Skip to main content

telemetry

Datadog telemetry for Bitfount.

Split into two pipelines:

  • bitfount.telemetry.logs — structured log events sent to the Datadog Logs API via telemetry_logger.
  • bitfount.telemetry.metrics — custom metrics sent to the Datadog Metrics API via metrics_logger / emit_metric.

The full public surface is re-exported here for convenience.

Module​

Submodules​

Functions​

clear_global_telemetry_context​

def clear_global_telemetry_context() ‑> None:

Empty the process-wide telemetry context.

clear_telemetry_context​

def clear_telemetry_context() ‑> None:

Drop everything set_telemetry_context put in the current context.

The counterpart to setting without a scope: a process that reuses a context across runs - and a test - needs a way back to empty. Values held by an open telemetry_context_scope are restored by that scope as usual.

current_telemetry_context​

def current_telemetry_context() ‑> dict[str, str | int | float | bool]:

Return the context to attach to an event, overlay over global.

Returns A new dict; mutating it does not affect either layer. Empty if nothing has been set, or if reading either layer failed.

emit_gpu_usage_event​

def emit_gpu_usage_event(    sampler: GpuUsageSampler,    *,    step_name: str,    model_ref: str | None = None,    model_class: str | None = None,    task_hash: str | None = None,    run_id: str | None = None,    exc: BaseException | None = None,) ‑> None:

Send a sampler's summary as one GpuUsageEvent.

Nothing is sent for a disabled sampler. Never raises.

Arguments

  • sampler: The sampler to report; stopped first if still running.
  • step_name: The step that was sampled.
  • model_ref: Which model ran.
  • model_class: Class name of the loaded model.
  • task_hash: The step's partition.
  • run_id: The flow run.
  • exc: The exception that ended the step, or None if it succeeded.

emit_metric​

def emit_metric(    metric_name: str,    value: float,    *,    timestamp: int | None = None,    tags: list[str] | None = None,    metric_type: MetricIntakeType | None = None,) ‑> None:

Emit a custom metric point to Datadog via metrics_logger.

Preferred public API for emitting metrics. Builds the log record contract expected by DatadogMetricsHandler (metric name as the message; point timestamp/value and options carried on extra) so call sites do not hand-build it. The timestamp/value are passed via extra rather than as logging args so the record's message never depends on %-formatting.

Arguments

  • metric_name: The metric name. Accepts a MetricName member or a raw string.
  • value: The metric value.
  • timestamp: Unix timestamp (seconds) for the point. Defaults to now.
  • tags: Per-call tags applied to this metric's series.
  • metric_type: The Datadog intake type for this metric. If omitted, the handler defaults the series to GAUGE.

flush_datadog_metrics​

def flush_datadog_metrics() ‑> None:

Flush the Datadog metrics buffer.

Should be called after each task completes, mirroring flush_datadog_telemetry.

flush_datadog_telemetry​

def flush_datadog_telemetry() ‑> None:

Flush the Datadog telemetry buffer.

This should be called to ensure all buffered logs are sent to Datadog.

record_dag_step_time​

def record_dag_step_time(    start: float,    *,    step_name: str,    task_hash: str | None = None,    files_processed: int | None = None,    patients_processed: int | None = None,    exc: BaseException | None = None,) ‑> None:

Emit the execution time of one DAG step, and how it ended.

Arguments

  • start: time.perf_counter() captured before the step began.
  • step_name: The step's name.
  • task_hash: Task hash of the run the step belongs to, when known.
  • files_processed: Files this invocation read, when the step reported it through report_step_progress. None when it reported nothing — a step that cannot answer leaves the field empty rather than having a number guessed for it.
  • patients_processed: Patients this invocation handled, same contract.
  • exc: The exception that ended the step, or None if it returned. Defaults to None so that a caller which does not track the outcome keeps working; callers that can tell should pass it.

record_execution_time​

def record_execution_time(start: float, *, tags: list[str] | None = None) ‑> None:

Emit elapsed time since start as the EXECUTION_TIME_SECONDS gauge.

Shared emission core for execution-time telemetry. Used by track_execution_time (which supplies function:/class: tags) and by call sites that need a dynamic tag the decorator cannot express — e.g. a per-DAG-step step:<name>. Only active when config.settings.enable_execution_time_metrics is True, and never raises: telemetry must not break execution.

Arguments

  • start: time.perf_counter() captured before the timed work began.
  • tags: Tags applied to the emitted gauge series.

report_step_progress​

def report_step_progress(    *, files_processed: int | None = None, patients_processed: int | None = None,) ‑> None:

Record how much work the running step did.

A no-op when no recorder is bound, so a step function stays callable outside a DAG run and in its own unit tests. Only the arguments given are written, so a step can report the two counts from different places.

Arguments

  • files_processed: Files this step actually read this run.
  • patients_processed: Distinct patients this step handled this run.

set_global_telemetry_context​

def set_global_telemetry_context(**fields: Any) ‑> None:

Merge fields into the process-wide telemetry context.

Call this as soon as a value is known - in particular immediately before setup_datadog_telemetry, so that events emitted from the moment the handler exists already carry it.

Arguments

  • **fields: Context entries. A None value is ignored rather than recorded, so a caller can pass an optional value unconditionally.

set_telemetry_context​

def set_telemetry_context(**fields: Any) ‑> None:

Add fields to the telemetry context without a scope to leave.

For a value that describes the whole of the work the current context is doing, set where that work begins - a DAG run naming its dataset. Use telemetry_context_scope instead wherever the value has an end as well as a beginning; this one stays until the context it was set in goes away, or until the next run overwrites it.

Arguments

  • **fields: Context entries. None values are ignored.

setup_datadog_metrics​

def setup_datadog_metrics(    dd_client_token: str | None = None,    dd_site: str | None = None,    hostname: str | None = None,    tags: list[str] | None = None,) ‑> None:

Setup Datadog metrics if credentials are available.

Mirrors setup_datadog_telemetry but routes to the Datadog Metrics API (/api/v2/series) rather than the Logs API. Must be called alongside setup_datadog_telemetry at pod startup.

This function is idempotent - calling it multiple times is safe.

Arguments

  • dd_client_token: The Datadog API key.
  • dd_site: The Datadog site (e.g. 'datadoghq.com', 'datadoghq.eu').
  • hostname: The host name attached to every metric series. Defaults to the system hostname.
  • tags: Static tags applied to every metric series (e.g. ['env:prod']).

Returns None

setup_datadog_telemetry​

def setup_datadog_telemetry(    dd_client_token: str | None = None,    dd_site: str | None = None,    service: str = 'pod',    hostname: str | None = None,    tags: list[str] | None = None,    log_level: str = 'INFO',) ‑> None:

Setup Datadog telemetry logging if credentials are available.

If credentials are not provided, the telemetry logger will exist but have no handlers, meaning all telemetry logs will be silently dropped.

This function is idempotent - calling it multiple times is safe.

Arguments

  • dd_client_token: The Datadog client token.
  • dd_site: The Datadog site to use (e.g., 'datadoghq.com', 'datadoghq.eu').
  • service: The service to use for the logs.
  • hostname: The hostname to use for the logs. Defaults to system hostname.
  • tags: The tags to use for the logs.
  • log_level: The log level for the Datadog handler (e.g., 'INFO', 'DEBUG').

Returns None

shutdown_datadog_metrics​

def shutdown_datadog_metrics() ‑> None:

