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
65 changes: 53 additions & 12 deletions apps/mobile/src/chat.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,13 @@ import { BackgroundUpdates } from "./background-updates";
import { BrowserRunContext, BrowserToolCard } from "./browser-tool-card";
import { BrowserThreadCard } from "./computer";
import { ConversationQueue, type QueuedMessage } from "./conversation-queue";
import { replayedRunError, runConversationTurn } from "./conversation-run";
import {
ConversationTurnError,
replayedRunError,
runConversationTurn,
showsRunError,
threadLocked,
} from "./conversation-run";
import { DesktopToolCard } from "./desktop-tool-card";
import { confirmedJevSelection, displayJevUserMessage, latestJevPanelId } from "./jev-actions";
import { JevInteractionContext, JevToolCard } from "./jev-tool-card";
Expand Down Expand Up @@ -222,6 +228,8 @@ export function ChatScreen({
const [busy, setBusy] = useState(false);
const [error, setError] = useState("");
const [loaded, setLoaded] = useState(false);
// Replayed messages show before connectAgent returns; new turns wait until it does.
const [syncing, setSyncing] = useState(false);
const [picking, setPicking] = useState(false);
const [attachments, setAttachments] = useState<string[]>([]);
const list = useRef<ScrollView>(null);
Expand All @@ -238,6 +246,8 @@ export function ChatScreen({
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);
// Set while a queued message is being sent, so a lock refusal holds it instead of failing.
const queuedTurn = useRef(false);
const [keyboardPadding, setKeyboardPadding] = useState(0);

useEffect(() => {
Expand All @@ -259,6 +269,7 @@ export function ChatScreen({
let active = true;
setHistoryError("");
setLoaded(false);
setSyncing(false);
const replay = agent.subscribe({
onMessagesChanged: ({ messages }) => {
if (active && richThreads && messages.length) setLoaded(true);
Expand All @@ -269,7 +280,10 @@ export function ChatScreen({
if (richThreads) {
if (selection.existing) {
// Replaying history re-emits past RUN_ERROR events; only connection failures block loading.
// connectAgent returns once the thread is idle, so it also waits for a reply still
// running from before a reload. Sending earlier fails with a thread lock.
replaying.current += 1;
setSyncing(true);
try {
await runConversationTurn(
agentId,
Expand All @@ -279,6 +293,7 @@ export function ChatScreen({
);
} finally {
replaying.current -= 1;
if (active) setSyncing(false);
}
}
} else {
Expand Down Expand Up @@ -306,11 +321,13 @@ export function ChatScreen({
if (!richThreads) await api.request("/api/conversation", { messages: agent.messages }, "PUT");
setSaveError("");
}, [agent, api, richThreads]);
/** Runs one turn; "held" means a queued message was refused by a lock and put back on hold. */
const run = useCallback(
async (message?: QueuedMessage) => {
if (runLock.current || agent.isRunning || !isReady || !loaded)
async (message?: QueuedMessage): Promise<"held" | undefined> => {
if (runLock.current || agent.isRunning || !isReady || !loaded || syncing)
throw new Error("The conversation is not ready yet.");
runLock.current = true;
queuedTurn.current = Boolean(message);
setBusy(true);
setError("");
if (message) agent.addMessage({ id: message.id, role: "user", content: message.text });
Expand All @@ -321,7 +338,18 @@ export function ChatScreen({
(onError) => copilotkit.subscribe({ onError }),
);
await Promise.all([refresh(), refreshAgent()]);
} catch (e) {
if (message && e instanceof ConversationTurnError && e.code === threadLocked) {
// The server refused the turn before running it (another reply holds the thread),
// so keep the message unsent and on hold instead of reporting a failed turn.
agent.setMessages(agent.messages.filter((m) => m.id !== message.id));
queue.restore(message);
queue.pause();
return "held";
}
throw e;
} finally {
queuedTurn.current = false;
try {
await saveHistory();
} catch (e) {
Expand All @@ -335,26 +363,39 @@ export function ChatScreen({
}
}
},
[agent, agentId, copilotkit, isReady, loaded, refresh, refreshAgent, saveHistory, queue],
[
agent,
agentId,
copilotkit,
isReady,
loaded,
syncing,
refresh,
refreshAgent,
saveHistory,
queue,
],
);
const runQueued = useCallback(
async (message: QueuedMessage) => {
let held = false;
try {
await run(message);
choiceCompletions.current.get(message.id)?.resolve();
held = (await run(message)) === "held";
if (!held) choiceCompletions.current.get(message.id)?.resolve();
} catch (error) {
choiceCompletions.current.get(message.id)?.reject(error);
throw error;
} finally {
choiceCompletions.current.delete(message.id);
// A held message is sent again later; its choice settles then.
if (!held) choiceCompletions.current.delete(message.id);
}
},
[run],
);
const flush = useCallback(() => {
if (!loaded || !isReady || runLock.current || agent.isRunning) return;
if (!loaded || syncing || !isReady || runLock.current || agent.isRunning) return;
void queue.flush(runQueued).catch((e) => setError(e instanceof Error ? e.message : String(e)));
}, [agent, isReady, loaded, queue, runQueued]);
}, [agent, isReady, loaded, syncing, queue, runQueued]);
const enqueue = useCallback(
(text: string) => {
queue.enqueue({ id: `user-${Date.now()}-${Math.random().toString(36).slice(2, 8)}`, text });
Expand Down Expand Up @@ -397,8 +438,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 during = { replaying: replaying.current > 0, queuedTurn: queuedTurn.current };
if (!showsRunError(event, during)) return;
const failure = event.error instanceof Error ? event.error : new Error(String(event.error));
setError(failure.message);
},
Expand Down Expand Up @@ -700,7 +741,7 @@ export function ChatScreen({
<Button
style={{ alignSelf: "flex-start" }}
icon={RotateCcw}
disabled={busy || agent.isRunning || !loaded || !isReady}
disabled={busy || agent.isRunning || !loaded || syncing || !isReady}
onPress={() => {
void run()
.then(() => {
Expand Down
6 changes: 6 additions & 0 deletions apps/mobile/src/conversation-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,12 @@ export class ConversationQueue {
remove(id: string) {
this.update({ pending: this.state.pending.filter((message) => message.id !== id) });
}
/** Put back, first in line, a message the server refused before running it. */
restore(message: QueuedMessage) {
this.update({
pending: [message, ...this.state.pending.filter((queued) => queued.id !== message.id)],
});
}
pause() {
this.update({ paused: true });
}
Expand Down
40 changes: 38 additions & 2 deletions apps/mobile/src/conversation-run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,18 +3,54 @@ 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";

/** The thread is busy with another run, so the server refused this one before starting it. */
export const threadLocked = "agent_thread_locked";

/** CopilotKit reports a lock as agent_thread_locked, then again as agent_run_failed. */
export function isThreadLocked(event: Pick<RunError, "error" | "code">) {
return (
event.code === threadLocked ||
(event.error instanceof Error && event.error.name === "AgentThreadLockedError")
);
}

/** A failed turn, carrying the CopilotKit error code that reported it. */
export class ConversationTurnError extends Error {
constructor(
readonly cause: Error,
readonly code?: string,
) {
super(cause.message);
this.name = cause.name;
}
}

/** Whether the chat's error banner shows an error reported while these are in progress. */
export function showsRunError(
event: Pick<RunError, "error" | "code">,
during: { replaying: boolean; queuedTurn: boolean },
) {
// A failed turn saved in history is already over; it is not a failure of this session.
if (during.replaying && event.code === replayedRunError) return false;
// A queued message refused by a lock goes back on hold. Retry is not a queued turn, so it shows.
if (during.queuedTurn && isThreadLocked(event)) return false;
return true;
}

/** 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;
let failure: ConversationTurnError | 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));
const error = event.error instanceof Error ? event.error : new Error(String(event.error));
// The first report is the specific one; CopilotKit follows it with a generic agent_run_failed.
failure ??= new ConversationTurnError(error, isThreadLocked(event) ? threadLocked : event.code);
});
try {
await execute();
Expand Down
11 changes: 11 additions & 0 deletions apps/mobile/test/conversation-queue.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,3 +72,14 @@ function deferred() {
});
return { promise, resolve };
}

test("a message the server refused is restored first in line, once", async () => {
const queue = new ConversationQueue();
queue.enqueue({ id: "later", text: "Later" });
queue.restore({ id: "refused", text: "Refused" });
queue.restore({ id: "refused", text: "Refused" });
assert.deepEqual(
queue.getSnapshot().pending.map((message) => message.id),
["refused", "later"],
);
});
64 changes: 62 additions & 2 deletions tests/conversation-sdk.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,16 @@ 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 { AgentThreadLockedError, CopilotKitCore } from "@copilotkit/core";
import { of, throwError } from "rxjs";
import { ConversationQueue } from "../apps/mobile/src/conversation-queue.ts";
import { replayedRunError, runConversationTurn } from "../apps/mobile/src/conversation-run.ts";
import {
ConversationTurnError,
replayedRunError,
runConversationTurn,
showsRunError,
threadLocked,
} 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 @@ -84,3 +90,57 @@ test("a connection failure still blocks loading the conversation", async () => {
/Thread service unavailable/,
);
});

test("a turn refused because the thread is locked reports the lock code", async () => {
class LockedAgent extends AbstractAgent {
run() {
return throwError(() => new AgentThreadLockedError("thread"));
}
}
const agent = new LockedAgent({ agentId: "default", threadId: "thread" });
const core = new CopilotKitCore({ agents__unsafe_dev_only: { default: agent } });
await assert.rejects(
runConversationTurn(
"default",
() => core.runAgent({ agent }),
(onError) => core.subscribe({ onError }),
),
(error: unknown) =>
error instanceof ConversationTurnError &&
error.code === threadLocked &&
error.message === "Thread thread is locked",
);
});

test("with replayed errors ignored, a lock refusal is still the reported failure", async () => {
const locked = new Error("Thread thread is locked");
locked.name = "AgentThreadLockedError";
await assert.rejects(
runConversationTurn(
"default",
async () => {},
(onError) => {
onError({ error: new Error("Saved failure"), code: replayedRunError });
onError({ error: locked, code: threadLocked });
onError({ error: new Error("Run failed"), code: "agent_run_failed" });
return { unsubscribe() {} };
},
[replayedRunError],
),
(error: unknown) => error instanceof ConversationTurnError && error.code === threadLocked,
);
});

test("the error banner hides replayed errors and held lock refusals, but not Retry's lock", () => {
const replayed = { error: new Error("Saved failure"), code: replayedRunError };
const locked = { error: new Error("Thread thread is locked"), code: threadLocked };
const idle = { replaying: false, queuedTurn: false };
assert.equal(showsRunError(replayed, { ...idle, replaying: true }), false);
assert.equal(showsRunError(replayed, idle), true);
assert.equal(showsRunError(locked, { ...idle, queuedTurn: true }), false);
assert.equal(showsRunError(locked, idle), true, "Retry is not a queued turn, so its lock shows");
assert.equal(
showsRunError({ error: new Error("Network down") }, { ...idle, queuedTurn: true }),
true,
);
});
Loading