Repository navigation
722 feature exit gracefully with pending when a checkpoint response has no checkpointtoken #757
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
07cd5a5
6689d27
0ab0d8c
c234aaf
aa2d825
88dffdc
c5af08b
387f58c
2a34ea5
4eeabba
50be8d4
f2f73f6
136d3cd
90607bc
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -38,9 +38,13 @@ | |
|
|
||
| class CheckpointResult(NamedTuple): | ||
| """Outcome of applying a checkpoint: the new token, the operations to | ||
| return to the handler this round, and the lifecycle effects raised.""" | ||
| return to the handler this round, and the lifecycle effects raised. | ||
|
|
||
| checkpoint_token: str | ||
| ``checkpoint_token`` is None when the execution is paused and the token | ||
| is withheld, telling the SDK this invocation may checkpoint no further. | ||
| """ | ||
|
|
||
| checkpoint_token: str | None | ||
| operations: list[Operation] | ||
| effects: list[CheckpointEffect] | ||
|
|
||
|
|
@@ -86,9 +90,11 @@ def apply( | |
| """Apply ``updates`` to ``execution`` and compute the response delta. | ||
|
|
||
| Advances ``token_sequence`` exactly once, returns the full set of | ||
| operations the handler has not yet seen, advances | ||
| ``handler_seen_seq`` to cover them, and records the | ||
| idempotency entry for a byte-identical replay of a retried call. | ||
| operations the handler has not yet seen unless paused, advances | ||
| ``handler_seen_seq`` only for returned operations, and records an | ||
| idempotency entry for a byte-identical replay of a retried call | ||
| only when returning a checkpoint token. | ||
|
|
||
| The caller is responsible for the invocation gate, locking, | ||
| persistence, and applying the returned effects. | ||
|
|
||
|
|
@@ -123,7 +129,10 @@ def apply( | |
| # The checkpoint response returns the full unseen delta in a single | ||
| # response. Advance handler_seen_seq to cover every returned op so | ||
| # the next delta carries only operations touched after this response. | ||
| response_ops: list[Operation] = paginator.unseen_operations() | ||
| # Paused checkpoints persist updates without delivering any state. | ||
| response_ops: list[Operation] = ( | ||
| [] if execution.is_paused else paginator.unseen_operations() | ||
| ) | ||
| if response_ops: | ||
| highest_delivered_seq: int = max( | ||
| execution.operation_last_touched_seq[op.operation_id] | ||
|
|
@@ -137,12 +146,26 @@ def apply( | |
| invocation_id=execution.current_invocation_id, | ||
| ).to_str() | ||
|
|
||
| execution.last_checkpoint = CheckpointIdempotencyRecord( | ||
| client_token=client_token or "", | ||
| inbound_checkpoint_token=checkpoint_token, | ||
| outbound_checkpoint_token=new_token, | ||
| operations=list(response_ops), | ||
| next_marker=None, | ||
| ) | ||
| # A paused execution registers this checkpoint's updates but withholds | ||
| # the token: a response without one tells the SDK this invocation may | ||
| # checkpoint no further, so it reports PENDING at its next checkpoint | ||
| # rather than continuing, and owes a re-invoke once resumed. | ||
| # | ||
| # Checkpoints that were interrupted by a pause are not added to the | ||
| # idempotency record. A retry of this checkpoint will fail rather than be | ||
| # answered again. | ||
|
|
||
| outbound_token: str | None = new_token | ||
| if execution.is_paused: | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Codex AI review · Finding [P2] Cache tokenless responses for idempotent retries An accepted paused checkpoint advances the token sequence but leaves |
||
| outbound_token = None | ||
| execution.defer_invocation() | ||
| else: | ||
| execution.last_checkpoint = CheckpointIdempotencyRecord( | ||
| client_token=client_token or "", | ||
| inbound_checkpoint_token=checkpoint_token, | ||
| outbound_checkpoint_token=outbound_token, | ||
| operations=list(response_ops), | ||
| next_marker=None, | ||
| ) | ||
|
hln33 marked this conversation as resolved.
|
||
|
|
||
| return CheckpointResult(new_token, response_ops, effects) | ||
| return CheckpointResult(outbound_token, response_ops, effects) | ||
|
hln33 marked this conversation as resolved.
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -5,6 +5,7 @@ | |
| import asyncio | ||
| import logging | ||
| import threading | ||
| import time | ||
| import uuid | ||
| from datetime import datetime | ||
| from typing import TYPE_CHECKING, assert_never | ||
|
|
@@ -351,7 +352,7 @@ def get_execution(self, execution_arn: str) -> Execution: | |
| try: | ||
| return self._store.load(execution_arn) | ||
| except KeyError as e: | ||
| msg: str = f"Execution {execution_arn} not found" | ||
| msg: str = "Durable Execution does not exist" | ||
| raise ResourceNotFoundException(msg) from e | ||
|
|
||
| def get_execution_details(self, execution_arn: str) -> GetDurableExecutionResponse: | ||
|
|
@@ -1454,15 +1455,22 @@ def _validate_invocation_response_and_store( | |
| ) | ||
|
|
||
| case InvocationStatus.PENDING: | ||
| # An operation the handler waited on may complete between | ||
| # the handler's return and this check. A change the | ||
| # handler has not seen, after the invocation's input was | ||
| # built, earns a re-invoke, so PENDING is valid; only a | ||
| # handler that waited on nothing is in error. | ||
| if not execution.has_pending_operations(execution) and not ( | ||
| invocation_seq is not None | ||
| and execution.has_changes_after( | ||
| max(invocation_seq, execution.handler_seen_seq) | ||
| # A paused execution answers the running invocation's | ||
| # checkpoint without a token and defers its next invocation; | ||
| # that invocation must stop as PENDING even with nothing | ||
| # pending, and resume re-invokes it. Any other PENDING needs | ||
| # pending operations or a change the handler has not seen | ||
| # since its input was built. The unseen-change case happens | ||
| # when an operation completes after this invocation's input | ||
| # was built but before this response is validated. | ||
| if ( | ||
| not execution.has_deferred_invocation | ||
| and not execution.has_pending_operations(execution) | ||
| and not ( | ||
| invocation_seq is not None | ||
| and execution.has_changes_after( | ||
| max(invocation_seq, execution.handler_seen_seq) | ||
| ) | ||
| ) | ||
| ): | ||
| msg_pending_ops: str = ( | ||
|
|
@@ -1521,6 +1529,15 @@ def _begin_invocation( | |
| ) | ||
| return None | ||
|
|
||
| if execution.is_paused: | ||
| execution.defer_invocation() | ||
| self._store.save(execution) | ||
| logger.debug( | ||
| "[%s] Holding back scheduled invocation while paused", | ||
| execution_arn, | ||
| ) | ||
| return None | ||
|
|
||
| # Claim the gate: at most one handler invocation per | ||
| # execution in flight. | ||
| self._set_invocation_gate(execution_arn, InvocationState.INVOKING) | ||
|
|
@@ -1776,6 +1793,76 @@ def _invoke_execution(self, execution_arn: str, delay: float = 0) -> None: | |
| completion_event=completion_event, | ||
| ) | ||
|
|
||
| def pause_execution(self, execution_arn: str) -> None: | ||
| """Make the local checkpoint server answer this execution's | ||
| checkpoints without a token, starting now. | ||
|
|
||
| Experimental; may change or be removed in a future release. | ||
|
|
||
| The invocation running now, if any, is answered without a token | ||
| on its next checkpoint. That checkpoint is accepted - its updates | ||
| stay durable - but the invocation reports PENDING, and no further | ||
| checkpoint of its is accepted. No new invocation starts until | ||
| resume_execution() is called. | ||
|
|
||
| Idempotent; a no-op once the execution has finished. | ||
| Resolves once no invocation of this execution is running. | ||
|
|
||
| Raises: | ||
| InvalidParameterValueException: If the ARN is blank. | ||
| ResourceNotFoundException: If the execution does not exist. | ||
| """ | ||
| self._validate_execution_arn(execution_arn) | ||
| self._registry.submit( | ||
This comment was marked as outdated.
Sorry, something went wrong.
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Codex AI review · Finding [P3] Validate execution existence before creating a worker A syntactically valid unknown ARN reaches |
||
| execution_arn, | ||
| CallableTask(lambda: self._set_paused(execution_arn)), | ||
| ).result() | ||
|
Comment on lines
+1815
to
+1819
This comment was marked as outdated.
Sorry, something went wrong. |
||
| self._wait_until_idle(execution_arn) | ||
|
|
||
| def resume_execution(self, execution_arn: str) -> None: | ||
| """Make the local checkpoint server answer this execution's | ||
| checkpoints with a token again. | ||
|
|
||
| Experimental; may change or be removed in a future release. | ||
|
|
||
| Starts the invocation that pause_execution() held back, if any. | ||
|
|
||
| Idempotent; a no-op once the execution has finished or if it was | ||
| not paused. | ||
|
|
||
| Raises: | ||
| InvalidParameterValueException: If the ARN is blank. | ||
| ResourceNotFoundException: If the execution does not exist. | ||
| """ | ||
| self._validate_execution_arn(execution_arn) | ||
| self._registry.submit( | ||
| execution_arn, | ||
| CallableTask(lambda: self._resume_execution(execution_arn)), | ||
| ).result() | ||
|
|
||
| def _set_paused(self, execution_arn: str) -> None: | ||
|
hln33 marked this conversation as resolved.
|
||
| execution = self.get_execution(execution_arn) | ||
| if execution.is_complete or execution.is_paused: | ||
| return | ||
| execution.pause() | ||
| self._store.save(execution) | ||
|
|
||
| def _resume_execution(self, execution_arn: str) -> None: | ||
| execution = self.get_execution(execution_arn) | ||
| if execution.is_complete or not execution.is_paused: | ||
| return | ||
|
|
||
| deferred = execution.resume() | ||
|
|
||
| self._store.save(execution) | ||
| if deferred: | ||
| self._invoke_execution(execution_arn) | ||
|
|
||
| def _wait_until_idle(self, execution_arn: str) -> None: | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This loop has no timeout, and two cases make it wait a long time:
Maybe take a |
||
| """Block until no invocation of ``execution_arn`` is running.""" | ||
| while self._invocation_gate(execution_arn) is InvocationState.INVOKING: | ||
| time.sleep(0.005) | ||
|
|
||
| def _complete_workflow( | ||
| self, execution_arn: str, result: str | None, error: ErrorObject | None | ||
| ): | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
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
Operationslist, and the SDK treats such a response as a finished execution (state.pylines 1064-1068). So the emulator reaches that branch only while paused. Maybe worth an issue, so an ordinary completion exercises it too.