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
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] | ||
|
|
||
|
|
@@ -137,12 +141,25 @@ def apply( | |
| invocation_id=execution.current_invocation_id, | ||
| ).to_str() | ||
|
|
||
| # 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. | ||
| # | ||
| # 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 | ||
| if execution.paused: | ||
| outbound_token = None | ||
| execution.deferred_invocation = True | ||
|
|
||
| execution.last_checkpoint = CheckpointIdempotencyRecord( | ||
| client_token=client_token or "", | ||
| inbound_checkpoint_token=checkpoint_token, | ||
| outbound_checkpoint_token=new_token, | ||
| outbound_checkpoint_token=outbound_token, | ||
|
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. 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.
Wanted: a retry of a checkpoint that was answered without a token, with or without a 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: |
||
| operations=list(response_ops), | ||
| next_marker=None, | ||
| ) | ||
|
|
||
| return CheckpointResult(new_token, response_ops, effects) | ||
| return CheckpointResult(outbound_token, response_ops, effects) | ||
|
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. suggestion: when the token is withheld, the response still carries every operation the handler hasn't seen, and
Wanted: when the token is withheld,
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -70,11 +70,15 @@ class CheckpointIdempotencyRecord: | |
| ``(client_token, inbound_checkpoint_token)`` pair is entitled to a | ||
| byte-identical response; this record is what we compare | ||
| against and replay from. | ||
|
|
||
| ``outbound_checkpoint_token`` is None when the response withheld the | ||
| token because the execution was paused, so the replay stays identical | ||
| after a resume. | ||
| """ | ||
|
|
||
| client_token: str | ||
| inbound_checkpoint_token: str | ||
| outbound_checkpoint_token: str | ||
| outbound_checkpoint_token: str | None | ||
| operations: list[Operation] | ||
| next_marker: str | None | ||
|
|
||
|
|
@@ -94,7 +98,7 @@ def from_json_dict(cls, data: dict[str, Any]) -> CheckpointIdempotencyRecord: | |
| return cls( | ||
| client_token=data["ClientToken"], | ||
| inbound_checkpoint_token=data["InboundCheckpointToken"], | ||
| outbound_checkpoint_token=data["OutboundCheckpointToken"], | ||
| outbound_checkpoint_token=data.get("OutboundCheckpointToken"), | ||
| operations=[ | ||
| Operation.from_json_dict(op_data) for op_data in data["Operations"] | ||
| ], | ||
|
|
@@ -163,6 +167,17 @@ def __init__( | |
| self.result: DurableExecutionInvocationOutput | None = None | ||
| self.consecutive_failed_invocation_attempts: int = 0 | ||
| self.close_status: ExecutionStatus | None = None | ||
| # While True, every checkpoint from pause_execution() until resume_execution() | ||
| # is answered without a checkpoint token. | ||
| # If responses have no token then no new invocation will start. | ||
| 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 | ||
|
Comment on lines
+173
to
+180
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. suggestion: I think these two flags would be safer as one enum, held in a private field and changed only through methods on
Something like this, next to 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 @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 Why |
||
|
|
||
| def touch_operation(self, operation_id: str) -> None: | ||
| """Record a state-affecting event on an operation. | ||
|
|
@@ -248,6 +263,8 @@ def to_json_dict(self) -> dict[str, Any]: | |
| "ConsecutiveFailedInvocationAttempts": self.consecutive_failed_invocation_attempts, | ||
| "CloseStatus": self.close_status.value if self.close_status else None, | ||
| "CurrentInvocationId": self.current_invocation_id, | ||
| "Paused": self.paused, | ||
| "DeferredInvocation": self.deferred_invocation, | ||
| } | ||
|
|
||
| @classmethod | ||
|
|
@@ -316,6 +333,8 @@ def from_json_dict(cls, data: dict[str, Any]) -> Execution: | |
| execution.close_status = ( | ||
| ExecutionStatus(close_status_str) if close_status_str else None | ||
| ) | ||
| execution.paused = data.get("Paused", False) | ||
| execution.deferred_invocation = data.get("DeferredInvocation", False) | ||
|
|
||
| return execution | ||
|
|
||
|
|
||
| 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 | ||
|
|
@@ -1454,19 +1455,29 @@ 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) | ||
| # PENDING is valid while paused because pause can make a | ||
| # checkpoint response omit the next checkpoint token, forcing | ||
| # the current handler invocation to stop as PENDING. | ||
| # | ||
| # Otherwise, PENDING requires either pending durable operations | ||
| # or a change the handler has not seen yet. The unseen-change | ||
| # case can happen when an operation completes after this | ||
| # invocation's input was built but before this response is | ||
| # validated, so a follow-up invocation is needed. | ||
| if ( | ||
| not execution.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. I think
That also makes main's message accurate again. Please could you keep it, |
||
| 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 = ( | ||
| "Cannot return PENDING status with no pending operations." | ||
| "Cannot return PENDING status unless execution is paused, " | ||
| "has pending durable operations, or has unseen changes " | ||
| "after invocation input was built." | ||
| ) | ||
| raise InvalidParameterValueException(msg_pending_ops) | ||
| logger.info("[%s] Execution pending async work", execution_arn) | ||
|
|
@@ -1521,6 +1532,15 @@ def _begin_invocation( | |
| ) | ||
| return None | ||
|
|
||
| if execution.paused: | ||
| execution.deferred_invocation = True | ||
| 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 +1796,68 @@ 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. | ||
| """ | ||
| self._validate_execution_arn(execution_arn) | ||
| self._registry.submit( | ||
| execution_arn, | ||
| CallableTask(lambda: self._set_paused(execution_arn)), | ||
| ).result() | ||
|
Comment on lines
+1814
to
+1818
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 Syntax validation allows an unknown ARN into |
||
| 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. | ||
| """ | ||
| 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: | ||
|
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. suggestion: what Wanted: both methods, with every store, raise That is the message the service returns for a missing execution, for example from Related, not for this PR: the emulator's |
||
| execution = self._store.load(execution_arn) | ||
| if execution.is_complete or execution.paused: | ||
| return | ||
| execution.paused = True | ||
| self._store.save(execution) | ||
|
|
||
| def _resume_execution(self, execution_arn: str) -> None: | ||
| execution = self._store.load(execution_arn) | ||
| if execution.is_complete or not execution.paused: | ||
| return | ||
| execution.paused = False | ||
| deferred = execution.deferred_invocation | ||
| execution.deferred_invocation = False | ||
| 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 | ||
| ): | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -714,6 +714,10 @@ def __init__( | |
| self._checkpoint_processor.set_chained_invoke_preflight( | ||
| self._executor.preflight_chained_invoke | ||
| ) | ||
| # 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 | ||
|
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. suggestion: please could you make
The default also adds state the test can't see:
JS's Concretely: def pause_execution(self, execution_arn: str) -> None:
def resume_execution(self, execution_arn: str) -> None:The rest is deletions: |
||
|
|
||
| def register_durable_function( | ||
| self, | ||
|
|
@@ -861,8 +865,59 @@ def run_async( | |
| if output.execution_arn is None: | ||
| msg_arn: str = "Execution ARN must exist to run test." | ||
| raise DurableFunctionsTestError(msg_arn) | ||
| self._default_execution_arn = output.execution_arn | ||
| return output.execution_arn | ||
|
|
||
| def pause_execution(self, execution_arn: str | None = None) -> 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. suggestion: the cloud runner has neither method, so a test that pauses and later runs against Please could you add both methods to
|
||
| """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. | ||
|
|
||
| Defaults to the execution run_async() last started. The | ||
| invocation running now, if any, is answered without a token on | ||
| its next checkpoint: that checkpoint is accepted, but the | ||
| invocation reports PENDING, and no new invocation starts until | ||
| resume_execution(). | ||
|
|
||
| Idempotent; a no-op once the execution has finished. | ||
| Resolves once no invocation of this execution is running. | ||
|
|
||
| Raises: | ||
| DurableFunctionsTestError: If no execution is in progress | ||
| (run_async() has not been called). | ||
| """ | ||
| self._executor.pause_execution(self._require_execution_arn(execution_arn)) | ||
|
|
||
| def resume_execution(self, execution_arn: str | None = None) -> 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. | ||
|
|
||
| Defaults to the execution run_async() last started. Starts the | ||
| invocation pause_execution() held back, if any. | ||
|
|
||
| Idempotent; a no-op once the execution has finished or if it was not paused. | ||
|
|
||
| Raises: | ||
| DurableFunctionsTestError: If no execution is in progress | ||
| (run_async() has not been called). | ||
| """ | ||
| self._executor.resume_execution(self._require_execution_arn(execution_arn)) | ||
|
|
||
| def _require_execution_arn(self, execution_arn: str | None) -> str: | ||
| arn = ( | ||
| execution_arn if execution_arn is not None else self._default_execution_arn | ||
| ) | ||
| if arn is None: | ||
| msg = ( | ||
| "No execution in progress. Call run_async() or run() first, " | ||
| "or pass execution_arn explicitly." | ||
| ) | ||
| raise DurableFunctionsTestError(msg) | ||
| return arn | ||
|
|
||
| def wait_for_result( | ||
| self, execution_arn: str, timeout: int = 60 | ||
| ) -> DurableFunctionTestResult: | ||
|
|
||
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.