Repository navigation
722 feature exit gracefully with pending when a checkpoint response has no checkpointtoken - #757
Conversation
A checkpoint response without a token means the current invocation must stop checkpointing. The background thread previously treated every such response as execution completion, which could report SUCCEEDED or FAILED while operation work was still in flight. It now reports PENDING. - Add ExecutionSuspendedByService and a token-revoked latch distinct from execution completion. - Exempt a batch carrying the execution's terminal update because that execution is already complete. - Reject newly queued operations and refreshes after token withdrawal, and wake accepted and queued waiters through the suspension path. - Return a bare PENDING response from every otherwise-terminal exit, including oversized-result paths and the CheckpointError exit, which now consults the latch for symmetry with the other exits. - Treat a refresh-only batch exactly like every other non-terminal batch when its response omits the token. - Document ExecutionSuspendedByService in the Raises section of create_checkpoint and schedule_refresh. - Drop ExecutionSuspendedByService's redundant constructor override, identical to SuspendExecution's. - Describe only the SDK-observed condition and resulting behavior in comments and docstrings. - Cover queue guards, deferred refreshes, completion versus withdrawal, oversized results, the CheckpointError exit, and every PENDING exit. Fixes #722
- Executor: skip PENDING-validator rejection when execution.paused - Add 9 new tests (6 unit + 3 e2e) covering pause/resume semantics - Fix 61 pre-existing executor_test.py mocks missing .paused=False (unconfigured Mock attrs are truthy, broke new paused-aware paths) - Update 3 mock-sequencing tests for new _invoke_execution store.load call
Use the existing SuspendExecution control flow for revoked checkpoint tokens and simplify the surrounding comments and tests. Clean up the local test runner pause and resume path, including deferred invocation naming and focused pause/resume coverage.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
8859615 to
aa2d825
Compare
…n-a-checkpoint-response-has-no-checkpointtoken
yaythomas
left a comment
There was a problem hiding this comment.
Thank you @hln33, this is a careful port of the JS fix, and the state_test.py coverage of the real checkpoint loop is excellent!
Thanks for the follow-up commit for the Codex findings. _invoke_execution no longer reads and saves outside the worker lane, a retried checkpoint no longer receives the withheld token, and the mid-step pause test now fails without the fix (I checked against main's core SDK).
Please could you rename the PR to the conventional commit format, e.g. fix(sdk): answer PENDING on revoked token (your first commit's subject)? The PR title becomes the squash commit message: https://github.com/aws/aws-durable-execution-sdk-python/blob/main/CONTRIBUTING.md#pull-request-title-and-commit-message-format
Maybe this behaviour also warrants a conformance requirement? If so, please could you note it in the description or open an issue in aws-durable-execution-conformance-tests?
| self.paused: bool = False | ||
| # Set while paused when progress was stopped and must be resumed later. | ||
| # This means either a new handler invocation was not started because the | ||
| # execution is paused, or the current handler was given no next checkpoint | ||
| # token and therefore must stop as PENDING. | ||
| # | ||
| # resume_execution() clears this flag and starts one new invocation. | ||
| self.deferred_invocation: bool = False |
There was a problem hiding this comment.
suggestion: I think these two flags would be safer as one enum, held in a private field and changed only through methods on Execution.
- Two bools allow four combinations, and only three are valid: not paused, paused with nothing owed, and paused with an invocation owed.
- The fourth combination, "not paused, invocation owed", does happen. The
_invoke_executionrace you just fixed produced it. The newprocessor_test.pytest also builds it: it setspaused = Falseafter a paused checkpoint has setdeferred_invocation = True. - One enum field cannot hold that combination.
Something like this, next to ExecutionStatus:
class PauseState(Enum):
"""Whether a test paused this execution, and whether resume must invoke."""
NOT_PAUSED = "NOT_PAUSED"
PAUSED = "PAUSED"
# While paused, an invocation was cut off or a wake-up was held back.
PAUSED_INVOCATION_DEFERRED = "PAUSED_INVOCATION_DEFERRED"with these on Execution:
@property
def is_paused(self) -> bool:
return self._pause_state is not PauseState.NOT_PAUSED
def pause(self) -> None:
if self._pause_state is PauseState.NOT_PAUSED:
self._pause_state = PauseState.PAUSED
def defer_invocation(self) -> None:
if self._pause_state is PauseState.PAUSED:
self._pause_state = PauseState.PAUSED_INVOCATION_DEFERRED
def resume(self) -> bool:
"""Clear the pause. Return whether resume must start an invocation."""
deferred: bool = self._pause_state is PauseState.PAUSED_INVOCATION_DEFERRED
self._pause_state = PauseState.NOT_PAUSED
return deferredWhy methods: code outside Execution assigns these flags on five lines, in CheckpointCore.apply, _begin_invocation, _set_paused and _resume_execution. With methods, the rules live in one class, the way start(), begin_new_invocation() and the complete_* methods already work. is_paused would also match is_complete, and to_json_dict would store one "PauseState" key instead of two.
Why PauseState: it matches InvocationState, the package's other internal state enum. In this package, ...Status names values that go out in API responses, like ExecutionStatus.
| # invocation's input was built but before this response is | ||
| # validated, so a follow-up invocation is needed. | ||
| if ( | ||
| not execution.paused |
There was a problem hiding this comment.
I think not execution.paused turns this check off for more invocations than the pause needs.
- While the execution is paused, the emulator answers the running invocation's next checkpoint without a token. The SDK then stops and answers PENDING, often with nothing pending. That PENDING is correct, so it needs an exception.
pausedalso exempts an invocation that answers PENDING with nothing pending although none of its checkpoints was answered without a token. That is the SDK bug this check catches.- I tried that second case on this head. The emulator accepts the answer and counts no failed attempt, and
resume_execution()starts no invocation, becausedeferred_invocationwas never set. So the execution makes no progress until it times out, and nothing names the cause.
CheckpointCore.apply sets deferred_invocation when it withholds a token, and it's now the only code that sets it while an invocation runs. With not execution.deferred_invocation in place of not execution.paused, the second case is rejected and retried as on main, and the pause/resume tests still pass. With the enum suggestion, the condition would check for PauseState.PAUSED_INVOCATION_DEFERRED.
That also makes main's message accurate again. Please could you keep it, Cannot return PENDING status with no pending operations.? Two tests in executor_test.py matched on it before this PR, so users' tests may too.
| if deferred: | ||
| self._invoke_execution(execution_arn) | ||
|
|
||
| def _wait_until_idle(self, execution_arn: str) -> None: |
There was a problem hiding this comment.
This loop has no timeout, and two cases make it wait a long time:
pause_execution()returns only after the running invocation ends, and the SDK ends it at its next checkpoint. So a step body that runs for a minute keepspause_execution()waiting for that minute.- The testing package accepts any core SDK from 1.0.0 (
aws-durable-execution-sdk-python>=1.0.0). A core SDK without this fix raisesOrphanedChildExceptioninstead of answering PENDING. I ran a two-step pause test against main's core SDK:pause_execution()returned only when the execution's 15 s timeout fired, and the execution ended timed out.
Maybe take a timeout, as wait_for_result(timeout=...) does, and raise DurableFunctionsTestError when it expires? It would also help to say in both docstrings that the pause takes effect at the running invocation's next checkpoint. When this releases, the testing package's core floor could move to the first core version with this fix.
| # ARN of the most recent execution started by run_async(), including calls | ||
| # through run(). Used as the default target for pause_execution() and | ||
| # resume_execution() when no ARN is provided. | ||
| self._default_execution_arn: str | None = None |
There was a problem hiding this comment.
suggestion: please could you make execution_arn a required parameter of pause_execution() and resume_execution(), and drop this field?
- The runner's other per-execution methods take the ARN:
wait_for_result,wait_for_callbackandget_execution_history. - A test can pause a running execution only after
run_async()returns, andrun_async()returns the ARN.run()returns only after the execution has finished. - So any test that can pause already holds the ARN. The default saves one argument.
The default also adds state the test can't see:
- It targets whichever execution
run_async()started last. A test that starts two executions pauses the second one. - It keeps that ARN after the execution finishes. A later
pause_execution()with no ARN then does nothing and raises nothing, because_set_pausedreturns early for a completed execution. The "No execution in progress" error fires only before the firstrun_async().
JS's pauseExecution() takes no argument, and that fits JS: its run() returns only a promise of the result, and pauseExecution() acts on the execution that run() is driving. Here, run_async() already hands the test the ARN.
Concretely:
def pause_execution(self, execution_arn: str) -> None:
def resume_execution(self, execution_arn: str) -> None:The rest is deletions: _default_execution_arn, its assignment in run_async(), _require_execution_arn, and in both docstrings the "Defaults to..." sentence and the Raises: entry. Worth doing before this ships, because removing a default later breaks every caller that relies on it, experimental or not.
| self._default_execution_arn = output.execution_arn | ||
| return output.execution_arn | ||
|
|
||
| def pause_execution(self, execution_arn: str | None = None) -> None: |
There was a problem hiding this comment.
suggestion: the cloud runner has neither method, so a test that pauses and later runs against DurableFunctionCloudTestRunner fails with a generic AttributeError: 'DurableFunctionCloudTestRunner' object has no attribute 'pause_execution'. JS gives both runners the same interface, and its cloud runner raises pauseExecution() is not implemented for CloudDurableTestRunner (aws-durable-execution-sdk-js#931).
Please could you add both methods to DurableFunctionCloudTestRunner, raising NotImplementedError with a message along the lines of:
pause_execution() is not implemented for DurableFunctionCloudTestRunner
resume_execution() is not implemented for DurableFunctionCloudTestRunner
NotImplementedError because the cloud runner can gain real pause and resume later, through the API.
| client_token=client_token or "", | ||
| inbound_checkpoint_token=checkpoint_token, | ||
| outbound_checkpoint_token=new_token, | ||
| outbound_checkpoint_token=outbound_token, |
There was a problem hiding this comment.
suggestion: I think a retry of a checkpoint that was answered without a token should fail like any other stale token, rather than be answered again.
- A response without a token means the invocation may not checkpoint again.
applyhas also advanced the token sequence, so the call's token is now stale. - On this head, a retry without a client token already fails that check with
InvalidParameterValueException("Invalid checkpoint token"). - A retry with the same client token matches this record and is answered again instead. A botocore retry over HTTP takes that path: the core SDK leaves
ClientTokenout, the model marks it as an idempotency token, so botocore fills one in and resends the same one on retry. - So whether the retry is rejected depends on whether it carries a client token.
Wanted: a retry of a checkpoint that was answered without a token, with or without a ClientToken, raises InvalidParameterValueException with the message exactly:
Invalid checkpoint token
That is the emulator's existing stale-token message, so nothing new is needed on the error side. Skipping this record when the token is withheld does it: last_checkpoint stays as it was, so both retries fall through to the existing token check. The withheld token still never reaches the handler. I tried it on this head: the e2e pause tests pass, and only the two test_paused_checkpoint_retries_without_a_token_even_after_resume tests, in executor_pause_resume_test.py and checkpoint/processor_test.py, need to expect the rejection instead of a second token-less answer.
| CallableTask(lambda: self._resume_execution(execution_arn)), | ||
| ).result() | ||
|
|
||
| def _set_paused(self, execution_arn: str) -> None: |
There was a problem hiding this comment.
suggestion: what pause_execution() and resume_execution() raise for an unknown ARN depends on the store. With the in-memory store, this load() raises a bare KeyError: 'arn:unknown' (I ran both). The filesystem and SQLite stores' load() raise ResourceNotFoundException with Execution arn:unknown not found instead.
Wanted: both methods, with every store, raise ResourceNotFoundException with the message exactly:
Durable Execution does not exist
That is the message the service returns for a missing execution, for example from StopDurableExecution on an unknown ARN. Catching KeyError and ResourceNotFoundException around the load() here and in _resume_execution does it. I tried it: both calls raise that, and no existing test depends on the KeyError.
Related, not for this PR: the emulator's get_execution and the stores say Execution <arn> not found for the same case. Aligning those with the service's text could be a follow-up.
| assert retry_after_resume.checkpoint_token is None | ||
|
|
||
|
|
||
| def test_invoke_execution_while_paused_still_schedules_with_its_delay(): |
There was a problem hiding this comment.
suggestion: a few tests would pin down the rules around a withheld token, including the exact error messages. The first two pass on this head. The last two fail on this head and pass with the suggestions on executor.py lines 1468 and 1838. They pass hatch fmt --check.
def test_checkpoint_with_the_withheld_invocations_token_is_rejected():
executor, store, execution, token_0 = _make_executor_with_started_execution()
execution.paused = True
store.save(execution)
executor.checkpoint_execution(
execution_arn=execution.durable_execution_arn,
checkpoint_token=token_0,
updates=[_step_start_update("step-A")],
)
with pytest.raises(InvalidParameterValueException) as exc_info:
executor.checkpoint_execution(
execution_arn=execution.durable_execution_arn,
checkpoint_token=token_0,
updates=[_step_start_update("step-B")],
)
assert str(exc_info.value) == "Invalid checkpoint token"
def test_resume_with_nothing_held_back_starts_no_invocation():
executor, _, execution, _ = _make_executor_with_started_execution()
arn = execution.durable_execution_arn
executor._set_invocation_gate(arn, InvocationState.PRE_INVOKE) # noqa: SLF001
executor.pause_execution(arn)
executor.resume_execution(arn)
executor._scheduler.call_later.assert_not_called() # noqa: SLF001
def test_paused_pending_is_accepted_only_when_a_token_was_withheld():
executor, store, execution, _ = _make_executor_with_started_execution()
execution.paused = True
store.save(execution)
pending = DurableExecutionInvocationOutput(status=InvocationStatus.PENDING)
with pytest.raises(InvalidParameterValueException) as exc_info:
executor._validate_invocation_response_and_store( # noqa: SLF001
execution.durable_execution_arn, pending, execution, execution.seq_counter
)
assert str(exc_info.value) == (
"Cannot return PENDING status with no pending operations."
)
execution.deferred_invocation = True
executor._validate_invocation_response_and_store( # noqa: SLF001
execution.durable_execution_arn, pending, execution, execution.seq_counter
)
def test_pause_and_resume_raise_for_an_unknown_execution():
executor, _, _, _ = _make_executor_with_started_execution()
for call in (executor.pause_execution, executor.resume_execution):
with pytest.raises(ResourceNotFoundException) as exc_info:
call("arn:unknown")
assert str(exc_info.value) == "Durable Execution does not exist"They need pytest, DurableExecutionInvocationOutput, InvocationStatus, InvalidParameterValueException and ResourceNotFoundException imported.
With all the suggestions applied, I ran the whole testing suite. These are the only existing tests that need a change:
test_paused_checkpoint_retries_without_a_token_even_after_resume, here and incheckpoint/processor_test.py: the retry now raisesInvalidParameterValueException("Invalid checkpoint token").test_checkpoint_while_paused_omits_token_but_registers_update: the response's operations list is[]; the check thatstep-Awas registered stays.executor_test.py, the two tests whosematch=this PR changed: back to main'smatch="no pending operations".executor_test.py,test_should_handle_pending_status_when_operations_existandtest_should_retry_when_pending_response_has_no_operations: addmock_execution.deferred_invocation = False, for the same reason this PR addedmock_execution.paused = False(an unsetMock()attribute is truthy).
| ) | ||
|
|
||
| return CheckpointResult(new_token, response_ops, effects) | ||
| return CheckpointResult(outbound_token, response_ops, effects) |
There was a problem hiding this comment.
suggestion: when the token is withheld, the response still carries every operation the handler hasn't seen, and apply advances handler_seen_seq over them.
- A token-less response from the service carries an empty
Operationslist. That is what the SDK gets today when a batch completes the execution. - For a batch without an
EXECUTIONupdate, the SDK reaches the revoked-token branch and stops before it readsnew_execution_state(state.pylines 1069-1071). So the handler never applies these operations. - The emulator still records them as delivered, because
handler_seen_seqmoved over them.
Wanted: when the token is withheld, new_execution_state.operations is [] and handler_seen_seq does not move.
[] if execution.paused else paginator.unseen_operations() does both, because the advance only runs over the returned list. I tried it: the e2e pause tests pass, and test_checkpoint_while_paused_omits_token_but_registers_update would then assert an empty list and keep its check that step-A was registered.
| # The withheld token is recorded as the idempotency record's outbound | ||
| # token so a retry of this call replays the same tokenless response, | ||
| # even after a resume has moved the execution on. | ||
| outbound_token: str | None = new_token |
There was a problem hiding this comment.
no-action note: when not paused, a checkpoint whose batch completes the execution still gets a token here. The service answers that batch without a token and with an empty Operations list, and the SDK treats such a response as a finished execution (state.py lines 1064-1068). So the emulator reaches that branch only while paused. Maybe worth an issue, so an ordinary completion exercises it too.
| else: | ||
| logger.exception("Invocation error. Must terminate.") | ||
| # Throw the error to trigger Lambda retry | ||
| raise |
There was a problem hiding this comment.
Codex AI review · Finding arf_v1_guweetjrqz6qymb4qyevo4qc4h
[P1] Defer retryable errors until cleanup checks revocation
This re-raises inside the with block. If an asynchronous checkpoint returns without a token while the context managers unwind, revocation is recorded after this check; the post-cleanup override is skipped and Lambda retries with the revoked invocation token. Capture the exception, finish checkpoint-thread cleanup, then return bare PENDING if revocation occurred; otherwise re-raise. Add a test racing a retryable InvocationError with a tokenless async checkpoint.
| "Checkpoint response omitted the token for an empty " | ||
| "checkpoint." | ||
| ) | ||
| self._handle_revoked_checkpoint_token(batch) |
There was a problem hiding this comment.
Codex AI review · Finding arf_v1_w75l3lpqi2d4aselqozh4pokc4
[P2] Apply accepted response state before suspending
The tokenless response accepted this batch and may contain its terminal operation state, but this exits before merging NewExecutionState or emitting plugin updates. After resume, instrumentation never receives on_operation_end, so OTel can omit the logical operation span and invocation-end state is stale. Merge the inline operations without pagination and emit updates before waking waiters with SuspendExecution; add plugin coverage for this path.
| self._validate_execution_arn(execution_arn) | ||
| self._registry.submit( | ||
| execution_arn, | ||
| CallableTask(lambda: self._set_paused(execution_arn)), | ||
| ).result() |
There was a problem hiding this comment.
Codex AI review · Finding arf_v1_wbuih3nwxewkn3oiqepguv2fp2
[P3] Validate execution existence before creating a worker
Syntax validation allows an unknown ARN into _registry.submit. The default memory store then exposes a raw KeyError while disk stores return a different exception, and the failed task leaves a worker lane registered for the nonexistent execution. Call get_execution() before submitting in both pause and resume, returning a consistent ResourceNotFoundException, and add unknown-ARN tests.
Codex AI reviewFound three issues: a cleanup-time revocation race, lost instrumentation state on tokenless checkpoints, and inconsistent unknown-ARN handling in pause/resume. Reviewed commit |
Issue #, if available:
Fixes #722
Description of changes:
This PR updates the SDK to return
PENDINGwhen a checkpoint response omits the next checkpoint token for non-terminal work, instead of treating that response as execution completion.Changes include:
SuspendExecutionpath.PENDINGwhen the token has been revoked.By submitting this pull request, I confirm that you can use, modify, copy, and redistribute this contribution, under the terms of your choice.