feat: public events/polls, friends graph, captcha, SSE hardening, and admin UX
Ship product gaps: register/vote pages, friend add/accept/decline/remove, email verify TTL, captcha on login/forgot, soft-fail user actions, SSE abort/shared client, Commando Centrum error toasts, admin delete for events/polls, and IT/NL i18n fills. Co-authored-by: Cursor <[email protected]>
This commit is contained in:
1 parent
2ff08e5127
commit
ed7db6e048
76 files changed
+4834
-1376
No files matched your search
@@ -40,4 +40,26 @@ describe("import/core/sse-batch", () => {
|
||||
failed: 1,
|
||||
});
|
||||
});
|
||||
|
||||
it("stops starting new chunks when AbortSignal fires", async () => {
|
||||
const ac = new AbortController();
|
||||
let started = 0;
|
||||
const res = runSseBatch({
|
||||
items: ["a", "b", "c", "d"],
|
||||
concurrency: 1,
|
||||
signal: ac.signal,
|
||||
labelOf: (item) => item,
|
||||
worker: async (item) => {
|
||||
started++;
|
||||
if (item === "a") ac.abort();
|
||||
// Slow enough that abort lands before the next chunk starts.
|
||||
await new Promise((r) => setTimeout(r, 20));
|
||||
return { ok: true };
|
||||
},
|
||||
});
|
||||
const events = await collect(res);
|
||||
expect(started).toBeLessThan(4);
|
||||
expect(events.some((e) => e.type === "batch_complete")).toBe(false);
|
||||
expect(events[0]).toMatchObject({ type: "batch_start", total: 4 });
|
||||
});
|
||||
});
|
||||
@@ -7,6 +7,8 @@ export interface SseWorkerResult {
|
||||
export interface RunSseBatchOptions<T> {
|
||||
items: T[];
|
||||
concurrency: number;
|
||||
/** Abort when the client disconnects (e.g. `request.signal`). */
|
||||
signal?: AbortSignal;
|
||||
/** Label used as the `classname` field on item_progress (kept for the existing client parser). */
|
||||
labelOf: (item: T) => string;
|
||||
/** Per-item worker. `report(status)` streams intermediate progress (e.g. 'downloading'). */
|
||||
@@ -20,21 +22,30 @@ export interface RunSseBatchOptions<T> {
|
||||
/**
|
||||
* Generic SSE batch runner. Emits the same event shape the furni client
|
||||
* parser consumes: batch_start / item_progress / batch_complete.
|
||||
* Stops starting new chunks when `signal` aborts or the client cancels the stream.
|
||||
*/
|
||||
export function runSseBatch<T>(opts: RunSseBatchOptions<T>): Response {
|
||||
const { items, labelOf, worker } = opts;
|
||||
const { items, labelOf, worker, signal } = opts;
|
||||
const concurrency = Math.min(Math.max(opts.concurrency || 3, 1), 5);
|
||||
const encoder = new TextEncoder();
|
||||
const ac = new AbortController();
|
||||
|
||||
if (signal) {
|
||||
if (signal.aborted) ac.abort();
|
||||
else signal.addEventListener("abort", () => ac.abort(), { once: true });
|
||||
}
|
||||
|
||||
const stream = new ReadableStream({
|
||||
async start(controller) {
|
||||
const send = (data: unknown) => {
|
||||
if (ac.signal.aborted) return;
|
||||
try {
|
||||
controller.enqueue(
|
||||
encoder.encode(`data: ${JSON.stringify(data)}\n\n`),
|
||||
);
|
||||
} catch {
|
||||
/* stream closed by client */
|
||||
ac.abort();
|
||||
}
|
||||
};
|
||||
|
||||
@@ -46,9 +57,12 @@ export function runSseBatch<T>(opts: RunSseBatchOptions<T>): Response {
|
||||
let withWarnings = 0;
|
||||
|
||||
for (let i = 0; i < items.length; i += concurrency) {
|
||||
if (ac.signal.aborted) break;
|
||||
|
||||
const chunk = items.slice(i, i + concurrency);
|
||||
await Promise.allSettled(
|
||||
chunk.map(async (item, chunkIdx) => {
|
||||
if (ac.signal.aborted) return;
|
||||
const index = i + chunkIdx;
|
||||
const classname = labelOf(item);
|
||||
send({
|
||||
@@ -61,6 +75,7 @@ export function runSseBatch<T>(opts: RunSseBatchOptions<T>): Response {
|
||||
const result = await worker(item, index, (status) =>
|
||||
send({ type: "item_progress", classname, status, index }),
|
||||
);
|
||||
if (ac.signal.aborted) return;
|
||||
if (result.ok) {
|
||||
succeeded++;
|
||||
if (result.warnings?.length) withWarnings++;
|
||||
@@ -84,6 +99,7 @@ export function runSseBatch<T>(opts: RunSseBatchOptions<T>): Response {
|
||||
});
|
||||
}
|
||||
} catch (err) {
|
||||
if (ac.signal.aborted) return;
|
||||
failed++;
|
||||
send({
|
||||
type: "item_progress",
|
||||
@@ -97,14 +113,23 @@ export function runSseBatch<T>(opts: RunSseBatchOptions<T>): Response {
|
||||
);
|
||||
}
|
||||
|
||||
send({
|
||||
type: "batch_complete",
|
||||
succeeded,
|
||||
failed,
|
||||
warnings: withWarnings,
|
||||
duration: Date.now() - startTime,
|
||||
});
|
||||
controller.close();
|
||||
if (!ac.signal.aborted) {
|
||||
send({
|
||||
type: "batch_complete",
|
||||
succeeded,
|
||||
failed,
|
||||
warnings: withWarnings,
|
||||
duration: Date.now() - startTime,
|
||||
});
|
||||
}
|
||||
try {
|
||||
controller.close();
|
||||
} catch {
|
||||
/* already closed */
|
||||
}
|
||||
},
|
||||
cancel() {
|
||||
ac.abort();
|
||||
},
|
||||
});
|
||||
|
||||
|
||||
Reference in new issue
Block a user