Skip to content

workflow: fix hook lifecycle event ordering - #432

Open
fantix wants to merge 1 commit into
mainfrom
fantix/workflow-hook-flush-order
Open

fantix wants to merge 1 commit into
mainfrom
fantix/workflow-hook-flush-order

Conversation

@fantix

@fantix fantix commented Oct 1, 2026 •

Copy link
Copy Markdown
Member

(The text below is approximately 20% AI-generated)

Before this PR, hook lifecycle events (hook_created / hook_disposed) were flushed to the world when the workflow suspended:

  1. Claim.wait() immediately added a hook to both context.hooks and context.suspensions, without requiring await or async for.
  2. On the next workflow suspension, the creation batch iterated context.suspensions and recorded hook_created events (or hook_conflict but let's not talk about that for now).

But hooks can also be manually dispose()-ed; async with also disposes them on exit. Disposal immediately removed the hook from context.suspensions, while keeping it in context.hooks and marking it as disposed. On the next suspension, a separate disposal batch iterated context.hooks and recorded hook_disposed for such marked hooks. This batch ran before the creation batch, as ordered in #357, so that token can be properly reused after disposal.

This is approximately correct, worked fine for happy paths like:

hook = Claim.wait(token="shared")
payload = await hook
hook.dispose()
await sleep("1h")

On the initial execution, await hook suspended the workflow while the hook was still live, so hook_created was recorded. After a payload arrived, replay let that await complete; the subsequent dispose() and suspension at sleep() caused hook_disposed to be recorded.

However, this lost lifecycle events when a hook was created and disposed before the next flush:

hook = Claim.wait(token="shared")
hook.dispose()
await sleep("1h")

The disposal batch attempted to dispose an unregistered hook and ignored the resulting HookNotFoundError. The creation batch no longer saw it because dispose() had already removed it from context.suspensions. Neither lifecycle event was recorded.

Missing lifecycle events could also leave a workflow stuck when a pending registration check overlapped with disposal and token reuse:

first = Claim.wait(token="shared")

async def replace():
    first.dispose()
    second = Claim.wait(token="shared")
    return await second.get_conflict()

await asyncio.gather(first.get_conflict(), replace())

Here, first.get_conflict() started waiting before replace() disposed first and reused its token. With the previous implementation, only second was registered. Its registration could succeed, but the pending check for first never received a hook_created or hook_conflict event to resolve it. The entire gather() is then blocked forever.

The TS SDK keeps disposed hooks in its pending invocation queue and flushes hooks in creation order within each token group. It registers each hook before applying its requested disposal, then proceeds to the next hook. Assuming the token is available, the example above produces:

hook_created(first) → hook_disposed(first) → hook_created(second)

Both registration checks can then resolve during replay, allowing the workflow to complete.

This PR aligns Python's hook lifecycle event-flushing with the TS SDK. In short, we now only iterate context.hooks to record all hook lifecycle events, depending on the dict insertion order of context.hooks. For parallelism, we group hooks by their tokens and record created/disposed events sequentially within each group.

@fantix
fantix added this pull request to stack #434 October 1, 2026 17:12
@fantix
fantix force-pushed the fantix/workflow-hook-flush-order branch from 8ef7704 to 51a4270 Compare October 5, 2026 17:40
@fantix
fantix force-pushed the fantix/workflow-hook-flush-order branch from 727a33c to 28d7528 Compare October 6, 2026 15:12
@fantix fantix changed the title workflow: flush hook lifecycles in token order workflow: fix hook lifecycle event ordering Oct 6, 2026
@fantix
fantix marked this pull request as ready for review October 6, 2026 17:00
@fantix
fantix requested a review from a team October 6, 2026 17:00
Record creation before disposal so hook registrations replay correctly
and tokens can be reused. Let independent token writes finish before
propagating failures.
@fantix
fantix force-pushed the fantix/workflow-hook-flush-order branch from 28d7528 to 2656a44 Compare October 7, 2026 18:26

@msullivan msullivan left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Looks good!


# Now that the workflow is fully suspended and old hook tokens have been
# released, create all pending events in parallel. Steps are enqueued only
# once every event is durable: a step may hand a hook's token to whoever

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The comment about "steps are enqueued only once every event is durable" is dropped here, but the code is still there.

I think we can drop the steps_to_queue code now, though, since all hooks are flushed in any earlier step.

Comment on lines +1836 to +1837
if failures:
raise failures[0]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't have a strong opinion here really but is there any particular reason we are specifically grabbing the first failure instead of all of them in an exceptiongroup, and for why we are catching and collecting them instead of letting the task group collect them (do we specifically want one failure to not cancel the others?)?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Let me add a comment; I went through exactly the same thinking process as you did.

I added ExceptionGroup previously but dropped it because of 1) TS parity and 2) Python 3.10 requires a polyfill package. Though neither is a blocking issue to use ExceptionGroup, I figured we should probably keep it simple at first, then iterate later.

The task group in anyio will cancel other tasks if one task raised an error, but we still want other tasks to flush as much events as possible, in the same way as the TS SDK does. And, such anyio behavior is not as configurable as asyncio.wait(..., return_when= FIRST_EXCEPTION) which allows the remaining tasks to run to completion. Therefore, it ended up the way in the PR.

This branch was successfully deployed

1 active deployment
ci — 2656a440 Deployed Oct 7, 2026 by fantix via Test (Windows, py3.10, anyio lowest) #1677
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants