/** Real Codex gate: cancellation isolation, RPC steer, and stored-thread resume. * Uses the existing Codex login unchanged and only disposable test threads. */ import assert from 'node:assert/strict'; import fs from 'node:fs'; import os from 'node:os'; import path from 'node:path'; import { randomUUID } from 'node:crypto'; const root = fs.mkdtempSync(path.join(os.tmpdir(), 'supen-runtime-e2e-')); process.env.SUPEN_HOME = path.join(root, 'supen'); const cwd = path.join(root, 'workspace'); fs.mkdirSync(cwd); const { getSharedCodexAppServerHost } = await import('../src/core/codex-app-server-host.js'); const { CodexAppServerDriver } = await import('../src/agent-sdk/drivers/codex-app-server-driver.js'); const { executeSupenRpc } = await import('../src/http/routes/rpc.js'); const host = getSharedCodexAppServerHost(); const driver = new CodexAppServerDriver(() => host); const events: Array> = []; const unsubscribe = host.subscribeNotifications((event) => events.push(event)); const ids = new Set(); const marker = `SUPEN-CONTROL-${Date.now()}`; const waitFor = async (predicate: () => boolean, label: string) => { const deadline = Date.now() + 120_000; while (!predicate()) { if (Date.now() > deadline) throw new Error(`Timed out: ${label}`); await new Promise((resolve) => setTimeout(resolve, 100)); } }; const create = (id: string, resume?: string) => driver.startThread({ agentId: 'codex', threadId: id, turnId: id, cwd, permissionMode: 'workspace-write', networkAccess: false, ...(resume ? { resume, resumeRequired: true } : {}), }); const consume = async (session: Awaited>) => { let threadId = ''; let turnId = ''; let text = ''; try { for await (const output of session.stream()) { const event = output as Record; if (event.method === 'turn/started') { threadId = event.params.threadId; turnId = event.params.turn.id; ids.add(threadId); } if (event.method === 'item/agentMessage/delta') text += event.params.delta; } return { threadId, turnId, text, error: null }; } catch (error) { return { threadId, turnId, text, error: String(error) }; } }; let a: Awaited> | undefined; let b: Awaited> | undefined; try { const missingId = randomUUID(); const missing = await create('missing-resume-canary', missingId); await missing.send('continue'); const missingResult = await consume(missing); assert.match(missingResult.error || '', /Failed to resume Codex thread/); assert(missingResult.error?.includes(missingId)); assert.equal(missingResult.threadId, ''); assert.equal(missingResult.turnId, ''); console.log(JSON.stringify({ requiredResumeRejectedMissingThread: true })); a = await create('cancel-canary'); await a.send('Use the execution tool to run sleep 60. After it finishes reply CANCEL-MISSED. Do not modify any files.'); const aDone = consume(a); const toolStarted = (threadId: string) => events.some((e) => e.method === 'item/started' && e.params?.threadId === threadId && ['commandExecution', 'mcpToolCall'].includes(e.params?.item?.type)); await waitFor(() => ids.size === 1 && toolStarted([...ids][0]), 'A executing a tool'); const aId = [...ids][0]; b = await create('steer-canary'); await b.send('Use the execution tool to run sleep 8. After it finishes reply BEFORE-STEER. Do not modify any files.'); const bDone = consume(b); await waitFor(() => ids.size === 2 && toolStarted([...ids][1]), 'B executing a tool'); const bId = [...ids][1]; const bTurn = events.find((e) => e.method === 'turn/started' && e.params?.threadId === bId)?.params.turn.id; a.close(); const steer = await executeSupenRpc({ id: 1, method: 'turn/steer', params: { agentId: 'codex', threadId: bId, expectedTurnId: bTurn, input: [{ type: 'text', text: `Change your final answer to exactly ${marker}.` }], } }, { adminContext: null, onMessage: () => { throw new Error('Unexpected legacy dispatch'); }, surface: 'cli' }); assert(!('error' in steer), JSON.stringify(steer)); const [aResult, bResult] = await Promise.all([aDone, bDone]); assert.match(aResult.error || '', /interrupted/); assert.equal(bResult.error, null); assert(bResult.text.includes(marker), JSON.stringify(bResult)); const aTerminal = events.find((e) => e.method === 'turn/completed' && e.params?.threadId === aId); assert.equal(aTerminal?.params.turn.status, 'interrupted'); assert.equal(events.filter((e) => e.method === 'turn/started' && e.params?.threadId === bId).length, 1); console.log(JSON.stringify({ aId, bId, nativeCancellation: true, siblingCompleted: true, steerStayedInTurn: true })); host.close(); const resumed = await create('resume-canary', bId); await resumed.send('What exact final answer did you just give? Repeat it only, without tools.'); const result = await consume(resumed); assert.equal(result.threadId, bId); assert.equal(result.error, null); assert(result.text.includes(marker)); const history = await host.request('thread/read', { threadId: bId, includeTurns: true }); assert(JSON.stringify(history.result).includes(marker)); console.log(JSON.stringify({ threadId: bId, resumedAfterHostRestart: true, historyVerified: true })); } finally { a?.close(); b?.close(); try { for (const threadId of ids) { await host.request('thread/archive', { threadId }); console.log(JSON.stringify({ threadId, archived: true })); } } finally { unsubscribe(); host.close(); fs.rmSync(root, { recursive: true, force: true }); } }