From cb9906e18b850fda5da879c2d374b110f0cbdae3 Mon Sep 17 00:00:00 2001 From: tanmay-sharma Date: Tue, 1 Sep 2026 10:51:43 +0530 Subject: [PATCH] Add scope-safe MMR retrieval to prevent alias crowding in recall. Authorization still precedes diversification; corpus recall and memory search can opt into deterministic source-diverse selection with adversarial tests and README proof. --- README.md | 129 ++++++++++++++- s13code/core/memory/store.py | 75 ++++++++- s13code/routes.py | 10 +- s13code/runtime.py | 22 +-- tests/test_document_memory.py | 4 +- tests/test_scope_safe_diverse_recall.py | 210 ++++++++++++++++++++++++ 6 files changed, 420 insertions(+), 30 deletions(-) create mode 100644 tests/test_scope_safe_diverse_recall.py diff --git a/README.md b/README.md index 7cf7872..6655f91 100644 --- a/README.md +++ b/README.md @@ -110,7 +110,7 @@ uv run pytest -q ## Student contribution -Fork the official [`theschoolofai/S13Code`](https://github.com/theschoolofai/S13Code) repository linked from Axiom, create a branch, implement one meaningful extension, and open one pull request against that repository. Do not open the Session 13 pull request against [`theschoolofai/glc_v3`](https://github.com/theschoolofai/glc_v3). +Fork the official `S13Code` repository linked from Axiom, create a branch, implement one meaningful extension, and open one pull request against that repository. Do not open the Session 13 pull request against `glc_v3`. Add one subsection to this README in the same pull request. It must contain: @@ -124,6 +124,133 @@ Add one subsection to this README in the same pull request. It must contain: Do not commit `.env`, credentials, personal memory, generated databases, unrestricted local paths, benchmark output containing private data, or provider responses containing secrets. Use synthetic identities in every proof. +### Scope-Safe Diverse Retrieval (Memory track) + +Cross-document questions used to return several chunks from the same aliased source before complementary papers appeared. The memory store now applies **scope-first maximal marginal relevance (MMR)** after authorization: tenant/project/user/agent/run filtering happens first, then an opt-in `diversify=true` pass penalizes near-duplicate embeddings and repeated source identities. Corpus recall in the live graph enables this automatically; direct memory search exposes it explicitly. + +#### 1. User-visible capability + +A caller can request diversified recall so near-identical document aliases no longer crowd out complementary sources. Authorization still precedes diversification — a foreign tenant’s records never enter the candidate pool. + +#### 2. Exact API request + +```bash +curl -s http://127.0.0.1:8113/v1/agent/memory/search \ + -H 'Content-Type: application/json' \ + -d '{ + "tenant_id": "course", + "project_id": "mmr-proof", + "user_id": "student-01", + "query": "credit assignment chain of thought reasoning", + "limit": 3, + "kinds": ["document_chunk"], + "diversify": true + }' +``` + +Corpus agent run (diversification enabled inside `memory_recall`): + +```bash +curl -s http://127.0.0.1:8113/v1/agent/runs \ + -H 'Content-Type: application/json' \ + -d '{ + "tenant_id": "course", + "project_id": "papers", + "user_id": "student-01", + "prompt": "Across the papers I have indexed, what do they say about chain-of-thought reasoning?" + }' +``` + +#### 3. Graph and ordered event trace + +Corpus query run (`run-*`): + +``` +run_started +graph_patched: add [recall] +task_started: recall +task_succeeded: recall +graph_patched: add [answer] +task_started: answer +task_succeeded: answer +graph_patched: finish +``` + +Live expansion floor (asyncio prompt, `run-f906d4c798ca`): + +``` +run_started → graph_patched add [search] → task_started search → task_succeeded search +→ graph_patched add [fetch_1, fetch_2, fetch_3] → parallel fetch_* → distill (failed) +→ answer succeeded +``` + +Fetch nodes appeared **only after** search returned URLs. + +A2A waiting/resume floor (`run_a2a_proof.py`): + +``` +run_started → graph_patched wait [remote_specialist] → a2a_task_completed +→ graph_patched resume [remote_specialist] → task_succeeded remote_specialist → finish +``` + +#### 4. Actual final result + +Memory search with `diversify=true` returned three distinct source URIs from the indexed alias/complementary fixture (`cot.txt` plus two DPO alias copies). Corpus query answer (Ollama provider): + +> Across the indexed papers, chain-of-thought reasoning is discussed in several contexts: improving model performance on GSM8K via intermediate reasoning steps, emergence in large models, … + +Birthday cross-run recall: + +> Mom's birthday is on 15 May 2026. … [source: api://agent/runs fact] + +#### 5. Evidence and provider/agent assignments + +| node | agent | provider | evidence | +|---|---|---|---| +| `recall` (corpus) | memory_recall | — | 8 hits spanning cot, react, attention, dpo, lora sources | +| `answer` (corpus) | answer_with_evidence | ollama | grounded synthesis from diversified recall hits | +| `remember` (birthday) | remember_explicit_fact | — | fact `Mom's birthday is 15 May 2026.` source `api://agent/runs` | +| `remote_specialist` (A2A) | remote-a2a-architecture-specialist | — | artifact: Agent Card grants no local-memory authority | + +#### 6. Adversarial failure and fix + +**Before:** six identical alias chunks indexed under different URIs filled every recall slot; a complementary ReAct chunk was excluded from the top-2 plain ranking (see `tests/test_scope_safe_diverse_recall.py::test_adversarial_alias_flood_blocked_only_when_diversify_enabled`). + +**After:** the same query with `diversify=true` returns the complementary `file:///papers/react.txt` chunk alongside one alias copy. Wrong-tenant records with higher novelty scores never appear because scope authorization filters them before MMR runs. + +#### 7. Reproduce from a fresh checkout + +```bash +# Terminal 1 — gateway (Ollama provider, no Gemini keys required) +cd glc_v3 && uv sync && cp .env.example .env +# set OLLAMA_MODEL to a local tag, e.g. qwen2.5:7b +uv run glc serve + +# Terminal 2 — S13Code +cd ../S13Code && uv sync +export GLC_BASE_URL=http://127.0.0.1:8111 +export S13_GATEWAY_PROVIDER=ollama +export S13_SANDBOX_ROOT="$PWD/sandbox" +export S13_CHUNK_MODEL=phi4-mini:latest +export S13_LIVE_SEMANTIC_CHUNKING=1 +uv run s13code serve + +# Deterministic proofs (no provider quota) +uv run ruff check . +uv run pytest -q +uv run pytest tests/test_scope_safe_diverse_recall.py -q + +cd ../S13Proof && uv sync && uv run pytest -q +uv run python run_a2a_proof.py --output a2a-proof.json + +# Live memory diversification API (after indexing six alias docs + one complementary doc) +curl -s http://127.0.0.1:8113/v1/agent/memory/search \ + -H 'Content-Type: application/json' \ + -d '{"tenant_id":"course","project_id":"mmr-proof","user_id":"student-01","query":"credit assignment chain of thought reasoning","limit":3,"kinds":["document_chunk"],"diversify":true}' +``` + +**Honest limitation:** with live Nomic embeddings, a genuinely top-ranked complementary chunk may already appear in plain recall; MMR’s value shows most clearly when many near-duplicate alias chunks would otherwise occupy every slot. The asyncio live run also exposed a `distill` worker failure (`RuntimeError` from the gateway) — the graph still finished with an `answer` node reporting fetched evidence, but specialist synthesis did not run. + ## License MIT. See `LICENSE`. diff --git a/s13code/core/memory/store.py b/s13code/core/memory/store.py index c71a9cb..e52ecf0 100644 --- a/s13code/core/memory/store.py +++ b/s13code/core/memory/store.py @@ -5,7 +5,8 @@ import sqlite3 import hashlib from pathlib import Path -from typing import Iterable, TYPE_CHECKING +from typing import Any, Iterable, TYPE_CHECKING +from urllib.parse import urlparse from .embeddings import DeterministicEmbedder, Embedder, cosine from .models import MemoryKind, MemoryRecord, MemoryScope, Principal, SourceRef, utcnow @@ -34,7 +35,9 @@ def __init__(self, path: str | Path = ":memory:", *, embedder: Embedder | None = self._create_schema() self.vectors = PersistentFaissIndex(self.path.parent if self.path else None) self._ensure_vector_index() - self.last_retrieval_stats: dict[str, int] = {"index_candidates": 0, "authorized_records": 0} + self.last_retrieval_stats: dict[str, Any] = { + "index_candidates": 0, "authorized_records": 0, "diversified": False, + } def close(self) -> None: self.db.close() @@ -211,7 +214,8 @@ def _scope_where(scope: MemoryScope) -> tuple[str, list[str | None]]: return " AND ".join(clauses), values def recall(self, query: str, scope: MemoryScope, *, kinds: Iterable[MemoryKind] | None = None, - limit: int = 8, include_history: bool = False, expand_neighbors: int = 0) -> list[MemoryRecord]: + limit: int = 8, include_history: bool = False, expand_neighbors: int = 0, + diversify: bool = False, diversify_lambda: float = 0.7) -> list[MemoryRecord]: where, values = self._scope_where(scope) if not include_history: where += " AND status='current' AND (valid_to IS NULL OR valid_to > ?)" @@ -235,7 +239,11 @@ def recall(self, query: str, scope: MemoryScope, *, kinds: Iterable[MemoryKind] # foreign tenant consuming the first few approximate candidates. vector_hits = self.vectors.search(query_embedding, max(128, limit * 16)) vector_scores = {record_id: score for record_id, score in vector_hits} - self.last_retrieval_stats = {"index_candidates": len(vector_hits), "authorized_records": len(rows)} + self.last_retrieval_stats = { + "index_candidates": len(vector_hits), + "authorized_records": len(rows), + "diversified": diversify, + } terms = [term for term in query.lower().replace("'", " ").split() if term.isalnum()] # We deliberately score only scope-authorized rows. FTS is a boost, not # an authorization mechanism and never gets to choose the tenant. @@ -246,7 +254,7 @@ def recall(self, query: str, scope: MemoryScope, *, kinds: Iterable[MemoryKind] lexical = {r[0] for r in self.db.execute("SELECT id FROM records_fts WHERE records_fts MATCH ?", (match,))} except sqlite3.OperationalError: pass - scored = [] + scored: list[tuple[float, MemoryRecord, list[float]]] = [] changed = False for row in rows: vector, stale = self._unpack_embedding(row["embedding_json"]) @@ -262,11 +270,18 @@ def recall(self, query: str, scope: MemoryScope, *, kinds: Iterable[MemoryKind] # remains the safe recall fallback for this authorized scope. score = vector_scores.get(row["id"], cosine(query_embedding, vector)) + (0.20 if row["id"] in lexical else 0.0) score += self._content_quality(row["text"]) - scored.append((score, self._to_record(row))) + scored.append((score, self._to_record(row), vector)) if changed: self.db.commit() - scored.sort(key=lambda pair: (-pair[0], pair[1].created_at)) - selected = [record for _, record in scored[:limit]] + scored.sort(key=lambda pair: (-pair[0], pair[1].created_at, pair[1].id)) + if diversify: + pool_size = min(len(scored), max(limit * 4, 32)) + selected = self._mmr_select(query_embedding, scored[:pool_size], limit=limit, + lambda_=diversify_lambda) + self.last_retrieval_stats["mmr_pool"] = pool_size + self.last_retrieval_stats["returned"] = len(selected) + else: + selected = [record for _, record, _ in scored[:limit]] if not expand_neighbors: return selected seen = {record.id for record in selected} @@ -289,6 +304,50 @@ def recall(self, query: str, scope: MemoryScope, *, kinds: Iterable[MemoryKind] expanded.append(neighbour) return expanded + @staticmethod + def _normalize_source_uri(uri: str) -> str: + """Collapse cosmetic URI differences that alias the same document.""" + parsed = urlparse(uri.strip().lower()) + path = parsed.path.rstrip("/") or "/" + return f"{parsed.scheme}://{parsed.netloc}{path}" + + @staticmethod + def _source_identity(record: MemoryRecord) -> str: + metadata = record.metadata + document_id = metadata.get("document_id") + if document_id is not None: + version = metadata.get("document_version", "current") + return f"doc:{document_id}:v{version}" + if record.sources: + return MemoryStore._normalize_source_uri(record.sources[0].uri) + return record.id + + def _mmr_select(self, query_embedding: list[float], + candidates: list[tuple[float, MemoryRecord, list[float]]], *, + limit: int, lambda_: float) -> list[MemoryRecord]: + """Deterministic maximal marginal relevance after scope authorization.""" + selected: list[MemoryRecord] = [] + selected_vectors: list[list[float]] = [] + selected_sources: set[str] = set() + pool = list(candidates) + while len(selected) < limit and pool: + best_idx = -1 + best_mmr = float("-inf") + best_tie = ("", "") + for index, (relevance, record, vector) in enumerate(pool): + redundancy = max((cosine(vector, chosen) for chosen in selected_vectors), default=0.0) + source_key = self._source_identity(record) + source_penalty = 0.12 if source_key in selected_sources else 0.0 + mmr = lambda_ * relevance - (1.0 - lambda_) * redundancy - source_penalty + tie = (record.created_at, record.id) + if mmr > best_mmr or (mmr == best_mmr and tie < best_tie): + best_mmr, best_idx, best_tie = mmr, index, tie + _, pick, pick_vector = pool.pop(best_idx) + selected.append(pick) + selected_vectors.append(pick_vector) + selected_sources.add(self._source_identity(pick)) + return selected + def _pack_embedding(self, values: list[float]) -> dict[str, object]: return {"fingerprint": getattr(self.embedder, "fingerprint", type(self.embedder).__name__), "values": values} diff --git a/s13code/routes.py b/s13code/routes.py index 1570088..7f616c9 100644 --- a/s13code/routes.py +++ b/s13code/routes.py @@ -49,6 +49,7 @@ class SearchBody(ScopeBody): query: str = Field(min_length=1, max_length=20_000) limit: int = Field(default=5, ge=1, le=50) kinds: list[MemoryKind] | None = None + diversify: bool = False @router.post("/runs") @@ -89,9 +90,12 @@ async def document(body: IndexBody, request: Request): @router.post("/memory/search") async def memory_search(body: SearchBody, request: Request): - hits = request.app.state.s13_runtime.memory.recall(body.query, body.scope(), - kinds=body.kinds, limit=body.limit) - return {"query": body.query, "hits": [{"id": hit.id, "kind": hit.kind.value, + store = request.app.state.s13_runtime.memory + hits = store.recall(body.query, body.scope(), kinds=body.kinds, limit=body.limit, + diversify=body.diversify) + return {"query": body.query, "diversified": body.diversify, + "retrieval_stats": store.last_retrieval_stats, + "hits": [{"id": hit.id, "kind": hit.kind.value, "text": hit.text, "sources": [source.uri for source in hit.sources], "metadata": hit.metadata} for hit in hits]} diff --git a/s13code/runtime.py b/s13code/runtime.py index 45c16e3..0c42a00 100644 --- a/s13code/runtime.py +++ b/s13code/runtime.py @@ -216,26 +216,14 @@ async def recall(task: TaskSpec) -> dict[str, Any]: corpus_query = bool(re.search(r"\b(papers?|documents?|indexed|corpus)\b", task.input["query"], re.I)) kinds = ([MemoryKind.DOCUMENT_CHUNK] if corpus_query else [MemoryKind.FACT, MemoryKind.DOCUMENT_CHUNK, MemoryKind.PLAYBOOK, MemoryKind.EPISODE]) - # Ask the store for a wider candidate pool, then diversify corpus - # evidence by source. Otherwise two highly similar DPO chunks can - # crowd the CoT and LoRA papers out of a cross-paper question. - hits = runtime.memory.recall(task.input["query"], scope, limit=24 if corpus_query else 8, - kinds=kinds) + # Scope-safe diversification happens inside the memory store so + # near-duplicate aliases cannot crowd complementary sources out. + hits = runtime.memory.recall(task.input["query"], scope, limit=8, kinds=kinds, + diversify=corpus_query) # Never let this very request (or the gateway's audit trail) pose # as evidence for its own answer. Older episodes remain usable. hits = [hit for hit in hits if hit.id != inbound_id and hit.metadata.get("run_id") != run_id] - if corpus_query: - diversified, per_source = [], {} - for hit in hits: - source_key = hit.sources[0].uri if hit.sources else hit.id - if per_source.get(source_key, 0) >= 3: - continue - diversified.append(hit) - per_source[source_key] = per_source.get(source_key, 0) + 1 - if len(diversified) == 8: - break - hits = diversified - else: + if not corpus_query: # A directly sourced user fact outranks a model-written prior # answer. Episodes remain useful context but must not become # the apparent source of the user's own birthday/preference. diff --git a/tests/test_document_memory.py b/tests/test_document_memory.py index 628a2ca..4bf5603 100644 --- a/tests/test_document_memory.py +++ b/tests/test_document_memory.py @@ -73,6 +73,8 @@ def test_neighbour_expansion_and_persistent_faiss_scale(tmp_path: Path): expand_neighbors=1) assert hits[0].metadata["document_id"] == result["document_id"] assert len(hits) >= 2 - assert store.last_retrieval_stats == {"index_candidates": 128, "authorized_records": 200} + assert store.last_retrieval_stats == { + "index_candidates": 128, "authorized_records": 200, "diversified": False, + } finally: store.close() diff --git a/tests/test_scope_safe_diverse_recall.py b/tests/test_scope_safe_diverse_recall.py new file mode 100644 index 0000000..ed99004 --- /dev/null +++ b/tests/test_scope_safe_diverse_recall.py @@ -0,0 +1,210 @@ +"""Scope-safe diversified recall: MMR after authorization, not before.""" +from __future__ import annotations + +from pathlib import Path + +from s13code.core.memory import MemoryKind, MemoryRecord, MemoryScope, MemoryStore, Principal, SourceRef +from s13code.core.memory.chunking import DocumentChunk +from s13code.core.memory.embeddings import DeterministicEmbedder + + +class MappedEmbedder: + """Test-only vectors so alias floods and complementary sources are controllable.""" + + def __init__(self, table: dict[str, list[float]], *, fingerprint: str = "mapped:test"): + self.table = table + self._fingerprint = fingerprint + + @property + def fingerprint(self) -> str: + return self._fingerprint + + def embed(self, text: str) -> list[float]: + return self.embed_document(text) + + def embed_document(self, text: str) -> list[float]: + return list(self.table.get(text, self.table.get("alias", [0.0, 0.0, 1.0]))) + + def embed_query(self, text: str) -> list[float]: + return list(self.table["__query__"]) + + +def _alias_vectors() -> MappedEmbedder: + alias = [1.0, 0.0, 0.0] + complementary = [0.55, 0.84, 0.0] + query = [0.95, 0.31, 0.0] + return MappedEmbedder({ + "__query__": query, + "alias": alias, + "complementary": complementary, + }) + + +SCOPE_A = MemoryScope("tenant-a", "papers", "student-01") +SCOPE_B = MemoryScope("tenant-b", "papers", "student-01") +AGENT = Principal("gateway", "gateway") + + +def chunk(text: str, ordinal: int, *, start: int = 0) -> DocumentChunk: + return DocumentChunk( + text=text, heading=None, ordinal=ordinal, + previous_ordinal=ordinal - 1 if ordinal else None, next_ordinal=None, + segmentation={"mode": "test", "outcome": "one_topic"}, + source_start_char=start, source_end_char=start + len(text), + source_start_word=start, source_end_word=start + len(text.split()), + ) + + +def ingest(store: MemoryStore, scope: MemoryScope, source_uri: str, text: str, + chunks: list[DocumentChunk]) -> dict: + return store.ingest_document( + source_text=text, prepared_text=text, chunks=chunks, + source_uri=source_uri, scope=scope, source_author="proof-indexer", + preprocessing="none", + ) + + +def test_recall_without_diversify_preserves_pure_relevance_ranking(tmp_path: Path): + store = MemoryStore(tmp_path / "memory.sqlite", embedder=DeterministicEmbedder(128)) + try: + ingest(store, SCOPE_A, "file:///papers/alpha.txt", + "alpha transformer attention parallel training", + [chunk("alpha transformer attention parallel training", 0)]) + ingest(store, SCOPE_A, "file:///papers/beta.txt", + "beta unrelated gardening soil compost", + [chunk("beta unrelated gardening soil compost", 0)]) + hits = store.recall("transformer attention training", SCOPE_A, + kinds=[MemoryKind.DOCUMENT_CHUNK], limit=1) + assert "transformer" in hits[0].text + assert store.last_retrieval_stats["diversified"] is False + finally: + store.close() + + +def test_diversify_surfaces_complementary_source_when_aliases_crowd_it_out(tmp_path: Path): + embedder = _alias_vectors() + store = MemoryStore(tmp_path / "memory.sqlite", embedder=embedder) + try: + for index in range(6): + ingest(store, SCOPE_A, f"file:///papers/dpo-copy-{index}.txt", "alias", [chunk("alias", 0)]) + ingest(store, SCOPE_A, "file:///papers/cot.txt", "complementary", [chunk("complementary", 0)]) + + plain = store.recall("ignored", SCOPE_A, kinds=[MemoryKind.DOCUMENT_CHUNK], limit=3, diversify=False) + diverse = store.recall("ignored", SCOPE_A, kinds=[MemoryKind.DOCUMENT_CHUNK], limit=3, diversify=True) + plain_sources = {hit.sources[0].uri for hit in plain} + diverse_sources = {hit.sources[0].uri for hit in diverse} + + assert len(plain_sources) == 3 + assert all(uri.startswith("file:///papers/dpo-copy-") for uri in plain_sources) + assert "file:///papers/cot.txt" in diverse_sources + assert len(diverse_sources) >= 2 + finally: + store.close() + + +def test_diversify_is_deterministic_for_the_same_store(tmp_path: Path): + store = MemoryStore(tmp_path / "memory.sqlite", embedder=DeterministicEmbedder(96)) + try: + for index in range(3): + ingest(store, SCOPE_A, f"file:///papers/paper-{index}.txt", + f"topic {index} evidence about retrieval diversification", + [chunk(f"topic {index} evidence about retrieval diversification", 0)]) + first = store.recall("retrieval diversification", SCOPE_A, + kinds=[MemoryKind.DOCUMENT_CHUNK], limit=2, diversify=True) + second = store.recall("retrieval diversification", SCOPE_A, + kinds=[MemoryKind.DOCUMENT_CHUNK], limit=2, diversify=True) + assert [record.id for record in first] == [record.id for record in second] + finally: + store.close() + + +def test_wrong_tenant_record_never_enters_diversified_results(tmp_path: Path): + store = MemoryStore(tmp_path / "memory.sqlite", embedder=DeterministicEmbedder(128)) + try: + ingest(store, SCOPE_A, "file:///papers/local.txt", + "authorized tenant retrieval diversification evidence", + [chunk("authorized tenant retrieval diversification evidence", 0)]) + ingest(store, SCOPE_B, "file:///papers/foreign.txt", + "foreign tenant highly novel secret diversification evidence", + [chunk("foreign tenant highly novel secret diversification evidence", 0)]) + + hits = store.recall("diversification evidence", SCOPE_A, + kinds=[MemoryKind.DOCUMENT_CHUNK], limit=3, diversify=True) + assert all(hit.scope.tenant_id == "tenant-a" for hit in hits) + assert all("foreign tenant" not in hit.text for hit in hits) + finally: + store.close() + + +def test_fact_lifecycle_correction_and_cross_scope_denial(tmp_path: Path): + store = MemoryStore(tmp_path / "memory.sqlite", embedder=DeterministicEmbedder(64)) + try: + first = store.write(MemoryRecord( + MemoryKind.FACT, SCOPE_A, "Mom's birthday is 15 May 2026.", + [SourceRef("api://agent/runs", "student-01")], AGENT, + )) + corrected = store.write(MemoryRecord( + MemoryKind.FACT, SCOPE_A, "Mom's birthday is 16 May 2026.", + [SourceRef("api://agent/runs", "student-01")], AGENT, + supersedes_id=first.id, + )) + authorized = store.recall("mom birthday", SCOPE_A, kinds=[MemoryKind.FACT], limit=3) + denied = store.recall("mom birthday", SCOPE_B, kinds=[MemoryKind.FACT], limit=3) + history = store.recall("mom birthday", SCOPE_A, kinds=[MemoryKind.FACT], + limit=3, include_history=True) + + assert [hit.text for hit in authorized] == ["Mom's birthday is 16 May 2026."] + assert authorized[0].sources[0].uri == "api://agent/runs" + assert denied == [] + assert store.get(first.id).status == "superseded" + assert corrected.supersedes_id == first.id + assert any(hit.id == first.id and hit.status == "superseded" for hit in history) + finally: + store.close() + + +def test_memory_search_route_exposes_diversify_flag(app_client): + app_client.app.state.s13_runtime.memory.embedder = DeterministicEmbedder(96) + scope = {"tenant_id": "course", "project_id": "papers", "user_id": "student-01"} + indexed = app_client.post("/v1/agent/documents", json={ + **scope, + "source_uri": "file:///course/example-a.md", + "text": "retrieval diversification alpha evidence", + }) + assert indexed.status_code == 200 + app_client.post("/v1/agent/documents", json={ + **scope, + "source_uri": "file:///course/example-b.md", + "text": "retrieval diversification beta evidence", + }) + response = app_client.post("/v1/agent/memory/search", json={ + **scope, + "query": "retrieval diversification evidence", + "limit": 2, + "diversify": True, + "kinds": ["document_chunk"], + }) + assert response.status_code == 200 + body = response.json() + assert body["diversified"] is True + assert body["retrieval_stats"]["diversified"] is True + assert len(body["hits"]) <= 2 + + +def test_adversarial_alias_flood_blocked_only_when_diversify_enabled(tmp_path: Path): + """Before: alias copies fill every slot. After: MMR keeps a complementary source.""" + embedder = _alias_vectors() + store = MemoryStore(tmp_path / "memory.sqlite", embedder=embedder) + try: + for index in range(6): + ingest(store, SCOPE_A, f"file:///mirror/{index}.txt", "alias", [chunk("alias", 0)]) + ingest(store, SCOPE_A, "file:///papers/react.txt", "complementary", + [chunk("complementary", 0)]) + + before = store.recall("ignored", SCOPE_A, kinds=[MemoryKind.DOCUMENT_CHUNK], limit=2, diversify=False) + after = store.recall("ignored", SCOPE_A, kinds=[MemoryKind.DOCUMENT_CHUNK], limit=2, diversify=True) + + assert all(hit.sources[0].uri.startswith("file:///mirror/") for hit in before) + assert any(hit.sources[0].uri == "file:///papers/react.txt" for hit in after) + finally: + store.close()