Files
EpicNext-Cms/src/lib/sse-client.ts
T
Simo b1ddda66ff
CI / check (push) Successful in 27s
CI / release (push) Skipped
CI / deploy (push) Successful in 43s
Revert "Merge pull request 'Complete Housekeeping migration and /ase cutover' (#52) from codex/housekeeping-complete into main"
This reverts commit 488b6e57c4, reversing
changes made to b506b4499a.
2026-08-30 21:31:34 +02:00

108 lines
2.5 KiB
TypeScript

import { toast } from "sonner";
import { adminFetch } from "@/lib/admin-fetch";
export type SseEvent = Record<string, unknown>;
/**
* Read an SSE response body and invoke `onEvent` for each `data:` JSON payload.
*/
export async function readSseStream(
body: ReadableStream<Uint8Array>,
onEvent: (event: SseEvent) => void,
signal?: AbortSignal,
): Promise<void> {
const reader = body.getReader();
const decoder = new TextDecoder();
let buf = "";
try {
while (true) {
if (signal?.aborted) {
await reader.cancel();
break;
}
const { value, done } = await reader.read();
if (done) break;
buf += decoder.decode(value, { stream: true });
const parts = buf.split("\n\n");
buf = parts.pop() ?? "";
for (const part of parts) {
if (!part.startsWith("data: ")) continue;
try {
onEvent(JSON.parse(part.slice(6)) as SseEvent);
} catch {
/* skip malformed events */
}
}
}
} finally {
reader.releaseLock();
}
}
/**
* POST JSON to an admin SSE import endpoint and drive the standard
* item_progress / batch_complete callbacks used by clothing & clone clients.
*/
export async function runSseImport(
url: string,
body: unknown,
onDone: (classname: string) => void,
onComplete: (succeeded: number, failed: number) => void,
signal?: AbortSignal,
onFailed?: (classname: string, error?: string) => void,
): Promise<void> {
const res = await adminFetch(url, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify(body),
signal,
});
if (!res.ok) {
const data = await res.json().catch(() => ({}));
toast.error(
typeof data.error === "string"
? data.error
: `Import failed (${res.status})`,
);
onComplete(0, 0);
return;
}
if (!res.body) {
toast.error("No response stream");
onComplete(0, 0);
return;
}
let succeeded = 0;
let failed = 0;
await readSseStream(
res.body,
(evt) => {
if (evt.type === "item_progress") {
const classname = String(evt.classname ?? "");
if (evt.status === "done") {
onDone(classname);
} else if (evt.status === "failed") {
onFailed?.(classname, String(evt.error ?? ""));
}
}
if (evt.type === "error") {
toast.error(
typeof evt.message === "string"
? evt.message
: "Import finished with errors",
);
}
if (evt.type === "batch_complete") {
succeeded = Number(evt.succeeded ?? 0);
failed = Number(evt.failed ?? 0);
}
},
signal,
);
onComplete(succeeded, failed);
}