fix(catalog): stream single imports and preserve response errors
CI / check (push) Successful in 1m31s
CI / deploy (push) Successful in 1m13s
CI / e2e (push) Successful in 21s

This commit is contained in:
Simo committed 2026-09-05 15:16:42 +02:00
1 parent f0dcbee162
commit 33b6d1520e
5 files changed
+351 -142

No files matched your search

+52
View File
@@ -0,0 +1,52 @@
import { readSseStream } from "@/lib/sse-client";
interface ImportResponse {
status: number;
data: Record<string, unknown>;
}
export async function readImportResponse(
response: Response,
): Promise<ImportResponse> {
if (
response.ok &&
response.headers.get("content-type")?.includes("text/event-stream")
) {
if (!response.body)
throw Error(
"Import response is empty. Check the furniture status before retrying.",
);
let result: ImportResponse | undefined;
await readSseStream(response.body, (event) => {
if (
event.type === "result" &&
typeof event.status === "number" &&
event.data &&
typeof event.data === "object"
)
result = {
status: event.status,
data: event.data as Record<string, unknown>,
};
});
if (!result)
throw Error(
"Import connection ended before confirmation. Some changes may already be saved. Check the furniture status before retrying.",
);
return result;
}
const text = await response.text();
let data: Record<string, unknown>;
try {
data = JSON.parse(text);
} catch {
const reason = [502, 503, 504, 524].includes(response.status)
? "Server or proxy unavailable, or request timed out"
: "Server returned an unexpected response";
throw Error(
`${reason} (HTTP ${response.status}). Check the furniture status before retrying.`,
);
}
if (!data || typeof data !== "object")
throw Error(`Invalid import response (HTTP ${response.status}).`);
return { status: response.status, data };
}
@@ -0,0 +1,73 @@
import { afterEach, expect, it, vi } from "vitest";
import { readImportResponse } from "@/lib/furni/import-response";
import { streamImportResponse } from "./single-response";
vi.mock("@/lib/server-log", () => ({ logServerError: vi.fn() }));
afterEach(() => vi.useRealTimers());
it("sends headers and heartbeats before a slow operation finishes", async () => {
vi.useFakeTimers();
let finish!: (response: Response) => void;
const operation = vi.fn(
() =>
new Promise<Response>((resolve) => {
finish = resolve;
}),
);
const response = streamImportResponse(operation);
expect(response.headers.get("X-Accel-Buffering")).toBe("no");
if (!response.body) throw Error("Missing stream");
const reader = response.body.getReader();
const decoder = new TextDecoder();
expect(decoder.decode((await reader.read()).value)).toContain('"started"');
await vi.advanceTimersByTimeAsync(10000);
expect(decoder.decode((await reader.read()).value)).toContain('"processing"');
finish(Response.json({ ok: true, id: 42 }));
expect(decoder.decode((await reader.read()).value)).toContain('"id":42');
expect((await reader.read()).done).toBe(true);
expect(operation).toHaveBeenCalledTimes(1);
expect(vi.getTimerCount()).toBe(0);
});
it("preserves failure status and the importer error", async () => {
const result = await readImportResponse(
streamImportResponse(async () =>
Response.json({ error: "Missing source asset" }, { status: 409 }),
),
);
expect(result).toEqual({
status: 409,
data: { error: "Missing source asset" },
});
});
it("delivers caught server failures instead of abruptly closing", async () => {
const result = await readImportResponse(
streamImportResponse(async () => {
throw Error("private database details");
}),
);
expect(result.status).toBe(500);
expect(result.data.error).toContain("Some changes may already be saved");
expect(result.data.error).not.toContain("private database");
});
it("reports proxy timeout status without rendering proxy HTML", async () => {
await expect(
readImportResponse(
new Response("<html>gateway error</html>", { status: 524 }),
),
).rejects.toThrow("HTTP 524");
});
it("does not treat a truncated stream as successful", async () => {
await expect(
readImportResponse(
new Response('data: {"type":"progress"}\n\n', {
headers: { "Content-Type": "text/event-stream" },
}),
),
).rejects.toThrow("before confirmation");
});
it("preserves JSON authentication failures", async () => {
expect(
await readImportResponse(
Response.json({ error: "Unauthorized" }, { status: 401 }),
),
).toEqual({ status: 401, data: { error: "Unauthorized" } });
});
@@ -0,0 +1,62 @@
import { logServerError } from "@/lib/server-log";
/** Return headers immediately and keep proxies alive while one import completes. */
export function streamImportResponse(
operation: () => Promise<Response>,
): Response {
const encoder = new TextEncoder();
let closed = false;
let timer: ReturnType<typeof setInterval> | undefined;
const body = new ReadableStream<Uint8Array>({
start(controller) {
const send = (data: unknown) => {
if (!closed)
controller.enqueue(
encoder.encode(`data: ${JSON.stringify(data)}\n\n`),
);
};
send({ type: "progress", stage: "started" });
timer = setInterval(
() => send({ type: "progress", stage: "processing" }),
10000,
);
void (async () => {
try {
const response = await operation();
send({
type: "result",
status: response.status,
data: await response.json(),
});
} catch (error) {
logServerError("furni.import_stream_failed", error);
send({
type: "result",
status: 500,
data: {
error:
"Import interrupted on the server. Some changes may already be saved. Check the furniture status before retrying.",
},
});
} finally {
clearInterval(timer);
if (!closed) {
closed = true;
controller.close();
}
}
})();
},
cancel() {
closed = true;
clearInterval(timer);
},
});
return new Response(body, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache, no-transform",
"X-Accel-Buffering": "no",
},
});
}