feat(todo): make the parallel in_progress policy configurable
Whether concurrent active tasks are legitimate depends on runtime concurrency the tool cannot observe, but whether a deployment's agents ever fan out is knowable at composition time. `allowParallelInProgress` (default true) therefore replaces the hardcoded policy: the flag moves the model-facing instruction and the accepted input together, so a deployment running strictly sequential agents can restore the single-active discipline from cordis.yml. The durable-log invariant does not follow the flag. A log written while parallel work was allowed must still replay after a deployment tightens the policy, so the invariant stays silent on the active count.
This commit is contained in:
@@ -0,0 +1,126 @@
|
||||
// Proves `allowParallelInProgress` is real configurability and not a constant:
|
||||
// the flag is set in a cordis.yml booted through the real Loader, and both faces
|
||||
// it controls — the model-facing description and the accepted input — follow it.
|
||||
import { mkdtemp, rm, writeFile } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { pathToFileURL } from 'node:url'
|
||||
import { afterEach, describe, expect, it } from 'vitest'
|
||||
import { Context } from 'cordis'
|
||||
import Loader from '@cordisjs/plugin-loader'
|
||||
import Include from '@cordisjs/plugin-include'
|
||||
import { CallId } from '@deepseek-ai/dsh-llm'
|
||||
import { Session, SessionId } from '@deepseek-ai/dsh-session'
|
||||
import AgentRegistry from '@deepseek-ai/dsh-agent'
|
||||
import type { Agent } from '@deepseek-ai/dsh-agent'
|
||||
import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
|
||||
import ToolRegistry from '@deepseek-ai/dsh-tools'
|
||||
import * as ToolTodo from '@deepseek-ai/dsh-tool-todo'
|
||||
|
||||
let root: string | undefined
|
||||
let context: Context | undefined
|
||||
|
||||
afterEach(async () => {
|
||||
await context?.fiber.dispose()
|
||||
context = undefined
|
||||
if (root !== undefined) await rm(root, { recursive: true, force: true })
|
||||
root = undefined
|
||||
})
|
||||
|
||||
function agent(ctx: Context): Agent {
|
||||
const scope = ctx.plugin(() => {})
|
||||
const id = SessionId('todo-loader-agent')
|
||||
const value: Agent = {
|
||||
id, options: {}, session: new Session(id), status: 'idle', acceptsNextStep: false, ctx: scope.ctx,
|
||||
followup: () => {}, steer: () => {}, inject: () => {}, send: () => {}, cancel() {}, whenIdle: () => Promise.resolve(),
|
||||
}
|
||||
ctx.agents.register(value)
|
||||
return value
|
||||
}
|
||||
|
||||
function resultText(result: { content: { type: string; text?: string }[] }): string {
|
||||
return result.content.filter(block => block.type === 'text').map(block => block.text).join('')
|
||||
}
|
||||
|
||||
/**
|
||||
* Boot a cordis.yml carrying the given tool-todo config block.
|
||||
* @param configLines - YAML lines nested under the tool's `config:` key.
|
||||
* @returns the booted context.
|
||||
*/
|
||||
async function boot(configLines: readonly string[]): Promise<Context> {
|
||||
root = await mkdtemp(join(tmpdir(), 'dsh-todo-loader-'))
|
||||
const configPath = join(root, 'cordis.yml')
|
||||
await writeFile(configPath, [
|
||||
"- name: '@deepseek-ai/dsh-agent'",
|
||||
"- name: '@deepseek-ai/dsh-system-prompt'",
|
||||
"- name: '@deepseek-ai/dsh-tools'",
|
||||
"- name: '@deepseek-ai/dsh-tool-todo'",
|
||||
...configLines.length > 0 ? [' config:', ...configLines] : [],
|
||||
'',
|
||||
].join('\n'))
|
||||
|
||||
const ctx = new Context()
|
||||
ctx.baseUrl = pathToFileURL(root).href + '/'
|
||||
await ctx.plugin(Loader)
|
||||
ctx.loader.builtins.include = Include
|
||||
const modules = new Map<string, unknown>([
|
||||
['@deepseek-ai/dsh-agent', AgentRegistry],
|
||||
['@deepseek-ai/dsh-system-prompt', SystemPrompt],
|
||||
['@deepseek-ai/dsh-tools', ToolRegistry],
|
||||
['@deepseek-ai/dsh-tool-todo', ToolTodo],
|
||||
])
|
||||
ctx.loader.internal = {
|
||||
version: 'v2',
|
||||
async import(specifier: string) {
|
||||
if (!modules.has(specifier)) throw new Error(`unexpected Loader import: ${specifier}`)
|
||||
return modules.get(specifier)
|
||||
},
|
||||
} as unknown as NonNullable<typeof ctx.loader.internal>
|
||||
await ctx.loader.create({ name: 'cordis:include', config: { path: pathToFileURL(configPath).href } })
|
||||
await ctx.loader.await()
|
||||
context = ctx
|
||||
return ctx
|
||||
}
|
||||
|
||||
const PARALLEL_TODOS = [
|
||||
{ content: 'run subagent a', status: 'in_progress' },
|
||||
{ content: 'run subagent b', status: 'in_progress' },
|
||||
]
|
||||
|
||||
describe('tool-todo real Loader composition through cordis.yml', () => {
|
||||
it('allowParallelInProgress: false narrows the description and rejects a parallel write', async () => {
|
||||
const ctx = await boot([' allowParallelInProgress: false'])
|
||||
const description = ctx.tools.schemas().find(s => s.name === 'todo_write')?.description ?? ''
|
||||
expect(description).toContain('Keep AT MOST ONE todo `in_progress`')
|
||||
expect(description).not.toContain('several at once')
|
||||
|
||||
const owner = agent(ctx)
|
||||
const result = await ctx.tools.execute({
|
||||
signal: new AbortController().signal,
|
||||
callId: CallId('parallel'),
|
||||
name: 'todo_write',
|
||||
arguments: { todos: PARALLEL_TODOS },
|
||||
agent: owner,
|
||||
})
|
||||
expect(result.isError).toBe(true)
|
||||
expect(resultText(result)).toContain('at most one task may be in_progress')
|
||||
expect(owner.session.events.some(e => e.type === 'todo/write')).toBe(false)
|
||||
}, 30_000)
|
||||
|
||||
it('the omitted default keeps the parallel policy end to end', async () => {
|
||||
const ctx = await boot([])
|
||||
const description = ctx.tools.schemas().find(s => s.name === 'todo_write')?.description ?? ''
|
||||
expect(description).toContain('several at once when work genuinely runs in parallel')
|
||||
|
||||
const owner = agent(ctx)
|
||||
const result = await ctx.tools.execute({
|
||||
signal: new AbortController().signal,
|
||||
callId: CallId('parallel-default'),
|
||||
name: 'todo_write',
|
||||
arguments: { todos: PARALLEL_TODOS },
|
||||
agent: owner,
|
||||
})
|
||||
expect(result.isError).toBe(false)
|
||||
expect(owner.session.events.findLast(e => e.type === 'todo/write')?.data.todos).toEqual(PARALLEL_TODOS)
|
||||
}, 30_000)
|
||||
})
|
||||
@@ -26,11 +26,11 @@ function agentWithSession(id = 'parent-1'): Agent & { session: Session } {
|
||||
return { id: SessionId(id), session } as unknown as Agent & { session: Session }
|
||||
}
|
||||
|
||||
async function setup(): Promise<Context> {
|
||||
async function setup(config: tool.Config = {}): Promise<Context> {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SystemPrompt)
|
||||
await ctx.plugin(ToolRegistry)
|
||||
await ctx.plugin(tool)
|
||||
await ctx.plugin(tool, config)
|
||||
return ctx
|
||||
}
|
||||
|
||||
@@ -140,6 +140,62 @@ describe('dsh-tool-todo', () => {
|
||||
expect(agent.session.events.findLast(e => e.type === 'todo/write')!.data.todos).toEqual(todos)
|
||||
})
|
||||
|
||||
describe('allowParallelInProgress: false', () => {
|
||||
const parallel = [
|
||||
{ content: 'run subagent a', status: 'in_progress' },
|
||||
{ content: 'run subagent b', status: 'in_progress' },
|
||||
]
|
||||
|
||||
it('rejects a call marking several items in_progress', async () => {
|
||||
const ctx = await setup({ allowParallelInProgress: false })
|
||||
const agent = agentWithSession('single-active')
|
||||
const result = await callTodo(ctx, { todos: parallel }, { agent })
|
||||
expect(result.isError).toBe(true)
|
||||
expect(text(result)).toContain('at most one task may be in_progress')
|
||||
// A rejected call must not reach the durable log.
|
||||
expect(agent.session.events.some(e => e.type === 'todo/write')).toBe(false)
|
||||
})
|
||||
|
||||
it('still accepts one active item', async () => {
|
||||
const ctx = await setup({ allowParallelInProgress: false })
|
||||
const todos: TodoItem[] = [
|
||||
{ content: 'run subagent a', status: 'in_progress' },
|
||||
{ content: 'run subagent b', status: 'pending' },
|
||||
]
|
||||
const result = await callTodo(ctx, { todos })
|
||||
expect(result.isError).toBe(false)
|
||||
})
|
||||
|
||||
it('an explicit true accepts a parallel write, like the omitted default', async () => {
|
||||
const ctx = await setup({ allowParallelInProgress: true })
|
||||
const result = await callTodo(ctx, { todos: parallel })
|
||||
expect(result.isError).toBe(false)
|
||||
})
|
||||
|
||||
it('defaults to parallel for a direct apply, which bypasses the schema default', async () => {
|
||||
// Composing through ctx.plugin lets schemastery fill the field; a caller
|
||||
// invoking apply() itself hands over a config object with it absent, so
|
||||
// the policy default has to hold on that path too.
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SystemPrompt)
|
||||
await ctx.plugin(ToolRegistry)
|
||||
tool.apply(ctx, {})
|
||||
const result = await callTodo(ctx, { todos: parallel })
|
||||
expect(result.isError).toBe(false)
|
||||
})
|
||||
|
||||
it('instructs the model to keep at most one active, and the default instructs parallel', async () => {
|
||||
const single = await setup({ allowParallelInProgress: false })
|
||||
const singleDesc = single.tools.schemas().find(s => s.name === 'todo_write')!.description
|
||||
expect(singleDesc).toContain('Keep AT MOST ONE todo `in_progress`')
|
||||
expect(singleDesc).not.toContain('several at once')
|
||||
|
||||
const parallelDesc = (await setup()).tools.schemas().find(s => s.name === 'todo_write')!.description
|
||||
expect(parallelDesc).toContain('several at once when work genuinely runs in parallel')
|
||||
expect(parallelDesc).not.toContain('AT MOST ONE')
|
||||
})
|
||||
})
|
||||
|
||||
it.each([
|
||||
{ label: 'empty content', todos: [{ content: ' ', status: 'pending' }], fragment: 'non-empty' },
|
||||
{ label: 'duplicate content', todos: [{ content: 'dup', status: 'pending' }, { content: 'dup', status: 'completed' }], fragment: 'duplicate' },
|
||||
|
||||
Reference in New Issue
Block a user