Shutdown Datadog metrics and flush any pending series.

Should be called at pod shutdown, mirroring shutdown_datadog_telemetry.

shutdown_datadog_telemetry​

def shutdown_datadog_telemetry() ‑> None:

Shutdown Datadog telemetry logging and flush any pending logs.

This should be called during application shutdown to ensure all buffered logs are sent to Datadog.

telemetry_context_scope​

def telemetry_context_scope(**fields: Any) ‑> collections.abc.Iterator[None]:

Add fields to the telemetry context for the duration of the block.

Nested scopes merge, innermost winning, and the overlay is restored on exit whether the block returned or raised.

Arguments

  • **fields: Context entries. None values are ignored.

track_dag_execution_time​

def track_dag_execution_time(func: F) ‑> F:

Decorator emitting a DAG entry point's execution time.

Emits a DAGExecutionTimeEvent naming the function and, for a method, its class. Supports sync and async callables.

Arguments

  • func: The function to wrap.

Returns The wrapped function.

track_execution_time​

def track_execution_time(func: F) ‑> F:

Decorator that emits function execution time as a Datadog gauge metric.

Supports both sync and async callables. Only active when config.settings.enable_execution_time_metrics is True. Emits the MetricName.EXECUTION_TIME_SECONDS gauge with the wrapped callable's class and function name attached as class:<X> / function:<Y> tags.

Arguments

  • func: The function to wrap.

Returns The wrapped function.

Classes​

AlgorithmProgressEvent​

class AlgorithmProgressEvent(**data: Any):

Emitted at algorithm lifecycle boundaries (run/epoch start and end).

Fired automatically by AlgorithmProgressHook for every algorithm via the BaseAlgorithmHook infrastructure — individual algorithm classes do not need any changes.

step is one of "run_start", "run_end", "epoch_start", or "epoch_end".

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static algorithm_name : str
  • static current_epoch : int | None
  • static datetime : str
  • static max_epochs : int | None
  • static model_config
  • static step : str
  • static step_info : str | None
  • static task_context : str

BackgroundDAGRunEvent​

class BackgroundDAGRunEvent(**data: Any):

Emitted at each observable moment of a background DAG deployment run.

A background run is submitted by the pod and executed minutes later by a subprocess the Prefect runner launches, so "how long did the DAG take?" is no longer the whole question — a run can be late because the runner was busy, because the project's concurrency limit was held, or because the DAG itself is long. Three phases separate those, and no other signal does: the Prefect UI shows the states but not the waits between them, and a duration alone cannot say which of the three it was.

Known gap, accepted: a run cancelled while still queued never executes, so it never emits any of these. The start-up sweep is what accounts for those.

Attributes

  • phase: "picked_up" when the subprocess began executing, "started" when the DAG body began, "finished" when the run ended.
  • seconds: Seconds since the run was created, for "picked_up"; seconds spent waiting for the concurrency slot, for "started"; the run's own duration, for "finished".
  • status: How the run ended. Only meaningful on "finished".
  • error_type: The exception class that ended the run, when one did.
  • task_hash: The (pod, datasource) task hash of the run.
  • project_id: The project the run belongs to.
  • datasource_name: The datasource the run processes.
  • trigger: Why the run was started.
  • concurrency_slot: The concurrency limit the run held, when it held one.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static concurrency_slot : str | None
  • static datasource_name : str | None
  • static error_type : str | None
  • static model_config
  • static phase : str
  • static project_id : str | None
  • static seconds : float
  • static task_hash : str | None
  • static trigger : str | None

BackgroundRunEmptyEvent​

class BackgroundRunEmptyEvent(**data: Any):

Emitted when a background DAG run ends without running any step.

Sent when the run's file selection is empty, either before any scope step (reason="no_files_selected") or after one kept none of the files it was given (reason="scope_selected_none"). The run completes rather than fails.

indexing_completed separates the expected case from the one worth an alert: an empty selection before the datasource's first index completes is a cold start, whereas one after it means the inventory matches nothing this task admits.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static dag_name : str | None
  • static datetime : str | None
  • static files_considered : int | None
  • static indexing_completed : bool | None
  • static model_config
  • static project_id : str | None
  • static reason : str | None
  • static run_id : str | None
  • static scope_step : str | None
  • static task_hash : str | None
  • static trigger : str | None

BatchProgressEvent​

class BatchProgressEvent(**data: Any):

Emitted by the orchestrator at the start of each protocol batch.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static batch_number : int
  • static model_config
  • static protocol_name : str

CriteriaMatchingSummaryEvent​

class CriteriaMatchingSummaryEvent(**data: Any):

Emitted once per criteria_matching run, covering every criterion.

One event rather than one per criterion: a run's whole verdict is then readable in a single record, which is what an investigation into an over-strict config actually needs.

The cost of that shape is that criteria and tree_nodes are nested arrays, which Datadog cannot facet or graph on. The numbers worth alerting on are therefore duplicated as flat fields here — the root verdict counts, and the single worst blocking node.

NO PHI. Every field is a count, a configured value, or a schema name. The observed value on a row, the codes a patient carried, the prose failure reason and the subject's identity are all excluded by construction, and no code string is emitted at all — configured or observed. A near-miss is reported as how many such codes exist, never which.

Attributes

  • rows_evaluated: Rows the run evaluated. Rows, not patients: on the file-system path the frame is scan-grain, so one patient's ten scans are ten rows. 0 with empty_reason set means the run had nothing to assess, which is not a verdict — see that field.
  • empty_reason: Why the run evaluated no rows, or None when it evaluated some. no_files_selected is expected and self-healing (a datasource still indexing). files_unreadable means files were selected and none could be loaded, which is unread data and needs someone to look. Both report criteria_yes: 0, and neither means nobody qualified.
  • criteria_yes: Rows where every criterion resolved to pass, the criteria_tree's own combined verdict among them. unknown does not match, so a row missing an input counts here exactly as a row that failed outright does — the two are told apart by rows_unknown on the criterion, not by this.
  • criteria_no: The rest.
  • criterion_count: Flat (non-tree) criteria evaluated, which is the length of criteria before truncation. A range field counts twice, once per bound. criteria_tree leaves are not included; they are counted by tree_leaf_count.
  • tree_node_count: Nodes in the criteria_tree, combinators included.
  • tree_leaf_count: Of those, leaves.
  • tree_depth: The tree's depth, counting a lone leaf as 1.
  • tree_shape: The tree's structure as a signature, e.g. and(or(code,code),not(column)). Config-derived, so a difference between two sites running the same trial is a config drift.
  • tree_shape_truncated: Whether tree_shape was cut short. A tree's breadth is uncapped, so a wide generated one would otherwise render a string large enough to have the whole event truncated.
  • tree_outline: The same tree drawn as an indented outline, each node's counts beside it.

tree_nodes is flat so that every node sits at one attribute depth and a single query reaches all of them; the cost is that reading one run's tree means reassembling it from parent_node_id. This is that reassembly, done once, for a reader rather than a query.

  • tree_outline_truncated: Whether tree_outline was cut short.
  • tree_root_pass: Rows the tree's root resolved to pass.
  • tree_root_fail: Rows it resolved to fail.
  • tree_root_unknown: Rows it could not resolve.
  • worst_blocker_key: The node_key of the node that alone kept the most rows out — the highest rows_sole_blocker in the tree, promoted here so the headline is alertable without unrolling tree_nodes.

