Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
129 changes: 128 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand All @@ -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`.
75 changes: 67 additions & 8 deletions s13code/core/memory/store.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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 > ?)"
Expand All @@ -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.
Expand All @@ -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"])
Expand All @@ -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}
Expand All @@ -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}

Expand Down
10 changes: 7 additions & 3 deletions s13code/routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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]}

Expand Down
22 changes: 5 additions & 17 deletions s13code/runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
4 changes: 3 additions & 1 deletion tests/test_document_memory.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Loading