From ba258ca093ff5c06982ba1c561bd5e60c689c0d5 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 20 Jun 2026 05:41:44 -0700 Subject: [PATCH] fix(mobile): structure thread outbox failures Co-authored-by: codex --- .../mobile/src/state/thread-outbox-manager.ts | 87 ++++++++++++++-- .../mobile/src/state/thread-outbox-storage.ts | 98 +++++++++++++++---- apps/mobile/src/state/thread-outbox.test.ts | 53 +++++++++- 3 files changed, 208 insertions(+), 30 deletions(-) diff --git a/apps/mobile/src/state/thread-outbox-manager.ts b/apps/mobile/src/state/thread-outbox-manager.ts index 477cb1273a3..7762e6cdf78 100644 --- a/apps/mobile/src/state/thread-outbox-manager.ts +++ b/apps/mobile/src/state/thread-outbox-manager.ts @@ -1,4 +1,5 @@ -import type { EnvironmentId, MessageId } from "@t3tools/contracts"; +import { EnvironmentId, MessageId, ThreadId } from "@t3tools/contracts"; +import * as Schema from "effect/Schema"; import { Atom, type AtomRegistry } from "effect/unstable/reactivity"; import { @@ -8,6 +9,27 @@ import { } from "./thread-outbox-model"; import type { ThreadOutboxStorage } from "./thread-outbox-storage"; +export class ThreadOutboxManagerError extends Schema.TaggedErrorClass()( + "ThreadOutboxManagerError", + { + operation: Schema.Literals([ + "load", + "enqueue", + "remove", + "clear-environment-load", + "clear-environment-remove", + ]), + environmentId: Schema.NullOr(EnvironmentId), + threadId: Schema.NullOr(ThreadId), + messageId: Schema.NullOr(MessageId), + cause: Schema.Defect(), + }, +) { + override get message(): string { + return `Thread outbox operation ${this.operation} failed for environment ${this.environmentId ?? "unknown"}, thread ${this.threadId ?? "unknown"}, message ${this.messageId ?? "unknown"}.`; + } +} + export interface ThreadOutboxManagerOptions { readonly registry: AtomRegistry.AtomRegistry; readonly storage: ThreadOutboxStorage; @@ -49,22 +71,51 @@ export function createThreadOutboxManager(options: ThreadOutboxManagerOptions) { loadPromise = serialize(async () => { const persistedMessages = await options.storage.load(); setMessages([...persistedMessages, ...currentMessages()]); - }).catch((error) => { + }).catch((cause) => { loadPromise = null; - warn("[thread-outbox] failed to load persisted messages", error); + warn( + "[thread-outbox] failed to load persisted messages", + new ThreadOutboxManagerError({ + operation: "load", + environmentId: null, + threadId: null, + messageId: null, + cause, + }), + ); }); return loadPromise; }; const enqueue = (message: QueuedThreadMessage): Promise => serialize(async () => { - await options.storage.write(message); + try { + await options.storage.write(message); + } catch (cause) { + throw new ThreadOutboxManagerError({ + operation: "enqueue", + environmentId: message.environmentId, + threadId: message.threadId, + messageId: message.messageId, + cause, + }); + } setMessages([...currentMessages(), message]); }); const remove = (message: QueuedThreadMessage): Promise => serialize(async () => { - await options.storage.remove(message); + try { + await options.storage.remove(message); + } catch (cause) { + throw new ThreadOutboxManagerError({ + operation: "remove", + environmentId: message.environmentId, + threadId: message.threadId, + messageId: message.messageId, + cause, + }); + } setMessages( currentMessages().filter((candidate) => candidate.messageId !== message.messageId), ); @@ -72,8 +123,17 @@ export function createThreadOutboxManager(options: ThreadOutboxManagerOptions) { const clearEnvironment = (environmentId: EnvironmentId): Promise => serialize(async () => { - const persisted = await options.storage.load().catch((error) => { - warn("[thread-outbox] failed to load messages while clearing environment", error); + const persisted = await options.storage.load().catch((cause) => { + warn( + "[thread-outbox] failed to load messages while clearing environment", + new ThreadOutboxManagerError({ + operation: "clear-environment-load", + environmentId, + threadId: null, + messageId: null, + cause, + }), + ); return []; }); const allMessages = flattenQueuedThreadMessages( @@ -88,8 +148,17 @@ export function createThreadOutboxManager(options: ThreadOutboxManagerOptions) { try { await options.storage.remove(message); removedMessageIds.add(message.messageId); - } catch (error) { - warn("[thread-outbox] failed to clear persisted message", error); + } catch (cause) { + warn( + "[thread-outbox] failed to clear persisted message", + new ThreadOutboxManagerError({ + operation: "clear-environment-remove", + environmentId: message.environmentId, + threadId: message.threadId, + messageId: message.messageId, + cause, + }), + ); } }), ); diff --git a/apps/mobile/src/state/thread-outbox-storage.ts b/apps/mobile/src/state/thread-outbox-storage.ts index e294aee4549..2003c220bad 100644 --- a/apps/mobile/src/state/thread-outbox-storage.ts +++ b/apps/mobile/src/state/thread-outbox-storage.ts @@ -1,4 +1,5 @@ -import type { MessageId } from "@t3tools/contracts"; +import { EnvironmentId, MessageId, ThreadId } from "@t3tools/contracts"; +import * as Schema from "effect/Schema"; import { decodeQueuedThreadMessage, @@ -8,6 +9,22 @@ import { const THREAD_OUTBOX_DIRECTORY = "thread-outbox"; +export class ThreadOutboxStorageError extends Schema.TaggedErrorClass()( + "ThreadOutboxStorageError", + { + operation: Schema.Literals(["load", "read-message", "write", "remove"]), + environmentId: Schema.NullOr(EnvironmentId), + threadId: Schema.NullOr(ThreadId), + messageId: Schema.NullOr(MessageId), + fileName: Schema.NullOr(Schema.String), + cause: Schema.Defect(), + }, +) { + override get message(): string { + return `Thread outbox storage operation ${this.operation} failed for environment ${this.environmentId ?? "unknown"}, thread ${this.threadId ?? "unknown"}, message ${this.messageId ?? "unknown"}, file ${this.fileName ?? "unknown"}.`; + } +} + export interface ThreadOutboxStorage { readonly load: () => Promise>; readonly write: (message: QueuedThreadMessage) => Promise; @@ -32,33 +49,78 @@ async function getMessageFile(messageId: MessageId) { export const expoThreadOutboxStorage: ThreadOutboxStorage = { load: async () => { - const { File } = await import("expo-file-system"); - const directory = await getOutboxDirectory(); const messages: QueuedThreadMessage[] = []; + try { + const { File } = await import("expo-file-system"); + const directory = await getOutboxDirectory(); - for (const entry of directory.list()) { - if (!(entry instanceof File) || !entry.name.endsWith(".json")) { - continue; - } - try { - messages.push(decodeQueuedThreadMessage(JSON.parse(await entry.text()) as unknown)); - } catch (error) { - console.warn("[thread-outbox] ignored invalid persisted message", entry.name, error); + for (const entry of directory.list()) { + if (!(entry instanceof File) || !entry.name.endsWith(".json")) { + continue; + } + try { + messages.push(decodeQueuedThreadMessage(JSON.parse(await entry.text()) as unknown)); + } catch (cause) { + console.warn( + "[thread-outbox] ignored invalid persisted message", + new ThreadOutboxStorageError({ + operation: "read-message", + environmentId: null, + threadId: null, + messageId: null, + fileName: entry.name, + cause, + }), + ); + } } + } catch (cause) { + throw new ThreadOutboxStorageError({ + operation: "load", + environmentId: null, + threadId: null, + messageId: null, + fileName: null, + cause, + }); } return messages; }, write: async (message) => { - const file = await getMessageFile(message.messageId); - if (!file.exists) { - file.create({ intermediates: true, overwrite: true }); + const fileName = messageFileName(message.messageId); + try { + const file = await getMessageFile(message.messageId); + if (!file.exists) { + file.create({ intermediates: true, overwrite: true }); + } + file.write(JSON.stringify(encodeQueuedThreadMessage(message))); + } catch (cause) { + throw new ThreadOutboxStorageError({ + operation: "write", + environmentId: message.environmentId, + threadId: message.threadId, + messageId: message.messageId, + fileName, + cause, + }); } - file.write(JSON.stringify(encodeQueuedThreadMessage(message))); }, remove: async (message) => { - const file = await getMessageFile(message.messageId); - if (file.exists) { - file.delete(); + const fileName = messageFileName(message.messageId); + try { + const file = await getMessageFile(message.messageId); + if (file.exists) { + file.delete(); + } + } catch (cause) { + throw new ThreadOutboxStorageError({ + operation: "remove", + environmentId: message.environmentId, + threadId: message.threadId, + messageId: message.messageId, + fileName, + cause, + }); } }, }; diff --git a/apps/mobile/src/state/thread-outbox.test.ts b/apps/mobile/src/state/thread-outbox.test.ts index d2634fb966f..d6b91c1c4f6 100644 --- a/apps/mobile/src/state/thread-outbox.test.ts +++ b/apps/mobile/src/state/thread-outbox.test.ts @@ -10,7 +10,7 @@ import { threadOutboxRetryDelayMs, type QueuedThreadMessage, } from "./thread-outbox-model"; -import { createThreadOutboxManager } from "./thread-outbox-manager"; +import { createThreadOutboxManager, ThreadOutboxManagerError } from "./thread-outbox-manager"; import type { ThreadOutboxStorage } from "./thread-outbox-storage"; function queuedMessage(input: { @@ -149,9 +149,48 @@ describe("thread outbox", () => { registry.dispose(); }); + it("reports structured load failures and permits a retry", async () => { + const registry = AtomRegistry.make(); + const loadCause = new Error("storage unavailable"); + const warnings: Array<{ message: string; error: unknown }> = []; + let loadCalls = 0; + const manager = createThreadOutboxManager({ + registry, + storage: { + load: async () => { + loadCalls += 1; + if (loadCalls === 1) throw loadCause; + return []; + }, + write: async () => undefined, + remove: async () => undefined, + }, + warn: (message, error) => warnings.push({ message, error }), + }); + + await manager.load(); + expect(warnings).toEqual([ + { + message: "[thread-outbox] failed to load persisted messages", + error: new ThreadOutboxManagerError({ + operation: "load", + environmentId: null, + threadId: null, + messageId: null, + cause: loadCause, + }), + }, + ]); + + await manager.load(); + expect(loadCalls).toBe(2); + registry.dispose(); + }); + it("keeps atom state aligned with durable writes and removals", async () => { const registry = AtomRegistry.make(); const stored = new Map(); + const removalCause = new Error("remove failed"); let failRemoval = true; const storage: ThreadOutboxStorage = { load: async () => [...stored.values()], @@ -160,7 +199,7 @@ describe("thread outbox", () => { }, remove: async (message) => { if (failRemoval) { - throw new Error("remove failed"); + throw removalCause; } stored.delete(message.messageId); }, @@ -176,7 +215,15 @@ describe("thread outbox", () => { "environment-1:thread-1": [message], }); - await expect(manager.remove(message)).rejects.toThrow("remove failed"); + await expect(manager.remove(message)).rejects.toEqual( + new ThreadOutboxManagerError({ + operation: "remove", + environmentId: message.environmentId, + threadId: message.threadId, + messageId: message.messageId, + cause: removalCause, + }), + ); expect(registry.get(manager.queuedMessagesByThreadKeyAtom)).toEqual({ "environment-1:thread-1": [message], });