import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { createAutoClient, createMainThreadClient, createRenderClient, RenderFailed, RenderTimeout, Superseded, WorkerUnavailable, type WorkerLike, } from "./client"; import type { AssembleJob, AssembleResult, RenderJob, RenderResult, WorkerMessage, WorkerRequest } from "./jobs"; class FakeWorker implements WorkerLike { onmessage: WorkerLike["onmessage"] = null; onerror: WorkerLike["onerror"] = null; sent: WorkerRequest[] = []; terminated = false; postMessage(message: WorkerRequest) { this.sent.push(message); } terminate() { this.terminated = true; } emit(data: WorkerMessage) { this.onmessage?.({ data }); } ready() { this.emit({ type: "ready" }); } reply(id: number) { this.emit({ type: "result", result: result(id) }); } replyAssembled(id: number) { this.emit({ type: "assembled", result: assembled(id) }); } } const assembled = (id: number): AssembleResult => ({ id, bytes: new Uint8Array([id]), ms: 3 }); const assembleJob = (id: number): AssembleJob => ({ id, title: "T", pages: [{ jpeg: new Uint8Array([1]), widthPt: 10, heightPt: 10 }] }); const kinds = (w: FakeWorker) => w.sent.map((m) => `${m.type}:${m.job.id}`); const result = (id: number): RenderResult => ({ id, bytes: new Uint8Array([id]), pages: 1, issues: [], fingerprint: `fp${id}`, ms: 5, }); const job = (id: number): RenderJob => ({ id, model: {} as RenderJob["model"], prefs: {} as RenderJob["prefs"] }); function setup(timeoutMs = 1000) { const workers: FakeWorker[] = []; const client = createRenderClient({ timeoutMs, workerFactory: () => { const w = new FakeWorker(); workers.push(w); return w; }, }); return { client, workers }; } const settled = (p: Promise) => p.then(() => "ok", (e) => e); describe("RenderClient", () => { beforeEach(() => vi.useFakeTimers()); afterEach(() => vi.useRealTimers()); it("sends nothing until the worker is ready, then resolves with the result", async () => { const { client, workers } = setup(); const p = client.render(job(1)); expect(workers).toHaveLength(1); expect(workers[0].sent).toHaveLength(0); workers[0].ready(); expect(workers[0].sent.map((m) => m.job.id)).toEqual([1]); workers[0].reply(1); expect((await p).fingerprint).toBe("fp1"); }); it("a new call supersedes the one in flight, and only the newest queued job runs next", async () => { const s = setup(); const a = s.client.render(job(1)); s.workers[0].ready(); const b = s.client.render(job(2)); const c = s.client.render(job(3)); expect(await settled(a)).toBeInstanceOf(Superseded); expect(await settled(b)).toBeInstanceOf(Superseded); // The worker is still busy with job 1: nothing else was posted. expect(s.workers[0].sent.map((m) => m.job.id)).toEqual([1]); s.workers[0].reply(1); // result of the superseded job is discarded expect(s.workers[0].sent.map((m) => m.job.id)).toEqual([1, 3]); s.workers[0].reply(3); expect((await c).id).toBe(3); }); it("rejects the superseded in-flight promise immediately", async () => { const { client, workers } = setup(); const a = client.render(job(1)); workers[0].ready(); const b = client.render(job(2)); expect(await settled(a)).toBeInstanceOf(Superseded); workers[0].reply(1); workers[0].reply(2); expect((await b).id).toBe(2); }); it("surfaces a render error from the worker and keeps going", async () => { const { client, workers } = setup(); const a = client.render(job(1)); workers[0].ready(); workers[0].emit({ type: "error", id: 1, message: "boom", stack: "at worker" }); const err = await settled(a); expect(err).toBeInstanceOf(RenderFailed); expect((err as Error).message).toBe("boom"); const b = client.render(job(2)); workers[0].reply(2); expect((await b).id).toBe(2); }); it("watchdog: terminates, respawns, rejects with RenderTimeout, then serves the next job", async () => { const { client, workers } = setup(1000); const a = client.render(job(1)); workers[0].ready(); vi.advanceTimersByTime(1000); expect(await settled(a)).toBeInstanceOf(RenderTimeout); expect(workers[0].terminated).toBe(true); expect(workers).toHaveLength(2); const b = client.render(job(2)); workers[1].ready(); expect(workers[1].sent.map((m) => m.job.id)).toEqual([2]); workers[1].reply(2); expect((await b).id).toBe(2); }); it("watchdog on a superseded job still replaces the worker and then runs the queued job", async () => { const { client, workers } = setup(1000); const a = client.render(job(1)); workers[0].ready(); const b = client.render(job(2)); expect(await settled(a)).toBeInstanceOf(Superseded); vi.advanceTimersByTime(1000); expect(workers[0].terminated).toBe(true); workers[1].ready(); expect(workers[1].sent.map((m) => m.job.id)).toEqual([2]); workers[1].reply(2); expect((await b).id).toBe(2); }); it("a start-up error before ready rejects with WorkerUnavailable, and so do later calls", async () => { const { client, workers } = setup(); const a = client.render(job(1)); workers[0].emit({ type: "error", id: null, message: "fonts missing" }); const err = await settled(a); expect(err).toBeInstanceOf(WorkerUnavailable); expect((err as Error).message).toContain("fonts missing"); expect(workers[0].terminated).toBe(true); expect(await settled(client.render(job(2)))).toBeInstanceOf(WorkerUnavailable); expect(workers).toHaveLength(1); }); it("a worker script error before ready is WorkerUnavailable; there is no timer fallback", async () => { const { client, workers } = setup(1000); const a = client.render(job(1)); // A silent worker that never says ready just waits: no timer turns that into a failure. vi.advanceTimersByTime(60_000); let done = false; void settled(a).then(() => (done = true)); await Promise.resolve(); expect(done).toBe(false); workers[0].onerror?.({ message: "SyntaxError" }); expect(await settled(a)).toBeInstanceOf(WorkerUnavailable); }); it("a factory that throws is WorkerUnavailable", async () => { const client = createRenderClient({ workerFactory: () => { throw new Error("no Worker here"); }, }); expect(await settled(client.render(job(1)))).toBeInstanceOf(WorkerUnavailable); }); it("dispose terminates the worker and rejects what is outstanding", async () => { const { client, workers } = setup(); const a = client.render(job(1)); workers[0].ready(); const b = client.render(job(2)); client.dispose(); expect(workers[0].terminated).toBe(true); expect(await settled(a)).toBeInstanceOf(Superseded); expect(String(await settled(b))).toContain("disposed"); expect(String(await settled(client.render(job(3))))).toContain("disposed"); }); }); describe("assemble lane", () => { beforeEach(() => vi.useFakeTimers()); afterEach(() => vi.useRealTimers()); it("queues behind a render in flight without superseding it", async () => { const { client, workers } = setup(); const r1 = client.render(job(1)); workers[0].ready(); const a = client.assemble(assembleJob(2)); expect(kinds(workers[0])).toEqual(["render:1"]); workers[0].reply(1); expect((await r1).id).toBe(1); expect(kinds(workers[0])).toEqual(["render:1", "assemble-images:2"]); workers[0].replyAssembled(2); expect((await a).id).toBe(2); }); it("runs ahead of a render that was queued before it", async () => { const { client, workers } = setup(); const r1 = client.render(job(1)); workers[0].ready(); const r3 = client.render(job(3)); // supersedes r1, waits for the worker const a = client.assemble(assembleJob(2)); expect(await settled(r1)).toBeInstanceOf(Superseded); workers[0].reply(1); expect(kinds(workers[0])).toEqual(["render:1", "assemble-images:2"]); workers[0].replyAssembled(2); expect((await a).id).toBe(2); expect(kinds(workers[0])).toEqual(["render:1", "assemble-images:2", "render:3"]); workers[0].reply(3); expect((await r3).id).toBe(3); }); it("is not superseded by renders that arrive while it runs, and the newest render still wins afterwards", async () => { const { client, workers } = setup(); const a = client.assemble(assembleJob(1)); workers[0].ready(); const r2 = client.render(job(2)); const r3 = client.render(job(3)); expect(await settled(r2)).toBeInstanceOf(Superseded); workers[0].replyAssembled(1); expect((await a).id).toBe(1); expect(kinds(workers[0])).toEqual(["assemble-images:1", "render:3"]); workers[0].reply(3); expect((await r3).id).toBe(3); }); it("runs several assembles in order and a render never rejects them", async () => { const { client, workers } = setup(); const r0 = client.render(job(0)); workers[0].ready(); workers[0].reply(0); await r0; const a = client.assemble(assembleJob(1)); const b = client.assemble(assembleJob(2)); void client.render(job(3)).catch(() => {}); void client.render(job(4)).catch(() => {}); workers[0].replyAssembled(1); expect((await a).id).toBe(1); expect(kinds(workers[0]).slice(-1)).toEqual(["assemble-images:2"]); workers[0].replyAssembled(2); expect((await b).id).toBe(2); }); it("surfaces an assemble error from the worker and keeps serving", async () => { const { client, workers } = setup(); const a = client.assemble(assembleJob(1)); workers[0].ready(); workers[0].emit({ type: "error", id: 1, message: "bad jpeg" }); const err = await settled(a); expect(err).toBeInstanceOf(RenderFailed); expect((err as Error).message).toBe("bad jpeg"); const b = client.assemble(assembleJob(2)); workers[0].replyAssembled(2); expect((await b).id).toBe(2); }); it("a render result cannot settle a running assemble with the same id", async () => { const { client, workers } = setup(); const a = client.assemble(assembleJob(5)); workers[0].ready(); workers[0].reply(5); let done = false; void settled(a).then(() => (done = true)); await Promise.resolve(); expect(done).toBe(false); workers[0].replyAssembled(5); expect((await a).id).toBe(5); }); it("dispose rejects queued assembles", async () => { const { client, workers } = setup(); void client.render(job(1)).catch(() => {}); workers[0].ready(); const a = client.assemble(assembleJob(2)); client.dispose(); expect(String(await settled(a))).toContain("disposed"); }); it("the main-thread client keeps the same lanes", async () => { const gates: Array<() => void> = []; const seen: string[] = []; const client = createMainThreadClient( (j) => new Promise((resolve) => { seen.push(`render:${j.id}`); gates.push(() => resolve(result(j.id))); }), (j) => new Promise((resolve) => { seen.push(`assemble:${j.id}`); gates.push(() => resolve(assembled(j.id))); }), ); const r1 = client.render(job(1)); const a = client.assemble(assembleJob(2)); gates[0](); expect((await r1).id).toBe(1); const r3 = client.render(job(3)); await Promise.resolve(); await Promise.resolve(); expect(seen).toEqual(["render:1", "assemble:2"]); gates[1](); expect((await a).id).toBe(2); await Promise.resolve(); await Promise.resolve(); expect(seen).toEqual(["render:1", "assemble:2", "render:3"]); gates[2](); expect((await r3).id).toBe(3); }); it("the auto client assembles on the main thread when the worker cannot start", async () => { const client = createAutoClient({ workerFactory: () => { throw new Error("no worker"); }, mainAssembleExecutor: async (j) => assembled(j.id), }); expect((await client.assemble(assembleJob(9))).id).toBe(9); expect(client.mode).toBe("main"); }); }); describe("main-thread client and fallback", () => { it("runs jobs serially with latest-wins", async () => { const gates: Array<() => void> = []; const seen: number[] = []; const client = createMainThreadClient( (j) => new Promise((resolve) => { seen.push(j.id); gates.push(() => resolve(result(j.id))); }), ); expect(client.mode).toBe("main"); const a = client.render(job(1)); const b = client.render(job(2)); const c = client.render(job(3)); expect(await settled(a)).toBeInstanceOf(Superseded); expect(await settled(b)).toBeInstanceOf(Superseded); gates[0](); await Promise.resolve(); await Promise.resolve(); expect(seen).toEqual([1, 3]); gates[1](); expect((await c).id).toBe(3); }); it("the auto client moves to the main thread when the worker cannot start", async () => { const client = createAutoClient({ workerFactory: () => { throw new Error("no worker"); }, mainExecutor: async (j) => result(j.id), }); expect(client.mode).toBe("worker"); expect((await client.render(job(7))).id).toBe(7); expect(client.mode).toBe("main"); expect((await client.render(job(8))).id).toBe(8); }); });