bridge-common.mjs raw

   1  // bridge-common.mjs - shared wasm bridge infrastructure.
   2  //
   3  // Exports utility functions and a factory for the common bridge functions
   4  // that appear identically in every wasm Worker host: SAB channel ops,
   5  // spawn_domain, is_worker_context, subtle_random_bytes, core helpers, WASI.
   6  
   7  // ── SAB ring buffer ────────────────────────────────────────────────────────
   8  
   9  export const SAB_HDR = 32;
  10  export const SAB_SLOT_HDR = 12;
  11  
  12  export function sabCreate(slotSize, slotCount) {
  13    const sab = new SharedArrayBuffer(SAB_HDR + slotCount * slotSize);
  14    const u32 = new Uint32Array(sab);
  15    u32[2] = slotSize;
  16    u32[3] = slotCount;
  17    return sab;
  18  }
  19  
  20  // forceClose sets the closed flag and wakes all blocked waiters.
  21  // No sentinel slot needed: send/recv wait loops check bit0 on every wakeup.
  22  export function sabForceClose(sab) {
  23    const u32 = new Uint32Array(sab);
  24    Atomics.or(u32, 4, 1);
  25    Atomics.notify(u32, 0, Atomics.MAX_VALUE);
  26    Atomics.notify(u32, 1, Atomics.MAX_VALUE);
  27  }
  28  
  29  export function sabSend(sab, bytes) {
  30    const u32 = new Uint32Array(sab);
  31    const slotSize = u32[2], slotCount = u32[3];
  32    const payloadCap = slotSize - SAB_SLOT_HDR;
  33    if (bytes.length > (slotCount - 1) * payloadCap)
  34      throw new Error('channel_send: overflow');
  35    let offset = 0, first = true;
  36    while (offset <= bytes.length) {
  37      const chunkLen = Math.min(bytes.length - offset, payloadCap);
  38      const isFinal = (offset + chunkLen >= bytes.length) ? 1 : 0;
  39      let wi = Atomics.load(u32, 0);
  40      while (wi - Atomics.load(u32, 1) >= slotCount) {
  41        Atomics.wait(u32, 1, Atomics.load(u32, 1));
  42        wi = Atomics.load(u32, 0);
  43        if (Atomics.load(u32, 4) & 1) throw new Error('channel_send: send on closed channel');
  44      }
  45      if (Atomics.load(u32, 4) & 1) throw new Error('channel_send: send on closed channel');
  46      const base = SAB_HDR + (wi & (slotCount - 1)) * slotSize;
  47      const dv = new DataView(sab, base, slotSize);
  48      dv.setUint8(0, isFinal);
  49      dv.setUint8(1, 0); dv.setUint8(2, 0); dv.setUint8(3, 0);
  50      dv.setUint32(4, first ? bytes.length : 0, true);
  51      dv.setUint32(8, chunkLen, true);
  52      if (chunkLen > 0)
  53        new Uint8Array(sab, base + SAB_SLOT_HDR, chunkLen)
  54          .set(bytes.subarray(offset, offset + chunkLen));
  55      Atomics.store(u32, 0, wi + 1);
  56      Atomics.notify(u32, 0, 1);
  57      offset += chunkLen;
  58      first = false;
  59      if (isFinal) break;
  60    }
  61  }
  62  
  63  export function sabRecv(sab, dstBuf) {
  64    const u32 = new Uint32Array(sab);
  65    const slotSize = u32[2], slotCount = u32[3];
  66    let totalLen = -1, dstOffset = 0, awaitingFirst = true;
  67    for (;;) {
  68      let ri = Atomics.load(u32, 1);
  69      while (Atomics.load(u32, 0) === ri) {
  70        Atomics.wait(u32, 0, ri);
  71        ri = Atomics.load(u32, 1);
  72        if (Atomics.load(u32, 4) & 1) return -1;
  73      }
  74      const base = SAB_HDR + (ri & (slotCount - 1)) * slotSize;
  75      const dv = new DataView(sab, base, slotSize);
  76      const isFinal = dv.getUint8(0);
  77      const slotTotalLen = dv.getUint32(4, true);
  78      const chunkLen = dv.getUint32(8, true);
  79      Atomics.store(u32, 1, ri + 1);
  80      Atomics.notify(u32, 1, 1);
  81      if (awaitingFirst) {
  82        if (slotTotalLen === 0)
  83          return (Atomics.load(u32, 4) & 1) ? -1 : 0;
  84        if (slotTotalLen > dstBuf.length)
  85          throw new Error('channel_recv: overflow');
  86        totalLen = slotTotalLen;
  87        awaitingFirst = false;
  88      }
  89      if (chunkLen > 0) {
  90        dstBuf.set(new Uint8Array(sab, base + SAB_SLOT_HDR, chunkLen), dstOffset);
  91        dstOffset += chunkLen;
  92      }
  93      if (isFinal) return totalLen;
  94    }
  95  }
  96  
  97  export function sabClose(sab) {
  98    const u32 = new Uint32Array(sab);
  99    Atomics.or(u32, 4, 1);
 100    const slotSize = u32[2], slotCount = u32[3];
 101    let wi = Atomics.load(u32, 0);
 102    while (wi - Atomics.load(u32, 1) >= slotCount) {
 103      Atomics.wait(u32, 1, Atomics.load(u32, 1));
 104      wi = Atomics.load(u32, 0);
 105    }
 106    const base = SAB_HDR + (wi & (slotCount - 1)) * slotSize;
 107    const dv = new DataView(sab, base, slotSize);
 108    dv.setUint8(0, 1); dv.setUint32(4, 0, true); dv.setUint32(8, 0, true);
 109    Atomics.store(u32, 0, wi + 1);
 110    Atomics.notify(u32, 0, 1);
 111  }
 112  
 113  // ── fromSlice ──────────────────────────────────────────────────────────────
 114  
 115  export function fromSlice(s) {
 116    if (s instanceof Uint8Array) return s;
 117    if (s && s.$array != null) {
 118      const u = new Uint8Array(s.$length);
 119      for (let i = 0; i < s.$length; i++) u[i] = s.$array[s.$offset + i];
 120      return u;
 121    }
 122    if (s && typeof s === 'string') return new TextEncoder().encode(s);
 123    return new Uint8Array(0);
 124  }
 125  
 126  // ── Core helpers factory ───────────────────────────────────────────────────
 127  
 128  // makeCoreHelpers returns memory-access helpers that close over getMem/getXp.
 129  // Call after building the bridge object; the getters provide the live wasm
 130  // instance so helpers always see the current memory even after wasm growth.
 131  export function makeCoreHelpers(getMem, getXp) {
 132    const enc = new TextEncoder();
 133    const dec = new TextDecoder();
 134  
 135    function readStr(ptr, len) {
 136      if (len <= 0) return '';
 137      return dec.decode(new Uint8Array(getMem().buffer, ptr >>> 0, len));
 138    }
 139    function readBytes(ptr, len) {
 140      if (len <= 0) return new Uint8Array(0);
 141      return new Uint8Array(getMem().buffer, ptr >>> 0, len);
 142    }
 143    function writeStr(s) {
 144      const bytes = enc.encode('' + s);
 145      const ptr = getXp().__alloc(bytes.length) >>> 0;
 146      new Uint8Array(getMem().buffer, ptr, bytes.length).set(bytes);
 147      return [ptr, bytes.length];
 148    }
 149    function writeI32(addr, val) {
 150      new DataView(getMem().buffer).setInt32(addr >>> 0, val, true);
 151    }
 152    function writeBytes(data) {
 153      const u = fromSlice(data);
 154      const ptr = getXp().__alloc(u.length) >>> 0;
 155      new Uint8Array(getMem().buffer, ptr, u.length).set(u);
 156      return [ptr, u.length];
 157    }
 158    function cb0(id) { getXp().__cb0(id); }
 159    function cbs(id, s) {
 160      const [ptr, len] = writeStr(s);
 161      getXp().__cbs(id, ptr, len);
 162    }
 163    function cbdata(id, data) {
 164      const [ptr, len] = writeBytes(data);
 165      getXp().__cbdata(id, ptr, len);
 166    }
 167  
 168    return { readStr, readBytes, writeStr, writeI32, writeBytes, cb0, cbs, cbdata };
 169  }
 170  
 171  // ── WASI ───────────────────────────────────────────────────────────────────
 172  
 173  // makeEnv returns the 'env' import object for the host shims pkg/mxutil
 174  // declares (write + the mxc_* file helpers). These are reached by app and
 175  // module wasm binaries that import mxutil, so every instantiate call must
 176  // provide the namespace or the module fails with
 177  // "Import #1 env: module is not an object or function".
 178  export function makeEnv(getMem) {
 179    const dec = new TextDecoder();
 180    const fail = () => -1;
 181    return {
 182      // No filesystem inside a worker: fail the file shims cleanly.
 183      mxc_readfile: fail,
 184      mxc_filesize: fail,
 185      mxc_listdir: fail,
 186      mxc_writefile: fail,
 187      mxc_mkdir: fail,
 188      mxc_getenv: fail,
 189      // write(fd, ptr, count): route 1/2 to the worker console.
 190      write(fd, ptr, count) {
 191        const n = count >>> 0;
 192        if ((fd === 1 || fd === 2) && n > 0) {
 193          const s = dec.decode(new Uint8Array(getMem().buffer, ptr >>> 0, n));
 194          if (fd === 1) console.log(s.replace(/\n$/, ''));
 195          else console.error(s.replace(/\n$/, ''));
 196        }
 197        return n;
 198      },
 199      'syscall.Exit'(code) {
 200        self.close();
 201        throw new Error('wasm exit ' + code);
 202      },
 203    };
 204  }
 205  
 206  export function makeWasi(getMem) {
 207    const dec = new TextDecoder();
 208    const _lineBuf = { 1: '', 2: '' };
 209    return {
 210      fd_write(fd, iovs, iovs_len, nwritten_ptr) {
 211        const dv = new DataView(getMem().buffer);
 212        let total = 0;
 213        const iovBase = iovs >>> 0;
 214        for (let i = 0; i < iovs_len; i++) {
 215          const ptr = dv.getUint32(iovBase + i * 8, true);
 216          const len = dv.getUint32(iovBase + i * 8 + 4, true);
 217          const s = dec.decode(new Uint8Array(getMem().buffer, ptr, len));
 218          if (fd === 1 || fd === 2) {
 219            _lineBuf[fd] += s;
 220            let nl;
 221            while ((nl = _lineBuf[fd].indexOf('\n')) !== -1) {
 222              const line = _lineBuf[fd].slice(0, nl);
 223              _lineBuf[fd] = _lineBuf[fd].slice(nl + 1);
 224              if (fd === 1) console.log(line);
 225              else console.error(line);
 226            }
 227          }
 228          total += len;
 229        }
 230        dv.setUint32(nwritten_ptr >>> 0, total, true);
 231        return 0;
 232      },
 233      // The app imports fd_read (os.Stdin read paths). There is no stdin in a
 234      // worker, so report EOF: nread = 0 and a success errno.
 235      fd_read(fd, iovs, iovs_len, nread_ptr) {
 236        const dv = new DataView(getMem().buffer);
 237        dv.setUint32(nread_ptr >>> 0, 0, true);
 238        return 0;
 239      },
 240  
 241      clock_time_get(clock_id, precision, time_ptr) {
 242        const dv = new DataView(getMem().buffer);
 243        const ms = (clock_id === 1) ? performance.now() : Date.now();
 244        const ns = BigInt(Math.trunc(ms * 1e6));
 245        dv.setBigUint64(time_ptr >>> 0, ns, true);
 246        return 0;
 247      },
 248    };
 249  }
 250  
 251  // ── Build hash ─────────────────────────────────────────────────────────────
 252  
 253  export async function computeBuildHash(wasmBytes) {
 254    const buf = await crypto.subtle.digest('SHA-256', wasmBytes);
 255    return Array.from(new Uint8Array(buf))
 256      .map(b => b.toString(16).padStart(2, '0')).join('');
 257  }
 258  
 259  // ── Spawned Worker lifecycle ───────────────────────────────────────────────
 260  
 261  export function createSpawnedWorker(wasmUrl, buildHash, fnIdx, argBytes, chanSabs, spawnedWorkerChans) {
 262    const worker = new Worker(new URL(wasmUrl), { type: 'module' });
 263    worker.postMessage({ type: 'init', mode: 'spawn', wasmUrl, buildHash, fnIdx, argBytes, chanSabs });
 264    spawnedWorkerChans.set(worker, chanSabs);
 265    worker.onmessage = function(e) {
 266      if (e.data && e.data.type === 'exit') spawnedWorkerChans.delete(worker);
 267    };
 268    worker.onerror = function(e) {
 269      const sabs = spawnedWorkerChans.get(worker);
 270      spawnedWorkerChans.delete(worker);
 271      console.error('[spawn] child Worker died:', e.message,
 272        'at', e.filename + ':' + e.lineno,
 273        e.error && e.error.stack ? '\n' + e.error.stack : '');
 274      if (sabs) sabs.forEach(sabForceClose);
 275    };
 276    worker.onmessageerror = function() {
 277      const sabs = spawnedWorkerChans.get(worker);
 278      spawnedWorkerChans.delete(worker);
 279      if (sabs) sabs.forEach(sabForceClose);
 280    };
 281  }
 282  
 283  // ── Common bridge factory ──────────────────────────────────────────────────
 284  
 285  // makeCommonBridge returns the bridge functions shared by all wasm Worker
 286  // hosts. The host merges these with its own specific functions (ext_*, dom_*,
 287  // etc.) to form the complete bridge object passed to WebAssembly.instantiate.
 288  //
 289  // h: the result of makeCoreHelpers(getMem, getXp)
 290  // sabTable: Map<int32, SharedArrayBuffer> - handle table (mutated by channel_create)
 291  // sabSeqRef: { value: int32 } - next handle sequence number (mutable)
 292  // wasmUrlRef: { value: string|null } - set during wasm boot
 293  // buildHashRef: { value: string|null } - set during wasm boot
 294  // spawnedWorkerChans: Map<Worker, SAB[]> - dead-Worker cleanup tracking
 295  export function makeCommonBridge(h, sabTable, sabSeqRef, wasmUrlRef, buildHashRef, spawnedWorkerChans) {
 296    const { readStr, readBytes, writeBytes, writeI32, cbs, cbdata } = h;
 297  
 298    return {
 299      // --- runtime context ---
 300      is_worker_context() { return 1; },
 301  
 302      // --- SAB channels ---
 303      channel_create(slotSize, slotCount) {
 304        const sab = sabCreate(slotSize, slotCount);
 305        const handle = sabSeqRef.value++;
 306        sabTable.set(handle, sab);
 307        return handle;
 308      },
 309      channel_send(handle, srcPtr, srcLen, srcCap) {
 310        const sab = sabTable.get(handle);
 311        if (!sab) throw new Error('channel_send: invalid handle ' + handle);
 312        // .slice() copies bytes out before the Atomics.wait in sabSend can
 313        // conflict with wasm advancing the memory pointer during a GC.
 314        sabSend(sab, readBytes(srcPtr, srcLen).slice());
 315      },
 316      channel_recv(handle, dstPtr, dstCap, cap) {
 317        const sab = sabTable.get(handle);
 318        if (!sab) return -1;
 319        const dst = new Uint8Array(dstCap);
 320        const n = sabRecv(sab, dst);
 321        if (n > 0) {
 322          // readBytes returns a writable view into wasm linear memory at dstPtr.
 323          // set() writes the received bytes directly to that address.
 324          readBytes(dstPtr, n).set(dst.subarray(0, n));
 325        }
 326        return n;
 327      },
 328      channel_close(handle) {
 329        const sab = sabTable.get(handle);
 330        if (sab) sabClose(sab);
 331      },
 332  
 333      // --- spawn ---
 334      spawn_domain(fnIdx, argPtr, argLen, chanHandlesPtr, nChans) {
 335        if (!wasmUrlRef.value) throw new Error('spawn_domain: wasmUrl not set');
 336        const chanSabs = [];
 337        for (let i = 0; i < nChans; i++) {
 338          // readBytes returns a view; DataView reads the int32 handle at each offset.
 339          const handleData = readBytes(chanHandlesPtr + i * 4, 4);
 340          const hh = new DataView(handleData.buffer, handleData.byteOffset, 4).getInt32(0, true);
 341          const sab = sabTable.get(hh);
 342          if (sab) chanSabs.push(sab);
 343        }
 344        const argBytes = argLen > 0 ? readBytes(argPtr, argLen).slice() : new Uint8Array(0);
 345        createSpawnedWorker(wasmUrlRef.value, buildHashRef.value, fnIdx, argBytes, chanSabs, spawnedWorkerChans);
 346      },
 347  
 348      // --- time ---
 349      timezone_offset_minutes() {
 350        return new Date().getTimezoneOffset();
 351      },
 352  
 353      // --- subtle ---
 354      subtle_random_bytes(ptr, len, cap) {
 355        crypto.getRandomValues(readBytes(ptr, len));
 356      },
 357    };
 358  }
 359