{
@Override
protected DurableConfig createConfiguration() {
- return DurableConfig.builder().withPlugins(info -> new LoggingPlugin()).build();
+ return DurableConfig.builder().withPlugins(new LoggingPlugin()).build();
}
@Override
diff --git a/insight-plugin/README.md b/insight-plugin/README.md
index 694cf0df9..75d4e5132 100644
--- a/insight-plugin/README.md
+++ b/insight-plugin/README.md
@@ -192,29 +192,24 @@ one event per record to `{logStreamPrefix}YYYY/MM/DD` (IAM: `logs:CreateLogStrea
logs the failure rather than silently dropping it or failing the execution.
- **`includeErrors`** gates **both** the execution-level error and each operation-level error; with
`includeErrors(false)` neither is emitted, so a sensitive failure message never reaches a record.
-- **Non-fatal plugin failures are isolated.** Record construction, snapshotting, transforms,
- truncation and exporter render/export/flush contain and log ordinary failures, including
- `AssertionError` and optional-dependency `NoClassDefFoundError`. `VirtualMachineError` and
- `ThreadDeath` propagate as their original error objects, including through future, reflection
- and Jackson transport wrappers. Arbitrary business-exception causes are not reclassified.
- A fatal failure stops the shared scheduler: queued records are released and drain/flush waiters
- fail, including when another exporter is blocked. Subsequent scheduler calls rethrow that fatal;
- a finished fatal pump cannot later be mistaken for a successful drain or flush. Already-running
- customer exporter code cannot be forcibly stopped, but it does not delay fatal delivery.
+- **Plugin failures never disrupt execution.** Every plugin-owned boundary — record construction,
+ input snapshotting, transforms, truncation, and each exporter's render/export/flush (including a
+ `NoClassDefFoundError` from an optional exporter's absent SDK) — is guarded against any `Throwable`
+ and logged, so one failing exporter cannot block the others and no plugin fault propagates into
+ the durable execution.
- **Per-exporter size truncation** (`Truncation`) drops, in order: operation results oldest-first,
then whole operations oldest-first, then execution input, then output — setting `truncated`,
`droppedOperations`, `droppedInput`, `droppedOutput` as applicable. The size is measured against
the exact shape each exporter emits (its `render`).
- **Exporter isolation.** Every exporter receives its own copy of each record, truncated to its own
- limit; ordinary exporter failures are logged without preventing the other exporters from running.
+ limit; a failing or slow exporter is logged and never blocks the others or the execution.
- **Export scheduling.** Exporter I/O never runs on the SDK threads that deliver plugin hooks.
Records are handed to a background worker that exports at most one record at a time; each
record is a complete snapshot of its execution, so while an export is in flight newer updates for
- the same invocation coalesce into that invocation's pending slot and only the latest is exported
- next; records of different invocations never displace each other. At invocation end the plugin waits
- for its latest record to reach every exporter, then requests a flush through the same worker.
- Concurrent flush requests may share one flush, and the final record is delivered before the
- invocation returns.
+ the same execution coalesce into that execution's pending slot and only the latest is exported
+ next; records of different executions never displace each other. At invocation end the plugin waits
+ for the queue to drain and then flushes every exporter once, so the final record is always
+ delivered before the invocation returns.
## Conformance
diff --git a/insight-plugin/pom.xml b/insight-plugin/pom.xml
index cc949030f..c2cbe8df5 100644
--- a/insight-plugin/pom.xml
+++ b/insight-plugin/pom.xml
@@ -7,7 +7,7 @@
software.amazon.lambda.durable
aws-durable-execution-sdk-java-parent
- 3.0.0-SNAPSHOT
+ 2.2.2-SNAPSHOT
aws-durable-execution-sdk-java-plugin-insight
diff --git a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/ContentConfig.java b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/ContentConfig.java
index 0ecc1e38b..3a8228ec9 100644
--- a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/ContentConfig.java
+++ b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/ContentConfig.java
@@ -21,8 +21,8 @@
* detached, JSON-compatible copy of the value, never the SDK's original Java object: a POJO is presented as a
* {@code Map}, a list as a {@code List}, and a Java-time type as its JSON representation (for example an
* {@code Instant} arrives as an ISO-8601 {@code String}). Mutating the argument is therefore safe — it cannot corrupt
- * the cached input snapshot or any later emission — and a non-fatal transform failure omits the field and is logged.
- * VM/thread-termination errors propagate.
+ * the cached input snapshot or any later emission — and a transform that throws omits the field (the failure is logged)
+ * rather than failing the execution.
*/
@Experimental
public final class ContentConfig {
diff --git a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/ExportScheduler.java b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/ExportScheduler.java
index 9731683a9..84e8be3c2 100644
--- a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/ExportScheduler.java
+++ b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/ExportScheduler.java
@@ -2,82 +2,37 @@
// SPDX-License-Identifier: Apache-2.0
package software.amazon.lambda.durable.insight;
-import java.util.ArrayDeque;
import java.util.ArrayList;
-import java.util.Deque;
-import java.util.LinkedHashSet;
+import java.util.Iterator;
+import java.util.LinkedHashMap;
import java.util.List;
-import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executor;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
-import java.util.concurrent.TimeUnit;
-import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicInteger;
-import java.util.concurrent.atomic.AtomicReference;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
-import software.amazon.lambda.durable.insight.internal.FatalErrors;
/**
- * Serializes record exports so that, at most, one export runs at a time, while keeping the records of concurrently
- * running executions independent.
+ * Serializes record exports so that, at most, one export runs at a time.
*
- * One scheduler serves the whole execution environment — the exporters it fans out to belong to the environment, and
- * an environment can host several durable executions at the same time (Lambda Managed Instances makes that routine).
- * The work, by contrast, is held on the invocation it belongs to: each {@link InsightPlugin} instance is one
- * invocation's plugin and owns that invocation's latest-record slot, drain signal, mid-export marker and drain-waiter
- * count. Coalescing happens only within one invocation, because the slot belongs to it.
+ *
Each {@link WorkflowInsightRecord} is a complete snapshot of its execution, so a newer record for the
+ * same execution fully supersedes any record of that execution still waiting to be exported. While an export is in
+ * flight, additional updates are coalesced into a single pending slot per execution — intermediate records are dropped
+ * because the latest one already contains all of their information. Records of different executions never displace each
+ * other: one execution's snapshot carries none of another's information, so each execution keeps its own slot and the
+ * slots are served in the order the executions first became pending. This prevents overlapping {@code export()} calls
+ * when updates arrive faster than the exporters can keep up, and it keeps exporter I/O off the SDK threads that deliver
+ * plugin hooks.
*
- *
That shape is not a tidying-up. The same facts once lived in five structures keyed by execution ARN — the plugin's
- * state map plus this class's {@code pending}, {@code settled}, {@code exporting} and {@code drainWaiters} — and two of
- * them could disagree about one execution. They did: "nothing queued for this ARN" was read as "nothing outstanding",
- * which is equally true of a record already taken and being exported, so an exiting pump completed another execution's
- * drain signal mid-export and that invocation returned before its record was delivered. Every fact now exists exactly
- * once, as a field of the one object the SDK gave that invocation, so the disagreement has no place to happen — and
- * since the SDK creates that object and drops it, the scheduler has no execution registry to keep in step with
- * anything.
- *
- *
Each {@link WorkflowInsightRecord} is a complete snapshot of its execution, so a newer record from the same
- * invocation fully supersedes one still waiting to be exported. While an export is in flight, additional updates from
- * that invocation are coalesced into its slot — intermediate records are dropped because the latest one already
- * contains all of their information. A record from a different invocation never displaces another's record.
- *
- *
The slot takes whichever record is handed to it last and compares nothing, so "newer" has to be established before
- * the hand-off. Customer code runs while a record is being built and can re-enter a hook of the same invocation, which
- * builds and hands over a newer record first; the build it re-entered from then hands over an older snapshot last.
- * {@code InsightPlugin}'s build revision identifies each build and {@link #scheduleIfNotSuperseded} drops a record
- * whose build has been overtaken, so the slot only ever advances.
- *
- *
A single pump exports the queued records one at a time, in the order the invocations first queued work
- * ({@link #queue}, which is ordering only — membership in it is the same fact as "this invocation has a record",
- * written in one place), so exporters still never see two exports at once and each record keeps its per-exporter
- * fan-out. {@link #flush()} requests are served by that same pump, between records, so an exporter never sees a
- * {@code flush()} overlap an {@code export()} either. Requests are served as a batch — the cadence is at most one flush
- * per requesting invocation end, not exactly one — and a flush is preceded by the queued records a drain is waiting
- * for, so a burst of invocation ends is covered by one flush rather than one each. A request made while a flush runs
- * waits for the next turn. Exports are otherwise fire-and-forget; {@link #drain(InsightPlugin)} is called before an
- * invocation returns and waits for that invocation's own latest record to reach every exporter.
- *
- *
Because there is one pump, that wait can also cover records other invocations had already queued ahead of this
- * one: a drain is not isolated from the queue's head-of-line cost. What per-invocation ownership guarantees is that
- * another invocation's record can never displace this one — records coalesce only within their own invocation,
- * and a drain cannot return until its own latest record has reached every exporter.
+ *
Exports are otherwise fire-and-forget; {@link #drain()} is called before the invocation returns to guarantee the
+ * final record is delivered.
*/
final class ExportScheduler {
private static final AtomicInteger THREAD_NUMBER = new AtomicInteger();
- /**
- * Upper bound on the passes {@link #drainAll()} makes over the outstanding invocations. Only reached if new work
- * keeps arriving for as long as the drain runs; a normal drain settles in two passes.
- */
- private static final int MAX_DRAIN_ALL_PASSES = 1_000;
-
- /** How long one {@link #drainAll()} pass waits for a running pump before taking another pass. */
- private static final long PUMP_WAIT_MILLIS = 50;
-
/** Shared for the process lifetime; idle daemon workers are reclaimed, so nothing keeps the runtime alive. */
private static final ExecutorService WORKERS = Executors.newCachedThreadPool(runnable -> {
var thread = new Thread(runnable, "workflow-insight-export-" + THREAD_NUMBER.incrementAndGet());
@@ -93,68 +48,11 @@ final class ExportScheduler {
/** Completes when the current pump finishes; {@code null} while idle. Guarded by {@code this}. */
private CompletableFuture inFlight;
- // A fatal worker failure is terminal for this shared scheduler. Keep it observable even
- // after its pump/owner signals have been released, and wake fan-out joins immediately.
- private volatile Error fatalFailure;
- private final CompletableFuture fatalSignal = new CompletableFuture<>();
- private InsightPlugin activeExecution; // guarded by this monitor
-
- /**
- * The thread serving the pump right now, or {@code null} while no pump is running. Deliberately not
- * guarded by {@code this}: it is read by {@link #flush()} and the drains before they touch anything else, and a
- * lock acquisition there would put the pump's own monitor on the path of every invocation end.
- *
- * Written only by the thread that enters {@link #pump} — a worker, or the caller's own thread on the
- * rejected-worker fallback, which is a case where the calling thread genuinely is the pump — and cleared
- * by that same thread on the way out, only if it is still the recorded one. A compare-and-set on the way out rather
- * than a blind clear: if some anomaly ever did leave two pumps running, the one that finishes first must not clear
- * the other, and must not leave a stale thread behind that a later, legitimate {@code flush()} from that same
- * thread would be mistaken for.
- */
- private final AtomicReference pumpThread = new AtomicReference<>();
-
- /**
- * Marks the thread currently running one exporter's share of a fan-out for this scheduler, so that a
- * {@code flush()} or {@code drain()} re-entered from an exporter callback can tell that the pump is waiting for it.
- *
- * With a single exporter the fan-out runs on the pump thread and {@link #pumpThread} already recognizes it. With
- * two or more, {@link #forEachExporterSettled} submits one task per exporter and the pump then joins them all, so
- * the callback runs on a thread that is not the pump but that the pump cannot outlive: a wait for the pump issued
- * from there is the same wait-for cycle, two threads wide instead of one. The pump parks in the join, so it never
- * reaches the point in its loop that would complete the future the worker is parked on.
- *
- *
An instance field rather than a static: a fan-out worker of one scheduler is not pump-dependent on any other
- * scheduler, and refusing its waits there would be a false positive. Set and cleared around each callback by the
- * thread that runs it, restoring whatever was there before rather than blindly removing, so a callback that the
- * pump ran inline (the rejected-worker fallback, where the fan-out thread is the pump) cannot clear a mark
- * an enclosing frame still needs.
- */
- private final ThreadLocal exporterFanOutThread = new ThreadLocal<>();
-
- /**
- * The invocations with a record no pump has picked up yet, in the order they first queued work. Ordering only — the
- * record itself lives on the invocation's plugin instance. Guarded by {@code this}.
- *
- * This is the only collection of per-invocation objects the scheduler has, and it holds an instance for exactly
- * as long as that instance has a record waiting: nothing here has to be cleaned up at an invocation boundary, and
- * an instance the SDK has dropped is unreachable from the scheduler the moment its last record is taken.
- *
- *
Invariant, and the only thing that could still be said twice: an invocation is in here exactly while its
- * {@link InsightPlugin#record} is non-null. Every record moves through {@link #queueRecord}, {@link #takeRecord} or
- * {@link #dropRecord}, which write both halves together, and a set makes a double entry impossible by construction.
- */
- private final Set queue = new LinkedHashSet<>();
-
/**
- * One entry per outstanding {@link #flush()} request, in request order, completed when a {@code flush()} that
- * started after that request was enqueued has reached every exporter. Guarded by {@code this}.
- *
- * A queue of requests rather than a single flag: the pump takes the requests that are queued when its turn
- * begins and satisfies all of them with one flush, so concurrent invocation ends share a flush; a request enqueued
- * while that flush runs stays in the queue for the next turn, because a flush already in progress cannot be shown
- * to have seen the new requester's records.
+ * The latest record of each execution not yet picked up by the pump, keyed by execution ARN and served in the order
+ * the executions first became pending. Guarded by {@code this}.
*/
- private final Deque> flushRequests = new ArrayDeque<>();
+ private final LinkedHashMap pending = new LinkedHashMap<>();
ExportScheduler(
List exporters,
@@ -174,214 +72,53 @@ final class ExportScheduler {
this.executor = executor;
}
- // --- Scheduling. ---
-
- /**
- * Queues the latest record of one invocation for export, with no ordering check. If an export is already running,
- * the record is held in that invocation's own slot (replacing only an earlier record of the same
- * invocation) and exported once the pump reaches it.
- *
- * The slot takes whichever record is handed over last and does not compare record ages, so this is the right
- * entry point only for a record that cannot be superseded. The plugin's RUNNING records go through
- * {@link #scheduleIfNotSuperseded} and its final record through {@link #closeAndSchedule}; both add the ordering
- * checks this one omits.
- */
- void schedule(InsightPlugin execution, WorkflowInsightRecord record) {
- CompletableFuture handle;
- synchronized (this) {
- queueRecord(execution, record);
- handle = claimPumpIfIdle();
- }
- startPump(handle);
- }
-
/**
- * Schedules a non-terminal record unless it has been superseded, which is two separate facts.
- *
- * The invocation may already have ended. No RUNNING snapshot may follow the final record, so
- * {@link InsightPlugin#closed} rejects it.
- *
- *
A newer build of this same invocation may already have started. Customer code runs inside a build — the
- * content transforms, an operation result transform, a serializer for a customer type — and can re-enter a hook, so
- * the build that hands its record over last is not necessarily the build that started last. Without the revision
- * check the slot would take that older snapshot and the newer one would be lost, or, if a pump had already taken
- * the newer one, an exporter would see the older snapshot after the newer one.
- *
- *
The superseded record is dropped rather than queued. Nothing is lost: a record is a complete snapshot of one
- * execution, so the record that superseded it carries everything it carries. That is the same property that makes
- * the slot's coalescing sound.
- *
- *
Both checks and the hand-off are one critical section, on the monitor that owns both fields, so a record
- * cannot pass the checks and then be queued after the record that supersedes it.
- *
- * @param buildRevision the revision the caller took before it started building this record
- * @return whether the record was queued
+ * Queues the latest record of its execution for export. If an export is already running, the record is held in that
+ * execution's pending slot (replacing any earlier pending record of the same execution) and exported once the
+ * in-flight export completes.
*/
- boolean scheduleIfNotSuperseded(InsightPlugin execution, WorkflowInsightRecord record, long buildRevision) {
+ void schedule(WorkflowInsightRecord record) {
CompletableFuture handle;
synchronized (this) {
- if (execution.closed || !execution.isNewestBuild(buildRevision)) {
- return false;
- }
- queueRecord(execution, record);
- handle = claimPumpIfIdle();
- }
- startPump(handle);
- return true;
- }
-
- /**
- * Marks the invocation ended and, when given a record, schedules it as the last one for that invocation.
- *
- * The final record is queued without the build-revision check {@link #scheduleIfNotSuperseded} makes. Customer
- * code running inside the final record's build can start a newer RUNNING build, which would leave the final
- * record's revision stale, and a checked hand-off would then drop it and leave a RUNNING snapshot as the
- * execution's last exported state. Exempting it cannot let a stale record win, because {@code closed} is set in
- * this same critical section and every RUNNING record handed over afterwards is rejected.
- */
- void closeAndSchedule(InsightPlugin execution, WorkflowInsightRecord finalRecord) {
- if (finalRecord == null && execution.closed) {
- // The idempotent second call from the hook's `finally`. A volatile read, so the common case of an
- // invocation end that already scheduled its final record does not take the lock again.
- return;
- }
- CompletableFuture handle = null;
- synchronized (this) {
- execution.closed = true;
- if (finalRecord != null) {
- queueRecord(execution, finalRecord);
- handle = claimPumpIfIdle();
+ // put() on an existing key keeps its position, so a chatty execution cannot jump ahead of a quieter one.
+ pending.put(executionKey(record), record);
+ if (inFlight != null) {
+ return;
}
- }
- startPump(handle);
- }
-
- /**
- * Claims the pump for the caller when none is running, returning the handle to run with, or {@code null} when a
- * pump already owns the scheduler and will pick the work up. Caller holds the lock.
- */
- private CompletableFuture claimPumpIfIdle() {
- if (inFlight != null) {
- return null;
- }
- CompletableFuture handle = new CompletableFuture<>();
- inFlight = handle;
- return handle;
- }
-
- /** Starts a claimed pump on a worker; a no-op when the caller claimed nothing. */
- private void startPump(CompletableFuture handle) {
- if (handle == null) {
- return;
+ handle = new CompletableFuture<>();
+ inFlight = handle;
}
try {
executor.execute(() -> pump(handle));
} catch (Throwable t) {
- rejectFatal(t);
- // No worker could be started. Keep the queued record and return to idle so a later schedule() retries, and
- // a drain runs whatever is still queued on the calling thread before the invocation returns. Complete the
- // handle too: a drain that already observed it must wake up and take that inline path.
+ // No worker could be started. Keep the pending record and return to idle so a later schedule() retries,
+ // and drain() runs whatever is still pending on the calling thread before the invocation returns. Complete
+ // the handle too: a drain() that already observed it must wake up and take that inline path.
synchronized (this) {
if (inFlight == handle) {
inFlight = null;
}
}
- completeSignal(handle);
+ handle.complete(null);
reportFailure(t);
}
}
- // --- The record slot: the two halves of "this invocation has a record queued", always written together. ---
-
- /** Puts this invocation's latest record in its slot and makes sure it has a drain signal. Caller holds the lock. */
- private void queueRecord(InsightPlugin execution, WorkflowInsightRecord record) {
- throwIfFailed();
- execution.record = record;
- queue.add(execution);
- if (execution.settled == null) {
- execution.settled = new CompletableFuture<>();
- }
- }
-
- /**
- * Takes this invocation's queued record for export and marks it mid-export. Caller holds the lock and has checked
- * that a record is there.
- *
- * The marking is not a separate step in a separate structure: leaving the slot and becoming "inside the
- * exporters" are one write of one object, so no reader can see the invocation between the two and conclude it has
- * nothing outstanding.
- */
- private WorkflowInsightRecord takeRecord(InsightPlugin execution) {
- WorkflowInsightRecord record = execution.record;
- execution.record = null;
- queue.remove(execution);
- execution.exporting = true;
- activeExecution = execution;
- return record;
- }
-
- /** Drops this invocation's queued record without exporting it. Caller holds the lock. */
- private void dropRecord(InsightPlugin execution) {
- execution.record = null;
- queue.remove(execution);
- }
-
- // --- Draining. ---
-
/**
- * Waits until the latest record of one invocation has been handed to every exporter. Safe to call when that
- * invocation has nothing outstanding. Used before the invocation returns to guarantee the final record is
- * delivered.
- *
- *
The wait is for that invocation's own latest record. Another invocation's record can never displace it, so
- * this always returns having delivered this invocation's latest snapshot; but since one pump exports serially, the
- * wait can also cover records other invocations had already queued ahead of it.
- *
- *
While this waits, the invocation counts a drain waiter — on the instance itself, so the count cannot come to
- * describe a different one. It tells the pump that this record gates an invocation return, so the pump exports it
- * before spending a flush fan-out. See {@link #exportRecordsADrainIsWaitingFor}.
- *
- *
Called from the pump thread itself, or from an exporter fan-out worker that pump is waiting for, the wait is
- * refused and reported instead of made: see {@link #refuseWaitThatWouldBlockThePump}.
+ * Waits for any in-flight and pending exports to complete. Safe to call when idle. Used before the invocation
+ * returns to guarantee the final record is delivered.
*/
- void drain(InsightPlugin execution) {
- throwIfFailed();
- // Re-entered from a thread the pump's progress depends on: waiting here would park on a signal only that pump
- // can settle. Refuse and return; the record stays queued and that same pump exports it once it resumes its
- // loop.
- if (refuseWaitThatWouldBlockThePump("drain(execution)")) {
- return;
- }
- synchronized (this) {
- throwIfFailed();
- if (execution.settled == null) {
- return;
- }
- execution.drainWaiters++;
- }
- try {
- drainUntilSettled(execution);
- } finally {
- synchronized (this) {
- execution.drainWaiters--;
- }
- }
- }
-
- private void drainUntilSettled(InsightPlugin execution) {
+ void drain() {
while (true) {
- CompletableFuture signal;
CompletableFuture handle;
boolean runInline = false;
synchronized (this) {
- throwIfFailed();
- signal = execution.settled;
- if (signal == null) {
- return;
- }
handle = inFlight;
if (handle == null) {
- // A record is outstanding with no pump running (a worker could not be started): export it here.
+ if (pending.isEmpty()) {
+ return;
+ }
+ // Records are pending with no pump running (a worker could not be started): export them here.
handle = new CompletableFuture<>();
inFlight = handle;
runInline = true;
@@ -389,447 +126,52 @@ private void drainUntilSettled(InsightPlugin execution) {
}
if (runInline) {
pump(handle);
- if (nothingCanSettle(execution, signal)) {
- // The pump this thread just ran found no record for this invocation and left none inside the
- // exporters, yet the signal survives: no later step can complete it, so waiting again would only
- // start empty pumps forever. Release it here and report — the pump's own exit does not sweep for
- // this any more, because it has no registry of invocations to sweep and does not need one: the
- // thread that would be stranded is this one, and it holds the instance.
- abandon(execution);
- reportFailure(new IllegalStateException(
- "a drain signal survived a pump that had nothing to export for it; the drain was released"
- + " rather than waiting for work nobody will do"));
- return;
- }
- continue;
- }
- // Wake either when this invocation's record has been exported or when the current pump ends — the pump may
- // have ended without taking this record (a rejected worker), in which case the loop re-evaluates and
- // exports it inline.
- try {
- CompletableFuture.anyOf(signal, handle).join();
- } catch (Throwable t) {
- // Never spin on an unexpected wait failure, and never let it escape into the execution. Abandon this
- // invocation's outstanding record instead of leaving it queued: WORKERS is a static, process-wide pool,
- // so a record left queued here would be exported later by some unrelated execution's pump — out of
- // order, and after this invocation has already returned. Completing the signal also releases any other
- // drain waiting on the same invocation rather than stranding it behind work nobody will do.
- abandon(execution);
- reportFailure(t);
- return;
- }
- }
- }
-
- /**
- * True when this invocation still holds the same drain signal but has no queued record and none inside the
- * exporters, so nothing that could complete the signal is left. No ordinary path produces that; this is the
- * liveness backstop for the unwinds that are hard to enumerate exhaustively.
- */
- private synchronized boolean nothingCanSettle(InsightPlugin execution, CompletableFuture signal) {
- return execution.settled == signal && execution.record == null && !execution.exporting;
- }
-
- /**
- * Waits for every outstanding record. Test seam for an environment-wide drain; the per-invocation path uses
- * {@link #drain(InsightPlugin)}.
- *
- * Returns once nothing is queued and no pump owns the scheduler, which is exactly "every record scheduled so far
- * has reached the exporters": a pump only returns to idle with its queue empty and the record it took settled.
- *
- *
Bounded by the number of passes, not by the set of invocations seen: an invocation that queues new work after
- * it was already drained must still be waited for (dropping it would silently weaken every assertion made after
- * this returns), while a producer that never stops cannot keep this spinning forever.
- */
- void drainAll() {
- throwIfFailed();
- // Every pass below is a drain, and each one would be refused; without this the loop spends all of its passes
- // reporting the same refusal.
- if (refuseWaitThatWouldBlockThePump("drainAll()")) {
- return;
- }
- for (int pass = 0; pass < MAX_DRAIN_ALL_PASSES; pass++) {
- List outstanding;
- CompletableFuture handle;
- synchronized (this) {
- throwIfFailed();
- outstanding = new ArrayList<>(queue);
- handle = inFlight;
- if (outstanding.isEmpty() && handle == null) {
- return;
- }
- }
- for (InsightPlugin execution : outstanding) {
- drain(execution);
- }
- if (outstanding.isEmpty()) {
- // Nothing is queued for anyone, but a pump still owns the scheduler: it may be inside an exporter with
- // a
- // record whose only reference is its own local variable, and with no registry of invocations there is
- // no
- // way to name that record and drain it. Waiting for the pump itself covers it — in bounded steps, so a
- // producer that keeps the pump permanently busy cannot make this unbounded, and so the wait is a real
- // wait rather than a re-poll.
- try {
- handle.get(PUMP_WAIT_MILLIS, TimeUnit.MILLISECONDS);
- } catch (TimeoutException e) {
- // Still running; take another pass.
- } catch (Throwable t) {
- reportFailure(t);
- return;
- }
+ } else {
+ handle.join();
}
}
}
- /**
- * Test seam: how many invocations the scheduler still holds a reference to.
- *
- * {@link #queue} is the only collection of per-invocation objects the scheduler has, so this is the whole of the
- * per-invocation state the environment retains. Zero means the environment — which outlives every invocation —
- * holds nothing belonging to any invocation it has served.
- */
- synchronized int retainedInvocationCount() {
- return queue.size();
- }
-
- /** Test seam: whether the scheduler still holds a reference to one particular invocation. */
- synchronized boolean retains(InsightPlugin execution) {
- return queue.contains(execution);
- }
-
- /** Gives up one invocation's outstanding work: drops its queued record and releases every drain waiting on it. */
- private void abandon(InsightPlugin execution) {
- CompletableFuture signal;
- synchronized (this) {
- dropRecord(execution);
- execution.exporting = false;
- signal = execution.settled;
- execution.settled = null;
- }
- if (signal != null) {
- completeSignal(signal);
- }
- }
-
- // --- The pump. ---
-
private void pump(CompletableFuture handle) {
- // Recorded for as long as this thread serves the pump — a worker, or a caller pumping inline after a rejected
- // worker — so that a flush() or drain() re-entered from anything the fan-out calls synchronously can tell that
- // it is asking itself. One atomic write per pump, and no lock: see the field.
- Thread self = Thread.currentThread();
- pumpThread.set(self);
- // The invocation this pump has taken a record from and not settled yet. Only this pump may release it, so an
- // abnormal unwind cannot strand a drain, and no other pump can mistake it for orphaned work.
- InsightPlugin taken = null;
- // Likewise for the flush requests this pump has taken out of the queue and not completed yet.
- List> takenFlushes = null;
try {
- // One record, then every flush request queued at that moment, alternating. Taking the record and returning
- // to idle both happen under the lock, so a record scheduled at any point is either exported by this pump or
- // starts the next one — never lost, and never displaced by another invocation's record. A flush therefore
- // waits at most one fan-out (it cannot be starved by a queue that never runs dry) and still never overlaps
- // an export, because this loop runs them one after the other.
- //
- // A loop, deliberately, not a pump that re-enters itself to pick up the next item: written that way, one
- // frame per queued item accumulates until the stack overflows, and the rest of the queue is dropped.
+ // Serve pending executions in order until none is left. Taking a record and returning to idle both happen
+ // under the lock, so an update scheduled at any point is either exported by this pump or starts the next
+ // one — never lost.
while (true) {
- InsightPlugin next = null;
- WorkflowInsightRecord record = null;
+ WorkflowInsightRecord record;
synchronized (this) {
- throwIfFailed();
- if (queue.isEmpty() && flushRequests.isEmpty()) {
- if (inFlight == handle) {
- inFlight = null;
- }
+ Iterator head = pending.values().iterator();
+ if (!head.hasNext()) {
+ inFlight = null;
return;
}
- if (!queue.isEmpty()) {
- next = queue.iterator().next();
- // Leaving the slot and being marked mid-export are one write of one object, so the invocation
- // is
- // never momentarily indistinguishable from one with nothing outstanding.
- record = takeRecord(next);
- }
- }
- if (next != null) {
- taken = next;
- try {
- exportToAll(record);
- } finally {
- signalSettled(next);
- taken = null;
- }
- }
- // Taken only now that the fan-out above has settled, and taken as a batch: every request queued at this
- // instant is satisfied by the single flush below, so invocation ends that ask together cost one flush
- // rather than one each. Sound because each requester drained its own record before asking, so a flush
- // that starts after the request was enqueued already has that record in the exporter's buffer.
- //
- // Emptying the queue here — rather than after the flush — is what keeps a request that arrives while
- // that flush runs out of this batch: it lands in the now-empty queue and is served by the next turn,
- // never credited to a flush that was already in progress when it was made.
- synchronized (this) {
- if (!flushRequests.isEmpty()) {
- takenFlushes = new ArrayList<>(flushRequests);
- flushRequests.clear();
- }
- }
- if (takenFlushes != null) {
- // Before spending the fan-out: export the queued records other invocations are still waiting on.
- // Those ends cannot have asked for their flush yet — they are inside a drain — so without this the
- // pump staggers them one record per turn, with a whole flush in between, and each pays for its own
- // flush however aggressively the queue is coalesced.
- exportRecordsADrainIsWaitingFor();
- // Re-take: the ends released above ask for their flush now, and one flush covers all of them since
- // it starts after every one of those records reached the exporters.
- synchronized (this) {
- if (!flushRequests.isEmpty()) {
- takenFlushes.addAll(flushRequests);
- flushRequests.clear();
- }
- }
- try {
- flushEveryExporter();
- } finally {
- // In `finally`: a Throwable from a customer's flush() — an Error, not just an exception — must
- // never leave the invocations waiting on these requests parked forever.
- completeAll(takenFlushes);
- takenFlushes = null;
- }
+ record = head.next();
+ head.remove();
}
+ exportToAll(record);
}
- } catch (Throwable failure) {
- rejectFatal(failure);
- throw failure;
} finally {
- try {
- // Before anything else, and before the handle below: whoever waits on these requests must be released
- // even
- // if this pump is unwinding for a reason none of the guards above anticipated.
- if (takenFlushes != null) {
- completeAll(takenFlushes);
- }
- CompletableFuture orphaned = null;
- synchronized (this) {
- if (inFlight == handle) {
- inFlight = null;
- }
- if (taken != null) {
- // Unwinding with a record still marked as being exported: this pump will never settle it.
- // Release
- // the marker and, unless a newer record for the same invocation is queued for a later pump to
- // export, complete the drain waiting on it — the instance is right here, so no sweep over other
- // invocations is needed to find it.
- taken.exporting = false;
- if (taken.record == null) {
- orphaned = taken.settled;
- taken.settled = null;
- }
- }
- }
- if (orphaned != null) {
- completeSignal(orphaned);
- }
- completeSignal(handle);
- // Last, because everything above is still this pump's work and a flush() re-entered from any of it
- // would
- // still have nobody to serve it. Conditional: a pump that recorded itself since must not be cleared
- // here.
- } finally {
- pumpThread.compareAndSet(self, null);
- }
- }
- }
-
- /**
- * Exports the queued records that a drain is waiting for, one at a time, and returns once they have all reached the
- * exporters. Called by the pump immediately before a flush.
- *
- * Those records are the last records of invocations that cannot return until they are exported, and their ends
- * cannot ask for their flush until then. Exporting them first is therefore what lets one flush serve a whole burst
- * of invocation ends: without it the pump interleaves one record and one flush fan-out, and each end pays for a
- * flush of its own even though every request is coalesced.
- *
- *
Bounded by the snapshot taken under the lock, so a producer that keeps scheduling for an invocation someone is
- * draining cannot hold a flush back indefinitely — and records nobody waits for are not exported here at all, so a
- * stream of {@code ON_CHANGE} snapshots still cannot starve a flush: it waits at most one ordinary fan-out plus
- * this pass over the invocations whose return is already blocked on their own record.
- */
- private void exportRecordsADrainIsWaitingFor() {
- List awaited = null;
- synchronized (this) {
- for (InsightPlugin execution : queue) {
- if (execution.drainWaiters > 0) {
- if (awaited == null) {
- awaited = new ArrayList<>();
- }
- awaited.add(execution);
- }
- }
- }
- if (awaited == null) {
- return;
- }
- for (InsightPlugin execution : awaited) {
- WorkflowInsightRecord record;
synchronized (this) {
- record = execution.record == null ? null : takeRecord(execution);
- }
- if (record == null) {
- continue;
- }
- try {
- exportToAll(record);
- } finally {
- try {
- signalSettled(execution);
- } catch (Throwable t) {
- // Nothing here is expected to throw, but a record left marked as being exported would strand the
- // drain that is waiting for it, so release it and that drain rather than leave the invocation
- // parked.
- abandon(execution);
- reportFailure(t);
+ if (inFlight == handle) {
+ inFlight = null;
}
}
+ handle.complete(null);
}
}
- /** Completes every taken flush request; one that cannot be completed must not stop the rest from being. */
- private void completeAll(List> requests) {
- for (CompletableFuture request : requests) {
- try {
- completeSignal(request);
- } catch (Throwable t) {
- reportFailure(t);
- }
- }
- }
-
- /**
- * Completes one invocation's drain signal now that its record has been exported, unless a newer record from the
- * same invocation arrived meanwhile — that one settles the signal instead, so a drain always waits for the latest.
- */
- private void signalSettled(InsightPlugin execution) {
- CompletableFuture signal;
- synchronized (this) {
- if (activeExecution == execution) activeExecution = null;
- if (execution.record != null) {
- // A newer record is queued for the same invocation. Leave it marked as being exported: it is still
- // outstanding, and the export of that newer record settles the signal.
- return;
- }
- execution.exporting = false;
- signal = execution.settled;
- execution.settled = null;
- }
- if (signal != null) {
- completeSignal(signal);
- }
- }
-
- // --- Flushing. ---
-
- /**
- * Flushes every exporter, serialized against exports: the request is queued and served by the pump between records,
- * so an exporter never sees {@code flush()} overlap {@code export()} — not even an export belonging to a different
- * execution running in the same environment. Returns once a flush that started after this request was enqueued has
- * reached every exporter.
- *
- * Requests are coalesced: the pump takes every request queued at the start of its turn, exports any queued
- * record a drain is still waiting for, re-takes the requests those ends make as they are released, and satisfies
- * them all with one flush. Invocation ends that overlap therefore share a flush instead of paying for one fan-out
- * each. That is sound because a caller drains its own record before asking, so a flush that starts after
- * the request was enqueued has that record in the exporter's buffer. A request enqueued while a flush is already
- * running is never satisfied by it — it waits for the next turn.
- *
- *
A queue that never runs dry cannot starve a request either: the pump alternates one record and one batch of
- * requests, so a flush waits at most one export fan-out.
- *
- *
Called from the pump thread itself — or from an exporter fan-out worker that pump is waiting for, which is
- * what a callback re-entering the scheduler does when two or more exporters are configured — the request is refused
- * and reported instead of made: see {@link #refuseWaitThatWouldBlockThePump}.
- */
- void flush() {
- throwIfFailed();
- // Re-entered from a thread the pump's progress depends on: the pump is the only thread that could serve the
- // request, and it cannot while this caller has not returned. Refuse rather than enqueue a request nobody
- // serves.
- if (refuseWaitThatWouldBlockThePump("flush()")) {
- return;
- }
- CompletableFuture request = new CompletableFuture<>();
- synchronized (this) {
- throwIfFailed();
- flushRequests.add(request);
- }
- while (true) {
- CompletableFuture handle;
- boolean startPump = false;
- synchronized (this) {
- throwIfFailed();
- if (request.isDone()) {
- return;
- }
- handle = inFlight;
- if (handle == null) {
- if (!flushRequests.contains(request)) {
- // Liveness backstop: a pump took this request and unwound without serving it, which its
- // `finally` is there to prevent. The request is no longer in the queue, so no future pump can
- // find it — release the caller here instead of spinning up pumps that have nothing to do.
- break;
- }
- // No pump is running (a worker could not be started earlier, or the pump went idle between the add
- // above and this check): start one.
- handle = new CompletableFuture<>();
- inFlight = handle;
- startPump = true;
- }
- }
- if (startPump) {
- CompletableFuture started = handle;
- try {
- executor.execute(() -> pump(started));
- } catch (Throwable t) {
- // No worker could be started. Serve the request on the calling thread, exactly as a drain exports a
- // queued record inline: this pump owns `inFlight`, so no export can run beside it.
- reportFailure(t);
- pump(started);
- continue;
- }
- }
- // Wake either when this request has been served or when the current pump ends — a pump can end without
- // serving it (a rejected worker), in which case the loop starts another one.
- try {
- CompletableFuture.anyOf(request, handle).join();
- } catch (Throwable t) {
- // Never spin on an unexpected wait failure, and never let it escape into the execution. Drop the
- // request rather than leaving it queued for some later, unrelated invocation's pump to serve.
- synchronized (this) {
- flushRequests.remove(request);
- }
- reportFailure(t);
- break;
- }
- }
- completeSignal(request);
+ private static String executionKey(WorkflowInsightRecord record) {
+ return record.executionArn() != null ? record.executionArn() : "";
}
/**
* Flushes every exporter, each on its own worker, and waits for all of them to settle. A slow or failing flush on
- * one exporter never delays or fails the others. Environment-wide, like the exporters themselves.
- *
- * Private and called only from the pump: routing every flush through the pump is what keeps a {@code flush()}
- * from overlapping an {@code export()}, so this must not be reachable from outside. The per-exporter fan-out below
- * is parallelism within one flush, not concurrency with an export.
+ * one exporter never delays or fails the others.
*/
- private void flushEveryExporter() {
+ void flushAll() {
forEachExporterSettled(InsightExporter::flush);
}
- // --- Exporting. ---
-
/**
* Exports one record to every exporter, each on its own worker, and waits for all of them to settle. One failing or
* slow exporter never blocks or fails the others, and an export error never propagates into the execution.
@@ -841,52 +183,26 @@ private void exportToAll(WorkflowInsightRecord record) {
/** Runs the action for every exporter concurrently and returns once all have settled, reporting each failure. */
private void forEachExporterSettled(Consumer action) {
if (exporters.size() == 1) {
- // On the pump thread itself, which the pump-thread check already refuses waits from.
runSafely(() -> action.accept(exporters.get(0)));
return;
}
- List> settledExporters = new ArrayList<>(exporters.size());
+ List> settled = new ArrayList<>(exporters.size());
for (InsightExporter exporter : exporters) {
- throwIfFailed();
- // Marked as a fan-out task: the pump joins every one of these below, so a wait for the pump issued from
- // inside one must be refused exactly as one issued from the pump itself.
- Runnable task = () -> runSafely(() -> runAsExporterFanOut(() -> action.accept(exporter)));
+ Runnable task = () -> runSafely(() -> action.accept(exporter));
try {
- settledExporters.add(CompletableFuture.runAsync(task, executor));
+ settled.add(CompletableFuture.runAsync(task, executor));
} catch (Throwable t) {
reportFailure(t);
task.run();
}
}
- for (CompletableFuture task : settledExporters) {
- // A fatal on another worker must reach the invocation even if this peer is blocked.
- runSafely(() -> CompletableFuture.anyOf(task, fatalSignal).join());
- }
- }
-
- /**
- * Runs one exporter's share of a fan-out with this thread marked pump-dependent, restoring the previous mark on the
- * way out. The mark is what makes {@link #refuseWaitThatWouldBlockThePump} recognize a fan-out worker.
- */
- private void runAsExporterFanOut(Runnable action) {
- Boolean previous = exporterFanOutThread.get();
- exporterFanOutThread.set(Boolean.TRUE);
- try {
- action.run();
- } finally {
- if (previous == null) {
- // Removed rather than set back to null: these run on a shared, process-wide pool, so a thread must not
- // keep an entry for this scheduler after its task ends.
- exporterFanOutThread.remove();
- } else {
- exporterFanOutThread.set(previous);
- }
+ for (CompletableFuture task : settled) {
+ runSafely(task::join);
}
}
private void runSafely(Runnable action) {
try {
- throwIfFailed();
action.run();
} catch (Throwable t) {
reportFailure(t);
@@ -894,102 +210,10 @@ private void runSafely(Runnable action) {
}
private void reportFailure(Throwable t) {
- rejectFatal(t);
try {
failureHandler.accept(t);
} catch (Throwable ignored) {
- rejectFatal(ignored);
- // Ordinary diagnostic failures remain isolated.
- }
- }
-
- void throwIfFailed() {
- Error failure = fatalFailure;
- if (failure != null) throw failure;
- }
-
- private void rejectFatal(Throwable failure) {
- Error fatal = FatalErrors.find(failure);
- if (fatal != null) {
- fail(fatal);
- throw fatal;
- }
- }
-
- private void completeSignal(CompletableFuture signal) {
- Error failure = fatalFailure;
- if (failure == null) signal.complete(null);
- else signal.completeExceptionally(failure);
- }
-
- /** Releases all owners and waiters without waiting for uncooperative exporter code. */
- synchronized void fail(Error failure) {
- if (fatalFailure != null) return;
- fatalFailure = failure;
- CompletableFuture handle = inFlight;
- inFlight = null;
- fatalSignal.completeExceptionally(failure);
- if (handle != null) handle.completeExceptionally(failure);
- if (activeExecution != null) {
- failOwner(activeExecution, failure);
- activeExecution = null;
- }
- while (!queue.isEmpty()) failOwner(queue.iterator().next(), failure);
- CompletableFuture request;
- while ((request = flushRequests.poll()) != null) request.completeExceptionally(failure);
- }
-
- private void failOwner(InsightPlugin execution, Error failure) {
- dropRecord(execution);
- execution.exporting = false;
- CompletableFuture signal = execution.settled;
- execution.settled = null;
- if (signal != null) signal.completeExceptionally(failure);
- }
-
- /**
- * Reports and refuses a wait for the pump that was issued from a thread the pump's own progress depends on. Returns
- * whether the caller is such a thread; when it is, the failure has already been reported and the caller must return
- * without waiting.
- *
- * Invariant: the thread that waits for the pump is never a thread the pump waits for. {@link #flush()} waits for
- * a request only a pump can complete, and a drain waits for a signal only a pump can complete or for the running
- * pump's own handle. All three are satisfied by the pump between records.
- *
- *
Two threads qualify. The pump thread itself: a wait issued from there is a wait-for cycle one thread wide —
- * the pump parks on the future it would itself have completed, so it never reaches the point in its loop that
- * completes it, and no other thread may take over because {@code inFlight} is this pump's. And an exporter fan-out
- * worker: with two or more exporters the pump submits one task per exporter and joins them all, so a wait issued
- * from a callback running on one of those workers is the same cycle two threads wide — the worker parks on a future
- * only the pump can complete, and the pump is parked in the join waiting for that worker. Neither is a monitor
- * deadlock, so the JVM's deadlock detection cannot see either one, and the invocation simply never returns.
- *
- *
Reachable through anything a fan-out calls synchronously: with a single exporter the fan-out runs on the pump
- * thread, so a customer exporter's {@code export()} or {@code flush()} that asks the scheduler for a flush, or a
- * non-conforming {@code exportOne}, is enough; with several it runs on a worker instead, and the same call is
- * refused for the same reason. A conforming production {@code exportOne} does not re-enter the scheduler, so this
- * is hardening.
- *
- *
So the call fails fast instead: the plugin's failure handler is told — it logs — and the caller returns as it
- * would from any other flush or drain, with nothing propagating into the execution. The queued work itself is not
- * dropped by refusing a drain: the record stays in the invocation's slot, and the pump that is waiting for this
- * caller exports it as soon as this caller returns and the fan-out it belongs to settles. Callers that are neither
- * — every SDK hook thread — never enter this branch and behave exactly as before; the check is a volatile read plus
- * a thread-local read, so no lock is added to that path.
- */
- private boolean refuseWaitThatWouldBlockThePump(String call) {
- if (pumpThread.get() == Thread.currentThread()) {
- reportFailure(new IllegalStateException(call
- + " was called from the export pump thread, the only thread able to serve it; the call was refused"
- + " rather than deadlocking the invocation"));
- return true;
- }
- if (Boolean.TRUE.equals(exporterFanOutThread.get())) {
- reportFailure(new IllegalStateException(call
- + " was called from an exporter fan-out worker the export pump is waiting for, so the pump cannot"
- + " serve it; the call was refused rather than deadlocking the invocation"));
- return true;
+ // A scheduler diagnostic must never disrupt durable execution.
}
- return false;
}
}
diff --git a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightExporter.java b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightExporter.java
index 59f04767e..69ff39739 100644
--- a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightExporter.java
+++ b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightExporter.java
@@ -11,34 +11,7 @@ public interface InsightExporter {
/** Emits one record to the destination. */
void export(WorkflowInsightRecord record);
- /**
- * Flushes any records this exporter has buffered. The default is a no-op; override it only if
- * {@link #export(WorkflowInsightRecord)} buffers rather than emitting immediately.
- *
- *
Called at most once per sampled-in invocation end, after that invocation's own record — if it emitted one —
- * has been handed to every exporter. An end that emits no record still flushes (a non-terminal suspend under
- * {@code ON_COMPLETE}, a success under {@code ON_FAILURE}), so records buffered by that execution's earlier
- * emissions are never left behind. Invocation ends that overlap may share a single flush: one flush is enough for
- * all of them, because it starts only after each of their records has been handed to every exporter. An execution
- * that is sampled out neither exports nor flushes.
- *
- *
Never called concurrently with {@link #export(WorkflowInsightRecord)} by the plugins one
- * {@link WorkflowInsight#workflowInsight} factory creates. That factory owns the scheduler serializing them, so the
- * guarantee is per factory rather than per environment: an exporter instance handed to two factories is served by
- * two schedulers, which can call its {@code export} and {@code flush} at the same time. Build the factory once per
- * handler — which is what a {@code DurableConfig} created once per handler does — and give each factory its own
- * exporter instances if a single exporter cannot tolerate concurrent calls.
- *
- *
May cover records belonging to other executions running in the same environment, so it is not a per-execution
- * barrier.
- *
- *
Must return promptly. No invocation whose end is waiting on this flush can return until it returns, and since
- * overlapping ends may share one flush, a slow flush is billed to every one of those invocations — not only to the
- * one that asked for it.
- *
- *
Non-fatal failures are reported and isolated without retries. Fatal VM/thread-termination errors propagate
- * through the scheduler to the invocation caller; pending scheduler work is released.
- */
+ /** Flushes any buffered records; no-op by default. */
default void flush() {}
/**
diff --git a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightPlugin.java b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightPlugin.java
deleted file mode 100644
index 922f675de..000000000
--- a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightPlugin.java
+++ /dev/null
@@ -1,458 +0,0 @@
-// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
-// SPDX-License-Identifier: Apache-2.0
-package software.amazon.lambda.durable.insight;
-
-import java.time.Instant;
-import java.util.ArrayList;
-import java.util.Comparator;
-import java.util.List;
-import java.util.Map;
-import java.util.concurrent.CompletableFuture;
-import java.util.concurrent.atomic.AtomicLong;
-import software.amazon.lambda.durable.insight.internal.FatalErrors;
-import software.amazon.lambda.durable.plugin.DurableExecutionPlugin;
-import software.amazon.lambda.durable.plugin.InvocationEndInfo;
-import software.amazon.lambda.durable.plugin.InvocationInfo;
-import software.amazon.lambda.durable.plugin.OperationChangeInfo;
-import software.amazon.lambda.durable.plugin.OperationChangeItemInfo;
-
-/**
- * The Workflow Insight plugin instance for one Lambda invocation: both the state that invocation's records are built
- * from and the slot the {@link ExportScheduler} exports them through.
- *
- *
The SDK creates one of these per invocation, from the {@link InvocationInfo} it is about to hand the first hook,
- * and drops it when the invocation returns. So everything about an execution is a plain field here — the parsed ARN,
- * the stable start time, the one-time sampling decision, the detached input snapshot, the latest queued record, the
- * drain signal, the mid-export marker and the drain-waiter count. There is nothing to key by execution ARN and nothing
- * to register or release: an instance is the registration, and its lifetime is the invocation's.
- *
- *
Two objects outlive the invocation and are shared by every instance the factory creates: the resolved
- * {@link InsightSettings} and the {@link ExportScheduler}. The scheduler is shared on purpose — serializing exports is
- * a property of the exporters, which belong to the environment, not to one invocation.
- *
- *
Ownership. Three groups of fields:
- *
- *
- * - Identity — {@link #executionArn}, {@link #arn}, {@link #startTime}, {@link #sampledIn} — is taken from
- * the {@link InvocationInfo} the factory receives and is {@code final}. It cannot be observed half-built, and
- * there is no second invocation that could change it.
- *
- The input snapshot — {@link #cachedInput} — is written by the thread that fires
- * {@code onInvocationStart} and read by the operation-change and invocation-end threads of the same invocation,
- * which the SDK does not promise are the same thread; {@code volatile} for that publication.
- *
- The build revision — {@link #buildRevision} — counts the record builds this invocation has started, so
- * that a build which was overtaken can be recognized at hand-off time and its record dropped. Atomic rather than
- * {@code volatile}, because the case it exists for is two builds running at once. See the field.
- *
- Scheduling state — {@link #record}, {@link #settled}, {@link #exporting}, {@link #drainWaiters} and
- * {@link #closed} — is shared with the export pump and guarded by the monitor of {@link #scheduler}. One monitor
- * for the whole environment, not one per invocation, so the {@code closed} check and the hand-off of a record are
- * a single critical section and there is no lock ordering between instances to get wrong.
- *
- *
- * {@link #closed} is additionally {@code volatile}: the hook threads read it without the lock as a fast pre-check.
- * That read only ever skips work — the authoritative check is made under the monitor by
- * {@link ExportScheduler#scheduleIfOpen}. It is written once, from false to true, and never back: an execution that
- * suspends and resumes gets a new instance rather than a reset one.
- */
-final class InsightPlugin implements DurableExecutionPlugin {
-
- /** Resolved configuration, shared by every instance of the environment. */
- private final InsightSettings settings;
-
- /**
- * Shared with every other instance: exports are serialized across the whole environment. Package-private because it
- * is the monitor this instance's scheduling fields are guarded by, which the tests in this package hold when they
- * read them.
- */
- final ExportScheduler scheduler;
-
- // --- Identity, from the InvocationInfo the factory was called with. ---
-
- /** The execution this instance observes. */
- final String executionArn;
-
- /** The parsed execution ARN, parsed once for every record this instance builds. */
- final ArnParser arn;
-
- /** Stable execution start time, from {@code InvocationInfo.executionStartTime()}. */
- final Instant startTime;
-
- /** The one-time sampling decision; deterministic in the ARN, so a resumed invocation decides the same way. */
- final boolean sampledIn;
-
- // --- The input snapshot. ---
-
- /**
- * Detached snapshot of the execution input, the single source of truth for {@code input} on every emission of this
- * invocation. Written by {@code onInvocationStart}, read by every later build.
- */
- private volatile Object cachedInput;
-
- // --- Build ordering. ---
-
- /**
- * Counts the record builds this invocation has started. The value a build takes identifies that build.
- *
- *
Customer code runs inside a build, on the hook thread: the input and output content transforms, an operation's
- * result transform, and any Jackson serializer registered for a customer type. That code can call back into a hook
- * of this same instance, and it runs before anything is scheduled, so a build can be overtaken by a newer build
- * that starts and finishes inside it. Two hook threads for one invocation would produce the same overlap.
- *
- *
The scheduler's slot holds one record per invocation and takes whichever record is handed to it last, with no
- * comparison of age. An overtaken build would therefore write its older snapshot over the newer one. Every build
- * takes the next value here before it starts, and the scheduler queues the record only while that value is still
- * the newest, so an overtaken build's record is dropped instead.
- *
- *
An {@link AtomicLong} rather than a {@code volatile long}: {@code ++} on a {@code volatile long} is a
- * read-modify-write, so two concurrent builds can take the same value and each conclude its own record is the
- * newest. That is the very case the check exists for, so a racy counter would guard nothing.
- */
- private final AtomicLong buildRevision = new AtomicLong();
-
- // --- Scheduling state: guarded by the scheduler's monitor. ---
-
- /**
- * The latest record for this invocation that no pump has picked up yet, or {@code null} when none is queued.
- *
- *
A newer record replaces an older one here — each record is a complete snapshot, so the older one carries
- * nothing the newer one lacks. That is the whole of coalescing: one slot, on the instance, which no other
- * invocation can reach.
- *
- *
Which record is newer is decided by {@link #buildRevision}, not by the order the records reach this slot. The
- * slot itself takes the last hand-off unconditionally, and the last hand-off is not the newest build when a build
- * was overtaken by one that customer code started from inside it.
- */
- WorkflowInsightRecord record;
-
- /**
- * Completes once this invocation's latest record has been handed to every exporter; {@code null} when nothing is
- * outstanding.
- */
- CompletableFuture settled;
-
- /**
- * Whether a pump has taken this invocation's record and is handing it to the exporters right now.
- *
- * Set and cleared in the same critical sections that move {@link #record}, so "no queued record" is never
- * mistaken for "nothing outstanding" while the record is inside the exporters.
- */
- boolean exporting;
-
- /**
- * How many {@code drain} calls are waiting for this invocation right now.
- *
- *
A record with a waiter gates an invocation return, so the pump exports it before it spends a flush fan-out.
- */
- int drainWaiters;
-
- /**
- * Set once invocation end begins; never cleared. Guarded by the scheduler's monitor — the same monitor that queues
- * the record, so the check and the hand-off are one critical section — and {@code volatile} for the hook-side
- * pre-check.
- *
- *
A checkpoint that completes while the end record is being drained still delivers an operation-change hook to
- * this same instance, and that RUNNING snapshot must not supersede the final record.
- *
- *
This orders RUNNING records against the final record; {@link #buildRevision} orders RUNNING records against
- * each other. Neither covers the other's case. A boolean cannot say which of two RUNNING builds is newer, and the
- * revision cannot reject a RUNNING record that follows the final one, because the final record is queued without a
- * revision check. See {@link ExportScheduler#closeAndSchedule}.
- */
- volatile boolean closed;
-
- /**
- * Creates the instance that serves one invocation. Identity comes from {@code info} rather than from the first
- * hook, so every field a record is keyed by exists before any hook can fire.
- *
- * @throws RuntimeException if the invocation has no usable execution ARN; the SDK contains that exactly as it
- * contains a hook failure, by skipping this plugin for the invocation
- */
- InsightPlugin(InsightSettings settings, ExportScheduler scheduler, InvocationInfo info) {
- this.settings = settings;
- this.scheduler = scheduler;
- this.executionArn = info.durableExecutionArn();
- this.arn = ArnParser.parse(executionArn);
- this.startTime = info.executionStartTime();
- this.sampledIn = WorkflowInsight.shouldSample(executionArn, settings.samplingRate);
- }
-
- /** Test seam: waits until every scheduled record has been handed to the exporters. */
- void drainExports() {
- scheduler.drainAll();
- }
-
- @Override
- public void onInvocationStart(InvocationInfo info) {
- try {
- scheduler.throwIfFailed();
- if (!sampledIn) {
- return;
- }
- // Detach the execution input from the live handler value immediately, before the user handler or any
- // content transform can mutate it. This raw, detached snapshot is the single source of truth for input
- // on every emission (start / change / end); each build hands transforms a separate defensive copy so a
- // mutating transform cannot corrupt it. Guard the snapshot: a Throwable here (e.g. a payload whose
- // serialization overflows the stack) must omit the captured input, never fail the user handler.
- try {
- cachedInput = Json.deepCopyContent(info.executionInput());
- } catch (Throwable t) {
- reportFailure("failed to snapshot execution input; omitting input", t);
- cachedInput = null;
- }
- if (settings.emitMode == WorkflowInsightConfig.EmitMode.ON_CHANGE) {
- // The revision is taken before the build, never after. Customer code runs inside buildRecord and can
- // re-enter a hook of this instance, which builds a newer record; a revision read afterwards would
- // already be that newer build's, and this older record would pass the check and overwrite it.
- long revision = beginBuild();
- scheduler.scheduleIfNotSuperseded(
- this, buildRecord("RUNNING", info.operations(), null, cachedInput, null, null), revision);
- }
- } catch (Throwable t) {
- reportFailure("onInvocationStart failed", t);
- }
- }
-
- @Override
- public void onOperationChange(OperationChangeInfo info) {
- try {
- scheduler.throwIfFailed();
- if (settings.emitMode != WorkflowInsightConfig.EmitMode.ON_CHANGE || !sampledIn) {
- return;
- }
- // Lock-free pre-check: this invocation's end may already have begun, in which case no RUNNING snapshot may
- // follow the final record. The authoritative check is made again under the scheduler's lock below.
- if (closed) {
- return;
- }
- long revision = beginBuild();
- scheduler.scheduleIfNotSuperseded(
- this, buildRecord("RUNNING", info.operations(), null, cachedInput, null, null), revision);
- } catch (Throwable t) {
- reportFailure("onOperationChange failed", t);
- }
- }
-
- // onInvocationEnd is the hook the SDK awaits, so it is where the export queue is drained before the invocation
- // returns; this guarantees the final record (scheduled above the drain) is delivered. The drain and flush run
- // in finally so they also cover the paths where record construction fails.
- @Override
- public void onInvocationEnd(InvocationEndInfo info) {
- try {
- scheduler.throwIfFailed();
- String status = WorkflowInsight.mapStatus(info.invocationStatus());
- boolean isTerminal = "SUCCEEDED".equals(status) || "FAILED".equals(status);
- boolean isFailure = "FAILED".equals(status);
- boolean shouldEmit;
- switch (settings.emitMode) {
- case ON_CHANGE:
- shouldEmit = true;
- break;
- case ON_FAILURE:
- shouldEmit = isFailure;
- break;
- case ON_COMPLETE:
- default:
- shouldEmit = isTerminal;
- break;
- }
-
- WorkflowInsightRecord finalRecord = null;
- if (sampledIn && shouldEmit) {
- // No build revision is taken here. Customer code running inside this build can start a newer RUNNING
- // build, which would make a revision taken here stale, and a checked hand-off would then drop the final
- // record and leave a RUNNING snapshot as this execution's last exported state. The final record is
- // instead ordered by `closed`, which closeAndSchedule sets in the same critical section that queues it.
- finalRecord = buildRecord(
- status,
- info.operations(),
- Instant.now(),
- cachedInput,
- info.executionResult(),
- info.executionError());
- }
- // Close before the drain below: an operation-change hook arriving from a checkpoint that completes
- // during the drain is rejected, so no RUNNING snapshot can follow (or replace) the final record.
- scheduler.closeAndSchedule(this, finalRecord);
- } catch (Throwable t) {
- // Ordinary hook/export failures remain isolated. Fatal VM/thread termination
- // marks the scheduler failed and escapes without further exporter work.
- reportFailure("onInvocationEnd failed", t);
- } finally {
- // If record construction failed above, this instance is still open: close it so a late change hook cannot
- // schedule into the drain. Idempotent when already closed.
- scheduler.closeAndSchedule(this, null);
- // Sampled-out invocations never schedule a record, so there is nothing to drain or flush. The drain needs
- // no lookup: this instance is the thing whose record it waits for.
- if (sampledIn) {
- drainAndFlush();
- }
- // Nothing is released here. There is no per-execution entry to remove — this instance is the state, the SDK
- // drops it when the invocation returns, and a suspended execution that resumes in the same container is
- // served by a new instance built from the resume's own InvocationInfo (same stable start time, same
- // deterministic sampling decision, its own input snapshot). A plugin failure therefore cannot turn into a
- // state leak, because there is no place a leak could accumulate.
- }
- }
-
- /**
- * Waits for this invocation's scheduled record to reach the exporters, then flushes each exporter once. The wait is
- * per invocation: another execution running in the same environment can never displace this record, so this always
- * returns having delivered this invocation's latest snapshot. It is not insulated from the queue, though — one pump
- * exports serially, so records another execution had already queued ahead of this one are exported first and this
- * drain waits for them too.
- *
- *
The flush goes through the scheduler's queue and is served by that same pump, between records, so no exporter
- * ever sees this invocation's {@code flush()} overlap another's {@code export()}. Invocation ends that overlap
- * share one flush: the cadence the exporter contract promises is at most one flush per sampled-in invocation end,
- * not exactly one.
- */
- private void drainAndFlush() {
- try {
- scheduler.drain(this);
- } catch (Throwable t) {
- reportFailure("failed to drain export scheduler", t);
- }
- try {
- scheduler.flush();
- } catch (Throwable t) {
- reportFailure("exporter flush failed", t);
- }
- }
-
- // --- Record building. ---
-
- /** Marks a fatal hook/logging failure before the invocation's cleanup can schedule more work. */
- private void reportFailure(String message, Throwable failure) {
- try {
- WorkflowInsight.logSafely(message, failure);
- } catch (Throwable reportingFailure) {
- Error fatal = FatalErrors.find(reportingFailure);
- if (fatal != null) scheduler.fail(fatal);
- throw reportingFailure;
- }
- }
-
- /** Starts a record build and returns the revision that identifies it. */
- private long beginBuild() {
- return buildRevision.incrementAndGet();
- }
-
- /**
- * Whether the identified build is still the newest one this invocation has started.
- *
- *
Read by the scheduler inside the critical section that queues the record, so a record that passes cannot be
- * queued after a record that supersedes it. A build that starts after the check passes still supersedes this one:
- * its record is handed over later and replaces this one in the slot, which is the order the slot should have.
- */
- boolean isNewestBuild(long revision) {
- return buildRevision.get() == revision;
- }
-
- private WorkflowInsightRecord buildRecord(
- String status,
- Map operations,
- Instant endTime,
- Object input,
- Object output,
- Throwable error) {
- ContentConfig content = settings.content;
- WorkflowInsightRecord record = new WorkflowInsightRecord();
- record.emittedAt = Instant.now().toString();
- record.executionArn = executionArn;
- record.executionName = WorkflowInsight.emptyToNull(arn.executionName());
- record.functionName = arn.functionName();
- record.functionQualifier = arn.qualifier();
- record.region = arn.region();
- record.accountId = arn.accountId();
- record.status = status;
- record.startTime = startTime != null ? startTime.toString() : null;
- if (endTime != null) {
- record.endTime = endTime.toString();
- if (startTime != null) {
- record.durationMs = endTime.toEpochMilli() - startTime.toEpochMilli();
- }
- }
- record.input = WorkflowInsight.applyDataContent(
- "input",
- input,
- content == null || content.includeInput(),
- content == null ? null : content.inputTransform());
- record.output = WorkflowInsight.applyDataContent(
- "output",
- output,
- content == null || content.includeOutput(),
- content == null ? null : content.outputTransform());
- // Honor ContentConfig.includeErrors for the execution-level error exactly as for operation-level errors
- // below: with includeErrors(false) no execution error is emitted, so a sensitive failure message never
- // reaches a record. Without this gate the execution error leaked even when errors were disabled.
- if (settings.includeErrors && error != null) {
- record.error = WorkflowInsight.toErrorInfo(error);
- }
- record.operations = buildOperationRecords(operations);
- return record;
- }
-
- private List buildOperationRecords(Map operations) {
- List out = new ArrayList<>();
- if (operations == null) {
- return out;
- }
- // The hook contract supplies a map with no iteration-order guarantee (the core snapshot originates from a
- // concurrent map). Sort by startTimestamp ascending (null timestamps last), then by a stable operation id
- // tie-breaker, so the emitted operations array is deterministic and OperationsIndex's "latest occurrence"
- // scalar fields reflect true chronological order rather than arbitrary map iteration order.
- List items = new ArrayList<>(operations.values());
- items.sort(Comparator.comparing(
- OperationChangeItemInfo::startTimestamp, Comparator.nullsLast(Comparator.naturalOrder()))
- .thenComparing(OperationChangeItemInfo::id, Comparator.nullsLast(Comparator.naturalOrder())));
- for (OperationChangeItemInfo item : items) {
- // The SDK core tracks the invocation/execution itself as a pseudo-entry of type EXECUTION; it is not a
- // customer operation and the record already carries the execution status/timing at top level.
- if ("EXECUTION".equals(item.type())) {
- continue;
- }
- // Unnamed operations can't be targeted or keyed — excluded by default (matches JS `if (!op.name)`).
- if (item.name() == null) {
- continue;
- }
- // top-level detail drops anything nested under a context (parallel branches, map items, nested steps).
- if (settings.topLevelOnly && item.parentId() != null) {
- continue;
- }
- OperationOverride override = settings.overridesByName.get(item.name());
- if (override != null && override.isExclude()) {
- continue;
- }
- OperationRecord rec = new OperationRecord()
- .id(item.id())
- .name(item.name())
- .type(item.type())
- .subType(item.subType())
- .parentId(item.parentId())
- .status(item.status() != null ? item.status().toString() : "UNKNOWN")
- .startTime(
- item.startTimestamp() != null
- ? item.startTimestamp().toString()
- : null)
- .endTime(item.endTimestamp() != null ? item.endTimestamp().toString() : null)
- .attempt(item.attempt());
- if (item.startTimestamp() != null && item.endTimestamp() != null) {
- rec.durationMs(item.endTimestamp().toEpochMilli()
- - item.startTimestamp().toEpochMilli());
- }
- if (settings.includeErrors && item.error() != null) {
- rec.error(WorkflowInsight.toErrorInfo(item.error()));
- }
- // Results are omitted unless an override explicitly opts in via a transform (matches JS).
- if (override != null && override.result() != null) {
- rec.result(WorkflowInsight.applyResultOverride(override.result(), item.result()));
- }
- out.add(rec);
- }
- return out;
- }
-
- @Override
- public String toString() {
- return "InsightPlugin[" + executionArn + "]";
- }
-}
diff --git a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightSettings.java b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightSettings.java
deleted file mode 100644
index 8729b8bec..000000000
--- a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightSettings.java
+++ /dev/null
@@ -1,59 +0,0 @@
-// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
-// SPDX-License-Identifier: Apache-2.0
-package software.amazon.lambda.durable.insight;
-
-import java.util.Collections;
-import java.util.LinkedHashMap;
-import java.util.List;
-import java.util.Map;
-import software.amazon.lambda.durable.insight.exporters.LambdaLogExporter;
-
-/**
- * The plugin's configuration, resolved once and then immutable: everything a record's shape depends on that does not
- * depend on which invocation is being observed.
- *
- * This belongs to the execution environment, not to an invocation. {@link WorkflowInsight#workflowInsight} resolves
- * it once and the factory it returns hands the same instance to every {@link InsightPlugin} it creates, so resolving
- * defaults, validating the sampling rate and indexing the operation overrides happen once per environment rather than
- * once per invocation.
- */
-final class InsightSettings {
-
- /** Sampling rate, clamped to [0, 1]; the per-invocation decision is derived from it and the execution ARN. */
- final double samplingRate;
-
- final WorkflowInsightConfig.EmitMode emitMode;
-
- /** True when nested operations (parallel branches, map items, nested steps) are dropped from the record. */
- final boolean topLevelOnly;
-
- final boolean includeErrors;
-
- /** May be null, which means "every default": include input, output and errors, with no transforms. */
- final ContentConfig content;
-
- /** Operation overrides indexed by operation name, in declaration order. */
- final Map overridesByName;
-
- /** The configured exporters, or the default single {@link LambdaLogExporter} when none were configured. */
- final List exporters;
-
- InsightSettings(WorkflowInsightConfig config) {
- this.samplingRate = WorkflowInsight.resolveSamplingRate(config.samplingRate());
- this.emitMode = config.emitMode() != null ? config.emitMode() : WorkflowInsightConfig.EmitMode.ON_COMPLETE;
- this.topLevelOnly = config.operationDetail() != WorkflowInsightConfig.OperationDetail.FULL_TREE;
- this.content = config.content();
- this.includeErrors = content == null || content.includeErrors();
- Map overrides = new LinkedHashMap<>();
- if (content != null) {
- for (OperationOverride override : content.overrides()) {
- overrides.put(override.operationName(), override);
- }
- }
- // Unmodifiable wrapper rather than Map.copyOf: declaration order is preserved and an override with a null
- // operation name is tolerated exactly as the mutable map tolerated it.
- this.overridesByName = Collections.unmodifiableMap(overrides);
- this.exporters =
- config.exporters().isEmpty() ? List.of(new LambdaLogExporter()) : List.copyOf(config.exporters());
- }
-}
diff --git a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/Json.java b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/Json.java
index fd843df59..dc72a7092 100644
--- a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/Json.java
+++ b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/Json.java
@@ -15,7 +15,6 @@
import java.util.List;
import java.util.Map;
import software.amazon.lambda.durable.annotations.Experimental;
-import software.amazon.lambda.durable.insight.internal.FatalErrors;
/** Minimal JSON helper for emitting insight records and measuring their serialized size. */
@Experimental
@@ -41,7 +40,6 @@ public static String stringify(Object value) {
try {
return MAPPER.writeValueAsString(value);
} catch (JsonProcessingException e) {
- FatalErrors.rethrow(e);
throw new IllegalStateException("failed to serialize insight record", e);
}
}
@@ -53,7 +51,6 @@ public static String prettyStringify(Object value) {
try {
return MAPPER.writer(PRETTY_PRINTER).writeValueAsString(value);
} catch (JsonProcessingException e) {
- FatalErrors.rethrow(e);
throw new IllegalStateException("failed to serialize insight record", e);
}
}
@@ -71,7 +68,6 @@ public static Integer byteSize(Object value) {
try {
return MAPPER.writeValueAsString(value).getBytes(StandardCharsets.UTF_8).length;
} catch (JsonProcessingException e) {
- FatalErrors.rethrow(e);
return null;
}
}
@@ -111,8 +107,6 @@ static Object deepCopyContent(Object value) {
try {
return MAPPER.convertValue(value, Object.class);
} catch (IllegalArgumentException e) {
- // convertValue transports Jackson serialization failures in IllegalArgumentException.
- if (e.getCause() instanceof JsonProcessingException) FatalErrors.rethrow(e.getCause());
throw new IllegalStateException("failed to copy insight record content", e);
}
}
diff --git a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/WorkflowInsight.java b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/WorkflowInsight.java
index ef50acaac..531bbc475 100644
--- a/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/WorkflowInsight.java
+++ b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/WorkflowInsight.java
@@ -2,7 +2,13 @@
// SPDX-License-Identifier: Apache-2.0
package software.amazon.lambda.durable.insight;
-import com.fasterxml.jackson.core.JsonProcessingException;
+import java.time.Instant;
+import java.util.ArrayList;
+import java.util.Comparator;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.function.Function;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -10,11 +16,12 @@
import software.amazon.lambda.durable.annotations.Experimental;
import software.amazon.lambda.durable.exception.DurableOperationException;
import software.amazon.lambda.durable.exception.UnrecoverableDurableExecutionException;
-import software.amazon.lambda.durable.insight.internal.FatalErrors;
-import software.amazon.lambda.durable.plugin.DurableExecutionPluginFactory;
+import software.amazon.lambda.durable.insight.exporters.LambdaLogExporter;
+import software.amazon.lambda.durable.plugin.DurableExecutionPlugin;
import software.amazon.lambda.durable.plugin.InvocationEndInfo;
import software.amazon.lambda.durable.plugin.InvocationInfo;
import software.amazon.lambda.durable.plugin.InvocationStatus;
+import software.amazon.lambda.durable.plugin.OperationChangeInfo;
import software.amazon.lambda.durable.plugin.OperationChangeItemInfo;
/**
@@ -28,18 +35,12 @@
* {@link OperationChangeItemInfo#result()}; these are the fields PR #618 surfaced on the hook records, so
* {@code input}, {@code output}, and operation {@code result} are now populated exactly as in the JS plugin.
*
- * {@link #workflowInsight} returns a {@link DurableExecutionPluginFactory}, so the SDK creates one
- * {@link InsightPlugin} per Lambda invocation and drops it when the invocation returns. Everything about an execution —
- * the stable start time, the parsed ARN, the one-time sampling decision, the detached input snapshot, the queued
- * record, the drain signal — is therefore a plain field of that instance. Nothing is keyed by execution ARN, and there
- * is no per-execution entry to remove at invocation end, so a suspended execution cannot leak one for the lifetime of a
- * warm container. Nothing is lost across a resume either: the resume's own {@link InvocationInfo} carries the same
- * stable start time, the sampling decision is deterministic in the ARN, and the input snapshot is taken again from
- * {@link InvocationInfo#executionInput()}.
- *
- *
What belongs to the execution environment rather than to an invocation stays in the factory: the resolved
- * {@link InsightSettings}, the exporters, and the {@link ExportScheduler} that serializes exports across every
- * execution the environment hosts.
+ *
Per-execution state (keyed by execution ARN) holds only the stable start time, the parsed ARN, the one-time
+ * sampling decision, and a detached snapshot of the execution input. State is removed on every {@code onInvocationEnd}
+ * — including non-terminal PENDING/RETRYING suspends — so a suspended execution never leaks a retained entry for the
+ * lifetime of a warm container. Nothing is lost across a resume: the next invocation recreates the same stable start
+ * time from {@link InvocationInfo#executionStartTime()}, the same sampling decision deterministically from the ARN, and
+ * the input snapshot from {@link InvocationInfo#executionInput()}.
*/
@Experimental
public final class WorkflowInsight {
@@ -48,51 +49,378 @@ public final class WorkflowInsight {
private WorkflowInsight() {}
- /**
- * Creates a Workflow Insight plugin factory from the given config. Mirrors the JS {@code workflowInsight(config)}.
- *
- *
The configuration is resolved once, here; the exporters and the scheduler that serializes exports across them
- * are created once, here. The returned factory then builds one plugin instance per invocation, which is what lets
- * that instance hold its execution's state in plain fields.
- *
- * @param config the plugin configuration
- * @return a factory to hand to {@code DurableConfig.Builder.withPlugins}
- */
- public static DurableExecutionPluginFactory workflowInsight(WorkflowInsightConfig config) {
- InsightSettings settings = new InsightSettings(config);
- ExportScheduler scheduler = new ExportScheduler(
- settings.exporters, WorkflowInsight::exportRecord, t -> logSafely("export scheduling failed", t));
- return info -> new InsightPlugin(settings, scheduler, info);
+ /** Creates a Workflow Insight plugin from the given config. Mirrors the JS {@code workflowInsight(config)}. */
+ public static DurableExecutionPlugin workflowInsight(WorkflowInsightConfig config) {
+ return new InsightPlugin(config);
}
- // --- helpers ---
+ /** Per-execution state, keyed by execution ARN, to prevent warm-container bleed and handle resume. */
+ private static final class ExecutionState {
+ final Instant startTime;
+ final ArnParser arn;
+ final boolean sampledIn;
+ volatile Object cachedInput;
- /** Shapes and exports one record to one exporter; runs on a scheduler worker, never on an SDK hook thread. */
- static void exportRecord(WorkflowInsightRecord record, InsightExporter exporter) {
- try {
- // Give each exporter its own deep copy: truncation returns the original record when it already fits,
- // so without this a custom exporter that mutates operations or nested content would corrupt every
- // other exporter's view of the same record.
- WorkflowInsightRecord isolated = record.deepCopy();
- WorkflowInsightRecord shaped =
- Truncation.truncateRecord(isolated, exporter.maxRecordSizeBytes(), exporter::render);
- exporter.export(shaped);
- } catch (Throwable t) {
- // Ordinary Errors (such as missing optional dependencies) are isolated too.
- // logSafely rethrows VM/thread termination before logging or continuing.
- logSafely("exporter failed", t);
+ /**
+ * Set once invocation end begins; guarded by {@code this}. A checkpoint that completes while the end record is
+ * being drained still delivers an operation-change hook, and that RUNNING snapshot must not supersede the final
+ * record.
+ */
+ boolean closed;
+
+ ExecutionState(Instant startTime, ArnParser arn, boolean sampledIn) {
+ this.startTime = startTime;
+ this.arn = arn;
+ this.sampledIn = sampledIn;
+ }
+
+ /** Schedules the record unless the invocation has already ended; the check and the hand-off are atomic. */
+ boolean scheduleIfOpen(ExportScheduler scheduler, WorkflowInsightRecord record) {
+ synchronized (this) {
+ if (closed) {
+ return false;
+ }
+ scheduler.schedule(record);
+ return true;
+ }
+ }
+
+ /** Marks the invocation ended and, when given a record, schedules it as the last one for this execution. */
+ void closeAndSchedule(ExportScheduler scheduler, WorkflowInsightRecord finalRecord) {
+ synchronized (this) {
+ closed = true;
+ if (finalRecord != null) {
+ scheduler.schedule(finalRecord);
+ }
+ }
}
}
+ static final class InsightPlugin implements DurableExecutionPlugin {
+ private final double samplingRate;
+ private final WorkflowInsightConfig.EmitMode emitMode;
+ private final boolean topLevelOnly;
+ private final boolean includeErrors;
+ private final ContentConfig content;
+ private final Map overridesByName = new LinkedHashMap<>();
+ private final List exporters;
+ private final ExportScheduler scheduler;
+
+ private final Map byArn = new ConcurrentHashMap<>();
+
+ /** Test seam: number of live per-execution state entries retained across invocations. */
+ int retainedStateCount() {
+ return byArn.size();
+ }
+
+ /** Test seam: waits until every scheduled record has been handed to the exporters. */
+ void drainExports() {
+ scheduler.drain();
+ }
+
+ InsightPlugin(WorkflowInsightConfig config) {
+ this.samplingRate = resolveSamplingRate(config.samplingRate());
+ this.emitMode = config.emitMode() != null ? config.emitMode() : WorkflowInsightConfig.EmitMode.ON_COMPLETE;
+ this.topLevelOnly = config.operationDetail() != WorkflowInsightConfig.OperationDetail.FULL_TREE;
+ this.content = config.content();
+ this.includeErrors = content == null || content.includeErrors();
+ if (content != null) {
+ for (OperationOverride o : content.overrides()) {
+ overridesByName.put(o.operationName(), o);
+ }
+ }
+ this.exporters =
+ config.exporters().isEmpty() ? List.of(new LambdaLogExporter()) : List.copyOf(config.exporters());
+ this.scheduler =
+ new ExportScheduler(exporters, this::exportRecord, t -> logSafely("export scheduling failed", t));
+ }
+
+ private ExecutionState getState(String arn, Instant startTime) {
+ return byArn.computeIfAbsent(
+ arn, a -> new ExecutionState(startTime, ArnParser.parse(a), shouldSample(a, samplingRate)));
+ }
+
+ @Override
+ public void onInvocationStart(InvocationInfo info) {
+ try {
+ ExecutionState state = getState(info.durableExecutionArn(), info.executionStartTime());
+ if (!state.sampledIn) {
+ return;
+ }
+ // Detach the execution input from the live handler value immediately, before the user handler or any
+ // content transform can mutate it. This raw, detached snapshot is the single source of truth for input
+ // on every emission (start / change / end); each build hands transforms a separate defensive copy so a
+ // mutating transform cannot corrupt it. Guard the snapshot: a Throwable here (e.g. a payload whose
+ // serialization overflows the stack) must omit the captured input, never fail the user handler.
+ try {
+ state.cachedInput = Json.deepCopyContent(info.executionInput());
+ } catch (Throwable t) {
+ logSafely("failed to snapshot execution input; omitting input", t);
+ state.cachedInput = null;
+ }
+ if (emitMode == WorkflowInsightConfig.EmitMode.ON_CHANGE) {
+ scheduler.schedule(buildRecord(
+ state,
+ info.durableExecutionArn(),
+ "RUNNING",
+ info.operations(),
+ null,
+ state.cachedInput,
+ null,
+ null));
+ }
+ } catch (Throwable t) {
+ logSafely("onInvocationStart failed", t);
+ }
+ }
+
+ @Override
+ public void onOperationChange(OperationChangeInfo info) {
+ try {
+ if (emitMode != WorkflowInsightConfig.EmitMode.ON_CHANGE) {
+ return;
+ }
+ ExecutionState state = byArn.get(info.durableExecutionArn());
+ if (state == null || !state.sampledIn) {
+ return;
+ }
+ state.scheduleIfOpen(
+ scheduler,
+ buildRecord(
+ state,
+ info.durableExecutionArn(),
+ "RUNNING",
+ info.operations(),
+ null,
+ state.cachedInput,
+ null,
+ null));
+ } catch (Throwable t) {
+ logSafely("onOperationChange failed", t);
+ }
+ }
+
+ // onInvocationEnd is the hook the SDK awaits, so it is where the export queue is drained before the invocation
+ // returns; this guarantees the final record (scheduled above the drain) is delivered. The drain and flush run
+ // in finally so they also cover the paths where record construction fails.
+ @Override
+ public void onInvocationEnd(InvocationEndInfo info) {
+ ExecutionState state = null;
+ try {
+ state = getState(info.durableExecutionArn(), info.executionStartTime());
+ String status = mapStatus(info.invocationStatus());
+ boolean isTerminal = "SUCCEEDED".equals(status) || "FAILED".equals(status);
+ boolean isFailure = "FAILED".equals(status);
+ boolean shouldEmit;
+ switch (emitMode) {
+ case ON_CHANGE:
+ shouldEmit = true;
+ break;
+ case ON_FAILURE:
+ shouldEmit = isFailure;
+ break;
+ case ON_COMPLETE:
+ default:
+ shouldEmit = isTerminal;
+ break;
+ }
+
+ WorkflowInsightRecord finalRecord = null;
+ if (state.sampledIn && shouldEmit) {
+ finalRecord = buildRecord(
+ state,
+ info.durableExecutionArn(),
+ status,
+ info.operations(),
+ Instant.now(),
+ state.cachedInput,
+ info.executionResult(),
+ info.executionError());
+ }
+ // Close before the drain below: an operation-change hook arriving from a checkpoint that completes
+ // during the drain is rejected, so no RUNNING snapshot can follow (or replace) the final record.
+ state.closeAndSchedule(scheduler, finalRecord);
+ } catch (Throwable t) {
+ // A plugin failure at end-of-invocation (record construction, transforms, truncation, export/flush,
+ // or optional exporter class linkage) must never disrupt durable execution.
+ logSafely("onInvocationEnd failed", t);
+ } finally {
+ // If record construction failed above, the state is still open: close it so a late change hook cannot
+ // schedule into the drain. Idempotent when already closed.
+ if (state != null) {
+ state.closeAndSchedule(scheduler, null);
+ }
+ // Sampled-out executions never schedule a record, so there is nothing to drain or flush. If the state
+ // lookup itself failed, drain anyway: it is a no-op when idle and otherwise delivers what is pending.
+ if (state == null || state.sampledIn) {
+ drainAndFlush();
+ }
+ // Remove per-execution state on EVERY invocation end, including non-terminal PENDING/RETRYING suspends,
+ // once any emission work above is done. Nothing durable is lost: the next invocation's onInvocation
+ // start recreates the stable startTime from InvocationInfo.executionStartTime() (stable across
+ // resumes),
+ // the one-time sampling decision deterministically from the ARN, and the input snapshot from
+ // InvocationInfo.executionInput(). Retaining state instead leaked one entry per suspended execution for
+ // the lifetime of the warm container. This runs even if emission above threw, so a plugin failure can
+ // never turn into a state leak.
+ byArn.remove(info.durableExecutionArn());
+ }
+ }
+
+ /** Waits for every scheduled record to reach the exporters, then flushes each exporter once, concurrently. */
+ private void drainAndFlush() {
+ try {
+ scheduler.drain();
+ } catch (Throwable t) {
+ logSafely("failed to drain export scheduler", t);
+ }
+ try {
+ scheduler.flushAll();
+ } catch (Throwable t) {
+ logSafely("exporter flush failed", t);
+ }
+ }
+
+ /** Shapes and exports one record to one exporter; runs on a scheduler worker, never on an SDK hook thread. */
+ private void exportRecord(WorkflowInsightRecord record, InsightExporter exporter) {
+ try {
+ // Give each exporter its own deep copy: truncation returns the original record when it already fits,
+ // so without this a custom exporter that mutates operations or nested content would corrupt every
+ // other exporter's view of the same record.
+ WorkflowInsightRecord isolated = record.deepCopy();
+ WorkflowInsightRecord shaped =
+ Truncation.truncateRecord(isolated, exporter.maxRecordSizeBytes(), exporter::render);
+ exporter.export(shaped);
+ } catch (Throwable t) {
+ // Catch Throwable, not just RuntimeException: deep copy, truncation, an exporter's render/export, or
+ // the linkage of an optional exporter class (a NoClassDefFoundError when the S3 / CloudWatch SDK is
+ // absent) can each fail with an Error. Isolating every Throwable here guarantees one failing exporter
+ // cannot affect the others, nor disrupt the execution.
+ logSafely("exporter failed", t);
+ }
+ }
+
+ private WorkflowInsightRecord buildRecord(
+ ExecutionState state,
+ String arn,
+ String status,
+ Map operations,
+ Instant endTime,
+ Object input,
+ Object output,
+ Throwable error) {
+ WorkflowInsightRecord record = new WorkflowInsightRecord();
+ ArnParser a = state.arn;
+ record.emittedAt = Instant.now().toString();
+ record.executionArn = arn;
+ record.executionName = emptyToNull(a.executionName());
+ record.functionName = a.functionName();
+ record.functionQualifier = a.qualifier();
+ record.region = a.region();
+ record.accountId = a.accountId();
+ record.status = status;
+ record.startTime = state.startTime != null ? state.startTime.toString() : null;
+ if (endTime != null) {
+ record.endTime = endTime.toString();
+ if (state.startTime != null) {
+ record.durationMs = endTime.toEpochMilli() - state.startTime.toEpochMilli();
+ }
+ }
+ record.input = applyDataContent(
+ "input",
+ input,
+ content == null || content.includeInput(),
+ content == null ? null : content.inputTransform());
+ record.output = applyDataContent(
+ "output",
+ output,
+ content == null || content.includeOutput(),
+ content == null ? null : content.outputTransform());
+ // Honor ContentConfig.includeErrors for the execution-level error exactly as for operation-level errors
+ // below: with includeErrors(false) no execution error is emitted, so a sensitive failure message never
+ // reaches a record. Without this gate the execution error leaked even when errors were disabled.
+ if (includeErrors && error != null) {
+ record.error = toErrorInfo(error);
+ }
+ record.operations = buildOperationRecords(operations);
+ return record;
+ }
+
+ private List buildOperationRecords(Map operations) {
+ List out = new ArrayList<>();
+ if (operations == null) {
+ return out;
+ }
+ // The hook contract supplies a map with no iteration-order guarantee (the core snapshot originates from a
+ // concurrent map). Sort by startTimestamp ascending (null timestamps last), then by a stable operation id
+ // tie-breaker, so the emitted operations array is deterministic and OperationsIndex's "latest occurrence"
+ // scalar fields reflect true chronological order rather than arbitrary map iteration order.
+ List items = new ArrayList<>(operations.values());
+ items.sort(Comparator.comparing(
+ OperationChangeItemInfo::startTimestamp, Comparator.nullsLast(Comparator.naturalOrder()))
+ .thenComparing(OperationChangeItemInfo::id, Comparator.nullsLast(Comparator.naturalOrder())));
+ for (OperationChangeItemInfo item : items) {
+ // The SDK core tracks the invocation/execution itself as a pseudo-entry of type EXECUTION; it is not a
+ // customer operation and the record already carries the execution status/timing at top level.
+ if ("EXECUTION".equals(item.type())) {
+ continue;
+ }
+ // Unnamed operations can't be targeted or keyed — excluded by default (matches JS `if (!op.name)`).
+ if (item.name() == null) {
+ continue;
+ }
+ // top-level detail drops anything nested under a context (parallel branches, map items, nested steps).
+ if (topLevelOnly && item.parentId() != null) {
+ continue;
+ }
+ OperationOverride override = overridesByName.get(item.name());
+ if (override != null && override.isExclude()) {
+ continue;
+ }
+ OperationRecord rec = new OperationRecord()
+ .id(item.id())
+ .name(item.name())
+ .type(item.type())
+ .subType(item.subType())
+ .parentId(item.parentId())
+ .status(item.status() != null ? item.status().toString() : "UNKNOWN")
+ .startTime(
+ item.startTimestamp() != null
+ ? item.startTimestamp().toString()
+ : null)
+ .endTime(
+ item.endTimestamp() != null
+ ? item.endTimestamp().toString()
+ : null)
+ .attempt(item.attempt());
+ if (item.startTimestamp() != null && item.endTimestamp() != null) {
+ rec.durationMs(item.endTimestamp().toEpochMilli()
+ - item.startTimestamp().toEpochMilli());
+ }
+ if (includeErrors && item.error() != null) {
+ rec.error(toErrorInfo(item.error()));
+ }
+ // Results are omitted unless an override explicitly opts in via a transform (matches JS).
+ if (override != null && override.result() != null) {
+ rec.result(applyResultOverride(override.result(), item.result()));
+ }
+ out.add(rec);
+ }
+ return out;
+ }
+ }
+
+ // --- helpers ---
+
/**
* Applies a user-supplied result transform to an operation's checkpointed (serialized JSON) result. Parses the JSON
* before handing it to the transform, so the transform always receives a detached, JSON-compatible value
* (a {@code Map} for a former POJO, a {@code List} for an array, or a scalar such as a {@code String} for a Java
* time value) — never the SDK's original Java object. The raw string is passed through only when the checkpointed
- * result is not valid JSON. Non-fatal transform failures omit the field rather than leaking the raw value and are
- * logged for diagnosis. VM/thread-termination errors propagate. Because the value is freshly parsed from the
- * immutable checkpoint string on every build, a transform that mutates its argument cannot corrupt any cached state
- * or a later emission.
+ * result is not valid JSON. User transforms are untrusted: a throwing transform (any {@link Throwable}) omits the
+ * field rather than leaking the raw value or failing the execution, and the failure is logged for diagnosis.
+ * Because the value is freshly parsed from the immutable checkpoint string on every build, a transform that mutates
+ * its argument cannot corrupt any cached state or a later emission.
*/
static Object applyResultOverride(Function