Skip to content

Commit 866cc78

Browse files
andystaplesCopilot
andauthored
Prevent timer callbacks after orchestration completion (#280)
Guard timer continuation at the executor callback boundary while preserving pre-terminal actions and trailing history processing. Cover terminal states, replay, native/chunked retries, cancellation, and Functions compatibility. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
1 parent f40ad56 commit 866cc78

6 files changed

Lines changed: 279 additions & 4 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,10 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
99

1010
FIXED
1111

12+
- Timer callbacks no longer schedule additional long-timer chunks, retry
13+
activities or sub-orchestrations, or resume orchestrator code after completion,
14+
failure, termination, or continue-as-new. Work scheduled before the terminal
15+
state is preserved.
1216
- `continue_as_new(..., save_events=True)` now preserves the global arrival
1317
order of unconsumed buffered external events across different event names,
1418
instead of grouping carryover events by name.

‎azure-functions-durable/CHANGELOG.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
99

1010
FIXED
1111

12+
- With the corresponding core `durabletask` SDK fix, timer callbacks no longer
13+
schedule additional long-timer chunks, retry activities or sub-orchestrations,
14+
or resume orchestrator code after completion, failure, termination, or
15+
continue-as-new. This applies to both native and compatibility orchestration APIs.
1216
- With the corresponding core `durabletask` SDK fix,
1317
`continue_as_new(..., save_events=True)` preserves the global arrival order
1418
of unconsumed buffered external events across different event names.

‎durabletask-azuremanaged/CHANGELOG.md‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,11 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
99

1010
FIXED
1111

12+
- With the corresponding core `durabletask` SDK fix, timer callbacks no longer
13+
retry activities or sub-orchestrations, or resume orchestrator code after
14+
completion, failure, termination, or continue-as-new. Azure Managed uses native
15+
long timers without chunking; the core long-timer chunking correction does not
16+
change that behavior.
1217
- With the corresponding core `durabletask` SDK fix,
1318
`continue_as_new(..., save_events=True)` preserves the global arrival order
1419
of unconsumed buffered external events across different event names.

‎durabletask/worker.py‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2576,6 +2576,10 @@ def process_event(
25762576
scheduled_time_ns=created_ns,
25772577
parent_trace_context=ctx._orchestration_trace_context or ctx._parent_trace_context, # pyright: ignore[reportPrivateUsage]
25782578
)
2579+
# A prior event in this batch may have ended the orchestration.
2580+
# Do not schedule another chunk, retry work, or resume user code.
2581+
if ctx._is_complete: # pyright: ignore[reportPrivateUsage]
2582+
return
25792583
next_fire_at = timer_task._handle_timer_fired(event.timerFired.fireAt.ToDatetime()) # pyright: ignore[reportPrivateUsage]
25802584
if next_fire_at is not None:
25812585
id = ctx.next_sequence_number()

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

Lines changed: 41 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@
1515
import json
1616
import threading
1717
from concurrent.futures import ThreadPoolExecutor
18-
from datetime import datetime
18+
from datetime import datetime, timedelta
1919
from types import SimpleNamespace
2020
from unittest.mock import AsyncMock, Mock
2121

@@ -27,6 +27,7 @@
2727
import azure.durable_functions as df
2828
from azure.durable_functions.internal import invocation, payloads
2929
from azure.durable_functions.worker import DurableFunctionsWorker
30+
from durabletask import task
3031
from durabletask.entities import EntityInstanceId
3132
from durabletask.payload import PayloadStore
3233

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

435436

437+
@pytest.mark.parametrize("native", [False, True])
438+
def test_long_timer_does_not_schedule_chunk_after_completion(native):
439+
def compatible_orchestrator(context):
440+
done = context.wait_for_external_event("done")
441+
timeout = context.create_timer(context.current_utc_datetime + timedelta(days=10))
442+
yield context.task_any([done, timeout])
443+
return "done"
444+
445+
def native_orchestrator(context: task.OrchestrationContext, _):
446+
done = context.wait_for_external_event("done")
447+
timeout = context.create_timer(timedelta(days=10))
448+
yield task.when_any([done, timeout])
449+
return "done"
450+
451+
start = datetime(2020, 1, 1)
452+
fire_at = start + timedelta(days=3)
453+
request = pb.OrchestratorRequest(
454+
instanceId=TEST_INSTANCE_ID,
455+
pastEvents=[
456+
helpers.new_orchestrator_started_event(start),
457+
helpers.new_execution_started_event("timer-race", TEST_INSTANCE_ID),
458+
helpers.new_timer_created_event(1, fire_at),
459+
],
460+
newEvents=[
461+
helpers.new_event_raised_event("done", json.dumps(True)),
462+
helpers.new_timer_fired_event(1, fire_at),
463+
],
464+
)
465+
encoded = base64.b64encode(request.SerializeToString()).decode("utf-8")
466+
orchestrator = native_orchestrator if native else compatible_orchestrator
467+
result = DurableFunctionsWorker().execute_orchestration_request(orchestrator, encoded)
468+
469+
response = _decode_orchestrator_response(result)
470+
assert len(response.actions) == 1
471+
completion = _get_completion_action(response)
472+
assert completion.orchestrationStatus == pb.ORCHESTRATION_STATUS_COMPLETED
473+
assert json.loads(completion.result.value) == "done"
474+
475+
436476
def test_execute_orchestration_request_registers_under_event_name():
437477
"""The orchestrator is registered under the name from the ExecutionStarted event."""
438478
def orchestrator(context):

‎tests/durabletask/test_orchestration_executor.py‎

Lines changed: 221 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -201,7 +201,8 @@ def delay_orchestrator(ctx: task.OrchestrationContext, _):
201201
assert actions[0].createTimer.fireAt.ToDatetime() == expected_fire_at
202202

203203

204-
def test_timer_fired_completion():
204+
@pytest.mark.parametrize("maximum_timer_interval", [timedelta(days=3), None])
205+
def test_timer_fired_completion(maximum_timer_interval):
205206
"""Tests the resumption of task using a timer_fired event"""
206207

207208
def delay_orchestrator(ctx: task.OrchestrationContext, _):
@@ -222,7 +223,10 @@ def delay_orchestrator(ctx: task.OrchestrationContext, _):
222223
new_events = [
223224
helpers.new_timer_fired_event(1, expected_fire_at)]
224225

225-
executor = worker._OrchestrationExecutor(registry, TEST_LOGGER, JsonDataConverter())
226+
executor = worker._OrchestrationExecutor(
227+
registry, TEST_LOGGER, JsonDataConverter(),
228+
maximum_timer_interval=maximum_timer_interval,
229+
)
226230
result = executor.execute(TEST_INSTANCE_ID, old_events, new_events)
227231
actions = result.actions
228232

@@ -352,6 +356,143 @@ def orchestrator(ctx: task.OrchestrationContext, _):
352356
assert complete_action.result.value == '"done"'
353357

354358

359+
@pytest.mark.parametrize("terminal_state", [
360+
"completed", "failed", "invalid_result", "terminated", "continued_as_new",
361+
])
362+
@pytest.mark.parametrize("replay_terminal_event", [False, True])
363+
def test_long_timer_does_not_continue_after_terminal_state(terminal_state, replay_terminal_event):
364+
def orchestrator(ctx: task.OrchestrationContext, _):
365+
done = ctx.wait_for_external_event("done")
366+
timeout = ctx.create_timer(timedelta(days=10))
367+
yield task.when_any([done, timeout])
368+
if terminal_state == "failed":
369+
raise ValueError("orchestrator failed")
370+
if terminal_state == "invalid_result":
371+
return object()
372+
if terminal_state == "continued_as_new":
373+
ctx.continue_as_new("next", save_events=True)
374+
return "done"
375+
376+
registry = worker._Registry()
377+
name = registry.add_orchestrator(orchestrator)
378+
executor = worker._OrchestrationExecutor(registry, TEST_LOGGER, JsonDataConverter())
379+
start = datetime(2020, 1, 1)
380+
fire_at = start + timedelta(days=3)
381+
old_events = [
382+
helpers.new_orchestrator_started_event(start),
383+
helpers.new_execution_started_event(name, TEST_INSTANCE_ID, encoded_input=None),
384+
helpers.new_timer_created_event(1, fire_at),
385+
]
386+
terminal_event = (
387+
helpers.new_terminated_event(encoded_output=json.dumps("terminated"))
388+
if terminal_state == "terminated"
389+
else helpers.new_event_raised_event("done", json.dumps(True))
390+
)
391+
new_events = [helpers.new_timer_fired_event(1, fire_at)]
392+
if replay_terminal_event:
393+
old_events.append(terminal_event)
394+
else:
395+
new_events.insert(0, terminal_event)
396+
new_events.append(helpers.new_event_raised_event("carryover", json.dumps("saved")))
397+
398+
result = executor.execute(TEST_INSTANCE_ID, old_events, new_events)
399+
400+
completion = get_and_validate_complete_orchestration_action_list(1, result.actions)
401+
if terminal_state in ("failed", "invalid_result"):
402+
assert completion.orchestrationStatus == pb.ORCHESTRATION_STATUS_FAILED
403+
assert completion.HasField("failureDetails")
404+
elif terminal_state == "terminated":
405+
assert completion.orchestrationStatus == pb.ORCHESTRATION_STATUS_TERMINATED
406+
assert completion.result.value == json.dumps("terminated")
407+
elif terminal_state == "continued_as_new":
408+
assert completion.orchestrationStatus == pb.ORCHESTRATION_STATUS_CONTINUED_AS_NEW
409+
assert completion.result.value == json.dumps("next")
410+
assert len(completion.carryoverEvents) == 1
411+
assert completion.carryoverEvents[0].eventRaised.name == "carryover"
412+
assert completion.carryoverEvents[0].eventRaised.input.value == json.dumps("saved")
413+
else:
414+
assert completion.orchestrationStatus == pb.ORCHESTRATION_STATUS_COMPLETED
415+
assert completion.result.value == json.dumps("done")
416+
417+
418+
@pytest.mark.parametrize("terminal_state", ["terminated", "continued_as_new"])
419+
def test_final_timer_callback_does_not_resume_terminal_orchestrator(terminal_state):
420+
def orchestrator(ctx: task.OrchestrationContext, _):
421+
done = ctx.wait_for_external_event("done")
422+
timeout = ctx.create_timer(timedelta(hours=1))
423+
yield task.when_any([done, timeout])
424+
if terminal_state == "continued_as_new":
425+
ctx.continue_as_new("next")
426+
yield timeout
427+
ctx.send_event("target", "unexpected")
428+
429+
registry = worker._Registry()
430+
name = registry.add_orchestrator(orchestrator)
431+
executor = worker._OrchestrationExecutor(registry, TEST_LOGGER, JsonDataConverter())
432+
start = datetime(2020, 1, 1)
433+
fire_at = start + timedelta(hours=1)
434+
old_events = [
435+
helpers.new_orchestrator_started_event(start),
436+
helpers.new_execution_started_event(name, TEST_INSTANCE_ID, encoded_input=None),
437+
helpers.new_timer_created_event(1, fire_at),
438+
]
439+
new_events = [
440+
helpers.new_terminated_event(encoded_output=None)
441+
if terminal_state == "terminated"
442+
else helpers.new_event_raised_event("done", json.dumps(True)),
443+
helpers.new_timer_fired_event(1, fire_at),
444+
]
445+
446+
result = executor.execute(TEST_INSTANCE_ID, old_events, new_events)
447+
448+
completion = get_and_validate_complete_orchestration_action_list(1, result.actions)
449+
expected_status = (
450+
pb.ORCHESTRATION_STATUS_TERMINATED
451+
if terminal_state == "terminated"
452+
else pb.ORCHESTRATION_STATUS_CONTINUED_AS_NEW
453+
)
454+
assert completion.orchestrationStatus == expected_status
455+
456+
457+
@pytest.mark.parametrize("cancel_timer", [False, True])
458+
def test_long_timer_preserves_work_scheduled_before_completion(cancel_timer):
459+
def orchestrator(ctx: task.OrchestrationContext, _):
460+
done = ctx.wait_for_external_event("done")
461+
timeout = ctx.create_timer(timedelta(days=10))
462+
yield task.when_any([done, timeout])
463+
if cancel_timer:
464+
assert timeout.cancel()
465+
ctx.send_event("target", "done", data=True)
466+
return "done"
467+
468+
registry = worker._Registry()
469+
name = registry.add_orchestrator(orchestrator)
470+
executor = worker._OrchestrationExecutor(registry, TEST_LOGGER, JsonDataConverter())
471+
start = datetime(2020, 1, 1)
472+
fire_at = start + timedelta(days=3)
473+
old_events = [
474+
helpers.new_orchestrator_started_event(start),
475+
helpers.new_execution_started_event(name, TEST_INSTANCE_ID, encoded_input=None),
476+
helpers.new_timer_created_event(1, fire_at),
477+
]
478+
new_events = [
479+
helpers.new_timer_fired_event(1, fire_at),
480+
helpers.new_event_raised_event("done", json.dumps(True)),
481+
]
482+
483+
result = executor.execute(TEST_INSTANCE_ID, old_events, new_events)
484+
485+
expected_actions = ["sendEvent", "completeOrchestration"]
486+
if not cancel_timer:
487+
expected_actions.insert(0, "createTimer")
488+
assert result.actions[0].id == 2
489+
assert result.actions[0].createTimer.fireAt.ToDatetime() == start + timedelta(days=6)
490+
assert [action.WhichOneof("orchestratorActionType") for action in result.actions] == expected_actions
491+
completion = result.actions[-1].completeOrchestration
492+
assert completion.orchestrationStatus == pb.ORCHESTRATION_STATUS_COMPLETED
493+
assert completion.result.value == json.dumps("done")
494+
495+
355496
def test_long_timer_can_be_cancelled_after_when_any_winner():
356497
"""Tests cancellation of a long timer after an external event wins when_any."""
357498

@@ -393,7 +534,10 @@ def orchestrator(ctx: task.OrchestrationContext, _):
393534
second = executor.execute(
394535
TEST_INSTANCE_ID,
395536
old_events,
396-
[helpers.new_event_raised_event("approval", json.dumps(True))],
537+
[
538+
helpers.new_event_raised_event("approval", json.dumps(True)),
539+
helpers.new_timer_fired_event(1, first_chunk_fire_at),
540+
],
397541
)
398542
complete_action = get_and_validate_complete_orchestration_action_list(1, second.actions)
399543
assert complete_action.orchestrationStatus == pb.ORCHESTRATION_STATUS_COMPLETED
@@ -1246,6 +1390,80 @@ def orchestrator(ctx: task.OrchestrationContext, orchestrator_input):
12461390
assert actions[-1].id == 1
12471391

12481392

1393+
@pytest.mark.parametrize("terminal_state", ["active", "completed", "terminated", "continued_as_new"])
1394+
@pytest.mark.parametrize("is_sub_orchestration", [False, True])
1395+
@pytest.mark.parametrize("retry_delay", [timedelta(hours=1), timedelta(days=10)])
1396+
@pytest.mark.parametrize("maximum_timer_interval", [timedelta(days=3), None])
1397+
def test_retry_timer_callback_respects_terminal_state(
1398+
terminal_state, is_sub_orchestration, retry_delay, maximum_timer_interval,
1399+
):
1400+
def orchestrator(ctx: task.OrchestrationContext, _):
1401+
retry_policy = task.RetryPolicy(
1402+
first_retry_interval=retry_delay, max_number_of_attempts=2,
1403+
)
1404+
if is_sub_orchestration:
1405+
work = ctx.call_sub_orchestrator("child", instance_id="child-instance", retry_policy=retry_policy)
1406+
else:
1407+
work = ctx.call_activity("activity", retry_policy=retry_policy)
1408+
done = ctx.wait_for_external_event("done")
1409+
yield task.when_any([done, work])
1410+
if terminal_state == "continued_as_new":
1411+
ctx.continue_as_new("next")
1412+
return "done"
1413+
1414+
registry = worker._Registry()
1415+
name = registry.add_orchestrator(orchestrator)
1416+
executor = worker._OrchestrationExecutor(
1417+
registry, TEST_LOGGER, JsonDataConverter(),
1418+
maximum_timer_interval=maximum_timer_interval,
1419+
)
1420+
start = datetime(2020, 1, 1)
1421+
interval = min(retry_delay, maximum_timer_interval) if maximum_timer_interval else retry_delay
1422+
fire_at = start + interval
1423+
old_events = [
1424+
helpers.new_orchestrator_started_event(start),
1425+
helpers.new_execution_started_event(name, TEST_INSTANCE_ID, encoded_input=None),
1426+
]
1427+
if is_sub_orchestration:
1428+
old_events.extend([
1429+
helpers.new_sub_orchestration_created_event(1, "child", "child-instance"),
1430+
helpers.new_sub_orchestration_failed_event(1, ValueError("retry me")),
1431+
])
1432+
else:
1433+
old_events.extend([
1434+
helpers.new_task_scheduled_event(1, "activity"),
1435+
helpers.new_task_failed_event(1, ValueError("retry me")),
1436+
])
1437+
old_events.append(helpers.new_timer_created_event(2, fire_at))
1438+
new_events = []
1439+
if terminal_state == "terminated":
1440+
new_events.append(helpers.new_terminated_event(encoded_output=None))
1441+
elif terminal_state != "active":
1442+
new_events.append(helpers.new_event_raised_event("done", json.dumps(True)))
1443+
new_events.append(helpers.new_timer_fired_event(2, fire_at))
1444+
1445+
result = executor.execute(TEST_INSTANCE_ID, old_events, new_events)
1446+
1447+
assert len(result.actions) == 1
1448+
if terminal_state == "active":
1449+
action = result.actions[0]
1450+
if interval < retry_delay:
1451+
assert action.HasField("createTimer")
1452+
assert action.id == 3
1453+
assert action.createTimer.fireAt.ToDatetime() == start + interval * 2
1454+
else:
1455+
assert action.HasField("createSubOrchestration" if is_sub_orchestration else "scheduleTask")
1456+
assert action.id == 1
1457+
else:
1458+
completion = get_and_validate_complete_orchestration_action_list(1, result.actions)
1459+
expected_status = {
1460+
"completed": pb.ORCHESTRATION_STATUS_COMPLETED,
1461+
"terminated": pb.ORCHESTRATION_STATUS_TERMINATED,
1462+
"continued_as_new": pb.ORCHESTRATION_STATUS_CONTINUED_AS_NEW,
1463+
}[terminal_state]
1464+
assert completion.orchestrationStatus == expected_status
1465+
1466+
12491467
def test_nondeterminism_expected_timer():
12501468
"""Tests the non-determinism detection logic when call_timer is expected but some other method (call_activity) is called instead"""
12511469
def dummy_activity(ctx, _):

0 commit comments

Comments
 (0)