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
71 changes: 71 additions & 0 deletions .changeset/clever-otters-show.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
---
"@voltagent/core": patch
---

fix: sub-agent stream error handling and propagation - #521

## What Changed

Fixed a critical issue where sub-agent stream errors were incorrectly reported as successful operations with empty responses. Supervisors now properly detect and handle sub-agent failures with configurable error handling behavior.

## The Problem (Before)

When a sub-agent's `streamText` encountered an error event:

- ❌ Supervisor received `status: "success"` with empty response
- ❌ No way to distinguish between empty success and failure
- ❌ Error details were lost, making debugging difficult
- ❌ Supervisors would continue as if the operation succeeded

```typescript
// Before: Sub-agent fails but supervisor doesn't know
const result = await subAgentManager.handoffTask({
task: "Process data",
targetAgent: failingAgent,
});

// result.status === "success" (WRONG!)
// result.result === "" (No error info)
```

## The Solution (After)

Sub-agent errors are now properly detected and reported:

- ✅ Stream errors return `status: "error"` with error details
- ✅ Error messages included in responses (configurable)
- ✅ Partial content preserved when errors occur after text generation
- ✅ New configuration options for flexible error handling

```typescript
// After: Proper error detection and handling
const result = await subAgentManager.handoffTask({
task: "Process data",
targetAgent: failingAgent,
});

// result.status === "error" (CORRECT!)
// result.result === "Error in FailingAgent: Stream processing failed"
// result.error === Error object with full details
```

## New Configuration Options

Added `SupervisorConfig` options for customizable error handling:

