diff --git a/docs/reference/protocols/periodic-report-v0.md b/docs/reference/protocols/periodic-report-v0.md index 9f61e01cd6..e61fbb63aa 100644 --- a/docs/reference/protocols/periodic-report-v0.md +++ b/docs/reference/protocols/periodic-report-v0.md @@ -352,7 +352,10 @@ does not parse rendered Markdown or infer titles from fingerprints. Its interaction semantics are `attention_kind=progress`, `interaction=inform`, `delivery=surface`, and `writable=false`. Compact counts distinguish facts added since the preceding verified publication from facts whose semantic -fingerprint changed. +fingerprint changed. "Already published" is a Goal-level property: every cursor +file the Goal recorded contributes its published fingerprints, so a fact another +lane delivered is not announced again by this lane, including by a lane that has +never published. Each lane's own baseline still decides what counts as changed. Generation alone does not expose this projection. The publication candidate binds its exact SHA-256, and the Goal Channel sink may advance that binding into the diff --git a/loopx/capabilities/periodic_report/incremental.py b/loopx/capabilities/periodic_report/incremental.py index 104f8690cc..c17a1c6596 100644 --- a/loopx/capabilities/periodic_report/incremental.py +++ b/loopx/capabilities/periodic_report/incremental.py @@ -114,19 +114,22 @@ def periodic_report_fact_fingerprint(item: Mapping[str, Any]) -> str: return _canonical_digest(identity) -def _cursor_path(runtime_root: Path, goal_id: str, agent_id: str) -> Path: +def _cursor_directory(runtime_root: Path, goal_id: str) -> Path: safe_goal = _identity_text(goal_id, "goal_id") - safe_agent = _identity_text(agent_id, "agent_id") return ( runtime_root.expanduser().resolve() / "goals" / safe_goal / "periodic_reports" / "publication-cursors" - / f"{safe_agent}.json" ) +def _cursor_path(runtime_root: Path, goal_id: str, agent_id: str) -> Path: + safe_agent = _identity_text(agent_id, "agent_id") + return _cursor_directory(runtime_root, goal_id) / f"{safe_agent}.json" + + def normalize_periodic_report_publication_cursor( raw: Mapping[str, Any], ) -> dict[str, Any]: @@ -205,6 +208,31 @@ def read_periodic_report_publication_cursor( ) +def read_periodic_report_goal_publication_cursors( + *, runtime_root: Path, goal_id: str +) -> list[dict[str, Any]]: + """Read every lane's published baseline that one Goal has recorded. + + A report covers the Goal, so what the Goal already announced is a Goal-level + fact even though each lane keeps its own cursor file. A cursor that cannot + be normalized raises rather than being skipped: narrowing the published set + silently would re-announce the facts that file was suppressing. + """ + directory = _cursor_directory(runtime_root, goal_id) + if not directory.is_dir(): + return [] + cursors: list[dict[str, Any]] = [] + for path in sorted(directory.glob("*.json")): + cursor = _read_periodic_report_publication_cursor_path( + path=path, + goal_id=goal_id, + agent_id=path.stem, + ) + if cursor is not None: + cursors.append(cursor) + return cursors + + def _read_periodic_report_publication_cursor_path( *, path: Path, goal_id: str, agent_id: str ) -> dict[str, Any] | None: @@ -260,8 +288,14 @@ def select_incremental_project_progress( snapshot: Mapping[str, Any], *, cursor: Mapping[str, Any] | None, + goal_cursors: Sequence[Mapping[str, Any]] | None = None, ) -> dict[str, Any] | None: - """Select only facts absent from, or changed since, a published cursor.""" + """Select only facts the Goal has not already published in this exact form. + + ``cursor`` is the calling lane's own baseline and decides ``changed``. + ``goal_cursors`` carries every lane's baseline for the same Goal, so a fact + another lane already announced does not become new once more. + """ items = snapshot.get("items") if not isinstance(items, list): @@ -271,6 +305,15 @@ def select_incremental_project_progress( if cursor is not None: normalized = normalize_periodic_report_publication_cursor(cursor) previous = {item["source_ref"]: item for item in normalized["fact_states"]} + published: set[tuple[str, str]] = set() + # Passing this lane's own cursor in is harmless: every fact it holds is + # already in ``previous``, so it cannot hide anything the checks below miss. + for sibling in goal_cursors or []: + sibling_cursor = normalize_periodic_report_publication_cursor(sibling) + published.update( + (item["source_ref"], item["fact_fingerprint"]) + for item in sibling_cursor["fact_states"] + ) selected: list[dict[str, Any]] = [] for raw in items: if not isinstance(raw, Mapping): @@ -281,6 +324,9 @@ def select_incremental_project_progress( prior = previous.get(source_ref) if prior and prior["fact_fingerprint"] == fingerprint: continue + # Another lane already announced this exact fact to the same Goal. + if prior is None and (source_ref, fingerprint) in published: + continue item["fact_fingerprint"] = fingerprint item["change_kind"] = "changed" if prior else "added" if prior: @@ -516,6 +562,7 @@ def commit_periodic_report_publication_cursor( "normalize_periodic_report_publication_cursor", "periodic_report_fact_fingerprint", "periodic_report_incremental_baseline", + "read_periodic_report_goal_publication_cursors", "read_periodic_report_publication_cursor", "select_incremental_project_progress", "write_periodic_report_publication_candidate", diff --git a/loopx/capabilities/periodic_report/post_writeback_hook.py b/loopx/capabilities/periodic_report/post_writeback_hook.py index 389f7dcdc4..8f5bdceee9 100644 --- a/loopx/capabilities/periodic_report/post_writeback_hook.py +++ b/loopx/capabilities/periodic_report/post_writeback_hook.py @@ -25,7 +25,7 @@ from .stage_completion import derive_periodic_report_stage_completion_from_runs from .presets import build_periodic_report_preset_activation from .project_progress_snapshot import build_project_progress_snapshot_from_state -from .incremental import read_periodic_report_publication_cursor +from .incremental import read_periodic_report_goal_publication_cursors from .machine_defaults import resolve_goal_periodic_report_subscription from .machine_store import read_periodic_report_machine_defaults from .triggers import build_periodic_report_trigger_decision @@ -282,10 +282,12 @@ def build_periodic_report_post_writeback_projection( if receipt is None: return {} result: dict[str, object] = {"stage_completion": receipt} - publication_cursor = read_periodic_report_publication_cursor( + goal_cursors = read_periodic_report_goal_publication_cursors( runtime_root=runtime_root, goal_id=goal_id, - agent_id=normalized_agent_id, + ) + publication_cursor = next( + (c for c in goal_cursors if c["agent_id"] == normalized_agent_id), None ) available_capabilities = payload.get("available_capabilities") if available_capabilities is None and isinstance(payload.get("turn"), Mapping): @@ -298,6 +300,7 @@ def build_periodic_report_post_writeback_projection( agent_id=normalized_agent_id, completed_at=str(receipt["completed_at"]), publication_cursor=publication_cursor, + goal_cursors=goal_cursors, available_capabilities=available_capabilities, rollout_events=load_rollout_events( rollout_event_log_path(runtime_root, goal_id), diff --git a/loopx/capabilities/periodic_report/project_progress_snapshot.py b/loopx/capabilities/periodic_report/project_progress_snapshot.py index 5147f8252b..4beb4b7c8c 100644 --- a/loopx/capabilities/periodic_report/project_progress_snapshot.py +++ b/loopx/capabilities/periodic_report/project_progress_snapshot.py @@ -1,6 +1,6 @@ from __future__ import annotations -from collections.abc import Mapping +from collections.abc import Mapping, Sequence from datetime import datetime from pathlib import Path from typing import Any @@ -62,6 +62,7 @@ def build_project_progress_snapshot( agent_id: str, completed_at: str, publication_cursor: Mapping[str, Any] | None = None, + goal_cursors: Sequence[Mapping[str, Any]] | None = None, available_capabilities: Any = None, rollout_events: list[dict[str, Any]] | None = None, ) -> dict[str, Any] | None: @@ -83,6 +84,7 @@ def build_project_progress_snapshot( agent_id=agent_id, completed_at=completed_at, publication_cursor=publication_cursor, + goal_cursors=goal_cursors, available_capabilities=available_capabilities, rollout_events=rollout_events, ) @@ -97,6 +99,7 @@ def build_project_progress_snapshot_from_state( agent_id: str, completed_at: str, publication_cursor: Mapping[str, Any] | None = None, + goal_cursors: Sequence[Mapping[str, Any]] | None = None, available_capabilities: Any = None, rollout_events: list[dict[str, Any]] | None = None, ) -> dict[str, Any] | None: @@ -214,10 +217,13 @@ def not_after_stage(item: Mapping[str, Any]) -> bool: "language": "zh-CN", "items": progress_items, } - if publication_cursor is not None: + if publication_cursor is not None or goal_cursors: + # A lane that has never published still must not re-announce what the + # Goal already announced, so the peer baseline applies without one. incremental = select_incremental_project_progress( snapshot, cursor=publication_cursor, + goal_cursors=goal_cursors, ) if incremental is None: return None diff --git a/tests/capabilities/test_periodic_report_incremental.py b/tests/capabilities/test_periodic_report_incremental.py index bcbfe995a4..098e5ec9d8 100644 --- a/tests/capabilities/test_periodic_report_incremental.py +++ b/tests/capabilities/test_periodic_report_incremental.py @@ -14,6 +14,7 @@ build_periodic_report_publication_candidate, commit_periodic_report_publication_cursor, periodic_report_incremental_baseline, + read_periodic_report_goal_publication_cursors, read_periodic_report_publication_cursor, select_incremental_project_progress, ) @@ -1102,3 +1103,144 @@ def test_snapshot_next_action_prefers_the_peer_lane_over_an_unowned_row( assert snapshot is not None assert _next_action_refs(snapshot) == ["todo:todo_peer"] + + +PEER_AGENT_ID = "peer-agent" + + +def _commit( + runtime: Path, + *, + agent_id: str, + facts: list[dict[str, object]], + generation_id: str, +) -> dict[str, object]: + candidate = build_periodic_report_publication_candidate( + goal_id=GOAL_ID, + agent_id=agent_id, + generation_id=generation_id, + trigger_receipt=_trigger(f"trigger_{generation_id}"), + facts=facts, + baseline=None, + ) + return commit_periodic_report_publication_cursor( + runtime_root=runtime, + candidate=candidate, + publication_id=f"goal-channel:{generation_id}", + delivered_at="2026-08-01T09:00:00Z", + covered_until="2026-08-01T08:00:00Z", + ) + + +def test_goal_cursors_drop_a_fact_another_lane_already_published( + tmp_path: Path, +) -> None: + """One Goal announces each fact once, whoever the reporting lane is.""" + + runtime = tmp_path / "runtime" + peer = _commit( + runtime, + agent_id=PEER_AGENT_ID, + generation_id="generation_peer", + facts=[_item("todo:a", title="A completed", summary="A is done.")], + ) + selected = select_incremental_project_progress( + _snapshot( + [ + _item("todo:a", title="A completed", summary="A is done."), + _item("todo:z", title="Z completed", summary="Z is new."), + ] + ), + cursor=None, + goal_cursors=[peer], + ) + + assert selected is not None + assert [item["source_ref"] for item in selected["items"]] == ["todo:z"] + assert selected["items"][0]["change_kind"] == "added" + + +def test_goal_cursors_keep_a_fact_a_peer_published_with_different_content( + tmp_path: Path, +) -> None: + runtime = tmp_path / "runtime" + peer = _commit( + runtime, + agent_id=PEER_AGENT_ID, + generation_id="generation_peer_text", + facts=[_item("todo:a", title="A completed", summary="First wording.")], + ) + selected = select_incremental_project_progress( + _snapshot([_item("todo:a", title="A completed", summary="Second wording.")]), + cursor=None, + goal_cursors=[peer], + ) + + assert selected is not None + assert [item["source_ref"] for item in selected["items"]] == ["todo:a"] + + +def test_goal_cursors_do_not_hide_this_lanes_own_changed_fact( + tmp_path: Path, +) -> None: + runtime = tmp_path / "runtime" + own = _commit( + runtime, + agent_id=AGENT_ID, + generation_id="generation_own", + facts=[_item("todo:a", title="A completed", summary="A is done.")], + ) + peer = _commit( + runtime, + agent_id=PEER_AGENT_ID, + generation_id="generation_peer_again", + facts=[_item("todo:a", title="A completed", summary="A is done.")], + ) + selected = select_incremental_project_progress( + _snapshot([_item("todo:a", title="A completed", summary="A reopened.")]), + cursor=own, + goal_cursors=[own, peer], + ) + + assert selected is not None + assert selected["items"][0]["change_kind"] == "changed" + assert selected["items"][0]["previous_fact_fingerprint"] == own["fact_states"][0][ + "fact_fingerprint" + ] + + +def test_goal_cursor_reader_returns_every_lane_and_fails_closed( + tmp_path: Path, +) -> None: + runtime = tmp_path / "runtime" + assert ( + read_periodic_report_goal_publication_cursors( + runtime_root=runtime, goal_id=GOAL_ID + ) + == [] + ) + + _commit( + runtime, + agent_id=AGENT_ID, + generation_id="generation_reader_own", + facts=[_item("todo:a", title="A completed", summary="A is done.")], + ) + _commit( + runtime, + agent_id=PEER_AGENT_ID, + generation_id="generation_reader_peer", + facts=[_item("todo:b", title="B completed", summary="B is done.")], + ) + cursors = read_periodic_report_goal_publication_cursors( + runtime_root=runtime, goal_id=GOAL_ID + ) + assert [cursor["agent_id"] for cursor in cursors] == [AGENT_ID, PEER_AGENT_ID] + + sibling = runtime / "goals" / GOAL_ID / "periodic_reports" / "publication-cursors" + corrupt = sibling / "third-agent.json" + corrupt.write_text('{"schema_version": "unexpected"}', encoding="utf-8") + with pytest.raises(ValueError, match="publication cursor must use"): + read_periodic_report_goal_publication_cursors( + runtime_root=runtime, goal_id=GOAL_ID + ) diff --git a/tests/control_plane/test_post_writeback_capability_hooks.py b/tests/control_plane/test_post_writeback_capability_hooks.py index 8fdfc8edfa..d6ce6cfeef 100644 --- a/tests/control_plane/test_post_writeback_capability_hooks.py +++ b/tests/control_plane/test_post_writeback_capability_hooks.py @@ -2606,10 +2606,65 @@ def test_periodic_report_projection_carries_peer_lane_progress( ) items = projection["project_progress"]["items"] - assert [ - (item["content_kind"], item["title"]) for item in items - ] == [ + assert [(item["content_kind"], item["title"]) for item in items] == [ ("outcome", "Land the reporting lane's change."), ("outcome", "Land a peer lane's change."), ("next_action", "Next action"), ] + + +def test_periodic_report_projection_drops_facts_a_peer_already_delivered( + tmp_path, +) -> None: + state_text = """# Goal + +## User Todo + +## Agent Todo + +- [x] Land the first bounded change. + +- [x] Land the second bounded change. + +- [ ] Continue the next bounded change. + +""" + runtime_root, registry_path = _projection_goal_fixture( + tmp_path, + state_text=state_text, + runs=[ + _successor_ack_run(), + _closed_vision_run(), + ], + ) + + def projection() -> dict[str, object]: + return build_periodic_report_post_writeback_projection( + payload={"state": {"path": str(tmp_path / "goal.md")}}, + registry_path=registry_path, + runtime_root=runtime_root, + goal_id="goal-1", + agent_id="agent-1", + ) + + first = projection() + facts = first["project_progress"]["items"] + assert len(facts) == 3 + + candidate = build_periodic_report_publication_candidate( + goal_id="goal-1", + agent_id="agent-2", + generation_id="generation-peer-delivery", + trigger_receipt={"coalesced_trigger_ids": ["trigger-peer-delivery"]}, + facts=facts, + baseline=None, + ) + commit_periodic_report_publication_cursor( + runtime_root=runtime_root, + candidate=candidate, + publication_id="goal-channel:peer-delivery", + delivered_at="2026-08-30T12:00:00Z", + covered_until="2026-08-30T11:00:00Z", + ) + + assert "project_progress" not in projection()