diff --git a/packages/paperclip-runner/src/drivers/acpx/runtime-host.test.ts b/packages/paperclip-runner/src/drivers/acpx/runtime-host.test.ts index 78d1e11e53..5fb1a709e1 100644 --- a/packages/paperclip-runner/src/drivers/acpx/runtime-host.test.ts +++ b/packages/paperclip-runner/src/drivers/acpx/runtime-host.test.ts @@ -25,6 +25,131 @@ const admissionControllers: AbortController[] = []; const pendingAdmissionOpenings = new Set>(); const pendingAdmissionCleanups = new Set>(); +// A credential or sandbox poll in this file can run behind a real retry +// envelope, not a mocked one. `stageManagedCodexCredential` first joins any +// in-flight quarantine recovery (codex-credentials.ts:183, :610-625). That +// recovery makes up to `MAX_AUTONOMOUS_CREDENTIAL_CLEANUP_ATTEMPTS` (8) +// attempts (codex-credentials.ts:19). The backoff between attempts is 10, +// 20, 40, 80, 160, 320, and 640 ms — 1,270 ms in total +// (codex-credentials.ts:558, :578-582). Each attempt can also run one real +// directory fsync. `DIRECTORY_SYNC_OPERATION_TIMEOUT_MS` (1,000 ms, +// codex-credentials.ts:18, :569, :1304) bounds that fsync. So one full +// recovery pass can cost up to 1,270 ms of backoff plus 8,000 ms of bounded +// fsync waits. +// +// When a pass does not clear the quarantine, +// `recoverQuarantinedCredentialCleanup` joins one more bounded attempt (up +// to 1,000 ms) before it gives up (codex-credentials.ts:616-619). So the +// full documented recovery path costs at least 8,000 + 1,270 + 1,000 = +// 10,270 ms. The vitest default `vi.waitFor` deadline is 1,000 ms, smaller +// than a single one of those inner bounds. So a poll can time out even +// though the call is still in progress. +// +// The helper deadline below adds margin on top of the 10,270 ms documented +// floor for the real filesystem work each attempt also does (two file +// removals and an intent-file delete, none of them bounded by +// `DIRECTORY_SYNC_OPERATION_TIMEOUT_MS`) and for this helper's own 50 ms +// poll granularity and Node event-loop scheduling jitter. A 30-second +// per-test budget keeps this wait reachable under the vitest per-test +// timeout, even for a test that runs the helper more than once. +const ACPX_OPERATION_WAIT_DEADLINE_MS = 15_000; +const ACPX_LONG_WAIT_TEST_TIMEOUT_MS = 30_000; + +/** + * Poll a credential or sandbox operation. Use a deadline derived from the + * real retry envelope described above. Unlike a bare `vi.waitFor`, report + * the last observed error when the deadline expires. Each attempt of + * `callback` is itself bounded by the remaining deadline, so a callback + * that stays pending cannot outlast the helper deadline and reach the + * enclosing vitest per-test timeout instead. + */ +async function waitForAcpxOperation( + callback: () => T | Promise, +): Promise { + const deadline = Date.now() + ACPX_OPERATION_WAIT_DEADLINE_MS; + let lastError: unknown = new Error( + "no attempt of this ACPX operation settled before the deadline", + ); + for (;;) { + try { + return await runAcpxOperationAttempt(callback, deadline); + } catch (error) { + lastError = error; + } + if (Date.now() >= deadline) { + const detail = + lastError instanceof Error + ? (lastError.stack ?? lastError.message) + : String(lastError); + throw new Error( + `ACPX operation did not settle within ${ACPX_OPERATION_WAIT_DEADLINE_MS}ms. Last observed error: ${detail}`, + { cause: lastError }, + ); + } + await new Promise((resolve) => setTimeout(resolve, 50)); + } +} + +/** + * Run one attempt of `callback`, bounded by the time remaining until + * `deadline`. A callback that is still pending when the remaining time + * runs out rejects with a timeout error instead of blocking the retry + * loop past the helper deadline. + * + * A rejected attempt does not cancel `callback`. If `callback` later + * resolves to a `ManagedCodexCredentialLease`, close that lease so its + * kernel lock and active lease generation do not stay allocated for the + * rest of the test run. + */ +async function runAcpxOperationAttempt( + callback: () => T | Promise, + deadline: number, +): Promise { + const remainingMs = Math.max(0, deadline - Date.now()); + let timer: ReturnType | undefined; + let timedOut = false; + const callbackResult = Promise.resolve().then(callback); + void callbackResult.then( + (value) => { + if (timedOut) void closeLateCredentialLease(value); + }, + () => undefined, + ); + try { + return await Promise.race([ + callbackResult, + new Promise((_resolve, reject) => { + timer = setTimeout(() => { + timedOut = true; + reject( + new Error( + `ACPX operation attempt did not settle within the remaining ${remainingMs}ms of the helper deadline`, + ), + ); + }, remainingMs); + }), + ]); + } finally { + if (timer) clearTimeout(timer); + } +} + +/** + * Close `value` if it is a `ManagedCodexCredentialLease` (or another lease + * with the same `close()` shape). Swallow a close failure so cleanup of one + * late lease cannot mask the original test failure. + */ +async function closeLateCredentialLease(value: unknown): Promise { + if ( + typeof value === "object" && + value !== null && + "close" in value && + typeof (value as { close: unknown }).close === "function" + ) { + await (value as { close(): Promise }).close().catch(() => undefined); + } +} + afterEach(async () => { for (const controller of admissionControllers.splice(0)) { if (!controller.signal.aborted) { @@ -211,15 +336,13 @@ describe("ACPX runtime host", () => { expect(createRuntime).not.toHaveBeenCalled(); expect(fixture.commandClose).toHaveBeenCalledOnce(); const authPath = join(credentialHome, "auth.json"); - const contender = await vi.waitFor( - () => - stageManagedCodexCredential({ - agentHomeDirectory: credentialHome, - environment: { - PAPERCLIP_ACPX_CODEX_AUTH_JSON_SECRET: '{"owner":"contender"}', - }, - }), - { timeout: 5_000 }, + const contender = await waitForAcpxOperation(() => + stageManagedCodexCredential({ + agentHomeDirectory: credentialHome, + environment: { + PAPERCLIP_ACPX_CODEX_AUTH_JSON_SECRET: '{"owner":"contender"}', + }, + }), ); await contender.close(); await expect(readFile(authPath)).rejects.toMatchObject({ code: "ENOENT" }); @@ -237,7 +360,7 @@ describe("ACPX runtime host", () => { ); await retryHost.close({ reason: "retry admission complete" }); expect(retryRuntime.close).toHaveBeenCalledOnce(); - }); + }, ACPX_LONG_WAIT_TEST_TIMEOUT_MS); it("composes admission, isolation, model verification, and cleanup", async () => { const fixture = await hostFixture(); @@ -585,7 +708,7 @@ describe("ACPX runtime host", () => { ), ).rejects.toThrow(/initialization and cleanup failed/); - await vi.waitFor(() => expect(runtime.close).toHaveBeenCalledTimes(2)); + await waitForAcpxOperation(() => expect(runtime.close).toHaveBeenCalledTimes(2)); await expect(readFile(authPath, "utf8")).resolves.toContain( "failed-admission", ); @@ -599,25 +722,23 @@ describe("ACPX runtime host", () => { ).rejects.toThrow("already has an active lease"); resolveRetryClose(); - await vi.waitFor(async () => { + await waitForAcpxOperation(async () => { await expect(readFile(authPath)).rejects.toMatchObject({ code: "ENOENT", }); }); // File removal precedes kernel lease release. Wait for the lease itself so // this assertion cannot race between those two ordered cleanup steps. - const contender = await vi.waitFor( - () => - stageManagedCodexCredential({ - agentHomeDirectory: credentialHome, - environment: { - PAPERCLIP_ACPX_CODEX_AUTH_JSON_SECRET: '{"owner":"contender"}', - }, - }), - { timeout: 5_000 }, + const contender = await waitForAcpxOperation(() => + stageManagedCodexCredential({ + agentHomeDirectory: credentialHome, + environment: { + PAPERCLIP_ACPX_CODEX_AUTH_JSON_SECRET: '{"owner":"contender"}', + }, + }), ); await contender.close(); - }); + }, ACPX_LONG_WAIT_TEST_TIMEOUT_MS); it("bounds post-handshake model verification and cleans the runtime", async () => { const fixture = await hostFixture(); @@ -732,14 +853,16 @@ describe("ACPX runtime host", () => { host.close({ reason: "retry close" }), ).resolves.toBeUndefined(); await expect(readFile(authPath)).rejects.toMatchObject({ code: "ENOENT" }); - const contender = await stageManagedCodexCredential({ - agentHomeDirectory: join(host.runtimeRoot(), "codex-home"), - environment: { - PAPERCLIP_ACPX_CODEX_AUTH_JSON_SECRET: '{"owner":"contender"}', - }, - }); + const contender = await waitForAcpxOperation(() => + stageManagedCodexCredential({ + agentHomeDirectory: join(host.runtimeRoot(), "codex-home"), + environment: { + PAPERCLIP_ACPX_CODEX_AUTH_JSON_SECRET: '{"owner":"contender"}', + }, + }), + ); await contender.close(); - }); + }, ACPX_LONG_WAIT_TEST_TIMEOUT_MS); it("scrubs credentials only after the exact pending runtime close resolves", async () => { const fixture = await hostFixture(); @@ -764,8 +887,8 @@ describe("ACPX runtime host", () => { const authPath = join(credentialHome, "auth.json"); const first = host.close({ reason: "runtime close pending" }); - await vi.waitFor(() => expect(runtime.close).toHaveBeenCalledOnce()); - await vi.waitFor(() => expect(fixture.commandClose).toHaveBeenCalledOnce()); + await waitForAcpxOperation(() => expect(runtime.close).toHaveBeenCalledOnce()); + await waitForAcpxOperation(() => expect(fixture.commandClose).toHaveBeenCalledOnce()); const second = host.close({ reason: "same exact close" }); await expect(readFile(authPath, "utf8")).resolves.toBe("{}"); await expect( @@ -784,14 +907,16 @@ describe("ACPX runtime host", () => { undefined, ]); await expect(readFile(authPath)).rejects.toMatchObject({ code: "ENOENT" }); - const contender = await stageManagedCodexCredential({ - agentHomeDirectory: credentialHome, - environment: { - PAPERCLIP_ACPX_CODEX_AUTH_JSON_SECRET: '{"owner":"contender"}', - }, - }); + const contender = await waitForAcpxOperation(() => + stageManagedCodexCredential({ + agentHomeDirectory: credentialHome, + environment: { + PAPERCLIP_ACPX_CODEX_AUTH_JSON_SECRET: '{"owner":"contender"}', + }, + }), + ); await contender.close(); - }); + }, ACPX_LONG_WAIT_TEST_TIMEOUT_MS); it("retains the exact pending cleanup while independent resources close", async () => { const fixture = await hostFixture(); @@ -814,7 +939,7 @@ describe("ACPX runtime host", () => { ); const first = host.close({ reason: "first close stalls" }); - await vi.waitFor(() => expect(runtime.close).toHaveBeenCalledOnce()); + await waitForAcpxOperation(() => expect(runtime.close).toHaveBeenCalledOnce()); const second = host.close({ reason: "same pending owner" }); let settled = false; void Promise.all([first, second]).finally(() => { @@ -824,7 +949,7 @@ describe("ACPX runtime host", () => { expect(settled).toBe(false); expect(runtime.close).toHaveBeenCalledOnce(); expect(fixture.commandClose).toHaveBeenCalledOnce(); - }); + }, ACPX_LONG_WAIT_TEST_TIMEOUT_MS); it("retries only after the exact close outcome settles with failure", async () => { const fixture = await hostFixture(); @@ -1095,9 +1220,9 @@ describe("ACPX runtime host", () => { close: lateCommandClose, }); - await vi.waitFor(() => expect(lateCommandClose).toHaveBeenCalledOnce()); + await waitForAcpxOperation(() => expect(lateCommandClose).toHaveBeenCalledOnce()); expect(openRuntime).not.toHaveBeenCalled(); - }); + }, ACPX_LONG_WAIT_TEST_TIMEOUT_MS); it("closes a credential lease that resolves after admission is aborted", async () => { const fixture = await hostFixture(); @@ -1158,10 +1283,10 @@ describe("ACPX runtime host", () => { close: lateCredentialClose, }); - await vi.waitFor(() => + await waitForAcpxOperation(() => expect(lateCredentialClose).toHaveBeenCalledTimes(2), ); - await vi.waitFor(async () => + await waitForAcpxOperation(async () => expect(readFile(lateCredentialPath)).rejects.toMatchObject({ code: "ENOENT", }), @@ -1174,7 +1299,7 @@ describe("ACPX runtime host", () => { }); expect(openRuntime).not.toHaveBeenCalled(); expect(fixture.commandClose).not.toHaveBeenCalled(); - }); + }, ACPX_LONG_WAIT_TEST_TIMEOUT_MS); it("retains managed credentials until an aborted late runtime is closed", async () => { const fixture = await hostFixture(); @@ -1237,7 +1362,7 @@ describe("ACPX runtime host", () => { ).rejects.toThrow("already has an active lease"); runtimeAdmission.resolve(lateRuntime); - await vi.waitFor(() => expect(lateRuntime.close).toHaveBeenCalledTimes(2)); + await waitForAcpxOperation(() => expect(lateRuntime.close).toHaveBeenCalledTimes(2)); expect(lateRuntime.close).toHaveBeenNthCalledWith(1, { reason: "ACPX runtime admission aborted", }); @@ -1252,25 +1377,23 @@ describe("ACPX runtime host", () => { ).rejects.toThrow("already has an active lease"); retryClose.resolve(undefined); - await vi.waitFor(async () => { + await waitForAcpxOperation(async () => { await expect(readFile(authPath)).rejects.toMatchObject({ code: "ENOENT", }); }); // File removal precedes kernel lease release. Wait for the lease itself so // this assertion cannot race between those two ordered cleanup steps. - const contender = await vi.waitFor( - () => - stageManagedCodexCredential({ - agentHomeDirectory: credentialHome, - environment: { - PAPERCLIP_ACPX_CODEX_AUTH_JSON_SECRET: '{"owner":"contender"}', - }, - }), - { timeout: 5_000 }, + const contender = await waitForAcpxOperation(() => + stageManagedCodexCredential({ + agentHomeDirectory: credentialHome, + environment: { + PAPERCLIP_ACPX_CODEX_AUTH_JSON_SECRET: '{"owner":"contender"}', + }, + }), ); await contender.close(); - }); + }, ACPX_LONG_WAIT_TEST_TIMEOUT_MS); it("scrubs credentials after rejected runtime cleanup is proven", async () => { const fixture = await hostFixture(); @@ -1319,8 +1442,8 @@ describe("ACPX runtime host", () => { expect(fixture.commandClose).toHaveBeenCalledOnce(); providerCleanup.resolve(undefined); - await vi.waitFor(() => expect(credentialClose).toHaveBeenCalledOnce()); - }); + await waitForAcpxOperation(() => expect(credentialClose).toHaveBeenCalledOnce()); + }, ACPX_LONG_WAIT_TEST_TIMEOUT_MS); }); function runtimePort( diff --git a/server/src/__tests__/secrets-service.test.ts b/server/src/__tests__/secrets-service.test.ts index 53e73b2ff8..dbc6050fdf 100644 --- a/server/src/__tests__/secrets-service.test.ts +++ b/server/src/__tests__/secrets-service.test.ts @@ -4,6 +4,7 @@ import { mkdir, rm } from "node:fs/promises"; import os from "node:os"; import path from "node:path"; import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "vitest"; +import type { MockInstance } from "vitest"; import { and, eq } from "drizzle-orm"; import { resolveCodexAuthCacheDir, withAccountHomeSecretMutationLock } from "@paperclipai/adapter-codex-local/server"; import { @@ -37,6 +38,112 @@ if (!embeddedPostgresSupport.supported) { ); } +// A deferred promise: a concurrency test resolves `resolve` from inside a +// mocked call, then a waiter `await`s `promise`. This proves the waiter's +// side reached a state, instead of guessing how long that state takes to +// reach. +function deferred() { + let resolve!: (value: T | PromiseLike) => void; + let reject!: (reason?: unknown) => void; + const promise = new Promise((promiseResolve, promiseReject) => { + resolve = promiseResolve; + reject = promiseReject; + }); + return { promise, resolve, reject }; +} + +// Waits for the given entry signal, but not blindly: if the operation +// itself settles first, the entry signal can never resolve, because the +// call never reached its mocked provider method. A plain `await` on the +// signal alone would then hang until the test's own timeout and hide the +// real error. Race the signal against the operation instead, so a create, +// rotate, or cleanup failure at setup surfaces immediately, at its own +// throw site. +async function awaitEntryOrOperationFailure( + entered: Promise, + operation: Promise, + label: string, +): Promise { + const failIfOperationSettlesFirst = operation.then(() => { + throw new Error(`${label}: the operation settled before it entered its mocked provider method`); + }); + // Attach a no-op handler so a later rejection here, once `entered` has + // already won the race below, never surfaces as an unhandled rejection. + failIfOperationSettlesFirst.catch(() => {}); + await Promise.race([entered, failIfOperationSettlesFirst]); +} + +// A wait built from a measured "uncontended entry" duration needs margin +// over that duration to absorb normal timing jitter, while it must still +// finish long before an unexcluded second operation could reach its own +// provider write. This multiple gives that margin. +const ENTRY_DETECTION_SAFETY_MULTIPLIER = 10; + +// The number of uncontended baseline calls to measure. One sample can be +// unusually fast by chance, which would understate real timing variance and +// let a broken lock slip past a too-short wait. The slowest of several +// samples gives a sturdier upper bound than any single sample alone. +const ENTRY_DETECTION_BASELINE_SAMPLE_COUNT = 3; + +// An absolute ceiling on the detection wait, independent of the measured +// baseline. A noisy baseline sample must never let this wait grow large +// enough to consume the test's own timeout. +const ENTRY_DETECTION_MAX_WAIT_MS = 3000; + +// Measures how long an uncontended call takes to reach a mocked provider +// method, by recording the time the mock is entered relative to the time +// the caller started. Repeats the measurement and keeps the slowest result, +// so the returned duration is a real, per-run upper bound, not a single +// possibly-lucky sample, and stays valid at any machine speed. +// +// Each baseline call must actually enter the mocked provider method. When +// one does not, that sample is meaningless, and a wait built from it would +// silently collapse toward its own one-millisecond floor instead of a real +// window. Throw here instead, so a broken baseline call fails loudly. +async function measureUncontendedEntryDurationMs( + spy: MockInstance<(...args: TArgs) => Promise>, + original: (...args: TArgs) => Promise, + triggerUncontendedCall: () => Promise, +): Promise { + let worstDurationMs = 0; + for (let sample = 0; sample < ENTRY_DETECTION_BASELINE_SAMPLE_COUNT; sample += 1) { + const startedAt = performance.now(); + let entered = false; + let enteredAt = startedAt; + spy.mockImplementationOnce(async (...args: TArgs) => { + entered = true; + enteredAt = performance.now(); + return original(...args); + }); + await triggerUncontendedCall(); + if (!entered) { + throw new Error( + "measureUncontendedEntryDurationMs: a baseline call never entered the mocked provider method, so it produced no valid measurement", + ); + } + worstDurationMs = Math.max(worstDurationMs, enteredAt - startedAt); + } + return worstDurationMs; +} + +// Waits long enough that a second operation, still queued behind a +// correctly excluding lock, cannot yet have reached its provider write — +// unless `violationSignal` resolves first. An unexcluded second operation +// resolves `violationSignal` itself, from inside its own mocked provider +// method, the instant it gets there, however long that takes: this ties the +// wait to the second operation's own confirmed progress, not to a blind +// sleep-then-check against a single guessed duration. The measured window +// below is only a ceiling on how long a correctly excluding lock is given +// to prove the second operation stayed queued. +function waitEntryDetectionWindow(uncontendedEntryDurationMs: number, violationSignal: Promise): Promise { + const waitMs = Math.min( + Math.max(uncontendedEntryDurationMs, 1) * ENTRY_DETECTION_SAFETY_MULTIPLIER, + ENTRY_DETECTION_MAX_WAIT_MS, + ); + const timeout = new Promise((resolve) => setTimeout(resolve, waitMs)); + return Promise.race([timeout, violationSignal]); +} + describeEmbeddedPostgres("secretService", () => { let stopDb: (() => Promise) | null = null; let db!: ReturnType; @@ -382,21 +489,42 @@ describeEmbeddedPostgres("secretService", () => { // whichever caller goes first. const companyId = await seedCompany(); const svc = secretService(db); + const originalCreateSecret = localEncryptedProvider.createSecret.bind(localEncryptedProvider); + const createSecretSpy = vi.spyOn(localEncryptedProvider, "createSecret"); + + // An uncontended create still crosses several asynchronous steps + // (directory checks, lock-root setup, database round trips) before it + // reaches its provider write. Measure that duration here, so the wait + // below can use a real measured value instead of a guessed sleep. + const uncontendedEntryDurationMs = await measureUncontendedEntryDurationMs( + createSecretSpy, + originalCreateSecret, + () => + svc.create(companyId, { + name: `baseline-${randomUUID()}`, + provider: "local_encrypted", + value: "/company/codex-home/acct-baseline", + }), + ); + const events: string[] = []; + const firstEntered = deferred(); + const secondEntered = deferred(); let releaseFirstWrite!: () => void; const firstWriteGate = new Promise((resolve) => { releaseFirstWrite = resolve; }); - const originalCreateSecret = localEncryptedProvider.createSecret.bind(localEncryptedProvider); - vi.spyOn(localEncryptedProvider, "createSecret") + createSecretSpy .mockImplementationOnce(async (input) => { events.push("first-provider-enter"); + firstEntered.resolve(); await firstWriteGate; events.push("first-provider-exit"); return originalCreateSecret(input); }) .mockImplementationOnce(async (input) => { events.push("second-provider-enter"); + secondEntered.resolve(); return originalCreateSecret(input); }); @@ -405,22 +533,49 @@ describeEmbeddedPostgres("secretService", () => { provider: "local_encrypted", value: "/company/codex-home/acct-a", }); - // Give the first call a chance to acquire the lock and enter its provider - // write before the second call starts racing for the same lock. - await new Promise((resolve) => setTimeout(resolve, 20)); - const secondCreate = svc.create(companyId, { - name: `hand-named-${randomUUID()}`, - provider: "local_encrypted", - value: "/company/codex-home/acct-a", - }); - // The second call must stay blocked on the lock while the first call - // still holds it: it must never enter its own provider write before the - // first call's provider write exits. - await new Promise((resolve) => setTimeout(resolve, 20)); - expect(events).toEqual(["first-provider-enter"]); - - releaseFirstWrite(); - await Promise.all([firstCreate, secondCreate]); + let secondCreate: ReturnType | undefined; + let outcomes: PromiseSettledResult[] = []; + try { + // Wait for the confirmed signal that the first call now holds the + // lock and sits inside its provider write. The lock stays held until + // we release it below, so the second call, once we start it, must + // contend for the same lock while the first call still holds it. Race + // against the call's own promise, so a setup failure that happens + // before the call ever reaches the lock surfaces immediately, at its + // own throw site, instead of hanging this wait until the test + // timeout. + await awaitEntryOrOperationFailure(firstEntered.promise, firstCreate, "firstCreate"); + secondCreate = svc.create(companyId, { + name: `hand-named-${randomUUID()}`, + provider: "local_encrypted", + value: "/company/codex-home/acct-a", + }); + // Wait a safety multiple of the measured uncontended entry duration, + // or until the second call itself confirms it reached its provider + // write, whichever comes first. A correctly excluding lock keeps the + // second call queued for the whole wait, so this cannot produce a + // false failure. A broken lock resolves `secondEntered` on its own, + // from inside the second call's mocked provider method, the instant + // it gets there. + await waitEntryDetectionWindow(uncontendedEntryDurationMs, secondEntered.promise); + // The lock is still held (we have not released it yet), so the second + // call must still be queued behind it and must not have entered its + // provider write. + expect(events).toEqual(["first-provider-enter"]); + } finally { + // Release and settle both calls even when the check above fails, so + // neither call stays parked inside the lock past this test and + // corrupts teardown. + releaseFirstWrite(); + outcomes = await Promise.allSettled([firstCreate, secondCreate]); + } + for (const outcome of outcomes) { + if (outcome.status === "rejected") throw outcome.reason; + } + // The lock enforces this order: the second call cannot start its + // provider write until the first call's whole locked operation + // completes. This final order is proof of mutual exclusion, not a + // timing guess. expect(events).toEqual(["first-provider-enter", "first-provider-exit", "second-provider-enter"]); }); @@ -435,7 +590,26 @@ describeEmbeddedPostgres("secretService", () => { provider: "local_encrypted", value: "/company/codex-home/acct-b", }); + const originalCreateSecret = localEncryptedProvider.createSecret.bind(localEncryptedProvider); + const createSecretSpy = vi.spyOn(localEncryptedProvider, "createSecret"); + + // The contended call below is a create, so measure how long an + // uncontended create takes to reach its own provider write. The wait + // later in this test uses that measured duration, not a guessed sleep. + const uncontendedEntryDurationMs = await measureUncontendedEntryDurationMs( + createSecretSpy, + originalCreateSecret, + () => + svc.create(companyId, { + name: `baseline-${randomUUID()}`, + provider: "local_encrypted", + value: "/company/codex-home/acct-baseline", + }), + ); + const events: string[] = []; + const rotateEntered = deferred(); + const createEntered = deferred(); let releaseRotateWrite!: () => void; const rotateWriteGate = new Promise((resolve) => { releaseRotateWrite = resolve; @@ -443,28 +617,60 @@ describeEmbeddedPostgres("secretService", () => { const originalCreateVersion = localEncryptedProvider.createVersion.bind(localEncryptedProvider); vi.spyOn(localEncryptedProvider, "createVersion").mockImplementationOnce(async (input) => { events.push("rotate-provider-enter"); + rotateEntered.resolve(); await rotateWriteGate; events.push("rotate-provider-exit"); return originalCreateVersion(input); }); - const originalCreateSecret = localEncryptedProvider.createSecret.bind(localEncryptedProvider); - vi.spyOn(localEncryptedProvider, "createSecret").mockImplementationOnce(async (input) => { + createSecretSpy.mockImplementationOnce(async (input) => { events.push("create-provider-enter"); + createEntered.resolve(); return originalCreateSecret(input); }); const rotateCall = svc.rotate(existing.id, { value: "/company/codex-home/acct-b-rotated" }); - await new Promise((resolve) => setTimeout(resolve, 20)); - const createCall = svc.create(companyId, { - name: `hand-named-${randomUUID()}`, - provider: "local_encrypted", - value: "/company/codex-home/acct-b", - }); - await new Promise((resolve) => setTimeout(resolve, 20)); - expect(events).toEqual(["rotate-provider-enter"]); - - releaseRotateWrite(); - await Promise.all([rotateCall, createCall]); + let createCall: ReturnType | undefined; + let outcomes: PromiseSettledResult[] = []; + try { + // Wait for the confirmed signal that the rotate now holds the lock + // and sits inside its provider write. The lock stays held until we + // release it below, so the create call, once we start it, must + // contend for the same lock while the rotate still holds it. Race + // against the call's own promise, so a setup failure that happens + // before the call ever reaches the lock surfaces immediately, at its + // own throw site, instead of hanging this wait until the test + // timeout. + await awaitEntryOrOperationFailure(rotateEntered.promise, rotateCall, "rotateCall"); + createCall = svc.create(companyId, { + name: `hand-named-${randomUUID()}`, + provider: "local_encrypted", + value: "/company/codex-home/acct-b", + }); + // Wait a safety multiple of the measured uncontended entry duration, + // or until the create call itself confirms it reached its provider + // write, whichever comes first. A correctly excluding lock keeps the + // create call queued for the whole wait, so this cannot produce a + // false failure. A broken lock resolves `createEntered` on its own, + // from inside the create call's mocked provider method, the instant + // it gets there. + await waitEntryDetectionWindow(uncontendedEntryDurationMs, createEntered.promise); + // The lock is still held (we have not released it yet), so the create + // call must still be queued behind it and must not have entered its + // provider write. + expect(events).toEqual(["rotate-provider-enter"]); + } finally { + // Release and settle both calls even when the check above fails, so + // neither call stays parked inside the lock past this test and + // corrupts teardown. + releaseRotateWrite(); + outcomes = await Promise.allSettled([rotateCall, createCall]); + } + for (const outcome of outcomes) { + if (outcome.status === "rejected") throw outcome.reason; + } + // The lock enforces this order: the create call cannot start its + // provider write until the rotate's whole locked operation completes. + // This final order is proof of mutual exclusion, not a timing guess. expect(events).toEqual(["rotate-provider-enter", "rotate-provider-exit", "create-provider-enter"]); }); @@ -478,32 +684,79 @@ describeEmbeddedPostgres("secretService", () => { const companyId = await seedCompany(); const svc = secretService(db); const accountHomeDir = await makeAccountHomeDir(companyId, "acct-queued-create"); + const originalCreateSecret = localEncryptedProvider.createSecret.bind(localEncryptedProvider); + const createSecretSpy = vi.spyOn(localEncryptedProvider, "createSecret"); + + // The contended call below is a create, so measure how long an + // uncontended create takes to reach its own provider write. The wait + // later in this test uses that measured duration, not a guessed sleep. + const uncontendedEntryDurationMs = await measureUncontendedEntryDurationMs( + createSecretSpy, + originalCreateSecret, + () => + svc.create(companyId, { + name: `baseline-${randomUUID()}`, + provider: "local_encrypted", + value: "/company/codex-home/acct-baseline", + }), + ); const events: string[] = []; + const cleanupEntered = deferred(); + const createEntered = deferred(); let releaseCleanup!: () => void; const cleanupGate = new Promise((resolve) => { releaseCleanup = resolve; }); const cleanupCall = withAccountHomeSecretMutationLock(undefined, companyId, async () => { events.push("cleanup-enter"); + cleanupEntered.resolve(); await cleanupGate; await rm(accountHomeDir, { recursive: true, force: true }); events.push("cleanup-exit"); }); - // Give the cleanup a chance to acquire the lock before the create starts - // racing for the same lock. - await new Promise((resolve) => setTimeout(resolve, 20)); - const createCall = svc.create(companyId, { - name: `account-home-${randomUUID()}`, - provider: "local_encrypted", - value: accountHomeDir, + createSecretSpy.mockImplementationOnce(async (input) => { + events.push("create-provider-enter"); + createEntered.resolve(); + return originalCreateSecret(input); }); - // The create must stay queued behind the held lock. - await new Promise((resolve) => setTimeout(resolve, 20)); - expect(events).toEqual(["cleanup-enter"]); - - releaseCleanup(); - await cleanupCall; + let createCall: ReturnType | undefined; + let outcomes: PromiseSettledResult[] = []; + try { + // Wait for the confirmed signal that the cleanup now holds the lock. + // The lock stays held until we release it below, so the create call, + // once we start it, must queue behind the cleanup. Race against the + // cleanup's own promise, so a setup failure that happens before the + // cleanup ever reaches the lock surfaces immediately, at its own + // throw site, instead of hanging this wait until the test timeout. + await awaitEntryOrOperationFailure(cleanupEntered.promise, cleanupCall, "cleanupCall"); + createCall = svc.create(companyId, { + name: `account-home-${randomUUID()}`, + provider: "local_encrypted", + value: accountHomeDir, + }); + // Wait a safety multiple of the measured uncontended entry duration, + // or until the create call itself confirms it reached its provider + // write, whichever comes first. A correctly excluding lock keeps the + // create call queued behind the cleanup's still-held lock for the + // whole wait, so this cannot produce a false failure. A broken lock + // resolves `createEntered` on its own, from inside the create call's + // mocked provider method, the instant it clears the directory check + // and gets there — the directory still exists until the cleanup + // (still paused on its own gate) actually removes it. + await waitEntryDetectionWindow(uncontendedEntryDurationMs, createEntered.promise); + // The lock is still held (we have not released it yet), so the create + // call must still be queued behind it and must not have entered its + // provider write. + expect(events).toEqual(["cleanup-enter"]); + } finally { + // Release and settle both calls even when the check above fails, so + // neither call stays parked inside the lock past this test and + // corrupts teardown. + releaseCleanup(); + outcomes = await Promise.allSettled([cleanupCall, createCall]); + } + if (outcomes[0]?.status === "rejected") throw outcomes[0].reason; await expect(createCall).rejects.toThrow(/no longer exists/); }); @@ -519,25 +772,75 @@ describeEmbeddedPostgres("secretService", () => { value: "/some/unrelated/placeholder/value", }); const accountHomeDir = await makeAccountHomeDir(companyId, "acct-queued-rotate"); + const originalCreateVersion = localEncryptedProvider.createVersion.bind(localEncryptedProvider); + const createVersionSpy = vi.spyOn(localEncryptedProvider, "createVersion"); + + // The contended call below is a rotate, so measure how long an + // uncontended rotate takes to reach its own provider write. The wait + // later in this test uses that measured duration, not a guessed sleep. + const baselineSecret = await svc.create(companyId, { + name: `baseline-${randomUUID()}`, + provider: "local_encrypted", + value: "/some/unrelated/placeholder/baseline", + }); + const uncontendedEntryDurationMs = await measureUncontendedEntryDurationMs( + createVersionSpy, + originalCreateVersion, + () => svc.rotate(baselineSecret.id, { value: "/some/unrelated/placeholder/baseline-rotated" }), + ); const events: string[] = []; + const cleanupEntered = deferred(); + const rotateEntered = deferred(); let releaseCleanup!: () => void; const cleanupGate = new Promise((resolve) => { releaseCleanup = resolve; }); const cleanupCall = withAccountHomeSecretMutationLock(undefined, companyId, async () => { events.push("cleanup-enter"); + cleanupEntered.resolve(); await cleanupGate; await rm(accountHomeDir, { recursive: true, force: true }); events.push("cleanup-exit"); }); - await new Promise((resolve) => setTimeout(resolve, 20)); - const rotateCall = svc.rotate(existing.id, { value: accountHomeDir }); - await new Promise((resolve) => setTimeout(resolve, 20)); - expect(events).toEqual(["cleanup-enter"]); - - releaseCleanup(); - await cleanupCall; + createVersionSpy.mockImplementationOnce(async (input) => { + events.push("rotate-provider-enter"); + rotateEntered.resolve(); + return originalCreateVersion(input); + }); + let rotateCall: ReturnType | undefined; + let outcomes: PromiseSettledResult[] = []; + try { + // Wait for the confirmed signal that the cleanup now holds the lock. + // The lock stays held until we release it below, so the rotate call, + // once we start it, must queue behind the cleanup. Race against the + // cleanup's own promise, so a setup failure that happens before the + // cleanup ever reaches the lock surfaces immediately, at its own + // throw site, instead of hanging this wait until the test timeout. + await awaitEntryOrOperationFailure(cleanupEntered.promise, cleanupCall, "cleanupCall"); + rotateCall = svc.rotate(existing.id, { value: accountHomeDir }); + // Wait a safety multiple of the measured uncontended entry duration, + // or until the rotate call itself confirms it reached its provider + // write, whichever comes first. A correctly excluding lock keeps the + // rotate call queued behind the cleanup's still-held lock for the + // whole wait, so this cannot produce a false failure. A broken lock + // resolves `rotateEntered` on its own, from inside the rotate call's + // mocked provider method, the instant it clears the directory check + // and gets there — the directory still exists until the cleanup + // (still paused on its own gate) actually removes it. + await waitEntryDetectionWindow(uncontendedEntryDurationMs, rotateEntered.promise); + // The lock is still held (we have not released it yet), so the + // rotate call must still be queued behind it and must not have + // entered its provider write. + expect(events).toEqual(["cleanup-enter"]); + } finally { + // Release and settle both calls even when the check above fails, so + // neither call stays parked inside the lock past this test and + // corrupts teardown. + releaseCleanup(); + outcomes = await Promise.allSettled([cleanupCall, rotateCall]); + } + if (outcomes[0]?.status === "rejected") throw outcomes[0].reason; await expect(rotateCall).rejects.toThrow(/no longer exists/); }); diff --git a/ui/public/sw.js b/ui/public/sw.js index e5997304dc..9a9d1a7ee6 100644 --- a/ui/public/sw.js +++ b/ui/public/sw.js @@ -1,4 +1,11 @@ -const CACHE_NAME = "paperclip-v2"; +// The build id is stamped into this file at production build time (see +// stampServiceWorkerBuildId in vite.config.ts), so a deploy that changes only +// the app bundle still changes sw.js byte-for-byte. That is what makes the +// browser install a new worker, which — via skipWaiting + controllerchange — +// reloads parked tabs onto the fresh bundle. Left as the literal placeholder in +// dev, where HMR (not the worker) drives refreshes. +const BUILD_ID = "__PAPERCLIP_BUILD_ID__"; +const CACHE_NAME = `paperclip-${BUILD_ID}`; self.addEventListener("install", () => { self.skipWaiting(); diff --git a/ui/src/lib/vite-sw-build-id.test.ts b/ui/src/lib/vite-sw-build-id.test.ts new file mode 100644 index 0000000000..5b978a9f84 --- /dev/null +++ b/ui/src/lib/vite-sw-build-id.test.ts @@ -0,0 +1,63 @@ +import fs from "node:fs"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; +import { describe, expect, it } from "vitest"; +import { + SERVICE_WORKER_BUILD_ID_PLACEHOLDER, + deriveBuildIdFromEntryFileName, + stampServiceWorkerBuildId, +} from "./vite-sw-build-id"; + +const swSource = () => + `const BUILD_ID = "${SERVICE_WORKER_BUILD_ID_PLACEHOLDER}";\n` + + "const CACHE_NAME = `paperclip-${BUILD_ID}`;\n"; + +describe("stampServiceWorkerBuildId", () => { + it("replaces the placeholder with the build id and leaves no placeholder", () => { + const out = stampServiceWorkerBuildId(swSource(), "index-abc123"); + expect(out).toContain("index-abc123"); + expect(out).not.toContain(SERVICE_WORKER_BUILD_ID_PLACEHOLDER); + }); + + it("produces different worker bytes for different build ids", () => { + // This is the whole point: a new bundle -> a new sw.js -> a new worker -> + // parked tabs reload. Identical build ids must stay byte-identical so the + // worker does not churn when the app did not change. + const a = stampServiceWorkerBuildId(swSource(), "index-aaaaaa"); + const b = stampServiceWorkerBuildId(swSource(), "index-bbbbbb"); + const again = stampServiceWorkerBuildId(swSource(), "index-aaaaaa"); + expect(a).not.toEqual(b); + expect(a).toEqual(again); + }); + + it("throws when the placeholder is missing so a drifted worker fails the build", () => { + expect(() => stampServiceWorkerBuildId("const CACHE_NAME = 'paperclip';", "x")).toThrow( + /placeholder/, + ); + }); + + it("throws on an empty build id rather than shipping a nameless cache", () => { + expect(() => stampServiceWorkerBuildId(swSource(), "")).toThrow(); + }); +}); + +describe("deriveBuildIdFromEntryFileName", () => { + it("uses the content-hashed entry file name", () => { + expect(deriveBuildIdFromEntryFileName("assets/index-BHbrFFmp.js")).toBe("index-BHbrFFmp"); + }); + + it("sanitizes characters that are unsafe in a cache name", () => { + expect(deriveBuildIdFromEntryFileName("assets/index @weird!.js")).toBe("index--weird-"); + }); +}); + +describe("public/sw.js contract", () => { + it("still contains the placeholder the plugin rewrites", () => { + const swPath = path.resolve( + path.dirname(fileURLToPath(import.meta.url)), + "../../public/sw.js", + ); + const source = fs.readFileSync(swPath, "utf8"); + expect(source).toContain(SERVICE_WORKER_BUILD_ID_PLACEHOLDER); + }); +}); diff --git a/ui/src/lib/vite-sw-build-id.ts b/ui/src/lib/vite-sw-build-id.ts new file mode 100644 index 0000000000..d9ea0a13d8 --- /dev/null +++ b/ui/src/lib/vite-sw-build-id.ts @@ -0,0 +1,81 @@ +import fs from "node:fs"; +import path from "node:path"; +import type { Plugin } from "vite"; + +/** + * Stamp the service worker with a per-build id so bundle-only deploys still + * refresh parked tabs. + * + * `sw.js` is a static public asset copied verbatim into the build, and its + * update machinery (`service-worker-updates.ts`) only reloads a parked tab when + * a *new* worker takes control — which happens only when `sw.js` changes + * byte-for-byte. Without this, a deploy that ships a new app bundle but the same + * `sw.js` installs no new worker, so an open tab keeps running the old bundle + * until someone reloads by hand. Rewriting the placeholder with a value derived + * from the bundle makes `sw.js` change exactly when the app does. + */ + +export const SERVICE_WORKER_BUILD_ID_PLACEHOLDER = "__PAPERCLIP_BUILD_ID__"; + +/** + * Replace the build-id placeholder in a service-worker source string. + * + * Throws when the placeholder is absent: that means the worker drifted away + * from the contract (renamed or removed placeholder) and would ship a service + * worker that never rotates — the exact bug this plugin exists to prevent — so + * a loud build failure beats a silent no-op. + */ +export function stampServiceWorkerBuildId(source: string, buildId: string): string { + if (!source.includes(SERVICE_WORKER_BUILD_ID_PLACEHOLDER)) { + throw new Error( + `service worker is missing the ${SERVICE_WORKER_BUILD_ID_PLACEHOLDER} placeholder; ` + + "the build cannot stamp a build id and parked tabs would not refresh after a deploy", + ); + } + if (!buildId) { + throw new Error("refusing to stamp the service worker with an empty build id"); + } + return source.split(SERVICE_WORKER_BUILD_ID_PLACEHOLDER).join(buildId); +} + +/** + * Derive a build id from the emitted bundle. The entry chunk's file name + * carries a content hash that changes whenever the app code changes and stays + * stable when it does not, so the worker rotates precisely with the app. + */ +export function deriveBuildIdFromEntryFileName(entryFileName: string): string { + const base = path.basename(entryFileName).replace(/\.js$/, ""); + // Keep only characters that are safe inside a Cache Storage name. + const sanitized = base.replace(/[^A-Za-z0-9_-]/g, "-"); + return sanitized || "build"; +} + +export function serviceWorkerBuildIdPlugin( + options: { serviceWorkerFileName?: string } = {}, +): Plugin { + const serviceWorkerFileName = options.serviceWorkerFileName ?? "sw.js"; + let buildId: string | null = null; + let outDir = "dist"; + + return { + name: "paperclip-sw-build-id", + apply: "build", + configResolved(config) { + outDir = config.build.outDir; + }, + generateBundle(_options, bundle) { + const entry = Object.values(bundle).find( + (chunk) => chunk.type === "chunk" && chunk.isEntry, + ); + if (entry) { + buildId = deriveBuildIdFromEntryFileName(entry.fileName); + } + }, + closeBundle() { + const swPath = path.resolve(outDir, serviceWorkerFileName); + const source = fs.readFileSync(swPath, "utf8"); + const stamped = stampServiceWorkerBuildId(source, buildId ?? "build"); + fs.writeFileSync(swPath, stamped); + }, + }; +} diff --git a/ui/vite.config.ts b/ui/vite.config.ts index d235797fec..3ac9f91485 100644 --- a/ui/vite.config.ts +++ b/ui/vite.config.ts @@ -4,11 +4,12 @@ import react from "@vitejs/plugin-react"; import tailwindcss from "@tailwindcss/vite"; import { createUiDevWatchOptions } from "./src/lib/vite-watch"; import { createApiProxy } from "./src/lib/vite-api-proxy"; +import { serviceWorkerBuildIdPlugin } from "./src/lib/vite-sw-build-id"; const apiProxy = createApiProxy(); export default defineConfig(({ mode }) => ({ - plugins: [react(), tailwindcss()], + plugins: [react(), tailwindcss(), serviceWorkerBuildIdPlugin()], build: { minify: "esbuild", },