Repository navigation
fix: export synthetic OTel roots and retain READY wait progress - #767
zhongkechen wants to merge 42 commits into
Conversation
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
| /** Shutdown the checkpoint batcher. */ | ||
| @Override | ||
| public void close() { | ||
| stopCheckpointContinuations(); |
This comment was marked as outdated.
This comment was marked as outdated.
Sorry, something went wrong.
There was a problem hiding this comment.
The queued-resumption race is addressed in ec27233 with a bounded admission cut before InvocationEndInfo is constructed: new checkpoint continuations are rejected and unfinished owned continuations are signaled. Explicit continuation/handler joins and checkpoint shutdown remain in manager close after End, avoiding a join of the current publisher or serialized callback path. Owner completion can still execute synchronous dependent callbacks; this is not claimed to be universally nonblocking.
Real DurableExecutor and OTel controls reproduce the former late starts for both an unawaited condition and supported early parallel completion. Queued losers now execute only their first predicate and may legitimately remain READY/incomplete, matching the End snapshot. The awaited control completes with SUCCEEDED. A separate already-running-predicate control proves that End can enter while that predicate is held, and normal close waits for it afterward; its later end hooks and completion are intentionally retained. This does not introduce an all-work-drained or universal no-hooks-after-End guarantee. The API Javadoc and README now state that boundary explicitly.
Full Java17 verification passed 2,416 tests (31 conditional skips, zero failures/errors), 74 focused Java25 controls and 16 artifact compatibility cases, including current-publisher/finally ordering, winner/loser identity, reentry and existing cleanup controls. New-head CI remains under observation.
This comment has been minimized.
This comment has been minimized.
| /** Shutdown the checkpoint batcher. */ | ||
| @Override | ||
| public void close() { | ||
| stopCheckpointContinuations(); |
This comment was marked as outdated.
This comment was marked as outdated.
Sorry, something went wrong.
There was a problem hiding this comment.
The queued-resumption portion is addressed in ec27233, with the tested bounded admission cut and its limits explained here: #767 (comment) . New admission closes before the End snapshot and unfinished owned continuations are signaled. Explicit task/handler draining remains in close afterward; already accepted/running work may still finish later. This does not replace the selected root outcome with a late failure or introduce an all-work-drained barrier. The API/README and real early-parallel/running-predicate controls state that boundary.
This comment has been minimized.
This comment has been minimized.
| public CompletableFuture<Void> runCheckpointContinuation(BaseDurableOperation owner, Runnable continuation) { | ||
| return runCheckpointContinuation(owner, continuation, InternalExecutor.INSTANCE); |
This comment was marked as outdated.
This comment was marked as outdated.
Sorry, something went wrong.
There was a problem hiding this comment.
Fixed in 2adc34c. Step retry resumption now enters the existing operation-owned checkpoint continuation, so a configured synchronous or submit-and-wait executor cannot occupy the serialized checkpoint callback. For a retry initiated by a live attempt, the continuation waits for that attempt's worker publication and completion before registering its replacement. The existing continuation activity lease spans that transition, and a stopped owner is checked before dispatch/body entry. Step failure classification, retry policy and checkpoint contents are unchanged.
The public regression starts with a real failed step, retains its PENDING history, advances the local backend to READY, and resumes through the configured executor. Before the fix, rejection returned PENDING with backend READY; submit-and-wait blocked the batcher until the probe's bounded release. Afterward, rejection reaches the caller as retryable control with the original cause, and async/direct/submit-and-wait modes complete. Three additional live-retry gates cover old-worker exit, executor-return/publication, and both together; the continuation holds no polling monitor, preserves activity, and completed replay does not repeat the body.
Validation: 2,423 full Java17 tests (31 expected skips), 26 focused Java25 tests, and all 16 installed-artifact compatibility cases passed. The same source defect was independently reproduced and the narrow fix validated on the other affected non-held branches; no lifecycle, sampler, workflow, or conformance-scenario port is included.
This comment has been minimized.
This comment has been minimized.
| * Recovery re-exports retain identity and timestamps; Workflow carries duration and outcome. Remote parents are | ||
| * never owned. | ||
| */ | ||
| static Span startExecutionRoot( |
This comment was marked as outdated.
This comment was marked as outdated.
Sorry, something went wrong.
There was a problem hiding this comment.
Fixed in 4091401. Successful handler output preparation now has a narrow failure boundary: if the configured output SerDes or oversized-result checkpoint throws, the existing completion thread dispatches one RETRYING End, then propagates the original normalized preparation failure. This releases the retained root before invocation return. The existing admission cut and subsequent manager drain are unchanged; this does not port the different handler-thread lifecycle from related PR #771 or cover later response-stream/runtime acknowledgment failures.
The real public tests reproduce both output-SerDes and greater-than-6MB checkpoint failures in both OTel views: before the change End/root-end/flush counts were zero and the root remained recording. They now each end and flush once, retain the caller error identity, and support another invocation. Combined preparation/End controls preserve the original failure and suppressed diagnostics, with the established direct JVM-fatal precedence. Full Java17: 2,462 tests (31 expected skips); 51 focused Java25 tests and all 16 compatibility cases passed.
| if (root != null) root.end(rootTimestamp); | ||
| if (shouldFlush && provider != null) { |
This comment was marked as outdated.
This comment was marked as outdated.
Sorry, something went wrong.
There was a problem hiding this comment.
Fixed in 4091401 for both views. A root-end Exception or LinkageError now still runs the flush through the existing cleanup combiner before the plugin boundary isolates that error. Unisolated Errors retain the previous escape/skip-flush behavior; this is not a blanket finally around JVM-fatal cleanup. Earlier-primary/later-fatal and same-object suppression rules remain intact.
The public control uses an actual throwing root SpanProcessor followed by BatchSpanProcessor with automatic export deferred. Before the fix, the invocation succeeded but no completed Invocation/Workflow spans had been exported; afterward they are flushed before return. A processor that throws need not export the root itself. Seven additional root/flush priority controls plus existing partial-start/reentry/reuse controls pass. Validation: 2,462 Java17 tests (31 expected skips), 51 focused Java25 tests, and 16 installed-artifact compatibility cases.
This comment has been minimized.
This comment has been minimized.
| Throwable error, | ||
| Object executionInput, | ||
| Object executionResult) { | ||
| executionManager.beginInvocationEnd(); |
This comment was marked as outdated.
This comment was marked as outdated.
Sorry, something went wrong.
There was a problem hiding this comment.
The selected root outcome is intentionally frozen before output preparation. A failure in unawaited work does not retroactively replace that outcome; await work whose success or failure must determine the handler result. The documented admission cut is before the End snapshot, after SDK output preparation, and accepted/running handlers retain their existing drain after End. It does not promise that all work stops immediately when the handler returns. An output serializer or oversized-result checkpoint failure itself still triggers one RETRYING End and propagates its original error, as covered by the new preparation tests. Moving admission earlier would add a different resource-policy boundary, rather than correct a violation of the current one. The retained early-parallel, already-running-handler and preparation-failure controls exercise these distinct guarantees; no global winner or earlier-cut rewrite is included.
| // Register before re-reading: another checkpoint may already have delivered READY before this poll existed. | ||
| known = getOperation(); | ||
| if (isReadyOrTerminal(known)) { | ||
| update.complete(known); |
This comment was marked as outdated.
This comment was marked as outdated.
Sorry, something went wrong.
There was a problem hiding this comment.
Fixed in 098c961. Before deciding whether an empty checkpoint batch needs a backend call, CheckpointManager now prunes polling futures that are already completed, failed or cancelled, under the existing pollingFutures monitor. Live consumers, including another future for the same operation ID, and real queued checkpoint updates are retained. This does not cancel an RPC that has already been admitted or change operation state/result semantics.
The deterministic component tests complete the exposed poll future exactly as the READY recheck does, then force the real delayed dispatcher to flush without closing CheckpointManager. Before the fix, all four already-done cases made an extra empty API call. They now make zero calls and acquire no checkpoint lease; four live-poll/real-update controls still make one balanced request and retain their results. Existing response-gate, READY, worker-handoff and shutdown controls also pass.
Validation on this branch: 2,474 full Java17 tests (31 expected skips), 56 focused Java25 tests and all 16 installed-artifact compatibility cases. The same source defect was independently reproduced and validated on the other affected non-held Java branches; workflow references, lifecycle ordering and the selected-outcome boundary are unchanged.
This comment has been minimized.
This comment has been minimized.
| DurableExecutionOutput output; | ||
| try { |
This comment was marked as outdated.
This comment was marked as outdated.
Sorry, something went wrong.
There was a problem hiding this comment.
The selected root outcome is intentionally frozen before output preparation. A failure in unawaited work does not retroactively replace that outcome; await work whose success or failure must determine the handler result. The documented admission cut is before the End snapshot, after SDK output preparation, and accepted/running handlers retain their existing drain after End. It does not promise that all work stops immediately when the handler returns. An output serializer or oversized-result checkpoint failure itself still triggers one RETRYING End and propagates its original error, as covered by the new preparation tests. Moving admission earlier would add a different resource-policy boundary, rather than correct a violation of the current one. The retained early-parallel, already-running-handler and preparation-failure controls exercise these distinct guarantees; no global winner or earlier-cut rewrite is included.
| // Register before re-reading: another checkpoint may already have delivered READY before this poll existed. | ||
| known = getOperation(); | ||
| if (isReadyOrTerminal(known)) { | ||
| update.complete(known); |
This comment was marked as outdated.
This comment was marked as outdated.
Sorry, something went wrong.
There was a problem hiding this comment.
Fixed in 098c961. Before deciding whether an empty checkpoint batch needs a backend call, CheckpointManager now prunes polling futures that are already completed, failed or cancelled, under the existing pollingFutures monitor. Live consumers, including another future for the same operation ID, and real queued checkpoint updates are retained. This does not cancel an RPC that has already been admitted or change operation state/result semantics.
The deterministic component tests complete the exposed poll future exactly as the READY recheck does, then force the real delayed dispatcher to flush without closing CheckpointManager. Before the fix, all four already-done cases made an extra empty API call. They now make zero calls and acquire no checkpoint lease; four live-poll/real-update controls still make one balanced request and retain their results. Existing response-gate, READY, worker-handoff and shutdown controls also pass.
Validation on this branch: 2,474 full Java17 tests (31 expected skips), 56 focused Java25 tests and all 16 installed-artifact compatibility cases. The same source defect was independently reproduced and validated on the other affected non-held Java branches; workflow references, lifecycle ordering and the selected-outcome boundary are unchanged.
This comment has been minimized.
This comment has been minimized.
| DurableExecutionOutput.success(handleLargePayload(executionManager, outputPayload)); | ||
| DurableExecutionOutput output; | ||
| try { | ||
| var outputPayload = config.getSerDes().serialize(result); |
There was a problem hiding this comment.
Codex AI review · Finding arf_v1_53bft2vg4vjxd3j53a3loalx2z
P2: The handler result has already won outcome selection here, but continuation admission remains open throughout serialization and oversized-result checkpointing. If an owned READY continuation fails concurrently, it selects retryable control through executionExceptionFuture, yet this callback continues to emit SUCCEEDED and return success. Call beginInvocationEnd() before output preparation, and add a test that blocks serialization or the large-payload checkpoint while a continuation fails.
Codex AI reviewFound one P2 correctness issue: a concurrent checkpoint-continuation failure can still be lost while preparing a successful response. Reviewed commit |
By submitting this pull request, I confirm that my contribution is made under the terms of the Apache 2.0 license.
Issue Link, if available
Fixes #756, including the first-invocation timing requirement tracked by aws/aws-durable-execution-sdk-js#938. Shared name/attribute match JS #935; this Java implementation is independently based on main.
Description
Fallback Workflow and Invocation spans referenced an SDK-created parent that was never exported. Both views now export
DurableExecutionRootwithdurable.execution.synthetic_root=truebefore invocation return, including the first PENDING or RETRYING response. Complete remote parents stay externally owned.The SDK-owned root now starts before descendants so a configured DurableSampler path uses its actual resolved parent trace state and flags. End captures and clears the root resource before callbacks, ends it once at the checkpointed start timestamp and then flushes. Partial Start, End failures, fatal precedence, reentrant cleanup and reuse are covered. A visible plain replacement retains its documented seeded-ancestor fallback and ordinary per-span sampling; deferred sampling first observes the actual root, without promising arbitrary name/input invariance.
The root retains the existing
generateExecutionRootSpanId(arn)identity and canonical trace resolution. Start and end timestamps both equal the checkpointed execution start, and attributes contain only the execution ARN and synthetic-root marker. It uses the existing provider resource and the same sampling intent as descendants. DurableSampler atomically consumes the per-span context holder and the one-shot sampling bridge before synchronous start processors run, and the remaining scope closes before end processors/exporters run, preserving sampling of unrelated spans those callbacks create. Workflow continues to carry the full execution duration and terminal outcome. Context carriers are used only when the actual public provider sampler is this plugin class copy of DurableSampler. Foreign or opaque samplers use the existing one-shot thread bridge, so provider visibility cannot leave an unconsumable app holder in the parent passed to processors.Every invocation may re-export the same anchor for recovery, including first-invocation redelivery. With stable sampling, re-exports retain identity, timestamps, and span fields. The configured provider resource is preserved; it can vary when recovery runs in another environment, and a deduplicating backend may retain either resource. No persistent exactly-once export mechanism is added. Different executions sharing a root-only trace ID own distinct ARN-derived anchors. No new hash algorithm, registration API, or factory migration is introduced.
The continuation follow-up prevents READY resumption failures from being lost in unobserved futures. Owned SDK continuation failures select retryable manager control before lease release; callers retain the original cause and persisted operation state remains available for retry. Direct JVM fatals also settle observation before escaping the coordinator. Rejected worker admission rolls back its activity registration, assuming rejection before task acceptance. Normal closing preserves unrelated operations and ordinary unowned helper behavior is unchanged. This repair retains this branch's existing invocation-hook threading, header/sampling behavior and workflow references.
Demo/Screenshots
Not applicable (exported trace hierarchy).
Checklist
READY retry progress
Wait-for-condition consumes known READY state and registers a poll before rechecking. Every retry attempt now retires its worker before the next attempt is submitted through the configured executor, including when READY is already available. An independent coordinator retains activity across prior-worker completion and next-worker registration; cancellation of the notification future cannot skip queued work or release a running continuation. This allows a sibling step queued by the root handler to unlock the condition on a bounded executor, and avoids dispatching a direct or submit-and-wait check on the serialized checkpoint batcher.
Normalized state is carried into the next attempt. True PENDING still suspends, and terminal live/replay outcomes preserve stored results and failures without repeating predicates or checkpoints. No sleeps, delay rounding, pool-size mandate or relaxed thresholds were added. Invocation lifecycle/finalization, OTel topology, operation identity, checkpoint formats and workflow refs remain unchanged.
Public old-code controls reproduce both direct/blocking checkpoint deadlocks and fixed-two-thread starvation while a cached executor completes the same root/sibling pattern. The tests retain three controlled worker/executor handoff windows, cancellation/rejection activity ownership, true PENDING and replay counts. Current verification: 318 focused controls on Java 17 and 25, full Java 17 with 2,326 tests (zero failures/errors; 31 existing cloud skips), and all 16 installed-artifact compatibility cases. These establish local defects without claiming a unique cause for the original cloud execution whose full history was unavailable. Original evidence: https://github.com/aws/aws-durable-execution-sdk-java/actions/runs/37566639113. Fresh cloud CI validates the current head.
Continuation registrations now participate in manager shutdown independently of their cancellable observation futures. Admission and closing are atomic; normal close stops only unfinished operations owned by admitted continuations and waits for actual work to release its registration before joining published handlers and shutting down checkpointing. Already-running predicates can finish their checkpoint, and stored outcomes replay without repeating side effects. This does not apply global stopAllOperations during normal close or move the root End/result boundary. Held-before-start, after-old-worker-join, cancellation, admission/rejection, early-parallel and running-checkpoint/replay controls pass.
The clock follow-up retains the actual parent Span while that operation is open, allowing the provider's normal SDK clock inheritance. Previously wrapping only its cached SpanContext created an independent clock anchor; a real initial-PENDING callback report failed strict containment by 10 microseconds despite child-first closure. Four deterministic public Clock controls reproduce that ordering issue before the fix and pass afterward, while a normal-provider ended/replayed-parent control retains the cached fallback, IDs/trace/flags and a single parent export. No timestamp clamping, rounding, shared/system-clock substitution, backend dates or conformance thresholds change. Opaque/cross-process clock behavior remains provider-defined; the new current-head cloud run validates the actual agent path.
Testing
Full Java 17
mvn -o -B clean verify: 2,326 tests, zero failures/errors, 31 existing cloud skips. The scoped coordination tests and all 16 released/current core-plugin artifact combinations pass. Formatting and diff checks pass. The sampler ownership follow-up changes only carrier selection; synthetic-root identity, topology, timestamps and resource semantics remain unchanged. Existing start/end callback, cross-classloader bridge and replay regressions remain in the full suite.Unit Tests
Both views cover absent/malformed context, valid trace without parent, invalid parent, remote-parent ownership, same-trace multiple executions, explicit Sampled=0/1, fallback always-off, first-invocation retry/redelivery, stable recovery exports, and unchanged deterministic identity. Existing raw span-count expectations include the newly exported ancestor. Synchronous onStart and onEnd processor regressions reproduce sampling leakage in both views before their fixes and pass afterward. Cross-classloader consumption, first-flush failure recovery in another environment, and distinct resource ownership are covered. A 32-case guarded onStart matrix additionally checks both views, shared/separate deterministic generators, SAMPLED/RECORD_ONLY root intents, always-on/off callback policies, two invocations, random non-colliding callback IDs, forwarded-parent and setNoParent callbacks, and before/after controls. The 32-case processor matrix and full reactor pass after context-carrier consumption; no ID-generator production change was needed.
A 12-case isolated-loader matrix covers both views, synthetic-root/Invocation parents, direct local/foreign samplers and opaque providers. The four direct-foreign cases fail before the fix; all twelve pass afterward, with unrelated always-off spans staying unrecorded, ambient span restoration and sampling-state cleanup across two invocations.
Integration Tests
Both views exercise an initial PENDING response before Workflow exists, then resume to success or failure. The anchor is already exported on PENDING, final Workflow connects to it, timestamps/resources remain stable, and the completed step side effect executes once. Cloud/LMI tests were not configured locally. Shared conformance coverage for early anchors belongs in the separate conformance repository.
Examples
No new example; README documents anchor ownership, timestamps, sampling, and recovery re-exports.
The close follow-up passed 198 focused controls on Java 17/25, the full Java 17 suite above and all 16 artifact combinations. No workflow or sampler change accompanies this close-only commit.
The clock production change passed full Java 17 verification with 2,330 tests (31 existing skips), then all 88 final focused controls on Java 17/25 including the added fallback case, and all 16 artifact combinations. Actual fresh cloud proof is pending: the failed predecessor run is https://github.com/aws/aws-durable-execution-sdk-java/actions/runs/37734642428.
Current maintenance validation
The latest foreign-sampler fix defers custom results to their owning sampler instead of reducing a foreign result to the property bridge decision enum; explicit upstream/ambient precedence and bridge format are unchanged. All 2,337 Java17 tests (31 existing skips), 140 focused Java17/25 controls and 16 artifact cases pass. Current CI uses validated workflow 1768 with case_count20 and the exact original 02d6 test ref. Runtime head 9760d9b passed Invocation/S3 20/20 including case10 with an actual initial PENDING and later callback completion; another suite in that run hit pre-deployment GetLayer throttling. Newest-head CI remains pending.
Deferred sampler decisions use the execution ARN plus canonical trace ID as a private cache key, retaining full sampling attributes and TraceState when separate executions share a trace. The256-entry LRU limit is unchanged; eviction may evaluate the delegate again. Public two-loader controls cover two executions and two invocations each. Validation:2,343 fullJava17 tests (31 existing skips,zero failures/errors),141 focused Java17/25 tests and16/16 compatibility cases passed.
Current continuation validation
Full Java17 clean verify: 2,393 tests, zero failures/errors, 31 conditional skips. The owned continuation repair passes 63 focused Java25 controls and all 16 artifact compatibility cases. Eight public controls failed on this branch before the repair (six hidden-failure assertions and two bounded rejection timeouts); they now verify caller/cause identity, End status, worker fatal escape, activity cleanup and subsequent successful resumption. Existing lifecycle and telemetry source outside the continuation/admission repair remains unchanged. Fresh CI, including pending core conformance workflows, is tracked separately.
Legacy outcome observer wakeup
A legacy End callback can execute synchronously inside continuation-failure publication. It now wakes waiters for the already-selected control before End can wait for handler cleanup, retaining the existing callback thread. An identity registry tracks active published controls, so a losing CompletableFuture publisher that drains the winner callback still wakes the winner cause. Exact entries are removed in finally; normal outcome diagnostics and threading remain unchanged. No scheduler handoff or timeout fallback is added.
Full Java17 verification passed 2,416 tests (zero failures/errors, 31 conditional skips), 82 focused Java25 controls and all 16 artifact compatibility cases. Fifteen public controls cover asynchronous, inline-root and submit-and-wait root dispatch with asynchronous children; ordinary/fatal/rejected continuations, cooperative End/finally ordering, replay, and unchanged normal End threads. Four of those controls fail on the published baseline. Two deterministic winner/loser cases fail with a thread-keyed candidate and pass with selected-control identity; cleanup/reentry controls also pass. Current-head cloud CI is tracked separately.
Bounded invocation-End admission cut
Before constructing InvocationEndInfo, the SDK now closes new checkpoint-continuation admission and signals unfinished owned continuations. This cancels queued READY resumption after early parallel success; those nonwinning operations may legitimately remain incomplete. Explicit task/handler draining and checkpoint shutdown stay in manager close after End. No scheduler/thread handoff or all-work-drained guarantee is added; owner completion can still run synchronous dependent callbacks. Already accepted/running predicates may finish after End, and their later end hooks need not appear in the earlier snapshot. The API Javadoc and README state these limits.
Full Java17: 2,416 tests, zero failures/errors, 31 conditional skips. All 74 focused Java25 controls and 16 artifact cases pass. The four real executor/OTel controls cover unawaited and early-parallel queued cancellation, the awaited successful path, and an already-running predicate: End can enter while it is held, then normal close waits for it afterward. Current-publisher, winner/loser, reentry and prior close controls remain enabled.