diff --git a/.github/workflows/pi-owner-runtime.yml b/.github/workflows/pi-owner-runtime.yml index 373afb7a539..9f9df4f9b0b 100644 --- a/.github/workflows/pi-owner-runtime.yml +++ b/.github/workflows/pi-owner-runtime.yml @@ -2,8 +2,9 @@ name: Pi owner runtime verification on: pull_request: paths: - - 'src/main/pi/agent-status-handler-source.ts' + - 'src/main/pi/**' - 'tests/tools/pi-owner-runtime-smoke.mjs' + - 'tests/tools/omp-completion-runtime-smoke.mjs' - '.github/workflows/pi-owner-runtime.yml' workflow_dispatch: permissions: @@ -27,3 +28,5 @@ jobs: run: npm install --prefix .cache/pi-owner --ignore-scripts --no-audit --no-fund @earendil-works/pi-coding-agent@0.83.0 - name: Verify real owner exit and hook delivery run: node tests/tools/pi-owner-runtime-smoke.mjs .cache/pi-owner/node_modules/@earendil-works/pi-coding-agent + - name: Verify OMP completion over native HTTP + run: node tests/tools/omp-completion-runtime-smoke.mjs diff --git a/src/main/pi/agent-status-completion-delivery.test.ts b/src/main/pi/agent-status-completion-delivery.test.ts new file mode 100644 index 00000000000..4c99563b6ed --- /dev/null +++ b/src/main/pi/agent-status-completion-delivery.test.ts @@ -0,0 +1,150 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' + +import { createAgentStatusExtensionHarness } from './agent-status-extension-test-harness' + +function events(mock: ReturnType): unknown[] { + return mock.mock.calls.map((call) => JSON.parse(String(call[1]?.body)).payload) +} + +describe('OMP completion delivery', () => { + beforeEach(() => vi.useFakeTimers()) + afterEach(() => vi.useRealTimers()) + + it.each(['rejection', 'HTTP failure'])( + 'retries a final %s without another turn', + async (failure) => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + await harness.callHook('agent_start') + await vi.advanceTimersByTimeAsync(0) + if (failure === 'rejection') { + harness.fetchMock.mockRejectedValueOnce(new Error('offline')) + } else { + harness.fetchMock.mockResolvedValueOnce({ ok: false, status: 503 }) + } + await harness.callHook('agent_end') + await vi.advanceTimersByTimeAsync(251) + expect(events(harness.fetchMock)).toEqual([ + { hook_event_name: 'agent_start' }, + { hook_event_name: 'agent_end' }, + { hook_event_name: 'agent_end' } + ]) + expect(vi.getTimerCount()).toBe(0) + } + ) + + it.each([22, null])('retries failed or timed-out WSL curl (exit %s)', async (curlExitCode) => { + const harness = createAgentStatusExtensionHarness({ + kind: 'omp', + env: { WSL_DISTRO_NAME: 'Ubuntu' }, + existsSync: (path) => path === '/mnt/c/Windows/System32/curl.exe', + curlExitCode, + fetchImpl: async () => { + throw new Error('guest unavailable') + } + }) + await harness.callHook('agent_end') + await vi.advanceTimersByTimeAsync(curlExitCode === null ? 11251 : 251) + expect(harness.spawnMock).toHaveBeenCalledTimes(2) + await harness.callHook('session_shutdown') + await vi.advanceTimersByTimeAsync(12000) + expect(harness.spawnMock).toHaveBeenCalledTimes(2) + expect(vi.getTimerCount()).toBe(0) + }) + + it('does not acknowledge a missing WSL curl bridge', async () => { + const harness = createAgentStatusExtensionHarness({ + kind: 'omp', + env: { WSL_DISTRO_NAME: 'Ubuntu' }, + fetchImpl: async () => { + throw new Error('guest unavailable') + } + }) + await harness.callHook('agent_end') + await vi.advanceTimersByTimeAsync(251) + expect(harness.fetchMock).toHaveBeenCalledTimes(2) + await harness.callHook('session_shutdown') + }) + + it('waits for WSL curl acknowledgment before draining a newer snapshot', async () => { + const harness = createAgentStatusExtensionHarness({ + kind: 'omp', + env: { WSL_DISTRO_NAME: 'Ubuntu' }, + existsSync: (path) => path === '/mnt/c/Windows/System32/curl.exe', + curlExitCode: null, + fetchImpl: async () => { + throw new Error('guest unavailable') + } + }) + await harness.callHook('agent_end') + await vi.advanceTimersByTimeAsync(0) + await harness.callHook('agent_start') + expect(harness.fetchMock).toHaveBeenCalledTimes(1) + harness.spawnedChildren[0]?.emit('close', 0) + await vi.advanceTimersByTimeAsync(0) + expect(harness.fetchMock).toHaveBeenCalledTimes(2) + harness.spawnedChildren[1]?.emit('close', 0) + await vi.advanceTimersByTimeAsync(1000) + expect(vi.getTimerCount()).toBe(0) + }) + + it.each(['before_agent_start', 'agent_start', 'session_shutdown', 'session_switch'])( + 'retires a failed completion at %s', + async (boundary) => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + harness.fetchMock.mockRejectedValueOnce(new Error('offline')) + await harness.callHook('agent_end') + await vi.advanceTimersByTimeAsync(0) + await harness.callHook(boundary, {}) + await vi.advanceTimersByTimeAsync(10_000) + expect( + events(harness.fetchMock).filter((event) => JSON.stringify(event).includes('agent_end')) + ).toHaveLength(1) + expect(vi.getTimerCount()).toBe(0) + } + ) + + it('does not schedule retries when a pending completion fails after a new start', async () => { + let rejectDelivery: ((error: Error) => void) | undefined + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + harness.fetchMock.mockImplementationOnce( + () => + new Promise((_resolve, reject) => { + rejectDelivery = reject + }) + ) + await harness.callHook('agent_end') + await harness.callHook('agent_start') + rejectDelivery?.(new Error('late failure')) + await vi.advanceTimersByTimeAsync(10_000) + expect(events(harness.fetchMock)).toEqual([ + { hook_event_name: 'agent_end' }, + { hook_event_name: 'agent_start' } + ]) + }) + + it('retries a timed out completion without blocking agent handlers', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + harness.fetchMock.mockImplementationOnce(() => new Promise(() => {})) + await harness.callHook('agent_end') + await vi.advanceTimersByTimeAsync(1251) + expect(events(harness.fetchMock)).toEqual([ + { hook_event_name: 'agent_end' }, + { hook_event_name: 'agent_end' } + ]) + expect(harness.fetchMock.mock.calls[0]?.[1]?.signal.aborted).toBe(true) + expect(vi.getTimerCount()).toBe(0) + }) + + it('bounds retries when Orca stays unreachable', async () => { + const harness = createAgentStatusExtensionHarness({ + kind: 'omp', + fetchImpl: async () => { + throw new Error('offline') + } + }) + await harness.callHook('agent_end') + await vi.advanceTimersByTimeAsync(30_000) + expect(harness.fetchMock).toHaveBeenCalledTimes(4) + expect(vi.getTimerCount()).toBe(0) + }) +}) diff --git a/src/main/pi/agent-status-extension-source.test.ts b/src/main/pi/agent-status-extension-source.test.ts index fa9823d76dd..7df089f13b9 100644 --- a/src/main/pi/agent-status-extension-source.test.ts +++ b/src/main/pi/agent-status-extension-source.test.ts @@ -349,6 +349,7 @@ describe('getPiAgentStatusExtensionSource', () => { expect(command).toBe('/mnt/c/Windows/System32/curl.exe') expect(args).toEqual([ '-sS', + '--fail', '--connect-timeout', '3', '--max-time', diff --git a/src/main/pi/agent-status-extension-source.ts b/src/main/pi/agent-status-extension-source.ts index 38775ca1973..ed2da5965ca 100644 --- a/src/main/pi/agent-status-extension-source.ts +++ b/src/main/pi/agent-status-extension-source.ts @@ -1,3 +1,4 @@ +import { getPiAgentStatusPostQueueSourceLines } from './agent-status-post-queue-source' // Why: pi has no settings.json hook surface — its extensibility is the // in-process TypeScript extension API (pi.on('agent_start'), 'tool_call', // etc.). To get pi panes into the unified agent-hooks pipeline alongside @@ -100,9 +101,8 @@ export function getPiAgentStatusExtensionSource(kind: PiAgentKind = 'pi'): strin '// critical path, and the latest-only pending slot prevents a stalled', '// Orca receiver from building an unbounded queue of obsolete snapshots.', 'const HOOK_POST_TIMEOUT_MS = 1000', - 'let activePost = false', + ...getPiAgentStatusPostQueueSourceLines(), ...(kind === 'pi' ? ['let piUiPromptDepth = 0', 'let piTurnInFlight = false'] : []), - 'let pendingPost: { hookEventName: string; extra: Record; metadata: Record; ompRuntime: boolean } | null = null', ...sessionMetadataSourceLines, '', '// Why: re-reading the endpoint file on every event is cheap (small file,', @@ -163,31 +163,26 @@ export function getPiAgentStatusExtensionSource(kind: PiAgentKind = 'pi'): strin '', 'function post(hookEventName: string, extra: Record = {}): void {', ' const ompRuntime = isOmpRuntime()', + ' cancelPostRetry()', + ' const metadata = getPostSessionMetadata(ompRuntime)', + '// Model changes must not erase an unacknowledged completion in the latest-only slot.', + " const previousCompletion = latestPost?.hookEventName === 'agent_end' && !latestPost.delivered && latestPost.metadata.session_id === metadata.session_id", ' pendingPost = {', - ' hookEventName,', + ' revision: ++postRevision,', + ' attempts: 0,', + ' delivered: false,', + " hookEventName: ompRuntime && hookEventName === 'model_select' && previousCompletion ? 'agent_end' : hookEventName,", // Why: every coalesced snapshot must retain an open modal, not just its start event. kind === 'pi' ? ' extra: { ...extra, ...(!ompRuntime && piUiPromptDepth > 0 ? { ui_prompt_active: true } : {}) },' : ' extra,', - ' metadata: getPostSessionMetadata(ompRuntime),', + ' metadata,', ' ompRuntime,', ' }', + ' latestPost = pendingPost', ' drainPosts()', '}', '', - 'function drainPosts(): void {', - ' if (activePost || !pendingPost) return', - ' const next = pendingPost', - ' pendingPost = null', - ' activePost = true', - ' void postOnce(next.hookEventName, next.extra, next.metadata, next.ompRuntime)', - ' .catch(() => {})', - ' .finally(() => {', - ' activePost = false', - ' drainPosts()', - ' })', - '}', - '', 'async function postOnce(', ' hookEventName: string,', ' extra: Record,', @@ -217,7 +212,7 @@ export function getPiAgentStatusExtensionSource(kind: PiAgentKind = 'pi'): strin " if (typeof timeout.unref === 'function') timeout.unref()", ' })', ' try {', - ' await Promise.race([', + ' const response = await Promise.race([', ' fetch(url, {', " method: 'POST',", ' headers: {', @@ -229,11 +224,12 @@ export function getPiAgentStatusExtensionSource(kind: PiAgentKind = 'pi'): strin ' }),', ' timeoutPromise,', ' ])', - ' } catch {', + " if (!response.ok) throw new Error('Orca hook HTTP ' + response.status)", + ' } catch (error) {', ' // Why: status reporting must never fail the pi run just because Orca', ' // is unavailable or the loopback request failed (e.g. Orca restart).', - ' if (!isWslRuntime()) return', - ' postViaWindowsCurl(body, ompRuntime)', + ' if (!isWslRuntime()) throw error', + ' await postViaWindowsCurl(body, ompRuntime)', ' } finally {', ' if (timeout) clearTimeout(timeout)', ' }', diff --git a/src/main/pi/agent-status-extension-test-harness.ts b/src/main/pi/agent-status-extension-test-harness.ts index 810bc3d04d5..42c5ca3cbe0 100644 --- a/src/main/pi/agent-status-extension-test-harness.ts +++ b/src/main/pi/agent-status-extension-test-harness.ts @@ -1,3 +1,4 @@ +import { EventEmitter } from 'node:events' import { runInNewContext } from 'node:vm' // TypeScript 7 is a native CLI; transpile tests still need the legacy JavaScript API. import ts from 'typescript-api' @@ -6,7 +7,10 @@ import { vi } from 'vitest' import { getPiAgentStatusExtensionSource } from './agent-status-extension-source' export type HookContext = { + hasUI?: boolean + ui?: { setEditorText?: (text: string) => void } isIdle?: () => boolean + model?: { provider?: unknown; id?: unknown } | null sessionManager?: { getSessionId?: () => unknown getSessionFile?: () => unknown @@ -16,6 +20,8 @@ export type HookContext = { export type HookHandler = (event?: unknown, context?: HookContext) => Promise | void type FakeCurlChild = { + kill: ReturnType + emit: (event: string, ...args: unknown[]) => boolean on: ReturnType stdin: { on: ReturnType @@ -66,6 +72,7 @@ export function createAgentStatusExtensionHarness(args: { existsSync?: (path: string) => boolean readFileSync?: (path: string, encoding: string) => string statSync?: (path: string) => { mtimeMs: number; size: number; ino: number } + curlExitCode?: number | null fetchImpl?: (...params: Parameters) => Promise }): AgentStatusExtensionHarness { const fetchMock = vi.fn( @@ -77,14 +84,20 @@ export function createAgentStatusExtensionHarness(args: { const spawnedChildren: FakeCurlChild[] = [] const spawnMock = vi.fn(() => { + const emitter = new EventEmitter() const child: FakeCurlChild = { - on: vi.fn(), + emit: emitter.emit.bind(emitter), + kill: vi.fn(() => emitter.emit('close', null)), + on: vi.fn(emitter.on.bind(emitter)), stdin: { on: vi.fn(), end: vi.fn() } } spawnedChildren.push(child) + if (args.curlExitCode !== null) { + void Promise.resolve().then(() => emitter.emit('close', args.curlExitCode ?? 0)) + } return child }) diff --git a/src/main/pi/agent-status-handler-source.ts b/src/main/pi/agent-status-handler-source.ts index 5a778a1c81f..8467ec9ea7e 100644 --- a/src/main/pi/agent-status-handler-source.ts +++ b/src/main/pi/agent-status-handler-source.ts @@ -114,6 +114,20 @@ export function getPiAgentStatusHandlerSourceLines(kind: PiAgentKind): string[] ' const selfPid = String(process.pid)', ' if (ownerPid && ownerPid !== selfPid && isStatusOwnerAlive(ownerPid)) return', ` process.env.${ownerEnv} = selfPid`, + ' resetPostQueue()', + ...(kind !== 'pi' + ? [" pi.on('session_shutdown', () => { resetPostQueue(); clearPendingAgentEndCheck() })"] + : []), + ...(kind !== 'prime-agent' + ? [ + " pi.on('session_switch', (_event, ctx) => {", + ' if (!isOmpRuntime()) return', + ' resetPostQueue()', + ' clearPendingAgentEndCheck()', + ' updateRuntimeOmpSessionMetadata(ctx)', + ' })' + ] + : []), ...sessionStartHandler, ` pi.on('before_agent_start', (event${ctxParam}) => {`, ...captureSessionMetadata, diff --git a/src/main/pi/agent-status-post-queue-source.ts b/src/main/pi/agent-status-post-queue-source.ts new file mode 100644 index 00000000000..512ed932339 --- /dev/null +++ b/src/main/pi/agent-status-post-queue-source.ts @@ -0,0 +1,52 @@ +export function getPiAgentStatusPostQueueSourceLines(): string[] { + return [ + 'type HookPost = { hookEventName: string; extra: Record; metadata: Record; ompRuntime: boolean; revision: number; attempts: number; delivered: boolean }', + 'let activePost = false', + 'let pendingPost: HookPost | null = null', + 'let latestPost: HookPost | null = null', + '// A newer snapshot or session boundary retires every older retry.', + 'let postRevision = 0', + 'let retryTimer: ReturnType | null = null', + '', + 'function cancelPostRetry(): void {', + ' if (retryTimer !== null) clearTimeout(retryTimer)', + ' retryTimer = null', + '}', + '', + 'function resetPostQueue(): void {', + ' cancelPostRetry()', + ' postRevision++', + ' pendingPost = null', + ' latestPost = null', + '}', + '', + 'function drainPosts(): void {', + ' if (activePost || !pendingPost) return', + ' const next = pendingPost', + ' pendingPost = null', + ' activePost = true', + ' void postOnce(next.hookEventName, next.extra, next.metadata, next.ompRuntime)', + ' .then(() => { next.delivered = true })', + ' .catch(() => {', + ' if (!next.ompRuntime || next.revision !== postRevision) return', + ' if (next.attempts >= 3) {', + " console.warn('[orca-pi-status] hook delivery failed after retries:', next.hookEventName)", + ' return', + ' }', + ' const delay = 250 * 2 ** next.attempts++', + ' retryTimer = setTimeout(() => {', + ' retryTimer = null', + ' if (next.revision !== postRevision) return', + ' pendingPost = next', + ' drainPosts()', + ' }, delay)', + " if (typeof retryTimer.unref === 'function') retryTimer.unref()", + ' })', + ' .finally(() => {', + ' activePost = false', + ' drainPosts()', + ' })', + '}', + '' + ] +} diff --git a/src/main/pi/agent-status-ui-prompt-source.ts b/src/main/pi/agent-status-ui-prompt-source.ts index 5790c5c30a7..b641a5b3328 100644 --- a/src/main/pi/agent-status-ui-prompt-source.ts +++ b/src/main/pi/agent-status-ui-prompt-source.ts @@ -32,6 +32,8 @@ export function getPiAgentStatusUiPromptHandlerSourceLines(kind: PiAgentKind): s ' })', '', " pi.on('session_shutdown', () => {", + ' resetPostQueue()', + ' clearPendingAgentEndCheck()', ' if (isOmpRuntime()) return', ' // Why: pi tears an open dialog down through resetExtensionUI without resolving its', ' // promise, so a replaced session never emits the matching ui_prompt_end and the wait', diff --git a/src/main/pi/agent-status-wsl-curl-source.ts b/src/main/pi/agent-status-wsl-curl-source.ts index 558e25a49db..4e9a5265587 100644 --- a/src/main/pi/agent-status-wsl-curl-source.ts +++ b/src/main/pi/agent-status-wsl-curl-source.ts @@ -43,20 +43,20 @@ export function getPiAgentStatusWslCurlSourceLines(): string[] { '}', '', '// Why: WSL loopback is not the Windows loopback, so use curl.exe on the host.', - 'function postViaWindowsCurl(body: string, ompRuntime: boolean): void {', + 'async function postViaWindowsCurl(body: string, ompRuntime: boolean): Promise {', ' const curlPath = resolveWindowsCurlPath()', ' const windowsPort = process.env.ORCA_AGENT_HOOK_PORT', ' const windowsToken = process.env.ORCA_AGENT_HOOK_TOKEN', - ' if (!curlPath || !windowsPort || !windowsToken) return', + " if (!curlPath || !windowsPort || !windowsToken) throw new Error('Orca WSL hook bridge unavailable')", ' // Why: a stale guest endpoint must fall back to current host coordinates.', ' const windowsUrl = `http://127.0.0.1:${windowsPort}${resolveHookPath(ompRuntime)}`', - ' try {', + ' await new Promise((resolve, reject) => {', " const { spawn } = require('child_process')", ' const child = spawn(', ' curlPath,', ' [', - " '-sS',", - ' // Why: detached delivery may take seconds under loaded WSL interop.', + " '-sS', '--fail',", + ' // WSL interop needs a bounded, acknowledged delivery window.', " '--connect-timeout', '3',", " '--max-time', '10',", " '--noproxy', '127.0.0.1',", @@ -69,12 +69,20 @@ export function getPiAgentStatusWslCurlSourceLines(): string[] { ' ],', " { stdio: ['pipe', 'ignore', 'ignore'] }", ' )', - " child.on('error', () => {})", - " child.stdin.on('error', () => {})", + ' const timer = setTimeout(() => {', + ' try { child.kill() } catch {}', + " reject(new Error('Orca WSL hook delivery timed out'))", + ' }, 11000)', + " if (typeof timer.unref === 'function') timer.unref()", + ' const fail = (error: Error): void => { clearTimeout(timer); reject(error) }', + " child.on('error', fail)", + " child.stdin.on('error', fail)", + " child.on('close', (code: number | null) => {", + ' clearTimeout(timer)', + " if (code === 0) resolve(); else reject(new Error('Orca WSL hook exit ' + code))", + ' })', ' child.stdin.end(body)', - ' } catch {', - ' // Why: status delivery must not surface inside the agent TUI.', - ' }', + ' })', '}', '' ] diff --git a/src/main/pi/titlebar-extension-source.test.ts b/src/main/pi/titlebar-extension-source.test.ts index eb826c1903e..b60eacc7be9 100644 --- a/src/main/pi/titlebar-extension-source.test.ts +++ b/src/main/pi/titlebar-extension-source.test.ts @@ -228,6 +228,21 @@ describe('getPiTitlebarExtensionSource', () => { expect(transferTitles[0]).toMatch(BRAILLE_RE) }) + it.each([{ kind: 'omp' as const }, { processTitle: 'omp' }])( + 'stops final OMP completion without waiting for isIdle: %j', + async (options) => { + const isIdle = vi.fn(() => false) + const harness = createHarness({ ...options, isIdle }) + await harness.callHook('agent_start') + await harness.callHook('agent_end', { willContinue: false }) + const completedTitle = harness.lastTitle() + expect(completedTitle).not.toMatch(BRAILLE_RE) + await vi.advanceTimersByTimeAsync(1000) + expect(harness.lastTitle()).toBe(completedTitle) + expect(isIdle).not.toHaveBeenCalled() + } + ) + it('keeps spinning across a non-terminal OMP agent_end', async () => { const harness = createHarness() diff --git a/src/main/pi/titlebar-extension-source.ts b/src/main/pi/titlebar-extension-source.ts index 570d23a15d3..bd77e169dd5 100644 --- a/src/main/pi/titlebar-extension-source.ts +++ b/src/main/pi/titlebar-extension-source.ts @@ -255,7 +255,7 @@ export function getPiTitlebarExtensionSource(kind: PiAgentKind = 'pi'): string { ' resetPromptState()', ' })', '', - ' // Why: modern Pi/OMP emit agent_end mid-run and only settle later, so settlement is the', + ' // Why: modern Pi emits agent_end mid-run and only settle later, so settlement is the', ' // authoritative completion boundary. Legacy runtimes never emit it, so agent_end stays.', " on('agent_settled', async (_event, ctx) => {", ' stopAnimation(ctx)', @@ -266,7 +266,7 @@ export function getPiTitlebarExtensionSource(kind: PiAgentKind = 'pi'): string { ' clearPendingAgentEndCheck()', ' return', ' }', - " if (!ctx || typeof ctx.isIdle !== 'function') {", + ` if (${kind === 'omp' ? 'true' : kind === 'pi' ? 'isOmpRuntime()' : 'false'} || !ctx || typeof ctx.isIdle !== 'function') {`, ' stopAnimation(ctx)', ' return', ' }', diff --git a/tests/e2e/omp-completion-retry-status.spec.ts b/tests/e2e/omp-completion-retry-status.spec.ts new file mode 100644 index 00000000000..e264c6ff3ae --- /dev/null +++ b/tests/e2e/omp-completion-retry-status.spec.ts @@ -0,0 +1,107 @@ +import { once } from 'node:events' +import { createServer } from 'node:http' +import { createRequire } from 'node:module' +import { join } from 'node:path' +import { runInNewContext } from 'node:vm' +import { transform } from 'esbuild' +import { getPiAgentStatusExtensionSource } from '../../src/main/pi/agent-status-extension-source' +import { test, expect } from './helpers/orca-app' +import { readHookEndpoint } from './helpers/agent-hook-endpoint' +import { ensureTerminalVisible, waitForActiveWorktree, waitForSessionReady } from './helpers/store' +import { waitForActivePaneHookDescriptor, waitForActiveTerminalManager } from './helpers/terminal' + +test('OMP completion retry clears the rendered working indicator', async ({ + orcaPage, + electronApp +}, testInfo) => { + await waitForSessionReady(orcaPage) + await waitForActiveWorktree(orcaPage) + await ensureTerminalVisible(orcaPage) + await waitForActiveTerminalManager(orcaPage, 30_000) + const endpoint = await readHookEndpoint(electronApp) + const { paneKey, worktreeId } = await waitForActivePaneHookDescriptor(orcaPage) + let completions = 0 + const proxy = createServer(async (request, response) => { + let body = '' + for await (const chunk of request) { + body += chunk + } + if (JSON.parse(body).payload.hook_event_name === 'agent_end' && ++completions === 1) { + response.writeHead(503).end() + return + } + const forwarded = await fetch(`http://127.0.0.1:${endpoint.port}/hook/omp`, { + method: 'POST', + headers: { 'Content-Type': 'application/json', 'X-Orca-Agent-Hook-Token': endpoint.token }, + body + }) + response.writeHead(forwarded.status).end() + }) + proxy.listen(0, '127.0.0.1') + await once(proxy, 'listening') + try { + const address = proxy.address() + if (!address || typeof address === 'string') { + throw new Error('Expected TCP listener') + } + type Handler = (event: Record, ctx: { isIdle: () => boolean }) => void + const handlers = new Map() + const module: { + exports: { default?: (api: { on: (name: string, fn: Handler) => void }) => void } + } = { exports: {} } + const { code } = await transform(getPiAgentStatusExtensionSource('omp'), { + loader: 'ts', + format: 'cjs' + }) + runInNewContext(code, { + module, + exports: module.exports, + require: createRequire(join(process.cwd(), 'package.json')), + process: { + pid: process.pid, + argv: [], + title: 'omp', + env: { + ORCA_PANE_KEY: paneKey, + ORCA_TAB_ID: paneKey.split(':')[0], + ORCA_WORKTREE_ID: worktreeId, + ORCA_AGENT_HOOK_PORT: String(address.port), + ORCA_AGENT_HOOK_TOKEN: endpoint.token, + ORCA_AGENT_HOOK_ENV: endpoint.env, + ORCA_AGENT_HOOK_VERSION: endpoint.version + } + }, + fetch, + AbortController, + Buffer, + console, + setTimeout, + clearTimeout + }) + expect(module.exports.default).toBeDefined() + module.exports.default?.({ on: (name, fn) => handlers.set(name, fn) }) + const working = orcaPage.locator('[aria-label="Working"]') + handlers.get('before_agent_start')?.( + { prompt: 'OMP completion recovery' }, + { isIdle: () => false } + ) + handlers.get('agent_start')?.({}, { isIdle: () => false }) + await expect(working.first()).toBeVisible() + await orcaPage.screenshot({ path: testInfo.outputPath('before-working.png') }) + handlers.get('agent_end')?.({ willContinue: false }, { isIdle: () => false }) + await expect.poll(() => completions).toBe(2) + await expect(working).toHaveCount(0) + await expect + .poll(() => + orcaPage.evaluate( + (key) => window.__store?.getState().agentStatusByPaneKey[key]?.state, + paneKey + ) + ) + .toBe('done') + await orcaPage.screenshot({ path: testInfo.outputPath('after-completed.png') }) + } finally { + proxy.closeAllConnections() + proxy.close() + } +}) diff --git a/tests/tools/omp-completion-runtime-smoke.mjs b/tests/tools/omp-completion-runtime-smoke.mjs new file mode 100644 index 00000000000..6c32c18db32 --- /dev/null +++ b/tests/tools/omp-completion-runtime-smoke.mjs @@ -0,0 +1,108 @@ +import assert from 'node:assert/strict' +import { once } from 'node:events' +import { mkdtemp, rm, writeFile } from 'node:fs/promises' +import { createServer } from 'node:http' +import { createRequire } from 'node:module' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { setTimeout as delay } from 'node:timers/promises' +import { build, transform } from 'esbuild' + +const scratch = await mkdtemp(join(tmpdir(), 'orca-omp-completion-')) +const received = [] +let rejectCompletion = true +const server = createServer(async (request, response) => { + let body = '' + for await (const chunk of request) { + body += chunk + } + const event = JSON.parse(body).payload.hook_event_name + received.push(event) + const reject = event === 'agent_end' && rejectCompletion + if (reject) { + rejectCompletion = false + } + response.writeHead(reject ? 503 : 204) + response.end() +}) + +async function waitForRequests(count) { + const deadline = Date.now() + 5000 + while (received.length < count && Date.now() < deadline) { + await delay(10) + } + assert.equal(received.length, count, 'Hook requests did not arrive before the deadline') +} + +try { + const bundle = join(scratch, 'source.cjs') + await build({ + entryPoints: ['src/main/pi/agent-status-extension-source.ts'], + bundle: true, + platform: 'node', + format: 'cjs', + outfile: bundle + }) + const require = createRequire(import.meta.url) + const { getPiAgentStatusExtensionSource } = require(bundle) + const generated = await transform(getPiAgentStatusExtensionSource('omp'), { + loader: 'ts', + format: 'cjs', + target: 'node24' + }) + const extension = join(scratch, 'extension.cjs') + await writeFile(extension, generated.code) + server.listen(0, '127.0.0.1') + await once(server, 'listening') + Object.assign(process.env, { + ORCA_BACKGROUND_LAUNCH: '1', + ORCA_PANE_KEY: 'completion-proof', + ORCA_TAB_ID: 'proof-tab', + ORCA_AGENT_HOOK_PORT: String(server.address().port), + ORCA_AGENT_HOOK_TOKEN: 'isolated-proof-token', + ORCA_AGENT_HOOK_ENDPOINT: '', + ORCA_PI_STATUS_OWNED: '', + WSL_DISTRO_NAME: '' + }) + const handlers = new Map() + require(extension).default({ on: (event, handler) => handlers.set(event, handler) }) + const emit = async (event) => { + assert.ok(handlers.has(event), `Missing lifecycle handler: ${event}`) + await handlers.get(event)({}, { isIdle: () => false }) + } + await emit('agent_start') + await waitForRequests(1) + await emit('agent_end') + await waitForRequests(3) + assert.deepEqual(received, ['agent_start', 'agent_end', 'agent_end']) + const recovered = [...received] + + for (const boundary of ['agent_start', 'session_switch', 'session_shutdown']) { + received.length = 0 + rejectCompletion = true + await emit('agent_start') + await waitForRequests(1) + await emit('agent_end') + await waitForRequests(2) + await emit(boundary) + await delay(600) + assert.equal( + received.filter((event) => event === 'agent_end').length, + 1, + `An obsolete completion retried after ${boundary}` + ) + } + console.log( + JSON.stringify({ + platform: process.platform, + node: process.version, + recovered, + cancelledAt: ['agent_start', 'session_switch', 'session_shutdown'], + scope: 'Generated OMP extension, real native HTTP 503, synthetic lifecycle callbacks' + }) + ) +} finally { + server.closeAllConnections() + server.close() + await rm(scratch, { recursive: true, force: true }) +}