Repository navigation
Conversation
8ef7704 to
51a4270
Compare
727a33c to
28d7528
Compare
Record creation before disposal so hook registrations replay correctly and tokens can be reused. Let independent token writes finish before propagating failures.
28d7528 to
2656a44
Compare
|
|
||
| # 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 |
There was a problem hiding this comment.
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.
| if failures: | ||
| raise failures[0] |
There was a problem hiding this comment.
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?)?
There was a problem hiding this comment.
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.
(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:Claim.wait()immediately added a hook to bothcontext.hooksandcontext.suspensions, without requiringawaitorasync for.context.suspensionsand recordedhook_createdevents (orhook_conflictbut let's not talk about that for now).But hooks can also be manually
dispose()-ed;async withalso disposes them on exit. Disposal immediately removed the hook fromcontext.suspensions, while keeping it incontext.hooksand marking it as disposed. On the next suspension, a separate disposal batch iteratedcontext.hooksand recordedhook_disposedfor 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:
On the initial execution,
await hooksuspended the workflow while the hook was still live, sohook_createdwas recorded. After a payload arrived, replay let thatawaitcomplete; the subsequentdispose()and suspension atsleep()causedhook_disposedto be recorded.However, this lost lifecycle events when a hook was created and disposed before the next flush:
The disposal batch attempted to dispose an unregistered hook and ignored the resulting
HookNotFoundError. The creation batch no longer saw it becausedispose()had already removed it fromcontext.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:
Here,
first.get_conflict()started waiting beforereplace()disposedfirstand reused its token. With the previous implementation, onlysecondwas registered. Its registration could succeed, but the pending check forfirstnever received ahook_createdorhook_conflictevent to resolve it. The entiregather()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:
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.hooksto record all hook lifecycle events, depending on the dict insertion order ofcontext.hooks. For parallelism, we group hooks by their tokens and record created/disposed events sequentially within each group.