Ties break towards the deepest node, so a leaf is named ahead of the AndGroup containing it: both blocked the same rows, but the leaf is the criterion someone would go and loosen. Chosen across every node before tree_nodes is truncated, so a tree over the cap still names its blocker.

None when no node scored, which is the normal state of a healthy config, and also what a tree of OrGroups reports since nothing under one is ever solely responsible (see TreeNodeSummary.rows_sole_blocker).

  • worst_blocker_node_id: That node's node_id.
  • worst_blocker_evidence_id: That node's evidence_id, or None where it set none.
  • worst_blocker_group: That node's resolved group, or None.
  • worst_blocker_rows: That node's rows_sole_blocker. 0 exactly when worst_blocker_key is None.
  • dead_code_list_count: Inclusion criteria where every configured code matched nobody — the config asked for codes this cohort does not carry. Exclusion criteria are excluded from this count: one that excludes nobody is working, not misconfigured.
  • near_miss_criteria_count: Code-list criteria with at least one near-miss code in the cohort — how many lists have a neighbouring code to look at, not how many codes there are.
  • criteria_truncated: Criteria omitted from criteria to bound the payload; 0 normally.
  • tree_nodes_truncated: Nodes omitted from tree_nodes for the same reason. Nodes are kept root-first, so what drops is the deepest detail rather than the structure above it.
  • criteria: Per-criterion aggregates for the flat criteria — one entry per comparison, so a range field appears twice.
  • tree_nodes: Per-node aggregates, in post-order (children before parents), so a reader meets a group's inputs before the group.

Flat rather than nested on purpose: every node sits at the same attribute depth, so one query reaches all of them across runs. tree_outline carries the shape for reading.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static criteria_no : int
  • static criteria_truncated : int
  • static criteria_yes : int
  • static criterion_count : int
  • static dead_code_list_count : int
  • static empty_reason : str | None
  • static model_config
  • static near_miss_criteria_count : int
  • static rows_evaluated : int
  • static tree_depth : int
  • static tree_leaf_count : int
  • static tree_node_count : int
  • static tree_nodes_truncated : int
  • static tree_outline : str | None
  • static tree_outline_truncated : bool
  • static tree_root_fail : int
  • static tree_root_pass : int
  • static tree_root_unknown : int
  • static tree_shape : str | None
  • static tree_shape_truncated : bool
  • static worst_blocker_evidence_id : str | None
  • static worst_blocker_group : str | None
  • static worst_blocker_key : str | None
  • static worst_blocker_node_id : str | None
  • static worst_blocker_rows : int

CriterionSummary​

class CriterionSummary(**data: Any):

One criterion's aggregate outcome across a criteria_matching run.

Counts only. The observed value a row carried, the codes it matched and the prose reason it failed are all deliberately absent — see CriteriaMatchingSummaryEvent for why.

Attributes

  • output_field: The criterion's logical field name, as the served evidence keys it. Not unique on its own: a range field contributes one entry per bound, both under this name.
  • column_name: The flat column name the filter writes. For a ColumnFilter this embeds the comparison — "Age (yrs) >= 50" — which is what keeps a range's two bounds apart as separate entries.
  • operator: The comparison operator as the filter was built with it. Either a symbol (">=") or the English spelling a task config may use ("greater than or equal"); both appear. A code-list criterion reports "code_match", which is not an ordering comparison and so never produces rows_failed_near_bound. None for a filter with no named comparison at all.
  • operand: The configured right-hand side of a scalar comparison. None for a code-list criterion, whose operand is the whole configured list: configured_code_count reports its length instead, so a few hundred codes are not repeated into every event.
  • grain: scan or patient — the cardinality axis.
  • provenance: scan, ehr, supplied or unknown — the source axis.
  • family: The semantic family (code_list, scan_metric, ...), or None when the producing filter declared none.
  • rows_pass: Rows the criterion passed. Rows, not patients: on the file-system path the frame is scan-grain, so a patient with ten scans counts ten times. The same holds for every rows_* field here.
  • rows_fail: Rows it failed — evaluated, and the answer was no.
  • rows_unknown: Rows it could not evaluate. For a code-list criterion this means the patient had no code data at all, which is missing input rather than a verdict — read it apart from rows_fail.
  • rows_absent_column: Rows whose target column was absent from the frame entirely, as opposed to present and null.

Only a ColumnFilter sets this. A code-list criterion is a MethodFilter and never does, so 0 there does not mean its column was present — check configured_codes_matching_nobody being None, which is how an absent code column shows up.

  • rows_failed_near_bound: Rows that failed an ordering comparison but came within 10% of passing it: abs(value - operand) <= 0.1 * abs(operand). With an operand of 0, where that fraction has no meaning, only an exact 0 counts.

Counts failures only — a row that could not be evaluated is not near anything. None for a criterion with no ordering comparison, which is different from 0 (an ordering comparison where nothing landed close).

A high count against a low rows_pass is the signature of a threshold set slightly too tight: rows_fail: 332 with rows_failed_near_bound: 301 says almost everyone who failed failed narrowly.

  • configured_code_count: How many codes the criterion's configured list holds. Patterns count as one each, so a single H35.3* covering twelve real codes counts once. Includes the _study_eye counterpart list where the criterion has one, since a code configured in either was genuinely configured. None for a non-code criterion.
  • configured_codes_matching_nobody: How many of those matched no code the cohort carried. Equal to configured_code_count means no patient carried any configured code — with the filters' OR semantics, that is a dead list rather than a demanding one.

Read it with the criterion's direction in mind. On an inclusion criterion it is the too-strict signal. On an exclusion criterion it means nobody was excludable, which is normal and is why dead_code_list_count counts inclusion criteria only.

None means it could not be established, never that nothing was

  • dead: the criterion's code column was absent from the frame (the data was never read), or the cohort scan did not run because enable_criteria_near_miss_counts is off, which leaves an exclusion criterion with no source for it at all.
  • near_miss_code_count: Distinct codes present in the cohort that matched no configured pattern but share an ICD-10 category (the first three characters) with one.

None when the pass did not run, and also on a procedure

  • criterion: those columns mix CPT/HCPCS with ICD-10-PCS, where a three-character prefix groups nothing, so no count is reported rather than an arbitrary one.
  • near_miss_row_count: Rows carrying at least one such code. Rows, not
  • patients: on the file-system path the frame is scan-grain, so a patient with ten scans counts ten times.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static column_name : str
  • static configured_code_count : int | None
  • static configured_codes_matching_nobody : int | None
  • static family : str | None
  • static grain : str
  • static model_config
  • static near_miss_code_count : int | None
  • static near_miss_row_count : int | None
  • static operand : Any
  • static operator : str | None
  • static output_field : str
  • static provenance : str
  • static rows_absent_column : int
  • static rows_fail : int
  • static rows_failed_near_bound : int | None
  • static rows_pass : int
  • static rows_unknown : int

DAGExecutionTimeEvent​

class DAGExecutionTimeEvent(**data: Any):

Emitted with the wall-clock time of a DAG step or executor entry point.

Replaces the execution_time_s Datadog custom metric for the DAG paths. Exactly one of step_name or function identifies what was timed: steps carry step_name, the executor's own entry points carry function (and class_name when the callable is a method).

