diff --git a/ouroboros/agent.py b/ouroboros/agent.py index 4451d5754..cb70ab503 100644 --- a/ouroboros/agent.py +++ b/ouroboros/agent.py @@ -352,6 +352,19 @@ class OuroborosAgent: """ try: started = getattr(self, "_task_started_ts", None) + # A queue row's focus is a REPLAY (retry clone, owner-wait restart + # handoff): it may be older than the focus the same task id already + # published durably, so it is accepted only when it is newer. + focus_kw: Dict[str, Any] = {} + if task.get("focus"): + from ouroboros.focus import compact_focus + from ouroboros.task_results import load_task_result + + incoming = compact_focus(task.get("focus")) + current = load_task_result(self.env.drive_root, str(task.get("id") or "")) + durable = compact_focus(current.get("focus")) if isinstance(current, dict) else None + if incoming and (durable is None or str(durable.get("authored_at") or "") < str(incoming.get("authored_at") or "")): + focus_kw = {"focus": incoming} running = write_task_result( self.env.drive_root, str(task.get("id") or ""), @@ -398,8 +411,9 @@ class OuroborosAgent: subagent_envelope=task.get("subagent_envelope"), configured_subagent=task.get("configured_subagent"), parent_cognitive_route=task.get("parent_cognitive_route"), subagent_availability=task.get("subagent_availability"), metadata=task.get("metadata") if isinstance(task.get("metadata"), dict) else {}, # A queue row without a focus (a retry clone, a fresh task) must not - # erase the durable focus the same task id already authored. - **({"focus": task["focus"]} if task.get("focus") else {}), + # erase the durable focus the same task id already authored, and a + # replayed older focus must not replace a newer durable one. + **focus_kw, # Ingress-captured owner-message identity (v6.73.0): persisted on the # durable record so a post-hoc "Turn into project" binds the start # message by value, never by content lookup. diff --git a/ouroboros/focus.py b/ouroboros/focus.py index ff71d00a6..172d7b686 100644 --- a/ouroboros/focus.py +++ b/ouroboros/focus.py @@ -49,7 +49,8 @@ def _safe_source(value: Any) -> Optional[Any]: # Provider-valid calls may spell an omitted optional field as "": treat # that exactly like omission so the schema and the handler agree. source = {key: item for key, item in source.items() - if not (isinstance(item, str) and not item.strip() and key != "reader")} + if not (isinstance(item, str) and not item.strip() and key != "reader") + and not (key == "offset" and item == 0 and not isinstance(item, bool))} reader = source.get("reader") if not isinstance(reader, str) or reader not in _SOURCE_FIELDS: return None diff --git a/ouroboros/room_consolidation.py b/ouroboros/room_consolidation.py index ed16cbca5..8338d2e5e 100644 --- a/ouroboros/room_consolidation.py +++ b/ouroboros/room_consolidation.py @@ -241,11 +241,20 @@ def summarize_source( # only with its corrected text: a draft whose correction failed # (and was then split) never entered the block, and a draft # nomination the correction dropped was never source-checked. - # Bound through the DRAFT call's read context: that call held the - # knowledge instruction and read the notes it nominates against, - # so its recorded reads attest the corrected block's revisions. - corrected, found = _extract_nominations(corrected, draft_knowledge if draft_knowledge is not None else knowledge) - entries.extend(found) + # The corrected block may DROP or FIX the draft's entries, never add + # topics: only the draft call read the notes it nominates against, so + # its recorded reads attest exactly those topics' revisions, and an + # entry the correction invented has no read behind it. + from ouroboros.reflection import _extract_trailing_json + + _, draft_raw = _extract_trailing_json(draft, "KNOWLEDGE_ENTRIES_JSON:") + draft_topics = {(str(e.get("topic") or ""), str(e.get("scope") or "")) + for e in (draft_raw if isinstance(draft_raw, list) else []) if isinstance(e, dict)} + corrected, raw = _extract_trailing_json(corrected, "KNOWLEDGE_ENTRIES_JSON:") + kept = [e for e in (raw if isinstance(raw, list) else []) if isinstance(e, dict) + and (str(e.get("topic") or ""), str(e.get("scope") or "")) in draft_topics] + binder = draft_knowledge if draft_knowledge is not None else knowledge + entries.extend(binder.bind_entries(kept) if binder is not None and kept else []) summaries.append(corrected.strip()) continue failure = usage["_consolidation_errors"][-1] diff --git a/tests/test_cross_focus_awareness.py b/tests/test_cross_focus_awareness.py index ee4a058ba..e419a123d 100644 --- a/tests/test_cross_focus_awareness.py +++ b/tests/test_cross_focus_awareness.py @@ -339,3 +339,29 @@ def test_focus_source_honours_the_task_contract_disabled_tools(tmp_path, monkeyp assert "FOCUS_SOURCE_UNRESOLVED" in refused and "withheld" in refused assert not list((tmp_path / "task_results" / "artifacts").glob("**/focus_source_*")) assert "focus" not in json.loads((tmp_path / "task_results" / "root.json").read_text()) + + +def test_schema_default_offset_zero_follows_the_omitted_path(): + from ouroboros.focus import normalize_focus + for ref in ({"reader": "workpad_read", "project_id": "alpha", "offset": 0, "snapshot": ""}, + {"reader": "get_task_result", "task_id": "task123", "offset": 0}): + assert "offset" not in normalize_focus("x", ref)["source_ref"] + with pytest.raises(ValueError): + normalize_focus("x", {"reader": "workpad_read", "project_id": "alpha", "offset": 3}) + + +def test_replayed_older_queue_focus_never_replaces_a_newer_durable_focus(tmp_path): + """A retry clone / restart handoff carries the queue row's focus, which may be + older than what the same task already published: the durable one wins.""" + from ouroboros.agent import OuroborosAgent + from ouroboros.focus import normalize_focus + + older = normalize_focus("older", {"reader": "recent_tasks"}, task_id="root", authored_at="2026-01-01T00:00:00+00:00") + newer = normalize_focus("newer", {"reader": "recent_tasks"}, task_id="root", authored_at="2026-01-02T00:00:00+00:00") + write_task_result(tmp_path, "root", STATUS_RUNNING, focus=newer) + agent = OuroborosAgent.__new__(OuroborosAgent) + agent.env = types.SimpleNamespace(drive_root=tmp_path, budget_drive_root=str(tmp_path)) + agent._persist_running_record({"id": "root", "description": "d", "focus": older}) + assert json.loads((tmp_path / "task_results" / "root.json").read_text())["focus"]["text"] == "newer" + agent._persist_running_record({"id": "root", "description": "d", "focus": {**newer, "text": "newest", "authored_at": "2026-01-03T00:00:00+00:00"}}) + assert json.loads((tmp_path / "task_results" / "root.json").read_text())["focus"]["text"] == "newest" diff --git a/tests/test_room_provenance_delta.py b/tests/test_room_provenance_delta.py index 096c352c7..6b7ed38b9 100644 --- a/tests/test_room_provenance_delta.py +++ b/tests/test_room_provenance_delta.py @@ -355,7 +355,8 @@ def test_nominations_come_from_the_corrected_response_not_the_draft(): if label == "Room summary": return ('draft memory\nKNOWLEDGE_ENTRIES_JSON: [{"topic":"leak","scope":"global","content":"owner approved"}]', usage, _Knowledge()) - return ('corrected memory\nKNOWLEDGE_ENTRIES_JSON: [{"topic":"kept","scope":"global","content":"owner asked"}]', + return ('corrected memory\nKNOWLEDGE_ENTRIES_JSON: [{"topic":"leak","scope":"global","content":"owner asked"},' + ' {"topic":"invented","scope":"global","content":"never read"}]', usage, _Knowledge()) content, usage = rc.summarize_source( @@ -366,5 +367,6 @@ def test_nominations_come_from_the_corrected_response_not_the_draft(): assert content == "corrected memory" # The draft's block reaches the correction under the same source check... assert "KNOWLEDGE_ENTRIES_JSON" in prompts[1][1] and "owner approved" in prompts[1][1] - # ...and only the corrected block is released. - assert [entry["topic"] for entry in usage["_knowledge_entries"]] == ["kept"] + # ...and only the corrected block is released: the draft's topic with the + # corrected content, never a topic the correction invented without a read. + assert [(e["topic"], e["content"]) for e in usage["_knowledge_entries"]] == [("leak", "owner asked")]