zephon.observability

Observability configuration and helpers used across Zephon.

class ExecutionTrackingMode(*values)[source]

Bases: str, Enum

Controls how much execution detail is collected.

OFF = 'off'
STAGES = 'stages'
NODES = 'nodes'
classmethod from_value(value)[source]
Return type:

ExecutionTrackingMode

property collects_nodes: bool
property collects_stages: bool
class FetchStageSummary(index, totals=<factory>, shard_totals=<factory>)[source]

Bases: object

Stage fetch timings keyed by (dataset_id, shard_id).

index: int
totals: FetchTimingTotals
shard_totals: MutableMapping[tuple[int, int], FetchTimingTotals]
apply(delta)[source]
merge(other)[source]
copy()[source]
Return type:

FetchStageSummary

to_records(*, plan_id, stage_name, tracking_mode, include_shards=True)[source]
Return type:

list[dict[str, Any]]

class FetchTimingDelta(stage_index, dataset_id, shard_id, samples, group_ns, resolve_ns, open_ns, read_ns, close_ns, retries, cache_hits, cache_misses, shard_reopens=0)[source]

Bases: object

Incremental fetch metrics emitted per shard group.

stage_index: int
dataset_id: int
shard_id: int
samples: int
group_ns: int
resolve_ns: int
open_ns: int
read_ns: int
close_ns: int
retries: int
cache_hits: int
cache_misses: int
shard_reopens: int
class FetchTimingSummary(plan_id=None, reporting_interval_s=None, tracking_mode=ExecutionTrackingMode.OFF, stages=<factory>)[source]

Bases: object

Aggregated fetch metrics for an entire pipeline.

plan_id: str | None = None
reporting_interval_s: float | None = None
tracking_mode: ExecutionTrackingMode = 'off'
stages: MutableMapping[int, FetchStageSummary]
apply(delta)[source]
merge(other)[source]
iter_stages()[source]
Return type:

Iterator[FetchStageSummary]

to_records(stage_names=None, *, include_shards=True)[source]
Return type:

list[dict[str, Any]]

clone()[source]
Return type:

FetchTimingSummary

has_samples()[source]
Return type:

bool

class FetchTimingTotals(samples=0, groups=0, group_ns=0, resolve_ns=0, open_ns=0, read_ns=0, close_ns=0, retries=0, cache_hits=0, cache_misses=0, shard_reopens=0)[source]

Bases: object

Aggregated timing totals accumulated across deltas.

samples: int
groups: int
group_ns: int
resolve_ns: int
open_ns: int
read_ns: int
close_ns: int
retries: int
cache_hits: int
cache_misses: int
shard_reopens: int
apply(delta)[source]
merge(other)[source]
copy()[source]
Return type:

FetchTimingTotals

property avg_group_ns: float
property avg_resolve_ns: float
property avg_open_ns: float
property avg_read_ns: float
property avg_close_ns: float
property cache_hit_ratio: float
class MTPQueueStats(depth, capacity, staged_bytes, prefetch_depth)[source]

Bases: object

Point-in-time occupancy of the MTP hand-off queue (data_q).

depth counts queued items, which may not have been pickled yet; staged_bytes is the already-pickled data in the kernel transport buffer, readable without the feeder running. Zero staged with depth near capacity means the feeder is starved. prefetch_depth counts items already drained into the main-process prefetch buffer; it is sampled without locking and may be slightly stale.

A field is -1 when unmeasurable: depth on macOS (sem_getvalue is unsupported there), staged_bytes once the queue is closed.

depth: int

Items put by the subprocess but not yet retrieved by the consumer.

capacity: int

Maximum depth — the resolved mtp_buffer.

staged_bytes: int

Pickled bytes staged in the transport buffer, consumer-ready.

prefetch_depth: int

Deserialized items in the prefetch buffer (0 when mtp_prefetch=0).

class MetricsSinkConfig(mode=MetricsSinkMode.LOG, namespace='zephon.pipeline', statsd_host='127.0.0.1', statsd_port=8125, default_tags=<factory>, json_logs=True, max_batch_size=32, flush_interval_s=5.0)[source]

Bases: object

Describe how observability data is exported.

mode: MetricsSinkMode = 'log'
namespace: str = 'zephon.pipeline'
statsd_host: str = '127.0.0.1'
statsd_port: int = 8125
default_tags: Mapping[str, str]
json_logs: bool = True
max_batch_size: int = 32
flush_interval_s: float = 5.0
should_log()[source]
Return type:

bool

should_emit_statsd()[source]
Return type:

bool

with_additional_tags(extra)[source]
Return type:

MetricsSinkConfig

class MetricsSinkMode(*values)[source]

Bases: str, Enum

Where metrics should be emitted.

LOG = 'log'
STATSD = 'statsd'
BOTH = 'both'
classmethod from_value(value)[source]
Return type:

MetricsSinkMode

includes_logs()[source]
Return type:

bool

includes_statsd()[source]
Return type:

bool

class PipelineSummary(plan_id=None, reporting_interval_s=None, tracking_mode=ExecutionTrackingMode.OFF, stages=<factory>, total_wait_ratio=None)[source]

Bases: object

Aggregated metrics for an entire pipeline invocation.

plan_id: str | None = None
reporting_interval_s: float | None = None
tracking_mode: ExecutionTrackingMode = 'off'
stages: MutableMapping[int, StageSummary]
total_wait_ratio: float | None = None
apply(delta)[source]
merge(other)[source]
iter_stages()[source]
Return type:

Iterator[StageSummary]

to_records()[source]
Return type:

list[dict[str, Any]]

clone()[source]
Return type:

PipelineSummary

compute_wait_ratios()[source]

Derive wait ratios for the pipeline based on processing and wait times.

class PrefetchTimingDelta(stage_index, batch_size, prefetch_requests, prefetch_succeeded, prefetch_failed)[source]

Bases: object

Incremental prefetch metrics emitted per batch.

stage_index: int
batch_size: int
prefetch_requests: int
prefetch_succeeded: int
prefetch_failed: int
class PrefetchTimingSummary(plan_id=None, reporting_interval_s=None, tracking_mode=ExecutionTrackingMode.OFF, stages=<factory>)[source]

Bases: object

Aggregated prefetch metrics for an entire pipeline.

plan_id: str | None = None
reporting_interval_s: float | None = None
tracking_mode: ExecutionTrackingMode = 'off'
stages: MutableMapping[int, PrefetchTimingTotals]
apply(delta)[source]
merge(other)[source]
clone()[source]
Return type:

PrefetchTimingSummary

has_samples()[source]
Return type:

bool

property success_rate: float
class PrefetchTimingTotals(batches=0, samples=0, prefetch_requests=0, prefetch_succeeded=0, prefetch_failed=0)[source]

Bases: object

Aggregated prefetch totals accumulated across deltas.

batches: int
samples: int
prefetch_requests: int
prefetch_succeeded: int
prefetch_failed: int
apply(delta)[source]
merge(other)[source]
copy()[source]
Return type:

PrefetchTimingTotals