Skip to content

Commit 3a2ac21

Browse files
andystaplesCopilotberndverst
authored
Preserve trailing external events after continue-as-new (#281)
Gate ordinary external-event delivery on the existing terminal context boundary while retaining entity routing and raw carryover buffering. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Co-authored-by: Bernd Verst <github@bernd.dev>
1 parent 5cdfeef commit 3a2ac21

7 files changed

Lines changed: 345 additions & 2 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,10 @@ state is preserved.
1818
- `continue_as_new(..., save_events=True)` now preserves the global arrival
1919
order of unconsumed buffered external events across different event names,
2020
instead of grouping carryover events by name.
21+
- Fixed external events arriving after `continue_as_new(..., save_events=True)`
22+
being lost to abandoned waits instead of carried into the next execution.
23+
Events already delivered to live waits are not carried over, and
24+
`save_events=False` still discards unprocessed events.
2125

2226
## v1.11.0
2327

‎azure-functions-durable/CHANGELOG.md‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,9 @@ continue-as-new. This applies to both native and compatibility orchestration API
2020
- With the corresponding core `durabletask` SDK fix,
2121
`continue_as_new(..., save_events=True)` preserves the global arrival order
2222
of unconsumed buffered external events across different event names.
23+
- With a corrected core `durabletask` SDK, native durabletask orchestrators
24+
preserve external events arriving after `continue_as_new(..., save_events=True)`
25+
instead of losing them to abandoned waits.
2326

2427
## v2.0.0rc2
2528

‎durabletask-azuremanaged/CHANGELOG.md‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,9 @@ change that behavior.
2020
- With the corresponding core `durabletask` SDK fix,
2121
`continue_as_new(..., save_events=True)` preserves the global arrival order
2222
of unconsumed buffered external events across different event names.
23+
- With a corrected core `durabletask` SDK, external events arriving after
24+
`continue_as_new(..., save_events=True)` are no longer lost to abandoned waits
25+
and are carried into the next execution.
2326

2427
## v1.11.0
2528

‎durabletask/task.py‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -393,12 +393,18 @@ def continue_as_new(self, new_input: Any, *, save_events: bool = False,
393393
new_version: str | None = None) -> None:
394394
"""Continue the orchestration execution as a new instance.
395395
396+
Orchestrators should return immediately after calling this method.
397+
Subsequent external events are no longer delivered to pending waits
398+
in the current execution.
399+
396400
Parameters
397401
----------
398402
new_input : Any
399403
The new input to use for the new orchestration instance.
400404
save_events : bool
401405
A flag indicating whether to add any unprocessed external events in the new orchestration history.
406+
Events already delivered to a waiting task are not saved, even if
407+
the orchestrator has not yielded that task.
402408
new_version : str | None
403409
An optional version to assign to the new orchestration instance.
404410
"""

‎durabletask/worker.py‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2852,7 +2852,9 @@ def _cancel_timer() -> None:
28522852
self._logger.info(f"{ctx.instance_id} Event raised: {event_name}")
28532853
task_list = ctx._pending_events.get(event_name, None) # pyright: ignore[reportPrivateUsage]
28542854
decoded_result: Any | None = None
2855-
if task_list:
2855+
# Completed executions leave abandoned waits behind. Buffer
2856+
# trailing events instead, so continue-as-new can carry them over.
2857+
if task_list and not ctx._is_complete: # pyright: ignore[reportPrivateUsage]
28562858
event_task = task_list.pop(0)
28572859
if not ph.is_empty(event.eventRaised.input):
28582860
decoded_result = self._data_converter.deserialize(
@@ -2879,7 +2881,7 @@ def _cancel_timer() -> None:
28792881
ctx._received_event_sequence += 1 # pyright: ignore[reportPrivateUsage]
28802882
if not ctx.is_replaying:
28812883
self._logger.info(
2882-
f"{ctx.instance_id}: Event '{event_name}' has been buffered as there are no tasks waiting for it."
2884+
f"{ctx.instance_id}: Event '{event_name}' has been buffered as there are no active tasks waiting for it."
28832885
)
28842886
elif event.HasField("executionSuspended"):
28852887
if not self._is_suspended and not ctx.is_replaying:

‎tests/azure-functions-durable/test_worker_compat.py‎

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -542,6 +542,37 @@ def orchestrator(context):
542542
assert "boom" in completion.failureDetails.errorMessage
543543

544544

545+
@pytest.mark.parametrize("save_events", [True, False])
546+
def test_continue_as_new_preserves_trailing_events_in_worker_response(save_events: bool):
547+
def orchestrator(ctx: task.OrchestrationContext, _):
548+
event_task = ctx.wait_for_external_event("event")
549+
timer = ctx.create_timer(timedelta(seconds=1))
550+
yield task.when_any([event_task, timer])
551+
ctx.continue_as_new(None, save_events=save_events)
552+
553+
started_at = datetime(2026, 1, 1)
554+
fire_at = started_at + timedelta(seconds=1)
555+
request = pb.OrchestratorRequest(instanceId=TEST_INSTANCE_ID)
556+
request.pastEvents.extend([
557+
helpers.new_orchestrator_started_event(started_at),
558+
helpers.new_execution_started_event("continue-events", TEST_INSTANCE_ID),
559+
helpers.new_timer_created_event(1, fire_at),
560+
])
561+
request.newEvents.extend([
562+
helpers.new_timer_fired_event(1, fire_at),
563+
helpers.new_event_raised_event("event", "1"),
564+
helpers.new_event_raised_event("event", "2"),
565+
])
566+
encoded = base64.b64encode(request.SerializeToString()).decode("utf-8")
567+
response = _decode_orchestrator_response(
568+
DurableFunctionsWorker().execute_orchestration_request(orchestrator, encoded))
569+
completion = _get_completion_action(response)
570+
571+
assert completion.orchestrationStatus == pb.ORCHESTRATION_STATUS_CONTINUED_AS_NEW
572+
assert [event.eventRaised.input.value for event in completion.carryoverEvents] == (
573+
["1", "2"] if save_events else [])
574+
575+
545576
def test_activity_retry_then_fan_out_uses_distinct_task_ids():
546577
"""Regression test for Azure/azure-functions-durable-python#603."""
547578
def orchestrator(context):

0 commit comments

Comments
 (0)