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 db04967bd4..9bd00f5637 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -2889,6 +2889,8 @@ source paths, authorize monitor writeback, or change provider/promotion holds. **D1 — qualify permanent projection delivery; may overlap T1/T2.** +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). + The D1 document-ownership slice gives readers, editors and projection one visible-region and Todo-block boundary. It fixes fenced examples becoming real tasks, narrative after an archive end marker entering history, and sparse imported ordinals or archived 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 9fd0058ef9..7dfb761310 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 @@ -2286,6 +2286,8 @@ Scoped fallback 的选择与门禁关系也已复用同一 TS decision owner, **D1 — 资格化永久投影交付,可与 T1/T2 重叠推进。** +Goal Channel 所有权观察先读取完整 provider revision,再限制展示;不修复 Markdown、不复活旧本地 lease,明确披露失败与截断。这是共用 TS 解释规则的 T3 读链路闭合,不完成 D1/D2 或 D3 切换,见 [coordination observation](../../reference/coordination-observation.md)。 + D1 的文档归属切片把读取、编辑与投影放到同一可见区域/Todo 行解码边界,修复 fenced 示例被当成真实任务、归档 end marker 后叙述进入历史、稀疏历史行号及归档 优先级阻塞读回的问题。投影复用普通状态的耐久原子写入;相同字节的重试仍完成 diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index a5fadfbdf7..795df45452 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -631,6 +631,8 @@ delivery. This does not finish all T2 commands or authorize whole-Goal promotion **T3 — close remaining structured consumers, then remove their old reads.** +Goal Channel ownership observation now reads a complete canonical Todo/lease revision and shares one TS batch policy with the legacy adapter. It retires display-layer lease time/generation/conflict decisions and local-file reads after promotion. Empty, unavailable and truncated observations remain distinct; see [coordination observation](../../reference/coordination-observation.md). This closes the Goal Channel ownership reader, not other channel panels or whole-Goal promotion. + The D1 document-ownership slice gives readers, editors and projection one visible-region and Todo-block boundary. It fixes fenced examples becoming real tasks, narrative after an archive end marker entering history, and sparse imported ordinals or archived @@ -664,7 +666,7 @@ derived inside acquire from the supplied owner/claim/exclusion/registration fact not from the old caller-provided `effective` hint. Other-Todo overlap facts still come from the existing complete execution snapshot; release retains its separate key/version cleanup fence. This closes one T3 reader and shared rule boundary, -not the remaining Goal-channel lease display, T1/T2 transactions or promotion. +not the remaining T1/T2 transactions or promotion. Goal Channel ownership display closes in the separate ownership-observation slice. Capability resolution now shares `agents/capability_gate.ts`: missing prerequisites, repair outputs, owner/agent resolution and blocked-Todo bindings have one typed 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 0ce2a496eb..cbe2a8d7e7 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 @@ -491,6 +491,8 @@ promotion 已完成。 **T3 — 闭合剩余 structured consumer,删除各自旧读路径。** +Goal Channel 所有权观察现从完整 canonical Todo/lease revision 读取,并与 legacy adapter 共用 TS 批量规则;删除展示层的时间/代数/冲突判断和晋升后的本地文件读路径。空值、不可用与截断分别披露,见 [coordination observation](../../reference/coordination-observation.md)。这只闭合所有权观察 reader,不宣称其余面板或整 Goal 晋升完成。 + D1 的文档归属切片把读取、编辑与投影放到同一可见区域/Todo 行解码边界,修复 fenced 示例被当成真实任务、归档 end marker 后叙述进入历史、稀疏历史行号及归档 优先级阻塞读回的问题。投影复用普通状态的耐久原子写入;相同字节的重试仍完成 @@ -515,8 +517,8 @@ handoff mode;canonical 无租约不复活本地旧文件,provider 失败不 供 acquire、lifecycle 与终态 fence 复用。当前租约是否有效由 acquire 内部根据同一输入 的 owner/claim/exclusion/注册事实推导,不再由旧 `effective` 派生提示覆盖。 其他 Todo 的 scope 冲突仍消费现有完整执行快照;release 保留独立的 key/version -清理门禁。这是一个 T3 reader 与共享规则边界的闭合,不代表 Goal-channel lease -展示、T1/T2 全部事务或 promotion 已完成。 +清理门禁。这是一个 T3 reader 与共享规则边界的闭合,不代表 T1/T2 全部事务或 +promotion 已完成;Goal Channel 所有权展示由独立的 observation 切片闭合。 Quota 的 scope/claim 消费者现通过每个 source 一次 `todo.quota_planning.project`, 组合选择、有限展示与既有 resume planner。`quota_selection.ts` 替代 Python diff --git a/docs/assets/coordination-observation-before.png b/docs/assets/coordination-observation-before.png new file mode 100644 index 0000000000..c9cd7e6bd9 Binary files /dev/null and b/docs/assets/coordination-observation-before.png differ diff --git a/docs/assets/coordination-observation-unavailable.png b/docs/assets/coordination-observation-unavailable.png new file mode 100644 index 0000000000..b4c623aa2b Binary files /dev/null and b/docs/assets/coordination-observation-unavailable.png differ diff --git a/docs/reference/coordination-observation.md b/docs/reference/coordination-observation.md new file mode 100644 index 0000000000..6366e9ca35 --- /dev/null +++ b/docs/reference/coordination-observation.md @@ -0,0 +1,87 @@ +# Goal Channel coordination observation + +Goal Channel's `active_leases` is a read-only display of Todo claims and +time-active lease records. A displayed lease is not an execution grant: actual +mutation still checks the owning operation's eligibility, mode and lease fence. + +```bash +loopx status --goal-id example-goal --format json +``` + +Inspect `goal_channel_projection.active_leases` and its `source_warnings` in the +status/attention item. No new capability activation or provider selection is +required. Existing Goal Channel HTML renders these rows and warnings; this read +never sends a message or changes channel bindings. + +Before canonical promotion, the source remains the supplied status Todo view +plus local lease files. Python adapts file records; one TS request evaluates +expiry, lease generation and claim conflicts for the whole batch. The legacy +Todo view can be incomplete and does not acquire canonical completeness by using +the shared rule. An explicitly supplied empty `active_leases=[]` does not cause +soft claims to be inferred; local hard-lease observation remains independent. + +After promotion, the selected provider supplies the complete Todo/lease snapshot +in one read. Stale or absent Markdown, legacy lease files and caller-supplied +claim/lease display overrides cannot replace canonical ownership facts. An empty +canonical result stays empty. Claims of completed/archived Todos do not re-enter +the active claim display. Canonical claim/lease conflict comparison happens +before output limits or text redaction. + +`coordination_observation` records the provider source/revision, observation time, +record counts, total observations, display limit and truncation. These fields +apply only to the ownership observation; other Goal Channel panels remain their +existing status/quota/history projections and need not share that revision. +The canonical display limit is 100 entries, after evaluating the complete source; +corrupt-lease and conflict diagnostics come first. A truncation warning means the +visible rows are not a complete work inventory. + +All leases use one observation time. Expiry equal to that time is expired. +Malformed active expiry, unknown schema, mismatched lease identity and invalid +generation produce `hard_lease_unreadable` / `corrupt_lease` observations instead +of silently disappearing or crashing the channel. A provider/protocol failure +produces an empty ownership list with `coordination_unavailable`; that empty list +is **not evidence of no ownership**. Raw errors, lease operation keys, write +scopes and arbitrary backend metadata are not copied into canonical display. +Existing channel text redaction still applies. + +The read performs no business mutation, receipt creation, Markdown repair, +promotion or fallback write. It does not change provider defaults or qualify a +PostgreSQL CLI selector, whole-Goal cutover or long-duration SQLite storage. +Disable/rollback follows the existing provider lifecycle; do not revive stale +local files to bypass an unavailable canonical source. + +## Recognize an unavailable source + +The same synthetic unavailable canonical source and stale local lease, rendered +by the pinned legacy implementation (left) and the provider-aware reader (right). +The corrected panel explicitly reports an unavailable observation rather than +presenting the local lease as current ownership. Source Warnings carries the +additional explanation. Desktop and mobile layouts use the existing renderer. + +| Legacy display | Provider-aware display | +| --- | --- | +| ![Stale local lease shown as ownership](../assets/coordination-observation-before.png) | ![Canonical ownership unavailable](../assets/coordination-observation-unavailable.png) | + +## 中文 + +Goal Channel 的 `active_leases` 展示认领与时间上有效的 lease 记录,不授予执行权; +实际修改仍受相应操作的资格、mode 和 lease 门禁约束。上面的 status 命令可读取 +现有 Goal Channel/attention 投影,不新增 activation,也不发送消息。 + +晋升前沿用传入的状态 Todo 视图与本地 lease 文件,由同一 TS 批量规则计算时间、 +代数和认领冲突;旧 Todo 视图仍可能不完整。明确传入空列表不再补出 soft claim, +本地 hard lease 仍独立观察。 + +晋升后从选定 provider 的完整同一 revision 读取所有权事实。陈旧/缺失 Markdown、 +旧 lease 文件或调用方的展示覆盖值不再成为事实来源;canonical 空值保持为空, +完成/归档 Todo 的认领不进入活动认领展示。先判断完整数据中的冲突,再脱敏与截断。 + +`coordination_observation` 披露来源、版本、观察时间、总数和截断情况,只描述 +所有权这一组数据,不声称整张 Goal Channel 的所有面板来自同一 revision。 +最多展示 100 条,异常/冲突优先;截断提示意味着不能把可见列表当作完整工作清单。 +所有 lease 使用同一个观察时间,到期时间恰好相等视为过期。损坏记录显示不可读提示; +provider 失败显示 `coordination_unavailable`,不会回退本地文件,也不把空列表说成无人负责。 +原始错误、操作密钥和任意后端字段不进入展示,现有文本脱敏继续生效。 + +这是一条只读链路,不修改权威状态、回执或 Markdown,不改变默认 provider,也不 +完成 SQLite 长程资格化、PostgreSQL CLI 选择入口或整 Goal 切换。 diff --git a/loopx/control_plane/coordination/local_authority_runtime.ts b/loopx/control_plane/coordination/local_authority_runtime.ts index 0e95a7a3d8..bc16684945 100644 --- a/loopx/control_plane/coordination/local_authority_runtime.ts +++ b/loopx/control_plane/coordination/local_authority_runtime.ts @@ -1,4 +1,5 @@ import {COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA} from "./todo_archive.ts"; +import {readCoordinationOwnership} from "./ownership_observation.ts"; import {executeTodoContinuation} from "./todo_continuation.ts"; import { withFileMutationLock } from "../effect_runtime_io.ts"; import { ShadowManagementError, requireShadowPrimaryWriteAllowed, shadowMaintenanceLockPath } from "./shadow_management.ts"; @@ -1207,3 +1208,22 @@ export async function continueLocalTodo(value: unknown): Promise { ...localAuthorityOpenFailure(error)}; } } + +/** Goal Channel observes a complete provider snapshot through one coarse read. */ +export async function observeLocalCoordinationOwnership(value: unknown): Promise { + let sourceAuthority = "canonical_unavailable"; + try { + const input = requireJsonObject(value, "local ownership observation"); + if (input.schema_version !== "loopx_local_ownership_observation_request_v0") throw new Error("ownership observation schema mismatch"); + const root = runtimeRoot(input.runtime_root); + const goalId = requireAuthorityStoreId(input.goal_id, "goal id"); + const store = await openLocalAuthorityStore(root, goalId); + sourceAuthority = sourceAuthorityFor(store); + return {...await readCoordinationOwnership(store, goalId, input.observed_at as string), + source_authority: sourceAuthority, decision_read_from_provider: true, legacy_fallback_used: false}; + } catch (error) { + return {schema_version: "loopx_ownership_observation_result_v0", status: "failed", + reason_code: "coordination_observation_unavailable", source_authority: sourceAuthority, + decision_read_from_provider: true, legacy_fallback_used: false, ...localAuthorityOpenFailure(error)}; + } +} diff --git a/loopx/control_plane/coordination/ownership_observation.ts b/loopx/control_plane/coordination/ownership_observation.ts new file mode 100644 index 0000000000..885d0943cb --- /dev/null +++ b/loopx/control_plane/coordination/ownership_observation.ts @@ -0,0 +1,95 @@ +/** Read-only ownership observations. These records describe claims/leases, never grant execution. */ +import type {JsonObject} from "../effect_program.ts"; +import type {AuthorityStore} from "./authority_store.ts"; +import {requireJsonObject} from "../runtime_decode.ts"; +import {parseIsoTimestamp} from "../runtime_timestamp.ts"; +import {requireAuthorityStoreId} from "./authority_store_codec.ts"; +import {indexCoordinationProjection, validateCoordinationTodoReadModel} from "./coordination_projection.ts"; +import {leaseEpoch, leaseIsActive, TASK_LEASE_SCHEMA_VERSION} from "../work_items/task_lease_acquire.ts"; + +export const OWNERSHIP_OBSERVATION_SCHEMA = "loopx_ownership_observation_request_v0"; +export const OWNERSHIP_OBSERVATION_RESULT = "loopx_ownership_observation_result_v0"; +export const CANONICAL_OWNERSHIP_DISPLAY_LIMIT = 100; + +type ObservationStatus = "soft_claim" | "hard_lease" | "hard_lease_unreadable"; +function objects(value: unknown, label: string): JsonObject[] { + if (!Array.isArray(value)) throw new Error(`${label} must be an array`); + return value.map(item => requireJsonObject(item, label)); +} +function text(value: unknown): string | null { + return typeof value === "string" && value.trim() ? value.trim() : null; +} +function note(todoId: string): JsonObject { + const status: ObservationStatus = "hard_lease_unreadable"; + return {todo_id: todoId, status, reason: "corrupt_lease"}; +} + +function displayEntry(item: JsonObject): JsonObject { + const entry: JsonObject = {}; + for (const key of ["todo_id", "owner_agent", "claimed_by", "lease_until", "expires_at", "status", "reason"]) { + const value = text(item[key]); + if (value) entry[key] = value; + } + for (const key of ["lease_version", "lease_epoch"]) { + if (typeof item[key] === "number" && Number.isSafeInteger(item[key])) entry[key] = item[key]; + } + return entry; +} + +/** Both source adapters share time/generation/conflict rules; explicit [] is authoritative. */ +export function projectOwnershipObservation(value: unknown): JsonObject { + const input = requireJsonObject(value, "ownership observation"); + if (input.schema_version !== OWNERSHIP_OBSERVATION_SCHEMA) throw new Error("ownership observation schema mismatch"); + const at = typeof input.observed_at === "string" ? parseIsoTimestamp(input.observed_at) : null; + if (at === null) throw new Error("observed_at must be a valid timestamp"); + const todos = objects(input.todos, "todos"); + const claims = new Map(todos.map(todo => [text(todo.todo_id), text(todo.claimed_by)])); + const explicit = input.explicit_entries == null ? null : objects(input.explicit_entries, "explicit_entries"); + const entries: JsonObject[] = explicit === null ? todos.filter(todo => text(todo.claimed_by)).map(todo => ({ + todo_id: todo.todo_id ?? null, owner_agent: todo.claimed_by!, status: "soft_claim" satisfies ObservationStatus, + })) : explicit.map(displayEntry).filter(entry => text(entry.todo_id) || text(entry.owner_agent) || text(entry.claimed_by)); + for (const row of objects(input.lease_rows, "lease_rows")) { + const todoId = text(row.todo_id); + if (!todoId) throw new Error("lease observation requires a Todo identity"); + if (row.unreadable === true) {entries.push(note(todoId)); continue;} + const lease = row.lease == null ? null : requireJsonObject(row.lease, "lease"); + if (lease === null) continue; + try { + if (lease.schema_version !== TASK_LEASE_SCHEMA_VERSION || lease.todo_id !== todoId) throw new Error("lease identity/schema mismatch"); + if (!leaseIsActive(lease, at)) continue; + const entry: JsonObject = {todo_id: todoId, status: "hard_lease" satisfies ObservationStatus, + lease_epoch: leaseEpoch(lease), expires_at: lease.expires_at!}; + const owner = text(lease.owner); + if (owner) entry.owner_agent = owner; + if (typeof lease.version === "number" && Number.isInteger(lease.version)) entry.lease_version = lease.version; + const claim = claims.get(todoId); + if (owner && claim && owner !== claim) {entry.reason = "owner_conflicts_with_claim"; entry.claimed_by = claim;} + entries.push(entry); + } catch { entries.push(note(todoId)); } + } + return {schema_version: OWNERSHIP_OBSERVATION_RESULT, status: "loaded", entries, + total_count: entries.length, observed_at: input.observed_at!}; +} + +/** One complete, validated revision; no display, local files, receipts or writes. */ +export async function readCoordinationOwnership(store: AuthorityStore, goalId: string, observedAt: string): Promise { + requireAuthorityStoreId(goalId, "goal id"); + const loaded = await store.loadAuthority(); + if (loaded.status !== "loaded") return {schema_version: OWNERSHIP_OBSERVATION_RESULT, ...loaded}; + const index = indexCoordinationProjection(loaded.head, goalId); + validateCoordinationTodoReadModel(loaded.head, goalId); + const todos = [...index.todos.values()].filter(todo => todo.archive_state === "active" && todo.done !== true); + const result = projectOwnershipObservation({schema_version: OWNERSHIP_OBSERVATION_SCHEMA, + observed_at: observedAt, todos, explicit_entries: todos.filter(todo => todo.role === "agent" && text(todo.claimed_by)) + .map(todo => ({todo_id: todo.todo_id, owner_agent: todo.claimed_by!, status: "soft_claim"})), + lease_rows: [...index.leases.entries()].sort(([a], [b]) => a < b ? -1 : a > b ? 1 : 0) + .map(([todo_id, lease]) => ({todo_id, lease})), + }); + // Evaluate the whole snapshot before bounding display; diagnostics are retained first. + const entries = result.entries as JsonObject[]; + const ordered = [...entries.filter(row => row.reason), ...entries.filter(row => !row.reason)]; + return {...result, entries: ordered.slice(0, CANONICAL_OWNERSHIP_DISPLAY_LIMIT), + truncated: entries.length > CANONICAL_OWNERSHIP_DISPLAY_LIMIT, display_limit: CANONICAL_OWNERSHIP_DISPLAY_LIMIT, + todo_count: index.todos.size, lease_count: index.leases.size, + provider_revision: loaded.provider_revision, cursor: loaded.cursor}; +} diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index d4fcbf13e0..0ff5a2276b 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -1,5 +1,7 @@ import {planHandoffMode} from "./coordination/handoff_mode_policy.ts"; import {setLocalHandoffMode} from "./coordination/handoff_mode_runtime.ts"; +import {projectOwnershipObservation} from "./coordination/ownership_observation.ts"; +import {observeLocalCoordinationOwnership} from "./coordination/local_authority_runtime.ts"; import {evaluateTaskLeaseOwnerEligibility} from "./work_items/task_lease_eligibility.ts"; import { evaluateSubagentContext, describeSubagentContext } from "./subagent_context.ts"; import { @@ -484,6 +486,8 @@ export function createEffectRuntimeHandlers( ["coordination.local_authority.todo_compatibility_edit", editLocalCoordinationTodo], ["coordination.local_authority.mutate", mutateLocalCoordinationAuthority], ["coordination.local_authority.todo_read", readLocalCoordinationTodo], + ["coordination.ownership_observation", projectOwnershipObservation], + ["coordination.local_authority.ownership_observation", observeLocalCoordinationOwnership], ["coordination.local_authority.todo_list", listLocalCoordinationTodos], [ "coordination.local_authority.legacy_writer_fence.engage", diff --git a/loopx/control_plane/goals/coordination_observation.py b/loopx/control_plane/goals/coordination_observation.py new file mode 100644 index 0000000000..76322575fb --- /dev/null +++ b/loopx/control_plane/goals/coordination_observation.py @@ -0,0 +1,60 @@ +"""Source adapters for the typed, read-only Goal Channel ownership observation.""" +from pathlib import Path +from typing import Any + +from ..coordination.local_authority import LOCAL_AUTHORITY_SOURCES, local_authority_is_promoted +from ..effect_runtime import effect_runtime_result +from ..runtime.time import now_local_iso + + +def observe_goal_coordination(*, runtime_root: Any, goal_id: str, + agent_todos: list[dict[str, Any]], + explicit_entries: list[dict[str, Any]] | None) -> dict[str, Any]: + canonical = False + try: + root = Path(runtime_root) if runtime_root is not None else None + if root is not None: + canonical = local_authority_is_promoted(runtime_root=root, goal_id=goal_id) + if canonical: + assert root is not None + result = effect_runtime_result('coordination.local_authority.ownership_observation', { + 'schema_version': 'loopx_local_ownership_observation_request_v0', + 'runtime_root': str(root.expanduser().resolve()), 'goal_id': goal_id, + 'observed_at': now_local_iso(), + }) + if (not isinstance(result, dict) or result.get('source_authority') not in LOCAL_AUTHORITY_SOURCES + or result.get('decision_read_from_provider') is not True or result.get('legacy_fallback_used') is not False + or not isinstance(result.get('provider_revision'), str)): + raise RuntimeError('canonical ownership observation unavailable') + else: + from ..work_items.local_lease_record import TaskLeaseError + from ..work_items import task_lease as task_lease_module + rows = [] + if root is not None: + for path in sorted(task_lease_module.task_lease_dir(runtime_root=root, goal_id=goal_id).glob('todo_*.json')): + try: + # task_lease re-exports this seam and existing callers/tests patch it there. + lease = task_lease_module.read_lease(path) # type: ignore[attr-defined] + except FileNotFoundError: + continue + except (TaskLeaseError, OSError): + rows.append({'todo_id': path.stem, 'unreadable': True}) + else: + if lease is not None: + rows.append({'todo_id': path.stem, 'lease': lease}) + if not rows and not agent_todos and not explicit_entries: + return {'status': 'loaded', 'entries': []} + result = effect_runtime_result('coordination.ownership_observation', { + 'schema_version': 'loopx_ownership_observation_request_v0', + 'todos': agent_todos, 'lease_rows': rows, 'explicit_entries': explicit_entries, + 'observed_at': now_local_iso(), + }) + if (not isinstance(result, dict) or result.get('schema_version') != 'loopx_ownership_observation_result_v0' + or result.get('status') != 'loaded' or not isinstance(result.get('entries'), list) + or any(not isinstance(row, dict) for row in result['entries'])): + raise RuntimeError('invalid ownership observation') + return result + except (OSError, RuntimeError, TypeError, ValueError): + # Observation failure is visible, but never activates a legacy fallback or leaks provider errors. + return {'status': 'unavailable', 'entries': [], 'legacy_fallback_used': False, + 'source_authority': 'canonical_unavailable' if canonical else 'unavailable'} diff --git a/loopx/control_plane/goals/goal_channel_projection.py b/loopx/control_plane/goals/goal_channel_projection.py index 203f24b5b9..4c303fea6e 100644 --- a/loopx/control_plane/goals/goal_channel_projection.py +++ b/loopx/control_plane/goals/goal_channel_projection.py @@ -255,137 +255,19 @@ def _open_gates( return gates -def _hard_lease_entries( - *, - runtime_root: Any, - goal_id: str, - agent_todos: Sequence[Mapping[str, Any]], -) -> list[dict[str, Any]]: - """Best-effort read of the on-disk hard task-lease store. - - Missing or empty lease directories stay silent because most goals never - use task leases. A corrupt or unreadable lease file degrades to a typed - note instead of failing the status projection. - """ - - if runtime_root is None: - return [] - from pathlib import Path - - from ..work_items.task_lease import ( - TaskLeaseError, - lease_epoch, - lease_is_active, - read_lease, - task_lease_dir, - ) - - try: - lease_dir = task_lease_dir( - runtime_root=Path(runtime_root), - goal_id=str(goal_id), - ) - lease_paths = sorted(lease_dir.glob("todo_*.json")) - except (TaskLeaseError, OSError, TypeError, ValueError): - return [] - if not lease_paths: - return [] - claim_by_todo: dict[str, str] = {} - for item in agent_todos: - todo_id = str(item.get("todo_id") or "").strip() - claimed_by = str(item.get("claimed_by") or "").strip() - if todo_id and claimed_by: - claim_by_todo[todo_id] = claimed_by - entries: list[dict[str, Any]] = [] - for path in lease_paths: - try: - lease = read_lease(path) - except FileNotFoundError: - # Lease released between directory listing and read: nothing left - # to surface, so skip instead of mislabeling it as corrupt. - continue - except (TaskLeaseError, OSError): - entries.append( - { - "todo_id": path.stem, - "status": "hard_lease_unreadable", - "reason": "corrupt_lease", - } - ) - continue - if lease is None or not lease_is_active(lease): - continue - todo_id = _text(lease.get("todo_id"), limit=120) or path.stem - entry: dict[str, Any] = {"todo_id": todo_id, "status": "hard_lease"} - owner = _text(lease.get("owner"), limit=120) - if owner: - entry["owner_agent"] = owner - version = lease.get("version") - if isinstance(version, int): - entry["lease_version"] = version - entry["lease_epoch"] = lease_epoch(lease) - expires_at = _text(lease.get("expires_at"), limit=80) - if expires_at: - entry["expires_at"] = expires_at - claimed_by = claim_by_todo.get(todo_id) - if owner and claimed_by and claimed_by != owner: - entry["reason"] = "owner_conflicts_with_claim" - entry["claimed_by"] = claimed_by - entries.append(entry) - return entries - - -def _active_leases( - *, - active_leases: Sequence[Mapping[str, Any]] | None, - agent_todos: Sequence[Mapping[str, Any]], - runtime_root: Any = None, - goal_id: str = "", -) -> list[dict[str, Any]]: - explicit = _as_mappings(active_leases) - if explicit: - source = explicit - else: - source = [ - { - "todo_id": item.get("todo_id"), - "owner_agent": item.get("claimed_by"), - "status": "soft_claim", - } - for item in agent_todos - if item.get("claimed_by") - ] - compact: list[dict[str, Any]] = [] - for item in source: - todo_id = _text(item.get("todo_id"), limit=120) - owner = _first_text(item.get("owner_agent"), item.get("claimed_by"), limit=120) - if not (todo_id or owner): - continue - lease: dict[str, Any] = {} - if todo_id: - lease["todo_id"] = todo_id - if owner: - lease["owner_agent"] = owner - for key in ("lease_until", "status"): - value = _text(item.get(key), limit=120) - if value: - lease[key] = value - write_scope = item.get("write_scope") - if isinstance(write_scope, list): - lease["write_scope"] = [ - scope - for scope in (_text(value, limit=120) for value in write_scope) - if scope - ] - compact.append(lease) - compact.extend( - _hard_lease_entries( - runtime_root=runtime_root, - goal_id=goal_id, - agent_todos=agent_todos, - ) - ) - return compact +def _compact_coordination_entry(item: Mapping[str, Any]) -> dict[str, Any]: + """The existing channel redaction boundary owns display text, not lease rules.""" + row: dict[str, Any] = {} + for key in ("todo_id", "owner_agent", "claimed_by", "lease_until", "expires_at", "status", "reason"): + value = _text(item.get(key), limit=120) + if value: + row[key] = value + for key in ("lease_version", "lease_epoch"): + if isinstance(item.get(key), int) and not isinstance(item[key], bool): + row[key] = item[key] + if isinstance(item.get("write_scope"), list): + row["write_scope"] = [text for value in item["write_scope"] if (text := _text(value, limit=120))] + return row def _compact_artifacts(artifacts: Sequence[Mapping[str, Any]] | None) -> list[dict[str, Any]]: @@ -468,6 +350,16 @@ def build_goal_channel_projection( project_asset = _project_asset(status_item_dict) user_todos = _compact_todos(project_asset, "user") agent_todos = _compact_todos(project_asset, "agent") + from .coordination_observation import observe_goal_coordination + + explicit = None if active_leases is None else [ + _compact_coordination_entry({**{key: item[key] for key in ("todo_id", "status", "lease_until", "write_scope") if key in item}, + "owner_agent": _first_text(item.get("owner_agent"), item.get("claimed_by"), limit=120)}) + for item in _as_mappings(active_leases) + if _text(item.get("todo_id"), limit=120) or _first_text(item.get("owner_agent"), item.get("claimed_by"), limit=120) + ] + observation = observe_goal_coordination(runtime_root=runtime_root, goal_id=str(goal_id), + agent_todos=agent_todos, explicit_entries=explicit) raw_keys = _raw_material_keys( status_item_dict, status_payload_dict, @@ -541,12 +433,7 @@ def build_goal_channel_projection( user_todos=user_todos, ), "artifacts": _compact_artifacts(artifacts), - "active_leases": _active_leases( - active_leases=active_leases, - agent_todos=agent_todos, - runtime_root=runtime_root, - goal_id=str(goal_id), - ), + "active_leases": [_compact_coordination_entry(row) for row in observation["entries"]], "recent_events": _recent_events(run_history_goal_dict), "source_warnings": _source_warnings(raw_keys), "truth_contract": { @@ -559,4 +446,14 @@ def build_goal_channel_projection( ), }, } + if "source_authority" in observation: + projection["coordination_observation"] = {key: observation[key] for key in ( + "status", "source_authority", "provider_revision", "observed_at", "legacy_fallback_used", + "total_count", "truncated", "display_limit", "todo_count", "lease_count") if key in observation} + if observation["status"] != "loaded": + projection["source_warnings"].append({"kind": "coordination_unavailable", + "message": "Task ownership could not be read; an empty list does not prove that no task is owned."}) + elif observation.get("truncated") is True: + projection["source_warnings"].append({"kind": "coordination_truncated", + "message": "Ownership display is limited to 100 entries; diagnostics are shown first. Read the full task/lease state before acting."}) return {key: value for key, value in projection.items() if value is not None} diff --git a/loopx/presentation/renderers/goal_channel_html.py b/loopx/presentation/renderers/goal_channel_html.py index 9500f3b1e6..560791194f 100644 --- a/loopx/presentation/renderers/goal_channel_html.py +++ b/loopx/presentation/renderers/goal_channel_html.py @@ -114,6 +114,7 @@ def render_goal_channel_projection_html(projection: Mapping[str, Any]) -> str: open_gates = _as_mappings(projection.get("open_gates")) artifacts = _as_mappings(projection.get("artifacts")) active_leases = _as_mappings(projection.get("active_leases")) + ownership_unavailable = _as_mapping(projection.get("coordination_observation")).get("status") == "unavailable" recent_events = _as_mappings(projection.get("recent_events")) source_warnings = _as_mappings(projection.get("source_warnings")) @@ -172,8 +173,8 @@ def render_goal_channel_projection_html(projection: Mapping[str, Any]) -> str: "reason", "claimed_by", ), - empty="No active claim or lease projected.", - tone="green", + empty="Task ownership is unavailable; see Source Warnings." if ownership_unavailable else "No active claim or lease projected.", + tone="red" if ownership_unavailable else "green", ), _html_item_panel( "artifacts", diff --git a/tests/control_plane/test_canonical_goal_channel_coordination.py b/tests/control_plane/test_canonical_goal_channel_coordination.py new file mode 100644 index 0000000000..048f64a546 --- /dev/null +++ b/tests/control_plane/test_canonical_goal_channel_coordination.py @@ -0,0 +1,90 @@ +"""Canonical ownership is authoritative even when the display/local files disagree.""" +import json + +import pytest + +from canonical_authority_fixture import initialize_canonical_authority, isolate_sqlite_runtime +from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection +from loopx.control_plane.goals.goal_channel_projection import build_goal_channel_projection +from test_goal_channel_hard_lease_visibility import GOAL_ID, _status_item, _write_active_lease + + +@pytest.mark.parametrize('provider', ['file', 'sqlite']) +def test_empty_canonical_ownership_never_revives_display_claim_or_local_lease(tmp_path, monkeypatch, provider): + if provider == 'sqlite': + isolate_sqlite_runtime(tmp_path, monkeypatch) + state = tmp_path / 'state.md' + state.write_text('---\nhandoff_mode: legacy\n---\n\n## Agent Todo\n') + _write_active_lease(tmp_path, owner='stale-agent') + projection = build_todo_runtime_shadow_projection(goal_id=GOAL_ID, handoff_mode='soft_claim', todos=[]) + initialize_canonical_authority(tmp_path, GOAL_ID, projection, state_path=state, provider=provider) + before = state.read_bytes() + observed = build_goal_channel_projection(goal_id=GOAL_ID, status_item=_status_item(claimed_by='stale-agent'), runtime_root=tmp_path) + assert observed['active_leases'] == [] + assert observed['coordination_observation']['source_authority'] == provider + '_v0' + assert state.read_bytes() == before + state.unlink() + assert build_goal_channel_projection(goal_id=GOAL_ID, runtime_root=tmp_path)['active_leases'] == [] + + +def test_explicit_empty_observation_is_not_unspecified(): + observed = build_goal_channel_projection(goal_id=GOAL_ID, status_item=_status_item(claimed_by='stale-agent'), active_leases=[]) + assert observed['active_leases'] == [] + + +@pytest.mark.parametrize('provider', ['file', 'sqlite']) +def test_canonical_failure_is_visible_and_never_discloses_or_falls_back(tmp_path, monkeypatch, provider): + from loopx.control_plane.goals import coordination_observation as adapter + from loopx.presentation.renderers.goal_channel_html import render_goal_channel_projection_html + if provider == 'sqlite': + isolate_sqlite_runtime(tmp_path, monkeypatch) + state = tmp_path / 'state.md' + state.write_text('## Agent Todo\n') + projection = build_todo_runtime_shadow_projection(goal_id=GOAL_ID, handoff_mode='legacy', todos=[]) + initialize_canonical_authority(tmp_path, GOAL_ID, projection, state_path=state, provider=provider) + _write_active_lease(tmp_path, owner='stale-agent') + def unavailable(*_args): + raise RuntimeError('PRIVATE_BACKEND_PATH_AND_CREDENTIAL') + monkeypatch.setattr(adapter, 'effect_runtime_result', unavailable) + observed = build_goal_channel_projection(goal_id=GOAL_ID, status_item=_status_item(claimed_by='stale-agent'), runtime_root=tmp_path) + assert observed['active_leases'] == [] + assert observed['coordination_observation']['status'] == 'unavailable' + assert observed['source_warnings'][-1]['kind'] == 'coordination_unavailable' + html = render_goal_channel_projection_html(observed) + assert 'PRIVATE_' not in html and 'PRIVATE_' not in json.dumps(observed) + assert 'Task ownership could not be read' in html + assert 'Task ownership is unavailable; see Source Warnings.' in html + assert 'No active claim or lease projected.' not in html + + +@pytest.mark.parametrize('provider', ['file', 'sqlite']) +def test_real_cli_export_uses_canonical_ownership(tmp_path, monkeypatch, provider): + import subprocess + import sys + if provider == 'sqlite': + isolate_sqlite_runtime(tmp_path, monkeypatch) + from test_goal_channel_hard_lease_visibility import _write_status_workspace + registry, repo, runtime = _write_status_workspace(tmp_path) + _write_active_lease(runtime, owner='stale-agent') + state = repo / 'ACTIVE_GOAL_STATE.md' + projection = build_todo_runtime_shadow_projection(goal_id=GOAL_ID, handoff_mode='legacy', todos=[]) + initialize_canonical_authority(runtime, GOAL_ID, projection, state_path=state, provider=provider) + before = state.read_bytes() + result = subprocess.run([sys.executable, '-m', 'loopx.cli', '--registry', str(registry), '--format', 'json', + 'status', '--goal-id', GOAL_ID], capture_output=True, text=True, timeout=90) + assert result.returncode == 0, result.stderr + payload = json.loads(result.stdout) + def channels(value): + if isinstance(value, dict): + if 'goal_channel_projection' in value: + yield value['goal_channel_projection'] + for child in value.values(): + yield from channels(child) + elif isinstance(value, list): + for child in value: + yield from channels(child) + found = list(channels(payload)) + assert found, payload.keys() + assert all(item['active_leases'] == [] for item in found) + assert all(item['coordination_observation']['source_authority'] == provider + '_v0' for item in found) + assert state.read_bytes() == before diff --git a/tests/control_plane_ts/authority_store_conformance.ts b/tests/control_plane_ts/authority_store_conformance.ts index 053d825401..f1bd232c44 100644 --- a/tests/control_plane_ts/authority_store_conformance.ts +++ b/tests/control_plane_ts/authority_store_conformance.ts @@ -1,6 +1,7 @@ import {registerAuthorityScanConformance} from "./authority_scan_conformance.ts"; import {executeCoordinationTodoArchiveCompleted} from "../../loopx/control_plane/coordination/todo_archive.ts"; import {registerHandoffModeConformance} from "./handoff_mode_conformance.ts"; +import {registerOwnershipObservationConformance} from "./ownership_observation_conformance.ts"; import assert from "node:assert/strict"; import { createHash } from "node:crypto"; import test from "node:test"; @@ -208,6 +209,7 @@ export function registerAuthorityStoreConformance( factory: AuthorityStoreConformanceFactory, ): void { registerAuthorityScanConformance(providerName, factory); + registerOwnershipObservationConformance(providerName, factory); registerNativePlanningUpdateConformance(providerName, factory); registerCoordinationReceiptConformance(providerName, factory); registerHandoffModeConformance(providerName, factory); diff --git a/tests/control_plane_ts/local_authority_provider.test.ts b/tests/control_plane_ts/local_authority_provider.test.ts index e2ab5ec459..6d692bebcd 100644 --- a/tests/control_plane_ts/local_authority_provider.test.ts +++ b/tests/control_plane_ts/local_authority_provider.test.ts @@ -98,6 +98,7 @@ function providerCalls(directory: string, revision: string, dryRun: boolean) { (value: unknown) => Promise ? K : never}[keyof typeof runtime]; // A new exported runtime action must deliberately enter this failure matrix. const requests = { + observeLocalCoordinationOwnership: [{...input, schema_version: "loopx_local_ownership_observation_request_v0"}], listLocalCoordinationTodos: [{...input, schema_version: runtime.LOCAL_COORDINATION_TODO_LIST_REQUEST_SCHEMA}], readLocalCoordinationTodo: [{...input, schema_version: runtime.LOCAL_COORDINATION_TODO_READ_REQUEST_SCHEMA}], mutateLocalCoordinationAuthority: [{...input, schema_version: runtime.LOCAL_COORDINATION_MUTATION_REQUEST_SCHEMA, diff --git a/tests/control_plane_ts/ownership_observation.test.ts b/tests/control_plane_ts/ownership_observation.test.ts new file mode 100644 index 0000000000..18271eb9fa --- /dev/null +++ b/tests/control_plane_ts/ownership_observation.test.ts @@ -0,0 +1,32 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import {projectOwnershipObservation, OWNERSHIP_OBSERVATION_SCHEMA} from "../../loopx/control_plane/coordination/ownership_observation.ts"; +const at = "2026-09-13T00:00:00Z"; +const request = {schema_version: OWNERSHIP_OBSERVATION_SCHEMA, observed_at: at, + todos: [{todo_id: "todo-a", claimed_by: "agent-a"}], explicit_entries: null, lease_rows: []}; +const lease = {schema_version: "task_lease_v0", todo_id: "todo-a", owner: "agent-b", + status: "active", expires_at: "2027-01-01T00:00:00Z", version: 2}; + +test("explicit empty differs from absent observation; source inputs are not mutated", () => { + assert.equal((projectOwnershipObservation(request).entries as unknown[]).length, 1); + const input = {...request, explicit_entries: []}; + assert.deepEqual(projectOwnershipObservation(input).entries, []); + assert.deepEqual(input.explicit_entries, []); +}); +test("explicit entries keep only the public ownership display fields", () => { + const result = projectOwnershipObservation({...request, explicit_entries: [{ + todo_id: "todo-a", owner_agent: "agent-a", status: "hard_lease", lease_epoch: 2, + write_scopes: ["private/**"], idempotency_key: "PRIVATE_OPERATION_KEY", private_backend: "PRIVATE_BACKEND", + }]}); + assert.deepEqual(result.entries, [{todo_id: "todo-a", owner_agent: "agent-a", status: "hard_lease", lease_epoch: 2}]); +}); +for (const patch of [{expires_at: "invalid"}, {schema_version: "unknown"}, {lease_epoch: true}, {todo_id: "wrong-id"}]) { + test(`invalid lease observation is visible and never crashes the channel: ${JSON.stringify(patch)}`, () => { + const result = projectOwnershipObservation({...request, lease_rows: [{todo_id: "todo-a", lease: {...lease, ...patch}}]}); + assert.deepEqual((result.entries as unknown[])[1], {todo_id: "todo-a", status: "hard_lease_unreadable", reason: "corrupt_lease"}); + }); +} +test("all leases use the same observation time; exact expiry is expired", () => { + assert.equal((projectOwnershipObservation({...request, lease_rows: [{todo_id: "todo-a", lease: {...lease, expires_at: at}}]}).entries as unknown[]).length, 1); + assert.throws(() => projectOwnershipObservation({...request, observed_at: "invalid"})); +}); diff --git a/tests/control_plane_ts/ownership_observation_conformance.ts b/tests/control_plane_ts/ownership_observation_conformance.ts new file mode 100644 index 0000000000..196c89a64e --- /dev/null +++ b/tests/control_plane_ts/ownership_observation_conformance.ts @@ -0,0 +1,53 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import {canonicalAuthoritySha256} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; +import {readCoordinationOwnership} from "../../loopx/control_plane/coordination/ownership_observation.ts"; +import type {AuthorityStoreConformanceFactory} from "./authority_store_conformance.ts"; +import {productionScaleCoordinationFixture} from "./production_scale_coordination_fixture.ts"; + +const at = "2026-09-13T00:00:00Z"; +export function registerOwnershipObservationConformance(provider: string, factory: AuthorityStoreConformanceFactory) { + for (const scenario of ["empty", "full_conflict", "bounded", "malformed_lease"] as const) { + test(`${provider}: readonly ownership observation ${scenario}`, async t => { + const {store} = await factory(t); + const goal = "ownership-goal"; + const projection = productionScaleCoordinationFixture(goal).projection; + const todos = projection.todos as JsonObject[]; + for (const todo of todos) delete todo.claimed_by; + const target = [...todos].reverse().find(todo => todo.role === "agent" && todo.done !== true)!; + const lease = {...(projection.leases as JsonObject[])[0]!, todo_id: target.todo_id, + status: "active", owner: "agent-a", expires_at: "2027-01-01T00:00:00Z", + private_diagnostic: "PRIVATE_LEASE_PAYLOAD", idempotency_key: "PRIVATE_OPERATION_KEY"}; + if (scenario !== "empty") target.claimed_by = "agent-b"; + if (scenario === "bounded") for (const todo of todos) if (todo.role === "agent" && todo.done !== true) todo.claimed_by = "agent-b"; + if (scenario === "malformed_lease") lease.expires_at = "invalid"; + projection.leases = scenario === "empty" ? [] : scenario === "bounded" + ? [...todos.slice(0, 140).map(todo => ({...lease, todo_id: todo.todo_id})), lease] + : [lease]; + (projection.todo_read_model as JsonObject).records_sha256 = canonicalAuthoritySha256(todos); + assert.equal((await store.commitAuthority({operation_id: "seed", expected_provider_revision: null, + next_projection: projection, events: [], receipts: []})).status, "applied"); + const before = await store.loadAuthority(); + const result = await readCoordinationOwnership(store, goal, at); + assert.equal(result.status, "loaded"); + assert.equal(result.todo_count, 464); + const entries = result.entries as JsonObject[]; + assert.equal(JSON.stringify(result).includes("PRIVATE_"), false); + if (scenario === "empty") assert.deepEqual(entries, []); + if (scenario === "full_conflict") { + assert.equal(entries[0]!.reason, "owner_conflicts_with_claim"); + assert.equal(entries[0]!.claimed_by, "agent-b"); + assert.equal(entries[0]!.todo_id, target.todo_id); + } + if (scenario === "bounded") { + assert.equal(entries.length, 100); assert.equal(result.truncated, true); + assert.ok(Number(result.total_count) > 100); + assert.equal(entries[0]!.reason, "owner_conflicts_with_claim"); + } + if (scenario === "malformed_lease") assert.equal(entries[0]!.status, "hard_lease_unreadable"); + assert.deepEqual(await store.loadAuthority(), before); + assert.equal((await store.readReceipt("observation")).status, "missing"); + }); + } +}