From 2b9ea06555cde5c4dfc5b76705a2c85683575e7e Mon Sep 17 00:00:00 2001 From: li872 <177593733+li872@users.noreply.github.com> Date: Mon, 14 Sep 2026 15:22:26 +0000 Subject: [PATCH] fix(core): honor AbortSignal during agent-team sleep waits Closes #85 Idle inbox waits and claim-slot backoff ignored cancellation, so Ctrl+C during delegation/handoff delays blocked until delayMs elapsed. Match agent-loop abort cleanup and wire the team abort signal into those call sites. --- packages/core/src/agent/agent-team.test.ts | 271 +++++++++++++++++++ packages/core/src/agent/agent-team.ts | 127 +++++++-- skills/builtin/src/agent-team-plugin.test.ts | 82 ++++++ skills/builtin/src/agent-team-plugin.ts | 1 + 4 files changed, 452 insertions(+), 29 deletions(-) create mode 100644 packages/core/src/agent/agent-team.test.ts create mode 100644 skills/builtin/src/agent-team-plugin.test.ts diff --git a/packages/core/src/agent/agent-team.test.ts b/packages/core/src/agent/agent-team.test.ts new file mode 100644 index 00000000..96eb00f7 --- /dev/null +++ b/packages/core/src/agent/agent-team.test.ts @@ -0,0 +1,271 @@ +import { describe, it, expect, vi } from "vitest"; +import { + AgentTeam, + type AgentTeamInboxStore, + type TeamMessage, +} from "./agent-team.js"; +import { AgentHarnessFactory } from "./harness.js"; +import type { MemoryConfig } from "./conversation-memory.js"; +import type { AgentRunConfig, CompletionResponse } from "@step-cli/protocol"; +import { createMutableRef } from "@step-cli/utils/mutable-ref.js"; + +class MemoryInboxStore implements AgentTeamInboxStore { + private readonly messages: TeamMessage[] = []; + + async append(message: TeamMessage): Promise { + this.messages.push(message); + } + + async read(inboxName: string, sessionId?: string): Promise { + return this.messages.filter((message) => { + if (message.to !== inboxName) { + return false; + } + if (!sessionId) { + return true; + } + return message.sessionId === sessionId; + }); + } +} + +function makeMemoryConfig(): MemoryConfig { + return { + maxContextTokens: 128_000, + reserveOutputTokens: 4096, + minRecentMessages: 4, + compressionTriggerRatio: 0.85, + compressionTargetRatio: 0.6, + maxSummaryChars: 2000, + compactedUserMessageTokenBudget: 2000, + maxCompactedUserMessages: 5, + compactedUserMessageMaxChars: 500, + maxDecisionEntries: 20, + decisionEntryMaxChars: 200, + microCompactKeepRecentToolMessages: 10, + microCompactToolContentChars: 2000, + }; +} + +function makeRunConfig(): AgentRunConfig { + return { + maxSteps: 4, + temperature: 0, + maxContextTokens: 128_000, + maxOutputTokens: 4096, + minOutputTokens: 256, + outputTokenSafetyMargin: 512, + parallelToolCalls: true, + maxToolCallsPerStep: 5, + repeatedToolCallLimit: 3, + maxToolResultCharsInContext: 25_000, + modelRequestRetries: 0, + toolExecutionRetries: 0, + }; +} + +function assistantReply(content: string): CompletionResponse { + return { + id: "cmpl-test", + object: "chat.completion", + created: 0, + model: "test-model", + choices: [ + { + index: 0, + message: { role: "assistant", content }, + finish_reason: "stop", + }, + ], + }; +} + +function createTeam(store = new MemoryInboxStore()) { + const factory = new AgentHarnessFactory({ + model: "test-model", + client: { + createChatCompletion: vi.fn().mockResolvedValue(assistantReply("done")), + }, + defaultSystemPrompt: "You are a teammate.", + memoryConfig: makeMemoryConfig(), + runConfig: makeRunConfig(), + commandTimeoutMs: 1_000, + commandOutputLimit: 1_000, + plugins: [], + interactionProfile: { surface: "headless", canAskUser: false }, + }); + const harnessFactoryRef = createMutableRef( + "AgentHarnessFactory", + ); + harnessFactoryRef.set(factory); + + return { + team: new AgentTeam({ + inboxStore: store, + harnessFactoryRef, + }), + store, + factory, + }; +} + +async function waitFor( + predicate: () => boolean, + timeoutMs = 3_000, +): Promise { + const startedAt = Date.now(); + while (Date.now() - startedAt < timeoutMs) { + if (predicate()) { + return; + } + await new Promise((resolve) => setTimeout(resolve, 15)); + } + throw new Error("Timed out waiting for condition"); +} + +describe("AgentTeam sleep abort", () => { + it("rejects promptly when the caller aborts during readInbox wait", async () => { + const { team } = createTeam(); + const controller = new AbortController(); + const startedAt = Date.now(); + const pending = team.readInbox({ + inboxName: "lead", + reader: "lead", + waitMs: 5_000, + signal: controller.signal, + }); + + setTimeout(() => controller.abort("Run interrupted by user."), 30); + + await expect(pending).rejects.toThrow("Run interrupted by user."); + expect(Date.now() - startedAt).toBeLessThan(1_000); + }); + + it("rejects immediately when the signal is already aborted", async () => { + const { team } = createTeam(); + const controller = new AbortController(); + controller.abort("Run interrupted by user."); + const startedAt = Date.now(); + + await expect( + team.readInbox({ + inboxName: "lead", + reader: "lead", + waitMs: 5_000, + signal: controller.signal, + }), + ).rejects.toThrow("Run interrupted by user."); + expect(Date.now() - startedAt).toBeLessThan(250); + }); + + it("removes the abort listener when sleep completes normally", async () => { + const { team } = createTeam(); + const add = vi.spyOn(AbortSignal.prototype, "addEventListener"); + const remove = vi.spyOn(AbortSignal.prototype, "removeEventListener"); + + try { + const result = await team.readInbox({ + inboxName: "lead", + reader: "lead", + waitMs: 80, + }); + + expect(result.messages).toEqual([]); + const abortAdds = add.mock.calls.filter(([type]) => type === "abort"); + const abortRemoves = remove.mock.calls.filter( + ([type]) => type === "abort", + ); + expect(abortAdds.length).toBeGreaterThan(0); + expect(abortRemoves.length).toBeGreaterThanOrEqual(abortAdds.length); + } finally { + add.mockRestore(); + remove.mockRestore(); + } + }); + + it("wakes idle worker inbox waits when the team closes", async () => { + const { team } = createTeam(); + await team.spawnTeammate({ + name: "researcher", + role: "researcher", + prompt: "Look this up", + requester: "lead", + parentId: "main", + parentDepth: 0, + workspaceRoot: "/tmp/workspace", + }); + + await waitFor(() => team.getTeammate("researcher")?.status === "idle"); + + const startedAt = Date.now(); + await team.close({ + abortRunning: true, + reason: "Agent team shutting down.", + }); + expect(Date.now() - startedAt).toBeLessThan(500); + }); + + it("interrupts worker claim-slot backoff when the team closes", async () => { + const store = new MemoryInboxStore(); + let modelCalls = 0; + const hangFactory = new AgentHarnessFactory({ + model: "test-model", + client: { + createChatCompletion: vi.fn().mockImplementation(() => { + modelCalls += 1; + if (modelCalls === 1) { + return Promise.resolve(assistantReply("spawn complete")); + } + return new Promise(() => {}); + }), + }, + defaultSystemPrompt: "You are a teammate.", + memoryConfig: makeMemoryConfig(), + runConfig: makeRunConfig(), + commandTimeoutMs: 1_000, + commandOutputLimit: 1_000, + plugins: [], + interactionProfile: { surface: "headless", canAskUser: false }, + }); + const harnessFactoryRef = createMutableRef( + "AgentHarnessFactory", + ); + harnessFactoryRef.set(hangFactory); + const team = new AgentTeam({ + inboxStore: store, + harnessFactoryRef, + }); + + await team.spawnTeammate({ + name: "coder", + role: "coder", + prompt: "Start work", + requester: "lead", + parentId: "main", + parentDepth: 0, + workspaceRoot: "/tmp/workspace", + }); + await waitFor(() => team.getTeammate("coder")?.status === "idle"); + + const hangingTurn = team.runTeammateTurn("coder", "Keep working"); + await waitFor(() => team.getTeammate("coder")?.status === "working"); + + await team.sendMessage({ + from: "lead", + to: "coder", + content: "Follow-up assignment", + sessionId: team.getTeammate("coder")?.sessionId, + }); + + // Worker idle wait is 800ms; after that claim-slot backoff sleeps up to 200ms. + await new Promise((resolve) => setTimeout(resolve, 850)); + + const startedAt = Date.now(); + await team.close({ + abortRunning: true, + reason: "Run interrupted by user.", + }); + expect(Date.now() - startedAt).toBeLessThan(300); + void hangingTurn.catch(() => undefined); + }); +}); diff --git a/packages/core/src/agent/agent-team.ts b/packages/core/src/agent/agent-team.ts index b0e27013..78f679eb 100644 --- a/packages/core/src/agent/agent-team.ts +++ b/packages/core/src/agent/agent-team.ts @@ -166,6 +166,7 @@ export class AgentTeam { private readonly activeRuns = new Map(); private readonly shutdownRequests = new Map(); private readonly planRequests = new Map(); + private readonly abortController = new AbortController(); private cursors: Record> = {}; private closePromise: Promise | null = null; @@ -325,6 +326,7 @@ export class AgentTeam { markRead?: boolean; limit?: number; waitMs?: number; + signal?: AbortSignal; }): Promise { const inboxName = normalizeInboxName(input.inboxName); const reader = normalizeInboxName(input.reader); @@ -333,27 +335,41 @@ export class AgentTeam { const limit = clamp(input.limit ?? MAX_READ_LIMIT, 1, MAX_READ_LIMIT); const waitMs = clamp(input.waitMs ?? 0, 0, 60_000); const deadline = Date.now() + waitMs; + const combined = combineAbortSignals( + input.signal, + this.abortController.signal, + ); - while (true) { - const messages = await this.loadInboxMessages(inboxName, sessionId); - const start = this.getCursor(reader, inboxName, sessionId); - const available = messages.slice(start); - - if (available.length > 0 || waitMs === 0 || Date.now() >= deadline) { - const selected = available.slice(0, limit); - if (markRead) { - this.setCursor(reader, inboxName, start + selected.length, sessionId); + try { + while (true) { + const messages = await this.loadInboxMessages(inboxName, sessionId); + const start = this.getCursor(reader, inboxName, sessionId); + const available = messages.slice(start); + + if (available.length > 0 || waitMs === 0 || Date.now() >= deadline) { + const selected = available.slice(0, limit); + if (markRead) { + this.setCursor( + reader, + inboxName, + start + selected.length, + sessionId, + ); + } + return { + messages: selected, + remaining: Math.max(0, available.length - selected.length), + total: messages.length, + }; } - return { - messages: selected, - remaining: Math.max(0, available.length - selected.length), - total: messages.length, - }; - } - await sleep( - Math.min(WORKER_IDLE_WAIT_MS, Math.max(50, deadline - Date.now())), - ); + await sleep( + Math.min(WORKER_IDLE_WAIT_MS, Math.max(50, deadline - Date.now())), + combined.signal, + ); + } + } finally { + combined.dispose(); } } @@ -386,6 +402,10 @@ export class AgentTeam { const reason = options.reason?.trim() || "Agent team shutting down."; this.closePromise = (async () => { + if (!this.abortController.signal.aborted) { + this.abortController.abort(reason); + } + const shutdownAt = new Date().toISOString(); for (const teammate of this.teammates.values()) { @@ -844,14 +864,22 @@ export class AgentTeam { return; } - const read = await this.readInbox({ - inboxName: name, - reader: name, - sessionId: teammate.sessionId, - markRead: false, - limit: 32, - waitMs: WORKER_IDLE_WAIT_MS, - }); + let read: TeamReadResult; + try { + read = await this.readInbox({ + inboxName: name, + reader: name, + sessionId: teammate.sessionId, + markRead: false, + limit: 32, + waitMs: WORKER_IDLE_WAIT_MS, + }); + } catch (error) { + if (this.shouldStopWorker(error)) { + return; + } + throw error; + } const latest = this.teammates.get(name); if (!latest) { @@ -872,7 +900,17 @@ export class AgentTeam { const run = this.claimRunSlot(name, "worker"); if (!run) { - await sleep(Math.min(200, WORKER_IDLE_WAIT_MS)); + try { + await sleep( + Math.min(200, WORKER_IDLE_WAIT_MS), + this.abortController.signal, + ); + } catch (error) { + if (this.shouldStopWorker(error)) { + return; + } + throw error; + } continue; } @@ -1073,6 +1111,10 @@ export class AgentTeam { } this.activeRuns.delete(name); } + + private shouldStopWorker(error: unknown): boolean { + return this.abortController.signal.aborted || isInterruptError(error); + } } function renderInboxPrompt( @@ -1374,8 +1416,35 @@ function encodeSessionStorageKey(sessionId: string): string { return Buffer.from(sessionId, "utf8").toString("base64url"); } -async function sleep(delayMs: number): Promise { - await new Promise((resolve) => setTimeout(resolve, delayMs)); +async function sleep(delayMs: number, signal?: AbortSignal): Promise { + await new Promise((resolve, reject) => { + const timer = setTimeout(() => { + signal?.removeEventListener("abort", abort); + resolve(); + }, delayMs); + + const abort = (): void => { + clearTimeout(timer); + signal?.removeEventListener("abort", abort); + reject(interruptError(signal)); + }; + + if (signal?.aborted) { + abort(); + return; + } + + signal?.addEventListener("abort", abort, { once: true }); + }); +} + +function interruptError(signal?: AbortSignal): Error { + const reason = signal?.reason; + return new Error( + typeof reason === "string" && reason.trim().length > 0 + ? reason + : "Run interrupted by user.", + ); } function finalizeTeammateHarness(teammate: LiveTeammate): void { diff --git a/skills/builtin/src/agent-team-plugin.test.ts b/skills/builtin/src/agent-team-plugin.test.ts new file mode 100644 index 00000000..4049d028 --- /dev/null +++ b/skills/builtin/src/agent-team-plugin.test.ts @@ -0,0 +1,82 @@ +import { describe, it, expect } from "vitest"; +import { createAgentTeamPlugin } from "./agent-team-plugin.js"; +import type { + AgentTeamInboxStore, + TeamMessage, +} from "@step-cli/core/agent/agent-team.js"; +import type { AgentHarnessFactory } from "@step-cli/core/agent/harness.js"; +import type { WorktreeManager } from "@step-cli/core/agent/worktree-manager.js"; +import type { ToolPluginContext } from "@step-cli/core/plugins/types.js"; +import { createMutableRef } from "@step-cli/utils/mutable-ref.js"; + +class MemoryInboxStore implements AgentTeamInboxStore { + private readonly messages: TeamMessage[] = []; + + async append(message: TeamMessage): Promise { + this.messages.push(message); + } + + async read(inboxName: string, sessionId?: string): Promise { + return this.messages.filter((message) => { + if (message.to !== inboxName) { + return false; + } + if (!sessionId) { + return true; + } + return message.sessionId === sessionId; + }); + } +} + +function mainPluginContext(): ToolPluginContext { + return { + workspaceRoot: "/tmp/workspace", + interactionProfile: { surface: "headless", canAskUser: false }, + harness: { + kind: "main", + name: "main", + depth: 0, + sessionId: "main-session", + goalId: "main:root", + executionProfile: { + workspaceMode: "shared", + memoryMode: "session", + priority: "interactive", + }, + }, + }; +} + +describe("agent-team-plugin read_inbox abort wiring", () => { + it("forwards the tool abort signal into team.readInbox waits", async () => { + const plugin = createAgentTeamPlugin( + createMutableRef("AgentHarnessFactory"), + new MemoryInboxStore(), + {} as WorktreeManager, + ); + const tools = plugin.register(mainPluginContext()); + const readInbox = tools.find( + (tool) => tool.definition.function.name === "read_inbox", + ); + expect(readInbox).toBeDefined(); + + const controller = new AbortController(); + const startedAt = Date.now(); + const pending = readInbox!.execute( + { waitMs: 5_000 }, + { + workspaceRoot: "/tmp/workspace", + commandTimeoutMs: 1_000, + commandOutputLimit: 1_000, + signal: controller.signal, + }, + {} as never, + ); + + setTimeout(() => controller.abort("Run interrupted by user."), 30); + + await expect(pending).rejects.toThrow("Run interrupted by user."); + expect(Date.now() - startedAt).toBeLessThan(1_000); + }); +}); diff --git a/skills/builtin/src/agent-team-plugin.ts b/skills/builtin/src/agent-team-plugin.ts index f5e85f28..35868f8b 100644 --- a/skills/builtin/src/agent-team-plugin.ts +++ b/skills/builtin/src/agent-team-plugin.ts @@ -541,6 +541,7 @@ function createReadInboxTool(team: AgentTeam): ToolSpec { markRead: args.markRead, limit: args.limit, waitMs: args.waitMs, + signal: ctx.signal, }); return {