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
1 change: 1 addition & 0 deletions .github/workflows/python-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,7 @@ jobs:
tests/test_file_lock_cross_process.py
tests/control_plane/test_coordination_file_provider.py
tests/control_plane/test_effect_runtime_integration.py
tests/control_plane/test_local_authority_shadow_outbox.py
tests/test_self_update_runtime_activation.py
tests/test_windows_install.py

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2134,14 +2134,14 @@ write that cannot preserve the contract.
fence, run `promote`, render the projections, and record the declaration of
question 11; refuse a rerun whose source digest changed unless the abandoned
store is explicitly discarded.
- A transaction-bound capture (question 14): the runtime shadow samples the
source after the commit, from outside the writer's lock, so a concurrent
writer can land inside the sampled projection and a crash between the commit
and the dispatch loses the mirror. The parity-half outbox closes both: the
prepared entry is written inside the lock the writer already holds, the
committed marker after the primary write returns, and a bounded drain turns
each entry into exactly one shadow transaction whose `operation_id` is the
entry id. `qualify` then counts entries, not samples.
- Full transaction-capture qualification (question 14): Todo add, update,
complete, supersede, and archive plus native lease acquire, renew, transfer,
release, auto-acquire, and fence-close now emit prepared/committed outbox
entries around their primary write. The bounded drain commits complete
versioned records into the existing `coordination.runtime_shadow` file-v0
lineage; it does not create the former second local-shadow candidate. Before
promotion, add sustained mixed-writer parity runs, event-only Todo coverage,
and the selected provider profile's recovery/capacity evidence.
- The provider-neutral authority binding, compatibility projection outbox,
and conformance rows for file, NoKV, and PostgreSQL. This does not require all
providers to promote together; each profile must pass the same contract
Expand Down Expand Up @@ -2172,6 +2172,6 @@ one document is acceptable only for promotion bootstrap and bounded test goals.
| Lane | May start | Scope and exit condition | Dependency |
| --- | --- | --- | --- |
| P. PostgreSQL provider plane | Now, from current `main` | Keep the existing `AuthorityStore` contract; finish schema migration/install ownership, authenticated service and tenant authorization, restore-incarnation rotation, pool/cancellation/failover behavior, and reviewed indexes, partitioning, retention, and measured capacity. Live PostgreSQL conformance remains mandatory. | Does not depend on #3870 and must not stack on its branch. This lane alone creates no runtime caller or promotion claim. |
| C. Canonical transaction capture | Now, by revising or replacing #3870 | Make the outbox transport complete versioned Todo/lease records into `coordination.runtime_shadow.commit`; retire the duplicate observation/local-shadow lineage and prove omission/explicit-clear behavior. | Can run in parallel with P, but both C and the selected provider profile must finish before parity or promotion integration. |
| C. Canonical transaction capture | In implementation, based on #3870 | Transaction-bound outbox capture now targets the one `coordination.runtime_shadow` lineage and retains complete versioned Todo/lease records. Finish sustained mixed-writer parity, explicit-clear/omission coverage, and event-only Todo recovery evidence. | Can run in parallel with P, but both C and the selected provider profile must finish before parity or promotion integration. |
| I. Binding and qualification integration | After P and C | Bind one exact provider lineage, field manifest, source revision, digest, and cursor; run sustained transaction parity and provider-specific recovery/capacity qualification without consulting legacy state for missing fields. | This is the merge point between provider work and capture work. |
| F. Promotion and cleanup | After I and explicit maintainer approval | Add provider-first CLI routing, the lock-owning promotion orchestrator, compatibility projection outbox, post-promotion fenced export/rollback, then delete duplicate reference aggregates and flip the reviewed stage/hold declarations. | No profile is eligible until its exact implementation and lineage pass P, C, and I. |
Original file line number Diff line number Diff line change
Expand Up @@ -1698,11 +1698,12 @@ integrity chain、确定性 scan 与 recovery readback。物理 profile 可以
片):取 Todo 与 lease 两把 legacy 锁,要求 `qualify` 在当前 revision 与 digest 上
为 `qualified`,engage fence,执行 `promote`,渲染投影,写入问题 11 的声明;源
digest 已变时拒绝重跑,除非显式丢弃被放弃的 store。
- 事务绑定的捕获(问题 14):runtime shadow 在提交之后、写者锁之外采样源,并发写者
可能落进被采样的投影,commit 与 dispatch 之间崩溃则丢失镜像。parity 半段的 outbox
同时关闭两者:prepared entry 在写者已持有的锁内写入,committed 标记在主写返回后
写入,有界 drain 把每条 entry 变成恰好一笔 `operation_id` 为 entry id 的 shadow
事务。此后 `qualify` 数的是 entry,不是采样。
- 完成事务捕获资格验证(问题 14):Todo add、update、complete、supersede、archive,
以及 native lease acquire、renew、transfer、release、auto-acquire、fence-close,现已在
主写前后生成 prepared/committed outbox entry。有界 drain 把完整、带版本的 record
提交到既有 `coordination.runtime_shadow` file-v0 lineage,不再创建第二套 local-shadow
candidate。晋升前仍需补持续 mixed-writer parity、event-only Todo 覆盖和所选 provider
profile 的 recovery/capacity 证据。
- provider-neutral authority binding、兼容投影 outbox,以及 file、NoKV、PostgreSQL
的 conformance row。三个 provider 不必同时晋升,但每个 profile 都必须先通过同一
合同才具备资格。
Expand All @@ -1727,6 +1728,6 @@ latency、response size 与 recovery。单文档全量保留只适用于 promoti
| Lane | 何时开始 | 范围与退出条件 | 依赖 |
| --- | --- | --- | --- |
| P. PostgreSQL provider plane | 现在,从当前 `main` 开始 | 保持既有 `AuthorityStore` 合同;完成 schema migration/install ownership、authenticated service 与 tenant authorization、restore-incarnation rotation、pool/cancellation/failover 行为,以及经评审的 index、partition、retention 与实测 capacity。真实 PostgreSQL conformance 始终是强制门禁。 | 不依赖 #3870,也不得叠在其分支上。仅完成本 lane 不产生 runtime caller 或 promotion 声明。 |
| C. Canonical transaction capture | 现在,通过修订或替代 #3870 | 让 outbox 把完整、带版本的 Todo/lease record 传给 `coordination.runtime_shadow.commit`;退役重复 observation/local-shadow lineage,并证明 omission/explicit-clear 行为。 | 可与 P 并行;但 C 与选定 provider profile 都完成后,才能进入 parity 或 promotion 集成。 |
| C. Canonical transaction capture | 实现中,基于 #3870 | transaction-bound outbox 已指向唯一 `coordination.runtime_shadow` lineage,并保留完整带版本的 Todo/lease record;继续完成 sustained mixed-writer parity、explicit-clear/omission 与 event-only Todo recovery 证据。 | 可与 P 并行;但 C 与选定 provider profile 都完成后,才能进入 parity 或 promotion 集成。 |
| I. Binding 与资格集成 | P 与 C 完成后 | 绑定一个精确 provider lineage、field manifest、source revision、digest 与 cursor;运行持续 transaction parity,以及 provider-specific recovery/capacity qualification;缺字段时不得查询 legacy state 补齐。 | 这是 provider 工作与 capture 工作的汇合点。 |
| F. Promotion 与清理 | I 完成且 maintainer 显式批准后 | 加入 provider-first CLI routing、持锁 promotion orchestrator、兼容投影 outbox、晋升后 fenced export/rollback;随后删除重复 reference aggregate,并翻转经评审的 stage/hold 声明。 | 一个 profile 的精确实现与 lineage 通过 P、C、I 前,不具备晋升资格。 |
13 changes: 13 additions & 0 deletions loopx/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@
handle_support_control_command,
handle_handoff_mode_command,
handle_task_lease_command,
handle_authority_shadow_command,
handle_version_command,
handle_host_mode_plan_command,
handle_worker_bridge_command,
Expand Down Expand Up @@ -140,6 +141,7 @@
register_support_control_commands,
register_handoff_mode_command,
register_task_lease_command,
register_authority_shadow_command,
register_todo_command,
register_version_command,
register_host_mode_plan_command,
Expand Down Expand Up @@ -325,6 +327,7 @@ def build_parser() -> LoopXArgumentParser:
register_todo_command(sub, add_subcommand_format)
register_coordination_shadow_command(sub, add_subcommand_format)
register_task_lease_command(sub, add_subcommand_format)
register_authority_shadow_command(sub, add_subcommand_format)
register_handoff_mode_command(sub, add_subcommand_format)
register_shared_goal_alignment_command(sub, add_subcommand_format)
register_quota_command(sub)
Expand Down Expand Up @@ -814,6 +817,16 @@ def main(argv: list[str] | None = None) -> int:
if task_lease_result is not None:
return task_lease_result

