From de6933f13d6d8abda68551c2def448132d63eb3c Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sat, 12 Sep 2026 02:10:46 +0800 Subject: [PATCH 1/5] refactor(frontier): unify typed revision and long-chain replan policy Signed-off-by: huangruiteng --- .../control_plane/effect_runtime_handlers.ts | 3 + .../goals/goal_frontier/__init__.py | 14 +- .../goals/goal_frontier/long_todo_chain.py | 167 +++------------- .../control_plane/todos/frontier_revision.py | 186 +++++------------- .../control_plane/todos/frontier_revision.ts | 154 +++++++++++++++ .../test_canonical_frontier_revision.py | 102 ++++++++++ .../test_frontier_revision_scope.py | 73 +++++++ .../frontier_revision.test.ts | 81 ++++++++ 8 files changed, 488 insertions(+), 292 deletions(-) create mode 100644 loopx/control_plane/todos/frontier_revision.ts create mode 100644 tests/control_plane/test_canonical_frontier_revision.py create mode 100644 tests/control_plane/test_frontier_revision_scope.py create mode 100644 tests/control_plane_ts/frontier_revision.test.ts diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 0122154cee..82c549a304 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -149,6 +149,7 @@ import {evaluateStandingDecisionProjection} from "./todos/standing_decision.ts"; import {evaluateDecisionScope} from "./todos/decision_scope.ts"; import {evaluateCapabilityGate} from "./agents/capability_gate.ts"; import {captureArchivedTodoDependencies} from "./todos/archive_capture.ts"; +import {projectAdvancementFrontier, evaluateLongTodoChain} from "./todos/frontier_revision.ts"; import { evaluateCoordinationTodoSuccessorDerivation } from "./coordination/todo_successor_derivation.ts"; import { checkLegacyCoordinationWriteAllowed, @@ -417,6 +418,8 @@ export function createEffectRuntimeHandlers( ["todo.resume_condition.evaluate", evaluateTodoResumeConditions], ["todo.resume_planning.project", projectTodoResumePlanning], ["todo.quota_planning.project", projectTodoQuotaPlanning], + ["todo.frontier_revision.project", projectAdvancementFrontier], + ["goal.long_todo_chain.evaluate", evaluateLongTodoChain], ["todo.external_wait.plan", planTodoExternalWaitTransition], ["scheduler.state_transition.evaluate", evaluateSchedulerStateTransition], ["scheduler.state.evaluate", evaluateSchedulerStateOperation], diff --git a/loopx/control_plane/goals/goal_frontier/__init__.py b/loopx/control_plane/goals/goal_frontier/__init__.py index 6127ea80ce..88bdbe311c 100644 --- a/loopx/control_plane/goals/goal_frontier/__init__.py +++ b/loopx/control_plane/goals/goal_frontier/__init__.py @@ -62,8 +62,7 @@ ) from .long_todo_chain import ( LONG_TODO_CHAIN_TRIGGER, - classify_long_todo_chain_ack, - observe_long_todo_chain, + evaluate_long_todo_chain, ) from .replan_rules import ( GoalFrontierReplanFacts, @@ -983,20 +982,13 @@ def derive_goal_frontier_replan_obligation_from_summaries( agent_todo_summary, agent_id=agent_id, ) - long_chain_observation = observe_long_todo_chain( + long_chain_observation, long_chain_ack_decision = evaluate_long_todo_chain( agent_todo_summary=agent_todo_summary, agent_counts=agent_counts, frontier_counts=frontier_counts, agent_id=agent_id, agent_todo_source_items=agent_todo_source_items, - ) - long_chain_ack_decision = ( - classify_long_todo_chain_ack( - long_chain_observation, - current_transition_replan_ack or latest_replan_ack, - ) - if long_chain_observation is not None - else None + latest_replan_ack=current_transition_replan_ack or latest_replan_ack, ) replan_rule = select_goal_frontier_replan_rule( GoalFrontierReplanFacts( 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 bb0300850a..73cad3581b 100644 --- a/loopx/control_plane/goals/goal_frontier/long_todo_chain.py +++ b/loopx/control_plane/goals/goal_frontier/long_todo_chain.py @@ -11,43 +11,17 @@ advancement_frontier_revision_from_index, selectable_advancement_frontier_revision, ) -from ...todos.contract import normalize_todo_replan_obligation_id +from ...effect_runtime import effect_runtime_result +from ...todos.frontier_revision import frontier_source_facts LONG_TODO_CHAIN_TRIGGER = "long_todo_chain" -LONG_TODO_CHAIN_ADVANCEMENT_THRESHOLD = 15 -LONG_TODO_CHAIN_OPEN_THRESHOLD = 20 TODO_TASK_CLASS_ADVANCEMENT = "advancement_task" LONG_TODO_CHAIN_FRONTIER_REVISION_SCHEMA_VERSION = ( TODO_FRONTIER_REVISION_SCHEMA_VERSION ) -def _safe_non_negative_int(value: Any) -> int: - try: - return max(0, int(value or 0)) - except (TypeError, ValueError): - return 0 - - -def _selectable_advancement_frontier_revision( - source_items: list[dict[str, Any]] | None, - *, - agent_id: str | None, -) -> tuple[str | None, str | None, bool]: - """Return a complete material revision for one selectable agent lane. - - Terminal advancement rows remain relevant because completion and pruning - mutate the frontier. Incomplete source revisions fail closed so a legacy - projection cannot silently suppress an obligation. - """ - - return selectable_advancement_frontier_revision( - source_items, - agent_id=agent_id, - ) - - @dataclass(frozen=True) class LongTodoChainObservation: trigger_count: int @@ -100,7 +74,7 @@ def long_todo_chain_source_checkpoint( frontier_revision, frontier_updated_at, revision_complete = ( projected if projected is not None - else _selectable_advancement_frontier_revision( + else selectable_advancement_frontier_revision( source_items, agent_id=agent_id, ) @@ -116,119 +90,28 @@ def long_todo_chain_source_checkpoint( ) -def observe_long_todo_chain( - *, - agent_todo_summary: dict[str, Any] | None, - agent_counts: dict[str, int], - frontier_counts: dict[str, int], - agent_id: str | None, - agent_todo_source_items: list[dict[str, Any]] | None = None, -) -> LongTodoChainObservation | None: - """Observe one agent-scoped long chain without inferring from prose.""" - - current_advancement = frontier_counts.get( - "current_agent_claimed_advancement_count", 0 - ) - unclaimed_advancement = frontier_counts.get("unclaimed_advancement_count", 0) - selectable_advancement = current_advancement + unclaimed_advancement - if isinstance(agent_todo_summary, dict): - current_open = _safe_non_negative_int( - agent_todo_summary.get("current_agent_claimed_open_count") - ) - unclaimed_open = _safe_non_negative_int( - agent_todo_summary.get("unclaimed_open_count") - ) - selectable_open = max( - current_open + unclaimed_open, - selectable_advancement, - ) - else: - selectable_open = max(agent_counts.get("open", 0), selectable_advancement) - threshold: int | None = None - trigger_count = 0 - count_kind = "" - if selectable_advancement >= LONG_TODO_CHAIN_ADVANCEMENT_THRESHOLD: - threshold = LONG_TODO_CHAIN_ADVANCEMENT_THRESHOLD - trigger_count = selectable_advancement - count_kind = "selectable_advancement_todos" - elif ( - selectable_open >= LONG_TODO_CHAIN_OPEN_THRESHOLD - and selectable_advancement > 0 - ): - threshold = LONG_TODO_CHAIN_OPEN_THRESHOLD - trigger_count = selectable_open - count_kind = "selectable_open_todos" - if threshold is None: - return None - projected = advancement_frontier_revision_from_index( - (agent_todo_summary or {}).get("advancement_frontier_revision_index"), - agent_id=agent_id, - ) - frontier_revision, _, revision_complete = ( - projected - if projected is not None - else _selectable_advancement_frontier_revision( - agent_todo_source_items, - agent_id=agent_id, - ) - ) - return LongTodoChainObservation( - trigger_count=trigger_count, - count_kind=count_kind, - selectable_open_count=selectable_open, - selectable_advancement_count=selectable_advancement, - current_agent_claimed_advancement_count=current_advancement, - unclaimed_advancement_count=unclaimed_advancement, - threshold=threshold, - agent_id=agent_id, - frontier_revision=frontier_revision, - frontier_revision_complete=revision_complete, - ) - - -def classify_long_todo_chain_ack( - observation: LongTodoChainObservation, - latest_replan_ack: dict[str, Any] | None, -) -> LongTodoChainAckDecision: - """Classify an accepted checkpoint against the current frontier revision.""" - - if ( - not isinstance(latest_replan_ack, dict) - or latest_replan_ack.get("recorded") is not True - ): - return LongTodoChainAckDecision(acknowledged=False) - semantic_delta_value = latest_replan_ack.get("semantic_delta") - semantic_delta: dict[str, Any] = ( - semantic_delta_value if isinstance(semantic_delta_value, dict) else {} - ) - trigger_kinds = { - str(value or "").strip() - for value in semantic_delta.get("trigger_kinds") or [] - if str(value or "").strip() - } - obligation_id = normalize_todo_replan_obligation_id( - semantic_delta.get("obligation_id") - ) - if ( - semantic_delta.get("accepted") is not True - or LONG_TODO_CHAIN_TRIGGER not in trigger_kinds - or not obligation_id - ): - return LongTodoChainAckDecision(acknowledged=False) - if not observation.frontier_revision_complete or not observation.frontier_revision: - return LongTodoChainAckDecision(acknowledged=False) - checkpoint_matches = any( - isinstance(checkpoint, dict) - and str(checkpoint.get("kind") or "").strip() == LONG_TODO_CHAIN_TRIGGER - and str(checkpoint.get("frontier_revision") or "").strip() - == observation.frontier_revision - for checkpoint in semantic_delta.get("trigger_checkpoints") or [] - ) - if checkpoint_matches: - return LongTodoChainAckDecision(acknowledged=True) - return LongTodoChainAckDecision( - acknowledged=False, - rearmed_after_obligation_id=obligation_id, +def evaluate_long_todo_chain( + *, agent_todo_summary: dict[str, Any] | None, + agent_counts: dict[str, int], frontier_counts: dict[str, int], + agent_id: str | None, agent_todo_source_items: list[dict[str, Any]] | None = None, + latest_replan_ack: dict[str, Any] | None = None, +) -> tuple[LongTodoChainObservation | None, LongTodoChainAckDecision | None]: + """One typed observation + checkpoint qualification, not two rule RPCs.""" + result = effect_runtime_result("goal.long_todo_chain.evaluate", { + "schema_version": "long_todo_chain_request_v0", "operation": "observe", + "summary": agent_todo_summary, "agent_counts": agent_counts, + "frontier_counts": frontier_counts, "agent_id": agent_id, + "rows": ( + None if isinstance((agent_todo_summary or {}).get("advancement_frontier_revision_index"), dict) + else frontier_source_facts(agent_todo_source_items) + ), + "ack": latest_replan_ack, + }) + observation = result["observation"] + decision = result["decision"] + return ( + LongTodoChainObservation(**observation) if observation is not None else None, + LongTodoChainAckDecision(**decision) if decision is not None else None, ) diff --git a/loopx/control_plane/todos/frontier_revision.py b/loopx/control_plane/todos/frontier_revision.py index f26f97528c..da4e2debf5 100644 --- a/loopx/control_plane/todos/frontier_revision.py +++ b/loopx/control_plane/todos/frontier_revision.py @@ -2,12 +2,11 @@ from __future__ import annotations -from hashlib import sha256 import json from typing import Any -from ..runtime.time import parse_timestamp -from .contract import normalize_todo_claimed_by +from ..effect_runtime import effect_runtime_result +from .contract import normalize_todo_claimed_by, normalize_todo_excluded_agents from .projection import todo_item_task_class @@ -50,159 +49,68 @@ ) -def _revision_for_claims( - source_items: list[dict[str, Any]] | None, - *, - included_claims: set[str | None] | None, -) -> tuple[str | None, str | None, bool]: +def frontier_source_facts(source_items: list[dict[str, Any]] | None) -> list[dict[str, Any]] | None: + """Legacy codecs only; TS selects lanes and builds complete revision identity.""" if not isinstance(source_items, list): - return None, None, False - revisions: list[dict[str, Any]] = [] - revision_times: list[tuple[Any, str]] = [] - relevant_count = 0 - for item in source_items: - if not isinstance(item, dict): - continue - if todo_item_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: - continue - claimed_by = normalize_todo_claimed_by(item.get("claimed_by")) - if included_claims is not None and claimed_by not in included_claims: - continue - relevant_count += 1 - todo_id = str(item.get("todo_id") or "").strip() - raw_revision = str( - item.get("updated_at") or item.get("completed_at") or "" - ).strip() - parsed_revision = parse_timestamp(raw_revision) - if not todo_id or parsed_revision is None: - return None, None, False - revisions.append( - { - field: item.get(field) - for field in FRONTIER_REVISION_FIELDS - if item.get(field) is not None - } - ) - revision_times.append((parsed_revision, raw_revision)) - if relevant_count == 0 or not revisions: - return None, None, False - encoded = json.dumps( - sorted(revisions, key=lambda item: str(item.get("todo_id") or "")), - ensure_ascii=True, - separators=(",", ":"), - sort_keys=True, - ).encode("utf-8") - return ( - f"{TODO_FRONTIER_REVISION_SCHEMA_VERSION}:" - f"{sha256(encoded).hexdigest()[:24]}", - max(revision_times, key=lambda item: item[0])[1], - True, - ) + return None + return [ + { + "id": str(item.get("todo_id") or "").strip(), + "claim": normalize_todo_claimed_by(item.get("claimed_by")), + "excluded": normalize_todo_excluded_agents(item.get("excluded_agents")), + "updated": str(item.get("updated_at") or item.get("completed_at") or "").strip(), + "advancement": todo_item_task_class(item) == TODO_TASK_CLASS_ADVANCEMENT, + "serialized": json.dumps( + {key: item[key] for key in FRONTIER_REVISION_FIELDS if item.get(key) is not None}, + ensure_ascii=True, separators=(",", ":"), sort_keys=True, + ), + } + for item in source_items if isinstance(item, dict) + ] -def selectable_advancement_frontier_revision( - source_items: list[dict[str, Any]] | None, - *, - agent_id: str | None, -) -> tuple[str | None, str | None, bool]: - """Return the material revision for one agent's selectable frontier.""" +def _request(operation: str, **facts: Any) -> dict[str, Any]: + result = effect_runtime_result("todo.frontier_revision.project", { + "schema_version": "todo_frontier_revision_request_v0", + "operation": operation, **facts, + }) + if not isinstance(result, dict): + raise TypeError("typed frontier revision response must be an object") + return result - normalized_agent_id = normalize_todo_claimed_by(agent_id) - included_claims = ( - {None, normalized_agent_id} if normalized_agent_id else None - ) - return _revision_for_claims(source_items, included_claims=included_claims) +def _checkpoint_tuple(value: Any) -> tuple[str | None, str | None, bool] | None: + if value is None: + return None + if value["complete"]: + return value["frontier_revision"], value["frontier_updated_at"], True + return None, None, False -def _checkpoint( - source_items: list[dict[str, Any]], - *, - included_claims: set[str | None] | None, -) -> dict[str, Any]: - revision, updated_at, complete = _revision_for_claims( - source_items, - included_claims=included_claims, - ) - checkpoint: dict[str, Any] = {"complete": complete} - if complete and revision and updated_at: - checkpoint["frontier_revision"] = revision - checkpoint["frontier_updated_at"] = updated_at - return checkpoint + +def selectable_advancement_frontier_revision( + source_items: list[dict[str, Any]] | None, *, agent_id: str | None, +) -> tuple[str | None, str | None, bool]: + result = _checkpoint_tuple(_request("select", rows=frontier_source_facts(source_items), + agent_id=normalize_todo_claimed_by(agent_id))["checkpoint"]) + assert result is not None + return result def build_advancement_frontier_revision_index( source_items: list[dict[str, Any]], ) -> dict[str, Any]: - """Project full frontier identity without retaining every Todo row.""" - - agent_ids = sorted( - { - claimed_by - for item in source_items - if isinstance(item, dict) - and todo_item_task_class(item) == TODO_TASK_CLASS_ADVANCEMENT - and (claimed_by := normalize_todo_claimed_by(item.get("claimed_by"))) - } - ) - return { - "schema_version": TODO_FRONTIER_REVISION_INDEX_SCHEMA_VERSION, - "all": _checkpoint(source_items, included_claims=None), - "unclaimed": _checkpoint(source_items, included_claims={None}), - "by_agent": [ - { - "agent_id": agent_id, - **_checkpoint(source_items, included_claims={None, agent_id}), - } - for agent_id in agent_ids - ], - } + return _request("index", rows=frontier_source_facts(source_items))["index"] def attach_advancement_frontier_revision_index( - summary: dict[str, Any], - source_items: list[dict[str, Any]], - *, - role: str | None, + summary: dict[str, Any], source_items: list[dict[str, Any]], *, role: str | None, ) -> None: - """Attach the complete decision checkpoint only to Agent Todo summaries.""" - if role == "agent": - summary["advancement_frontier_revision_index"] = ( - build_advancement_frontier_revision_index(source_items) - ) + summary["advancement_frontier_revision_index"] = build_advancement_frontier_revision_index(source_items) def advancement_frontier_revision_from_index( - value: Any, - *, - agent_id: str | None, + value: Any, *, agent_id: str | None, ) -> tuple[str | None, str | None, bool] | None: - """Read one complete lane checkpoint; invalid indexes fail closed.""" - - if not isinstance(value, dict): - return None - if value.get("schema_version") != TODO_FRONTIER_REVISION_INDEX_SCHEMA_VERSION: - return None, None, False - normalized_agent_id = normalize_todo_claimed_by(agent_id) - checkpoint: Any = value.get("all") - if normalized_agent_id: - by_agent = value.get("by_agent") - if not isinstance(by_agent, list): - return None, None, False - checkpoint = next( - ( - row - for row in by_agent - if isinstance(row, dict) - and normalize_todo_claimed_by(row.get("agent_id")) - == normalized_agent_id - ), - value.get("unclaimed"), - ) - if not isinstance(checkpoint, dict) or checkpoint.get("complete") is not True: - return None, None, False - revision = str(checkpoint.get("frontier_revision") or "").strip() - updated_at = str(checkpoint.get("frontier_updated_at") or "").strip() - if not revision or parse_timestamp(updated_at) is None: - return None, None, False - return revision, updated_at, True + return _checkpoint_tuple(_request("read", index=value, + agent_id=normalize_todo_claimed_by(agent_id))["checkpoint"]) diff --git a/loopx/control_plane/todos/frontier_revision.ts b/loopx/control_plane/todos/frontier_revision.ts new file mode 100644 index 0000000000..ded2e2236f --- /dev/null +++ b/loopx/control_plane/todos/frontier_revision.ts @@ -0,0 +1,154 @@ +/** Complete advancement frontier identity and long-chain checkpoint policy. + * Python supplies normalized legacy facts and the exact v0 serialization codec; + * selection, completeness, hashing, thresholds and ACK authority live here. + */ +import { createHash } from "node:crypto"; +import type { JsonObject } from "../effect_program.ts"; +import { requireJsonObject } from "../runtime_decode.ts"; +import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { parseTodoTimestampMicros } from "../runtime_timestamp.ts"; +import { normalizeTodoAgent, stripPythonWhitespace } from "../coordination/todo_agents.ts"; +import { AuthorityStoreProtocolError } from "../coordination/authority_store_codec.ts"; + +const REVISION = "todo_frontier_revision_v0"; +const INDEX = "todo_frontier_revision_index_v0"; +const TRIGGER = "long_todo_chain"; +type Checkpoint = { complete: false } | { + complete: true; frontier_revision: string; frontier_updated_at: string; +}; +type Row = { + id: string; claim: string | null; excluded: string[]; + updated: string; serialized: string; advancement: boolean; +}; +type LongChainObservation = { + trigger_count: number; + count_kind: "selectable_advancement_todos" | "selectable_open_todos"; + selectable_open_count: number; selectable_advancement_count: number; + current_agent_claimed_advancement_count: number; unclaimed_advancement_count: number; + threshold: 15 | 20; agent_id: string | null; + frontier_revision: string | null; frontier_revision_complete: boolean; +}; +type AckDecision = {acknowledged: boolean; rearmed_after_obligation_id: string | null}; +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) : ""; +function agentId(value: unknown): string | null { + try { return normalizeTodoAgent(value, "agent_id"); } + catch (error) { + if (error instanceof AuthorityStoreProtocolError) return null; + throw error; + } +} +const strings = (value: unknown): string[] => Array.isArray(value) + ? value.filter((item): item is string => typeof item === "string").map(item => item.trim()).filter(Boolean) : []; +const count = (value: unknown): number => { + const number = Number(value ?? 0); + return Number.isFinite(number) ? Math.max(0, Math.trunc(number)) : 0; +}; + +function decodeRows(value: unknown): Row[] | null { + if (value == null) return null; + if (!Array.isArray(value)) throw new EffectRuntimeRequestError("frontier source must be an array"); + return value.map(raw => { + const row = requireJsonObject(raw, "frontier row"); + if (typeof row.serialized !== "string" || typeof row.advancement !== "boolean") { + throw new EffectRuntimeRequestError("frontier row codec facts are missing"); + } + return {id: text(row.id), claim: text(row.claim) || null, + excluded: strings(row.excluded), updated: text(row.updated), + serialized: row.serialized, advancement: row.advancement}; + }); +} + +function checkpoint(rows: Row[] | null, agent: string | null, unclaimedOnly = false): Checkpoint { + if (rows === null) return {complete: false}; + const selected = rows.filter(row => row.advancement && + (!unclaimedOnly || row.claim === null) && + (!agent || ((row.claim === null || row.claim === agent) && !row.excluded.includes(agent)))); + if (selected.length === 0) return {complete: false}; + const ids = new Set(); + let latest: bigint | null = null; + let updated = ""; + for (const row of selected) { + const instant = parseTodoTimestampMicros(row.updated); + if (!row.id || ids.has(row.id) || instant === null) return {complete: false}; + ids.add(row.id); + if (latest === null || instant > latest) { latest = instant; updated = row.updated; } + } + selected.sort((a, b) => a.id < b.id ? -1 : a.id > b.id ? 1 : 0); + const digest = createHash("sha256").update(`[${selected.map(row => row.serialized).join(",")}]`).digest("hex"); + return {complete: true, frontier_revision: `${REVISION}:${digest.slice(0, 24)}`, + frontier_updated_at: updated}; +} + +function readIndex(value: unknown, agent: string | null): Checkpoint | null { + if (value === null || value === undefined || typeof value !== "object" || Array.isArray(value)) return null; + const index = object(value); + if (index.schema_version !== INDEX) return {complete: false}; + let raw = index.all; + if (agent) { + if (!Array.isArray(index.by_agent)) return {complete: false}; + const matches = index.by_agent.filter(row => agentId(object(row).agent_id) === agent); + if (matches.length > 1) return {complete: false}; + raw = matches[0] ?? index.unclaimed; + } + const entry = object(raw); + const revision = text(entry.frontier_revision), updated = text(entry.frontier_updated_at); + if (entry.complete !== true || !revision || parseTodoTimestampMicros(updated) === null) return {complete: false}; + return {complete: true, frontier_revision: revision, frontier_updated_at: updated}; +} + +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 === "read") return {checkpoint: readIndex(request.index, agent)}; + const rows = decodeRows(request.rows); + if (request.operation === "select") return {checkpoint: checkpoint(rows, agent)}; + if (request.operation !== "index") throw new EffectRuntimeRequestError("unsupported frontier revision operation"); + // An excluded agent can have no claimed rows. It still needs its own lane; + // falling back to the global unclaimed checkpoint would include excluded work. + const agents = [...new Set((rows ?? []).filter(row => row.advancement) + .flatMap(row => [...(row.claim ? [row.claim] : []), ...row.excluded]))].sort(); + return {index: {schema_version: INDEX, all: checkpoint(rows, null), + unclaimed: checkpoint(rows, null, true), + by_agent: agents.map(agent_id => ({agent_id, ...checkpoint(rows, agent_id)}))}}; +} + +function classifyAck(observation: LongChainObservation, value: unknown): AckDecision { + const ack = object(value), delta = object(ack.semantic_delta); + const id = text(delta.obligation_id); + const rejected = {acknowledged: false, rearmed_after_obligation_id: null}; + if (ack.recorded !== true || delta.accepted !== true || + !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); + return text(row.kind) === TRIGGER && text(row.frontier_revision) === observation.frontier_revision; + }); + return {acknowledged: matches, rearmed_after_obligation_id: matches ? null : id}; +} + +export function evaluateLongTodoChain(value: unknown): JsonObject { + const request = requireJsonObject(value, "long chain request"); + if (request.schema_version !== "long_todo_chain_request_v0") throw new EffectRuntimeRequestError("long chain schema mismatch"); + if (request.operation !== "observe") throw new EffectRuntimeRequestError("unsupported long chain operation"); + const summary = object(request.summary), frontier = object(request.frontier_counts); + const current = count(frontier.current_agent_claimed_advancement_count); + const unclaimed = count(frontier.unclaimed_advancement_count); + const advancement = current + unclaimed; + const open = Math.max(advancement, request.summary == null ? count(object(request.agent_counts).open) : + count(summary.current_agent_claimed_open_count) + count(summary.unclaimed_open_count)); + const threshold = advancement >= 15 ? 15 : open >= 20 && advancement > 0 ? 20 : null; + if (threshold === null) return {observation: null, decision: null}; + const agent = agentId(request.agent_id); + const revision = readIndex(summary.advancement_frontier_revision_index, agent) + ?? checkpoint(decodeRows(request.rows), agent); + const observation: LongChainObservation = {trigger_count: threshold === 15 ? advancement : open, + count_kind: threshold === 15 ? "selectable_advancement_todos" : "selectable_open_todos", + selectable_open_count: open, selectable_advancement_count: advancement, + current_agent_claimed_advancement_count: current, unclaimed_advancement_count: unclaimed, + threshold, agent_id: agent, frontier_revision: revision.complete ? revision.frontier_revision : null, + frontier_revision_complete: revision.complete}; + return {observation, decision: classifyAck(observation, request.ack)}; +} diff --git a/tests/control_plane/test_canonical_frontier_revision.py b/tests/control_plane/test_canonical_frontier_revision.py new file mode 100644 index 0000000000..f87f7172fd --- /dev/null +++ b/tests/control_plane/test_canonical_frontier_revision.py @@ -0,0 +1,102 @@ +"""Full-source frontier replay on the shared complex fixture and real file store.""" +from copy import deepcopy +import json +from pathlib import Path +import subprocess + +import pytest + +from canonical_authority_fixture import initialize_canonical_authority +from test_goal_amendment_proposal import _write_fixture, GOAL_ID +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.status import active_state_todo_fields + + +def _fixture(): + module = (Path(__file__).resolve().parents[2] / + "tests/control_plane_ts/production_scale_coordination_fixture.ts").as_uri() + process = subprocess.run([ + "node", "--no-warnings", "--experimental-strip-types", "--input-type=module", "-e", + f"import {{productionScaleCoordinationFixture}} from {json.dumps(module)};" + "process.stdout.write(JSON.stringify(productionScaleCoordinationFixture(process.argv[1])));", GOAL_ID, + ], check=True, capture_output=True, text=True, timeout=30) + return json.loads(process.stdout)["projection"] + + +def _commit_variant(paths, operation, todo_id): + module = (Path(__file__).resolve().parents[2] / "loopx/control_plane/coordination") + script = ( + f"import {{FileAuthorityStore}} from {json.dumps((module / 'file_authority_store.ts').as_uri())};" + f"import {{canonicalAuthoritySha256}} from {json.dumps((module / 'authority_store_codec.ts').as_uri())};" + "const s=new FileAuthorityStore(process.argv[1],process.argv[2]);const h=await s.loadAuthority();" + "const t=h.head.todos.find(t=>t.todo_id===process.argv[4]);" + "if(process.argv[3]==='remove-exclusion')delete t.excluded_agents;else t.priority='P0';" + "h.head.todo_read_model.records_sha256=canonicalAuthoritySha256(h.head.todos);" + "const r=await s.commitAuthority({expected_provider_revision:h.provider_revision," + "operation_id:process.argv[3],events:[],receipts:[],next_projection:h.head});" + "if(r.status!=='applied')throw Error(JSON.stringify(r));" + ) + subprocess.run(["node", "--no-warnings", "--experimental-strip-types", "--input-type=module", "-e", + script, str(paths["runtime"] / "authority/file-v0"), GOAL_ID, operation, todo_id], + check=True, capture_output=True, text=True, timeout=30) + + +@pytest.mark.parametrize("display", ["stale", "missing"]) +def test_complex_canonical_frontier_ack_tracks_only_selectable_material_changes(tmp_path, display): + paths = _write_fixture(tmp_path) + projection = _fixture() + # The shared fixture intentionally includes incomplete historical timestamps. + # This variant exercises a complete checkpoint without changing that fixture's contract. + for record in projection["todos"]: + record.setdefault("updated_at", "2026-09-01T00:00:00.000001Z") + for index in range(30): + projection["todos"].append({ + "schema_version": "todo_item_v0", "todo_id": f"todo_zz_frontier_{index:03}", + "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", + **({"excluded_agents": ["agent-a"]} if index == 29 else {}), + }) + from loopx.control_plane.coordination.local_authority_shadow_projection import canonical_bytes + from hashlib import sha256 + projection["todo_read_model"].update(todo_count=len(projection["todos"]), + records_sha256=sha256(canonical_bytes(projection["todos"])).hexdigest()) + initialize_canonical_authority(paths["runtime"], GOAL_ID, projection, state_path=paths["state_file"]) + if display == "missing": + paths["state_file"].unlink() + before = paths["state_file"].read_bytes() if paths["state_file"].exists() else None + goal = json.loads(paths["registry"].read_text())["goals"][0] + + def observe(ack=None): + summary = active_state_todo_fields(goal, runtime_root=paths["runtime"])["agent_todos"] + # The material edit starts outside the hot-path executable backlog. + if ack is None: + assert not any(item.get("todo_id") == "todo_zz_frontier_028" + for item in summary["executable_backlog_items"]) + alignment = project_shared_goal_alignment(goal_id=GOAL_ID, agent_id="agent-a", project=paths["project"]) + return evaluate_long_todo_chain(agent_todo_summary=summary, agent_counts={}, + frontier_counts=alignment["frontier_counts"], agent_id="agent-a", latest_replan_ack=ack) + + initial, _ = observe() + assert initial is not None and initial.frontier_revision_complete + 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}]}} + 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 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 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"]) + assert code == 0, result + assert (paths["state_file"].read_bytes() if paths["state_file"].exists() else None) == before diff --git a/tests/control_plane/test_frontier_revision_scope.py b/tests/control_plane/test_frontier_revision_scope.py new file mode 100644 index 0000000000..d31edc3ea1 --- /dev/null +++ b/tests/control_plane/test_frontier_revision_scope.py @@ -0,0 +1,73 @@ +from copy import deepcopy +from hashlib import sha256 +import json + +from loopx.control_plane.todos.frontier_revision import ( + advancement_frontier_revision_from_index, + build_advancement_frontier_revision_index, + selectable_advancement_frontier_revision, +) + + +def rows(): + return [ + {"todo_id": "todo_a", "task_class": "advancement_task", "status": "open", + "text": "验证边界", "updated_at": "2026-09-01T00:00:00.000001Z"}, + {"todo_id": "todo_b", "task_class": "advancement_task", "status": "open", + "text": "Independent work", "excluded_agents": ["worker-a"], + "updated_at": "2026-09-01T00:00:00.000002Z"}, + ] + + +def test_excluded_unclaimed_work_does_not_rearm_this_agents_frontier(): + before = rows() + after = deepcopy(before) + after[1]["priority"] = "P0" + after[1]["updated_at"] = "2026-09-02T00:00:00Z" + for source in (before, after): + direct = selectable_advancement_frontier_revision(source, agent_id="worker-a") + indexed = advancement_frontier_revision_from_index( + build_advancement_frontier_revision_index(source), agent_id="worker-a", + ) + assert direct == indexed + assert selectable_advancement_frontier_revision(before, agent_id="worker-a") == ( + selectable_advancement_frontier_revision(after, agent_id="worker-a") + ) + assert selectable_advancement_frontier_revision(before, agent_id="worker-b") != ( + selectable_advancement_frontier_revision(after, agent_id="worker-b") + ) + + +def test_duplicate_identity_cannot_be_a_complete_frontier(): + source = rows() + source[1]["todo_id"] = source[0]["todo_id"] + assert selectable_advancement_frontier_revision(source, agent_id=None) == (None, None, False) + + +def test_removing_exclusion_rearms_newly_available_work(): + source = rows() + before = selectable_advancement_frontier_revision(source, agent_id="worker-a") + source[1].pop("excluded_agents") + assert selectable_advancement_frontier_revision(source, agent_id="worker-a") != before + + +def test_v0_codec_preserves_unicode_hash_and_microsecond_ordering(): + # Explicit protocol records, not output generated by the implementation. + material = [{"status": "open", "task_class": "advancement_task", "text": "边界 🧭", + "todo_id": "todo_a"}, + {"status": "done", "task_class": "advancement_task", "todo_id": "todo_b"}] + encoded = json.dumps(material, ensure_ascii=True, separators=(",", ":"), sort_keys=True) + expected = "todo_frontier_revision_v0:" + sha256(encoded.encode()).hexdigest()[:24] + source = [{**material[1], "completed_at": "2026-09-01T08:00:00.000001+08:00"}, + {**material[0], "updated_at": "2026-09-01T00:00:00.000002Z"}] + assert selectable_advancement_frontier_revision(source, agent_id=None) == ( + expected, "2026-09-01T00:00:00.000002Z", True) + + +def test_legacy_agent_spelling_in_index_remains_readable(): + source = rows()[:1] + source[0]["claimed_by"] = "Worker A" + index = build_advancement_frontier_revision_index(source) + index["by_agent"][0]["agent_id"] = " Worker\u0085A " + assert advancement_frontier_revision_from_index(index, agent_id="worker-a") == ( + selectable_advancement_frontier_revision(source, agent_id="Worker A")) diff --git a/tests/control_plane_ts/frontier_revision.test.ts b/tests/control_plane_ts/frontier_revision.test.ts new file mode 100644 index 0000000000..4c82235094 --- /dev/null +++ b/tests/control_plane_ts/frontier_revision.test.ts @@ -0,0 +1,81 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import {projectAdvancementFrontier, evaluateLongTodoChain} from "../../loopx/control_plane/todos/frontier_revision.ts"; + +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})}; +} +function project(rows: ReturnType[], agent_id: string | null = null) { + return projectAdvancementFrontier({schema_version: "todo_frontier_revision_request_v0", + operation: "select", rows, agent_id}).checkpoint as Record; +} +function observe(overrides: Record = {}) { + return evaluateLongTodoChain({schema_version: "long_todo_chain_request_v0", operation: "observe", + agent_id: "worker-a", summary: {current_agent_claimed_open_count: 15, unclaimed_open_count: 0}, + frontier_counts: {current_agent_claimed_advancement_count: 15, unclaimed_advancement_count: 0}, + rows: [row("todo_a", "worker-a")], ...overrides}); +} + +test("revision is order-independent and maintenance timestamps do not change material identity", () => { + const rows = [row("todo_b"), row("todo_a")]; + assert.deepEqual(project(rows), project([...rows].reverse())); + const changed = [{...rows[0], updated: "2026-09-02T00:00:00Z"}, rows[1]]; + assert.equal(project(rows).frontier_revision, project(changed).frontier_revision); + assert.notEqual(project(rows).frontier_updated_at, project(changed).frontier_updated_at); +}); + +test("excluded-only agents receive an explicit checkpoint instead of the unclaimed fallback", () => { + const rows = [row("todo_a"), row("todo_b", null, ["worker-a"])]; + const index = projectAdvancementFrontier({schema_version: "todo_frontier_revision_request_v0", + operation: "index", rows}).index; + const read = projectAdvancementFrontier({schema_version: "todo_frontier_revision_request_v0", + operation: "read", index, agent_id: "worker-a"}).checkpoint; + assert.deepEqual(read, project([rows[0]], "worker-a")); + assert.notDeepEqual(read, project(rows)); +}); + +test("duplicate identities and malformed or absent revision facts cannot authorize an ACK", () => { + for (const rows of [[], [row("todo_a"), row("todo_a")], [{...row("todo_a"), updated: "bad"}]]) { + assert.deepEqual(project(rows), {complete: false}); + } + const result = observe({summary: {advancement_frontier_revision_index: {schema_version: "bad"}}}); + assert.equal((result.observation as Record).frontier_revision_complete, false); + assert.deepEqual(result.decision, {acknowledged: false, rearmed_after_obligation_id: null}); +}); + +test("long-chain thresholds remain 15 advancement or 20 selectable with advancement", () => { + assert.notEqual(observe().observation, null); + assert.equal(observe({frontier_counts: {current_agent_claimed_advancement_count: 14}}).observation, null); + const open = {current_agent_claimed_open_count: 20, unclaimed_open_count: 0}; + assert.notEqual(observe({summary: open, frontier_counts: {unclaimed_advancement_count: 1}}).observation, null); + assert.equal(observe({summary: open, frontier_counts: {}}).observation, null); +}); + +test("only exact accepted checkpoint suppresses a repeated trigger; material change rearms", () => { + 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_revision: observation.frontier_revision}]}}; + assert.deepEqual(observe({ack}).decision, {acknowledged: true, rearmed_after_obligation_id: null}); + assert.deepEqual(observe({ack, rows: [row("todo_new", "worker-a")]}).decision, + {acknowledged: false, rearmed_after_obligation_id: "replan-0123456789abcdef"}); + for (const invalid of [null, {...ack, recorded: false}, {...ack, semantic_delta: {...ack.semantic_delta, + accepted: false}}, {...ack, semantic_delta: {...ack.semantic_delta, trigger_kinds: "long_todo_chain"}}]) { + assert.equal((observe({ack: invalid}).decision as Record).acknowledged, false); + } +}); + +test("duplicate agent checkpoints fail closed instead of selecting the first receipt", () => { + const entry = {agent_id: "worker-a", ...project([row("todo_a")])}; + const index = {schema_version: "todo_frontier_revision_index_v0", by_agent: [entry, entry]}; + assert.deepEqual(projectAdvancementFrontier({schema_version: "todo_frontier_revision_request_v0", + operation: "read", index, agent_id: "worker-a"}).checkpoint, {complete: false}); +}); + +test("the typed transport rejects unsupported operations and missing source codec facts", () => { + assert.throws(() => projectAdvancementFrontier({schema_version: "old"})); + assert.throws(() => projectAdvancementFrontier({schema_version: "todo_frontier_revision_request_v0", + operation: "index", rows: [{}]})); + assert.throws(() => evaluateLongTodoChain({schema_version: "long_todo_chain_request_v0", operation: "unknown"})); +}); From b22e2cdb6facb05b054b33011ade9f8903db2468 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sat, 12 Sep 2026 02:10:46 +0800 Subject: [PATCH 2/5] docs(rfc): record bounded frontier checkpoint consumer closure Signed-off-by: huangruiteng --- .../shared-goal-authority-state-provider-v0.md | 8 ++++++++ ...d-goal-authority-state-provider-v0.zh-CN.md | 7 +++++++ .../typescript-control-plane-migration-v0.md | 18 ++++++++++++++++++ ...escript-control-plane-migration-v0.zh-CN.md | 13 +++++++++++++ 4 files changed, 46 insertions(+) diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md index ba09e6543f..dbda6b8fa8 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -2655,6 +2655,14 @@ and retired Python selectors. Real FileAuthorityStore CLI tests cover missing an stale display without writing it back. This is consumer-rule consolidation, not a transaction/store change, provider qualification or whole-Goal cutover. +Long-chain checkpoint reads now use one typed frontier revision/ACK policy across +legacy and canonical sources (TS RFC T3). The index is built before display limits; +excluded work cannot spuriously rearm another Agent, and an incomplete or ambiguous +checkpoint cannot acknowledge the chain. Python keeps the persisted v0 codec, not +a second revision/threshold policy. This is a consumer change: it adds no provider, +commit receipt, promotion route or Markdown writer. Existing CAS/replay, permanent +projection and D1–D3 qualification remain unchanged. + The original direction remains; execution cards expand these stages rather than cancel them: 1. **Close TS transactions and consumers.** Follow [T0–T3](typescript-control-plane-migration-v0.md#execution-cards-after-the-current-stack) to consolidate rules and delete duplicate decisions. diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md index a05b7d9dc2..4e026e3cbd 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md @@ -2102,6 +2102,13 @@ selector 见 TS RFC 的 T3 卡。真实 FileAuthorityStore CLI 测试覆盖展 不写回。这是 consumer 规则收拢,不是 transaction/store 改造、provider 资格化或 整 Goal cutover。 +长链 checkpoint 读取现由同一 typed frontier revision/ACK 策略处理 legacy 与 canonical +来源(TS RFC T3)。Index 在展示限制之前生成;被排除工作不能误触发该 Agent, +不完整或有歧义的 checkpoint 不能确认长链已处理。Python 保留持久 v0 codec, +不再持有第二套 revision/threshold 策略。这是 consumer 改造,不新增 provider、 +commit receipt、promotion 路由或 Markdown writer。既有 CAS/replay、永久投影与 +D1–D3 资格化要求保持不变。 + 以下规划保留原有方向;执行卡是它们的展开,不是替代或取消: 1. **闭合 TS 事务与 consumer。** 按 [T0–T3](typescript-control-plane-migration-v0.zh-CN.md#当前-stack-合入后的执行卡) 收口规则并删除重复决策。 diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index 8c6f53b36a..41d52f491c 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -588,6 +588,24 @@ read or resume evaluation is added. Remaining T3 work includes consumers that reconstruct diagnostics from compact summaries; do not call those migrated. This does not close T1/T2, all T3 consumers, or any durability/promotion hold. +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. + The list-filter consumer now uses `compact_evaluated_todo_group` instead of re-running resume evaluation on active-only rows. Initial parsing/canonical reads still evaluate against the full source through the TS owner; filtering requires 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 6aefd2a769..bae6db2dea 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 @@ -449,6 +449,19 @@ Todo,同一 Todo 的不同展示不重复计算,权威空 backlog 不再复 target capability 是修复产出,不是安装或授权。没有新 provider/inventory/enablement/ promotion;压缩候选来源的上限和其余 T3 consumer 仍需分别闭合。 +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 完成。 + 列表过滤现改用 `compact_evaluated_todo_group`,不再用仅活动项重算 resume。 初始解析/canonical 读取仍通过 TS owner 在完整来源上求值;过滤要求匹配的已求值 条件,不能把归档中的已完成依赖变成丢失。共享合成 fixture 增补“有 scope 无 outcome” From 1949e8532bcc1a5fea1e2df5ba61faa0a480f0cb Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sat, 12 Sep 2026 12:49:16 +0800 Subject: [PATCH 3/5] fix(frontier): reject malformed typed index responses Signed-off-by: huangruiteng --- loopx/control_plane/todos/frontier_revision.py | 5 ++++- tests/control_plane/test_frontier_revision_scope.py | 12 ++++++++++++ 2 files changed, 16 insertions(+), 1 deletion(-) diff --git a/loopx/control_plane/todos/frontier_revision.py b/loopx/control_plane/todos/frontier_revision.py index da4e2debf5..bdbeca13e6 100644 --- a/loopx/control_plane/todos/frontier_revision.py +++ b/loopx/control_plane/todos/frontier_revision.py @@ -99,7 +99,10 @@ def selectable_advancement_frontier_revision( def build_advancement_frontier_revision_index( source_items: list[dict[str, Any]], ) -> dict[str, Any]: - return _request("index", rows=frontier_source_facts(source_items))["index"] + index = _request("index", rows=frontier_source_facts(source_items)).get("index") + if not isinstance(index, dict): + raise TypeError("typed frontier revision response index must be an object") + return index def attach_advancement_frontier_revision_index( diff --git a/tests/control_plane/test_frontier_revision_scope.py b/tests/control_plane/test_frontier_revision_scope.py index d31edc3ea1..4fff2ed23e 100644 --- a/tests/control_plane/test_frontier_revision_scope.py +++ b/tests/control_plane/test_frontier_revision_scope.py @@ -2,6 +2,10 @@ from hashlib import sha256 import json +import pytest + +from loopx.control_plane.todos import frontier_revision + from loopx.control_plane.todos.frontier_revision import ( advancement_frontier_revision_from_index, build_advancement_frontier_revision_index, @@ -19,6 +23,14 @@ def rows(): ] +@pytest.mark.parametrize("index", [None, [], "invalid", 1, True]) +def test_index_adapter_rejects_non_object_typed_response(monkeypatch, index): + monkeypatch.setattr(frontier_revision, "effect_runtime_result", + lambda *_args: {"index": index}) + with pytest.raises(TypeError, match="typed frontier revision response index must be an object"): + build_advancement_frontier_revision_index(rows()) + + def test_excluded_unclaimed_work_does_not_rearm_this_agents_frontier(): before = rows() after = deepcopy(before) From 0ed96be86cfed9c0b41d022e7c81fbab62c971bb Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sat, 12 Sep 2026 12:53:42 +0800 Subject: [PATCH 4/5] fix(frontier): preserve large source reads with lossless facts transport Signed-off-by: huangruiteng --- .../control_plane/todos/frontier_revision.py | 15 +++++++++++-- .../control_plane/todos/frontier_revision.ts | 15 +++++++++++++ .../test_frontier_revision_scope.py | 16 ++++++++++++++ .../frontier_revision.test.ts | 22 +++++++++++++++++++ 4 files changed, 66 insertions(+), 2 deletions(-) diff --git a/loopx/control_plane/todos/frontier_revision.py b/loopx/control_plane/todos/frontier_revision.py index bdbeca13e6..eb732f983d 100644 --- a/loopx/control_plane/todos/frontier_revision.py +++ b/loopx/control_plane/todos/frontier_revision.py @@ -2,7 +2,9 @@ from __future__ import annotations +import base64 import json +import zlib from typing import Any from ..effect_runtime import effect_runtime_result @@ -49,11 +51,13 @@ ) -def frontier_source_facts(source_items: list[dict[str, Any]] | None) -> list[dict[str, Any]] | None: +def frontier_source_facts( + source_items: list[dict[str, Any]] | None, +) -> list[dict[str, Any]] | dict[str, str] | None: """Legacy codecs only; TS selects lanes and builds complete revision identity.""" if not isinstance(source_items, list): return None - return [ + rows = [ { "id": str(item.get("todo_id") or "").strip(), "claim": normalize_todo_claimed_by(item.get("claimed_by")), @@ -67,6 +71,13 @@ def frontier_source_facts(source_items: list[dict[str, Any]] | None) -> list[dic } for item in source_items if isinstance(item, dict) ] + # Lossless transport codec only: never truncate material identity or raise + # the shared Effect request limit for large history/frontier reads. + raw = json.dumps(rows, ensure_ascii=True, separators=(",", ":")).encode() + if len(raw) < 512 * 1024: + return rows + return {"encoding": "deflate-base64-json-v0", + "data": base64.b64encode(zlib.compress(raw)).decode("ascii")} def _request(operation: str, **facts: Any) -> dict[str, Any]: diff --git a/loopx/control_plane/todos/frontier_revision.ts b/loopx/control_plane/todos/frontier_revision.ts index ded2e2236f..71d7771911 100644 --- a/loopx/control_plane/todos/frontier_revision.ts +++ b/loopx/control_plane/todos/frontier_revision.ts @@ -3,6 +3,7 @@ * selection, completeness, hashing, thresholds and ACK authority live here. */ import { createHash } from "node:crypto"; +import {inflateSync} from "node:zlib"; import type { JsonObject } from "../effect_program.ts"; import { requireJsonObject } from "../runtime_decode.ts"; import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; @@ -48,6 +49,20 @@ const count = (value: unknown): number => { function decodeRows(value: unknown): Row[] | null { if (value == null) return null; + if (!Array.isArray(value)) { + const encoded = requireJsonObject(value, "frontier rows transport"); + if (encoded.encoding !== "deflate-base64-json-v0" || typeof encoded.data !== "string" || + encoded.data.length > 2 * 1024 * 1024 || !/^[A-Za-z0-9+/]+={0,2}$/.test(encoded.data)) { + throw new EffectRuntimeRequestError("invalid frontier rows transport"); + } + try { + value = JSON.parse(inflateSync(Buffer.from(encoded.data, "base64"), { + maxOutputLength: 64 * 1024 * 1024, + }).toString("utf8")); + } catch { + throw new EffectRuntimeRequestError("frontier rows must be valid compressed JSON within 64 MiB"); + } + } if (!Array.isArray(value)) throw new EffectRuntimeRequestError("frontier source must be an array"); return value.map(raw => { const row = requireJsonObject(raw, "frontier row"); diff --git a/tests/control_plane/test_frontier_revision_scope.py b/tests/control_plane/test_frontier_revision_scope.py index 4fff2ed23e..8afd88feb5 100644 --- a/tests/control_plane/test_frontier_revision_scope.py +++ b/tests/control_plane/test_frontier_revision_scope.py @@ -83,3 +83,19 @@ def test_legacy_agent_spelling_in_index_remains_readable(): index["by_agent"][0]["agent_id"] = " Worker\u0085A " assert advancement_frontier_revision_from_index(index, agent_id="worker-a") == ( selectable_advancement_frontier_revision(source, agent_id="Worker A")) + + +def test_large_frontier_transport_preserves_exact_v0_identity_and_tail_edits(): + material = [{"todo_id": f"todo_{index:04}", "task_class": "advancement_task", + "text": "Complete synthetic description " * 40} for index in range(2500)] + encoded = json.dumps(material, ensure_ascii=True, separators=(",", ":"), sort_keys=True) + assert len(encoded.encode()) > 2 * 1024 * 1024 + expected = "todo_frontier_revision_v0:" + sha256(encoded.encode()).hexdigest()[:24] + source = [{**item, "updated_at": "2026-09-01T00:00:00Z"} for item in material] + assert selectable_advancement_frontier_revision(source, agent_id=None) == ( + expected, "2026-09-01T00:00:00Z", True) + assert advancement_frontier_revision_from_index( + build_advancement_frontier_revision_index(source), agent_id=None, + ) == (expected, "2026-09-01T00:00:00Z", True) + source[-1]["text"] += " Material edit at the tail." + assert selectable_advancement_frontier_revision(source, agent_id=None)[0] != expected diff --git a/tests/control_plane_ts/frontier_revision.test.ts b/tests/control_plane_ts/frontier_revision.test.ts index 4c82235094..35ddd5384f 100644 --- a/tests/control_plane_ts/frontier_revision.test.ts +++ b/tests/control_plane_ts/frontier_revision.test.ts @@ -1,5 +1,6 @@ import assert from "node:assert/strict"; import test from "node:test"; +import {deflateSync} from "node:zlib"; import {projectAdvancementFrontier, evaluateLongTodoChain} from "../../loopx/control_plane/todos/frontier_revision.ts"; function row(id: string, claim: string | null = null, excluded: string[] = []) { @@ -17,6 +18,27 @@ function observe(overrides: Record = {}) { rows: [row("todo_a", "worker-a")], ...overrides}); } +test("lossless compressed rows preserve index and ACK semantics and reject malformed transport", () => { + const rows = [row("todo_a", "worker-a"), row("todo_b", null, ["worker-a"])]; + const compressed = {encoding: "deflate-base64-json-v0", + data: deflateSync(JSON.stringify(rows)).toString("base64")}; + for (const operation of ["index", "select"]) { + const request = {schema_version: "todo_frontier_revision_request_v0", operation, agent_id: "worker-a"}; + assert.deepEqual(projectAdvancementFrontier({...request, rows: compressed}), + projectAdvancementFrontier({...request, rows})); + } + assert.deepEqual(observe({rows: compressed}), observe({rows})); + for (const invalid of [ + {...compressed, encoding: "unknown"}, {...compressed, data: "!"}, + {...compressed, data: "AAAA"}, + {...compressed, data: deflateSync("{}").toString("base64")}, + {...compressed, data: deflateSync("x".repeat(64 * 1024 * 1024 + 1)).toString("base64")}, + ]) { + assert.throws(() => projectAdvancementFrontier({schema_version: "todo_frontier_revision_request_v0", + operation: "index", rows: invalid})); + } +}); + test("revision is order-independent and maintenance timestamps do not change material identity", () => { const rows = [row("todo_b"), row("todo_a")]; assert.deepEqual(project(rows), project([...rows].reverse())); From e2a0eedbc5ca152f575c157e046e214555f58eed Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sat, 12 Sep 2026 12:53:42 +0800 Subject: [PATCH 5/5] docs(frontier): clarify bounded lossless transport Signed-off-by: huangruiteng --- .../rfcs/typescript-control-plane-migration-v0.md | 6 ++++++ .../rfcs/typescript-control-plane-migration-v0.zh-CN.md | 5 +++++ 2 files changed, 11 insertions(+) diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index 41d52f491c..d006c13e27 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -606,6 +606,12 @@ 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. +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 +rejects malformed payloads and inflation beyond 64 MiB; it never truncates rows +or silently falls back to Python decisions. Real completed-history HTTP reads +and complete-checkpoint tail edits guard against transport-size regressions. + The list-filter consumer now uses `compact_evaluated_todo_group` instead of re-running resume evaluation on active-only rows. Initial parsing/canonical reads still evaluate against the full source through the TS owner; filtering requires 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 bae6db2dea..67194c7fc6 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 @@ -462,6 +462,11 @@ replan。这些是明确的只读语义修正,不是执行授权。既有 cano 新可用工作及陈旧/缺失展示;私有快照只读对照结果不公开原始数据。 本批闭合一个 T3 规则组,不代表其余 consumer 或 T1/T2/D1–D3 完成。 +来源 facts 超过 512 KiB 时使用无损 deflate/base64 传输,保留精确 v0 内容和共享 +2 MiB 请求边界。TS 拒绝畸形载荷及解压超过 64 MiB 的输入,不截断 Todo,也不 +静默退回 Python 决策。真实 completed-history HTTP 和完整 checkpoint 尾项变更 +回归保护传输容量语义。 + 列表过滤现改用 `compact_evaluated_todo_group`,不再用仅活动项重算 resume。 初始解析/canonical 读取仍通过 TS owner 在完整来源上求值;过滤要求匹配的已求值 条件,不能把归档中的已完成依赖变成丢失。共享合成 fixture 增补“有 scope 无 outcome”