Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Original file line number Diff line number Diff line change
Expand Up @@ -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,不发布原始证据。
18 changes: 11 additions & 7 deletions loopx/control_plane/coordination/authority_store_codec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
15 changes: 12 additions & 3 deletions loopx/control_plane/coordination/file_authority_journal.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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<FileAuthorityJournal> {
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");
Expand All @@ -87,6 +88,10 @@ export class FileAuthorityJournal {
parseAuthorityCursor(cursor) !== BigInt(value.committed.length)) return invalid("lineage is invalid");
const rows: StoredCommit[] = [], operations = new Set<string>();
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"])) {
Expand All @@ -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");
Expand All @@ -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<FileAuthorityJournal> {
let previous: JsonObject | null = null;
const compact = rows.map(transaction => {
const encoded = retain(transaction, previous);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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};
}
Expand All @@ -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) {
Expand Down
21 changes: 17 additions & 4 deletions loopx/control_plane/coordination/file_authority_store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, Promise<FileAuthorityJournal>>();

function documentDigest(raw: Uint8Array): string {
return createHash("sha256").update(raw).digest("hex");
Expand Down Expand Up @@ -153,7 +156,7 @@ function decodeDocument(
value: unknown,
goalId: string,
storeIdentity: string,
): FileAuthorityJournal {
): Promise<FileAuthorityJournal> {
return FileAuthorityJournal.decode(value, goalId, storeIdentity, (previous, transaction) =>
fileAuthorityRevision(goalId, storeIdentity, previous, transaction));
}
Expand Down Expand Up @@ -199,7 +202,7 @@ export class FileAuthorityStore implements AuthorityStore {
protected async archiveRenamed(): Promise<void> {}

/** 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<FileAuthorityJournal> {
return decodeDocument(value, this.goalId, identity);
}

Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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,
Expand Down
1 change: 1 addition & 0 deletions skills/loopx-self-repair/references/repair-patterns.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Expand Down
10 changes: 10 additions & 0 deletions skills/loopx-self-repair/references/targeted-diagnostics.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
29 changes: 29 additions & 0 deletions tests/control_plane_ts/authority_store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
});
Loading
Loading