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"; }