Schema-default offset follows the omitted path; corrected nominations may drop or fix draft topics, never add; a replayed older queue focus never replaces a newer durable one

This commit is contained in:
Ouroboros 2026-09-22 23:27:49 +03:00
parent cd6467f4cd
commit 1bbed4e9aa
5 changed files with 63 additions and 11 deletions

View file

@ -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.

View file

@ -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

View file

@ -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]

View file

@ -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"

View file

@ -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")]