```typescript
const supervisor = new Agent({
name: "Supervisor",
subAgents: [agent1, agent2],
supervisorConfig: {
// Throw exceptions on stream errors instead of returning error results
throwOnStreamError: false, // default: false

// Include error messages in empty responses
includeErrorInEmptyResponse: true, // default: true

// Custom guidelines for error handling
customGuidelines: ["When a sub-agent fails, provide alternative solutions"],
},
});
```
1 change: 1 addition & 0 deletions packages/core/src/agent/index.ts
Original file line number Diff line number Diff line change
@@ -1 +1,2 @@
export { Agent } from "./agent";
export type { SupervisorConfig } from "./types";
171 changes: 168 additions & 3 deletions packages/core/src/agent/subagent/index.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -885,7 +885,7 @@ describe("SubAgentManager", () => {
// Create subAgentManager with error type in configuration
const supervisorConfig = {
fullStreamEventForwarding: {
types: ["tool-call", "tool-result", "error"],
types: ["tool-call", "tool-result", "error"] as ("tool-call" | "tool-result" | "error")[],
},
};
const localSubAgentManager = new SubAgentManager("Main Agent", [], supervisorConfig);
Expand Down Expand Up @@ -943,6 +943,127 @@ describe("SubAgentManager", () => {
expect(result.status).toBe("success");
});

it("should return error status when stream error occurs with no text content", async () => {
const mockAgent = new MockAgent("error-only-agent", "Error Only Agent");

// Mock streamText to only emit an error (no text content)
mockAgent.streamText = vi.fn().mockReturnValue({
fullStream: (async function* () {
yield { type: "error", error: new Error("Critical stream error") };
})(),
textStream: (async function* () {
// No text emitted
})(),
});

const options: AgentHandoffOptions = {
task: "Task that will fail immediately",
targetAgent: mockAgent,
context: {},
sharedContext: [],
forwardEvent: vi.fn(),
};

const result = await subAgentManager.handoffTask(options);

// Should return error status
expect(result.status).toBe("error");
expect(result.error).toBeInstanceOf(Error);
expect((result.error as Error).message).toBe("Critical stream error");

// By default, includeErrorInEmptyResponse is true, so error message should be in result
expect(result.result).toContain("Error in Error Only Agent");
expect(result.result).toContain("Critical stream error");
});

it("should return success with partial content when error occurs after text", async () => {
const mockAgent = new MockAgent("partial-content-agent", "Partial Content Agent");

// Mock streamText to emit some text then error
mockAgent.streamText = vi.fn().mockReturnValue({
fullStream: (async function* () {
yield { type: "text-delta", textDelta: "Partial response" };
yield { type: "error", error: new Error("Stream interrupted") };
})(),
});

const options: AgentHandoffOptions = {
task: "Task with partial completion",
targetAgent: mockAgent,
context: {},
sharedContext: [],
forwardEvent: vi.fn(),
};

const result = await subAgentManager.handoffTask(options);

// Should return success with the partial content
expect(result.status).toBe("success");
expect(result.result).toBe("Partial response");
expect(result.error).toBeUndefined();
});

it("should throw error when throwOnStreamError is true", async () => {
const supervisorConfig = {
throwOnStreamError: true,
};
const localSubAgentManager = new SubAgentManager("Main Agent", [], supervisorConfig);

const mockAgent = new MockAgent("throw-error-agent", "Throw Error Agent");

// Mock streamText to only emit an error
mockAgent.streamText = vi.fn().mockReturnValue({
fullStream: (async function* () {
yield { type: "error", error: new Error("Stream error to throw") };
})(),
});

const options: AgentHandoffOptions = {
task: "Task that should throw",
targetAgent: mockAgent,
context: {},
sharedContext: [],
forwardEvent: vi.fn(),
};

// Should throw the error
await expect(localSubAgentManager.handoffTask(options)).rejects.toThrow(
"Stream error in Throw Error Agent: Stream error to throw",
);
});

it("should not include error text when includeErrorInEmptyResponse is false", async () => {
const supervisorConfig = {
includeErrorInEmptyResponse: false,
};
const localSubAgentManager = new SubAgentManager("Main Agent", [], supervisorConfig);

const mockAgent = new MockAgent("no-error-text-agent", "No Error Text Agent");

// Mock streamText to only emit an error
mockAgent.streamText = vi.fn().mockReturnValue({
fullStream: (async function* () {
yield { type: "error", error: new Error("Hidden error message") };
})(),
});

const options: AgentHandoffOptions = {
task: "Task with hidden error",
targetAgent: mockAgent,
context: {},
sharedContext: [],
forwardEvent: vi.fn(),
};

const result = await localSubAgentManager.handoffTask(options);

// Should return error status but empty result
expect(result.status).toBe("error");
expect(result.result).toBe("");
expect(result.error).toBeInstanceOf(Error);
expect((result.error as Error).message).toBe("Hidden error message");
});

it("should forward events through delegate tool", async () => {
const forwardEventSpy = vi.fn();
const mockAgent = new MockAgent("delegate-agent", "Delegate Agent");
Expand Down Expand Up @@ -979,6 +1100,45 @@ describe("SubAgentManager", () => {
}
});

it("should propagate stream errors through delegate_task tool", async () => {
const forwardEventSpy = vi.fn();
const mockAgent = new MockAgent("error-delegate-agent", "Error Delegate Agent");

// Mock streamText to only emit an error
mockAgent.streamText = vi.fn().mockReturnValue({
fullStream: (async function* () {
yield { type: "error", error: new Error("Delegation error") };
})(),
});

subAgentManager.addSubAgent(mockAgent as any);

const tool = subAgentManager.createDelegateTool({
sourceAgent: { id: "supervisor-agent" } as any,
operationContext: { userContext: new Map(), systemContext: new Map() } as any,
currentHistoryEntryId: "history-error",
forwardEvent: forwardEventSpy,
});

const result = await tool.execute({
task: "Task that will fail in delegation",
targetAgents: ["Error Delegate Agent"],
context: {},
});

// Should return structured results with error status
expect(Array.isArray(result)).toBe(true);
expect(result[0]).toMatchObject({
agentName: "Error Delegate Agent",
status: "error",
error: "Delegation error",
});

// Error message should be included in response by default
expect(result[0].response).toContain("Error in Error Delegate Agent");
expect(result[0].response).toContain("Delegation error");
});

it("should handle multiple agents with event forwarding", async () => {
const forwardEventSpy = vi.fn();
const mockAgent1 = new MockAgent("multi-agent-1", "Multi Agent 1");
Expand Down Expand Up @@ -1089,7 +1249,12 @@ describe("SubAgentManager", () => {
it("should use custom event types from supervisor configuration", async () => {
const supervisorConfig = {
fullStreamEventForwarding: {
types: ["tool-call", "tool-result", "text-delta", "reasoning"],
types: ["tool-call", "tool-result", "text-delta", "reasoning"] as (
| "tool-call"
| "tool-result"
| "text-delta"
| "reasoning"
)[],
},
};
subAgentManager = new SubAgentManager("Main Agent", [], supervisorConfig);
Expand All @@ -1115,7 +1280,7 @@ describe("SubAgentManager", () => {
it("should respect addSubAgentPrefix configuration", async () => {
const supervisorConfig = {
fullStreamEventForwarding: {
types: ["tool-call"],
types: ["tool-call"] as "tool-call"[],
addSubAgentPrefix: false,
},
};
Expand Down
59 changes: 59 additions & 0 deletions packages/core/src/agent/subagent/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -295,6 +295,9 @@ ${guidelinesText}
// Use the provided conversationId or generate a new one
const handoffConversationId = conversationId || crypto.randomUUID();

// Track if we should rethrow stream errors
let streamErrorToThrow: Error | null = null;

try {
// Call onHandoff hook if source agent is provided
if (sourceAgent && targetAgent.hooks) {
Expand Down Expand Up @@ -391,6 +394,10 @@ ${task}\n\nContext: ${safeStringify(context, { indentation: 2 })}`;
// Collect all stream chunks for final result
finalResult = "";

// Track stream errors and whether we received any text content
let streamError: Error | null = null;
let hasTextContent = false;

if (streamResponse.fullStream && forwardEvent) {
// Get event forwarding configuration
const eventForwardingConfig = {
Expand All @@ -410,6 +417,7 @@ ${task}\n\nContext: ${safeStringify(context, { indentation: 2 })}`;
switch (part.type) {
case "text-delta": {
finalResult += part.textDelta;
hasTextContent = true;

const eventData = {
type: "text-delta",
Expand Down Expand Up @@ -491,6 +499,9 @@ ${task}\n\nContext: ${safeStringify(context, { indentation: 2 })}`;
}

case "error": {
// Capture the error for proper handling
streamError = part.error;

const eventData = {
type: "error",
data: {
Expand All @@ -513,9 +524,52 @@ ${task}\n\nContext: ${safeStringify(context, { indentation: 2 })}`;
} else {
for await (const part of streamResponse.textStream) {
finalResult += part;
hasTextContent = true;
}
}

// Handle stream errors based on configuration
if (streamError && !hasTextContent) {
const errorMessage =
streamError instanceof Error ? streamError.message : String(streamError);

// Check if we should throw the error
if (this.supervisorConfig?.throwOnStreamError) {
// Store the error to throw after the try-catch
streamErrorToThrow = new Error(`Stream error in ${targetAgent.name}: ${errorMessage}`);
// Still throw here to exit the try block
throw streamErrorToThrow;
}

// Check if we should include error message in empty response
const includeErrorInResponse = this.supervisorConfig?.includeErrorInEmptyResponse ?? true;

return {
result: includeErrorInResponse ? `Error in ${targetAgent.name}: ${errorMessage}` : "",
conversationId: handoffConversationId,
messages: [
taskMessage,
{
role: "system" as const,
content: `Stream error occurred: ${errorMessage}`,
},
],
status: "error",
error: streamError,
};
}

// If we have partial content despite an error, log warning but return the content
if (streamError && hasTextContent) {
const logger =
options.parentOperationContext?.logger ||
getGlobalLogger().child({ component: "subagent-manager" });
logger.warn(`Stream error occurred after partial content in ${targetAgent.name}`, {
error: streamError,
partialContent: finalResult,
});
}

finalMessages = [taskMessage, { role: "assistant", content: finalResult }];
}

Expand All @@ -526,6 +580,11 @@ ${task}\n\nContext: ${safeStringify(context, { indentation: 2 })}`;
status: "success",
};
} catch (error) {
// If this is the stream error we marked for rethrowing, rethrow it
if (streamErrorToThrow && error === streamErrorToThrow) {
throw error;
}

const logger =
options.parentOperationContext?.logger ||
getGlobalLogger().child({ component: "subagent-manager" });
Expand Down
Loading