From a38902f517da39f617a7ffc0a4645b7adc7eb525 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Mon, 28 Sep 2026 00:37:12 -0700 Subject: [PATCH 01/16] test(authority): characterize independent JSON snapshots Signed-off-by: Lihua <1017343802@qq.com> --- .../authority_journal_scan.test.ts | 29 ++++++++++++++ .../authority_state_log.test.ts | 34 ++++++++++++++++ .../sqlite_authority_store.test.ts | 40 +++++++++++++++++++ 3 files changed, 103 insertions(+) diff --git a/tests/control_plane_ts/authority_journal_scan.test.ts b/tests/control_plane_ts/authority_journal_scan.test.ts index bf3dd6056d..7c89328c34 100644 --- a/tests/control_plane_ts/authority_journal_scan.test.ts +++ b/tests/control_plane_ts/authority_journal_scan.test.ts @@ -60,6 +60,35 @@ test("scan requests reject malformed types before a provider effect", () => { } }); +test("scan pages copy every mutable JSON container and retain row metadata", () => { + const rows = ["1", "2"].map(cursor => ({ + ...row(cursor), + projection: {...row(cursor).projection, padding: "p".repeat(1024 * 1024), + ["__proto__"]: {marker: "data"}, negative: -0, holes: Array(2)}, + events: [{sequence: cursor, detail: {marker: "event"}}], + receipts: [{accepted: true, detail: {marker: "receipt"}}], + provenance: {marker: "retained metadata"}, + })); + const observed = {...head("2"), head: rows[1]!.projection}; + const page = scan(null, 2).page(rows, observed); + assert.equal(page.status, "page"); + if (page.status !== "page") return; + assert.deepEqual(page.transactions, rows); + assert.deepEqual(Object.keys(page.transactions[0]!), Object.keys(rows[0]!)); + assert.deepEqual(Object.keys(page.transactions[0]!.projection), Object.keys(rows[0]!.projection)); + const first = page.transactions[0]!; + assert.equal(Object.is(first.projection.negative, -0), true); + assert.equal(0 in (first.projection.holes as unknown[]), false); + (first.projection["__proto__"] as {marker: string}).marker = "edited"; + (first.events[0]!.detail as {marker: string}).marker = "edited"; + (first.receipts[0]!.detail as {marker: string}).marker = "edited"; + assert.deepEqual(rows[0]!.projection["__proto__"], {marker: "data"}); + assert.deepEqual(rows[0]!.events[0]!.detail, {marker: "event"}); + assert.deepEqual(rows[0]!.receipts[0]!.detail, {marker: "receipt"}); + assert.deepEqual(page.transactions[1], rows[1]); + assert.equal(({} as Record).marker, undefined); +}); + test("transaction decoder requires a canonical positive decimal cursor", () => { for (const cursor of ["0", "01", "-1", "1.0", "unknown", 1, null]) { assert.throws(() => decodeAuthorityTransaction({...row("1"), cursor}), /cursor/); diff --git a/tests/control_plane_ts/authority_state_log.test.ts b/tests/control_plane_ts/authority_state_log.test.ts index 979c153533..200ca22b8c 100644 --- a/tests/control_plane_ts/authority_state_log.test.ts +++ b/tests/control_plane_ts/authority_state_log.test.ts @@ -238,6 +238,40 @@ test("private replay keeps exact proofs while copying only changed paths", async assert.equal(replay.canonicalJson(), stable); }); +test("replay snapshots isolate mutable JSON while retaining primitive edge cases", async () => { + const {AuthorityStateReplay} = await import("../../loopx/control_plane/coordination/authority_state_log.ts"); + const initial = { + padding: "p".repeat(1024 * 1024), + negative: -0, + holes: Array(2), + unicode: "\ud800🎛", + ["__proto__"]: {marker: "data"}, + nested: {rows: [{value: 1}, {value: 2}]}, + }; + const replay = new AuthorityStateReplay(initial); + const first = replay.snapshot(), second = replay.snapshot(); + assert.deepEqual(first, initial); + assert.equal(Object.is(first.negative, -0), true); + assert.equal(0 in (first.holes as unknown[]), false); + assert.equal(Object.hasOwn(first, "__proto__"), true); + assert.equal(Object.getPrototypeOf(first), Object.prototype); + const rows = (value: Record) => + (value.nested as {rows: {value: number}[]}).rows; + rows(first)[0]!.value = 99; + rows(first).push({value: 3}); + (first["__proto__"] as {marker: string}).marker = "edited"; + assert.deepEqual(rows(second), [{value: 1}, {value: 2}]); + assert.deepEqual(second["__proto__"], {marker: "data"}); + assert.deepEqual(replay.snapshot(), initial); + replay.apply({schema_version: AUTHORITY_STATE_DELTA_SCHEMA, operations: [ + {op: "set", path: ["nested", "rows"], value: [{value: 4}]}, + ]}); + assert.deepEqual(rows(second), [{value: 1}, {value: 2}]); + assert.deepEqual(rows(replay.snapshot()), [{value: 4}]); + assert.equal(second.padding, initial.padding); + assert.equal(({} as Record).marker, undefined); +}); + test("replay encoding preserves canonical Unicode, sparse values, numeric keys and strict boundaries", async () => { const {AuthorityStateReplay} = await import("../../loopx/control_plane/coordination/authority_state_log.ts"); const {createHash} = await import("node:crypto"); diff --git a/tests/control_plane_ts/sqlite_authority_store.test.ts b/tests/control_plane_ts/sqlite_authority_store.test.ts index 769eb815dc..9eeb3f68bc 100644 --- a/tests/control_plane_ts/sqlite_authority_store.test.ts +++ b/tests/control_plane_ts/sqlite_authority_store.test.ts @@ -111,6 +111,46 @@ test("SQLite first empty projection still derives its root digest", async t => { assert.equal((await store.readReceipt("empty-root")).status, "found"); }); +test("SQLite scan snapshots stay independent across checkpoint boundaries", async t => { + const {store} = await fixture(t); + let revision: string | null = null; + const projection = (ordinal: number) => ({ + padding: "p".repeat(8192), + ["__proto__"]: {marker: "data"}, + nested: {rows: [{ordinal}]}, + }); + for (let ordinal = 1; ordinal <= 70; ordinal++) { + const committed = await store.commitAuthority({expected_provider_revision: revision, + operation_id: "snapshot-" + ordinal, next_projection: projection(ordinal), + events: [{ordinal}], receipts: [{ordinal}]}); + assert.equal(committed.status, "applied"); + if (committed.status !== "applied") return; + revision = committed.provider_revision; + } + const page = await store.scanCommitted("62", 5); + assert.equal(page.status, "page"); + if (page.status !== "page") return; + assert.deepEqual(page.transactions.map(row => row.cursor), ["63", "64", "65", "66", "67"]); + for (const [index, row] of page.transactions.entries()) { + assert.deepEqual(row.projection, projection(63 + index)); + assert.deepEqual(row.receipts, [{ordinal: 63 + index}]); + } + const first = page.transactions[0]!.projection; + (first.nested as {rows: {ordinal: number}[]}).rows[0]!.ordinal = 900; + (first["__proto__"] as {marker: string}).marker = "edited"; + assert.deepEqual(page.transactions[1]!.projection, projection(64)); + const again = await store.scanCommitted("62", 5); + assert.equal(again.status, "page"); + if (again.status === "page") { + assert.deepEqual(again.transactions[0]!.projection, projection(63)); + assert.equal(Object.getPrototypeOf(again.transactions[0]!.projection), Object.prototype); + } + const head = await store.loadAuthority(); + assert.equal(head.status, "loaded"); + if (head.status === "loaded") assert.deepEqual(head.head, projection(70)); + assert.equal(({} as Record).marker, undefined); +}); + test("SQLite head continuity is independent of retained history", {timeout: 30000}, async t => { const {store} = await fixture(t); assert.equal((await store.storeIdentity()).status, "available"); From aabf559796c44b29a8bc96fc64110fdf12abfee8 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Mon, 28 Sep 2026 00:37:12 -0700 Subject: [PATCH 02/16] perf(authority): avoid repeated primitive copies in journal pages Signed-off-by: Lihua <1017343802@qq.com> --- .../coordination/authority_journal_scan.ts | 6 +++-- .../coordination/authority_state_log.ts | 7 +++++- .../coordination/authority_store_codec.ts | 23 ++++++++++++++++--- .../coordination/sqlite_authority_store.ts | 2 +- 4 files changed, 31 insertions(+), 7 deletions(-) diff --git a/loopx/control_plane/coordination/authority_journal_scan.ts b/loopx/control_plane/coordination/authority_journal_scan.ts index 97bbb752ba..49c4f9a322 100644 --- a/loopx/control_plane/coordination/authority_journal_scan.ts +++ b/loopx/control_plane/coordination/authority_journal_scan.ts @@ -2,7 +2,7 @@ * Storage effects and snapshot acquisition remain with each provider. */ import type {AuthorityStoreCommittedTransaction, AuthorityStoreHead, AuthorityStoreReadFailure, AuthorityStoreScanResult} from "./authority_store.ts"; -import {AuthorityStoreProtocolError, canonicalAuthorityBytes, parseAuthorityCursor} from "./authority_store_codec.ts"; +import {AuthorityStoreProtocolError, canonicalAuthorityBytes, copyAuthorityJson, parseAuthorityCursor} from "./authority_store_codec.ts"; export class AuthorityJournalScan { readonly after: string | null; @@ -57,7 +57,9 @@ export class AuthorityJournalScan { throw new AuthorityStoreProtocolError("committed scan head lineage is invalid"); } } - const transactions = structuredClone(rows.slice(0, this.limit)); + // Each caller owns the mutable JSON containers; immutable primitive values + // need no additional byte copy for every retained projection in the page. + const transactions = copyAuthorityJson(rows.slice(0, this.limit)) as AuthorityStoreCommittedTransaction[]; return {status: "page", transactions, next_cursor: transactions.at(-1)?.cursor ?? this.after, has_more: rows.length > this.limit}; } diff --git a/loopx/control_plane/coordination/authority_state_log.ts b/loopx/control_plane/coordination/authority_state_log.ts index df299d435e..154c6e113b 100644 --- a/loopx/control_plane/coordination/authority_state_log.ts +++ b/loopx/control_plane/coordination/authority_state_log.ts @@ -22,6 +22,7 @@ import { canonicalAuthorityJson, canonicalAuthorityObject, canonicalAuthoritySha256, + copyAuthorityJson, isAuthorityJsonObject, } from "./authority_store_codec.ts"; @@ -215,7 +216,11 @@ export class AuthorityStateReplay { this.#state = next; } - snapshot(): JsonObject { return structuredClone(this.#state); } + snapshot(): JsonObject { + // Copy every mutable JSON container; immutable strings need no new byte + // allocation for each historical row carrying the same large value. + return copyAuthorityJson(this.#state) as JsonObject; + } /** Immutable canonical text; storage readers may parse it into independent rows. */ canonicalJson(): string { return this.#encode(this.#state); } diff --git a/loopx/control_plane/coordination/authority_store_codec.ts b/loopx/control_plane/coordination/authority_store_codec.ts index 06ec9e553f..a99c6fbddb 100644 --- a/loopx/control_plane/coordination/authority_store_codec.ts +++ b/loopx/control_plane/coordination/authority_store_codec.ts @@ -46,6 +46,15 @@ export function canonicalAuthorityJson( value: unknown, stack = new Set(), ): unknown { + return cloneAuthorityJson(value, stack, true); +} + +/** Own JSON containers without copying immutable primitives or reordering keys. */ +export function copyAuthorityJson(value: unknown): unknown { + return cloneAuthorityJson(value, new Set(), false); +} + +function cloneAuthorityJson(value: unknown, stack: Set, canonicalKeys: boolean): unknown { if (value === null || typeof value === "string" || typeof value === "boolean") { return value; } @@ -59,7 +68,13 @@ export function canonicalAuthorityJson( if (stack.has(value)) throw new AuthorityStoreProtocolError("JSON value must be acyclic"); stack.add(value); try { - return value.map((item) => canonicalAuthorityJson(item, stack)); + if (canonicalKeys) return value.map((item) => cloneAuthorityJson(item, stack, true)); + const copy: unknown[] = Array(value.length); + for (const key of Object.keys(value)) Object.defineProperty(copy, key, { + value: cloneAuthorityJson(Reflect.get(value, key), stack, false), + writable: true, enumerable: true, configurable: true, + }); + return copy; } finally { stack.delete(value); } @@ -74,10 +89,12 @@ export function canonicalAuthorityJson( if (stack.has(value)) throw new AuthorityStoreProtocolError("JSON value must be acyclic"); stack.add(value); try { + const keys = Object.keys(value); + if (canonicalKeys) keys.sort(authorityUnicodeCompare); return Object.fromEntries( - Object.keys(value).sort(authorityUnicodeCompare).map((key) => [ + keys.map((key) => [ key, - canonicalAuthorityJson(value[key], stack), + cloneAuthorityJson(value[key], stack, canonicalKeys), ]), ); } finally { diff --git a/loopx/control_plane/coordination/sqlite_authority_store.ts b/loopx/control_plane/coordination/sqlite_authority_store.ts index 024222d73d..86c147d971 100644 --- a/loopx/control_plane/coordination/sqlite_authority_store.ts +++ b/loopx/control_plane/coordination/sqlite_authority_store.ts @@ -594,7 +594,7 @@ export class SqliteAuthorityStore implements AuthorityStore { if (rows.length) this.verifiedRange(db, BigInt(String(rows[0]!.sequence)), BigInt(String(rows[rows.length - 1]!.sequence)), (row, replay, identity) => { verified.push({cursor: row.cursor.toString(), provider_revision: `${identity}:${row.cursor}`, - operation_id: row.operation_id, projection: JSON.parse(replay.canonicalJson()) as JsonObject, events: row.events, receipts: row.receipts}); + operation_id: row.operation_id, projection: replay.snapshot(), events: row.events, receipts: row.receipts}); }); return scan.page(verified, head === null ? null : {cursor: head.state.cursor.toString(), provider_revision: head.provider_revision, head: head.state.projection}); From 64a75ca94136e756594d910fb6b2d6c23d162323 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Mon, 28 Sep 2026 01:52:03 -0700 Subject: [PATCH 03/16] chore(semantics): refresh contract registry I/O location Signed-off-by: Lihua <1017343802@qq.com> --- loopx/semantics/project_registry_io_manifest_v1.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 09c271ee70..f11501031e 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -887,7 +887,7 @@ }, { "site": "loopx/contract.py::.check_contract::codec_read:load_registry#1", - "line": 1014, + "line": 1027, "column": 20, "kind": "codec_read", "api": "load_registry", From 1df14ecc6ad60ecd5a6ea0d30e496ba874afb85e Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Mon, 28 Sep 2026 19:36:31 -0700 Subject: [PATCH 04/16] test(quota): align canonical fixtures and host ACK settlement Signed-off-by: Lihua <1017343802@qq.com> --- .../test_quota_plan_observation_payload.py | 4 +- tests/control_plane/test_quota_selection.py | 4 +- .../test_quota_settlement_cli.py | 57 +++++++++++-------- 3 files changed, 40 insertions(+), 25 deletions(-) diff --git a/tests/control_plane/test_quota_plan_observation_payload.py b/tests/control_plane/test_quota_plan_observation_payload.py index c6df3768a6..d9e5f26ac6 100644 --- a/tests/control_plane/test_quota_plan_observation_payload.py +++ b/tests/control_plane/test_quota_plan_observation_payload.py @@ -67,7 +67,8 @@ def test_real_cli_compact_and_full_detail_preserve_canonical_todos(tmp_path, mon "role": "agent" if i < 40 else "user", "status": "open" if i % 3 else "done", "done": i % 3 == 0, "text": f"Retained work {i}", "note": "exact metadata🙂" * 100, "archive_state": "active", "source_section": "Agent Todo" if i < 40 else "User Todo", - "index": i + 1, "task_class": "advancement_task"} for i in range(60)] + "index": i + 1, "task_class": "advancement_task" if i < 40 else "user_action", + **({"goal_bound": True} if i >= 40 else {})} for i in range(60)] projection = build_todo_runtime_shadow_projection(goal_id="example", todos=records, leases=[], handoff_mode="soft_claim") initialize_canonical_authority(runtime, "example", projection, state_path=state, provider=provider) @@ -97,6 +98,7 @@ def call(mode, *details): if record["role"] == role: assert actual[record["todo_id"]]["note"] == record["note"] assert actual[record["todo_id"]]["status"] == record["status"] + assert actual[record["todo_id"]]["task_class"] == record["task_class"] assert len(json.dumps(compact)) < len(json.dumps(full)) / 2 assert not list((runtime / "goals" / "example" / "runs").glob("*.json*")) finally: diff --git a/tests/control_plane/test_quota_selection.py b/tests/control_plane/test_quota_selection.py index b10c6f481b..642bd2361c 100644 --- a/tests/control_plane/test_quota_selection.py +++ b/tests/control_plane/test_quota_selection.py @@ -59,10 +59,12 @@ def test_public_quota_keeps_gate_from_real_canonical_provider(tmp_path: Path, di runtime, registry, state = tmp_path / "runtime", tmp_path / "registry.json", tmp_path / "state.md" history = "\n".join(f"- [x] [P2] Completed synthetic work {i}.\n" f" " for i in range(12)) + # A goal-wide gate cannot also bind continuation through a legacy agent claim. + claim = " claimed_by=agent-a" if scope == "blocks_agent=agent-b" else "" state.write_text("---\nstatus: active\n---\n# Goal\n## Objective\nDeliver a checked change.\n\n" "## Agent Todo\n" + history + "\n\n## User Todo\n" "- [ ] [P0] Owner approval is required.\n" - f" \n") + f" \n") write_fixture_registry(project=tmp_path, runtime_root=runtime, registry_path=registry, goal_id="goal-scope", domain="quota-scope", adapter_kind="generic_project_goal_v0", state_file=str(state), registered_agents=["agent-a", "agent-b"], quota_allowed_slots=None) diff --git a/tests/control_plane/test_quota_settlement_cli.py b/tests/control_plane/test_quota_settlement_cli.py index 70ac39c78b..26fadf4d94 100644 --- a/tests/control_plane/test_quota_settlement_cli.py +++ b/tests/control_plane/test_quota_settlement_cli.py @@ -2290,21 +2290,7 @@ def test_standard_codex_app_settlement_is_receipted_and_idempotent( assert refresh_replay["idempotent_replay"] is True assert _classification_count(runtime, "validated_progress") == 1 - fresh_guard_rc, fresh_guard = _run_cli( - registry_path, - runtime, - "quota", - "should-run", - "--codex-app", - "--goal-id", - GOAL_ID, - "--agent-id", - AGENT_ID, - "--scan-path", - str(project), - ) - assert fresh_guard_rc == 0, fresh_guard - assert fresh_guard["selected_todo"]["todo_id"] == successor_id + # Finish this Turn and its ACK before the next guard owns scheduler authority. spend_args = ( "quota", @@ -2392,6 +2378,14 @@ def test_standard_codex_app_settlement_is_receipted_and_idempotent( ) assert fresh_turn_rc == 0, fresh_turn assert fresh_turn["selected_todo"]["todo_id"] == successor_id + stale_ack_rc, stale_ack = _run_cli( + registry_path, runtime, *settled_ack_hint["cli_args"], + ) + assert stale_ack_rc == 1, stale_ack + assert stale_ack["error_code"] == "SCHEDULER_FOLLOWUP_HEARTBEAT_RECEIPT_STALE" + assert stale_ack["write_performed"] is False + assert stale_ack["scheduler_state_mutated"] is False + assert _spend_run_count(runtime) == 1 def _assert_material_monitor_writeback_can_add_workspace_before_spend( @@ -4931,7 +4925,7 @@ def test_pending_action_selection_does_not_commit_after_new_user_gate( assert all(not event["details"].get("settlement_effect_id") for event in events) -def test_todoless_autonomous_replan_settles_quota_refresh_spend_chain( +def test_todoless_replan_keeps_scheduler_ack_outside_agent_settlement( tmp_path: Path, ) -> None: project, runtime, registry_path = _write_fixture(tmp_path) @@ -4975,6 +4969,7 @@ def test_todoless_autonomous_replan_settles_quota_refresh_spend_chain( assert identity["replan_obligation_id"] == obligation_id assert "todo_id" not in identity cli_channel = guard["interaction_contract"]["cli_channel"] + assert cli_channel["settlement_plan"]["host_handoff"]["inside_agent_settlement"] is False plan_identity = cli_channel["settlement_plan"]["identity"] assert plan_identity["binding_kind"] == identity["binding_kind"] assert plan_identity["binding_id"] == identity["binding_id"] @@ -4986,7 +4981,9 @@ def test_todoless_autonomous_replan_settles_quota_refresh_spend_chain( original_scheduler_ack_args = original_scheduler_ack_args[ original_scheduler_ack_args.index("quota"): ] - assert original_scheduler_ack_args[:2] == ["quota", "scheduler-ack-current"] + # The bounded hint may return an explicit ACK instead of resolving current state. + assert original_scheduler_ack_args[0] == "quota" + assert original_scheduler_ack_args[1] in {"scheduler-ack", "scheduler-ack-current"} assert "--turn-instance-id" in original_scheduler_ack_args assert turn_instance_id in original_scheduler_ack_args actions = cli_channel["next_cli_actions"] @@ -5082,6 +5079,11 @@ def test_todoless_autonomous_replan_settles_quota_refresh_spend_chain( ] == "autonomous_replan" assert _spend_run_count(runtime) == 1 + # Scheduler handoff is outside agent settlement. Its ACK may update only + # scheduler state; repeated execution cannot change Goal state or accounting. + state_path = project / ".codex" / "goals" / GOAL_ID / "ACTIVE_GOAL_STATE.md" + run_index = runtime / "goals" / GOAL_ID / "runs" / "index.jsonl" + before_state, before_runs = state_path.read_bytes(), run_index.read_bytes() ack_rc, ack = _run_cli( registry_path, runtime, @@ -5089,13 +5091,22 @@ def test_todoless_autonomous_replan_settles_quota_refresh_spend_chain( ) assert ack_rc == 0, ack assert ack["ok"] is True - assert ack["mode"] == "scheduler-ack-current" - assert ack["status"] == "heartbeat_settled_skip" - assert ack["idempotent_replay"] is True - assert ack["write_performed"] is False - assert ack["scheduler_state_mutated"] is False - assert ack["quota_spend_performed"] is False + assert ack["mode"] in {"scheduler-ack", "scheduler-ack-current"} + assert ack["registry_mutated"] is False assert ack["appended"] is False + scheduler_path = ack.get("scheduler_state_path") + scheduler_bytes = Path(scheduler_path).read_bytes() if scheduler_path else None + ack_replay_rc, ack_replay = _run_cli( + registry_path, runtime, *original_scheduler_ack_args, + ) + assert ack_replay_rc == 0, ack_replay + assert ack_replay["ok"] is True + assert ack_replay["registry_mutated"] is False + assert ack_replay["appended"] is False + assert state_path.read_bytes() == before_state + assert run_index.read_bytes() == before_runs + if scheduler_path: + assert Path(scheduler_path).read_bytes() == scheduler_bytes assert _spend_run_count(runtime) == 1 fresh_turn_id = "turn-autonomous-replan-settlement-2" From 53fab7f098fbe48736d1df008c5c9c74ceeba9c2 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Mon, 28 Sep 2026 21:10:30 -0700 Subject: [PATCH 05/16] test(attached-session): prove claim wait releases lifetime guard Signed-off-by: Lihua <1017343802@qq.com> --- tests/test_attached_session_goal_instance.py | 72 ++++++++++++++++---- 1 file changed, 59 insertions(+), 13 deletions(-) diff --git a/tests/test_attached_session_goal_instance.py b/tests/test_attached_session_goal_instance.py index 4a71910d64..86658812d9 100644 --- a/tests/test_attached_session_goal_instance.py +++ b/tests/test_attached_session_goal_instance.py @@ -4,10 +4,13 @@ import threading import time from pathlib import Path +from types import SimpleNamespace from typing import Any import pytest +import loopx.attached_session as attached_session_module + from loopx.attached_session import ( bind_attached_agent_session, claim_attached_agent_turn, @@ -525,12 +528,32 @@ def test_missing_registry_without_exact_history_preserves_legacy_mutation( assert "admitted_goal_instance_id" not in completed_turn -def test_claim_wait_does_not_hold_goal_lifetime_guard(tmp_path: Path) -> None: +def test_claim_wait_does_not_hold_goal_lifetime_guard( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: registry_path, _runtime_root, store, instance_a = _register(tmp_path) session_a = str(_bind(store, registry_path)["session"]["session_id"]) # type: ignore[index] - finished = threading.Event() + wait_entered = threading.Event() + release_wait = threading.Event() + claim_finished = threading.Event() + recreation_finished = threading.Event() + recreated: list[str] = [] + stale_claims: list[str] = [] errors: list[BaseException] = [] + def blocked_poll_wait(_seconds: float) -> None: + wait_entered.set() + if not release_wait.wait(timeout=10): + raise TimeoutError("claim poll wait was not released") + + # Replace only this module's wait hook, retaining real admission, storage, + # lifetime locks and recreation. Event order is the oracle, not wall time. + monkeypatch.setattr( + attached_session_module, "time", + SimpleNamespace(monotonic=time.monotonic, sleep=blocked_poll_wait), + ) + def wait_for_claim() -> None: try: claim_attached_agent_turn( @@ -543,22 +566,45 @@ def wait_for_claim() -> None: wait_seconds=0.5, ) except ValueError as exc: - if "stale_goal_instance" not in str(exc): + if "stale_goal_instance" in str(exc): + stale_claims.append(str(exc)) + else: errors.append(exc) + except BaseException as exc: + errors.append(exc) finally: - finished.set() + claim_finished.set() - thread = threading.Thread(target=wait_for_claim) - thread.start() - time.sleep(0.1) - started = time.monotonic() - _recreate(registry_path, instance_a) - elapsed = time.monotonic() - started - thread.join(timeout=2) + def recreate() -> None: + try: + recreated.append(_recreate(registry_path, instance_a)) + except BaseException as exc: + errors.append(exc) + finally: + recreation_finished.set() + + claim_thread = threading.Thread(target=wait_for_claim) + recreate_thread = threading.Thread(target=recreate) + claim_thread.start() + try: + assert wait_entered.wait(timeout=5) + assert not claim_finished.is_set() + recreate_thread.start() + assert recreation_finished.wait(timeout=5) + assert not release_wait.is_set() + assert not claim_finished.is_set() + finally: + release_wait.set() + claim_thread.join(timeout=5) + if recreate_thread.ident is not None: + recreate_thread.join(timeout=5) - assert elapsed < 0.4 - assert finished.is_set() + assert not claim_thread.is_alive() + assert not recreate_thread.is_alive() assert errors == [] + assert recreated and recreated[0] != instance_a + assert claim_finished.is_set() + assert len(stale_claims) == 1 def test_resume_writeback_holds_goal_lifetime_guard( From 8e5c3dca325d3986773c5afa914ada1a4afd7aff Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Mon, 28 Sep 2026 22:37:33 -0700 Subject: [PATCH 06/16] test(authority): verify bootstrap receipts after ambiguous response Signed-off-by: Lihua <1017343802@qq.com> --- .../test_local_authority_shadow_outbox.py | 113 ++++++++++++++++-- 1 file changed, 106 insertions(+), 7 deletions(-) diff --git a/tests/control_plane/test_local_authority_shadow_outbox.py b/tests/control_plane/test_local_authority_shadow_outbox.py index 0a1bf6d831..e5c5946ccd 100644 --- a/tests/control_plane/test_local_authority_shadow_outbox.py +++ b/tests/control_plane/test_local_authority_shadow_outbox.py @@ -8,20 +8,34 @@ from loopx.control_plane.coordination import local_authority_shadow_adapter as adapter from loopx.control_plane.coordination import local_authority_shadow_outbox as outbox +from loopx.control_plane.coordination.coordination_state_contract_generated import ( + COORDINATION_RUNTIME_SHADOW_BOOTSTRAP_RESULT_SCHEMA, +) from loopx.control_plane.coordination.local_authority_shadow_projection import ( ProjectionValueError, canonical_bytes, lease_partition_projection, partition_digest, sha256_digest, + source_effect_runtime_result, text_digest, todo_partition_projection, ) from loopx.control_plane.coordination.runtime_shadow import ( bootstrap_coordination_runtime_shadow, build_runtime_shadow_source_snapshot, + RuntimeInvoker, +) +from loopx.control_plane.coordination.shadow_management import ( + read_shadow_bootstrap_source_path, + read_shadow_management_state, + require_shadow_primary_write_allowed, + shadow_maintenance_lock_target, + shadow_management_directory, + shadow_management_state_path, ) -from loopx.control_plane.coordination.shadow_management import require_shadow_primary_write_allowed +from loopx.control_plane.effect_runtime import EffectRuntimeResponseAmbiguous, EffectRuntimeRejected +from loopx.file_lock import exclusive_mutation_file_lock from loopx.history import load_registry from loopx.registry import find_registry_goal @@ -29,7 +43,10 @@ GOAL_ID = "goal-outbox" -def _fixture(tmp_path: Path, *, bootstrap: bool = True) -> tuple[Path, Path, Path]: +def _fixture( + tmp_path: Path, *, bootstrap: bool = True, + runtime_invoker: RuntimeInvoker = source_effect_runtime_result, +) -> tuple[Path, Path, Path]: repo = tmp_path / "repo" repo.mkdir() state = repo / "ACTIVE_GOAL_STATE.md" @@ -75,15 +92,35 @@ def _fixture(tmp_path: Path, *, bootstrap: bool = True) -> tuple[Path, Path, Pat "schema_version": "loopx_coordination_runtime_shadow_config_v0", "enabled": True, "provider": "file_v0", }}} + def bootstrap_or_receipt(method: str, request: dict) -> dict: + try: + return runtime_invoker(method, request) + except EffectRuntimeResponseAmbiguous: + # The RPC deadline still applies. Do not resend an uncertain + # write; read its exact completed receipt after the existing + # maintenance lock releases, using the normal lock budget. + with exclusive_mutation_file_lock(shadow_maintenance_lock_target(runtime_root, GOAL_ID)): + journal = read_shadow_management_state(runtime_root, GOAL_ID) + if journal is None: + raise + assert journal["status"] == "active", journal + assert journal["operation"]["operation_id"] == request["operation_id"], journal + assert journal["operation"]["request_digest"] == sha256_digest(request), journal + binding = require_shadow_primary_write_allowed(runtime_root, GOAL_ID) + assert binding is not None, journal + assert read_shadow_bootstrap_source_path(runtime_root, GOAL_ID, binding) == state.resolve() + receipt = journal["result"] + assert receipt["operation_id"] == request["operation_id"], receipt + assert {key: receipt[key] for key in binding} == binding, receipt + return {"schema_version": COORDINATION_RUNTIME_SHADOW_BOOTSTRAP_RESULT_SCHEMA, **receipt} + result = bootstrap_coordination_runtime_shadow( goal=enabled_goal, runtime_root=runtime_root, goal_id=GOAL_ID, operation_id="bootstrap:outbox-test", source_version="source:initial", - projection=projection, source_snapshot=snapshot, + projection=projection, source_snapshot=snapshot, runtime_invoker=bootstrap_or_receipt, ) - # The managed Effect runtime may lose the first response after the - # durable bootstrap commit and retry the same operation. Windows CI is - # slow enough to exercise that path, so the public success contract is - # applied/recovered/replayed rather than applied-only. + # Outbox assertions require a verified bootstrap, including an exact + # durable result when the managed RPC response is ambiguous. assert result["status"] in {"applied", "recovered", "replayed"}, result assert result["operation_id"] == "bootstrap:outbox-test", result assert result["cursor"] == "1", result @@ -161,6 +198,68 @@ def _files(directory: Path) -> dict[str, bytes]: return {str(path.relative_to(directory)): path.read_bytes() for path in directory.rglob("*") if path.is_file()} +def test_outbox_fixture_reads_exact_bootstrap_receipt_after_lost_response(tmp_path: Path) -> None: + calls = [] + committed = {} + + def lose_response(method: str, request: dict) -> dict: + result = source_effect_runtime_result(method, request) + assert result["status"] in {"applied", "recovered", "replayed"}, result + calls.append(request) + path = shadow_management_state_path(Path(request["runtime_root"]), GOAL_ID) + committed["journal_bytes"] = path.read_bytes() + raise EffectRuntimeResponseAmbiguous(method, timeout=10) + + registry, state, runtime_root = _fixture(tmp_path, runtime_invoker=lose_response) + assert len(calls) == 1 # A receipt readback must never dispatch another write. + assert shadow_management_state_path(runtime_root, GOAL_ID).read_bytes() == committed["journal_bytes"] + _record_change(registry, state, runtime_root, "Continue after verified bootstrap readback.") + assert _drain(registry, runtime_root).reason_code is None + + +@pytest.mark.parametrize("damage", ["missing", "pending", "other_operation", "changed_request", "changed_manifest"]) +def test_outbox_fixture_rejects_unverified_bootstrap_after_lost_response(tmp_path: Path, damage: str) -> None: + calls = [] + + def lose_response(method: str, request: dict) -> dict: + calls.append(request) + if damage != "missing": + source_effect_runtime_result(method, request) + root = Path(request["runtime_root"]) + path = shadow_management_state_path(root, GOAL_ID) + journal = json.loads(path.read_text(encoding="utf-8")) + if damage == "pending": + journal.update(status="bootstrapping", binding=None, result=None) + journal["operation"]["phase"] = "prepared" + elif damage == "other_operation": + journal["operation"]["operation_id"] = "bootstrap:another-operation" + elif damage == "changed_request": + journal["operation"]["request_digest"] = "sha256:" + "0" * 64 + else: + manifest = next((shadow_management_directory(root, GOAL_ID) / "operations").glob("*/manifest.json")) + manifest.write_text("{}", encoding="utf-8") + path.write_text(json.dumps(journal), encoding="utf-8") + raise EffectRuntimeResponseAmbiguous(method, timeout=10) + + with pytest.raises(AssertionError): + _fixture(tmp_path, runtime_invoker=lose_response) + assert len(calls) == 1 + assert not _todo_dir(tmp_path / "runtime").exists() + + +def test_outbox_fixture_preserves_semantic_rejection_even_with_a_bootstrap_receipt(tmp_path: Path) -> None: + calls = [] + + def reject_response(method: str, request: dict) -> dict: + source_effect_runtime_result(method, request) + calls.append(request) + raise EffectRuntimeRejected("deliberate semantic rejection", diagnostic_code="invalid_request") + + with pytest.raises(AssertionError, match="deliberate semantic rejection"): + _fixture(tmp_path, runtime_invoker=reject_response) + assert len(calls) == 1 + + def test_capture_records_prepared_then_committed_and_skips_prose_only_writes(tmp_path: Path) -> None: registry, state, runtime_root = _fixture(tmp_path) original = state.read_text(encoding="utf-8") From 0097931ea349c404f0f430f3b088b4f669da0120 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Tue, 29 Sep 2026 01:07:24 -0700 Subject: [PATCH 07/16] chore(semantics): align CLI registry census with integrated main Signed-off-by: Lihua <1017343802@qq.com> --- loopx/semantics/project_registry_io_manifest_v1.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 05cb12cb81..9ed1c447c8 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -487,7 +487,7 @@ }, { "site": "loopx/cli.py::.main::codec_read:load_project_registry#1", - "line": 846, + "line": 835, "column": 17, "kind": "codec_read", "api": "load_project_registry", From 106b923a57c60d411dba9276c37d1f927e0023bc Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Tue, 29 Sep 2026 19:15:26 -0700 Subject: [PATCH 08/16] test: align CI guards with generated digest and hook dispatch Signed-off-by: Lihua <1017343802@qq.com> --- tests/architecture/test_turn_contract_generation.py | 4 +++- tests/control_plane/test_prompt_upgrade_hook.py | 10 +++++++++- 2 files changed, 12 insertions(+), 2 deletions(-) diff --git a/tests/architecture/test_turn_contract_generation.py b/tests/architecture/test_turn_contract_generation.py index e14adf24c5..8741eeace2 100644 --- a/tests/architecture/test_turn_contract_generation.py +++ b/tests/architecture/test_turn_contract_generation.py @@ -260,7 +260,9 @@ def test_new_independent_twin_cannot_hide_behind_generated_pair(monkeypatch): ) assert counts is not None raw, generated, maintained, budget = map(int, counts.groups()) - assert generated == 1 and raw == maintained + generated + # The Turn contract and content digest are the two verified generated + # twins. A new authored twin still consumes the independent module budget. + assert generated == 2 and raw == maintained + generated from loopx.semantics.inventory import SourceFile target = smoke["check_dual_runtime_twins"].__globals__ diff --git a/tests/control_plane/test_prompt_upgrade_hook.py b/tests/control_plane/test_prompt_upgrade_hook.py index a8ade637cb..935879f1b7 100644 --- a/tests/control_plane/test_prompt_upgrade_hook.py +++ b/tests/control_plane/test_prompt_upgrade_hook.py @@ -167,8 +167,16 @@ def test_live_decision_adds_only_existing_required_read_channel(tmp_path, monkey assert envelope["compaction"]["budget_bytes"] == 8192 + 1536 assert envelope["compaction"]["hook_prompt_budget_bytes"] == 1536 assert build_turn_envelope(baseline)["compaction"]["budget_bytes"] == 8192 + assert baseline.get("turn_start_capability_hook_dispatch") is None + dispatch = pending["turn_start_capability_hook_dispatch"] + assert len(dispatch["required_reads"]) == 1 + dispatched_read = dispatch["required_reads"][0] + assert dispatched_read["hook_id"] == "heartbeat.prompt_upgrade" + assert dispatched_read["capability_id"] == "automation-prompt-upgrade" + assert {key: dispatched_read[key] for key in hint} == hint for key in baseline.keys() | pending.keys(): - if key not in {"required_reads", "interaction_contract", "protocol_action_packet"}: + if key not in {"required_reads", "interaction_contract", "protocol_action_packet", + "turn_start_capability_hook_dispatch"}: assert pending.get(key) == baseline.get(key), key _set_fixture_prompt(path, database, desired) assert build_live_quota_should_run_decision(status, **kwargs) == baseline From e4e754982883306d98e9ebb943ff52b198c35af1 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Tue, 29 Sep 2026 20:05:25 -0700 Subject: [PATCH 09/16] test: publish Host descendant marker atomically Signed-off-by: Lihua <1017343802@qq.com> --- tests/control_plane_ts/host_process.test.ts | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/tests/control_plane_ts/host_process.test.ts b/tests/control_plane_ts/host_process.test.ts index cbea7db0d4..23b235a8e2 100644 --- a/tests/control_plane_ts/host_process.test.ts +++ b/tests/control_plane_ts/host_process.test.ts @@ -49,10 +49,12 @@ for (const mode of ["timeout", "abort", "leader_exit", "closed_pipes"] as const) t.after(() => rm(root, {recursive: true, force: true})); const marker = join(root, "counter"); // Ignore TERM so the test proves escalation and does not merely observe a - // cooperative child. Its marker is the semantic oracle, not a PID lookup. + // cooperative child. Publish the marker atomically: a kill during a write + // must not look like a surviving child, but each completed tick stays visible. const child = `const fs=require('fs');let n=0;process.on('SIGTERM',()=>{}); - fs.writeFileSync(${JSON.stringify(marker)},String(n)); - setInterval(()=>fs.writeFileSync(${JSON.stringify(marker)},String(++n)),10)`; + const marker=${JSON.stringify(marker)}, staged=marker+'.next'; + const publish=()=>{fs.writeFileSync(staged,String(n));fs.renameSync(staged,marker)}; + publish();setInterval(()=>{n++;publish()},10)`; const script = `const{spawn}=require('child_process');const fs=require('fs'); spawn(process.execPath,['-e',${JSON.stringify(child)}],{stdio:${JSON.stringify(mode === "closed_pipes" ? "ignore" : "inherit")}}); const timer=setInterval(()=>{if(fs.existsSync(${JSON.stringify(marker)})){ From 70f4d3684e23217be4e7ef3af144c3f7ae3aecf3 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Wed, 30 Sep 2026 14:09:14 -0700 Subject: [PATCH 10/16] test(runtime): make source churn race deterministic Signed-off-by: Lihua <1017343802@qq.com> --- .../test_turn_journal_runtime_readiness.py | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/tests/control_plane/test_turn_journal_runtime_readiness.py b/tests/control_plane/test_turn_journal_runtime_readiness.py index cb0a2b0165..1c0c76ea7b 100644 --- a/tests/control_plane/test_turn_journal_runtime_readiness.py +++ b/tests/control_plane/test_turn_journal_runtime_readiness.py @@ -1,6 +1,7 @@ from __future__ import annotations from pathlib import Path +from threading import Event from typing import Any import pytest @@ -117,12 +118,16 @@ def test_runtime_fingerprint_rescans_when_a_snapshotted_file_disappears_while_re later.write_text("export const later = true;\n", encoding="utf-8") original_read_bytes = Path.read_bytes reads: list[str] = [] + removed = Event() def remove_later_after_first_read(path: Path) -> bytes: + if path == later: + assert removed.wait(timeout=5), "first source read did not remove later file" reads.append(path.name) content = original_read_bytes(path) if path == first and later.exists(): later.unlink() + removed.set() return content monkeypatch.setattr(effect_runtime, "_control_plane_root", lambda: tmp_path) @@ -142,17 +147,22 @@ def _install_persistent_stat_read_churn( original_scan = effect_runtime._scan_runtime_source_files original_read_bytes = Path.read_bytes scans: list[tuple[str, ...]] = [] + removed = Event() def restore_then_scan(root: Path) -> tuple[str, ...]: + removed.clear() later.write_text("export const later = true;\n", encoding="utf-8") files = original_scan(root) scans.append(files) return files def remove_later_after_first_read(path: Path) -> bytes: + if path == later: + assert removed.wait(timeout=5), "first source read did not remove later file" content = original_read_bytes(path) if path == first: later.unlink() + removed.set() return content monkeypatch.setattr(effect_runtime, "_control_plane_root", lambda: tmp_path) From 6e897e48948e7588535db0463c3a1505dcdd1939 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Wed, 30 Sep 2026 14:43:55 -0700 Subject: [PATCH 11/16] fix(work-items): reuse canonical SHA-256 matcher for lease workspace Signed-off-by: Lihua <1017343802@qq.com> --- loopx/control_plane/work_items/task_lease_workspace.ts | 3 ++- tests/control_plane_ts/content_digest_single_owner.test.ts | 1 + 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/loopx/control_plane/work_items/task_lease_workspace.ts b/loopx/control_plane/work_items/task_lease_workspace.ts index bf7114f257..3c6778dafe 100644 --- a/loopx/control_plane/work_items/task_lease_workspace.ts +++ b/loopx/control_plane/work_items/task_lease_workspace.ts @@ -5,6 +5,7 @@ import {realpath, stat, readFile, lstat, opendir} from "node:fs/promises"; import {platform} from "node:os"; import {createHash} from "node:crypto"; import {isAbsolute, resolve} from "node:path"; +import {BARE_SHA256_PATTERN} from "../content_digest.ts"; import type {JsonObject} from "../effect_program.ts"; import {EffectRuntimeRequestError} from "../effect_runtime_errors.ts"; import {requireJsonObject} from "../runtime_decode.ts"; @@ -23,7 +24,7 @@ export function leaseWorkspace(value: unknown): LeaseWorkspace | null { if (value == null) return null; const row = requireJsonObject(value, "lease worktree identity"); if (Object.keys(row).sort().join(",") !== "common_directory,host,repository,worktree" || - [row.host, row.common_directory, row.worktree].some(v => typeof v !== "string" || !/^[a-f0-9]{64}$/u.test(v))) { + [row.host, row.common_directory, row.worktree].some(v => typeof v !== "string" || !BARE_SHA256_PATTERN.test(v))) { throw new EffectRuntimeRequestError("invalid lease worktree identity"); } const repository = leaseWriteRepository(row.repository); diff --git a/tests/control_plane_ts/content_digest_single_owner.test.ts b/tests/control_plane_ts/content_digest_single_owner.test.ts index 47b3c0d4c0..e46d52ab2b 100644 --- a/tests/control_plane_ts/content_digest_single_owner.test.ts +++ b/tests/control_plane_ts/content_digest_single_owner.test.ts @@ -126,6 +126,7 @@ const CANONICAL_CONSUMERS = [ "control_plane/work_items/task_lease_acquire.ts", "control_plane/work_items/task_lease_lifecycle.ts", "control_plane/work_items/task_lease_lifecycle_request.ts", + "control_plane/work_items/task_lease_workspace.ts", ]; function packageFiles(dir: string, base = ""): string[] { From 786b5edbf36b3fe80b39337581a9adf8b8785b87 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Wed, 30 Sep 2026 15:42:30 -0700 Subject: [PATCH 12/16] test(cli): align source-first usage notice with v5 contract Signed-off-by: Lihua <1017343802@qq.com> --- tests/control_plane/test_source_cli_entrypoint.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/control_plane/test_source_cli_entrypoint.py b/tests/control_plane/test_source_cli_entrypoint.py index fdb5fcd150..c0cce900d5 100644 --- a/tests/control_plane/test_source_cli_entrypoint.py +++ b/tests/control_plane/test_source_cli_entrypoint.py @@ -311,7 +311,7 @@ def test_source_first_usage_disclosure_keeps_json_pure_and_does_not_send(tmp_pat assert json.loads(result.stdout)["ok"] is True assert "random installation ID" in result.stderr stored = json.loads((state / "usage-ping.json").read_text()) - assert stored["notice"]["version"] == 4 + assert stored["notice"]["version"] == 5 assert "last_attempt_day" not in stored and "counters" not in stored From e9a5d348bf7fae13379f61a6c55439185b4123f1 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Wed, 30 Sep 2026 16:39:34 -0700 Subject: [PATCH 13/16] test(ci): exclude disposable probe subprocess from source coverage Signed-off-by: Lihua <1017343802@qq.com> --- tests/architecture/test_semantic_development_probe.py | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/tests/architecture/test_semantic_development_probe.py b/tests/architecture/test_semantic_development_probe.py index 94987444c4..1a80dcf430 100644 --- a/tests/architecture/test_semantic_development_probe.py +++ b/tests/architecture/test_semantic_development_probe.py @@ -2,6 +2,7 @@ from __future__ import annotations +import os from pathlib import Path import subprocess import sys @@ -344,6 +345,12 @@ def probe_cli(repository: Path) -> Path: def _run_probe_cli(repository: Path) -> subprocess.CompletedProcess[str]: + # The copied CLI runs in a disposable repository. Do not merge its coverage + # into the source checkout: pytest removes that repository before CI reports. + env = { + key: value for key, value in os.environ.items() + if not key.startswith(("COV_CORE_", "COVERAGE_")) + } return subprocess.run( [ sys.executable, @@ -352,6 +359,7 @@ def _run_probe_cli(repository: Path) -> subprocess.CompletedProcess[str]: "HEAD", ], cwd=repository, + env=env, capture_output=True, text=True, check=False, From 96657f9e1fb8a10630725bd9e33d1ccd5e14fcec Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Thu, 1 Oct 2026 18:39:35 +0800 Subject: [PATCH 14/16] perf(authority): avoid per-field tuples in strict JSON copies Signed-off-by: huangruiteng --- .../coordination/authority_store_codec.ts | 9 ++++++- .../authority_store_codec.test.ts | 24 ++++++++++++++++++- 2 files changed, 31 insertions(+), 2 deletions(-) diff --git a/loopx/control_plane/coordination/authority_store_codec.ts b/loopx/control_plane/coordination/authority_store_codec.ts index a99c6fbddb..8c5ac6fd1f 100644 --- a/loopx/control_plane/coordination/authority_store_codec.ts +++ b/loopx/control_plane/coordination/authority_store_codec.ts @@ -90,7 +90,14 @@ function cloneAuthorityJson(value: unknown, stack: Set, canonicalKeys: b stack.add(value); try { const keys = Object.keys(value); - if (canonicalKeys) keys.sort(authorityUnicodeCompare); + if (!canonicalKeys) { + // Fill without inherited setters (including __proto__), then normalize + // to the same ordinary object as Object.fromEntries. No per-field tuples. + const copy: JsonObject = Object.create(null); + for (const key of keys) copy[key] = cloneAuthorityJson(value[key], stack, false); + return Object.setPrototypeOf(copy, Object.prototype); + } + keys.sort(authorityUnicodeCompare); return Object.fromEntries( keys.map((key) => [ key, diff --git a/tests/control_plane_ts/authority_store_codec.test.ts b/tests/control_plane_ts/authority_store_codec.test.ts index af7bb7cb30..3e7ae4dfdb 100644 --- a/tests/control_plane_ts/authority_store_codec.test.ts +++ b/tests/control_plane_ts/authority_store_codec.test.ts @@ -1,7 +1,7 @@ import assert from "node:assert/strict"; import {createHash} from "node:crypto"; import test from "node:test"; -import {authorityUnicodeCompare, canonicalAuthorityBytes, canonicalAuthoritySha256} +import {authorityUnicodeCompare, canonicalAuthorityBytes, canonicalAuthoritySha256, copyAuthorityJson} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; // Frozen pre-optimization comparator: persisted revisions depend on code points, @@ -24,6 +24,28 @@ test("canonical key ordering preserves the persisted Unicode comparator", () => assert.deepEqual([...keys].sort(authorityUnicodeCompare), [...keys].sort(referenceCompare)); }); +test("copying JSON defines own data even when a key names an inherited setter", () => { + const key = "authorityCopySentinel"; + const previous = Object.getOwnPropertyDescriptor(Object.prototype, key); + let setterCalls = 0; + Object.defineProperty(Object.prototype, key, { + configurable: true, set: () => { setterCalls++; }, + }); + try { + const input = JSON.parse('{"authorityCopySentinel":{"nested":1},"constructor":{"marker":2}}'); + const copy = copyAuthorityJson(input) as Record; + assert.equal(setterCalls, 0); + assert.deepEqual(copy, input); + assert.equal(Object.getPrototypeOf(copy), Object.prototype); + assert.equal(Object.hasOwn(copy, key), true); + (copy[key] as {nested: number}).nested = 99; + assert.deepEqual(input[key], {nested: 1}); + } finally { + if (previous) Object.defineProperty(Object.prototype, key, previous); + else Reflect.deleteProperty(Object.prototype, key); + } +}); + test("canonical bytes and digest retain JSON enumeration, scalar and special-key semantics", () => { const input = JSON.parse('{"😀":4,"\\ue000":3,"__proto__":{"z":2,"a":1},"2":2,"10":10,"":0}'); input.scalars = [-0, 1e30, "\ud800", true, null]; From 69fd2de6d2b54577664813bfe9ee684c8f287451 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Thu, 1 Oct 2026 18:45:19 +0800 Subject: [PATCH 15/16] test(quota): use original execution hint for settled Turn ACK Signed-off-by: huangruiteng --- tests/control_plane/test_quota_settlement_cli.py | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/tests/control_plane/test_quota_settlement_cli.py b/tests/control_plane/test_quota_settlement_cli.py index c65c32fcda..1b53ad6232 100644 --- a/tests/control_plane/test_quota_settlement_cli.py +++ b/tests/control_plane/test_quota_settlement_cli.py @@ -2219,6 +2219,7 @@ def test_standard_codex_app_settlement_is_receipted_and_idempotent( ) assert guard_rc == 0, guard + original_ack_hint = guard["scheduler_hint"]["app_automation"]["ack_hint"] identity = guard["heartbeat_receipt"]["settlement_identity"] assert identity["todo_id"] == TODO_ID assert identity["effect_id"] == (f"{GOAL_ID}:{AGENT_ID}:{TODO_ID}:{TURN_ID}") @@ -2351,9 +2352,9 @@ def test_standard_codex_app_settlement_is_receipted_and_idempotent( ) assert _spend_run_count(runtime) == 1 - settled_ack_hint = settled_replay["scheduler_hint"]["codex_app"]["ack_hint"] - assert settled_ack_hint["args"]["turn_instance_id"] == TURN_ID - assert settled_ack_hint["cli_args"][-3:] == [ + # Use the original execution hint: the settled skip packet has no host work. + assert original_ack_hint["args"]["turn_instance_id"] == TURN_ID + assert original_ack_hint["cli_args"][-3:] == [ "--turn-instance-id", TURN_ID, "--execute", @@ -2361,7 +2362,7 @@ def test_standard_codex_app_settlement_is_receipted_and_idempotent( ack_rc, ack = _run_cli( registry_path, runtime, - *settled_ack_hint["cli_args"], + *original_ack_hint["cli_args"], ) # A settled Turn remains current until a newer heartbeat is admitted. assert ack_rc == 0, ack @@ -2394,7 +2395,7 @@ def test_standard_codex_app_settlement_is_receipted_and_idempotent( assert fresh_turn_rc == 0, fresh_turn assert fresh_turn["selected_todo"]["todo_id"] == successor_id stale_ack_rc, stale_ack = _run_cli( - registry_path, runtime, *settled_ack_hint["cli_args"], + registry_path, runtime, *original_ack_hint["cli_args"], ) assert stale_ack_rc == 1, stale_ack assert stale_ack["error_code"] == "SCHEDULER_FOLLOWUP_HEARTBEAT_RECEIPT_STALE" From 4a944ca19de494797152453a0c3cc42336f6311c Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Thu, 1 Oct 2026 18:45:19 +0800 Subject: [PATCH 16/16] docs(authority): distinguish merge, cohort trial and default qualification Signed-off-by: huangruiteng --- .../2026-09-28-retirement-cadence.md | 21 +++++++++++-- .../2026-09-28-retirement-cadence.zh-CN.md | 13 ++++++-- ...shared-goal-authority-state-provider-v0.md | 31 +++++++++++++++++-- 3 files changed, 57 insertions(+), 8 deletions(-) diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.md index a3425b9bef..093cebd517 100644 --- a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.md +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.md @@ -257,9 +257,24 @@ Scale characterization with 4,101 synthetic Agent Todos still hits the existing repair for File/SQLite. The repaired contract API can read that collection; this does not qualify the remaining whole-command payload boundary. -The next B work is SQLite admission on a frozen source/runtime profile: rerun -the existing reference capacity axes, reconcile concurrency/recovery/consumer-lag -evidence, and verify the applicability of retained natural-time soak results. +B distinguishes a bounded provider PR, a recoverable opt-in developer cohort, +and the release default, as specified in Section 7.2 of the owning RFC. Proposed +absolute latency budgets do not veto every merge or trial: use matched current +release measurements, disclose absolute and relative regressions, and inspect +their consumer impact. Correctness, original receipts, complete metadata and +recoverable migrations remain hard requirements. Frozen reports retain their +original budgets and failed/missing rows; revising a decision does not rewrite +past evidence. + +[PR #5251](https://github.com/loopx-project/loopx/pull/5251) refines strict JSON +materialization behind the existing codec owner. It preserves persisted +canonical encoding while avoiding repeated immutable primitive allocation in +historical projections. Its historical formal reports are source-specific: +the author reports 8 passed / 6 failed / 10 missing on `d767b06f1`, then +14 passed / 0 failed / 10 missing on `02d3dee83`; neither report qualifies a +later head, nor completes the missing axes. The next qualification reconciles +changed-path measurements and affected real providers/callers, then verifies +concurrency, recovery, consumer lag and retained natural-time soak applicability. The comparison runner's former conflict expectation contradicted merged #5169: an identical historical intent must return its original applied revision/cursor. The runner now checks that result, independently rejects projection/event/receipt diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.zh-CN.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.zh-CN.md index e5fd3dd600..86fc81f2ae 100644 --- a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.zh-CN.md +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.zh-CN.md @@ -200,8 +200,17 @@ Todo 写入时的业务校验,也不重审完成/deferred 历史的授权。 仍触及既有 `todo.succession.project` RPC 响应预算;修复后的合同 API 能读取该集合, 不代表剩余整命令包体边界已完成验收。 -B 下一步聚焦冻结 source/runtime profile 下的 SQLite 准入:重新跑已有 reference -容量轴,对齐并发/恢复/consumer lag 证据,并核对保留的自然时间 soak 适用性。 +B 按主 RFC 7.2 区分有界 provider PR、可恢复的小范围开发者试用和发布默认值。 +提议的绝对延迟预算不否决每次合入或试用:在匹配负载下比较当前受支持版本, +公开绝对增量与相对变化,并核对消费者影响。正确性、原始回执、完整 metadata 和 +可恢复迁移仍是硬要求。冻结报告保留原预算及失败/缺失项,调整决策不改写旧证据。 + +[PR #5251](https://github.com/loopx-project/loopx/pull/5251) 在已有 strict JSON codec +owner 中优化历史数据物化,保留持久化 canonical 编码,减少历史投影中不可变值的 +重复分配。旧正式报告各自绑定 source:作者报告 `d767b06f1` 为 8 通过/6 失败/10 +缺失,`02d3dee83` 为 14 通过/0 失败/10 缺失;这些结果不验证后续 head,也不补齐 +缺失轴。下一步核对变更路径的配对测量与真实 provider/调用方,再对齐并发、恢复、 +consumer lag 及保留的自然时间 soak 适用性。 比较 runner 原先要求历史重试返回 conflict,与已合并 #5169 矛盾:相同完整意图应 返回原 applied revision/cursor。现在核对原结果,分别拒绝 projection/event/receipt 漂移,在重试前后分页验证全部历史,不保留所有预期快照。不变量失败就不发布成功 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 ea1826b3fc..aa0c3bebe2 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -1284,9 +1284,34 @@ For the 64 KiB live-state axis at 10,000 versus 100,000 commits: 1 MiB axis and 300,000-commit headroom separately; failures narrow the supported profile rather than disappearing into averaged results. -These thresholds are proposed engineering budgets, not current measurements. -Review them against the first matched baseline before activation; do not relax -correctness, silently change the workload, or advertise an unqualified horizon. +These thresholds are proposed engineering budgets, not current measurements or +a universal prerequisite for merging an improvement or trying an opt-in provider. +Use the current supported implementation as the performance control, on the same +host, runtime, complete data, history, durability and command mix. Retain the +original pre-migration baseline as a product comparison; it cannot hide a +regression against the current supported release. Report both absolute latency +and relative change. A percentage increase on a short store operation is not, +by itself, a user-visible regression; trace it through the affected command or +consumer before deciding whether the tradeoff is acceptable. Conversely, parity +with an already unusable baseline is not sufficient. + +Make three distinct decisions using the existing qualification evidence: + +| Decision | Required evidence and scope | +| --- | --- | +| Merge a bounded provider improvement | Exact-source correctness, metadata/history/receipt preservation, affected real backends and callers, and matched measurements of the changed paths. Disclose local regressions and their consumer impact; a missed proposed latency target alone is not a merge blocker. | +| Invite a small opt-in developer cohort | A recoverable installed workflow: verified backup, migration, restart, ordinary commands, new writes and return to the previous provider without losing those writes. Relevant concurrency and interruption controls must pass. Bound the advertised workload to demonstrated evidence, expose actionable failures, and observe natural operation. A ten-day certificate is not required to begin this reversible trial. | +| Select the release default | Installation and runtime support, new-Goal creation, upgrade, migration, recovery and rollback must work on the supported profiles. Representative sustained operation must show no material degradation of common user journeys against the current release, with acceptable growth, resource use and recovery space. Reconcile existing soak results with changed boundaries; do not restart elapsed-time evidence for unrelated changes. | + +Data loss, duplicate effects, altered original receipts, incorrect decisions, +broken fencing or unrecoverable migration remain blockers at the affected +boundary. Performance targets may be calibrated from matched evidence with an +explicit tradeoff; correctness is not calibrated away. Preserve frozen workload +and budget identities in existing reports, including failures and missing rows. +If a budget is revised, record a new declared comparison rather than relabeling +an old failed report as passing. The formal ten-day/100,000-commit qualification +and its unmeasured axes remain explicit; neither a merged PR nor a successful +cohort trial advertises that horizon or settles the release-default decision. #### Retention, recovery and delivery gates