Skip to content
Open
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
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,12 @@
# Changelog

## Unreleased

- Preserve classified provider errors when streamed responses fail, including
aborts racing late message writes (#320). `DeltaStreamer.getOrCreateStreamId`
now accepts `{ ifAborted: "returnUndefined" }` for abort-aware writes; its
existing no-argument behavior is unchanged.

## 0.7.1

- Persists oversized streamed files (#307)
Expand Down
72 changes: 69 additions & 3 deletions src/component/messages.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -263,9 +263,7 @@ describe("agent", () => {
const { messages } = await t.mutation(api.messages.addMessages, {
threadId: thread._id as Id<"threads">,
order: "next",
messages: [
{ message: { role: "assistant", content: "separate reply" } },
],
messages: [{ message: { role: "assistant", content: "separate reply" } }],
});

expect(messages[0]).toMatchObject({ order: 1, stepOrder: 0 });
Expand Down Expand Up @@ -1262,3 +1260,71 @@ describe("agent", () => {
});
});
});

describe("late saves racing a failed pending message (issue #320)", () => {
const PROVIDER_ERROR = "invalid_prompt: Invalid prompt: flagged by policy.";

test("keeps the first durable failure authoritative", async () => {
const t = initConvexTest();
const thread = await t.mutation(api.threads.createThread, {
userId: "u1",
});
const threadId = thread._id as Id<"threads">;

const { messages: seeded } = await t.mutation(api.messages.addMessages, {
threadId,
messages: [
{ message: { role: "user", content: "hello" } },
{ message: { role: "assistant", content: [] }, status: "pending" },
],
});
const pending = seeded.at(-1)!;
expect(pending.status).toBe("pending");

const streamId = await t.mutation(api.streams.create, {
threadId,
order: pending.order,
stepOrder: pending.stepOrder,
format: "UIMessageChunk",
});

await t.mutation(api.messages.finalizeMessage, {
messageId: pending._id as Id<"messages">,
result: { status: "failed", error: PROVIDER_ERROR },
});
await t.mutation(api.streams.abort, { streamId, reason: PROVIDER_ERROR });

const { messages: late } = await t.mutation(api.messages.addMessages, {
threadId,
pendingMessageId: pending._id as Id<"messages">,
finishStreamId: streamId,
failPendingSteps: false,
messages: [
{ message: { role: "assistant", content: "partial response" } },
],
});

const assistants = (
await t.run(async (ctx) =>
ctx.db
.query("messages")
.withIndex("threadId_status_tool_order_stepOrder", (q) =>
q.eq("threadId", threadId),
)
.collect(),
)
).filter((message) => message.message?.role === "assistant");

expect(late).toHaveLength(1);
expect(assistants).toHaveLength(1);
expect(assistants[0]!._id).toBe(pending._id);
expect(assistants[0]!.status).toBe("failed");
expect(assistants[0]!.error).toBe(PROVIDER_ERROR);
expect(assistants[0]!.text).toBe("partial response");

const stream = await t.run((ctx) =>
ctx.db.get("streamingMessages", streamId),
);
expect(stream?.state.kind).toBe("aborted");
});
});
4 changes: 2 additions & 2 deletions src/component/messages.ts
Original file line number Diff line number Diff line change
Expand Up @@ -327,8 +327,8 @@ async function addMessagesHandler(
if (pendingMessage.status === "failed") {
fail = true;
error =
`Trying to update a message that failed: ${pendingMessageId}, ` +
`error: ${pendingMessage.error ?? error}`;
pendingMessage.error ??
`Trying to update a message that failed: ${pendingMessageId}`;
messageDoc.status = "failed";
messageDoc.error = error;
}
Expand Down
76 changes: 76 additions & 0 deletions src/errors.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
import { APICallError } from "@ai-sdk/provider";
import { describe, expect, test } from "vitest";
import { errorToString } from "./errors.js";

describe("errorToString", () => {
test("preserves provider error classifications", () => {
const details = {
error: {
code: "invalid_prompt",
message: "Invalid prompt: flagged by policy",
},
};
const apiError = new APICallError({
message: "Invalid prompt: flagged by policy",
url: "https://api.example.test",
requestBodyValues: {},
statusCode: 400,
data: details,
});

expect(errorToString(details)).toBe(
"invalid_prompt: Invalid prompt: flagged by policy",
);
expect(errorToString(apiError)).toBe(
"invalid_prompt: Invalid prompt: flagged by policy",
);
expect(errorToString(new Error())).toBe("Error");
expect(errorToString(new TypeError())).toBe("TypeError");
const systemError = Object.assign(new Error("socket hang up"), {
code: "ECONNRESET",
});
expect(errorToString(systemError)).toBe("socket hang up");
const codeOnly = Object.assign(new Error("Request failed"), {
data: { code: "rate_limit" },
});
expect(errorToString(codeOnly)).toBe("rate_limit: Request failed");
});

test("serializes objects without mistaking shared values for cycles", () => {
const shared = { detail: "provider disconnected" };
const circular: Record<string, unknown> = { shared };
circular.self = circular;

expect(errorToString({ x: shared, y: shared })).toBe(
'{"x":{"detail":"provider disconnected"},"y":{"detail":"provider disconnected"}}',
);
expect(errorToString(circular)).toBe(
'{"shared":{"detail":"provider disconnected"},"self":"[Circular]"}',
);
});

test("bounds stored error text without splitting surrogate pairs", () => {
const serialized = errorToString(`${"x".repeat(1022)}😀tail`);

expect(serialized.length).toBeLessThanOrEqual(1024);
expect(serialized.endsWith("x…")).toBe(true);
});

test("does not throw when Error properties are hostile accessors", () => {
const error = new Error();
Object.defineProperties(error, {
message: {
get() {
throw new Error("message getter failed");
},
},
name: {
get() {
throw new Error("name getter failed");
},
},
});

expect(errorToString(error)).toBe("Unknown error");
});
});
113 changes: 113 additions & 0 deletions src/errors.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
const MAX_ERROR_LENGTH = 1024;

export function errorToString(error: unknown): string {
return truncateError(describeError(error));
}

function describeError(error: unknown): string {
if (typeof error === "string") return error;
if (error instanceof Error) {
const message = property(error, "message");
if (typeof message !== "string" || message.length === 0) {
const name = property(error, "name");
return typeof name === "string" && name.length > 0
? name
: safeString(error);
}
const nested = errorDetails(
property(error, "error") ?? property(error, "data"),
);
return (
formatDetails({
message: nested.message ?? message,
code: nested.code,
}) ?? message
);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

const details = formatDetails(errorDetails(error));
if (details) return details;

if (error && typeof error === "object") {
try {
const ancestors: object[] = [];
const serialized = JSON.stringify(error, function (_key, value: unknown) {
if (typeof value === "bigint") return value.toString();
if (!value || typeof value !== "object") return value;
while (ancestors.length > 0 && ancestors.at(-1) !== this) {
ancestors.pop();
}
if (ancestors.includes(value)) return "[Circular]";
ancestors.push(value);
return value;
});
if (serialized) return serialized;
} catch {
return safeString(error);
}
}

return safeString(error);
}

function safeString(error: unknown): string {
try {
return String(error);
} catch {
return "Unknown error";
}
}

function errorDetails(error: unknown): { message?: string; code?: string } {
let current = error;
let message: string | undefined;
let code: string | undefined;
for (let depth = 0; depth < 3; depth++) {
if (typeof current === "string") {
message ??= current;
break;
}
if (!current || typeof current !== "object") break;

const currentMessage = property(current, "message");
if (typeof currentMessage === "string" && currentMessage.length > 0) {
message ??= currentMessage;
}
const currentCode = property(current, "code");
if (typeof currentCode === "string" || typeof currentCode === "number") {
code ??= String(currentCode);
}
if (message && code) break;
current = property(current, "error") ?? property(current, "data");
}
return { message, code };
}

function property(value: object, key: string): unknown {
try {
return (value as Record<string, unknown>)[key];
} catch {
return undefined;
}
}

function formatDetails({
message,
code,
}: {
message?: string;
code?: string;
}): string | undefined {
if (message && code) {
return message.startsWith(`${code}:`) ? message : `${code}: ${message}`;
}
return message ?? code;
}

function truncateError(error: string): string {
if (error.length <= MAX_ERROR_LENGTH) return error;
let truncated = error.slice(0, MAX_ERROR_LENGTH - 1);
const last = truncated.charCodeAt(truncated.length - 1);
if (last >= 0xd800 && last <= 0xdbff) truncated = truncated.slice(0, -1);
return `${truncated}…`;
}
Loading
Loading