status says how the timed work ended - "success", "error" or "cancelled" - because the timing is emitted whether it succeeded or not, and a duration alone cannot tell a step that finished from one that died partway through. error_type names the exception class for the two non-success outcomes, so a dashboard can separate a timeout from a crash without a second query. It defaults to success rather than being required, because an event must never raise when constructed. It serialises as execution_status: Datadog reads a status attribute as the log's severity, which the handler sets from the record's level, and two producers for one attribute means one of them loses.

The volume fields divide by who can answer. A step timing carries files_processed / patients_processed - work done by that one invocation, reported by the step itself. The executor's own entry-point timings carry files_in_scope instead, the run's denominator: summing the per-step counts would double-count, since every file-scoped step sees the same files.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static class_name : str | None
  • static duration_seconds : float
  • static error_type : str | None
  • static files_in_scope : int | None
  • static files_processed : int | None
  • static function : str | None
  • static model_config
  • static patients_processed : int | None
  • static step_name : str | None
  • static task_hash : str | None

DAGSweepSummaryEvent​

class DAGSweepSummaryEvent(**data: Any):

Emitted once when a waved background DAG run finishes its sweep.

Attributes

  • dag_name: The DAG being run.
  • waves_planned: Waves the run planned.
  • waves_banked: Waves that completed.
  • waves_failed: Waves attempted and abandoned.
  • waves_unstarted: Waves left for a later run.
  • stop_reason: Why the loop stopped early (budget, standing aside), or "no_waves" when nothing was planned; None when it ran every planned wave.
  • duration_seconds: How long the sweep took, from its first wave to its last.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static dag_name : str
  • static duration_seconds : float
  • static model_config
  • static stop_reason : str | None
  • static waves_banked : int
  • static waves_failed : int
  • static waves_planned : int
  • static waves_unstarted : int

DAGWaveEvent​

class DAGWaveEvent(**data: Any):

Emitted when one attempt at a background DAG wave ends.

outcome is "banked" for a wave that completed, "attempt_failed" for a failed attempt that will be retried, and "abandoned" once the retries are spent (its files stay outstanding for the next run).

Attributes

  • dag_name: The DAG being run.
  • outcome: How the attempt ended.
  • wave_index: The wave's position in its sweep, from 1.
  • attempt: This attempt's number, from 1.
  • attempts: How many attempts the wave is allowed.
  • waves_planned: Waves the sweep planned, so wave_index reads as n of N.
  • 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; False when creating one failed and the wave ran without it.
  • error_type: The exception class, for a failed attempt.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static attempt : int
  • static attempts : int
  • static dag_name : str
  • static duration_seconds : float
  • static error_type : str | None
  • static files : int
  • static model_config
  • static outcome : str
  • static patient_count : int
  • static prefect_subflow : bool
  • static wave_index : int
  • static wave_kind : str
  • static waves_planned : int

DatadogLogsHandler​

class DatadogLogsHandler(    api_instance: LogsApi,    source: str,    hostname: str,    service: str,    tags: list[str] | None = None,    capacity: int = 1000,    flush_interval_seconds: float = 0.0,):

A MemoryHandler that sends logs to Datadog.

This handler automatically flushes when the buffer approaches a 5MB limit. These numbers are taken from Datadog's Logs API documentation: https://docs.datadoghq.com/api/latest/logs/

We piggyback off Python's MemoryHandler class since we do not want to implement our own buffer management system, including what happens when there are records left in the buffer when the handler is closed.

Initialize the DatadogMemoryHandler.

Arguments

  • api_instance: The Datadog API instance.
  • source: The source of the logs.
  • hostname: The hostname of the logs.
  • service: The service of the logs.
  • tags: The tags of the logs.
  • capacity: The capacity of the buffer, defaulted to 1000 according to Datadog's documentation.
  • flush_interval_seconds: How often to flush the buffer regardless of how full it is. Defaults to no periodic flush; the production path is setup_datadog_telemetry, which passes the configured interval.

Variables​

  • static BUFFER_PERCENTAGE
  • static MAX_BUFFER_SIZE

Methods​


close​

def close(self) ‑> None:

Stop the periodic flush, then flush and close the handler.

The flusher is signalled but not waited for: logging.shutdown calls this while holding the handler, and a flush already in flight cannot finish until that caller returns. shutdown_datadog_* stops the flusher properly first, so the orderly path loses nothing.

emit​

def emit(self, record: logging.LogRecord) ‑> None:

Emit a record to the buffer.

Note that we immediately build the HTTPLogItem object and add it to a separate buffer, rather than waiting for the flush operation to do so.

Arguments

  • record: The log record.

flush​

def flush(self) ‑> None:

Flush the buffer by sending all records to Datadog in a single request.

Override the parent's flush method to flush records instead of calling emit() per record.

shouldFlush​

def shouldFlush(self, record: logging.LogRecord, item_size: int = 0) ‑> bool:

Determine if we should flush the buffer.

Arguments

  • record: The log record.
  • item_size: The size of the item to add to the buffer.

Returns True if we should flush the buffer, False otherwise.

DatadogMetricsHandler​

class DatadogMetricsHandler(    api_instance: MetricsApi,    hostname: str,    tags: list[str] | None = None,    flush_interval_seconds: float = 0.0,):

A logging handler that sends custom metrics to Datadog's Metrics API.

Buffers points per unique (metric_name, tags) combination, so each combination becomes its own MetricSeries. Each buffer is monitored independently and flushed when it approaches 90% of the 5 MB uncompressed limit, and again on flush()/close().

Call sites do not use this handler directly; they go through emit_metric, which forwards to metrics_logger:

metrics_logger.info(metric_name, extra={...})

where metric_name maps to record.msg and the point's timestamp/value and options are carried as record attributes via extra:

  • metric_timestamp / metric_value: the point. Carried on extra (not as logging args) so the message never depends on %-formatting.
  • tags: per-call tags appended to this series' tags.
  • metric_type: the Datadog intake type for the series. It is set per call site; calls that omit it default the series to GAUGE.

See: https://docs.datadoghq.com/api/latest/metrics/

Initialise the handler.

Arguments

  • api_instance: A configured Datadog MetricsApi instance.
  • hostname: Attached to every MetricSeries via resources so metrics are filterable by host in Datadog.
  • tags: Static tags applied to every series (e.g. ['env:prod']).
  • flush_interval_seconds: How often to flush every series buffer regardless of how full it is. Defaults to no periodic flush; the production path is setup_datadog_metrics, which passes the configured interval.

Variables​

  • static MAX_SERIES_SIZE

Methods​


close​

def close(self) ‑> None:

Stop the periodic flush, then flush and close the handler.

The flusher is signalled but not waited for: logging.shutdown calls this while holding the handler, and a flush already in flight cannot finish until that caller returns. shutdown_datadog_* stops the flusher properly first, so the orderly path loses nothing.

emit​

def emit(self, record: logging.LogRecord) ‑> None:

Buffer a metric point.

Expects record.msg to be the metric name, with the point's timestamp and value carried as the metric_timestamp / metric_value attributes (set via extra=...). Per-call tags can be supplied via extra={"tags": [...]}. Records that do not match this shape are silently dropped. The timestamp/value are deliberately not passed as logging args so the record's message never depends on %-formatting and cannot break a foreign formatter.

Each unique (metric_name, tags) combination is buffered independently and flushed as its own MetricSeries when it approaches the size limit.

Arguments

  • record: The log record.

flush​

def flush(self) ‑> None:

Flush all series buffers.

DatasetConnectedEvent​

class DatasetConnectedEvent(**data: Any):

