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 1/9] 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 2/9] 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 3/9] 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 4/9] 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 5/9] 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 6/9] 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 7/9] 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 8/9] 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 9/9] 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)})){