Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions docs/architecture/rfcs/typescript-control-plane-migration-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,20 @@ not the supported legacy writer. The absent arm also removes all four shadow
producer modules using the qualified CLI startup isolation above.
[Retained duties and rollback](../../reference/canonical-todo-completion-update.md#monitor-source-writer-isolation).

Shadow diagnostics now have an independent read-only boundary: retained outbox
disk codecs live in `local_authority_shadow_history.py`, and status/candidate
read transport in `local_authority_shadow_readback.py`. The public status CLI
does not load source producers; drain still uses its existing orchestration.
Original entry identity, malformed evidence, cursor residue and missing-store
behavior remain unchanged. Python retains historical filesystem decoding and
presentation IO; TS retains history proof, delivery and cleanup authority.
Existing producer imports are aliases for supported callers, not duplicate
implementations. Source and normal-wheel File/SQLite recipient checks compare
complete status results and unchanged original bytes with producers present or
absent. This is a historical-reader exit prerequisite, not writer deletion,
outbox disposition, installed cohort adoption or release-default qualification.
See [operator readback and rollback](../../reference/authority-observation-retirement.md#retained-readback-without-source-producers).

Existing Goal migration, two-policy ownership retirement and storage-format
upgrade retain separate receipts and exits. Original-receipt recovery does not
justify retaining `legacy` as a live policy. Required migration readers remain.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,17 @@ wheel 删源检查须保留后来写入及原 successor;晋升后重放保持
producer module。
见[保留义务及回滚](../../reference/canonical-todo-completion-update.md#monitor-source-writer-isolation)。

shadow 诊断现在有独立只读边界:历史 outbox 磁盘 codec 位于
`local_authority_shadow_history.py`,status/candidate read transport 位于
`local_authority_shadow_readback.py`。公开 status CLI 不加载源 producer;drain
仍走现有 orchestration。原 entry 身份、损坏证据、cursor 残留和缺失 store 行为
保持。Python 保留历史文件解码及展示 IO;TS 保留历史证明、投递与清理权限。
旧 producer 的已有导入为支持调用方保留 alias,不复制实现。源码及正常 wheel
的 File/SQLite 接收端检查,在 producer 存在/缺席两臂比较完整 status 响应和
原文件字节。这是历史 reader 的退役前置,不代表 writer 删除、outbox 处置、
安装 cohort 采用或发布默认资格。
见[操作读回与回滚](../../reference/authority-observation-retirement.md#retained-readback-without-source-producers)。

存量 Goal 迁移、两策略退役和格式升级仍各有独立回执及出口;原回执恢复不能成为
保留 legacy 活跃策略的理由,必要迁移 reader 保留。

Expand Down
18 changes: 18 additions & 0 deletions docs/reference/authority-observation-retirement.md
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,24 @@ can restart the old observer if its setting is still present, so clear the
setting on installations that must not resume observation after rollback.
No format rewrite or destructive data cleanup accompanies retirement.

## Retained readback without source producers

`authority-shadow status` reads retained outbox entries and candidate history
through independent read-only modules. It remains available when the old source
capture producers are absent. Pending prepared/committed files, cursor residue,
invalid evidence and the original management binding remain visible. Residue
counts are diagnostic hints: they do not authorize cleanup or prove delivery.
A missing candidate stays missing; status never creates a replacement identity,
repairs bytes, drains an outbox or grants a writer permission.

The retained Python filesystem codecs preserve the existing wire and errors;
candidate proof and drain effects remain in their TypeScript owners. Supported
unpromoted writes and `authority-shadow drain` still need their current producer
and orchestration modules. This separation permits eventual producer retirement
without discarding historical readback, but does not itself authorize deletion
or disposition of pending originals. Code rollback restores the previous import
layout without changing data or restoring retired write authority.

## Why two matching reads were not a capture guarantee

The observer sampled after the writer released its lock. Another writer could
Expand Down
9 changes: 6 additions & 3 deletions loopx/cli_commands/authority_shadow.py
Original file line number Diff line number Diff line change
Expand Up @@ -147,17 +147,20 @@ def handle_authority_shadow_command(
action = str(getattr(args, "authority_shadow_command", None))
payload: dict[str, object]
try:
from ..control_plane.coordination.local_authority_shadow_adapter import (
drain_local_authority_shadow_outbox,
from ..control_plane.coordination.local_authority_shadow_readback import (
local_authority_shadow_status,
)
from ..control_plane.coordination.local_authority_shadow_outbox import OutboxError
from ..control_plane.coordination.local_authority_shadow_history import OutboxError

try:
# The same resolver every writer hook uses, so drain and status address
# the lineage those hooks wrote.
runtime_root = effective_runtime_root(registry_path, runtime_root_arg)
if action == "drain":
from ..control_plane.coordination.local_authority_shadow_adapter import (
drain_local_authority_shadow_outbox,
)

if args.max_entries < 1:
raise ValueError("--max-entries must be at least 1")
result = drain_local_authority_shadow_outbox(
Expand Down
174 changes: 7 additions & 167 deletions loopx/control_plane/coordination/local_authority_shadow_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@
from pathlib import Path
from typing import Any

from ..projects.registry_codec import load_registry
from ...paths import effective_runtime_root as effective_runtime_root, resolve_runtime_root
from ...registry import find_registry_goal
from ..effect_runtime import effect_runtime_result
Expand All @@ -22,13 +21,15 @@
LOCAL_AUTHORITY_SHADOW_READ_RESULT_SCHEMA,
LOCAL_AUTHORITY_SHADOW_TRANSACTION_EVIDENCE_SCHEMA,
)
from .local_authority_shadow_projection import (
head_digest,
)
from .configuration import resolve_coordination_runtime_shadow_config, local_authority_shadow_summary
from .configuration import resolve_coordination_runtime_shadow_config
from .authority_source_capture import todo_partition_projector as todo_partition_projector
from .shadow_management import read_shadow_capture_binding, shadow_management_state_path
from .shadow_management import shadow_management_state_path
from .shadow_goal_scope import shadow_goal_scope
from .local_authority_shadow_readback import (
RETENTION_PRESSURE_BYTES as RETENTION_PRESSURE_BYTES,
read_local_authority_shadow as read_local_authority_shadow,
local_authority_shadow_status as local_authority_shadow_status,
)



Expand All @@ -51,8 +52,6 @@
INLINE_DRAIN_LOCK_TIMEOUT_SECONDS = 0.25
SHADOW_DRAIN_SCHEMA = "loopx_shadow_drain_v0"
SHADOW_EXACT_DRAIN_SCHEMA = "loopx_shadow_drain_v1"
RETENTION_PRESSURE_BYTES = 8 * 1024 * 1024
_SEED_WRITE_CLASSES = {"seed", "reseed_after_crash_gap"}


@dataclass
Expand Down Expand Up @@ -108,41 +107,6 @@ def to_payload(self) -> dict[str, Any]:
return payload


def read_local_authority_shadow(
*,
runtime_root: Path,
goal_id: str,
store_kind: str = "runtime_shadow",
scan_after_cursor: str | None = None,
scan_limit: int = 0,
receipt_operation_id: str | None = None,
read_model: str = "full",
) -> dict[str, Any]:
"""Read-only candidate view through the TypeScript store boundary."""

result = effect_runtime_result(
"coordination.runtime_shadow.outbox_read",
{
"schema_version": LOCAL_AUTHORITY_SHADOW_READ_REQUEST_SCHEMA,
"runtime_root": str(runtime_root),
"goal_id": goal_id,
"store_kind": store_kind,
"scan_after_cursor": scan_after_cursor,
"scan_limit": scan_limit,
"receipt_operation_id": receipt_operation_id,
"read_model": read_model,
},
timeout=15.0,
)
if (
not isinstance(result, dict)
or result.get("schema_version") != LOCAL_AUTHORITY_SHADOW_READ_RESULT_SCHEMA
or result.get("goal_id") != goal_id
):
raise RuntimeError("local authority shadow read result is invalid")
return dict(result)


def _drain_prelude(
result: DrainResult,
*,
Expand Down Expand Up @@ -226,130 +190,6 @@ def drain_local_authority_shadow_outbox(
return result


class _CandidateMissing(Exception):
"""The candidate store directory does not exist yet."""


def _store_bytes(runtime_root: Path, goal_id: str, *, legacy_observation: bool) -> int:
directory = (
runtime_root / "authority-shadow" / "file" / goal_id
if legacy_observation
else runtime_root / "authority-shadow" / "file-v0"
)
if not directory.is_dir():
return 0
return sum(path.stat().st_size for path in directory.iterdir() if path.is_file())


def local_authority_shadow_status(
*,
registry_path: Path,
runtime_root: Path | None,
goal_id: str,
) -> dict[str, Any]:
"""Operator readback: configuration, outbox backlog, and candidate head facts."""

registry = load_registry(registry_path)
goal = find_registry_goal(registry, goal_id)
if not isinstance(goal, dict):
raise ValueError(f"goal {goal_id!r} is not registered")
if runtime_root is None:
runtime_root = resolve_runtime_root(registry, None, registry_path=registry_path)
config = local_authority_shadow_summary(goal)
runtime_config = resolve_coordination_runtime_shadow_config(goal)
management = read_shadow_capture_binding(runtime_root, goal_id)
legacy_observation = (
config.get("configured") is True
and not runtime_config.enabled
and management["status"] == "missing"
)
backlog = outbox.outbox_summary(runtime_root, goal_id)
candidate: dict[str, Any]
try:
directory = (
runtime_root / "authority-shadow" / "file" / goal_id
if legacy_observation
else runtime_root / "authority-shadow" / "file-v0"
)
if not directory.is_dir() or not any(directory.glob("authority-store-*.json")):
# Reading through the store boundary would mint a store identity;
# a status probe must not create candidate lineage.
raise _CandidateMissing
view = read_local_authority_shadow(
runtime_root=runtime_root,
goal_id=goal_id,
store_kind=(
"legacy_observation" if legacy_observation else "runtime_shadow"
),
)
head = view.get("head") if isinstance(view.get("head"), dict) else None
candidate = {
"status": view.get("status"),
"reason_code": view.get("reason_code"),
"store_identity": view.get("store_identity"),
"provider_revision": view.get("provider_revision"),
"cursor": view.get("cursor"),
"head_digest": view.get("head_digest"),
"head_schema_version": head.get("schema_version") if head else None,
"partitions": view.get("partitions"),
"codec_agreement": (head_digest(head) == view.get("head_digest"))
if head
else None,
}
except _CandidateMissing:
candidate = {
"status": "missing",
"reason_code": None,
"store_identity": None,
"provider_revision": None,
"cursor": None,
"head_digest": None,
"head_schema_version": None,
"partitions": None,
"codec_agreement": None,
}
except Exception:
candidate = {
"status": "unavailable",
"reason_code": "shadow_read_failed",
"store_identity": None,
"provider_revision": None,
"cursor": None,
"head_digest": None,
"head_schema_version": None,
"partitions": None,
"codec_agreement": None,
}
candidate["store_kind"] = "legacy_observation" if legacy_observation else "runtime_shadow"
candidate["historical_only"] = legacy_observation
try:
store_bytes = _store_bytes(
runtime_root, goal_id, legacy_observation=legacy_observation
)
storage_error = None
except OSError:
store_bytes = None
storage_error = "shadow_store_unavailable"
return {
"ok": all(item["invalid"] is None for item in backlog.values())
and storage_error is None
and management["status"] != "hold",
"action": "status",
"goal_id": goal_id,
"config": config,
"runtime_config": asdict(runtime_config),
"management": management,
"storage_error": storage_error,
"runtime_root_digest": outbox.runtime_root_digest(runtime_root),
"outbox": backlog,
"candidate": candidate,
"store_bytes": store_bytes,
"retention_pressure": store_bytes > RETENTION_PRESSURE_BYTES
if store_bytes is not None
else None,
}


def capture_evidence(
*,
goal_id: str,
Expand Down
Loading
Loading