Emitted once per datasource when a pod successfully connects it.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static connection_datetime : str
  • static dataset_name : str
  • static datasource_type : str
  • static model_config

EHRQueryCompleteEvent​

class EHRQueryCompleteEvent(**data: Any):

Emitted after an EHR patient query batch completes.

Fires once per protocol batch in EHR-enabled algorithms. Provides visibility into the per-batch EHR query phase, which is otherwise the longest silent stretch during task execution.

The patients_* outcome counts are per input row, like the per-patient loop that produces them; patients_queried and records_found count unique patients.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static config_type : str
  • static model_config
  • static patients_auth_failed : int
  • static patients_errored : int
  • static patients_not_found : int
  • static patients_queried : int
  • static records_found : int

EHRQueryOutcomeEvent​

class EHRQueryOutcomeEvent(**data: Any):

Emitted once per background ehr_criteria_query step run.

run_outcome is the step's EHRRunOutcome value: "primed", "degraded_partial", "degraded_total" or "not_configured".

Attributes

  • run_outcome: How well the run primed the EHR cache.
  • ehr_availability: What the preflight probe found, when one ran.
  • records_stored: EHR rows this run persisted.
  • skipped_total_failure: Patients served from an older row.
  • skipped_unidentified: Patients with no identity to query by.
  • skipped_fresh: Patients whose stored row was still fresh.
  • task_hash: The step's partition.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static ehr_availability : str | None
  • static model_config
  • static records_stored : int
  • static run_outcome : str
  • static skipped_fresh : int
  • static skipped_total_failure : int
  • static skipped_unidentified : int
  • static task_hash : str | None

EHRScreeningBatchEvent​

class EHRScreeningBatchEvent(**data: Any):

Emitted after each EHR screening page is queried and eligibility-matched.

Fires once per page, so a screening run sends one per page it queries.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static model_config
  • static page_eligible : int
  • static page_number : int
  • static page_size : int
  • static project_id : str | None
  • static task_id : str
  • static total_eligible : int

EHRSessionInitialisedEvent​

class EHRSessionInitialisedEvent(**data: Any):

Emitted after a NextGen or FHIR R4 EHR session is set up.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static config_type : str
  • static model_config

EHRTokenRefreshedEvent​

class EHRTokenRefreshedEvent(**data: Any):

Emitted after a FHIR client access token is refreshed.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static config_type : str
  • static model_config

ExecutionStatus​

class ExecutionStatus(*args, **kwds):

How a timed piece of DAG work ended.

CANCELLED is kept apart from ERROR because the pod cancels in-flight work on shutdown and redeploy, and counting those as failures would bury the real ones.

Inherits from str so values serialise directly as JSON strings.

Variables​

  • static CANCELLED
  • static ERROR
  • static SUCCESS

FileMultiSeriesSummaryEvent​

class FileMultiSeriesSummaryEvent(**data: Any):

Aggregated count of files whose matching series were cut to one row.

Replaces one event per reduced file, and is sent at the same points as FileSkipSummaryEvent.

Attributes

  • datasource_type: Class name of the datasource that reduced the files.
  • dataset_name: The datasource's configured name.
  • trigger: What closed the window: "pass_end", "task_end" or "interval".
  • window_seconds: Length of the window the counts cover.
  • files_reduced: Files that produced more than one matching row.
  • rows_dropped: Rows discarded across those files.
  • by_laterality: Files reduced, keyed by configured laterality.
  • by_series_protocol: Files reduced, keyed by configured series protocol.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static by_laterality : dict[str, int]
  • static by_series_protocol : dict[str, int]
  • static dataset_name : str | None
  • static datasource_type : str
  • static files_reduced : int
  • static model_config
  • static rows_dropped : int
  • static trigger : str
  • static window_seconds : float

FileSkipSummaryEvent​

class FileSkipSummaryEvent(**data: Any):

Aggregated count of files a datasource skipped over one window.

Replaces one event per skipped file. Emitted at the end of each file-listing pass, at each task boundary, and periodically during a long pass, so a window covers whatever was skipped since the previous summary.

Attributes

  • datasource_type: Class name of the datasource that skipped the files.
  • dataset_name: The datasource's configured name.
  • trigger: What closed the window: "pass_end", "task_end" or "interval".
  • window_seconds: Length of the window the counts cover.
  • total_skipped: Files skipped in the window.
  • by_reason: Counts keyed by FileSkipReason name.
  • by_source: Counts keyed by skip source: "datasource" for skips persisted to the cache, "task" for task-scoped skips.
  • by_extension: Counts keyed by file extension ("" when none).
  • sample_files: Up to a few redacted file names per reason.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static by_extension : dict[str, int]
  • static by_reason : dict[str, int]
  • static by_source : dict[str, int]
  • static dataset_name : str | None
  • static datasource_type : str
  • static model_config
  • static sample_files : dict[str, list[str]]
  • static total_skipped : int
  • static trigger : str
  • static window_seconds : float

FilteringCompleteEvent​

class FilteringCompleteEvent(**data: Any):

Emitted after RecordFilterAlgorithm.setup_run() finishes selecting files.

Fires once per task, before the first batch begins. Provides visibility into how many files were considered and how many passed the filter criteria.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static files_selected : int
  • static model_config
  • static total_files : int

GpuUsageEvent​

class GpuUsageEvent(**data: Any):

Emitted once per model_inference invocation that ran inference.

Summarises the GPU over the step's working window, from the model load to the end of the chunk loop: whether one is present, whether torch put work on it, and the mean and maximum of its compute and memory utilisation. One event per invocation rather than per sample, so the volume does not grow with how long inference takes.

Utilisation and device memory come from NVML and describe the whole device, not this process. The torch_* fields say what torch itself did. A fully-cached invocation does no GPU work and emits nothing.

Attributes

  • step_name: The step that was sampled.
  • model_ref: Which model ran.
  • model_class: Class name of the loaded model, which separates a torch model from one on another runtime (e.g. ONNX) when NVML shows load but torch did not use the GPU.
  • task_hash: The step's partition.
  • run_id: The flow run.
  • duration_seconds: Length of the sampled window.
  • status: How the step ended.
  • error_type: Exception class name when the step did not succeed.
  • gpu_present: Whether a GPU was found. None when it could not be determined, e.g. no NVIDIA driver library.
  • gpu_count: GPUs NVML reports.
  • gpu_name: Name of the sampled device.
  • gpu_memory_total_bytes: Total memory of the sampled device.
  • nvml_driver_version: NVIDIA driver version.
  • gpu_probe_error: Exception class name when probing for a GPU failed.
  • torch_cuda_build: CUDA version torch was built with; None for a CPU-only build.
  • torch_cuda_available: torch.cuda.is_available().
  • torch_device: Device the SDK selects for torch models: cuda, mps or cpu.
  • torch_used_gpu: Whether torch allocated GPU memory during the window.
  • torch_gpu_memory_peak_bytes: Peak GPU memory torch allocated during the window; on MPS, the driver's allocation at the window's end.
  • gpu_utilisation_mean_pct: Mean device compute utilisation.
  • gpu_utilisation_max_pct: Maximum device compute utilisation.
  • gpu_memory_used_mean_bytes: Mean device memory in use.
  • gpu_memory_used_max_bytes: Maximum device memory in use.
  • sample_count: Samples the utilisation figures are built from.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static duration_seconds : float | None
  • static error_type : str | None
  • static gpu_count : int | None
  • static gpu_memory_total_bytes : int | None
  • static gpu_memory_used_max_bytes : int | None
  • static gpu_memory_used_mean_bytes : float | None
  • static gpu_name : str | None
  • static gpu_present : bool | None
  • static gpu_probe_error : str | None
  • static gpu_utilisation_max_pct : float | None
  • static gpu_utilisation_mean_pct : float | None
  • static model_class : str | None
  • static model_config
  • static model_ref : str | None
  • static nvml_driver_version : str | None
  • static run_id : str | None
  • static sample_count : int
  • static step_name : str
  • static task_hash : str | None
  • static torch_cuda_available : bool | None
  • static torch_cuda_build : str | None
  • static torch_device : str | None
  • static torch_gpu_memory_peak_bytes : int | None
  • static torch_used_gpu : bool | None

