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.