diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0-evidence.zh-CN.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0-evidence.zh-CN.md index f0b5811bab..8146a99668 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0-evidence.zh-CN.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0-evidence.zh-CN.md @@ -195,3 +195,48 @@ live NoKV 测试。确定性 fake 的通过结果与这项静态 API 核对必 因绑定具体部署形态保留在本地忽略状态,不入公共树。回执增长、并发包络 等数字是单节点 dev 栈的量化观测,按第 5 节规则不构成生产结论;其含义 已写入 RFC 第 11 节 Stage 2 状态小节。 + +## 8. Stage ladder E2E 复跑记录(2026-09-03) + +按第 5 节的复跑要求逐项记录本轮 stage ladder 证据: + +1. **精确 commit**:LoopX 侧为 `test/shared-authority-e2e-ladder` 分支(基于 + PR #3818 第三轮修订头,含单一 effective runtime root 修复;ladder 报告的 `bindings.loopx_commit` 与 + `bindings.loopx_tree_dirty` 记录实际运行的树)。NoKV 侧本轮在本机单节点栈上 + 执行(`nokv` 83971e62ab,Python SDK 0.11.0,`nokv serve` 以静态 etcd 路由 + + RustFS 对象存储运行;报告只记录 `bindings.nokv_client_config_sha256` 与 SDK + 版本,不记录任何连接值);PostgreSQL 侧无可达栈,对应行按设计报告 `unverified`。 +2. **探针源码**:`loopx/control_plane/testing/authority_e2e_ladder.py`、 + `loopx/control_plane/testing/authority_e2e_fixtures.py` 与只读 TypeScript + 探针 `tests/control_plane_ts/authority_store_readback_probe.ts`(随分支评审; + 报告 `bindings.probe_sha256[]` 记录其 digest)。入口为 + `examples/shared-goal-authority-e2e/ladder.py`,pytest 投影为 + `tests/control_plane/test_shared_goal_authority_e2e.py`。各行驱动的是真实 + `python -m loopx.cli` 与生产 `FileAuthorityStore`,不是参考实现。 +3. **断言**:本轮 9 个 deterministic 行全部 pass,2 个 NoKV live 行在本机栈上 + pass:`s0.nokv_live_matrix` 十三个 NoKV 场景行全为 true 且与 file provider 的 + 十二个共享行逐行一致;`s2a.nokv_live_qualification` 对既有 workbench 以新铸 + tenant/goal 运行已合并的 live 资格探针,13 项 check 全部 `passed`,final + generation `3`,SDK `0.11.0` / API `1`,且报告声明未改变 authority source、未 + 证明可用性。deterministic 行:`s0.file_matrix_twelve_rows` + 十二个 file provider 场景行全为 true;`s1.cli_document_decodes_through_ts_store` + 三次 CLI 写入经 TypeScript store 回读 cursor `3`、operation id 按序一致、首条 + receipt found;七个 `s2c1.*` 行(configure 往返、12 个 writer family 全部 + captured 且候选 cursor `12`、default-off 隔离、候选失败保主写、SIGKILL 崩溃 + 间隙只丢失一次 observation、`--runtime-root` 与 `common_runtime_root` 不同时 + 五次写入落入单一 store identity 且候选 cursor `5`、`migrate-state` 新 lineage + cursor `1`)。 +4. **负例**:幂等 re-acquire 不携带 `authority_shadow`;default-off goal 无 + `authority-shadow/` 目录且响应字段与 observed goal 一致;候选目录被占用时 + `outcome=failed` / `reason_code=shadow_observation_failed` 但主写已提交; + 迁移后的候选序列化中不含旧 store identity、legacy revision、源路径与私有 + 字节;报告隐私扫描把注入的临时路径改写为 `fail/privacy_violation`,仅泄露到 + bindings 块时 `summary.privacy_violations=1` 且退出码为 1,且 live 变量清空时 + ladder 退出码为 1(均由 pytest 钉住)。 +5. **未执行与限定**:`s2b.postgresql_conformance_live`(`postgres_url_missing`) + 本轮 unverified;在没有 NoKV 与 PostgreSQL 栈的 CI 环境里三个 live 行都报告 + `unverified`。9 个 `s2c2.*` 行以 pending 声明,未宣称;选中任一 pending 行而不传 + `--allow-pending` 时 ladder 以非零退出,零执行的报告不可能显示为 green。默认 + 全量运行退出码为 1,`--allow-unverified --allow-pending` 才为 0。Stage 2C + parity、outbox、drain 与增长门槛均未验证;本记录不构成任何 provider 晋升 + 或生产结论。 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 272820c937..dcb6c2652a 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -1563,6 +1563,79 @@ flip, rollback, and retention decisions below - not a diagnostic CLI that creates a second writer. This status section claims proven contracts, not a shipped production capability. +#### Stage-ladder end-to-end evidence (2026-09-03) + +What exists on this branch is one incremental end-to-end "stage ladder" that +exercises every completed stage claim of this RFC through the real +`python -m loopx.cli` and reports a machine-checkable verdict per row: +`loopx/control_plane/testing/authority_e2e_ladder.py` (row registry, runners, +the `loopx_shared_goal_authority_e2e_report_v0` JSON report, exit policy, and +privacy scan), `loopx/control_plane/testing/authority_e2e_fixtures.py` (goal +workspaces, CLI runners, the observation-lock window, candidate read-back), the +read-only TypeScript probe +`tests/control_plane_ts/authority_store_readback_probe.ts`, the pytest +projection `tests/control_plane/test_shared_goal_authority_e2e.py`, and the +entry point `examples/shared-goal-authority-e2e/ladder.py`. + +Per stage, this increment implements: + +- Stage 0: `s0.file_matrix_twelve_rows` runs the retained live matrix script + and requires exactly the twelve shared scenario rows to be true on the file + provider; `s0.nokv_live_matrix` requires the same rows plus + `restored_lineage_fails_closed` and identical file/NoKV outcomes on a live + NoKV stack. +- Stage 1: `s1.cli_document_decodes_through_ts_store` writes three + observations through the product CLI (`todo add`, `task-lease acquire`, + `todo update`) and reads them back through `FileAuthorityStore` with + `loadAuthority`, paged `scanCommitted`, and `readReceipt`: cursor `3`, the + three operation ids in order, and the first receipt found. +- Stage 2A: `s2a.nokv_live_qualification` runs the merged live qualification + probe (`examples/nokv-authority-store/live-qualification.ts --execute-live`) + against an existing workbench with a fresh tenant/goal pair and requires + `ok=true`, the single-node store-conformance scope, every check passed, NoKV + SDK `0.11.0` / API `1`, and no promotion or availability claim. +- Stage 2B: `s2b.postgresql_conformance_live` runs the PostgreSQL integration + test file under node's TAP reporter and requires at least nine passes, zero + failures, and zero skips. +- Stage 2C observation foundation: seven `s2c1.*` rows port the local-shadow CLI + E2E and migration assertions and pin the single-lineage guarantee. The configure round trip previews, enables, + reads back, and disables the observer; every writer family (handoff-mode, + todo add/update/complete/supersede/capture-followups/archive-completed, + task-lease acquire/renew/transfer) captures with + `primary_writeback_preserved`, `provider_to_local_writes=false`, and + `candidate_read_for_decision=false`, while an idempotent re-acquire does not + observe; default-off goals stay isolated; candidate failure preserves the + primary commit; a POSIX SIGKILL in the crash gap loses only that + observation; a `--runtime-root` override that differs from + `common_runtime_root` keeps todo add, task-lease acquire, todo update, + follow-up capture, and a leased completion in one store identity while the + registry root gains neither a candidate lineage nor lease state; and + `migrate-state` seeds a fresh lineage without legacy bytes. + +Live rows are environment-gated (`LOOPX_TEST_POSTGRES_URL`; +`NOKV_COORDINATION_LIVE=1` plus the `NOKV_*` stack variables; +`LOOPX_NOKV_AUTHORITY_LIVE=1` plus the `LOOPX_NOKV_AUTHORITY_*` inputs). +Without a stack they report `unverified`, and the ladder exits non-zero unless +`--allow-unverified` is passed; an unverified row is never counted as green. +A pending row is an unmet obligation as well: selecting one exits non-zero +unless `--allow-pending` is passed, so a report cannot read as green while it +executed nothing. +The report binds the LoopX commit, tree dirtiness, probe digests, and hashed +connection facts, and its privacy scan turns any leak of a temporary root, +home directory, connection URL, or configuration path into +`fail/privacy_violation`; a leak confined to the bindings block is redacted +and still fails the run through `summary.privacy_violations`, which no flag +relaxes. + +Delivery boundary: test-only. No production entry point constructs any store; +the ladder adds no product path and reads the candidate only through the +retained TypeScript store. The Stage 2C parity half +(`s2c2.*`: outbox entries, idempotent drain, SIGKILL before and during drain, +rollback with pending entries, parity equal and divergent, +migration seed-and-drain, growth measurement) are declared as pending rows, +not claimed. This subsection records executable evidence for the stages above; +it does not promote any provider or complete the Stage 2C promotion. + ### P0: contract and deterministic proof - this ownership matrix and explicit shared-mode boundary; 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 40d8f25cb2..900c0972ab 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 @@ -1250,6 +1250,66 @@ visible governance 台账中属 coverage-only。该接受允许内聚的 referen retention 决策,不能用一个会制造第二 writer 的诊断 CLI 代替。本状态节声明的是已 证明的合同,不是已 ship 的生产能力。 +#### Stage ladder 端到端证据(2026-09-03) + +本分支上存在一条增量式端到端 "stage ladder":它通过真实的 +`python -m loopx.cli` 逐行演练本 RFC 每个已完成阶段的声明,并按行给出可机器 +判定的结论:`loopx/control_plane/testing/authority_e2e_ladder.py`(行注册表、 +runner、`loopx_shared_goal_authority_e2e_report_v0` JSON 报告、退出策略与隐私 +扫描)、`loopx/control_plane/testing/authority_e2e_fixtures.py`(goal 工作区、 +CLI runner、observation-lock 窗口、候选回读)、只读 TypeScript 探针 +`tests/control_plane_ts/authority_store_readback_probe.ts`、pytest 投影 +`tests/control_plane/test_shared_goal_authority_e2e.py`,以及入口 +`examples/shared-goal-authority-e2e/ladder.py`。 + +按阶段,本增量实现: + +- Stage 0:`s0.file_matrix_twelve_rows` 运行保留的 live matrix 脚本,要求 file + provider 上恰好十二个共享场景行全为 true;`s0.nokv_live_matrix` 要求 live + NoKV 栈上同样的行加 `restored_lineage_fails_closed` 全为 true,且 file/NoKV + 逐行结果一致。 +- Stage 1:`s1.cli_document_decodes_through_ts_store` 通过产品 CLI 写入三次 + observation(`todo add`、`task-lease acquire`、`todo update`),再经 + `FileAuthorityStore` 的 `loadAuthority`、分页 `scanCommitted` 与 + `readReceipt` 回读:cursor 为 `3`、三个 operation id 按序一致、首条 receipt + 可找到。 +- Stage 2A:`s2a.nokv_live_qualification` 对一个已存在的 workbench 以新铸的 + tenant/goal 运行已合并的 live 资格探针 + (`examples/nokv-authority-store/live-qualification.ts --execute-live`),要求 + `ok=true`、单节点 store conformance 范围、每项 check 通过、NoKV SDK `0.11.0` + / API `1`,且不宣称晋升或可用性。 +- Stage 2B:`s2b.postgresql_conformance_live` 在 node TAP reporter 下运行 + PostgreSQL 集成测试文件,要求至少九个 pass、零 fail、零 skip。 +- Stage 2C 观察基础:七个 `s2c1.*` 行移植本地 shadow CLI E2E 与迁移断言,并钉住单一 + lineage 保证。 + configure 往返先预览、再开启、回读、最后关闭 observer;每个 writer family + (handoff-mode、todo add/update/complete/supersede/capture-followups/ + archive-completed、task-lease acquire/renew/transfer)都以 + `primary_writeback_preserved`、`provider_to_local_writes=false`、 + `candidate_read_for_decision=false` 完成 capture,而幂等 re-acquire 不产生 + observation;default-off goal 保持隔离;候选失败不推翻主写;POSIX SIGKILL + 落在崩溃间隙时只丢失该次 observation;`--runtime-root` 与 `common_runtime_root` + 不同时,todo add、task-lease acquire、todo update、follow-up 捕获与带 lease 的 + complete 仍落入同一个 store identity,registry root 既不产生候选 lineage 也不 + 产生 lease 状态;`migrate-state` 在不携带 legacy 字节的前提下建立新 lineage。 + +Live 行按环境门控(`LOOPX_TEST_POSTGRES_URL`;`NOKV_COORDINATION_LIVE=1` 加 +`NOKV_*` 栈变量;`LOOPX_NOKV_AUTHORITY_LIVE=1` 加 `LOOPX_NOKV_AUTHORITY_*` 输入)。 +没有栈时它们报告 `unverified`,除非传入 `--allow-unverified`,否则 ladder 以非零 +退出;unverified 行永不计为 green。pending 行同样是未兑现的义务:选中它而不传 +`--allow-pending` 就非零退出,所以一份零执行的报告不可能显示为 green。 +报告绑定 LoopX commit、工作树是否 dirty、探针 digest 与经哈希的连接事实;其隐私 +扫描会把任何临时根目录、home 目录、连接 URL 或配置路径的泄露改写为 +`fail/privacy_violation`;仅出现在 bindings 块的泄露会被抹除并同样判定为失败, +`summary.privacy_violations` 阻止 green 退出,任何开关都不能放宽。 + +交付边界:test-only。没有任何生产入口构造任何 store;ladder 不新增产品路径, +只经保留的 TypeScript store 读取候选。Stage 2C parity 后半段 +(`s2c2.*`:outbox 条目、幂等 drain、drain 前与 drain 中的 SIGKILL、带 pending +条目的 rollback、parity 相等与分歧、迁移 seed-and-drain、增长 +度量)以 pending 行声明,而非宣称已完成。本小节记录的是上述阶段的可执行证据; +它不晋升任何 provider,也不完成 Stage 2C promotion。 + ### P0:合同与 deterministic proof - 本 ownership matrix 与显式 shared-mode boundary; diff --git a/examples/shared-goal-authority-e2e/README.md b/examples/shared-goal-authority-e2e/README.md new file mode 100644 index 0000000000..82cd454934 --- /dev/null +++ b/examples/shared-goal-authority-e2e/README.md @@ -0,0 +1,106 @@ +# Shared-goal-authority E2E stage ladder + +One incremental end-to-end ladder for +[RFC: LoopX shared control-plane authority and pluggable state providers v0](../../docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md). +Every completed RFC stage claim is one row; every row drives the product +through the real `python -m loopx.cli` (`real_cli`) or runs a retained +store-level probe (`store_direct`). The ladder adds no product path, reads the +candidate only through the production TypeScript `FileAuthorityStore`, and +never reports green while a selected row is unverified. + +```bash +python examples/shared-goal-authority-e2e/ladder.py # exit 1 here: live rows unverified, parity rows pending +python examples/shared-goal-authority-e2e/ladder.py --allow-unverified --allow-pending +python examples/shared-goal-authority-e2e/ladder.py --stage 2c1 --report-json ladder-report.json +python examples/shared-goal-authority-e2e/ladder.py --list +``` + +The pytest projection is `tests/control_plane/test_shared_goal_authority_e2e.py`; +there, an unverified row skips as `unverified: ` and a POSIX-only row +skips on Windows. Five `s2c1.*` rows whose assertions +`tests/control_plane/test_local_authority_shadow_cli_e2e.py` already pins +through the same product path (configure round trip, default-off isolation, +candidate failure, crash gap, dual runtime root) are skipped in the default CI +projection to stay within the pytest job budget; `LOOPX_LADDER_FULL=1` runs +them in pytest, and the example runner always runs every row. + +## Rows + +| Row | Stage | Path | Gate | Asserts | +| --- | --- | --- | --- | --- | +| `s0.file_matrix_twelve_rows` | 0 | store_direct | deterministic | `examples/nokv-shadow-provider/live_e2e.py` reports exactly the twelve known file-provider scenario rows, all true | +| `s0.nokv_live_matrix` | 0 | store_direct | env:nokv_legacy | the same twelve rows plus `restored_lineage_fails_closed` are true on a live NoKV stack and file/NoKV outcomes are identical | +| `s1.cli_document_decodes_through_ts_store` | 1 | real_cli | deterministic | three CLI writes (`todo add`, `task-lease acquire`, `todo update`) read back through `FileAuthorityStore`: `loadAuthority` loaded at cursor `3`, paged `scanCommitted` yields the three `observation_id`s in order, `readReceipt` finds the first | +| `s2a.nokv_live_qualification` | 2a | store_direct | env:nokv_authority | runs the merged `examples/nokv-authority-store/live-qualification.ts --execute-live` against an existing workbench with a fresh tenant/goal pair; requires `ok=true`, the single-node store-conformance scope, every check `passed`, NoKV SDK `0.11.0` / API `1`, and no promotion or availability claim; evidence carries check ids, counts, and config and workbench digest prefixes, never a configuration value or the workbench name | +| `s2b.postgresql_conformance_live` | 2b | store_direct | env:postgresql | `postgresql_authority_store.integration.test.ts` under node's TAP reporter: `# pass >= 9`, `# fail 0`, `# skipped 0` | +| `s2c1.configure_enable_disable_roundtrip` | 2c1 | real_cli | deterministic | `configure-goal` preview does not write, enable writes, captured observations for a todo and a lease, read-back summary `enabled/file_one_way`, disable writes and later writes neither observe nor touch candidate bytes | +| `s2c1.every_writer_family_captures` | 2c1 | real_cli | deterministic | handoff-mode set, todo add/update/complete/supersede/capture-followups/archive-completed, task-lease acquire/renew/transfer each carry `outcome in {captured, replayed, ambiguous_reconciled}`, `primary_writeback_preserved=true`, `provider_to_local_writes=false`, `candidate_read_for_decision=false`; an idempotent re-acquire carries no `authority_shadow`; candidate `cursor == captured count`, operation ids equal observation ids, no time-active lease in the head, head todos equal `todo list` | +| `s2c1.default_off_isolation` | 2c1 | real_cli | deterministic | a default-off goal returns the same response fields as an observed goal, carries no `authority_shadow`, and creates no `authority-shadow/` directory | +| `s2c1.candidate_failure_preserves_primary` | 2c1 | real_cli | deterministic | a blocked candidate directory yields `outcome=failed`, `reason_code=shadow_observation_failed`, and the committed todo is in the primary state | +| `s2c1.crash_gap_loses_observation` | 2c1 | real_cli | deterministic (POSIX) | a writer SIGKILLed while the observation lock is held commits its todo but leaves no candidate document; the next write captures the full two-todo snapshot without claiming an outbox or correlation | +| `s2c1.dual_runtime_root_consistency` | 2c1 | real_cli | deterministic | with `common_runtime_root` different from `--runtime-root`, todo add, task-lease acquire, todo update, capture-followups, and a leased completion all observe into one store identity; the head holds both todos and the released lease; the registry root gains neither a candidate lineage nor lease state | +| `s2c1.migration_seeds_new_lineage` | 2c1 | real_cli | deterministic | `migrate-state` dry run plans the seed without writing; execute seeds one fresh `file:` lineage at cursor `1` that carries no legacy identity, revision, source path, or private byte | + +Pending rows are declared in the report as `pending`, never counted as pass, +and they block a green exit unless `--allow-pending` is passed. The Stage 2C +parity rows are pending: `s2c2.outbox_prepared_then_committed_entries`, `s2c2.drain_idempotent`, +`s2c2.sigkill_between_primary_write_and_drain`, `s2c2.sigkill_mid_drain`, +`s2c2.rollback_with_pending_entries`, `s2c2.parity_equal`, `s2c2.parity_divergent_detects_foreign_edit`, +`s2c2.migration_seeds_and_drains`, `s2c2.growth_measurement_gate` (until the +Stage 2C parity PRs land). + +## Gates and environment variables + +| Gate | Requirement | Unverified reason when absent | +| --- | --- | --- | +| `deterministic` | none (needs `node` on `PATH` for the CLI's TypeScript runtime and the read-back probe) | `node_missing` when the probe cannot run | +| `env:postgresql` | `LOOPX_TEST_POSTGRES_URL` plus `node_modules/pg` (`npm ci`) | `postgres_url_missing`, `pg_dependency_missing`, `node_missing` | +| `env:nokv_legacy` | `NOKV_COORDINATION_LIVE=1` and `NOKV_ETCD`, `NOKV_ETCD_PREFIX`, `NOKV_ROOT_ID`, `NOKV_BUCKET`, `NOKV_OBJECT_ENDPOINT`, `NOKV_OBJECT_ROOT`, `NOKV_OBJECT_KEY`, `NOKV_OBJECT_SECRET`; the `nokv` SDK importable | `nokv_live_env_missing`, `nokv_coordination_live_not_enabled`, `nokv_sdk_missing` | +| `env:nokv_authority` | `LOOPX_NOKV_AUTHORITY_LIVE=1` (the probe writes durable test data), `LOOPX_NOKV_AUTHORITY_CONFIG_JSON` (absolute path to the ignored NoKV client configuration), `LOOPX_NOKV_AUTHORITY_PYTHON` (absolute path to the Python executable that resolves NoKV SDK 0.11.0), `LOOPX_NOKV_AUTHORITY_WORKBENCH` (an existing workbench); `node` on `PATH` | `nokv_authority_env_missing`, `loopx_nokv_authority_live_not_enabled`, `nokv_authority_config_missing`, `nokv_authority_python_missing`, `node_missing` | + +POSIX-only rows report `unverified/posix_only` on Windows. + +## Report and exit policy + +The report schema is `loopx_shared_goal_authority_e2e_report_v0`: +`rows[]` (`status in {pass, fail, unverified}`, `reason_code`, public-safe +`evidence`, `duration_ms`), `pending[]`, `summary{pass, fail, unverified, +pending, executed, privacy_violations}`, `bindings{loopx_commit, loopx_tree_dirty, probe_sha256[], +nokv_client_config_sha256, nokv_sdk_version, postgres_url_sha256_prefix, +pg_package_version}` (`null` when unknown), and `exit_policy`. + +Exit code is `0` iff `fail == 0` and `privacy_violations == 0` and +(`unverified == 0` or `--allow-unverified`) and (`pending == 0` or +`--allow-pending`): a selected +row that never executed, whether gated or declared pending, is an unmet +obligation, so `--row s2c2.parity_equal` exits 1 with zero executions, and a +mixed selection exits 1 even when its executable rows pass. `--list` only +prints the registry and never claims verification. A privacy scan runs over +the finished report: any occurrence of a temporary root, the home directory, +the repository path, the PostgreSQL URL, a NoKV configuration value, or a NoKV +authority input path rewrites that row to `fail/privacy_violation`; a leak +confined to the `bindings` block nulls every binding, marks +`bindings.privacy_violation`, and still exits `1` through +`summary.privacy_violations`, which no flag relaxes. Evidence therefore +carries counters, cursors, outcome tokens, and sha256 prefixes only. + +## Test seams later PRs must provide + +The pending `s2c2.*` rows will be implemented against these seams; a Stage 2C +parity PR that does not expose them cannot be ladder-verified: + +- a drain lock file at `/authority-shadow/outbox//drain` so the + ladder can hold the drain window with `loopx.file_lock.exclusive_file_lock` + exactly as it holds `/authority-shadow/file//observation` + today, then SIGKILL a writer before or during drain; +- one file per outbox entry under `/authority-shadow/outbox//` + with a prepared-then-committed marker, so pending entries are countable and + a rollback with pending entries is observable from disk; +- `drain` and `verify` product commands that emit JSON with `drained_count`, + `cursor_before`, `cursor_after`, `parity_verdict`, and the source and + candidate digests, so parity-equal and foreign-edit rows can assert on + typed fields rather than prose; +- the same commands must resolve the runtime root the way `todo` and + `task-lease` do (`effective_runtime_root`), so the one-lineage guarantee that + `s2c1.dual_runtime_root_consistency` proves for the observation hooks also + holds for drain and verify. diff --git a/examples/shared-goal-authority-e2e/ladder.py b/examples/shared-goal-authority-e2e/ladder.py new file mode 100644 index 0000000000..792eba4a0d --- /dev/null +++ b/examples/shared-goal-authority-e2e/ladder.py @@ -0,0 +1,22 @@ +#!/usr/bin/env python3 +"""Run the shared-goal-authority E2E stage ladder against this checkout. + +Thin entry point: the row registry, runners, report, and exit policy live in +``loopx.control_plane.testing.authority_e2e_ladder``. See README.md next to +this file for rows, gates, environment variables, and the exit policy. +""" + +from __future__ import annotations + +import sys +from pathlib import Path + +REPO_ROOT = Path(__file__).resolve().parents[2] +if str(REPO_ROOT) not in sys.path: + sys.path.insert(0, str(REPO_ROOT)) + +from loopx.control_plane.testing.authority_e2e_ladder import main # noqa: E402 + + +if __name__ == "__main__": + raise SystemExit(main(sys.argv[1:])) diff --git a/loopx/control_plane/testing/authority_e2e_fixtures.py b/loopx/control_plane/testing/authority_e2e_fixtures.py new file mode 100644 index 0000000000..4befa5abde --- /dev/null +++ b/loopx/control_plane/testing/authority_e2e_fixtures.py @@ -0,0 +1,631 @@ +"""Workspace and process fixtures for the shared-goal-authority E2E ladder. + +Every helper drives the product through ``python -m loopx.cli`` in a child +process, or reads candidate bytes back through the TypeScript +``FileAuthorityStore`` probe. Nothing here imports a LoopX writer, so the +ladder can never become a second authority over the local goal state. +""" + +from __future__ import annotations + +import json +import os +import shutil +import subprocess +import sys +import time +import uuid +from collections.abc import Callable, Iterator, Mapping, Sequence +from contextlib import contextmanager +from dataclasses import dataclass +from pathlib import Path +from typing import Any, Protocol + +from ...file_lock import exclusive_file_lock + +REPO_ROOT = Path(__file__).resolve().parents[3] +TS_READBACK_PROBE = Path("tests") / "control_plane_ts" / "authority_store_readback_probe.ts" +DEFAULT_REGISTERED_AGENTS: tuple[str, ...] = ("agent-a", "agent-b") +RUNTIME_ROOT_BINDINGS: tuple[str, ...] = ("registry", "cli_override", "cli_override_divergent") +HANDOFF_MODES: tuple[str, ...] = ("legacy", "soft_claim", "hard_lease") +LOCAL_AUTHORITY_SHADOW_CONFIG = { + "schema_version": "loopx_local_authority_shadow_config_v0", + "mode": "file_one_way", +} + +JsonObject = dict[str, Any] + + +class CliOutputError(AssertionError): + """The CLI did not print exactly one JSON object.""" + + +class CliCommandError(AssertionError): + """The CLI exited non-zero while the row expected a committed response.""" + + def __init__( + self, + *, + command: Sequence[str], + returncode: int, + payload: Mapping[str, Any], + ) -> None: + verb = " ".join(command[:2]) + super().__init__( + f"loopx {verb} exited {returncode}: " + f"error_code={payload.get('error_code')!r} error={payload.get('error')!r}" + ) + self.command = tuple(command) + self.returncode = returncode + self.payload = dict(payload) + + +class ProbeError(AssertionError): + """The TypeScript read-back probe did not complete.""" + + +class CliWorkspace(Protocol): + """Anything the CLI runners can address: a home plus global CLI arguments.""" + + @property + def home(self) -> Path: ... + + def cli_prefix(self) -> list[str]: ... + + +@dataclass(frozen=True) +class GoalWorkspace: + """One registered goal with its own registry, repo, runtime root, and home.""" + + goal_id: str + root: Path + repo: Path + registry_path: Path + state_path: Path + runtime_root: Path + home: Path + runtime_root_binding: str + registry_runtime_root: Path + + @property + def shadow_directory(self) -> Path: + return self.runtime_root / "authority-shadow" / "file" / self.goal_id + + @property + def observation_lock_target(self) -> Path: + return self.shadow_directory / "observation" + + def cli_prefix(self) -> list[str]: + prefix = ["--registry", str(self.registry_path)] + if self.runtime_root_binding in ("cli_override", "cli_override_divergent"): + prefix.extend(["--runtime-root", str(self.runtime_root)]) + prefix.extend(["--format", "json"]) + return prefix + + +@dataclass(frozen=True) +class LegacyMigrationSource: + """A legacy registry/runtime pair whose old shadow lineage must never migrate.""" + + old_goal_id: str + new_goal_id: str + old_store_identity: str + legacy_revision: str + private_marker: str + legacy_registry: Path + target_registry: Path + legacy_runtime: Path + target_runtime: Path + source_repo: Path + target_repo: Path + home: Path + + @property + def target_shadow_directory(self) -> Path: + return self.target_runtime / "authority-shadow" / "file" / self.new_goal_id + + def cli_prefix(self) -> list[str]: + return ["--registry", str(self.target_registry), "--format", "json"] + + +@dataclass(frozen=True) +class CandidateDocument: + """Stable fields of the single ``authority-store-*.json`` candidate document.""" + + path: Path + document: JsonObject + + @property + def cursor(self) -> str: + return str(self.document.get("cursor")) + + @property + def store_identity(self) -> str: + return str(self.document.get("store_identity")) + + @property + def head(self) -> JsonObject: + head = self.document.get("head") + return dict(head) if isinstance(head, dict) else {} + + @property + def operation_ids(self) -> list[str]: + committed = self.document.get("committed") + if not isinstance(committed, list): + return [] + return [ + str(entry.get("operation_id")) + for entry in committed + if isinstance(entry, dict) + ] + + @property + def todo_ids(self) -> list[str]: + return [ + str(todo.get("todo_id")) + for todo in self.head.get("todos") or [] + if isinstance(todo, dict) + ] + + @property + def leases(self) -> list[JsonObject]: + return [ + dict(lease) + for lease in self.head.get("leases") or [] + if isinstance(lease, dict) + ] + + +@dataclass(frozen=True) +class TapSummary: + """The ``# pass`` / ``# fail`` / ``# skipped`` trailer of a node TAP run.""" + + returncode: int + tests: int | None + passed: int | None + failed: int | None + skipped: int | None + + +def unique_goal_id(prefix: str) -> str: + """Return a single-segment goal id that is unique across xdist workers.""" + + return f"{prefix}-{uuid.uuid4().hex[:10]}" + + +def parse_json_object(text: str) -> JsonObject: + """Decode one JSON object from CLI stdout; anything else is a contract break.""" + + try: + decoded = json.loads(text) + except json.JSONDecodeError as exc: + raise CliOutputError(f"CLI output is not JSON: {exc.msg}") from None + if not isinstance(decoded, dict): + raise CliOutputError("CLI output is not a JSON object") + return {str(key): value for key, value in decoded.items()} + + +def _write_active_state(state_path: Path, *, goal_id: str, handoff_mode: str) -> None: + state_path.write_text( + "---\n" + f"goal_id: {goal_id}\n" + f"handoff_mode: {handoff_mode}\n" + "updated_at: 2026-09-02T00:00:00+00:00\n" + "---\n\n" + "## Agent Todo\n\n", + encoding="utf-8", + ) + + +def build_goal_workspace( + root: Path, + *, + goal_id: str, + handoff_mode: str = "legacy", + shadow_enabled: bool = False, + runtime_root_binding: str = "registry", + registered_agents: Sequence[str] = DEFAULT_REGISTERED_AGENTS, +) -> GoalWorkspace: + """Materialize one goal exactly as the local-shadow CLI E2E fixture does. + + ``runtime_root_binding`` selects how the CLI learns the runtime root: + ``registry`` relies on ``common_runtime_root`` alone, ``cli_override`` also + passes the same directory as ``--runtime-root``, and + ``cli_override_divergent`` registers a different ``common_runtime_root`` + than the ``--runtime-root`` override so a row can prove that every writer + hook of one CLI call shares the override root. + """ + + if handoff_mode not in HANDOFF_MODES: + raise ValueError(f"unsupported handoff_mode {handoff_mode!r}") + if runtime_root_binding not in RUNTIME_ROOT_BINDINGS: + raise ValueError(f"unsupported runtime_root_binding {runtime_root_binding!r}") + repo = root / goal_id + repo.mkdir() + state_path = repo / "ACTIVE_GOAL_STATE.md" + _write_active_state(state_path, goal_id=goal_id, handoff_mode=handoff_mode) + runtime_root = root / f"{goal_id}-runtime" + registry_runtime_root = ( + root / f"{goal_id}-registry-runtime" + if runtime_root_binding == "cli_override_divergent" + else runtime_root + ) + home = root / f"{goal_id}-home" + home.mkdir() + coordination: JsonObject = { + "agent_model": "peer_v1", + "registered_agents": list(registered_agents), + } + if shadow_enabled: + coordination["authority_shadow"] = dict(LOCAL_AUTHORITY_SHADOW_CONFIG) + registry_path = root / f"{goal_id}-registry.json" + registry_path.write_text( + json.dumps( + { + "common_runtime_root": str(registry_runtime_root), + "goals": [ + { + "id": goal_id, + "status": "active", + "repo": str(repo), + "state_file": state_path.name, + "coordination": coordination, + } + ], + } + ), + encoding="utf-8", + ) + return GoalWorkspace( + goal_id=goal_id, + root=root, + repo=repo, + registry_path=registry_path, + state_path=state_path, + runtime_root=runtime_root, + home=home, + runtime_root_binding=runtime_root_binding, + registry_runtime_root=registry_runtime_root, + ) + + +def build_legacy_migration_source( + root: Path, + *, + old_goal_id: str, + new_goal_id: str, + old_store_identity: str = "file:11111111111111111111111111111111", + registered_agents: Sequence[str] = DEFAULT_REGISTERED_AGENTS, +) -> LegacyMigrationSource: + """Lift the state-migration shadow fixture: a legacy goal with a stale lineage.""" + + legacy_runtime = root / "legacy-runtime" + target_runtime = root / "target-runtime" + source_repo = root / "legacy-repo" + target_repo = root / "target-repo" + home = root / "migration-home" + for directory in (source_repo, target_repo, home): + directory.mkdir() + source_state = source_repo / "ACTIVE_GOAL_STATE.md" + source_state.write_text( + "---\n" + f"goal_id: {old_goal_id}\n" + "handoff_mode: soft_claim\n" + "updated_at: 2026-09-02T00:00:00+10:00\n" + "---\n\n" + "## Agent Todo\n\n" + "- [ ] Preserve the new local authority only.\n", + encoding="utf-8", + ) + legacy_registry = root / "legacy-registry.json" + legacy_registry.write_text( + json.dumps( + { + "schema_version": "0.1", + "common_runtime_root": str(legacy_runtime), + "goals": [ + { + "id": old_goal_id, + "status": "active", + "repo": str(source_repo), + "state_file": source_state.name, + "coordination": { + "agent_model": "peer_v1", + "registered_agents": list(registered_agents), + "authority_shadow": dict(LOCAL_AUTHORITY_SHADOW_CONFIG), + }, + } + ], + } + ), + encoding="utf-8", + ) + lease_dir = legacy_runtime / "goals" / old_goal_id / "task-leases" + lease_dir.mkdir(parents=True) + (lease_dir / "safe-local.json").write_text( + json.dumps( + { + "goal_id": old_goal_id, + "todo_id": "safe-local", + "owner": registered_agents[0], + "version": 1, + "lease_epoch": 1, + "status": "released", + } + ), + encoding="utf-8", + ) + legacy_revision = "file:99:legacy-lineage" + private_marker = "must-never-migrate" + source_shadow = legacy_runtime / "authority-shadow" / "file" / old_goal_id + source_shadow.mkdir(parents=True) + (source_shadow / "store-identity").write_text(old_store_identity, encoding="utf-8") + (source_shadow / "authority-store-legacy.json").write_text( + json.dumps( + { + "goal_id": old_goal_id, + "store_identity": old_store_identity, + "provider_revision": legacy_revision, + "cursor": "99", + "private_provider_byte": private_marker, + "source_path": str(source_repo), + } + ), + encoding="utf-8", + ) + return LegacyMigrationSource( + old_goal_id=old_goal_id, + new_goal_id=new_goal_id, + old_store_identity=old_store_identity, + legacy_revision=legacy_revision, + private_marker=private_marker, + legacy_registry=legacy_registry, + target_registry=root / "target-registry.json", + legacy_runtime=legacy_runtime, + target_runtime=target_runtime, + source_repo=source_repo, + target_repo=target_repo, + home=home, + ) + + +def cli_env(workspace: CliWorkspace) -> dict[str, str]: + """Child environment: this checkout on ``PYTHONPATH`` and an isolated home.""" + + env = os.environ.copy() + env["PYTHONPATH"] = str(REPO_ROOT) + env["HOME"] = str(workspace.home) + if os.name == "nt": + env["USERPROFILE"] = str(workspace.home) + return env + + +def cli_command(workspace: CliWorkspace, *args: str) -> list[str]: + return [sys.executable, "-m", "loopx.cli", *workspace.cli_prefix(), *args] + + +def run_cli( + workspace: CliWorkspace, + *args: str, + timeout: float = 60, + check: bool = True, +) -> JsonObject: + """Run one product CLI command and return its JSON object response.""" + + command = cli_command(workspace, *args) + completed = subprocess.run( + command, + cwd=REPO_ROOT, + env=cli_env(workspace), + capture_output=True, + text=True, + timeout=timeout, + check=False, + ) + payload = parse_json_object(completed.stdout) + if check and completed.returncode != 0: + raise CliCommandError( + command=args, + returncode=completed.returncode, + payload=payload, + ) + return payload + + +def spawn_cli(workspace: CliWorkspace, *args: str) -> subprocess.Popen[str]: + """Start one product CLI command without waiting for it.""" + + return subprocess.Popen( + cli_command(workspace, *args), + cwd=REPO_ROOT, + env=cli_env(workspace), + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + ) + + +def wait_until( + predicate: Callable[[], bool], + timeout: float, + *, + interval: float = 0.01, +) -> bool: + """Poll ``predicate`` until it holds or ``timeout`` seconds elapse.""" + + deadline = time.monotonic() + timeout + while True: + if predicate(): + return True + if time.monotonic() >= deadline: + return False + time.sleep(interval) + + +def kill_now(process: subprocess.Popen[str]) -> None: + """SIGKILL (or TerminateProcess) the child and reap it.""" + + process.kill() + process.communicate(timeout=5) + + +@contextmanager +def hold_observation_lock(workspace: GoalWorkspace) -> Iterator[Path]: + """Hold the observer's own lock so a primary commit cannot be observed.""" + + with exclusive_file_lock( + workspace.observation_lock_target, + operation="e2e_window", + ) as lock_path: + yield lock_path + + +def candidate_store_paths(workspace: GoalWorkspace) -> list[Path]: + return sorted(workspace.shadow_directory.glob("authority-store-*.json")) + + +def candidate_document(workspace: GoalWorkspace) -> CandidateDocument: + """Return the single candidate document; zero or several is a failure.""" + + paths = candidate_store_paths(workspace) + if len(paths) != 1: + raise AssertionError(f"expected exactly one candidate document, found {len(paths)}") + return CandidateDocument( + path=paths[0], + document=parse_json_object(paths[0].read_text(encoding="utf-8")), + ) + + +def node_executable() -> str | None: + return shutil.which("node") + + +def ts_readback( + workspace: GoalWorkspace, + *, + receipt: str | None = None, + page_size: int = 2, +) -> JsonObject | None: + """Read the candidate back through ``FileAuthorityStore``; ``None`` without node.""" + + node = node_executable() + if node is None: + return None + command = [ + node, + "--no-warnings", + "--experimental-strip-types", + str(REPO_ROOT / TS_READBACK_PROBE), + "--directory", + str(workspace.shadow_directory), + "--goal-id", + workspace.goal_id, + "--page-size", + str(page_size), + ] + if receipt is not None: + command.extend(["--receipt", receipt]) + completed = subprocess.run( + command, + cwd=REPO_ROOT, + capture_output=True, + text=True, + timeout=60, + check=False, + ) + if completed.returncode != 0: + raise ProbeError(f"read-back probe exited {completed.returncode}") + return parse_json_object(completed.stdout) + + +def _tap_counter(line: str, label: str) -> int | None: + prefix = f"# {label} " + if not line.startswith(prefix): + return None + try: + return int(line[len(prefix):].strip()) + except ValueError: + return None + + +def parse_tap_summary(output: str, *, returncode: int) -> TapSummary: + """Extract the node TAP reporter trailer counters.""" + + counters: dict[str, int | None] = { + "tests": None, + "pass": None, + "fail": None, + "skipped": None, + } + aliases = {"tests": ("tests",), "pass": ("pass",), "fail": ("fail",), "skipped": ("skipped", "skip")} + for raw in output.splitlines(): + line = raw.strip() + for key, labels in aliases.items(): + for label in labels: + value = _tap_counter(line, label) + if value is not None: + counters[key] = value + return TapSummary( + returncode=returncode, + tests=counters["tests"], + passed=counters["pass"], + failed=counters["fail"], + skipped=counters["skipped"], + ) + + +def tap_summary( + argv: Sequence[str], + *, + cwd: Path = REPO_ROOT, + env: Mapping[str, str] | None = None, + timeout: float = 600, +) -> TapSummary: + """Run a node TAP command and summarize its trailer.""" + + completed = subprocess.run( + list(argv), + cwd=cwd, + env=dict(env) if env is not None else None, + capture_output=True, + text=True, + timeout=timeout, + check=False, + ) + return parse_tap_summary(completed.stdout, returncode=completed.returncode) + + +__all__ = [ + "CandidateDocument", + "CliCommandError", + "CliOutputError", + "CliWorkspace", + "DEFAULT_REGISTERED_AGENTS", + "GoalWorkspace", + "HANDOFF_MODES", + "JsonObject", + "LOCAL_AUTHORITY_SHADOW_CONFIG", + "LegacyMigrationSource", + "ProbeError", + "REPO_ROOT", + "RUNTIME_ROOT_BINDINGS", + "TS_READBACK_PROBE", + "TapSummary", + "build_goal_workspace", + "build_legacy_migration_source", + "candidate_document", + "candidate_store_paths", + "cli_command", + "cli_env", + "hold_observation_lock", + "kill_now", + "node_executable", + "parse_json_object", + "parse_tap_summary", + "run_cli", + "spawn_cli", + "tap_summary", + "ts_readback", + "unique_goal_id", + "wait_until", +] diff --git a/loopx/control_plane/testing/authority_e2e_ladder.py b/loopx/control_plane/testing/authority_e2e_ladder.py new file mode 100644 index 0000000000..6d8a9af8c4 --- /dev/null +++ b/loopx/control_plane/testing/authority_e2e_ladder.py @@ -0,0 +1,1144 @@ +"""Incremental end-to-end stage ladder for the shared-goal-authority RFC. + +One ladder, one row per completed RFC stage claim, one exit policy: the run +is green only when every selected row passed. A row that cannot run in this +environment (no PostgreSQL, no NoKV stack, no POSIX signals) reports +``unverified`` and, unless ``--allow-unverified`` is given, the ladder exits +non-zero. Rows for stages that later PRs will land are declared as +``pending`` so the report never silently claims them. + +Rows drive the product through ``python -m loopx.cli`` (``real_cli``) or run +the retained store-level probes (``store_direct``); this module adds no +product path of its own. +""" + +from __future__ import annotations + +import argparse +import json +import os +import subprocess +import sys +import tempfile +import time +from collections.abc import Callable, Iterable, Mapping, Sequence +from dataclasses import dataclass +from datetime import datetime, timezone +from importlib import metadata +from pathlib import Path + +from .authority_e2e_row_support import ( + AGENT_A, + RowAssertionError, + RowContext, + RowOutcome, + acquire_lease, + add_todo, + committed_observation, + expect, + passed, + sha256_hex, + unverified, +) +from .authority_e2e_rows_stage2c import ( + row_candidate_failure_preserves_primary, + row_configure_enable_disable_roundtrip, + row_crash_gap_loses_observation, + row_default_off_isolation, + row_dual_runtime_root_consistency, + row_every_writer_family_captures, + row_migration_seeds_new_lineage, +) +from .authority_e2e_fixtures import ( + REPO_ROOT, + CliOutputError, + TS_READBACK_PROBE, + JsonObject, + build_goal_workspace, + candidate_document, + node_executable, + parse_json_object, + run_cli, + tap_summary, + ts_readback, + unique_goal_id, +) + +REPORT_SCHEMA = "loopx_shared_goal_authority_e2e_report_v0" +LIST_SCHEMA = "loopx_shared_goal_authority_e2e_rows_v0" +STAGES: tuple[str, ...] = ("0", "1", "2a", "2b", "2c1", "2c2") +PRODUCT_PATHS: tuple[str, ...] = ("real_cli", "store_direct") +GATES: tuple[str, ...] = ( + "deterministic", + "env:postgresql", + "env:nokv_authority", + "env:nokv_legacy", +) +ROW_STATUSES: tuple[str, ...] = ("pass", "fail", "unverified") +EXIT_POLICY_RULE = ( + "exit 0 iff fail == 0 and privacy_violations == 0 " + "and (unverified == 0 or allow_unverified) and (pending == 0 or allow_pending)" +) + +POSTGRES_URL_VARIABLE = "LOOPX_TEST_POSTGRES_URL" +NOKV_LIVE_FLAG = "NOKV_COORDINATION_LIVE" +NOKV_STACK_VARIABLES: tuple[str, ...] = ( + "NOKV_ETCD", + "NOKV_ETCD_PREFIX", + "NOKV_ROOT_ID", + "NOKV_BUCKET", + "NOKV_OBJECT_ENDPOINT", + "NOKV_OBJECT_ROOT", + "NOKV_OBJECT_KEY", + "NOKV_OBJECT_SECRET", +) +NOKV_SECRET_VARIABLES: tuple[str, ...] = ("NOKV_OBJECT_KEY", "NOKV_OBJECT_SECRET") +# Stage 2A qualification inputs: the probe writes durable test data into an +# existing workbench, so it needs an explicit opt-in flag plus the ignored +# client configuration file, the Python executable that resolves the qualified +# NoKV SDK, and the workbench name. Tenant and goal ids are minted per run. +NOKV_AUTHORITY_LIVE_FLAG = "LOOPX_NOKV_AUTHORITY_LIVE" +NOKV_AUTHORITY_CONFIG_VARIABLE = "LOOPX_NOKV_AUTHORITY_CONFIG_JSON" +NOKV_AUTHORITY_PYTHON_VARIABLE = "LOOPX_NOKV_AUTHORITY_PYTHON" +NOKV_AUTHORITY_WORKBENCH_VARIABLE = "LOOPX_NOKV_AUTHORITY_WORKBENCH" +NOKV_AUTHORITY_VARIABLES: tuple[str, ...] = ( + NOKV_AUTHORITY_CONFIG_VARIABLE, + NOKV_AUTHORITY_PYTHON_VARIABLE, + NOKV_AUTHORITY_WORKBENCH_VARIABLE, +) +LIVE_OPT_IN_FLAGS: tuple[str, ...] = (NOKV_LIVE_FLAG, NOKV_AUTHORITY_LIVE_FLAG) +GATE_REQUIREMENTS: dict[str, tuple[str, ...]] = { + "deterministic": (), + "env:postgresql": (POSTGRES_URL_VARIABLE,), + "env:nokv_authority": (NOKV_AUTHORITY_LIVE_FLAG, *NOKV_AUTHORITY_VARIABLES), + "env:nokv_legacy": (NOKV_LIVE_FLAG, *NOKV_STACK_VARIABLES), +} +GATE_UNVERIFIED_REASON: dict[str, str] = { + "env:postgresql": "postgres_url_missing", + "env:nokv_authority": "nokv_authority_env_missing", + "env:nokv_legacy": "nokv_live_env_missing", +} + +LIVE_E2E_SCRIPT = Path("examples") / "nokv-shadow-provider" / "live_e2e.py" +PG_INTEGRATION_TEST = ( + Path("tests") / "control_plane_ts" / "postgresql_authority_store.integration.test.ts" +) +NOKV_QUALIFICATION_SCRIPT = Path("examples") / "nokv-authority-store" / "live-qualification.ts" +NOKV_HELPER = Path("loopx") / "control_plane" / "coordination" / "nokv_jsonl_helper.py" +NOKV_QUALIFICATION_REPORT_SCHEMA = "loopx_nokv_authority_live_qualification_v0" +NOKV_QUALIFICATION_SCOPE = "stage_2a_single_node_store_conformance" +QUALIFIED_NOKV_SDK_VERSION = "0.11.0" +QUALIFIED_NOKV_API_VERSION = 1 +PROBE_SOURCES: tuple[Path, ...] = ( + LIVE_E2E_SCRIPT, + TS_READBACK_PROBE, + PG_INTEGRATION_TEST, + NOKV_QUALIFICATION_SCRIPT, + NOKV_HELPER, + Path("loopx") / "control_plane" / "testing" / "authority_e2e_ladder.py", + Path("loopx") / "control_plane" / "testing" / "authority_e2e_fixtures.py", + Path("loopx") / "control_plane" / "testing" / "authority_e2e_row_support.py", + Path("loopx") / "control_plane" / "testing" / "authority_e2e_rows_stage2c.py", +) +FILE_MATRIX_ROWS: tuple[str, ...] = ( + "same_todo_one_winner", + "independent_todo_applies", + "replay_returns_original_receipt", + "identity_mismatch_rejected", + "stale_revision_conflicts", + "lost_response_recovers_receipt", + "receipts_retained", + "authority_revision_advanced_twice", + "renew_extends_the_active_lease", + "expired_lease_reclaimed_with_new_epoch", + "superseded_executor_cannot_write_back", + "complete_creates_claimable_successor_atomically", +) +NOKV_ONLY_MATRIX_ROW = "restored_lineage_fails_closed" +MINIMUM_POSTGRES_TAP_PASSES = 9 +@dataclass(frozen=True) +class LadderRow: + id: str + stage: str + title: str + product_path: str + gate: str + posix_only: bool + run: Callable[[RowContext], RowOutcome] + + def __post_init__(self) -> None: + if self.stage not in STAGES: + raise ValueError(f"row {self.id}: unsupported stage {self.stage!r}") + if self.product_path not in PRODUCT_PATHS: + raise ValueError(f"row {self.id}: unsupported product_path {self.product_path!r}") + if self.gate not in GATES: + raise ValueError(f"row {self.id}: unsupported gate {self.gate!r}") + + def describe(self) -> JsonObject: + return { + "id": self.id, + "stage": self.stage, + "title": self.title, + "product_path": self.product_path, + "gate": self.gate, + "posix_only": self.posix_only, + } + + +@dataclass(frozen=True) +class PendingRow: + id: str + stage: str + pending_until: str + + def __post_init__(self) -> None: + if self.stage not in STAGES: + raise ValueError(f"pending row {self.id}: unsupported stage {self.stage!r}") + + def describe(self) -> JsonObject: + return { + "id": self.id, + "stage": self.stage, + "status": "pending", + "pending_until": self.pending_until, + } + + +@dataclass(frozen=True) +class RowResult: + row: LadderRow + status: str + reason_code: str | None + evidence: JsonObject + duration_ms: int + + def as_dict(self) -> JsonObject: + return { + **self.row.describe(), + "status": self.status, + "reason_code": self.reason_code, + "evidence": dict(self.evidence), + "duration_ms": self.duration_ms, + } + + +def _run_live_matrix_script(environ: Mapping[str, str], *, live: bool) -> JsonObject: + env = dict(environ) + env["PYTHONPATH"] = str(REPO_ROOT) + if not live: + env.pop(NOKV_LIVE_FLAG, None) + completed = subprocess.run( + [sys.executable, str(REPO_ROOT / LIVE_E2E_SCRIPT)], + cwd=REPO_ROOT, + env=env, + capture_output=True, + text=True, + timeout=900, + check=False, + ) + matrix = parse_json_object(completed.stdout) + matrix["_exit_code"] = completed.returncode + return matrix + + +def _matrix_rows(matrix: Mapping[str, object], key: str) -> dict[str, object]: + rows = matrix.get(key) + expect(isinstance(rows, dict), f"live matrix must report {key}") + assert isinstance(rows, dict) + return {str(name): value for name, value in rows.items()} + + +def _false_rows(rows: Mapping[str, object]) -> list[str]: + return sorted(name for name, value in rows.items() if value is not True) + + +# --------------------------------------------------------------------------- +# Stage 0: recoverable reference foundation (store_direct) +# --------------------------------------------------------------------------- + + +def _row_file_matrix_twelve_rows(context: RowContext) -> RowOutcome: + matrix = _run_live_matrix_script(context.environ, live=False) + file_rows = _matrix_rows(matrix, "file_provider") + expect( + set(file_rows) == set(FILE_MATRIX_ROWS), + "file provider matrix must contain exactly the twelve known rows", + ) + expect(not _false_rows(file_rows), "every file provider matrix row must be true") + expect(matrix["_exit_code"] == 0, "live matrix script must exit 0 without a stack") + return passed(matrix_rows=len(file_rows), script_exit_code=matrix["_exit_code"]) + + +def _row_nokv_live_matrix(context: RowContext) -> RowOutcome: + matrix = _run_live_matrix_script(context.environ, live=True) + nokv_rows = _matrix_rows(matrix, "nokv_provider") + if "unverified" in nokv_rows: + reason = str(nokv_rows["unverified"]) + code = "nokv_sdk_missing" if "SDK" in reason else "nokv_matrix_unverified" + return unverified(code) + expected = {*FILE_MATRIX_ROWS, NOKV_ONLY_MATRIX_ROW} + expect(set(nokv_rows) == expected, "NoKV matrix must contain the shared rows plus the lineage row") + expect(not _false_rows(nokv_rows), "every NoKV matrix row must be true") + parity = _matrix_rows(matrix, "file_nokv_parity") + expect(parity.get("identical_row_outcomes") is True, "file and NoKV rows must be identical") + expect(parity.get("rows") == len(FILE_MATRIX_ROWS), "parity must cover the twelve shared rows") + expect(matrix["_exit_code"] == 0, "live matrix script must exit 0") + return passed( + nokv_rows=len(nokv_rows), + parity_rows=parity.get("rows"), + restored_lineage_fails_closed=True, + script_exit_code=matrix["_exit_code"], + ) + + +# --------------------------------------------------------------------------- +# Stage 1: provider-neutral boundary, read back through the TypeScript store +# --------------------------------------------------------------------------- + + +def _row_cli_document_decodes_through_ts_store(context: RowContext) -> RowOutcome: + if node_executable() is None: + return unverified("node_missing") + workspace = build_goal_workspace( + context.root, + goal_id=unique_goal_id("ladder-s1"), + handoff_mode="hard_lease", + shadow_enabled=True, + runtime_root_binding="registry", + ) + added = add_todo(workspace, "Decode this observation through the TypeScript store.") + todo_id = str(added["todo_id"]) + acquired = acquire_lease( + workspace, + todo_id=todo_id, + owner=AGENT_A, + idempotency_key="ladder-s1-lease", + ) + updated = run_cli( + workspace, + "todo", + "update", + "--goal-id", + workspace.goal_id, + "--todo-id", + todo_id, + "--note", + "A public-safe note recorded after the lease.", + "--agent-id", + AGENT_A, + ) + observations = [ + committed_observation(payload, label=label) + for label, payload in (("todo add", added), ("task-lease acquire", acquired), ("todo update", updated)) + ] + expect( + all(evidence["outcome"] == "captured" for evidence in observations), + "three distinct CLI writes must each be captured", + ) + observation_ids = [str(evidence["observation_id"]) for evidence in observations] + probe = ts_readback(workspace, receipt=observation_ids[0]) + expect(probe is not None, "read-back probe requires node") + assert probe is not None + load = probe.get("load") + scan = probe.get("scan") + receipt = probe.get("receipt") + expect(isinstance(load, dict) and load.get("status") == "loaded", "TS store must load the CLI document") + expect(isinstance(load, dict) and load.get("cursor") == "3", "TS store cursor must be 3 after three writes") + expect( + isinstance(load, dict) + and load.get("provider_revision") == observations[-1]["provider_revision"], + "TS store head revision must equal the last observation revision", + ) + expect( + isinstance(scan, dict) and scan.get("operation_ids") == observation_ids, + "scanCommitted must page through the three observation ids in order", + ) + expect(isinstance(receipt, dict) and receipt.get("status") == "found", "readReceipt must find the first observation") + document = candidate_document(workspace) + expect(document.cursor == "3" and document.operation_ids == observation_ids, "candidate bytes must match the probe") + assert isinstance(scan, dict) + return passed(cursor="3", scan_pages=scan.get("pages"), observation_count=len(observation_ids)) + + +# --------------------------------------------------------------------------- +# Stage 2A: NoKV candidate conformance against a live single-node stack +# --------------------------------------------------------------------------- + + +def _absolute_existing_path(value: str | None, *, must_be_file: bool) -> Path | None: + if not value: + return None + path = Path(value) + if not path.is_absolute(): + return None + if must_be_file and not path.is_file(): + return None + if not must_be_file and not path.exists(): + return None + return path + + +def _nokv_authority_config_sha256(path: Path) -> str: + """Digest of the canonical client configuration; the values never leave the file.""" + + document = json.loads(path.read_text(encoding="utf-8")) + return sha256_hex(json.dumps(document, sort_keys=True, separators=(",", ":"))) + + +def _row_nokv_live_qualification(context: RowContext) -> RowOutcome: + node = node_executable() + if node is None: + return unverified("node_missing") + config_path = _absolute_existing_path( + context.environ.get(NOKV_AUTHORITY_CONFIG_VARIABLE), must_be_file=True + ) + if config_path is None: + return unverified("nokv_authority_config_missing", variable=NOKV_AUTHORITY_CONFIG_VARIABLE) + python = _absolute_existing_path( + context.environ.get(NOKV_AUTHORITY_PYTHON_VARIABLE), must_be_file=True + ) + if python is None: + return unverified("nokv_authority_python_missing", variable=NOKV_AUTHORITY_PYTHON_VARIABLE) + workbench = str(context.environ.get(NOKV_AUTHORITY_WORKBENCH_VARIABLE) or "") + tenant_id = unique_goal_id("ladder-tenant") + goal_id = unique_goal_id("ladder-2a") + completed = subprocess.run( + [ + node, + "--no-warnings", + "--experimental-strip-types", + str(REPO_ROOT / NOKV_QUALIFICATION_SCRIPT), + "--execute-live", + "--config-json", + str(config_path), + "--python-executable", + str(python), + "--tenant-id", + tenant_id, + "--goal-id", + goal_id, + "--workbench", + workbench, + ], + cwd=REPO_ROOT, + env=dict(context.environ), + capture_output=True, + text=True, + timeout=600, + check=False, + ) + if completed.returncode != 0: + failure: JsonObject = {} + try: + failure = parse_json_object(completed.stderr.strip().splitlines()[-1]) + except (CliOutputError, IndexError): + pass + raise RowAssertionError( + f"qualification probe exited {completed.returncode}: " + f"{failure.get('reason_code') or 'no typed failure on stderr'}" + ) + report = parse_json_object(completed.stdout) + expect( + report.get("schema_version") == NOKV_QUALIFICATION_REPORT_SCHEMA + and report.get("ok") is True, + "qualification report must carry the reviewed schema and ok=true", + ) + expect( + report.get("qualification_scope") == NOKV_QUALIFICATION_SCOPE, + "qualification scope must be single-node store conformance", + ) + checks = report.get("checks") + expect(isinstance(checks, list) and len(checks) > 0, "qualification must report its checks") + assert isinstance(checks, list) + check_ids = [str(check.get("id")) for check in checks if isinstance(check, dict)] + expect( + len(check_ids) == len(checks) + and all(check.get("status") == "passed" for check in checks if isinstance(check, dict)), + "every qualification check must have passed", + ) + expect( + report.get("nokv_sdk_version") == QUALIFIED_NOKV_SDK_VERSION + and report.get("nokv_api_version") == QUALIFIED_NOKV_API_VERSION, + "qualification must name the qualified NoKV SDK and API versions", + ) + expect( + report.get("authority_source_changed") is False + and report.get("availability_or_ha_proven") is False, + "qualification must not claim promotion or availability", + ) + return passed( + qualification_scope=NOKV_QUALIFICATION_SCOPE, + check_count=len(check_ids), + check_ids=check_ids, + final_generation=report.get("final_generation"), + final_cursor=report.get("final_cursor"), + nokv_sdk_version=report.get("nokv_sdk_version"), + nokv_api_version=report.get("nokv_api_version"), + config_sha256_prefix=_nokv_authority_config_sha256(config_path)[:12], + workbench_sha256_prefix=sha256_hex(workbench)[:12], + tenant_id=tenant_id, + goal_id=goal_id, + durable_test_data_left=report.get("durable_test_data_left") is True, + ) + + +# --------------------------------------------------------------------------- +# Stage 2B: PostgreSQL candidate conformance against a live database +# --------------------------------------------------------------------------- + + +def _pg_package_version() -> str | None: + manifest = REPO_ROOT / "node_modules" / "pg" / "package.json" + if not manifest.exists(): + return None + try: + version = json.loads(manifest.read_text(encoding="utf-8")).get("version") + except (OSError, ValueError, AttributeError): + return None + return str(version) if version else None + + +def _row_postgresql_conformance_live(context: RowContext) -> RowOutcome: + node = node_executable() + if node is None: + return unverified("node_missing") + pg_version = _pg_package_version() + if pg_version is None: + return unverified("pg_dependency_missing") + summary = tap_summary( + [ + node, + "--no-warnings", + "--experimental-strip-types", + "--test", + "--test-reporter=tap", + str(PG_INTEGRATION_TEST), + ], + env=context.environ, + timeout=900, + ) + expect(summary.failed == 0, "PostgreSQL conformance must report zero TAP failures") + expect(summary.skipped == 0, "PostgreSQL conformance must not skip; the URL was not honoured") + expect( + summary.passed is not None and summary.passed >= MINIMUM_POSTGRES_TAP_PASSES, + "PostgreSQL conformance must pass the conformance suite plus provider tests", + ) + expect(summary.returncode == 0, "node test runner must exit 0") + return passed( + tap_pass=summary.passed, + tap_fail=summary.failed, + tap_skipped=summary.skipped, + postgres_url_sha256_prefix=sha256_hex(context.environ.get(POSTGRES_URL_VARIABLE, ""))[:12], + pg_package_version=pg_version, + ) + + +# --------------------------------------------------------------------------- +# Registry +# --------------------------------------------------------------------------- + + +LADDER_ROWS: tuple[LadderRow, ...] = ( + LadderRow( + id="s0.file_matrix_twelve_rows", + stage="0", + title="Twelve shared lifecycle scenarios pass on the file coordination provider", + product_path="store_direct", + gate="deterministic", + posix_only=False, + run=_row_file_matrix_twelve_rows, + ), + LadderRow( + id="s0.nokv_live_matrix", + stage="0", + title="The same matrix passes on a live NoKV stack with identical outcomes", + product_path="store_direct", + gate="env:nokv_legacy", + posix_only=False, + run=_row_nokv_live_matrix, + ), + LadderRow( + id="s1.cli_document_decodes_through_ts_store", + stage="1", + title="CLI-written candidate documents decode through the TypeScript FileAuthorityStore", + product_path="real_cli", + gate="deterministic", + posix_only=False, + run=_row_cli_document_decodes_through_ts_store, + ), + LadderRow( + id="s2a.nokv_live_qualification", + stage="2a", + title="The merged NoKV candidate passes its live single-node qualification", + product_path="store_direct", + gate="env:nokv_authority", + posix_only=False, + run=_row_nokv_live_qualification, + ), + LadderRow( + id="s2b.postgresql_conformance_live", + stage="2b", + title="PostgreSQL candidate passes the conformance suite against a live database", + product_path="store_direct", + gate="env:postgresql", + posix_only=False, + run=_row_postgresql_conformance_live, + ), + LadderRow( + id="s2c1.configure_enable_disable_roundtrip", + stage="2c1", + title="configure-goal previews, enables, reads back, and disables the observer", + product_path="real_cli", + gate="deterministic", + posix_only=False, + run=row_configure_enable_disable_roundtrip, + ), + LadderRow( + id="s2c1.every_writer_family_captures", + stage="2c1", + title="Every local writer family records a post-commit observation", + product_path="real_cli", + gate="deterministic", + posix_only=False, + run=row_every_writer_family_captures, + ), + LadderRow( + id="s2c1.default_off_isolation", + stage="2c1", + title="Default-off goals produce identical responses and no candidate storage", + product_path="real_cli", + gate="deterministic", + posix_only=False, + run=row_default_off_isolation, + ), + LadderRow( + id="s2c1.candidate_failure_preserves_primary", + stage="2c1", + title="Candidate construction failure never reverses the primary commit", + product_path="real_cli", + gate="deterministic", + posix_only=False, + run=row_candidate_failure_preserves_primary, + ), + LadderRow( + id="s2c1.crash_gap_loses_observation", + stage="2c1", + title="A SIGKILL between primary commit and observer loses only that observation", + product_path="real_cli", + gate="deterministic", + posix_only=True, + run=row_crash_gap_loses_observation, + ), + LadderRow( + id="s2c1.dual_runtime_root_consistency", + stage="2c1", + title="A --runtime-root override that differs from common_runtime_root keeps one candidate lineage", + product_path="real_cli", + gate="deterministic", + posix_only=False, + run=row_dual_runtime_root_consistency, + ), + LadderRow( + id="s2c1.migration_seeds_new_lineage", + stage="2c1", + title="migrate-state excludes the legacy lineage and seeds a fresh candidate", + product_path="real_cli", + gate="deterministic", + posix_only=False, + run=row_migration_seeds_new_lineage, + ), +) + +PENDING_ROWS: tuple[PendingRow, ...] = ( + PendingRow("s2c2.outbox_prepared_then_committed_entries", "2c2", "Stage 2C parity PRs"), + PendingRow("s2c2.drain_idempotent", "2c2", "Stage 2C parity PRs"), + PendingRow("s2c2.sigkill_between_primary_write_and_drain", "2c2", "Stage 2C parity PRs"), + PendingRow("s2c2.sigkill_mid_drain", "2c2", "Stage 2C parity PRs"), + PendingRow("s2c2.rollback_with_pending_entries", "2c2", "Stage 2C parity PRs"), + PendingRow("s2c2.parity_equal", "2c2", "Stage 2C parity PRs"), + PendingRow("s2c2.parity_divergent_detects_foreign_edit", "2c2", "Stage 2C parity PRs"), + PendingRow("s2c2.migration_seeds_and_drains", "2c2", "Stage 2C parity PRs"), + PendingRow("s2c2.growth_measurement_gate", "2c2", "Stage 2C parity PRs"), +) + + +def row_by_id(row_id: str) -> LadderRow: + for row in LADDER_ROWS: + if row.id == row_id: + return row + raise KeyError(row_id) + + +# --------------------------------------------------------------------------- +# Gates, redaction, and row execution +# --------------------------------------------------------------------------- + + +def gate_unverified_reason(gate: str, environ: Mapping[str, str]) -> tuple[str, JsonObject] | None: + """Return the unverified reason for an env gate whose stack is absent.""" + + required = GATE_REQUIREMENTS[gate] + missing = sorted(name for name in required if not environ.get(name)) + if missing: + return GATE_UNVERIFIED_REASON[gate], {"missing_variables": missing} + for flag in LIVE_OPT_IN_FLAGS: + if flag in required and environ.get(flag) != "1": + return f"{flag.lower()}_not_enabled", {"flag": flag} + return None + + +def default_forbidden_tokens(roots: Iterable[Path], environ: Mapping[str, str]) -> list[str]: + """Substrings whose presence in a report is a privacy leak.""" + + tokens: set[str] = set() + for root in roots: + tokens.add(str(root)) + tokens.add(str(Path(root).resolve())) + temp_root = Path(tempfile.gettempdir()) + tokens.update({str(temp_root), str(temp_root.resolve())}) + tokens.add(environ.get("HOME") or str(Path.home())) + tokens.update({str(REPO_ROOT), str(REPO_ROOT.resolve())}) + for name in (POSTGRES_URL_VARIABLE, *NOKV_STACK_VARIABLES, *NOKV_AUTHORITY_VARIABLES): + value = environ.get(name) + if value: + tokens.add(value) + tokens.update(_nokv_authority_config_tokens(environ)) + return sorted((token for token in tokens if len(token) >= 4), key=len, reverse=True) + + +def _nokv_authority_config_tokens(environ: Mapping[str, str]) -> set[str]: + """Every string leaf of the ignored client configuration is a forbidden token.""" + + path = _absolute_existing_path(environ.get(NOKV_AUTHORITY_CONFIG_VARIABLE), must_be_file=True) + if path is None: + return set() + try: + document = json.loads(path.read_text(encoding="utf-8")) + except (OSError, ValueError): + return set() + tokens: set[str] = set() + + def walk(value: object) -> None: + if isinstance(value, str): + tokens.add(value) + elif isinstance(value, dict): + for item in value.values(): + walk(item) + elif isinstance(value, list): + for item in value: + walk(item) + + walk(document) + return tokens + + +def redact(text: str, forbidden: Sequence[str]) -> str: + for token in forbidden: + text = text.replace(token, "") + return text + + +def _failure_result( + row: LadderRow, + exc: BaseException, + *, + forbidden: Sequence[str], + started: float, +) -> RowResult: + reason_code = "assertion_failed" if isinstance(exc, AssertionError) else "row_exception" + return RowResult( + row=row, + status="fail", + reason_code=reason_code, + evidence={ + "error_type": type(exc).__name__, + "message": redact(str(exc), forbidden)[:500], + }, + duration_ms=int((time.monotonic() - started) * 1000), + ) + + +def run_row( + row: LadderRow, + *, + root: Path, + environ: Mapping[str, str] | None = None, +) -> RowResult: + """Run one row in ``root`` and classify its outcome without raising.""" + + started = time.monotonic() + env: Mapping[str, str] = dict(os.environ) if environ is None else environ + forbidden = default_forbidden_tokens([root], env) + if row.posix_only and os.name == "nt": + return RowResult(row, "unverified", "posix_only", {"platform": os.name}, 0) + gate = gate_unverified_reason(row.gate, env) + if gate is not None: + return RowResult(row, "unverified", gate[0], gate[1], 0) + try: + outcome = row.run(RowContext(root=root, environ=env)) + except Exception as exc: + # Every row failure becomes a typed, redacted result; nothing escapes. + return _failure_result(row, exc, forbidden=forbidden, started=started) + if outcome.status not in {"pass", "unverified"}: + return _failure_result( + row, + RowAssertionError(f"row returned unsupported status {outcome.status!r}"), + forbidden=forbidden, + started=started, + ) + return RowResult( + row=row, + status=outcome.status, + reason_code=outcome.reason_code, + evidence=outcome.evidence, + duration_ms=int((time.monotonic() - started) * 1000), + ) + + +# --------------------------------------------------------------------------- +# Report, bindings, privacy scan, exit policy +# --------------------------------------------------------------------------- + + +def _git_output(*arguments: str) -> str | None: + try: + completed = subprocess.run( + ["git", *arguments], + cwd=REPO_ROOT, + capture_output=True, + text=True, + timeout=30, + check=False, + ) + except (OSError, subprocess.SubprocessError): + return None + return completed.stdout if completed.returncode == 0 else None + + +def _probe_digests() -> list[JsonObject]: + digests: list[JsonObject] = [] + for relative in PROBE_SOURCES: + path = REPO_ROOT / relative + if path.exists(): + digests.append({"path": relative.as_posix(), "sha256": sha256_hex(path.read_bytes())}) + return digests + + +def _nokv_client_config_digest(environ: Mapping[str, str]) -> str | None: + config_path = _absolute_existing_path( + environ.get(NOKV_AUTHORITY_CONFIG_VARIABLE), must_be_file=True + ) + if config_path is not None: + try: + return _nokv_authority_config_sha256(config_path) + except (OSError, ValueError): + return None + public = { + name: environ[name] + for name in NOKV_STACK_VARIABLES + if name not in NOKV_SECRET_VARIABLES and environ.get(name) + } + if len(public) != len(NOKV_STACK_VARIABLES) - len(NOKV_SECRET_VARIABLES): + return None + return sha256_hex(json.dumps(public, sort_keys=True, separators=(",", ":"))) + + +def _nokv_sdk_version() -> str | None: + for distribution in ("nokv", "nokv-python"): + try: + return metadata.version(distribution) + except metadata.PackageNotFoundError: + continue + return None + + +def collect_bindings(environ: Mapping[str, str]) -> JsonObject: + """Pin what the report was produced against; ``None`` means unknown.""" + + commit = _git_output("rev-parse", "HEAD") + status = _git_output("status", "--porcelain") + postgres_url = environ.get(POSTGRES_URL_VARIABLE) + return { + "loopx_commit": commit.strip() if commit else None, + "loopx_tree_dirty": bool(status.strip()) if status is not None else None, + "probe_sha256": _probe_digests(), + "nokv_client_config_sha256": _nokv_client_config_digest(environ), + "nokv_sdk_version": _nokv_sdk_version(), + "postgres_url_sha256_prefix": sha256_hex(postgres_url)[:12] if postgres_url else None, + "pg_package_version": _pg_package_version(), + } + + +def exit_code_for( + summary: Mapping[str, int], + *, + allow_unverified: bool, + allow_pending: bool = False, +) -> int: + """Green means every selected obligation was verified, not that a report exists. + + A pending row is a selected obligation with no executable evidence, so it + blocks a green exit exactly like an unverified row unless the caller + explicitly allows it. A privacy violation anywhere in the report, in a row + or confined to the bindings, is a failed evidence run even after the leak + was redacted; no flag relaxes it. + """ + + if summary["fail"] != 0: + return 1 + if summary.get("privacy_violations", 0) != 0: + return 1 + if summary["unverified"] != 0 and not allow_unverified: + return 1 + if summary.get("pending", 0) != 0 and not allow_pending: + return 1 + return 0 + + +def _finalize_report( + rows: Sequence[JsonObject], + pending: Sequence[JsonObject], + bindings: JsonObject, + *, + allow_unverified: bool, + allow_pending: bool, + generated_at: str, +) -> JsonObject: + summary = {status: sum(1 for row in rows if row["status"] == status) for status in ROW_STATUSES} + summary["pending"] = len(pending) + summary["executed"] = len(rows) + # Row leaks are already failures; a leak confined to the bindings has no + # row to fail, so the count is what the exit policy consumes. + summary["privacy_violations"] = sum( + 1 for row in rows if row.get("reason_code") == "privacy_violation" + ) + (1 if bindings.get("privacy_violation") else 0) + return { + "schema_version": REPORT_SCHEMA, + "generated_at": generated_at, + "rows": list(rows), + "pending": list(pending), + "summary": summary, + "bindings": bindings, + "exit_policy": { + "allow_unverified": allow_unverified, + "allow_pending": allow_pending, + "exit_code": exit_code_for( + summary, + allow_unverified=allow_unverified, + allow_pending=allow_pending, + ), + "rule": EXIT_POLICY_RULE, + }, + } + + +def _leaked_tokens(value: object, forbidden: Sequence[str]) -> list[str]: + serialized = json.dumps(value, sort_keys=True, ensure_ascii=False) + return [token for token in forbidden if token in serialized] + + +def assert_public_safe(report: JsonObject, *, forbidden: Sequence[str]) -> JsonObject: + """Turn any leak of a forbidden substring into ``fail/privacy_violation``. + + A leak confined to the bindings nulls every binding and still fails the + run through ``summary.privacy_violations``. + """ + + rows: list[JsonObject] = [] + for row in report["rows"]: + leaked = _leaked_tokens(row, forbidden) + if leaked: + rows.append( + { + **row, + "status": "fail", + "reason_code": "privacy_violation", + "evidence": {"leaked_token_count": len(leaked)}, + } + ) + else: + rows.append(dict(row)) + bindings = dict(report["bindings"]) + if _leaked_tokens(bindings, forbidden): + bindings = {key: None for key in bindings} + bindings["privacy_violation"] = True + return _finalize_report( + rows, + report["pending"], + bindings, + allow_unverified=bool(report["exit_policy"]["allow_unverified"]), + allow_pending=bool(report["exit_policy"].get("allow_pending", False)), + generated_at=str(report["generated_at"]), + ) + + +def build_report( + results: Sequence[RowResult], + *, + pending: Sequence[PendingRow], + allow_unverified: bool, + allow_pending: bool = False, + environ: Mapping[str, str] | None = None, + forbidden: Sequence[str] | None = None, +) -> JsonObject: + """Assemble the report and run the privacy scan over it.""" + + env: Mapping[str, str] = dict(os.environ) if environ is None else environ + tokens = list(forbidden) if forbidden is not None else default_forbidden_tokens([], env) + report = _finalize_report( + [result.as_dict() for result in results], + [row.describe() for row in pending], + collect_bindings(env), + allow_unverified=allow_unverified, + allow_pending=allow_pending, + generated_at=datetime.now(timezone.utc).isoformat(timespec="seconds"), + ) + return assert_public_safe(report, forbidden=tokens) + + +# --------------------------------------------------------------------------- +# Command line +# --------------------------------------------------------------------------- + + +def _parser() -> argparse.ArgumentParser: + parser = argparse.ArgumentParser( + prog="ladder", + description="Run the shared-goal-authority stage ladder against this checkout.", + ) + parser.add_argument("--stage", action="append", choices=STAGES, help="Only run rows of this stage; repeatable.") + parser.add_argument("--row", action="append", help="Only run this row id; repeatable.") + parser.add_argument( + "--allow-unverified", + action="store_true", + help="Exit 0 even when env-gated rows could not run; failures still exit 1.", + ) + parser.add_argument( + "--allow-pending", + action="store_true", + help=( + "Exit 0 even when the selection includes declared-but-unimplemented " + "rows; they are still reported as pending, never as pass." + ), + ) + parser.add_argument("--report-json", help="Write the JSON report to this path as well as stdout.") + parser.add_argument("--list", action="store_true", help="Print the row registry and pending rows, then exit.") + return parser + + +def _selected( + stages: Sequence[str] | None, + row_ids: Sequence[str] | None, +) -> tuple[list[LadderRow], list[PendingRow]]: + rows = [row for row in LADDER_ROWS if not stages or row.stage in stages] + pending = [row for row in PENDING_ROWS if not stages or row.stage in stages] + if row_ids: + wanted = set(row_ids) + rows = [row for row in rows if row.id in wanted] + pending = [row for row in pending if row.id in wanted] + return rows, pending + + +def _run_selected(rows: Sequence[LadderRow], environ: Mapping[str, str]) -> tuple[list[RowResult], list[str]]: + results: list[RowResult] = [] + tokens: set[str] = set() + for row in rows: + with tempfile.TemporaryDirectory(prefix="loopx-authority-ladder-") as scratch: + root = Path(scratch) + tokens.update(default_forbidden_tokens([root], environ)) + results.append(run_row(row, root=root, environ=environ)) + return results, sorted(tokens, key=len, reverse=True) + + +def _print_summary(report: JsonObject) -> None: + for status in ("fail", "unverified"): + ids = [f"{row['id']} ({row['reason_code']})" for row in report["rows"] if row["status"] == status] + if ids: + print(f"{status} rows: {', '.join(ids)}", file=sys.stderr) + pending_ids = [f"{row['id']} ({row['pending_until']})" for row in report["pending"]] + if pending_ids: + print(f"pending rows (not verified): {', '.join(pending_ids)}", file=sys.stderr) + if report["bindings"].get("privacy_violation"): + print("privacy violation confined to report bindings: redacted, the run fails", file=sys.stderr) + print( + "summary: " + ", ".join(f"{key}={value}" for key, value in sorted(report["summary"].items())) + + f"; exit_code={report['exit_policy']['exit_code']}", + file=sys.stderr, + ) + + +def main(argv: Sequence[str] | None = None) -> int: + args = _parser().parse_args(list(argv) if argv is not None else None) + rows, pending = _selected(args.stage, args.row) + if args.list: + listing = { + "schema_version": LIST_SCHEMA, + "rows": [row.describe() for row in rows], + "pending": [row.describe() for row in pending], + } + print(json.dumps(listing, indent=2, sort_keys=True)) + return 0 + if not rows and not pending: + print("no rows match the requested selection", file=sys.stderr) + return 2 + environ = dict(os.environ) + results, tokens = _run_selected(rows, environ) + report = build_report( + results, + pending=pending, + allow_unverified=bool(args.allow_unverified), + allow_pending=bool(args.allow_pending), + environ=environ, + forbidden=tokens or default_forbidden_tokens([], environ), + ) + rendered = json.dumps(report, indent=2, sort_keys=True) + if args.report_json: + Path(args.report_json).write_text(rendered + "\n", encoding="utf-8") + print(rendered) + _print_summary(report) + return int(report["exit_policy"]["exit_code"]) + + +__all__ = [ + "EXIT_POLICY_RULE", + "FILE_MATRIX_ROWS", + "GATES", + "GATE_REQUIREMENTS", + "LADDER_ROWS", + "LIVE_OPT_IN_FLAGS", + "NOKV_AUTHORITY_CONFIG_VARIABLE", + "NOKV_AUTHORITY_LIVE_FLAG", + "NOKV_AUTHORITY_PYTHON_VARIABLE", + "NOKV_AUTHORITY_VARIABLES", + "NOKV_AUTHORITY_WORKBENCH_VARIABLE", + "NOKV_LIVE_FLAG", + "NOKV_STACK_VARIABLES", + "PENDING_ROWS", + "POSTGRES_URL_VARIABLE", + "PRODUCT_PATHS", + "REPORT_SCHEMA", + "ROW_STATUSES", + "STAGES", + "LadderRow", + "PendingRow", + "RowAssertionError", + "RowContext", + "RowOutcome", + "RowResult", + "assert_public_safe", + "build_report", + "collect_bindings", + "default_forbidden_tokens", + "exit_code_for", + "gate_unverified_reason", + "main", + "passed", + "redact", + "row_by_id", + "run_row", + "unverified", +] + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/loopx/control_plane/testing/authority_e2e_row_support.py b/loopx/control_plane/testing/authority_e2e_row_support.py new file mode 100644 index 0000000000..b64a38de6b --- /dev/null +++ b/loopx/control_plane/testing/authority_e2e_row_support.py @@ -0,0 +1,192 @@ +"""Shared row vocabulary and CLI helpers for the shared-goal-authority ladder. + +Every row runner returns a ``RowOutcome`` or raises ``RowAssertionError``; the +helpers here drive the product only through ``run_cli`` and never import a +LoopX writer, so a row cannot become a second authority over the goal state. +""" + +from __future__ import annotations + +import hashlib +from collections.abc import Mapping +from dataclasses import dataclass +from pathlib import Path + +from .authority_e2e_fixtures import GoalWorkspace, JsonObject, run_cli + +AGENT_A = "agent-a" + + +AGENT_B = "agent-b" + + +PRIMARY_VISIBILITY_TIMEOUT_SECONDS = 15.0 + + +COMMITTED_OBSERVATION_OUTCOMES = frozenset({"captured", "replayed", "ambiguous_reconciled"}) + + +LOCAL_SHADOW_SUMMARY_ENABLED = { + "enabled": True, + "mode": "file_one_way", + "status": "enabled", +} + + +DEFAULT_OFF_PARITY_FIELDS: tuple[str, ...] = ( + "ok", + "added", + "already_exists", + "metadata_updated", + "status_changed", + "role", + "status", + "task_class", + "action_kind", + "continuation_policy", +) + + +MIGRATION_SEED_SCHEMA = "loopx_state_migration_shadow_seed_evidence_v0" + + +class RowAssertionError(AssertionError): + """A row invariant failed; the message is written to be public-safe.""" + + +@dataclass(frozen=True) +class RowContext: + """Per-row scratch root and the environment the row may consult.""" + + root: Path + environ: Mapping[str, str] + + +@dataclass(frozen=True) +class RowOutcome: + """What a row runner returns when it did not raise.""" + + status: str + reason_code: str | None + evidence: JsonObject + + +def passed(**evidence: object) -> RowOutcome: + return RowOutcome(status="pass", reason_code=None, evidence=dict(evidence)) + + +def unverified(reason_code: str, **evidence: object) -> RowOutcome: + return RowOutcome(status="unverified", reason_code=reason_code, evidence=dict(evidence)) + + +def expect(condition: bool, message: str) -> None: + if not condition: + raise RowAssertionError(message) + + +def sha256_hex(value: str | bytes) -> str: + payload = value.encode("utf-8") if isinstance(value, str) else value + return hashlib.sha256(payload).hexdigest() + + +def shadow_evidence(payload: Mapping[str, object], *, label: str) -> JsonObject: + evidence = payload.get("authority_shadow") + expect(isinstance(evidence, dict), f"{label} must carry authority_shadow evidence") + assert isinstance(evidence, dict) + return {str(key): value for key, value in evidence.items()} + + +def committed_observation(payload: Mapping[str, object], *, label: str) -> JsonObject: + evidence = shadow_evidence(payload, label=label) + expect( + evidence.get("outcome") in COMMITTED_OBSERVATION_OUTCOMES, + f"{label} observation outcome must be captured, replayed, or ambiguous_reconciled", + ) + expect( + evidence.get("primary_writeback_preserved") is True, + f"{label} must preserve the primary writeback", + ) + expect( + evidence.get("provider_to_local_writes") is False, + f"{label} must never write from provider to local state", + ) + expect( + evidence.get("candidate_read_for_decision") is False, + f"{label} must never read the candidate for a decision", + ) + return evidence + + +def add_todo(workspace: GoalWorkspace, text: str) -> JsonObject: + return run_cli( + workspace, + "todo", + "add", + "--goal-id", + workspace.goal_id, + "--role", + "agent", + "--text", + text, + "--task-class", + "advancement_task", + ) + + +def acquire_lease( + workspace: GoalWorkspace, + *, + todo_id: str, + owner: str, + idempotency_key: str, +) -> JsonObject: + return run_cli( + workspace, + "task-lease", + "acquire", + "--goal-id", + workspace.goal_id, + "--todo-id", + todo_id, + "--owner", + owner, + "--idempotency-key", + idempotency_key, + "--ttl-seconds", + "120", + ) + + +def configure_shadow(workspace: GoalWorkspace, *flags: str) -> JsonObject: + return run_cli(workspace, "configure-goal", "--goal-id", workspace.goal_id, *flags) + + +def lease_version(payload: Mapping[str, object], *, label: str) -> str: + lease = payload.get("lease") + expect(isinstance(lease, dict), f"{label} must return a lease record") + assert isinstance(lease, dict) + return str(int(lease["version"])) + + +__all__ = [ + "AGENT_A", + "AGENT_B", + "COMMITTED_OBSERVATION_OUTCOMES", + "DEFAULT_OFF_PARITY_FIELDS", + "LOCAL_SHADOW_SUMMARY_ENABLED", + "MIGRATION_SEED_SCHEMA", + "PRIMARY_VISIBILITY_TIMEOUT_SECONDS", + "RowAssertionError", + "RowContext", + "RowOutcome", + "acquire_lease", + "add_todo", + "committed_observation", + "configure_shadow", + "expect", + "lease_version", + "passed", + "shadow_evidence", + "sha256_hex", + "unverified", +] diff --git a/loopx/control_plane/testing/authority_e2e_rows_stage2c.py b/loopx/control_plane/testing/authority_e2e_rows_stage2c.py new file mode 100644 index 0000000000..8a3385b7ff --- /dev/null +++ b/loopx/control_plane/testing/authority_e2e_rows_stage2c.py @@ -0,0 +1,500 @@ +"""Stage 2C observation-foundation rows of the shared-goal-authority ladder. + +These rows port the local-shadow CLI E2E and migration assertions onto the +ladder: every write goes through the real ``python -m loopx.cli`` and the +candidate is read back only from its retained bytes. +""" + +from __future__ import annotations + +import json +from collections.abc import Mapping +from dataclasses import dataclass, field + +from .authority_e2e_fixtures import ( + GoalWorkspace, + JsonObject, + LegacyMigrationSource, + build_goal_workspace, + build_legacy_migration_source, + candidate_document, + candidate_store_paths, + hold_observation_lock, + kill_now, + parse_json_object, + run_cli, + spawn_cli, + unique_goal_id, + wait_until, +) +from .authority_e2e_row_support import ( + AGENT_A, + AGENT_B, + DEFAULT_OFF_PARITY_FIELDS, + LOCAL_SHADOW_SUMMARY_ENABLED, + MIGRATION_SEED_SCHEMA, + PRIMARY_VISIBILITY_TIMEOUT_SECONDS, + RowContext, + RowOutcome, + acquire_lease, + add_todo, + committed_observation, + configure_shadow, + expect, + lease_version, + passed, + shadow_evidence, +) + +def shadow_workspace(context: RowContext, prefix: str, *, shadow_enabled: bool) -> GoalWorkspace: + return build_goal_workspace( + context.root, + goal_id=unique_goal_id(prefix), + handoff_mode="hard_lease", + shadow_enabled=shadow_enabled, + runtime_root_binding="cli_override", + ) + + +def row_configure_enable_disable_roundtrip(context: RowContext) -> RowOutcome: + workspace = shadow_workspace(context, "ladder-configure", shadow_enabled=False) + preview = configure_shadow(workspace, "--local-authority-shadow-file") + expect(preview.get("dry_run") is True and preview.get("written") is False, "preview must not write") + enabled = configure_shadow(workspace, "--local-authority-shadow-file", "--execute") + expect(enabled.get("written") is True, "enable must write the registry") + + observed = add_todo(workspace, "Capture one post-commit observation through the product CLI.") + evidence = committed_observation(observed, label="todo add") + expect(evidence["outcome"] == "captured", "first observation must be captured") + expect(evidence["parity_verdict"] == "not_evaluated", "observation must not claim parity") + lease = acquire_lease( + workspace, + todo_id=str(observed["todo_id"]), + owner=AGENT_A, + idempotency_key="ladder-configure-lease", + ) + expect(lease.get("acquired") is True, "lease must be acquired") + expect(committed_observation(lease, label="task-lease acquire")["outcome"] == "captured", "lease observation must be captured") + document = candidate_document(workspace) + expect(len(document.todo_ids) == 1 and len(document.leases) == 1, "candidate head must hold one todo and one lease") + + inspected = configure_shadow(workspace) + after = inspected.get("after") + expect( + isinstance(after, dict) and after.get("local_authority_shadow") == LOCAL_SHADOW_SUMMARY_ENABLED, + "configure-goal must read back the enabled shadow summary", + ) + disabled = configure_shadow(workspace, "--clear-local-authority-shadow", "--execute") + expect(disabled.get("written") is True, "disable must write the registry") + before_disabled_write = document.path.read_bytes() + after_disable = add_todo(workspace, "This local lifecycle write must not execute the observer.") + expect(after_disable.get("ok") is True and after_disable.get("added") is True, "disabled write must still commit") + expect("authority_shadow" not in after_disable, "disabled write must not observe") + expect(document.path.read_bytes() == before_disabled_write, "candidate bytes must not change once disabled") + return passed(candidate_cursor=document.cursor, head_todos=1, head_leases=1) + + +def row_default_off_isolation(context: RowContext) -> RowOutcome: + enabled = shadow_workspace(context, "ladder-enabled", shadow_enabled=True) + baseline = shadow_workspace(context, "ladder-baseline", shadow_enabled=False) + text = "Capture one post-commit observation through the product CLI." + observed = add_todo(enabled, text) + committed_observation(observed, label="enabled todo add") + default_off = add_todo(baseline, text) + expect("authority_shadow" not in default_off, "default-off write must not carry observation evidence") + differing = [field for field in DEFAULT_OFF_PARITY_FIELDS if observed.get(field) != default_off.get(field)] + expect(not differing, f"default-off response fields must match the observed response: {differing}") + expect(not (baseline.runtime_root / "authority-shadow").exists(), "default-off must not create candidate storage") + return passed(compared_fields=len(DEFAULT_OFF_PARITY_FIELDS)) + + +def row_candidate_failure_preserves_primary(context: RowContext) -> RowOutcome: + workspace = shadow_workspace(context, "ladder-failure", shadow_enabled=False) + configure_shadow(workspace, "--local-authority-shadow-file", "--execute") + workspace.runtime_root.mkdir(parents=True, exist_ok=True) + (workspace.runtime_root / "authority-shadow").write_text("block candidate directory", encoding="utf-8") + result = add_todo(workspace, "The primary write survives a candidate construction failure.") + expect(result.get("ok") is True and result.get("added") is True, "primary write must commit") + evidence = shadow_evidence(result, label="todo add") + expect(evidence.get("outcome") == "failed", "candidate failure must be reported as failed") + expect(evidence.get("reason_code") == "shadow_observation_failed", "candidate failure must carry its typed reason") + expect(evidence.get("primary_writeback_preserved") is True, "candidate failure must preserve the primary writeback") + expect( + str(result["todo_id"]) in workspace.state_path.read_text(encoding="utf-8"), + "the committed todo must be present in the primary state", + ) + return passed(outcome="failed", reason_code="shadow_observation_failed") + + +def row_crash_gap_loses_observation(context: RowContext) -> RowOutcome: + workspace = shadow_workspace(context, "ladder-crash-gap", shadow_enabled=False) + configure_shadow(workspace, "--local-authority-shadow-file", "--execute") + first_text = "Primary commit that loses its post-commit observation." + with hold_observation_lock(workspace): + process = spawn_cli( + workspace, + "todo", + "add", + "--goal-id", + workspace.goal_id, + "--role", + "agent", + "--text", + first_text, + "--task-class", + "advancement_task", + ) + try: + visible = wait_until( + lambda: first_text in workspace.state_path.read_text(encoding="utf-8"), + PRIMARY_VISIBILITY_TIMEOUT_SECONDS, + ) + expect(visible, "primary Todo commit did not become visible") + expect(process.poll() is None, "writer must still be blocked on the observation lock") + finally: + kill_now(process) + expect(not candidate_store_paths(workspace), "a killed writer must leave no candidate document") + recovered = add_todo(workspace, "A later primary commit refreshes the current full snapshot.") + evidence = committed_observation(recovered, label="recovering todo add") + expect(evidence["outcome"] == "captured", "recovery observation must be captured") + expect(evidence["durable_source_outbox"] is False, "no durable outbox may be claimed") + expect(evidence["source_transaction_correlated"] is False, "no transaction correlation may be claimed") + expect(evidence["parity_verdict"] == "not_evaluated", "recovery must not claim parity") + document = candidate_document(workspace) + expect(len(document.todo_ids) == 2, "the refreshed snapshot must include both primary commits") + return passed(lost_observations=1, refreshed_head_todos=2, candidate_cursor=document.cursor) + + +@dataclass +class _WriterSequence: + workspace: GoalWorkspace + committed: list[tuple[str, JsonObject]] = field(default_factory=list) + + def commit(self, label: str, payload: JsonObject, *, flag: str) -> JsonObject: + expect(payload.get(flag) is True, f"{label} must report {flag}=true") + self.committed.append((label, payload)) + return payload + + def cli(self, *args: str) -> JsonObject: + return run_cli(self.workspace, *args, "--goal-id", self.workspace.goal_id) + + +def _writer_sequence_lease_lifecycle(sequence: _WriterSequence) -> str: + """Add a todo, then acquire, update, renew, transfer, and complete it.""" + + first = sequence.commit( + "todo add", + add_todo(sequence.workspace, "Deliver one bounded control-plane change."), + flag="added", + ) + todo_id = str(first["todo_id"]) + acquired = sequence.commit( + "task-lease acquire", + acquire_lease(sequence.workspace, todo_id=todo_id, owner=AGENT_A, idempotency_key="ladder-lease-a"), + flag="acquired", + ) + sequence.commit( + "todo update", + sequence.cli("todo", "update", "--todo-id", todo_id, "--note", "A public-safe update.", "--agent-id", AGENT_A), + flag="changed", + ) + renewed = sequence.commit( + "task-lease renew", + sequence.cli( + "task-lease", "renew", "--todo-id", todo_id, "--owner", AGENT_A, + "--idempotency-key", "ladder-lease-a", + "--expected-version", lease_version(acquired, label="acquire"), + "--ttl-seconds", "120", + ), + flag="renewed", + ) + transferred = sequence.commit( + "task-lease transfer", + sequence.cli( + "task-lease", "transfer", "--todo-id", todo_id, "--owner", AGENT_A, + "--idempotency-key", "ladder-lease-a", "--new-owner", AGENT_B, + "--new-idempotency-key", "ladder-lease-b", + "--expected-version", lease_version(renewed, label="renew"), + "--ttl-seconds", "120", + ), + flag="transferred", + ) + sequence.commit( + "todo complete", + sequence.cli( + "todo", "complete", "--todo-id", todo_id, "--agent-id", AGENT_B, + "--task-lease-idempotency-key", "ladder-lease-b", + "--task-lease-expected-version", lease_version(transferred, label="transfer"), + "--evidence", "validation://ladder-complete", "--no-follow-up", + ), + flag="completed", + ) + return todo_id + + +def _writer_sequence_supersede_and_hygiene(sequence: _WriterSequence) -> str: + """Add a second todo, replay its acquire, supersede it, then run hygiene writers.""" + + second = sequence.commit( + "todo add (second)", + add_todo(sequence.workspace, "Replace this bounded work with a successor."), + flag="added", + ) + todo_id = str(second["todo_id"]) + acquired = sequence.commit( + "task-lease acquire (second)", + acquire_lease(sequence.workspace, todo_id=todo_id, owner=AGENT_A, idempotency_key="ladder-lease-c"), + flag="acquired", + ) + replayed = acquire_lease(sequence.workspace, todo_id=todo_id, owner=AGENT_A, idempotency_key="ladder-lease-c") + expect(replayed.get("idempotent") is True, "re-acquire with the same key must be idempotent") + expect("authority_shadow" not in replayed, "an idempotent re-acquire must not observe") + sequence.commit( + "todo supersede", + sequence.cli( + "todo", "supersede", "--todo-id", todo_id, "--agent-id", AGENT_A, + "--reason", "Replace obsolete work.", + "--next-agent-todo", "Carry the bounded work forward.", + "--task-lease-idempotency-key", "ladder-lease-c", + "--task-lease-expected-version", lease_version(acquired, label="acquire (second)"), + ), + flag="superseded", + ) + sequence.commit( + "todo capture-followups", + sequence.cli( + "todo", "capture-followups", + "--follow-up", "Verify the migrated authority projection.", + "--evidence", "validation://ladder-followup", + ), + flag="changed", + ) + sequence.commit( + "todo archive-completed", + sequence.cli("todo", "archive-completed", "--max-active-done", "0", "--execute"), + flag="changed", + ) + return todo_id + + +def row_every_writer_family_captures(context: RowContext) -> RowOutcome: + workspace = build_goal_workspace( + context.root, + goal_id=unique_goal_id("ladder-writers"), + handoff_mode="legacy", + shadow_enabled=True, + runtime_root_binding="cli_override", + ) + sequence = _WriterSequence(workspace) + sequence.commit( + "handoff-mode set", + sequence.cli("handoff-mode", "set", "--mode", "hard_lease"), + flag="changed", + ) + first_todo = _writer_sequence_lease_lifecycle(sequence) + second_todo = _writer_sequence_supersede_and_hygiene(sequence) + + observations = [committed_observation(payload, label=label) for label, payload in sequence.committed] + captured = [evidence for evidence in observations if evidence["outcome"] == "captured"] + document = candidate_document(workspace) + expect(document.cursor == str(len(captured)), "candidate cursor must equal the number of captured observations") + expect( + set(document.operation_ids) == {str(evidence["observation_id"]) for evidence in observations}, + "candidate operation ids must be exactly the observation ids", + ) + expect( + {str(lease.get("todo_id")) for lease in document.leases} == {first_todo, second_todo}, + "candidate head must retain the lease records of both leased todos", + ) + expect( + all(lease.get("status") == "released" for lease in document.leases), + "no time-active lease may remain in the candidate head", + ) + listed = sequence.cli("todo", "list") + listed_ids = sorted(str(todo.get("todo_id")) for todo in listed.get("todos") or [] if isinstance(todo, dict)) + expect(sorted(document.todo_ids) == listed_ids, "candidate head todos must equal the projected todo list") + return passed( + writer_families=len(sequence.committed), + captured=len(captured), + outcomes=sorted({str(evidence["outcome"]) for evidence in observations}), + candidate_cursor=document.cursor, + head_todos=len(document.todo_ids), + head_leases=len(document.leases), + ) + + +def _migration_arguments(source: LegacyMigrationSource) -> list[str]: + return [ + "migrate-state", + "--legacy-registry", str(source.legacy_registry), + "--legacy-runtime-root", str(source.legacy_runtime), + "--target-runtime-root", str(source.target_runtime), + "--goal-id", source.old_goal_id, + "--goal-id-map", f"{source.old_goal_id}={source.new_goal_id}", + "--path-map", f"{source.source_repo}={source.target_repo}", + "--copy-active-state", + "--copy-runtime", + "--no-global-sync", + ] + + +def _first_entry(payload: Mapping[str, object], key: str, *, label: str) -> JsonObject: + entries = payload.get(key) + expect(isinstance(entries, list) and len(entries) == 1, f"{label} must report exactly one {key} entry") + assert isinstance(entries, list) + entry = entries[0] + expect(isinstance(entry, dict), f"{label} {key} entry must be an object") + assert isinstance(entry, dict) + return {str(field): value for field, value in entry.items()} + + +def _assert_migration_preview(source: LegacyMigrationSource, preview: JsonObject, sentinel: bytes) -> None: + expect(preview.get("ok") is True and preview.get("dry_run") is True, "migration preview must be a dry run") + expect(source.target_registry.read_bytes() == sentinel, "dry run must not write the target registry") + expect(not source.target_runtime.exists(), "dry run must not create the target runtime") + runtime_result = _first_entry(preview, "runtime_goals", label="migration preview") + expect(runtime_result.get("copied_file_count") == 0, "dry run must copy no runtime files") + seed = _first_entry(preview, "authority_shadow_seeds", label="migration preview") + expect( + seed + == { + "schema_version": MIGRATION_SEED_SCHEMA, + "goal_id": source.new_goal_id, + "attempted": False, + "outcome": "planned", + "reason_code": None, + }, + "dry run must plan, not attempt, the shadow seed", + ) + expect(source.private_marker not in json.dumps(preview, sort_keys=True), "preview must not leak private provider bytes") + + +def _assert_migration_executed(source: LegacyMigrationSource, executed: JsonObject) -> JsonObject: + expect(executed.get("ok") is True and executed.get("wrote_project_registry") is True, "execute must write the registry") + runtime_result = _first_entry(executed, "runtime_goals", label="migration execute") + expect(runtime_result.get("copied") is True, "execute must copy the runtime goal directory") + lease_path = source.target_runtime / "goals" / source.new_goal_id / "task-leases" / "safe-local.json" + copied_lease = parse_json_object(lease_path.read_text(encoding="utf-8")) + expect(copied_lease.get("goal_id") == source.new_goal_id, "copied lease must carry the migrated goal id") + identity = (source.target_shadow_directory / "store-identity").read_text(encoding="utf-8") + expect(identity.startswith("file:") and identity != source.old_store_identity, "target lineage must be fresh") + store_paths = sorted(source.target_shadow_directory.glob("authority-store-*.json")) + expect(len(store_paths) == 1, "execute must seed exactly one candidate document") + store = parse_json_object(store_paths[0].read_text(encoding="utf-8")) + expect(store.get("goal_id") == source.new_goal_id and store.get("store_identity") == identity, "seeded store must bind the new lineage") + committed = store.get("committed") + expect(store.get("cursor") == "1" and isinstance(committed, list) and len(committed) == 1, "seed must be the first and only commit") + serialized = json.dumps(store, sort_keys=True) + for forbidden in (source.old_store_identity, source.legacy_revision, str(source.source_repo), source.private_marker): + expect(forbidden not in serialized, "seeded store must not carry any legacy lineage or private byte") + expect(not (source.target_shadow_directory / "authority-store-legacy.json").exists(), "legacy document must not migrate") + seed = _first_entry(executed, "authority_shadow_seeds", label="migration execute") + expect(seed.get("goal_id") == source.new_goal_id and seed.get("attempted") is True, "seed must target the migrated goal") + expect(seed.get("outcome") == "captured", "seed must be captured") + return store + + +def row_migration_seeds_new_lineage(context: RowContext) -> RowOutcome: + source = build_legacy_migration_source( + context.root, + old_goal_id=unique_goal_id("legacy"), + new_goal_id=unique_goal_id("migrated"), + ) + sentinel = b'{"schema_version":"existing","goals":[]}\n' + source.target_registry.write_bytes(sentinel) + arguments = _migration_arguments(source) + _assert_migration_preview(source, run_cli(source, *arguments), sentinel) + store = _assert_migration_executed(source, run_cli(source, *arguments, "--execute")) + return passed(seed_outcome="captured", seeded_cursor=str(store.get("cursor")), legacy_lineage_excluded=True) + + +def row_dual_runtime_root_consistency(context: RowContext) -> RowOutcome: + """``--runtime-root`` differs from ``common_runtime_root``: one lineage per goal.""" + + workspace = build_goal_workspace( + context.root, + goal_id=unique_goal_id("ladder-one-root"), + handoff_mode="hard_lease", + shadow_enabled=True, + runtime_root_binding="cli_override_divergent", + ) + expect( + workspace.registry_runtime_root != workspace.runtime_root, + "fixture must register a different common_runtime_root than the override", + ) + added = add_todo(workspace, "Every hook of one CLI call shares one runtime root.") + todo_id = str(added["todo_id"]) + acquired = acquire_lease( + workspace, + todo_id=todo_id, + owner=AGENT_A, + idempotency_key="ladder-one-root-lease", + ) + updated = run_cli( + workspace, "todo", "update", "--goal-id", workspace.goal_id, "--todo-id", todo_id, + "--note", "Observed under the override root.", "--agent-id", AGENT_A, + ) + followups = run_cli( + workspace, "todo", "capture-followups", "--goal-id", workspace.goal_id, + "--follow-up", "Keep one candidate lineage per goal.", + "--evidence", "validation://ladder-one-root", + ) + completed = run_cli( + workspace, "todo", "complete", "--goal-id", workspace.goal_id, "--todo-id", todo_id, + "--agent-id", AGENT_A, "--task-lease-idempotency-key", "ladder-one-root-lease", + "--task-lease-expected-version", lease_version(acquired, label="acquire"), + "--evidence", "validation://ladder-one-root-complete", "--no-follow-up", + ) + observations = [ + committed_observation(payload, label=label) + for label, payload in ( + ("todo add", added), + ("task-lease acquire", acquired), + ("todo update", updated), + ("todo capture-followups", followups), + ("todo complete", completed), + ) + ] + identities = {str(evidence.get("store_identity")) for evidence in observations} + expect(len(identities) == 1, "every writer family must observe into one store identity") + document = candidate_document(workspace) + expect(document.store_identity in identities, "candidate bytes must carry the observed identity") + expect(document.cursor == str(len(observations)), "candidate cursor must equal the observation count") + expect( + todo_id in document.todo_ids and len(document.todo_ids) == 2, + "head must hold the completed todo and its captured follow-up", + ) + expect( + [lease.get("todo_id") for lease in document.leases] == [todo_id] + and document.leases[0].get("status") == "released", + "head must hold exactly the released lease of the completed todo", + ) + lease_path = workspace.runtime_root / "goals" / workspace.goal_id / "task-leases" / f"{todo_id}.json" + expect(lease_path.exists(), "lease state must live under the override root") + expect( + not (workspace.registry_runtime_root / "authority-shadow").exists(), + "the registry root must not gain a candidate lineage", + ) + expect( + not (workspace.registry_runtime_root / "goals").exists(), + "the registry root must not gain lease state", + ) + return passed( + observations=len(observations), + store_identities=len(identities), + candidate_cursor=document.cursor, + head_todos=len(document.todo_ids), + head_leases=len(document.leases), + ) + + +__all__ = [ + "row_candidate_failure_preserves_primary", + "row_configure_enable_disable_roundtrip", + "row_crash_gap_loses_observation", + "row_default_off_isolation", + "row_dual_runtime_root_consistency", + "row_every_writer_family_captures", + "row_migration_seeds_new_lineage", + "shadow_workspace", +] diff --git a/pyproject.toml b/pyproject.toml index d5c94b1f9f..5880ca09a7 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -132,6 +132,10 @@ files = [ "loopx/control_plane/runtime/time.py", "loopx/control_plane/runtime/trajectory_hygiene.py", "loopx/control_plane/scheduler/time.py", + "loopx/control_plane/testing/authority_e2e_fixtures.py", + "loopx/control_plane/testing/authority_e2e_ladder.py", + "loopx/control_plane/testing/authority_e2e_row_support.py", + "loopx/control_plane/testing/authority_e2e_rows_stage2c.py", "loopx/control_plane/work_items/delivery_batch_scale.py", "loopx/control_plane/work_items/delivery_outcome.py", "loopx/control_plane/work_items/lifecycle.py", diff --git a/tests/control_plane/test_shared_goal_authority_e2e.py b/tests/control_plane/test_shared_goal_authority_e2e.py new file mode 100644 index 0000000000..009dd9024c --- /dev/null +++ b/tests/control_plane/test_shared_goal_authority_e2e.py @@ -0,0 +1,382 @@ +"""Pytest projection of the shared-goal-authority E2E stage ladder. + +Each ladder row becomes one parametrized test. Rows whose environment gate is +closed skip with ``unverified: `` so a green pytest run never implies +that a live provider was exercised; the ladder's own exit policy is pinned +separately so the standalone runner cannot report green while unverified. +""" + +from __future__ import annotations + +import ast +import json +import os +from collections.abc import Iterator +from pathlib import Path + +import pytest + +from loopx.control_plane.testing import authority_e2e_ladder as ladder + + +LIVE_ENVIRONMENT_VARIABLES = ( + ladder.POSTGRES_URL_VARIABLE, + ladder.NOKV_LIVE_FLAG, + *ladder.NOKV_STACK_VARIABLES, + ladder.NOKV_AUTHORITY_LIVE_FLAG, + *ladder.NOKV_AUTHORITY_VARIABLES, +) +GATED_ROW_IDS = ( + "s0.nokv_live_matrix", + "s2a.nokv_live_qualification", + "s2b.postgresql_conformance_live", +) +PENDING_ONLY_ROW_ID = "s2c2.parity_equal" +CHEAP_DETERMINISTIC_ROW_ID = "s0.file_matrix_twelve_rows" +FULL_LADDER_VARIABLE = "LOOPX_LADDER_FULL" +# Rows whose assertions the in-repo CLI E2E suite already pins through the same +# product path. The pytest job runs close to its time budget, so the default CI +# projection keeps the rows that only the ladder exercises and defers these to +# the example runner (or LOOPX_LADDER_FULL=1). +# Rows the default CI projection skips, each named with the product-CLI E2E +# test that pins the same assertions; a guard below fails when a twin is +# renamed or deleted so the skip cannot outlive its coverage. +CLI_E2E_TWIN_FILE = Path(__file__).with_name("test_local_authority_shadow_cli_e2e.py") +CLI_E2E_COVERAGE = { + "s2c1.configure_enable_disable_roundtrip": ( + "test_product_cli_configure_capture_readback_disable_and_default_off_lifecycle_isolation" + ), + "s2c1.default_off_isolation": ( + "test_product_cli_configure_capture_readback_disable_and_default_off_lifecycle_isolation" + ), + "s2c1.candidate_failure_preserves_primary": ( + "test_product_cli_candidate_failure_preserves_the_primary_lifecycle_commit" + ), + "s2c1.crash_gap_loses_observation": ( + "test_product_cli_loses_capture_between_commit_and_observer_then_refreshes_snapshot" + ), + "s2c1.dual_runtime_root_consistency": ( + "test_product_cli_runtime_root_override_keeps_one_candidate_lineage" + ), +} +CLI_E2E_COVERED_ROW_IDS = tuple(CLI_E2E_COVERAGE) +CLI_E2E_COVERAGE_REASON = ( + "pinned by tests/control_plane/test_local_authority_shadow_cli_e2e.py; " + "run examples/shared-goal-authority-e2e/ladder.py or set LOOPX_LADDER_FULL=1" +) +REQUIRED_PENDING_ROW_IDS = ( + "s2c2.outbox_prepared_then_committed_entries", + "s2c2.drain_idempotent", + "s2c2.sigkill_between_primary_write_and_drain", + "s2c2.sigkill_mid_drain", + "s2c2.rollback_with_pending_entries", + "s2c2.parity_equal", + "s2c2.parity_divergent_detects_foreign_edit", + "s2c2.migration_seeds_and_drains", + "s2c2.growth_measurement_gate", +) + + +def _row_parameters() -> Iterator[object]: + full_ladder = os.environ.get(FULL_LADDER_VARIABLE) == "1" + for row in ladder.LADDER_ROWS: + marks = [] + if row.posix_only: + marks.append( + pytest.mark.skipif(os.name == "nt", reason="requires POSIX cross-process flock and SIGKILL") + ) + if row.id in CLI_E2E_COVERED_ROW_IDS and not full_ladder: + marks.append(pytest.mark.skip(reason=CLI_E2E_COVERAGE_REASON)) + yield pytest.param(row, id=row.id, marks=marks) + + +@pytest.mark.parametrize("row", list(_row_parameters())) +def test_ladder_row_passes_or_is_declared_unverified( + row: ladder.LadderRow, + tmp_path: Path, +) -> None: + result = ladder.run_row(row, root=tmp_path, environ=os.environ) + if result.status == "unverified": + pytest.skip(f"unverified: {result.reason_code}") + report = ladder.build_report( + [result], + pending=(), + allow_unverified=False, + environ=os.environ, + forbidden=ladder.default_forbidden_tokens([tmp_path], os.environ), + ) + reported = report["rows"][0] + assert reported["status"] == "pass", (reported["reason_code"], reported["evidence"]) + assert report["summary"] == {"pass": 1, "fail": 0, "unverified": 0, "pending": 0, "executed": 1, "privacy_violations": 0} + assert report["exit_policy"]["exit_code"] == 0 + + +def test_registry_vocabulary_and_pending_rows_are_declared_not_claimed() -> None: + row_ids = [row.id for row in ladder.LADDER_ROWS] + assert set(CLI_E2E_COVERED_ROW_IDS) < set(row_ids) + pending_ids = [row.id for row in ladder.PENDING_ROWS] + assert len(set(row_ids)) == len(row_ids) + assert set(row_ids).isdisjoint(pending_ids) + assert set(REQUIRED_PENDING_ROW_IDS) <= set(pending_ids) + assert {row.stage for row in ladder.LADDER_ROWS} == {"0", "1", "2a", "2b", "2c1"} + assert {row.stage for row in ladder.PENDING_ROWS} == {"2c2"} + assert all("#3819" not in row.pending_until for row in ladder.PENDING_ROWS) + assert all(row.gate in ladder.GATES for row in ladder.LADDER_ROWS) + assert all(row.product_path in ladder.PRODUCT_PATHS for row in ladder.LADDER_ROWS) + with pytest.raises(ValueError): + ladder.LadderRow( + id="bad.stage", + stage="9", + title="invalid", + product_path="real_cli", + gate="deterministic", + posix_only=False, + run=lambda _context: ladder.passed(), + ) + + +def test_main_never_reports_green_while_unverified( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + capsys: pytest.CaptureFixture[str], +) -> None: + for name in LIVE_ENVIRONMENT_VARIABLES: + monkeypatch.delenv(name, raising=False) + report_path = tmp_path / "report.json" + argv = ["--report-json", str(report_path)] + for row_id in GATED_ROW_IDS: + argv.extend(["--row", row_id]) + + assert ladder.main(argv) == 1 + report = json.loads(report_path.read_text(encoding="utf-8")) + assert report["schema_version"] == ladder.REPORT_SCHEMA + assert report["summary"] == { + "pass": 0, + "fail": 0, + "unverified": 3, + "pending": 0, + "executed": 3, + "privacy_violations": 0, + } + assert {row["status"] for row in report["rows"]} == {"unverified"} + assert {row["reason_code"] for row in report["rows"]} == { + "nokv_live_env_missing", + "nokv_authority_env_missing", + "postgres_url_missing", + } + assert report["exit_policy"] == { + "allow_unverified": False, + "allow_pending": False, + "exit_code": 1, + "rule": ladder.EXIT_POLICY_RULE, + } + assert report["bindings"]["postgres_url_sha256_prefix"] is None + assert report["bindings"]["nokv_client_config_sha256"] is None + assert report["bindings"]["loopx_commit"] is None or len(report["bindings"]["loopx_commit"]) == 40 + captured = capsys.readouterr() + assert "unverified rows:" in captured.err + + assert ladder.main([*argv, "--allow-unverified"]) == 0 + relaxed = json.loads(report_path.read_text(encoding="utf-8")) + assert relaxed["summary"]["unverified"] == 3 + assert relaxed["exit_policy"]["allow_unverified"] is True + assert relaxed["exit_policy"]["exit_code"] == 0 + capsys.readouterr() + + +def test_stage_2a_row_reports_specific_unverified_reasons_for_each_missing_input( + tmp_path: Path, +) -> None: + row = ladder.row_by_id("s2a.nokv_live_qualification") + assert row.gate == "env:nokv_authority" + assert row.stage == "2a" + base = {name: value for name, value in os.environ.items() if name not in LIVE_ENVIRONMENT_VARIABLES} + + gated = ladder.run_row(row, root=tmp_path, environ=base) + assert gated.status == "unverified" + assert gated.reason_code == "nokv_authority_env_missing" + assert gated.evidence["missing_variables"] == sorted( + [ladder.NOKV_AUTHORITY_LIVE_FLAG, *ladder.NOKV_AUTHORITY_VARIABLES] + ) + + config = tmp_path / "nokv-client.json" + config.write_text(json.dumps({"root_id": "0" * 32, "object_store": {"kind": "memory"}}), encoding="utf-8") + inputs = { + **base, + ladder.NOKV_AUTHORITY_LIVE_FLAG: "0", + ladder.NOKV_AUTHORITY_CONFIG_VARIABLE: str(config), + ladder.NOKV_AUTHORITY_PYTHON_VARIABLE: "relative/python", + ladder.NOKV_AUTHORITY_WORKBENCH_VARIABLE: "ladder-workbench", + } + not_enabled = ladder.run_row(row, root=tmp_path, environ=inputs) + assert not_enabled.status == "unverified" + assert not_enabled.reason_code == "loopx_nokv_authority_live_not_enabled" + + inputs[ladder.NOKV_AUTHORITY_LIVE_FLAG] = "1" + relative_python = ladder.run_row(row, root=tmp_path, environ=inputs) + assert relative_python.status == "unverified" + assert relative_python.reason_code == "nokv_authority_python_missing" + + inputs[ladder.NOKV_AUTHORITY_CONFIG_VARIABLE] = str(tmp_path / "absent.json") + missing_config = ladder.run_row(row, root=tmp_path, environ=inputs) + assert missing_config.status == "unverified" + assert missing_config.reason_code == "nokv_authority_config_missing" + + # Configuration values are secrets: every string leaf becomes a forbidden token. + inputs[ladder.NOKV_AUTHORITY_CONFIG_VARIABLE] = str(config) + tokens = ladder.default_forbidden_tokens([tmp_path], inputs) + assert str(config) in tokens + assert "0" * 32 in tokens + assert "memory" in tokens + assert ladder.collect_bindings(inputs)["nokv_client_config_sha256"] is not None + + +def test_pending_rows_never_exit_green_without_allow_pending( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + capsys: pytest.CaptureFixture[str], +) -> None: + for name in LIVE_ENVIRONMENT_VARIABLES: + monkeypatch.delenv(name, raising=False) + report_path = tmp_path / "report.json" + + # Pending-only selection: zero executions must not read as success. + assert ladder.main(["--row", PENDING_ONLY_ROW_ID, "--report-json", str(report_path)]) == 1 + report = json.loads(report_path.read_text(encoding="utf-8")) + assert report["summary"] == {"pass": 0, "fail": 0, "unverified": 0, "pending": 1, "executed": 0, "privacy_violations": 0} + assert report["rows"] == [] + assert report["pending"][0]["id"] == PENDING_ONLY_ROW_ID + assert report["pending"][0]["status"] == "pending" + assert report["exit_policy"]["exit_code"] == 1 + assert "pending rows (not verified):" in capsys.readouterr().err + + assert ladder.main(["--row", PENDING_ONLY_ROW_ID, "--allow-pending", "--report-json", str(report_path)]) == 0 + allowed = json.loads(report_path.read_text(encoding="utf-8")) + assert allowed["exit_policy"] == { + "allow_unverified": False, + "allow_pending": True, + "exit_code": 0, + "rule": ladder.EXIT_POLICY_RULE, + } + assert allowed["summary"]["pending"] == 1 + capsys.readouterr() + + # A whole pending stage behaves the same way. + assert ladder.main(["--stage", "2c2", "--report-json", str(report_path)]) == 1 + capsys.readouterr() + + # Mixed selection: one executable pass does not excuse a pending obligation. + mixed = ["--row", CHEAP_DETERMINISTIC_ROW_ID, "--row", PENDING_ONLY_ROW_ID, "--report-json", str(report_path)] + assert ladder.main(mixed) == 1 + report = json.loads(report_path.read_text(encoding="utf-8")) + assert report["summary"] == {"pass": 1, "fail": 0, "unverified": 0, "pending": 1, "executed": 1, "privacy_violations": 0} + capsys.readouterr() + assert ladder.main([*mixed, "--allow-pending"]) == 0 + capsys.readouterr() + + # List-only never claims verification: it prints the registry and exits 0. + assert ladder.main(["--list", "--row", PENDING_ONLY_ROW_ID]) == 0 + listing = json.loads(capsys.readouterr().out) + assert listing["schema_version"] == ladder.LIST_SCHEMA + assert listing["rows"] == [] + assert [row["id"] for row in listing["pending"]] == [PENDING_ONLY_ROW_ID] + assert "exit_policy" not in listing + + +def test_privacy_scan_turns_leaks_into_failures(tmp_path: Path) -> None: + row = ladder.row_by_id("s0.file_matrix_twelve_rows") + leaking = ladder.RowResult( + row=row, + status="pass", + reason_code=None, + evidence={"pointer": str(tmp_path / "leaked-root")}, + duration_ms=1, + ) + clean = ladder.RowResult( + row=ladder.row_by_id("s1.cli_document_decodes_through_ts_store"), + status="pass", + reason_code=None, + evidence={"cursor": "3"}, + duration_ms=1, + ) + + report = ladder.build_report( + [leaking, clean], + pending=ladder.PENDING_ROWS, + allow_unverified=True, + environ=os.environ, + forbidden=ladder.default_forbidden_tokens([tmp_path], os.environ), + ) + + statuses = {row["id"]: row for row in report["rows"]} + assert statuses[leaking.row.id]["status"] == "fail" + assert statuses[leaking.row.id]["reason_code"] == "privacy_violation" + assert str(tmp_path) not in json.dumps(report) + assert statuses[clean.row.id]["status"] == "pass" + assert report["summary"] == { + "pass": 1, + "fail": 1, + "unverified": 0, + "pending": len(ladder.PENDING_ROWS), + "executed": 2, + "privacy_violations": 1, + } + assert report["exit_policy"]["exit_code"] == 1 + assert {row["status"] for row in report["pending"]} == {"pending"} + + +def test_privacy_leak_confined_to_bindings_still_fails_the_run() -> None: + row = ladder.row_by_id("s1.cli_document_decodes_through_ts_store") + clean = ladder.RowResult(row=row, status="pass", reason_code=None, evidence={"cursor": "3"}, duration_ms=1) + # The probe manifest names this module, so the token can only leak through + # the bindings; the row itself stays clean. + token = "authority_e2e_ladder.py" + assert token not in json.dumps(clean.as_dict()) + + report = ladder.build_report( + [clean], + pending=(), + allow_unverified=True, + allow_pending=True, + environ=os.environ, + forbidden=[token], + ) + + assert report["rows"][0]["status"] == "pass" + assert report["bindings"]["privacy_violation"] is True + assert all(value is None for key, value in report["bindings"].items() if key != "privacy_violation") + assert token not in json.dumps(report) + assert report["summary"] == { + "pass": 1, + "fail": 0, + "unverified": 0, + "pending": 0, + "executed": 1, + "privacy_violations": 1, + } + assert report["exit_policy"]["exit_code"] == 1 + # No relaxation flag reaches a privacy violation. + assert ladder.exit_code_for(report["summary"], allow_unverified=True, allow_pending=True) == 1 + + +def test_ci_projection_skips_only_rows_whose_cli_e2e_twin_still_exists() -> None: + tree = ast.parse(CLI_E2E_TWIN_FILE.read_text(encoding="utf-8")) + defined = {node.name for node in ast.walk(tree) if isinstance(node, ast.FunctionDef)} + for row_id, twin in CLI_E2E_COVERAGE.items(): + ladder.row_by_id(row_id) + assert twin in defined, f"{row_id} is skipped by default but its twin {twin} no longer exists" + + +def test_list_prints_rows_and_pending_declarations( + capsys: pytest.CaptureFixture[str], +) -> None: + assert ladder.main(["--list"]) == 0 + listing = json.loads(capsys.readouterr().out) + assert [row["id"] for row in listing["rows"]] == [row.id for row in ladder.LADDER_ROWS] + assert [row["id"] for row in listing["pending"]] == [row.id for row in ladder.PENDING_ROWS] + assert all(row["status"] == "pending" for row in listing["pending"]) + + assert ladder.main(["--list", "--stage", "2c2"]) == 0 + stage_listing = json.loads(capsys.readouterr().out) + assert stage_listing["rows"] == [] + assert {row["stage"] for row in stage_listing["pending"]} == {"2c2"} diff --git a/tests/control_plane_ts/authority_store_readback_probe.ts b/tests/control_plane_ts/authority_store_readback_probe.ts new file mode 100644 index 0000000000..beb57cbd33 --- /dev/null +++ b/tests/control_plane_ts/authority_store_readback_probe.ts @@ -0,0 +1,198 @@ +/** + * Read-only probe over the file-backed AuthorityStore for the E2E stage ladder. + * + * The Python ladder writes through `python -m loopx.cli`; this probe proves the + * resulting candidate bytes decode through the production TypeScript store + * (`loadAuthority`, paged `scanCommitted`, `readReceipt`). It prints one JSON + * line and never writes: `storeIdentity()` is deliberately not called because + * it would mint an identity for a store that has none. + * + * Usage: + * node --experimental-strip-types tests/control_plane_ts/authority_store_readback_probe.ts \ + * --directory --goal-id [--receipt ] [--page-size ] + */ + +import type { JsonObject } from "../../loopx/control_plane/effect_program.ts"; +import type { AuthorityStore } from "../../loopx/control_plane/coordination/authority_store.ts"; +import { FileAuthorityStore } from "../../loopx/control_plane/coordination/file_authority_store.ts"; + +interface ProbeArguments { + directory: string; + goalId: string; + receipt: string | null; + pageSize: number; +} + +interface LoadSummary { + status: string; + reason_code: string | null; + cursor: string | null; + provider_revision: string | null; + head_handoff_mode: string | null; + head_todo_ids: string[]; + head_lease_count: number | null; +} + +interface ScanSummary { + status: string; + reason_code: string | null; + page_size: number; + pages: number; + operation_ids: string[]; + cursors: string[]; +} + +interface ReceiptSummary { + status: string; + reason_code: string | null; + operation_id: string; + cursor: string | null; + provider_revision: string | null; + receipt_count: number | null; + observation_id: string | null; +} + +class UsageError extends Error {} + +function readOption(argv: readonly string[], name: string): string | null { + const index = argv.indexOf(name); + if (index === -1) return null; + const value = argv[index + 1]; + if (value === undefined || value.startsWith("--")) { + throw new UsageError(`${name} requires a value`); + } + return value; +} + +function parseArguments(argv: readonly string[]): ProbeArguments { + const directory = readOption(argv, "--directory"); + const goalId = readOption(argv, "--goal-id"); + if (directory === null || goalId === null) { + throw new UsageError("--directory and --goal-id are required"); + } + const rawPageSize = readOption(argv, "--page-size"); + const pageSize = rawPageSize === null ? 2 : Number.parseInt(rawPageSize, 10); + if (!Number.isSafeInteger(pageSize) || pageSize < 1) { + throw new UsageError("--page-size must be a positive integer"); + } + return { directory, goalId, receipt: readOption(argv, "--receipt"), pageSize }; +} + +function stringList(value: unknown, key: string): string[] { + if (!Array.isArray(value)) return []; + const ids: string[] = []; + for (const entry of value) { + if (entry !== null && typeof entry === "object") { + const field = (entry as Record)[key]; + if (typeof field === "string") ids.push(field); + } + } + return ids; +} + +async function summarizeLoad(store: AuthorityStore): Promise { + const loaded = await store.loadAuthority(); + const summary: LoadSummary = { + status: loaded.status, + reason_code: null, + cursor: null, + provider_revision: null, + head_handoff_mode: null, + head_todo_ids: [], + head_lease_count: null, + }; + if (loaded.status === "loaded") { + const head: JsonObject = loaded.head; + summary.cursor = loaded.cursor; + summary.provider_revision = loaded.provider_revision; + summary.head_handoff_mode = typeof head.handoff_mode === "string" ? head.handoff_mode : null; + summary.head_todo_ids = stringList(head.todos, "todo_id"); + summary.head_lease_count = Array.isArray(head.leases) ? head.leases.length : null; + } else if (loaded.status !== "missing") { + summary.reason_code = loaded.reason_code; + } + return summary; +} + +async function summarizeScan(store: AuthorityStore, pageSize: number): Promise { + const summary: ScanSummary = { + status: "page", + reason_code: null, + page_size: pageSize, + pages: 0, + operation_ids: [], + cursors: [], + }; + let after: string | null = null; + for (;;) { + const page = await store.scanCommitted(after, pageSize); + if (page.status !== "page") { + summary.status = page.status; + summary.reason_code = page.reason_code; + return summary; + } + summary.pages += 1; + for (const transaction of page.transactions) { + summary.operation_ids.push(transaction.operation_id); + summary.cursors.push(transaction.cursor); + } + if (!page.has_more || page.next_cursor === after) return summary; + after = page.next_cursor; + } +} + +async function summarizeReceipt(store: AuthorityStore, operationId: string): Promise { + const receipt = await store.readReceipt(operationId); + const summary: ReceiptSummary = { + status: receipt.status, + reason_code: null, + operation_id: operationId, + cursor: null, + provider_revision: null, + receipt_count: null, + observation_id: null, + }; + if (receipt.status === "found") { + summary.cursor = receipt.cursor; + summary.provider_revision = receipt.provider_revision; + summary.receipt_count = receipt.receipts.length; + const first = receipt.receipts[0] as Record | undefined; + summary.observation_id = typeof first?.observation_id === "string" ? first.observation_id : null; + } else if (receipt.status !== "missing") { + summary.reason_code = receipt.reason_code; + } + return summary; +} + +async function main(argv: readonly string[]): Promise { + let parsed: ProbeArguments; + try { + parsed = parseArguments(argv); + } catch (error) { + if (error instanceof UsageError) { + process.stderr.write(`${error.message}\n`); + return 2; + } + throw error; + } + const store = new FileAuthorityStore(parsed.directory, parsed.goalId); + const result = { + schema_version: "loopx_authority_store_readback_probe_v0", + goal_id: parsed.goalId, + load: await summarizeLoad(store), + scan: await summarizeScan(store, parsed.pageSize), + receipt: parsed.receipt === null ? null : await summarizeReceipt(store, parsed.receipt), + }; + process.stdout.write(`${JSON.stringify(result)}\n`); + return 0; +} + +main(process.argv.slice(2)).then( + (code) => { + process.exitCode = code; + }, + (error: unknown) => { + process.stderr.write(`${error instanceof Error ? error.message : String(error)}\n`); + process.exitCode = 1; + }, +); diff --git a/tsconfig.control-plane.json b/tsconfig.control-plane.json index 427768a5d6..342253c704 100644 --- a/tsconfig.control-plane.json +++ b/tsconfig.control-plane.json @@ -70,6 +70,7 @@ "tests/control_plane_ts/nokv_authority_store.test.ts", "tests/control_plane_ts/nokv_jsonl_transport.test.ts", "tests/control_plane_ts/nokv_stage2a_qualification_harness.test.ts", + "tests/control_plane_ts/authority_store_readback_probe.ts", "tests/control_plane_ts/postgresql_authority_store.integration.test.ts", "tests/control_plane_ts/delivery_continuity.test.ts", "tests/control_plane_ts/delivery_workspace.test.ts",