From c98a4d09493350b546f8449295c15f5938aecda2 Mon Sep 17 00:00:00 2001 From: pku-xht Date: Tue, 11 Aug 2026 17:12:53 +0800 Subject: [PATCH] refactor(tasks): simplify admission configuration --- .../background-task-admission/session.jsonl | 2 +- .../examples/acp-demo/tests/acp-agent.spec.ts | 18 +++++++++++++--- .../agent-spine-demo/tests/agent-core.spec.ts | 19 +++++++++++++---- packages/tasks/tasks-local/src/index.ts | 21 +++++++------------ .../tests/loader-composition.spec.ts | 18 ++++++++++++++-- .../tasks/tasks-local/tests/tasks.spec.ts | 16 ++++---------- 6 files changed, 58 insertions(+), 36 deletions(-) diff --git a/examples/acp-agent/tests/snapshots/background-task-admission/session.jsonl b/examples/acp-agent/tests/snapshots/background-task-admission/session.jsonl index 0a303e1866..d9d40b2539 100644 --- a/examples/acp-agent/tests/snapshots/background-task-admission/session.jsonl +++ b/examples/acp-agent/tests/snapshots/background-task-admission/session.jsonl @@ -25,7 +25,7 @@ {"type":"assistant/chunk","seq":23,"time":1786434813605,"data":{"turn":1,"step":2,"chunk":{"type":"finish","reason":{"kind":"tool-calls"}}}} {"type":"assistant/message","seq":24,"time":1786434813605,"data":{"turn":1,"step":2,"message":{"role":"assistant","content":[{"type":"tool-call","id":"bounded-task-second","name":"bash","arguments":"{\"command\":\"printf SHOULD_NOT_RUN > second-task-ran.txt; while :; do sleep 60; done\",\"description\":\"Attempt a second background task\",\"run_in_background\":true}"}],"source":{"kind":"model","provider":"deepseek-official","model":"deepseek-v4-flash"},"id":"48c909b3-5651-462f-b0d3-09198d119a2f"},"usage":{"inputTokens":10,"outputTokens":5}},"sourceEventSeqs":[19,20,21,22,23],"surfaceOp":"append"} {"type":"tool/call","seq":25,"time":1786434813606,"data":{"turn":1,"step":2,"callId":"bounded-task-second","name":"bash","arguments":"{\"command\":\"printf SHOULD_NOT_RUN > second-task-ran.txt; while :; do sleep 60; done\",\"description\":\"Attempt a second background task\",\"run_in_background\":true}"}} -{"type":"tool/result","seq":26,"time":1786434813609,"data":{"turn":1,"step":2,"message":{"source":{"kind":"tool","callId":"bounded-task-second"},"content":[{"type":"tool-result","toolCallId":"bounded-task-second","content":[{"type":"text","text":"Error: background task limit reached for this owner (1/1 active); use task_kill to stop an unneeded task, wait for it to finish, then retry"}],"isError":true}],"role":"user","id":"982adb5c-13ef-44c5-a53e-b726c986ac68"}},"sourceEventSeqs":[25],"surfaceOp":"append"} +{"type":"tool/result","seq":26,"time":1786434813609,"data":{"turn":1,"step":2,"message":{"source":{"kind":"tool","callId":"bounded-task-second"},"content":[{"type":"tool-result","toolCallId":"bounded-task-second","content":[{"type":"text","text":"Error: background task limit reached for this owner (limit: 1); use task_kill to stop an unneeded task, wait for it to finish, then retry"}],"isError":true}],"role":"user","id":"c0386bdf-df3c-4d2b-af8e-04ee28682214"}},"sourceEventSeqs":[25],"surfaceOp":"append"} {"type":"step/end","seq":27,"time":1786434813609,"data":{"turn":1,"step":2}} {"type":"step/start","seq":28,"time":1786434813614,"data":{"turn":1,"step":3}} {"type":"assistant/chunk","seq":29,"time":1786434813618,"data":{"turn":1,"step":3,"chunk":{"type":"block-start","index":0,"blockType":"tool-call"}}} diff --git a/packages/examples/acp-demo/tests/acp-agent.spec.ts b/packages/examples/acp-demo/tests/acp-agent.spec.ts index 65479f0489..3a64404704 100644 --- a/packages/examples/acp-demo/tests/acp-agent.spec.ts +++ b/packages/examples/acp-demo/tests/acp-agent.spec.ts @@ -186,12 +186,24 @@ describe('dsh-acp-demo composition', () => { const ctx = await mount({ provider: 'mock', model: 'mock', - tasks: { maxConcurrentTasksPerOwner: 2 }, + tasks: { maxConcurrentTasksPerOwner: 1 }, skills: await isolatedSkillsConfig(), workspaceContext: false, }) - expect((ctx.tasks as unknown as { config: { maxConcurrentTasksPerOwner: number } }) - .config.maxConcurrentTasksPerOwner).toBe(2) + let settle!: (outcome: { status: 'killed' }) => void + ctx.tasks.start({ + kind: 'bash', + label: 'hold configured slot', + run: () => ({ + cancel: () => { settle({ status: 'killed' }) }, + done: new Promise((resolve) => { settle = resolve }), + }), + }) + expect(() => ctx.tasks.start({ + kind: 'bash', + label: 'blocked configured task', + run: () => ({ cancel: () => {}, done: Promise.resolve({ status: 'completed' }) }), + })).toThrow('(limit: 1)') await ctx.fiber.dispose() }) diff --git a/packages/examples/agent-spine-demo/tests/agent-core.spec.ts b/packages/examples/agent-spine-demo/tests/agent-core.spec.ts index d3b02ae582..9f738e3b56 100644 --- a/packages/examples/agent-spine-demo/tests/agent-core.spec.ts +++ b/packages/examples/agent-spine-demo/tests/agent-core.spec.ts @@ -10,7 +10,6 @@ import { agentEvents, type Agent } from '@deepseek-ai/dsh-agent' import { SessionId } from '@deepseek-ai/dsh-session' import LocalBashExecutor from '@deepseek-ai/dsh-bash-local' import LocalFileSystem from '@deepseek-ai/dsh-fs-local' -import LocalTaskService from '@deepseek-ai/dsh-tasks-local' import * as ToolFs from '@deepseek-ai/dsh-tool-fs' import { MockAdapter, textResponse, toolCallResponse } from '../../../core/agent-loop/tests/mock-adapter.ts' import { @@ -158,7 +157,6 @@ describe('dsh-agent-spine-demo bundle', () => { expect(ctx.get('skills')).toBeDefined() expect(ctx.get('agents')).toBeDefined() expect(ctx.get('tasks')).toBeDefined() - expect((ctx.tasks as LocalTaskService).config.maxConcurrentTasksPerOwner).toBe(10) expect(ctx.get('invariants')).toBeDefined() expect(ctx.get('agentLoop')).toBeDefined() expect(ctx.get('goals')).toBeUndefined() @@ -304,10 +302,23 @@ describe('dsh-agent-spine-demo bundle', () => { it('forwards task admission config to the process-local provider', async () => { const ctx = await mount({ - tasks: { maxConcurrentTasksPerOwner: 3 }, + tasks: { maxConcurrentTasksPerOwner: 1 }, workspaceContext: false, }) - expect((ctx.tasks as LocalTaskService).config.maxConcurrentTasksPerOwner).toBe(3) + let settle!: (outcome: { status: 'killed' }) => void + ctx.tasks.start({ + kind: 'probe', + label: 'hold configured slot', + run: () => ({ + cancel: () => { settle({ status: 'killed' }) }, + done: new Promise((resolve) => { settle = resolve }), + }), + }) + expect(() => ctx.tasks.start({ + kind: 'probe', + label: 'blocked configured task', + run: () => ({ cancel: () => {}, done: Promise.resolve({ status: 'completed' }) }), + })).toThrow('(limit: 1)') await ctx.fiber.dispose() }) diff --git a/packages/tasks/tasks-local/src/index.ts b/packages/tasks/tasks-local/src/index.ts index c4f90808a6..cf60c416fc 100644 --- a/packages/tasks/tasks-local/src/index.ts +++ b/packages/tasks/tasks-local/src/index.ts @@ -33,9 +33,6 @@ export interface Config { maxConcurrentTasksPerOwner?: number } -/** Configuration after defaults and load-time validation. */ -type ResolvedConfig = Required - /** The registry's mutable per-task record (never handed out — see {@link LocalTaskService.snapshot}). */ interface TrackedTask { id: TaskId @@ -97,8 +94,8 @@ export class LocalTaskService extends TaskService { .default(DEFAULT_MAX_CONCURRENT_TASKS_PER_OWNER), }) - /** Validated registry configuration. */ - readonly config: ResolvedConfig + /** Schemastery-defaulted active-task limit. */ + private readonly maxConcurrentTasksPerOwner: number private store = new Map() private counters = new Map() /** @@ -120,14 +117,10 @@ export class LocalTaskService extends TaskService { /** Service context used by detached settlement continuations and teardown. */ private readonly selfCtx: Context - constructor(ctx: Context, config: Config = {}) { + constructor(ctx: Context, config: Config) { super(ctx) - const maxConcurrentTasksPerOwner = config.maxConcurrentTasksPerOwner - ?? DEFAULT_MAX_CONCURRENT_TASKS_PER_OWNER - if (!Number.isSafeInteger(maxConcurrentTasksPerOwner) || maxConcurrentTasksPerOwner <= 0) { - throw new TypeError('tasks-local: maxConcurrentTasksPerOwner must be a positive safe integer') - } - this.config = { maxConcurrentTasksPerOwner } + // Schemastery validates and fills the default before constructing the service. + this.maxConcurrentTasksPerOwner = (config as Required).maxConcurrentTasksPerOwner this.selfCtx = ctx ctx.effect(() => () => this.disposeAll(), 'tasks teardown') } @@ -145,9 +138,9 @@ export class LocalTaskService extends TaskService { if (spec.owner !== undefined) this.ensureOwnerCleanup(spec.owner) const active = this.activeTaskCount(spec.owner) - if (active >= this.config.maxConcurrentTasksPerOwner) { + if (active >= this.maxConcurrentTasksPerOwner) { throw new Error( - `background task limit reached for this owner (${active}/${this.config.maxConcurrentTasksPerOwner} active); use task_kill to stop an unneeded task, wait for it to finish, then retry`, + `background task limit reached for this owner (limit: ${this.maxConcurrentTasksPerOwner}); use task_kill to stop an unneeded task, wait for it to finish, then retry`, ) } diff --git a/packages/tasks/tasks-local/tests/loader-composition.spec.ts b/packages/tasks/tasks-local/tests/loader-composition.spec.ts index 0cd0cac575..9d4b09b6ff 100644 --- a/packages/tasks/tasks-local/tests/loader-composition.spec.ts +++ b/packages/tasks/tasks-local/tests/loader-composition.spec.ts @@ -25,7 +25,7 @@ describe('tasks-local through a real Loader composition', () => { await writeFile(configPath, [ "- name: '@deepseek-ai/dsh-tasks-local'", ' config:', - ' maxConcurrentTasksPerOwner: 2', + ' maxConcurrentTasksPerOwner: 1', '', ].join('\n')) @@ -47,6 +47,20 @@ describe('tasks-local through a real Loader composition', () => { await context.loader.await() expect(context.tasks).toBeInstanceOf(LocalTaskService) - expect((context.tasks as LocalTaskService).config.maxConcurrentTasksPerOwner).toBe(2) + context.tasks.attachController('loader-test') + let settle!: (outcome: { status: 'killed' }) => void + context.tasks.start({ + kind: 'bash', + label: 'hold loader slot', + run: () => ({ + cancel: () => { settle({ status: 'killed' }) }, + done: new Promise((resolve) => { settle = resolve }), + }), + }) + expect(() => context!.tasks.start({ + kind: 'bash', + label: 'blocked loader task', + run: () => ({ cancel: () => {}, done: Promise.resolve({ status: 'completed' }) }), + })).toThrow('(limit: 1)') }) }) diff --git a/packages/tasks/tasks-local/tests/tasks.spec.ts b/packages/tasks/tasks-local/tests/tasks.spec.ts index 12cf6f917c..239fca4fb3 100644 --- a/packages/tasks/tasks-local/tests/tasks.spec.ts +++ b/packages/tasks/tasks-local/tests/tasks.spec.ts @@ -172,35 +172,27 @@ describe('LocalTaskService.start', () => { const ctx = new Context() await expect(ctx.plugin(LocalTaskService, { maxConcurrentTasksPerOwner })) .rejects.toThrow() - expect(() => new LocalTaskService(new Context(), { maxConcurrentTasksPerOwner })) - .toThrow('maxConcurrentTasksPerOwner must be a positive safe integer') }, ) it('accepts the largest safe integer limit', async () => { const ctx = await harness({ maxConcurrentTasksPerOwner: Number.MAX_SAFE_INTEGER }) - expect((ctx.tasks as LocalTaskService).config.maxConcurrentTasksPerOwner) - .toBe(Number.MAX_SAFE_INTEGER) + expect(ctx.tasks).toBeInstanceOf(LocalTaskService) }) it('defaults each owner bucket to ten active tasks', async () => { const ctx = await harness() - expect((ctx.tasks as LocalTaskService).config.maxConcurrentTasksPerOwner).toBe(10) const live = Array.from({ length: 10 }, () => producer()) for (const task of live) ctx.tasks.start(task.spec) const blocked = producer() const run = vi.fn(() => blocked.spec.run()) expect(() => ctx.tasks.start({ ...blocked.spec, run })) - .toThrow('background task limit reached for this owner (10/10 active)') + .toThrow('background task limit reached for this owner (limit: 10)') expect(run).not.toHaveBeenCalled() for (const task of live) task.settle({ status: 'completed' }) }) - it('defaults direct construction when the config schema is bypassed', () => { - expect(new LocalTaskService(new Context()).config.maxConcurrentTasksPerOwner).toBe(10) - }) - it('rejects before producer start and id allocation, then admits immediately after settlement', async () => { const ctx = await harness({ maxConcurrentTasksPerOwner: 1 }) const first = producer() @@ -224,7 +216,7 @@ describe('LocalTaskService.start', () => { expect(ctx.tasks.kill(id)).toBe('requested') const replacement = producer() - expect(() => ctx.tasks.start(replacement.spec)).toThrow('(1/1 active)') + expect(() => ctx.tasks.start(replacement.spec)).toThrow('(limit: 1)') first.settle({ status: 'killed' }) await tick() @@ -260,7 +252,7 @@ describe('LocalTaskService.start', () => { expect(() => ctx.tasks.start(producer({ owner: replacement }).spec)).not.toThrow() ctx.tasks.start(producer().spec) - expect(() => ctx.tasks.start(producer().spec)).toThrow('(1/1 active)') + expect(() => ctx.tasks.start(producer().spec)).toThrow('(limit: 1)') expect(() => ctx.tasks.start(producer({ owner: oldOwner }).spec)) .toThrow('is not the registered agent instance')