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
52 changes: 35 additions & 17 deletions docs/architecture/rfcs/typescript-control-plane-migration-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,也不
Expand Down
2 changes: 1 addition & 1 deletion loopx/control_plane/goals/goal_frontier/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
82 changes: 36 additions & 46 deletions loopx/control_plane/goals/goal_frontier/ack_policy.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
)


Expand Down Expand Up @@ -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"))
Expand All @@ -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,
Expand All @@ -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"
),
}
Expand Down
90 changes: 18 additions & 72 deletions loopx/control_plane/goals/goal_frontier/long_todo_chain.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -35,72 +30,37 @@ 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:
acknowledged: bool
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(
Expand All @@ -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)
Loading