Skip to main content

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_PARENT or STARTUP_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_PARENT or STARTUP_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.bump for 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 for spawn to 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.