Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion src/core/explain/evidence.py
Original file line number Diff line number Diff line change
Expand Up @@ -146,7 +146,8 @@ def _metric_onsets(db: Session, scope: str, window_start: datetime, window_end:

lookback = window_end - window_start
rows = db.execute(
select(MetricSample.service, MetricSample.metric, MetricSample.value, MetricSample.ts).where(
select(MetricSample.service, MetricSample.metric, MetricSample.value, MetricSample.ts,
MetricSample.attributes).where(
MetricSample.scope == scope,
MetricSample.ts >= window_start - lookback,
MetricSample.ts <= window_end,
Expand Down
11 changes: 10 additions & 1 deletion src/core/ingestion/telemetry.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
"""
from __future__ import annotations

import json
import uuid
from dataclasses import dataclass
from datetime import datetime
Expand Down Expand Up @@ -62,8 +63,16 @@ def span_pk(scope: str, s: ParsedSpan) -> uuid.UUID:


def metric_pk(scope: str, m: ParsedMetricSample) -> uuid.UUID:
"""Deterministic id for a metric sample — (scope, service, metric, ts)."""
"""Deterministic id for a metric sample — (scope, service, metric, ts[, series attributes]).

When the source recorded series identity (``attributes`` is a dict — datapoint attributes such as
``cpu.mode`` plus the reporting ``service.instance.id``), it is part of the key, so two series of
one metric sampled at the same timestamp are two rows rather than one silently dropped on
conflict. Sources without series identity (``attributes`` None) keep the original key, so their
idempotent re-ingest is unchanged."""
key = f"{scope}|{m.service}|{m.metric}|{m.ts}"
if m.attributes is not None:
key += "|" + json.dumps(m.attributes, sort_keys=True, default=str)
Comment thread
cursor[bot] marked this conversation as resolved.
return uuid.uuid5(_NS, key)


Expand Down
17 changes: 12 additions & 5 deletions src/core/rca/features.py
Original file line number Diff line number Diff line change
Expand Up @@ -301,17 +301,24 @@ def trace_symptoms(
def metric_features(
samples, baseline_start: datetime, incident_start: datetime, incident_end: datetime
) -> dict[str, float]:
"""Per service, the max over its metrics of the incident-vs-baseline mean
change ratio. ``samples`` have ``service`` / ``metric`` / ``value`` / ``ts``."""
base: dict[tuple[str, str], list[float]] = defaultdict(list)
inc: dict[tuple[str, str], list[float]] = defaultdict(list)
"""Per service, the max over its metric **series** of the incident-vs-baseline mean change ratio.
``samples`` have ``service`` / ``metric`` / ``value`` / ``ts`` (and optionally ``attributes``).

A series is ``(service, metric, series_key(attributes))`` (#209 M2b): an instrument that emits
several series (per ``cpu.mode``, per replica, …) is reduced per series, never averaged into a
number that describes none of them. Samples without recorded identity form one group per
``(service, metric)``, exactly as before identity existed."""
from src.core.rca.metric_series import series_key

base: dict[tuple, list[float]] = defaultdict(list)
inc: dict[tuple, list[float]] = defaultdict(list)
for m in samples:
s = m.service
ts = m.ts
v = m.value
if not s or ts is None or v is None:
continue
key = (s, m.metric)
key = (s, m.metric, series_key(getattr(m, "attributes", None)))
if baseline_start <= ts < incident_start:
base[key].append(float(v))
elif incident_start <= ts <= incident_end:
Expand Down
65 changes: 65 additions & 0 deletions src/core/rca/metric_series.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
"""Metric series identity (#209 M2b) — one definition shared by every reducer of metric samples.

A metric *name* is not a series. One instrument can emit several series — per datapoint attribute
(``cpu.mode``, ``system.memory.state``, a route) and per reporting instance (replicas) — and averaging
across them produces a number that describes no real series. Ingestion records a sample's series
identity in ``attributes``; every consumer that reduces metric samples must group by it.

Three levels of knowledge, deliberately kept distinct:

* ``attributes is None`` — the source never recorded identity (legacy corpora, RCAEval). All such
samples of a ``(service, metric)`` form one legacy group, exactly as before identity existed.
* a dict **with** a named instance — the series is fully identified.
* a dict **without** a named instance (``{}`` included) — the datapoint attributes are known, but the
reporting process is not: two replicas look like one series. That is *not* a verified series.
"""
from __future__ import annotations

import json
from typing import Optional

# Resource attributes that each name one reporting process (OpenTelemetry semantic conventions).
INSTANCE_ATTRS: tuple[str, ...] = ("service.instance.id", "k8s.pod.uid", "container.id")
# Attributes that name one process only together (a pid alone repeats across hosts).
INSTANCE_ATTR_PAIRS: tuple[tuple[str, str], ...] = (("host.name", "process.pid"),)
_INSTANCE_KEYS = frozenset(INSTANCE_ATTRS) | {k for pair in INSTANCE_ATTR_PAIRS for k in pair}


def series_key(attributes) -> Optional[str]:
"""Canonical full series identity (datapoint attributes + instance), or ``None`` when the source
never recorded identity. Reducers group by ``(service, metric, series_key)``."""
if attributes is None:
return None
return json.dumps(attributes, sort_keys=True, default=str)


def instance_identity(attributes) -> Optional[str]:
"""The reporting instance named by the recorded attributes, or ``None`` if none is named."""
if not isinstance(attributes, dict):
return None
for key in INSTANCE_ATTRS:
if attributes.get(key):
return f"{key}={attributes[key]}"
for a, b in INSTANCE_ATTR_PAIRS:
if attributes.get(a) and attributes.get(b):
return f"{a}={attributes[a]},{b}={attributes[b]}"
return None


def datapoint_signature(attributes) -> Optional[str]:
"""The datapoint-attribute part of the identity (instance attributes removed), or ``None`` when
identity was never recorded. Two samples with different signatures are different *modes* of the
instrument (e.g. ``cpu.mode=idle`` vs ``user``), not different replicas."""
if attributes is None:
return None
return json.dumps({k: v for k, v in attributes.items() if k not in _INSTANCE_KEYS},
sort_keys=True, default=str)


def instance_attributes(resource_attrs: dict[str, str]) -> dict[str, str]:
"""The instance-naming attributes present on a resource, for the converter to record."""
out = {k: resource_attrs[k] for k in INSTANCE_ATTRS if resource_attrs.get(k)}
for a, b in INSTANCE_ATTR_PAIRS:
if resource_attrs.get(a) and resource_attrs.get(b):
out[a], out[b] = resource_attrs[a], resource_attrs[b]
return out
6 changes: 3 additions & 3 deletions src/core/rca/structural.py
Original file line number Diff line number Diff line change
Expand Up @@ -66,12 +66,12 @@ def build_structural_view(
# The one shared structural-model builder (#209 M1): span-derived sig (latency median + explicit
# OTLP status only) merged with the metric sig, and the INCIDENT call graph (edges whose child span
# starts at/after window_start — the baseline half of this load feeds latency ratios, never edges).
signals, edges, edge_signals = structural_signals(span_rows, metric_rows, window_start)
hypotheses = build_hypotheses(signals, edges, edge_signals)
inp = structural_signals(span_rows, metric_rows, window_start)
hypotheses = build_hypotheses(inp.signals, inp.edges, inp.edge_signals, inp.util_signals)
if not hypotheses:
return None # no candidate -> abstain, never a fabricated structural result

observations = build_observables(signals, edge_signals)
observations = build_observables(inp.signals, inp.edge_signals, inp.util_signals)
part = partition(hypotheses, observations)
result = resolve(part, observations=observations)
ranking = rank_classes(part)
Expand Down
Loading
Loading