diff --git a/extensions/codex/media-understanding-provider.test.ts b/extensions/codex/media-understanding-provider.test.ts index bf5fee328923..69aec650f068 100644 --- a/extensions/codex/media-understanding-provider.test.ts +++ b/extensions/codex/media-understanding-provider.test.ts @@ -11,10 +11,14 @@ const sharedClientMocks = vi.hoisted(() => ({ createIsolatedCodexAppServerClient: vi.fn(), })); -vi.mock("./src/app-server/shared-client.js", () => ({ - createIsolatedCodexAppServerClient: sharedClientMocks.createIsolatedCodexAppServerClient, - retireSharedCodexAppServerClientIfCurrent: () => undefined, -})); +vi.mock("./src/app-server/shared-client.js", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + createIsolatedCodexAppServerClient: sharedClientMocks.createIsolatedCodexAppServerClient, + retireSharedCodexAppServerClientIfCurrent: () => undefined, + }; +}); function codexModel(inputModalities: string[] = ["text", "image"]) { return { @@ -93,6 +97,7 @@ function createFakeClient(options?: { approvalRequestMethod?: string; responseText?: string; onTurnStart?: () => void; + beforeRequest?: (method: string) => Promise; }) { const notifications = new Set<(notification: CodexServerNotification) => void>(); type RequestHandler = Parameters[0]; @@ -101,6 +106,9 @@ function createFakeClient(options?: { const approvalResponses: JsonValue[] = []; const request = vi.fn(async (method: string, params?: JsonValue) => { requests.push({ method, params }); + if (options?.beforeRequest) { + await options.beforeRequest(method); + } if (method === "model/list") { return { data: [codexModel(options?.inputModalities)], @@ -521,10 +529,132 @@ describe("codex media understanding provider", () => { expect(closeAndWait).toHaveBeenCalledOnce(); }); + it.each([0, 200, -200])( + "keeps one media-turn budget through selection retry after a %s ms clock step", + async (clockStepMs) => { + vi.useFakeTimers({ toFake: ["Date", "performance", "setTimeout", "clearTimeout"] }); + let wallOffsetMs = 0; + vi.spyOn(Date, "now").mockImplementation( + () => 1_700_000_000_000 + performance.now() + wallOffsetMs, + ); + + const firstStartup = createDeferred(); + const releaseFirstStartup = createDeferred(); + const selectionCheck = createDeferred(); + const releaseSelectionCheck = createDeferred(); + const retryStartup = createDeferred(); + const releaseRetryStartup = createDeferred(); + const turnStarted = createDeferred(); + const caller = new AbortController(); + const selectionChanged = Object.assign(new Error("managed executable selection changed"), { + code: "CODEX_APP_SERVER_START_SELECTION_CHANGED", + }); + const first = createFakeClient({ + beforeRequest: async (method) => { + if (method === "thread/start") { + selectionCheck.resolve(); + await releaseSelectionCheck.promise; + throw selectionChanged; + } + }, + }); + const second = createFakeClient({ + deferTurnCompletion: true, + onTurnStart: () => turnStarted.resolve(), + }); + sharedClientMocks.createIsolatedCodexAppServerClient + .mockImplementationOnce(async () => { + firstStartup.resolve(); + await releaseFirstStartup.promise; + return first.client; + }) + .mockImplementationOnce(async () => { + retryStartup.resolve(); + await releaseRetryStartup.promise; + return second.client; + }); + + const provider = buildCodexMediaUnderstandingProvider(); + if (!provider.describeImage) { + throw new Error("media provider must expose image understanding"); + } + let settled = false; + const observed = provider + .describeImage({ + buffer: Buffer.from("image-bytes"), + fileName: "image.png", + mime: "image/png", + provider: "codex", + model: "gpt-5.4", + timeoutMs: 1_000, + signal: caller.signal, + cfg: {}, + agentDir: "/tmp/openclaw-agent", + }) + .then( + (value) => { + settled = true; + return value; + }, + (error: unknown) => { + settled = true; + return error; + }, + ); + const reach = async (boundary: Promise) => { + const reached = await Promise.race([boundary.then(() => true), observed.then(() => false)]); + expect(reached, "operation settled before the required scenario boundary").toBe(true); + }; + + try { + await reach(firstStartup.promise); + await vi.advanceTimersByTimeAsync(100); + releaseFirstStartup.resolve(); + await reach(selectionCheck.promise); + await vi.advanceTimersByTimeAsync(200); + wallOffsetMs = clockStepMs; + releaseSelectionCheck.resolve(); + + await reach(retryStartup.promise); + expect(first.closeAndWait).toHaveBeenCalledOnce(); + await vi.advanceTimersByTimeAsync(100); + releaseRetryStartup.resolve(); + await reach(turnStarted.promise); + + await vi.advanceTimersByTimeAsync(599); + expect(settled).toBe(false); + expect(second.requests.filter(({ method }) => method === "turn/interrupt")).toEqual([]); + await vi.advanceTimersByTimeAsync(1); + expect(settled).toBe(true); + expect(await observed).toMatchObject({ + name: "TimeoutError", + message: "codex app-server image understanding turn timed out after 1s", + }); + expect(second.requests.filter(({ method }) => method === "turn/interrupt")).toEqual([ + { method: "turn/interrupt", params: { threadId: "thread-1", turnId: "turn-1" } }, + ]); + expect(first.requests.some(({ method }) => method === "turn/start")).toBe(false); + expect(first.requests.some(({ method }) => method === "turn/interrupt")).toBe(false); + expect(second.closeAndWait).toHaveBeenCalledOnce(); + expect(caller.signal.aborted).toBe(false); + expect(sharedClientMocks.createIsolatedCodexAppServerClient).toHaveBeenCalledTimes(2); + expect(vi.getTimerCount()).toBe(0); + } finally { + releaseFirstStartup.resolve(); + releaseSelectionCheck.resolve(); + releaseRetryStartup.resolve(); + caller.abort("retry proof cleanup"); + await observed; + vi.restoreAllMocks(); + vi.useRealTimers(); + } + }, + ); + it("clamps oversized image understanding turn timeouts", async () => { // The bounded timer subtracts startup time from its clamped deadline. // Freeze the clock so the clamp assertion cannot lose a real millisecond. - const dateNowSpy = vi.spyOn(Date, "now").mockReturnValue(1_700_000_000_000); + const clockSpy = vi.spyOn(performance, "now").mockReturnValue(1_000); const setTimeoutSpy = vi.spyOn(globalThis, "setTimeout"); try { const { client } = createFakeClient(); @@ -546,7 +676,7 @@ describe("codex media understanding provider", () => { expect(result?.text).toBe("A red square."); expect(setTimeoutSpy).toHaveBeenCalledWith(expect.any(Function), MAX_TIMER_TIMEOUT_MS); } finally { - dateNowSpy.mockRestore(); + clockSpy.mockRestore(); vi.restoreAllMocks(); vi.clearAllTimers(); vi.useRealTimers(); diff --git a/extensions/codex/src/app-server/bounded-turn.clock.test.ts b/extensions/codex/src/app-server/bounded-turn.clock.test.ts new file mode 100644 index 000000000000..6cb7cde1e5da --- /dev/null +++ b/extensions/codex/src/app-server/bounded-turn.clock.test.ts @@ -0,0 +1,183 @@ +import { once } from "node:events"; +import { describe, expect, it, vi } from "vitest"; +import { WebSocketServer } from "ws"; +import { runBoundedCodexAppServerTurn } from "./bounded-turn.js"; +import { CodexAppServerClient } from "./client.js"; +import { threadStartResult, turnStartResult } from "./codex-app-server.test-fixtures.js"; +import type { RpcRequest } from "./protocol.js"; +import { CODEX_APP_SERVER_VERSION } from "./version.js"; + +describe("bounded Codex turn elapsed deadlines over WebSocket", () => { + it.each([0, 60_000, -60_000])( + "keeps its real timeout after a %s ms startup clock correction", + async (clockStepMs) => { + // Only the protocol peer is controlled: client, socket, and timers are real. + const server = new WebSocketServer({ host: "127.0.0.1", port: 0 }); + const methods: string[] = []; + const thread = threadStartResult("clock-thread", process.cwd()); + const turn = turnStartResult("clock-turn"); + server.on("connection", (socket) => { + socket.on("message", (data) => { + const buffer = Array.isArray(data) + ? Buffer.concat(data) + : Buffer.isBuffer(data) + ? data + : Buffer.from(data); + const request = JSON.parse(buffer.toString("utf8")) as RpcRequest; + methods.push(request.method); + const respond = (result: unknown) => { + socket.send(JSON.stringify({ id: request.id, result })); + }; + switch (request.method) { + case "initialize": + respond({ userAgent: `codex/${CODEX_APP_SERVER_VERSION}` }); + break; + case "initialized": + break; + case "model/list": + respond({ + data: [ + { + id: thread.model, + model: thread.model, + upgrade: null, + upgradeInfo: null, + availabilityNux: null, + displayName: "Clock proof model", + description: "Controlled protocol fixture; no provider calls", + hidden: false, + isDefault: true, + inputModalities: ["text"], + supportedReasoningEfforts: [{ reasoningEffort: "low", description: "Low" }], + defaultReasoningEffort: "low", + supportsPersonality: false, + multiAgentVersion: null, + additionalSpeedTiers: [], + serviceTiers: [], + defaultServiceTier: null, + }, + ], + nextCursor: null, + }); + break; + case "thread/start": + respond(thread); + break; + case "turn/start": + respond(turn); + break; + case "turn/interrupt": + respond({}); + socket.send( + JSON.stringify({ + method: "turn/completed", + params: { + threadId: thread.thread.id, + ...turnStartResult(turn.turn.id, "interrupted"), + }, + }), + ); + break; + default: + socket.send( + JSON.stringify({ + id: request.id, + error: { code: -32601, message: `Unexpected method: ${request.method}` }, + }), + ); + } + }); + }); + + const realDateNow = Date.now; + const wallClock = vi.spyOn(Date, "now"); + const caller = new AbortController(); + let client: CodexAppServerClient | undefined; + let initialized = false; + let run: Promise | undefined; + let watchdog: ReturnType | undefined; + try { + await once(server, "listening"); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("expected a loopback WebSocket port"); + } + const startedAt = performance.now(); + // Bound the old backward-clock bug without simulating its expiry. + watchdog = setTimeout(() => { + caller.abort(new Error("clock proof watchdog")); + if (!initialized) { + client?.close(); + } + }, 4_000); + run = runBoundedCodexAppServerTurn({ + model: { mode: "required", id: thread.model }, + timeoutMs: 1_000, + signal: caller.signal, + options: { + clientFactory: async () => { + client = await CodexAppServerClient.start({ + transport: "websocket", + url: `ws://127.0.0.1:${address.port}`, + headers: {}, + authToken: undefined, + }); + await client.initialize(); + initialized = true; + // Apply one process-local offset after startup, never an OS clock change. + wallClock.mockImplementation(() => realDateNow() + clockStepMs); + return client; + }, + }, + taskLabel: "clock proof", + developerInstructions: "Wait for cancellation.", + input: [{ type: "text", text: "Clock proof.", text_elements: [] }], + requiredModalities: ["text"], + isolation: "configured-transport", + }).catch((error: unknown) => error); + const outcome = await run; + const elapsedMs = performance.now() - startedAt; + console.info("bounded-turn-clock", { + clockStepMs, + elapsedMs: Math.round(elapsedMs), + outcome: outcome instanceof Error ? outcome.name : "completed", + watchdogAborted: caller.signal.aborted, + methods, + }); + expect(methods).toEqual([ + "initialize", + "initialized", + "model/list", + "thread/start", + "turn/start", + "turn/interrupt", + ]); + expect(outcome).toMatchObject({ + name: "TimeoutError", + message: "codex app-server clock proof turn timed out after 1s", + }); + expect(caller.signal.aborted).toBe(false); + expect(elapsedMs).toBeGreaterThanOrEqual(900); + expect(elapsedMs).toBeLessThan(4_000); + } finally { + clearTimeout(watchdog); + wallClock.mockRestore(); + caller.abort(); + await run; + await client?.closeAndWait(); + await Promise.all( + [...server.clients].map((socket) => { + const closed = once(socket, "close"); + socket.terminate(); + return closed; + }), + ); + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + expect(server.clients.size).toBe(0); + expect(server.address()).toBeNull(); + } + }, + ); +}); diff --git a/extensions/codex/src/app-server/bounded-turn.ts b/extensions/codex/src/app-server/bounded-turn.ts index 61f969dc4da2..582572ea20a6 100644 --- a/extensions/codex/src/app-server/bounded-turn.ts +++ b/extensions/codex/src/app-server/bounded-turn.ts @@ -154,8 +154,9 @@ async function runBoundedCodexAppServerTurnInWorkspace( ): Promise { const totalTimeoutMs = timing?.timeoutMs ?? resolveTimerTimeoutMs(params.timeoutMs, 100, 100); const timeoutError = new CodexBoundedTurnTimeoutError(params.taskLabel, totalTimeoutMs); - const deadline = timing?.deadline ?? Date.now() + totalTimeoutMs; - const timeoutMs = deadline - Date.now(); + // Startup and selection retries share an elapsed budget, not a wall-clock deadline. + const deadline = timing?.deadline ?? performance.now() + totalTimeoutMs; + const timeoutMs = deadline - performance.now(); if (timeoutMs <= 0) { throw timeoutError; } @@ -215,7 +216,7 @@ async function runBoundedCodexAppServerTurnInWorkspace( } else { params.signal?.addEventListener("abort", abortFromCaller, { once: true }); } - const remainingRunMs = deadline - Date.now(); + const remainingRunMs = deadline - performance.now(); if (remainingRunMs <= 0) { abortRun(timeoutError); }