Skip to content
Merged
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
5 changes: 4 additions & 1 deletion docs/reference/protocols/periodic-report-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
55 changes: 51 additions & 4 deletions loopx/capabilities/periodic_report/incremental.py
Original file line number Diff line number Diff line change
Expand Up @@ -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]:
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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):
Expand All @@ -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):
Expand All @@ -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:
Expand Down Expand Up @@ -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",
Expand Down
9 changes: 6 additions & 3 deletions loopx/capabilities/periodic_report/post_writeback_hook.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand All @@ -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),
Expand Down
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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:
Expand All @@ -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,
)
Expand All @@ -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:
Expand Down Expand Up @@ -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
Expand Down
142 changes: 142 additions & 0 deletions tests/capabilities/test_periodic_report_incremental.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
)
Expand Down Expand Up @@ -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
)
Loading
Loading