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
26 changes: 19 additions & 7 deletions apps/mobile/src/chat.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ import { BackgroundUpdates } from "./background-updates";
import { BrowserRunContext, BrowserToolCard } from "./browser-tool-card";
import { BrowserThreadCard } from "./computer";
import { ConversationQueue, type QueuedMessage } from "./conversation-queue";
import { runConversationTurn } from "./conversation-run";
import { replayedRunError, runConversationTurn } from "./conversation-run";
import { MailToolCard } from "./mail-tool-card";
import { FileThreadCard, TaskThreadCard } from "./thread-artifacts";
import { type Selection, useMuseThread } from "./threads";
Expand Down Expand Up @@ -200,6 +200,8 @@ export function ChatScreen({
const [saveError, setSaveError] = useState("");
const [historyError, setHistoryError] = useState("");
const [historyAttempt, setHistoryAttempt] = useState(0);
// Loads still replaying history. A count, so an old load finishing does not end a newer one.
const replaying = useRef(0);
useEffect(() => {
if (!isReady) return;
let active = true;
Expand All @@ -213,12 +215,20 @@ export function ChatScreen({
async function hydrate() {
try {
if (richThreads) {
if (selection.existing)
await runConversationTurn(
agentId,
() => copilotkit.connectAgent({ agent }),
(onError) => copilotkit.subscribe({ onError }),
);
if (selection.existing) {
// Replaying history re-emits past RUN_ERROR events; only connection failures block loading.
replaying.current += 1;
try {
await runConversationTurn(
agentId,
() => copilotkit.connectAgent({ agent }),
(onError) => copilotkit.subscribe({ onError }),
[replayedRunError],
);
} finally {
replaying.current -= 1;
}
}
} else {
const { messages } = await api.request<{ messages: Message[] }>("/api/conversation");
if (active) agent.setMessages(messages);
Expand Down Expand Up @@ -299,6 +309,8 @@ export function ChatScreen({
const subscription = copilotkit.subscribe({
onError: (event) => {
if (event.context?.agentId && event.context.agentId !== agentId) return;
// A failed turn saved in history is already over; it is not a failure of this session.
if (replaying.current > 0 && event.code === replayedRunError) return;
const failure = event.error instanceof Error ? event.error : new Error(String(event.error));
setError(failure.message);
},
Expand Down
7 changes: 6 additions & 1 deletion apps/mobile/src/conversation-run.ts
Original file line number Diff line number Diff line change
@@ -1,14 +1,19 @@
type RunError = { error: unknown; context?: { agentId?: string } };
type RunError = { error: unknown; code?: string; context?: { agentId?: string } };

/** A RUN_ERROR event from the thread's saved history, not a failure of this call. */
export const replayedRunError = "agent_run_error_event";

/** CopilotKit emits run failures through onError even when runAgent resolves. */
export async function runConversationTurn(
agentId: string,
execute: () => Promise<unknown>,
subscribe: (listener: (event: RunError) => void) => { unsubscribe: () => void },
ignore: readonly string[] = [],
) {
let failure: Error | undefined;
const subscription = subscribe((event) => {
if (event.context?.agentId && event.context.agentId !== agentId) return;
if (event.code && ignore.includes(event.code)) return;
failure = event.error instanceof Error ? event.error : new Error(String(event.error));
});
try {
Expand Down
52 changes: 50 additions & 2 deletions tests/conversation-sdk.test.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
import assert from "node:assert/strict";
import { test } from "node:test";
import { AbstractAgent } from "@ag-ui/client";
import { EventType } from "@ag-ui/core";
import { CopilotKitCore } from "@copilotkit/core";
import { throwError } from "rxjs";
import { of, throwError } from "rxjs";
import { ConversationQueue } from "../apps/mobile/src/conversation-queue.ts";
import { runConversationTurn } from "../apps/mobile/src/conversation-run.ts";
import { replayedRunError, runConversationTurn } from "../apps/mobile/src/conversation-run.ts";

test("an emitted CopilotKit run error stops the queue even when runAgent resolves", async () => {
let attempts = 0;
Expand Down Expand Up @@ -36,3 +37,50 @@ test("an emitted CopilotKit run error stops the queue even when runAgent resolve
["second"],
);
});

test("a failed turn saved in thread history does not block loading the conversation", async () => {
class ReplayingAgent extends AbstractAgent {
run() {
return throwError(() => new Error("not used"));
}
connect() {
return of(
{ type: EventType.RUN_STARTED, threadId: "thread", runId: "old-run" },
{ type: EventType.RUN_ERROR, message: "Missing Authentication header" },
);
}
}
const agent = new ReplayingAgent({ agentId: "default", threadId: "thread" });
const core = new CopilotKitCore({ agents__unsafe_dev_only: { default: agent } });
const connect = (ignore?: string[]) =>
runConversationTurn(
"default",
() => core.connectAgent({ agent }),
(onError) => core.subscribe({ onError }),
ignore,
);
await assert.rejects(connect(), /Missing Authentication header/);
await connect([replayedRunError]);
});

test("a connection failure still blocks loading the conversation", async () => {
class UnreachableAgent extends AbstractAgent {
run() {
return throwError(() => new Error("not used"));
}
connect() {
return throwError(() => new Error("Thread service unavailable"));
}
}
const agent = new UnreachableAgent({ agentId: "default", threadId: "thread" });
const core = new CopilotKitCore({ agents__unsafe_dev_only: { default: agent } });
await assert.rejects(
runConversationTurn(
"default",
() => core.connectAgent({ agent }),
(onError) => core.subscribe({ onError }),
[replayedRunError],
),
/Thread service unavailable/,
);
});
Loading