From aff9318bf46beaf05cc7155b428d3f0b8711efd2 Mon Sep 17 00:00:00 2001 From: Wout Stiens <71498452+StiensWout@users.noreply.github.com> Date: Tue, 22 Sep 2026 11:25:14 +0200 Subject: [PATCH] fix(web): retry failed attachment uploads after reconnect (#10338) --- .../lib/composerAttachmentUploadQueue.test.ts | 25 ++++++ .../web/src/lib/attachmentUploadQueue.test.ts | 85 ++++++++++++++++++- apps/web/src/lib/attachmentUploadQueue.ts | 37 +++++++- 3 files changed, 145 insertions(+), 2 deletions(-) diff --git a/apps/mobile/src/lib/composerAttachmentUploadQueue.test.ts b/apps/mobile/src/lib/composerAttachmentUploadQueue.test.ts index 8bf67456ce53..79be3bbe8e2f 100644 --- a/apps/mobile/src/lib/composerAttachmentUploadQueue.test.ts +++ b/apps/mobile/src/lib/composerAttachmentUploadQueue.test.ts @@ -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> = {}; + 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(); const started = Promise.withResolvers(); diff --git a/apps/web/src/lib/attachmentUploadQueue.test.ts b/apps/web/src/lib/attachmentUploadQueue.test.ts index f376203cce6b..10b1e2e30a43 100644 --- a/apps/web/src/lib/attachmentUploadQueue.test.ts +++ b/apps/web/src/lib/attachmentUploadQueue.test.ts @@ -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, @@ -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(), @@ -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 }, @@ -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); @@ -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"), diff --git a/apps/web/src/lib/attachmentUploadQueue.ts b/apps/web/src/lib/attachmentUploadQueue.ts index 2a98f38133ac..b435c486a34f 100644 --- a/apps/web/src/lib/attachmentUploadQueue.ts +++ b/apps/web/src/lib/attachmentUploadQueue.ts @@ -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, @@ -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"; @@ -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(); const queue: UploadJob[] = []; const activeUploadsByEnvironment = new Map(); @@ -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; @@ -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", @@ -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) {