authority_shadow_result = handle_authority_shadow_command(
args,
registry_path=registry_path,
runtime_root_arg=args.runtime_root,
output_format=output_format,
print_payload=print_payload,
)
if authority_shadow_result is not None:
return authority_shadow_result

handoff_mode_result = handle_handoff_mode_command(
args,
registry_path=registry_path,
Expand Down
6 changes: 6 additions & 0 deletions loopx/cli_commands/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,10 @@ def _load_exports() -> None:
handle_support_control_command,
register_support_control_commands,
)
from .authority_shadow import (
handle_authority_shadow_command,
register_authority_shadow_command,
)
from .task_lease import handle_task_lease_command, register_task_lease_command
from .todo import handle_todo_command, register_todo_command
from .version import handle_version_command, register_version_command
Expand Down Expand Up @@ -221,6 +225,7 @@ def _load_exports() -> None:
"handle_starter_visible_pilot_command",
"handle_summary_all_command",
"handle_support_control_command",
"handle_authority_shadow_command",
"handle_task_lease_command",
"handle_todo_command",
"handle_version_command",
Expand Down Expand Up @@ -268,6 +273,7 @@ def _load_exports() -> None:
"register_summary_all_command",
"register_status_commands",
"register_support_control_commands",
"register_authority_shadow_command",
"register_task_lease_command",
"register_todo_command",
"register_version_command",
Expand Down
219 changes: 219 additions & 0 deletions loopx/cli_commands/authority_shadow.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,219 @@
from __future__ import annotations

