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
14 changes: 14 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,20 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## Unreleased

FIXED

- Timer callbacks no longer schedule additional long-timer chunks, retry
activities or sub-orchestrations, or resume orchestrator code after completion,
failure, termination, or continue-as-new. Work scheduled before the terminal
state is preserved.
- `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.
- Fixed external events arriving after `continue_as_new(..., save_events=True)`
being lost to abandoned waits instead of carried into the next execution.
Events already delivered to live waits are not carried over, and
`save_events=False` still discards unprocessed events.

## v1.11.0

FIXED
Expand Down
13 changes: 13 additions & 0 deletions azure-functions-durable/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,19 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## Unreleased

FIXED

- With the corresponding core `durabletask` SDK fix, timer callbacks no longer
schedule additional long-timer chunks, retry activities or sub-orchestrations,
or resume orchestrator code after completion, failure, termination, or
continue-as-new. This applies to both native and compatibility orchestration APIs.
- 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.
- With a corrected core `durabletask` SDK, native durabletask orchestrators
preserve external events arriving after `continue_as_new(..., save_events=True)`
instead of losing them to abandoned waits.

## v2.0.0rc2

ADDED
Expand Down
14 changes: 14 additions & 0 deletions durabletask-azuremanaged/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,20 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## Unreleased

FIXED

- With the corresponding core `durabletask` SDK fix, timer callbacks no longer
retry activities or sub-orchestrations, or resume orchestrator code after
completion, failure, termination, or continue-as-new. Azure Managed uses native
long timers without chunking; the core long-timer chunking correction does not
change that behavior.
- 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.
- With a corrected core `durabletask` SDK, external events arriving after
`continue_as_new(..., save_events=True)` are no longer lost to abandoned waits
and are carried into the next execution.

## v1.11.0

ADDED
Expand Down
6 changes: 6 additions & 0 deletions durabletask/task.py
Original file line number Diff line number Diff line change
Expand Up @@ -393,12 +393,18 @@ def continue_as_new(self, new_input: Any, *, save_events: bool = False,
new_version: str | None = None) -> None:
"""Continue the orchestration execution as a new instance.

Orchestrators should return immediately after calling this method.
Subsequent external events are no longer delivered to pending waits
in the current execution.

Parameters
----------
new_input : Any
The new input to use for the new orchestration instance.
save_events : bool
A flag indicating whether to add any unprocessed external events in the new orchestration history.
Events already delivered to a waiting task are not saved, even if
the orchestrator has not yielded that task.
new_version : str | None
An optional version to assign to the new orchestration instance.
"""
Expand Down
36 changes: 24 additions & 12 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 @@ -2571,6 +2576,10 @@ def process_event(
scheduled_time_ns=created_ns,
parent_trace_context=ctx._orchestration_trace_context or ctx._parent_trace_context, # pyright: ignore[reportPrivateUsage]
)
# A prior event in this batch may have ended the orchestration.
# Do not schedule another chunk, retry work, or resume user code.
if ctx._is_complete: # pyright: ignore[reportPrivateUsage]
return
next_fire_at = timer_task._handle_timer_fired(event.timerFired.fireAt.ToDatetime()) # pyright: ignore[reportPrivateUsage]
if next_fire_at is not None:
id = ctx.next_sequence_number()
Expand Down Expand Up @@ -2843,7 +2852,9 @@ def _cancel_timer() -> None:
self._logger.info(f"{ctx.instance_id} Event raised: {event_name}")
task_list = ctx._pending_events.get(event_name, None) # pyright: ignore[reportPrivateUsage]
decoded_result: Any | None = None
if task_list:
# Completed executions leave abandoned waits behind. Buffer
# trailing events instead, so continue-as-new can carry them over.
if task_list and not ctx._is_complete: # pyright: ignore[reportPrivateUsage]
event_task = task_list.pop(0)
if not ph.is_empty(event.eventRaised.input):
decoded_result = self._data_converter.deserialize(
Expand All @@ -2866,10 +2877,11 @@ 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."
f"{ctx.instance_id}: Event '{event_name}' has been buffered as there are no active tasks waiting for it."
)
elif event.HasField("executionSuspended"):
if not self._is_suspended and not ctx.is_replaying:
Expand Down
73 changes: 72 additions & 1 deletion tests/azure-functions-durable/test_worker_compat.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
import json
import threading
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime
from datetime import datetime, timedelta
from types import SimpleNamespace
from unittest.mock import AsyncMock, Mock

Expand All @@ -27,6 +27,7 @@
import azure.durable_functions as df
from azure.durable_functions.internal import invocation, payloads
from azure.durable_functions.worker import DurableFunctionsWorker
from durabletask import task
from durabletask.entities import EntityInstanceId
from durabletask.payload import PayloadStore

Expand Down Expand Up @@ -433,6 +434,45 @@ def orchestrator(context):
assert json.loads(completion.result.value) == {"echo": {"n": 5}}


@pytest.mark.parametrize("native", [False, True])
def test_long_timer_does_not_schedule_chunk_after_completion(native):
def compatible_orchestrator(context):
done = context.wait_for_external_event("done")
timeout = context.create_timer(context.current_utc_datetime + timedelta(days=10))
yield context.task_any([done, timeout])
return "done"

def native_orchestrator(context: task.OrchestrationContext, _):
done = context.wait_for_external_event("done")
timeout = context.create_timer(timedelta(days=10))
yield task.when_any([done, timeout])
return "done"

start = datetime(2020, 1, 1)
fire_at = start + timedelta(days=3)
request = pb.OrchestratorRequest(
instanceId=TEST_INSTANCE_ID,
pastEvents=[
helpers.new_orchestrator_started_event(start),
helpers.new_execution_started_event("timer-race", TEST_INSTANCE_ID),
helpers.new_timer_created_event(1, fire_at),
],
newEvents=[
helpers.new_event_raised_event("done", json.dumps(True)),
helpers.new_timer_fired_event(1, fire_at),
],
)
encoded = base64.b64encode(request.SerializeToString()).decode("utf-8")
orchestrator = native_orchestrator if native else compatible_orchestrator
result = DurableFunctionsWorker().execute_orchestration_request(orchestrator, encoded)

response = _decode_orchestrator_response(result)
assert len(response.actions) == 1
completion = _get_completion_action(response)
assert completion.orchestrationStatus == pb.ORCHESTRATION_STATUS_COMPLETED
assert json.loads(completion.result.value) == "done"


def test_execute_orchestration_request_registers_under_event_name():
"""The orchestrator is registered under the name from the ExecutionStarted event."""
def orchestrator(context):
Expand Down Expand Up @@ -502,6 +542,37 @@ def orchestrator(context):
assert "boom" in completion.failureDetails.errorMessage


@pytest.mark.parametrize("save_events", [True, False])
def test_continue_as_new_preserves_trailing_events_in_worker_response(save_events: bool):
def orchestrator(ctx: task.OrchestrationContext, _):
event_task = ctx.wait_for_external_event("event")
timer = ctx.create_timer(timedelta(seconds=1))
yield task.when_any([event_task, timer])
ctx.continue_as_new(None, save_events=save_events)

started_at = datetime(2026, 1, 1)
fire_at = started_at + timedelta(seconds=1)
request = pb.OrchestratorRequest(instanceId=TEST_INSTANCE_ID)
request.pastEvents.extend([
helpers.new_orchestrator_started_event(started_at),
helpers.new_execution_started_event("continue-events", TEST_INSTANCE_ID),
helpers.new_timer_created_event(1, fire_at),
])
request.newEvents.extend([
helpers.new_timer_fired_event(1, fire_at),
helpers.new_event_raised_event("event", "1"),
helpers.new_event_raised_event("event", "2"),
])
encoded = base64.b64encode(request.SerializeToString()).decode("utf-8")
response = _decode_orchestrator_response(
DurableFunctionsWorker().execute_orchestration_request(orchestrator, encoded))
completion = _get_completion_action(response)

assert completion.orchestrationStatus == pb.ORCHESTRATION_STATUS_CONTINUED_AS_NEW
assert [event.eventRaised.input.value for event in completion.carryoverEvents] == (
["1", "2"] if save_events else [])


def test_activity_retry_then_fan_out_uses_distinct_task_ids():
"""Regression test for Azure/azure-functions-durable-python#603."""
def orchestrator(context):
Expand Down
Loading
Loading