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
25 changes: 25 additions & 0 deletions apps/mobile/src/lib/composerAttachmentUploadQueue.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,31 @@ describe("composer attachment upload queue", () => {
queue.dispose();
});

it("retries a failed attachment when the worker restores requests after reconnect", async () => {
let states: Readonly<Record<string, ComposerAttachmentUploadState>> = {};
const upload = vi.fn().mockRejectedValueOnce(new Error("Disconnected")).mockResolvedValue(true);
const queue = createComposerAttachmentUploadQueue({
upload,
onChange: (next) => {
states = next;
},
});
const local = request("failed-upload");
const key = composerAttachmentUploadKey(environmentId, local.attachment.id);
queue.sync([local]);
await queue.settled();
expect(states[key]?.status).toBe("failed");
queue.sync([local]);
await queue.settled();
expect(upload).toHaveBeenCalledTimes(1);
queue.sync([]);
queue.sync([local]);
await queue.settled();
expect(upload).toHaveBeenCalledTimes(2);
expect(states[key]?.status).toBe("ready");
queue.dispose();
});

it("ignores a late completion after removal or environment switch", async () => {
const gate = Promise.withResolvers<boolean>();
const started = Promise.withResolvers<void>();
Expand Down
85 changes: 84 additions & 1 deletion apps/web/src/lib/attachmentUploadQueue.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
import { EnvironmentId } from "@t3tools/contracts";
import { afterEach, beforeEach, describe, expect, it, vi } from "vite-plus/test";
import { Atom, AsyncResult } from "effect/unstable/reactivity";
import { appAtomRegistry } from "../rpc/atomRegistry";

import {
composerFileNeedsReattach,
Expand All @@ -10,6 +12,7 @@ import {
} from "../composerDraftStore";

const mocks = vi.hoisted(() => ({
connectionStateAtom: vi.fn(),
createAssetUrl: vi.fn(),
createUploadUrl: Symbol("create-upload-url"),
executeAtomQuery: vi.fn(),
Expand All @@ -24,7 +27,14 @@ vi.mock("@t3tools/client-runtime/state/runtime", () => ({
squashAtomCommandFailure: (result: { readonly error: unknown }) => result.error,
}));

vi.mock("../rpc/atomRegistry", () => ({ appAtomRegistry: {} }));
vi.mock("../rpc/atomRegistry", async () => {
const { AtomRegistry } = await import("effect/unstable/reactivity");
return { appAtomRegistry: AtomRegistry.make() };
});

vi.mock("../connection/catalog", () => ({
environmentCatalog: { stateAtom: mocks.connectionStateAtom },
}));

vi.mock("../state/assets", () => ({
assetEnvironment: { createUrl: mocks.createAssetUrl },
Expand Down Expand Up @@ -140,8 +150,22 @@ function makeFile(id: string): ComposerFileAttachment {
};
}

const connectionStates = Atom.family((_environmentId: EnvironmentId) =>
Atom.make(AsyncResult.success({ phase: "connected" })),
);

function setConnected(environmentId: EnvironmentId, connected: boolean) {
appAtomRegistry.set(
connectionStates(environmentId),
AsyncResult.success({ phase: connected ? "connected" : "backoff" }),
);
}

describe("attachmentUploadQueue", () => {
beforeEach(() => {
mocks.connectionStateAtom.mockImplementation(connectionStates);
setConnected(firstEnvironment, true);
setConnected(secondEnvironment, true);
TestXmlHttpRequest.requests = [];
mocks.createAssetUrl.mockReset();
mocks.createAssetUrl.mockImplementation((target: unknown) => target);
Expand Down Expand Up @@ -183,6 +207,65 @@ describe("attachmentUploadQueue", () => {
vi.unstubAllGlobals();
});

it.each([false, true])(
"retries a failed file once after reconnect, including a late HTTP failure: %s",
async (lateFailure) => {
const image = makeFile("reconnect");
startAttachmentUpload({ environmentId: firstEnvironment, image });
await Promise.resolve();
const firstSettled = awaitAttachmentUploads([image.id]);
setConnected(firstEnvironment, false);
if (lateFailure) setConnected(firstEnvironment, true);
TestXmlHttpRequest.requests[0]!.complete(503);
await firstSettled;
if (!lateFailure) setConnected(firstEnvironment, true);
await Promise.resolve();
await Promise.resolve();
expect(TestXmlHttpRequest.requests).toHaveLength(2);
const retrySettled = awaitAttachmentUploads([image.id]);
TestXmlHttpRequest.requests[1]!.complete(503);
await retrySettled;
setConnected(firstEnvironment, true);
await Promise.resolve();
expect(TestXmlHttpRequest.requests).toHaveLength(2);
expect(readAttachmentUpload(image.id)?.status).toBe("failed");

setConnected(firstEnvironment, false);
setConnected(firstEnvironment, true);
await Promise.resolve();
await Promise.resolve();
const finalSettled = awaitAttachmentUploads([image.id]);
TestXmlHttpRequest.requests[2]!.complete();
await finalSettled;
expect(
getUploadedAttachments({ environmentId: firstEnvironment, images: [image] }),
).not.toBeNull();
setConnected(firstEnvironment, false);
setConnected(firstEnvironment, true);
await Promise.resolve();
expect(TestXmlHttpRequest.requests).toHaveLength(3);
},
);

it("does not retry for another environment or after the attachment is removed", async () => {
const image = makeFile("removed");
startAttachmentUpload({ environmentId: firstEnvironment, image });
await Promise.resolve();
const settled = awaitAttachmentUploads([image.id]);
TestXmlHttpRequest.requests[0]!.complete(503);
await settled;
setConnected(secondEnvironment, false);
setConnected(secondEnvironment, true);
await Promise.resolve();
expect(TestXmlHttpRequest.requests).toHaveLength(1);
setConnected(firstEnvironment, false);
setConnected(firstEnvironment, true);
releaseAttachmentUpload(image.id);
await Promise.resolve();
expect(TestXmlHttpRequest.requests).toHaveLength(1);
expect(readAttachmentUpload(image.id)).toBeUndefined();
});

it("uploads images immediately and sends attachment references", async () => {
const image = {
...makeImage("image-1"),
Expand Down
37 changes: 36 additions & 1 deletion apps/web/src/lib/attachmentUploadQueue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@ import {
type PersistedAttachmentVerification,
} from "@t3tools/client-runtime/state/attachments";
import { create } from "zustand";
import * as Option from "effect/Option";
import { AsyncResult } from "effect/unstable/reactivity";

import {
DraftId,
Expand All @@ -21,6 +23,7 @@ import {
type ComposerThreadTarget,
} from "../composerDraftStore";
import { appAtomRegistry } from "../rpc/atomRegistry";
import { environmentCatalog } from "../connection/catalog";
import { assetEnvironment } from "../state/assets";
import { attachmentEnvironment } from "../state/attachments";
import { readPreparedConnection } from "../state/session";
Expand Down Expand Up @@ -58,8 +61,10 @@ interface UploadJob {
attachmentId: string | null;
cancelled: boolean;
abort: (() => void) | null;
stopWatchingConnection: () => void;
}

// Failed jobs retain their source and connection subscription until retry or release.
const jobsByImageId = new Map<string, UploadJob>();
const queue: UploadJob[] = [];
const activeUploadsByEnvironment = new Map<EnvironmentId, number>();
Expand Down Expand Up @@ -358,7 +363,11 @@ function pumpUploads(): void {
}
})
.finally(() => {
if (jobsByImageId.get(job.image.id) === job) {
if (
jobsByImageId.get(job.image.id) === job &&
readAttachmentUpload(job.image.id)?.status !== "failed"
) {
job.stopWatchingConnection();
jobsByImageId.delete(job.image.id);
}
const remaining = (activeUploadsByEnvironment.get(job.environmentId) ?? 1) - 1;
Expand Down Expand Up @@ -427,9 +436,34 @@ export function startAttachmentUpload(input: {
attachmentId: null,
cancelled: false,
abort: null,
stopWatchingConnection: () => {},
};

jobsByImageId.set(input.image.id, job);
const connectionAtom = environmentCatalog.stateAtom(job.environmentId);
const isConnected = () =>
Option.exists(
AsyncResult.value(appAtomRegistry.get(connectionAtom)),
(state) => state.phase === "connected",
);
let wasConnected = isConnected();
job.stopWatchingConnection = appAtomRegistry.subscribe(connectionAtom, () => {
const connected = isConnected();
const reconnected = connected && !wasConnected;
wasConnected = connected;
if (!reconnected) return;
// The HTTP failure can arrive after the socket has already reconnected.
// Wait for that attempt, then retry only if this job still owns the file.
void job.settled.then(() => {
if (
jobsByImageId.get(job.image.id) === job &&
readAttachmentUpload(job.image.id)?.status === "failed" &&
isConnected()
) {
retryAttachmentUpload(input);
}
});
});
queue.push(job);
setUploadState(input.image.id, {
status: "uploading",
Expand All @@ -451,6 +485,7 @@ function cancelAttachmentUpload(imageId: string): void {
return;
}
job.cancelled = true;
job.stopWatchingConnection();
jobsByImageId.delete(imageId);
const queuedIndex = queue.indexOf(job);
if (queuedIndex !== -1) {
Expand Down
Loading