Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 6 additions & 4 deletions src/node/services/agentSession.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3608,8 +3608,7 @@ export class AgentSession {
workspaceGoalService: this.workspaceGoalService,
experiments: options?.experiments,
disableWorkspaceAgents: options?.disableWorkspaceAgents,
hasQueuedMessage: () =>
!this.messageQueue.isEmpty() && this.messageQueue.getQueueDispatchMode() === "tool-end",
hasQueuedMessages: this.hasQueuedMessages.bind(this),
openaiTruncationModeOverride,
});

Expand Down Expand Up @@ -5022,8 +5021,11 @@ export class AgentSession {
this.backgroundProcessManager.setMessageQueued(this.workspaceId, false);
}

hasQueuedMessages(): boolean {
return !this.messageQueue.isEmpty();
hasQueuedMessages(dispatchMode?: "tool-end" | "turn-end"): boolean {
return (
!this.messageQueue.isEmpty() &&
(dispatchMode == null || this.messageQueue.getQueueDispatchMode() === dispatchMode)
);
}

async waitForPendingStreamErrorRecoveryDecision(): Promise<void> {
Expand Down
23 changes: 23 additions & 0 deletions src/node/services/agentSession.waitForIdle.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,29 @@ describe("AgentSession.waitForIdle", () => {
}
});

test("tracks only tool-end queued messages for foreground wait backgrounding", async () => {
const { session, cleanup } = await createAgentSessionHarness({
workspaceId: "wait-for-idle-tool-end-queue",
});

try {
session.queueMessage("later", {
model: "anthropic:claude-sonnet-4-5",
agentId: "exec",
queueDispatchMode: "turn-end",
});
expect(session.hasQueuedMessages()).toBe(true);
expect(session.hasQueuedMessages("tool-end")).toBe(false);

session.clearQueue();
session.queueMessage("now");
expect(session.hasQueuedMessages("tool-end")).toBe(true);
} finally {
session.dispose();
await cleanup();
}
});

test("removes repeated aborted idle waiters during one busy turn", async () => {
const { session, cleanup } = await createAgentSessionHarness({
workspaceId: "wait-for-idle-repeated-aborts",
Expand Down
6 changes: 3 additions & 3 deletions src/node/services/aiService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -250,7 +250,7 @@ export interface StreamMessageOptions {
allowAgentSetGoal?: boolean;
workspaceGoalService?: WorkspaceGoalService;
disableWorkspaceAgents?: boolean;
hasQueuedMessage?: () => boolean;
hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean;
muxMetadata?: MuxMessageMetadata;
openaiTruncationModeOverride?: "auto" | "disabled";
}
Expand Down Expand Up @@ -1003,7 +1003,7 @@ export class AIService extends EventEmitter {
allowAgentSetGoal,
workspaceGoalService,
disableWorkspaceAgents,
hasQueuedMessage,
hasQueuedMessages,
openaiTruncationModeOverride,
muxMetadata,
} = opts;
Expand Down Expand Up @@ -2694,7 +2694,7 @@ export class AIService extends EventEmitter {
maxOutputTokens,
effectiveToolPolicy,
streamToken, // Pass the pre-generated stream token
hasQueuedMessage,
hasQueuedMessages,
metadata.name,
effectiveThinkingLevel,
requestHeaders,
Expand Down
8 changes: 4 additions & 4 deletions src/node/services/streamManager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -416,7 +416,7 @@ describe("StreamManager - cleanupStreamTempDir", () => {
describe("StreamManager - stopWhen configuration", () => {
type StopWhenCondition = (options: { steps: unknown[] }) => boolean;
type BuildStopWhenCondition = (request: {
hasQueuedMessage?: () => boolean;
hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean;
toolPolicy?: ToolPolicy;
}) => StopWhenCondition[];

Expand All @@ -429,7 +429,7 @@ describe("StreamManager - stopWhen configuration", () => {

function requiredToolConditionForTests(toolPolicy: ToolPolicy): StopWhenCondition {
const [, , requiredToolCondition] = buildStopWhenForTests()({
hasQueuedMessage: () => false,
hasQueuedMessages: () => false,
toolPolicy,
});
return requiredToolCondition;
Expand All @@ -441,7 +441,7 @@ describe("StreamManager - stopWhen configuration", () => {

test("returns step-cap and queued-message conditions with no policy", () => {
let queued = false;
const stopWhen = buildStopWhenForTests()({ hasQueuedMessage: () => queued });
const stopWhen = buildStopWhenForTests()({ hasQueuedMessages: () => queued });
expect(stopWhen).toHaveLength(3);

const [maxStepCondition, queuedMessageCondition, requiredToolCondition] = stopWhen;
Expand Down Expand Up @@ -666,7 +666,7 @@ describe("StreamManager - sequential tool execution", () => {
headers?: Record<string, string | undefined>;
maxOutputTokens?: number;
streamCallSettings?: Record<string, unknown>;
hasQueuedMessage?: () => boolean;
hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean;
toolPolicy?: ToolPolicy;
toolChoice?: { type: "tool"; toolName: string };
}
Expand Down
20 changes: 10 additions & 10 deletions src/node/services/streamManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -156,7 +156,7 @@ interface StreamRequestConfig {
headers?: Record<string, string | undefined>;
maxOutputTokens?: number;
streamCallSettings?: Omit<ResolvedCallSettingsOverrides, "maxOutputTokens">;
hasQueuedMessage?: () => boolean;
hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean;
/** Optional hook for callers that need chunk-level visibility during streaming. */
onChunk?: StreamTextOnChunk;
/** Optional hook for callers that need the live prepared step transcript. */
Expand Down Expand Up @@ -1392,7 +1392,7 @@ export class StreamManager extends EventEmitter {
maxOutputTokens?: number,
callSettingsOverrides?: ResolvedCallSettingsOverrides,
toolPolicy?: ToolPolicy,
hasQueuedMessage?: () => boolean,
hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean,
headers?: Record<string, string | undefined>,
anthropicCacheTtlOverride?: AnthropicCacheTtl,
onChunk?: StreamTextOnChunk,
Expand Down Expand Up @@ -1449,15 +1449,15 @@ export class StreamManager extends EventEmitter {
maxOutputTokens: effectiveMaxOutputTokens,
streamCallSettings:
Object.keys(streamCallSettings).length > 0 ? streamCallSettings : undefined,
hasQueuedMessage,
hasQueuedMessages,
onChunk,
onStepMessages,
toolPolicy,
};
}

private createStopWhenCondition(
request: Pick<StreamRequestConfig, "hasQueuedMessage" | "toolPolicy">
request: Pick<StreamRequestConfig, "hasQueuedMessages" | "toolPolicy">
): Array<ReturnType<typeof stepCountIs>> {
// Completion-tool stop check: completion/routing tools use explicit
// success/ok markers (agent_report, propose_plan).
Expand Down Expand Up @@ -1507,7 +1507,7 @@ export class StreamManager extends EventEmitter {

return [
stepCountIs(100000),
() => request.hasQueuedMessage?.() ?? false,
() => request.hasQueuedMessages?.("tool-end") ?? false,
hasSuccessfulRequiredToolResult,
];
}
Expand Down Expand Up @@ -1570,7 +1570,7 @@ export class StreamManager extends EventEmitter {
maxOutputTokens?: number,
toolPolicy?: ToolPolicy,
callSettingsOverrides?: ResolvedCallSettingsOverrides,
hasQueuedMessage?: () => boolean,
hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean,
workspaceName?: string,
thinkingLevel?: string,
headers?: Record<string, string | undefined>,
Expand All @@ -1593,7 +1593,7 @@ export class StreamManager extends EventEmitter {
maxOutputTokens,
callSettingsOverrides,
toolPolicy,
hasQueuedMessage,
hasQueuedMessages,
headers,
anthropicCacheTtlOverride,
onChunk,
Expand Down Expand Up @@ -2249,7 +2249,7 @@ export class StreamManager extends EventEmitter {
fallbackState.original.maxOutputTokens,
prepared.data.callSettingsOverrides,
streamInfo.request.toolPolicy,
streamInfo.request.hasQueuedMessage,
streamInfo.request.hasQueuedMessages,
prepared.data.headers,
prepared.data.anthropicCacheTtl,
streamInfo.request.onChunk,
Expand Down Expand Up @@ -3638,7 +3638,7 @@ export class StreamManager extends EventEmitter {
maxOutputTokens?: number,
toolPolicy?: ToolPolicy,
providedStreamToken?: StreamToken,
hasQueuedMessage?: () => boolean,
hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean,
workspaceName?: string,
thinkingLevel?: string,
headers?: Record<string, string | undefined>,
Expand Down Expand Up @@ -3723,7 +3723,7 @@ export class StreamManager extends EventEmitter {
maxOutputTokens,
toolPolicy,
callSettingsOverrides,
hasQueuedMessage,
hasQueuedMessages,
workspaceName,
thinkingLevel,
headers,
Expand Down
74 changes: 74 additions & 0 deletions src/node/services/taskService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -363,6 +363,7 @@ function createWorkspaceServiceMocks(
sendMessage: ReturnType<typeof mock>;
resumeStream: ReturnType<typeof mock>;
clearQueue: ReturnType<typeof mock>;
hasQueuedMessages: ReturnType<typeof mock>;
hasPendingQueuedOrPreparingTurn: ReturnType<typeof mock>;
waitForPendingStreamErrorRecoveryDecision: ReturnType<typeof mock>;
remove: ReturnType<typeof mock>;
Expand All @@ -380,6 +381,7 @@ function createWorkspaceServiceMocks(
sendMessage: ReturnType<typeof mock>;
resumeStream: ReturnType<typeof mock>;
clearQueue: ReturnType<typeof mock>;
hasQueuedMessages: ReturnType<typeof mock>;
hasPendingQueuedOrPreparingTurn: ReturnType<typeof mock>;
waitForPendingStreamErrorRecoveryDecision: ReturnType<typeof mock>;
remove: ReturnType<typeof mock>;
Expand All @@ -398,6 +400,7 @@ function createWorkspaceServiceMocks(
overrides?.resumeStream ??
mock((): Promise<Result<{ started: boolean }>> => Promise.resolve(Ok({ started: true })));
const clearQueue = overrides?.clearQueue ?? mock((): Result<void> => Ok(undefined));
const hasQueuedMessages = overrides?.hasQueuedMessages ?? mock(() => false);
const hasPendingQueuedOrPreparingTurn =
overrides?.hasPendingQueuedOrPreparingTurn ?? mock(() => false);
const waitForPendingStreamErrorRecoveryDecision =
Expand Down Expand Up @@ -429,6 +432,7 @@ function createWorkspaceServiceMocks(
sendMessage,
resumeStream,
clearQueue,
hasQueuedMessages,
hasPendingQueuedOrPreparingTurn,
waitForPendingStreamErrorRecoveryDecision,
remove,
Expand All @@ -444,6 +448,7 @@ function createWorkspaceServiceMocks(
sendMessage,
resumeStream,
clearQueue,
hasQueuedMessages,
hasPendingQueuedOrPreparingTurn,
waitForPendingStreamErrorRecoveryDecision,
remove,
Expand Down Expand Up @@ -521,6 +526,7 @@ describe("TaskService", () => {
sendMessage?: ReturnType<typeof mock>;
remove?: ReturnType<typeof mock>;
isStreaming?: ReturnType<typeof mock>;
hasQueuedMessages?: ReturnType<typeof mock>;
hasPendingQueuedOrPreparingTurn?: ReturnType<typeof mock>;
waitForPendingStreamErrorRecoveryDecision?: ReturnType<typeof mock>;
} = {}
Expand Down Expand Up @@ -553,6 +559,9 @@ describe("TaskService", () => {
create: createWorkspace,
...(options.sendMessage != null ? { sendMessage: options.sendMessage } : {}),
...(options.remove != null ? { remove: options.remove } : {}),
...(options.hasQueuedMessages != null
? { hasQueuedMessages: options.hasQueuedMessages }
: {}),
...(options.hasPendingQueuedOrPreparingTurn != null
? { hasPendingQueuedOrPreparingTurn: options.hasPendingQueuedOrPreparingTurn }
: {}),
Expand Down Expand Up @@ -1891,6 +1900,23 @@ describe("TaskService", () => {
expect(taskService.backgroundForegroundWaitsForWorkspace(parentId)).toBe(0);
});

test("waitForWorkspaceTurn backgrounds when tool-end message was already queued", async () => {
const hasQueuedMessages = mock(() => true);
const { parentId, taskService } = await startWorkspaceTurnForTest({ hasQueuedMessages });

const waitError = await taskService
.waitForWorkspaceTurn("wst_handle", {
requestingWorkspaceId: parentId,
timeoutMs: 1_000,
backgroundOnMessageQueued: true,
})
.catch((error: unknown) => error);

expect(waitError).toBeInstanceOf(ForegroundWaitBackgroundedError);
expect(hasQueuedMessages).toHaveBeenCalledWith(parentId, "tool-end");
expect(taskService.backgroundForegroundWaitsForWorkspace(parentId)).toBe(0);
});

test("disposable workspace turns are removed after completion, error, or interruption", async () => {
const completedRemove = mock((): Promise<Result<void>> => Promise.resolve(Ok(undefined)));
const completed = await startWorkspaceTurnForTest({
Expand Down Expand Up @@ -7449,6 +7475,54 @@ describe("TaskService", () => {
expect(count2).toBe(0);
});

test("backgrounds waiters when tool-end message was already queued", async () => {
const config = await createTestConfig(rootDir);

const parentId = "parent-ws";
const childId = "child-task-ws";
const projectPath = "/test/project";

await saveWorkspaces(
config,
projectPath,
[
{ path: `${projectPath}/parent`, id: parentId, name: "parent" },
{
path: `${projectPath}/child`,
id: childId,
name: "agent_explore_child",
parentWorkspaceId: parentId,
agentType: "explore",
taskStatus: "queued",
},
],
testTaskSettings(2, 3)
);

const hasQueuedMessages = mock(() => true);
const { workspaceService } = createWorkspaceServiceMocks({ hasQueuedMessages });
const { taskService } = createTaskServiceHarness(config, { workspaceService });
const internal = taskService as unknown as {
backgroundableForegroundWaitersByWorkspaceId: Map<string, Set<unknown>>;
pendingStartWaitersByTaskId: Map<string, unknown[]>;
pendingWaitersByTaskId: Map<string, unknown[]>;
};

const waitError = await taskService
.waitForAgentReport(childId, {
requestingWorkspaceId: parentId,
backgroundOnMessageQueued: true,
})
.catch((error: unknown) => error);

expect(waitError).toBeInstanceOf(ForegroundWaitBackgroundedError);
expect(hasQueuedMessages).toHaveBeenCalledWith(parentId, "tool-end");
expect(taskService.backgroundForegroundWaitsForWorkspace(parentId)).toBe(0);
expect(internal.backgroundableForegroundWaitersByWorkspaceId.has(parentId)).toBe(false);
expect(internal.pendingStartWaitersByTaskId.has(childId)).toBe(false);
expect(internal.pendingWaitersByTaskId.has(childId)).toBe(false);
});

test("defaults to queue-backgroundable when requestingWorkspaceId is present", async () => {
const config = await createTestConfig(rootDir);

Expand Down
15 changes: 15 additions & 0 deletions src/node/services/taskService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4030,6 +4030,13 @@ export class TaskService {
timeoutMs
);

if (
shouldBackgroundOnQueuedMessage &&
this.workspaceService.hasQueuedMessages(options.requestingWorkspaceId, "tool-end")
) {
this.backgroundForegroundWaitsForWorkspace(options.requestingWorkspaceId);
}

void (async () => {
const record = await this.taskHandleStore.getWorkspaceTurn(
options.requestingWorkspaceId,
Expand Down Expand Up @@ -4375,6 +4382,14 @@ export class TaskService {
};
options.abortSignal.addEventListener("abort", abortListener, { once: true });
}

if (
shouldBackgroundOnQueuedMessage &&
requestingWorkspaceId &&
this.workspaceService.hasQueuedMessages(requestingWorkspaceId, "tool-end")
) {
this.backgroundForegroundWaitsForWorkspace(requestingWorkspaceId);
}
})().catch((error: unknown) => {
reject(error instanceof Error ? error : new Error(String(error)));
});
Expand Down
4 changes: 4 additions & 0 deletions src/node/services/workspaceService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7578,6 +7578,10 @@ export class WorkspaceService extends EventEmitter {
}
}

hasQueuedMessages(workspaceId: string, dispatchMode?: "tool-end" | "turn-end"): boolean {
return this.sessions.get(workspaceId.trim())?.hasQueuedMessages(dispatchMode) ?? false;
}

async waitForPendingStreamErrorRecoveryDecision(workspaceId: string): Promise<void> {
const session = this.sessions.get(workspaceId.trim());
await session?.waitForPendingStreamErrorRecoveryDecision();
Expand Down
Loading