diff --git a/src/main/agent-hooks/server-pi-normalization.test.ts b/src/main/agent-hooks/server-pi-normalization.test.ts index 6c341f62815..d5e83d69f4f 100644 --- a/src/main/agent-hooks/server-pi-normalization.test.ts +++ b/src/main/agent-hooks/server-pi-normalization.test.ts @@ -1,6 +1,6 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' -import { agentHookServer, _internals } from './server' -import { buildBody } from './server.test-fixtures' +import { AgentHookServer, agentHookServer, _internals } from './server' +import { buildBody, PANE, postHookEvent } from './server.test-fixtures' const { getCohortAtEmitMock, trackMock } = vi.hoisted(() => ({ getCohortAtEmitMock: vi.fn(), @@ -279,3 +279,183 @@ describe('Pi hook normalization', () => { expect(result).toBeNull() }) }) + +const SCOUT = { id: 'run-scout', state: 'working', startedAt: 1_000, agentType: 'scout' } +const REVIEWER = { + id: 'run-reviewer', + state: 'working', + startedAt: 2_000, + agentType: 'reviewer', + description: 'Review the diff' +} + +describe('Pi-family child rows through the hook lane', () => { + let server: AgentHookServer + + beforeEach(async () => { + server = new AgentHookServer() + await server.start({ env: 'production' }) + }) + + afterEach(() => { + server.stop() + }) + + async function post(source: 'pi' | 'omp', payload: Record): Promise { + const response = await postHookEvent(server, buildBody(payload), `/hook/${source}`) + expect(response.status).toBe(204) + } + + it.each(['pi', 'omp'] as const)( + 'publishes the %s extension roster as the row subagents', + async (source) => { + await post(source, { + hook_event_name: 'before_agent_start', + prompt: 'fan out', + subagents: [SCOUT, REVIEWER] + }) + + expect(server.getStatusSnapshot()).toEqual([ + expect.objectContaining({ + state: 'working', + agentType: source, + prompt: 'fan out', + subagents: [SCOUT, REVIEWER] + }) + ]) + } + ) + + it('restates the last row with the new roster when a child ends between lead events', async () => { + await post('pi', { + hook_event_name: 'before_agent_start', + prompt: 'fan out', + subagents: [SCOUT, REVIEWER] + }) + await post('pi', { hook_event_name: 'tool_execution_start', tool_name: 'bash' }) + await post('pi', { hook_event_name: 'subagents_update', subagents: [REVIEWER] }) + + const [row] = server.getStatusSnapshot() + expect(row).toMatchObject({ + state: 'working', + prompt: 'fan out', + toolName: 'bash', + subagents: [REVIEWER] + }) + + // Why: the roster is a full restatement, so an update without one clears every child row. + await post('pi', { hook_event_name: 'subagents_update' }) + expect(server.getStatusSnapshot()[0]).toMatchObject({ state: 'working', prompt: 'fan out' }) + expect(server.getStatusSnapshot()[0]?.subagents).toBeUndefined() + }) + + it('keeps an update from inventing a row or unhiding a resume placeholder', async () => { + await post('pi', { hook_event_name: 'subagents_update', subagents: [SCOUT] }) + expect(server.getStatusSnapshot()).toEqual([]) + + await post('pi', { + hook_event_name: 'session_start', + session_id: 'pi-session-1', + session_file: '/home/dev/.pi/agent/sessions/pi-session-1.jsonl' + }) + await post('pi', { hook_event_name: 'subagents_update', subagents: [SCOUT] }) + const [placeholder] = server.getStatusSnapshot() + expect(placeholder).toMatchObject({ providerSessionOnly: true, state: 'done' }) + expect(placeholder?.subagents).toBeUndefined() + + // Why: the same guard covers Pi's model_select, which shares this path. + await post('pi', { hook_event_name: 'model_select', model: 'anthropic/claude-opus-5' }) + expect(server.getStatusSnapshot()[0]).toMatchObject({ providerSessionOnly: true }) + expect(server.getStatusSnapshot()[0]?.model).toBeUndefined() + }) + + it('records a run its session ended as a session boundary, not a completion', async () => { + await post('pi', { hook_event_name: 'before_agent_start', prompt: 'fan out' }) + await post('pi', { hook_event_name: 'agent_end', session_boundary: true }) + expect(server.getStatusSnapshot()[0]).toMatchObject({ state: 'done', sessionBoundary: true }) + + await post('pi', { hook_event_name: 'before_agent_start', prompt: 'again' }) + await post('pi', { hook_event_name: 'agent_end' }) + expect(server.getStatusSnapshot()[0]?.sessionBoundary).toBeUndefined() + }) + + it('never lets one agent restate another agent row', async () => { + await post('pi', { hook_event_name: 'before_agent_start', prompt: 'pi turn' }) + await post('omp', { hook_event_name: 'subagents_update', subagents: [SCOUT] }) + + const [row] = server.getStatusSnapshot() + expect(row).toMatchObject({ agentType: 'pi', prompt: 'pi turn' }) + expect(row?.subagents).toBeUndefined() + }) +}) + +// Why: publishing a roster for Pi puts its panes on the same child-work guard Claude and Codex +// already sit behind, which changes what Ctrl+C records. Pinned here so the next change to the +// guard — or to the extension, once it reports a main-agent state of its own — has to face it. +describe('a Pi cancel beside a live child row', () => { + let server: AgentHookServer + + beforeEach(async () => { + server = new AgentHookServer() + await server.start({ env: 'production' }) + }) + + afterEach(() => { + server.stop() + }) + + async function startTurn(subagents?: Record[]): Promise { + const response = await postHookEvent( + server, + buildBody({ + hook_event_name: 'before_agent_start', + prompt: 'fan out', + ...(subagents ? { subagents } : {}) + }), + '/hook/pi' + ) + expect(response.status).toBe(204) + } + + function pressCtrlC(): boolean { + const baseline = server.getStatusSnapshotForPane(PANE)[0] + if (!baseline) { + throw new Error('the pane has no row') + } + return server.inferInterrupt({ + paneKey: PANE, + baselineUpdatedAt: baseline.receivedAt, + baselineStateStartedAt: baseline.stateStartedAt, + baselinePrompt: baseline.prompt, + baselineAgentType: 'pi', + intent: 'ctrl-c' + }) + } + + it('still settles a stopped row when the pane has no children', async () => { + await startTurn() + + expect(pressCtrlC()).toBe(true) + expect(server.getStatusSnapshotForPane(PANE)[0]).toMatchObject({ + state: 'done', + interrupted: true, + mainAgent: { state: 'done', outcome: 'cancellation' } + }) + }) + + // Why: without a main-agent state from the extension, Orca cannot tell a cancelled turn from + // Ctrl+C at the idle prompt of a lead that children alone hold open — where it cancels nothing. + // It keeps the live row rather than claiming a cancellation the children contradict. + it('leaves the working row alone while a child still runs', async () => { + await startTurn([{ id: 'run-scout', state: 'working', startedAt: 1_000, agentType: 'scout' }]) + + expect(pressCtrlC()).toBe(false) + const row = server.getStatusSnapshotForPane(PANE)[0] + expect(row).toMatchObject({ + state: 'working', + subagents: [expect.objectContaining({ id: 'run-scout' })] + }) + expect(row?.interrupted).toBeUndefined() + expect(row?.mainAgent).toBeUndefined() + }) +}) diff --git a/src/main/pi/agent-status-completion-delivery.test.ts b/src/main/pi/agent-status-completion-delivery.test.ts index 4c99563b6ed..0354235c938 100644 --- a/src/main/pi/agent-status-completion-delivery.test.ts +++ b/src/main/pi/agent-status-completion-delivery.test.ts @@ -46,8 +46,8 @@ describe('OMP completion delivery', () => { 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) + await vi.advanceTimersByTimeAsync(60_000) + expect(harness.spawnMock).toHaveBeenCalledTimes(4) expect(vi.getTimerCount()).toBe(0) }) @@ -87,7 +87,7 @@ describe('OMP completion delivery', () => { expect(vi.getTimerCount()).toBe(0) }) - it.each(['before_agent_start', 'agent_start', 'session_shutdown', 'session_switch'])( + it.each(['before_agent_start', 'agent_start'])( 'retires a failed completion at %s', async (boundary) => { const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) @@ -122,6 +122,59 @@ describe('OMP completion delivery', () => { ]) }) + it.each(['scheduled retry', 'in-flight failure', 'in-flight success'] as const)( + 'preserves completion delivery across a session boundary after %s', + async (delivery) => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + let id = 'A' + const ctx = { + isIdle: () => true, + sessionManager: { getSessionId: () => id, getSessionFile: () => `/sessions/${id}.jsonl` } + } + await harness.callHook('agent_start', {}, ctx) + await vi.advanceTimersByTimeAsync(0) + let acknowledge: ((value: { ok: boolean }) => void) | undefined + let fail: ((error: Error) => void) | undefined + if (delivery === 'scheduled retry') { + harness.fetchMock.mockRejectedValueOnce(new Error('offline')) + } else { + harness.fetchMock.mockImplementationOnce( + () => + new Promise((resolve, reject) => { + acknowledge = resolve + fail = reject + }) + ) + } + await harness.callHook('agent_end', {}, ctx) + await vi.advanceTimersByTimeAsync(0) + id = 'B' + await harness.callHook('session_switch', { reason: 'new' }, ctx) + await harness.callHook('agent_start', {}, ctx) + if (delivery === 'in-flight failure') { + fail?.(new Error('late failure')) + } else { + acknowledge?.({ ok: true }) + } + await vi.advanceTimersByTimeAsync(5_000) + + const payloads = events(harness.fetchMock) + const completions = payloads.filter( + (event) => + event && + typeof event === 'object' && + 'hook_event_name' in event && + event.hook_event_name === 'agent_end' + ) + expect(completions).toHaveLength(delivery === 'in-flight success' ? 1 : 2) + for (const completion of completions) { + expect(completion).toMatchObject({ session_id: 'A' }) + } + expect(payloads.at(-1)).toMatchObject({ hook_event_name: 'agent_start', session_id: 'B' }) + expect(vi.getTimerCount()).toBe(0) + } + ) + it('retries a timed out completion without blocking agent handlers', async () => { const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) harness.fetchMock.mockImplementationOnce(() => new Promise(() => {})) diff --git a/src/main/pi/agent-status-extension-async-subagents.test.ts b/src/main/pi/agent-status-extension-async-subagents.test.ts index 89d5eebcf72..00e375060ea 100644 --- a/src/main/pi/agent-status-extension-async-subagents.test.ts +++ b/src/main/pi/agent-status-extension-async-subagents.test.ts @@ -1,56 +1,20 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { AGENT_STATUS_MAX_SUBAGENTS } from '../../shared/agent-status-types' +import { createAgentStatusExtensionHarness } from './agent-status-extension-test-harness' import { - AGENT_STATUS_EXTENSION_SELF_PID, - createAgentStatusExtensionHarness, - type AgentStatusExtensionHarness -} from './agent-status-extension-test-harness' - -// Event shapes and orderings mirror traces recorded from pi-subagents 0.71.0. -const WORKFLOW = 'workflow-1' -const idle = { isIdle: () => true } - -function postedHookNames(harness: AgentStatusExtensionHarness): string[] { - return harness.fetchMock.mock.calls.map((call) => { - const body: { payload?: { hook_event_name?: unknown } } = JSON.parse(String(call[1]?.body)) - return String(body.payload?.hook_event_name) - }) -} - -function agentEndCount(harness: AgentStatusExtensionHarness): number { - return postedHookNames(harness).filter((name) => name === 'agent_end').length -} - -function startWorkflow(harness: AgentStatusExtensionHarness): void { - harness.emitPiEvent('subagent:async-started', { - id: WORKFLOW, - mode: 'workflow', - pid: AGENT_STATUS_EXTENSION_SELF_PID - }) -} - -function startChild(harness: AgentStatusExtensionHarness, id: string, parent = WORKFLOW): void { - harness.emitPiEvent('subagent:async-started', { - id, - mode: 'single', - pid: 4000, - parentWorkflowRunId: parent - }) -} - -function exitRunner(harness: AgentStatusExtensionHarness, runId: string): void { - harness.emitPiEvent('subagent:process-terminal', { runId, state: 'observed' }) -} - -function complete(harness: AgentStatusExtensionHarness, id: string): void { - harness.emitPiEvent('subagent:async-complete', { id, runId: id, state: 'complete' }) -} - -async function endTurn(harness: AgentStatusExtensionHarness): Promise { - await harness.callHook('agent_end', {}, idle) - await harness.callHook('agent_settled', undefined, idle) - await vi.advanceTimersByTimeAsync(0) -} + agentEndCount, + childIds, + complete, + endTurn, + exitRunner, + postedHookNames, + posts, + startAsync, + startChild, + startWorkflow, + WORKFLOW +} from './agent-status-subagent-event-fixtures' describe('Pi async subagent roster', () => { beforeEach(() => { @@ -190,16 +154,242 @@ describe('Pi async subagent roster', () => { expect(agentEndCount(harness)).toBe(1) }) +}) - it('keeps one runner-exit subscription and the roster across reloads', async () => { +// Roster shapes also mirror OMP 18.3.2's task:subagent:lifecycle. +describe('Pi child rows', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + afterEach(() => { + vi.useRealTimers() + }) + + it('posts each running pi-subagents child with its agent name, never its redacted task', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'pi' }) + startAsync(harness, 'run-scout', 'scout') + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness)).toEqual([ + { + hook_event_name: 'agent_start', + subagents: [ + { id: 'run-scout', state: 'working', startedAt: expect.any(Number), agentType: 'scout' } + ] + } + ]) + }) + + it('posts an OMP task child with its description', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + await harness.callHook('session_start') + harness.emitPiEvent('task:subagent:lifecycle', { + id: '0-explore', + agent: 'explore', + description: 'Map the auth module', + detached: true, + status: 'started', + index: 0 + }) + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness).at(-1)?.subagents).toEqual([ + { + id: '0-explore', + state: 'working', + startedAt: expect.any(Number), + agentType: 'explore', + description: 'Map the auth module' + } + ]) + }) + + it('restates the roster when a child ends mid-turn', async () => { const harness = createAgentStatusExtensionHarness({ kind: 'pi' }) await harness.callHook('agent_start') - startChild(harness, 'child-a', 'tool-call-1') - harness.reload() - expect(harness.piEventListenerCount('subagent:process-terminal')).toBe(1) + startAsync(harness, 'run-a', 'scout') + startAsync(harness, 'run-b', 'reviewer') + await vi.advanceTimersByTimeAsync(0) + complete(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) - exitRunner(harness, 'child-a') + const last = posts(harness).at(-1) + expect(last?.hook_event_name).toBe('subagents_update') + expect(childIds(last)).toEqual(['run-b']) + }) + + it('restates the roster while the lead waits, then completes once the last child ends', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'pi' }) + await harness.callHook('agent_start') + startAsync(harness, 'run-a', 'scout') + startAsync(harness, 'run-b', 'reviewer') await endTurn(harness) - expect(agentEndCount(harness)).toBe(1) + complete(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'subagents_update' }) + expect(childIds(posts(harness).at(-1))).toEqual(['run-b']) + + complete(harness, 'run-b') + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toEqual({ hook_event_name: 'agent_end' }) + expect( + posts(harness).filter((post) => post.hook_event_name === 'subagents_update') + ).toHaveLength(1) + }) + + it('drops a child as soon as its runner exits, once per exit', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'pi' }) + await harness.callHook('agent_start') + startAsync(harness, 'run-a', 'scout') + startAsync(harness, 'run-b', 'reviewer') + await vi.advanceTimersByTimeAsync(0) + exitRunner(harness, 'run-a') + exitRunner(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + + const updates = posts(harness).filter((post) => post.hook_event_name === 'subagents_update') + expect(updates).toHaveLength(1) + expect(childIds(updates[0])).toEqual(['run-b']) + }) + + it('lets a queued post carry a roster change instead of adding an update behind it', async () => { + const releases: (() => void)[] = [] + const harness = createAgentStatusExtensionHarness({ + kind: 'pi', + fetchImpl: () => + new Promise((resolve) => { + releases.push(() => resolve({ ok: true })) + }) + }) + await harness.callHook('agent_start') + startAsync(harness, 'run-a', 'scout') + startAsync(harness, 'run-b', 'reviewer') + complete(harness, 'run-a') + + releases.shift()?.() + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).map((post) => post.hook_event_name)).toEqual([ + 'agent_start', + 'agent_start' + ]) + expect(childIds(posts(harness)[1])).toEqual(['run-b']) + }) + + // Why: the generated cap is interpolated from the host's, so a drift would show up here + // as an over-long roster the host would silently truncate on arrival. + it('caps the posted roster at the host limit', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'pi' }) + await harness.callHook('agent_start') + for (let index = 0; index < AGENT_STATUS_MAX_SUBAGENTS + 4; index++) { + startAsync(harness, `run-${index}`, 'scout') + } + await vi.advanceTimersByTimeAsync(0) + + expect(childIds(posts(harness).at(-1))).toHaveLength(AGENT_STATUS_MAX_SUBAGENTS) + expect(childIds(posts(harness).at(-1))?.at(-1)).toBe(`run-${AGENT_STATUS_MAX_SUBAGENTS - 1}`) + }) + + it('posts nothing for the end of a run it is not tracking', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'pi' }) + await harness.callHook('agent_start') + complete(harness, 'unknown-run') + await vi.advanceTimersByTimeAsync(0) + + expect(postedHookNames(harness)).toEqual(['agent_start']) + }) + + it('labels a reused child id from its latest start, and keeps a running child’s first label', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + const start = (agent: string) => + harness.emitPiEvent('task:subagent:lifecycle', { id: '0-task', agent, status: 'started' }) + const labels = () => + posts(harness) + .at(-1) + ?.subagents?.map((child) => child.agentType) + await harness.callHook('agent_start') + start('explore') + start('review') + await vi.advanceTimersByTimeAsync(0) + expect(labels()).toEqual(['explore']) + + harness.emitPiEvent('task:subagent:lifecycle', { id: '0-task', status: 'completed' }) + start('review') + await vi.advanceTimersByTimeAsync(0) + expect(labels()).toEqual(['review']) + + await harness.callHook('session_switch', { reason: 'new' }, {}) + start('plan') + await vi.advanceTimersByTimeAsync(0) + expect(labels()).toEqual(['plan']) + }) + + it('posts no subagents field for a pane without children', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'pi' }) + await harness.callHook('before_agent_start', { prompt: 'plain turn' }) + await harness.callHook('tool_execution_start', { toolName: 'bash', args: {} }) + await endTurn(harness) + + expect(posts(harness).some((post) => 'subagents' in post)).toBe(false) + }) + + it('posts one roster update when a runner exit precedes its completion', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'pi' }) + await harness.callHook('agent_start') + startAsync(harness, 'run-a', 'scout') + startAsync(harness, 'run-b', 'reviewer') + await vi.advanceTimersByTimeAsync(0) + exitRunner(harness, 'run-a') + await vi.advanceTimersByTimeAsync(150) + complete(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + + const updates = posts(harness).filter((post) => post.hook_event_name === 'subagents_update') + expect(updates.map(childIds)).toEqual([['run-b']]) + }) + + it('leaves a scheduled OMP retry to carry a roster change', async () => { + let attempts = 0 + const harness = createAgentStatusExtensionHarness({ + kind: 'omp', + fetchImpl: async () => { + attempts += 1 + if (attempts === 2) { + throw new Error('Orca restarting') + } + return { ok: true } + } + }) + await harness.callHook('agent_start') + await vi.advanceTimersByTimeAsync(0) + harness.emitPiEvent('task:subagent:lifecycle', { id: 'c1', agent: 'task', status: 'started' }) + await vi.advanceTimersByTimeAsync(0) + harness.emitPiEvent('task:subagent:lifecycle', { id: 'c1', status: 'completed' }) + await vi.advanceTimersByTimeAsync(250) + + expect(posts(harness).map((post) => post.hook_event_name)).toEqual([ + 'agent_start', + 'agent_start', + 'agent_start' + ]) + expect(childIds(posts(harness)[1])).toEqual(['c1']) + expect(posts(harness)[2]?.subagents).toBeUndefined() + }) + + it('keeps a workflow run out of the rows while it still holds the pane', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'pi' }) + await harness.callHook('agent_start') + startWorkflow(harness) + startAsync(harness, 'run-a', 'scout') + await endTurn(harness) + expect(childIds(posts(harness).at(-1))).toEqual(['run-a']) + + complete(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toEqual({ hook_event_name: 'subagents_update' }) + + complete(harness, WORKFLOW) + await vi.advanceTimersByTimeAsync(0) + expect(postedHookNames(harness).slice(-1)).toEqual(['agent_end']) + expect(posts(harness).some((post) => childIds(post)?.includes(WORKFLOW))).toBe(false) }) }) diff --git a/src/main/pi/agent-status-extension-omp-lifecycle.test.ts b/src/main/pi/agent-status-extension-omp-lifecycle.test.ts index 03408a3ba6b..fe410842ba1 100644 --- a/src/main/pi/agent-status-extension-omp-lifecycle.test.ts +++ b/src/main/pi/agent-status-extension-omp-lifecycle.test.ts @@ -44,9 +44,16 @@ describe('OMP agent_end contract', () => { ) }) + it('subscribes once when the factory runs again on the same bus', () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + harness.reload() + expect(harness.piEventListenerCount('task:subagent:lifecycle')).toBe(1) + expect(harness.piEventListenerCount('subagent:process-terminal')).toBe(1) + }) + it('keeps one lifecycle subscription across extension reloads', async () => { const harness = createAgentStatusExtensionHarness({ kind: 'pi' }) - harness.reload() + await harness.reloadPi() expect(harness.piEventListenerCount('task:subagent:lifecycle')).toBe(1) expect(harness.piEventListenerCount('subagent:async-started')).toBe(1) expect(harness.piEventListenerCount('subagent:async-complete')).toBe(1) diff --git a/src/main/pi/agent-status-extension-session-change.test.ts b/src/main/pi/agent-status-extension-session-change.test.ts new file mode 100644 index 00000000000..1122349ae10 --- /dev/null +++ b/src/main/pi/agent-status-extension-session-change.test.ts @@ -0,0 +1,768 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' + +import { + createAgentStatusExtensionHarness, + type AgentStatusExtensionHarness, + type HookContext +} from './agent-status-extension-test-harness' +import { + agentEndCount, + childIds, + complete, + endTurn, + exitRunner, + idle, + postedHookNames, + posts, + startAsync, + startChild, + startWorkflow, + WORKFLOW +} from './agent-status-subagent-event-fixtures' + +// Orderings mirror runs recorded from Pi 0.87.1 with pi-subagents 0.71.0: a session change or +// /reload shuts the old registration down and runs the factory again on a fresh `pi.events`. +function session(id: string) { + return { + isIdle: () => true, + sessionManager: { getSessionId: () => id, getSessionFile: () => `/sessions/${id}.jsonl` } + } +} + +function createPi(fetchImpl?: () => Promise): AgentStatusExtensionHarness { + return createAgentStatusExtensionHarness({ kind: 'pi', existsSync: () => true, fetchImpl }) +} + +async function holdRunOpen(harness: AgentStatusExtensionHarness, sessionId = 'A'): Promise { + await harness.callHook('session_start', { reason: 'startup' }, session(sessionId)) + await harness.callHook('agent_start', {}, session(sessionId)) + startAsync(harness, 'run-a', 'scout') + await endTurn(harness) +} + +describe('Pi session changes', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + afterEach(() => { + vi.useRealTimers() + }) + + it.each([ + ['new', undefined], + ['resume', '/sessions/B.jsonl'], + ['fork', '/sessions/B.jsonl'] + ] as const)( + 'ends the run its children held open under the old session on %s, before the next one starts', + async (reason, target) => { + const harness = createPi() + await holdRunOpen(harness) + expect(agentEndCount(harness)).toBe(0) + + await harness.replacePiSession(reason, target) + await harness.callHook('session_start', { reason }, session('B')) + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness).slice(-2)).toEqual([ + { + hook_event_name: 'agent_end', + session_id: 'A', + session_file: '/sessions/A.jsonl', + session_boundary: true + }, + { hook_event_name: 'session_start', session_id: 'B', session_file: '/sessions/B.jsonl' } + ]) + } + ) + + it('ends the run for a Pi too old to say why it shut the session down', async () => { + const harness = createPi() + await holdRunOpen(harness) + await harness.callHook('session_shutdown') + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end', session_id: 'A' }) + }) + + it('posts nothing on quit, even with children still running', async () => { + const harness = createPi() + await holdRunOpen(harness) + const sent = posts(harness).length + await harness.callHook('session_shutdown', { reason: 'quit' }) + await vi.advanceTimersByTimeAsync(5_000) + + expect(posts(harness)).toHaveLength(sent) + }) + + it('re-opens the run for a child of the next session after a turn the old one cut off', async () => { + const harness = createPi() + await harness.callHook('session_start', { reason: 'startup' }, session('A')) + await harness.callHook('agent_start', {}, session('A')) + await harness.replacePiSession('new') + await harness.callHook('session_start', { reason: 'new' }, session('B')) + startAsync(harness, 'run-b', 'scout') + complete(harness, 'run-b') + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end', session_id: 'B' }) + expect(posts(harness).at(-1)).not.toHaveProperty('session_boundary') + }) + + it('keeps a completion that was still waiting to be sent when the session changed', async () => { + const releases: (() => void)[] = [] + const harness = createPi( + () => + new Promise((resolve) => { + releases.push(() => resolve({ ok: true })) + }) + ) + await harness.callHook('session_start', { reason: 'startup' }, session('A')) + await harness.callHook('agent_start', {}, session('A')) + // Pi aborts the turn before it shuts the session down; that completion queues behind the post in flight. + await endTurn(harness) + await harness.replacePiSession('new') + await harness.callHook('session_start', { reason: 'new' }, session('B')) + + while (releases.length > 0) { + releases.shift()?.() + await vi.advanceTimersByTimeAsync(0) + } + expect(posts(harness).slice(-2)).toMatchObject([ + { hook_event_name: 'agent_end', session_id: 'A' }, + { hook_event_name: 'session_start', session_id: 'B' } + ]) + // A turn that really finished is a completion, not a session boundary. + expect(posts(harness).at(-2)).not.toHaveProperty('session_boundary') + }) + + it('gives up on a completion Orca keeps refusing, then sends what came after', async () => { + let attempts = 0 + const harness = createPi(async () => { + attempts += 1 + if (attempts >= 3 && attempts <= 6) { + throw new Error('Orca down') + } + return { ok: true } + }) + await holdRunOpen(harness) + await harness.replacePiSession('new') + await harness.callHook('session_start', { reason: 'new' }, session('B')) + await vi.advanceTimersByTimeAsync(5_000) + + expect(agentEndCount(harness)).toBe(4) + expect(posts(harness).at(-1)).toMatchObject({ + hook_event_name: 'session_start', + session_id: 'B' + }) + }) + + it('posts nothing for an old session whose run already ended', async () => { + const harness = createPi() + await harness.callHook('agent_start', {}, session('A')) + await endTurn(harness) + expect(agentEndCount(harness)).toBe(1) + + await harness.replacePiSession('fork', '/sessions/B.jsonl') + await vi.advanceTimersByTimeAsync(0) + + expect(agentEndCount(harness)).toBe(1) + }) + + it('delivers the old session’s completion ahead of a new session posted behind an in-flight delivery', async () => { + const releases: (() => void)[] = [] + const harness = createPi( + () => + new Promise((resolve) => { + releases.push(() => resolve({ ok: true })) + }) + ) + await holdRunOpen(harness) + await harness.replacePiSession('new') + await harness.callHook('session_start', { reason: 'new' }, session('B')) + // A child of the new session must not ride on the old session's last post. + startAsync(harness, 'run-b', 'scout') + + while (releases.length > 0) { + releases.shift()?.() + await vi.advanceTimersByTimeAsync(0) + } + const completion = posts(harness).find((post) => post.hook_event_name === 'agent_end') + expect(completion).toMatchObject({ session_id: 'A' }) + expect(completion?.subagents).toBeUndefined() + expect(postedHookNames(harness).slice(-2)).toEqual(['agent_end', 'agent_start']) + expect(posts(harness).at(-1)).toMatchObject({ session_id: 'B' }) + expect(childIds(posts(harness).at(-1))).toEqual(['run-b']) + }) + + it('retries the old session’s completion before sending anything newer', async () => { + let attempts = 0 + const harness = createPi(async () => { + attempts += 1 + // The close-out follows session_start and agent_start; fail its first delivery. + if (attempts === 3) { + throw new Error('Orca restarting') + } + return { ok: true } + }) + await holdRunOpen(harness) + await harness.replacePiSession('new') + await harness.callHook('session_start', { reason: 'new' }, session('B')) + await vi.advanceTimersByTimeAsync(0) + // The new session's post waits out the backoff behind the old session's completion. + expect(postedHookNames(harness).at(-1)).toBe('agent_end') + expect(agentEndCount(harness)).toBe(1) + await vi.advanceTimersByTimeAsync(250) + + expect(posts(harness).slice(-3)).toMatchObject([ + { hook_event_name: 'agent_end', session_id: 'A' }, + { hook_event_name: 'agent_end', session_id: 'A' }, + { hook_event_name: 'session_start', session_id: 'B' } + ]) + }) + + it('stops listening once the old session is closed out', async () => { + const harness = createPi() + await holdRunOpen(harness) + await harness.callHook('session_shutdown', { reason: 'new' }) + await vi.advanceTimersByTimeAsync(0) + const sent = posts(harness).length + + // Another extension's shutdown handler can still be running while pi-subagents emits here. + startAsync(harness, 'run-late', 'scout') + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness)).toHaveLength(sent) + + await harness.replacePiSession('new') + await harness.callHook('agent_start', {}, session('B')) + await endTurn(harness) + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end' }) + expect( + posts(harness) + .slice(sent) + .some((post) => post.subagents) + ).toBe(false) + }) + + it('leaves no timer of the old session running', async () => { + const harness = createPi() + await holdRunOpen(harness) + exitRunner(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + expect(vi.getTimerCount()).toBe(1) + + await harness.replacePiSession('new') + await vi.advanceTimersByTimeAsync(0) + expect(vi.getTimerCount()).toBe(0) + }) + + it('does not bring back a child whose runner had already exited', async () => { + const harness = createPi() + await holdRunOpen(harness) + exitRunner(harness, 'run-a') + await harness.replacePiSession('new') + await harness.callHook('session_start', { reason: 'new' }, session('B')) + await harness.replacePiSession('resume', '/sessions/A.jsonl') + await harness.callHook('session_start', { reason: 'resume' }, session('A')) + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness).at(-1)).toMatchObject({ + hook_event_name: 'session_start', + session_id: 'A' + }) + }) + + it('brings a session’s children back when it is resumed, and ends the run when they finish', async () => { + const harness = createPi() + await holdRunOpen(harness) + await harness.replacePiSession('new') + await harness.callHook('session_start', { reason: 'new' }, session('B')) + await harness.replacePiSession('resume', '/sessions/A.jsonl') + await harness.callHook('session_start', { reason: 'resume' }, session('A')) + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_start', session_id: 'A' }) + expect(posts(harness).at(-1)?.subagents).toEqual([ + { id: 'run-a', state: 'working', startedAt: expect.any(Number), agentType: 'scout' } + ]) + + // pi-subagents reports the run's completion to the resumed session, finished or not. + complete(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toEqual({ + hook_event_name: 'agent_end', + session_id: 'A', + session_file: '/sessions/A.jsonl' + }) + }) + + it('restores a session’s children once, not on every later resume', async () => { + const harness = createPi() + await holdRunOpen(harness) + await harness.replacePiSession('new') + await harness.callHook('session_start', { reason: 'new' }, session('B')) + await harness.replacePiSession('resume', '/sessions/A.jsonl') + await harness.callHook('session_start', { reason: 'resume' }, session('A')) + complete(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + await harness.replacePiSession('new') + await harness.callHook('session_start', { reason: 'new' }, session('B')) + await harness.replacePiSession('resume', '/sessions/A.jsonl') + await harness.callHook('session_start', { reason: 'resume' }, session('A')) + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness).at(-1)).toMatchObject({ + hook_event_name: 'session_start', + session_id: 'A' + }) + expect(posts(harness).at(-1)?.subagents).toBeUndefined() + }) + + it('keeps the children across a resume into the same session file', async () => { + const harness = createPi() + await holdRunOpen(harness) + await harness.replacePiSession('resume', '/sessions/A.jsonl') + await harness.callHook('session_start', { reason: 'resume' }, session('A')) + await vi.advanceTimersByTimeAsync(0) + expect(agentEndCount(harness)).toBe(0) + + complete(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + expect(agentEndCount(harness)).toBe(1) + }) +}) + +describe('Pi /reload', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + afterEach(() => { + vi.useRealTimers() + }) + + it('keeps the children and the hold across a reload', async () => { + const harness = createPi() + await holdRunOpen(harness) + await harness.reloadPi() + await harness.callHook('session_start', { reason: 'reload' }, session('A')) + expect(harness.piEventListenerCount('subagent:async-complete')).toBe(1) + expect(harness.piEventListenerCount('subagent:process-terminal')).toBe(1) + + // A turn that ends while the pre-reload child still runs must not report done. + await harness.callHook('agent_start', {}, session('A')) + await vi.advanceTimersByTimeAsync(0) + expect(childIds(posts(harness).at(-1))).toEqual(['run-a']) + await endTurn(harness) + expect(agentEndCount(harness)).toBe(0) + + complete(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + expect(postedHookNames(harness).at(-1)).toBe('agent_end') + expect(agentEndCount(harness)).toBe(1) + }) + + it('keeps the turn counters across a reload, so a child starting afterwards is not read as late', async () => { + const harness = createPi() + await harness.callHook('agent_start', {}, session('A')) + await harness.reloadPi() + startAsync(harness, 'run-a', 'scout') + complete(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + expect(agentEndCount(harness)).toBe(0) + + await endTurn(harness) + expect(agentEndCount(harness)).toBe(1) + }) + + it('releases a workflow’s pre-reload children when the workflow ends', async () => { + const harness = createPi() + await harness.callHook('agent_start', {}, session('A')) + startWorkflow(harness) + startChild(harness, 'child-a') + await endTurn(harness) + // Pi drops the old registration's runner-exit events, so child-a never reports its end. + await harness.reloadPi() + complete(harness, WORKFLOW) + await vi.advanceTimersByTimeAsync(0) + + expect(agentEndCount(harness)).toBe(1) + expect(posts(harness).at(-1)?.subagents).toBeUndefined() + }) + + it('drops the rows of pre-reload children released while the turn is still running', async () => { + const harness = createPi() + await harness.callHook('agent_start', {}, session('A')) + startWorkflow(harness) + startChild(harness, 'child-a') + await harness.reloadPi() + complete(harness, WORKFLOW) + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness).at(-1)).toEqual({ hook_event_name: 'subagents_update' }) + }) + + it.each(['reload', 'resume'] as const)( + 'settles an exited runner after same-session %s without another turn', + async (reason) => { + const harness = createPi() + await harness.callHook('session_start', {}, session('A')) + await harness.callHook('agent_start', {}, session('A')) + startChild(harness, 'child-a', 'tool-call-1') + await endTurn(harness) + exitRunner(harness, 'child-a') + await vi.advanceTimersByTimeAsync(1_000) + await (reason === 'reload' + ? harness.reloadPi() + : harness.replacePiSession('resume', '/sessions/A.jsonl')) + await harness.callHook('session_start', { reason }, session('A')) + await vi.advanceTimersByTimeAsync(999) + expect(agentEndCount(harness)).toBe(0) + await vi.advanceTimersByTimeAsync(1) + + expect(agentEndCount(harness)).toBe(1) + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end', session_id: 'A' }) + expect(posts(harness).at(-1)?.subagents).toBeUndefined() + expect(vi.getTimerCount()).toBe(0) + } + ) + + it.each(['new', 'quit'] as const)( + 'clears runner grace when the session ends on %s', + async (reason) => { + const harness = createPi() + await holdRunOpen(harness) + exitRunner(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + expect(vi.getTimerCount()).toBe(1) + await harness.callHook('session_shutdown', { reason }) + await vi.advanceTimersByTimeAsync(0) + const completionCount = agentEndCount(harness) + await vi.advanceTimersByTimeAsync(5_000) + + expect(agentEndCount(harness)).toBe(completionCount) + expect(vi.getTimerCount()).toBe(0) + } + ) +}) + +describe('children that start outside a turn', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + afterEach(() => { + vi.useRealTimers() + }) + + it('re-opens a finished OMP run for a late child and ends it when the child does', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + await harness.callHook('agent_start') + await harness.callHook('agent_end', {}) + await vi.advanceTimersByTimeAsync(0) + harness.emitPiEvent('task:subagent:lifecycle', { id: 'w1', agent: 'task', status: 'started' }) + await vi.advanceTimersByTimeAsync(0) + harness.emitPiEvent('task:subagent:lifecycle', { id: 'w1', status: 'completed' }) + await vi.advanceTimersByTimeAsync(0) + + expect(postedHookNames(harness)).toEqual([ + 'agent_start', + 'agent_end', + 'agent_start', + 'agent_end' + ]) + }) + + it('ends a run for a child that started before any turn in this session', async () => { + const harness = createPi() + startAsync(harness, 'run-a', 'scout') + await vi.advanceTimersByTimeAsync(0) + complete(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + + expect(postedHookNames(harness)).toEqual(['agent_start', 'agent_end']) + }) + + it('keeps working when a dialog closes while a child runs, whether or not the turn has ended', async () => { + const harness = createPi() + await harness.callHook('agent_start') + startAsync(harness, 'run-a', 'scout') + await harness.callHook('ui_prompt_start', {}) + await harness.callHook('ui_prompt_end', {}, idle) + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toMatchObject({ + hook_event_name: 'ui_prompt_end', + is_idle: false + }) + + await endTurn(harness) + await harness.callHook('ui_prompt_start', {}) + await harness.callHook('ui_prompt_end', {}, idle) + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toMatchObject({ + hook_event_name: 'ui_prompt_end', + is_idle: false + }) + expect(childIds(posts(harness).at(-1))).toEqual(['run-a']) + + complete(harness, 'run-a') + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toEqual({ hook_event_name: 'agent_end' }) + }) +}) + +describe('OMP session switches', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + afterEach(() => { + vi.useRealTimers() + }) + + // OMP keeps one registration and one SessionManager across a switch. + function ompSession() { + let id = 'A' + const context = { + sessionManager: { getSessionId: () => id, getSessionFile: () => `/sessions/${id}.jsonl` } + } + return { context, switchTo: (next: string) => (id = next) } + } + + async function holdOmpRunOpen(harness: AgentStatusExtensionHarness, context: HookContext) { + await harness.callHook('agent_start', {}, context) + harness.emitPiEvent('task:subagent:lifecycle', { id: 'c1', agent: 'task', status: 'started' }) + await harness.callHook('agent_end', {}, context) + await vi.advanceTimersByTimeAsync(0) + } + + it('ends the held run under the session that ran it and forgets its children', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + const { context, switchTo } = ompSession() + await holdOmpRunOpen(harness, context) + expect(agentEndCount(harness)).toBe(0) + + switchTo('B') + await harness.callHook( + 'session_switch', + { reason: 'new', previousSessionFile: '/sessions/A.jsonl' }, + context + ) + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end', session_id: 'A' }) + expect(posts(harness).at(-1)?.subagents).toBeUndefined() + + await harness.callHook('agent_start', {}, context) + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_start', session_id: 'B' }) + expect(posts(harness).at(-1)?.subagents).toBeUndefined() + }) + + it('ends a turn the switch cut off, under the session that ran it', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + const { context, switchTo } = ompSession() + await harness.callHook('agent_start', {}, context) + switchTo('B') + await harness.callHook( + 'session_switch', + { reason: 'new', previousSessionFile: '/sessions/A.jsonl' }, + context + ) + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end', session_id: 'A' }) + }) + + it('keeps the children when OMP reloads the same session', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + const { context } = ompSession() + await holdOmpRunOpen(harness, context) + await harness.callHook( + 'session_switch', + { reason: 'resume', previousSessionFile: '/sessions/A.jsonl' }, + context + ) + await vi.advanceTimersByTimeAsync(0) + expect(agentEndCount(harness)).toBe(0) + + harness.emitPiEvent('task:subagent:lifecycle', { id: 'c1', status: 'completed' }) + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end', session_id: 'A' }) + }) + + it.each(['fork', 'resume'] as const)( + 'keeps the children on a %s into another session, where OMP leaves them running', + async (reason) => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + const { context, switchTo } = ompSession() + await holdOmpRunOpen(harness, context) + switchTo('B') + await harness.callHook( + 'session_switch', + { reason, previousSessionFile: '/sessions/A.jsonl' }, + context + ) + await vi.advanceTimersByTimeAsync(0) + expect(agentEndCount(harness)).toBe(0) + + harness.emitPiEvent('task:subagent:lifecycle', { id: 'c1', status: 'completed' }) + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end' }) + } + ) + + it('ends the held run when OMP branches the session', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + const { context, switchTo } = ompSession() + await holdOmpRunOpen(harness, context) + switchTo('B') + await harness.callHook('session_branch', { previousSessionFile: '/sessions/A.jsonl' }, context) + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end', session_id: 'A' }) + expect(posts(harness).at(-1)?.subagents).toBeUndefined() + }) + + it('keeps a completion that was still waiting to be sent when OMP starts a new session', async () => { + const releases: (() => void)[] = [] + const harness = createAgentStatusExtensionHarness({ + kind: 'omp', + fetchImpl: () => + new Promise((resolve) => { + releases.push(() => resolve({ ok: true })) + }) + }) + const { context, switchTo } = ompSession() + await harness.callHook('agent_start', {}, context) + await harness.callHook('agent_end', {}, context) + switchTo('B') + await harness.callHook( + 'session_switch', + { reason: 'new', previousSessionFile: '/sessions/A.jsonl' }, + context + ) + await harness.callHook('agent_start', {}, context) + + while (releases.length > 0) { + releases.shift()?.() + await vi.advanceTimersByTimeAsync(0) + } + expect(posts(harness).slice(-2)).toMatchObject([ + { hook_event_name: 'agent_end', session_id: 'A' }, + { hook_event_name: 'agent_start', session_id: 'B' } + ]) + }) + + it.each(['omp', 'pi'] as const)( + 'takes over a roster and its subscriptions that an older build left on the %s bus', + async (kind) => { + const harness = createAgentStatusExtensionHarness({ + kind, + sharedEventBus: true, + seedEventBus: (bus) => { + Object.assign(bus, { + __orcaPiSubagents: { + active: new Set(['old-child']), + waiting: false, + listener: () => {}, + runnerExitListener: () => {} + } + }) + } + }) + expect(harness.piEventListenerCount('task:subagent:lifecycle')).toBe(0) + expect(harness.piEventListenerCount('subagent:process-terminal')).toBe(0) + + await harness.callHook('agent_start') + await vi.advanceTimersByTimeAsync(0) + expect(childIds(posts(harness).at(-1))).toEqual(['old-child']) + } + ) + + it('posts nothing for a resume before any turn has run', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + const { context } = ompSession() + await harness.callHook( + 'session_switch', + { reason: 'resume', previousSessionFile: '/sessions/A.jsonl' }, + context + ) + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness)).toEqual([]) + }) + + it('keeps the next session’s children off the old session’s last post', async () => { + const releases: (() => void)[] = [] + const harness = createAgentStatusExtensionHarness({ + kind: 'omp', + fetchImpl: () => + new Promise((resolve) => { + releases.push(() => resolve({ ok: true })) + }) + }) + const { context, switchTo } = ompSession() + await holdOmpRunOpen(harness, context) + switchTo('B') + await harness.callHook( + 'session_switch', + { reason: 'new', previousSessionFile: '/sessions/A.jsonl' }, + context + ) + harness.emitPiEvent('task:subagent:lifecycle', { id: 'c2', agent: 'task', status: 'started' }) + + while (releases.length > 0) { + releases.shift()?.() + await vi.advanceTimersByTimeAsync(0) + } + const completion = posts(harness).find((post) => post.hook_event_name === 'agent_end') + expect(completion).toMatchObject({ session_id: 'A' }) + expect(completion?.subagents).toBeUndefined() + expect(childIds(posts(harness).at(-1))).toEqual(['c2']) + }) + + it('ends a turn that a reload of the same session cut off', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + const { context } = ompSession() + await harness.callHook('agent_start', {}, context) + await harness.callHook( + 'session_switch', + { reason: 'resume', previousSessionFile: '/sessions/A.jsonl' }, + context + ) + await vi.advanceTimersByTimeAsync(0) + + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end', session_id: 'A' }) + + // The turn is over, so a child that starts now re-opens the run. + harness.emitPiEvent('task:subagent:lifecycle', { id: 'w1', agent: 'task', status: 'started' }) + harness.emitPiEvent('task:subagent:lifecycle', { id: 'w1', status: 'completed' }) + await vi.advanceTimersByTimeAsync(0) + expect(postedHookNames(harness)).toEqual([ + 'agent_start', + 'agent_end', + 'agent_start', + 'agent_end' + ]) + }) + + it('keeps describing the lead’s children after a task child registers on its own bus', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + await harness.callHook('session_start') + harness.emitPiEvent('task:subagent:lifecycle', { id: 'c1', agent: 'task', status: 'started' }) + await vi.advanceTimersByTimeAsync(0) + harness.registerTaskChild() + harness.emitPiEvent('task:subagent:lifecycle', { id: 'c2', agent: 'task', status: 'started' }) + await vi.advanceTimersByTimeAsync(0) + + expect(childIds(posts(harness).at(-1))).toEqual(['c1', 'c2']) + }) + + it('ignores a task child’s own subagents, which run on that child’s bus', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'omp' }) + await harness.callHook('agent_start') + harness.emitPiEvent('task:subagent:lifecycle', { id: 'c1', agent: 'task', status: 'started' }) + await vi.advanceTimersByTimeAsync(0) + const emitOnChildBus = harness.registerTaskChild() + emitOnChildBus('task:subagent:lifecycle', { id: 'g1', agent: 'task', status: 'started' }) + emitOnChildBus('task:subagent:lifecycle', { id: 'g1', status: 'completed' }) + await vi.advanceTimersByTimeAsync(0) + + expect(postedHookNames(harness)).toEqual(['agent_start', 'agent_start']) + }) +}) diff --git a/src/main/pi/agent-status-extension-source.ts b/src/main/pi/agent-status-extension-source.ts index 58a18b78fa5..5a3d23f47a7 100644 --- a/src/main/pi/agent-status-extension-source.ts +++ b/src/main/pi/agent-status-extension-source.ts @@ -15,6 +15,7 @@ import type { PiAgentKind } from '../../shared/pi-agent-kind' import { getPiAgentStatusHandlerSourceLines } from './agent-status-handler-source' import { getPiAgentStatusRuntimeDetectionSourceLines } from './agent-status-runtime-detection-source' import { getPiAgentStatusWslCurlSourceLines } from './agent-status-wsl-curl-source' +import { getPiSubagentSnapshotSourceLines } from './agent-status-subagent-roster-source' export const ORCA_PI_AGENT_STATUS_EXTENSION_FILE = 'orca-agent-status.ts' @@ -53,14 +54,14 @@ export function getPiAgentStatusExtensionSource(kind: PiAgentKind = 'pi'): strin ' return ompRuntime ? { ...runtimeOmpSessionMetadata, ...modelMetadata } : sessionMetadata', '}', '', - 'function getPersistedSessionMetadata(): Record {', - ' const sessionFile = sessionMetadata.session_file', + 'function getPersistedSessionMetadata(metadata: Record): Record {', + ' const sessionFile = metadata.session_file', " if (typeof sessionFile !== 'string' || !sessionFile) return {}", ' try {', " const fs = require('fs')", ' // Why: Pi publishes its planned path before creating the transcript;', ' // recheck on every post so the first completed turn becomes resumable.', - ' return fs.existsSync(sessionFile) ? sessionMetadata : {}', + ' return fs.existsSync(sessionFile) ? metadata : {}', ' } catch {', ' return {}', ' }', @@ -123,8 +124,8 @@ export function getPiAgentStatusExtensionSource(kind: PiAgentKind = 'pi'): strin // Why: Pi resumes from an existing transcript; OMP resumes directly by session id (#8962). const payloadLine = kind !== 'omp' - ? ' payload: { hook_event_name: hookEventName, ...(ompRuntime ? metadata : getPersistedSessionMetadata()), ...extra },' - : ' payload: { hook_event_name: hookEventName, ...metadata, ...extra },' + ? ' payload: { hook_event_name: hookEventName, ...(ompRuntime ? metadata : getPersistedSessionMetadata(metadata)), ...(final ? {} : subagentPayload()), ...extra },' + : ' payload: { hook_event_name: hookEventName, ...metadata, ...(final ? {} : subagentPayload()), ...extra },' // Why: keep this string self-contained — it runs inside the pi process, // so it cannot import from Orca's main bundle. fs/http coords come from @@ -140,10 +141,11 @@ 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', - ...getPiAgentStatusPostQueueSourceLines(), + ...getPiAgentStatusPostQueueSourceLines(kind), ...(kind === 'pi' ? ['let piUiPromptDepth = 0', 'let piTurnInFlight = false'] : []), ...modelMetadataSourceLines, '', + ...getPiSubagentSnapshotSourceLines(), ...sessionMetadataSourceLines, '', '// Why: re-reading the endpoint file on every event is cheap (small file,', @@ -202,14 +204,23 @@ export function getPiAgentStatusExtensionSource(kind: PiAgentKind = 'pi'): strin '', ...getPiAgentStatusRuntimeDetectionSourceLines(kind), '', - 'function post(hookEventName: string, extra: Record = {}): void {', + // `final` marks the last post of a session that is being closed. + 'function post(hookEventName: string, extra: Record = {}, final = false): void {', ' const ompRuntime = isOmpRuntime()', - ' cancelPostRetry()', + // Why: the body is built at delivery, so the session is pinned here, where the event happened. ' const metadata = getPostSessionMetadata(ompRuntime)', + ' if (final) {', + ' postQueue.finalPosts.push({ revision: postQueue.postRevision, attempts: 0, delivered: false, hookEventName, extra, metadata, ompRuntime, final })', + ' drainPosts()', + ' return', + ' }', + ' cancelPostRetry()', + // Why: a new turn supersedes a same-session completion even when /reload preserved its delivery. + " if (hookEventName === 'before_agent_start' || hookEventName === 'agent_start') retireTurnCompletionPosts(metadata)", '// 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 = {', - ' revision: ++postRevision,', + " const previousCompletion = postQueue.latestPost?.hookEventName === 'agent_end' && !postQueue.latestPost.delivered && postQueue.latestPost.metadata.session_id === metadata.session_id", + ' postQueue.pendingPost = {', + ' revision: ++postQueue.postRevision,', ' attempts: 0,', ' delivered: false,', " hookEventName: ompRuntime && hookEventName === 'model_select' && previousCompletion ? 'agent_end' : hookEventName,", @@ -220,7 +231,7 @@ export function getPiAgentStatusExtensionSource(kind: PiAgentKind = 'pi'): strin ' metadata,', ' ompRuntime,', ' }', - ' latestPost = pendingPost', + ' postQueue.latestPost = postQueue.pendingPost', ' drainPosts()', '}', '', @@ -228,7 +239,8 @@ export function getPiAgentStatusExtensionSource(kind: PiAgentKind = 'pi'): strin ' hookEventName: string,', ' extra: Record,', ' metadata: Record,', - ' ompRuntime: boolean', + ' ompRuntime: boolean,', + ' final: boolean', '): Promise {', ' const coords = resolveHookCoords()', ' const paneKey = process.env.ORCA_PANE_KEY', diff --git a/src/main/pi/agent-status-extension-test-harness.ts b/src/main/pi/agent-status-extension-test-harness.ts index 0c07caf8256..c9cfe952697 100644 --- a/src/main/pi/agent-status-extension-test-harness.ts +++ b/src/main/pi/agent-status-extension-test-harness.ts @@ -1,5 +1,5 @@ import { EventEmitter } from 'node:events' -import { runInNewContext } from 'node:vm' +import { createContext, runInContext } from 'node:vm' // TypeScript 7 is a native CLI; transpile tests still need the legacy JavaScript API. import ts from 'typescript-api' import { vi } from 'vitest' @@ -19,6 +19,10 @@ export type HookContext = { } } +type PiEventBus = { + on: (name: string, listener: (event: unknown) => void) => unknown +} + export type HookHandler = (event?: unknown, context?: HookContext) => Promise | void type FakeCurlChild = { @@ -48,9 +52,16 @@ export type AgentStatusExtensionHarness = { callHook: (name: string, event?: unknown, context?: HookContext) => Promise emitPiEvent: (name: string, event: unknown) => void piEventListenerCount: (name: string) => number - // Re-invoke the extension factory in the same process (as Pi does on an - // in-process extension reload), swapping in the freshly registered handlers. + // Re-run the factory on the same event bus and module state. reload: () => void + // What Pi does for /new, resume and fork: shut the old registration down, drop its bus + // subscriptions, then run the factory again on a fresh `pi.events` (module state kept). + replacePiSession: (reason: 'new' | 'resume' | 'fork', targetSessionFile?: string) => Promise + // What Pi does for /reload: as above, but the module is evaluated again and `globalThis` survives. + reloadPi: () => Promise + // An OMP in-process task child: the same module's factory, run again on the child's own bus. + // Returns an emitter for events on that child's bus. + registerTaskChild: () => (name: string, event: unknown) => void } const BASE_ENV = { @@ -80,6 +91,10 @@ export function createAgentStatusExtensionHarness(args: { statSync?: (path: string) => { mtimeMs: number; size: number; ino: number } curlExitCode?: number | null fetchImpl?: (...params: Parameters) => Promise + // Runs before the first registration, to leave state an older build would have put on the bus. + seedEventBus?: (bus: EventEmitter) => void + // Pi before 0.84 handed every registration the one shared bus. + sharedEventBus?: boolean }): AgentStatusExtensionHarness { const fetchMock = vi.fn( args.fetchImpl ?? @@ -132,7 +147,7 @@ export function createAgentStatusExtensionHarness(args: { command: { handler: (args: string, context: HookContext) => Promise } ) => void setModel: (model: unknown) => Promise - events?: EventEmitter + events?: PiEventBus }) => void } } = { exports: {} } @@ -186,30 +201,71 @@ export function createAgentStatusExtensionHarness(args: { target: ts.ScriptTarget.ES2020 } }).outputText - runInNewContext(output, context) - - const register = module.exports.default - if (!register) { - throw new Error('expected default export from generated source') + createContext(context) + // Why: a function scope per evaluation, so a Pi reload can evaluate the module again in one realm. + const evaluateModule = (): NonNullable => { + runInContext(`(function () {\n${output}\n})()`, context) + const factory = module.exports.default + if (!factory) { + throw new Error('expected default export from generated source') + } + return factory } + let register = evaluateModule() const handlers: Record = {} - const piEvents = new EventEmitter() + // Why: Pi calls every handler an extension registers for an event, in registration order. + let handlerLists: Record = {} + let piEvents = new EventEmitter() + let busSubscriptions: [string, (event: unknown) => void][] = [] const commands: AgentStatusExtensionHarness['commands'] = {} const setModelMock = vi.fn(async (_model: unknown) => true) - const registerInto = (target: Record): void => { + const registerInto = ( + target: Record, + events: PiEventBus = piEvents + ): void => { + handlerLists = {} register({ registerCommand: (name, command) => { commands[name] = command }, setModel: setModelMock, - events: piEvents, + events, on(name: string, handler: HookHandler) { target[name] = handler + ;(handlerLists[name] ??= []).push(handler) } }) } - registerInto(handlers) + const callHook: AgentStatusExtensionHarness['callHook'] = async (name, event, hookContext) => { + for (const handler of handlerLists[name] ?? []) { + await handler(event, hookContext) + } + } + // Why: each Pi registration gets its own `pi.events` object, and Pi removes its subscriptions when + // the registration is replaced; anything an extension stores on that object goes with it. + const registerLikePi = (): void => { + for (const [name, listener] of busSubscriptions) { + piEvents.off(name, listener) + } + busSubscriptions = [] + for (const key of Object.keys(handlers)) { + delete handlers[key] + } + const bus = piEvents + registerInto(handlers, { + on(name: string, listener: (event: unknown) => void) { + bus.on(name, listener) + busSubscriptions.push([name, listener]) + } + }) + } + args.seedEventBus?.(piEvents) + if (args.kind === 'pi' && !args.sharedEventBus) { + registerLikePi() + } else { + registerInto(handlers) + } return { setModelMock, @@ -221,9 +277,7 @@ export function createAgentStatusExtensionHarness(args: { fsMock, handlers, processEnv: processMock.env, - callHook: async (name, event, hookContext) => { - await handlers[name]?.(event, hookContext) - }, + callHook, emitPiEvent: (name, event) => { piEvents.emit(name, event) }, @@ -233,6 +287,25 @@ export function createAgentStatusExtensionHarness(args: { delete handlers[key] } registerInto(handlers) + }, + replacePiSession: async (reason, targetSessionFile) => { + await callHook('session_shutdown', { reason, targetSessionFile }) + piEvents = new EventEmitter() + registerLikePi() + }, + registerTaskChild: () => { + const leadHandlerLists = handlerLists + const childBus = new EventEmitter() + registerInto({}, childBus) + handlerLists = leadHandlerLists + return (name, event) => { + childBus.emit(name, event) + } + }, + reloadPi: async () => { + await callHook('session_shutdown', { reason: 'reload' }) + register = evaluateModule() + registerLikePi() } } } diff --git a/src/main/pi/agent-status-handler-source.ts b/src/main/pi/agent-status-handler-source.ts index 76eff628731..af2735159fe 100644 --- a/src/main/pi/agent-status-handler-source.ts +++ b/src/main/pi/agent-status-handler-source.ts @@ -8,6 +8,10 @@ import { getPiSubagentRosterEventSourceLines, getPiSubagentRosterSetupSourceLines } from './agent-status-subagent-roster-source' +import { + getAgentStatusRunCloseOutSourceLines, + getAgentStatusSessionBoundaryHandlerSourceLines +} from './agent-status-session-boundary-source' // Why: keep the generated handler registrations separate from hook transport; // both are independently sizeable and the installed extension concatenates them. @@ -22,6 +26,13 @@ export function getPiAgentStatusHandlerSourceLines(kind: PiAgentKind): string[] ' // turn boundary and must not clear the visible status or unread state.', " if (event.reason === 'reload') return", " post('session_start')", + ...(kind === 'pi' + ? [ + ' restoreParkedSubagents()', + // Why: the host reads session_start as an idle session; children still hold this one. + " if (lifecycleState.waiting) post('agent_start')" + ] + : []), ' })', '' ] @@ -133,27 +144,9 @@ export function getPiAgentStatusHandlerSourceLines(kind: PiAgentKind): string[] ' if (ownerPid && ownerPid !== selfPid && isStatusOwnerAlive(ownerPid)) return', ` process.env.${ownerEnv} = selfPid`, ' resetPostQueue()', - ...getPiSubagentRosterSetupSourceLines(), - ...(kind !== 'pi' - ? [ - " pi.on('session_shutdown', () => { lifecycleState.active.clear(); lifecycleState.exited?.clear(); lifecycleState.waiting = false; lifecycleState.rootRunInFlight = false; resetPostQueue(); clearPendingAgentEndCheck() })" - ] - : []), - ...(kind !== 'prime-agent' - ? [ - " pi.on('session_switch', (_event, ctx) => {", - ' if (!isOmpRuntime()) return', - ' lifecycleState.active.clear()', - ' lifecycleState.exited?.clear()', - ' lifecycleState.waiting = false', - ' lifecycleState.rootRunInFlight = false', - ' resetPostQueue()', - ' clearPendingAgentEndCheck()', - ' updateRuntimeOmpSessionMetadata(ctx)', - ' })' - ] - : []), + ...getPiSubagentRosterSetupSourceLines(kind), ...getOmpSessionOwnerHandlerSourceLines(), + ...getAgentStatusSessionBoundaryHandlerSourceLines(kind), ...getOmpModelCommandSourceLines(), ...sessionStartHandler, ...(kind === 'omp' ? getPiPrefillHandlerSourceLines('omp', true) : []), @@ -166,8 +159,7 @@ export function getPiAgentStatusHandlerSourceLines(kind: PiAgentKind): string[] ...captureSessionMetadata, ' clearPendingAgentEndCheck()', ' lifecycleState.waiting = false', - ' lifecycleState.rootRunInFlight = true', - ' runGeneration += 1', + ' lifecycleState.runGeneration += 1', // Why: a turn cannot begin under a dialog holding input focus, so this is the one // boundary that can recover a modal whose close never arrived. ...(kind === 'pi' ? [' piUiPromptDepth = 0', ' piTurnInFlight = true'] : []), @@ -219,15 +211,6 @@ export function getPiAgentStatusHandlerSourceLines(kind: PiAgentKind): string[] ' const AGENT_END_IDLE_RECHECK_MS = 25', ' const AGENT_END_IDLE_RECHECK_MAX_MS = 250', ' let agentSettledSupported = false', - // Why: completion is a per-RUN fact. A sibling extension (the memory reminder is one) - // can start the next run from inside its own agent_settled handler, and Pi dispatches - // handlers in registration order, so this extension sees that run's agent_start - // BEFORE its own agent_settled for the run that just ended. A boolean "already - // posted" latch reset on agent_start then eats the newer run's completion and leaves - // the host stuck on that run's last working event. - ' let runGeneration = 0', - ' let endedRunGeneration = 0', - ' let completionPostedGeneration = -1', ' let agentEndIdleRecheckMs = AGENT_END_IDLE_RECHECK_MS', ' let pendingAgentEndCheck: ReturnType | null = null', ' let pendingAgentEndContext: { isIdle: () => boolean } | null = null', @@ -238,26 +221,37 @@ export function getPiAgentStatusHandlerSourceLines(kind: PiAgentKind): string[] ' pendingAgentEndContext = null', ' }', ...getPiSubagentRosterEventSourceLines(), - ' function postAgentEndOnce(): void {', - ' for (const id of lifecycleState.exited ?? []) lifecycleState.active.delete(id)', - ' lifecycleState.exited?.clear()', - ' if (lifecycleState.active.size > 0) {', + // The one answer to "do children still hold this pane"; an exited runner holds through its grace. + ' function isHeldByChildren(): boolean {', + ' return lifecycleState.active.size > 0', + ' }', + '', + ' function isTurnInFlight(): boolean {', + ' return lifecycleState.runGeneration !== lifecycleState.endedRunGeneration', + ' }', + '', + ' function postAgentEndOnce(final = false): boolean {', + ' for (const id of lifecycleState.exited ?? []) forgetSubagent(id)', + ' if (isHeldByChildren()) {', ' lifecycleState.waiting = true', - ' return', + ' return false', ' }', ' lifecycleState.waiting = false', - ' if (completionPostedGeneration === endedRunGeneration) return', - ' completionPostedGeneration = endedRunGeneration', + ' if (lifecycleState.completionPostedGeneration === lifecycleState.endedRunGeneration) return false', + ' lifecycleState.completionPostedGeneration = lifecycleState.endedRunGeneration', // Why: distinct from the completion guard, which holds the generation of the posted run // and so starts clean on a pane that has not run a turn yet — that pane is idle, not busy. ...(kind === 'pi' ? [' piTurnInFlight = false'] : []), - " post('agent_end')", + // Why: a run that ends because its session did has not completed; the host records it quietly. + " post('agent_end', final ? { session_boundary: true } : {}, final)", + ' return true', ' }', '', + ...getAgentStatusRunCloseOutSourceLines(), ' function checkPendingAgentEnd(): void {', ' pendingAgentEndCheck = null', ' const ctx = pendingAgentEndContext', - ' if (!ctx || agentSettledSupported || completionPostedGeneration === endedRunGeneration) {', + ' if (!ctx || agentSettledSupported || lifecycleState.completionPostedGeneration === lifecycleState.endedRunGeneration) {', ' pendingAgentEndContext = null', ' return', ' }', @@ -289,8 +283,7 @@ export function getPiAgentStatusHandlerSourceLines(kind: PiAgentKind): string[] ' clearPendingAgentEndCheck()', ' return', ' }', - ' lifecycleState.rootRunInFlight = false', - ' endedRunGeneration = runGeneration', + ' lifecycleState.endedRunGeneration = lifecycleState.runGeneration', ' if (isOmpRuntime()) {', ' postAgentEndOnce()', ' return', diff --git a/src/main/pi/agent-status-post-queue-source.ts b/src/main/pi/agent-status-post-queue-source.ts index 512ed932339..2775cb75f14 100644 --- a/src/main/pi/agent-status-post-queue-source.ts +++ b/src/main/pi/agent-status-post-queue-source.ts @@ -1,49 +1,97 @@ -export function getPiAgentStatusPostQueueSourceLines(): string[] { +import type { PiAgentKind } from '../../shared/pi-agent-kind' + +export function getPiAgentStatusPostQueueSourceLines(kind: PiAgentKind): 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', + 'type HookPost = { hookEventName: string; extra: Record; metadata: Record; ompRuntime: boolean; revision: number; attempts: number; delivered: boolean; final?: boolean }', + 'type HookPostQueue = { activePost: HookPost | null; pendingPost: HookPost | null; latestPost: HookPost | null; finalPosts: HookPost[]; finalRetryTimer: ReturnType | null; postRevision: number; retryTimer: ReturnType | null }', + ...(kind === 'pi' + ? [ + 'declare global { var __orcaPiStatusPostQueue: HookPostQueue | undefined }', + // Why: old delivery continuations and a reloaded registration must serialize through one queue. + 'const postQueue: HookPostQueue = globalThis.__orcaPiStatusPostQueue ??= { activePost: null, pendingPost: null, latestPost: null, finalPosts: [], finalRetryTimer: null, postRevision: 0, retryTimer: null }' + ] + : [ + 'const postQueue: HookPostQueue = { activePost: null, pendingPost: null, latestPost: null, finalPosts: [], finalRetryTimer: null, postRevision: 0, retryTimer: null }' + ]), '', 'function cancelPostRetry(): void {', - ' if (retryTimer !== null) clearTimeout(retryTimer)', - ' retryTimer = null', + ' if (postQueue.retryTimer !== null) clearTimeout(postQueue.retryTimer)', + ' postQueue.retryTimer = null', + '}', + '', + '// Why: a queued or scheduled post builds its body later, so it already carries newer state.', + 'function hasQueuedPost(): boolean {', + ' return postQueue.pendingPost !== null || postQueue.retryTimer !== null', '}', '', 'function resetPostQueue(): void {', ' cancelPostRetry()', - ' postRevision++', - ' pendingPost = null', - ' latestPost = null', + ' postQueue.postRevision++', + // Why: preserve the original object so an in-flight acknowledgment also retires its final queue entry. + ' for (const post of [postQueue.activePost, postQueue.pendingPost, postQueue.latestPost]) {', + " if (!post || post.hookEventName !== 'agent_end' || post.delivered || post.final) continue", + ' post.final = true', + ' postQueue.finalPosts.push(post)', + ' }', + ' postQueue.pendingPost = null', + ' postQueue.latestPost = null', + ' drainPosts()', + '}', + '', + 'function retireTurnCompletionPosts(metadata: Record): void {', + ' const first = postQueue.finalPosts[0]', + ' for (const post of postQueue.finalPosts) {', + ' if (post.extra.session_boundary !== true && post.metadata.session_id === metadata.session_id) post.final = false', + ' }', + ' postQueue.finalPosts = postQueue.finalPosts.filter((post) => post.final)', + ' if (first && !first.final && postQueue.finalRetryTimer !== null) {', + ' clearTimeout(postQueue.finalRetryTimer)', + ' postQueue.finalRetryTimer = 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 })', + // Why: a final post waiting to retry holds newer posts back, so the old session's row is never written last. + ' if (postQueue.activePost || postQueue.finalRetryTimer !== null) return', + ' const next = postQueue.finalPosts[0] ?? postQueue.pendingPost', + ' if (!next) return', + ' if (!next.final) postQueue.pendingPost = null', + ' postQueue.activePost = next', + ' void postOnce(next.hookEventName, next.extra, next.metadata, next.ompRuntime, next.final === true)', + ' .then(() => {', + ' next.delivered = true', + ' if (next.final) postQueue.finalPosts.shift()', + ' })', ' .catch(() => {', - ' if (!next.ompRuntime || next.revision !== postRevision) return', + ' if (next.final) {', + ' if (next.attempts >= 3) {', + ' postQueue.finalPosts.shift()', + " console.warn('[orca-pi-status] hook delivery failed after retries:', next.hookEventName)", + ' return', + ' }', + ' postQueue.finalRetryTimer = setTimeout(() => {', + ' postQueue.finalRetryTimer = null', + ' drainPosts()', + ' }, 250 * 2 ** next.attempts++)', + " if (typeof postQueue.finalRetryTimer.unref === 'function') postQueue.finalRetryTimer.unref()", + ' return', + ' }', + ' if (!next.ompRuntime || next.revision !== postQueue.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', + ' postQueue.retryTimer = setTimeout(() => {', + ' postQueue.retryTimer = null', + ' if (next.revision !== postQueue.postRevision) return', + ' postQueue.pendingPost = next', ' drainPosts()', ' }, delay)', - " if (typeof retryTimer.unref === 'function') retryTimer.unref()", + " if (typeof postQueue.retryTimer.unref === 'function') postQueue.retryTimer.unref()", ' })', ' .finally(() => {', - ' activePost = false', + ' postQueue.activePost = null', ' drainPosts()', ' })', '}', diff --git a/src/main/pi/agent-status-reload-delivery.test.ts b/src/main/pi/agent-status-reload-delivery.test.ts new file mode 100644 index 00000000000..d8b4b9e4f7a --- /dev/null +++ b/src/main/pi/agent-status-reload-delivery.test.ts @@ -0,0 +1,87 @@ +import { afterEach, beforeEach, expect, it, vi } from 'vitest' +import { createAgentStatusExtensionHarness } from './agent-status-extension-test-harness' +import { agentEndCount, endTurn, posts } from './agent-status-subagent-event-fixtures' + +beforeEach(() => vi.useFakeTimers()) +afterEach(() => vi.useRealTimers()) + +it('does not overwrite the next turn with an old completion retried across reload', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'pi', existsSync: () => true }) + const ctx = { + isIdle: () => true, + sessionManager: { getSessionId: () => 'A', getSessionFile: () => '/sessions/A.jsonl' } + } + await harness.callHook('session_start', { reason: 'startup' }, ctx) + await harness.callHook('agent_start', {}, ctx) + await vi.advanceTimersByTimeAsync(0) + harness.fetchMock.mockRejectedValueOnce(new Error('offline')) + await endTurn(harness) + harness.fetchMock.mockRejectedValueOnce(new Error('still offline')) + await harness.reloadPi() + await vi.advanceTimersByTimeAsync(0) + await harness.callHook('session_start', { reason: 'reload' }, ctx) + const completionCount = agentEndCount(harness) + await harness.callHook('agent_start', {}, ctx) + await vi.advanceTimersByTimeAsync(1_000) + expect(agentEndCount(harness)).toBe(completionCount) + + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_start', session_id: 'A' }) +}) + +it.each(['failure', 'success'] as const)( + 'serializes the next turn behind an in-flight completion on reload (%s)', + async (delivery) => { + const harness = createAgentStatusExtensionHarness({ kind: 'pi', existsSync: () => true }) + const ctx = { + isIdle: () => true, + sessionManager: { getSessionId: () => 'A', getSessionFile: () => '/sessions/A.jsonl' } + } + await harness.callHook('session_start', { reason: 'startup' }, ctx) + await harness.callHook('agent_start', {}, ctx) + await vi.advanceTimersByTimeAsync(0) + let acknowledge: ((value: { ok: boolean }) => void) | undefined + let fail: ((error: Error) => void) | undefined + harness.fetchMock.mockImplementationOnce( + () => + new Promise((resolve, reject) => { + acknowledge = resolve + fail = reject + }) + ) + await endTurn(harness) + await harness.reloadPi() + await harness.callHook('session_start', { reason: 'reload' }, ctx) + await harness.callHook('agent_start', {}, ctx) + await vi.advanceTimersByTimeAsync(0) + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end', session_id: 'A' }) + const completionCount = agentEndCount(harness) + if (delivery === 'failure') { + fail?.(new Error('late failure')) + } else { + acknowledge?.({ ok: true }) + } + await vi.advanceTimersByTimeAsync(5_000) + expect(agentEndCount(harness)).toBe(completionCount) + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_start', session_id: 'A' }) + expect(vi.getTimerCount()).toBe(0) + } +) + +it('recovers a completion on reload without requiring a new turn', async () => { + const harness = createAgentStatusExtensionHarness({ kind: 'pi', existsSync: () => true }) + const ctx = { + isIdle: () => true, + sessionManager: { getSessionId: () => 'A', getSessionFile: () => '/sessions/A.jsonl' } + } + await harness.callHook('session_start', { reason: 'startup' }, ctx) + await harness.callHook('agent_start', {}, ctx) + await vi.advanceTimersByTimeAsync(0) + harness.fetchMock.mockRejectedValueOnce(new Error('offline')) + await endTurn(harness) + await harness.reloadPi() + await harness.callHook('session_start', { reason: 'reload' }, ctx) + await vi.advanceTimersByTimeAsync(5_000) + expect(agentEndCount(harness)).toBe(2) + expect(posts(harness).at(-1)).toMatchObject({ hook_event_name: 'agent_end', session_id: 'A' }) + expect(vi.getTimerCount()).toBe(0) +}) diff --git a/src/main/pi/agent-status-session-boundary-source.ts b/src/main/pi/agent-status-session-boundary-source.ts new file mode 100644 index 00000000000..b958ba5b346 --- /dev/null +++ b/src/main/pi/agent-status-session-boundary-source.ts @@ -0,0 +1,110 @@ +import type { PiAgentKind } from '../../shared/pi-agent-kind' + +// What the generated extension does when a session ends or is replaced. A run belongs to one +// session: once that session is gone its children never report here again, so the run ends with it. + +// Pi replaces the registration for /new, resume and fork, and again for /reload, which keeps the session. +function getPiSessionShutdownHandlerSourceLines(): string[] { + return [ + " onStatus('session_shutdown', (event) => {", + ' clearPendingAgentEndCheck()', + ' if (isOmpRuntime()) {', + ' clearRunnerExitCheck()', + ' resetPostQueue()', + ' return', + ' }', + // Why: pi tears an open dialog down without resolving its promise, so no ui_prompt_end follows. + ' piUiPromptDepth = 0', + ' const reason = (event as { reason?: unknown } | null)?.reason', + ' const target = (event as { targetSessionFile?: unknown } | null)?.targetSessionFile', + // Why: /reload and a resume of the file already open keep the session, and its children still report to it. + " const keepsSession = reason === 'reload' || (reason === 'resume' && typeof target === 'string' && target === sessionMetadata.session_file)", + ' if (!keepsSession) clearRunnerExitCheck()', + // Why: on quit the PTY's exit clears the pane, and a done here would notify on every quit. + " if (keepsSession || reason === 'quit') {", + // Why: this registration's queue outlives it and would deliver stale posts after the next one's. + ' resetPostQueue()', + ' return', + ' }', + // Also a Pi too old to give a reason: its children would otherwise hold the pane for good. + ' closeOutRun(sessionMetadata.session_file)', + // Why: pi-subagents can still emit here until Pi invalidates this registration. + ' lifecycleState.onEvent = undefined', + ' })', + '' + ] +} + +// Registered after onStatus exists; expects the roster and closeOutRun() from the handler scope. +export function getAgentStatusSessionBoundaryHandlerSourceLines(kind: PiAgentKind): string[] { + return [ + ...(kind !== 'pi' + ? [ + " pi.on('session_shutdown', () => { resetSubagentRoster(); resetPostQueue(); clearPendingAgentEndCheck() })" + ] + : []), + ...(kind !== 'prime-agent' + ? [ + // Why: OMP keeps this registration and bus across a switch. /new and a branch cancel the + // session's own jobs; fork and resume (which is also how it reloads) leave them running + // and reporting here. + ' const onOmpSessionChange = (event, ctx) => {', + ' if (!isOmpRuntime()) return', + " if (event?.reason === 'fork' || event?.reason === 'resume') {", + // Why: a resume cuts a running turn off without an agent_end. + ' if (isTurnInFlight()) {', + ' lifecycleState.endedRunGeneration = lifecycleState.runGeneration', + ' postAgentEndOnce()', + ' }', + ' } else {', + ' closeOutRun()', + ' }', + ' updateRuntimeOmpSessionMetadata(ctx)', + ' }', + " pi.on('session_switch', onOmpSessionChange)", + " pi.on('session_branch', onOmpSessionChange)" + ] + : []), + ...(kind === 'pi' ? getPiSessionShutdownHandlerSourceLines() : []) + ] +} + +export function getAgentStatusRunCloseOutSourceLines(): string[] { + return [ + // Why: pi-subagents re-attaches a resumed session's runs and reports their completion to it. + ' function restoreParkedSubagents(): void {', + ' const file = sessionMetadata.session_file', + " if (typeof file !== 'string') return", + ' const parked = lifecycleState.parked?.get(file)', + ' if (!parked) return', + ' lifecycleState.parked?.delete(file)', + ' for (const [id, detail] of parked) {', + ' lifecycleState.active.add(id)', + ' subagentDetails.set(id, detail)', + ' }', + ' if (!isHeldByChildren()) return', + ' lifecycleState.waiting = true', + ' lifecycleState.completionPostedGeneration = -1', + ' }', + '', + // `parkUnder` keeps the children for a resume into that session file, where pi-subagents + // announces each one's completion again. + ' function closeOutRun(parkUnder?: unknown): void {', + ' const unsettled = lifecycleState.waiting || isTurnInFlight()', + " if (typeof parkUnder === 'string' && isHeldByChildren()) {", + ' const parked = new Map()', + ' for (const id of lifecycleState.active) {', + ' const detail = subagentDetails.get(id)', + ' if (detail && !lifecycleState.exited?.has(id)) parked.set(id, detail)', + ' }', + ' ;(lifecycleState.parked ??= new Map()).set(parkUnder, parked)', + ' }', + ' resetSubagentRoster()', + ' resetPostQueue()', + ' if (!unsettled) return', + ' lifecycleState.endedRunGeneration = lifecycleState.runGeneration', + ' postAgentEndOnce(true)', + ' }', + '' + ] +} diff --git a/src/main/pi/agent-status-subagent-event-fixtures.ts b/src/main/pi/agent-status-subagent-event-fixtures.ts new file mode 100644 index 00000000000..92e92fd7bf8 --- /dev/null +++ b/src/main/pi/agent-status-subagent-event-fixtures.ts @@ -0,0 +1,84 @@ +import { vi } from 'vitest' + +import { + AGENT_STATUS_EXTENSION_SELF_PID, + type AgentStatusExtensionHarness +} from './agent-status-extension-test-harness' + +// Event shapes and orderings mirror traces recorded from pi-subagents 0.71.0. +export const WORKFLOW = 'workflow-1' +export const idle = { isIdle: () => true } + +type PostedChild = { id: string; state: string; startedAt: number; agentType?: string } +type PostedPayload = { + hook_event_name: string + session_id?: string + session_file?: string + subagents?: PostedChild[] +} + +export function posts(harness: AgentStatusExtensionHarness): PostedPayload[] { + return harness.fetchMock.mock.calls.map((call) => { + const body: { payload: PostedPayload } = JSON.parse(String(call[1]?.body)) + return body.payload + }) +} + +export function postedHookNames(harness: AgentStatusExtensionHarness): string[] { + return posts(harness).map((post) => post.hook_event_name) +} + +export function agentEndCount(harness: AgentStatusExtensionHarness): number { + return postedHookNames(harness).filter((name) => name === 'agent_end').length +} + +export function startWorkflow(harness: AgentStatusExtensionHarness): void { + harness.emitPiEvent('subagent:async-started', { + id: WORKFLOW, + mode: 'workflow', + agent: 'workflow', + pid: AGENT_STATUS_EXTENSION_SELF_PID + }) +} + +export function startChild( + harness: AgentStatusExtensionHarness, + id: string, + parent = WORKFLOW +): void { + harness.emitPiEvent('subagent:async-started', { + id, + mode: 'single', + pid: 4000, + parentWorkflowRunId: parent + }) +} + +export function exitRunner(harness: AgentStatusExtensionHarness, runId: string): void { + harness.emitPiEvent('subagent:process-terminal', { runId, state: 'observed' }) +} + +export function complete(harness: AgentStatusExtensionHarness, id: string): void { + harness.emitPiEvent('subagent:async-complete', { id, runId: id, state: 'complete' }) +} + +export function childIds(payload: PostedPayload | undefined): string[] | undefined { + return payload?.subagents?.map((child) => child.id) +} + +export function startAsync(harness: AgentStatusExtensionHarness, id: string, agent: string): void { + harness.emitPiEvent('subagent:async-started', { + id, + mode: 'single', + agent, + task: '[REDACTED]', + goal: '[REDACTED]', + pid: 4000 + }) +} + +export async function endTurn(harness: AgentStatusExtensionHarness): Promise { + await harness.callHook('agent_end', {}, idle) + await harness.callHook('agent_settled', undefined, idle) + await vi.advanceTimersByTimeAsync(0) +} diff --git a/src/main/pi/agent-status-subagent-roster-source.ts b/src/main/pi/agent-status-subagent-roster-source.ts index 106d49a7533..ac3d294542b 100644 --- a/src/main/pi/agent-status-subagent-roster-source.ts +++ b/src/main/pi/agent-status-subagent-roster-source.ts @@ -1,75 +1,161 @@ // Why: Pi settles its own turn while pi-subagents children keep running, so the // generated extension holds the pane's completion until every child it saw start is gone. -// The roster lives on pi.events so an in-process /reload keeps children and listeners. -export function getPiSubagentRosterSetupSourceLines(): string[] { +import type { PiAgentKind } from '../../shared/pi-agent-kind' +import { AGENT_STATUS_MAX_SUBAGENTS } from '../../shared/agent-status-types' + +// Module scope: post() reads the roster when a body is built, so a coalesced or +// retried post always carries the children live at delivery. +export function getPiSubagentSnapshotSourceLines(): string[] { return [ - ' const piEventBus = (pi as { events?: { on?: (name: string, handler: (event: unknown) => void) => void } }).events', - ' const lifecycleState = (piEventBus as { __orcaPiSubagents?: { active: Set; exited?: Set; waiting: boolean; ownsPane?: boolean; rootRunInFlight?: boolean; onEvent?: (event: unknown, forcedStatus?: string) => void; listener?: (event: unknown) => void; onRunnerExit?: (event: unknown) => void; runnerExitListener?: (event: unknown) => void } } | undefined)?.__orcaPiSubagents ?? { active: new Set(), waiting: false }', - ' if (piEventBus) (piEventBus as { __orcaPiSubagents?: unknown }).__orcaPiSubagents = lifecycleState', - ' if (piEventBus?.on && !(lifecycleState as { listener?: unknown }).listener) {', - ' const listener = (event: unknown) => lifecycleState.onEvent?.(event)', - ' lifecycleState.listener = listener', - " piEventBus.on('task:subagent:lifecycle', listener)", + 'type SubagentDetail = { agentType?: string; description?: string; startedAt: number; workflow?: boolean; parent?: string; registration?: object }', + 'type SubagentRoster = { active: Set; exited?: Set; details?: Map; waiting: boolean; ownsPane?: boolean; runGeneration?: number; endedRunGeneration?: number; completionPostedGeneration?: number; parked?: Map>; onEvent?: (event: unknown, forcedStatus?: string) => void; listener?: (event: unknown) => void; onRunnerExit?: (event: unknown) => void; runnerExitListener?: (event: unknown) => void; runnerExitCheck?: ReturnType | null; onRunnerExitSettled?: () => void }', + // Why: interpolated, not re-typed, so the extension cap cannot drift from the host's. + `const MAX_SUBAGENT_SNAPSHOT = ${AGENT_STATUS_MAX_SUBAGENTS}`, + 'let subagentRoster: SubagentRoster | null = null', + '', + // Why: an exited runner is already gone, and a pi-subagents workflow run is the lead + // coordinating children that post their own rows; both still hold the pane. + 'function isVisibleSubagent(roster: SubagentRoster, id: string): boolean {', + ' return roster.active.has(id) && !roster.exited?.has(id) && roster.details?.get(id)?.workflow !== true', + '}', + '', + 'function subagentPayload(): Record {', + ' if (!subagentRoster) return {}', + ' const subagents: Record[] = []', + ' for (const id of subagentRoster.active) {', + ' if (!isVisibleSubagent(subagentRoster, id)) continue', + ' const detail = subagentRoster.details?.get(id)', + " subagents.push({ id, state: 'working', startedAt: detail?.startedAt ?? 0, ...(detail?.agentType ? { agentType: detail.agentType } : {}), ...(detail?.description ? { description: detail.description } : {}) })", + ' if (subagents.length >= MAX_SUBAGENT_SNAPSHOT) break', + ' }', + ' return subagents.length > 0 ? { subagents } : {}', + '}', + '' + ] +} + +// The run state (children, the hold, the turn counters) has to outlive a registration: Pi hands +// each one a fresh `pi.events` and evaluates this module again on /reload, so only globalThis +// survives there. OMP and Prime keep one bus for the session. +export function getPiSubagentRosterSetupSourceLines(kind: PiAgentKind): string[] { + return [ + ' const piEventBus = (pi as { events?: { __orcaPiSubagents?: SubagentRoster; __orcaPiSubagentsHeard?: boolean; __orcaPiRunnerExitsHeard?: boolean; on?: (name: string, handler: (event: unknown) => void) => void } }).events', + kind === 'pi' + ? ' const runStateHome: { __orcaPiSubagents?: SubagentRoster } | undefined = isOmpRuntime() ? piEventBus : (globalThis as { __orcaPiSubagents?: SubagentRoster })' + : ' const runStateHome: { __orcaPiSubagents?: SubagentRoster } | undefined = piEventBus', + // Why: a roster an older in-process build left on this bus is adopted, with the subscriptions it made. + ' const busRoster = piEventBus?.__orcaPiSubagents', + ' const lifecycleState: SubagentRoster = busRoster ?? runStateHome?.__orcaPiSubagents ?? { active: new Set(), waiting: false }', + ' if (runStateHome) runStateHome.__orcaPiSubagents = lifecycleState', + // Why: optional on the shared roster so one created by an older in-process build gains them on /reload. + ' const subagentDetails = (lifecycleState.details ??= new Map())', + // Why: completion is a per-RUN fact. A sibling extension (the memory reminder is one) can + // start the next run from inside its own agent_settled handler, so this extension sees that + // run's agent_start BEFORE its own agent_settled for the run that just ended; a boolean + // "already posted" latch would eat the newer run's completion. + ' lifecycleState.runGeneration ??= 0', + ' lifecycleState.endedRunGeneration ??= 0', + ' lifecycleState.completionPostedGeneration ??= -1', + // Why: an OMP task child runs this factory again on its own bus. Posts keep describing the + // lead's children, and the child's own subagents must not settle the lead's pane. + ' subagentRoster ??= lifecycleState', + ' const ownsPaneRoster = subagentRoster === lifecycleState', + // Tells the children this registration saw start from the ones a /reload handed it. + ' const registration = {}', + ' function resetSubagentRoster(): void {', + ' clearRunnerExitCheck()', + ' lifecycleState.active.clear()', + ' lifecycleState.exited?.clear()', + ' subagentDetails.clear()', + ' lifecycleState.waiting = false', + ' }', + // Why: one subscription per bus object; Pi drops a replaced registration's own. + ' if (ownsPaneRoster && piEventBus?.on && !piEventBus.__orcaPiSubagentsHeard && !busRoster?.listener) {', + ' piEventBus.__orcaPiSubagentsHeard = true', + " piEventBus.on('task:subagent:lifecycle', (event: unknown) => lifecycleState.onEvent?.(event))", " piEventBus.on('subagent:async-started', (event: unknown) => lifecycleState.onEvent?.(event, 'started'))", " piEventBus.on('subagent:async-complete', (event: unknown) => lifecycleState.onEvent?.(event, 'completed'))", ' }', - // Why: separate guard so a roster created by an older in-process build still subscribes. - ' if (piEventBus?.on && !lifecycleState.runnerExitListener) {', - ' const runnerExitListener = (event: unknown) => lifecycleState.onRunnerExit?.(event)', - ' lifecycleState.runnerExitListener = runnerExitListener', - " piEventBus.on('subagent:process-terminal', runnerExitListener)", + ' if (ownsPaneRoster && piEventBus?.on && !piEventBus.__orcaPiRunnerExitsHeard && !busRoster?.runnerExitListener) {', + ' piEventBus.__orcaPiRunnerExitsHeard = true', + " piEventBus.on('subagent:process-terminal', (event: unknown) => lifecycleState.onRunnerExit?.(event))", ' }' ] } -// Expects post(), the run generations and postAgentEndOnce() from the handler scope; -// the latter prunes exited runners before deciding whether children still hold the pane. +// Expects post() and postAgentEndOnce() from the handler scope; the latter prunes +// exited runners, then returns whether it settled the pane. export function getPiSubagentRosterEventSourceLines(): string[] { return [ // Why: a run that reports its own completion does so ~150ms after its runner exits; // the grace lets that path (and the wake turn it triggers) settle the pane first. ' const RUNNER_EXIT_GRACE_MS = 2000', - ' let runnerExitCheck: ReturnType | null = null', + ' function clearRunnerExitCheck(): void {', + ' if (lifecycleState.runnerExitCheck != null) clearTimeout(lifecycleState.runnerExitCheck)', + ' lifecycleState.runnerExitCheck = null', + ' }', + ' function forgetSubagent(id: string): void {', + ' lifecycleState.active.delete(id)', + ' lifecycleState.exited?.delete(id)', + ' subagentDetails.delete(id)', + ' }', + // Why: a child can end with no lead event to carry it; a queued post already reads the new roster. + ' function postSubagentsUpdate(): void {', + " if (!hasQueuedPost()) post('subagents_update')", + ' }', + " const readLabel = (value: unknown): string | undefined => typeof value === 'string' && value.trim() ? value : undefined", ' lifecycleState.onEvent = (event: unknown, forcedStatus?: string): void => {', " if (!event || typeof event !== 'object') return", - ' const record = event as { id?: unknown; runId?: unknown }', + ' const record = event as { id?: unknown; runId?: unknown; agent?: unknown; description?: unknown; mode?: unknown; parentWorkflowRunId?: unknown }', " const id = typeof record.id === 'string' && record.id ? record.id : typeof record.runId === 'string' ? record.runId : ''", ' const status = forcedStatus ?? (event as { status?: unknown }).status', ' if (!id) return', - // Why: each OMP task session runs its own copy on its own bus; only the pane's bus tracks children. ' if (isOmpRuntime() && !lifecycleState.ownsPane) return', " if (status === 'started') {", ' lifecycleState.active.add(id)', - // Why: a child starting after the run's done owes a fresh done. Under OMP the same holds - // before the first turn of this factory run (a resumed root, or a reload), but only when no - // root run is in flight -- a reload mid-turn resets these counters while the root still works. - ' if (completionPostedGeneration === runGeneration || (isOmpRuntime() && runGeneration === 0 && !lifecycleState.rootRunInFlight)) {', + // Why: pi-subagents redacts task prompts, so only the agent name and OMP's short label are shown. + " if (!subagentDetails.has(id)) subagentDetails.set(id, { agentType: readLabel(record.agent), description: readLabel(record.description), startedAt: Date.now(), workflow: record.mode === 'workflow', parent: readLabel(record.parentWorkflowRunId), registration })", + // Why: Pi re-opens only a posted completion; an earlier child must leave the idle check intact. + ' if (!isTurnInFlight() && (isOmpRuntime() || lifecycleState.runGeneration === 0 || lifecycleState.completionPostedGeneration === lifecycleState.runGeneration)) {', ' lifecycleState.waiting = true', - ' completionPostedGeneration = -1', + ' lifecycleState.completionPostedGeneration = -1', ' }', " post('agent_start')", ' return', ' }', " if (status !== 'completed' && status !== 'failed' && status !== 'aborted') return", - ' lifecycleState.active.delete(id)', - ' lifecycleState.exited?.delete(id)', - ' if (lifecycleState.waiting) postAgentEndOnce()', + ' let wasVisible = isVisibleSubagent(lifecycleState, id)', + ' forgetSubagent(id)', + // Why: Pi drops the runner-exit events of runs started before a /reload, so a finished run + // takes the children it launched back then with it. + ' for (const [childId, detail] of subagentDetails) {', + ' if (detail.parent !== id || detail.registration === registration) continue', + ' wasVisible ||= isVisibleSubagent(lifecycleState, childId)', + ' forgetSubagent(childId)', + ' }', + ' if (lifecycleState.waiting && postAgentEndOnce()) return', + ' if (wasVisible) postSubagentsUpdate()', ' }', // Why: awaited workflow children never get subagent:async-complete; their runner // exiting is the only end signal pi-subagents publishes for them. + ' lifecycleState.onRunnerExitSettled = (): void => {', + ' if (lifecycleState.waiting) postAgentEndOnce()', + ' }', ' lifecycleState.onRunnerExit = (event: unknown): void => {', " const runId = event && typeof event === 'object' ? (event as { runId?: unknown }).runId : undefined", " if (typeof runId !== 'string' || !lifecycleState.active.has(runId)) return", ' if (!lifecycleState.exited) lifecycleState.exited = new Set()', + ' const wasVisible = isVisibleSubagent(lifecycleState, runId)', ' lifecycleState.exited.add(runId)', + ' if (wasVisible) postSubagentsUpdate()', ' if (!lifecycleState.waiting) return', - ' if (runnerExitCheck !== null) clearTimeout(runnerExitCheck)', - ' runnerExitCheck = setTimeout(() => {', - ' runnerExitCheck = null', - ' if (lifecycleState.waiting) postAgentEndOnce()', + ' clearRunnerExitCheck()', + ' lifecycleState.runnerExitCheck = setTimeout(() => {', + ' lifecycleState.runnerExitCheck = null', + ' lifecycleState.onRunnerExitSettled?.()', ' }, RUNNER_EXIT_GRACE_MS)', - " if (typeof runnerExitCheck.unref === 'function') runnerExitCheck.unref()", + " if (typeof lifecycleState.runnerExitCheck.unref === 'function') lifecycleState.runnerExitCheck.unref()", ' }' ] } diff --git a/src/main/pi/agent-status-ui-prompt-source.ts b/src/main/pi/agent-status-ui-prompt-source.ts index d961864cd41..9ff30bd888d 100644 --- a/src/main/pi/agent-status-ui-prompt-source.ts +++ b/src/main/pi/agent-status-ui-prompt-source.ts @@ -29,19 +29,10 @@ export function getPiAgentStatusUiPromptHandlerSourceLines(kind: PiAgentKind): s ' } catch {', ' // Why: a runner this very modal invalidated cannot answer; keep the local verdict.', ' }', + ' // Why: Pi reports idle once its own turn settles, but children still hold the run open.', + ' isIdle &&= !isHeldByChildren()', " post('ui_prompt_end', { is_idle: isIdle })", ' })', - '', - " onStatus('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', - ' // would stick forever. Reset without posting: shutdown is not a turn boundary, and', - ' // the session_start that follows republishes the corrected state.', - ' piUiPromptDepth = 0', - ' })', '' ] } diff --git a/src/shared/agent-hook-listener/providers/pi-family-events.ts b/src/shared/agent-hook-listener/providers/pi-family-events.ts index 9786f09fb56..bf2acdc89ef 100644 --- a/src/shared/agent-hook-listener/providers/pi-family-events.ts +++ b/src/shared/agent-hook-listener/providers/pi-family-events.ts @@ -10,7 +10,8 @@ import { readString } from '../tool-input-preview' /** Maps a Pi-family hook event (Pi, OMP, Prime) onto a pane status: lifecycle * events become `working` / `done`, an ask tool becomes `blocked`, and OMP's - * `model` stamp rides along. Returns null for events that carry no status. */ + * `model` stamp and the extension's live `subagents` ride along. Returns null + * for events that carry no status. */ export function normalizePiCompatibleEvent( state: HookListenerState, agentType: 'pi' | 'omp' | 'prime-agent', @@ -32,15 +33,26 @@ export function normalizePiCompatibleEvent( const model = readString(hookPayload, 'model') const modelSwitchCommand = hookPayload.model_switch_command === 'orca-model' ? 'orca-model' : undefined - if (eventName === 'model_select') { - // Why: a model switch happens between turns, so it must ride on the pane's last - // known status instead of inventing a state — and before any status exists there - // is nothing for a model to describe. - const previous = state.lastStatusByPaneKey.get(paneKey)?.payload - if (!model || !previous || previous.agentType !== agentType) { + // Why: every post restates the extension's whole live roster, so an absent list means none. + const subagents = hookPayload.subagents + if (eventName === 'model_select' || eventName === 'subagents_update') { + // Why: a model switch or a child ending happens between lead events, so it rides on the + // pane's last visible row instead of inventing a state. Before any row exists there is + // nothing to describe, and a providerSessionOnly placeholder is a hidden resume record. + const previous = state.lastStatusByPaneKey.get(paneKey) + if ( + !previous || + previous.providerSessionOnly === true || + previous.payload.agentType !== agentType || + (eventName === 'model_select' && !model) + ) { return null } - return normalizeAgentStatusPayload({ ...previous, model, modelSwitchCommand }) + return normalizeAgentStatusPayload({ + ...previous.payload, + ...(eventName === 'model_select' ? { model, modelSwitchCommand } : {}), + subagents + }) } // Why: gate on the event's own tool_name so a stale cached question can't re-enter blocked. @@ -107,7 +119,12 @@ export function normalizePiCompatibleEvent( toolInput: snapshot.toolInput, interactivePrompt: snapshot.interactivePrompt, lastAssistantMessage: snapshot.lastAssistantMessage, - lastAssistantMessageIsToolOutput: snapshot.lastAssistantMessageIsToolOutput + lastAssistantMessageIsToolOutput: snapshot.lastAssistantMessageIsToolOutput, + subagents, + // Why: the extension ends a run whose session was replaced; that is not a completed turn. + ...(eventName === 'agent_end' && hookPayload.session_boundary === true + ? { sessionBoundary: true } + : {}) }) } diff --git a/tests/tools/omp-completion-runtime-smoke.mjs b/tests/tools/omp-completion-runtime-smoke.mjs index 6c32c18db32..6e9a2da65ad 100644 --- a/tests/tools/omp-completion-runtime-smoke.mjs +++ b/tests/tools/omp-completion-runtime-smoke.mjs @@ -16,8 +16,9 @@ const server = createServer(async (request, response) => { for await (const chunk of request) { body += chunk } - const event = JSON.parse(body).payload.hook_event_name - received.push(event) + const payload = JSON.parse(body).payload + const event = payload.hook_event_name + received.push(payload) const reject = event === 'agent_end' && rejectCompletion if (reject) { rejectCompletion = false @@ -66,38 +67,68 @@ try { }) const handlers = new Map() require(extension).default({ on: (event, handler) => handlers.set(event, handler) }) - const emit = async (event) => { + let sessionId = 'A' + const context = { + isIdle: () => false, + sessionManager: { + getSessionId: () => sessionId, + getSessionFile: () => join(scratch, `${sessionId}.jsonl`) + } + } + const emit = async (event, payload = {}) => { assert.ok(handlers.has(event), `Missing lifecycle handler: ${event}`) - await handlers.get(event)({}, { isIdle: () => false }) + await handlers.get(event)(payload, context) } await emit('agent_start') await waitForRequests(1) await emit('agent_end') await waitForRequests(3) - assert.deepEqual(received, ['agent_start', 'agent_end', 'agent_end']) + assert.deepEqual( + received.map((event) => event.hook_event_name), + ['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('agent_start') + await delay(600) + assert.equal( + received.filter((event) => event.hook_event_name === 'agent_end').length, + 1, + 'An obsolete completion retried after the next turn started' + ) + + for (const boundary of ['session_switch', 'session_shutdown']) { received.length = 0 rejectCompletion = true + sessionId = 'A' await emit('agent_start') await waitForRequests(1) await emit('agent_end') await waitForRequests(2) - await emit(boundary) + sessionId = 'B' + await emit(boundary, { reason: 'new' }) + await emit('agent_start') + await waitForRequests(4) await delay(600) - assert.equal( - received.filter((event) => event === 'agent_end').length, - 1, - `An obsolete completion retried after ${boundary}` - ) + const completions = received.filter((event) => event.hook_event_name === 'agent_end') + assert.equal(completions.length, 2, `Unsent completion was lost at ${boundary}`) + assert.ok(completions.every((event) => event.session_id === 'A')) + assert.equal(received.at(-1).hook_event_name, 'agent_start') + assert.equal(received.at(-1).session_id, 'B', 'Old completion overwrote the new session') } console.log( JSON.stringify({ platform: process.platform, node: process.version, recovered, - cancelledAt: ['agent_start', 'session_switch', 'session_shutdown'], + cancelledAt: ['agent_start'], + preservedAt: ['session_switch', 'session_shutdown'], scope: 'Generated OMP extension, real native HTTP 503, synthetic lifecycle callbacks' }) )