"""Note usage telemetry — was a surfaced note ever actually pulled? Two event streams, deliberately independent: - SURFACED: we put this note's title in front of the agent (auto-inject, or either arm of the write-path prior-art trigger). - PULLED: someone then opened it in full (get_snippet / get_note / the REST detail route). The ratio between them is the signal. A snippet surfaced forty times and never pulled is not neutral — it occupies the injection budget on every future turn and dilutes the menu — so this is what makes dead weight visible and prunable. Design notes (mirrors retrieval_telemetry, for the same reasons): - Writes are fire-and-forget. `record_surfaced` / `record_pulled` extract plain ints synchronously and schedule the insert as a background task, so telemetry never adds latency to — or can break — the surface it observes. - Failures degrade, but they must not degrade SILENTLY. The original version swallowed everything into logger.debug, and the deployed instance ran with every counter at zero for weeks while surfacing demonstrably fired — a total outage indistinguishable from "nobody uses this" (#2663). A subsystem whose every failure mode is invisible cannot report its own death, so failures now log at WARNING and drop one AppLog error row per process per site, where the admin UI shows it. - Reads (`usage_for_notes`) are NOT fire-and-forget — a readout the caller awaits, aggregated in one round-trip for a whole page of snippets rather than per row. """ from __future__ import annotations import asyncio import logging from collections.abc import Sequence from sqlalchemy import case, func, select from scribe.models import async_session from scribe.models.note import Note from scribe.models.note_usage import PULLED, SURFACED, NoteUsageEvent from scribe.models.base import iso from scribe.services.background import report_telemetry_failure logger = logging.getLogger(__name__) # Strong references to in-flight inserts. The event loop keeps only a WEAK # reference to a task, so a fire-and-forget create_task with no other holder # can be garbage-collected before it completes — a write that never errors and # never lands. The done-callback discard keeps the set from growing. _pending: set[asyncio.Task] = set() async def _report_failure(site: str) -> None: """This subsystem's canary, now the shared one. The per-site dedup, the WARNING and the single AppLog row all moved to `background.report_telemetry_failure` unchanged when `rule_usage` needed the identical behaviour — two hand-kept copies of a thing whose whole job is to be reliable is the wrong number. `retrieval_telemetry` deliberately still has its own: its canary is a different shape (one process-wide flag, no AppLog row), so repointing it would change behaviour rather than consolidate it. """ await report_telemetry_failure("note_usage", site) async def _insert_events(rows: list[dict]) -> None: """Persist usage rows. Best-effort: failures degrade, visibly.""" try: async with async_session() as session: session.add_all([NoteUsageEvent(**row) for row in rows]) await session.commit() except Exception: await _report_failure("write") def _schedule(rows: list[dict]) -> None: if not rows: return try: task = asyncio.get_running_loop().create_task(_insert_events(rows)) except RuntimeError: # No running loop (sync context outside the app) — skip rather than # block. Every app path runs on the loop. logger.debug("note usage telemetry skipped — no running event loop") return _pending.add(task) task.add_done_callback(_pending.discard) def _project_or_none(project_id: int | None) -> int | None: """0 and None both mean "no project reported" — store one of them. Callers reach this from two conventions at once: the MCP tools spell "no project" as `0` (it is an int parameter with an int default), while the column is nullable. Folding them here keeps every call site from having to remember which one this function wants, and stops a row claiming it was read on project #0. """ try: pid = int(project_id or 0) except (TypeError, ValueError): return None return pid or None def record_surfaced( *, user_id: int | None, note_ids: list[int] | set[int], source: str, project_id: int | None = None, ) -> None: """Fire-and-forget: record that these notes were shown to the agent. Takes the whole menu at once — one insert per surfacing event, not per note — because a menu is a single decision and its rows should land together. `project_id` is where the READER was, not where the record lives. Every surfacing arm knows it — it is the scope it just searched — so pass it; it is what makes "surfaced away from home" answerable without joining RetrievalLog. 0 is normalised to None: a project id of zero means "no project" everywhere else in this codebase, and storing it would read as project #0. """ try: rows = [ { "user_id": user_id, "note_id": int(nid), "event": SURFACED, "source": source, "project_id": _project_or_none(project_id), } for nid in note_ids ] except Exception: logger.debug("note usage payload build failed", exc_info=True) return _schedule(rows) def record_pulled( *, user_id: int | None, note_id: int, source: str, project_id: int | None = None, ) -> None: """Fire-and-forget: record that a note was opened in full. `project_id` is where the READER was — the caller's active project, never the record's own. Compared against the record's `project_id`, it answers whether this was opened somewhere other than where it was written, which is the evidence that a record TRANSFERRED (milestone 385, #3735). Unlike a surfacing arm, a getter only knows what it was handed, so this stays optional and null is an ordinary answer meaning "not reported". A pull with no project is still a pull: it counts toward dead-weight detection and simply cannot speak to transfer. """ try: rows = [ { "user_id": user_id, "note_id": int(note_id), "event": PULLED, "source": source, "project_id": _project_or_none(project_id), } ] except Exception: logger.debug("note usage payload build failed", exc_info=True) return _schedule(rows) # Surfacings that are NOT a ranked choice. enter_project returns whatever the # top-N-by-recency happen to be; the skill sync installs every Process the # operator can reach. Counting those alongside auto-inject would make a note's # surfaced_count dominated by "it was recently updated in a project you # opened", and dead-weight detection would read that as popularity (#2477). # They still matter — a pull that follows one must not float unattributed — so # they land in their own bucket rather than not landing at all. AMBIENT_SOURCES = ("enter_project", "process_skill_sync") def empty_usage() -> dict: """The zero readout — what a note with no recorded events looks like. Callers render this shape unconditionally, so a note predating the table reads as "never surfaced, never pulled" rather than as a missing key. `surfaced_count` is RANKED surfacings only — a scored surface chose this record. `ambient_count` is the rest (see AMBIENT_SOURCES). The split is the readout half of #2477: the "high surfaced, zero pulls → dead weight" reading is only valid over surfacings that were choices. `surfaced_away_count` / `pulled_away_count` are the subsets that happened on a project other than the one the record was written on — the evidence that a record TRANSFERRED, which is the claim the lesson kind rests on (milestone 385, #3735). Counted only where both projects are known: an event with no reader project, or a record with no project of its own, cannot speak to "away" and is left out rather than guessed. """ return { "surfaced_count": 0, "ambient_count": 0, "pull_count": 0, "last_surfaced_at": None, "last_pulled_at": None, "surfaced_away_count": 0, "pulled_away_count": 0, } async def usage_for_notes(note_ids: list[int]) -> dict[int, dict]: """Aggregate usage for a set of notes: {note_id: {counts + timestamps}}. One GROUP BY for the whole page rather than a query per row — this feeds a list view, so the per-row shape would be N+1 by construction. Notes with no events are returned with `empty_usage()` so the caller never has to distinguish "no events" from "not in the result". """ ids = [int(n) for n in note_ids] out: dict[int, dict] = {nid: empty_usage() for nid in ids} if not ids: return out # Classified in SQL so the group count stays small: per note we get at most # (surfaced-ranked, surfaced-ambient, pulled) rather than one row per # distinct source. ONE labelled expression, grouped by its label — a second # case() instance in GROUP BY renders with its own expanding-IN bind names # under asyncpg, so the database sees two DIFFERENT expressions and rejects # the query with a GroupingError. That rejection was swallowed, which is # how every counter read zero in production while the writes were landing # fine (#2663). ambient = case( (NoteUsageEvent.source.in_(AMBIENT_SOURCES), True), else_=False, ).label("ambient") try: async with async_session() as session: rows = ( await session.execute( select( NoteUsageEvent.note_id, NoteUsageEvent.event, func.count().label("n"), func.max(NoteUsageEvent.created_at).label("last_at"), ambient, ) .where(NoteUsageEvent.note_id.in_(ids)) .group_by( NoteUsageEvent.note_id, NoteUsageEvent.event, ambient, ) ) ).all() # Away from home: the reader's project against the record's own. # A second aggregate rather than a column on the first, because # it needs the join to `notes` and the first must keep counting # events whose note has since been deleted (the table is FK-free # so that evidence outlives the row). Ranked surfacings only, for # the same reason surfaced_count is (#2477). away_rows = ( await session.execute( select( NoteUsageEvent.note_id, NoteUsageEvent.event, func.count().label("n"), ) .join(Note, Note.id == NoteUsageEvent.note_id) .where( NoteUsageEvent.note_id.in_(ids), NoteUsageEvent.project_id.is_not(None), Note.project_id.is_not(None), NoteUsageEvent.project_id != Note.project_id, NoteUsageEvent.source.not_in(AMBIENT_SOURCES), ) .group_by(NoteUsageEvent.note_id, NoteUsageEvent.event) ) ).all() except Exception: # A telemetry readout must not be able to break the list it decorates — # but it must say it failed, or a broken readout is indistinguishable # from a corpus nobody uses (#2663). await _report_failure("readout") return out for note_id, event, n, last_at, ambient in rows: slot = out.get(int(note_id)) if slot is None: continue if event == SURFACED and ambient: slot["ambient_count"] = int(n) elif event == SURFACED: slot["surfaced_count"] = int(n) slot["last_surfaced_at"] = iso(last_at) elif event == PULLED: # Pulls are pulls regardless of what surfaced the record — the # question a pull answers ("did anyone ever open this?") doesn't # depend on how it was found. slot["pull_count"] = slot["pull_count"] + int(n) latest = iso(last_at) if latest and (slot["last_pulled_at"] or "") < latest: slot["last_pulled_at"] = latest for note_id, event, n in away_rows: slot = out.get(int(note_id)) if slot is None: continue if event == SURFACED: slot["surfaced_away_count"] = int(n) elif event == PULLED: slot["pulled_away_count"] = int(n) return out def _row_id(row: dict, key: str) -> int | None: """The note id on a payload row, or None when there is not one to read. Skipping is deliberate: an id this cannot parse is not a reason to fail a whole list, and GUESSING one would credit another record's counts to this row — a wrong chip is worse than no chip, because it reads as a measurement. `bool` is excluded explicitly because `int(True)` is 1, which would quietly attach note #1's usage to a row carrying a flag. """ raw = row.get(key) if raw is None or isinstance(raw, bool): return None try: return int(raw) except (TypeError, ValueError): return None async def attach_usage(rows: Sequence[dict], *, key: str = "id") -> None: """Add `usage` to every row of a payload a door is about to return (#4230). The one seam both doors and every record kind share. Before this, four call sites carried their own copy of the same lines — two list routes and two detail routes — and `/api/knowledge`, which is the list a person ACTUALLY browses notes and lessons in, was about to become a fifth. That is how the chip came to reach two record kinds out of four while a service named `usage_for_notes` worked on all of them: each door read fine on its own, and nobody was comparing them. ONE AGGREGATE FOR THE WHOLE PAGE. `usage_for_notes` is a single GROUP BY over the id set; calling it per row would be N+1 by construction, which is the one shape a list route must not have. EVERY ROW GETS THE KEY, zero-filled, so a record predating the table reads as "never surfaced, never pulled" rather than making the UI treat a missing field as a state. `UsageBadge` then renders nothing at all below one surfacing, because "0/0" would look like a verdict where there is only an absence of evidence. NO try/except HERE, deliberately — it is not an oversight. The fail-open already lives one layer down: `usage_for_notes` catches its own failure, reports it through `_report_failure("readout")` and returns the zero-filled map, so a broken readout degrades without breaking the list it decorates. Wrapping it again would swallow the REPORT along with the error, and a silently-swallowed readout failure is exactly #2663 — every counter reading zero in production for weeks while the writes landed fine. Mutates in place and returns None, matching how the call sites already used it: these rows are the payload, not a copy of it. A detail payload is just a one-row list — `await attach_usage([data])` — so the single-record doors share this seam rather than keeping a second shape that could drift from it. """ pairs = [(row, _row_id(row, key)) for row in rows] usage = await usage_for_notes([nid for _, nid in pairs if nid is not None]) for row, nid in pairs: if nid is not None: row["usage"] = usage.get(nid, empty_usage())