process_supervision
Spawn, supervise and instrument a compute child process.
Nothing here knows what the child computes: the caller supplies the target, its arguments and its result type. Both the v8 protocol worker and the v9 background DAG are driven through it.
Liveness is a shared multiprocessing.Value holding a monotonic timestamp that
the child bumps and the parent ages, not process liveness alone: a child that is
alive but wedged has to be distinguishable from one that is working.
Every call that blocks — spawn_with_retry, and join on the returned process
— is the caller's to place. An asyncio caller runs them on a worker thread.
Module
Functions
emit_startup_phase
def emit_startup_phase( phase: str, duration_seconds: float, *, process: str, task_id: str, pod_name: str | None = None,) ‑> None:Record how long one phase of child startup took.
Emits two ways: a local INFO log (so it lands in pod logs on disk) and a
WorkerStartupPhaseEvent to the Datadog Logs pipeline.
Never raises - instrumentation must not be able to fail a task. Note that Datadog telemetry is only set up part-way through child setup, so phases completing before that point are reported once it is available.
Arguments
phase: Stable snake_case phase name, e.g."protocol_pickle".duration_seconds: How long the phase took.process:STARTUP_PHASE_PARENTorSTARTUP_PHASE_CHILD.task_id: The task this startup belongs to.pod_name: The pod identifier, where the caller has it.
log_startup_phase
def log_startup_phase( phase: str, *, process: str, task_id: str, pod_name: str | None = None,) ‑> collections.abc.Generator[None, None, None]:Time a child startup phase and report it via emit_startup_phase().
Reports on the way out whether or not the phase raised, so a phase that dies part-way through still tells us how long it burned first.
Arguments
phase: Stable snake_case phase name, e.g."protocol_pickle".process:STARTUP_PHASE_PARENTorSTARTUP_PHASE_CHILD.task_id: The task this startup belongs to.pod_name: The pod identifier, where the caller has it.
poll_result_queue
async def poll_result_queue( result_queue: multiprocessing.Queue[Any], proc: SpawnProcess, *, heartbeat: Heartbeat, stale_timeout: float, cancel_event: threading.Event | None = None,) ‑> Any:Poll result_queue without blocking the event loop.
Runs result_queue.get(timeout=POLL_INTERVAL_SECONDS) in a thread-pool
executor on a tight loop so that the asyncio event loop stays responsive for
cancellation checks and other coroutines.
Arguments
result_queue: The multiprocessing queue the child writes its result to.proc: The child process (used to detect early death).heartbeat: The child's liveness timestamp.stale_timeout: Maximum allowed heartbeat age in seconds before the child is treated as unresponsive.cancel_event: Set when the caller wants to give up waiting. Omit it when nothing can cancel the wait.
Returns Whatever the child put on the queue.
Raises
ChildCancelled: When cancel_event is set before a result arrives.ChildTimeout: When the child heartbeat becomes stale.ChildDiedEmpty: When the child has exited but left the queue empty.
run_heartbeat_refresher
def run_heartbeat_refresher( stop_event: threading.Event, bump: Callable[[], None], interval_seconds: float = 30.0,) ‑> None:Refresh a heartbeat on a fixed interval until the work stops.
Arguments
stop_event: Set by the child when its work finishes.bump: Refreshes the shared timestamp;Heartbeat.bumpfor a child that has one, or any callable for a caller that keeps its own.interval_seconds: How often to refresh.
spawn_with_retry
def spawn_with_retry( ctx: SpawnContext, target: Callable[..., Any], args: Sequence[Any], *, process_name: str,) ‑> multiprocessing.context.SpawnProcess:Spawn a child process with retry and backoff.
Uses daemon=False so the child is properly joined and reaped by the
parent rather than killed at interpreter exit.
Arguments
ctx: The multiprocessing spawn context.target: The child entry point. Must be importable at module level forspawnto pickle it.args: Positional arguments passed to target.process_name: Name for the process, used in logs and in the spawn-error hook.
Returns The started child process.
Raises
WorkerProcessSpawnError: After all retry attempts are exhausted.
Classes
ChildCancelled
class ChildCancelled(*args, **kwargs):Raised by poll_result_queue when the cancel event fires.
ChildDiedEmpty
class ChildDiedEmpty(*args, **kwargs):Raised by poll_result_queue when the child died without putting a result.
ChildTimeout
class ChildTimeout(*args, **kwargs):Raised by poll_result_queue when the child heartbeat goes stale.
Heartbeat
class Heartbeat(ctx: SpawnContext):A child's liveness timestamp, shared with its parent.
The child bumps it; the parent reads its age. Carries the shared value
itself so callers do not handle a bare multiprocessing.Value, and crosses
a spawn boundary as part of the child's arguments.
Methods
age
def age(self) ‑> float:Return seconds since the child last bumped this.
Returns The age in seconds.
bump
def bump(self) ‑> None:Record that the child is still working.