GpuUsageSampler​

class GpuUsageSampler(interval_seconds: float = 1.0):

Samples GPU utilisation over a window of work.

Use as a context manager, or call start and stop directly when the window does not match a block; both are idempotent. Never raises: every probe and sample failure is logged at debug level and leaves the affected fields empty.

A no-op when settings.enable_gpu_telemetry is False.

Arguments

  • interval_seconds: Time between samples.

Variables​

  • duration_seconds : float | None - Length of the sampled window, once stopped.
  • enabled : bool - Whether this sampler is collecting, as decided at start.

Methods​


start​

def start(self) ‑> None:

Probe for a GPU and begin sampling.

stop​

def stop(self) ‑> None:

Stop sampling and read torch's own figures for the window.

summary​

def summary(self) ‑> GpuUsageSummary:

Summarise what was observed so far.

Returns The summary; all fields empty when the sampler is disabled.

GpuUsageSummary​

class GpuUsageSummary(    gpu_present: bool | None = None,    gpu_count: int | None = None,    gpu_name: str | None = None,    gpu_memory_total_bytes: int | None = None,    nvml_driver_version: str | None = None,    gpu_probe_error: str | None = None,    torch_cuda_build: str | None = None,    torch_cuda_available: bool | None = None,    torch_device: str | None = None,    torch_used_gpu: bool | None = None,    torch_gpu_memory_peak_bytes: int | None = None,    gpu_utilisation_mean_pct: float | None = None,    gpu_utilisation_max_pct: float | None = None,    gpu_memory_used_mean_bytes: float | None = None,    gpu_memory_used_max_bytes: int | None = None,    sample_count: int = 0,):

What a GpuUsageSampler observed over its window.

Field meanings match the same-named fields on GpuUsageEvent. A field the sampler could not determine is None.

Variables​

  • static gpu_count : int | None
  • static gpu_memory_total_bytes : int | None
  • static gpu_memory_used_max_bytes : int | None
  • static gpu_memory_used_mean_bytes : float | None
  • static gpu_name : str | None
  • static gpu_present : bool | None
  • static gpu_probe_error : str | None
  • static gpu_utilisation_max_pct : float | None
  • static gpu_utilisation_mean_pct : float | None
  • static nvml_driver_version : str | None
  • static sample_count : int
  • static torch_cuda_available : bool | None
  • static torch_cuda_build : str | None
  • static torch_device : str | None
  • static torch_gpu_memory_peak_bytes : int | None
  • static torch_used_gpu : bool | None

GroupingSummaryEvent​

class GroupingSummaryEvent(**data: Any):

Aggregate Datadog event for grouping-aware batching behaviour.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static batch_size : int | None
  • static cached_files : int | None
  • static group_count : int | None
  • static include_non_new_group_files : bool | None
  • static maybe_final : bool | None
  • static model_config
  • static oversized_cohort_size : int | None
  • static selected_files : int | None
  • static step : str

IndexingProgressEvent​

class IndexingProgressEvent(**data: Any):

Emitted by a metadata runtime each time it persists a chunk.

Totals are cumulative for the run, so consecutive events give both the rate and the running position of a pass that may take hours.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static bytes_read : int
  • static dataset_identifier : str | None
  • static elapsed_seconds : float
  • static files_keyed : int
  • static files_per_second : float
  • static files_reused : int
  • static files_seen : int
  • static files_unchanged : int
  • static io_concurrency : int | None
  • static model_config
  • static process_read_bytes : int | None
  • static rows_written : int
  • static run_type : str
  • static task_hash : str

IndexingSummaryEvent​

class IndexingSummaryEvent(**data: Any):

Emitted once when a metadata runtime finishes, successfully or not.

status mirrors the run's bookkeeping status ("success", "skipped", "partial" or "error"). It serialises as run_status, because Datadog reads a status attribute as the log's severity. The per-phase seconds are what separate a pass that is bound by hashing from one bound by header parsing or by cache writes.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static bytes_read : int
  • static dataset_identifier : str | None
  • static elapsed_seconds : float
  • static files_keyed : int
  • static files_per_second : float
  • static files_reused : int
  • static files_seen : int
  • static files_unchanged : int
  • static flush_seconds : float
  • static io_concurrency : int | None
  • static key_seconds : float
  • static model_config
  • static parse_seconds : float
  • static process_read_bytes : int | None
  • static root_probe_seconds : float
  • static root_probes : int
  • static rows_written : int
  • static run_type : str
  • static status : str
  • static task_hash : str
  • static transient_failures : int

MetadataKickoffEvent​

class MetadataKickoffEvent(**data: Any):

Emitted once each time the pod decides which datasources to index.

submission_order is the order the file_metadata_runtime runs are scheduled in. The health lists partition every datasource considered; skipped_archived_only names those left out because every project they are linked to is archived. hub_filters_failed names the Hub calls that failed and were therefore not applied.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static healthy : list[str]
  • static hub_filters_applied : list[str]
  • static hub_filters_failed : list[str]
  • static model_config
  • static pod_name : str
  • static skipped_archived_only : list[str]
  • static submission_order : list[str]
  • static trigger : str
  • static unknown : list[str]

MetricName​

class MetricName(*args, **kwds):

Canonical names for all Datadog custom metrics.

Inherits from str so values serialise directly as metric names. The intake type for each metric is set at the call site (see emit_metric), defaulting to GAUGE when omitted.

Variables​

  • static EHR_REQUEST_COUNT
  • static EHR_RESPONSE_BODY_SIZE_BYTES
  • static EHR_RESPONSE_TIME_SECONDS
  • static EXECUTION_TIME_SECONDS
  • static TIMING_WINDOW_MAX_SECONDS
  • static TIMING_WINDOW_MEAN_SECONDS
  • static TIMING_WINDOW_SAMPLES
  • static TIMING_WINDOW_SECONDS_TOTAL

ModelInferenceProgressEvent​

class ModelInferenceProgressEvent(**data: Any):

Emitted while model_inference is still working.

Modelled on IndexingProgressEvent, and for the same reason: a pass that takes hours otherwise says nothing between its first log line and its last, which is indistinguishable from a hang. Consecutive events give both the rate and the running position.

Emitted on a wall clock and a file floor rather than per batch, so the volume is bounded both by how long inference takes and by how many files it covers — the pair that keeps a large datasource from paying a point per file even when one file takes longer to infer than the interval. One trailing event closes out each invocation.

The counts are scoped to this invocation of the step, which under waving is one wave rather than the whole datasource. The wave's own banner in bitfount.flows.dag.executor carries the outer denominator.

