feat(subagent): add explicit child reports
This commit is contained in:
@@ -0,0 +1,369 @@
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import { mkdtempSync, rmSync } from 'node:fs'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { Context } from 'cordis'
|
||||
import type { Agent } from '@deepseek-ai/dsh-agent'
|
||||
import AgentLoop from '@deepseek-ai/dsh-agent-loop'
|
||||
import { mountAgentLoopTestDependencies } from '@deepseek-ai/dsh-agent-loop-testkit'
|
||||
import { CallId, LlmAdapter } from '@deepseek-ai/dsh-llm'
|
||||
import type { GenerateOptions, StreamChunk } from '@deepseek-ai/dsh-llm'
|
||||
import { SessionId } from '@deepseek-ai/dsh-session'
|
||||
import type { SessionEvent } from '@deepseek-ai/dsh-session'
|
||||
import JsonlSessionPersistence from '@deepseek-ai/dsh-session-persistence-jsonl'
|
||||
import SubagentService from '@deepseek-ai/dsh-subagent'
|
||||
import * as SubagentSpawn from '@deepseek-ai/dsh-subagent-spawn'
|
||||
import * as control from '@deepseek-ai/dsh-tool-subagent-control'
|
||||
import { textResponse } from '../../../core/agent-loop/tests/mock-adapter.ts'
|
||||
import * as tool from '../src/index.ts'
|
||||
|
||||
const testSignal = new AbortController().signal
|
||||
|
||||
/** Adapter that keeps child Activations resident until released. */
|
||||
class HeldAdapter extends LlmAdapter {
|
||||
readonly requests: GenerateOptions[] = []
|
||||
private readonly gate = Promise.withResolvers<undefined>()
|
||||
|
||||
async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
|
||||
this.requests.push(options)
|
||||
await this.gate.promise
|
||||
for (const chunk of textResponse('held answer')) {
|
||||
if (options.signal?.aborted) throw new Error('aborted')
|
||||
yield chunk
|
||||
}
|
||||
}
|
||||
|
||||
release(): void {
|
||||
this.gate.resolve(undefined)
|
||||
}
|
||||
}
|
||||
|
||||
const cleanups: (() => Promise<void>)[] = []
|
||||
afterEach(async () => {
|
||||
for (const cleanup of cleanups.splice(0).reverse()) await cleanup()
|
||||
})
|
||||
|
||||
/** Boot the real continuation graph with optional report installation. */
|
||||
async function setup(options: { load?: boolean; config?: tool.Config } = {}) {
|
||||
const ctx = new Context()
|
||||
await mountAgentLoopTestDependencies(ctx)
|
||||
const root = mkdtempSync(join(tmpdir(), 'dsh-tool-subagent-report-'))
|
||||
await ctx.plugin(JsonlSessionPersistence, { root })
|
||||
await ctx.plugin(AgentLoop, { agents: [] })
|
||||
await ctx.plugin(SubagentService)
|
||||
await ctx.plugin(SubagentSpawn, { providerName: 'spawn' })
|
||||
const fiber = options.load === false
|
||||
? undefined
|
||||
: await ctx.plugin(tool, options.config ?? { reportDelivery: 'quiet' })
|
||||
const adapter = new HeldAdapter()
|
||||
ctx.llm.registerAdapter(['mock'], adapter)
|
||||
const parent = ctx.agentLoop.create(SessionId('parent'), { provider: 'mock', model: 'mock' })
|
||||
cleanups.push(async () => {
|
||||
adapter.release()
|
||||
await ctx.fiber.dispose()
|
||||
rmSync(root, { recursive: true, force: true })
|
||||
})
|
||||
return { ctx, parent, adapter, fiber }
|
||||
}
|
||||
|
||||
/** Start and resolve one resident continuable child. */
|
||||
async function startChild(ctx: Context, parent: Agent, prompt = 'child task') {
|
||||
const started = await ctx.subagents.startContinuable({
|
||||
provider: 'spawn',
|
||||
label: prompt,
|
||||
request: {
|
||||
prompt: [{ type: 'text', text: prompt }],
|
||||
parent,
|
||||
},
|
||||
signal: testSignal,
|
||||
})
|
||||
const child = await vi.waitFor(() => {
|
||||
const live = ctx.agents.get(started.childId)
|
||||
expect(live).toBeDefined()
|
||||
return live as Agent
|
||||
})
|
||||
return { started, child }
|
||||
}
|
||||
|
||||
let calls = 0
|
||||
function callReport(ctx: Context, child: Agent, output: string, signal = testSignal) {
|
||||
return ctx.tools.execute({
|
||||
signal,
|
||||
callId: CallId(`report-${++calls}`),
|
||||
name: 'report',
|
||||
arguments: { output },
|
||||
agent: child,
|
||||
})
|
||||
}
|
||||
|
||||
/** Reports durably visible in one Agent's Session. */
|
||||
function reports(agent: Agent): { id: string; text: string; sender: string }[] {
|
||||
return agent.session.events.flatMap((event) => {
|
||||
if (event.type !== 'user/message' || event.data.source.kind !== 'subagent-report') return []
|
||||
return [{
|
||||
id: event.data.id,
|
||||
text: event.data.content.flatMap(block => block.type === 'text' ? [block.text] : []).join('\n'),
|
||||
sender: event.data.source.senderSessionId,
|
||||
}]
|
||||
})
|
||||
}
|
||||
|
||||
function renderedText(result: { content: { type: string; text?: string }[] }): string {
|
||||
return result.content.flatMap(block => block.type === 'text' ? [block.text ?? ''] : []).join('')
|
||||
}
|
||||
|
||||
describe('dsh-tool-subagent-report', () => {
|
||||
it('registers report only in continuable child scopes', async () => {
|
||||
const { ctx, parent } = await setup()
|
||||
expect(ctx.tools.schemas().map(schema => schema.name)).not.toContain('report')
|
||||
expect(ctx.tools.schemas(parent).map(schema => schema.name)).not.toContain('report')
|
||||
|
||||
const { child } = await startChild(ctx, parent)
|
||||
const schemas = ctx.tools.schemas(child).filter(schema => schema.name === 'report')
|
||||
expect(schemas).toHaveLength(1)
|
||||
const properties = (schemas[0]?.parameters as { properties: Record<string, unknown> }).properties
|
||||
expect(Object.keys(properties)).toEqual(['output'])
|
||||
})
|
||||
|
||||
it('adds no implicit capability when the package is absent', async () => {
|
||||
const { ctx, parent } = await setup({ load: false })
|
||||
const { child } = await startChild(ctx, parent)
|
||||
expect(ctx.tools.schemas(child).map(schema => schema.name)).not.toContain('report')
|
||||
expect((await callReport(ctx, child, 'missing')).isError).toBe(true)
|
||||
})
|
||||
|
||||
it('does not imply parent controls and survives a global-tool allow-list', async () => {
|
||||
const { ctx, parent } = await setup()
|
||||
expect(ctx.tools.schemas().map(schema => schema.name)).not.toContain('send_message')
|
||||
await ctx.plugin(control)
|
||||
expect(ctx.tools.schemas().map(schema => schema.name)).toContain('send_message')
|
||||
|
||||
const started = await ctx.subagents.startContinuable({
|
||||
provider: 'spawn',
|
||||
label: 'restricted child',
|
||||
request: {
|
||||
prompt: [{ type: 'text', text: 'restricted child' }],
|
||||
parent,
|
||||
toolFilter: { allow: [] },
|
||||
},
|
||||
signal: testSignal,
|
||||
})
|
||||
const child = await vi.waitFor(() => {
|
||||
const live = ctx.agents.get(started.childId)
|
||||
expect(live).toBeDefined()
|
||||
return live as Agent
|
||||
})
|
||||
const names = ctx.tools.schemas(child).map(schema => schema.name)
|
||||
expect(names).toContain('report')
|
||||
expect(names).not.toContain('send_message')
|
||||
})
|
||||
|
||||
it('delivers quiet reports with stable identity and provenance without waking', async () => {
|
||||
const { ctx, parent, adapter } = await setup()
|
||||
const { started, child } = await startChild(ctx, parent)
|
||||
const parentRequests = adapter.requests.filter(request => request.sessionId === parent.id).length
|
||||
const enqueues: string[] = []
|
||||
ctx.on('agent/inbox/enqueue', (agent, item) => {
|
||||
if (agent === parent) enqueues.push(item.placement)
|
||||
})
|
||||
|
||||
const result = await callReport(ctx, child, 'CHILD_FINDING')
|
||||
|
||||
expect(result.isError).toBe(false)
|
||||
if (result.isError) throw new Error('report unexpectedly failed')
|
||||
const messageId = (result.value as { messageId: string }).messageId
|
||||
expect(renderedText(result)).toContain(messageId)
|
||||
expect(reports(parent)).toEqual([{
|
||||
id: messageId,
|
||||
text: `Background subagent ${started.childId} reported:\nCHILD_FINDING`,
|
||||
sender: started.childId,
|
||||
}])
|
||||
expect(enqueues).toEqual([])
|
||||
expect(parent.status).toBe('idle')
|
||||
expect(adapter.requests.filter(request => request.sessionId === parent.id)).toHaveLength(parentRequests)
|
||||
})
|
||||
|
||||
it('queues wakeup reports as one later parent turn', async () => {
|
||||
const { ctx, parent, adapter } = await setup({ config: { reportDelivery: 'wakeup' } })
|
||||
const { child } = await startChild(ctx, parent)
|
||||
const enqueues: string[] = []
|
||||
ctx.on('agent/inbox/enqueue', (agent, item) => {
|
||||
if (agent === parent) enqueues.push(item.placement)
|
||||
})
|
||||
|
||||
const result = await callReport(ctx, child, 'WAKE_UP')
|
||||
expect(result.isError).toBe(false)
|
||||
expect(enqueues).toEqual(['queued'])
|
||||
await vi.waitFor(() => {
|
||||
expect(adapter.requests.some(request => request.sessionId === parent.id)).toBe(true)
|
||||
})
|
||||
})
|
||||
|
||||
it('preserves accepted order across repeated reports', async () => {
|
||||
const { ctx, parent } = await setup()
|
||||
const { child } = await startChild(ctx, parent)
|
||||
|
||||
expect((await callReport(ctx, child, 'FIRST')).isError).toBe(false)
|
||||
expect((await callReport(ctx, child, 'SECOND')).isError).toBe(false)
|
||||
expect(reports(parent).map(report => report.text.split('\n').at(-1))).toEqual(['FIRST', 'SECOND'])
|
||||
})
|
||||
|
||||
it('keeps an accepted report after the child settles', async () => {
|
||||
const { ctx, parent, adapter } = await setup()
|
||||
const { started, child } = await startChild(ctx, parent)
|
||||
expect((await callReport(ctx, child, 'DURABLE_SELECTION')).isError).toBe(false)
|
||||
|
||||
adapter.release()
|
||||
await vi.waitFor(() => { expect(ctx.agents.get(started.childId)).toBeUndefined() })
|
||||
expect(reports(parent).map(report => report.text)).toEqual([
|
||||
`Background subagent ${started.childId} reported:\nDURABLE_SELECTION`,
|
||||
])
|
||||
})
|
||||
|
||||
it('routes nested reports exactly one edge upward', async () => {
|
||||
const { ctx, parent, adapter } = await setup()
|
||||
const { child } = await startChild(ctx, parent, 'outer task')
|
||||
const { started: grandchildStart, child: grandchild } = await startChild(ctx, child, 'inner task')
|
||||
|
||||
expect((await callReport(ctx, grandchild, 'FROM_GRANDCHILD')).isError).toBe(false)
|
||||
expect(reports(parent)).toEqual([])
|
||||
// The intermediate parent's turn is open, so quiet context is staged until
|
||||
// that turn reaches its next safe log boundary.
|
||||
expect(reports(child)).toEqual([])
|
||||
adapter.release()
|
||||
await vi.waitFor(() => { expect(reports(child)).toHaveLength(1) })
|
||||
expect(reports(child)[0]?.sender).toBe(grandchildStart.childId)
|
||||
expect(reports(child)[0]?.text).toContain('FROM_GRANDCHILD')
|
||||
})
|
||||
|
||||
it('accounts wakeup reports delivered to a resident continuable parent', async () => {
|
||||
const { ctx, parent, adapter } = await setup({ config: { reportDelivery: 'wakeup' } })
|
||||
const { child } = await startChild(ctx, parent, 'outer task')
|
||||
const { started: grandchildStart, child: grandchild } = await startChild(ctx, child, 'inner task')
|
||||
|
||||
expect((await callReport(ctx, grandchild, 'WAKE_PARENT_CHILD')).isError).toBe(false)
|
||||
expect(ctx.agents.get(child.id)).toBe(child)
|
||||
|
||||
adapter.release()
|
||||
await vi.waitFor(() => { expect(reports(child)).toHaveLength(1) })
|
||||
expect(reports(child)[0]?.sender).toBe(grandchildStart.childId)
|
||||
expect(reports(child)[0]?.text).toContain('WAKE_PARENT_CHILD')
|
||||
})
|
||||
|
||||
it('normalizes a direct parent send rejection', async () => {
|
||||
const { ctx, parent } = await setup()
|
||||
const { child } = await startChild(ctx, parent)
|
||||
vi.spyOn(parent, 'inject').mockImplementationOnce(() => {
|
||||
throw new Error('parent closed during delivery')
|
||||
})
|
||||
|
||||
await expect(ctx.subagents.reportFrom(child, [{ type: 'text', text: 'rejected' }], {
|
||||
delivery: 'quiet',
|
||||
signal: testSignal,
|
||||
})).rejects.toMatchObject({ code: 'PARENT_UNAVAILABLE' })
|
||||
expect(reports(parent)).toEqual([])
|
||||
})
|
||||
|
||||
it('rejects roots, forged same-id senders, absent parents, cancellation, and drain', async () => {
|
||||
const { ctx, parent, adapter } = await setup()
|
||||
await expect(ctx.subagents.reportFrom(parent, [{ type: 'text', text: 'root' }], {
|
||||
delivery: 'quiet',
|
||||
signal: testSignal,
|
||||
})).rejects.toMatchObject({ code: 'UNAUTHORIZED' })
|
||||
|
||||
const disposable = await ctx.agents.create({
|
||||
sessionId: SessionId('disposable-parent'),
|
||||
agentOptions: { provider: 'mock', model: 'mock' },
|
||||
})
|
||||
const { child } = await startChild(ctx, disposable.agent)
|
||||
const forged = { ...child } as Agent
|
||||
await expect(ctx.subagents.reportFrom(forged, [{ type: 'text', text: 'forged' }], {
|
||||
delivery: 'quiet',
|
||||
signal: testSignal,
|
||||
})).rejects.toMatchObject({ code: 'UNAUTHORIZED' })
|
||||
|
||||
const aborted = new AbortController()
|
||||
aborted.abort()
|
||||
expect((await callReport(ctx, child, 'cancelled', aborted.signal)).isError).toBe(true)
|
||||
|
||||
await disposable.dispose()
|
||||
expect((await callReport(ctx, child, 'orphaned')).isError).toBe(true)
|
||||
|
||||
adapter.release()
|
||||
const draining = ctx.subagents.drainContinuableDescendants([child])
|
||||
await expect(ctx.subagents.reportFrom(child, [{ type: 'text', text: 'draining' }], {
|
||||
delivery: 'quiet',
|
||||
signal: testSignal,
|
||||
})).rejects.toMatchObject({ code: 'DRAINING' })
|
||||
await draining
|
||||
})
|
||||
|
||||
it('revokes resident installations and defers later grants to the next Activation', async () => {
|
||||
const { ctx, parent, fiber } = await setup()
|
||||
const { child } = await startChild(ctx, parent)
|
||||
expect(ctx.tools.schemas(child).map(schema => schema.name)).toContain('report')
|
||||
|
||||
await fiber?.dispose()
|
||||
expect(ctx.tools.schemas(child).map(schema => schema.name)).not.toContain('report')
|
||||
expect((await callReport(ctx, child, 'revoked')).isError).toBe(true)
|
||||
|
||||
const late = await ctx.plugin(tool, { reportDelivery: 'quiet' })
|
||||
expect(ctx.tools.schemas(child).map(schema => schema.name)).not.toContain('report')
|
||||
await late.dispose()
|
||||
})
|
||||
|
||||
it('rolls back materialization when a setup contribution revokes itself', async () => {
|
||||
const { ctx, parent } = await setup({ load: false })
|
||||
const self: { revoke?: () => void } = {}
|
||||
self.revoke = ctx.subagents.registerContinuableSetup((childCtx) => {
|
||||
const dispose = childCtx.tools.register({
|
||||
name: 'racing-report',
|
||||
description: 'racing setup',
|
||||
parameters: { type: 'object', properties: {} },
|
||||
output: { schema: { type: 'object', properties: {} }, render: () => [] },
|
||||
execute: () => Promise.resolve({}),
|
||||
})
|
||||
self.revoke?.()
|
||||
return dispose
|
||||
})
|
||||
|
||||
await expect(ctx.subagents.startContinuable({
|
||||
provider: 'spawn',
|
||||
label: 'racing child',
|
||||
request: {
|
||||
prompt: [{ type: 'text', text: 'racing child' }],
|
||||
parent,
|
||||
},
|
||||
signal: testSignal,
|
||||
})).rejects.toMatchObject({ code: 'ACTIVATION_SETUP_REVOKED' })
|
||||
expect(ctx.agents.list().map(agent => agent.id)).toEqual([parent.id])
|
||||
})
|
||||
|
||||
it('keeps the namespace plugin shape and validates its default', () => {
|
||||
expect('default' in tool).toBe(false)
|
||||
expect(tool.name).toBe('tool-subagent-report')
|
||||
expect(tool.inject).toEqual(['subagents', 'tools'])
|
||||
expect(tool.Config({}).reportDelivery).toBe('quiet')
|
||||
expect(() => tool.Config({ reportDelivery: 'shout' } as never)).toThrow()
|
||||
})
|
||||
})
|
||||
|
||||
/** Prove report delivery uses ordinary logged user messages. */
|
||||
function userTexts(events: readonly SessionEvent[]): string[] {
|
||||
return events.flatMap(event => event.type === 'user/message'
|
||||
? event.data.content.flatMap(block => block.type === 'text' ? [block.text] : [])
|
||||
: [])
|
||||
}
|
||||
|
||||
describe('dsh-tool-subagent-report result independence', () => {
|
||||
it('does not report a final assistant answer automatically or create Tasks', async () => {
|
||||
const { ctx, parent, adapter } = await setup()
|
||||
const { started } = await startChild(ctx, parent)
|
||||
adapter.release()
|
||||
await vi.waitFor(() => { expect(ctx.agents.get(started.childId)).toBeUndefined() })
|
||||
|
||||
expect(reports(parent)).toEqual([])
|
||||
expect(userTexts((await ctx.sessionPersistence.load(started.childId)).events)).toEqual(['child task'])
|
||||
expect(ctx.get('tasks')).toBeUndefined()
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user