import { strToU8, zipSync } from "fflate"; import { describe, expect, it, vi } from "vitest"; import { encodeToolcraftModelDocument } from "../canonical/model-document-codec"; import { digestToolcraftModelDocument } from "../canonical/model-document-digest"; import { createValidModelDocument } from "../canonical/model-document-test-support"; import { TOOLCRAFT_DEFAULT_MODEL_IMPORT_LIMITS } from "../model-import-limits"; import { encodeToolcraftModelRepairPlanEnvelope } from "../topology/model-repair-plan-codec"; import { createToolcraftModelWorkerClient } from "./model-import-worker-client"; import { completeDecode, deferred, draftFor, emitDecodeLifecycle, FakeWorker, forgedDigest, importRequest, resultFor, } from "./model-import-worker-client-test-support"; import { createToolcraftModelWorkerReceiptDigest, digestToolcraftModelWorkerBytes, } from "./model-import-worker-receipt"; import { fakeRepairPlan } from "./model-import-worker-runtime-test-support"; describe("createToolcraftModelWorkerClient", () => { it("copies one archive into the worker and verifies bounded package output", async () => { const worker = new FakeWorker(); const archive = zipSync({ "scene.gltf": strToU8("{}") }); const archiveDigest = digestToolcraftModelWorkerBytes(archive); const onPackageResult = vi.fn(); const client = createToolcraftModelWorkerClient({ sha256: async (bytes) => digestToolcraftModelWorkerBytes(bytes), workerFactory: () => worker, }); const pending = client.extractModelPackage({ archive, archiveDigest, jobId: "extract-package", limits: TOOLCRAFT_DEFAULT_MODEL_IMPORT_LIMITS, }, { onPackageResult }); const [{ message, transfer }] = worker.posts; if (message.kind !== "extract-model-package") { throw new Error("Expected a package extraction request."); } expect(transfer).toEqual([message.payload.archive]); expect(message.payload.archive).not.toBe(archive.buffer); expect(archive.byteLength).toBeGreaterThan(0); worker.emitMessage({ generation: message.generation, jobId: message.jobId, kind: "progress", phase: "extracting", progress: 0.5, }); const entryBytes = strToU8("{}").buffer; const resultWithoutReceipt = { archiveDigest, entries: [{ bytes: entryBytes, compressedBytes: 4, contentDigest: digestToolcraftModelWorkerBytes(entryBytes), mimeType: "model/gltf+json", path: "scene.gltf", uncompressedBytes: entryBytes.byteLength, }], operation: "extract-model-package" as const, }; worker.emitMessage({ generation: message.generation, jobId: message.jobId, kind: "package-result", result: { ...resultWithoutReceipt, receiptDigest: createToolcraftModelWorkerReceiptDigest({ generation: message.generation, jobId: message.jobId, result: resultWithoutReceipt, }), }, }); await expect(pending).resolves.toMatchObject({ kind: "package-result", result: { archiveDigest }, }); expect(onPackageResult).toHaveBeenCalledTimes(1); }); it("rejects a package result whose entry digest does not match its bytes", async () => { const worker = new FakeWorker(); const archive = new Uint8Array([1, 2, 3]); const archiveDigest = digestToolcraftModelWorkerBytes(archive); const onPackageResult = vi.fn(); const client = createToolcraftModelWorkerClient({ sha256: async (bytes) => digestToolcraftModelWorkerBytes(bytes), workerFactory: () => worker, }); const pending = client.extractModelPackage({ archive, archiveDigest, jobId: "forged-package", limits: TOOLCRAFT_DEFAULT_MODEL_IMPORT_LIMITS, }, { onPackageResult }); const request = worker.posts[0]!.message; if (request.kind !== "extract-model-package") { throw new Error("Expected a package extraction request."); } worker.emitMessage({ generation: request.generation, jobId: request.jobId, kind: "progress", phase: "extracting", progress: 0.5, }); const entryBytes = strToU8("{}").buffer; const resultWithoutReceipt = { archiveDigest, entries: [{ bytes: entryBytes, compressedBytes: 4, contentDigest: `sha256:${"0".repeat(64)}`, mimeType: "model/gltf+json", path: "scene.gltf", uncompressedBytes: entryBytes.byteLength, }], operation: "extract-model-package" as const, }; worker.emitMessage({ generation: request.generation, jobId: request.jobId, kind: "package-result", result: { ...resultWithoutReceipt, receiptDigest: createToolcraftModelWorkerReceiptDigest({ generation: request.generation, jobId: request.jobId, result: resultWithoutReceipt, }), }, }); await expect(pending).resolves.toMatchObject({ feedback: { code: "model-worker-invalid-message" }, kind: "error", }); expect(onPackageResult).not.toHaveBeenCalled(); expect(worker.terminate).toHaveBeenCalledTimes(1); }); it("copies and transfers source buffers exactly once without detaching caller-owned originals", () => { const worker = new FakeWorker(); const original = new Uint8Array([7, 8, 9]).buffer; const client = createToolcraftModelWorkerClient({ workerFactory: () => worker }); void client.importModel(importRequest("job-1", original)); expect(new Uint8Array(original)).toEqual(new Uint8Array([7, 8, 9])); expect(worker.posts).toHaveLength(1); const [{ message, transfer }] = worker.posts; if (message.kind !== "decode-and-analyze") throw new Error("Expected decode request."); expect(transfer).toHaveLength(1); expect(new Set(transfer).size).toBe(transfer.length); expect(transfer[0]).toBe(message.payload.bundle.sourceFiles[0]!.bytes); expect(transfer[0]).not.toBe(original); }); it("hard-terminates active decode on replacement, rejects stale publication, and uses a fresh worker", async () => { const workers: FakeWorker[] = []; const factory = vi.fn(() => { const worker = new FakeWorker(); workers.push(worker); return worker; }); const onResult = vi.fn(); const onDiagnostic = vi.fn(); const client = createToolcraftModelWorkerClient({ workerFactory: factory }); const first = client.importModel(importRequest("job-1"), { onDiagnostic, onResult }); const firstRequest = workers[0]!.posts[0]!.message; workers[0]!.emitMessage({ generation: firstRequest.generation, jobId: firstRequest.jobId, kind: "progress", phase: "decoding", progress: 0.2, }); const second = client.importModel(importRequest("job-2"), { onResult }); expect(workers[0]!.terminate).toHaveBeenCalledTimes(1); expect(workers).toHaveLength(2); expect(workers[1]).not.toBe(workers[0]); await expect(first).resolves.toMatchObject({ jobId: "job-1", kind: "cancelled" }); workers[0]!.emitRemovedMessage({ diagnostic: { affectedCount: 1, code: "stale", explanation: "Stale diagnostic.", severity: "warning", }, generation: firstRequest.generation, jobId: firstRequest.jobId, kind: "diagnostic", }); workers[0]!.emitRemovedMessage(resultFor(firstRequest, forgedDigest)); expect(onDiagnostic).not.toHaveBeenCalled(); expect(onResult).not.toHaveBeenCalled(); const secondRequest = workers[1]!.posts[0]!.message; expect(secondRequest.generation).toBeGreaterThan(firstRequest.generation); completeDecode(workers[1]!, secondRequest); await expect(second).resolves.toMatchObject({ kind: "result", result: { canonicalDocumentDigest: digestToolcraftModelDocument(createValidModelDocument()) }, }); expect(onResult).toHaveBeenCalledTimes(1); }); it("hard-terminates an unresponsive compressed decode on explicit cancellation", async () => { const workers: FakeWorker[] = []; const client = createToolcraftModelWorkerClient({ workerFactory: () => { const worker = new FakeWorker(); workers.push(worker); return worker; }, }); const pending = client.importModel(importRequest("compressed")); const request = workers[0]!.posts[0]!.message; workers[0]!.emitMessage({ generation: request.generation, jobId: request.jobId, kind: "progress", phase: "decoding", progress: 0.5, }); client.cancel("compressed"); expect(workers[0]!.terminate).toHaveBeenCalledTimes(1); await expect(pending).resolves.toMatchObject({ kind: "cancelled" }); expect(workers[0]!.posts.some(({ message }) => message.kind === "cancel")).toBe(false); void client.importModel(importRequest("retry")); expect(workers).toHaveLength(2); }); it("publishes only monotonic current-generation progress below one", async () => { const worker = new FakeWorker(); const onProgress = vi.fn(); const client = createToolcraftModelWorkerClient({ workerFactory: () => worker }); const pending = client.importModel(importRequest("job-1"), { onProgress }); const request = worker.posts[0]!.message; for (const progress of [0.2, 0.4, 0.8]) { worker.emitMessage({ generation: request.generation, jobId: request.jobId, kind: "progress", phase: "decoding", progress, }); } worker.emitMessage(draftFor(request)); worker.emitMessage({ generation: request.generation, jobId: request.jobId, kind: "progress", phase: "analyzing", progress: 0.8, }); worker.emitMessage(resultFor(request)); await pending; expect(onProgress.mock.calls.map(([event]) => event.progress)).toEqual([0.2, 0.4, 0.8, 0.8]); }); it.each(["onDraft", "onProgress", "onDiagnostic", "onResult"] as const)( "isolates a throwing %s observer, settles once, and keeps the worker reusable", async (observerName) => { const worker = new FakeWorker(); const factory = vi.fn(() => worker); const observer = vi.fn(() => { throw new Error(`${observerName} must not control worker lifecycle`); }); const client = createToolcraftModelWorkerClient({ sha256: async (bytes) => digestToolcraftModelWorkerBytes(bytes), workerFactory: factory, }); const first = client.importModel(importRequest("throwing-observer"), { [observerName]: observer, }); const firstRequest = worker.posts[0]!.message; if (observerName === "onProgress") { expect(() => worker.emitMessage({ generation: firstRequest.generation, jobId: firstRequest.jobId, kind: "progress", phase: "decoding", progress: 0.1, })).not.toThrow(); worker.emitMessage({ generation: firstRequest.generation, jobId: firstRequest.jobId, kind: "progress", phase: "decoding", progress: 0.9, }); worker.emitMessage(draftFor(firstRequest)); worker.emitMessage({ generation: firstRequest.generation, jobId: firstRequest.jobId, kind: "progress", phase: "analyzing", progress: 0.9, }); } else { emitDecodeLifecycle(worker, firstRequest); } if (observerName === "onDiagnostic") { expect(() => worker.emitMessage({ diagnostic: { affectedCount: 1, code: "observer-test", explanation: "Observer exceptions are isolated.", severity: "info", }, generation: firstRequest.generation, jobId: firstRequest.jobId, kind: "diagnostic", })).not.toThrow(); } worker.emitMessage(resultFor(firstRequest)); await expect(first).resolves.toMatchObject({ kind: "result" }); expect(observer).toHaveBeenCalledTimes(observerName === "onProgress" ? 3 : 1); expect(worker.terminate).not.toHaveBeenCalled(); const retry = client.importModel(importRequest("observer-retry")); completeDecode(worker, worker.posts[1]!.message); await expect(retry).resolves.toMatchObject({ kind: "result" }); expect(factory).toHaveBeenCalledTimes(1); }, ); it("cancels without publication when superseded during asynchronous digest verification", async () => { const workers: FakeWorker[] = []; const verification = deferred(); const sha256 = vi.fn() .mockImplementationOnce(() => verification.promise) .mockImplementation(async (bytes) => digestToolcraftModelWorkerBytes(bytes)); const firstOnDraft = vi.fn(); const firstOnResult = vi.fn(); const client = createToolcraftModelWorkerClient({ sha256, workerFactory: () => { const worker = new FakeWorker(); workers.push(worker); return worker; }, }); const first = client.importModel(importRequest("verifying-first"), { onDraft: firstOnDraft, onResult: firstOnResult, }); const firstRequest = workers[0]!.posts[0]!.message; emitDecodeLifecycle(workers[0]!, firstRequest); workers[0]!.emitMessage(resultFor(firstRequest)); expect(sha256).toHaveBeenCalledTimes(1); const second = client.importModel(importRequest("verifying-second")); await expect(first).resolves.toMatchObject({ jobId: "verifying-first", kind: "cancelled", }); expect(workers[0]!.terminate).toHaveBeenCalledTimes(1); expect(workers).toHaveLength(2); verification.resolve(digestToolcraftModelDocument(createValidModelDocument())); await Promise.resolve(); await Promise.resolve(); expect(firstOnDraft).not.toHaveBeenCalled(); expect(firstOnResult).not.toHaveBeenCalled(); const secondRequest = workers[1]!.posts[0]!.message; completeDecode(workers[1]!, secondRequest); await expect(second).resolves.toMatchObject({ kind: "result" }); expect(sha256).toHaveBeenCalledTimes(5); }); it.each(["error", "messageerror"] as const)( "maps worker %s, cleans up listeners, and creates a fresh worker for retry", async (eventType) => { const workers: FakeWorker[] = []; const client = createToolcraftModelWorkerClient({ workerFactory: () => { const worker = new FakeWorker(); workers.push(worker); return worker; }, }); const pending = client.importModel(importRequest("job-1")); workers[0]!.emit(eventType); const outcome = await pending; expect(outcome).toMatchObject({ kind: "error", feedback: { category: "resource-unavailable", code: eventType === "error" ? "model-worker-error" : "model-worker-message-error", }, }); if (outcome.kind !== "error") throw new Error("Expected error."); expect(outcome.feedback.message.length).toBeLessThanOrEqual(240); expect(workers[0]!.terminate).toHaveBeenCalledTimes(1); expect([...workers[0]!.listeners.values()].every((set) => set.size === 0)).toBe(true); void client.importModel(importRequest("job-2")); expect(workers).toHaveLength(2); }, ); it("keeps the live worker reusable after a typed job error", async () => { const worker = new FakeWorker(); const factory = vi.fn(() => worker); const client = createToolcraftModelWorkerClient({ workerFactory: factory }); const first = client.importModel(importRequest("job-1")); const firstRequest = worker.posts[0]!.message; worker.emitMessage({ feedback: { category: "format", code: "model-decode-failed", message: "Decode failed.", }, generation: firstRequest.generation, jobId: firstRequest.jobId, kind: "error", }); await expect(first).resolves.toMatchObject({ kind: "error" }); const second = client.importModel(importRequest("job-2")); const secondRequest = worker.posts[1]!.message; completeDecode(worker, secondRequest); await expect(second).resolves.toMatchObject({ kind: "result" }); expect(factory).toHaveBeenCalledTimes(1); expect(worker.terminate).not.toHaveBeenCalled(); }); it("transfers a repair snapshot and disposes idempotently", async () => { const worker = new FakeWorker(); const client = createToolcraftModelWorkerClient({ workerFactory: () => worker }); const original = encodeToolcraftModelDocument(createValidModelDocument()); const plan = fakeRepairPlan(); const planBytes = encodeToolcraftModelRepairPlanEnvelope(plan); const pending = client.repair({ canonicalDocument: original, canonicalDocumentDigest: digestToolcraftModelDocument(createValidModelDocument()), jobId: "repair-1", limits: TOOLCRAFT_DEFAULT_MODEL_IMPORT_LIMITS, repairPlanEnvelope: { bytes: planBytes.buffer, envelopeDigest: digestToolcraftModelWorkerBytes(planBytes), planDigest: plan.planDigest, version: 1, }, }); await vi.waitFor(() => expect(worker.posts).toHaveLength(1)); const posted = worker.posts[0]!; if (posted.message.kind !== "repair") throw new Error("Expected repair request."); expect(posted.transfer).toEqual([ posted.message.payload.canonicalDocument, posted.message.payload.repairPlanEnvelope.bytes, ]); expect(posted.message.payload.canonicalDocument).not.toBe(original.buffer); expect(posted.message.payload.repairPlanEnvelope.bytes).not.toBe(planBytes.buffer); expect(original.byteLength).toBeGreaterThan(0); expect(planBytes.byteLength).toBeGreaterThan(0); client.dispose(); client.dispose(); expect(worker.terminate).toHaveBeenCalledTimes(1); await expect(pending).resolves.toMatchObject({ kind: "cancelled" }); expect(() => client.importModel(importRequest("after-dispose"))).toThrow(/disposed/u); }); });