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
87 changes: 78 additions & 9 deletions apps/mobile/src/state/thread-outbox-manager.ts
Original file line number Diff line number Diff line change
@@ -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 {
Expand All @@ -8,6 +9,27 @@ import {
} from "./thread-outbox-model";
import type { ThreadOutboxStorage } from "./thread-outbox-storage";

export class ThreadOutboxManagerError extends Schema.TaggedErrorClass<ThreadOutboxManagerError>()(
"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;
Expand Down Expand Up @@ -49,31 +71,69 @@ 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<void> =>
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<void> =>
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),
);
});

const clearEnvironment = (environmentId: EnvironmentId): Promise<void> =>
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(
Expand All @@ -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,
}),
);
}
}),
);
Expand Down
98 changes: 80 additions & 18 deletions apps/mobile/src/state/thread-outbox-storage.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -8,6 +9,22 @@ import {

const THREAD_OUTBOX_DIRECTORY = "thread-outbox";

export class ThreadOutboxStorageError extends Schema.TaggedErrorClass<ThreadOutboxStorageError>()(
"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<ReadonlyArray<QueuedThreadMessage>>;
readonly write: (message: QueuedThreadMessage) => Promise<void>;
Expand All @@ -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,
});
}
},
};
53 changes: 50 additions & 3 deletions apps/mobile/src/state/thread-outbox.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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: {
Expand Down Expand Up @@ -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<MessageId, QueuedThreadMessage>();
const removalCause = new Error("remove failed");
let failRemoval = true;
const storage: ThreadOutboxStorage = {
load: async () => [...stored.values()],
Expand All @@ -160,7 +199,7 @@ describe("thread outbox", () => {
},
remove: async (message) => {
if (failRemoval) {
throw new Error("remove failed");
throw removalCause;
}
stored.delete(message.messageId);
},
Expand All @@ -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],
});
Expand Down
Loading