diff --git a/specs/011-summarise-analysis-results/research.md b/specs/011-summarise-analysis-results/research.md index 4da52d7..b2c3681 100644 --- a/specs/011-summarise-analysis-results/research.md +++ b/specs/011-summarise-analysis-results/research.md @@ -87,6 +87,28 @@ than surprising. rejected as a blocker disproportionate to the first increment; noting it as follow-up work is enough. +**Built 2026-09-20, and the follow-up is now open rather than implied.** +`SummaryStore` holds summaries in process, bounded at 512 entries, evicting +the least recently used. So: + +- a summary survives a reload and not a deploy +- beta deploys often, so most readers will regenerate at some point +- `cached` on `start` is what makes that honest rather than surprising, and + is the reason FR-015 exists alongside FR-014 + +**The durable store is the open piece.** It needs somewhere to put it — +`POSTGRES_LANGGRAPH_DB` is unset on beta and LangGraph already falls back to +`MemorySaver`, so there is no existing home to reuse. Anyone picking this up +should treat "summaries vanish on deploy" as a known state, not a bug. + +Two decisions inside the store worth not reversing by accident. An **empty +summary is never stored**: a failed or abandoned generation leaves no text, +and storing it would serve the emptiness back forever as though it were the +answer, indistinguishable from a result with nothing to say. And eviction +**drops the oldest rather than refusing the newest**, because a reader whose +summary was evicted simply regenerates, where refusing new entries would +quietly stop the feature working for everyone after the first few hundred. + ## D5 — What "an option that discloses no identifiers" means, precisely **Decision**: Two disclosure tiers, defined by field rather than by intention. diff --git a/specs/011-summarise-analysis-results/tasks.md b/specs/011-summarise-analysis-results/tasks.md index 52bc52e..d1c9432 100644 --- a/specs/011-summarise-analysis-results/tasks.md +++ b/specs/011-summarise-analysis-results/tasks.md @@ -106,10 +106,10 @@ exists to prevent. Both signals are checked now. ## Phase 7: Stability and transparency (FR-014, FR-015) -- [ ] T027 Store summaries keyed `(token, release, tier)` in `src/analysis/store.py`; an aggregate summary and a disclosing one are different artefacts and must not be interchanged -- [ ] T028 [P] Test in `tests/analysis/test_store.py` that the same token returns byte-identical text on a second request, and that a release change discards the stored summary — the second half matters because the Analysis Service deletes the underlying result on a release (research D2, D3) -- [ ] T029 Report `cached` on the `start` event in `src/api/analysis_summary.py`, so the interface can say a summary was reused rather than implying the generator is deterministic (FR-015) -- [ ] T030 Record in [research.md](./research.md) that the first increment's store is in-process and lost on deploy, and open follow-up work for a durable store — beta sets no `POSTGRES_LANGGRAPH_DB` today (research D4) +- [x] T027 Store summaries keyed `(token, release, tier)` in `src/analysis/store.py`; an aggregate summary and a disclosing one are different artefacts and must not be interchanged. **Keyed on the tier that *applied*, not the one requested** — a disclosure that could not be honoured produced an aggregate summary, and storing it under `identifiers` would serve it back later as though the identifiers had been used +- [x] T028 [P] Test in `tests/analysis/test_store.py` that the same token returns byte-identical text on a second request, and that a release change discards the stored summary — the second half matters because the Analysis Service deletes the underlying result on a release (research D2, D3) +- [x] T029 Report `cached` on the `start` event in `src/api/analysis_summary.py`, so the interface can say a summary was reused rather than implying the generator is deterministic (FR-015) +- [x] T030 Record in [research.md](./research.md) that the first increment's store is in-process and lost on deploy, and open follow-up work for a durable store — beta sets no `POSTGRES_LANGGRAPH_DB` today (research D4) ## Phase 8: Human presence (FR-013) — UNBLOCKED 2026-09-19 diff --git a/src/analysis/store.py b/src/analysis/store.py new file mode 100644 index 0000000..1e2a9ac --- /dev/null +++ b/src/analysis/store.py @@ -0,0 +1,83 @@ +"""Summaries kept so the same analysis reads the same way twice. + +FR-014 wants a token to yield the same summary request after request. +Generation cannot provide that: measured on this repository, the same +question through the same surface twice scores 0.33 similarity. Storage can, +because an analysis result is a fixed artefact -- so stability here is +**reuse**, not determinism, and FR-015 requires the interface to say which. + +**Keyed by `(token, release, tier)`, and all three matter.** + +`release`, because the Analysis Service *deletes results on a new release*. +Without it a stored summary outlives the result it describes and we serve a +confident account of an analysis that no longer exists. + +`tier`, because an aggregate summary and a disclosing one are different +artefacts. Sharing a key would let a reader who chose the default be served +a summary built from their identifiers, or the reverse -- one of them a +disclosure nobody asked for. + +**In process, and lost on deploy.** Beta sets no `POSTGRES_LANGGRAPH_DB` and +LangGraph already falls back to `MemorySaver`, so nothing on that host +persists across a restart. This satisfies FR-014 within a process lifetime +and is honest only because `cached` tells the reader when a summary was +regenerated. A durable store is follow-up work, recorded in research D4. +""" + +from collections import OrderedDict +from dataclasses import dataclass, field +from time import time + +#: Bounded because this lives for the life of the process. Summaries are a +#: couple of kilobytes, so this is a few megabytes at worst, and the oldest +#: is dropped rather than the newest refused -- a reader whose summary was +#: evicted regenerates, which is the documented behaviour anyway. +DEFAULT_MAX_ENTRIES = 512 + + +@dataclass(frozen=True) +class Stored: + text: str + citations: tuple[tuple[str, str], ...] + generated_at: float + + +@dataclass +class SummaryStore: + max_entries: int = DEFAULT_MAX_ENTRIES + _entries: OrderedDict[tuple[str, str, str], Stored] = field( + default_factory=OrderedDict + ) + + def get(self, token: str, release: str, tier: str) -> Stored | None: + key = (token, release, tier) + found = self._entries.get(key) + if found is not None: + self._entries.move_to_end(key) + return found + + def put( + self, + token: str, + release: str, + tier: str, + text: str, + citations: tuple[tuple[str, str], ...], + ) -> Stored: + """Store a summary. An empty one is never stored. + + A failed or abandoned generation leaves no text, and storing that + would serve the emptiness back forever as though it were the answer. + """ + stored = Stored(text=text, citations=citations, generated_at=time()) + if not text.strip(): + return stored + key = (token, release, tier) + self._entries[key] = stored + self._entries.move_to_end(key) + while len(self._entries) > self.max_entries: + self._entries.popitem(last=False) + return stored + + def __len__(self) -> int: + return len(self._entries) diff --git a/src/api/analysis_summary.py b/src/api/analysis_summary.py index 65dd884..006610e 100644 --- a/src/api/analysis_summary.py +++ b/src/api/analysis_summary.py @@ -28,6 +28,7 @@ from agent.models import get_llm from analysis.client import current_release, fetch_not_found, fetch_result from analysis.disclosure import Tier, for_tier +from analysis.store import SummaryStore from analysis.summarise import ( INEXACT_COUNT_INSTRUCTION, NAMED_UNMATCHED_INSTRUCTION, @@ -51,6 +52,10 @@ #: not hold a connection open. _limiter = limiter_from_env() +#: Process-local, lost on deploy (research D4). Module state so it outlives +#: a request, as the limiter does. +_store = SummaryStore() + SYSTEM_PROMPT = """ You explain a completed Reactome pathway-analysis result to the researcher who ran it. @@ -169,15 +174,19 @@ async def stream() -> AsyncIterator[str]: "lookup returned nothing; serving aggregate" ) release = await current_release() + # Keyed on the tier that *applied*, not the one requested: a + # disclosure that could not be honoured produced an aggregate + # summary, and storing it under `identifiers` would serve it + # back later as though the identifiers had been used. + cached = _store.get(body.token, release, applied) if release else None yield _sse( "start", { "release": release, "analysis_type": model_input.get("analysis_type"), - # No store yet (Phase 7), so nothing is ever reused. - # Reported rather than omitted, because the interface - # must not imply a determinism this does not have. - "cached": False, + # Stability is reuse, not determinism (FR-015). This + # is how the interface knows which it is looking at. + "cached": cached is not None, # What the summary was built from. Equal to the # request's `disclosure` except when the disclosing # tier could not be honoured. @@ -188,15 +197,26 @@ async def stream() -> AsyncIterator[str]: # From the result, never from the model's prose. An invented # or mismatched identifier is impossible by construction # rather than by checking afterwards (SC-003). - for pathway in model_input["pathways"]: - if pathway.get("st_id"): - yield _sse( - "citation", - { - "st_id": pathway["st_id"], - "display_name": pathway.get("name") or pathway["st_id"], - }, - ) + citations: tuple[tuple[str, str], ...] = ( + cached.citations + if cached + else tuple( + (p["st_id"], p.get("name") or p["st_id"]) + for p in model_input["pathways"] + if p.get("st_id") + ) + ) + for st_id, display_name in citations: + yield _sse( + "citation", {"st_id": st_id, "display_name": display_name} + ) + + if cached: + # Byte-identical, and in one event: re-streaming it token + # by token would imitate generation that is not happening. + yield _sse("token", {"text": cached.text}) + yield _done("summarised", time.monotonic() - started) + return provider, model, base_url = resolve_llm_model(None) llm = get_llm(provider, model, base_url=base_url, request_timeout=90.0) @@ -220,10 +240,16 @@ async def stream() -> AsyncIterator[str]: f"Data:\n{json.dumps(model_input, default=str)}", ), ] + produced: list[str] = [] async for chunk in llm.astream(messages): text = getattr(chunk, "content", "") if isinstance(text, str) and text: + produced.append(text) yield _sse("token", {"text": text}) + if release: + _store.put( + body.token, release, applied, "".join(produced), citations + ) state = "summarised" except (asyncio.CancelledError, GeneratorExit): logger.info( diff --git a/tests/analysis/test_store.py b/tests/analysis/test_store.py new file mode 100644 index 0000000..b7a322d --- /dev/null +++ b/tests/analysis/test_store.py @@ -0,0 +1,79 @@ +"""Stability by reuse, and the two keys that stop it being wrong.""" + +from analysis.store import SummaryStore + +CITES = (("R-HSA-109581", "Apoptosis"),) + + +def test_the_same_request_returns_byte_identical_text() -> None: + # FR-014. Generation cannot provide this -- the same question through the + # same surface scores 0.33 similarity twice -- so stability is reuse. + store = SummaryStore() + store.put("tok", "97", "aggregate", "Four pathways pass correction.", CITES) + first = store.get("tok", "97", "aggregate") + second = store.get("tok", "97", "aggregate") + assert first is not None + assert second is not None + assert first.text == second.text == "Four pathways pass correction." + assert first.citations == CITES + + +def test_a_release_change_discards_the_summary() -> None: + # The Analysis Service deletes results on a new release, so without the + # release in the key a stored summary outlives the result it describes + # and we serve a confident account of an analysis that no longer exists. + store = SummaryStore() + store.put("tok", "97", "aggregate", "from release 97", CITES) + assert store.get("tok", "98", "aggregate") is None + assert store.get("tok", "97", "aggregate") is not None + + +def test_the_two_tiers_are_different_artefacts() -> None: + # Sharing a key would serve a reader who chose the default a summary + # built from their identifiers, or the reverse. One of those is a + # disclosure nobody asked for. + store = SummaryStore() + store.put("tok", "97", "aggregate", "no identifiers named", CITES) + store.put("tok", "97", "identifiers", "ABCA1_TYPO was not matched", CITES) + aggregate = store.get("tok", "97", "aggregate") + identifiers = store.get("tok", "97", "identifiers") + assert aggregate is not None + assert identifiers is not None + assert aggregate.text != identifiers.text + assert "ABCA1_TYPO" not in aggregate.text + + +def test_an_empty_summary_is_never_stored() -> None: + # A failed or abandoned generation leaves no text. Storing it would serve + # the emptiness back forever as though it were the answer, and nothing + # downstream could tell it from a result with nothing to say. + store = SummaryStore() + store.put("tok", "97", "aggregate", "", CITES) + store.put("tok", "97", "identifiers", " \n ", CITES) + assert store.get("tok", "97", "aggregate") is None + assert store.get("tok", "97", "identifiers") is None + assert len(store) == 0 + + +def test_the_oldest_is_dropped_rather_than_the_newest_refused() -> None: + # Bounded because this lives for the life of the process. A reader whose + # summary was evicted regenerates, which is the documented behaviour + # anyway; refusing to store new ones would silently stop the feature + # working for everyone after the first few hundred readers. + store = SummaryStore(max_entries=2) + for n in ("a", "b", "c"): + store.put(n, "97", "aggregate", f"summary {n}", CITES) + assert len(store) == 2 + assert store.get("a", "97", "aggregate") is None + assert store.get("c", "97", "aggregate") is not None + + +def test_reading_a_summary_keeps_it_from_being_evicted() -> None: + # The one people actually reload is the one worth keeping. + store = SummaryStore(max_entries=2) + store.put("a", "97", "aggregate", "summary a", CITES) + store.put("b", "97", "aggregate", "summary b", CITES) + store.get("a", "97", "aggregate") + store.put("c", "97", "aggregate", "summary c", CITES) + assert store.get("a", "97", "aggregate") is not None + assert store.get("b", "97", "aggregate") is None diff --git a/tests/api/test_analysis_summary.py b/tests/api/test_analysis_summary.py index 56043d2..86e202b 100644 --- a/tests/api/test_analysis_summary.py +++ b/tests/api/test_analysis_summary.py @@ -22,6 +22,7 @@ from fastapi.testclient import TestClient from analysis.client import Fetched +from analysis.store import SummaryStore from api.analysis_summary import router from util.caller_token import DEFAULT_AUDIENCE from util.rate_limit import SlidingWindowLimiter @@ -85,6 +86,11 @@ def wired(monkeypatch: pytest.MonkeyPatch) -> _Counter: "api.analysis_summary._limiter", SlidingWindowLimiter(limit=10_000, window=600.0), ) + # A private store per test. It is module state that outlives a request by + # design, so without this one test serves another its cached summary and + # the model is never called -- which looks like the feature being broken + # and is the tests interfering. + monkeypatch.setattr("api.analysis_summary._store", SummaryStore()) monkeypatch.setattr("api.analysis_summary.get_llm", lambda *a, **k: counter) monkeypatch.setattr( "api.analysis_summary.resolve_llm_model", lambda _c: ("openai", "m", None) @@ -470,3 +476,74 @@ async def astream(self, messages: Any) -> AsyncIterator[Any]: prompt = json.dumps(sent, default=str) assert expected in prompt, f"{analysis_type} was not given its own reading" assert forbidden not in prompt, f"{analysis_type} was given another's" + + +def _prose(response: Any) -> str: + return "".join( + payload["text"] for kind, payload in _events(response.text) if kind == "token" + ) + + +def test_a_second_request_is_byte_identical_and_calls_no_model( + keys: tuple[str, str], wired: _Counter +) -> None: + # FR-014 on the served path. The generator is not deterministic, so if + # the second request reached it the text would differ -- which is why + # the assertion is on the bytes and the call count together. + private, public = keys + first = _post(public, caller_token=_token(private)) + calls_after_first = wired.calls + second = _post(public, caller_token=_token(private)) + + assert _prose(first) == _prose(second) + assert _prose(first), "no prose was produced, so this proves nothing" + assert wired.calls == calls_after_first, "the model ran again" + assert _events(first.text)[0][1]["cached"] is False + assert _events(second.text)[0][1]["cached"] is True + + +def test_a_release_change_regenerates( + keys: tuple[str, str], wired: _Counter, monkeypatch: pytest.MonkeyPatch +) -> None: + # The service deletes results on a release, so a summary from the old one + # describes something that no longer exists. + private, public = keys + _post(public, caller_token=_token(private)) + before = wired.calls + + async def _next_release() -> str: + return "98" + + monkeypatch.setattr("api.analysis_summary.current_release", _next_release) + second = _post(public, caller_token=_token(private)) + assert _events(second.text)[0][1]["cached"] is False + assert wired.calls == before + 1 + + +def test_the_two_tiers_do_not_share_a_cached_summary( + keys: tuple[str, str], monkeypatch: pytest.MonkeyPatch +) -> None: + # Serving the aggregate summary to someone who chose to disclose, or the + # disclosing one to someone who did not, are both failures -- and the + # second is a disclosure nobody asked for. + _with_not_found_spy(monkeypatch) + private, public = keys + _post(public, caller_token=_token(private), disclosure="aggregate") + start = _events( + _post(public, caller_token=_token(private), disclosure="identifiers").text + )[0][1] + assert start["cached"] is False, "the disclosing tier reused the aggregate summary" + + +def test_citations_come_back_with_a_cached_summary(keys: tuple[str, str]) -> None: + # A reused summary that lost its chips would look like a summary citing + # nothing, and the caller cannot tell that from a result with no pathways. + private, public = keys + first = _post(public, caller_token=_token(private)) + second = _post(public, caller_token=_token(private)) + + def cited(response: Any) -> list[str]: + return [p["st_id"] for k, p in _events(response.text) if k == "citation"] + + assert cited(second) == cited(first) + assert cited(second), "no citations at all"