From 6c7b98e5cec4972c4ed0b86d9c8a77167be2d7df Mon Sep 17 00:00:00 2001 From: Andy Staples Date: Tue, 29 Sep 2026 11:51:52 -0600 Subject: [PATCH] Preserve global carryover event arrival order Track buffered event arrival indexes independently of task IDs and retain per-name FIFO consumption and raw payloads. Add replay, selective consumption, cancellation, converter, and save-events regression coverage. Fixes #277 Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- CHANGELOG.md | 6 + azure-functions-durable/CHANGELOG.md | 6 + durabletask-azuremanaged/CHANGELOG.md | 6 + durabletask/worker.py | 26 +-- .../test_orchestration_executor.py | 152 ++++++++++++++++++ 5 files changed, 186 insertions(+), 10 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 9a320ea7..4e99fdd5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,12 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). ## Unreleased +FIXED + +- `continue_as_new(..., save_events=True)` now preserves the global arrival +order of unconsumed buffered external events across different event names, +instead of grouping carryover events by name. + ## v1.11.0 FIXED diff --git a/azure-functions-durable/CHANGELOG.md b/azure-functions-durable/CHANGELOG.md index 5944f4bb..845d5a03 100644 --- a/azure-functions-durable/CHANGELOG.md +++ b/azure-functions-durable/CHANGELOG.md @@ -7,6 +7,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## Unreleased +FIXED + +- With the corresponding core `durabletask` SDK fix, +`continue_as_new(..., save_events=True)` preserves the global arrival order +of unconsumed buffered external events across different event names. + ## v2.0.0rc2 ADDED diff --git a/durabletask-azuremanaged/CHANGELOG.md b/durabletask-azuremanaged/CHANGELOG.md index d26c18ff..24b45f5e 100644 --- a/durabletask-azuremanaged/CHANGELOG.md +++ b/durabletask-azuremanaged/CHANGELOG.md @@ -7,6 +7,12 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). ## Unreleased +FIXED + +- With the corresponding core `durabletask` SDK fix, +`continue_as_new(..., save_events=True)` preserves the global arrival order +of unconsumed buffered external events across different event names. + ## v1.11.0 ADDED diff --git a/durabletask/worker.py b/durabletask/worker.py index f4f31a32..0fec1b4e 100644 --- a/durabletask/worker.py +++ b/durabletask/worker.py @@ -1534,7 +1534,9 @@ def __init__(self, self._version: str | None = None self._parent_instance_id: str | None = None self._completion_status: pb.OrchestrationStatus | None = None - self._received_events: dict[str, list[str | None]] = {} + # Keep arrival indexes separate from task IDs and retain raw payloads. + self._received_events: dict[str, list[tuple[int, str | None]]] = {} + self._received_event_sequence = 0 self._pending_events: dict[str, list[task.CancellableTask[Any]]] = {} self._new_input: Any | None = None self._new_version: str | None = None @@ -1674,13 +1676,16 @@ def get_actions(self) -> list[pb.OrchestratorAction]: carryover_events = [] # We need to save the current set of pending events so that they can be # replayed when the new instance starts. - for event_name, values in self._received_events.items(): - for event_value in values: - # Buffered events are stored as their raw JSON payload - # (or None), so carry them over as-is without re-encoding. - carryover_events.append( - ph.new_event_raised_event(event_name, event_value) - ) + buffered_events = sorted( + (index, event_name, value) + for event_name, values in self._received_events.items() + for index, value in values + ) + for _, event_name, event_value in buffered_events: + # Carry over raw JSON payloads (or None) without re-encoding. + carryover_events.append( + ph.new_event_raised_event(event_name, event_value) + ) action = ph.new_complete_orchestration_action( self.next_sequence_number(), pb.ORCHESTRATION_STATUS_CONTINUED_AS_NEW, @@ -2109,7 +2114,7 @@ def wait_for_external_event(self, name: str, *, event_name = name.casefold() event_list = self._received_events.get(event_name, None) if event_list: - event_data = event_list.pop(0) + _, event_data = event_list.pop(0) if not event_list: del self._received_events[event_name] external_event_task.complete(self._data_converter.deserialize(event_data, data_type)) @@ -2866,7 +2871,8 @@ def _cancel_timer() -> None: buffered_payload: str | None = None if not ph.is_empty(event.eventRaised.input): buffered_payload = event.eventRaised.input.value - event_list.append(buffered_payload) + event_list.append((ctx._received_event_sequence, buffered_payload)) # pyright: ignore[reportPrivateUsage] + ctx._received_event_sequence += 1 # pyright: ignore[reportPrivateUsage] if not ctx.is_replaying: self._logger.info( f"{ctx.instance_id}: Event '{event_name}' has been buffered as there are no tasks waiting for it." diff --git a/tests/durabletask/test_orchestration_executor.py b/tests/durabletask/test_orchestration_executor.py index 4c213403..b78df02b 100644 --- a/tests/durabletask/test_orchestration_executor.py +++ b/tests/durabletask/test_orchestration_executor.py @@ -1857,6 +1857,158 @@ def orchestrator(ctx: task.OrchestrationContext, input: int): assert event.eventRaised.input.value == json.dumps(42 + i) +@pytest.mark.parametrize("save_events", [True, False]) +@pytest.mark.parametrize("replayed_events", [0, 1, 3]) +@pytest.mark.parametrize("timer_fires_first", [True, False]) +def test_continue_as_new_preserves_global_event_order( + save_events: bool, replayed_events: int, timer_fires_first: bool): + def orchestrator(ctx: task.OrchestrationContext, _): + yield ctx.create_timer(ctx.current_utc_datetime + timedelta(seconds=1)) + ctx.continue_as_new(None, save_events=save_events) + + registry = worker._Registry() + name = registry.add_orchestrator(orchestrator) + start_time = datetime(2026, 1, 1) + fire_at = start_time + timedelta(seconds=1) + raised_events = [ + helpers.new_event_raised_event("A", "1"), + helpers.new_event_raised_event("b", "2"), + helpers.new_event_raised_event("a", "3"), + ] + old_events = [ + helpers.new_orchestrator_started_event(start_time), + helpers.new_execution_started_event(name, TEST_INSTANCE_ID), + helpers.new_timer_created_event(1, fire_at), + *raised_events[:replayed_events], + ] + timer_fired = helpers.new_timer_fired_event(1, fire_at) + new_events = raised_events[replayed_events:] + new_events.insert(0 if timer_fires_first else len(new_events), timer_fired) + + executor = worker._OrchestrationExecutor(registry, TEST_LOGGER, JsonDataConverter()) + result = executor.execute(TEST_INSTANCE_ID, old_events, new_events) + + complete_action = get_and_validate_complete_orchestration_action_list(1, result.actions) + assert complete_action.orchestrationStatus == pb.ORCHESTRATION_STATUS_CONTINUED_AS_NEW + assert result.actions[0].id == 2 + expected = [("a", "1"), ("b", "2"), ("a", "3")] if save_events else [] + assert [ + (event.eventRaised.name, event.eventRaised.input.value) + for event in complete_action.carryoverEvents + ] == expected + + +@pytest.mark.parametrize("replay_consumption", [True, False]) +def test_continue_as_new_preserves_order_after_selective_event_consumption( + replay_consumption: bool): + def orchestrator(ctx: task.OrchestrationContext, _): + live_result = yield ctx.wait_for_external_event("LIVE") + cancelled_wait = ctx.wait_for_external_event("C") + cancelled_wait.cancel() + yield ctx.create_timer(ctx.current_utc_datetime + timedelta(seconds=1)) + first = yield ctx.wait_for_external_event("a", data_type=int) + second = yield ctx.wait_for_external_event("A", data_type=int) + third = yield ctx.wait_for_external_event("c", data_type=int) + yield ctx.create_timer(ctx.current_utc_datetime + timedelta(seconds=2)) + ctx.continue_as_new([live_result, first, second, third], save_events=True) + + registry = worker._Registry() + name = registry.add_orchestrator(orchestrator) + start_time = datetime(2026, 1, 1) + first_fire_at = start_time + timedelta(seconds=1) + second_fire_at = start_time + timedelta(seconds=2) + history = [ + helpers.new_orchestrator_started_event(start_time), + helpers.new_execution_started_event(name, TEST_INSTANCE_ID), + helpers.new_event_raised_event("live", "0"), + helpers.new_timer_created_event(1, first_fire_at), + helpers.new_event_raised_event("A", "1"), + helpers.new_event_raised_event("b", "2"), + helpers.new_event_raised_event("a", "3"), + helpers.new_event_raised_event("C", "4"), + helpers.new_event_raised_event("B", "5"), + helpers.new_event_raised_event("A", "6"), + helpers.new_timer_fired_event(1, first_fire_at), + helpers.new_timer_created_event(2, second_fire_at), + ] + trailing_events = [ + helpers.new_event_raised_event("c", "7"), + helpers.new_event_raised_event("a", "8"), + helpers.new_event_raised_event("b", "9"), + helpers.new_timer_fired_event(2, second_fire_at), + ] + old_events = history if replay_consumption else [] + new_events = trailing_events if replay_consumption else history + trailing_events + executor = worker._OrchestrationExecutor(registry, TEST_LOGGER, JsonDataConverter()) + result = executor.execute(TEST_INSTANCE_ID, old_events, new_events) + + complete_action = get_and_validate_complete_orchestration_action_list(1, result.actions) + assert complete_action.orchestrationStatus == pb.ORCHESTRATION_STATUS_CONTINUED_AS_NEW + assert result.actions[0].id == 3 + assert json.loads(complete_action.result.value) == [0, 1, 3, 4] + assert [ + (event.eventRaised.name, event.eventRaised.input.value) + for event in complete_action.carryoverEvents + ] == [("b", "2"), ("b", "5"), ("a", "6"), ("c", "7"), ("a", "8"), ("b", "9")] + + +@pytest.mark.parametrize("consumed_payload", [None, "null", "23"]) +def test_continue_as_new_preserves_order_and_raw_event_payloads(consumed_payload: str | None): + class RecordingConverter(JsonDataConverter): + def __init__(self): + self.serialized: list[Any] = [] + self.deserialized: list[tuple[str | None, type | None]] = [] + + def serialize(self, value: Any) -> str | None: + self.serialized.append(value) + return None if value is None else json.dumps({"wrapped": value}) + + def deserialize(self, data: str | None, target_type: type | None = None) -> Any: + self.deserialized.append((data, target_type)) + return super().deserialize(data, target_type) + + def orchestrator(ctx: task.OrchestrationContext, _): + yield ctx.create_timer(ctx.current_utc_datetime + timedelta(seconds=1)) + consumed = yield ctx.wait_for_external_event("Consume", data_type=int) + ctx.continue_as_new(consumed, save_events=True) + + registry = worker._Registry() + name = registry.add_orchestrator(orchestrator) + start_time = datetime(2026, 1, 1) + fire_at = start_time + timedelta(seconds=1) + expected = [ + ("a", '{ "value" : [1, 2] }'), + ("b", None), + ("a", "null"), + ("b", '"hello\\u0020world"'), + ("a", '{ "value" : [1, 2] }'), + ("c", "false"), + ("b", "0"), + ("a", '""'), + ] + old_events = [ + helpers.new_orchestrator_started_event(start_time), + helpers.new_execution_started_event(name, TEST_INSTANCE_ID), + helpers.new_timer_created_event(1, fire_at), + helpers.new_event_raised_event("consume", consumed_payload), + *[helpers.new_event_raised_event(name, payload) for name, payload in expected], + ] + converter = RecordingConverter() + executor = worker._OrchestrationExecutor(registry, TEST_LOGGER, converter) + result = executor.execute( + TEST_INSTANCE_ID, old_events, [helpers.new_timer_fired_event(1, fire_at)]) + + complete_action = get_and_validate_complete_orchestration_action_list(1, result.actions) + assert complete_action.orchestrationStatus == pb.ORCHESTRATION_STATUS_CONTINUED_AS_NEW + assert [ + (event.eventRaised.name, + event.eventRaised.input.value if event.eventRaised.HasField("input") else None) + for event in complete_action.carryoverEvents + ] == expected + assert converter.deserialized == [(consumed_payload, int)] + assert converter.serialized == [json.loads(consumed_payload) if consumed_payload else None] + + def test_fan_out(): """Tests that a fan-out pattern correctly schedules N tasks""" def hello(_, name: str):