Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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]

Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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]
Expand All @@ -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

Copy link
Copy Markdown
Contributor

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 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.

if execution.is_paused:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Codex AI review · Finding arf_v1_qyavuzqaxgovssjakqmvosmm35

[P2] Cache tokenless responses for idempotent retries

An accepted paused checkpoint advances the token sequence but leaves last_checkpoint unchanged. If its response is lost, retrying the same ClientToken and inbound token is rejected as stale instead of replaying the tokenless response, breaking the existing checkpoint idempotency contract. Persist the record with outbound_checkpoint_token=None and update both checkpoint-path tests to expect identical replay while paused and after resume.

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,
)
Comment thread
hln33 marked this conversation as resolved.

return CheckpointResult(new_token, response_ops, effects)
return CheckpointResult(outbound_token, response_ops, effects)
Comment thread
hln33 marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -62,19 +62,33 @@ class ExecutionStatus(Enum):
TIMED_OUT = "TIMED_OUT"


class PauseState(Enum):
"""Whether a test paused this execution, and whether resume must invoke."""

NOT_PAUSED = "NOT_PAUSED"
# Paused while idle: no handler was running and no invocation was due.
# Resume lifts the pause and the next trigger invokes as usual.
PAUSED = "PAUSED"
# Paused while progress was owed. Either the running handler checkpointed
# and was answered without a token, so it stopped as PENDING before it
# finished, or a scheduled invocation came due and was held back.
# Resume must start one new invocation to make up for it.
PAUSED_INVOCATION_DEFERRED = "PAUSED_INVOCATION_DEFERRED"


@dataclass(frozen=True)
class CheckpointIdempotencyRecord:
"""Single-slot cache of the most recent accepted checkpoint response.
"""Single-slot cache of the most recent accepted checkpoint response
with a checkpoint token.

Single-slot cache of the most recent accepted checkpoint response.
``(client_token, inbound_checkpoint_token)`` pair is entitled to a
byte-identical response; this record is what we compare
against and replay from.
A matching ``(client_token, inbound_checkpoint_token)`` pair replays
this response without applying updates again. Pause-interrupted
checkpoints leave this record unchanged so their retries are rejected.
"""

client_token: str
inbound_checkpoint_token: str
outbound_checkpoint_token: str
outbound_checkpoint_token: str | None
operations: list[Operation]
next_marker: str | None

Expand All @@ -94,7 +108,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"]
],
Expand Down Expand Up @@ -163,6 +177,33 @@ def __init__(
self.result: DurableExecutionInvocationOutput | None = None
self.consecutive_failed_invocation_attempts: int = 0
self.close_status: ExecutionStatus | None = None
self._pause_state: PauseState = PauseState.NOT_PAUSED

@property
def pause_state(self) -> PauseState:
return self._pause_state

@property
def is_paused(self) -> bool:
return self._pause_state is not PauseState.NOT_PAUSED

@property
def has_deferred_invocation(self) -> bool:
return self._pause_state is PauseState.PAUSED_INVOCATION_DEFERRED

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 and return whether the caller must schedule an invocation."""
deferred: bool = self._pause_state is PauseState.PAUSED_INVOCATION_DEFERRED
self._pause_state = PauseState.NOT_PAUSED
return deferred

def touch_operation(self, operation_id: str) -> None:
"""Record a state-affecting event on an operation.
Expand Down Expand Up @@ -248,6 +289,7 @@ 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,
"PauseState": self._pause_state.value,
}

@classmethod
Expand Down Expand Up @@ -316,6 +358,9 @@ def from_json_dict(cls, data: dict[str, Any]) -> Execution:
execution.close_status = (
ExecutionStatus(close_status_str) if close_status_str else None
)
execution._pause_state = PauseState(
data.get("PauseState", PauseState.NOT_PAUSED.value)
)

return execution

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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 = (
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Codex AI review · Finding arf_v1_wbuih3nwxewkn3oiqepguv2fp2

[P3] Validate execution existence before creating a worker

A syntactically valid unknown ARN reaches _registry.submit, which creates a daemon worker before _store.load fails. Memory stores expose a raw KeyError, disk stores return a different exception, and the failed task leaves its worker registered. Call get_execution() before submitting in both pause and resume so unknown ARNs consistently raise ResourceNotFoundException without allocating a lane, and add coverage for each method.

execution_arn,
CallableTask(lambda: self._set_paused(execution_arn)),
).result()
Comment on lines +1815 to +1819

This comment was marked as outdated.

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:
Comment thread
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:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The 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:

  1. 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 keeps pause_execution() waiting for that minute.
  2. 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 raises OrphanedChildException instead 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.

"""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
):
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3222,9 +3222,12 @@ def to_dict(self) -> dict[str, Any]:

@dataclass(frozen=True)
class CheckpointDurableExecutionResponse:
"""Response from checkpointing a durable execution."""
"""Response from checkpointing a durable execution.

checkpoint_token: str
``checkpoint_token`` is None when this invocation may checkpoint no further
"""

checkpoint_token: str | None
new_execution_state: CheckpointUpdatedExecutionState | None = None

@classmethod
Expand All @@ -3234,12 +3237,14 @@ def from_dict(cls, data: dict) -> CheckpointDurableExecutionResponse:
new_execution_state = CheckpointUpdatedExecutionState.from_dict(state_data)

return cls(
checkpoint_token=data["CheckpointToken"],
checkpoint_token=data.get("CheckpointToken"),
new_execution_state=new_execution_state,
)

def to_dict(self) -> dict[str, Any]:
result: dict[str, Any] = {"CheckpointToken": self.checkpoint_token}
result: dict[str, Any] = {}
if self.checkpoint_token is not None:
result["CheckpointToken"] = self.checkpoint_token
if self.new_execution_state is not None:
result["NewExecutionState"] = self.new_execution_state.to_dict()
return result
Expand Down
Loading
Loading