docs(subagent): update package READMEs for the activation lifecycle
Rewrites the service API table, authority-versus-provenance contract, residency routing, and deferred-work list; scopes the in-process driver README to one-shot runs; and restates both model-facing tools' outputs, which no longer carry a task id.
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { Context } from 'cordis'
|
||||
import { createUserMessage, CallId, type ContentBlock, type GenerateOptions } from '@deepseek-ai/dsh-llm'
|
||||
import { createUserMessage, CallId, type ContentBlock, type GenerateOptions } from '@deepseek-ai/dsh-llm'
|
||||
import { SessionId } from '@deepseek-ai/dsh-session'
|
||||
import AgentLoop from '@deepseek-ai/dsh-agent-loop'
|
||||
import { mountAgentLoopTestDependencies } from '@deepseek-ai/dsh-agent-loop-testkit'
|
||||
@@ -120,28 +120,6 @@ describe('in-process structured output', () => {
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('confirmed steering rejects delivery once the structured result is captured', async () => {
|
||||
const { ctx, parent } = await setup([
|
||||
toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 7 }),
|
||||
])
|
||||
// oxlint-disable-next-line prefer-const -- single assignment follows listener registration so pre-fulfillment events remain guardable.
|
||||
let run: Awaited<ReturnType<typeof ctx.subagents.start>> | undefined
|
||||
let delivery: Promise<void> | undefined
|
||||
ctx.on('session/event', (session, event) => {
|
||||
if (session.header.parentSession === undefined || run === undefined
|
||||
|| event.type !== 'tool/result' || delivery !== undefined) return
|
||||
delivery = run.steer?.([{ type: 'text', text: 'one more thing' }], { kind: 'user' })
|
||||
void delivery?.catch(() => undefined)
|
||||
})
|
||||
run = await ctx.subagents.start('spawn', structuredRequest(parent))
|
||||
const result = await run.result
|
||||
if (delivery === undefined) throw new Error('structured result did not submit steering')
|
||||
await expect(delivery)
|
||||
.rejects.toThrow(/already reported its structured result; the message was not delivered/)
|
||||
expect(result.structured).toEqual({ answer: 7 })
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('denies tool calls that FOLLOW the capture in the same response — terminal means terminal', async () => {
|
||||
// One model response carrying structured_output FIRST and a side-effecting
|
||||
// call after it: the continuation veto only fires at step end, so without
|
||||
|
||||
@@ -2,17 +2,16 @@ import { createUserMessage } from '@deepseek-ai/dsh-llm'
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { Context } from 'cordis'
|
||||
import { type Agent, type AgentOptions } from '@deepseek-ai/dsh-agent'
|
||||
import { Session, SessionId } from '@deepseek-ai/dsh-session'
|
||||
import { SessionId } from '@deepseek-ai/dsh-session'
|
||||
import AgentLoop from '@deepseek-ai/dsh-agent-loop'
|
||||
import { mountAgentLoopTestDependencies } from '@deepseek-ai/dsh-agent-loop-testkit'
|
||||
import InvariantService from '@deepseek-ai/dsh-invariants'
|
||||
import * as SessionInvariant from '@deepseek-ai/dsh-session/invariant'
|
||||
import * as AgentInvariant from '@deepseek-ai/dsh-agent/invariant'
|
||||
import * as AgentLoopInvariant from '@deepseek-ai/dsh-agent-loop/invariant'
|
||||
import SubagentService, { SUBAGENT_DESCRIPTOR_VERSION, SubagentError } from '@deepseek-ai/dsh-subagent'
|
||||
import { defineContentToolFixture } from '@deepseek-ai/dsh-tools'
|
||||
import { maxTokensResponse, MockAdapter, textResponse, toolCallResponse } from '../../../core/agent-loop/tests/mock-adapter.ts'
|
||||
import { resumeInProcessRun, startInProcessRun } from '../src/index.ts'
|
||||
import SubagentService from '@deepseek-ai/dsh-subagent'
|
||||
import { maxTokensResponse, MockAdapter, textResponse } from '../../../core/agent-loop/tests/mock-adapter.ts'
|
||||
import { startInProcessRun } from '../src/index.ts'
|
||||
|
||||
type Script = ConstructorParameters<typeof MockAdapter>[0]
|
||||
|
||||
@@ -39,22 +38,6 @@ function request(parent: Agent, signal = new AbortController().signal) {
|
||||
return { prompt: [{ type: 'text' as const, text: 'child task' }], parent, signal }
|
||||
}
|
||||
|
||||
function continuableRequest(parent: Agent) {
|
||||
const sessionId = SessionId('continuable-child')
|
||||
return {
|
||||
...request(parent),
|
||||
continuation: {
|
||||
sessionId,
|
||||
descriptor: {
|
||||
version: SUBAGENT_DESCRIPTOR_VERSION,
|
||||
provider: 'spawn',
|
||||
agentProvider: 'mock',
|
||||
agentModel: 'mock',
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
function text(blocks: readonly { type: string; text?: string }[]): string {
|
||||
return blocks.filter(block => block.type === 'text').map(block => block.text).join('')
|
||||
}
|
||||
@@ -88,110 +71,7 @@ describe('startInProcessRun', () => {
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('rejects a continuable child when no durability listener is registered', async () => {
|
||||
const { parent } = await setup([textResponse('driver answer')])
|
||||
|
||||
const run = await startInProcessRun(continuableRequest(parent), {})
|
||||
const caught: unknown = await run.result.catch((error: unknown) => error)
|
||||
|
||||
expect(caught).toBeInstanceOf(SubagentError)
|
||||
const durabilityError = caught as SubagentError
|
||||
expect(durabilityError.code).toBe('DURABILITY_FAILED')
|
||||
expect(durabilityError.message).toContain('required durability checkpoint has no registered listener')
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('rejects when the durability listener disappears before final confirmation', async () => {
|
||||
const { ctx, parent } = await setup([textResponse('driver answer')])
|
||||
let flushes = 0
|
||||
let detach = (): void => {}
|
||||
detach = ctx.on('session/flush', (session) => {
|
||||
if (session.header.parentSession === undefined) return
|
||||
flushes++
|
||||
if (flushes === 1) detach()
|
||||
})
|
||||
|
||||
const run = await startInProcessRun(continuableRequest(parent), {})
|
||||
const caught: unknown = await run.result.catch((error: unknown) => error)
|
||||
|
||||
expect(caught).toBeInstanceOf(SubagentError)
|
||||
const durabilityError = caught as SubagentError
|
||||
expect(durabilityError.code).toBe('DURABILITY_FAILED')
|
||||
expect(durabilityError.message).toContain('required durability checkpoint has no registered listener')
|
||||
expect(flushes).toBe(1)
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('requires a final durability checkpoint for a continuable child', async () => {
|
||||
const { ctx, parent } = await setup([textResponse('driver answer')])
|
||||
const failure = new Error('disk full')
|
||||
let flushes = 0
|
||||
ctx.on('session/flush', (session) => {
|
||||
if (session.header.parentSession === undefined) return
|
||||
flushes++
|
||||
throw failure
|
||||
})
|
||||
|
||||
const run = await startInProcessRun(continuableRequest(parent), {})
|
||||
const caught: unknown = await run.result.catch((error: unknown) => error)
|
||||
expect(caught).toBeInstanceOf(SubagentError)
|
||||
const durabilityError = caught as SubagentError
|
||||
expect(durabilityError.code).toBe('DURABILITY_FAILED')
|
||||
expect(durabilityError.cause).toBe(failure)
|
||||
expect(durabilityError.message).toContain(
|
||||
'the latest child state was not confirmed persisted and may be unavailable or stale on resume: disk full',
|
||||
)
|
||||
expect(flushes).toBe(2)
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('completes a continuable child when the final checkpoint retries a transient flush failure', async () => {
|
||||
const { ctx, parent } = await setup([textResponse('driver answer')])
|
||||
let flushes = 0
|
||||
ctx.on('session/flush', (session) => {
|
||||
if (session.header.parentSession === undefined) return
|
||||
flushes++
|
||||
if (flushes === 1) throw new Error('temporary append failure')
|
||||
})
|
||||
|
||||
const run = await startInProcessRun(continuableRequest(parent), {})
|
||||
await expect(run.result).resolves.toMatchObject({ stopReason: 'completed' })
|
||||
expect(flushes).toBe(2)
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it.each([
|
||||
{ checkpoint: 'succeeds', failure: undefined },
|
||||
{ checkpoint: 'fails', failure: new Error('disk full') },
|
||||
])('lets cancellation own the result when the final durability checkpoint $checkpoint', async ({ failure }) => {
|
||||
const { ctx, parent } = await setup([textResponse('driver answer')])
|
||||
const checkpointStarted = Promise.withResolvers<undefined>()
|
||||
const releaseCheckpoint = Promise.withResolvers<undefined>()
|
||||
let flushes = 0
|
||||
ctx.on('session/flush', async (session) => {
|
||||
if (session.header.parentSession === undefined) return
|
||||
flushes++
|
||||
if (flushes !== 2) return
|
||||
checkpointStarted.resolve(undefined)
|
||||
await releaseCheckpoint.promise
|
||||
if (failure !== undefined) throw failure
|
||||
})
|
||||
const controller = new AbortController()
|
||||
|
||||
const run = await startInProcessRun({
|
||||
...continuableRequest(parent),
|
||||
signal: controller.signal,
|
||||
}, {})
|
||||
await checkpointStarted.promise
|
||||
controller.abort()
|
||||
releaseCheckpoint.resolve(undefined)
|
||||
|
||||
await expect(run.result).resolves.toMatchObject({ stopReason: 'aborted' })
|
||||
expect(flushes).toBe(2)
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('keeps foreground runs best-effort when their turn checkpoint fails', async () => {
|
||||
it('does not add a final durability checkpoint to a foreground run', async () => {
|
||||
const { ctx, parent } = await setup([textResponse('driver answer')])
|
||||
let flushes = 0
|
||||
ctx.on('session/flush', (session) => {
|
||||
@@ -336,69 +216,18 @@ describe('startInProcessRun', () => {
|
||||
expect(ctx.sessions.list()).toHaveLength(beforeSessions)
|
||||
})
|
||||
|
||||
it('rejects an already-aborted resume before publication', async () => {
|
||||
const { parent } = await setup([])
|
||||
const controller = new AbortController()
|
||||
controller.abort('too late')
|
||||
await expect(resumeInProcessRun({
|
||||
sessionId: SessionId('resumed-child'),
|
||||
prompt: [{ type: 'text', text: 'continue' }],
|
||||
source: { kind: 'user' },
|
||||
parent,
|
||||
signal: controller.signal,
|
||||
descriptor: { version: SUBAGENT_DESCRIPTOR_VERSION, provider: 'spawn' },
|
||||
})).rejects.toThrow('aborted before child publication')
|
||||
})
|
||||
|
||||
it('resumes without inventing undeclared agent model options', async () => {
|
||||
const childId = SessionId('resumed-child')
|
||||
let flushes = 0
|
||||
const child = {
|
||||
id: childId,
|
||||
options: {},
|
||||
session: new Session(childId),
|
||||
status: 'idle',
|
||||
acceptsNextStep: false,
|
||||
ctx: {
|
||||
sessions: {
|
||||
flush: () => {
|
||||
flushes++
|
||||
return Promise.resolve(true)
|
||||
},
|
||||
},
|
||||
} as unknown as Context,
|
||||
send(): void {},
|
||||
reserveTurnAdmission: () => undefined,
|
||||
updateInbox: () => 'not-found',
|
||||
followup(): void {},
|
||||
steer() { return { outcome: Promise.resolve({ status: 'rejected' as const }) } },
|
||||
inject(): void {},
|
||||
cancel(): void {},
|
||||
whenIdle: () => Promise.resolve(),
|
||||
} as Agent
|
||||
let resumedOptions: unknown
|
||||
const parent = {
|
||||
ctx: {
|
||||
agents: {
|
||||
resume: (options: { agentOptions: unknown }) => {
|
||||
resumedOptions = options.agentOptions
|
||||
return Promise.resolve({ agent: child, dispose: () => Promise.resolve() })
|
||||
},
|
||||
},
|
||||
},
|
||||
} as unknown as Agent
|
||||
|
||||
const run = await resumeInProcessRun({
|
||||
sessionId: childId,
|
||||
prompt: [{ type: 'text', text: 'continue' }],
|
||||
source: { kind: 'plugin', plugin: 'test-coordinator' },
|
||||
parent,
|
||||
signal: new AbortController().signal,
|
||||
descriptor: { version: SUBAGENT_DESCRIPTOR_VERSION, provider: 'spawn' },
|
||||
})
|
||||
expect(resumedOptions).toEqual({})
|
||||
it('stamps only the resolved depth when neither parent nor request declares a model route', async () => {
|
||||
// The one-shot analogue of the deleted resume coverage ("resumes without
|
||||
// inventing undeclared agent model options"): a bare parent with no request
|
||||
// agentOptions yields a child whose options carry ONLY the stamped depth —
|
||||
// no provider/model is fabricated, so the child's turn errors for want of a
|
||||
// route rather than silently adopting one.
|
||||
const { ctx } = await setup([])
|
||||
const parent = ctx.agentLoop.create(SessionId('routeless-parent'), {})
|
||||
const run = await startInProcessRun(request(parent), {})
|
||||
const child = ctx.agents.get(run.id)!
|
||||
expect(child.options).toEqual({ subagentDepth: 1 })
|
||||
await expect(run.result).resolves.toMatchObject({ stopReason: 'error' })
|
||||
expect(flushes).toBe(1)
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
@@ -461,104 +290,4 @@ describe('startInProcessRun', () => {
|
||||
expect(ctx.agents.list()).toHaveLength(beforeAgents)
|
||||
expect(ctx.sessions.list()).toHaveLength(beforeSessions)
|
||||
})
|
||||
|
||||
it('confirmed steering rejects a settled child instead of queueing an untracked turn', async () => {
|
||||
const { ctx, parent } = await setup([textResponse('done')])
|
||||
const run = await startInProcessRun(request(parent), {})
|
||||
await run.result
|
||||
await expect(run.steer!([{ type: 'text', text: 'late' }], { kind: 'user' }))
|
||||
.rejects.toThrow(/not running; the message was not delivered/)
|
||||
const child = ctx.agents.get(run.id)!
|
||||
expect(child.session.events.some(event => event.type === 'steering/message')).toBe(false)
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('confirmed steering rejects when a concluding tool prevents request admission', async () => {
|
||||
const { ctx, parent } = await setup([toolCallResponse('c1', 'finalize', {})])
|
||||
const enteredTool = Promise.withResolvers<undefined>()
|
||||
const releaseTool = Promise.withResolvers<undefined>()
|
||||
ctx.tools.register(defineContentToolFixture({
|
||||
name: 'finalize',
|
||||
description: 'Finish the child run.',
|
||||
parameters: {},
|
||||
async execute(_args, exec) {
|
||||
enteredTool.resolve(undefined)
|
||||
await releaseTool.promise
|
||||
exec.concludeTurn()
|
||||
return [{ type: 'text', text: 'final' }]
|
||||
},
|
||||
}))
|
||||
const run = await startInProcessRun(request(parent), {})
|
||||
const child = ctx.agents.get(run.id)!
|
||||
await enteredTool.promise
|
||||
|
||||
const delivery = run.steer!([{ type: 'text', text: 'terminal race' }], { kind: 'user' })
|
||||
releaseTool.resolve(undefined)
|
||||
await expect(delivery).rejects.toThrow(/stopped before steering admission; the message was not delivered/)
|
||||
await run.result
|
||||
expect(child.session.events.some(event => event.type === 'steering/message')).toBe(false)
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('confirmed steering fulfills only after the next request snapshot admits it', async () => {
|
||||
const { ctx, parent, adapter } = await setup([textResponse('first'), textResponse('second')])
|
||||
const enteredStopping = Promise.withResolvers<undefined>()
|
||||
const releaseStopping = Promise.withResolvers<undefined>()
|
||||
let held = false
|
||||
ctx.on('agent/turn-stopping', (agent) => {
|
||||
if (agent.session.header.parentSession === undefined || held) return
|
||||
held = true
|
||||
enteredStopping.resolve(undefined)
|
||||
return releaseStopping.promise
|
||||
})
|
||||
|
||||
const run = await startInProcessRun(request(parent), {})
|
||||
const child = ctx.agents.get(run.id)!
|
||||
await enteredStopping.promise
|
||||
|
||||
let settled = false
|
||||
const delivery = run.steer!([{ type: 'text', text: 'after the first step' }], { kind: 'user' })
|
||||
.then(() => { settled = true })
|
||||
await Promise.resolve()
|
||||
expect(settled).toBe(false)
|
||||
releaseStopping.resolve(undefined)
|
||||
await delivery
|
||||
|
||||
const result = await run.result
|
||||
expect(adapter.requests).toHaveLength(2)
|
||||
expect(JSON.stringify(adapter.requests[1]?.messages)).toContain('after the first step')
|
||||
expect((result.output[0] as { text?: string }).text).toBe('second')
|
||||
const steering = child.session.events.find(event => event.type === 'steering/message')
|
||||
expect(steering?.type === 'steering/message' && steering.data.message.source).toEqual({ kind: 'user' })
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('carries steering from a non-terminal flush window into a tracked next turn', async () => {
|
||||
const { ctx, parent, adapter } = await setup([textResponse('first'), textResponse('second')])
|
||||
const enteredFlush = Promise.withResolvers<undefined>()
|
||||
const releaseFlush = Promise.withResolvers<undefined>()
|
||||
let held = false
|
||||
ctx.on('session/flush', (session) => {
|
||||
if (session.header.parentSession === undefined || held) return
|
||||
if (!session.events.some(event => event.type === 'turn/end')) return
|
||||
held = true
|
||||
enteredFlush.resolve(undefined)
|
||||
return releaseFlush.promise
|
||||
})
|
||||
|
||||
const run = await startInProcessRun(request(parent), {})
|
||||
const child = ctx.agents.get(run.id)!
|
||||
await enteredFlush.promise
|
||||
expect(child.status).toBe('running')
|
||||
|
||||
const delivery = run.steer!([{ type: 'text', text: 'next tracked turn' }], { kind: 'user' })
|
||||
releaseFlush.resolve(undefined)
|
||||
await delivery
|
||||
const result = await run.result
|
||||
expect(adapter.requests).toHaveLength(2)
|
||||
expect(child.session.events.filter(event => event.type === 'turn/start')).toHaveLength(2)
|
||||
expect(child.session.events.some(event => event.type === 'steering/message')).toBe(false)
|
||||
expect((result.output[0] as { text?: string }).text).toBe('second')
|
||||
await run.dispose()
|
||||
})
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user