Spaces:
Running
Running
File size: 9,947 Bytes
cc08b9f | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 | /* rows.js — rows of a large lookup table (an embedding-style table) served from OPFS into a small GPU-resident
* cache. Extracted from LocalMind's ple-opfs.js (2026-10-05), where it serves Gemma 4 E2B/E4B's per-layer embedding
* table from disk; the Gemma-specific helpers stay in LocalMind. `root` is a new option (LocalMind passes
* 'localmind-ssd').
*
* RowFile fixed-size rows in one OPFS file plus manifest.json; synchronous reads through opfs-reader.js in
* inline mode, so it must live in a dedicated worker.
* RowCache `slots` rows on the GPU, split into planes (byte ranges of a row that live in separate GPU buffers),
* an O(1) LRU, and a CPU id→slot mirror (optionally a GPU u32 map). lookup(ids) makes every id
* resident and returns its slot, writing missing rows with queue.writeBuffer.
*/
import { OpfsReaderPool, OpfsWriter, canInline, readOpfsText, writeOpfsText, removeOpfs } from './opfs-reader.js';
export const ROOT = 'diskformer';
const FORMAT = 'localmind-rows/1'; // unchanged from LocalMind, so a store it wrote stays valid
// FNV-1a over a few samples of each source buffer and their lengths: cheap enough to run on every
// load, and it changes when the model's weights change.
export function fingerprint(buffers, sample = 1 << 16) {
let h = 0x811c9dc5;
const mix = (b) => { h ^= b; h = Math.imul(h, 0x01000193) >>> 0; };
for (const u8 of buffers) {
for (const v of [u8.length & 0xff, (u8.length >>> 8) & 0xff, (u8.length >>> 16) & 0xff, (u8.length >>> 24) & 0xff]) mix(v);
const starts = [0, Math.max(0, (u8.length >> 1) - (sample >> 1)), Math.max(0, u8.length - sample)];
for (const s of starts) for (let i = s, e = Math.min(u8.length, s + sample); i < e; i++) mix(u8[i]);
}
return h.toString(16).padStart(8, '0');
}
// Fixed-size rows in one OPFS file, <root>/<key>/<name>, with manifest.json beside it.
// The file I/O is opfs-reader.js in inline mode: the sync handle lives in this dedicated worker,
// so a row read is a plain synchronous call with no worker hop.
export class RowFile {
#reader = null;
constructor({ key, name, rowBytes, rows, root = ROOT }) {
this.path = `${root}/${key}/${name}`;
this.manifestPath = `${root}/${key}/manifest.json`;
this.name = name;
this.rowBytes = rowBytes;
this.rows = rows;
this.manifest = null;
}
static async open({ key, name = 'rows.bin', rowBytes, rows, root = ROOT }) {
if (!canInline()) throw new Error('RowFile needs a dedicated worker (FileSystemSyncAccessHandle)');
const f = new RowFile({ key, name, rowBytes, rows, root });
try { f.manifest = JSON.parse(await readOpfsText(f.manifestPath)); } catch (_) { f.manifest = null; }
return f;
}
matches(fp) {
const m = this.manifest;
return !!m && m.format === FORMAT && m.complete === true && m.name === this.name &&
m.rowBytes === this.rowBytes && m.rows === this.rows && m.fingerprint === fp;
}
// fill(row0, count, dst) writes `count` rows starting at `row0` into dst (count * rowBytes bytes).
// The manifest is removed first and written last, so an interrupted write is never taken as valid.
async write(fp, fill, { batchRows = 8192, extra = {} } = {}) {
this.close();
await removeOpfs(this.manifestPath).catch(() => {});
this.manifest = null;
const total = this.rows * this.rowBytes;
const t0 = performance.now();
const w = await OpfsWriter.open(this.path, { truncate: true, inline: true });
try {
const buf = new Uint8Array(batchRows * this.rowBytes);
for (let row0 = 0; row0 < this.rows; row0 += batchRows) {
const n = Math.min(batchRows, this.rows - row0);
const view = buf.subarray(0, n * this.rowBytes);
fill(row0, n, view);
await w.write(view, row0 * this.rowBytes);
}
} finally { await w.close(); }
const writeMs = performance.now() - t0;
this.manifest = { format: FORMAT, name: this.name, rowBytes: this.rowBytes, rows: this.rows, bytes: total,
fingerprint: fp, complete: true, writtenAt: new Date().toISOString(), writeMs: Math.round(writeMs), ...extra };
await writeOpfsText(this.manifestPath, JSON.stringify(this.manifest, null, 1));
return writeMs;
}
async openRead() {
if (this.#reader) return;
const r = await OpfsReaderPool.open(this.path, { inline: true });
if (r.size !== this.rows * this.rowBytes) {
await r.close();
throw new Error(`RowFile ${this.name}: size ${r.size} != ${this.rows * this.rowBytes}`);
}
this.#reader = r;
}
// Reads rows [row, row + count) to the start of dst (a Uint8Array that starts its buffer).
readRows(row, count, dst) {
if (dst.byteOffset !== 0) throw new Error('RowFile.readRows: dst must start at offset 0 of its buffer');
const len = count * this.rowBytes;
const got = this.#reader.readSync(row * this.rowBytes, len, dst.buffer);
if (got !== len) throw new Error(`RowFile ${this.name}: read ${got} of ${len} bytes at row ${row}`);
}
// Closes the read handle. In inline mode the handle closes synchronously inside this call.
close() {
if (this.#reader) { this.#reader.close(); this.#reader = null; }
}
}
export class RowCache {
// file: an open RowFile. slots: rows held on the GPU. planes: [{ offset, bytes, buffer }], each
// row byte range [offset, offset + bytes) is stored at buffer[slot * bytes]. queue: GPUQueue.
// mapBuffer (optional): a GPU u32[file.rows] copy of the id->slot map, 0xFFFFFFFF when absent.
constructor({ file, slots, planes, queue, mapBuffer = null }) {
this.file = file;
this.slots = slots;
this.planes = planes;
this.queue = queue;
this.mapBuffer = mapBuffer;
this.slotOf = new Int32Array(file.rows).fill(-1);
this.idOf = new Int32Array(slots).fill(-1);
// LRU as a doubly linked list over slots; head = least recently used. Starts as 0..slots-1.
this.prev = new Int32Array(slots);
this.next = new Int32Array(slots);
for (let s = 0; s < slots; s++) { this.prev[s] = s - 1; this.next[s] = s + 1 < slots ? s + 1 : -1; }
this.head = 0;
this.tail = slots - 1;
this.stamp = new Uint32Array(slots);
this.clock = 0;
this.row = new Uint8Array(file.rowBytes);
this.u32 = new Uint32Array(1);
this.stats = { lookups: 0, ids: 0, misses: 0, readMs: 0, replays: 0 };
this.missLog = []; // the first 4,096 ids installed after load, for analysis
}
#touch(s) {
if (s === this.tail) return;
const p = this.prev[s], n = this.next[s];
if (p >= 0) this.next[p] = n; else this.head = n;
this.prev[n] = p;
this.prev[s] = this.tail;
this.next[s] = -1;
this.next[this.tail] = s;
this.tail = s;
}
#mapWrite(id, slot) {
if (!this.mapBuffer) return;
this.u32[0] = slot >>> 0;
this.queue.writeBuffer(this.mapBuffer, id * 4, this.u32);
}
#install(id, s) {
const old = this.idOf[s];
if (old >= 0) { this.slotOf[old] = -1; this.#mapWrite(old, 0xffffffff); }
const t0 = performance.now();
this.file.readRows(id, 1, this.row);
this.stats.readMs += performance.now() - t0;
for (const p of this.planes) this.queue.writeBuffer(p.buffer, s * p.bytes, this.row, p.offset, p.bytes);
this.idOf[s] = id;
this.slotOf[id] = s;
this.#mapWrite(id, s);
this.stats.misses++;
if (this.missLog.length < 4096) this.missLog.push(id);
}
// Makes every id resident and returns its slot. Rows one call uses are never evicted by the same
// call, so a call may name at most `slots` distinct ids.
lookup(ids, out = new Uint32Array(ids.length)) {
const call = ++this.clock;
this.stats.lookups++;
this.stats.ids += ids.length;
for (let i = 0; i < ids.length; i++) {
const id = ids[i];
let s = this.slotOf[id];
if (s < 0) {
s = this.head;
if (this.stamp[s] === call) throw new Error(`RowCache: one lookup needs more than ${this.slots} rows`);
this.#install(id, s);
}
this.stamp[s] = call;
this.#touch(s);
out[i] = s;
}
return out;
}
has(id) { return this.slotOf[id] >= 0; }
// Loads rows [row0, row0 + count) into the least recently used slots in large reads (for a warm
// set at load), then uploads the whole id->slot map once. Returns the milliseconds spent.
warm(row0, count, batch = 4096) {
const t0 = performance.now();
count = Math.min(count, this.slots);
const rb = this.file.rowBytes;
const buf = new Uint8Array(batch * rb);
const planeBufs = this.planes.map((p) => new Uint8Array(batch * p.bytes));
for (let r = row0; r < row0 + count; r += batch) {
const n = Math.min(batch, row0 + count - r);
this.file.readRows(r, n, buf);
// Fresh slots come from the LRU head in order; collect them and write contiguous runs.
const slots = new Int32Array(n);
for (let i = 0; i < n; i++) {
const id = r + i;
let s = this.slotOf[id];
if (s < 0) {
s = this.head;
const old = this.idOf[s];
if (old >= 0) this.slotOf[old] = -1;
this.idOf[s] = id;
this.slotOf[id] = s;
}
this.#touch(s);
slots[i] = s;
}
for (let pi = 0; pi < this.planes.length; pi++) {
const p = this.planes[pi], dst = planeBufs[pi];
for (let i = 0; i < n; i++) dst.set(buf.subarray(i * rb + p.offset, i * rb + p.offset + p.bytes), i * p.bytes);
let i = 0;
while (i < n) {
let j = i + 1;
while (j < n && slots[j] === slots[j - 1] + 1) j++;
this.queue.writeBuffer(p.buffer, slots[i] * p.bytes, dst, i * p.bytes, (j - i) * p.bytes);
i = j;
}
}
}
// slotOf's bytes are the GPU map: -1 as an Int32 is 0xFFFFFFFF as a u32.
if (this.mapBuffer) this.queue.writeBuffer(this.mapBuffer, 0, this.slotOf);
return performance.now() - t0;
}
}
|