The deployed instance ran with every usage counter at zero while surfacing demonstrably fired. Every unit test was green because every unit test mocked either _schedule or the session — the two functions that touch the database ran against real Postgres nowhere. Both telemetry writers also held no reference to their fire-and-forget tasks (the loop keeps only weak ones), and swallowed every failure into logger.debug, so a total outage was indistinguishable from an unused corpus. - note_usage + retrieval_telemetry keep strong task references until done - failures log at WARNING; note_usage additionally drops one AppLog error row per process per site, so the admin UI shows the outage without host access - integration tests cover _insert_events -> usage_for_notes and the full record_pulled chain on a running loop, splitting the write and read halves so a failure names its side Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
138 lines
4.8 KiB
Python
138 lines
4.8 KiB
Python
"""Retrieval telemetry — one RetrievalLog row per semantic-retrieval call.
|
|
|
|
This is the empirical basis for KB-injection tuning: it records what each query
|
|
asked for, the score distribution of what came back, and the effective params,
|
|
so the similarity threshold and top-k can be tuned from data rather than guessed.
|
|
|
|
Design notes:
|
|
- Fire-and-forget, mirroring upsert_note_embedding: `record_retrieval` extracts
|
|
the primitives it needs SYNCHRONOUSLY (while the caller's Note objects are
|
|
still valid) and schedules the DB insert as a background task, so logging
|
|
never adds latency to — or can break — the search response.
|
|
- Result objects are reduced to {id, score, rank} before scheduling; the
|
|
background writer touches only plain data, never a possibly-detached ORM row.
|
|
- Every failure path is swallowed: telemetry must never take down retrieval.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import logging
|
|
|
|
from scribe.models import async_session
|
|
from scribe.models.note import Note
|
|
from scribe.models.retrieval_log import RetrievalLog
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Strong references to in-flight inserts — the loop holds tasks only weakly,
|
|
# and an unreferenced fire-and-forget task can be collected before it runs
|
|
# (same guard as note_usage, found via #2663).
|
|
_pending: set[asyncio.Task] = set()
|
|
|
|
# Whether this process already dropped its one warning about failing writes.
|
|
_reported = False
|
|
|
|
|
|
def _build_payload(
|
|
*,
|
|
user_id: int | None,
|
|
source: str,
|
|
query: str | None,
|
|
threshold: float | None,
|
|
limit: int | None,
|
|
project_id: int | None,
|
|
is_task: bool | None,
|
|
results: list[tuple[float, Note]],
|
|
duration_ms: float | None,
|
|
) -> dict:
|
|
"""Reduce a retrieval call to a flat, JSON-safe RetrievalLog payload.
|
|
|
|
Pure and synchronous (no DB, no event loop) so it is unit-testable and safe
|
|
to run inline before scheduling the write. `results` is the
|
|
`(score, Note)` list from semantic_search_notes, already highest-first.
|
|
"""
|
|
items = [
|
|
{"id": int(note.id), "score": round(float(score), 5), "rank": rank}
|
|
for rank, (score, note) in enumerate(results)
|
|
]
|
|
scores = [it["score"] for it in items]
|
|
return {
|
|
"user_id": user_id,
|
|
"source": source,
|
|
"query": query,
|
|
"threshold": threshold,
|
|
"limit_n": limit,
|
|
"project_id": project_id,
|
|
"is_task": is_task,
|
|
"result_count": len(items),
|
|
"top_score": (scores[0] if scores else None),
|
|
"min_score": (scores[-1] if scores else None),
|
|
"result_ids": items,
|
|
"duration_ms": (round(duration_ms, 2) if duration_ms is not None else None),
|
|
}
|
|
|
|
|
|
async def _insert_retrieval_log(payload: dict) -> None:
|
|
"""Persist one RetrievalLog row. Best-effort: failures degrade, visibly.
|
|
|
|
WARNING rather than debug — this table is the empirical basis for threshold
|
|
tuning, and a silent write outage yields a dataset that looks complete while
|
|
covering only part of the traffic (#2663's shape). Once per process is
|
|
enough to be found; per-call would flood the log with what it already said.
|
|
"""
|
|
global _reported
|
|
try:
|
|
async with async_session() as session:
|
|
session.add(RetrievalLog(**payload))
|
|
await session.commit()
|
|
except Exception:
|
|
if not _reported:
|
|
_reported = True
|
|
logger.warning("retrieval telemetry write failed", exc_info=True)
|
|
else:
|
|
logger.debug("retrieval telemetry write skipped", exc_info=True)
|
|
|
|
|
|
def record_retrieval(
|
|
*,
|
|
user_id: int | None,
|
|
source: str,
|
|
query: str | None,
|
|
threshold: float | None,
|
|
limit: int | None,
|
|
project_id: int | None,
|
|
is_task: bool | None,
|
|
results: list[tuple[float, Note]],
|
|
duration_ms: float | None = None,
|
|
) -> None:
|
|
"""Fire-and-forget: record one retrieval call.
|
|
|
|
Builds the payload inline (synchronously) then schedules the insert so the
|
|
caller returns immediately. Never raises — telemetry must not affect search.
|
|
"""
|
|
try:
|
|
payload = _build_payload(
|
|
user_id=user_id,
|
|
source=source,
|
|
query=query,
|
|
threshold=threshold,
|
|
limit=limit,
|
|
project_id=project_id,
|
|
is_task=is_task,
|
|
results=results,
|
|
duration_ms=duration_ms,
|
|
)
|
|
except Exception:
|
|
logger.debug("retrieval telemetry payload build failed", exc_info=True)
|
|
return
|
|
|
|
try:
|
|
task = asyncio.get_running_loop().create_task(_insert_retrieval_log(payload))
|
|
except RuntimeError:
|
|
# No running loop (e.g. called from sync context outside the app) —
|
|
# skip rather than block. The app paths always run on the loop.
|
|
logger.debug("retrieval telemetry skipped — no running event loop")
|
|
return
|
|
_pending.add(task)
|
|
task.add_done_callback(_pending.discard)
|