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
18 changes: 7 additions & 11 deletions apps/mobile/src/state/thread-outbox-manager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -276,17 +276,13 @@ export function createThreadOutboxManager(options: ThreadOutboxManagerOptions) {
const revisionsAtRequest = new Map(revisions);
return serialize(async () => {
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 [];
throw new ThreadOutboxManagerError({
operation: "clear-environment-load",
environmentId,
threadId: null,
messageId: null,
cause,
});
});
const allMessages = flattenQueuedThreadMessages(
groupQueuedThreadMessages([...persisted, ...currentMessages()]),
Expand Down
22 changes: 11 additions & 11 deletions apps/mobile/src/state/thread-outbox-storage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -80,20 +80,20 @@ export const expoThreadOutboxStorage: ThreadOutboxStorage = {
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,
}),
);
// A partial queue hides attachment owners from cleanup. Keep all
// records untouched until every persisted message can be read.
throw new ThreadOutboxStorageError({
operation: "read-message",
environmentId: null,
threadId: null,
messageId: null,
fileName: entry.name,
cause,
});
}
}
} catch (cause) {
if (cause instanceof ThreadOutboxStorageError) throw cause;
throw new ThreadOutboxStorageError({
operation: "load",
environmentId: null,
Expand Down
96 changes: 95 additions & 1 deletion apps/mobile/src/state/thread-outbox.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,34 @@ import {
ThreadId,
} from "@t3tools/contracts";
import { AtomRegistry } from "effect/unstable/reactivity";
import { onTestFinished, vi } from "vite-plus/test";

const outboxFiles = vi.hoisted(() => new Map<string, string | Error>());

vi.mock("expo-file-system", () => {
class File {
constructor(readonly name: string) {}

async text(): Promise<string> {
const contents = outboxFiles.get(this.name);
if (contents instanceof Error) throw contents;
if (contents === undefined) throw new Error("Missing file");
return contents;
}
}

return {
File,
Directory: class {
create() {}

list() {
return Array.from(outboxFiles.keys(), (name) => new File(name));
}
},
Paths: { document: "/documents" },
};
});

import {
decodeQueuedThreadMessage,
Expand All @@ -24,7 +52,7 @@ import {
type QueuedThreadMessage,
} from "./thread-outbox-model";
import { createThreadOutboxManager, ThreadOutboxManagerError } from "./thread-outbox-manager";
import type { ThreadOutboxStorage } from "./thread-outbox-storage";
import { expoThreadOutboxStorage, type ThreadOutboxStorage } from "./thread-outbox-storage";

function queuedMessage(input: {
readonly environmentId?: string;
Expand All @@ -44,6 +72,72 @@ function queuedMessage(input: {
}

describe("thread outbox", () => {
it.each(["read", "json", "schema"] as const)(
"does not load a partial outbox after a record %s failure",
async (failure) => {
onTestFinished(() => outboxFiles.clear());
const first = queuedMessage({
messageId: "message-1",
createdAt: "2026-06-08T10:00:01.000Z",
});
const second = queuedMessage({
messageId: "message-2",
createdAt: "2026-06-08T10:00:02.000Z",
});
outboxFiles.set("message-1.json", JSON.stringify(encodeQueuedThreadMessage(first)));
outboxFiles.set(
"message-2.json",
failure === "read"
? new Error("storage unavailable")
: failure === "json"
? "{"
: JSON.stringify({ ...second, schemaVersion: 999 }),
);

await expect(expoThreadOutboxStorage.load()).rejects.toMatchObject({
operation: "read-message",
fileName: "message-2.json",
});

outboxFiles.set("message-2.json", JSON.stringify(encodeQueuedThreadMessage(second)));
await expect(expoThreadOutboxStorage.load()).resolves.toEqual([first, second]);
},
);

it("preserves queued messages when environment cleanup cannot read the outbox", async () => {
const registry = AtomRegistry.make();
onTestFinished(() => registry.dispose());
const message = queuedMessage({
messageId: "message-1",
createdAt: "2026-06-08T10:00:01.000Z",
});
const stored = new Map<MessageId, QueuedThreadMessage>();
const manager = createThreadOutboxManager({
registry,
warn: () => {},
storage: {
load: async () => {
throw new Error("storage unavailable");
},
write: async (entry) => {
stored.set(entry.messageId, entry);
},
remove: async (entry) => {
stored.delete(entry.messageId);
},
},
});
await manager.enqueue(message);

await expect(manager.clearEnvironment(message.environmentId)).rejects.toMatchObject({
operation: "clear-environment-load",
});
expect([...stored.values()]).toEqual([message]);
expect(registry.get(manager.queuedMessagesByThreadKeyAtom)).toEqual({
"environment-1:thread-1": [message],
});
});

it("groups messages by scoped thread and preserves creation order", () => {
const later = queuedMessage({
messageId: "message-2",
Expand Down
69 changes: 66 additions & 3 deletions apps/mobile/src/state/use-composer-drafts.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,8 @@ import {
import { onTestFinished, vi } from "vite-plus/test";

const composerDraftFileMocks = vi.hoisted(() => {
let document = "";
let document = JSON.stringify({ schemaVersion: 1, drafts: {} });
let readError: Error | null = null;
let writeError: Error | null = null;
let releaseRead: (() => void) | null = null;
let readBarrier = Promise.resolve();
Expand All @@ -33,6 +34,9 @@ const composerDraftFileMocks = vi.hoisted(() => {
setDocument(value: unknown) {
document = JSON.stringify(value);
},
setReadError(error: Error | null) {
readError = error;
},
setWriteError(error: Error | null) {
writeError = error;
},
Expand All @@ -50,6 +54,10 @@ const composerDraftFileMocks = vi.hoisted(() => {
},
Directory: class {
create() {}

list() {
return [];
}
},
File: class {
exists = true;
Expand All @@ -61,6 +69,7 @@ const composerDraftFileMocks = vi.hoisted(() => {

async text() {
await readBarrier;
if (readError) throw readError;
return document;
}

Expand Down Expand Up @@ -156,7 +165,8 @@ const DRAFT: ComposerDraft = {
afterEach(() => {
vi.useRealTimers();
resetComposerDraftsLoadState();
composerDraftFileMocks.setDocument("");
composerDraftFileMocks.setDocument({ schemaVersion: 1, drafts: {} });
composerDraftFileMocks.setReadError(null);
composerDraftFileMocks.setWriteError(null);
composerDraftFileMocks.setNextWriteBarrier(null);
composerDraftFileMocks.setOnWrite(null);
Expand Down Expand Up @@ -1043,9 +1053,62 @@ describe("mobile composer drafts", () => {
});
});

it.each(["read", "decode"] as const)(
"preserves saved drafts and attachment files when the draft %s fails",
async (failure) => {
vi.useFakeTimers();
const file = {
id: "saved-file",
type: "file" as const,
name: "report.pdf",
mimeType: "application/pdf",
sizeBytes: 42,
fileUri: "file:///documents/t3-composer-attachments/report.pdf",
};
composerDraftFileMocks.setDocument({
schemaVersion: failure === "decode" ? 999 : 1,
drafts: { "environment-1:saved": { text: "Saved draft", attachments: [file] } },
});
const original = composerDraftFileMocks.getDocument();
if (failure === "read") {
composerDraftFileMocks.setReadError(new Error("storage unavailable"));
}

await expect(releaseUnusedComposerAttachmentFiles([file])).rejects.toMatchObject({
operation: failure,
});
setComposerDraftText("environment-1:new", "Keep my new edits too");
await expect(flushComposerDrafts()).rejects.toMatchObject({ operation: failure });

expect(composerDraftFileMocks.getDocument()).toBe(original);
expect(composerAttachmentCleanupMocks.remove).not.toHaveBeenCalled();
},
);

it("retries a failed debounced read on final flush without dropping saved drafts or new edits", async () => {
vi.useFakeTimers();
composerDraftFileMocks.setDocument({
schemaVersion: 1,
drafts: { "environment-1:saved": DRAFT },
});
const original = composerDraftFileMocks.getDocument();
composerDraftFileMocks.setReadError(new Error("storage unavailable"));
setComposerDraftText("environment-1:new", "New edits");
await vi.advanceTimersByTimeAsync(200);
expect(composerDraftFileMocks.getDocument()).toBe(original);

composerDraftFileMocks.setReadError(null);
await flushComposerDrafts();

expect(JSON.parse(composerDraftFileMocks.getDocument()).drafts).toEqual({
"environment-1:saved": DRAFT,
"environment-1:new": { text: "New edits", attachments: [] },
});
});

it("serializes environment cleanup after an older queued write", async () => {
vi.useFakeTimers();
composerDraftFileMocks.setDocument(JSON.stringify({ schemaVersion: 1, drafts: {} }));
composerDraftFileMocks.setDocument({ schemaVersion: 1, drafts: {} });
composerDraftFileMocks.resetWrites();
let releaseFirstWrite!: () => void;
const firstWriteBarrier = new Promise<void>((resolve) => {
Expand Down
Loading
Loading