Conversation
There was a problem hiding this comment.
🟡 Changes recommended
Existing cancellation details can be overwritten, and the regression assertion does not guarantee the new fallback ran.
Once you've addressed the issues Copilot identified, you can request another Copilot review.
Pull request overview
Prevents worker shutdown from hanging on activities no longer tracked by Core.
Changes:
- Cancels orphaned activities after polling stops.
- Adds a regression test and changelog entry.
File summaries
| File | Description |
|---|---|
temporalio/worker/_activity.py |
Cancels remaining activities during shutdown. |
tests/worker/test_workflow.py |
Tests orphaned local-activity shutdown. |
CHANGELOG.md |
Documents the fix. |
Review details
- Files reviewed: 3/3 changed files
- Comments generated: 3
- Review effort level: Balanced
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
b2b1bc7 to
2c5b821
Compare
ea32d2d to
807fd47
Compare
|
I don't really want this to be lang specific code. I believe having If "On a workflow eviction Core queues a cancel for a running local activity, then invalidates the run and drops the activity from its outstanding set, so the queued cancel is discarded as untracked while Python keeps running the activity." is accurate (which, seems it can be under certain conditions). Then we should just fix that in core. All cancel requests for all open LAs should get delivered before we send out the eviction request. That way they can't get lost. If lang chooses to ignore the cancels, fine, cool, but at least they all get acked before we evict. |
When a workflow run is evicted while one of its local activities is executing, Core queues a cancel for the activity and then, once the eviction activation completes, invalidates the run, which removes the activity from Core's outstanding set. If the cancel is still queued at that point Core discards it as no longer tracked, while the Python worker keeps running the activity. Core then reports activity polling as finished and wait_all_completed waits for that task forever, so worker shutdown hangs. Once Core's activity poll has shut down it tracks no activity, so anything still executing can only be finished from here. Cancel it with worker_shutdown cancellation details, log a warning, and let its completion be ignored by Core as untracked. This is the shutdown hang behind the frequent macOS CI timeouts of test_workflow_cancel_activity[True], whose captured logs show a workflow task eviction followed by two local activity starts for one run and one fewer cancel. The regression test forces the same race by delaying the activity poll that follows a local activity start and terminating the workflow so the workflow task heartbeat fails.
The regression test forced the lost-cancel race with wall-clock timing: a 3s delay on the activity poll that follows a local activity start and a 4s wait after terminating the workflow. On the slow fips CI job shutdown still hung for its 20s limit twice, and the fallback cancel never ran, so the worker was still waiting for its pollers to exit when the test gave up. Gate the test on events instead. The activity poll after the local activity start is held until the workflow poller has returned shutdown: Core only ends the workflow stream once it has processed the eviction completion that invalidates the run, and it marks workflows as shut down for the local activity manager before returning, so the cancel queued at eviction can never reach the worker and shutdown has to cancel the activity itself. The assertion now requires worker_shutdown details. A second test covers the other ordering: the cancel is delivered while Core still tracks the activity, the activity ignores it, and eviction is held until then so the run is invalidated afterwards. Shutdown must still end the activity without replacing its cancel_requested details, which are set once, so the fallback now only sets worker_shutdown details when none were recorded and the changelog entry says so. The workflow task timeout is 3s so that on a loaded runner the task cannot time out before the terminate lands, which replays the workflow and starts a second local activity.
…ding The new cancel in wait_all_completed can land while an orphaned activity is still encoding its failure through a codec that yields, which ends the activity task with CancelledError. gather() re-raised it, Worker._run skipped finalize_shutdown, and shutdown() hung, the very symptom this branch fixes. Collect task exceptions instead, logging anything other than cancellation, so wait_all_completed keeps its no-raise contract. Add a regression test with a slow payload codec, share the poll-hold setup between the orphan tests, and move the changelog entry from the released 1.33.0 section to Unreleased.
807fd47 to
2fffae5
Compare
|
Closing in favour of a Core fix, per Spencer's comment above: the lost local-activity cancel will be handled in sdk-core so every open LA's cancel is delivered before the eviction request goes out, and lang stays out of it. The Python-only problem this surfaced (a cancel landing while an activity's result is still being encoded) is fixed separately in #1910. I'll link the sdk-core PR here once it's open |
|
The Core side is up as a draft: temporalio/sdk-rust#1620 (withholds the eviction activation until every queued local-activity cancel has been polled) |
When a run was evicted, Core queued a cancel for each of its in-flight local activities and produced the eviction activation in the same pass. If lang completed the eviction before its activity poll reached the cancel, invalidating the run made the LA manager drop the cancel as untracked, so lang kept running the local activity while Core reported local activity shutdown complete (temporalio/sdk-python#1837). The eviction activation is now withheld while any of the run's in-flight local activities has a queued, not yet polled cancel, and the LA manager wakes the run once the last one is handed to lang. Cancels are queued once per attempt, a local activity still waiting for dispatch resolves as cancelled instead of starting, and an attempt that fails after its cancel was requested is not retried locally. Local activity results that arrive while the eviction is withheld are discarded, as they already were once the eviction activation was outstanding, so the cancellation the eviction caused never reaches workflow code as an activation. CancelAllInRun yields no immediate resolutions for the same reason: it is only sunk for runs that are terminating or evicting, and applying them surfaced spurious cancellations for backing-off or never-dispatched local activities.
What was changed
After activity polling has shut down, the worker cancels any activity still running, with a warning and, if it has no cancellation details yet,
worker_shutdowncancellation details, before waiting for completions. Regression tests and CHANGELOG entry.wait_all_completedalso collects the activity tasks' exceptions instead of letting them escape: the new cancel can land while an orphaned activity is still encoding its failure through a codec that yields, and the escapingCancelledErrorused to skipfinalize_shutdownand hangshutdown(). A regression test with a slow payload codec covers that path.Why
On a workflow eviction Core queues a cancel for a running local activity, then invalidates the run and drops the activity from its outstanding set, so the queued cancel is discarded as untracked while Python keeps running the activity. Core then reports shutdown complete and
wait_all_completedwaits forever; this is the 60s timeout intest_workflow_cancel_activity[True]across 17 CI attempts. Once polling is down Core tracks nothing, so anything still running is orphaned. The lost cancel itself is a Core race worth an sdk-core issue.Testing
Two tests force each ordering with events instead of sleeps. One holds the activity poll until the workflow poller has shut down, so the queued cancel is never delivered and shutdown must cancel the activity with
worker_shutdowndetails. The other delivers the cancel first, has the activity ignore it, and checks that the shutdown cancel keeps itscancel_requesteddetails. Both hang for the 20s limit without the fix; 100/100 and 60/60 runs under heavy load on Python 3.14 and 3.10. Lint clean.