import argparse
from collections.abc import Callable
from pathlib import Path

from ..control_plane.coordination.local_authority_shadow_adapter import (
CLI_DRAIN_LOCK_TIMEOUT_SECONDS,
drain_local_authority_shadow_outbox,
effective_runtime_root,
local_authority_shadow_status,
)
from ..control_plane.coordination.local_authority_shadow_outbox import OutboxError
from ..file_lock import LockAcquireTimeoutError


AUTHORITY_SHADOW_CLI_SCHEMA = "loopx_authority_shadow_cli_v0"
CLI_DRAIN_MAX_ENTRIES = 256
CLI_DRAIN_BUDGET_SECONDS = 30.0

PrintPayload = Callable[
[dict[str, object], str, Callable[[dict[str, object]], str]],
None,
]


_DRAIN_FIELDS = (
"outcome",
"reason_code",
"delivered",
"replayed",
"reconciled",
"no_op",
"reseeded",
"reclaimed_residue",
"pending_after",
"prepared_only_after",
"budget_exhausted",
"last_cursor",
"candidate_readback_verified",
)


def _drain_markdown_lines(payload: dict[str, object]) -> list[str]:
lines = [f"- {key}: `{payload.get(key)}`" for key in _DRAIN_FIELDS]
stopped_at = payload.get("stopped_at")
if isinstance(stopped_at, dict):
lines.append(
"- stopped_at: "
f"`{stopped_at.get('partition')}#{stopped_at.get('seq')}` "
f"→ `{stopped_at.get('outcome')}` ({stopped_at.get('reason_code')})"
)
return lines


def _status_markdown_lines(payload: dict[str, object]) -> list[str]:
lines: list[str] = []
config = payload.get("config")
if isinstance(config, dict):
lines.append(f"- config: `{config.get('status')}`")
backlog = payload.get("outbox")
partitions = backlog.items() if isinstance(backlog, dict) else []
for partition, facts in partitions:
if isinstance(facts, dict):
lines.append(
f"- outbox.{partition}: committed_pending="
f"`{facts.get('committed_pending')}` prepared_only="
f"`{facts.get('prepared_only')}` cursor_last_seq="
f"`{facts.get('cursor_last_seq')}`"
)
candidate = payload.get("candidate")
if isinstance(candidate, dict):
lines.append(
f"- candidate: status=`{candidate.get('status')}` cursor="
f"`{candidate.get('cursor')}` store_identity=`{candidate.get('store_identity')}`"
)
lines.append(f"- store_bytes: `{payload.get('store_bytes')}`")
lines.append(f"- retention_pressure: `{payload.get('retention_pressure')}`")
return lines


