From 4a2f9da8a57d23bcdc95a792042eca2b9b173b05 Mon Sep 17 00:00:00 2001 From: Ulysse Pence Date: Wed, 2 Sep 2026 14:14:50 -0200 Subject: [PATCH] Comments --- src/backlog.py | 12 +++++++----- src/telemetry/prometheus/metrics.py | 10 ++++++++++ tests/telemetry/test_metric_zero_init.py | 5 ++++- tests/test_deriver_metrics.py | 20 ++++++++++++++++++++ 4 files changed, 41 insertions(+), 6 deletions(-) diff --git a/src/backlog.py b/src/backlog.py index c23899b4..1180be57 100644 --- a/src/backlog.py +++ b/src/backlog.py @@ -112,6 +112,9 @@ class DeriverMetricsPoller: stats = await crud.get_deriver_metrics(db) if self._dream_poll_due(): self._dreams_due = await count_due_dreams(db) + self._next_dream_poll = ( + time.monotonic() + settings.DREAM.DUE_POLL_INTERVAL_SECONDS + ) signal = outstanding_work_seconds(stats, dreams_due=self._dreams_due) measured_at = time.time() @@ -130,6 +133,7 @@ class DeriverMetricsPoller: pending_items=stats.pending_items, oldest_pending_age_seconds=stats.oldest_pending_age_seconds, embeddings_pending=stats.embeddings_pending, + embeddings_pending_due=stats.embeddings_pending_due, ) metrics.set_dreams_due(count=self._dreams_due) metrics.set_deriver_outstanding_work(seconds=signal) @@ -137,8 +141,6 @@ class DeriverMetricsPoller: def _dream_poll_due(self) -> bool: """The dream query is far more expensive, so it runs on its own spacing.""" - now = time.monotonic() - if self._next_dream_poll is not None and now < self._next_dream_poll: - return False - self._next_dream_poll = now + settings.DREAM.DUE_POLL_INTERVAL_SECONDS - return True + return ( + self._next_dream_poll is None or time.monotonic() >= self._next_dream_poll + ) diff --git a/src/telemetry/prometheus/metrics.py b/src/telemetry/prometheus/metrics.py index c8394228..6859c7cd 100644 --- a/src/telemetry/prometheus/metrics.py +++ b/src/telemetry/prometheus/metrics.py @@ -199,6 +199,14 @@ message_embeddings_pending_gauge = NamespacedGauge( ["namespace"], ) +message_embeddings_pending_due_gauge = NamespacedGauge( + "message_embeddings_pending_due", + "Pending MessageEmbedding rows past their retry backoff, so a sync attempt " + + "is due. Service-wide DB count, reported independently by every API " + + "replica — aggregate with max() or avg(), never sum()", + ["namespace"], +) + deriver_outstanding_work_seconds_gauge = NamespacedGauge( "deriver_outstanding_work_seconds", "Seconds of outstanding deriver work, 0 when a deriver has nothing to do. " @@ -615,6 +623,7 @@ class PrometheusMetrics: pending_items: int = 0, oldest_pending_age_seconds: float = 0.0, embeddings_pending: int = 0, + embeddings_pending_due: int = 0, ) -> None: try: deriver_queue_work_units_eligible_gauge.labels().set(eligible_work_units) @@ -624,6 +633,7 @@ class PrometheusMetrics: oldest_pending_age_seconds ) message_embeddings_pending_gauge.labels().set(embeddings_pending) + message_embeddings_pending_due_gauge.labels().set(embeddings_pending_due) except Exception as e: self._handle_metric_error("set_deriver_metrics", e) diff --git a/tests/telemetry/test_metric_zero_init.py b/tests/telemetry/test_metric_zero_init.py index 69b8f7ec..eb412287 100644 --- a/tests/telemetry/test_metric_zero_init.py +++ b/tests/telemetry/test_metric_zero_init.py @@ -138,8 +138,11 @@ _API_DERIVER_METRIC_GAUGES = ( "deriver_queue_items_pending", "deriver_queue_oldest_pending_age_seconds", "dreams_due", + "message_embeddings_pending_due", ) +_SHARED_DERIVER_METRIC_GAUGES = ("message_embeddings_pending",) + # --------------------------------------------------------------------------- # API-process zero-init @@ -171,7 +174,7 @@ def test_api_init_materializes_dialectic_and_embed(): ) assert sample("embed_now_tasks_shed_total") is not None assert sample("embed_now_tasks_in_flight") == 0.0 # gauge, explicit .set(0) - for gauge in _API_DERIVER_METRIC_GAUGES: + for gauge in (*_API_DERIVER_METRIC_GAUGES, *_SHARED_DERIVER_METRIC_GAUGES): assert sample(gauge) == 0.0, f"{gauge} was not zero-initialized" diff --git a/tests/test_deriver_metrics.py b/tests/test_deriver_metrics.py index a5d59c3a..acd97330 100644 --- a/tests/test_deriver_metrics.py +++ b/tests/test_deriver_metrics.py @@ -102,6 +102,26 @@ class TestPoller: assert dream_count.await_count == 1 assert poller.snapshot.dreams_due == 1 + async def test_a_failed_dream_query_is_retried_on_the_next_pass(self): + """Advancing the deadline first would republish the old count for a whole interval.""" + stats = schemas.DeriverMetrics() + poller = DeriverMetricsPoller() + dream_count = AsyncMock(side_effect=[RuntimeError("db down"), 4]) + + with ( + patch( + "src.backlog.crud.get_deriver_metrics", + AsyncMock(return_value=stats), + ), + patch("src.backlog.count_due_dreams", dream_count), + ): + with pytest.raises(RuntimeError): + await poller.refresh() + await poller.refresh() + + assert dream_count.await_count == 2 + assert poller.snapshot.dreams_due == 4 + async def test_a_failed_pass_leaves_the_previous_snapshot_alone(self): """A half-finished pass must never be published as a measurement.""" stats = schemas.DeriverMetrics(eligible_work_units=1)