Skip to content
Open
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 @@ -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

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

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.

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.

  1. A response without a token means the invocation may not checkpoint again. apply has also advanced the token sequence, so the call's token is now stale.
  2. On this head, a retry without a client token already fails that check with InvalidParameterValueException("Invalid checkpoint token").
  3. 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 ClientToken out, the model marks it as an idempotency token, so botocore fills one in and resends the same one on retry.
  4. 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.

operations=list(response_ops),
next_marker=None,
)

return CheckpointResult(new_token, response_ops, effects)
return CheckpointResult(outbound_token, response_ops, effects)

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.

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.

  1. A token-less response from the service carries an empty Operations list. That is what the SDK gets today when a batch completes the execution.
  2. For a batch without an EXECUTION update, the SDK reaches the revoked-token branch and stops before it reads new_execution_state (state.py lines 1069-1071). So the handler never applies these operations.
  3. The emulator still records them as delivered, because handler_seen_seq moved 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.

Original file line number Diff line number Diff line change
Expand Up @@ -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

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

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.

suggestion: I think these two flags would be safer as one enum, held in a private field and changed only through methods on Execution.

  1. Two bools allow four combinations, and only three are valid: not paused, paused with nothing owed, and paused with an invocation owed.
  2. The fourth combination, "not paused, invocation owed", does happen. The _invoke_execution race you just fixed produced it. The new processor_test.py test also builds it: it sets paused = False after a paused checkpoint has set deferred_invocation = True.
  3. 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 deferred

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


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

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

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.

I think not execution.paused turns this check off for more invocations than the pause needs.

  1. 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.
  2. paused also 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.
  3. I tried that second case on this head. The emulator accepts the answer and counts no failed attempt, and resume_execution() starts no invocation, because deferred_invocation was 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.

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

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

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.

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:

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.

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.

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:

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
Original file line number Diff line number Diff line change
Expand Up @@ -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

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.

suggestion: please could you make execution_arn a required parameter of pause_execution() and resume_execution(), and drop this field?

  1. The runner's other per-execution methods take the ARN: wait_for_result, wait_for_callback and get_execution_history.
  2. A test can pause a running execution only after run_async() returns, and run_async() returns the ARN. run() returns only after the execution has finished.
  3. 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_paused returns early for a completed execution. The "No execution in progress" error fires only before the first run_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.


def register_durable_function(
self,
Expand Down Expand Up @@ -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:

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.

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.

"""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:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -408,3 +408,52 @@ def test_process_checkpoint_delivers_due_wait_completion() -> None:
op for op in persisted.operations if op.operation_id == "wait-1"
)
assert persisted_wait.status is OperationStatus.SUCCEEDED


def test_paused_checkpoint_retries_without_a_token_even_after_resume():
"""A retry of a checkpoint answered while paused replays the same
tokenless response, so the invocation cannot keep checkpointing. The
withholding is recorded on the idempotency entry, so a resume in
between does not hand the retry a live token."""
store = InMemoryExecutionStore()
scheduler = Mock(spec=Scheduler)
processor = CheckpointProcessor(store, scheduler)

start_input = StartDurableExecutionInput(
account_id="123456789012",
function_name="test-function",
function_qualifier="$LATEST",
execution_name="test-execution",
execution_timeout_seconds=300,
execution_retention_period_days=7,
invocation_id="inv-paused-idem",
)
execution = Execution.new(start_input)
execution.start()
execution.paused = True
store.save(execution)

inbound = CheckpointToken(
execution_arn=execution.durable_execution_arn, token_sequence=0
).to_str()
updates = [
OperationUpdate(
operation_id="step-A",
operation_type=OperationType.STEP,
action=OperationAction.START,
name="step-A",
)
]

first = processor.process_checkpoint(inbound, updates, "c1")
assert first.checkpoint_token is None

retry = processor.process_checkpoint(inbound, updates, "c1")
assert retry.checkpoint_token is None

resumed = store.load(execution.durable_execution_arn)
resumed.paused = False
store.save(resumed)

retry_after_resume = processor.process_checkpoint(inbound, updates, "c1")
assert retry_after_resume.checkpoint_token is None
Loading
Loading