diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md index 3944614ac..1b8f81653 100644 --- a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md @@ -96,3 +96,36 @@ PostgreSQL 16 server passed, with no skipped checks in these suites. A detached 129-commit real-source snapshot passed File and SQLite restore and independent audit. Recovery of the earlier timed-out SQLite destination passed without reissuing its committed operations. These checks do not claim active cutover. + +## Shared-runtime latency reconciliation (`96a3b90f4`) + +The recovery delivery above is now merged as #5140. #5144 is the open managed +Host process supervision slice; attached hosts still need their declared +cancellation boundary. Whole-Goal activation/rollback and default entrypoint +cutover remain the two planned subsequent implementation PRs. Existing #5054 +(retirement) and #4931 (SQLite proof encoding) remain open and are not new work. +Thus the inventory is two planned implementation PRs plus those three existing +PRs, **before this newly reproduced latency repair**. This is an inventory, not +an unconditional completion count: #4224 D2 capacity/soak and D1/D3 evidence +remain gates, and failures may require additional scoped fixes. + +The latency repair does not retire another Python owner or close D2. Isolated +fixed File snapshots reproduce 9.9–10.5 second cold history verification, +including a ping timeout at the original 10 second budget. Warm reads hide the +problem; alternating Goals evict the single verified read cache. An isolated +CPU profile attributes about 56% of samples to allocating code-point arrays in +the shared key comparator. An allocation-free comparator preserves ordering; +File verification yields between complete transactions and coalesces identical +in-flight proofs keyed by path, store identity and exact byte digest. A late +corrupt row still rejects an early receipt; failed proofs never become cache +entries. No schema, revision formula, timeout or selector changes. + +On the same snapshots, cold reads take about 2.9 seconds and concurrent light +requests 18–52 ms. These are local observations, not formal capacity/p95 or +cross-platform qualification. One very large transaction, JSON parse, other +synchronous handlers and SQLite replay can still occupy the event loop; this +is not general worker isolation. The common comparator benefits every provider; +the yielding and in-flight proof lifecycle belong to File. #4931's digest +window remains a separate optimization. The regression uses a private real +server and the existing mixed Todo/lease/decision fixture; production locators, +active Goals and raw evidence are never modified or published. diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md index 8964d3291..6c2fbf770 100644 --- a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md @@ -74,3 +74,26 @@ Goal。公开材料排除私有 Goal 内容及原始日志。 PostgreSQL 16 的 4 项跨 provider 检查全部通过,这些套件没有跳过项。129 笔历史的 独立真实来源快照通过 File、SQLite 恢复及独立审计;此前超时的 SQLite 目标也成功 恢复,没有重发已提交操作。这些结果不表示已经完成活跃 Goal 切换。 + +## 共享运行时延迟核对(`96a3b90f4`) + +上文恢复交付已作为 #5140 合入。#5144 是待合并的受管 Host 进程监督切片;attached +Host 仍需明确取消能力边界。整 Goal 激活/回退、默认入口切换仍是后续两个规划实现 +PR。已有 #5054(退役)、#4931(SQLite 证明编码)仍开放,不重复实现。因此清单是 +**两个规划实现 PR,加三个已有 PR,再加本次新复现的延迟修复**。这是工作清单, +不是无条件完成倒计时:#4224 D2 容量/soak、D1/D3 证据仍需验收;失败可产生新的 +有界修复,必须指出具体缺陷,不能重新复述一个固定区间。 + +本次不退役额外 Python owner,也不关闭 D2。隔离固定 File 快照复现了 9.9–10.5 秒 +的冷历史校验,期间 ping 在原 10 秒预算内超时。热缓存掩盖问题,交替读取 Goal 又 +会淘汰唯一的验证缓存。隔离 CPU 采样中,约 56% 样本落在为键比较分配码点数组。 +改用不分配数组的同序比较;File 在完整事务的校验之间让出事件循环,并按路径、 +store identity、精确字节摘要合并进行中的相同校验。尾部损坏仍拒绝返回早期回执, +失败证明不会成为缓存。不修改格式、revision 算法、超时或 selector。 + +相同快照冷读约 2.9 秒,并发轻请求约 18–52 毫秒。这是本机观察,不是正式容量、 +p95 或跨平台资格。单笔巨大事务、JSON 解析、其他同步 handler 和 SQLite 重放仍 +可能占用事件循环;本次没有实现通用 worker 隔离。公共比较器惠及各 provider, +让步及进行中证明的生命周期归 File;#4931 的 digest window 仍是独立优化。 +回归使用私有真实 server 和既有混合 Todo/lease/decision fixture,不改生产 locator +及活跃 Goal,不发布原始证据。 diff --git a/loopx/control_plane/coordination/authority_store_codec.ts b/loopx/control_plane/coordination/authority_store_codec.ts index a386352b7..06ec9e553 100644 --- a/loopx/control_plane/coordination/authority_store_codec.ts +++ b/loopx/control_plane/coordination/authority_store_codec.ts @@ -11,14 +11,18 @@ export function isAuthorityJsonObject(value: unknown): value is JsonObject { } export function authorityUnicodeCompare(left: string, right: string): number { - const leftPoints = Array.from(left, (item) => item.codePointAt(0) ?? 0); - const rightPoints = Array.from(right, (item) => item.codePointAt(0) ?? 0); - const shared = Math.min(leftPoints.length, rightPoints.length); - for (let index = 0; index < shared; index += 1) { - const difference = leftPoints[index] - rightPoints[index]; - if (difference !== 0) return difference; + // Walk code points without allocating two arrays for every sort comparison. + // JS's default sort compares UTF-16 units, which would change persisted + // revisions for supplementary characters relative to BMP characters. + let leftIndex = 0, rightIndex = 0; + while (leftIndex < left.length && rightIndex < right.length) { + const leftPoint = left.codePointAt(leftIndex)!; + const rightPoint = right.codePointAt(rightIndex)!; + if (leftPoint !== rightPoint) return leftPoint - rightPoint; + leftIndex += leftPoint > 0xffff ? 2 : 1; + rightIndex += rightPoint > 0xffff ? 2 : 1; } - return leftPoints.length - rightPoints.length; + return leftIndex < left.length ? 1 : rightIndex < right.length ? -1 : 0; } export function hasExactAuthorityKeys( diff --git a/loopx/control_plane/coordination/file_authority_journal.ts b/loopx/control_plane/coordination/file_authority_journal.ts index 36caca25f..fa94135f9 100644 --- a/loopx/control_plane/coordination/file_authority_journal.ts +++ b/loopx/control_plane/coordination/file_authority_journal.ts @@ -1,5 +1,6 @@ /** File's physical journal codec. Logical revisions, receipts and transactions * stay unchanged; only repeated projections become checkpoints and deltas. */ +import {setImmediate as yieldToRuntime} from "node:timers/promises"; import type {JsonObject} from "../effect_program.ts"; import type {AuthorityStoreCommit, AuthorityStoreCommittedTransaction} from "./authority_store.ts"; import {AuthorityStoreProtocolError, canonicalAuthorityBytes, canonicalAuthorityObject, @@ -73,7 +74,7 @@ export class FileAuthorityJournal { this.operations = new Map(rows.map(row => [row.operation_id, row])); } - static decode(value: unknown, goal: string, identity: string, revisionFor: JournalRevision): FileAuthorityJournal { + static async decode(value: unknown, goal: string, identity: string, revisionFor: JournalRevision): Promise { if (!isAuthorityJsonObject(value) || !hasExactAuthorityKeys(value, HEADER_KEYS) || value.schema_version !== FILE_AUTHORITY_JOURNAL_SCHEMA) { return invalid("schema mismatch; run loopx authority-archive upgrade --execute before opening this store"); @@ -87,6 +88,10 @@ export class FileAuthorityJournal { parseAuthorityCursor(cursor) !== BigInt(value.committed.length)) return invalid("lineage is invalid"); const rows: StoredCommit[] = [], operations = new Set(); let previous: JsonObject | null = null, previousRevision: string | null = null; + // Historical verification is CPU work inside the shared Effect server. + // Yield between complete transactions, never publish a partially verified + // journal. Promise.resolve() would only drain microtasks and starve sockets. + let sliceStart = performance.now(); for (const [index, raw] of value.committed.entries()) { if (!isAuthorityJsonObject(raw) || !hasExactAuthorityKeys(raw, ["cursor", "provider_revision", "operation_id", "events", "receipts", "state"])) { @@ -104,6 +109,10 @@ export class FileAuthorityJournal { const {projection, ...entry} = transaction; rows.push({...entry, state}); previous = projection; previousRevision = transaction.provider_revision; + if (performance.now() - sliceStart >= 8) { + await yieldToRuntime(); + sliceStart = performance.now(); + } } if (rows.at(-1)!.cursor !== cursor || previousRevision !== revision || !canonicalAuthorityBytes(previous).equals(canonicalAuthorityBytes(head))) return invalid("head lineage is invalid"); @@ -123,8 +132,8 @@ export class FileAuthorityJournal { /** Migration boundary: callers supply fully verified logical transactions. * Decode the resulting wire format again before it can be adopted. */ - static fromTransactions(goal: string, identity: string, rows: readonly AuthorityStoreCommittedTransaction[], - revisionFor: JournalRevision): FileAuthorityJournal { + static async fromTransactions(goal: string, identity: string, rows: readonly AuthorityStoreCommittedTransaction[], + revisionFor: JournalRevision): Promise { let previous: JsonObject | null = null; const compact = rows.map(transaction => { const encoded = retain(transaction, previous); diff --git a/loopx/control_plane/coordination/file_authority_migration.ts b/loopx/control_plane/coordination/file_authority_migration.ts index 6444e6875..3c2fc5da6 100644 --- a/loopx/control_plane/coordination/file_authority_migration.ts +++ b/loopx/control_plane/coordination/file_authority_migration.ts @@ -39,7 +39,7 @@ export async function migrateFileAuthorityStore(directory: string, goal: string, const value: unknown = JSON.parse(source.toString("utf8")); if (!isAuthorityJsonObject(value)) throw new Error("Invalid authority document"); if (value.schema_version === FILE_AUTHORITY_JOURNAL_SCHEMA) { - const current = FileAuthorityJournal.decode(value, goal, identity, revisionFor); + const current = await FileAuthorityJournal.decode(value, goal, identity, revisionFor); return {status: "already_current", provider: "file", cursor: current.cursor, provider_revision: current.provider_revision}; } @@ -48,7 +48,7 @@ export async function migrateFileAuthorityStore(directory: string, goal: string, throw new Error("Unsupported file authority format or mismatched lineage; source was not changed"); } const legacy = decodeRetainedAuthorityJournal(value, "file migration source", revisionFor); - const compact = FileAuthorityJournal.fromTransactions(goal, identity, legacy.committed, revisionFor); + const compact = await FileAuthorityJournal.fromTransactions(goal, identity, legacy.committed, revisionFor); // Compare complete logical history, not only the head or receipt count. const logicalDigest = canonicalAuthoritySha256(legacy.committed); if (canonicalAuthoritySha256(compact.scan(0, legacy.committed.length)) !== logicalDigest) { diff --git a/loopx/control_plane/coordination/file_authority_store.ts b/loopx/control_plane/coordination/file_authority_store.ts index 6eff136fd..1793858bb 100644 --- a/loopx/control_plane/coordination/file_authority_store.ts +++ b/loopx/control_plane/coordination/file_authority_store.ts @@ -50,6 +50,9 @@ interface VerifiedDocument { document?: FileAuthorityJournal; } let verifiedDocument: VerifiedDocument | null = null; +// Only identical immutable input bytes share in-flight verification. Failed +// proofs are removed too; neither a path nor a pending promise grants authority. +const pendingVerification = new Map>(); function documentDigest(raw: Uint8Array): string { return createHash("sha256").update(raw).digest("hex"); @@ -153,7 +156,7 @@ function decodeDocument( value: unknown, goalId: string, storeIdentity: string, -): FileAuthorityJournal { +): Promise { return FileAuthorityJournal.decode(value, goalId, storeIdentity, (previous, transaction) => fileAuthorityRevision(goalId, storeIdentity, previous, transaction)); } @@ -199,7 +202,7 @@ export class FileAuthorityStore implements AuthorityStore { protected async archiveRenamed(): Promise {} /** Full-history verification seam; unchanged byte-identical reads may reuse it. */ - protected decodeStoredDocument(value: unknown, identity: string): FileAuthorityJournal { + protected decodeStoredDocument(value: unknown, identity: string): Promise { return decodeDocument(value, this.goalId, identity); } @@ -272,7 +275,17 @@ export class FileAuthorityStore implements AuthorityStore { (!requireHistory || verifiedDocument.document !== undefined)) { return verifiedDocument; } - const document = this.decodeStoredDocument(JSON.parse(raw.toString("utf8")), identity); + const key = JSON.stringify([this.path, identity, digest]); + let proof = pendingVerification.get(key); + if (!proof) { + proof = this.decodeStoredDocument(JSON.parse(raw.toString("utf8")), identity); + pendingVerification.set(key, proof); + } + let document: FileAuthorityJournal; + try { document = await proof; } + finally { + if (pendingVerification.get(key) === proof) pendingVerification.delete(key); + } return rememberVerifiedDocument(this.path, identity, raw, digest, document, this.fullDocumentCacheLimitBytes()); } catch (error) { @@ -473,7 +486,7 @@ export class FileAuthorityStore implements AuthorityStore { const identity = await this.readStoreIdentity(false); let archived: FileAuthorityJournal | null = null; try { - archived = decodeDocument( + archived = await decodeDocument( JSON.parse(await readFile(archivePath, "utf8")), this.goalId, identity, diff --git a/skills/loopx-self-repair/references/repair-patterns.md b/skills/loopx-self-repair/references/repair-patterns.md index 537149f14..7a7f5a91e 100644 --- a/skills/loopx-self-repair/references/repair-patterns.md +++ b/skills/loopx-self-repair/references/repair-patterns.md @@ -5,6 +5,7 @@ teaches a reusable control-plane lesson. | Pattern | Symptoms | Evidence To Read | Likely Root | Durable Repair | | --- | --- | --- | --- | --- | +| `authority_cold_read_starves_runtime` | Unrelated pure Todo rules and ping time out while warm reads are fast. | Same-byte isolated cold/warm/alternating-Goal probes, original response budget, CPU profile and exact provider revision. | Synchronous retained-history verification monopolizes the shared event loop; allocation-heavy canonical key sorting amplifies it. | Preserve byte-level proof semantics, reduce codec allocations, yield between complete transaction proofs and share only identical in-flight proofs. Test real socket concurrency and late-history corruption; do not raise timeouts, weaken integrity or replay ambiguous writes. | | `native_todo_status_index_schema_gap` | A Goal shows “status load failed / invalid response” after native Todos appear, while the scoped status endpoint returns valid JSON. | Exact scoped status response, Zod issue paths, Todo `todo_id` and `index` fields, native presentation contract, packaged Goal load. | The dashboard still requires every Todo to have a numeric source index, but native Todos deliberately use stable `todo_id` with absent or null index. One row rejects the entire Goal snapshot. | Accept nullable/absent index only when a nonempty stable Todo ID exists; keep numeric legacy indexes and reject anonymous or malformed rows. Render and act by Todo ID, then prove a mixed native/legacy scoped snapshot loads in the packaged UI. | | `capability_catalog_editor_kind_drift` | Machine or Goal settings report an empty capability list even though the configuration API returns registered capabilities. | Live API catalog IDs and editor kinds, the dashboard's accepted field-kind schema, and the page's load-error state. | One new descriptor emits an unsupported field kind; strict validation rejects the shared catalog and the machine page presents the failed load as an empty registry. | Keep the published editor vocabulary aligned with the browser contract, check every built-in descriptor together, and show a retryable error when catalog loading or validation fails. Only a successfully loaded empty catalog may show the empty state. | | `acceptance_scope_capture` | A bounded validation experiment leaves unrelated existing/new work unbound; a recorded blocker quiets replan without repairing admission. | Canonical contract scope/bindings, exact held generation, authorized configuration source, ordinary task validators and post-correction claim/lease readback. | Omitted scope silently imposed Goal-wide acceptance; repeated per-task binding masked the missing scope contract. | Require explicit scope on new owner configuration, preserve legacy persisted semantics/replay, and enforce one typed scope across admission, completion and verification freshness. Expose scope on existing read surfaces. Diagnose scope before proposing rebinding; a blocker ACK is neither a repair nor a handoff. Apply authorized corrections through CAS and validate independent work resumes while selected holds and ordinary validation remain. | diff --git a/skills/loopx-self-repair/references/targeted-diagnostics.md b/skills/loopx-self-repair/references/targeted-diagnostics.md index 382342634..8ded4ec29 100644 --- a/skills/loopx-self-repair/references/targeted-diagnostics.md +++ b/skills/loopx-self-repair/references/targeted-diagnostics.md @@ -57,6 +57,16 @@ and synthetic fixture or authorized read-only snapshot. Preserve integrity, receipt recovery and lease/CAS semantics; do not benchmark by mutating an active Goal. Check existing PRs before starting an overlapping store refactor. +When unrelated lightweight rules and `runtime.ping` slow down together, test +shared event-loop starvation before attributing the timeout to the named rule. +Compare cold, warm and alternating-Goal reads in a separate runtime using fixed +snapshot bytes; capture a CPU profile there, not by restarting a shared live +service. A hot cache can hide full-history CPU work. Keep the original response +budget, preserve uncertain-write recovery, and distinguish lower CPU cost from +cooperative scheduling. Yielding between verified transactions must not publish +an incomplete proof; concurrent identical reads may share only an exact-input +in-flight proof, with failures removed so a later read can revalidate. + Searchable reference lookup has no Goal authority and can stay in the skill. Runtime admission, recovery decisions and provider integrity remain in their existing typed owners. This is the S10 diagnostic-efficiency boundary alongside diff --git a/tests/control_plane_ts/authority_store.test.ts b/tests/control_plane_ts/authority_store.test.ts index b6613fa6b..eaf4c7b1c 100644 --- a/tests/control_plane_ts/authority_store.test.ts +++ b/tests/control_plane_ts/authority_store.test.ts @@ -249,3 +249,32 @@ test("ambiguous file commits reconcile only from durable receipt readback", asyn assert.equal(loaded.status, "loaded"); if (loaded.status === "loaded") assert.equal(loaded.head.authority_revision, 1); }); + +test("concurrent cold reads share only the same exact-byte proof and recover after failure", async t => { + const {root, store} = await fixture(t); + assert.equal((await store.commitAuthority(commit(null, "operation-shared", 1, 1))).status, "applied"); + const valid = await readFile(store.path, "utf8"); + await writeFile(store.path, valid + "\n"); + class CountingStore extends FileAuthorityStore { + static validations = 0; + protected override async decodeStoredDocument(value: unknown, identity: string) { + CountingStore.validations++; + // Hold the asynchronous proof open while sibling handles enter the read. + await new Promise(resolve => setTimeout(resolve, 25)); + return super.decodeStoredDocument(value, identity); + } + } + const readers = Array.from({length: 6}, () => new CountingStore(root, "goal-a")); + const first = await Promise.all(readers.map(s => s.readReceipt("operation-shared"))); + assert.ok(first.every(r => r.status === "found")); + assert.equal(CountingStore.validations, 1); + const corrupt = JSON.parse(valid); corrupt.committed[0].provider_revision = "corrupt"; + await writeFile(store.path, JSON.stringify(corrupt)); + assert.ok((await Promise.all(readers.map(s => s.loadAuthority()))).every(r => r.status === "failed")); + assert.equal(CountingStore.validations, 2); + assert.equal((await readers[0]!.loadAuthority()).status, "failed"); + assert.equal(CountingStore.validations, 3, "a rejected promise must not remain in the in-flight registry"); + await writeFile(store.path, valid + "\n\n"); + assert.equal((await readers[0]!.loadAuthority()).status, "loaded"); + assert.equal(CountingStore.validations, 4); +}); diff --git a/tests/control_plane_ts/authority_store_codec.test.ts b/tests/control_plane_ts/authority_store_codec.test.ts new file mode 100644 index 000000000..af7bb7cb3 --- /dev/null +++ b/tests/control_plane_ts/authority_store_codec.test.ts @@ -0,0 +1,39 @@ +import assert from "node:assert/strict"; +import {createHash} from "node:crypto"; +import test from "node:test"; +import {authorityUnicodeCompare, canonicalAuthorityBytes, canonicalAuthoritySha256} + from "../../loopx/control_plane/coordination/authority_store_codec.ts"; + +// Frozen pre-optimization comparator: persisted revisions depend on code points, +// including unpaired surrogates, rather than the default UTF-16 sort order. +function referenceCompare(a: string, b: string): number { + const x = Array.from(a, c => c.codePointAt(0)!); + const y = Array.from(b, c => c.codePointAt(0)!); + for (let i = 0; i < Math.min(x.length, y.length); i++) { + if (x[i] !== y[i]) return x[i]! - y[i]!; + } + return x.length - y.length; +} + +test("canonical key ordering preserves the persisted Unicode comparator", () => { + const alphabet = ["", "a", "\0", "é", "中", "\ud800", "\udfff", "\ue000", "😀", "𐀀", "\u{10ffff}"]; + const keys = [...alphabet, ...alphabet.flatMap(a => alphabet.map(b => a + b))]; + for (const a of keys) for (const b of keys) { + assert.equal(Math.sign(authorityUnicodeCompare(a, b)), Math.sign(referenceCompare(a, b))); + } + assert.deepEqual([...keys].sort(authorityUnicodeCompare), [...keys].sort(referenceCompare)); +}); + +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]; + input.sparse = Array(2); + const expected = '{"2":2,"10":10,"":0,"__proto__":{"a":1,"z":2},"scalars":[0,1e+30,"\\ud800",true,null],"sparse":[null,null],"\ue000":3,"😀":4}'; + assert.equal(canonicalAuthorityBytes(input).toString(), expected); + assert.equal(canonicalAuthoritySha256(input), createHash("sha256").update(expected).digest("hex")); + assert.equal(Object.getPrototypeOf(input), Object.prototype); + const cyclic: Record = {}; cyclic.self = cyclic; + for (const invalid of [cyclic, {n: NaN}, {n: Infinity}, {x: undefined}, {x: new Date()}, [undefined]]) { + assert.throws(() => canonicalAuthorityBytes(invalid)); + } +}); diff --git a/tests/control_plane_ts/effect_runtime_fairness.test.ts b/tests/control_plane_ts/effect_runtime_fairness.test.ts new file mode 100644 index 000000000..8a4fc69f6 --- /dev/null +++ b/tests/control_plane_ts/effect_runtime_fairness.test.ts @@ -0,0 +1,117 @@ +import assert from "node:assert/strict"; +import {spawn} from "node:child_process"; +import {once} from "node:events"; +import {mkdtemp, readFile, rm, writeFile} from "node:fs/promises"; +import {connect} from "node:net"; +import {tmpdir} from "node:os"; +import {join} from "node:path"; +import {setTimeout as delay} from "node:timers/promises"; +import test from "node:test"; +import {FileAuthorityJournal} from "../../loopx/control_plane/coordination/file_authority_journal.ts"; +import {FileAuthorityStore, fileAuthorityRevision} from "../../loopx/control_plane/coordination/file_authority_store.ts"; +import {canonicalAuthorityBytes} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; +import {productionScaleCoordinationFixture} from "./production_scale_coordination_fixture.ts"; + +// A private child process and explicitly rooted store avoid touching the user's +// shared runtime. No tempfile cache, production locator or installed CLI is used. +test("a real runtime serves lightweight work during cold File history verification", {timeout: 60_000}, async t => { + const root = await mkdtemp(join(tmpdir(), "loopx-runtime-fairness-")); + t.after(() => rm(root, {recursive: true, force: true})); + const goal = "fairness-fixture"; + const store = new FileAuthorityStore(join(root, "authority", "file-v0"), goal); + const identity = await store.storeIdentity(); + assert.equal(identity.status, "available"); + if (identity.status !== "available") throw new Error("fixture identity unavailable"); + const revision = (previous: string | null, transaction: Parameters[3]) => + fileAuthorityRevision(goal, identity.store_identity, previous, transaction); + const projection = productionScaleCoordinationFixture(goal).projection; + let journal: FileAuthorityJournal | null = null; + // Retain the actual mixed Todo/lease/decision shapes. Repeated observations + // have distinct receipts and revisions even when they publish unchanged state. + for (let i = 0; i < 96; i++) { + journal = FileAuthorityJournal.append(journal, goal, identity.store_identity, { + expected_provider_revision: journal?.provider_revision ?? null, + operation_id: `observation-${i}`, events: [{kind: "observed", sequence: i}], + next_projection: projection, receipts: [{sequence: i}], + }, revision); + } + await writeFile(store.path, canonicalAuthorityBytes(journal!.toDocument())); + const infoPath = join(root, "runtime.json"); + const child = spawn(process.execPath, ["--no-warnings", "--experimental-sqlite", "--experimental-strip-types", + "loopx/control_plane/effect_runtime_server.ts", "--info", infoPath, "--fingerprint", "fairness-fixture"], { + env: {...process.env, LOOPX_EFFECT_RUNTIME_TOKEN: "isolated-test-token"}, stdio: "ignore", + }); + const closed = once(child, "close"); + t.after(async () => { if (child.exitCode === null) child.kill(); await closed; }); + let info: {port: number; token: string} | undefined; + for (let i = 0; i < 200; i++) { + try { info = JSON.parse(await readFile(infoPath, "utf8")); break; } + catch { if (child.exitCode !== null) throw new Error("isolated runtime exited"); await delay(20); } + } + assert.ok(info, "isolated runtime published its locator"); + const port = info.port, token = info.token; + let sequence = 0; + function rpc(method: string, params: Record = {}): Promise { + return new Promise((resolve, reject) => { + const socket = connect(port, "127.0.0.1"); + socket.setEncoding("utf8"); + socket.setTimeout(10_000, () => socket.destroy(new Error("original 10s response budget exceeded"))); + socket.on("error", reject); + socket.on("connect", () => socket.write(JSON.stringify({schema_version: "loopx_effect_runtime_request_v0", + token, request_id: `request-${sequence++}`, method, params}) + "\n")); + let raw = ""; + socket.on("data", chunk => { raw += chunk; }); + socket.on("end", () => { + try { const result = JSON.parse(raw); assert.equal(result.ok, true); resolve(result.result); } + catch (error) { reject(error); } + }); + }); + } + let heavyFinished = false; + const heavy = rpc("coordination.local_authority.operation_receipt", { + schema_version: "loopx_local_coordination_operation_receipt_request_v0", runtime_root: root, + goal_id: goal, operation_id: "observation-0", + }).finally(() => { heavyFinished = true; }); + void heavy.catch(() => {}); + // Give the dedicated server time to enter its cold read, then exercise a + // separate socket while its historical proof is still in flight. + await delay(100); + const binding = {id: "review", agent_id: "reviewer", todo_id: "todo_review", workspace: root, + requesters: ["coordinator"], host_args: ["--host", "dsh"], timeout_seconds: 60, output_refs: ["output.json"]}; + const [ping, decisions, scope, ownership, selectedBinding] = await Promise.all([rpc("runtime.ping"), rpc("todo.standing_decision.project", { + schema_version: "standing_decision_projection_request_v0", items: [], legacy_source_order: false, + }), rpc("todo.authoring_scope.plan", { + schema_version: "todo_authoring_scope_request_v0", command: "class", role: "agent", + todo: {task_class: "continuous_monitor"}, intent: {}, + }), rpc("todo.ownership_gate.decide", { + handoff_mode: "hard_lease", ownership_mutation: true, authority_mode: "registered_peer_actor", + }), rpc("collaboration.delegation.binding", { + agent_id: "coordinator", binding_id: "review", + config: {schema_version: "loopx_local_delegation_v0", bindings: [binding]}, + })]); + assert.deepEqual(selectedBinding, binding, "grant selection progresses without launching a Host"); + assert.equal(scope.schema_version, "todo_authoring_scope_result_v0"); + assert.equal(ownership.ownership_gate, "require_holder"); + assert.ok(ping.pid); + assert.equal(decisions, null, "no standing decisions means no authority projection"); + assert.equal(heavyFinished, false, "a full history proof must not monopolize other requests"); + const receipt = await heavy; + assert.equal(receipt.status, "found"); + assert.deepEqual(receipt.receipts, [{sequence: 0}]); + assert.equal(receipt.cursor, "1"); + // A valid early receipt is not publishable if a later historical row fails. + const corrupt = journal!.toDocument(); + const rows = corrupt.committed as Record[]; + rows.at(-1)!.provider_revision = "tampered-last-row"; + const damagedBytes = canonicalAuthorityBytes(corrupt); + await writeFile(store.path, damagedBytes); + const refused = await rpc("coordination.local_authority.operation_receipt", { + schema_version: "loopx_local_coordination_operation_receipt_request_v0", runtime_root: root, + goal_id: goal, operation_id: "observation-0", + }); + assert.equal(refused.status, "failed"); + assert.equal(refused.reason_code, "provider_protocol_violation"); + assert.deepEqual(await readFile(store.path), damagedBytes, "read failure must not rewrite history"); + await rpc("runtime.shutdown"); + await closed; +});