diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 4b54e10..3f62760 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -14,6 +14,8 @@ jobs: python-version: ["3.11", "3.12"] steps: - uses: actions/checkout@v4 + with: + fetch-depth: 0 - uses: actions/setup-python@v5 with: python-version: ${{ matrix.python-version }} @@ -28,12 +30,14 @@ jobs: runs-on: ubuntu-latest steps: - uses: actions/checkout@v4 + with: + fetch-depth: 0 - uses: actions/setup-python@v5 with: python-version: "3.12" - uses: actions/setup-node@v4 with: - node-version: "20" + node-version: "22" cache: npm cache-dependency-path: extension/package-lock.json - name: Scrub check diff --git a/engine/env.example b/engine/env.example index 09d972c..b8b8bfe 100644 --- a/engine/env.example +++ b/engine/env.example @@ -39,3 +39,8 @@ PRODUCER_B_DIM=128 # Live Cursor SDK scenarios (optional): set the Cursor integrations API key in a # local .env that is never committed. Do not place that key in chat-compressor.env. + +# Cross-turn inject dedup (default on). When StateNode.meta.recipient_id is set, +# the inject ledger is partitioned per recipient (CC-2); hops reset suppression +# and warmup (CC-3..CC-5). Absent recipient_id keeps session-scoped 0.2.0 behavior. +# CHAT_COMPRESSOR_CROSS_TURN_DEDUP=1 diff --git a/engine/src/chat_compressor/handle.py b/engine/src/chat_compressor/handle.py index d22cfad..170c309 100644 --- a/engine/src/chat_compressor/handle.py +++ b/engine/src/chat_compressor/handle.py @@ -175,17 +175,58 @@ def flush_graph(self) -> str | None: self._last_graph_path = str(snap) return self._last_graph_path - def sample_for(self, target: str, query: str | None = None) -> SampledPayload: + def sample_for( + self, + target: str, + query: str | None = None, + *, + recipient_id: str | None = None, + ) -> SampledPayload: """cursor-sdk => packed HOT_SET/typed/ranked text. local: may return C_B floats.""" node = self.latest() q = (query or "").strip() or self._last_user_query() hot = self.graph.hot_set(query=q or None) window = self.graph.window_text() typed = self.graph.typed_projection(q or None, hot_set=hot) - history = load_inject_history(self._agent_dir()) + + # Resolve recipient: explicit arg, else latest StateNode.meta (CC-1). + rid = recipient_id + if rid is None and node is not None: + meta_rid = (node.meta or {}).get("recipient_id") + if meta_rid is not None and str(meta_rid).strip(): + rid = str(meta_rid).strip() + elif rid is not None: + rid = str(rid).strip() or None + + # CC-2: partition inject ledger by recipient_id; absent ⇒ session ledger. + history = load_inject_history(self._agent_dir(), recipient_id=rid) t = int(node.t) if node is not None else 0 + + # CC-3/CC-5: previous recipient from parent lineage node. + prev_rid: str | None = None + if node is not None and node.parent_id: + try: + parent = self.store.load(node.parent_id) + except KeyError: + parent = None + if parent is not None: + raw_prev = (parent.meta or {}).get("recipient_id") + if raw_prev is not None and str(raw_prev).strip(): + prev_rid = str(raw_prev).strip() + + if rid is None: + # Absent recipient_id ⇒ exact 0.2.0 session-scoped behavior. + recipient_changed = False + recipient_continued = True + recipient_t = t + else: + recipient_changed = prev_rid != rid + recipient_continued = prev_rid == rid + recipient_t = self._recipient_turn_count(rid) + novelty = rolling_novelty(history, k=3) - budget = adaptive_budget(t, novelty, cap=forward_budget()) + # CC-5: warmup against per-recipient turn counter, not session t alone. + budget = adaptive_budget(recipient_t, novelty, cap=forward_budget()) if not cross_turn_dedup_enabled(): budget = forward_budget() last = history[-1] if history else {} @@ -197,7 +238,12 @@ def sample_for(self, target: str, query: str | None = None) -> SampledPayload: openitem_changed = self.graph.openitem_signature() != prev_sig node_superseded = self.graph.supersede_count() > int(last.get("supersede_count") or 0) recent = recent_line_hashes(history, k=3) - allow_skip = bool(cross_turn_dedup_enabled() and t > WARMUP_TURNS) + # CC-4: never allow_skip on a recipient's first turn (or hop). + allow_skip = bool( + cross_turn_dedup_enabled() + and recipient_continued + and recipient_t > WARMUP_TURNS + ) pack_kwargs = { "hot_set": hot, "window_text": window, @@ -208,6 +254,7 @@ def sample_for(self, target: str, query: str | None = None) -> SampledPayload: "recent_hashes": recent, "openitem_changed": openitem_changed, "node_superseded": node_superseded, + "recipient_changed": recipient_changed, "allow_skip": allow_skip, } t0 = time.perf_counter() @@ -224,6 +271,7 @@ def sample_for(self, target: str, query: str | None = None) -> SampledPayload: { "state_id": None if node is None else node.state_id, "t": t, + "recipient_t": recipient_t, "hashes": list(sampled.line_hashes), "text": (sampled.text or "")[:8000], "openitem_sig": self.graph.openitem_signature(), @@ -232,6 +280,7 @@ def sample_for(self, target: str, query: str | None = None) -> SampledPayload: "novel_tokens": int(sampled.novel_tokens), "dup_suppressed_tokens": int(sampled.dup_suppressed_tokens), }, + recipient_id=rid, ) return sampled if target.startswith("local:"): @@ -244,6 +293,17 @@ def sample_for(self, target: str, query: str | None = None) -> SampledPayload: finally: self.last_sample_ms = (time.perf_counter() - t0) * 1000.0 + def _recipient_turn_count(self, recipient_id: str) -> int: + """How many lineage nodes record this recipient_id (CC-5 warmup key).""" + rid = str(recipient_id).strip() + if not rid: + return 0 + return sum( + 1 + for n in self.store.lineage(self.agent_id) + if str((n.meta or {}).get("recipient_id") or "").strip() == rid + ) + def _last_user_query(self) -> str: turns = [ n diff --git a/engine/src/chat_compressor/pack.py b/engine/src/chat_compressor/pack.py index 405eac1..516b5e9 100644 --- a/engine/src/chat_compressor/pack.py +++ b/engine/src/chat_compressor/pack.py @@ -82,6 +82,7 @@ def pack_forward( recent_hashes: set[str] | None = None, openitem_changed: bool = True, node_superseded: bool = False, + recipient_changed: bool = False, allow_skip: bool = False, marginal_jaccard: float = MARGINAL_JACCARD, skip_floor_tokens: int = SKIP_FLOOR_TOKENS, @@ -95,7 +96,8 @@ def pack_forward( method = "hot_set" dup_suppressed = 0 suppress = set(recent_hashes or ()) - if node_superseded or not cross_turn_dedup_enabled(): + # CC-3: recipient change clears suppression like supersede (hop safety). + if node_superseded or recipient_changed or not cross_turn_dedup_enabled(): suppress = set() def _blocked(text: str) -> bool: @@ -188,6 +190,7 @@ def _fits(block: str) -> bool: and cross_turn_dedup_enabled() and not openitem_changed and not node_superseded + and not recipient_changed and packed < skip_floor_tokens ) if skip: diff --git a/engine/src/chat_compressor/store.py b/engine/src/chat_compressor/store.py index ee1b45b..1c1d4d0 100644 --- a/engine/src/chat_compressor/store.py +++ b/engine/src/chat_compressor/store.py @@ -224,40 +224,102 @@ def inject_history_path(agent_dir: str | Path) -> Path: return Path(agent_dir) / INJECT_HISTORY_NAME -def load_inject_history(agent_dir: str | Path) -> list[dict[str, Any]]: +def _empty_inject_doc() -> dict[str, Any]: + return {"turns": [], "recipients": {}} + + +def _read_inject_doc(agent_dir: str | Path) -> dict[str, Any]: + """Load inject ledger document. Supports legacy list and turns-only shapes.""" path = inject_history_path(agent_dir) if not path.is_file(): - return [] + return _empty_inject_doc() try: raw = json.loads(path.read_text(encoding="utf-8")) except (OSError, json.JSONDecodeError): - return [] - if isinstance(raw, dict): - rows = raw.get("turns") - return list(rows) if isinstance(rows, list) else [] + return _empty_inject_doc() if isinstance(raw, list): - return raw - return [] - - -def save_inject_history(agent_dir: str | Path, turns: list[dict[str, Any]]) -> Path: + return {"turns": list(raw), "recipients": {}} + if not isinstance(raw, dict): + return _empty_inject_doc() + turns = raw.get("turns") + turns_list = list(turns) if isinstance(turns, list) else [] + recipients_raw = raw.get("recipients") + recipients: dict[str, list[dict[str, Any]]] = {} + if isinstance(recipients_raw, dict): + for key, val in recipients_raw.items(): + if isinstance(val, list): + recipients[str(key)] = list(val) + return {"turns": turns_list, "recipients": recipients} + + +def _write_inject_doc(agent_dir: str | Path, doc: dict[str, Any]) -> Path: dest = inject_history_path(agent_dir) dest.parent.mkdir(parents=True, exist_ok=True) - kept = turns[-INJECT_HISTORY_KEEP:] + turns = list(doc.get("turns") or [])[-INJECT_HISTORY_KEEP:] + recipients_in = doc.get("recipients") or {} + recipients: dict[str, list[dict[str, Any]]] = {} + if isinstance(recipients_in, dict): + for key, val in recipients_in.items(): + if isinstance(val, list) and val: + recipients[str(key)] = list(val)[-INJECT_HISTORY_KEEP:] + payload: dict[str, Any] = {"turns": turns} + if recipients: + payload["recipients"] = recipients dest.write_text( - json.dumps({"turns": kept}, ensure_ascii=False, indent=2) + "\n", + json.dumps(payload, ensure_ascii=False, indent=2) + "\n", encoding="utf-8", ) return dest +def load_inject_history( + agent_dir: str | Path, + recipient_id: str | None = None, +) -> list[dict[str, Any]]: + """Return inject turns. Absent recipient_id ⇒ legacy session-scoped ledger (0.2.0).""" + doc = _read_inject_doc(agent_dir) + if recipient_id is None: + return list(doc.get("turns") or []) + rid = str(recipient_id).strip() + if not rid: + return list(doc.get("turns") or []) + recipients = doc.get("recipients") or {} + rows = recipients.get(rid) + return list(rows) if isinstance(rows, list) else [] + + +def save_inject_history( + agent_dir: str | Path, + turns: list[dict[str, Any]], + recipient_id: str | None = None, +) -> Path: + """Persist inject turns. With recipient_id, write that partition only.""" + doc = _read_inject_doc(agent_dir) + kept = list(turns)[-INJECT_HISTORY_KEEP:] + if recipient_id is None or not str(recipient_id).strip(): + doc["turns"] = kept + else: + recipients = dict(doc.get("recipients") or {}) + recipients[str(recipient_id).strip()] = kept + doc["recipients"] = recipients + return _write_inject_doc(agent_dir, doc) + + def append_inject_history( agent_dir: str | Path, row: dict[str, Any], + recipient_id: str | None = None, ) -> list[dict[str, Any]]: - turns = load_inject_history(agent_dir) - turns.append(row) - save_inject_history(agent_dir, turns) + """Append one inject row to the session ledger or a recipient partition (CC-2).""" + rid = str(recipient_id).strip() if recipient_id is not None else None + if rid == "": + rid = None + turns = load_inject_history(agent_dir, recipient_id=rid) + entry = dict(row) + if rid is not None: + entry.setdefault("recipient_id", rid) + turns.append(entry) + save_inject_history(agent_dir, turns, recipient_id=rid) return turns diff --git a/engine/src/chat_compressor/translate/vocab_bridge.py b/engine/src/chat_compressor/translate/vocab_bridge.py index f66d085..58203b7 100644 --- a/engine/src/chat_compressor/translate/vocab_bridge.py +++ b/engine/src/chat_compressor/translate/vocab_bridge.py @@ -240,6 +240,7 @@ def sample_text( recent_hashes: set[str] | None = None, openitem_changed: bool = True, node_superseded: bool = False, + recipient_changed: bool = False, allow_skip: bool = False, ) -> SampledPayload: """Primary forward channel: HOT_SET → typed → ranked chunks; P1 debug-only.""" @@ -272,6 +273,7 @@ def sample_text( recent_hashes=recent_hashes, openitem_changed=openitem_changed, node_superseded=node_superseded, + recipient_changed=recipient_changed, allow_skip=allow_skip, ) method = packed.method diff --git a/engine/tests/test_handle.py b/engine/tests/test_handle.py index 7c6a495..2cbfa64 100644 --- a/engine/tests/test_handle.py +++ b/engine/tests/test_handle.py @@ -103,3 +103,30 @@ def test_step_persists_recipient_meta_through_lineage(tmp_path) -> None: reloaded = store.load(latest.state_id) assert reloaded.meta == latest.meta + + +def test_sample_for_recipient_hop_resets_dedup_and_budget(tmp_path, monkeypatch) -> None: + """CC-2..CC-5 smoke: hop clears suppress path and restores full budget.""" + monkeypatch.setenv("CHAT_COMPRESSOR_CROSS_TURN_DEDUP", "1") + monkeypatch.setenv("CHAT_COMPRESSOR_FORWARD_BUDGET", "1024") + store = StateStore(tmp_path / "state") + handle = PersistentAgentHandle( + agent_id="hop-h", + store=store, + producer=EmbeddingProducer(d=64, k_max=8), + k_max=8, + ) + for i in range(5): + handle.step( + f'Create todo "item-{i}" and keep milk bread groceries on the list. substance {i}.', + recipient_id="model-a", + ) + handle.sample_for("cursor-sdk") + handle.step( + 'Hop turn: keep milk bread groceries visible for the new model.', + recipient_id="model-b", + ) + hop = handle.sample_for("cursor-sdk") + assert hop.method != "skip" + assert hop.budget == 1024 + assert hop.packed_tokens > 0 diff --git a/engine/tests/test_m1_hop_safety.py b/engine/tests/test_m1_hop_safety.py new file mode 100644 index 0000000..ba67ddf --- /dev/null +++ b/engine/tests/test_m1_hop_safety.py @@ -0,0 +1,174 @@ +"""M1 exit: hop at turn 20 delivers full unsuppressed full-budget payload; no-hop matches 0.2.0.""" + +from __future__ import annotations + +from pathlib import Path + +import pytest + +from chat_compressor.handle import PersistentAgentHandle +from chat_compressor.pack import WARMUP_TURNS, adaptive_budget, forward_budget +from chat_compressor.producer import EmbeddingProducer +from chat_compressor.store import ( + load_inject_history, + recent_line_hashes, + StateStore, +) + + +def _handle(tmp_path: Path, agent_id: str) -> PersistentAgentHandle: + store = StateStore(tmp_path / "state") + return PersistentAgentHandle( + agent_id=agent_id, + store=store, + producer=EmbeddingProducer(d=64, k_max=8), + k_max=8, + ) + + +def _turn_text(i: int) -> str: + # Stable open items early, then distinctive content in late turns for hash checks. + return ( + f'Turn {i}: Create todo "buy groceries" and add "milk" and "bread". ' + f'Also note unique-marker-turn-{i} for dedup tracking with substance.' + ) + + +def test_m1_hop_at_turn_20_full_unsuppressed_full_budget( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """Scripted hop: 19 turns model-a, turn 20 model-b ⇒ full / unsuppressed / full budget.""" + monkeypatch.setenv("CHAT_COMPRESSOR_CROSS_TURN_DEDUP", "1") + monkeypatch.delenv("CHAT_COMPRESSOR_INJECT_P1", raising=False) + monkeypatch.setenv("CHAT_COMPRESSOR_FORWARD_BUDGET", "1024") + + handle = _handle(tmp_path, "hop-session") + model_a = "model-a" + model_b = "model-b" + + late_a_hashes: set[str] = set() + for i in range(1, 20): + handle.step(_turn_text(i), recipient_id=model_a, recipient_version="v1") + payload = handle.sample_for("cursor-sdk") + assert payload.kind == "text" + if i >= 17 and payload.line_hashes: + late_a_hashes.update(payload.line_hashes) + + assert late_a_hashes, "expected inject hashes from late model-a turns" + + # Confirm model-a ledger would suppress those hashes on a continued turn. + hist_a = load_inject_history(handle._agent_dir(), recipient_id=model_a) + suppress_a = recent_line_hashes(hist_a, k=3) + assert late_a_hashes & suppress_a, "late A hashes should sit in A's recent suppress set" + + # Turn 20: hop to model-b. + handle.step(_turn_text(20), recipient_id=model_b, recipient_version="v1") + hop = handle.sample_for("cursor-sdk") + + assert hop.method != "skip", "CC-4: first turn for recipient must not skip" + assert hop.text, "hop payload must be non-empty" + assert hop.budget == forward_budget(), "CC-5: late joiner gets full (warmup) budget" + assert hop.packed_tokens > 0 + + # CC-2/CC-3: B's ledger was empty / suppress cleared ⇒ content not hole-punched. + # Markers from recent A turns should still be packable for B. + body = hop.text.lower() + assert ( + "unique-marker-turn-17" in body + or "unique-marker-turn-18" in body + or "unique-marker-turn-19" in body + or "groceries" in body + or "bread" in body + ), "hop payload should include context A had already injected" + + # B partition is independent of A. + hist_b = load_inject_history(handle._agent_dir(), recipient_id=model_b) + assert hist_b, "model-b should have its own inject rows after hop sample" + assert load_inject_history(handle._agent_dir(), recipient_id=model_a), "model-a ledger preserved" + + # recipient_t for B is 1 ⇒ adaptive_budget equals full cap. + assert adaptive_budget(1, 0.0, cap=1024) == 1024 + assert handle._recipient_turn_count(model_b) == 1 + assert handle._recipient_turn_count(model_a) == 19 + + +def test_m1_no_hop_matches_legacy_token_accounting( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """No-hop with recipient_id tracks same budgets/tokens as absent-recipient (0.2.0) path.""" + monkeypatch.setenv("CHAT_COMPRESSOR_CROSS_TURN_DEDUP", "1") + monkeypatch.delenv("CHAT_COMPRESSOR_INJECT_P1", raising=False) + monkeypatch.setenv("CHAT_COMPRESSOR_FORWARD_BUDGET", "1024") + + legacy = _handle(tmp_path, "legacy") + tagged = _handle(tmp_path, "tagged") + + legacy_rows: list[dict] = [] + tagged_rows: list[dict] = [] + + for i in range(1, 21): + text = _turn_text(i) + legacy.step(text) + tagged.step(text, recipient_id="model-a", recipient_version="v1") + + lp = legacy.sample_for("cursor-sdk") + tp = tagged.sample_for("cursor-sdk") + + legacy_rows.append( + { + "t": i, + "method": lp.method, + "budget": lp.budget, + "packed_tokens": lp.packed_tokens, + "novel_tokens": lp.novel_tokens, + "dup_suppressed_tokens": lp.dup_suppressed_tokens, + } + ) + tagged_rows.append( + { + "t": i, + "method": tp.method, + "budget": tp.budget, + "packed_tokens": tp.packed_tokens, + "novel_tokens": tp.novel_tokens, + "dup_suppressed_tokens": tp.dup_suppressed_tokens, + } + ) + + # Budget schedule must match (session t == per-recipient t when single recipient). + assert [r["budget"] for r in legacy_rows] == [r["budget"] for r in tagged_rows] + # Token accounting within tight equality for identical inputs. + assert [r["packed_tokens"] for r in legacy_rows] == [r["packed_tokens"] for r in tagged_rows] + assert [r["method"] for r in legacy_rows] == [r["method"] for r in tagged_rows] + + # After warmup, budgets may decay; first WARMUP_TURNS stay at full cap. + for r in legacy_rows[:WARMUP_TURNS]: + assert r["budget"] == forward_budget() + + +def test_cc4_first_recipient_sample_never_skip( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("CHAT_COMPRESSOR_CROSS_TURN_DEDUP", "1") + handle = _handle(tmp_path, "first") + # Build a session that would be skip-eligible for a continued recipient. + for i in range(1, 6): + handle.step(_turn_text(i), recipient_id="model-a") + handle.sample_for("cursor-sdk") + handle.step(_turn_text(6), recipient_id="model-b") + payload = handle.sample_for("cursor-sdk") + assert payload.method != "skip" + assert payload.budget == forward_budget() + + +def test_absent_recipient_keeps_session_ledger(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + """Absent recipient_id ⇒ 0.2.0 session-scoped inject_history.turns.""" + monkeypatch.setenv("CHAT_COMPRESSOR_CROSS_TURN_DEDUP", "1") + handle = _handle(tmp_path, "legacy-ledger") + handle.step(_turn_text(1)) + handle.sample_for("cursor-sdk") + handle.step(_turn_text(2)) + handle.sample_for("cursor-sdk") + legacy = load_inject_history(handle._agent_dir()) + assert len(legacy) >= 1 + assert load_inject_history(handle._agent_dir(), recipient_id="nope") == [] diff --git a/engine/tests/test_pack.py b/engine/tests/test_pack.py index 62bae6a..a5a5fab 100644 --- a/engine/tests/test_pack.py +++ b/engine/tests/test_pack.py @@ -30,3 +30,61 @@ def test_typed_then_chunks_order() -> None: chunk_idx = packed.text.index("ranked chunk about SQL joins") assert hot_idx < path_idx < chunk_idx assert packed.method == "query-pack" + + +def test_recipient_changed_clears_suppression_and_blocks_skip() -> None: + """CC-3/CC-4: recipient change resets suppress; never skip on hop.""" + from chat_compressor.pack import line_hash, pack_forward + + line = "OpenItem: open: milk" + h = line_hash(line) + # Same hashes would suppress without recipient_changed. + suppressed = pack_forward( + hot_set="", + typed_lines=[line], + budget=1024, + recent_hashes={h}, + openitem_changed=False, + node_superseded=False, + recipient_changed=False, + allow_skip=True, + skip_floor_tokens=64, + ) + # With typed content suppressed and packed small, skip is allowed. + assert suppressed.method == "skip" or h not in suppressed.line_hashes + + hopped = pack_forward( + hot_set="", + typed_lines=[line], + budget=1024, + recent_hashes={h}, + openitem_changed=False, + node_superseded=False, + recipient_changed=True, + allow_skip=True, + skip_floor_tokens=64, + ) + assert hopped.method != "skip" + assert line.lower() in hopped.text.lower() or any( + line_hash(x) == h for x in hopped.text.splitlines() if x.strip() + ) + assert h in hopped.line_hashes + + +def test_first_recipient_turn_never_skips() -> None: + """CC-4: allow_skip false / recipient_changed true ⇒ no skip even under floor.""" + from chat_compressor.pack import pack_forward + + packed = pack_forward( + hot_set="x", + typed_lines=[], + budget=1024, + recent_hashes=set(), + openitem_changed=False, + node_superseded=False, + recipient_changed=True, + allow_skip=True, + skip_floor_tokens=10_000, + ) + assert packed.method != "skip" + assert packed.text diff --git a/engine/tests/test_store.py b/engine/tests/test_store.py index c5ac537..ba4c5da 100644 --- a/engine/tests/test_store.py +++ b/engine/tests/test_store.py @@ -77,3 +77,48 @@ def test_recipient_meta_roundtrip_and_lineage(tmp_path) -> None: assert chain[0].meta["route_decision_id"] == "urn:mg:routedecision:a1b2c3" assert chain[1].meta["recipient_version"] == "v2" + + +def test_per_recipient_inject_ledger_partitions(tmp_path) -> None: + """CC-2: inject ledger partitions by recipient_id; absent ⇒ session ledger.""" + from chat_compressor.store import ( + append_inject_history, + load_inject_history, + recent_line_hashes, + ) + + agent = tmp_path / "agent" + agent.mkdir() + + # Legacy / no-recipient path (0.2.0). + append_inject_history(agent, {"t": 1, "hashes": ["aaaa"], "packed_tokens": 10, "novel_tokens": 10}) + append_inject_history(agent, {"t": 2, "hashes": ["bbbb"], "packed_tokens": 8, "novel_tokens": 4}) + legacy = load_inject_history(agent) + assert len(legacy) == 2 + assert recent_line_hashes(legacy, k=3) == {"aaaa", "bbbb"} + + # Recipient A and B are isolated. + append_inject_history( + agent, + {"t": 3, "hashes": ["hashA1", "hashA2"], "packed_tokens": 12, "novel_tokens": 12}, + recipient_id="model-a", + ) + append_inject_history( + agent, + {"t": 4, "hashes": ["hashA3"], "packed_tokens": 6, "novel_tokens": 2}, + recipient_id="model-a", + ) + append_inject_history( + agent, + {"t": 5, "hashes": ["hashB1"], "packed_tokens": 9, "novel_tokens": 9}, + recipient_id="model-b", + ) + + hist_a = load_inject_history(agent, recipient_id="model-a") + hist_b = load_inject_history(agent, recipient_id="model-b") + assert [h for row in hist_a for h in row["hashes"]] == ["hashA1", "hashA2", "hashA3"] + assert [h for row in hist_b for h in row["hashes"]] == ["hashB1"] + # Legacy session ledger untouched by recipient partitions. + assert len(load_inject_history(agent)) == 2 + # New recipient starts empty. + assert load_inject_history(agent, recipient_id="model-c") == []