mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 00:02:29 +00:00
feat(pi): show Pi and OMP background children as rows (#23201)
* feat(pi): show Pi and OMP background children as rows Refs #22868 * fix(pi): derive the extension's roster cap from the host's The generated extension re-typed the 32-row limit as its own literal, so a change to AGENT_STATUS_MAX_SUBAGENTS would leave the extension posting rosters the host silently truncates on arrival. Interpolate the host constant instead, and cover the cap with a test that posts past it. * test(pi): pin the Pi cancel verdict beside a live child row Publishing a subagent roster for Pi puts its panes behind the same child-work guard Claude and Codex sit behind, so Ctrl+C beside a live child no longer settles a stopped row. The extension reports no main-agent state, so Orca cannot tell a cancelled turn from Ctrl+C at the idle prompt of a lead children alone hold open, and keeps the live row. Both halves of that are now pinned, as is the resume-placeholder guard on the model_select arm of the same code path. * fix(pi): end a run with its session and keep children across reload Pi re-runs the extension on a fresh event bus for /new, resume, fork and /reload, so the roster kept on that bus was lost and nothing ended the run the old session's children held open. * fix(pi): preserve pending completion and OMP pane ownership * test(pi): use the hook owner context type * fix(pi): retain completion timers and delivery across session boundaries * test(omp): preserve unsent completion at session boundaries * test(journal): retain provider handle import across main integration * test(startup): account for encoded Windows prompt lines * fix(pi): serialize status delivery across module reloads * test(compatibility): retain the released journal parser dependency --------- Co-authored-by: Neil <neil@stably.ai> Co-authored-by: Neil <4138956+nwparker@users.noreply.github.com>
This commit is contained in:
co-authored by
Neil
Neil
parent
1b52be6255
commit
a259616f7c
@@ -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<string, unknown>): Promise<void> {
|
||||
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<string, unknown>[]): Promise<void> {
|
||||
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()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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(() => {}))
|
||||
|
||||
@@ -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<void> {
|
||||
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)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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<unknown>): AgentStatusExtensionHarness {
|
||||
return createAgentStatusExtensionHarness({ kind: 'pi', existsSync: () => true, fetchImpl })
|
||||
}
|
||||
|
||||
async function holdRunOpen(harness: AgentStatusExtensionHarness, sessionId = 'A'): Promise<void> {
|
||||
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'])
|
||||
})
|
||||
})
|
||||
@@ -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<string, unknown> {',
|
||||
' const sessionFile = sessionMetadata.session_file',
|
||||
'function getPersistedSessionMetadata(metadata: Record<string, unknown>): Record<string, unknown> {',
|
||||
' 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<string, unknown> = {}): void {',
|
||||
// `final` marks the last post of a session that is being closed.
|
||||
'function post(hookEventName: string, extra: Record<string, unknown> = {}, 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<string, unknown>,',
|
||||
' metadata: Record<string, unknown>,',
|
||||
' ompRuntime: boolean',
|
||||
' ompRuntime: boolean,',
|
||||
' final: boolean',
|
||||
'): Promise<void> {',
|
||||
' const coords = resolveHookCoords()',
|
||||
' const paneKey = process.env.ORCA_PANE_KEY',
|
||||
|
||||
@@ -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> | void
|
||||
|
||||
type FakeCurlChild = {
|
||||
@@ -48,9 +52,16 @@ export type AgentStatusExtensionHarness = {
|
||||
callHook: (name: string, event?: unknown, context?: HookContext) => Promise<void>
|
||||
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<void>
|
||||
// What Pi does for /reload: as above, but the module is evaluated again and `globalThis` survives.
|
||||
reloadPi: () => Promise<void>
|
||||
// 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<typeof fetch>) => Promise<unknown>
|
||||
// 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> }
|
||||
) => void
|
||||
setModel: (model: unknown) => Promise<boolean>
|
||||
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<typeof module.exports.default> => {
|
||||
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<string, HookHandler> = {}
|
||||
const piEvents = new EventEmitter()
|
||||
// Why: Pi calls every handler an extension registers for an event, in registration order.
|
||||
let handlerLists: Record<string, HookHandler[]> = {}
|
||||
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<string, HookHandler>): void => {
|
||||
const registerInto = (
|
||||
target: Record<string, HookHandler>,
|
||||
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()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<typeof setTimeout> | 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',
|
||||
|
||||
@@ -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<string, unknown>; metadata: Record<string, unknown>; 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<typeof setTimeout> | null = null',
|
||||
'type HookPost = { hookEventName: string; extra: Record<string, unknown>; metadata: Record<string, unknown>; 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<typeof setTimeout> | null; postRevision: number; retryTimer: ReturnType<typeof setTimeout> | 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<string, unknown>): 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()',
|
||||
' })',
|
||||
'}',
|
||||
|
||||
@@ -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)
|
||||
})
|
||||
@@ -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<string, SubagentDetail>()',
|
||||
' 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)',
|
||||
' }',
|
||||
''
|
||||
]
|
||||
}
|
||||
@@ -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<void> {
|
||||
await harness.callHook('agent_end', {}, idle)
|
||||
await harness.callHook('agent_settled', undefined, idle)
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
}
|
||||
@@ -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<string>; exited?: Set<string>; 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<string>(), 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<string>; exited?: Set<string>; details?: Map<string, SubagentDetail>; waiting: boolean; ownsPane?: boolean; runGeneration?: number; endedRunGeneration?: number; completionPostedGeneration?: number; parked?: Map<string, Map<string, SubagentDetail>>; onEvent?: (event: unknown, forcedStatus?: string) => void; listener?: (event: unknown) => void; onRunnerExit?: (event: unknown) => void; runnerExitListener?: (event: unknown) => void; runnerExitCheck?: ReturnType<typeof setTimeout> | 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<string, unknown> {',
|
||||
' if (!subagentRoster) return {}',
|
||||
' const subagents: Record<string, unknown>[] = []',
|
||||
' 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<string>(), 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<string, SubagentDetail>())',
|
||||
// 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<typeof setTimeout> | 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<string>()',
|
||||
' 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()",
|
||||
' }'
|
||||
]
|
||||
}
|
||||
|
||||
@@ -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',
|
||||
' })',
|
||||
''
|
||||
]
|
||||
}
|
||||
|
||||
@@ -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 }
|
||||
: {})
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -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'
|
||||
})
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user