{
@Override
protected DurableConfig createConfiguration() {
- return DurableConfig.builder().withPlugins(new LoggingPlugin()).build();
+ return DurableConfig.builder().withPlugins(info -> new LoggingPlugin()).build();
}
@Override
diff --git a/insight-plugin/README.md b/insight-plugin/README.md
index 75d4e5132..694cf0df9 100644
--- a/insight-plugin/README.md
+++ b/insight-plugin/README.md
@@ -192,24 +192,29 @@ 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.
-- **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.
+- **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.
- **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; a failing or slow exporter is logged and never blocks the others or the execution.
+ limit; ordinary exporter failures are logged without preventing the other exporters from running.
- **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 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.
+ 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.
## Conformance
diff --git a/insight-plugin/pom.xml b/insight-plugin/pom.xml
index c2cbe8df5..cc949030f 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
- 2.2.2-SNAPSHOT
+ 3.0.0-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 3a8228ec9..0ecc1e38b 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 transform that throws omits the field (the failure is logged)
- * rather than failing the execution.
+ * 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.
*/
@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 84e8be3c2..9731683a9 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,37 +2,82 @@
// SPDX-License-Identifier: Apache-2.0
package software.amazon.lambda.durable.insight;
+import java.util.ArrayDeque;
import java.util.ArrayList;
-import java.util.Iterator;
-import java.util.LinkedHashMap;
+import java.util.Deque;
+import java.util.LinkedHashSet;
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.
+ * Serializes record exports so that, at most, one export runs at a time, while keeping the records of concurrently
+ * running executions independent.
*
- * 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.
+ *
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.
*
- *
Exports are otherwise fire-and-forget; {@link #drain()} is called before the invocation returns to guarantee the
- * final record is delivered.
+ *
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.
*/
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());
@@ -48,11 +93,68 @@ 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<>();
+
/**
- * 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}.
+ * 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.
*/
- private final LinkedHashMap pending = new LinkedHashMap<>();
+ private final Deque> flushRequests = new ArrayDeque<>();
ExportScheduler(
List exporters,
@@ -72,53 +174,214 @@ final class ExportScheduler {
this.executor = executor;
}
+ // --- Scheduling. ---
+
/**
- * 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.
+ * 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(WorkflowInsightRecord record) {
+ void schedule(InsightPlugin execution, WorkflowInsightRecord record) {
CompletableFuture handle;
synchronized (this) {
- // 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;
+ 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
+ */
+ boolean scheduleIfNotSuperseded(InsightPlugin execution, WorkflowInsightRecord record, long buildRevision) {
+ CompletableFuture handle;
+ synchronized (this) {
+ if (execution.closed || !execution.isNewestBuild(buildRevision)) {
+ return false;
}
- handle = new CompletableFuture<>();
- inFlight = handle;
+ 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();
+ }
+ }
+ 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;
}
try {
executor.execute(() -> pump(handle));
} catch (Throwable t) {
- // 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.
+ 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.
synchronized (this) {
if (inFlight == handle) {
inFlight = null;
}
}
- handle.complete(null);
+ completeSignal(handle);
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 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.
+ * 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}.
*/
- void drain() {
+ 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) {
while (true) {
+ CompletableFuture signal;
CompletableFuture handle;
boolean runInline = false;
synchronized (this) {
+ throwIfFailed();
+ signal = execution.settled;
+ if (signal == null) {
+ return;
+ }
handle = inFlight;
if (handle == null) {
- if (pending.isEmpty()) {
- return;
- }
- // Records are pending with no pump running (a worker could not be started): export them here.
+ // A record is outstanding with no pump running (a worker could not be started): export it here.
handle = new CompletableFuture<>();
inFlight = handle;
runInline = true;
@@ -126,52 +389,447 @@ void drain() {
}
if (runInline) {
pump(handle);
- } else {
- handle.join();
+ 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;
+ }
+ }
+ }
+ }
+
+ /**
+ * 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 {
- // 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.
+ // 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.
while (true) {
- WorkflowInsightRecord record;
+ InsightPlugin next = null;
+ WorkflowInsightRecord record = null;
synchronized (this) {
- Iterator head = pending.values().iterator();
- if (!head.hasNext()) {
- inFlight = null;
+ throwIfFailed();
+ if (queue.isEmpty() && flushRequests.isEmpty()) {
+ if (inFlight == handle) {
+ inFlight = null;
+ }
return;
}
- record = head.next();
- head.remove();
+ 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;
+ }
}
- 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) {
- if (inFlight == handle) {
- inFlight = null;
+ 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);
}
}
- handle.complete(null);
}
}
- private static String executionKey(WorkflowInsightRecord record) {
- return record.executionArn() != null ? record.executionArn() : "";
+ /** 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);
}
/**
* 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.
+ * 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.
*/
- void flushAll() {
+ private void flushEveryExporter() {
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.
@@ -183,26 +841,52 @@ 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> settled = new ArrayList<>(exporters.size());
+ List> settledExporters = new ArrayList<>(exporters.size());
for (InsightExporter exporter : exporters) {
- Runnable task = () -> runSafely(() -> action.accept(exporter));
+ 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)));
try {
- settled.add(CompletableFuture.runAsync(task, executor));
+ settledExporters.add(CompletableFuture.runAsync(task, executor));
} catch (Throwable t) {
reportFailure(t);
task.run();
}
}
- for (CompletableFuture task : settled) {
- runSafely(task::join);
+ 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);
+ }
}
}
private void runSafely(Runnable action) {
try {
+ throwIfFailed();
action.run();
} catch (Throwable t) {
reportFailure(t);
@@ -210,10 +894,102 @@ private void runSafely(Runnable action) {
}
private void reportFailure(Throwable t) {
+ rejectFatal(t);
try {
failureHandler.accept(t);
} catch (Throwable ignored) {
- // A scheduler diagnostic must never disrupt durable execution.
+ 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;
}
+ 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 69ff39739..59f04767e 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,7 +11,34 @@ public interface InsightExporter {
/** Emits one record to the destination. */
void export(WorkflowInsightRecord record);
- /** Flushes any buffered records; no-op by default. */
+ /**
+ * 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.
+ */
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
new file mode 100644
index 000000000..922f675de
--- /dev/null
+++ b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightPlugin.java
@@ -0,0 +1,458 @@
+// 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
new file mode 100644
index 000000000..8729b8bec
--- /dev/null
+++ b/insight-plugin/src/main/java/software/amazon/lambda/durable/insight/InsightSettings.java
@@ -0,0 +1,59 @@
+// 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 dc72a7092..fd843df59 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,6 +15,7 @@
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
@@ -40,6 +41,7 @@ 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);
}
}
@@ -51,6 +53,7 @@ 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);
}
}
@@ -68,6 +71,7 @@ public static Integer byteSize(Object value) {
try {
return MAPPER.writeValueAsString(value).getBytes(StandardCharsets.UTF_8).length;
} catch (JsonProcessingException e) {
+ FatalErrors.rethrow(e);
return null;
}
}
@@ -107,6 +111,8 @@ 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 531bbc475..ef50acaac 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,13 +2,7 @@
// 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.LinkedHashMap;
-import java.util.List;
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
+import com.fasterxml.jackson.core.JsonProcessingException;
import java.util.function.Function;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -16,12 +10,11 @@
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.exporters.LambdaLogExporter;
-import software.amazon.lambda.durable.plugin.DurableExecutionPlugin;
+import software.amazon.lambda.durable.insight.internal.FatalErrors;
+import software.amazon.lambda.durable.plugin.DurableExecutionPluginFactory;
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;
/**
@@ -35,12 +28,18 @@
* {@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.
*
- * 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()}.
+ *
{@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.
*/
@Experimental
public final class WorkflowInsight {
@@ -49,378 +48,51 @@ public final class WorkflowInsight {
private WorkflowInsight() {}
- /** 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);
- }
-
- /** 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;
-
- /**
- * 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);
- }
- }
- }
+ /**
+ * 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);
}
- 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;
- }
+ // --- helpers ---
- 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;
+ /** 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);
}
}
- // --- 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. 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.
+ * 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.
*/
static Object applyResultOverride(Function