Attributes

  • model_ref: Which model is running.
  • batches_done: Batches finished so far, of batches_total.
  • batches_total: Batches this invocation planned.
  • files_in_window: Files covered since the previous report, which at the one-file batches every shipped v9 template configures is how many batches the window held.
  • files_done: Files carried so far, of files_total.
  • files_total: Files this invocation is inferring over, after the cached ones were subtracted.
  • elapsed_seconds: Since the step's loop started.
  • files_per_second: Mean rate so far, for projecting what is left.
  • task_hash: The step's partition.
  • run_id: The flow run.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static batches_done : int
  • static batches_total : int
  • static elapsed_seconds : float
  • static files_done : int
  • static files_in_window : int
  • static files_per_second : float
  • static files_total : int
  • static model_config
  • static model_ref : str
  • static run_id : str | None
  • static task_hash : str | None

PatientEligibilityEvent​

class PatientEligibilityEvent(**data: Any):

Emitted after eligible patient counts are calculated.

Two producers, whose fields do not overlap fully, so every field either producer cannot supply is optional — a required field would raise ValidationError in the other, inside a try that also guards unrelated work:

  • the orchestrator's on_task_end hook, counting eligible patients over a task's merged output; it has the datasource and no denominator.
  • the patient_eligibility reduce step, which has the denominator, task_hash and run_id, and reads stored verdicts rather than a datasource.

From the reduce step, patients_considered counts the patients rolled up in that run's partition — the pipeline is scan-driven, so that is the patients with at least one stored scan, not every patient in the datasource — and eligible_count is how many of those came out ELIGIBLE. Both describe the rows the run wrote, not what a reader is served: the pointer cutover can be withheld, leaving the previous partition published.

ineligible_count and unknown_count complete the verdict breakdown from the reduce step, and are None from the orchestrator hook, which counts only the eligible. Read them together: the three move independently, and a fall in eligible_count means something different depending on which of the other two rose.

unknown_count is not an EHR-outage signal on its own. An unreachable EHR leaves its criteria unevaluated and those patients roll up to unknown, but so does a patient with no patient-level row at all. The run-scoped EHRRunOutcome that would tell the two apart does not reach the step that emits this; the ehr_coverage tag published beside these counts is the discriminator.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static datasource : str | None
  • static datetime : str
  • static eligible_count : int
  • static ineligible_count : int | None
  • static model_config
  • static patients_considered : int | None
  • static project_id : str
  • static run_id : str | None
  • static task_hash : str | None
  • static unknown_count : int | None

PodErrorEvent​

class PodErrorEvent(**data: Any):

Emitted when the pod fails outside a task.

phase is "init" or "startup" for an exception escaping Pod.__init__ or Pod.start, and "process_spawn" when a background process could not be started after all retries. Frames come from the record's exc_info where the exception is in hand; the message is never sent.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static attempts : int | None
  • static error_type : str
  • static model_config
  • static phase : str
  • static process_name : str | None

RecoverySummaryEvent​

class RecoverySummaryEvent(**data: Any):

Emitted once per start-of-process metadata run reconciliation.

Attributes

  • backlog_cancelled: Queued duplicate metadata runs cancelled.
  • runs_resubmitted: Replacement runs submitted for lost ones.
  • refresh_runs_ended: Refresh runs ended after wedging.
  • slots_released: Stranded concurrency slots released.
  • failed_steps: Reconciliation steps that raised, by name.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static backlog_cancelled : int
  • static failed_steps : list[str]
  • static model_config
  • static refresh_runs_ended : int
  • static runs_resubmitted : int
  • static slots_released : int

ScanReduceEvent​

class ScanReduceEvent(**data: Any):

Emitted after the scan_reduce step narrows a run's file selection.

How much a reduction saved is otherwise invisible: covered and the published coverage figure both count the inventory, so a run that skipped most of it looks the same from the outside as one that processed everything.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static datetime : str | None
  • static files_considered : int | None
  • static files_selected : int | None
  • static files_unidentified : int | None
  • static model_config
  • static patients_considered : int | None
  • static patients_dropped : int | None
  • static run_id : str | None
  • static strategy : str | None
  • static task_hash : str | None
  • static unsatisfiable : bool | None

SchemaGenerationEndEvent​

class SchemaGenerationEndEvent(**data: Any):

Emitted when the Prefect schema worker finishes generating a schema.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static dataset_id : str
  • static datetime : str
  • static model_config

SchemaGenerationErrorEvent​

class SchemaGenerationErrorEvent(**data: Any):

Emitted when the Prefect schema worker fails.

stage is "generation" for a failure while building the schema, "upload" for a failed full upload, and "partial_upload" for a failed upload after an interruption.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static dataset_id : str
  • static error_type : str
  • static model_config
  • static stage : str

SchemaGenerationStartEvent​

class SchemaGenerationStartEvent(**data: Any):

Emitted when the Prefect schema worker begins generating a schema.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static dataset_id : str
  • static datetime : str
  • static model_config

SchemaInitiationSource​

class SchemaInitiationSource(*args, **kwds):

Source that triggered schema generation.

Inherits from str so values serialise directly as JSON strings.

Variables​

  • static PREFECT_SCHEMA_WORKER

SchemaUploadSuccessEvent​

class SchemaUploadSuccessEvent(**data: Any):

Emitted after a full schema is successfully uploaded to the Hub.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static dataset_id : str
  • static datetime : str
  • static model_config
  • static number_of_records : int

StepProgress​

class StepProgress(    files_processed: int | None = None, patients_processed: int | None = None,):

How much work one DAG step invocation did.

Mutable by design. contextvars.copy_context() copies bindings, not the objects they point at, so a step mutating this on its worker thread is seen by the executor that created it - which is how a count reported inside a step reaches the timing event emitted around it.

Both fields mean work done by this invocation and nothing else. A step that cannot answer honestly leaves them None; a partition size or a denominator does not belong here. files_in_scope on the executor's own timing events is the field for a denominator.

Attributes

  • files_processed: Files this step actually read this run, excluding ones skipped because a cached result already covered them.
  • patients_processed: Distinct patients this step handled this run.

Variables​

  • static files_processed : int | None
  • static patients_processed : int | None

TaskAcceptedByPodEvent​

class TaskAcceptedByPodEvent(**data: Any):

Emitted by the modeller when a pod accepts a task request.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static app_version : str
  • static datetime : str
  • static model_config
  • static pod_identifier : str
  • static task_id : str

TaskAcceptedEvent​

class TaskAcceptedEvent(**data: Any):

Emitted by the worker when it accepts and begins a task.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static app_version : str
  • static datasource_type : str
  • static datetime : str
  • static model_config
  • static modeller_username : str
  • static pod_name : str
  • static task_id : str

TaskCompleteEvent​

class TaskCompleteEvent(**data: Any):

Emitted by the worker when a task finishes successfully.

Symmetric counterpart to TaskAcceptedEvent — together they bracket the full task execution window in Datadog logs.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static datetime : str
  • static model_config
  • static modeller_username : str
  • static pod_name : str
  • static task_id : str

TaskCompleteTimeoutEvent​

class TaskCompleteTimeoutEvent(**data: Any):

Emitted when the worker times out waiting for TASK_COMPLETE from the modeller.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static datetime : str
  • static model_config
  • static modeller_username : str
  • static pod_name : str
  • static task_id : str

TaskErrorEvent​

class TaskErrorEvent(**data: Any):

Emitted by the worker when a task does not run to completion.

