CI & Build / Python lint (push) Successful in 3s
CI & Build / Plugin hooks (push) Successful in 8s
CI & Build / integration (push) Successful in 30s
CI & Build / TypeScript typecheck (push) Successful in 33s
CI & Build / Python tests (push) Successful in 1m5s
CI & Build / Build & push image (push) Successful in 25s
The readout already grouped usage by source — `group_by(event, source)` — and the loop directly below it threw the source away, collapsing every surface into one corpus-wide ratio. So the question a threshold is actually tuned against, "is THIS surface worth its noise", could not be asked of any surface, while the data to answer it sat in the table. `usage.by_source` reports notes_surfaced / notes_pulled / pull_through per surface. The grain is the note, not the call: a pull records the door it came through, not the surface that led there, so grouping the pulled rows by source would answer a different question. Joining surfaced rows to pulled rows on note_id answers this one without the session identity #2085 declined to invent — at the cost of being an upper bound per surface, which the docstring says where it is read. Ambient surfaces report counts and a null ratio: nothing chose those records, so "surfaced often, opened never" is not a judgment about them. A surface that genuinely produced nothing reports 0.0, which must not look like the null. The join is guarded separately from the two reads above it. #2663 was a novel SQL shape the database rejected inside a broad except; this is the novel shape here, and it must not take down two readouts that work. Tests are integration for that same reason — a mock passes on a query Postgres refuses. They pin the distinct-first property (three surfacings of one note are one note), the ambient null, and the LIKE escape, since an unescaped `mcp_%` also matches `mcpXget_note` and nothing else in the payload would show the difference.
362 lines
16 KiB
Python
362 lines
16 KiB
Python
"""Tests for services.retrieval_telemetry.
|
|
|
|
_build_payload is pure (no DB, no loop) and gets unit coverage. The persistence
|
|
path (_insert_retrieval_log + the RetrievalLog model / JSONB roundtrip) is an
|
|
integration test against real Postgres.
|
|
"""
|
|
from types import SimpleNamespace
|
|
|
|
import pytest
|
|
|
|
from scribe.services.retrieval_telemetry import (
|
|
_build_payload,
|
|
record_retrieval,
|
|
)
|
|
|
|
|
|
def _note(nid):
|
|
"""Minimal stand-in — _build_payload only reads .id."""
|
|
return SimpleNamespace(id=nid)
|
|
|
|
|
|
# ─── _build_payload (pure) ───────────────────────────────────────────────────
|
|
|
|
|
|
def test_build_payload_ranks_and_score_bounds():
|
|
results = [(0.91, _note(11)), (0.72, _note(22)), (0.55, _note(33))]
|
|
p = _build_payload(
|
|
user_id=7, source="mcp_search", query="hello", threshold=0.45,
|
|
limit=10, project_id=3, is_task=None, results=results, duration_ms=12.345,
|
|
)
|
|
assert p["result_count"] == 3
|
|
assert p["top_score"] == 0.91
|
|
assert p["min_score"] == 0.55
|
|
assert [it["rank"] for it in p["result_ids"]] == [0, 1, 2]
|
|
assert [it["id"] for it in p["result_ids"]] == [11, 22, 33]
|
|
assert p["duration_ms"] == 12.35 # rounded to 2dp
|
|
assert p["user_id"] == 7 and p["project_id"] == 3 and p["threshold"] == 0.45
|
|
|
|
|
|
def test_build_payload_empty_results():
|
|
p = _build_payload(
|
|
user_id=1, source="rest_search", query="x", threshold=0.3,
|
|
limit=5, project_id=None, is_task=False, results=[], duration_ms=None,
|
|
)
|
|
assert p["result_count"] == 0
|
|
assert p["top_score"] is None and p["min_score"] is None
|
|
assert p["result_ids"] == []
|
|
assert p["duration_ms"] is None
|
|
|
|
|
|
def test_build_payload_rounds_scores_to_5dp():
|
|
p = _build_payload(
|
|
user_id=1, source="mcp_search", query="q", threshold=0.45,
|
|
limit=1, project_id=None, is_task=None,
|
|
results=[(0.123456789, _note(1))], duration_ms=0.0,
|
|
)
|
|
assert p["result_ids"][0]["score"] == 0.12346
|
|
|
|
|
|
def test_record_retrieval_without_event_loop_is_safe():
|
|
"""Called from a sync context (no running loop) it must swallow and return,
|
|
never raise — telemetry can't be allowed to break a caller."""
|
|
# No event loop running in this plain sync test.
|
|
assert record_retrieval(
|
|
user_id=1, source="mcp_search", query="q", threshold=0.45,
|
|
limit=10, project_id=None, is_task=None,
|
|
results=[(0.9, _note(1))],
|
|
) is None
|
|
|
|
|
|
# ─── persistence (integration) ───────────────────────────────────────────────
|
|
|
|
|
|
@pytest.mark.integration
|
|
@pytest.mark.asyncio
|
|
async def test_insert_retrieval_log_roundtrip(_dispose_engine):
|
|
from sqlalchemy import delete, select
|
|
|
|
from scribe.models import async_session
|
|
from scribe.models.retrieval_log import RetrievalLog
|
|
from scribe.services.retrieval_telemetry import _insert_retrieval_log
|
|
|
|
payload = _build_payload(
|
|
user_id=990001, source="mcp_search", query="pgvector tuning",
|
|
threshold=0.45, limit=10, project_id=None, is_task=None,
|
|
results=[(0.88, _note(501)), (0.61, _note(502))], duration_ms=9.9,
|
|
)
|
|
await _insert_retrieval_log(payload)
|
|
|
|
async with async_session() as s:
|
|
row = (
|
|
await s.execute(
|
|
select(RetrievalLog).where(RetrievalLog.user_id == 990001)
|
|
)
|
|
).scalars().first()
|
|
assert row is not None
|
|
assert row.source == "mcp_search"
|
|
assert row.result_count == 2
|
|
assert row.top_score == 0.88
|
|
# JSONB roundtrips as a list of dicts with the expected shape.
|
|
assert row.result_ids[0] == {"id": 501, "score": 0.88, "rank": 0}
|
|
assert row.created_at is not None # server_default now()
|
|
await s.execute(delete(RetrievalLog).where(RetrievalLog.user_id == 990001))
|
|
await s.commit()
|
|
|
|
|
|
# ─── the read half: retrieval_summary (integration) ──────────────────────────
|
|
# Integration, not mocked, and deliberately so. #2663 is the bug where a
|
|
# GROUP BY the database rejected was swallowed by a broad except, so every
|
|
# counter read zero in production while the writes landed fine and the mocked
|
|
# tests passed. `retrieval_summary` runs a grouped aggregate with
|
|
# percentile_cont ... WITHIN GROUP and a two-label CASE — precisely the shape
|
|
# that failed then. Only a real Postgres can say it parses.
|
|
|
|
|
|
@pytest.mark.integration
|
|
@pytest.mark.asyncio
|
|
async def test_retrieval_summary_reads_what_the_writer_wrote(_dispose_engine):
|
|
from sqlalchemy import delete
|
|
|
|
from scribe.models import async_session
|
|
from scribe.models.note_usage import NoteUsageEvent
|
|
from scribe.models.retrieval_log import RetrievalLog
|
|
from scribe.services.retrieval_telemetry import (
|
|
_insert_retrieval_log, retrieval_summary,
|
|
)
|
|
|
|
UID = 990002
|
|
# Three auto_inject calls at a 0.55 bar: two clear it, one does not.
|
|
# Plus one call that returned nothing at all — a different failure from a
|
|
# low-scoring one, and the readout must not blend them.
|
|
for score in (0.91, 0.72, 0.40):
|
|
await _insert_retrieval_log(_build_payload(
|
|
user_id=UID, source="auto_inject", query="q", threshold=0.55,
|
|
limit=3, project_id=None, is_task=None,
|
|
results=[(score, _note(1))], duration_ms=5.0,
|
|
))
|
|
await _insert_retrieval_log(_build_payload(
|
|
user_id=UID, source="auto_inject", query="q", threshold=0.55,
|
|
limit=3, project_id=None, is_task=None, results=[], duration_ms=5.0,
|
|
))
|
|
# A second surface, so the GROUP BY has something to separate.
|
|
await _insert_retrieval_log(_build_payload(
|
|
user_id=UID, source="mcp_search", query="q", threshold=0.45,
|
|
limit=10, project_id=None, is_task=None,
|
|
results=[(0.80, _note(2))], duration_ms=11.0,
|
|
))
|
|
# Corpus side: two ranked surfacings, one ambient, one pull.
|
|
async with async_session() as s:
|
|
s.add_all([
|
|
NoteUsageEvent(user_id=UID, note_id=1, event="surfaced", source="auto_inject"),
|
|
NoteUsageEvent(user_id=UID, note_id=2, event="surfaced", source="auto_inject"),
|
|
NoteUsageEvent(user_id=UID, note_id=3, event="surfaced", source="enter_project"),
|
|
NoteUsageEvent(user_id=UID, note_id=1, event="pulled", source="mcp_get_note"),
|
|
])
|
|
await s.commit()
|
|
|
|
try:
|
|
out = await retrieval_summary(UID, days=30)
|
|
|
|
assert out["read_failed"] is False, "the aggregate did not execute"
|
|
ai = out["sources"]["auto_inject"]
|
|
assert ai["calls"] == 4
|
|
assert ai["zero_result_calls"] == 1
|
|
assert ai["cleared_threshold"] == 2 # 0.91 and 0.72, not 0.40
|
|
# p50 over the three scored calls; the empty one contributes no score.
|
|
assert ai["top_score"]["p50"] == pytest.approx(0.72, abs=1e-4)
|
|
assert ai["top_score"]["min"] == pytest.approx(0.40, abs=1e-4)
|
|
assert ai["top_score"]["max"] == pytest.approx(0.91, abs=1e-4)
|
|
assert out["sources"]["mcp_search"]["calls"] == 1
|
|
|
|
u = out["usage"]
|
|
assert u["surfaced"] == 2 and u["ambient"] == 1 and u["pulled"] == 1
|
|
assert u["distinct_notes_surfaced"] == 2
|
|
# The pull came from `mcp_get_note`, so it counts as an AGENT pull
|
|
# and drives pull_through; a human `rest_*` pull would not.
|
|
assert u["pulled_by_agent"] == 1 and u["pulled_by_human"] == 0
|
|
assert u["pull_through"] == pytest.approx(0.5)
|
|
finally:
|
|
async with async_session() as s:
|
|
await s.execute(delete(RetrievalLog).where(RetrievalLog.user_id == UID))
|
|
await s.execute(delete(NoteUsageEvent).where(NoteUsageEvent.user_id == UID))
|
|
await s.commit()
|
|
|
|
|
|
@pytest.mark.integration
|
|
@pytest.mark.asyncio
|
|
async def test_retrieval_summary_is_empty_not_broken_for_a_fresh_install(_dispose_engine):
|
|
"""Rule #115: an install with no telemetry gets a coherent zero readout,
|
|
and `read_failed` stays False — the distinction #2663 says must exist."""
|
|
from scribe.services.retrieval_telemetry import retrieval_summary
|
|
|
|
out = await retrieval_summary(990003, days=30)
|
|
assert out["read_failed"] is False
|
|
assert out["sources"] == {}
|
|
assert out["usage"]["pull_through"] is None # no division by zero
|
|
assert out["usage"]["surfaced"] == 0
|
|
# An empty dict, not a missing key and not a failure flag — the same
|
|
# "no rows" / "read broke" distinction the rest of this readout keeps.
|
|
assert out["usage"]["by_source"] == {}
|
|
assert "by_source_failed" not in out["usage"]
|
|
|
|
|
|
@pytest.mark.integration
|
|
@pytest.mark.asyncio
|
|
async def test_retrieval_summary_sees_only_its_own_users_telemetry(_dispose_engine):
|
|
"""A retrieval log records what one user's agent asked for, query text
|
|
included. The owner filter is the access rule, so it gets a test."""
|
|
from sqlalchemy import delete
|
|
|
|
from scribe.models import async_session
|
|
from scribe.models.retrieval_log import RetrievalLog
|
|
from scribe.services.retrieval_telemetry import (
|
|
_insert_retrieval_log, retrieval_summary,
|
|
)
|
|
|
|
await _insert_retrieval_log(_build_payload(
|
|
user_id=990004, source="auto_inject", query="theirs", threshold=0.55,
|
|
limit=3, project_id=None, is_task=None, results=[(0.9, _note(1))],
|
|
duration_ms=1.0,
|
|
))
|
|
try:
|
|
assert (await retrieval_summary(990005, days=30))["sources"] == {}
|
|
assert (await retrieval_summary(990004, days=30))["sources"]["auto_inject"]["calls"] == 1
|
|
finally:
|
|
async with async_session() as s:
|
|
await s.execute(delete(RetrievalLog).where(RetrievalLog.user_id == 990004))
|
|
await s.commit()
|
|
|
|
|
|
# ─── per-source pull-through (#3311) ─────────────────────────────────────────
|
|
# Integration for the same reason the block above is: this is a self-join with
|
|
# two DISTINCT subqueries and a LIKE escape, which is a new SQL shape in a
|
|
# module whose one production outage (#2663) was a new SQL shape the database
|
|
# rejected inside a broad except. A mock would pass on a query Postgres refuses.
|
|
|
|
|
|
@pytest.mark.integration
|
|
@pytest.mark.asyncio
|
|
async def test_by_source_separates_a_surface_that_earns_its_noise_from_one_that_does_not(
|
|
_dispose_engine,
|
|
):
|
|
"""The whole point: the corpus average cannot say WHICH surface is working.
|
|
|
|
Two ranked surfaces, identical volume, opposite outcomes — and a top-level
|
|
ratio that describes neither of them.
|
|
"""
|
|
from sqlalchemy import delete
|
|
|
|
from scribe.models import async_session
|
|
from scribe.models.note_usage import NoteUsageEvent
|
|
from scribe.services.retrieval_telemetry import retrieval_summary
|
|
|
|
UID = 990010
|
|
async with async_session() as s:
|
|
s.add_all([
|
|
# auto_inject chose note 1 three times and note 2 once. Three
|
|
# surfacings of one note is ONE note surfaced — the DISTINCT that
|
|
# keeps the join from multiplying rows is what this pins.
|
|
NoteUsageEvent(user_id=UID, note_id=1, event="surfaced", source="auto_inject"),
|
|
NoteUsageEvent(user_id=UID, note_id=1, event="surfaced", source="auto_inject"),
|
|
NoteUsageEvent(user_id=UID, note_id=1, event="surfaced", source="auto_inject"),
|
|
NoteUsageEvent(user_id=UID, note_id=2, event="surfaced", source="auto_inject"),
|
|
# write_path_semantic chose two notes and got nothing opened.
|
|
NoteUsageEvent(user_id=UID, note_id=3, event="surfaced", source="write_path_semantic"),
|
|
NoteUsageEvent(user_id=UID, note_id=4, event="surfaced", source="write_path_semantic"),
|
|
# One agent pull, of a note only auto_inject surfaced.
|
|
NoteUsageEvent(user_id=UID, note_id=1, event="pulled", source="mcp_get_note"),
|
|
])
|
|
await s.commit()
|
|
|
|
try:
|
|
out = await retrieval_summary(UID, days=30)
|
|
assert out["read_failed"] is False
|
|
by_source = out["usage"]["by_source"]
|
|
assert "by_source_failed" not in out["usage"], "the join did not execute"
|
|
|
|
ai = by_source["auto_inject"]
|
|
assert ai["notes_surfaced"] == 2, "three surfacings of note 1 are one note"
|
|
assert ai["notes_pulled"] == 1
|
|
assert ai["pull_through"] == pytest.approx(0.5)
|
|
|
|
wp = by_source["write_path_semantic"]
|
|
assert wp["notes_surfaced"] == 2
|
|
assert wp["notes_pulled"] == 0
|
|
# 0.0, NOT None. "This surface produced nothing" is a finding; None is
|
|
# what a surface with no data reads as, and they must not look alike.
|
|
assert wp["pull_through"] == 0.0
|
|
|
|
# And the number that exists today, which is true of neither surface:
|
|
# one agent pull over six ranked surfacings.
|
|
assert out["usage"]["pull_through"] == pytest.approx(1 / 6, abs=1e-4)
|
|
finally:
|
|
async with async_session() as s:
|
|
await s.execute(delete(NoteUsageEvent).where(NoteUsageEvent.user_id == UID))
|
|
await s.commit()
|
|
|
|
|
|
@pytest.mark.integration
|
|
@pytest.mark.asyncio
|
|
async def test_an_ambient_surface_reports_its_counts_but_no_ratio(_dispose_engine):
|
|
"""`enter_project` bulk-loads records; nothing CHOSE them. "Surfaced often,
|
|
opened never" is not a judgment about a record that was never picked, so the
|
|
counts stay visible and the ratio that would be misread is null."""
|
|
from sqlalchemy import delete
|
|
|
|
from scribe.models import async_session
|
|
from scribe.models.note_usage import NoteUsageEvent
|
|
from scribe.services.retrieval_telemetry import retrieval_summary
|
|
|
|
UID = 990011
|
|
async with async_session() as s:
|
|
s.add_all([
|
|
NoteUsageEvent(user_id=UID, note_id=1, event="surfaced", source="enter_project"),
|
|
NoteUsageEvent(user_id=UID, note_id=2, event="surfaced", source="enter_project"),
|
|
])
|
|
await s.commit()
|
|
|
|
try:
|
|
row = (await retrieval_summary(UID, days=30))["usage"]["by_source"]["enter_project"]
|
|
assert row["ambient"] is True
|
|
assert row["notes_surfaced"] == 2
|
|
assert row["pull_through"] is None
|
|
finally:
|
|
async with async_session() as s:
|
|
await s.execute(delete(NoteUsageEvent).where(NoteUsageEvent.user_id == UID))
|
|
await s.commit()
|
|
|
|
|
|
@pytest.mark.integration
|
|
@pytest.mark.asyncio
|
|
async def test_the_agent_pull_filter_does_not_treat_its_underscore_as_a_wildcard(
|
|
_dispose_engine,
|
|
):
|
|
"""`_` is a LIKE wildcard, so an unescaped `LIKE 'mcp_%'` also matches
|
|
`mcpXsomething`. The Python half of this readout uses str.startswith and
|
|
cannot have the bug; the SQL half needs autoescape to match it, and nothing
|
|
else in the payload would reveal the difference."""
|
|
from sqlalchemy import delete
|
|
|
|
from scribe.models import async_session
|
|
from scribe.models.note_usage import NoteUsageEvent
|
|
from scribe.services.retrieval_telemetry import retrieval_summary
|
|
|
|
UID = 990012
|
|
async with async_session() as s:
|
|
s.add_all([
|
|
NoteUsageEvent(user_id=UID, note_id=1, event="surfaced", source="auto_inject"),
|
|
# Not an agent pull: the door is `mcpXget_note`, not `mcp_get_note`.
|
|
NoteUsageEvent(user_id=UID, note_id=1, event="pulled", source="mcpXget_note"),
|
|
])
|
|
await s.commit()
|
|
|
|
try:
|
|
row = (await retrieval_summary(UID, days=30))["usage"]["by_source"]["auto_inject"]
|
|
assert row["notes_pulled"] == 0, "a wildcard match counted a non-agent pull"
|
|
assert row["pull_through"] == 0.0
|
|
finally:
|
|
async with async_session() as s:
|
|
await s.execute(delete(NoteUsageEvent).where(NoteUsageEvent.user_id == UID))
|
|
await s.commit()
|