Skip to main content

dag

Telemetry for background DAG wave and sweep outcomes.

Module​

Functions​

report_dag_sweep​

def report_dag_sweep(    *,    dag_name: str,    waves_planned: int,    waves_banked: int,    waves_failed: int,    waves_unstarted: int,    duration_seconds: float,    stop_reason: str | None,) ‑> None:

Send a DAGSweepSummaryEvent; ERROR when every wave failed. Never raises.

Arguments

  • dag_name: The DAG that was run.
  • waves_planned: Waves the run planned.
  • waves_banked: Waves that completed.
  • waves_failed: Waves abandoned.
  • waves_unstarted: Waves left for a later run.
  • duration_seconds: How long the sweep took.
  • stop_reason: Why the loop stopped early, if it did.

report_dag_wave​

def report_dag_wave(    *,    dag_name: str,    outcome: str,    wave_index: int,    attempt: int,    attempts: int,    waves_planned: int,    files: int,    patient_count: int,    wave_kind: str,    duration_seconds: float,    prefect_subflow: bool,    exc: BaseException | None = None,) ‑> None:

Send a DAGWaveEvent for one ended wave attempt. Never raises.

Arguments

  • dag_name: The DAG being run.
  • outcome: "banked", "attempt_failed" or "abandoned".
  • wave_index: The wave's position in its sweep.
  • attempt: This attempt's number.
  • attempts: Attempts the wave is allowed.
  • waves_planned: Waves the sweep planned.
  • files: Files in the wave.
  • patient_count: Distinct patients the wave's files belong to.
  • wave_kind: "first_carry", "refresh" or "orphan".
  • duration_seconds: How long the attempt took.
  • prefect_subflow: Whether the wave ran as its own Prefect subflow run.
  • exc: The failure, for a failed attempt.