def render_authority_shadow_markdown(payload: dict[str, object]) -> str:
lines = [
"# LoopX Authority Shadow",
"",
f"- ok: `{payload.get('ok')}`",
f"- action: `{payload.get('action')}`",
f"- goal_id: `{payload.get('goal_id')}`",
]
if payload.get("error"):
lines.append(f"- error: {payload.get('error')}")
if payload.get("error_code"):
lines.append(f"- error_code: `{payload.get('error_code')}`")
if payload.get("action") == "drain":
lines.extend(_drain_markdown_lines(payload))
elif payload.get("action") == "status":
lines.extend(_status_markdown_lines(payload))
return "\n".join(lines) + "\n"


def register_authority_shadow_command(
subparsers: argparse._SubParsersAction,
add_subcommand_format: Callable[[argparse.ArgumentParser], None],
) -> None:
parser = subparsers.add_parser(
"authority-shadow",
help=(
"Drain or inspect the transaction-bound local authority shadow outbox "
"for one goal. The candidate store is evidence only; it never decides."
),
)
add_subcommand_format(parser)
parser.add_argument(
"authority_shadow_command",
choices=["drain", "status"],
help="drain delivers pending outbox entries; status reports backlog and candidate facts.",
)
parser.add_argument("--goal-id", required=True, help="Goal id whose shadow outbox to operate on.")
parser.add_argument(
"--max-entries",
type=int,
default=CLI_DRAIN_MAX_ENTRIES,
help=f"Maximum entries one drain pass delivers (default {CLI_DRAIN_MAX_ENTRIES}).",
)
parser.add_argument(
"--budget-seconds",
type=float,
default=CLI_DRAIN_BUDGET_SECONDS,
help=f"Wall-clock budget for one drain pass (default {CLI_DRAIN_BUDGET_SECONDS:g}).",
)
parser.add_argument(
"--lock-timeout-seconds",
type=float,
default=CLI_DRAIN_LOCK_TIMEOUT_SECONDS,
help=(
"How long to wait for the per-goal drain lock before reporting drain_deferred "
f"(default {CLI_DRAIN_LOCK_TIMEOUT_SECONDS:g})."
),
)


def handle_authority_shadow_command(
args: argparse.Namespace,
*,
registry_path: Path,
runtime_root_arg: str | None,
output_format: Callable[..., str],
print_payload: PrintPayload,
) -> int | None:
if args.command != "authority-shadow":
return None
action = str(getattr(args, "authority_shadow_command", None))
payload: dict[str, object]
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":
if args.max_entries < 1:
raise ValueError("--max-entries must be at least 1")
result = drain_local_authority_shadow_outbox(
registry_path=registry_path,
runtime_root=runtime_root,
goal_id=args.goal_id,
max_entries=args.max_entries,
budget_seconds=args.budget_seconds,
lock_timeout_seconds=args.lock_timeout_seconds,
)
payload = {
"schema_version": AUTHORITY_SHADOW_CLI_SCHEMA,
"action": "drain",
**result.to_payload(),
}
else:
payload = {
"schema_version": AUTHORITY_SHADOW_CLI_SCHEMA,
**local_authority_shadow_status(
registry_path=registry_path,
runtime_root=runtime_root,
goal_id=args.goal_id,
),
}
except OutboxError as exc:
payload = {
"ok": False,
"schema_version": AUTHORITY_SHADOW_CLI_SCHEMA,
"action": action,
"goal_id": args.goal_id,
"error": str(exc),
"error_code": exc.reason_code,
}
except LockAcquireTimeoutError as exc:
payload = {
"ok": False,
"schema_version": AUTHORITY_SHADOW_CLI_SCHEMA,
"action": action,
"goal_id": args.goal_id,
"error": str(exc),
**exc.to_payload(),
}
except Exception as exc:
payload = {
"ok": False,
"schema_version": AUTHORITY_SHADOW_CLI_SCHEMA,
"action": action,
"goal_id": args.goal_id,
"error": str(exc),
"error_code": exc.__class__.__name__,
}
print_payload(payload, output_format(args), render_authority_shadow_markdown)
return 0 if payload.get("ok") else 1


__all__ = [
"AUTHORITY_SHADOW_CLI_SCHEMA",
"handle_authority_shadow_command",
"register_authority_shadow_command",
"render_authority_shadow_markdown",
]
Loading
Loading