feat(web): rewrite subagent conversations for FIFO activation
This commit is contained in:
@@ -71,7 +71,14 @@ describe('subagent gateway', () => {
|
||||
expect(response.rpcId).toBe('subagent-rpc')
|
||||
expect(response.result).toMatchObject({
|
||||
ok: true,
|
||||
value: { parentAvailable: false, entries: [{ kind: 'child' }, { kind: 'diagnostic' }] },
|
||||
value: {
|
||||
parentAvailable: false,
|
||||
entries: [
|
||||
{ kind: 'child', mode: 'continuable' },
|
||||
{ kind: 'child', mode: 'one-shot' },
|
||||
{ kind: 'diagnostic' },
|
||||
],
|
||||
},
|
||||
})
|
||||
expect(listChildren).toHaveBeenCalledWith(PARENT, undefined)
|
||||
})
|
||||
@@ -79,7 +86,7 @@ describe('subagent gateway', () => {
|
||||
it('reads a healthy direct child without looking up or activating any Agent', async () => {
|
||||
const { api, getAgent, readSession } = bench()
|
||||
const response = await api.subagents.history(request({
|
||||
parentSessionId: PARENT, childSessionId: CHILD, maxMessages: 10,
|
||||
parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', maxMessages: 10,
|
||||
}))
|
||||
expect(response.result).toMatchObject({
|
||||
ok: true,
|
||||
@@ -89,12 +96,26 @@ describe('subagent gateway', () => {
|
||||
expect(getAgent).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('reads one-shot history and rejects an address with the wrong mode', async () => {
|
||||
const oneShot = {
|
||||
kind: 'child', id: CHILD, mode: 'one-shot', label: 'batch', activity: 'inactive',
|
||||
}
|
||||
const { api, readSession } = bench({ entries: [oneShot] })
|
||||
expect((await api.subagents.history(request({
|
||||
parentSessionId: PARENT, childSessionId: CHILD, mode: 'one-shot',
|
||||
}))).result).toMatchObject({ ok: true })
|
||||
expect((await api.subagents.history(request({
|
||||
parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
|
||||
}))).result).toMatchObject({ ok: false, error: { code: 'subagent-not-found' } })
|
||||
expect(readSession).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('rejects a diagnostic address before reading history', async () => {
|
||||
const { api, readSession } = bench({ entries: [
|
||||
{ kind: 'diagnostic', id: CHILD, reason: 'unsupported' },
|
||||
] })
|
||||
const response = await api.subagents.history(request({
|
||||
parentSessionId: PARENT, childSessionId: CHILD,
|
||||
parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
|
||||
}))
|
||||
expect(response.result).toMatchObject({
|
||||
ok: false,
|
||||
@@ -109,30 +130,36 @@ describe('subagent gateway', () => {
|
||||
it('routes human content through the exact live parent with rpc attribution', async () => {
|
||||
const { api, parent, followup } = bench()
|
||||
const content = [{ type: 'text' as const, text: '继续' }]
|
||||
const signal = new AbortController().signal
|
||||
const response = await api.subagents.prompt(request({
|
||||
parentSessionId: PARENT, childSessionId: CHILD, content,
|
||||
}))
|
||||
parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content,
|
||||
}), signal)
|
||||
expect(response.result).toMatchObject({
|
||||
ok: true, value: { messageId: 'message-1' },
|
||||
})
|
||||
expect(followup).toHaveBeenCalledTimes(1)
|
||||
const [actualParent, actualChild, actualContent, delivery] = followup.mock.calls[0]!
|
||||
expect([actualParent, actualChild, actualContent]).toEqual([parent, CHILD, content])
|
||||
expect(delivery.source).toEqual({ kind: 'user', rpcId: RpcId('subagent-rpc') })
|
||||
expect(delivery.signal).toBeInstanceOf(AbortSignal)
|
||||
expect(followup).toHaveBeenCalledWith(
|
||||
parent,
|
||||
CHILD,
|
||||
content,
|
||||
{ source: { kind: 'user', rpcId: RpcId('subagent-rpc') }, signal },
|
||||
)
|
||||
})
|
||||
|
||||
it('fails before delivery when the parent is absent and maps continuation failures', async () => {
|
||||
const absent = bench({ parentLive: false })
|
||||
expect((await absent.api.subagents.prompt(request({
|
||||
parentSessionId: PARENT, childSessionId: CHILD, content: [],
|
||||
}))).result).toMatchObject({ ok: false, error: { code: 'subagent-parent-unavailable' } })
|
||||
parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content: [],
|
||||
}), new AbortController().signal)).result).toMatchObject({
|
||||
ok: false, error: { code: 'subagent-parent-unavailable' },
|
||||
})
|
||||
expect(absent.listChildren).not.toHaveBeenCalled()
|
||||
|
||||
const failed = bench({ followupError: new SubagentError('not delivered', 'DRAINING') })
|
||||
const failed = bench({ followupError: new SubagentError('draining', 'DRAINING') })
|
||||
expect((await failed.api.subagents.prompt(request({
|
||||
parentSessionId: PARENT, childSessionId: CHILD, content: [],
|
||||
}))).result).toMatchObject({ ok: false, error: { code: 'subagent-not-delivered' } })
|
||||
parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content: [],
|
||||
}), new AbortController().signal)).result).toMatchObject({
|
||||
ok: false, error: { code: 'subagent-delivery-unavailable' },
|
||||
})
|
||||
})
|
||||
|
||||
it('maps history disappearance and hides unexpected backend details', async () => {
|
||||
@@ -140,7 +167,7 @@ describe('subagent gateway', () => {
|
||||
readError: new SessionQueryError('secret path', 'SESSION_QUERY_SESSION_NOT_FOUND'),
|
||||
})
|
||||
expect((await disappeared.api.subagents.history(request({
|
||||
parentSessionId: PARENT, childSessionId: CHILD,
|
||||
parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
|
||||
}))).result).toMatchObject({
|
||||
ok: false,
|
||||
error: {
|
||||
@@ -160,8 +187,8 @@ describe('subagent gateway', () => {
|
||||
|
||||
const prompt = bench({ followupError: new Error('secret provider') })
|
||||
expect((await prompt.api.subagents.prompt(request({
|
||||
parentSessionId: PARENT, childSessionId: CHILD, content: [],
|
||||
}))).result).toMatchObject({
|
||||
parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content: [],
|
||||
}), new AbortController().signal)).result).toMatchObject({
|
||||
ok: false,
|
||||
error: { code: 'internal', message: 'subagent prompt failed' },
|
||||
})
|
||||
|
||||
@@ -110,7 +110,18 @@ function fakeApi(overrides: Partial<{ muxFrames: MuxFrame[]; hostFrames: HostFra
|
||||
async history(request) {
|
||||
return { rpcId: request.rpcId, result: { ok: true, value: { events: [], hasMore: false } } }
|
||||
},
|
||||
async prompt(request) {
|
||||
async prompt(request, signal) {
|
||||
if (request.payload.content.some(block => block.type === 'text' && block.text === 'hang')) {
|
||||
if (!signal.aborted) {
|
||||
await new Promise<void>((resolve) => {
|
||||
signal.addEventListener('abort', () => { resolve() }, { once: true })
|
||||
})
|
||||
}
|
||||
return {
|
||||
rpcId: request.rpcId,
|
||||
result: { ok: false, error: { code: 'cancelled' as const, message: 'aborted', details: {} } },
|
||||
}
|
||||
}
|
||||
return {
|
||||
rpcId: request.rpcId,
|
||||
result: { ok: true, value: { messageId: 'message-1' as never } },
|
||||
@@ -400,6 +411,23 @@ describe('unary round trip (handler ⇄ client, no network)', () => {
|
||||
}
|
||||
})
|
||||
|
||||
it('round-trips the subagent domain through the wire form', async () => {
|
||||
const c = client()
|
||||
expect((await c.subagents.list({ parentSessionId: 'parent' as never })).result)
|
||||
.toEqual({ ok: true, value: { entries: [], parentAvailable: false } })
|
||||
expect((await c.subagents.history({
|
||||
parentSessionId: 'parent' as never,
|
||||
childSessionId: 'child' as never,
|
||||
mode: 'one-shot',
|
||||
})).result).toEqual({ ok: true, value: { events: [], hasMore: false } })
|
||||
expect((await c.subagents.prompt({
|
||||
parentSessionId: 'parent' as never,
|
||||
childSessionId: 'child' as never,
|
||||
mode: 'continuable',
|
||||
content: [],
|
||||
})).result).toEqual({ ok: true, value: { messageId: 'message-1' } })
|
||||
})
|
||||
|
||||
it('keeps caller and connection aborts on command.execute', async () => {
|
||||
const api = fakeApi()
|
||||
const started = Promise.withResolvers<AbortSignal>()
|
||||
@@ -465,6 +493,34 @@ describe('unary round trip (handler ⇄ client, no network)', () => {
|
||||
expect(parsed.result.error?.code).toBe('cancelled')
|
||||
})
|
||||
|
||||
it('propagates the carrier Request signal into subagent.prompt', async () => {
|
||||
const handler = toFetchHandler(fakeApi())
|
||||
const controller = new AbortController()
|
||||
const body = JSON.stringify({
|
||||
type: 'client-request',
|
||||
rpcId: 'r-subagent-sig',
|
||||
method: 'subagent.prompt',
|
||||
payload: {
|
||||
parentSessionId: 'parent',
|
||||
childSessionId: 'child',
|
||||
mode: 'continuable',
|
||||
content: [{ type: 'text', text: 'hang' }],
|
||||
},
|
||||
})
|
||||
const pending = handler.fetch(new Request(
|
||||
'http://x/api/subagent.prompt',
|
||||
{ method: 'POST', headers: { 'content-type': 'application/json' }, body, signal: controller.signal },
|
||||
))
|
||||
controller.abort()
|
||||
const response = await pending
|
||||
const parsed = await response.json() as {
|
||||
rpcId: string
|
||||
result: { error?: { code: string } }
|
||||
}
|
||||
expect(parsed.rpcId).toBe('r-subagent-sig')
|
||||
expect(parsed.result.error?.code).toBe('cancelled')
|
||||
})
|
||||
|
||||
it('propagates the carrier Request signal into host.pickDirectory', async () => {
|
||||
const api = fakeApi()
|
||||
api.host.pickDirectory = async (request, signal) => {
|
||||
|
||||
@@ -85,7 +85,7 @@ describe('rpcErrorSchema', () => {
|
||||
expect(rpcErrorSchema.parse({ code: 'subagent-catalog-diagnostic', message: 'm', details: { parentSessionId: 'p', childSessionId: 'c', reason: 'corrupt' } }).code).toBe('subagent-catalog-diagnostic')
|
||||
expect(rpcErrorSchema.parse({ code: 'subagent-not-resumable', message: 'm', details: { childSessionId: 'c' } }).code).toBe('subagent-not-resumable')
|
||||
expect(rpcErrorSchema.parse({ code: 'subagent-unauthorized', message: 'm', details: { childSessionId: 'c' } }).code).toBe('subagent-unauthorized')
|
||||
expect(rpcErrorSchema.parse({ code: 'subagent-not-delivered', message: 'm', details: { childSessionId: 'c' } }).code).toBe('subagent-not-delivered')
|
||||
expect(rpcErrorSchema.parse({ code: 'subagent-delivery-unavailable', message: 'm', details: { childSessionId: 'c' } }).code).toBe('subagent-delivery-unavailable')
|
||||
expect(rpcErrorSchema.parse({ code: 'internal', message: 'm', details: {} }).code).toBe('internal')
|
||||
})
|
||||
|
||||
@@ -291,29 +291,34 @@ describe('sessions domain schemas', () => {
|
||||
|
||||
describe('subagent domain schemas', () => {
|
||||
it('validates the direct catalog and addressed history pair', () => {
|
||||
const child = { kind: 'child', id: 'c', label: 'worker', activity: 'running' }
|
||||
const child = {
|
||||
kind: 'child', id: 'c', mode: 'continuable', label: 'worker', activity: 'running',
|
||||
}
|
||||
const oneShot = { kind: 'child', id: 'o', mode: 'one-shot', activity: 'inactive' }
|
||||
const diagnostic = { kind: 'diagnostic', id: 'bad', reason: 'unsupported' }
|
||||
expect(subagentListEntrySchema.parse(child)).toEqual(child)
|
||||
expect(subagentListEntrySchema.parse(oneShot)).toEqual(oneShot)
|
||||
expect(subagentListEntrySchema.parse(diagnostic)).toEqual(diagnostic)
|
||||
expect(subagentListRequestSchema.parse({ parentSessionId: 'p' })).toEqual({ parentSessionId: 'p' })
|
||||
expect(subagentListValueSchema.parse({
|
||||
entries: [child, diagnostic], parentAvailable: true,
|
||||
}).entries).toHaveLength(2)
|
||||
entries: [child, oneShot, diagnostic], parentAvailable: true,
|
||||
}).entries).toHaveLength(3)
|
||||
expect(subagentHistoryRequestSchema.parse({
|
||||
parentSessionId: 'p', childSessionId: 'c', beforeSeq: 4, maxMessages: 2,
|
||||
parentSessionId: 'p', childSessionId: 'c', mode: 'continuable', beforeSeq: 4, maxMessages: 2,
|
||||
}).beforeSeq).toBe(4)
|
||||
expect(() => subagentHistoryRequestSchema.parse({
|
||||
parentSessionId: 'p', childSessionId: 'c', maxMessages: 0,
|
||||
parentSessionId: 'p', childSessionId: 'c', mode: 'continuable', maxMessages: 0,
|
||||
})).toThrow()
|
||||
expect(subagentHistoryValueSchema.parse({ events: [], hasMore: false }).hasMore).toBe(false)
|
||||
})
|
||||
|
||||
it('validates prompt content and the accepted inbox identity', () => {
|
||||
it('validates continuable prompt content and the accepted inbox identity', () => {
|
||||
expect(subagentPromptRequestSchema.parse({
|
||||
parentSessionId: 'p', childSessionId: 'c', content: [{ type: 'text', text: '继续' }],
|
||||
parentSessionId: 'p', childSessionId: 'c', mode: 'continuable',
|
||||
content: [{ type: 'text', text: '继续' }],
|
||||
}).childSessionId).toBe('c')
|
||||
expect(subagentPromptValueSchema.parse({ messageId: 'm1' }).messageId).toBe('m1')
|
||||
expect(() => subagentPromptValueSchema.parse({ taskId: 't1' })).toThrow()
|
||||
expect(() => subagentPromptValueSchema.parse({ route: 'started', taskId: 't2' })).toThrow()
|
||||
})
|
||||
})
|
||||
|
||||
|
||||
Reference in New Issue
Block a user