telemetry
Datadog telemetry for Bitfount.
Split into two pipelines:
bitfount.telemetry.logs— structured log events sent to the Datadog Logs API viatelemetry_logger.bitfount.telemetry.metrics— custom metrics sent to the Datadog Metrics API viametrics_logger/emit_metric.
The full public surface is re-exported here for convenience.
Module
Submodules
- bitfount.telemetry.algorithm - Telemetry hook reporting algorithm progress.
- bitfount.telemetry.context - Run context carried onto every Datadog telemetry event.
- bitfount.telemetry.criteria - Datadog telemetry for the
criteria_matchingstep. - bitfount.telemetry.dag - Telemetry for background DAG wave and sweep outcomes.
- bitfount.telemetry.datasource - Aggregated telemetry for per-file datasource outcomes.
- bitfount.telemetry.ehr - Telemetry for background EHR query step outcomes.
- bitfount.telemetry.errors - Turn an exception into something safe to send to Datadog.
- bitfount.telemetry.execution - Execution-time telemetry for the DAG paths.
- bitfount.telemetry.gpu - GPU usage telemetry for
model_inference. - bitfount.telemetry.indexing - Per-run counters for the background metadata runtimes.
- bitfount.telemetry.lifecycle - Telemetry for pod and task lifecycle events.
- bitfount.telemetry.logs - Datadog log telemetry: handler, lifecycle, logger, and event models.
- bitfount.telemetry.metrics - Datadog metrics telemetry: handler, lifecycle, logger, and metric registry.
- bitfount.telemetry.periodic_flush - A background thread that flushes a telemetry handler on a fixed interval.
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, orNoneif 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 aMetricNamemember 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 toGAUGE.
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 throughreport_step_progress.Nonewhen 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, orNoneif it returned. Defaults toNoneso 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. ANonevalue 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.Nonevalues 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.Nonevalues 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- static
model_config
- static
phase : str
- static
project_id : str | None
- static
seconds : float
- static
status : ExecutionStatus
- 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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.0withempty_reasonset means the run had nothing to assess, which is not a verdict — see that field.empty_reason: Why the run evaluated no rows, orNonewhen it evaluated some.no_files_selectedis expected and self-healing (a datasource still indexing).files_unreadablemeans files were selected and none could be loaded, which is unread data and needs someone to look. Both reportcriteria_yes: 0, and neither means nobody qualified.criteria_yes: Rows where every criterion resolved topass, thecriteria_tree's own combined verdict among them.unknowndoes not match, so a row missing an input counts here exactly as a row that failed outright does — the two are told apart byrows_unknownon the criterion, not by this.criteria_no: The rest.criterion_count: Flat (non-tree) criteria evaluated, which is the length ofcriteriabefore truncation. A range field counts twice, once per bound.criteria_treeleaves are not included; they are counted bytree_leaf_count.tree_node_count: Nodes in thecriteria_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: Whethertree_shapewas 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: Whethertree_outlinewas cut short.tree_root_pass: Rows the tree's root resolved topass.tree_root_fail: Rows it resolved tofail.tree_root_unknown: Rows it could not resolve.worst_blocker_key: Thenode_keyof the node that alone kept the most rows out — the highestrows_sole_blockerin the tree, promoted here so the headline is alertable without unrollingtree_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'snode_id.worst_blocker_evidence_id: That node'sevidence_id, orNonewhere it set none.worst_blocker_group: That node's resolvedgroup, orNone.worst_blocker_rows: That node'srows_sole_blocker.0exactly whenworst_blocker_keyisNone.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 fromcriteriato bound the payload;0normally.tree_nodes_truncated: Nodes omitted fromtree_nodesfor 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 : list[CriterionSummary]
- 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
event : TelemetryEventName
- 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 : list[TreeNodeSummary]
- 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 aColumnFilterthis 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 producesrows_failed_near_bound.Nonefor a filter with no named comparison at all.operand: The configured right-hand side of a scalar comparison.Nonefor a code-list criterion, whose operand is the whole configured list:configured_code_countreports its length instead, so a few hundred codes are not repeated into every event.grain:scanorpatient— the cardinality axis.provenance:scan,ehr,suppliedorunknown— the source axis.family: The semantic family (code_list,scan_metric, ...), orNonewhen 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 everyrows_*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 fromrows_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 of0, where that fraction has no meaning, only an exact0counts.
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 singleH35.3*covering twelve real codes counts once. Includes the_study_eyecounterpart list where the criterion has one, since a code configured in either was genuinely configured.Nonefor a non-code criterion.configured_codes_matching_nobody: How many of those matched no code the cohort carried. Equal toconfigured_code_countmeans 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 becauseenable_criteria_near_miss_countsis 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, notpatients: 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
event : TelemetryEventName
- static
files_in_scope : int | None
- static
files_processed : int | None
- static
function : str | None
- static
model_config
- static
patients_processed : int | None
- static
status : ExecutionStatus
- 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;Nonewhen 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
event : TelemetryEventName
- 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, sowave_indexreads 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;Falsewhen 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
event : TelemetryEventName
- 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 issetup_datadog_telemetry, which passes the configured interval.
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 onextra(not as loggingargs) 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 toGAUGE.
See: https://docs.datadoghq.com/api/latest/metrics/
Initialise the handler.
Arguments
api_instance: A configured DatadogMetricsApiinstance.hostname: Attached to everyMetricSeriesviaresourcesso 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 issetup_datadog_metrics, which passes the configured interval.
Ancestors
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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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.
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.
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.
Ancestors
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
event : TelemetryEventName
- 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 byFileSkipReasonname.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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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.Nonewhen 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;Nonefor a CPU-only build.torch_cuda_available:torch.cuda.is_available().torch_device: Device the SDK selects for torch models:cuda,mpsorcpu.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
event : TelemetryEventName
- 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
status : ExecutionStatus
- 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 atstart.
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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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
unhealthy : list[UnhealthyDatasource]
- 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.
Ancestors
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, ofbatches_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, offiles_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
event : TelemetryEventName
- 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_endhook, counting eligible patients over a task's merged output; it has thedatasourceand no denominator. - the
patient_eligibilityreduce step, which has the denominator,task_hashandrun_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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- static
model_config
- static
source_of_initiation : SchemaInitiationSource
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
event : TelemetryEventName
- 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
event : TelemetryEventName
- static
model_config
- static
source_of_initiation : SchemaInitiationSource
SchemaInitiationSource
class SchemaInitiationSource(*args, **kwds):Source that triggered schema generation.
Inherits from str so values serialise directly as JSON strings.
Ancestors
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
event : TelemetryEventName
- 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.
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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- 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.
Ancestors
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_idwhen the author set one, else this node's owngroupwhen it is a region root, elsenode_id. Readnode_key_stableto 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: Whethernode_keycame from an author-set value —evidence_id, or a region root's owngroup— rather than falling back to the positionalnode_id. A tree whose nodes are mostlyFalsehere cannot be tracked across a config edit; settingevidence_idon the nodes that matter fixes it.is_region_root: Whether this node set thegroupit 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-assignednode_<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'snode_id, orNonefor the root.depth: Distance from the root, which is0.child_count: How many children the node has;0for a leaf.group: The node's resolved (inherited)groupreporting label, orNonewhen neither it nor any ancestor set one.evidence_id: The node's ownevidence_id. Not inherited, and unique tree-wide, so it identifies exactly one node.min_count: AnAtLeastGroup'smin_count;Noneon every other node.combine: AnAndGroup'scombinemode;Noneon every other node.rows_pass: Rows whose evaluation resolved this node topass. For a combinator that is its own combined verdict, not a tally of its children.rows_fail: Rows it resolved tofail.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 sameAndGrouppassed.
Three conditions, all required, because the point of the field is that loosening this one node would have admitted those rows:
- The row's tree verdict was
fail. A row the tree accepted was kept out by nobody, whatever an individual branch did. - The node's parent is an
AndGroup, where every child must pass, so a lone failing child decided it.0for a node under any other combinator — anOrGroupsibling can carry the branch regardless, and under aNotGroupthe failure is what let the row through. - Every node between that parent and the root is also an
AndGroup, so the group's failure reached the verdict rather than being absorbed. AnAtLeastGroupwhosemin_countequals 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
event : TelemetryEventName
- 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
event : TelemetryEventName
- static
model_config
- static
phase : str
- static
pod_name : str | None
- static
process : str
- static
task_id : str