import { Worker } from "node:worker_threads"; import { performance } from "node:perf_hooks"; import { generateMap } from "../src/mapPipeline.js"; import { createWorldMap } from "../src/worldMap.js"; function assert(condition, message) { if (!condition) throw new Error(message); console.log(`OK: ${message}`); } function isPlainObject(value) { if (!value || typeof value !== "object" || Array.isArray(value) || ArrayBuffer.isView(value)) return false; const prototype = Object.getPrototypeOf(value); return prototype === Object.prototype || prototype === null; } function sourceEntryNeedsChunking(value) { return isPlainObject(value) && Object.values(value).some((entry) => ArrayBuffer.isView(entry) || entry instanceof ArrayBuffer); } function manifestFor(world) { const sourceKeys = Object.keys(world.sourceMap || {}); const expandedSourceObjects = {}; for (const key of sourceKeys) if (sourceEntryNeedsChunking(world.sourceMap[key])) expandedSourceObjects[key] = Object.keys(world.sourceMap[key]); return { rootKeys: Object.keys(world).filter((key) => key !== "fields" && key !== "sourceMap" && key !== "generatedMask"), fieldKeys: Object.keys(world.fields || {}), sourceKeys, expandedSourceObjects, hasGeneratedMask: !!world.generatedMask, }; } function transferCopy(value) { if (ArrayBuffer.isView(value)) { if (value instanceof DataView) { const buffer = value.buffer.slice(value.byteOffset, value.byteOffset + value.byteLength); return { value: new DataView(buffer), transfer: [buffer], bytes: buffer.byteLength }; } const copy = new value.constructor(value); return { value: copy, transfer: [copy.buffer], bytes: copy.byteLength }; } if (value instanceof ArrayBuffer) { const copy = value.slice(0); return { value: copy, transfer: [copy], bytes: copy.byteLength }; } return { value, transfer: [], bytes: 0 }; } function waitFor(worker, predicate, timeoutMs = 120_000) { return new Promise((resolve, reject) => { const timer = setTimeout(() => done(new Error("Worker response timed out.")), timeoutMs); const onMessage = (message) => { if (predicate(message)) done(null, message); }; const onError = (error) => done(error); const done = (error, value) => { clearTimeout(timer); worker.off("message", onMessage); worker.off("error", onError); if (error) reject(error); else resolve(value); }; worker.on("message", onMessage); worker.on("error", onError); }); } async function sendSync(worker, envelope, type, payload = {}, transfer = []) { const sequence = ++envelope.sequence; const reply = waitFor(worker, (message) => message?.type === "patch-mirror-sync-ack" && message.id === envelope.id && message.syncId === envelope.syncId && message.sequence === sequence); worker.postMessage({ id: envelope.id, type, syncId: envelope.syncId, sequence, ...payload }, transfer); const ack = await reply; if (!ack.ok) throw new Error(ack.error || `${type} failed`); return ack; } async function main() { const initial = generateMap(114514); const world = createWorldMap(initial); const originalElevationBuffer = world.fields.elevation.buffer; const worker = new Worker(new URL("./browser-worker-node-shim.mjs", import.meta.url), { type: "module" }); const envelope = { id: 91, syncId: "node-cold-sync:1", sequence: 0 }; const manifest = manifestFor(world); let maxBinaryChunkBytes = 0; const syncStartedAt = performance.now(); try { await sendSync(worker, envelope, "patch-mirror-sync-start", { committedRevision: 1, manifest }); for (const key of manifest.rootKeys) { const prepared = transferCopy(world[key]); maxBinaryChunkBytes = Math.max(maxBinaryChunkBytes, prepared.bytes); await sendSync(worker, envelope, "patch-mirror-sync-root", { key, value: prepared.value }, prepared.transfer); } for (const key of manifest.fieldKeys) { const prepared = transferCopy(world.fields[key]); maxBinaryChunkBytes = Math.max(maxBinaryChunkBytes, prepared.bytes); await sendSync(worker, envelope, "patch-mirror-sync-field", { key, value: prepared.value }, prepared.transfer); } if (manifest.hasGeneratedMask) { const prepared = transferCopy(world.generatedMask); maxBinaryChunkBytes = Math.max(maxBinaryChunkBytes, prepared.bytes); await sendSync(worker, envelope, "patch-mirror-sync-generated-mask", { value: prepared.value }, prepared.transfer); } for (const key of manifest.sourceKeys) { const value = world.sourceMap[key]; const childKeys = manifest.expandedSourceObjects[key]; if (Array.isArray(childKeys)) { await sendSync(worker, envelope, "patch-mirror-sync-source-object-start", { key }); for (const childKey of childKeys) { const prepared = transferCopy(value[childKey]); maxBinaryChunkBytes = Math.max(maxBinaryChunkBytes, prepared.bytes); await sendSync(worker, envelope, "patch-mirror-sync-source-object-entry", { key, childKey, value: prepared.value }, prepared.transfer); } } else { const prepared = transferCopy(value); maxBinaryChunkBytes = Math.max(maxBinaryChunkBytes, prepared.bytes); await sendSync(worker, envelope, "patch-mirror-sync-source", { key, value: prepared.value }, prepared.transfer); } } const finish = await sendSync(worker, envelope, "patch-mirror-sync-finish"); const syncMs = performance.now() - syncStartedAt; assert(finish.mirrorCommittedRevision === 1, "cold mirror bootstrap installs the requested committed revision only after finish"); assert(world.fields.elevation.buffer === originalElevationBuffer && world.fields.elevation.byteLength > 0, "cold mirror bootstrap transfers copies and never detaches the main committed rasters"); assert(maxBinaryChunkBytes < 1024 * 1024, `cold mirror binary synchronization remains sub-megabyte per main-thread dispatch (${maxBinaryChunkBytes} bytes max)`); assert(syncMs < 10_000, `cold mirror synchronization completes without a whole-world structured-clone stall (${Math.round(syncMs)} ms)`); const rect = { x0: world.originX + 72, y0: world.originY + 58, x1: world.originX + 132, y1: world.originY + 118, }; const candidates = [1, 2, 3].map((variant, index) => ({ candidateId: `node-sync:${variant}`, candidateOrdinal: index + 1, variant, seed: (0x51a7c3d3 + Math.imul(index, 0x9e3779b9)) >>> 0, })); const resultPromise = waitFor(worker, (message) => message?.id === 92 && message?.type !== "progress", 180_000); worker.postMessage({ id: 92, world: null, rect, options: { patchMode: "regeneration", terrainType: "auto", variant: candidates[0].variant, seed: candidates[0].seed, maxQualityRetries: 0, qualityTerrainAttempts: 1, acceptBestAvailableQuality: true, includeSeamVisualization: false, }, search: { searchId: "node-sync-search", operationId: "node-sync-search", committedRevision: 1, workerEpoch: 1, executionAttempt: 1, totalCandidateCount: candidates.length, reuseCommittedMirror: true, resolvedPatchMode: "regeneration", selectBestCandidate: true, draftSelection: true, candidatePlan: candidates, }, }); const result = await resultPromise; assert(result.ok === true && result.result?.ok === true, "a cold-synchronized mirror can run draft-ranked production without retransmitting world"); const draftDebug = result.result?.draftSelection; const exactBoundBehavior = draftDebug?.fullCandidateCount < 3 ? draftDebug?.branchBoundPrunedCount > 0 : draftDebug?.fullCandidateCount === 3 && draftDebug?.branchBoundPrunedCount === 0 && draftDebug?.fullComparedCandidateOrdinals?.length === 3; assert(draftDebug?.enabled === true && draftDebug?.draftCount === 3 && draftDebug?.fullCandidateCount >= 1 && draftDebug?.fullCandidateCount <= 3 && draftDebug?.reusedWinningDraft === false && draftDebug?.fullProductionFromTerrainOnly === true && exactBoundBehavior, "real Worker production ranks three terrain scouts, prunes only mathematically dominated candidates, and fully compares all candidates when admissible bounds overlap"); assert(result.result?.acceptedWorldHash && result.result?.applyToken, "draft-ranked finalist returns transactional hash and Apply token"); const ackId = "node-sync-apply"; const applyPromise = waitFor(worker, (message) => message?.type === "patch-apply-ack-result" && message.ackId === ackId); worker.postMessage({ type: "patch-apply-ack", ackId, applyToken: result.result.applyToken, baseCommittedRevision: 1, committedRevision: 2, }); const applyAck = await applyPromise; assert(applyAck.ok === true && applyAck.mirrorCommittedRevision === 2, "Apply ACK advances the synchronized persistent mirror by exactly one revision"); assert(applyAck.mirrorHash === result.result.acceptedWorldHash, "Apply ACK mirror hash matches the accepted candidate hash"); } finally { await worker.terminate(); } } await main(); console.log("All patch Worker mirror synchronization tests passed.");