diff --git a/apps/mobile/src/connection/platform.ts b/apps/mobile/src/connection/platform.ts
index 90a4d12226a2..7f1e3e5c344d 100644
--- a/apps/mobile/src/connection/platform.ts
+++ b/apps/mobile/src/connection/platform.ts
@@ -27,6 +27,7 @@ import * as MobileStorage from "../persistence/mobile-storage";
import { appAtomRegistry } from "../state/atom-registry";
import { clearThreadOutboxEnvironment } from "../state/thread-outbox-removal";
import { clearComposerDraftsEnvironment } from "../state/use-composer-drafts";
+import { clearThreadComposerErrorsForEnvironment } from "../state/thread-composer-error";
import { mobileApplicationActiveWakeup } from "./app-state-wakeups";
import { connectionStorageLayer } from "./storage";
@@ -264,6 +265,7 @@ const environmentOwnedDataCleanupLayer = Layer.succeed(
[
Effect.promise(() => clearThreadOutboxEnvironment(environmentId)),
Effect.promise(() => clearComposerDraftsEnvironment(environmentId)),
+ Effect.sync(() => clearThreadComposerErrorsForEnvironment(environmentId)),
],
{ concurrency: "unbounded", discard: true },
).pipe(
diff --git a/apps/mobile/src/features/threads/ComposerErrorNotice.tsx b/apps/mobile/src/features/threads/ComposerErrorNotice.tsx
new file mode 100644
index 000000000000..37ea6078f2ab
--- /dev/null
+++ b/apps/mobile/src/features/threads/ComposerErrorNotice.tsx
@@ -0,0 +1,55 @@
+import { useEffect } from "react";
+import { AccessibilityInfo, Platform, Pressable, View } from "react-native";
+
+import { AppText as Text } from "../../components/AppText";
+import { SymbolView } from "../../components/AppSymbol";
+
+/** Why the thread's last message did not send, above the composer until dismissed. */
+export function ComposerErrorNotice({
+ message,
+ onDismiss,
+}: {
+ readonly message: string;
+ readonly onDismiss: () => void;
+}) {
+ // accessibilityLiveRegion below only reaches TalkBack; VoiceOver needs an
+ // explicit announcement.
+ useEffect(() => {
+ if (Platform.OS === "ios") {
+ AccessibilityInfo.announceForAccessibility(message);
+ }
+ }, [message]);
+ return (
+
+
+
+
+ {message}
+
+
+
+
+
+
+ );
+}
diff --git a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx
index 013bf900912d..886c717fabc8 100644
--- a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx
+++ b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx
@@ -1,3 +1,4 @@
+import { useAtomValue } from "@effect/atom-react";
import { useThreadReportedModelSelection } from "../../state/entities";
import { UsageLimitRecoveryCard } from "./UsageLimitRecoveryCard";
import { useNavigation } from "@react-navigation/native";
@@ -93,6 +94,10 @@ import { useEnvironmentQuery } from "../../state/query";
import { threadDevicePreviews } from "../devices/threadDevicePreviews";
import type { QueuedThreadMessage } from "../../state/thread-outbox-model";
import { scopedThreadKey } from "../../lib/scopedEntities";
+import {
+ clearThreadComposerError,
+ threadComposerErrorsAtom,
+} from "../../state/thread-composer-error";
import { threadEnvironment } from "../../state/threads";
import { useAtomCommand } from "../../state/use-atom-command";
import { useDelayedStatus } from "../../lib/useDelayedStatus";
@@ -104,6 +109,7 @@ import type {
ThreadFeedLatestRun,
} from "../../lib/threadActivity";
import { PendingApprovalCard } from "./PendingApprovalCard";
+import { ComposerErrorNotice } from "./ComposerErrorNotice";
import { ComposerFeedback } from "./ComposerFeedback";
import { ComposerUsageLimits } from "./ComposerUsageLimits";
import { PendingUserInputCard } from "./PendingUserInputCard";
@@ -362,6 +368,7 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread
const navigationHeaderHeight = useContext(HeaderHeightContext) || insets.top + 44;
const agentLabel = `${props.selectedThread.modelSelection.instanceId} agent`;
const selectedThreadKey = scopedThreadKey(props.environmentId, props.selectedThread.id);
+ const composerError = useAtomValue(threadComposerErrorsAtom)[selectedThreadKey]?.message ?? null;
const queuedCount = useThreadQueuedCount({
environmentId: props.environmentId,
threadId: props.selectedThread.id,
@@ -1190,6 +1197,18 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread
onDismiss={() => props.onDismissFeedback(submission.id)}
/>
))}
+ {composerError !== null ? (
+
+ clearThreadComposerError(selectedThreadKey)}
+ />
+
+ ) : null}
{usageLimitsReport && activeUserInputRequestId === null ? (
{
+ appAtomRegistry.set(threadComposerErrorsAtom, {});
+});
+
+describe("thread composer errors", () => {
+ it("clears a message-scoped error only for that message", () => {
+ setThreadComposerError("environment-1:thread-1", "rejected", "message-1");
+
+ clearThreadComposerError("environment-1:thread-1", "message-2");
+ expect(appAtomRegistry.get(threadComposerErrorsAtom)["environment-1:thread-1"]?.message).toBe(
+ "rejected",
+ );
+
+ clearThreadComposerError("environment-1:thread-1");
+ expect(appAtomRegistry.get(threadComposerErrorsAtom)).toEqual({});
+ });
+
+ it("clears every thread's error for a removed environment and no others", () => {
+ setThreadComposerError("environment-1:thread-1", "one");
+ setThreadComposerError("environment-1:thread-2", "two");
+ setThreadComposerError("environment-10:thread-1", "other environment");
+
+ clearThreadComposerErrorsForEnvironment("environment-1");
+
+ expect(Object.keys(appAtomRegistry.get(threadComposerErrorsAtom))).toEqual([
+ "environment-10:thread-1",
+ ]);
+ });
+});
diff --git a/apps/mobile/src/state/thread-composer-error.ts b/apps/mobile/src/state/thread-composer-error.ts
new file mode 100644
index 000000000000..bf51139ede97
--- /dev/null
+++ b/apps/mobile/src/state/thread-composer-error.ts
@@ -0,0 +1,56 @@
+import { Atom } from "effect/unstable/reactivity";
+
+import { appAtomRegistry } from "./atom-registry";
+
+interface ThreadComposerError {
+ readonly message: string;
+ /** The queued message this error is about, when it is about one. */
+ readonly messageId: string | null;
+}
+
+/**
+ * Why a thread's last message did not go out, shown above that thread's
+ * composer. The outbox drain can reject a message after the user has left the
+ * thread, so the reason is kept per thread until they dismiss it, send again,
+ * or the message it describes is delivered after all. Keyed by `scopedThreadKey`.
+ */
+export const threadComposerErrorsAtom = Atom.make>>(
+ {},
+).pipe(Atom.keepAlive, Atom.withLabel("mobile:thread-composer-errors"));
+
+export function setThreadComposerError(
+ threadKey: string,
+ message: string,
+ messageId: string | null = null,
+): void {
+ appAtomRegistry.set(threadComposerErrorsAtom, {
+ ...appAtomRegistry.get(threadComposerErrorsAtom),
+ [threadKey]: { message, messageId },
+ });
+}
+
+/** With `messageId`, clears only an error that describes that message. */
+export function clearThreadComposerError(threadKey: string, messageId?: string): void {
+ const current = appAtomRegistry.get(threadComposerErrorsAtom);
+ const error = current[threadKey];
+ if (!error || (messageId !== undefined && error.messageId !== messageId)) {
+ return;
+ }
+ const next = { ...current };
+ delete next[threadKey];
+ appAtomRegistry.set(threadComposerErrorsAtom, next);
+}
+
+export function clearThreadComposerErrorsForEnvironment(environmentId: string): void {
+ const current = appAtomRegistry.get(threadComposerErrorsAtom);
+ const prefix = `${environmentId}:`;
+ const keys = Object.keys(current).filter((key) => key.startsWith(prefix));
+ if (keys.length === 0) {
+ return;
+ }
+ const next = { ...current };
+ for (const key of keys) {
+ delete next[key];
+ }
+ appAtomRegistry.set(threadComposerErrorsAtom, next);
+}
diff --git a/apps/mobile/src/state/use-thread-composer-state.ts b/apps/mobile/src/state/use-thread-composer-state.ts
index 1f9eed3d675c..9688eddc65a4 100644
--- a/apps/mobile/src/state/use-thread-composer-state.ts
+++ b/apps/mobile/src/state/use-thread-composer-state.ts
@@ -96,6 +96,7 @@ import {
useQueuedRunEdit,
} from "./queued-run-edit";
import { setPendingConnectionError } from "../state/use-remote-environment-registry";
+import { clearThreadComposerError, setThreadComposerError } from "./thread-composer-error";
import {
useSelectedThreadProjection,
useSelectedThreadVisibleTurnItems,
@@ -452,7 +453,8 @@ export function useThreadComposerState() {
});
}
endQueuedRunEdit(selectedThreadKey, { deferAttachmentCleanup: keepable });
- setPendingConnectionError(
+ setThreadComposerError(
+ selectedThreadKey,
keepable
? "That message already started. Your edit is back in the composer."
: "That message already started, so the edit was discarded.",
@@ -670,6 +672,8 @@ export function useThreadComposerState() {
const metadata = makeQueuedMessageMetadata();
const messageId = MessageId.make(metadata.messageId);
+ // A new send supersedes the reason the previous one bounced back.
+ clearThreadComposerError(threadKey);
// Enqueue publishes the queued atom synchronously (the durable write
// happens behind it), so clearing the draft here gives send feedback on
// the tap frame instead of after file I/O. If the write fails the message
@@ -706,7 +710,8 @@ export function useThreadComposerState() {
attachments: [],
});
appendComposerDraftAttachments(threadKey, attachments, { allowOverflow: true });
- setPendingConnectionError(
+ setThreadComposerError(
+ threadKey,
error instanceof Error ? error.message : "Failed to save the queued message.",
);
},
diff --git a/apps/mobile/src/state/use-thread-outbox-drain.test.ts b/apps/mobile/src/state/use-thread-outbox-drain.test.ts
index 890b29f9c3e8..a7b807e6267a 100644
--- a/apps/mobile/src/state/use-thread-outbox-drain.test.ts
+++ b/apps/mobile/src/state/use-thread-outbox-drain.test.ts
@@ -19,7 +19,6 @@ const harness = vi.hoisted(() => ({
removePersistedFile: vi.fn(async () => undefined),
removeOutboxMessage: vi.fn(async (_message: QueuedThreadMessage) => undefined),
prepareTurnAttachments: vi.fn(),
- setPendingConnectionError: vi.fn(),
draftFile: (() => {
let document = "";
let writeError: Error | null = null;
@@ -106,7 +105,6 @@ vi.mock("./use-thread-outbox", async () => {
});
vi.mock("./use-remote-environment-registry", () => ({
- setPendingConnectionError: harness.setPendingConnectionError,
useRemoteConnectionStatus: () => ({ connectedEnvironments: [] }),
}));
@@ -137,6 +135,11 @@ import {
clearPendingThreadCreationOutcome,
pendingThreadCreationOutcomesAtom,
} from "./pending-thread-creation";
+import {
+ clearThreadComposerError,
+ setThreadComposerError,
+ threadComposerErrorsAtom,
+} from "./thread-composer-error";
import type { QueuedThreadMessage } from "./thread-outbox-model";
import * as composerDrafts from "./use-composer-drafts";
import { recoverFailedThreadDraft } from "./recover-failed-thread-draft";
@@ -209,11 +212,11 @@ afterEach(() => {
appAtomRegistry.set(composerDrafts.composerCloudDraftsAtom, { accountId: null, signedOut: {} });
appAtomRegistry.set(editingQueuedMessageIdsAtom, {});
appAtomRegistry.set(pendingThreadCreationOutcomesAtom, {});
+ appAtomRegistry.set(threadComposerErrorsAtom, {});
harness.draftFile.setWriteError(null);
harness.removePersistedFile.mockClear();
harness.removeOutboxMessage.mockClear();
harness.prepareTurnAttachments.mockReset();
- harness.setPendingConnectionError.mockClear();
});
describe("thread outbox attachment preparation", () => {
@@ -442,6 +445,25 @@ describe("thread outbox drain delivery cleanup", () => {
expect(appAtomRegistry.get(acknowledgedThreadMessagesAtom)).toEqual([message]);
});
+ it("clears an error about the delivered message but keeps one about another message", async () => {
+ const threadKey = "environment-1:thread-1";
+ const retried = queuedMessage({ messageId: "message-retried", text: "retried" });
+ await harness.manager.enqueue(retried);
+ // A failed recovery left this message queued with an error about it.
+ setThreadComposerError(threadKey, "could not be restored", retried.messageId);
+
+ await completeQueuedMessageDelivery(retried, harness.manager.revisionOf(retried.messageId));
+ expect(appAtomRegistry.get(threadComposerErrorsAtom)[threadKey]).toBeUndefined();
+
+ const other = queuedMessage({ messageId: "message-other", text: "other" });
+ await harness.manager.enqueue(other);
+ // The thread's error explains a different, rejected message still in the draft.
+ setThreadComposerError(threadKey, "rejected", "message-rejected");
+
+ await completeQueuedMessageDelivery(other, harness.manager.revisionOf(other.messageId));
+ expect(appAtomRegistry.get(threadComposerErrorsAtom)[threadKey]?.message).toBe("rejected");
+ });
+
it("keeps a delivered message when its editor opens during storage removal", async () => {
const message = queuedMessage({
messageId: "message-editor-removal-race",
@@ -697,7 +719,8 @@ describe("thread outbox recovery rollback", () => {
},
});
expect(remainingMessages()).toEqual([]);
- expect(harness.setPendingConnectionError).toHaveBeenCalledWith("rejected by server");
+ // The creation's failure card shows the reason; the composer is hidden.
+ expect(appAtomRegistry.get(threadComposerErrorsAtom)).toEqual({});
// The thread screen opened for this creation reads the failure from here.
expect(
appAtomRegistry.get(pendingThreadCreationOutcomesAtom)[
@@ -706,6 +729,37 @@ describe("thread outbox recovery rollback", () => {
).toEqual({ kind: "failed", message, reason: "rejected by server" });
});
+ it("drops an earlier recovery error once a rejected new task is restored", async () => {
+ const message: QueuedThreadMessage = {
+ ...queuedMessage({ messageId: "message-creation-retried", text: "new task text" }),
+ modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5.6-sol" },
+ creation: {
+ projectId: ProjectId.make("project-1"),
+ workspaceMode: "local",
+ branch: null,
+ worktreePath: null,
+ },
+ };
+ await harness.manager.enqueue(message);
+ harness.draftFile.setWriteError(new Error("disk full"));
+ await expect(restoreRejectedQueuedMessage(message, "rejected by server")).resolves.toBe(
+ "retry",
+ );
+ const threadKey = `${message.environmentId}:${message.threadId}`;
+ expect(appAtomRegistry.get(threadComposerErrorsAtom)[threadKey]?.messageId).toBe(
+ message.messageId,
+ );
+
+ harness.draftFile.setWriteError(null);
+ await expect(restoreRejectedQueuedMessage(message, "rejected by server")).resolves.toBe(
+ "restored",
+ );
+
+ // The failure card carries the reason now; nothing stale sits above it.
+ expect(appAtomRegistry.get(threadComposerErrorsAtom)).toEqual({});
+ expect(appAtomRegistry.get(pendingThreadCreationOutcomesAtom)[threadKey]?.kind).toBe("failed");
+ });
+
it("keeps a failed outcome until its thread screen consumes it", async () => {
const message: QueuedThreadMessage = {
...queuedMessage({ messageId: "message-creation-kept", text: "new task text" }),
@@ -734,6 +788,46 @@ describe("thread outbox recovery rollback", () => {
await expect(restoreRejectedQueuedMessage(message, "rejected")).resolves.toBe("restored");
expect(appAtomRegistry.get(pendingThreadCreationOutcomesAtom)).toEqual({});
+ // The thread screen shows why the message came back into the composer.
+ expect(
+ appAtomRegistry.get(threadComposerErrorsAtom)[`${message.environmentId}:${message.threadId}`],
+ ).toEqual({ message: "rejected", messageId: message.messageId });
+ });
+
+ it("leaves no error when the restored text is resent before recovery finishes", async () => {
+ const message = queuedMessage({ messageId: "message-resent", text: "resend me" });
+ const threadKey = `${message.environmentId}:${message.threadId}`;
+ await harness.manager.enqueue(message);
+ // The user sees the restored text as soon as the merge publishes it and
+ // sends it again while the recovery is still awaiting persistence.
+ const unsubscribe = appAtomRegistry.subscribe(composerDrafts.composerDraftsAtom, (drafts) => {
+ if (drafts[threadKey]?.text === "resend me") {
+ unsubscribe();
+ clearThreadComposerError(threadKey);
+ void composerDrafts.clearComposerDraftContent(threadKey);
+ }
+ });
+
+ await expect(restoreRejectedQueuedMessage(message, "rejected")).resolves.toBe("restored");
+
+ expect(appAtomRegistry.get(threadComposerErrorsAtom)).toEqual({});
+ });
+
+ it("withdraws the error when an edit makes the recovery back out", async () => {
+ const message = queuedMessage({ messageId: "message-edited-mid-recovery", text: "edit me" });
+ const threadKey = `${message.environmentId}:${message.threadId}`;
+ await harness.manager.enqueue(message);
+ const unsubscribe = appAtomRegistry.subscribe(composerDrafts.composerDraftsAtom, (drafts) => {
+ if (drafts[threadKey]?.text === "edit me") {
+ unsubscribe();
+ appAtomRegistry.set(editingQueuedMessageIdsAtom, { [message.messageId]: true });
+ }
+ });
+
+ await expect(restoreRejectedQueuedMessage(message, "rejected")).resolves.toBe("deferred");
+
+ expect(appAtomRegistry.get(threadComposerErrorsAtom)).toEqual({});
+ expect(remainingMessages()).toEqual([message]);
});
it("rolls a failed recovery merge back so the retry cannot duplicate the text", async () => {
@@ -759,6 +853,6 @@ describe("thread outbox recovery rollback", () => {
"typed offline\n\nqueued text",
);
expect(remainingMessages()).toEqual([]);
- expect(harness.setPendingConnectionError).toHaveBeenCalledWith("too large");
+ expect(appAtomRegistry.get(threadComposerErrorsAtom)[draftKey]?.message).toBe("too large");
});
});
diff --git a/apps/mobile/src/state/use-thread-outbox-drain.ts b/apps/mobile/src/state/use-thread-outbox-drain.ts
index 454980d404a5..db4d5b14f590 100644
--- a/apps/mobile/src/state/use-thread-outbox-drain.ts
+++ b/apps/mobile/src/state/use-thread-outbox-drain.ts
@@ -80,10 +80,8 @@ import {
useThreadOutboxMessages,
useThreadOutboxShellStatuses,
} from "./use-thread-outbox";
-import {
- setPendingConnectionError,
- useRemoteConnectionStatus,
-} from "./use-remote-environment-registry";
+import { clearThreadComposerError, setThreadComposerError } from "./thread-composer-error";
+import { useRemoteConnectionStatus } from "./use-remote-environment-registry";
// Ordinary offline behavior (a socket dropping mid-request, a retryable
// attachment upload failure) must not spam `console.warn` on every backoff
@@ -282,6 +280,12 @@ export async function completeQueuedMessageDelivery(
queuedMessage: QueuedThreadMessage,
deliveryRevision: number,
): Promise<"removed" | "edited" | "failed"> {
+ // The server took it after all: an error left by an earlier failed recovery
+ // of this same message no longer applies.
+ clearThreadComposerError(
+ scopedThreadKey(queuedMessage.environmentId, queuedMessage.threadId),
+ queuedMessage.messageId,
+ );
try {
await removeDeliveredCloudQueuedMessage(queuedMessage).catch((error) => {
console.warn("[thread-outbox] could not update sign-out snapshot after delivery", {
@@ -438,6 +442,7 @@ export async function restoreRejectedQueuedMessage(
message: string,
): Promise<"restored" | "deferred" | "blocked" | "retry"> {
const draftKey = recoveryDraftKey(queuedMessage);
+ const threadKey = scopedThreadKey(queuedMessage.environmentId, queuedMessage.threadId);
// Set once the merge publishes, cleared once the queued message is removed.
// The catch below uses it to take the merged content back out, so a retry
// after a mid-recovery failure cannot append the recovered text again.
@@ -467,12 +472,21 @@ export async function restoreRejectedQueuedMessage(
(attachment) => !existingAttachmentIds.has(attachment.id),
).length;
if (existingAttachmentIds.size + addedAttachmentCount > PROVIDER_SEND_TURN_MAX_ATTACHMENTS) {
- setPendingConnectionError(
+ setThreadComposerError(
+ threadKey,
`Remove attachments from the draft before restoring this message. Messages can contain at most ${PROVIDER_SEND_TURN_MAX_ATTACHMENTS} attachments.`,
+ queuedMessage.messageId,
);
return "blocked";
}
+ // Shown before the merge publishes the text, so a resend of that text,
+ // which can happen while this recovery still awaits persistence, clears
+ // it. Withdrawn below wherever the recovery backs out.
+ const withdrawError = () => clearThreadComposerError(threadKey, queuedMessage.messageId);
+ if (!queuedMessage.creation) {
+ setThreadComposerError(threadKey, message, queuedMessage.messageId);
+ }
let mergedDraft: ComposerDraft;
try {
stampRecoveryDraftProject(queuedMessage, draftKey);
@@ -492,6 +506,7 @@ export async function restoreRejectedQueuedMessage(
rollback = { snapshot: originalDraft, merged: mergedDraft };
}
if (appAtomRegistry.get(editingQueuedMessageIdsAtom)[queuedMessage.messageId]) {
+ withdrawError();
await undoComposerDraftMerge(draftKey, originalDraft, mergedDraft);
return "deferred";
}
@@ -520,6 +535,7 @@ export async function restoreRejectedQueuedMessage(
!(await confirmThreadOutboxMessageQueued(queuedMessage)) ||
appAtomRegistry.get(editingQueuedMessageIdsAtom)[queuedMessage.messageId]
) {
+ withdrawError();
await undoComposerDraftMerge(draftKey, originalDraft, restoredDraft);
return "deferred";
}
@@ -532,6 +548,7 @@ export async function restoreRejectedQueuedMessage(
() => !appAtomRegistry.get(editingQueuedMessageIdsAtom)[queuedMessage.messageId],
))
) {
+ withdrawError();
await undoComposerDraftMerge(draftKey, originalDraft, restoredDraft);
return "deferred";
}
@@ -539,15 +556,17 @@ export async function restoreRejectedQueuedMessage(
// must never be rolled back.
rollback = null;
if (queuedMessage.creation) {
+ // The failure card shows the reason, so an error left by an earlier
+ // failed attempt at this recovery no longer applies.
+ withdrawError();
// The thread screen for this creation is likely open; it reads the
- // outcome to offer reopening the restored draft.
+ // outcome to offer reopening the restored draft, and shows the reason.
recordPendingThreadCreationOutcome({
kind: "failed",
message: queuedMessage,
reason: message,
});
}
- setPendingConnectionError(message);
return "restored";
} catch (error) {
if (rollback !== null) {
@@ -561,8 +580,10 @@ export async function restoreRejectedQueuedMessage(
);
}
console.warn("[thread-outbox] failed to restore an undeliverable message", error);
- setPendingConnectionError(
+ setThreadComposerError(
+ threadKey,
error instanceof Error ? error.message : "The unsent message could not be restored.",
+ queuedMessage.messageId,
);
return "retry";
}