This commit is contained in:
Ulysse Pence 2026-09-02 14:14:50 -02:00
parent 144862924b
commit 4a2f9da8a5
4 changed files with 41 additions and 6 deletions

View File

@ -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
)

View File

@ -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)

View File

@ -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"

View File

@ -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)