From 7220306c766e75654c701a758b1118988095222b Mon Sep 17 00:00:00 2001 From: Leonardo Araujo Date: Sat, 26 Sep 2026 07:07:49 -0300 Subject: [PATCH 1/4] feat(rca): utilization observables for the structural model (#209 M2b) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds util:{service}:{resource} as its own observable family (never folded into sig:{service}) to recover locally silent RESOURCE faults — a CPU- or event-loop- saturated service whose requests still look normal in traces. Trace evidence plateaued at 4/8 (M2a); the ceiling survey found clean, attributable metric separation for 2 of the 4 remaining misses. Frozen semantics (posted to #209 before implementing): - Metrics are classified by OpenTelemetry semantic-convention NAME SHAPE, not by names picked from the corpus (util_metric_class): a utilization gauge (last segment ends in "utilization"; resource = preceding segment, e.g. cpu, eventloop, memory) or a cumulative CPU-time counter (*.cpu.time -> cpu). Collector/SDK self-telemetry (otel.sdk.*, otelcol*) is excluded — it moves when a service merely emits more spans (the product-catalog false-positive trap). - summarize_utilization: gauges compare incident vs baseline MEAN (periodic aggregates); a malformed sample leaves the metric unmeasured. CPU-time counters compare incident vs baseline RATE, trusted only for a verifiably single-series, monotonic stream — the converter flattens attribute dimensions, so interleaved series (e.g. process.cpu.time state=user/system: 2 samples/timestamp, decreases) are UNKNOWN rather than guessed. Unattributed samples (service=None) are not evidence. PRESENT iff ratio >= 2x; ABSENT if measured below (a drop is not a saturation anomaly); else UNKNOWN. The (service, resource) observable takes the witness metric (PRESENT > ABSENT > UNKNOWN). - build_hypotheses: a service with a PRESENT util observable joins the candidates and process:{S} expects its util coordinates PRESENT. Soft only; without util signals the M2a model is byte-identical. - structural_signals now returns a StructuralInputs NamedTuple (signals, edges, edge_signals, util_signals); both runners use it — still one model. 19 new unit tests (name classes, self-telemetry exclusion, gauge/counter rules, interleaved + reset counters, witness, model wiring, end-to-end locally silent CPU fault). Full unit suite 1515; structural/trigger/absence integration green. Co-Authored-By: Claude Opus 5.5 --- src/core/rca/structural.py | 6 +- src/core/rca/structural_model.py | 187 ++++++++++++++++++--- src/eval/structural_shadow.py | 9 +- tests/unit/test_structural_edge_signals.py | 4 +- tests/unit/test_structural_span_signals.py | 3 +- tests/unit/test_structural_util_signals.py | 161 ++++++++++++++++++ 6 files changed, 340 insertions(+), 30 deletions(-) create mode 100644 tests/unit/test_structural_util_signals.py diff --git a/src/core/rca/structural.py b/src/core/rca/structural.py index a439732..90ecbd1 100644 --- a/src/core/rca/structural.py +++ b/src/core/rca/structural.py @@ -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) diff --git a/src/core/rca/structural_model.py b/src/core/rca/structural_model.py index 7eaa4d6..881af1f 100644 --- a/src/core/rca/structural_model.py +++ b/src/core/rca/structural_model.py @@ -24,7 +24,7 @@ from collections import Counter, defaultdict from dataclasses import dataclass from datetime import datetime -from typing import Optional +from typing import NamedTuple, Optional from src.core.rca.expectations import ExpectedObservation, ObservationModel from src.core.rca.hypothesis import Hypothesis, Kind, from_observation_model @@ -36,6 +36,11 @@ def _median(xs: list[float]) -> float: return statistics.median(xs) if xs else 0.0 +def _mean(xs: list[float]) -> float: + """For periodic gauge samples (already aggregates), unlike raw per-request durations.""" + return sum(xs) / len(xs) if xs else 0.0 + + # OTLP span status codes: 0 = UNSET, 1 = OK, 2 = ERROR. Only an EXPLICIT status is a measurement: # UNSET (or null / missing) says nothing about whether the request succeeded, so it is excluded from # the error-rate denominator entirely — it is neither a success nor a failure. Counting UNSET as @@ -372,27 +377,163 @@ def summarize_edges( return out -def structural_signals( - spans: list, metrics: list, window_start: datetime -) -> tuple[dict[str, ServiceSignal], set[tuple[str, str]], dict[tuple[str, str], EdgeSignal]]: +@dataclass(frozen=True) +class UtilSignal: + """A service's resource-utilization summary (#209 M2b) — the observable ``util:{service}:{resource}``, + its own coordinate family, never folded into ``sig:{service}``. It recovers *locally silent resource + faults*: a CPU-saturated or event-loop-saturated service whose requests still look normal in traces. + ``ratio`` is the witness metric's incident/baseline ratio; ``sig_state`` is three-valued.""" + + service: str + resource: str + ratio: float + measured: bool + sig_state: Optional[str] + + @property + def id(self) -> str: + return util_observable_id(self.service, self.resource) + + +def util_observable_id(service: str, resource: str) -> str: + return f"util:{service}:{resource}" + + +_SELF_TELEMETRY_PREFIXES = ("otel.sdk.", "otelcol") + + +def util_metric_class(metric: Optional[str]) -> Optional[tuple[str, str]]: + """Classify a metric by OpenTelemetry semantic-convention **name shape** (never by names picked from + a corpus). Returns ``(resource, kind)`` or ``None``: + + * a utilization **gauge** — last dotted segment ends in ``utilization``; the resource is the segment + before it (``jvm.cpu.recent_utilization`` → ``cpu``, ``nodejs.eventloop.utilization`` → + ``eventloop``, ``system.memory.utilization`` → ``memory``); + * a cumulative **CPU-time counter** — name ends ``.cpu.time`` (semconv seconds); resource ``cpu``. + + Collector/SDK self-telemetry (``otel.sdk.*``, ``otelcol*``) is excluded: it measures the telemetry + pipeline, not the service, and moves whenever a service simply emits more spans.""" + if not metric or metric.startswith(_SELF_TELEMETRY_PREFIXES): + return None + parts = metric.split(".") + if len(parts) >= 2 and parts[-1].endswith("utilization"): + return parts[-2], "gauge" + if metric.endswith(".cpu.time"): + return "cpu", "cpu_time" + return None + + +def _counter_rate(points: list[tuple[datetime, float]]) -> Optional[float]: + if len(points) < 2: + return None + dt = (points[-1][0] - points[0][0]).total_seconds() + return (points[-1][1] - points[0][1]) / dt if dt > 0 else None + + +def summarize_utilization(samples: list, window_start: datetime) -> dict[tuple[str, str], UtilSignal]: + """Per ``(service, resource)`` utilization observable from metric-sample records (#209 M2b), duck-typed + on ``service`` / ``metric`` / ``value`` / ``ts``. + + Per metric (only classes recognized by :func:`util_metric_class`; unattributed ``service=None`` + samples contribute nothing): + + * **gauge** — incident mean / baseline mean; measured iff both windows have samples, every raw value + is finite and ≥ 0, and the baseline mean is > 0 (a malformed sample leaves the metric unmeasured, + never averaged into false-normal evidence — as in :func:`summarize_metrics`). + * **cpu_time** — incident rate / baseline rate, rate = (last − first) / Δt. Measured **only for a + verifiably single-series, monotonic stream**: at most one sample per timestamp and never decreasing + across the whole load. A converter that flattened attribute dimensions interleaves several + cumulative series (e.g. ``state=user``/``system``) into one stream, and a rate over that is + meaningless — so such a stream is UNKNOWN, not guessed. + + A metric is PRESENT if its ratio ≥ 2× (``discretize_ratio`` HIGH), ABSENT if measured and below, + else UNKNOWN. The ``(service, resource)`` signal takes the witness: PRESENT if any metric of that + class is PRESENT, else ABSENT if any is measured-normal, else UNKNOWN. A utilization *drop* is ABSENT + (not saturated), not an anomaly.""" + series: dict[tuple[str, str], dict[str, list[tuple[datetime, float]]]] = defaultdict( + lambda: {"b": [], "i": []} + ) + invalid: set[tuple[str, str]] = set() + for s in samples: + svc, metric, value, ts = (getattr(s, "service", None), getattr(s, "metric", None), + getattr(s, "value", None), getattr(s, "ts", None)) + if svc is None or value is None or ts is None or util_metric_class(metric) is None: + continue + key = (svc, metric) + v = float(value) + if not math.isfinite(v) or v < 0: + invalid.add(key) + continue + series[key]["i" if ts >= window_start else "b"].append((ts, v)) + + per_class: dict[tuple[str, str], list[tuple[Optional[str], float, bool]]] = defaultdict(list) + for (svc, metric), d in series.items(): + resource, kind = util_metric_class(metric) + ratio, measured = 1.0, False + if (svc, metric) not in invalid and d["b"] and d["i"]: + if kind == "gauge": + base = _mean([v for _, v in d["b"]]) + if base > 0: + ratio, measured = _mean([v for _, v in d["i"]]) / base, True + else: # cpu_time: single-series + monotonic, or UNKNOWN + pts = sorted(d["b"] + d["i"]) + single = len({t for t, _ in pts}) == len(pts) + monotonic = all(b >= a for (_, a), (_, b) in zip(pts, pts[1:])) + rb, ri = _counter_rate(sorted(d["b"])), _counter_rate(sorted(d["i"])) + if single and monotonic and rb is not None and rb > 0 and ri is not None: + ratio, measured = ri / rb, True + if not measured: + state: Optional[str] = None + elif discretize_ratio(ratio) == State.HIGH: + state = State.PRESENT + else: + state = State.ABSENT + per_class[(svc, resource)].append((state, ratio, measured)) + + out: dict[tuple[str, str], UtilSignal] = {} + for (svc, resource), metrics in sorted(per_class.items()): + witness = (next((m for m in metrics if m[0] == State.PRESENT), None) + or next((m for m in metrics if m[0] == State.ABSENT), None) + or metrics[0]) + out[(svc, resource)] = UtilSignal(svc, resource, witness[1], witness[2], witness[0]) + return out + + +class StructuralInputs(NamedTuple): + """Everything the structural model reads for one ``(scope, window)`` (#209): per-service ``sig``, + the incident call graph, per-edge observables (M2a) and per-resource utilization observables (M2b).""" + + signals: dict[str, ServiceSignal] + edges: set[tuple[str, str]] + edge_signals: dict[tuple[str, str], EdgeSignal] + util_signals: dict[tuple[str, str], UtilSignal] + + +def structural_signals(spans: list, metrics: list, window_start: datetime) -> StructuralInputs: """The single structural-model input builder shared by the product path (:func:`~src.core.rca.structural.build_structural_view`) and the shadow eval (``src.eval.structural_shadow.shadow_result``), so both run one model. ``spans`` and ``metrics`` span - the baseline *and* incident windows. Returns ``(signals, edges, edge_signals)``: span-derived and - metric-derived ``sig`` merged by :func:`combine_signals`, the **incident** call graph, and the - per-edge observables (#209 M2a).""" + the baseline *and* incident windows. Returns :class:`StructuralInputs`: span-derived and + metric-derived ``sig`` merged by :func:`combine_signals`, the **incident** call graph, the per-edge + observables (M2a) and the utilization observables (M2b).""" signals = combine_signals(summarize_spans(spans, window_start), summarize_metrics(metrics, window_start)) - return signals, call_edges(spans, since=window_start), summarize_edges(spans, window_start) + return StructuralInputs( + signals=signals, + edges=call_edges(spans, since=window_start), + edge_signals=summarize_edges(spans, window_start), + util_signals=summarize_utilization(metrics, window_start), + ) def build_observables( signals: dict[str, ServiceSignal], edge_signals: Optional[dict[tuple[str, str], EdgeSignal]] = None, + util_signals: Optional[dict[tuple[str, str], UtilSignal]] = None, ) -> list[Observable]: """One ``sig:{service}`` observable per service whose combined signal is *proven* PRESENT or ABSENT - (three-valued), plus one ``edge:{caller}->{callee}`` observable per proven call edge. Anything - UNKNOWN (a branch unmeasured, or an edge with no incident calls) is omitted — its coordinate stays - UNKNOWN, never a fabricated OBSERVED ABSENT.""" + (three-valued), plus one ``edge:{caller}->{callee}`` per proven call edge and one + ``util:{service}:{resource}`` per proven utilization class. Anything UNKNOWN is omitted — its + coordinate stays UNKNOWN, never a fabricated OBSERVED ABSENT.""" obs = [ observed(f"sig:{svc}", sig.sig_state) for svc, sig in sorted(signals.items()) @@ -401,6 +542,9 @@ def build_observables( for _key, e in sorted((edge_signals or {}).items()): if e.sig_state is not None: obs.append(observed(e.id, e.sig_state)) + for _key, u in sorted((util_signals or {}).items()): + if u.sig_state is not None: + obs.append(observed(u.id, u.sig_state)) return obs @@ -418,31 +562,36 @@ def build_hypotheses( signals: dict[str, ServiceSignal], edges: set[tuple[str, str]], edge_signals: Optional[dict[tuple[str, str], EdgeSignal]] = None, + util_signals: Optional[dict[tuple[str, str], UtilSignal]] = None, ) -> list[Hypothesis]: """Candidate ``process`` hypotheses: anomalous services, plus a service's callees when its fault is not already explained by a visible (anomalous) callee — so a silent downstream root is still - generated — plus (#209 M2a) the **callee of every PRESENT edge**: a failing or slow call toward a - service is evidence about that service even when it emits nothing itself (the unreachable-callee - case). Each hypothesis predicts its own ``sig`` present and, when edge signals are given, every - incoming ``edge:{caller}->{svc}`` present. All expectations are soft; dependency direction is **not** - a hard rule (an anomalous callee doesn't exclude its caller as root — that stays soft ranking).""" + generated — plus (#209 M2a) the **callee of every PRESENT edge**, plus (#209 M2b) **every service with + a PRESENT utilization observable** (a locally silent resource fault). Each hypothesis predicts its own + ``sig`` present and, when given, every incoming ``edge:{caller}->{svc}`` and every + ``util:{svc}:{resource}`` present. All expectations are soft; dependency direction is **not** a hard + rule (an anomalous callee doesn't exclude its caller as root — that stays soft ranking).""" anomalous = {s for s, sig in signals.items() if sig.anomalous} candidates = set(anomalous) for s in anomalous: callee_set = _callees(s, edges) if not any(c in anomalous for c in callee_set): # fault unexplained by a visible callee candidates |= callee_set - incoming: dict[str, list[str]] = defaultdict(list) + expected_extra: dict[str, list[str]] = defaultdict(list) for e in (edge_signals or {}).values(): - incoming[e.callee].append(e.id) + expected_extra[e.callee].append(e.id) if e.sig_state == State.PRESENT: candidates.add(e.callee) + for u in (util_signals or {}).values(): + expected_extra[u.service].append(u.id) + if u.sig_state == State.PRESENT: + candidates.add(u.service) return [ from_observation_model( f"process:{svc}", Kind.PROCESS, svc, ObservationModel(expected=( ExpectedObservation(f"sig:{svc}", State.PRESENT), - *(ExpectedObservation(eid, State.PRESENT) for eid in sorted(incoming[svc])), + *(ExpectedObservation(oid, State.PRESENT) for oid in sorted(expected_extra[svc])), )), ) for svc in sorted(candidates) diff --git a/src/eval/structural_shadow.py b/src/eval/structural_shadow.py index 44a30c4..7f77319 100644 --- a/src/eval/structural_shadow.py +++ b/src/eval/structural_shadow.py @@ -102,19 +102,18 @@ def shadow_result(case: EvalCase) -> ShadowResult: # The same structural-model builder the product path uses (#209 M1): span+metric sig and the # incident call graph — one model, not a divergent shadow copy. spans = load_spans_jsonl(case.spans_path) - signals, edges, edge_signals = structural_signals( - spans, load_metrics_jsonl(case.metrics_path), case.window_start - ) + inp = structural_signals(spans, load_metrics_jsonl(case.metrics_path), case.window_start) + signals = inp.signals # The service universe is every service seen in telemetry — metric-bearing services and all span # services (incl. root-only spans with no edge) — so candidate_ratio's "1.0 = enumerate all" holds. n_services = len(service_universe(signals, {sp.service for sp in spans if sp.service})) - hypotheses = build_hypotheses(signals, edges, edge_signals) + hypotheses = build_hypotheses(signals, inp.edges, inp.edge_signals, inp.util_signals) if not hypotheses: # No hypothesis returned -> an abstention, not a zero-abstention success. return ShadowResult(case.id, truth, "no_candidates", (), (), False, False, True, n_candidates=0, n_services=n_services) - result = resolve(partition(hypotheses, build_observables(signals, edge_signals))) + result = resolve(partition(hypotheses, build_observables(signals, inp.edge_signals, inp.util_signals))) classes = tuple(ClassInfo(c.localizations, c.signature, c.d_missing) for c in result.classes) localizations = result.localization return ShadowResult( diff --git a/tests/unit/test_structural_edge_signals.py b/tests/unit/test_structural_edge_signals.py index 0a8f39e..ae40f30 100644 --- a/tests/unit/test_structural_edge_signals.py +++ b/tests/unit/test_structural_edge_signals.py @@ -148,13 +148,13 @@ def _spans(self): return spans def test_m2a_retains_the_unreachable_callee(self): - signals, edges, edge_signals = structural_signals(self._spans(), [], _W) + signals, edges, edge_signals, _ = structural_signals(self._spans(), [], _W) assert "payment" not in {c for _, c in edges} # no incident payment span -> not in call graph res = resolve(partition(build_hypotheses(signals, edges, edge_signals), build_observables(signals, edge_signals))) assert "payment" in res.localization def test_m1_model_alone_cannot(self): - signals, edges, _ = structural_signals(self._spans(), [], _W) + signals, edges, _, _ = structural_signals(self._spans(), [], _W) hyps = build_hypotheses(signals, edges) # no edge signals = the M1 model assert "payment" not in {h.localization for h in hyps} diff --git a/tests/unit/test_structural_span_signals.py b/tests/unit/test_structural_span_signals.py index 83ff6a9..e376819 100644 --- a/tests/unit/test_structural_span_signals.py +++ b/tests/unit/test_structural_span_signals.py @@ -156,6 +156,7 @@ def test_baseline_neighbour_is_not_generated_as_a_candidate(self): class TestSharedBuilder: def test_structural_signals_merges_modalities_and_uses_the_incident_graph(self): spans = _stream("checkout", inc_ms=40.0) + TestIncidentCallGraph()._spans() - signals, edges, _edge_signals = structural_signals(spans, [], _W) + inp = structural_signals(spans, [], _W) + signals, edges = inp.signals, inp.edges assert signals["checkout"].sig_state == State.PRESENT assert ("checkout", "fraud") not in edges diff --git a/tests/unit/test_structural_util_signals.py b/tests/unit/test_structural_util_signals.py new file mode 100644 index 0000000..d5897d7 --- /dev/null +++ b/tests/unit/test_structural_util_signals.py @@ -0,0 +1,161 @@ +"""#209 M2b — utilization observables util:{service}:{resource}. Pure, no DB. + +Pins the frozen semantics: metrics are classified by OpenTelemetry semantic-convention name shape (not +by names picked from a corpus), collector/SDK self-telemetry is excluded, a cumulative CPU-time rate is +only trusted for a verifiably single-series monotonic stream, unattributed or malformed samples never +become evidence, and the observable family stays separate from sig:{service}.""" +from datetime import datetime, timedelta, timezone + +from src.core.rca.observable import State +from src.core.rca.outcome import resolve +from src.core.rca.partition import partition +from src.core.rca.structural_model import ( + ServiceSignal, + UtilSignal, + build_hypotheses, + build_observables, + structural_signals, + summarize_utilization, + util_metric_class, +) + +_W = datetime(2026, 1, 1, 12, 0, 0, tzinfo=timezone.utc) + + +class _M: + def __init__(self, service, metric, value, ts): + self.service, self.metric, self.value, self.ts = service, metric, value, ts + + +def _gauge(service, metric, base, inc, n=5): + out = [_M(service, metric, base, _W - timedelta(seconds=10 * (i + 1))) for i in range(n)] + out += [_M(service, metric, inc, _W + timedelta(seconds=10 * (i + 1))) for i in range(n)] + return out + + +def _counter(service, metric, base_rate, inc_rate, n=5): + """A single monotonic cumulative series: `base_rate`/s before the window, `inc_rate`/s after.""" + out, v = [], 0.0 + for i in range(n, 0, -1): + out.append(_M(service, metric, v, _W - timedelta(seconds=10 * i))) + v += base_rate * 10 + for i in range(1, n + 1): + out.append(_M(service, metric, v, _W + timedelta(seconds=10 * i))) + v += inc_rate * 10 + return out + + +class TestMetricClass: + def test_semantic_convention_shapes(self): + assert util_metric_class("jvm.cpu.recent_utilization") == ("cpu", "gauge") + assert util_metric_class("system.cpu.utilization") == ("cpu", "gauge") + assert util_metric_class("nodejs.eventloop.utilization") == ("eventloop", "gauge") + assert util_metric_class("system.memory.utilization") == ("memory", "gauge") + assert util_metric_class("process.cpu.time") == ("cpu", "cpu_time") + + def test_self_telemetry_and_other_metrics_are_not_classified(self): + assert util_metric_class("otel.sdk.span.started") is None + assert util_metric_class("otelcol_scraper_scraped_metric_points") is None + assert util_metric_class("process.memory.usage") is None + assert util_metric_class(None) is None + + +class TestGauges: + def test_saturation_is_present(self): + u = summarize_utilization(_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9), _W)[("ad", "cpu")] + assert u.sig_state == State.PRESENT and u.id == "util:ad:cpu" + + def test_normal_is_absent_and_a_drop_is_not_an_anomaly(self): + assert summarize_utilization(_gauge("ad", "jvm.cpu.recent_utilization", 0.2, 0.25), _W)[ + ("ad", "cpu")].sig_state == State.ABSENT + assert summarize_utilization(_gauge("fe", "nodejs.eventloop.utilization", 0.4, 0.1), _W)[ + ("fe", "eventloop")].sig_state == State.ABSENT + + def test_missing_window_or_zero_baseline_is_unknown(self): + only_incident = [m for m in _gauge("ad", "jvm.cpu.recent_utilization", 0.1, 0.9) if m.ts >= _W] + assert summarize_utilization(only_incident, _W)[("ad", "cpu")].sig_state is None + zero_base = _gauge("ad", "jvm.cpu.recent_utilization", 0.0, 0.9) + assert summarize_utilization(zero_base, _W)[("ad", "cpu")].sig_state is None + + def test_malformed_sample_leaves_the_metric_unmeasured(self): + samples = _gauge("ad", "jvm.cpu.recent_utilization", 0.1, 0.12) + samples.append(_M("ad", "jvm.cpu.recent_utilization", -1.0, _W + timedelta(seconds=99))) + assert summarize_utilization(samples, _W)[("ad", "cpu")].sig_state is None + + def test_unattributed_samples_are_not_evidence(self): + assert summarize_utilization(_gauge(None, "container.cpu.utilization", 0.01, 0.9), _W) == {} + + def test_self_telemetry_spike_is_not_evidence(self): + # the product-catalog trap: SDK counters rise because the service emits more spans. + assert summarize_utilization(_gauge("product-catalog", "otel.sdk.span.started", 10, 20), _W) == {} + + +class TestCpuTimeCounters: + def test_single_monotonic_series_rate_increase_is_present(self): + u = summarize_utilization(_counter("ad", "jvm.cpu.time", 0.01, 0.5), _W)[("ad", "cpu")] + assert u.sig_state == State.PRESENT and u.measured + + def test_steady_rate_is_absent(self): + assert summarize_utilization(_counter("ad", "jvm.cpu.time", 0.1, 0.11), _W)[ + ("ad", "cpu")].sig_state == State.ABSENT + + def test_interleaved_series_is_unknown(self): + # two flattened attribute series (e.g. state=user/system) share timestamps -> rate meaningless. + a = _counter("rec", "process.cpu.time", 0.1, 0.5) + b = [_M(m.service, m.metric, m.value * 3 + 7, m.ts) for m in _counter("rec", "process.cpu.time", 0.1, 0.1)] + assert summarize_utilization(a + b, _W)[("rec", "cpu")].sig_state is None + + def test_counter_reset_is_unknown(self): + samples = _counter("ad", "jvm.cpu.time", 0.1, 0.5) + samples.append(_M("ad", "jvm.cpu.time", 0.0, _W + timedelta(seconds=60))) # restart + assert summarize_utilization(samples, _W)[("ad", "cpu")].sig_state is None + + +class TestWitnessAcrossMetricsOfAClass: + def test_present_metric_wins_the_class(self): + samples = (_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9) + + _counter("ad", "jvm.cpu.time", 0.1, 0.11)) + u = summarize_utilization(samples, _W)[("ad", "cpu")] + assert u.sig_state == State.PRESENT and u.ratio > 2.0 + + +def _util(service, resource, state): + return UtilSignal(service, resource, 5.0 if state == State.PRESENT else 1.0, True, state) + + +class TestModel: + def test_present_util_generates_its_service_with_a_util_expectation(self): + hyps = build_hypotheses({}, set(), None, {("ad", "cpu"): _util("ad", "cpu", State.PRESENT)}) + (h,) = hyps + assert h.localization == "ad" + assert h.predictions == {"sig:ad": State.PRESENT, "util:ad:cpu": State.PRESENT} + + def test_absent_util_does_not_generate(self): + assert build_hypotheses({}, set(), None, {("ad", "cpu"): _util("ad", "cpu", State.ABSENT)}) == [] + + def test_without_util_signals_the_m2a_model_is_unchanged(self): + signals = {"checkout": ServiceSignal("checkout", 0.5, 1.0, True, True, State.PRESENT)} + (h,) = build_hypotheses(signals, set(), {}) + assert h.predictions == {"sig:checkout": State.PRESENT} + + def test_util_observable_is_its_own_coordinate(self): + obs = build_observables({"ad": ServiceSignal("ad", 0.0, 1.0, False, False, None)}, None, + {("ad", "cpu"): _util("ad", "cpu", State.PRESENT)}) + assert [o.id for o in obs] == ["util:ad:cpu"] # sig:ad stays UNKNOWN + + +class TestLocallySilentResourceFaultEndToEnd: + def _metrics(self): + # ad's requests look normal in traces, but its JVM CPU saturates. + return _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9) + _gauge("cart", "jvm.cpu.recent_utilization", 0.2, 0.2) + + def test_m2b_retains_the_saturated_service(self): + inp = structural_signals([], self._metrics(), _W) + res = resolve(partition( + build_hypotheses(inp.signals, inp.edges, inp.edge_signals, inp.util_signals), + build_observables(inp.signals, inp.edge_signals, inp.util_signals))) + assert res.localization == ("ad",) + + def test_without_util_signals_nothing_is_generated(self): + inp = structural_signals([], self._metrics(), _W) + assert build_hypotheses(inp.signals, inp.edges, inp.edge_signals) == [] From e68bcae21eba8cbcd1834b4cb5a5405ab9195727 Mon Sep 17 00:00:00 2001 From: Leonardo Araujo Date: Sat, 26 Sep 2026 07:31:09 -0300 Subject: [PATCH 2/4] fix(rca,ingest): utilization is measured only on one verified series (#212 review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Addresses both points of the #212 review. Our metric samples had no series identity anywhere in the pipeline, so M2b's "single-series" guarantee was not real: 1. BLOCKING — per-state gauges were averaged into OBSERVED ABSENT. system.cpu.utilization (per cpu.mode) / system.memory.utilization (per state): with no series key, a saturated CPU (idle 0.9->0.1, user 0.1->0.9) averaged to a flat 0.5 and was emitted as measured-normal. 2. MAJOR — the counter "single-series" check was not a series check, and the product path could not see collisions: metric_pk = scope|service|metric|ts with ON CONFLICT DO NOTHING silently DROPPED a second series at the same timestamp (whichever row landed first became "the" series); offset-timestamp series with a monotonic blend passed; metric_type was never read. Fix — series identity end to end: - Converter (parse_otlp_metrics) records each datapoint's series identity: datapoint attributes + the resource service.instance.id, always a dict for OTLP-derived samples ({} = attribute-free datapoint; None = never recorded). - metric_pk includes the canonical series identity when present, so two series at one timestamp are two rows (a general data-loss bug). Identity-less sources keep their original key, so their idempotent re-ingest is unchanged. - The jsonl corpus writer/loader round-trips attributes (omitted when None). - summarize_utilization measures a metric ONLY when every sample carries identity, they form exactly one series, no two samples share a timestamp, and the declared instrument (metric_type) is the one the reducer needs (gauge; cumulative counter for *.cpu.time — the name gives resource + required instrument, never the instrument itself). Anything else is UNKNOWN, never a measured ABSENT. A decrease inside a counter series is a reset -> UNKNOWN. Consequence (stated in the PR): the otel-fresh corpus carries no series identity, so M2b measures nothing on it and the earlier 6/8 claim is withdrawn — it rested on an unverified single-series assumption. M2b becomes measurable on M3's fresh capture with the fixed converter. Tests: util suite rewritten around series identity incl. the reviewer's per-mode CPU and offset-interleaved-counter cases, multiple instances, no-identity, instrument mismatch; converter/key/jsonl identity unit tests; product-path integration test (two series at one ts both persist -> UNKNOWN; single verified series through the DB -> measured; identity-less rows -> nothing claimed). Unit 1527, integration green. Co-Authored-By: Claude Opus 5.5 --- src/core/ingestion/telemetry.py | 11 +- src/core/rca/structural_model.py | 122 +++++++------ src/eval/otlp.py | 29 ++- src/eval/rcaeval.py | 6 + .../integration/test_util_series_identity.py | 96 ++++++++++ tests/unit/test_metric_series_identity.py | 84 +++++++++ tests/unit/test_structural_util_signals.py | 170 +++++++++++------- 7 files changed, 401 insertions(+), 117 deletions(-) create mode 100644 tests/integration/test_util_series_identity.py create mode 100644 tests/unit/test_metric_series_identity.py diff --git a/src/core/ingestion/telemetry.py b/src/core/ingestion/telemetry.py index 3a234fe..a3cab36 100644 --- a/src/core/ingestion/telemetry.py +++ b/src/core/ingestion/telemetry.py @@ -13,6 +13,7 @@ """ from __future__ import annotations +import json import uuid from dataclasses import dataclass from datetime import datetime @@ -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) return uuid.uuid5(_NS, key) diff --git a/src/core/rca/structural_model.py b/src/core/rca/structural_model.py index 881af1f..da93482 100644 --- a/src/core/rca/structural_model.py +++ b/src/core/rca/structural_model.py @@ -19,6 +19,7 @@ """ from __future__ import annotations +import json import math import statistics from collections import Counter, defaultdict @@ -403,13 +404,16 @@ def util_observable_id(service: str, resource: str) -> str: def util_metric_class(metric: Optional[str]) -> Optional[tuple[str, str]]: - """Classify a metric by OpenTelemetry semantic-convention **name shape** (never by names picked from - a corpus). Returns ``(resource, kind)`` or ``None``: + """Map a metric name to ``(resource, required_instrument)`` by OpenTelemetry semantic-convention + **name shape**, or ``None``. The name only says *which resource* the metric is about and *which + instrument it must be* for the reducer to apply; whether it actually **is** that instrument is read + from the sample's ``metric_type``, never inferred from the name: - * a utilization **gauge** — last dotted segment ends in ``utilization``; the resource is the segment - before it (``jvm.cpu.recent_utilization`` → ``cpu``, ``nodejs.eventloop.utilization`` → - ``eventloop``, ``system.memory.utilization`` → ``memory``); - * a cumulative **CPU-time counter** — name ends ``.cpu.time`` (semconv seconds); resource ``cpu``. + * last dotted segment ends in ``utilization`` → the preceding segment is the resource + (``jvm.cpu.recent_utilization`` → ``cpu``, ``nodejs.eventloop.utilization`` → ``eventloop``, + ``system.memory.utilization`` → ``memory``); required instrument ``"gauge"``; + * name ends ``.cpu.time`` → resource ``cpu``; required instrument ``"counter"`` (a cumulative + monotonic sum). Collector/SDK self-telemetry (``otel.sdk.*``, ``otelcol*``) is excluded: it measures the telemetry pipeline, not the service, and moves whenever a service simply emits more spans.""" @@ -419,7 +423,7 @@ def util_metric_class(metric: Optional[str]) -> Optional[tuple[str, str]]: if len(parts) >= 2 and parts[-1].endswith("utilization"): return parts[-2], "gauge" if metric.endswith(".cpu.time"): - return "cpu", "cpu_time" + return "cpu", "counter" return None @@ -430,57 +434,75 @@ def _counter_rate(points: list[tuple[datetime, float]]) -> Optional[float]: return (points[-1][1] - points[0][1]) / dt if dt > 0 else None +def _series_key(attributes) -> Optional[str]: + """Canonical series identity of a sample, or ``None`` when the source never recorded it.""" + if attributes is None: + return None + return json.dumps(attributes, sort_keys=True, default=str) + + +def _instrument(metric_type: Optional[str]) -> str: + """``MetricSample`` convention: a null ``metric_type`` is a gauge.""" + return metric_type or "gauge" + + def summarize_utilization(samples: list, window_start: datetime) -> dict[tuple[str, str], UtilSignal]: """Per ``(service, resource)`` utilization observable from metric-sample records (#209 M2b), duck-typed - on ``service`` / ``metric`` / ``value`` / ``ts``. - - Per metric (only classes recognized by :func:`util_metric_class`; unattributed ``service=None`` - samples contribute nothing): - - * **gauge** — incident mean / baseline mean; measured iff both windows have samples, every raw value - is finite and ≥ 0, and the baseline mean is > 0 (a malformed sample leaves the metric unmeasured, - never averaged into false-normal evidence — as in :func:`summarize_metrics`). - * **cpu_time** — incident rate / baseline rate, rate = (last − first) / Δt. Measured **only for a - verifiably single-series, monotonic stream**: at most one sample per timestamp and never decreasing - across the whole load. A converter that flattened attribute dimensions interleaves several - cumulative series (e.g. ``state=user``/``system``) into one stream, and a rate over that is - meaningless — so such a stream is UNKNOWN, not guessed. - - A metric is PRESENT if its ratio ≥ 2× (``discretize_ratio`` HIGH), ABSENT if measured and below, - else UNKNOWN. The ``(service, resource)`` signal takes the witness: PRESENT if any metric of that - class is PRESENT, else ABSENT if any is measured-normal, else UNKNOWN. A utilization *drop* is ABSENT - (not saturated), not an anomaly.""" - series: dict[tuple[str, str], dict[str, list[tuple[datetime, float]]]] = defaultdict( - lambda: {"b": [], "i": []} - ) - invalid: set[tuple[str, str]] = set() + on ``service`` / ``metric`` / ``value`` / ``ts`` / ``metric_type`` / ``attributes``. + + **A measurement requires one verified series.** A metric's samples form a series only when the + source recorded series identity (``attributes`` — datapoint attributes such as ``cpu.mode`` plus + the reporting instance). A metric is **measured only if** every one of its samples carries identity, + they all belong to **exactly one** series, no two samples share a timestamp, and the declared + instrument (``metric_type``) is the one the reducer requires. Otherwise it is UNKNOWN: + + * no identity — the samples cannot be told apart, so a per-state instrument + (``system.cpu.utilization`` per ``cpu.mode``, ``system.memory.utilization`` per state) or several + reporting instances would be averaged (or, after a lossy ingest, reduced to whichever row landed + first) into a number that is not one resource's utilization — never emitted as ABSENT; + * several series — there are no per-mode semantics yet, so the reducer neither averages nor picks; + * wrong instrument — a ``*.cpu.time`` that is not declared a cumulative counter is not rated. + + Reducers: a **gauge** compares incident mean / baseline mean (periodic aggregates; every raw value + finite and ≥ 0, baseline mean > 0); a **counter** compares incident rate / baseline rate, and a + decrease inside the series is a **reset**, so that series is UNKNOWN. Unattributed samples + (``service`` None) and self-telemetry contribute nothing. + + Per metric: PRESENT if ratio ≥ 2× (``discretize_ratio`` HIGH), ABSENT if measured and below, else + UNKNOWN. The ``(service, resource)`` signal takes the witness (PRESENT > ABSENT > UNKNOWN). A + utilization *drop* is ABSENT (not saturated), not an anomaly.""" + groups: dict[tuple[str, str], list] = defaultdict(list) for s in samples: - svc, metric, value, ts = (getattr(s, "service", None), getattr(s, "metric", None), - getattr(s, "value", None), getattr(s, "ts", None)) - if svc is None or value is None or ts is None or util_metric_class(metric) is None: + svc, metric = getattr(s, "service", None), getattr(s, "metric", None) + if svc is None or getattr(s, "value", None) is None or getattr(s, "ts", None) is None: continue - key = (svc, metric) - v = float(value) - if not math.isfinite(v) or v < 0: - invalid.add(key) + if util_metric_class(metric) is None: continue - series[key]["i" if ts >= window_start else "b"].append((ts, v)) + groups[(svc, metric)].append(s) per_class: dict[tuple[str, str], list[tuple[Optional[str], float, bool]]] = defaultdict(list) - for (svc, metric), d in series.items(): - resource, kind = util_metric_class(metric) + for (svc, metric), rows in groups.items(): + resource, required = util_metric_class(metric) + keys = {_series_key(getattr(r, "attributes", None)) for r in rows} + values = [(r.ts, float(r.value)) for r in rows] + single_series = None not in keys and len(keys) == 1 + distinct_ts = len({t for t, _ in values}) == len(values) + right_instrument = all(_instrument(getattr(r, "metric_type", None)) == required for r in rows) + well_formed = all(math.isfinite(v) and v >= 0 for _, v in values) + base = sorted(p for p in values if p[0] < window_start) + inc = sorted(p for p in values if p[0] >= window_start) + ratio, measured = 1.0, False - if (svc, metric) not in invalid and d["b"] and d["i"]: - if kind == "gauge": - base = _mean([v for _, v in d["b"]]) - if base > 0: - ratio, measured = _mean([v for _, v in d["i"]]) / base, True - else: # cpu_time: single-series + monotonic, or UNKNOWN - pts = sorted(d["b"] + d["i"]) - single = len({t for t, _ in pts}) == len(pts) - monotonic = all(b >= a for (_, a), (_, b) in zip(pts, pts[1:])) - rb, ri = _counter_rate(sorted(d["b"])), _counter_rate(sorted(d["i"])) - if single and monotonic and rb is not None and rb > 0 and ri is not None: + if single_series and distinct_ts and right_instrument and well_formed and base and inc: + if required == "gauge": + bmean = _mean([v for _, v in base]) + if bmean > 0: + ratio, measured = _mean([v for _, v in inc]) / bmean, True + else: + pts = sorted(values) + no_reset = all(b >= a for (_, a), (_, b) in zip(pts, pts[1:])) + rb, ri = _counter_rate(base), _counter_rate(inc) + if no_reset and rb is not None and rb > 0 and ri is not None: ratio, measured = ri / rb, True if not measured: state: Optional[str] = None diff --git a/src/eval/otlp.py b/src/eval/otlp.py index 04b0bda..3c9997c 100644 --- a/src/eval/otlp.py +++ b/src/eval/otlp.py @@ -39,6 +39,24 @@ def _service(resource: dict) -> Optional[str]: return _attr((resource or {}).get("attributes"), "service.name") +def _attr_value(v: dict): + for k in ("stringValue", "stringvalue", "intValue", "intvalue", "doubleValue", "doublevalue", + "boolValue", "boolvalue"): + if k in v: + return str(v[k]) + return None + + +def _attrs(attributes) -> dict[str, str]: + """Flatten an OTLP attribute list to ``{key: str(value)}`` (scalar values only).""" + out: dict[str, str] = {} + for a in attributes or []: + key, val = a.get("key"), _attr_value(a.get("value") or {}) + if key and val is not None: + out[key] = val + return out + + def _dt(unix_nano: object) -> Optional[datetime]: """OTLP timestamps are unsigned nanoseconds since epoch (JSON strings).""" if unix_nano is None: @@ -119,6 +137,7 @@ def parse_otlp_metrics(objs: list[dict], window: Window = None) -> list[ParsedMe for obj in objs: for rm in obj.get("resourceMetrics") or []: service = _service(rm.get("resource") or {}) + instance = _attr((rm.get("resource") or {}).get("attributes"), "service.instance.id") for sm in rm.get("scopeMetrics") or []: for metric in sm.get("metrics") or []: name = metric.get("name") @@ -132,8 +151,16 @@ def parse_otlp_metrics(objs: list[dict], window: Window = None) -> list[ParsedMe value = float(dp["count"]) if mtype == "histogram" and dp.get("count") is not None else _num(dp) if value is None: continue + # Series identity (#209 M2b): the datapoint attributes (e.g. cpu.mode, + # system.memory.state) plus the reporting instance. Always a dict for + # OTLP-derived samples — {} means "this datapoint had no attributes", + # which is different from a source that never recorded them (None). + series = _attrs(dp.get("attributes")) + if instance: + series["service.instance.id"] = instance samples.append(ParsedMetricSample( - service=service, metric=name, value=value, ts=ts, metric_type=mtype + service=service, metric=name, value=value, ts=ts, metric_type=mtype, + attributes=series, )) return samples diff --git a/src/eval/rcaeval.py b/src/eval/rcaeval.py index 4cc0dcd..f4fa164 100644 --- a/src/eval/rcaeval.py +++ b/src/eval/rcaeval.py @@ -387,6 +387,11 @@ def _metric_to_jsonl(m) -> dict: mtype = getattr(m, "metric_type", None) if mtype is not None: d["metric_type"] = mtype + # Series identity (#209 M2b): only emitted when the source recorded it, so identity-less + # sources stay byte-identical and a loaded sample never pretends to have identity it lacks. + attrs = getattr(m, "attributes", None) + if attrs is not None: + d["attributes"] = attrs return d @@ -434,6 +439,7 @@ def load_metrics_jsonl(path: Path) -> list: value=d.get("value"), ts=datetime.fromisoformat(ts) if ts else None, metric_type=d.get("metric_type"), + attributes=d.get("attributes"), )) return samples diff --git a/tests/integration/test_util_series_identity.py b/tests/integration/test_util_series_identity.py new file mode 100644 index 0000000..c38eabc --- /dev/null +++ b/tests/integration/test_util_series_identity.py @@ -0,0 +1,96 @@ +"""Integration: #209 M2b series identity on the PRODUCT path (the metric_samples rows +build_structural_view reads). Two series of one metric at the same timestamp must both persist, and a +per-mode gauge must come out UNKNOWN — never a measured ABSENT from whichever row landed first. +Skipped without Postgres.""" +from __future__ import annotations + +import os +from datetime import datetime, timedelta, timezone + +import pytest + +pytestmark = pytest.mark.skipif( + not os.getenv("DB_URL") and not os.getenv("INTEGRATION_TESTS"), + reason="Integration tests require DB_URL environment variable", +) + +W = datetime(2026, 1, 1, 12, 0, 0, tzinfo=timezone.utc) +BASE = W - timedelta(seconds=300) +END = W + timedelta(seconds=300) +SCOPE = "eval:util-series" + + +@pytest.fixture +def db_session(): + from sqlalchemy import text + + from src.db.models import Base + from src.db.session import check_connection, get_db, get_engine + + if not check_connection(): + pytest.skip("Cannot connect to database") + engine = get_engine() + with engine.connect() as conn: + conn.execute(text("CREATE EXTENSION IF NOT EXISTS vector")) + conn.commit() + Base.metadata.drop_all(engine) + Base.metadata.create_all(engine) + with get_db() as db: + yield db + + +def _series(service, metric, base, inc, attributes, metric_type="gauge"): + from src.core.ingestion.telemetry import ParsedMetricSample + out = [] + for i in range(1, 6): + out.append(ParsedMetricSample(service=service, metric=metric, value=base, ts=W - timedelta(seconds=30 * i), + metric_type=metric_type, attributes=attributes)) + out.append(ParsedMetricSample(service=service, metric=metric, value=inc, ts=W + timedelta(seconds=30 * i), + metric_type=metric_type, attributes=attributes)) + return out + + +def _rows(db): + from sqlalchemy import select + + from src.db.models import MetricSample + return db.execute(select(MetricSample).where(MetricSample.scope == SCOPE)).scalars().all() + + +def test_per_mode_gauge_persists_both_series_and_is_unknown(db_session): + from src.core.ingestion.telemetry import persist_metric_samples + from src.core.rca.structural_model import structural_signals + + # saturated CPU split by cpu.mode: idle 0.9 -> 0.1, user 0.1 -> 0.9 (both means stay 0.5) + inst = {"service.instance.id": "pod-1"} + persist_metric_samples(db_session, _series("host", "system.cpu.utilization", 0.9, 0.1, {"cpu.mode": "idle", **inst}) + + _series("host", "system.cpu.utilization", 0.1, 0.9, {"cpu.mode": "user", **inst}), + scope=SCOPE) + db_session.flush() + rows = _rows(db_session) + assert len(rows) == 20 # both series survive at every shared timestamp (none dropped on conflict) + + util = structural_signals([], rows, W).util_signals + assert util[("host", "cpu")].sig_state is None # UNKNOWN, not an averaged "measured normal" + + +def test_single_verified_series_through_the_db_is_measured(db_session): + from src.core.ingestion.telemetry import persist_metric_samples + from src.core.rca.structural import build_structural_view + + persist_metric_samples(db_session, _series("ad", "jvm.cpu.recent_utilization", 0.01, 0.9, + {"service.instance.id": "ad-1"}), scope=SCOPE) + db_session.flush() + built = build_structural_view(db_session, SCOPE, W, END, BASE) + assert built is not None + result, _ranking = built + assert "ad" in result.localization + + +def test_identity_less_rows_are_never_measured(db_session): + from src.core.ingestion.telemetry import persist_metric_samples + from src.core.rca.structural import build_structural_view + + persist_metric_samples(db_session, _series("ad", "jvm.cpu.recent_utilization", 0.01, 0.9, None), scope=SCOPE) + db_session.flush() + assert build_structural_view(db_session, SCOPE, W, END, BASE) is None # no evidence claimed diff --git a/tests/unit/test_metric_series_identity.py b/tests/unit/test_metric_series_identity.py new file mode 100644 index 0000000..dcf8763 --- /dev/null +++ b/tests/unit/test_metric_series_identity.py @@ -0,0 +1,84 @@ +"""#209 M2b — metric series identity through ingestion. Pure, no DB. + +The OTLP converter records each datapoint's series identity (datapoint attributes + the reporting +service.instance.id), the row key includes it so two series at one timestamp are two rows, and the +jsonl corpus format round-trips it — while identity-less sources keep their original key and format.""" +import json +from datetime import datetime, timezone +from pathlib import Path + +from src.core.ingestion.telemetry import ParsedMetricSample, metric_pk +from src.eval.otlp import parse_otlp_metrics +from src.eval.rcaeval import load_metrics_jsonl + +_TS = 1_767_268_800_000_000_000 # 2026-01-01T12:00:00Z + + +def _dp(value, **attrs): + return {"timeUnixNano": str(_TS), "asDouble": value, + "attributes": [{"key": k, "value": {"stringValue": v}} for k, v in attrs.items()]} + + +def _otlp(metric, instance="pod-1"): + resource = [{"key": "service.name", "value": {"stringValue": "host"}}] + if instance: + resource.append({"key": "service.instance.id", "value": {"stringValue": instance}}) + return [{"resourceMetrics": [{"resource": {"attributes": resource}, + "scopeMetrics": [{"metrics": [metric]}]}]}] + + +class TestConverterRecordsSeriesIdentity: + def test_per_mode_gauge_datapoints_are_distinct_series(self): + metric = {"name": "system.cpu.utilization", + "gauge": {"dataPoints": [_dp(0.9, **{"cpu.mode": "idle"}), _dp(0.1, **{"cpu.mode": "user"})]}} + samples = parse_otlp_metrics(_otlp(metric)) + assert [s.attributes for s in samples] == [ + {"cpu.mode": "idle", "service.instance.id": "pod-1"}, + {"cpu.mode": "user", "service.instance.id": "pod-1"}, + ] + assert {s.metric_type for s in samples} == {"gauge"} + + def test_attribute_free_datapoint_still_records_identity(self): + metric = {"name": "jvm.cpu.recent_utilization", "gauge": {"dataPoints": [_dp(0.5)]}} + (s,) = parse_otlp_metrics(_otlp(metric, instance=None)) + assert s.attributes == {} # "this datapoint had no attributes" — not None ("never recorded") + + def test_cumulative_monotonic_sum_is_declared_a_counter(self): + metric = {"name": "jvm.cpu.time", + "sum": {"isMonotonic": True, "aggregationTemporality": 2, "dataPoints": [_dp(12.0)]}} + (s,) = parse_otlp_metrics(_otlp(metric)) + assert s.metric_type == "counter" + + +class TestRowKey: + def _m(self, attributes): + return ParsedMetricSample(service="host", metric="system.cpu.utilization", value=0.5, + ts=datetime(2026, 1, 1, 12, tzinfo=timezone.utc), attributes=attributes) + + def test_two_series_at_one_timestamp_get_two_keys(self): + assert metric_pk("s", self._m({"cpu.mode": "idle"})) != metric_pk("s", self._m({"cpu.mode": "user"})) + + def test_key_is_independent_of_attribute_order(self): + assert (metric_pk("s", self._m({"a": "1", "b": "2"})) + == metric_pk("s", self._m({"b": "2", "a": "1"}))) + + def test_identity_less_sources_keep_their_original_key(self): + import uuid + + from src.core.ingestion.telemetry import _NS + m = self._m(None) + assert metric_pk("s", m) == uuid.uuid5(_NS, f"s|{m.service}|{m.metric}|{m.ts}") + + +class TestJsonlRoundTrip: + def test_identity_round_trips_and_is_omitted_when_unknown(self, tmp_path: Path): + from src.eval.rcaeval import _metric_to_jsonl as _metric_to_dict + ts = datetime(2026, 1, 1, 12, tzinfo=timezone.utc) + with_id = ParsedMetricSample(service="h", metric="m", value=1.0, ts=ts, metric_type="gauge", + attributes={"cpu.mode": "user"}) + without = ParsedMetricSample(service="h", metric="m", value=1.0, ts=ts) + assert "attributes" not in _metric_to_dict(without) + path = tmp_path / "metrics.jsonl" + path.write_text("\n".join(json.dumps(_metric_to_dict(m)) for m in (with_id, without)) + "\n") + a, b = load_metrics_jsonl(path) + assert a.attributes == {"cpu.mode": "user"} and b.attributes is None diff --git a/tests/unit/test_structural_util_signals.py b/tests/unit/test_structural_util_signals.py index d5897d7..cb6bb66 100644 --- a/tests/unit/test_structural_util_signals.py +++ b/tests/unit/test_structural_util_signals.py @@ -1,9 +1,10 @@ """#209 M2b — utilization observables util:{service}:{resource}. Pure, no DB. -Pins the frozen semantics: metrics are classified by OpenTelemetry semantic-convention name shape (not -by names picked from a corpus), collector/SDK self-telemetry is excluded, a cumulative CPU-time rate is -only trusted for a verifiably single-series monotonic stream, unattributed or malformed samples never -become evidence, and the observable family stays separate from sig:{service}.""" +Pins the semantics: a metric is measured only when it is ONE verified series — every sample carries +series identity (datapoint attributes + reporting instance), all samples belong to exactly one series, +no two share a timestamp, and the declared instrument (metric_type) is the one the reducer needs. +Anything else is UNKNOWN, never a measured ABSENT: in particular a per-state gauge +(system.cpu.utilization by cpu.mode) is never averaged into "normal".""" from datetime import datetime, timedelta, timezone from src.core.rca.observable import State @@ -20,38 +21,46 @@ ) _W = datetime(2026, 1, 1, 12, 0, 0, tzinfo=timezone.utc) +_ONE = {"service.instance.id": "pod-1"} # series identity of an attribute-free instrument class _M: - def __init__(self, service, metric, value, ts): + def __init__(self, service, metric, value, ts, *, attributes=None, metric_type=None): self.service, self.metric, self.value, self.ts = service, metric, value, ts + self.attributes, self.metric_type = attributes, metric_type -def _gauge(service, metric, base, inc, n=5): - out = [_M(service, metric, base, _W - timedelta(seconds=10 * (i + 1))) for i in range(n)] - out += [_M(service, metric, inc, _W + timedelta(seconds=10 * (i + 1))) for i in range(n)] +def _gauge(service, metric, base, inc, n=5, *, attributes=_ONE, offset=0): + kw = {"attributes": attributes, "metric_type": "gauge"} + out = [_M(service, metric, base, _W - timedelta(seconds=10 * (i + 1) + offset), **kw) for i in range(n)] + out += [_M(service, metric, inc, _W + timedelta(seconds=10 * (i + 1) + offset), **kw) for i in range(n)] return out -def _counter(service, metric, base_rate, inc_rate, n=5): - """A single monotonic cumulative series: `base_rate`/s before the window, `inc_rate`/s after.""" +def _counter(service, metric, base_rate, inc_rate, n=5, *, attributes=_ONE, metric_type="counter", offset=0): out, v = [], 0.0 for i in range(n, 0, -1): - out.append(_M(service, metric, v, _W - timedelta(seconds=10 * i))) + out.append(_M(service, metric, v, _W - timedelta(seconds=10 * i - offset), + attributes=attributes, metric_type=metric_type)) v += base_rate * 10 for i in range(1, n + 1): - out.append(_M(service, metric, v, _W + timedelta(seconds=10 * i))) + out.append(_M(service, metric, v, _W + timedelta(seconds=10 * i + offset), + attributes=attributes, metric_type=metric_type)) v += inc_rate * 10 return out +def _state(samples, key): + u = summarize_utilization(samples, _W).get(key) + return u.sig_state if u else "no-signal" + + class TestMetricClass: - def test_semantic_convention_shapes(self): + def test_name_shape_gives_resource_and_required_instrument(self): assert util_metric_class("jvm.cpu.recent_utilization") == ("cpu", "gauge") - assert util_metric_class("system.cpu.utilization") == ("cpu", "gauge") assert util_metric_class("nodejs.eventloop.utilization") == ("eventloop", "gauge") assert util_metric_class("system.memory.utilization") == ("memory", "gauge") - assert util_metric_class("process.cpu.time") == ("cpu", "cpu_time") + assert util_metric_class("process.cpu.time") == ("cpu", "counter") def test_self_telemetry_and_other_metrics_are_not_classified(self): assert util_metric_class("otel.sdk.span.started") is None @@ -60,61 +69,90 @@ def test_self_telemetry_and_other_metrics_are_not_classified(self): assert util_metric_class(None) is None -class TestGauges: +class TestOneVerifiedSeriesGauge: def test_saturation_is_present(self): - u = summarize_utilization(_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9), _W)[("ad", "cpu")] - assert u.sig_state == State.PRESENT and u.id == "util:ad:cpu" + assert _state(_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9), ("ad", "cpu")) == State.PRESENT def test_normal_is_absent_and_a_drop_is_not_an_anomaly(self): - assert summarize_utilization(_gauge("ad", "jvm.cpu.recent_utilization", 0.2, 0.25), _W)[ - ("ad", "cpu")].sig_state == State.ABSENT - assert summarize_utilization(_gauge("fe", "nodejs.eventloop.utilization", 0.4, 0.1), _W)[ - ("fe", "eventloop")].sig_state == State.ABSENT - - def test_missing_window_or_zero_baseline_is_unknown(self): - only_incident = [m for m in _gauge("ad", "jvm.cpu.recent_utilization", 0.1, 0.9) if m.ts >= _W] - assert summarize_utilization(only_incident, _W)[("ad", "cpu")].sig_state is None - zero_base = _gauge("ad", "jvm.cpu.recent_utilization", 0.0, 0.9) - assert summarize_utilization(zero_base, _W)[("ad", "cpu")].sig_state is None - - def test_malformed_sample_leaves_the_metric_unmeasured(self): - samples = _gauge("ad", "jvm.cpu.recent_utilization", 0.1, 0.12) - samples.append(_M("ad", "jvm.cpu.recent_utilization", -1.0, _W + timedelta(seconds=99))) - assert summarize_utilization(samples, _W)[("ad", "cpu")].sig_state is None - - def test_unattributed_samples_are_not_evidence(self): + assert _state(_gauge("ad", "jvm.cpu.recent_utilization", 0.2, 0.25), ("ad", "cpu")) == State.ABSENT + assert _state(_gauge("fe", "nodejs.eventloop.utilization", 0.4, 0.1), ("fe", "eventloop")) == State.ABSENT + + def test_per_mode_gauge_is_unknown_not_averaged_into_absent(self): + # The reviewer's case: cpu.mode idle 0.9->0.1 and user 0.1->0.9. Averaged, both windows read 0.5 + # and the saturated CPU would be "measured normal". Two series -> UNKNOWN. + samples = (_gauge("h", "system.cpu.utilization", 0.9, 0.1, attributes={"cpu.mode": "idle", **_ONE}) + + _gauge("h", "system.cpu.utilization", 0.1, 0.9, attributes={"cpu.mode": "user", **_ONE})) + assert _state(samples, ("h", "cpu")) is None + assert build_observables({}, None, summarize_utilization(samples, _W)) == [] + + def test_per_state_memory_gauge_is_unknown(self): + samples = (_gauge("h", "system.memory.utilization", 0.3, 0.9, attributes={"system.memory.state": "used"}) + + _gauge("h", "system.memory.utilization", 0.7, 0.1, attributes={"system.memory.state": "free"})) + assert _state(samples, ("h", "memory")) is None + + def test_two_reporting_instances_are_unknown(self): + samples = (_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9, attributes={"service.instance.id": "a"}) + + _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.02, + attributes={"service.instance.id": "b"}, offset=3)) + assert _state(samples, ("ad", "cpu")) is None + + def test_samples_without_series_identity_are_never_measured(self): + assert _state(_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9, attributes=None), ("ad", "cpu")) is None + + def test_same_timestamp_collision_within_a_series_is_unknown(self): + samples = _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9) + samples.append(_M("ad", "jvm.cpu.recent_utilization", 0.02, samples[0].ts, + attributes=_ONE, metric_type="gauge")) + assert _state(samples, ("ad", "cpu")) is None + + def test_declared_non_gauge_utilization_is_unknown(self): + samples = _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9) + for s in samples: + s.metric_type = "sum" + assert _state(samples, ("ad", "cpu")) is None + + def test_malformed_or_zero_baseline_is_unknown(self): + bad = _gauge("ad", "jvm.cpu.recent_utilization", 0.1, 0.12) + bad.append(_M("ad", "jvm.cpu.recent_utilization", -1.0, _W + timedelta(seconds=99), + attributes=_ONE, metric_type="gauge")) + assert _state(bad, ("ad", "cpu")) is None + assert _state(_gauge("ad", "jvm.cpu.recent_utilization", 0.0, 0.9), ("ad", "cpu")) is None + + def test_unattributed_and_self_telemetry_samples_are_not_evidence(self): assert summarize_utilization(_gauge(None, "container.cpu.utilization", 0.01, 0.9), _W) == {} - - def test_self_telemetry_spike_is_not_evidence(self): - # the product-catalog trap: SDK counters rise because the service emits more spans. - assert summarize_utilization(_gauge("product-catalog", "otel.sdk.span.started", 10, 20), _W) == {} + assert summarize_utilization(_gauge("pc", "otel.sdk.span.started", 10, 20), _W) == {} -class TestCpuTimeCounters: - def test_single_monotonic_series_rate_increase_is_present(self): - u = summarize_utilization(_counter("ad", "jvm.cpu.time", 0.01, 0.5), _W)[("ad", "cpu")] - assert u.sig_state == State.PRESENT and u.measured +class TestOneVerifiedSeriesCounter: + def test_declared_cumulative_counter_rate_increase_is_present(self): + assert _state(_counter("ad", "jvm.cpu.time", 0.01, 0.5), ("ad", "cpu")) == State.PRESENT def test_steady_rate_is_absent(self): - assert summarize_utilization(_counter("ad", "jvm.cpu.time", 0.1, 0.11), _W)[ - ("ad", "cpu")].sig_state == State.ABSENT - - def test_interleaved_series_is_unknown(self): - # two flattened attribute series (e.g. state=user/system) share timestamps -> rate meaningless. - a = _counter("rec", "process.cpu.time", 0.1, 0.5) - b = [_M(m.service, m.metric, m.value * 3 + 7, m.ts) for m in _counter("rec", "process.cpu.time", 0.1, 0.1)] - assert summarize_utilization(a + b, _W)[("rec", "cpu")].sig_state is None - - def test_counter_reset_is_unknown(self): + assert _state(_counter("ad", "jvm.cpu.time", 0.1, 0.11), ("ad", "cpu")) == State.ABSENT + + def test_name_shape_alone_does_not_make_a_counter(self): + # a *.cpu.time not declared a cumulative counter (null = gauge, or a delta "sum") is not rated. + assert _state(_counter("ad", "jvm.cpu.time", 0.01, 0.5, metric_type=None), ("ad", "cpu")) is None + assert _state(_counter("ad", "jvm.cpu.time", 0.01, 0.5, metric_type="sum"), ("ad", "cpu")) is None + + def test_offset_interleaved_series_with_monotonic_blend_is_unknown(self): + # The reviewer's case: two cumulative series on offset timestamps whose merged sequence happens to + # be non-decreasing. With identity they are two series -> UNKNOWN, not a blended rate. + samples = (_counter("r", "process.cpu.time", 0.1, 0.5, attributes={"cpu.mode": "user", **_ONE}) + + _counter("r", "process.cpu.time", 0.1, 0.5, attributes={"cpu.mode": "system", **_ONE}, + offset=5)) + assert _state(samples, ("r", "cpu")) is None + + def test_reset_is_unknown(self): samples = _counter("ad", "jvm.cpu.time", 0.1, 0.5) - samples.append(_M("ad", "jvm.cpu.time", 0.0, _W + timedelta(seconds=60))) # restart - assert summarize_utilization(samples, _W)[("ad", "cpu")].sig_state is None + samples.append(_M("ad", "jvm.cpu.time", 0.0, _W + timedelta(seconds=61), + attributes=_ONE, metric_type="counter")) + assert _state(samples, ("ad", "cpu")) is None class TestWitnessAcrossMetricsOfAClass: def test_present_metric_wins_the_class(self): - samples = (_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9) - + _counter("ad", "jvm.cpu.time", 0.1, 0.11)) + samples = _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9) + _counter("ad", "jvm.cpu.time", 0.1, 0.11) u = summarize_utilization(samples, _W)[("ad", "cpu")] assert u.sig_state == State.PRESENT and u.ratio > 2.0 @@ -125,8 +163,7 @@ def _util(service, resource, state): class TestModel: def test_present_util_generates_its_service_with_a_util_expectation(self): - hyps = build_hypotheses({}, set(), None, {("ad", "cpu"): _util("ad", "cpu", State.PRESENT)}) - (h,) = hyps + (h,) = build_hypotheses({}, set(), None, {("ad", "cpu"): _util("ad", "cpu", State.PRESENT)}) assert h.localization == "ad" assert h.predictions == {"sig:ad": State.PRESENT, "util:ad:cpu": State.PRESENT} @@ -141,13 +178,13 @@ def test_without_util_signals_the_m2a_model_is_unchanged(self): def test_util_observable_is_its_own_coordinate(self): obs = build_observables({"ad": ServiceSignal("ad", 0.0, 1.0, False, False, None)}, None, {("ad", "cpu"): _util("ad", "cpu", State.PRESENT)}) - assert [o.id for o in obs] == ["util:ad:cpu"] # sig:ad stays UNKNOWN + assert [o.id for o in obs] == ["util:ad:cpu"] class TestLocallySilentResourceFaultEndToEnd: def _metrics(self): - # ad's requests look normal in traces, but its JVM CPU saturates. - return _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9) + _gauge("cart", "jvm.cpu.recent_utilization", 0.2, 0.2) + return (_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9) + + _gauge("cart", "jvm.cpu.recent_utilization", 0.2, 0.2)) def test_m2b_retains_the_saturated_service(self): inp = structural_signals([], self._metrics(), _W) @@ -156,6 +193,9 @@ def test_m2b_retains_the_saturated_service(self): build_observables(inp.signals, inp.edge_signals, inp.util_signals))) assert res.localization == ("ad",) - def test_without_util_signals_nothing_is_generated(self): - inp = structural_signals([], self._metrics(), _W) - assert build_hypotheses(inp.signals, inp.edges, inp.edge_signals) == [] + def test_without_series_identity_nothing_is_claimed(self): + # the same evidence from a source that never recorded series identity: no util observable at all. + stripped = [_M(m.service, m.metric, m.value, m.ts, metric_type="gauge") for m in self._metrics()] + inp = structural_signals([], stripped, _W) + assert inp.util_signals[("ad", "cpu")].sig_state is None + assert build_hypotheses(inp.signals, inp.edges, inp.edge_signals, inp.util_signals) == [] From 63edde4d0cb5411a5e0340a18512aed5a5c44a7a Mon Sep 17 00:00:00 2001 From: Leonardo Araujo Date: Sat, 26 Sep 2026 08:23:42 -0300 Subject: [PATCH 3/4] fix(rca,ingest): name the instance, and reduce every metric consumer per series (#212 review r2) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Both points of the second #212 review: 1. BLOCKING — an empty identity was treated as a verified series. The converter stamped {} when a datapoint had no attributes and the resource no instance id, and summarize_utilization accepted it: two unnamed replicas on offset timestamps (A 0.40->0.45, B 0.40->0.90) blended to 1.69x and read as measured-normal while B was saturated. Now a series is verified only if its instance is NAMED. New shared module src/core/rca/metric_series defines identity once: series_key, datapoint signature, and instance_identity from the semconv resource attributes that name one process (service.instance.id, k8s.pod.uid, container.id, or host.name + process.pid). The converter records whichever exist. summarize_utilization: unnamed instance -> UNKNOWN; more than one datapoint signature (modes, e.g. cpu.mode idle/user) -> UNKNOWN; under one mode each named instance is measured separately — PRESENT if any replica is saturated, ABSENT only if every replica is measured-normal. 2. MAJOR — other reducers still averaged the series metric_pk now retains. metric_anomaly_onsets (trigger path; _metric_onsets now also SELECTs attributes), metric_features (ranker features) and summarize_metrics (structural sig) now group by (service, metric, series_key): features take the max over series, onsets are detected per series, the sig branches are anomalous if ANY series is and measured-normal only if EVERY series is. Identity-less samples form one group per (service, metric), exactly as before. Eval delta (measured, not assumed): raglogs eval on trace-loc with the learned ranker + rare_event trigger path, main vs this branch: 0/24 cases differ in root cause, predicted services, top trigger or confidence; aggregates identical. All existing corpora carry no series identity, so the default path is byte-identical on them; series-aware reduction applies to identity-carrying captures. Tests: reviewer's unnamed-replica case (unit + from the DB), named replicas (calm+pegged -> PRESENT, both calm -> ABSENT, one unmeasured -> UNKNOWN), each semconv instance attribute (host.name alone does not name), per-series reducer tests with identity-less equivalence, and a DB test that _metric_onsets and compute_features see every cpu.mode series. Unit 1538, integration 52 green. Co-Authored-By: Claude Opus 5.5 --- src/core/explain/evidence.py | 3 +- src/core/rca/features.py | 17 +- src/core/rca/metric_series.py | 65 +++++++ src/core/rca/structural_model.py | 181 ++++++++++-------- src/core/rca/triggers.py | 10 +- src/eval/otlp.py | 14 +- .../integration/test_util_series_identity.py | 44 +++++ tests/unit/test_metric_reducers_per_series.py | 82 ++++++++ tests/unit/test_metric_series_identity.py | 18 +- tests/unit/test_structural_util_signals.py | 36 +++- 10 files changed, 368 insertions(+), 102 deletions(-) create mode 100644 src/core/rca/metric_series.py create mode 100644 tests/unit/test_metric_reducers_per_series.py diff --git a/src/core/explain/evidence.py b/src/core/explain/evidence.py index 2a15459..ebf1a7d 100644 --- a/src/core/explain/evidence.py +++ b/src/core/explain/evidence.py @@ -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, diff --git a/src/core/rca/features.py b/src/core/rca/features.py index bf7856a..b34d80f 100644 --- a/src/core/rca/features.py +++ b/src/core/rca/features.py @@ -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: diff --git a/src/core/rca/metric_series.py b/src/core/rca/metric_series.py new file mode 100644 index 0000000..82d8c2c --- /dev/null +++ b/src/core/rca/metric_series.py @@ -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 diff --git a/src/core/rca/structural_model.py b/src/core/rca/structural_model.py index da93482..6cb045e 100644 --- a/src/core/rca/structural_model.py +++ b/src/core/rca/structural_model.py @@ -19,7 +19,6 @@ """ from __future__ import annotations -import json import math import statistics from collections import Counter, defaultdict @@ -93,14 +92,22 @@ def summarize_metrics(samples: list, window_start: datetime) -> dict[str, Servic ``error_rate`` exists; a latency branch only when both incident **and** baseline ``latency_ms`` exist (a ratio needs both). ``sig`` is ``PRESENT`` if a measured branch is anomalous, ``ABSENT`` only if both branches are measured-and-normal, else UNKNOWN. Records are duck-typed on - ``service`` / ``value`` / ``ts`` / ``metric``.""" - inc: dict[str, dict[str, list[float]]] = defaultdict(lambda: defaultdict(list)) - base: dict[str, dict[str, list[float]]] = defaultdict(lambda: defaultdict(list)) + ``service`` / ``value`` / ``ts`` / ``metric`` (and optionally ``attributes``). + + Each branch is reduced **per series** — ``(service, metric, series_key(attributes))`` (#209 M2b) — + never averaged across the series of one instrument (per route, per replica, …): a branch is + anomalous if **any** series is, and measured-normal only if **every** series is measured and + normal. Samples without recorded identity form one series per ``(service, metric)``, so for them + this is exactly the single-series computation.""" + from src.core.rca.metric_series import series_key + + inc: dict[str, dict[tuple, list[float]]] = defaultdict(lambda: defaultdict(list)) + base: dict[str, dict[tuple, list[float]]] = defaultdict(lambda: defaultdict(list)) for s in samples: if s.service is None or s.value is None or s.ts is None: continue bucket = inc if s.ts >= window_start else base - bucket[s.service][s.metric].append(float(s.value)) + bucket[s.service][(s.metric, series_key(getattr(s, "attributes", None)))].append(float(s.value)) def mean(xs: list[float]) -> float: return sum(xs) / len(xs) if xs else 0.0 @@ -113,31 +120,34 @@ def all_nonneg(xs: list[float]) -> bool: signals: dict[str, ServiceSignal] = {} for svc in sorted(set(inc) | set(base)): - inc_err = inc[svc].get("error_rate", []) - inc_lat = inc[svc].get("latency_ms", []) - base_lat = base[svc].get("latency_ms", []) - base_lat_mean = mean(base_lat) - + keys = set(inc[svc]) | set(base[svc]) # Validate every RAW sample, not just the aggregate — an in-range mean does not prove valid # inputs (e.g. incident latency [-10, 30] averages to a normal-looking 10). Any malformed - # sample leaves that branch UNKNOWN (unmeasured), never averaged into false-normal evidence. - # A rate must be finite in [0,1]; a latency needs valid non-negative incident + positive - # baseline samples (30/0 is UNKNOWN, not normal). - error_measured = bool(inc_err) and all_valid_rates(inc_err) - latency_measured = ( - bool(inc_lat) and all_nonneg(inc_lat) - and bool(base_lat) and all_nonneg(base_lat) and base_lat_mean > 0 - ) - err = mean(inc_err) - ratio = mean(inc_lat) / base_lat_mean if latency_measured else 1.0 - - error_present = error_measured and discretize_rate(err) == State.PRESENT - latency_high = latency_measured and discretize_ratio(ratio) == State.HIGH - if error_present or latency_high: # a measured branch proves the anomaly + # sample leaves that series unmeasured, never averaged into false-normal evidence. A rate must + # be finite in [0,1]; a latency needs valid non-negative incident + positive baseline samples + # (30/0 is UNKNOWN, not normal). + err_series: list[tuple[bool, float]] = [] # (measured, rate) + for key in sorted(k for k in keys if k[0] == "error_rate"): + xs = inc[svc].get(key, []) + err_series.append((bool(xs) and all_valid_rates(xs), mean(xs))) + lat_series: list[tuple[bool, float]] = [] # (measured, ratio) + for key in sorted(k for k in keys if k[0] == "latency_ms"): + ix, bx = inc[svc].get(key, []), base[svc].get(key, []) + ok = bool(ix) and all_nonneg(ix) and bool(bx) and all_nonneg(bx) and mean(bx) > 0 + lat_series.append((ok, mean(ix) / mean(bx) if ok else 1.0)) + + err_hit = [r for m, r in err_series if m and discretize_rate(r) == State.PRESENT] + lat_hit = [r for m, r in lat_series if m and discretize_ratio(r) == State.HIGH] + error_measured = bool(err_series) and all(m for m, _ in err_series) + latency_measured = bool(lat_series) and all(m for m, _ in lat_series) + err = max(err_hit) if err_hit else max((r for m, r in err_series if m), default=0.0) + ratio = max(lat_hit) if lat_hit else max((r for m, r in lat_series if m), default=1.0) + + if err_hit or lat_hit: # a measured series proves the anomaly sig_state: Optional[str] = State.PRESENT - elif error_measured and latency_measured: # both measured and normal -> proven absent + elif error_measured and latency_measured: # every series measured and normal -> proven absent sig_state = State.ABSENT - else: # some branch unmeasured, nothing proves present + else: # some series unmeasured, nothing proves present sig_state = None signals[svc] = ServiceSignal(svc, err, ratio, error_measured, latency_measured, sig_state) return signals @@ -434,43 +444,61 @@ def _counter_rate(points: list[tuple[datetime, float]]) -> Optional[float]: return (points[-1][1] - points[0][1]) / dt if dt > 0 else None -def _series_key(attributes) -> Optional[str]: - """Canonical series identity of a sample, or ``None`` when the source never recorded it.""" - if attributes is None: - return None - return json.dumps(attributes, sort_keys=True, default=str) - - def _instrument(metric_type: Optional[str]) -> str: """``MetricSample`` convention: a null ``metric_type`` is a gauge.""" return metric_type or "gauge" +def _measure_series(points: list[tuple[datetime, float]], required: str, + window_start: datetime) -> tuple[Optional[str], float, bool]: + """``(state, ratio, measured)`` for ONE fully identified series (a single instance of a single mode).""" + values = sorted(points) + base = [p for p in values if p[0] < window_start] + inc = [p for p in values if p[0] >= window_start] + ok = (len({t for t, _ in values}) == len(values) # no timestamp collision in a series + and all(math.isfinite(v) and v >= 0 for _, v in values) # well-formed + and bool(base) and bool(inc)) + ratio, measured = 1.0, False + if ok and required == "gauge": + bmean = _mean([v for _, v in base]) + if bmean > 0: + ratio, measured = _mean([v for _, v in inc]) / bmean, True + elif ok: # counter: a decrease inside the series is a reset -> not rated + no_reset = all(b >= a for (_, a), (_, b) in zip(values, values[1:])) + rb, ri = _counter_rate(base), _counter_rate(inc) + if no_reset and rb is not None and rb > 0 and ri is not None: + ratio, measured = ri / rb, True + if not measured: + return None, ratio, False + return (State.PRESENT if discretize_ratio(ratio) == State.HIGH else State.ABSENT), ratio, True + + def summarize_utilization(samples: list, window_start: datetime) -> dict[tuple[str, str], UtilSignal]: """Per ``(service, resource)`` utilization observable from metric-sample records (#209 M2b), duck-typed on ``service`` / ``metric`` / ``value`` / ``ts`` / ``metric_type`` / ``attributes``. - **A measurement requires one verified series.** A metric's samples form a series only when the - source recorded series identity (``attributes`` — datapoint attributes such as ``cpu.mode`` plus - the reporting instance). A metric is **measured only if** every one of its samples carries identity, - they all belong to **exactly one** series, no two samples share a timestamp, and the declared - instrument (``metric_type``) is the one the reducer requires. Otherwise it is UNKNOWN: - - * no identity — the samples cannot be told apart, so a per-state instrument - (``system.cpu.utilization`` per ``cpu.mode``, ``system.memory.utilization`` per state) or several - reporting instances would be averaged (or, after a lossy ingest, reduced to whichever row landed - first) into a number that is not one resource's utilization — never emitted as ABSENT; - * several series — there are no per-mode semantics yet, so the reducer neither averages nor picks; - * wrong instrument — a ``*.cpu.time`` that is not declared a cumulative counter is not rated. - - Reducers: a **gauge** compares incident mean / baseline mean (periodic aggregates; every raw value - finite and ≥ 0, baseline mean > 0); a **counter** compares incident rate / baseline rate, and a - decrease inside the series is a **reset**, so that series is UNKNOWN. Unattributed samples - (``service`` None) and self-telemetry contribute nothing. - - Per metric: PRESENT if ratio ≥ 2× (``discretize_ratio`` HIGH), ABSENT if measured and below, else - UNKNOWN. The ``(service, resource)`` signal takes the witness (PRESENT > ABSENT > UNKNOWN). A - utilization *drop* is ABSENT (not saturated), not an anomaly.""" + **A measurement requires fully identified series** (see :mod:`src.core.rca.metric_series`). For one + ``(service, metric)``, the metric is UNKNOWN — never a measured ABSENT — unless: + + * every sample carries recorded identity (``attributes`` is not None); + * every sample **names its reporting instance** (``service.instance.id``, ``k8s.pod.uid``, + ``container.id``, or ``host.name`` + ``process.pid``) — an empty or instance-less identity cannot + tell two replicas apart, so they would be blended (a calm replica plus a pegged one reading as + normal); + * all samples share **one datapoint signature** — several signatures are several *modes* of the + instrument (``cpu.mode`` idle/user, ``system.memory.state`` used/free) whose direction differs, + and there are no per-mode semantics yet, so the reducer neither averages nor picks; + * the declared instrument (``metric_type``) is the one the reducer needs (gauge; cumulative counter + for ``*.cpu.time``) — the name gives the resource, never the instrument. + + Under one mode, each named instance is its own series, measured separately (:func:`_measure_series`: + gauge incident/baseline mean ratio, or counter incident/baseline rate ratio with reset detection; no + timestamp collisions). The metric is PRESENT if **any** instance is saturated (ratio ≥ 2×), ABSENT + only if **every** instance is measured and below, else UNKNOWN. A utilization *drop* is ABSENT (not + saturated). Unattributed samples (``service`` None) and self-telemetry contribute nothing. The + ``(service, resource)`` signal takes the witness metric (PRESENT > ABSENT > UNKNOWN).""" + from src.core.rca.metric_series import datapoint_signature, instance_identity + groups: dict[tuple[str, str], list] = defaultdict(list) for s in samples: svc, metric = getattr(s, "service", None), getattr(s, "metric", None) @@ -483,34 +511,25 @@ def summarize_utilization(samples: list, window_start: datetime) -> dict[tuple[s per_class: dict[tuple[str, str], list[tuple[Optional[str], float, bool]]] = defaultdict(list) for (svc, metric), rows in groups.items(): resource, required = util_metric_class(metric) - keys = {_series_key(getattr(r, "attributes", None)) for r in rows} - values = [(r.ts, float(r.value)) for r in rows] - single_series = None not in keys and len(keys) == 1 - distinct_ts = len({t for t, _ in values}) == len(values) - right_instrument = all(_instrument(getattr(r, "metric_type", None)) == required for r in rows) - well_formed = all(math.isfinite(v) and v >= 0 for _, v in values) - base = sorted(p for p in values if p[0] < window_start) - inc = sorted(p for p in values if p[0] >= window_start) - - ratio, measured = 1.0, False - if single_series and distinct_ts and right_instrument and well_formed and base and inc: - if required == "gauge": - bmean = _mean([v for _, v in base]) - if bmean > 0: - ratio, measured = _mean([v for _, v in inc]) / bmean, True - else: - pts = sorted(values) - no_reset = all(b >= a for (_, a), (_, b) in zip(pts, pts[1:])) - rb, ri = _counter_rate(base), _counter_rate(inc) - if no_reset and rb is not None and rb > 0 and ri is not None: - ratio, measured = ri / rb, True - if not measured: - state: Optional[str] = None - elif discretize_ratio(ratio) == State.HIGH: - state = State.PRESENT + attrs = [getattr(r, "attributes", None) for r in rows] + instances = [instance_identity(a) for a in attrs] + identified = (all(a is not None for a in attrs) and None not in instances + and len({datapoint_signature(a) for a in attrs}) == 1 + and all(_instrument(getattr(r, "metric_type", None)) == required for r in rows)) + if not identified: + per_class[(svc, resource)].append((None, 1.0, False)) + continue + by_instance: dict[str, list[tuple[datetime, float]]] = defaultdict(list) + for r, inst in zip(rows, instances): + by_instance[inst].append((r.ts, float(r.value))) + results = [_measure_series(pts, required, window_start) for _, pts in sorted(by_instance.items())] + hit = [x for x in results if x[0] == State.PRESENT] + if hit: + per_class[(svc, resource)].append(max(hit, key=lambda x: x[1])) + elif all(m for _, _, m in results): + per_class[(svc, resource)].append((State.ABSENT, max(x[1] for x in results), True)) else: - state = State.ABSENT - per_class[(svc, resource)].append((state, ratio, measured)) + per_class[(svc, resource)].append((None, 1.0, False)) out: dict[tuple[str, str], UtilSignal] = {} for (svc, resource), metrics in sorted(per_class.items()): diff --git a/src/core/rca/triggers.py b/src/core/rca/triggers.py index 38924d9..051bd0b 100644 --- a/src/core/rca/triggers.py +++ b/src/core/rca/triggers.py @@ -136,13 +136,19 @@ def metric_anomaly_onsets( """Earliest per-(service, metric) deviation from baseline, as onset candidates. ``samples`` are duck-typed rows with ``service`` / ``metric`` / ``value`` / - ``ts``. Baseline = points with ``ts < incident_start`` (needs ``min_baseline`` + ``ts`` (and optionally ``attributes``). A series is ``(service, metric, + series_key(attributes))`` (#209 M2b): several series of one instrument (per + ``cpu.mode``, per replica, …) are detected separately, never averaged into a + flat line; samples without recorded identity form one series per + ``(service, metric)``, as before. Baseline = points with ``ts < incident_start`` (needs ``min_baseline`` of them); a point in ``[incident_start, incident_end]`` is anomalous when it is more than ``z_threshold`` baseline sigmas from the baseline mean, or — when the baseline is flat (sigma ~ 0) — more than ``rel_threshold`` of |mean| away. The earliest anomalous point per series is its onset; ranked by earliness x magnitude. This surfaces the injection time on faults that never log. """ + from src.core.rca.metric_series import series_key + baseline: dict[tuple, list[float]] = {} incident: dict[tuple, list[tuple[datetime, float]]] = {} for m in samples: @@ -152,7 +158,7 @@ def metric_anomaly_onsets( value = getattr(m, "value", None) if metric is None or ts is None or value is None: continue - key = (service, metric) + key = (service, metric, series_key(getattr(m, "attributes", None))) if ts < incident_start: baseline.setdefault(key, []).append(float(value)) elif incident_start <= ts <= incident_end: diff --git a/src/eval/otlp.py b/src/eval/otlp.py index 3c9997c..2e2c0c0 100644 --- a/src/eval/otlp.py +++ b/src/eval/otlp.py @@ -21,6 +21,7 @@ from typing import Optional from src.core.ingestion.telemetry import ParsedMetricSample, ParsedSpan +from src.core.rca.metric_series import instance_attributes from src.eval.rcaeval import _infer_level Window = Optional[tuple[datetime, datetime]] @@ -137,7 +138,7 @@ def parse_otlp_metrics(objs: list[dict], window: Window = None) -> list[ParsedMe for obj in objs: for rm in obj.get("resourceMetrics") or []: service = _service(rm.get("resource") or {}) - instance = _attr((rm.get("resource") or {}).get("attributes"), "service.instance.id") + instance = instance_attributes(_attrs((rm.get("resource") or {}).get("attributes"))) for sm in rm.get("scopeMetrics") or []: for metric in sm.get("metrics") or []: name = metric.get("name") @@ -152,12 +153,11 @@ def parse_otlp_metrics(objs: list[dict], window: Window = None) -> list[ParsedMe if value is None: continue # Series identity (#209 M2b): the datapoint attributes (e.g. cpu.mode, - # system.memory.state) plus the reporting instance. Always a dict for - # OTLP-derived samples — {} means "this datapoint had no attributes", - # which is different from a source that never recorded them (None). - series = _attrs(dp.get("attributes")) - if instance: - series["service.instance.id"] = instance + # system.memory.state) plus whichever resource attributes name the + # reporting instance (see metric_series). Always a dict for OTLP-derived + # samples; None is reserved for sources that never recorded identity. A + # dict with no instance attribute ({} included) is NOT a verified series. + series = {**_attrs(dp.get("attributes")), **instance} samples.append(ParsedMetricSample( service=service, metric=name, value=value, ts=ts, metric_type=mtype, attributes=series, diff --git a/tests/integration/test_util_series_identity.py b/tests/integration/test_util_series_identity.py index c38eabc..4f11c99 100644 --- a/tests/integration/test_util_series_identity.py +++ b/tests/integration/test_util_series_identity.py @@ -94,3 +94,47 @@ def test_identity_less_rows_are_never_measured(db_session): persist_metric_samples(db_session, _series("ad", "jvm.cpu.recent_utilization", 0.01, 0.9, None), scope=SCOPE) db_session.flush() assert build_structural_view(db_session, SCOPE, W, END, BASE) is None # no evidence claimed + + +def _replicas_without_instance(offset_s=5): + """Two replicas whose resource names no instance (identity {}), offset timestamps: A calm, B pegged.""" + from src.core.ingestion.telemetry import ParsedMetricSample + out = [] + for value_base, value_inc, off in ((0.40, 0.45, 0), (0.40, 0.90, offset_s)): + for i in range(1, 6): + out.append(ParsedMetricSample(service="ad", metric="jvm.cpu.recent_utilization", value=value_base, + ts=W - timedelta(seconds=30 * i - off), metric_type="gauge", + attributes={})) + out.append(ParsedMetricSample(service="ad", metric="jvm.cpu.recent_utilization", value=value_inc, + ts=W + timedelta(seconds=30 * i + off), metric_type="gauge", + attributes={})) + return out + + +def test_replicas_without_a_named_instance_are_unknown_from_the_db(db_session): + from src.core.ingestion.telemetry import persist_metric_samples + from src.core.rca.structural_model import structural_signals + + persist_metric_samples(db_session, _replicas_without_instance(), scope=SCOPE) + db_session.flush() + rows = _rows(db_session) + assert len(rows) == 20 # both replicas stored (distinct timestamps) + assert structural_signals([], rows, W).util_signals[("ad", "cpu")].sig_state is None + + +def test_default_path_reducers_see_every_series_from_the_db(db_session): + """The trigger onsets and ranker features (default explain path) read the rows metric_pk now keeps + per series; they must reduce per series, not average cpu.mode idle/user into a flat 0.5.""" + from src.core.explain.evidence import _metric_onsets + from src.core.ingestion.telemetry import persist_metric_samples + from src.core.rca.features import compute_features + + inst = {"service.instance.id": "pod-1"} + persist_metric_samples(db_session, _series("host", "system.cpu.utilization", 0.9, 0.1, {"cpu.mode": "idle", **inst}) + + _series("host", "system.cpu.utilization", 0.1, 0.9, {"cpu.mode": "user", **inst}), + scope=SCOPE) + db_session.flush() + onsets = _metric_onsets(db_session, SCOPE, W, END) + assert any(o.service == "host" and o.metric == "system.cpu.utilization" for o in onsets) + table = compute_features(db_session, SCOPE, incident_start=W, incident_end=END, baseline_start=BASE) + assert {f.service: f.met_anom for f in table.services}["host"] > 1.0 diff --git a/tests/unit/test_metric_reducers_per_series.py b/tests/unit/test_metric_reducers_per_series.py new file mode 100644 index 0000000..c80307b --- /dev/null +++ b/tests/unit/test_metric_reducers_per_series.py @@ -0,0 +1,82 @@ +"""#209 M2b — every metric-sample reducer groups by series, not by metric name. Pure, no DB. + +Ingestion now retains every series of an instrument (per cpu.mode, per route, per replica). The +reducers that already owned metric samples — ranker features, trigger onsets and the structural metric +sig — must reduce per series instead of averaging them into a line that describes none of them, and +samples without recorded identity must behave exactly as before.""" +from datetime import datetime, timedelta, timezone + +from src.core.rca.features import metric_features +from src.core.rca.observable import State +from src.core.rca.structural_model import summarize_metrics +from src.core.rca.triggers import metric_anomaly_onsets + +_W = datetime(2026, 1, 1, 12, 0, 0, tzinfo=timezone.utc) +_B = _W - timedelta(seconds=300) +_E = _W + timedelta(seconds=300) +_INST = {"service.instance.id": "pod-1"} + + +class _M: + def __init__(self, service, metric, value, ts, attributes=None): + self.service, self.metric, self.value, self.ts, self.attributes = service, metric, value, ts, attributes + + +def _series(service, metric, base, inc, attributes, n=6): + out = [_M(service, metric, base, _W - timedelta(seconds=30 * (i + 1)), attributes) for i in range(n)] + out += [_M(service, metric, inc, _W + timedelta(seconds=30 * (i + 1)), attributes) for i in range(n)] + return out + + +def _cpu_modes(): + # saturated CPU split by cpu.mode: idle 0.9 -> 0.1, user 0.1 -> 0.9 (the blended mean stays 0.5) + return (_series("h", "system.cpu.utilization", 0.9, 0.1, {"cpu.mode": "idle", **_INST}) + + _series("h", "system.cpu.utilization", 0.1, 0.9, {"cpu.mode": "user", **_INST})) + + +def _strip(samples): + return [_M(m.service, m.metric, m.value, m.ts, None) for m in samples] + + +class TestRankerFeatures: + def test_per_mode_change_is_not_averaged_away(self): + assert metric_features(_cpu_modes(), _B, _W, _E)["h"] > 1.0 # user series moved 9x + + def test_identity_less_samples_are_one_series_as_before(self): + # the same values without identity blend exactly like the pre-identity reducer did + assert metric_features(_strip(_cpu_modes()), _B, _W, _E)["h"] < 1e-6 + + +class TestTriggerOnsets: + def test_per_mode_onset_is_detected(self): + onsets = metric_anomaly_onsets(_cpu_modes(), _W, _E) + assert onsets and onsets[0].service == "h" and onsets[0].metric == "system.cpu.utilization" + + def test_identity_less_samples_are_one_series_as_before(self): + assert metric_anomaly_onsets(_strip(_cpu_modes()), _W, _E) == [] + + +class TestStructuralMetricSig: + def _routes(self): + # route A latency 10ms -> 30ms (3x); route B 100ms -> 100ms. Blended: 55 -> 65 (1.18x, "normal"). + return (_series("api", "latency_ms", 10.0, 30.0, {"http.route": "/a", **_INST}) + + _series("api", "latency_ms", 100.0, 100.0, {"http.route": "/b", **_INST}) + + _series("api", "error_rate", 0.0, 0.0, {"http.route": "/a", **_INST}) + + _series("api", "error_rate", 0.0, 0.0, {"http.route": "/b", **_INST})) + + def test_one_slow_route_is_present_not_averaged_away(self): + sig = summarize_metrics(self._routes(), _W)["api"] + assert sig.sig_state == State.PRESENT and abs(sig.latency_ratio - 3.0) < 1e-9 + + def test_absent_needs_every_series_measured_and_normal(self): + normal = (_series("api", "latency_ms", 10.0, 11.0, {"http.route": "/a", **_INST}) + + _series("api", "error_rate", 0.0, 0.0, {"http.route": "/a", **_INST})) + assert summarize_metrics(normal, _W)["api"].sig_state == State.ABSENT + # a second route seen only in the incident has no baseline -> latency not measured everywhere + extra = [m for m in _series("api", "latency_ms", 10.0, 11.0, {"http.route": "/b", **_INST}) + if m.ts >= _W] + assert summarize_metrics(normal + extra, _W)["api"].sig_state is None + + def test_identity_less_samples_are_one_series_as_before(self): + sig = summarize_metrics(_strip(self._routes()), _W)["api"] + assert sig.sig_state == State.ABSENT and abs(sig.latency_ratio - 65 / 55) < 1e-9 diff --git a/tests/unit/test_metric_series_identity.py b/tests/unit/test_metric_series_identity.py index dcf8763..5a65b30 100644 --- a/tests/unit/test_metric_series_identity.py +++ b/tests/unit/test_metric_series_identity.py @@ -38,10 +38,24 @@ def test_per_mode_gauge_datapoints_are_distinct_series(self): ] assert {s.metric_type for s in samples} == {"gauge"} - def test_attribute_free_datapoint_still_records_identity(self): + def test_attribute_free_datapoint_without_an_instance_is_recorded_but_unverified(self): + from src.core.rca.metric_series import instance_identity metric = {"name": "jvm.cpu.recent_utilization", "gauge": {"dataPoints": [_dp(0.5)]}} (s,) = parse_otlp_metrics(_otlp(metric, instance=None)) - assert s.attributes == {} # "this datapoint had no attributes" — not None ("never recorded") + assert s.attributes == {} # identity recorded (not None) ... + assert instance_identity(s.attributes) is None # ... but it names no instance: not verified + + def test_every_semconv_instance_attribute_is_recorded(self): + metric = {"name": "jvm.cpu.recent_utilization", "gauge": {"dataPoints": [_dp(0.5)]}} + objs = _otlp(metric, instance=None) + objs[0]["resourceMetrics"][0]["resource"]["attributes"] += [ + {"key": "container.id", "value": {"stringValue": "c1"}}, + {"key": "host.name", "value": {"stringValue": "h"}}, + {"key": "process.pid", "value": {"intValue": "42"}}, + {"key": "deployment.environment", "value": {"stringValue": "prod"}}, # not an instance attr + ] + (s,) = parse_otlp_metrics(objs) + assert s.attributes == {"container.id": "c1", "host.name": "h", "process.pid": "42"} def test_cumulative_monotonic_sum_is_declared_a_counter(self): metric = {"name": "jvm.cpu.time", diff --git a/tests/unit/test_structural_util_signals.py b/tests/unit/test_structural_util_signals.py index cb6bb66..0c50d23 100644 --- a/tests/unit/test_structural_util_signals.py +++ b/tests/unit/test_structural_util_signals.py @@ -90,12 +90,40 @@ def test_per_state_memory_gauge_is_unknown(self): + _gauge("h", "system.memory.utilization", 0.7, 0.1, attributes={"system.memory.state": "free"})) assert _state(samples, ("h", "memory")) is None - def test_two_reporting_instances_are_unknown(self): - samples = (_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9, attributes={"service.instance.id": "a"}) - + _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.02, - attributes={"service.instance.id": "b"}, offset=3)) + def test_replicas_without_a_named_instance_are_unknown(self): + # The reviewer's case: two replicas whose resource names no instance (identity {}), offset by + # 5s. Replica A 0.40->0.45, replica B 0.40->0.90: blended 0.40->0.675 (1.69x) would read as + # measured-normal while B is saturated. An unnamed instance is not a verified series -> UNKNOWN. + samples = (_gauge("ad", "jvm.cpu.recent_utilization", 0.40, 0.45, attributes={}) + + _gauge("ad", "jvm.cpu.recent_utilization", 0.40, 0.90, attributes={}, offset=5)) assert _state(samples, ("ad", "cpu")) is None + def test_named_replicas_are_measured_per_instance(self): + a = {"service.instance.id": "a"} + b = {"service.instance.id": "b"} + calm_plus_pegged = (_gauge("ad", "jvm.cpu.recent_utilization", 0.40, 0.45, attributes=a) + + _gauge("ad", "jvm.cpu.recent_utilization", 0.40, 0.90, attributes=b, offset=5)) + assert _state(calm_plus_pegged, ("ad", "cpu")) == State.PRESENT # replica B alone is 2.25x + both_calm = (_gauge("ad", "jvm.cpu.recent_utilization", 0.40, 0.45, attributes=a) + + _gauge("ad", "jvm.cpu.recent_utilization", 0.40, 0.42, attributes=b, offset=5)) + assert _state(both_calm, ("ad", "cpu")) == State.ABSENT + + def test_absent_needs_every_replica_measured(self): + a = {"service.instance.id": "a"} + b = {"service.instance.id": "b"} + calm = _gauge("ad", "jvm.cpu.recent_utilization", 0.40, 0.45, attributes=a) + b_incident_only = [m for m in _gauge("ad", "jvm.cpu.recent_utilization", 0.4, 0.4, attributes=b, offset=5) + if m.ts >= _W] + assert _state(calm + b_incident_only, ("ad", "cpu")) is None + + def test_instance_can_be_named_by_any_semconv_instance_attribute(self): + for attrs in ({"k8s.pod.uid": "u1"}, {"container.id": "c1"}, {"host.name": "h", "process.pid": "7"}): + assert _state(_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9, attributes=attrs), + ("ad", "cpu")) == State.PRESENT, attrs + # a host name alone repeats across processes on that host -> not an instance + assert _state(_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9, attributes={"host.name": "h"}), + ("ad", "cpu")) is None + def test_samples_without_series_identity_are_never_measured(self): assert _state(_gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9, attributes=None), ("ad", "cpu")) is None From 219b28444d6e13c759485e69514b5a0d7f89643d Mon Sep 17 00:00:00 2001 From: Leonardo Araujo Date: Sat, 26 Sep 2026 08:34:20 -0300 Subject: [PATCH 4/4] fix(trigger): one onset per (service, metric) so series cannot flood the budget (#212 review r3) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review r3 (merge acceptable for M2b; one default-path MAJOR + a stale docstring): - metric_anomaly_onsets detected per series but emitted onsets that carry only (service, metric), then cut the list at max_candidates=5. One multi-mode instrument (eight cpu.mode series of system.cpu.utilization) could spend the whole budget on indistinguishable copies and push a different instrument on the cause service off the list. Detection stays per series; the candidate is now ONE onset per (service, metric) — its strongest series (highest score; the max, never a mean) — before ranking and the cap. Docstring's first line corrected. For identity-less data each (service, metric) is already one series, so this is a no-op there: default-path eval on trace-loc (ranker + rare_event) vs main still 0/24 cases differ, aggregates identical. - test_structural_util_signals header now states the actual contract (named instance, one datapoint signature, any saturated replica PRESENT) instead of the superseded single-series refusal. New test: eight moving cpu.mode series + a different instrument on another service with max_candidates=5 -> exactly one system.cpu.utilization onset and the other instrument present (fails on the previous code). Unit 1539. Co-Authored-By: Claude Opus 5.5 --- src/core/rca/triggers.py | 34 +++++++++++++------ tests/unit/test_metric_reducers_per_series.py | 13 +++++++ tests/unit/test_structural_util_signals.py | 12 ++++--- 3 files changed, 44 insertions(+), 15 deletions(-) diff --git a/src/core/rca/triggers.py b/src/core/rca/triggers.py index 051bd0b..5c789fa 100644 --- a/src/core/rca/triggers.py +++ b/src/core/rca/triggers.py @@ -133,19 +133,26 @@ def metric_anomaly_onsets( rel_threshold: float = 0.5, max_candidates: int = 5, ) -> list[AnomalyOnset]: - """Earliest per-(service, metric) deviation from baseline, as onset candidates. + """At most one onset per ``(service, metric)``, as trigger candidates. ``samples`` are duck-typed rows with ``service`` / ``metric`` / ``value`` / - ``ts`` (and optionally ``attributes``). A series is ``(service, metric, - series_key(attributes))`` (#209 M2b): several series of one instrument (per - ``cpu.mode``, per replica, …) are detected separately, never averaged into a - flat line; samples without recorded identity form one series per - ``(service, metric)``, as before. Baseline = points with ``ts < incident_start`` (needs ``min_baseline`` + ``ts`` (and optionally ``attributes``). Detection runs **per series** — + ``(service, metric, series_key(attributes))`` (#209 M2b) — so several series of + one instrument (per ``cpu.mode``, per replica, …) are never averaged into a flat + line; samples without recorded identity form one series per ``(service, metric)``, + as before. Baseline = points with ``ts < incident_start`` (needs ``min_baseline`` of them); a point in ``[incident_start, incident_end]`` is anomalous when it is more than ``z_threshold`` baseline sigmas from the baseline mean, or — when the baseline is flat (sigma ~ 0) — more than ``rel_threshold`` of |mean| away. The - earliest anomalous point per series is its onset; ranked by earliness x - magnitude. This surfaces the injection time on faults that never log. + earliest anomalous point per series is that series' onset, scored by earliness x + magnitude. + + The candidate is then **one per (service, metric)**: its strongest series (highest + score — the max, never a mean). An onset carries only service and metric, so + several series of one instrument would otherwise be indistinguishable copies that + spend the whole ``max_candidates`` budget (eight ``cpu.mode`` series of + ``system.cpu.utilization``) and push a different instrument off the list. This + surfaces the injection time on faults that never log. """ from src.core.rca.metric_series import series_key @@ -188,5 +195,12 @@ def metric_anomaly_onsets( )) break # earliest anomalous point per series is the onset - onsets.sort(key=lambda o: (-o.score, _epoch(o.onset))) - return onsets[:max_candidates] + # One candidate per (service, metric): keep the strongest series' onset. + best: dict[tuple, AnomalyOnset] = {} + for o in onsets: + k = (o.service, o.metric) + cur = best.get(k) + if cur is None or (-o.score, _epoch(o.onset)) < (-cur.score, _epoch(cur.onset)): + best[k] = o + ranked = sorted(best.values(), key=lambda o: (-o.score, _epoch(o.onset))) + return ranked[:max_candidates] diff --git a/tests/unit/test_metric_reducers_per_series.py b/tests/unit/test_metric_reducers_per_series.py index c80307b..92dc056 100644 --- a/tests/unit/test_metric_reducers_per_series.py +++ b/tests/unit/test_metric_reducers_per_series.py @@ -55,6 +55,19 @@ def test_per_mode_onset_is_detected(self): def test_identity_less_samples_are_one_series_as_before(self): assert metric_anomaly_onsets(_strip(_cpu_modes()), _W, _E) == [] + def test_one_candidate_per_metric_so_series_cannot_flood_the_budget(self): + # Eight cpu.mode series all move; a different instrument on another service moves less. With + # max_candidates=5 the list must hold ONE system.cpu.utilization onset, not five copies of it. + modes = ("idle", "user", "system", "nice", "iowait", "irq", "softirq", "steal") + samples = [] + for i, mode in enumerate(modes): + samples += _series("h", "system.cpu.utilization", 0.1, 0.5 + 0.01 * i, {"cpu.mode": mode, **_INST}) + samples += _series("ad", "jvm.cpu.recent_utilization", 0.2, 0.9, _INST) + onsets = metric_anomaly_onsets(samples, _W, _E, max_candidates=5) + keys = [(o.service, o.metric) for o in onsets] + assert keys.count(("h", "system.cpu.utilization")) == 1 + assert ("ad", "jvm.cpu.recent_utilization") in keys + class TestStructuralMetricSig: def _routes(self): diff --git a/tests/unit/test_structural_util_signals.py b/tests/unit/test_structural_util_signals.py index 0c50d23..99a6cce 100644 --- a/tests/unit/test_structural_util_signals.py +++ b/tests/unit/test_structural_util_signals.py @@ -1,10 +1,12 @@ """#209 M2b — utilization observables util:{service}:{resource}. Pure, no DB. -Pins the semantics: a metric is measured only when it is ONE verified series — every sample carries -series identity (datapoint attributes + reporting instance), all samples belong to exactly one series, -no two share a timestamp, and the declared instrument (metric_type) is the one the reducer needs. -Anything else is UNKNOWN, never a measured ABSENT: in particular a per-state gauge -(system.cpu.utilization by cpu.mode) is never averaged into "normal".""" +Pins the semantics: a metric is measured only when every sample carries series identity that NAMES +its reporting instance, all samples share ONE datapoint signature (one mode), and the declared +instrument (metric_type) is the one the reducer needs. Each named instance is then its own series +(no timestamp collisions within it): PRESENT if ANY replica is saturated, ABSENT only if EVERY replica +is measured and below. Anything else is UNKNOWN, never a measured ABSENT — in particular a per-mode +gauge (system.cpu.utilization by cpu.mode) is never averaged into "normal", and unnamed replicas are +never blended.""" from datetime import datetime, timedelta, timezone from src.core.rca.observable import State