Skip to main content

context

Run context carried onto every Datadog telemetry event.

Identity (username, pod_identifier) and run coordinates (dataset_name, task_id, run_id) are known at different times and in different places from where telemetry is configured: setup_datadog_telemetry runs at process start, before login in several entry points, and its tags are fixed for the life of the handler. This module holds those values instead, and DatadogLogsHandler.emit reads them per record.

Two layers, because a contextvars value set after a worker thread copied its context is invisible to that thread:

  • The global layer is process-wide and lock-guarded, for values true of the whole process - who is logged in, which pod this is.
  • The overlay is a ContextVar, for values true of one run or one step. run_on_daemon_thread copies the context, so a DAG step inherits the overlay its executor set.

Every read is total: telemetry must not raise into the code it observes.

Module​

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_files_in_scope​

def current_files_in_scope() ‑> int | None:

Return the file count noted for the running DAG or wave, if any.

current_step_progress​

def current_step_progress() ‑> StepProgress | None:

Return the recorder bound to the current context, if any.

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.

files_in_scope_slot​

def files_in_scope_slot(    ,) ‑> collections.abc.Iterator[FilesInScopeSlot]:

Bind a fresh slot for the entry point about to run to note into.

The slot starts from the enclosing slot's value, so an entry point that notes nothing reports the denominator it was called with rather than nothing at all.

note_files_in_scope​

def note_files_in_scope(count: int | None) ‑> None:

Record how many files the running DAG or wave was given.

Written into the slot the timing decorator bound, not into a value scoped around the entry point: the timing event is emitted after the body returns, so a scope closed inside it would already have been restored. Writing into the innermost slot is what keeps a wave's count off the run's timing.

Arguments

  • count: Files in scope, or None when the run has no file scope to report - an unwaved run does not narrow the datasource, so there is no cheap count to take and none is invented.

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.

step_progress_scope​

def step_progress_scope(progress: StepProgress) ‑> collections.abc.Iterator[None]:

Make progress the recorder report_step_progress writes to.

Arguments

  • progress: The recorder to bind for the duration of the block.

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.

Classes​

FilesInScopeSlot​

class FilesInScopeSlot(value: int | None = None):

Where one entry point's file denominator is held while it runs.

A slot per entry point rather than a single value per context: a wave runs inside the run that scheduled it, and both are timed, so a shared value would leave the run's timing reporting the last wave's count. The slot is an object the caller keeps a reference to, so the timing decorator can read it after the body has returned and the binding has been restored.

Attributes

  • value: Files this entry point was given, or None when it has none to report.

Variables​

  • static value : int | None

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