Skip to content
Merged
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
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions azure-functions-durable/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions durabletask-azuremanaged/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
26 changes: 16 additions & 10 deletions durabletask/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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."
Expand Down
152 changes: 152 additions & 0 deletions tests/durabletask/test_orchestration_executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
Loading