From 7d08b68b3a99984f50a969d6274910c6396bbd58 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 21 Sep 2026 18:54:22 +0800 Subject: [PATCH 1/2] fix(periodic-report): read canonical progress and retry decisions Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../periodic_report/pending_intent.py | 76 +++---- .../periodic_report/post_writeback_hook.py | 66 +++--- .../project_progress_snapshot.py | 191 +++++++----------- .../periodic_report/todo_source.py | 70 +++++++ .../capabilities/periodic_report_progress.ts | 96 +++++++++ .../coordination/local_authority.py | 2 + .../control_plane/effect_runtime_handlers.ts | 3 + 7 files changed, 302 insertions(+), 202 deletions(-) create mode 100644 loopx/capabilities/periodic_report/todo_source.py create mode 100644 loopx/control_plane/capabilities/periodic_report_progress.ts diff --git a/loopx/capabilities/periodic_report/pending_intent.py b/loopx/capabilities/periodic_report/pending_intent.py index ce3721f657..f9678c85b8 100644 --- a/loopx/capabilities/periodic_report/pending_intent.py +++ b/loopx/capabilities/periodic_report/pending_intent.py @@ -15,12 +15,12 @@ POST_WRITEBACK_HOOK_RECEIPT_SCHEMA_VERSION, InteractionProjectionHookRegistration, ) -from ...control_plane.todos.active_state_todo_parser import parse_active_state_todos +from ...control_plane.effect_runtime import effect_runtime_result +from .todo_source import read_report_todo_source from ...registry import ( atomic_write_json, find_registry_goal, read_json, - resolve_state_file, ) from ...todos import add_goal_todo from ...file_lock import LockAcquisitionPolicy, exclusive_file_lock @@ -177,57 +177,45 @@ def _editorial_response_path( ) -def _decision_scope_text(value: object) -> str: - if isinstance(value, Mapping): - return ":".join( - str(value.get(field) or "") - for field in ("kind", "granularity", "scope_key") - ) - return str(value or "") - - def _superseding_approval_revision( *, registry_path: Path, goal_id: str, agent_id: str, receipt: Mapping[str, Any], + runtime_root: Path | None = None, ) -> str | None: approval_scope = str(receipt.get("approval_scope") or "") if receipt.get("status") != "approval_pending" or not approval_scope: return None - registry = read_json(registry_path) - goal = find_registry_goal(registry, goal_id) - if not isinstance(goal, Mapping): - return None - repo = Path(str(goal.get("repo") or "")).expanduser() - state_path = resolve_state_file(repo, str(goal.get("state_file") or "")) - if state_path is None or not state_path.is_file(): - return None - parsed = parse_active_state_todos( - state_path.read_text(encoding="utf-8"), - goal=dict(goal), - state_path=state_path, - item_limit=None, + _, items = read_report_todo_source( + registry_path=registry_path, goal_id=goal_id, runtime_root=runtime_root ) - user_summary = parsed.get("user_todos") - items = user_summary.get("items") if isinstance(user_summary, Mapping) else [] - superseding = [ - item - for item in items or [] - if isinstance(item, Mapping) - and item.get("status") == "done" - and item.get("action_kind") - in {"approve_periodic_report_payload", "cancel_periodic_report_payload"} - and item.get("decision_outcome") in {"reject", "cancel"} - and _decision_scope_text(item.get("decision_scope")) == approval_scope - and str(item.get("bound_agent") or item.get("blocks_agent") or "") == agent_id - ] - if not superseding: - return None - latest = max(superseding, key=lambda item: str(item.get("updated_at") or "")) - revision = f"{latest.get('todo_id')}:{latest.get('updated_at')}" - return hashlib.sha256(revision.encode("utf-8")).hexdigest()[:16] + keys = ( + "todo_id", + "status", + "action_kind", + "decision_outcome", + "decision_scope", + "bound_agent", + "blocks_agent", + "updated_at", + ) + result = effect_runtime_result( + "capabilities.periodic_report.approval_retry.select", + { + "schema_version": "periodic_report_approval_retry_request_v0", + "agent_id": agent_id, + "approval_scope": approval_scope, + "items": [{key: item.get(key) for key in keys} for item in items], + }, + ) + if ( + not isinstance(result, dict) + or result.get("schema_version") != "periodic_report_approval_retry_result_v0" + ): + raise ValueError("periodic-report approval retry selection result mismatch") + return result["revision"] def _load_consumption_receipt( @@ -292,7 +280,7 @@ def _next_attempt_revision( candidate_revisions: list[str] = [] for receipt in receipts: revision = _superseding_approval_revision( - registry_path=registry_path, + runtime_root=runtime_root, registry_path=registry_path, goal_id=goal_id, agent_id=agent_id, receipt=receipt, @@ -540,7 +528,7 @@ def _progress_facts( from ...rollout_event_log import load_rollout_events, rollout_event_log_path snapshot = build_project_progress_snapshot( - registry_path=registry_path, + runtime_root=runtime_root, registry_path=registry_path, goal_id=goal_id, agent_id=agent_id, completed_at=completed_at, diff --git a/loopx/capabilities/periodic_report/post_writeback_hook.py b/loopx/capabilities/periodic_report/post_writeback_hook.py index 8f5bdceee9..d2306ee128 100644 --- a/loopx/capabilities/periodic_report/post_writeback_hook.py +++ b/loopx/capabilities/periodic_report/post_writeback_hook.py @@ -14,7 +14,7 @@ from ...control_plane.goals.goal_frontier import ( build_goal_frontier_projection_from_summaries, ) -from ...control_plane.todos.active_state_todo_parser import parse_active_state_todos +from .todo_source import read_report_todo_source from ...control_plane.todos.quota_summary import summarize_user_todos_for_quota from ...control_plane.todos.todo_index import MAX_TODO_INDEX_ROLLOUT_EVENTS_PER_GOAL from ...history import collect_history, load_registry @@ -24,7 +24,7 @@ from .stage_completion import STAGE_COMPLETION_RECEIPT_SCHEMA from .stage_completion import derive_periodic_report_stage_completion_from_runs from .presets import build_periodic_report_preset_activation -from .project_progress_snapshot import build_project_progress_snapshot_from_state +from .project_progress_snapshot import build_project_progress_snapshot_from_fields from .incremental import read_periodic_report_goal_publication_cursors from .machine_defaults import resolve_goal_periodic_report_subscription from .machine_store import read_periodic_report_machine_defaults @@ -144,18 +144,10 @@ def periodic_report_post_writeback_hooks_for_goal( def _frontier_projection( *, - state_text: str, - goal: Mapping[str, Any], - state_path: Path, + todos: Mapping[str, Any], goal_id: str, agent_id: str, ) -> tuple[dict[str, Any], dict[str, Any], dict[str, Any]]: - todos = parse_active_state_todos( - state_text, - goal=dict(goal), - state_path=state_path, - item_limit=None, - ) raw_user_summary = ( dict(todos.get("user_todos")) if isinstance(todos.get("user_todos"), Mapping) @@ -197,29 +189,31 @@ def build_periodic_report_post_writeback_projection( """Reduce private runtime state to one bounded public-safe stage receipt.""" normalized_agent_id = str(agent_id or "").strip() + if not normalized_agent_id: + return {} + available_capabilities = payload.get("available_capabilities") + if available_capabilities is None and isinstance(payload.get("turn"), Mapping): + available_capabilities = payload["turn"].get("available_capabilities") + events = load_rollout_events( + rollout_event_log_path(runtime_root, goal_id), + limit=MAX_TODO_INDEX_ROLLOUT_EVENTS_PER_GOAL, + ) state = payload.get("state") - state_path_value = ( + state_path = ( state.get("path") if isinstance(state, Mapping) else None ) or payload.get("state_file") - state_path = Path(str(state_path_value or "")).expanduser() - if not normalized_agent_id or not state_path.is_file(): - return {} - registry = load_registry(registry_path) - goal = next( - ( - item - for item in registry_goals(registry) - if str(item.get("id") or "").strip() == str(goal_id or "").strip() - ), - {}, + fields, _ = read_report_todo_source( + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=goal_id, + state_path=Path(state_path) + if isinstance(state_path, str) and state_path + else None, + rollout_events=events, + available_capabilities=available_capabilities, ) - state_text = state_path.read_text(encoding="utf-8") projection, _user_summary, _agent_summary = _frontier_projection( - state_text=state_text, - goal=goal, - state_path=state_path, - goal_id=goal_id, - agent_id=normalized_agent_id, + todos=fields, goal_id=goal_id, agent_id=normalized_agent_id ) history = collect_history( registry_path=registry_path, @@ -289,23 +283,13 @@ def build_periodic_report_post_writeback_projection( publication_cursor = next( (c for c in goal_cursors if c["agent_id"] == normalized_agent_id), None ) - available_capabilities = payload.get("available_capabilities") - if available_capabilities is None and isinstance(payload.get("turn"), Mapping): - available_capabilities = payload["turn"].get("available_capabilities") - project_progress = build_project_progress_snapshot_from_state( - state_text=state_text, - goal=goal, - state_path=state_path, + project_progress = build_project_progress_snapshot_from_fields( + fields=fields, goal_id=goal_id, agent_id=normalized_agent_id, completed_at=str(receipt["completed_at"]), publication_cursor=publication_cursor, goal_cursors=goal_cursors, - available_capabilities=available_capabilities, - rollout_events=load_rollout_events( - rollout_event_log_path(runtime_root, goal_id), - limit=MAX_TODO_INDEX_ROLLOUT_EVENTS_PER_GOAL, - ), ) if publication_cursor is not None and project_progress is None: return {} diff --git a/loopx/capabilities/periodic_report/project_progress_snapshot.py b/loopx/capabilities/periodic_report/project_progress_snapshot.py index 4beb4b7c8c..9e395c05ab 100644 --- a/loopx/capabilities/periodic_report/project_progress_snapshot.py +++ b/loopx/capabilities/periodic_report/project_progress_snapshot.py @@ -1,64 +1,21 @@ from __future__ import annotations from collections.abc import Mapping, Sequence -from datetime import datetime from pathlib import Path from typing import Any from ...control_plane.todos.active_state_todo_parser import parse_active_state_todos from ...control_plane.todos.todo_semantics import todo_item_is_actionable_open -from ...registry import find_registry_goal, read_json, resolve_state_file +from ...control_plane.effect_runtime import EffectRuntimeRejected, effect_runtime_result +from .todo_source import read_report_todo_source from .incremental import select_incremental_project_progress -def _stage_timestamp(value: str) -> datetime | None: - """Parse an offset-aware ISO-8601 timestamp or reject the value. - - Sibling validations (``_validated_snapshot_timestamp``, - ``_actual_work_window``, ``incremental._timestamp``) all reject naive - timestamps, so stage filtering must not compare offset-naive and - offset-aware datetimes either. - """ - - try: - parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) - except (TypeError, ValueError): - return None - if parsed.tzinfo is None: - return None - return parsed - - -def _outcome_completed_at( - item: Mapping[str, Any], *, stage_time: datetime -) -> str | None: - """Return the durable completion timestamp an outcome fact may carry. - - A done todo without a trustworthy completion timestamp (handwritten - agent-lane entries may omit them) must not become an outcome fact: the - frozen fact would fail timestamp validation on every consumption retry. - """ - - raw = str(item.get("completed_at") or item.get("updated_at") or "").strip() - parsed = _stage_timestamp(raw) - if parsed is None or parsed > stage_time: - return None - return raw - - -_META_ACTION_KINDS = frozenset( - { - "consume_periodic_report_intent", - "repair_periodic_report_intent_consumption", - "repair_periodic_report_editorial", - } -) - - def build_project_progress_snapshot( *, registry_path: Path, goal_id: str, + runtime_root: Path | None = None, agent_id: str, completed_at: str, publication_cursor: Mapping[str, Any] | None = None, @@ -68,25 +25,20 @@ def build_project_progress_snapshot( ) -> dict[str, Any] | None: """Build a bounded public-safe progress snapshot at a stage boundary.""" - registry = read_json(registry_path) - goal = find_registry_goal(registry, goal_id) - if not isinstance(goal, Mapping): - raise ValueError("periodic-report Goal is not registered") - repo = Path(str(goal.get("repo") or "")).expanduser() - state_path = resolve_state_file(repo, str(goal.get("state_file") or "")) - if state_path is None or not state_path.is_file(): - raise ValueError("periodic-report active state is unavailable") - return build_project_progress_snapshot_from_state( - state_text=state_path.read_text(encoding="utf-8"), - goal=dict(goal), - state_path=state_path, + fields, _ = read_report_todo_source( + registry_path=registry_path, + goal_id=goal_id, + runtime_root=runtime_root, + rollout_events=rollout_events, + available_capabilities=available_capabilities, + ) + return build_project_progress_snapshot_from_fields( + fields=fields, goal_id=goal_id, agent_id=agent_id, completed_at=completed_at, publication_cursor=publication_cursor, goal_cursors=goal_cursors, - available_capabilities=available_capabilities, - rollout_events=rollout_events, ) @@ -103,7 +55,7 @@ def build_project_progress_snapshot_from_state( available_capabilities: Any = None, rollout_events: list[dict[str, Any]] | None = None, ) -> dict[str, Any] | None: - """Build a progress snapshot from one already-read authoritative state. + """Legacy text adapter; provider-aware callers use the shared Todo source. Evidence is selected for the Goal, not for the calling lane: every Agent's eligible rows are reportable and ``agent_id`` only ranks the reporting @@ -126,44 +78,64 @@ def build_project_progress_snapshot_from_state( rollout_events=rollout_events, available_capabilities=available_capabilities, ) - agent_summary = parsed.get("agent_todos") - items = agent_summary.get("items") if isinstance(agent_summary, Mapping) else [] - stage_time = _stage_timestamp(completed_at) - if stage_time is None: - raise ValueError("periodic-report stage completion timestamp is invalid") - - def produced_by(item: Mapping[str, Any]) -> str: - """Return the Agent whose lane produced this row, or empty for none.""" - - return str(item.get("claimed_by") or "").strip() + return build_project_progress_snapshot_from_fields( + fields=parsed, + goal_id=goal_id, + agent_id=agent_id, + completed_at=completed_at, + publication_cursor=publication_cursor, + goal_cursors=goal_cursors, + ) - def not_after_stage(item: Mapping[str, Any]) -> bool: - raw = str(item.get("updated_at") or item.get("completed_at") or "").strip() - if not raw: - return True - parsed = _stage_timestamp(raw) - if parsed is None: - return False - return parsed <= stage_time - done = [ - dict(item) - for item in items or [] - if isinstance(item, Mapping) - and item.get("status") == "done" - and produced_by(item) - and not_after_stage(item) - and str(item.get("action_kind") or "") not in _META_ACTION_KINDS - ] - # Tier order must not disturb recency order inside a tier, so the lane sort - # runs last over the already-newest-first list. - done.sort(key=lambda item: str(item.get("updated_at") or ""), reverse=True) - done.sort(key=lambda item: produced_by(item) != agent_id) +def build_project_progress_snapshot_from_fields( + *, + fields: Mapping[str, Any], + goal_id: str, + agent_id: str, + completed_at: str, + publication_cursor: Mapping[str, Any] | None = None, + goal_cursors: Sequence[Mapping[str, Any]] | None = None, +) -> dict[str, Any] | None: + """Select facts from one complete evaluated snapshot, then apply publication history.""" + items = list((fields.get("agent_todos") or {}).get("items") or []) + selection_fields = ( + "todo_id", + "status", + "claimed_by", + "updated_at", + "completed_at", + "action_kind", + "task_class", + ) + try: + selected = effect_runtime_result( + "capabilities.periodic_report.progress.select", + { + "schema_version": "periodic_report_progress_selection_request_v0", + "agent_id": agent_id, + "completed_at": completed_at, + "items": [ + { + **{key: item.get(key) for key in selection_fields}, + "actionable": todo_item_is_actionable_open(item), + } + for item in items + ], + }, + ) + except EffectRuntimeRejected as error: + raise ValueError(str(error)) from error + if ( + not isinstance(selected, dict) + or selected.get("schema_version") + != "periodic_report_progress_selection_result_v0" + ): + raise ValueError("periodic-report progress selection result mismatch") progress_items: list[dict[str, Any]] = [] - for index, item in enumerate(done): - outcome_completed_at = _outcome_completed_at(item, stage_time=stage_time) - if outcome_completed_at is None: - continue + for outcome in selected["outcomes"]: + index = outcome["rank"] + item = items[outcome["index"]] summary = " ".join( str( item.get("evidence") or item.get("note") or item.get("text") or "" @@ -177,35 +149,20 @@ def not_after_stage(item: Mapping[str, Any]) -> bool: "summary": summary[:360] or "Validated completion is durably recorded.", "content_kind": "outcome", "value_rank": 10 + index, - "source_ref": f"todo:{item.get('todo_id')}", - "completed_at": outcome_completed_at, + "source_ref": f"todo:{item['todo_id']}", + "completed_at": outcome["completed_at"], } ) - open_items = [ - dict(item) - for item in items or [] - if isinstance(item, Mapping) - and todo_item_is_actionable_open(dict(item)) - and produced_by(item) - and not_after_stage(item) - and item.get("task_class") != "continuous_monitor" - and item.get("action_kind") - not in { - "consume_periodic_report_intent", - "repair_periodic_report_intent_consumption", - } - ] - open_items.sort(key=lambda item: produced_by(item) != agent_id) - if open_items: - next_item = open_items[0] + if selected["next_index"] is not None: + item = items[selected["next_index"]] progress_items.append( { "item_id": "next_action", "title": "Next action", - "summary": " ".join(str(next_item.get("text") or "").split())[:360], + "summary": " ".join(str(item.get("text") or "").split())[:360], "content_kind": "next_action", "value_rank": 90, - "source_ref": f"todo:{next_item.get('todo_id')}", + "source_ref": f"todo:{item['todo_id']}", } ) if not progress_items: diff --git a/loopx/capabilities/periodic_report/todo_source.py b/loopx/capabilities/periodic_report/todo_source.py new file mode 100644 index 0000000000..a938fa5089 --- /dev/null +++ b/loopx/capabilities/periodic_report/todo_source.py @@ -0,0 +1,70 @@ +"""One read-only Todo source for report staging and approval retry. + +Only the absence of a promotion fence permits Markdown parsing. A caller may +reuse the returned fields for frontier and fact selection without mixing heads. +""" + +from __future__ import annotations + +from collections.abc import Mapping +from pathlib import Path +from typing import Any + +from ...control_plane.coordination.local_authority import ( + canonical_todo_summary_fields, + read_canonical_todos_if_promoted, +) +from ...control_plane.todos.active_state_todo_parser import parse_active_state_todos +from ...paths import resolve_runtime_root +from ...registry import find_registry_goal, read_json, resolve_state_file + + +def read_report_todo_source( + *, + registry_path: Path, + goal_id: str, + runtime_root: Path | None = None, + state_path: Path | None = None, + rollout_events: list[dict[str, Any]] | None = None, + available_capabilities: Any = None, +) -> tuple[dict[str, Any], list[dict[str, Any]]]: + """Return full evaluated summaries and retained User decision records. + + Archived decisions can still supersede an approval-pending receipt; they + must not disappear just because a display stopped showing them. + """ + registry = read_json(registry_path) + goal = find_registry_goal(registry, goal_id) + if not isinstance(goal, Mapping): + raise ValueError("periodic-report Goal is not registered") + runtime_root = runtime_root or resolve_runtime_root( + registry, None, registry_path=registry_path + ) + canonical = read_canonical_todos_if_promoted( + runtime_root=runtime_root, goal_id=goal_id + ) + if canonical is not None: + fields = canonical_todo_summary_fields( + canonical["todos"], + rollout_events=rollout_events, + available_capabilities=available_capabilities, + goal_acceptance_contract=canonical.get("goal_acceptance_contract"), + goal_acceptance_work_guards=canonical.get("goal_acceptance_work_guards"), + ) + return fields, [row for row in canonical["todos"] if row.get("role") == "user"] + state_path = state_path or resolve_state_file( + Path(str(goal.get("repo") or "")).expanduser(), + str(goal.get("state_file") or ""), + ) + if state_path is None: + raise ValueError("periodic-report active state is unavailable") + fields = parse_active_state_todos( + state_path.read_text(encoding="utf-8"), + goal=dict(goal), + state_path=state_path, + item_limit=None, + rollout_events=rollout_events, + available_capabilities=available_capabilities, + ) + # Legacy decision selection retains its existing active-section boundary. + return fields, list((fields.get("user_todos") or {}).get("items") or []) diff --git a/loopx/control_plane/capabilities/periodic_report_progress.ts b/loopx/control_plane/capabilities/periodic_report_progress.ts new file mode 100644 index 0000000000..48a5aa5e5d --- /dev/null +++ b/loopx/control_plane/capabilities/periodic_report_progress.ts @@ -0,0 +1,96 @@ +/** Report selection owns report policy; evaluated Todo facts retain their owner. + * Ordinals address this exact input snapshot and are never durable identities. */ +import {createHash} from 'node:crypto'; +import type {JsonObject} from '../effect_program.ts'; +import {EffectRuntimeRequestError} from '../effect_runtime_errors.ts'; +import {requireJsonObject, requireNonEmptyString, requireBoolean, requireStringLiteral} from '../runtime_decode.ts'; +import {parseTodoTimestampMicros} from '../runtime_timestamp.ts'; +import {authorityUnicodeCompare} from '../coordination/authority_store_codec.ts'; + +const META_ACTIONS = new Set(['consume_periodic_report_intent', 'repair_periodic_report_intent_consumption', 'repair_periodic_report_editorial']); +const text = (value: unknown): string => typeof value === 'string' ? value.trim() : ''; +function timestamp(value: unknown): bigint | null { + const raw = text(value); + // Report facts require an explicit offset. The shared Todo codec otherwise + // accepts naive timestamps as UTC for legacy compatibility. + if (!/[T ].*(?:Z|[+-]\d{2}(?::?\d{2})?(?::?\d{2}(?:[.,]\d+)?)?)$/u.test(raw)) return null; + return parseTodoTimestampMicros(raw); +} +function request(value: unknown, schema: string): JsonObject { + const input = requireJsonObject(value, 'periodic-report read request'); + if (input.schema_version !== schema || !Array.isArray(input.items)) { + throw new EffectRuntimeRequestError('periodic-report read request schema/items mismatch'); + } + const seen = new Set(); + for (const raw of input.items) { + const row = requireJsonObject(raw, 'periodic-report Todo'); + const id = requireNonEmptyString(row.todo_id, 'todo_id'); + if (seen.has(id)) throw new EffectRuntimeRequestError('periodic-report source contains duplicate Todo identities'); + seen.add(id); + } + return input; +} + +interface ProgressRow { + index: number; + status: 'open' | 'blocked' | 'done' | 'deferred'; + owner: string; + action: string; + taskClass: string; + actionable: boolean; + updatedAt: bigint | null; + completedAt: string; + completedTime: bigint | null; + observedTime: bigint | null; +} + +export function selectPeriodicReportProgress(value: unknown): JsonObject { + const input = request(value, 'periodic_report_progress_selection_request_v0'); + const agent = requireNonEmptyString(input.agent_id, 'agent_id'); + const stage = timestamp(input.completed_at); + if (stage === null) throw new EffectRuntimeRequestError('periodic-report stage completion timestamp is invalid'); + const rows: ProgressRow[] = (input.items as JsonObject[]).map((row, index) => { + const updated = text(row.updated_at), completed = text(row.completed_at) || updated; + return {index, status: requireStringLiteral(row.status, ['open','blocked','done','deferred'], 'Todo status'), + owner: text(row.claimed_by), action: text(row.action_kind), taskClass: text(row.task_class), + actionable: requireBoolean(row.actionable, 'evaluated Todo actionable'), + updatedAt: timestamp(updated), completedAt: completed, completedTime: timestamp(completed), + observedTime: updated || completed ? timestamp(updated || completed) : stage}; + }).filter(row => row.owner !== '' && row.observedTime !== null && row.observedTime <= stage); + const ownFirst = (a: ProgressRow, b: ProgressRow): number => Number(a.owner !== agent) - Number(b.owner !== agent); + const outcomes = rows.filter(row => row.status === 'done' && !META_ACTIONS.has(row.action)) + .sort((a, b) => { + const tier = ownFirst(a, b); + if (tier) return tier; + const left = a.updatedAt ?? a.completedTime ?? 0n; + const right = b.updatedAt ?? b.completedTime ?? 0n; + return left === right ? a.index - b.index : left > right ? -1 : 1; + }).map((row, rank) => ({row, rank})) + .filter(({row}) => row.completedTime !== null && row.completedTime <= stage) + .map(({row, rank}) => ({index: row.index, completed_at: row.completedAt, rank})); + const next = rows.filter(row => row.actionable && row.status === 'open' && + row.taskClass !== 'continuous_monitor' && !META_ACTIONS.has(row.action)) + .sort((a, b) => ownFirst(a, b) || a.index - b.index)[0]; + return {schema_version: 'periodic_report_progress_selection_result_v0', outcomes, next_index: next?.index ?? null}; +} + +export function selectPeriodicReportApprovalRetry(value: unknown): JsonObject { + const input = request(value, 'periodic_report_approval_retry_request_v0'); + const agent = requireNonEmptyString(input.agent_id, 'agent_id'); + const scope = requireNonEmptyString(input.approval_scope, 'approval_scope'); + const rows = (input.items as JsonObject[]).filter(row => { + const decisionScope = row.decision_scope; + const actualScope = typeof decisionScope === 'object' && decisionScope !== null && !Array.isArray(decisionScope) + ? ['kind', 'granularity', 'scope_key'].map(key => text((decisionScope as JsonObject)[key])).join(':') : text(decisionScope); + return row.status === 'done' && ['approve_periodic_report_payload', 'cancel_periodic_report_payload'].includes(text(row.action_kind)) && + ['reject', 'cancel'].includes(text(row.decision_outcome)) && actualScope === scope && + text(row.bound_agent || row.blocks_agent) === agent && timestamp(row.updated_at) !== null; + }).sort((a, b) => { + const left = timestamp(a.updated_at)!, right = timestamp(b.updated_at)!; + return left === right ? authorityUnicodeCompare(text(a.todo_id), text(b.todo_id)) : left > right ? -1 : 1; + }); + const row = rows[0]; + // Preserve the existing durable retry key; time parsing only selects a row. + const revision = row ? createHash('sha256').update(`${row.todo_id}:${row.updated_at}`).digest('hex').slice(0, 16) : null; + return {schema_version: 'periodic_report_approval_retry_result_v0', revision}; +} diff --git a/loopx/control_plane/coordination/local_authority.py b/loopx/control_plane/coordination/local_authority.py index c1b9c7b655..c68faf38a6 100644 --- a/loopx/control_plane/coordination/local_authority.py +++ b/loopx/control_plane/coordination/local_authority.py @@ -283,6 +283,7 @@ def canonical_todo_summary_fields( todos: list[dict[str, Any]], *, rollout_events: list[dict[str, Any]] | None = None, + available_capabilities: Any = None, goal_acceptance_contract: dict[str, Any] | None = None, goal_acceptance_work_guards: dict[str, Any] | None = None, ) -> dict[str, Any]: @@ -340,6 +341,7 @@ def canonical_todo_summary_fields( include_empty_source=True, resume_source_items=todos, rollout_events=rollout_events, + available_capabilities=available_capabilities, item_limit=None, ) if summary: diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 2c16d92cf2..f537f62c7c 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -1,3 +1,4 @@ +import {selectPeriodicReportProgress, selectPeriodicReportApprovalRetry} from "./capabilities/periodic_report_progress.ts"; import {inspectTaskLease} from "./work_items/task_lease_inspection.ts"; import {evaluateTodoPriority} from "./todos/priority.ts"; import {evaluateUserCompletion} from "./todos/user_completion.ts"; @@ -423,6 +424,8 @@ export function createEffectRuntimeHandlers( ["todo.public_update.plan", planPublicTodoUpdate], ["todo.standing_decision.project", evaluateStandingDecisionProjection], ["todo.summary_lanes.project", projectTodoSummaryLanes], + ["capabilities.periodic_report.progress.select", selectPeriodicReportProgress], + ["capabilities.periodic_report.approval_retry.select", selectPeriodicReportApprovalRetry], ["todo.succession.project", projectTodoSuccession], ["todo.succession.closure", projectTodoClosure], ["todo.work_counts.project", projectLegacyTodoWorkCounts], From d1dcf38dd889a90b60e1652e90020ea9646ce795 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 21 Sep 2026 18:58:07 +0800 Subject: [PATCH 2/2] test(periodic-report): qualify provider-backed reporting consumers Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- ...shared-goal-authority-state-provider-v0.md | 12 ++ .../typescript-control-plane-migration-v0.md | 12 ++ ...script-control-plane-migration-v0.zh-CN.md | 9 + loopx/capabilities/periodic_report/README.md | 39 ++++ .../test_periodic_report_pending_intent.py | 67 +++++++ ...periodic_report_intent_capability_chain.py | 56 ++++++ .../test_periodic_report_authority.py | 168 ++++++++++++++++++ .../authority_store_conformance.ts | 2 + .../periodic_report_conformance.ts | 46 +++++ .../periodic_report_progress.test.ts | 64 +++++++ 10 files changed, 475 insertions(+) create mode 100644 tests/control_plane/test_periodic_report_authority.py create mode 100644 tests/control_plane_ts/periodic_report_conformance.ts create mode 100644 tests/control_plane_ts/periodic_report_progress.test.ts 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 c3b88bb46a..b206041f94 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -3036,6 +3036,18 @@ source paths, authorize monitor writeback, or change provider/promotion holds. **D1 — qualify permanent projection delivery; may overlap T1/T2.** +Periodic-report staging, live editorial fallback and approval retry now share +one canonical-first Todo source. Frontier and progress selection reuse the same +complete evaluated snapshot; missing/stale display cannot invent or hide work. +`capabilities/periodic_report_progress.ts` owns report selection and rejection +retry ordering, retiring Python selection/sorting loops. Offset-aware instants +retain microseconds, canonical archived rejection records remain effective, and +explicit runtime-root applies to both intent and Todo IO. Frozen editorial +requests retain their original basis. See [operation and boundaries](../../../loopx/capabilities/periodic_report/README.md#todo-authority-and-report-retries). +This closes that T3/L5 consumer family, not D1 permanent display freshness, +D2 durability, D3 whole-Goal qualification or default-provider selection. The +conditional 5–8 remaining delivery-package estimate is unchanged. + Summary/work-lane counts now remain independent of display limits and retain incomplete-source knowledge through Agent scoping; canonical list acceptance holds match status. This closes one L5 read consumer, not permanent projection freshness or D1–D3. See [count semantics](../../reference/todo-work-counts.md). The Goal Channel ownership observation consumes one complete provider revision before bounding display. It never repairs Markdown or revives old local leases; provider failures and truncation stay visible. This is a T3 read closure with shared TS interpretation, not D1/D2 qualification or D3 cutover. See [coordination observation](../../reference/coordination-observation.md). diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index c1a49d4d5d..ae097032df 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -789,6 +789,18 @@ all T2 commands or authorize whole-Goal promotion. **T3 — close remaining structured consumers, then remove their old reads.** +Periodic-report staging, live editorial fallback and approval retry now share +one canonical-first Todo source. Frontier and progress selection reuse the same +complete evaluated snapshot; missing/stale display cannot invent or hide work. +`capabilities/periodic_report_progress.ts` owns report selection and rejection +retry ordering, retiring Python selection/sorting loops. Offset-aware instants +retain microseconds, canonical archived rejection records remain effective, and +explicit runtime-root applies to both intent and Todo IO. Frozen editorial +requests retain their original basis. See [operation and boundaries](../../../loopx/capabilities/periodic_report/README.md#todo-authority-and-report-retries). +This closes that T3/L5 consumer family, not D1 permanent display freshness, +D2 durability, D3 whole-Goal qualification or default-provider selection. The +conditional 5–8 remaining delivery-package estimate is unchanged. + Todo summary lanes and pre-limit work counts now share `todos/summary_lanes.ts`. Python's lane classification and hidden-work inference loops are removed; quota recomputes counts after scope selection and carries incomplete source knowledge 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 e697c68e15..02df208907 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 @@ -599,6 +599,15 @@ delivery pending;这不代表全部 T2 命令或整 Goal promotion 已完成 **T3 — 闭合剩余 structured consumer,删除各自旧读路径。** +Periodic-report 的阶段判断、实时编辑输入回退与审批重试现共用 canonical-first +Todo 来源;frontier 和报告事实复用同一完整已求值快照。展示缺失、过期或损坏不再 +隐藏/复活工作。`capabilities/periodic_report_progress.ts` 拥有报告选择及拒绝重试 +排序,删除 Python 对应循环;时间按带偏移的实际时刻比较并保留微秒,canonical +归档拒绝记录仍有效,显式 runtime-root 同时约束 intent 和 Todo IO。已冻结的编辑 +请求沿用原始依据,不因重试刷新。见[操作边界](../../../loopx/capabilities/periodic_report/README.md#todo-authority-and-report-retries)。 +这闭合一组 T3/L5 消费者,不代表 D1 永久展示新鲜度、D2 耐久性、D3 整 Goal +资格或默认 provider 已完成;条件性的 5–8 个后续完整交付批次估算保持不变。 + Todo 摘要 lane 与裁剪前工作计数现共用 `todos/summary_lanes.ts`,删除 Python 的 lane 分类和隐藏任务推断循环。quota 在作用域筛选后重新计数,不完整来源状态贯穿 压缩与重复投影;公开 canonical Todo 列表保留同版本 acceptance 限制。见 diff --git a/loopx/capabilities/periodic_report/README.md b/loopx/capabilities/periodic_report/README.md index 14c68817c3..b0e045cc02 100644 --- a/loopx/capabilities/periodic_report/README.md +++ b/loopx/capabilities/periodic_report/README.md @@ -95,6 +95,45 @@ delivery-Todo pipeline. It acknowledges the provider source only after `delivery_ready` durability; failed ACKs become settlement-only retries and do not duplicate delivery work. +## Todo authority and report retries + +Stage-boundary frontier evaluation and project-progress selection share one +complete evaluated Todo snapshot. After Goal promotion, File/SQLite authority +owns those facts even when the Markdown display is stale, absent or malformed. +An empty canonical graph is empty; a failed canonical read remains a failed +observation and never falls back to Markdown. Reading does not repair the display. +Before promotion, the existing Markdown source remains in use. + +The TS report selector consumes bounded decision fields, not report prose. It +keeps the reporter's outcomes before peer outcomes, excludes unowned work, and +uses the existing resume/acceptance evaluation for the next action. Report +consumption and all report-repair actions are excluded from next-action facts. +Publication-history filtering still runs before the six-outcome cap so already +published facts do not hide an unseen seventh outcome. + +Completion and retry ordering compares offset-aware instants at microsecond +precision, rather than timestamp strings. Invalid/naive completion timestamps +are not reportable; an invalid stage timestamp is an error. A retained canonical +User rejection/cancellation can supersede an approval-pending receipt after +archival. It must match the exact decision scope and addressed Agent; this is +permission to reconsider that pending attempt, not approval to publish. Existing +retry-key bytes remain unchanged for the same selected decision. + +Use the existing local command to consume pending work: + +```bash +loopx periodic-report consume-pending --goal-id --agent-id --execute --format json +``` + +An explicit global `--runtime-root` routes both intent state and Todo reads. +When the result is `editorial_required`, inspect its local editorial-request +artifact: facts identify their source Todos. An identical retry reuses the +frozen editorial input; it does not rebuild it from newer Todo state. Provider +outages leave the intent for retry after authority recovery. Keep existing +editorial artifacts and receipts when rolling back; do not delete them to force +re-generation. No provider default, subscription, scheduler, external-delivery +permission or publication cursor is changed by this read-path refactor. + ## Customize or schedule The capability remains **inactive for background work and external writes by diff --git a/tests/capabilities/test_periodic_report_pending_intent.py b/tests/capabilities/test_periodic_report_pending_intent.py index 515243372c..e8af9caf6b 100644 --- a/tests/capabilities/test_periodic_report_pending_intent.py +++ b/tests/capabilities/test_periodic_report_pending_intent.py @@ -1908,3 +1908,70 @@ def test_cross_agent_or_malformed_intent_fails_closed(tmp_path: Path) -> None: ) == [] ) + + +@pytest.mark.parametrize('provider', ['file', 'sqlite']) +def test_real_cli_editorial_fallback_reads_canonical_work_without_display( + tmp_path, monkeypatch, provider +): + """The public command freezes canonical facts and sends nothing externally.""" + import subprocess + import sys + from tests.control_plane.canonical_authority_fixture import ( + initialize_canonical_authority, + isolate_sqlite_runtime, + ) + from loopx.control_plane.coordination.runtime_shadow import ( + build_todo_runtime_shadow_projection, + ) + + if provider == "sqlite": + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry, runtime = _fixture(tmp_path) + state = registry.parent / "ACTIVE_GOAL_STATE.md" + records = parse_active_state_todos(state.read_text(), item_limit=None)[ + "agent_todos" + ]["items"] + projection = build_todo_runtime_shadow_projection( + goal_id=GOAL_ID, todos=records, handoff_mode="soft_claim" + ) + initialize_canonical_authority( + runtime, GOAL_ID, projection, state_path=state, provider=provider + ) + # A runtime override must route both intent IO and the Todo snapshot. + config = json.loads(registry.read_text()) + config["common_runtime_root"] = str(tmp_path / "unused-runtime") + registry.write_text(json.dumps(config)) + state.unlink() + command = [ + sys.executable, + "-m", + "loopx.cli", + "--registry", + str(registry), + "--runtime-root", + str(runtime), + "--format", + "json", + "periodic-report", + "consume-pending", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--execute", + ] + run = subprocess.run(command, capture_output=True, text=True, timeout=90) + assert run.returncode == 0, run.stderr + result = json.loads(run.stdout) + assert result["status"] == "editorial_required", result + frozen = Path(result["editorial_request_path"]) + before = frozen.read_bytes() + request = json.loads(before) + assert {fact["source_ref"] for fact in request["facts"]} == {"todo:todo_finished"} + assert not state.exists() + retry = subprocess.run(command, capture_output=True, text=True, timeout=90) + assert retry.returncode == 0, retry.stderr + assert frozen.read_bytes() == before, ( + "retry must retain the same authored-input basis" + ) diff --git a/tests/cli_commands/test_periodic_report_intent_capability_chain.py b/tests/cli_commands/test_periodic_report_intent_capability_chain.py index b56fda8540..ce8b04f5b9 100644 --- a/tests/cli_commands/test_periodic_report_intent_capability_chain.py +++ b/tests/cli_commands/test_periodic_report_intent_capability_chain.py @@ -124,3 +124,59 @@ def test_pending_intent_fallback_fails_closed_without_producer_evidence( _claim_gated_successors(registry.parent) assert f"todo:{GATED_TODO}" not in _editorial_fact_sources(registry, runtime) assert f"todo:{PLAIN_TODO}" in _editorial_fact_sources(registry, runtime) + + +def test_post_writeback_frontier_and_progress_share_one_canonical_snapshot( + tmp_path, monkeypatch +): + from tests.control_plane.canonical_authority_fixture import ( + initialize_canonical_authority, + ) + from loopx.control_plane.coordination.runtime_shadow import ( + build_todo_runtime_shadow_projection, + ) + from loopx.control_plane.todos.active_state_todo_parser import ( + parse_active_state_todos, + ) + from loopx.capabilities.periodic_report.post_writeback_hook import ( + build_periodic_report_post_writeback_projection, + ) + from loopx.capabilities.periodic_report import todo_source + + _captured, registry, runtime = complete_todo_via_cli( + tmp_path, + journal_capabilities=["network"], + write_state=_write_unclaimed_frontier_state, + ) + _claim_gated_successors(registry.parent) + state = registry.parent / "goal.md" + records = parse_active_state_todos(state.read_text(), item_limit=None)[ + "agent_todos" + ]["items"] + source = build_todo_runtime_shadow_projection( + goal_id=GOAL_ID, todos=records, handoff_mode="soft_claim" + ) + initialize_canonical_authority(runtime, GOAL_ID, source, state_path=state) + state.unlink() + reads = [] + original = todo_source.read_canonical_todos_if_promoted + + def observe(**kwargs): + result = original(**kwargs) + reads.append(result["provider_revision"]) + return result + + monkeypatch.setattr(todo_source, "read_canonical_todos_if_promoted", observe) + result = build_periodic_report_post_writeback_projection( + payload={"available_capabilities": ["network"]}, + registry_path=registry, + runtime_root=runtime, + goal_id=GOAL_ID, + agent_id=AGENT_ID, + ) + assert len(reads) == 1 + assert result["stage_completion"]["acceptance"] == "validated" + assert f"todo:{GATED_TODO}" in { + r["source_ref"] for r in result["project_progress"]["items"] + } + assert not state.exists() diff --git a/tests/control_plane/test_periodic_report_authority.py b/tests/control_plane/test_periodic_report_authority.py new file mode 100644 index 0000000000..845263f135 --- /dev/null +++ b/tests/control_plane/test_periodic_report_authority.py @@ -0,0 +1,168 @@ +"""The reporting consumer never treats a stale display as Todo authority.""" + +import json + +import pytest + +from canonical_authority_fixture import ( + initialize_canonical_authority, + isolate_sqlite_runtime, +) +from loopx.capabilities.periodic_report.project_progress_snapshot import ( + build_project_progress_snapshot, +) +from loopx.control_plane.coordination.runtime_shadow import ( + build_todo_runtime_shadow_projection, +) +from loopx.control_plane.todos.active_state_todo_parser import parse_active_state_todos + +GOAL = "report-authority" +STAGE = "2026-09-21T12:00:00Z" + + +def workspace(tmp_path, monkeypatch, provider, *, empty=False, decision=False): + if provider == "sqlite": + isolate_sqlite_runtime(tmp_path, monkeypatch) + state = tmp_path / "state.md" + text = ( + "# Report fixture\n\n## Agent Todo\n" + "- [x] Canonical finished work.\n" + " \n" + ) + state.write_text(text) + records = parse_active_state_todos(text, item_limit=None)["agent_todos"]["items"] + if decision: + user_text = ( + "## User Todo\n- [x] Reject the report payload.\n" + " \n" + ) + user = parse_active_state_todos(user_text, item_limit=None)["user_todos"][ + "items" + ][0] + user.update(archive_state="archive", source_section="Completed Work Archive") + records.append(user) + projection = build_todo_runtime_shadow_projection( + goal_id=GOAL, todos=[] if empty else records, handoff_mode="soft_claim" + ) + runtime = tmp_path / "runtime" + initialize_canonical_authority( + runtime, GOAL, projection, state_path=state, provider=provider + ) + registry = tmp_path / "registry.json" + registry.write_text( + json.dumps( + { + "common_runtime_root": str(runtime), + "goals": [ + {"id": GOAL, "repo": str(tmp_path), "state_file": str(state)} + ], + } + ) + ) + return registry, state, runtime + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +@pytest.mark.parametrize("display", ["missing", "stale", "malformed"]) +def test_snapshot_uses_authority_without_repairing_display( + tmp_path, monkeypatch, provider, display +): + registry, state, _ = workspace(tmp_path, monkeypatch, provider) + if display == "missing": + state.unlink() + elif display == "malformed": + state.write_bytes(b"\xff\xfe") + else: + state.write_text("## Agent Todo\n- [ ] Stale display work.\n") + before = state.read_bytes() if state.exists() else None + snapshot = build_project_progress_snapshot( + registry_path=registry, goal_id=GOAL, agent_id="reporter", completed_at=STAGE + ) + assert [row["source_ref"] for row in snapshot["items"]] == ["todo:todo_real"] + assert (state.read_bytes() if state.exists() else None) == before + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_empty_authority_does_not_resurrect_markdown_outcome( + tmp_path, monkeypatch, provider +): + registry, _, _ = workspace(tmp_path, monkeypatch, provider, empty=True) + assert ( + build_project_progress_snapshot( + registry_path=registry, + goal_id=GOAL, + agent_id="reporter", + completed_at=STAGE, + ) + is None + ) + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_unavailable_authority_is_not_empty_or_legacy(tmp_path, monkeypatch, provider): + from loopx.capabilities.periodic_report import todo_source + from loopx.control_plane.coordination.local_authority import ( + LocalCoordinationAuthorityUnavailable, + ) + + registry, _, _ = workspace(tmp_path, monkeypatch, provider) + + def unavailable(**kwargs): + raise LocalCoordinationAuthorityUnavailable( + "unavailable", code="test_outage", payload={} + ) + + monkeypatch.setattr(todo_source, "read_canonical_todos_if_promoted", unavailable) + with pytest.raises(LocalCoordinationAuthorityUnavailable): + build_project_progress_snapshot( + registry_path=registry, + goal_id=GOAL, + agent_id="reporter", + completed_at=STAGE, + ) + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_approval_retry_reads_retained_canonical_decision_and_explicit_runtime( + tmp_path, monkeypatch, provider +): + from loopx.capabilities.periodic_report.pending_intent import ( + _superseding_approval_revision, + ) + from loopx.control_plane.coordination.local_authority import ( + read_canonical_todos_if_promoted, + ) + import hashlib + + registry, state, runtime = workspace(tmp_path, monkeypatch, provider, decision=True) + config = json.loads(registry.read_text()) + todo_id = "todo_decision" + row = next( + t + for t in read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL)[ + "todos" + ] + if t["todo_id"] == todo_id + ) + config["common_runtime_root"] = str(tmp_path / "wrong-runtime") + registry.write_text(json.dumps(config)) + state.unlink() + expected = hashlib.sha256(f"{todo_id}:{row['updated_at']}".encode()).hexdigest()[ + :16 + ] + assert ( + _superseding_approval_revision( + registry_path=registry, + runtime_root=runtime, + goal_id=GOAL, + agent_id="reporter", + receipt={ + "status": "approval_pending", + "approval_scope": "other:action:report-1", + }, + ) + == expected + ) diff --git a/tests/control_plane_ts/authority_store_conformance.ts b/tests/control_plane_ts/authority_store_conformance.ts index 77d097b325..766fd6ac60 100644 --- a/tests/control_plane_ts/authority_store_conformance.ts +++ b/tests/control_plane_ts/authority_store_conformance.ts @@ -1,3 +1,4 @@ +import {registerPeriodicReportConformance} from "./periodic_report_conformance.ts"; import {registerTodoConsumerScopeConformance} from "./todo_consumer_scope_conformance.ts"; import {registerProjectionConfirmationConformance} from "./projection_confirmation_conformance.ts"; import {registerUserCompletionFollowthroughConformance} from "./user_completion_followthrough_conformance.ts"; @@ -244,6 +245,7 @@ export function registerAuthorityStoreConformance( factory: AuthorityStoreConformanceFactory, ): void { registerProjectionConfirmationConformance(providerName, factory); + registerPeriodicReportConformance(providerName, factory); registerLeaseLifecycleConformance(providerName, factory); registerClaimTransferConformance(providerName, factory); registerLeaseAcquisitionConformance(providerName, factory); diff --git a/tests/control_plane_ts/periodic_report_conformance.ts b/tests/control_plane_ts/periodic_report_conformance.ts new file mode 100644 index 0000000000..cf1ae6f4f5 --- /dev/null +++ b/tests/control_plane_ts/periodic_report_conformance.ts @@ -0,0 +1,46 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; +import type {AuthorityStoreConformanceFactory} from './authority_store_conformance.ts'; +import {productionScaleCoordinationFixture} from './production_scale_coordination_fixture.ts'; +import {coordinationTodoReadModel} from '../../loopx/control_plane/coordination/coordination_projection.ts'; +import {listLocalCoordinationTodos} from '../../loopx/control_plane/coordination/local_authority_read.ts'; +import {LOCAL_COORDINATION_TODO_LIST_REQUEST_SCHEMA} from '../../loopx/control_plane/coordination/coordination_state_contract.generated.ts'; +import {selectPeriodicReportProgress,selectPeriodicReportApprovalRetry} from '../../loopx/control_plane/capabilities/periodic_report_progress.ts'; +import type {JsonObject} from '../../loopx/control_plane/effect_program.ts'; + +export function registerPeriodicReportConformance(name: string, factory: AuthorityStoreConformanceFactory): void { + for (const shape of ['native','legacy'] as const) test(`${name}: periodic report full source (${shape})`, async context => { + const {store} = await factory(context); + const goal = 'report-source', fixture = productionScaleCoordinationFixture(goal,shape); + const todos = fixture.projection.todos as JsonObject[]; + const own = todos.find(row => row.role === 'agent' && row.status === 'done')!; + const peer = todos.find(row => row.role === 'agent' && row.status === 'done' && row !== own)!; + Object.assign(own,{claimed_by:'reporter',updated_at:'2026-09-21T18:00:00+08:00',completed_at:'2026-09-21T18:00:00+08:00'}); + Object.assign(peer,{claimed_by:'peer',updated_at:'2026-09-21T11:00:00.000001Z',completed_at:'2026-09-21T11:00:00.000001Z'}); + const decision = todos.find(row => row.role === 'user' && row.status === 'done')!; + Object.assign(decision,{action_kind:'approve_periodic_report_payload',decision_outcome:'reject', + bound_agent:'reporter',decision_scope:{kind:'other',granularity:'action',scope_key:'report-1'}, + updated_at:'2026-09-21T11:00:00Z',archive_state:'archive', + ...(shape === 'legacy' ? {source_section:'Completed Work Archive'} : {})}); + fixture.projection.todo_read_model = coordinationTodoReadModel(todos,(fixture.projection.todo_read_model as JsonObject).schema_version); + const seeded = await store.commitAuthority({operation_id:'report-source',expected_provider_revision:null, + next_projection:fixture.projection,events:[],receipts:[]}); + assert.equal(seeded.status,'applied'); + const before = await store.loadAuthority(); + const source = await listLocalCoordinationTodos({schema_version:LOCAL_COORDINATION_TODO_LIST_REQUEST_SCHEMA, + goal_id:goal,runtime_root:'/synthetic-runtime'}, {createStore:()=>store}); + assert.equal(source.status,'loaded'); + const records = source.todos as JsonObject[]; + assert.equal(records.length,fixture.expected_initial_todo_count); + const agents = records.filter(row => row.role === 'agent' && row.archive_state === 'active'); + const selected = selectPeriodicReportProgress({schema_version:'periodic_report_progress_selection_request_v0', + agent_id:'reporter',completed_at:'2026-09-21T12:00:00Z',items:agents.map(row=>({...row,actionable:false}))}); + const outcomeIds = (selected.outcomes as JsonObject[]).map(item=>agents[Number(item.index)].todo_id); + assert.equal(outcomeIds[0],own.todo_id); + assert.ok(outcomeIds.includes(peer.todo_id)); + const retry = selectPeriodicReportApprovalRetry({schema_version:'periodic_report_approval_retry_request_v0', + agent_id:'reporter',approval_scope:'other:action:report-1',items:records.filter(row=>row.role==='user')}); + assert.match(String(retry.revision),/^[a-f0-9]{16}$/); + assert.deepEqual(await store.loadAuthority(),before,'report reads never mutate provider state'); + }); +} diff --git a/tests/control_plane_ts/periodic_report_progress.test.ts b/tests/control_plane_ts/periodic_report_progress.test.ts new file mode 100644 index 0000000000..610bc9981c --- /dev/null +++ b/tests/control_plane_ts/periodic_report_progress.test.ts @@ -0,0 +1,64 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; +import {createHash} from 'node:crypto'; +import {selectPeriodicReportProgress, selectPeriodicReportApprovalRetry} from '../../loopx/control_plane/capabilities/periodic_report_progress.ts'; +import type {JsonObject} from '../../loopx/control_plane/effect_program.ts'; + +const row = (id: string, extra: JsonObject = {}): JsonObject => ({todo_id:id, claimed_by:'writer', + status:'done', updated_at:'2026-09-21T10:00:00Z', task_class:'advancement_task', actionable:false, ...extra}); +const progress = (items: JsonObject[], completed_at='2026-09-21T12:00:00Z') => selectPeriodicReportProgress({ + schema_version:'periodic_report_progress_selection_request_v0', agent_id:'writer', completed_at, items}); + +test('report selection orders actual instants and preserves sub-millisecond ordering', () => { + const result = progress([row('a',{updated_at:'2026-09-21T18:00:00+08:00'}), + row('b',{updated_at:'2026-09-21T10:00:00.000002Z'}), row('c',{updated_at:'2026-09-21T10:00:00.000001Z'})]); + assert.deepEqual((result.outcomes as JsonObject[]).map(r => r.index), [1,2,0]); +}); + +test('report stage is inclusive at the exact microsecond, not a millisecond bucket', () => { + const result = progress([row('a',{updated_at:'2026-09-21T10:00:00.000001Z'}), + row('b',{updated_at:'2026-09-21T10:00:00.000002Z'})], '2026-09-21T10:00:00.000001Z'); + assert.deepEqual((result.outcomes as JsonObject[]).map(r => r.index), [0]); +}); + +test('report selection retains peer progress and prioritizes reporter, never unowned work', () => { + const result = progress([row('peer',{claimed_by:'peer', updated_at:'2026-09-21T11:00:00Z'}), row('own'), row('unowned',{claimed_by:null})]); + assert.deepEqual((result.outcomes as JsonObject[]).map(r => r.index), [1,0]); +}); + +test('report next action obeys evaluated authority and excludes all report repair tasks', () => { + const result = progress([row('held',{status:'open',actionable:false}), + row('monitor',{status:'open',actionable:true,task_class:'continuous_monitor'}), + row('editorial',{status:'open',actionable:true,action_kind:'repair_periodic_report_editorial'}), + row('next',{status:'open',actionable:true})]); + assert.equal(result.next_index,3); +}); + +for (const value of ['2026-09-21T10:00:00','2026-02-30T10:00:00Z','bad']) { + test(`report rejects untrustworthy stage ${value}`, () => assert.throws(() => progress([],value), /timestamp/)); + test(`report excludes untrustworthy completion ${value}`, () => assert.deepEqual(progress([row('bad',{completed_at:value})]).outcomes, [])); +} + +test('report rejects duplicate identities before selecting facts', () => assert.throws(() => progress([row('a'),row('a')]), /duplicate/)); + +const scope = {kind:'other',granularity:'action',scope_key:'report-1'}; +const decision = (id: string, extra: JsonObject = {}) => row(id,{action_kind:'approve_periodic_report_payload', + decision_outcome:'reject',decision_scope:scope,bound_agent:'writer',...extra}); +const approval = (items: JsonObject[]) => selectPeriodicReportApprovalRetry({ + schema_version:'periodic_report_approval_retry_request_v0',agent_id:'writer',approval_scope:'other:action:report-1',items}); + +test('approval retry reads the latest real instant and preserves the durable key bytes', () => { + const result = approval([decision('old',{updated_at:'2026-09-21T18:00:00+08:00'}), + decision('new',{updated_at:'2026-09-21T11:00:00Z',archive_state:'archive'})]); + assert.equal(result.revision,createHash('sha256').update('new:2026-09-21T11:00:00Z').digest('hex').slice(0,16)); +}); + +test('approval retry does not authorize other agents, scopes, approvals or malformed time', () => { + assert.equal(approval([decision('foreign',{bound_agent:'peer'}),decision('scope',{decision_scope:'other'}), + decision('approve',{decision_outcome:'approve'}),decision('invalid',{updated_at:'invalid'})]).revision,null); +}); + +test('report selection rejects malformed evaluated state rather than coercing it', () => { + assert.throws(() => progress([row('a',{actionable:'true'})]),/boolean/); + assert.throws(() => progress([row('a',{status:'finished'})]),/status/); +});