Skip to content

Commit 976856d

Browse files
committed
fix(metrics): initialize counters before first scrape
1 parent ae50736 commit 976856d

8 files changed

Lines changed: 109 additions & 3 deletions

File tree

openviking/metrics/collectors/base.py

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -122,6 +122,27 @@ def inc_counter(
122122
label_names=final_label_names,
123123
)
124124

125+
def initialize_counter(
126+
self,
127+
name: str,
128+
*,
129+
labels: dict[str, str] | None = None,
130+
label_names: tuple[str, ...] = (),
131+
account_id: str | None = None,
132+
) -> None:
133+
"""Initialize a zero-valued counter after normalizing account-aware labels."""
134+
final_labels, final_label_names = self._with_account_labels(
135+
metric_name=name,
136+
labels=labels,
137+
label_names=label_names,
138+
explicit_account_id=account_id,
139+
)
140+
self._registry.initialize_counter(
141+
name,
142+
labels=final_labels,
143+
label_names=final_label_names,
144+
)
145+
125146
def set_gauge(
126147
self,
127148
name: str,

openviking/metrics/collectors/queue.py

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,16 @@ def collect_hook(self, registry, metric_input) -> None:
5656
statuses = metric_input
5757
for queue_name, status in statuses.items():
5858
labels = {"queue": str(queue_name)}
59+
registry.initialize_counter(
60+
self.PROCESSED_TOTAL,
61+
labels=labels,
62+
label_names=("queue",),
63+
)
64+
registry.initialize_counter(
65+
self.ERRORS_TOTAL,
66+
labels=labels,
67+
label_names=("queue",),
68+
)
5969
registry.set_gauge(
6070
self.PENDING,
6171
float(status.pending),

openviking/metrics/collectors/resource.py

Lines changed: 14 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@
2020

2121
from openviking.metrics.core.base import MetricCollector
2222

23-
from .base import EventMetricCollector
23+
from .base import CollectorMetricWriter, EventMetricCollector
2424

2525

2626
@dataclass
@@ -48,10 +48,21 @@ class ResourceIngestionCollector(EventMetricCollector):
4848
)
4949

5050
SUPPORTED_EVENTS: ClassVar[frozenset[str]] = frozenset({"resource.stage", "resource.wait"})
51+
STAGES: ClassVar[tuple[str, ...]] = ("parse", "finalize", "persist", "summarize")
52+
STAGE_STATUSES: ClassVar[tuple[str, ...]] = ("ok", "error")
5153

5254
def collect(self, registry=None) -> None:
53-
"""Implement the collector interface as a no-op because resource metrics are push-driven."""
54-
return None
55+
"""Initialize bounded stage outcome series before the first resource event."""
56+
if registry is None:
57+
return
58+
writer = CollectorMetricWriter(registry)
59+
for stage in self.STAGES:
60+
for status in self.STAGE_STATUSES:
61+
writer.initialize_counter(
62+
self.STAGE_TOTAL,
63+
labels={"stage": stage, "status": status},
64+
label_names=("stage", "status"),
65+
)
5566

5667
def receive_hook(self, event_name: str, payload: dict, registry) -> None:
5768
"""

openviking/metrics/core/registry.py

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -106,6 +106,23 @@ def inc_counter(
106106
"""
107107
self.counter(name, label_names=label_names).inc(labels=labels, amount=float(amount))
108108

109+
def initialize_counter(
110+
self,
111+
name: str,
112+
*,
113+
labels: Mapping[str, str] | None = None,
114+
label_names: Sequence[str] = (),
115+
) -> None:
116+
"""
117+
Materialize a Counter series at zero without treating zero as an increment.
118+
119+
Args:
120+
name: Prometheus metric name.
121+
labels: Optional label dict. Keys must exactly match `label_names`.
122+
label_names: Ordered label key tuple for this metric name.
123+
"""
124+
self.counter(name, label_names=label_names).initialize(labels=labels)
125+
109126
def set_gauge(
110127
self,
111128
name: str,
@@ -376,6 +393,17 @@ def inc(self, *, labels: Mapping[str, str] | None, amount: float) -> None:
376393
return
377394
self._values[key] = self._values.get(key, 0.0) + float(amount)
378395

396+
def initialize(self, *, labels: Mapping[str, str] | None) -> None:
397+
"""Materialize one counter series at zero while preserving any existing value."""
398+
key = self._normalize_and_validate(labels)
399+
with self._lock:
400+
if key in self._values:
401+
return
402+
if len(self._values) >= self._max_series:
403+
self._on_drop(self.name)
404+
return
405+
self._values[key] = 0.0
406+
379407
def copy_values(self) -> dict[tuple[tuple[str, str], ...], float]:
380408
"""Return a detached copy of all series values in this family."""
381409
with self._lock:
@@ -620,6 +648,10 @@ def inc(self, amount: float = 1.0, *, labels: Mapping[str, str] | None = None) -
620648
"""Increment one series in the bound counter family using the public wrapper API."""
621649
self._family.inc(labels=labels, amount=amount)
622650

651+
def initialize(self, *, labels: Mapping[str, str] | None = None) -> None:
652+
"""Materialize one series at zero without weakening positive-increment validation."""
653+
self._family.initialize(labels=labels)
654+
623655

624656
class _Gauge:
625657
"""Public lightweight handle used by callers to mutate one gauge family."""
@@ -636,6 +668,7 @@ def inc(self, amount: float = 1.0, *, labels: Mapping[str, str] | None = None) -
636668
"""Increase one series in the bound gauge family by a positive delta."""
637669
self._family.add(labels=labels, delta=amount)
638670

671+
639672
class _Histogram:
640673
"""Public lightweight handle used by callers to mutate one histogram family."""
641674

openviking/metrics/global_api.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -301,6 +301,7 @@ def _build_event_router(registry: MetricRegistry) -> EventCollectorRouter:
301301
vlm_collector = VLMCollector()
302302
session_collector = SessionCollector()
303303
resource_collector = ResourceIngestionCollector()
304+
resource_collector.collect(registry)
304305
retrieval_collector = RetrievalCollector()
305306
encryption_collector = EncryptionCollector()
306307
telemetry_bridge_collector = TelemetryBridgeCollector()

tests/metrics/collectors/test_state_collectors.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,7 @@ async def check_status(self):
4545
assert 'openviking_queue_in_progress{queue="semantic"} 1.0' in text
4646
assert 'openviking_queue_processed_total{queue="semantic"} 10' in text
4747
assert 'openviking_queue_errors_total{queue="semantic"} 2' in text
48+
assert 'openviking_queue_errors_total{queue="embedding"} 0' in text
4849

4950

5051
def test_task_tracker_collector_maps_counts(monkeypatch):

tests/metrics/core/test_registry.py

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,22 @@ def test_counter_only_increases():
6464
c.inc(amount=-1)
6565

6666

67+
def test_counter_can_initialize_labeled_series_at_zero(render_prometheus):
68+
registry = MetricRegistry()
69+
c = registry.counter("openviking_initialized_total", label_names=("status",))
70+
71+
c.initialize(labels={"status": "error"})
72+
assert 'openviking_initialized_total{status="error"} 0' in render_prometheus(registry)
73+
74+
c.inc(amount=2, labels={"status": "error"})
75+
c.initialize(labels={"status": "error"})
76+
77+
text = render_prometheus(registry)
78+
assert 'openviking_initialized_total{status="error"} 2' in text
79+
with pytest.raises(ValueError):
80+
c.inc(amount=0, labels={"status": "error"})
81+
82+
6783
def test_histogram_boundary_bucket(registry, render_prometheus):
6884
h = registry.histogram("openviking_latency_seconds", buckets=(0.05, 0.1))
6985
h.observe(0.05)

tests/metrics/integration/test_bootstrap.py

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
import pytest
99

1010
import openviking.metrics.bootstrap as bootstrap
11+
import openviking.metrics.global_api as global_api
1112
from openviking.metrics.account_dimension import (
1213
configure_metric_account_dimension,
1314
reset_metric_account_dimension,
@@ -96,6 +97,18 @@ def _boom():
9697
bootstrap.create_default_collector_manager(app=None, service=None)
9798

9899

100+
def test_event_router_initializes_resource_stage_counters():
101+
registry = MetricRegistry()
102+
103+
global_api._build_event_router(registry)
104+
105+
text = PrometheusExporter(registry=registry).render()
106+
assert (
107+
'openviking_resource_stage_total{account_id="__unknown__",stage="parse",status="error"} 0'
108+
in text
109+
)
110+
111+
99112
def test_optional_cache_datasource_instrumentation_is_wired_into_key_call_sites():
100113
project_root = Path(__file__).resolve().parents[3]
101114
targets = [

0 commit comments

Comments
 (0)