outcome says how it ended: "error" or "abort" for a failure inside the task, "rejected", "no_data", "data_not_available", "cancelled", "timeout" or "child_died" otherwise.

stacktrace_frames carries frame strings only - never the exception message or .args, which can name a patient. The handler fills it from the record's exc_info when the emitting code has the exception itself; it stays None on the paths that only learn a child process failed, because the traceback the child sent back is format_exception output and includes the message. error_type is available on both paths.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static error_type : str
  • static model_config
  • static modeller_username : str
  • static outcome : str
  • static pod_name : str
  • static stacktrace_frames : str | None
  • static task_id : str

TelemetryEventName​

class TelemetryEventName(*args, **kwds):

Canonical names for all Datadog telemetry events.

Inherits from str so values serialise directly as JSON strings.

Variables​

  • static ALGORITHM_PROGRESS
  • static BACKGROUND_DAG_RUN
  • static BACKGROUND_RUN_EMPTY
  • static BATCH_PROGRESS
  • static CRITERIA_MATCHING_SUMMARY
  • static DAG_EXECUTION_TIME
  • static DAG_SWEEP_SUMMARY
  • static DAG_WAVE
  • static DATASET_CONNECTED
  • static EHR_QUERY_COMPLETE
  • static EHR_QUERY_OUTCOME
  • static EHR_SCREENING_BATCH
  • static EHR_SESSION_INITIALISED
  • static EHR_TOKEN_REFRESHED
  • static FILE_MULTI_SERIES_SUMMARY
  • static FILE_SKIP_SUMMARY
  • static FILTERING_COMPLETE
  • static GPU_USAGE
  • static GROUPING_SUMMARY
  • static INDEXING_PROGRESS
  • static INDEXING_SUMMARY
  • static METADATA_KICKOFF
  • static MODEL_INFERENCE_PROGRESS
  • static PATIENT_API_REQUEST_ERROR
  • static PATIENT_API_SERVER
  • static PATIENT_ELIGIBILITY
  • static POD_ERROR
  • static RECOVERY_SUMMARY
  • static SCAN_REDUCE
  • static SCHEMA_GENERATION_END
  • static SCHEMA_GENERATION_ERROR
  • static SCHEMA_GENERATION_START
  • static SCHEMA_UPLOAD_SUCCESS
  • static TASK_ACCEPTED
  • static TASK_ACCEPTED_BY_POD
  • static TASK_COMPLETE
  • static TASK_COMPLETE_TIMEOUT
  • static TASK_ERROR
  • static TELEMETRY_DROPPED
  • static TRIAL_FILTER_SUMMARY
  • static WORKER_STARTUP_PHASE

TreeNodeSummary​

class TreeNodeSummary(**data: Any):

One criteria_tree node's aggregate outcome, leaf or combinator alike.

parent_node_id and depth are carried on every node so the whole tree is rebuildable from the node list alone, without a separate structure payload.

Attributes

  • node_key: The identifier to key a dashboard on, in precedence order: evidence_id when the author set one, else this node's own group when it is a region root, else node_id. Read node_key_stable to know whether it resolved to an author-set value or fell back.

An inherited group is never used. A region root's own group is unique tree-wide and author-set (_check_group_inheritance), so it identifies that one node; the descendants sharing the string do not set it and are not identified by it.

  • node_key_stable: Whether node_key came from an author-set value — evidence_id, or a region root's own group — rather than falling back to the positional node_id. A tree whose nodes are mostly False here cannot be tracked across a config edit; setting evidence_id on the nodes that matter fixes it.
  • is_region_root: Whether this node set the group it reports, rather than inheriting it. A region root's own verdict is the region's verdict, so filtering on this gives one row per reporting group without unrolling the subtree beneath it.
  • node_id: The node's auto-assigned node_<n> id.

Positional, and assigned in tree order, so inserting a node renumbers every node after it. Present for correlating against a single run's evidence; prefer node_key for anything that spans runs.

  • node_type: The node's discriminator — and, or, at_least, not, or a leaf type (column, code, observation, medication, allergy, device, appointment).
  • parent_node_id: The parent's node_id, or None for the root.
  • depth: Distance from the root, which is 0.
  • child_count: How many children the node has; 0 for a leaf.
  • group: The node's resolved (inherited) group reporting label, or None when neither it nor any ancestor set one.
  • evidence_id: The node's own evidence_id. Not inherited, and unique tree-wide, so it identifies exactly one node.
  • min_count: An AtLeastGroup's min_count; None on every other node.
  • combine: An AndGroup's combine mode; None on every other node.
  • rows_pass: Rows whose evaluation resolved this node to pass. For a combinator that is its own combined verdict, not a tally of its children.
  • rows_fail: Rows it resolved to fail.
  • rows_unknown: Rows it could not resolve.
  • rows_sole_blocker: Rows this node alone kept out of the cohort: rows where it failed while every sibling under the same AndGroup passed.

Three conditions, all required, because the point of the field is that loosening this one node would have admitted those rows:

  1. The row's tree verdict was fail. A row the tree accepted was kept out by nobody, whatever an individual branch did.
  2. The node's parent is an AndGroup, where every child must pass, so a lone failing child decided it. 0 for a node under any other combinator — an OrGroup sibling can carry the branch regardless, and under a NotGroup the failure is what let the row through.
  3. Every node between that parent and the root is also an AndGroup, so the group's failure reached the verdict rather than being absorbed. An AtLeastGroup whose min_count equals its child count is conjunctive in effect but does not count here — conservative, and it keeps the rule checkable by eye.

A nested AndGroup scores alongside the child that sank it: both statements are true, and the group's own count is worth having. worst_blocker_key breaks that tie towards the deeper node, being the criterion someone would actually loosen.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static child_count : int
  • static combine : str | None
  • static depth : int
  • static evidence_id : str | None
  • static group : str | None
  • static is_region_root : bool
  • static min_count : int | None
  • static model_config
  • static node_id : str
  • static node_key : str
  • static node_key_stable : bool
  • static node_type : str
  • static parent_node_id : str | None
  • static rows_fail : int
  • static rows_pass : int
  • static rows_sole_blocker : int
  • static rows_unknown : int

TrialFilterSummaryEvent​

class TrialFilterSummaryEvent(**data: Any):

Aggregate Datadog event for trial filter outcomes within a batch.

rows_missing counts rows with missing inputs required by the filter. It is not mutually exclusive with rows_matched or rows_failed.

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static filter_name : str
  • static matching_column_count : int | None
  • static missing_columns : list[str] | None
  • static model_config
  • static rows_failed : int
  • static rows_matched : int
  • static rows_missing : int
  • static step : str

WorkerStartupPhaseEvent​

class WorkerStartupPhaseEvent(**data: Any):

Emitted after each phase of worker startup, with how long that phase took.

Together these bracket the window between the worker reporting "Configuring task" and its first exchange with the modeller — protocol deserialisation and model download, protocol pickling, spawning the child interpreter, child setup, and algorithm initialisation. That window is otherwise silent, and a message arriving from the modeller during it is only tolerated for handler_register_grace_period seconds, so the split between these phases is what decides whether a task starts or stalls.

process is "parent" (the pod) or "child" (the spawned worker).

Create a new model by parsing and validating input data from keyword arguments.

Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.

self is explicitly positional-only to allow self as a field name.

Variables​

  • static duration_seconds : float
  • static model_config
  • static phase : str
  • static pod_name : str | None
  • static process : str
  • static task_id : str