From 00fae001fb994ef4d29a1aaa202589ed22f14597 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 17 Sep 2026 19:32:34 +0800 Subject: [PATCH 1/2] fix(replan): unify typed checkpoints and preserve successor causality Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../goals/goal_frontier/__init__.py | 2 +- .../goals/goal_frontier/ack_policy.py | 82 ++++++++--------- .../goals/goal_frontier/long_todo_chain.py | 90 ++++-------------- .../control_plane/todos/frontier_revision.py | 27 ------ .../control_plane/todos/frontier_revision.ts | 75 ++++++++++++++- loopx/control_plane/todos/summary_item.py | 4 +- .../work_items/autonomous_replan_ack.py | 3 + .../work_items/progress_observation.py | 24 ++--- .../test_canonical_frontier_revision.py | 16 +++- .../test_goal_frontier_replan_rules.py | 48 +++++++++- .../test_replan_successor_durable_ack.py | 91 +++++++++++++++++++ ...est_replan_successor_frontier_causality.py | 65 +++++++++++++ .../frontier_revision.test.ts | 84 ++++++++++++++++- 13 files changed, 442 insertions(+), 169 deletions(-) create mode 100644 tests/control_plane/test_replan_successor_frontier_causality.py diff --git a/loopx/control_plane/goals/goal_frontier/__init__.py b/loopx/control_plane/goals/goal_frontier/__init__.py index b403a96fb0..610894e389 100644 --- a/loopx/control_plane/goals/goal_frontier/__init__.py +++ b/loopx/control_plane/goals/goal_frontier/__init__.py @@ -1204,7 +1204,7 @@ def derive_goal_frontier_replan_obligation_from_summaries( if replan_rule.rule is GoalFrontierReplanRule.LONG_TODO_CHAIN: assert long_chain_observation is not None assert long_chain_ack_decision is not None - long_chain_trigger = long_chain_observation.to_trigger() + long_chain_trigger = long_chain_observation.trigger return build_autonomous_replan_obligation_payload( schema_version=AUTONOMOUS_REPLAN_OBLIGATION_SCHEMA_VERSION, agent_id=agent_id, diff --git a/loopx/control_plane/goals/goal_frontier/ack_policy.py b/loopx/control_plane/goals/goal_frontier/ack_policy.py index 3447ebb995..19f372d15e 100644 --- a/loopx/control_plane/goals/goal_frontier/ack_policy.py +++ b/loopx/control_plane/goals/goal_frontier/ack_policy.py @@ -16,10 +16,10 @@ replan_obligation_trigger_kinds, required_semantic_outcomes, ) +from ...work_items.autonomous_replan_obligation import ensure_replan_novelty_policy from .long_todo_chain import ( LONG_TODO_CHAIN_TRIGGER, - long_todo_chain_source_checkpoint, - long_todo_chain_transition_is_fresh, + long_todo_chain_successor_checkpoints, ) @@ -101,48 +101,45 @@ def replan_successor_transition_ack( item for item in agent_todo_items if isinstance(item, dict) ] trigger_kinds = replan_obligation_trigger_kinds(replan_obligation or {}) - long_chain_checkpoint: dict[str, str] | None = None - long_chain_frontier_updated_at: str | None = None + candidates = [item for item in source_items + if todo_item_is_actionable_open(item) + and todo_item_task_class(item) == TODO_TASK_CLASS_ADVANCEMENT + and normalize_todo_claimed_by(item.get("claimed_by")) == safe_agent_id + and normalize_todo_replan_obligation_id(item.get("replan_obligation_id")) + and normalize_todo_id(item.get("todo_id")) + and replan_successor_semantic_binding(action_kind=item.get("action_kind"), + target_key=item.get("target_key"), explore_result_node_refs=item.get("explore_result_node_refs"))] + if not candidates: + return None + trigger_checkpoints: list[dict[str, str]] | None = None + eligible_ids = {item["todo_id"] for item in candidates if item["replan_obligation_id"] == obligation_id} if LONG_TODO_CHAIN_TRIGGER in trigger_kinds: - source_checkpoint = long_todo_chain_source_checkpoint( + source_checkpoint = long_todo_chain_successor_checkpoints( source_items, agent_id=safe_agent_id, + triggers=(replan_obligation or {}).get("triggers") or [], + obligation_id=obligation_id, + candidates=[{"todo_id": item["todo_id"], "updated_at": item.get("updated_at"), + "origin_obligation_id": item["replan_obligation_id"]} for item in candidates], frontier_revision_index=(agent_todo_summary or {}).get( "advancement_frontier_revision_index" ), ) if source_checkpoint is None: return None - long_chain_checkpoint, long_chain_frontier_updated_at = source_checkpoint - successor = next( - ( - item - for item in source_items - if isinstance(item, dict) - and todo_item_is_actionable_open(item) - and todo_item_task_class(item) == TODO_TASK_CLASS_ADVANCEMENT - and normalize_todo_claimed_by(item.get("claimed_by")) - == safe_agent_id - and normalize_todo_replan_obligation_id( - item.get("replan_obligation_id") - ) - == obligation_id - and normalize_todo_id(item.get("todo_id")) - and replan_successor_semantic_binding( - action_kind=item.get("action_kind"), - target_key=item.get("target_key"), - explore_result_node_refs=item.get("explore_result_node_refs"), - ) - and ( - LONG_TODO_CHAIN_TRIGGER not in trigger_kinds - or long_todo_chain_transition_is_fresh( - frontier_updated_at=long_chain_frontier_updated_at, - transition_generated_at=item.get("updated_at"), - ) - ) - ), - None, - ) + trigger_checkpoints = source_checkpoint["trigger_checkpoints"] + eligible_ids = set() + origins = {item["todo_id"]: item["replan_obligation_id"] for item in candidates} + for binding in source_checkpoint["bindings"]: + if binding["kind"] == "exact": + eligible_ids.add(binding["todo_id"]) + elif binding["kind"] == "predecessor": + prior = ensure_replan_novelty_policy({**(replan_obligation or {}), "triggers": [ + {**trigger, "frontier_revision": binding["frontier_revision"]} + for trigger in (replan_obligation or {}).get("triggers") or []]}) + if prior["obligation_id"] == origins[binding["todo_id"]]: + eligible_ids.add(binding["todo_id"]) + successor = next((item for item in candidates if item["todo_id"] in eligible_ids), None) if successor is None: return None successor_todo_id = normalize_todo_id(successor.get("todo_id")) @@ -153,16 +150,8 @@ def replan_successor_transition_ack( ) if successor_binding is None: return None - trigger_checkpoints = replan_obligation_trigger_checkpoints( - replan_obligation or {} - ) - if LONG_TODO_CHAIN_TRIGGER in trigger_kinds: - assert long_chain_checkpoint is not None - trigger_checkpoints = [ - checkpoint - for checkpoint in trigger_checkpoints - if checkpoint.get("kind") != LONG_TODO_CHAIN_TRIGGER - ] + [long_chain_checkpoint] + if trigger_checkpoints is None: + trigger_checkpoints = replan_obligation_trigger_checkpoints(replan_obligation or {}) semantic_delta = { "schema_version": "replan_semantic_delta_v0", "accepted": True, @@ -173,9 +162,10 @@ def replan_successor_transition_ack( "trigger_checkpoints": trigger_checkpoints, "obligation_id": obligation_id, "successor_todo_id": successor_todo_id, + "successor_origin_obligation_id": successor["replan_obligation_id"], "successor_binding": successor_binding, "reason": ( - "an exact current-obligation Todo transition created a runnable " + "an exact or source-proven predecessor Todo transition created a runnable " "successor" ), } diff --git a/loopx/control_plane/goals/goal_frontier/long_todo_chain.py b/loopx/control_plane/goals/goal_frontier/long_todo_chain.py index dd506e0fff..99a2e5ce31 100644 --- a/loopx/control_plane/goals/goal_frontier/long_todo_chain.py +++ b/loopx/control_plane/goals/goal_frontier/long_todo_chain.py @@ -5,16 +5,11 @@ from dataclasses import dataclass from typing import Any -from ...runtime.time import parse_timestamp from ...todos.frontier_revision import ( TODO_FRONTIER_REVISION_SCHEMA_VERSION, - advancement_frontier_owned_identity, - advancement_frontier_revision_from_index, - selectable_advancement_frontier_owned_identity, - selectable_advancement_frontier_revision, + frontier_source_facts, ) from ...effect_runtime import effect_runtime_result -from ...todos.frontier_revision import frontier_source_facts LONG_TODO_CHAIN_TRIGGER = "long_todo_chain" @@ -35,25 +30,9 @@ class LongTodoChainObservation: agent_id: str | None frontier_revision: str | None frontier_revision_complete: bool + trigger: dict[str, Any] frontier_owned_identity: str | None = None - def to_trigger(self) -> dict[str, Any]: - trigger: dict[str, Any] = { - "trigger_count": self.trigger_count, - "count_kind": self.count_kind, - "selectable_open_count": self.selectable_open_count, - "selectable_advancement_count": self.selectable_advancement_count, - "current_agent_claimed_advancement_count": ( - self.current_agent_claimed_advancement_count - ), - "unclaimed_advancement_count": self.unclaimed_advancement_count, - "threshold": self.threshold, - "agent_id": self.agent_id, - } - if self.frontier_revision_complete and self.frontier_revision: - trigger["frontier_revision"] = self.frontier_revision - return trigger - @dataclass(frozen=True) class LongTodoChainAckDecision: @@ -61,46 +40,27 @@ class LongTodoChainAckDecision: rearmed_after_obligation_id: str | None = None -def long_todo_chain_source_checkpoint( +def long_todo_chain_successor_checkpoints( source_items: list[dict[str, Any]], *, agent_id: str | None, + triggers: list[dict[str, Any]], + obligation_id: str, + candidates: list[dict[str, Any]], frontier_revision_index: Any = None, -) -> tuple[dict[str, str], str] | None: - """Return the revision and ordering fence for an exact Todo source.""" +) -> dict[str, Any] | None: + """Resolve successor checkpoints and fresh causal bindings in one TS read.""" - projected = advancement_frontier_revision_from_index( - frontier_revision_index, - agent_id=agent_id, - ) - frontier_revision, frontier_updated_at, revision_complete = ( - projected - if projected is not None - else selectable_advancement_frontier_revision( - source_items, - agent_id=agent_id, - ) - ) - if not revision_complete or not frontier_revision or not frontier_updated_at: - return None - owned_identity = ( - advancement_frontier_owned_identity( - frontier_revision_index, agent_id=agent_id - ) - if projected is not None - else selectable_advancement_frontier_owned_identity( - source_items, agent_id=agent_id - ) - ) - checkpoint = { - "kind": LONG_TODO_CHAIN_TRIGGER, - "frontier_revision": frontier_revision, - } - if owned_identity: - # The ACK stays valid while this agent's own selectable rows are - # unchanged, even when another lane claims work this agent can still see. - checkpoint["frontier_owned_identity"] = owned_identity - return (checkpoint, frontier_updated_at) + needs_predecessor_source = any(row["origin_obligation_id"] != obligation_id for row in candidates) + result: dict[str, Any] | None = effect_runtime_result("todo.frontier_revision.project", { + "schema_version": "todo_frontier_revision_request_v0", + "operation": "successor_checkpoints", "agent_id": agent_id, + "triggers": triggers, + "obligation_id": obligation_id, "candidates": candidates, + "index": frontier_revision_index, + "rows": frontier_source_facts(source_items) if needs_predecessor_source or not isinstance(frontier_revision_index, dict) else None, + })["source_checkpoint"] + return result def evaluate_long_todo_chain( @@ -126,17 +86,3 @@ def evaluate_long_todo_chain( LongTodoChainObservation(**observation) if observation is not None else None, LongTodoChainAckDecision(**decision) if decision is not None else None, ) - - -def long_todo_chain_transition_is_fresh( - *, - frontier_updated_at: Any, - transition_generated_at: Any, -) -> bool: - """Fence a successor Todo against the authoritative source revision.""" - - frontier_time = parse_timestamp(frontier_updated_at) - transition_time = parse_timestamp(transition_generated_at) - if frontier_time is None or transition_time is None: - return False - return bool(transition_time >= frontier_time) diff --git a/loopx/control_plane/todos/frontier_revision.py b/loopx/control_plane/todos/frontier_revision.py index 6bf1a659ad..4c2f590d50 100644 --- a/loopx/control_plane/todos/frontier_revision.py +++ b/loopx/control_plane/todos/frontier_revision.py @@ -112,33 +112,6 @@ def selectable_advancement_frontier_revision( return result -def _owned_identity(checkpoint: Any) -> str | None: - """Read the agent-owned identity a checkpoint carries, when it has one.""" - - if not isinstance(checkpoint, dict) or checkpoint.get("complete") is not True: - return None - identity = str(checkpoint.get("frontier_owned_identity") or "").strip() - return identity or None - - -def advancement_frontier_owned_identity( - value: Any, *, agent_id: str | None, -) -> str | None: - """Identity over the rows this agent owns inside an indexed checkpoint.""" - - return _owned_identity(_request("read", index=value, - agent_id=normalize_todo_claimed_by(agent_id)).get("checkpoint")) - - -def selectable_advancement_frontier_owned_identity( - source_items: list[dict[str, Any]] | None, *, agent_id: str | None, -) -> str | None: - """Identity over the rows this agent owns in a freshly selected checkpoint.""" - - return _owned_identity(_request("select", rows=frontier_source_facts(source_items), - agent_id=normalize_todo_claimed_by(agent_id)).get("checkpoint")) - - def build_advancement_frontier_revision_index( source_items: list[dict[str, Any]], ) -> dict[str, Any]: diff --git a/loopx/control_plane/todos/frontier_revision.ts b/loopx/control_plane/todos/frontier_revision.ts index 745a7fe6a1..cfc225eda3 100644 --- a/loopx/control_plane/todos/frontier_revision.ts +++ b/loopx/control_plane/todos/frontier_revision.ts @@ -32,6 +32,11 @@ type LongChainObservation = { frontier_owned_identity: string | null; }; type AckDecision = {acknowledged: boolean; rearmed_after_obligation_id: string | null}; +type SuccessorBinding = {kind: "exact"; todo_id: string} | + {kind: "predecessor"; todo_id: string; frontier_revision: string}; +type TriggerCheckpoint = { + kind: string; frontier_revision: string; frontier_owned_identity?: string; +}; const object = (value: unknown): JsonObject => value !== null && typeof value === "object" && !Array.isArray(value) ? value as JsonObject : {}; const text = (value: unknown): string => typeof value === "string" ? stripPythonWhitespace(value) : ""; @@ -49,6 +54,21 @@ const count = (value: unknown): number => { return Number.isFinite(number) ? Math.max(0, Math.trunc(number)) : 0; }; +/** One receipt shape for observation, successor and semantic-writeback paths. + * Historical revision-only checkpoints remain valid. An owned identity cannot + * stand alone or confer the long-chain matching rule on another trigger kind. + */ +function triggerCheckpoint(value: unknown): TriggerCheckpoint | null { + const row = object(value), kind = text(row.kind), revision = text(row.frontier_revision); + if (!kind || !revision || row.frontier_revision_complete === false) return null; + const owned = kind === TRIGGER ? text(row.frontier_owned_identity) : ""; + return {kind, frontier_revision: revision, ...(owned ? {frontier_owned_identity: owned} : {})}; +} + +function triggerCheckpoints(value: unknown): TriggerCheckpoint[] { + return (Array.isArray(value) ? value : []).map(triggerCheckpoint).filter(row => row !== null); +} + function decodeRows(value: unknown): Row[] | null { if (value == null) return null; if (!Array.isArray(value)) { @@ -124,10 +144,53 @@ function readIndex(value: unknown, agent: string | null): Checkpoint | null { frontier_owned_identity: text(entry.frontier_owned_identity) || null}; } +function successorCheckpoints(request: JsonObject, agent: string | null): JsonObject | null { + const indexed = readIndex(request.index, agent); + const rows = decodeRows(request.rows); + const source = indexed ?? checkpoint(rows, agent); + if (!source.complete) return null; + const latest = parseTodoTimestampMicros(source.frontier_updated_at)!; + const candidates = (Array.isArray(request.candidates) ? request.candidates : []).map(object) + .filter(row => { + const updated = parseTodoTimestampMicros(text(row.updated_at)); + return text(row.todo_id) && updated !== null && updated >= latest; + }); + const bindings: SuccessorBinding[] = candidates.filter(row => row.origin_obligation_id === request.obligation_id) + .map(row => ({kind: "exact", todo_id: text(row.todo_id)})); + const triggers = Array.isArray(request.triggers) ? request.triggers : []; + const trigger = object(triggers[0]); + // A new successor changes the revision it was created to settle. Reconstruct + // only a unique fresh insertion, using a complete source matching the index. + // The existing obligation-id owner still verifies the predecessor revision. + const priorAdvancement = count(trigger.selectable_advancement_count) - 1; + const priorOpen = count(trigger.selectable_open_count) - 1; + if (bindings.length === 0 && candidates.length === 1 && triggers.length === 1 && + trigger.kind === TRIGGER && trigger.frontier_revision === source.frontier_revision && (priorAdvancement >= 15 || priorOpen >= 20 && priorAdvancement > 0)) { + const completeSource = indexed === null ? source : checkpoint(rows, agent); + if (completeSource.complete && completeSource.frontier_revision === source.frontier_revision) { + const candidate = candidates[0]; + const prior = checkpoint(rows === null ? null : rows.filter(row => row.id !== candidate.todo_id), agent); + if (prior.complete && prior.frontier_revision !== source.frontier_revision) { + bindings.push({kind: "predecessor", todo_id: text(candidate.todo_id), frontier_revision: prior.frontier_revision}); + } + } + } + return {trigger_checkpoints: [ + ...triggerCheckpoints(request.triggers).filter(row => row.kind !== TRIGGER), + triggerCheckpoint({kind: TRIGGER, ...source}), + ], bindings}; +} + export function projectAdvancementFrontier(value: unknown): JsonObject { const request = requireJsonObject(value, "frontier revision request"); if (request.schema_version !== "todo_frontier_revision_request_v0") throw new EffectRuntimeRequestError("frontier revision schema mismatch"); const agent = agentId(request.agent_id); + if (request.operation === "trigger_checkpoints") { + return {trigger_checkpoints: triggerCheckpoints(request.triggers)}; + } + if (request.operation === "successor_checkpoints") { + return {source_checkpoint: successorCheckpoints(request, agent)}; + } if (request.operation === "read") return {checkpoint: readIndex(request.index, agent)}; const rows = decodeRows(request.rows); if (request.operation === "select") return {checkpoint: checkpoint(rows, agent)}; @@ -149,9 +212,9 @@ function classifyAck(observation: LongChainObservation, value: unknown): AckDeci !strings(delta.trigger_kinds).includes(TRIGGER) || !/^replan-[a-f0-9]{16}$/.test(id) || observation.frontier_revision_complete !== true || !text(observation.frontier_revision)) return rejected; const matches = Array.isArray(delta.trigger_checkpoints) && delta.trigger_checkpoints.some(raw => { - const row = object(raw); - if (text(row.kind) !== TRIGGER) return false; - if (text(row.frontier_revision) === observation.frontier_revision) return true; + const row = triggerCheckpoint(raw); + if (row === null || row.kind !== TRIGGER) return false; + if (row.frontier_revision === observation.frontier_revision) return true; // Another lane claiming or editing an unclaimed row moves the revision but // leaves this agent's own selectable rows untouched; that is not new // evidence about this agent's chain, so it must not re-arm the obligation. @@ -183,5 +246,9 @@ export function evaluateLongTodoChain(value: unknown): JsonObject { threshold, agent_id: agent, frontier_revision: revision.complete ? revision.frontier_revision : null, frontier_revision_complete: revision.complete, frontier_owned_identity: revision.complete ? revision.frontier_owned_identity : null}; - return {observation, decision: classifyAck(observation, request.ack)}; + const {frontier_revision, frontier_revision_complete, frontier_owned_identity, ...counts} = observation; + const receipt = triggerCheckpoint({kind: TRIGGER, frontier_revision, + frontier_revision_complete, frontier_owned_identity}); + return {observation: {...observation, trigger: {...counts, ...receipt}}, + decision: classifyAck(observation, request.ack)}; } diff --git a/loopx/control_plane/todos/summary_item.py b/loopx/control_plane/todos/summary_item.py index a9239025cf..9f4bbcee61 100644 --- a/loopx/control_plane/todos/summary_item.py +++ b/loopx/control_plane/todos/summary_item.py @@ -19,6 +19,7 @@ from .handoff_gate import handoff_ready_successor_todo_ids from .handoff_note import attach_todo_handoff_note, compact_todo_continuation_hint from .todo_semantics import todo_item_task_class +from .frontier_revision import FRONTIER_REVISION_FIELDS TODO_SUMMARY_COMPACT_FIELDS = ( "schema_version", @@ -273,7 +274,8 @@ def todo_planning_source_items( seen.add(todo_id) compact = compact_todo_summary_item(item, text=text) if include_terminal: - for field in ("evidence", "note", "last_actor_agent_id"): + # Terminal-inclusive planning also proves exact frontier causality. + for field in (*FRONTIER_REVISION_FIELDS, "evidence", "note", "last_actor_agent_id"): if item.get(field) is not None: compact[field] = item[field] planning_items.append(compact) diff --git a/loopx/control_plane/work_items/autonomous_replan_ack.py b/loopx/control_plane/work_items/autonomous_replan_ack.py index 777a34f3ba..8e00d7d9fc 100644 --- a/loopx/control_plane/work_items/autonomous_replan_ack.py +++ b/loopx/control_plane/work_items/autonomous_replan_ack.py @@ -108,6 +108,9 @@ def compact_autonomous_replan_ack(run: dict[str, Any] | None) -> dict[str, Any] "trigger_kinds", "trigger_checkpoints", "obligation_id", + "successor_todo_id", + "successor_origin_obligation_id", + "successor_binding", "observation_fingerprint", "reason", ) diff --git a/loopx/control_plane/work_items/progress_observation.py b/loopx/control_plane/work_items/progress_observation.py index 512521688e..9e781463ca 100644 --- a/loopx/control_plane/work_items/progress_observation.py +++ b/loopx/control_plane/work_items/progress_observation.py @@ -7,6 +7,7 @@ from typing import Any from ...turn_identity import normalize_turn_instance_id +from ..effect_runtime import effect_runtime_result from ..goals.goal_vision_state import normalize_goal_vision_state from ..todos.contract import ( normalize_todo_task_domain, @@ -396,20 +397,15 @@ def replan_obligation_trigger_checkpoints( ) -> list[dict[str, str]]: """Bind an ACK to typed trigger revisions supplied by the obligation.""" - checkpoints: list[dict[str, str]] = [] - for trigger in obligation.get("triggers") or []: - if not isinstance(trigger, Mapping): - continue - kind = str(trigger.get("kind") or "").strip() - frontier_revision = str(trigger.get("frontier_revision") or "").strip() - if not kind or not frontier_revision: - continue - checkpoints.append( - { - "kind": kind, - "frontier_revision": frontier_revision, - } - ) + triggers = obligation.get("triggers") + if not triggers: + return [] + checkpoints = effect_runtime_result("todo.frontier_revision.project", { + "schema_version": "todo_frontier_revision_request_v0", + "operation": "trigger_checkpoints", "triggers": triggers, + })["trigger_checkpoints"] + if not isinstance(checkpoints, list): + raise TypeError("typed frontier checkpoint response must be a list") return checkpoints diff --git a/tests/control_plane/test_canonical_frontier_revision.py b/tests/control_plane/test_canonical_frontier_revision.py index f87f7172fd..5a60ef2081 100644 --- a/tests/control_plane/test_canonical_frontier_revision.py +++ b/tests/control_plane/test_canonical_frontier_revision.py @@ -11,6 +11,7 @@ from loopx.control_plane.goals.goal_frontier.long_todo_chain import evaluate_long_todo_chain from loopx.control_plane.goals.shared_goal_alignment import project_shared_goal_alignment from loopx.control_plane.testing.canary_harness import run_json_cli_result +from loopx.control_plane.work_items.progress_observation import semantic_delta_from_writeback from loopx.status import active_state_todo_fields @@ -44,7 +45,8 @@ def _commit_variant(paths, operation, todo_id): @pytest.mark.parametrize("display", ["stale", "missing"]) -def test_complex_canonical_frontier_ack_tracks_only_selectable_material_changes(tmp_path, display): +@pytest.mark.parametrize("ack_scope", ["revision_only", "owned"]) +def test_complex_canonical_frontier_ack_tracks_only_selectable_material_changes(tmp_path, display, ack_scope): paths = _write_fixture(tmp_path) projection = _fixture() # The shared fixture intentionally includes incomplete historical timestamps. @@ -57,6 +59,7 @@ def test_complex_canonical_frontier_ack_tracks_only_selectable_material_changes( "role": "agent", "status": "open", "done": False, "task_class": "advancement_task", "text": "Synthetic independent work", "archive_state": "active", "source_section": "Agent Todo", "index": len(projection["todos"]) + 1, "updated_at": "2026-09-01T00:00:00.000002Z", + **({"claimed_by": "agent-a"} if index == 0 else {}), **({"excluded_agents": ["agent-a"]} if index == 29 else {}), }) from loopx.control_plane.coordination.local_authority_shadow_projection import canonical_bytes @@ -84,17 +87,26 @@ def observe(ack=None): ack = {"recorded": True, "semantic_delta": {"accepted": True, "obligation_id": "replan-0123456789abcdef", "trigger_kinds": ["long_todo_chain"], "trigger_checkpoints": [{"kind": "long_todo_chain", "frontier_revision": initial.frontier_revision}]}} + if ack_scope == "owned": + ack["semantic_delta"] = semantic_delta_from_writeback(obligation={ + "obligation_id": "replan-0123456789abcdef", "triggers": [initial.trigger], + }, progress_observation={"result_class": "advanced", "surface_id": "canonical-frontier", + "evidence_ids": ["evidence:source-readback"]}) + assert ack["semantic_delta"]["accepted"] and initial.frontier_owned_identity assert observe(ack)[1].acknowledged _commit_variant(paths, "excluded-edit", "todo_zz_frontier_029") assert observe(ack)[1].acknowledged _commit_variant(paths, "eligible-edit", "todo_zz_frontier_028") changed, decision = observe(ack) - assert not decision.acknowledged and decision.rearmed_after_obligation_id == "replan-0123456789abcdef" + assert decision.acknowledged is (ack_scope == "owned") + assert decision.rearmed_after_obligation_id == (None if ack_scope == "owned" else "replan-0123456789abcdef") assert changed.frontier_revision != initial.frontier_revision next_ack = deepcopy(ack) next_ack["semantic_delta"]["trigger_checkpoints"][0]["frontier_revision"] = changed.frontier_revision assert observe(next_ack)[1].acknowledged _commit_variant(paths, "remove-exclusion", "todo_zz_frontier_029") + assert observe(next_ack)[1].acknowledged is (ack_scope == "owned") + _commit_variant(paths, "owned-edit", "todo_zz_frontier_000") assert not observe(next_ack)[1].acknowledged code, result = run_json_cli_result("quota", "should-run", "--goal-id", GOAL_ID, "--agent-id", "agent-a", registry_path=paths["registry"]) diff --git a/tests/control_plane/test_goal_frontier_replan_rules.py b/tests/control_plane/test_goal_frontier_replan_rules.py index fb5572a078..7bbd6f829a 100644 --- a/tests/control_plane/test_goal_frontier_replan_rules.py +++ b/tests/control_plane/test_goal_frontier_replan_rules.py @@ -29,6 +29,7 @@ build_interaction_contract, interaction_next_cli_actions, ) +from loopx.control_plane.work_items.progress_observation import semantic_delta_from_writeback @pytest.mark.parametrize( @@ -228,12 +229,13 @@ def _accepted_long_chain_ack(obligation: dict[str, object]) -> dict[str, object] def _derive_long_chain( source_items: list[dict[str, object]], *, + agent_todo_summary: dict[str, object] | None = None, latest_replan_ack: dict[str, object] | None = None, current_transition_replan_ack: dict[str, object] | None = None, ) -> dict[str, object] | None: return derive_goal_frontier_replan_obligation_from_summaries( user_todo_summary={"open_count": 0}, - agent_todo_summary=_long_chain_summary(source_items), + agent_todo_summary=agent_todo_summary or _long_chain_summary(source_items), agent_todo_source_items=source_items, work_lane_contract={"lane": "advancement_task", "must_attempt_work": True}, agent_id="current-agent", @@ -271,6 +273,50 @@ def test_long_todo_chain_checkpoint_is_edge_triggered_and_rearms_on_change() -> ) +@pytest.mark.parametrize("change,rearms", [ + ("peer_claim", False), ("unclaimed_priority", False), + ("owned_priority", True), ("maintenance_timestamp", False), +]) +def test_writeback_ack_preserves_owned_material_basis(change: str, rearms: bool) -> None: + items = [*_long_chain_source_items(), { + **_advancement("todo_unclaimed", ""), + "updated_at": "2026-08-22T09:00:00+08:00", + }] + + def derive(rows, ack=None): + owned = [row for row in rows if row.get("claimed_by") == "current-agent"] + unclaimed = [row for row in rows if not row.get("claimed_by")] + return _derive_long_chain(rows, latest_replan_ack=ack, agent_todo_summary={ + "open_count": len(rows), "current_agent_claimed_open_count": len(owned), + "current_agent_claimed_advancement_count": len(owned), + "unclaimed_open_count": len(unclaimed), + "unclaimed_priority_open_items": unclaimed, + "executable_backlog_items": owned + unclaimed, + "claim_scope": {"other_agent_claimed_items": [ + row for row in rows if row not in owned + unclaimed]}, + }) + + original = derive(items) + assert original is not None + delta = semantic_delta_from_writeback(obligation=original, progress_observation={ + "result_class": "advanced", "surface_id": "dependency-recovery", + "evidence_ids": ["evidence:independent-acceptance"], + }) + assert delta["accepted"] + ack = {"recorded": True, "semantic_delta": delta} + assert derive(items, ack) is None + changed = deepcopy(items) + if change == "peer_claim": + changed[-1]["claimed_by"] = "peer-agent" + elif change == "unclaimed_priority": + changed[-1]["priority"] = "P0" + elif change == "owned_priority": + changed[0]["priority"] = "P0" + else: + changed[0]["updated_at"] = "2026-08-22T10:00:00+08:00" + assert (derive(changed, ack) is not None) is rearms + + def test_frontier_revision_index_preserves_complete_agent_lane_semantics() -> None: source_items = [ { diff --git a/tests/control_plane/test_replan_successor_durable_ack.py b/tests/control_plane/test_replan_successor_durable_ack.py index c43bb86905..9ab2de040c 100644 --- a/tests/control_plane/test_replan_successor_durable_ack.py +++ b/tests/control_plane/test_replan_successor_durable_ack.py @@ -104,3 +104,94 @@ def test_transition_settles_only_exact_guard_and_owner() -> None: with pytest.raises(ReplanWritebackRejected): enforce_open_replan_writeback( state_text=successor_state(obligation["obligation_id"], owner="another-agent"), **kwargs) + + +@pytest.mark.parametrize("route", ["writeback", "successor", "canonical-successor"]) +def test_long_chain_ack_survives_real_cli_history_and_peer_claim(tmp_path: Path, route: str) -> None: + import subprocess + import sys + + project = tmp_path / "project" + project.mkdir() + runtime = tmp_path / "runtime" + state = project / "ACTIVE_GOAL_STATE.md" + rows = [ + f"- [ ] [P1] Synthetic work {i}.\n" + f" \n" + for i in range(15) + ] + state.write_text("---\nstatus: active\n---\n\n# Goal\n\n## Agent Todo\n\n" + "".join(rows) + + "- [ ] [P1] Synthetic shared work.\n" + " \n") + registry = tmp_path / "registry.json" + registry.write_text(json.dumps({"common_runtime_root": str(runtime), "goals": [{ + "id": GOAL, "status": "active", "repo": str(project), "state_file": state.name, + "domain": "synthetic-replan", + "adapter": {"kind": "fixture_connected_delivery_v0", "status": "connected-delivery"}, + "quota": {"compute": 1.0, "window_hours": 24}, + "coordination": {"agent_model": "peer_v1", "registered_agents": [AGENT, "peer-agent"]}, + }]})) + + if route == "canonical-successor": + from canonical_authority_fixture import initialize_canonical_authority + from loopx.control_plane.todos.active_state_todo_parser import parse_active_state_todos + todos = parse_active_state_todos(state.read_text(), item_limit=None)["agent_todos"]["items"] + module = (Path(__file__).resolve().parents[2] / "tests/control_plane_ts/authority_projection_fixture.ts").as_uri() + built = subprocess.run(["node", "--no-warnings", "--experimental-strip-types", "--input-type=module", "-e", + f"import {{authorityProjectionFixture}} from {json.dumps(module)};" + "let s='';for await(const c of process.stdin)s+=c;" + "process.stdout.write(JSON.stringify(authorityProjectionFixture(process.argv[1],JSON.parse(s),[], 'legacy')));", GOAL], + input=json.dumps(todos), capture_output=True, text=True, check=True, timeout=30) + initialize_canonical_authority(runtime, GOAL, json.loads(built.stdout), state_path=state) + + def call(*args): + completed = subprocess.run([sys.executable, "-m", "loopx.cli", "--registry", str(registry), + "--runtime-root", str(runtime), "--format", "json", *args], + capture_output=True, text=True, timeout=60) + assert completed.returncode == 0, completed.stdout + completed.stderr + return json.loads(completed.stdout) + + def guard(): + result = call("quota", "should-run", "--goal-id", GOAL, "--agent-id", AGENT, + "--runtime-profile", "generic_cli") + return result.get("autonomous_replan_obligation") + + original = guard() + assert original and "long_todo_chain" in [t["kind"] for t in original["triggers"]] + progress_args = ["--progress-result-class", "advanced", "--progress-surface-id", "artifact-adoption", + "--progress-evidence-id", "evidence:independent-acceptance"] + if route.endswith("successor"): + added = call("todo", "add", "--goal-id", GOAL, "--role", "agent", "--claimed-by", AGENT, + "--text", "Validate dependent artifact adoption", "--task-class", "advancement_task", + "--action-kind", "validate", "--target-key", "artifact-adoption", + "--replan-obligation-id", original["obligation_id"]) + progress_args = [] + refreshed = call("refresh-state", "--goal-id", GOAL, "--agent-id", AGENT, + "--classification", "bounded_replan_progress", "--delivery-outcome", "surface_only", + *progress_args, + "--vision-unchanged-reason", "The synthetic acceptance remains unchanged.", + "--no-global-sync", "--suppress-external-sinks") + persisted = json.loads(Path(refreshed["json_path"]).read_text()) + ack = persisted["autonomous_replan_ack"] + assert ack["recorded"] and ack["semantic_delta"]["accepted"] + if route.endswith("successor"): + assert ack["semantic_delta"]["successor_todo_id"] == added["todo_id"] + assert ack["semantic_delta"]["successor_origin_obligation_id"] == original["obligation_id"] + assert ack["semantic_delta"]["trigger_checkpoints"][0]["frontier_owned_identity"] + compact = json.loads((runtime / "goals" / GOAL / "runs" / "index.jsonl").read_text().splitlines()[-1]) + assert compact["autonomous_replan_ack"]["semantic_delta"]["trigger_checkpoints"] == ( + ack["semantic_delta"]["trigger_checkpoints"]) + if route.endswith("successor"): + for field in ("successor_todo_id", "successor_origin_obligation_id", "successor_binding"): + assert compact["autonomous_replan_ack"]["semantic_delta"][field] == ack["semantic_delta"][field] + assert guard() is None + call("todo", "claim", "--goal-id", GOAL, "--todo-id", "todo_shared_work", + "--agent-id", "peer-agent", "--claimed-by", "peer-agent") + assert guard() is None + call("todo", "update", "--goal-id", GOAL, "--todo-id", "todo_owned_00", + "--agent-id", AGENT, "--text", "Changed material acceptance for this lane") + rearmed = guard() + assert rearmed and rearmed["rearmed_after_obligation_id"] == ack["semantic_delta"]["obligation_id"] diff --git a/tests/control_plane/test_replan_successor_frontier_causality.py b/tests/control_plane/test_replan_successor_frontier_causality.py new file mode 100644 index 0000000000..1f3e89cd7e --- /dev/null +++ b/tests/control_plane/test_replan_successor_frontier_causality.py @@ -0,0 +1,65 @@ +"""A new successor may close only the exact frontier from which it was created.""" +from copy import deepcopy + +import pytest + +from loopx.control_plane.goals.goal_frontier.ack_policy import replan_successor_transition_ack +from loopx.control_plane.goals.goal_frontier.long_todo_chain import evaluate_long_todo_chain +from loopx.control_plane.todos.frontier_revision import build_advancement_frontier_revision_index +from loopx.control_plane.work_items.autonomous_replan_obligation import ensure_replan_novelty_policy + + +def obligation(rows): + observation, _ = evaluate_long_todo_chain( + agent_todo_summary={"current_agent_claimed_open_count": len(rows)}, + agent_counts={}, frontier_counts={"current_agent_claimed_advancement_count": len(rows)}, + agent_id="worker-a", agent_todo_source_items=rows) + assert observation is not None + return ensure_replan_novelty_policy({"schema_version": "autonomous_replan_obligation_v0", + "agent_id": "worker-a", "triggers": [observation.trigger]}) + + +@pytest.mark.parametrize("indexed", [False, True]) +@pytest.mark.parametrize("mutation", [None, "wrong-origin", "material-edit", "stale", "ambiguous", "truncated"]) +def test_successor_proves_exact_predecessor_before_closing(indexed, mutation): + rows = [{"todo_id": f"todo_owned_{i:03}", "text": "Synthetic work", "status": "open", + "done": False, "task_class": "advancement_task", "claimed_by": "worker-a", + "updated_at": "2026-08-01T00:00:00Z"} for i in range(15)] + origin = obligation(rows) + successor = {**rows[0], "todo_id": "todo_successor", "action_kind": "validate", + "target_key": "artifact-adoption", "updated_at": "2026-08-02T00:00:00Z", + "replan_obligation_id": origin["obligation_id"]} + rows.append(successor) + if mutation == "wrong-origin": + successor["replan_obligation_id"] = "replan-0123456789abcdef" + elif mutation == "material-edit": + rows[0]["text"] = "A different acceptance condition" + elif mutation == "stale": + rows[0]["updated_at"] = "2026-08-03T00:00:00Z" + elif mutation == "ambiguous": + rows.append({**successor, "todo_id": "todo_other_successor"}) + current = obligation(rows) + summary = {"advancement_frontier_revision_index": build_advancement_frontier_revision_index(rows)} if indexed else {} + source = deepcopy(rows) + if mutation == "truncated": + source.pop(0) + ack = replan_successor_transition_ack(summary, agent_id="worker-a", + replan_obligation=current, agent_todo_items=source) + if mutation: + assert ack is None + else: + delta = ack["semantic_delta"] + assert delta["obligation_id"] == current["obligation_id"] + assert delta["successor_origin_obligation_id"] == origin["obligation_id"] + assert delta["successor_todo_id"] == successor["todo_id"] + assert delta["trigger_checkpoints"][0]["frontier_owned_identity"] + + +def test_planning_compaction_preserves_material_frontier_bytes(): + from loopx.control_plane.todos.frontier_revision import frontier_source_facts + from loopx.control_plane.todos.summary_item import todo_planning_source_items + rows = [{"todo_id": "todo_with_dependency", "text": "Validate adoption", "done": False, + "status": "open", "task_class": "advancement_task", "claimed_by": "worker-a", + "depends_on_todo_ids": ["todo_dependency"], "updated_at": "2026-08-01T00:00:00Z"}] + planned = todo_planning_source_items({"items": rows}, include_terminal=True) + assert frontier_source_facts(planned) == frontier_source_facts(rows) diff --git a/tests/control_plane_ts/frontier_revision.test.ts b/tests/control_plane_ts/frontier_revision.test.ts index 66d86323c0..f8337c2001 100644 --- a/tests/control_plane_ts/frontier_revision.test.ts +++ b/tests/control_plane_ts/frontier_revision.test.ts @@ -5,7 +5,8 @@ import {projectAdvancementFrontier, evaluateLongTodoChain} from "../../loopx/con function row(id: string, claim: string | null = null, excluded: string[] = []) { return {id, claim, excluded, advancement: true, updated: "2026-09-01T00:00:00.000001Z", - serialized: JSON.stringify({task_class: "advancement_task", todo_id: id})}; + serialized: JSON.stringify({task_class: "advancement_task", todo_id: id, + ...(claim ? {claimed_by: claim} : {})})}; } function project(rows: ReturnType[], agent_id: string | null = null) { return projectAdvancementFrontier({schema_version: "todo_frontier_revision_request_v0", @@ -123,6 +124,8 @@ test("another lane taking over an unclaimed row does not re-arm this lane's ACK" row("todo_c", "worker-a")]}); assert.deepEqual(ownChange.decision, {acknowledged: false, rearmed_after_obligation_id: "replan-0123456789abcdef"}); + assert.equal((observe({ack, rows: [row("todo_a"), row("todo_b")]}).decision as + Record).acknowledged, false); // An ACK recorded before the owned identity existed still matches on revision. const legacy = {...ack, semantic_delta: {...ack.semantic_delta, trigger_checkpoints: [ @@ -130,3 +133,82 @@ test("another lane taking over an unclaimed row does not re-arm this lane's ACK" assert.deepEqual(observe({ack: legacy, rows: before}).decision, {acknowledged: true, rearmed_after_obligation_id: null}); }); + +test("observation, source and writeback preserve one complete checkpoint", () => { + const rows = [row("todo_a", "worker-a"), row("todo_shared")]; + const result = observe({rows}).observation as Record; + const request = {schema_version: "todo_frontier_revision_request_v0", agent_id: "worker-a"}; + const source = projectAdvancementFrontier({...request, operation: "successor_checkpoints", rows}) + .source_checkpoint as Record; + const writeback = projectAdvancementFrontier({...request, operation: "trigger_checkpoints", + triggers: [result.trigger]}).trigger_checkpoints; + assert.deepEqual(writeback, source.trigger_checkpoints); + const replaced = projectAdvancementFrontier({...request, operation: "successor_checkpoints", rows, + triggers: [{kind: "other", frontier_revision: "r"}, + {kind: "long_todo_chain", frontier_revision: "old"}, result.trigger]}).source_checkpoint as + Record; + assert.deepEqual(replaced.trigger_checkpoints, [ + {kind: "other", frontier_revision: "r"}, ...(source.trigger_checkpoints as object[])]); + assert.equal((source.trigger_checkpoints as Record[])[0].frontier_owned_identity, + result.frontier_owned_identity); + const index = projectAdvancementFrontier({...request, operation: "index", rows}).index; + assert.deepEqual(projectAdvancementFrontier({...request, operation: "successor_checkpoints", index}) + .source_checkpoint, source); + // A present incomplete index is authoritative; a valid fallback cannot heal it. + assert.equal(projectAdvancementFrontier({...request, operation: "successor_checkpoints", rows, + index: {schema_version: "invalid"}}).source_checkpoint, null); +}); + +test("checkpoint transport preserves legacy revisions and rejects incomplete identity authority", () => { + const checkpoints = projectAdvancementFrontier({schema_version: "todo_frontier_revision_request_v0", + operation: "trigger_checkpoints", triggers: [null, {}, + {kind: "long_todo_chain", frontier_owned_identity: "owned"}, + {kind: "long_todo_chain", frontier_revision: "r", frontier_revision_complete: false, + frontier_owned_identity: "owned"}, + {kind: " long_todo_chain ", frontier_revision: " r ", frontier_owned_identity: " owned "}, + {kind: "long_todo_chain", frontier_revision: "legacy"}, + {kind: "other", frontier_revision: "other-revision", frontier_owned_identity: "owned"}, + ]}).trigger_checkpoints; + assert.deepEqual(checkpoints, [ + {kind: "long_todo_chain", frontier_revision: "r", frontier_owned_identity: "owned"}, + {kind: "long_todo_chain", frontier_revision: "legacy"}, + {kind: "other", frontier_revision: "other-revision"}, + ]); + const observation = observe().observation as Record; + const ack = {recorded: true, semantic_delta: {accepted: true, obligation_id: "replan-0123456789abcdef", + trigger_kinds: ["long_todo_chain"], trigger_checkpoints: [{kind: "long_todo_chain", + frontier_owned_identity: observation.frontier_owned_identity}]}}; + assert.equal((observe({ack}).decision as Record).acknowledged, false); + const incomplete = observe({rows: null}).observation as Record; + assert.equal((incomplete.trigger as Record).frontier_owned_identity, undefined); + assert.equal((incomplete.trigger as Record).frontier_revision, undefined); +}); + +test("an entirely unclaimed chain cannot borrow the owned-work exemption", () => { + const observation = observe({rows: [row("todo_unclaimed")]}).observation as Record; + const ack = {recorded: true, semantic_delta: {accepted: true, obligation_id: "replan-0123456789abcdef", + trigger_kinds: ["long_todo_chain"], trigger_checkpoints: [observation.trigger]}}; + assert.equal(observation.frontier_owned_identity, null); + assert.equal((observe({ack, rows: [row("todo_replacement")]}).decision as Record) + .acknowledged, false); +}); + +test("successor reconstruction requires a complete matching source and an already-long predecessor", () => { + const rows = Array.from({length: 16}, (_, i) => row(`todo_${i}`, "worker-a")); + rows[15].updated = "2026-09-02T00:00:00Z"; + const current = project(rows, "worker-a"); + const request = {schema_version: "todo_frontier_revision_request_v0", operation: "successor_checkpoints", + agent_id: "worker-a", rows, obligation_id: "replan-current", + candidates: [{todo_id: "todo_15", updated_at: rows[15].updated, origin_obligation_id: "replan-prior"}], + triggers: [{kind: "long_todo_chain", ...current, selectable_advancement_count: 16, selectable_open_count: 16}]}; + const result = projectAdvancementFrontier(request).source_checkpoint as Record; + assert.deepEqual(result.bindings, [{kind: "predecessor", todo_id: "todo_15", + frontier_revision: project(rows.slice(0, 15), "worker-a").frontier_revision}]); + for (const triggers of [ + [{...request.triggers[0], selectable_advancement_count: 15, selectable_open_count: 15}], + [{...request.triggers[0], frontier_revision: "different-current-source"}], + [...request.triggers, {kind: "periodic_review"}], + ]) { + assert.deepEqual((projectAdvancementFrontier({...request, triggers}).source_checkpoint as Record).bindings, []); + } +}); From acf44344b823537a680bc815884a9531edd8f717 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 17 Sep 2026 19:32:34 +0800 Subject: [PATCH 2/2] docs(control-plane): reconcile frontier checkpoint ownership Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../typescript-control-plane-migration-v0.md | 52 +++++++++++++------ ...script-control-plane-migration-v0.zh-CN.md | 38 +++++++++----- 2 files changed, 61 insertions(+), 29 deletions(-) diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index a890870191..6845400f59 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -824,23 +824,41 @@ editors govern optional feature configuration, not host tool observations, so th configuration owner and fields are unchanged. See [operating semantics](../../quota-allocation.md). -Advancement-frontier checkpoint closure: `todos/frontier_revision.ts` now owns -agent selection, completeness, material hashing, long-chain thresholds and exact -ACK/rearm classification. Python retains the v0 field manifest and legacy JSON/ -metadata codecs so unchanged legal frontiers retain their persisted fingerprints; -the old Python revision builder, index selector and two-step long-chain decision -are retired. Terminal advancement rows still affect material identity, while -timestamp-only maintenance does not rearm it. Thresholds remain 15 advancement -Todos or 20 selectable open Todos with advancement work. Excluded unclaimed work -no longer changes that Agent's checkpoint, including Agents with no claimed rows; -removing the exclusion makes that work relevant again. Duplicate identities in a -selected frontier, duplicate matching index lanes and incomplete timestamps cannot -provide a complete checkpoint or suppress replanning. These are explicit read -corrections, not new execution permissions. The existing canonical source feeds -the index before display truncation. Complex-fixture tests replay accepted ACKs, -excluded/eligible edits and newly available work through a real provider with -stale/missing display; a read-only private-snapshot comparison remains private. -This closes one T3 rule group, not the remaining consumers or T1/T2/D1–D3. +Advancement-frontier checkpoint closure: `todos/frontier_revision.ts` owns +agent selection, completeness, material hashing, long-chain thresholds, +checkpoint construction and ACK/rearm classification. Observation, semantic +writeback and runnable-successor receipts now use the same typed checkpoint +constructor. Successor projection resolves the complete source, owned identity +and replacement checkpoint list in one request; Python no longer assembles +receipts or fetches the same frontier separately for each identity field. +Python retains the v0 field manifest, legacy JSON/metadata codecs, successor +eligibility and the existing obligation-id derivation. TS owns timestamp +ordering and reconstructs only a unique fresh successor insertion against the +complete current source; Python verifies its predecessor obligation id. +Compaction preserves the material `done` field, and history retains successor +lineage. Ambiguous, stale, truncated or unrelated material changes cannot close +the current obligation. This closes one T3 rule group, not +the remaining consumers or T1/T2/D1–D3. + +Thresholds remain 15 advancement Todos or 20 selectable open Todos with +advancement work. Full material revisions include terminal advancement rows; +timestamp-only maintenance does not rearm them. A complete agent-owned identity +also keeps an accepted long-chain ACK valid when peers change shared unclaimed +work. Owned material edits still rearm; an entirely unclaimed chain cannot use +that exemption. Historical revision-only ACKs keep exact-revision matching. +Intentional corrections: semantic writeback now preserves the owned identity; +an identity without a revision or an explicitly incomplete checkpoint cannot +suppress replanning. Other trigger kinds cannot borrow long-chain identity +matching. No threshold, write authority or obligation-id rule changes. + +The canonical index is built before display truncation. Exclusions, duplicate +ids/index lanes, incomplete timestamps and authoritative incomplete indexes +retain fail-closed behavior. Real CLI tests cover both ACK routes through +persisted run/history readback, a peer claim and an owned material edit. The +complex fixture also exercises revision-only and owned ACKs through the real +File provider with stale/missing display. Frontend/Lark configuration is +unchanged: this is the shared quota/recovery checkpoint path, not a new control +or user confirmation. No provider promotion is implied. Large source facts use lossless deflate/base64 transport above 512 KiB, retaining the exact v0 material bytes and the shared 2 MiB request boundary. The TS decoder diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md index 0be3a36319..3f79f44022 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md @@ -642,18 +642,32 @@ observation 调用。Python 只适配 registry、宿主、hook 组合和 CLI, 混入功能开关。详见[操作语义](../../quota-allocation.md)。 -Advancement-frontier checkpoint 闭合:`todos/frontier_revision.ts` 现统一 Agent -选择、完整度、实质内容哈希、长链阈值与精确 ACK/rearm 分类。Python 保留 v0 字段清单 -与 legacy JSON/metadata codec,使合法且未变化的 frontier 保持已有指纹;删除旧 Python -revision builder、index selector 和分两步执行的长链决策。终态 advancement 仍影响 -实质身份;仅更新时间不重新触发。阈值仍是 15 项 advancement,或存在 advancement -时的 20 项可选 open Todo。被排除的 unclaimed 工作不再改变该 Agent 的 checkpoint, -包括没有 claimed Todo 的 Agent;取消 exclusion 后,该工作重新相关。选中 frontier -中的重复 ID、重复匹配的 index lane 与不完整时间不能提供完整 checkpoint 或压制 -replan。这些是明确的只读语义修正,不是执行授权。既有 canonical source 在展示截断前 -生成 index。复杂 fixture 经真实 provider 验证 accepted ACK、excluded/eligible 修改、 -新可用工作及陈旧/缺失展示;私有快照只读对照结果不公开原始数据。 -本批闭合一个 T3 规则组,不代表其余 consumer 或 T1/T2/D1–D3 完成。 +Advancement-frontier checkpoint 闭合:`todos/frontier_revision.ts` 统一 Agent +选择、完整度、实质内容哈希、长链阈值、checkpoint 构造与 ACK/rearm 分类。 +Observation、语义写回和 runnable-successor 回执现在共用一个 typed checkpoint +构造器。后继路径在一次请求内解析完整来源、owned identity 和替换后的 checkpoint +列表;Python 不再自行拼装回执,也不为每个身份字段重复读取同一 frontier。 +Python 保留 v0 字段清单、legacy JSON/metadata codec、后继资格与既有 obligation-id +推导。TS 统一时间顺序,并仅对唯一新鲜后继插入从完整当前来源重建前置 revision; +Python 验证前置 obligation id。压缩保留实质字段 `done`,历史保留后继来源关系。 +多后继歧义、过期、来源截断或无关实质变化都不能关闭当前 obligation。 +这闭合一个 T3 规则组,不代表其余 consumer 或 T1/T2/D1–D3 完成。 + +阈值仍是 15 项 advancement,或存在 advancement 时的 20 项可选 open Todo。 +完整实质 revision 包含终态 advancement;仅更新时间不重新触发。完整的 Agent-owned +identity 还能在同伴改变共享 unclaimed 工作时保持既有 long-chain ACK 有效。 +自己的实质工作变化仍重新触发;全是未认领工作的链不能使用 owned 豁免。 +历史 revision-only ACK 仍按精确 revision 匹配。明确修正:语义写回不再丢失 +owned identity;只有 identity 而没有 revision、或明确不完整的 checkpoint 不能 +压制 replan;其他 trigger kind 不能借用长链身份匹配。阈值、写权限和 obligation-id +规则均未改变。 + +Canonical index 仍在展示截断前生成。Exclusion、重复 ID/index lane、不完整时间与 +权威 index 不完整时均保持 fail-closed。真实 CLI 验证两条 ACK 路径经过运行记录及 +历史回读后,同伴 claim 不重新触发、自己的实质修改重新触发;复杂 fixture 还通过 +真实 File provider,在展示陈旧/缺失时覆盖 revision-only 与 owned ACK。 +前端/Lark 配置未改变:这是共享 quota/recovery checkpoint 路径,没有新增控制项 +或用户确认,也不代表 provider promotion。 来源 facts 超过 512 KiB 时使用无损 deflate/base64 传输,保留精确 v0 内容和共享 2 MiB 请求边界。TS 拒绝畸形载荷及解压超过 64 MiB 的输入,不截断 Todo,也不