lifecycle
Lifecycle signalling for DAG execution.
A single LifecycleNotifier signals a DAG's lifecycle transitions — accepted
(start), completed, failed — to whichever destinations are wired for the run:
- the cache
Runrow (marked complete / failed), whencache+run_idare supplied; and - the task initiator's
_WorkerMailbox(protoJOB_ACCEPT/TASK_COMPLETE/TASK_ABORT), when amailboxis supplied — used forworker_onlyinteractive DAGs triggered by a gRPC JOB_REQUEST, where the initiator (the envelope sender) has no modeller process listening.
Construct with the destinations available for a run:
cache+run_id→ the cache Run row is closed out;mailbox→ lifecycle messages are sent to the initiator;- neither → no-op (tests, or a run with nothing to signal);
- both → both destinations are updated.
Each destination self-guards on being present and swallows + logs its own delivery / IO errors, so one destination failing never skips the other and a signalling hiccup never re-raises into the DAG.
The mailbox is bound to the event loop the DAG runs on, so the async hooks can
only be awaited from that loop. Steps run on worker threads, off the loop, so
the mid-task progress signals also have synchronous twins — notify_num_batches
and notify_current_batch_id — which bridge back to the loop recorded by
bind_loop. Call those from a step; call the on_* hooks from the loop.
Classes
LifecycleNotifier
class LifecycleNotifier( *, cache: CacheProtocol | None = None, run_id: str | None = None, mailbox: _LifecycleMailbox | None = None,):Signal DAG lifecycle transitions to whichever destinations are wired.
Construct with the destinations available for a run (cache + run_id
closes the cache Run row; mailbox sends proto messages to the initiator).
A LifecycleNotifier() with no arguments is the no-op notifier.
Add a new lifecycle signal by adding one on_* method that touches the
relevant destination(s); add a new destination by adding one optional
constructor argument and branching on it inside the hooks. A signal a step
needs to send from its worker thread also needs a notify_* twin — see the
module docstring.
Arguments
cache: The cache holding theRuntable. Combined withrun_id, enables closing out the run row.run_id: The run UUID returned bymark_run_started.mailbox: A_WorkerMailboxaddressed to the task initiator. Enables protoJOB_ACCEPT/TASK_COMPLETE/TASK_ABORTreplies.
Initialise with whichever destinations are available for the run.
Methods
bind_loop
def bind_loop(self, loop: asyncio.AbstractEventLoop) ‑> None:Record the event loop that lifecycle sends must be delivered on.
The mailbox is bound to the loop the DAG runs on, so a step executing on
a worker thread cannot await it directly — notify_num_batches and
notify_current_batch_id bridge back to this loop instead. Call once
per run, from the loop itself.
Arguments
loop: The event loop the DAG is running on.
notify_current_batch_id
def notify_current_batch_id(self, batch_id: int) ‑> None:Signal progression to a batch, callable from any thread.
The synchronous twin of on_current_batch_id, for steps that run off
the event loop. Requires bind_loop; without it the signal is dropped
with a warning rather than raised.
Arguments
batch_id: The zero-based id of the batch now being processed.
notify_num_batches
def notify_num_batches(self, num_batches: int) ‑> None:Signal the number of batches to expect, callable from any thread.
The synchronous twin of on_num_batches, for steps that run off the
event loop. Requires bind_loop; without it the signal is dropped with
a warning rather than raised.
Arguments
num_batches: The number of batches the initiator should expect.
on_accept
async def on_accept(self) ‑> None:Signal task acceptance (proto JOB_ACCEPT) at DAG start.
Delivery failures are logged as warnings but never re-raised.
on_batches_complete
async def on_batches_complete(self, completion_state: str = 'BATCHES_ONLY') ‑> None:Signal all batches are complete (proto BATCHES_COMPLETE).
Message-only: no cache Run-row write. Delivery failures are logged as warnings but never re-raised.
Arguments
completion_state: The batches-complete body. Defaults to"BATCHES_ONLY"so a stray call never sends the"TASK_COMPLETE"variant — the terminalTASK_COMPLETEis owned byon_success.
on_configuring_task
async def on_configuring_task(self) ‑> None:Signal the worker is configuring the task (proto CONFIGURING_TASK).
Message-only: no cache Run-row write. Delivery failures are logged as warnings but never re-raised.
on_current_batch_id
async def on_current_batch_id(self, batch_id: int) ‑> None:Signal progression to a batch (proto CURRENT_BATCH_ID).
Message-only: no cache Run-row write. Delivery failures are logged as warnings but never re-raised.
Arguments
batch_id: The zero-based id of the batch now being processed.
on_failure
async def on_failure( self, exc: BaseException, *, reason: Reason | None = None, user_readable_message: str | None = None,) ‑> None:Signal failure — mark the run failed and send TASK_ABORT.
Arguments
exc: The exception that caused the DAG to fail.reason: Machine-readableTASK_ABORTreason. Defaults toReason.WORKER_ERRORwhen unset — pass a more specific reason when the failure context is known (e.g.NO_NEW_DATA).user_readable_message: Human-readable abort message. Defaults tostr(exc)when unset.
on_num_batches
async def on_num_batches(self, num_batches: int) ‑> None:Signal the number of batches to expect (proto NUMBER_OF_BATCHES).
Message-only: no cache Run-row write. Delivery failures are logged as warnings but never re-raised.
Arguments
num_batches: The number of batches the initiator should expect.
on_preparing_data
async def on_preparing_data(self) ‑> None:Signal the worker is preparing data (proto PREPARING_DATA).
Message-only: no cache Run-row write. Delivery failures are logged as warnings but never re-raised.
on_preparing_data_batches
async def on_preparing_data_batches(self) ‑> None:Signal the worker is preparing data batches (PREPARING_DATA_BATCHES).
Message-only: no cache Run-row write. Delivery failures are logged as warnings but never re-raised.
on_success
async def on_success(self) ‑> None:Signal successful completion — close the run row, send TASK_COMPLETE.
Delivery of the TASK_COMPLETE reply is logged as a warning on failure
but never re-raised.
on_task_start
async def on_task_start(self) ‑> None:Signal the task is about to run (proto TASK_START).
Message-only: no cache Run-row write. Delivery failures are logged as warnings but never re-raised.