Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
210 changes: 192 additions & 18 deletions dag/gunbc/spark/patches/d2d649e6/engram_file_backed.py.patch
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
diff --git a/vllm/models/deepseek_v4_1/nvidia/engram_file_backed.py b/vllm/models/deepseek_v4_1/nvidia/engram_file_backed.py
new file mode 100644
index 0000000..b5c332c
index 0000000..4026664
--- /dev/null
+++ b/vllm/models/deepseek_v4_1/nvidia/engram_file_backed.py
@@ -0,0 +1,115 @@
@@ -0,0 +1,289 @@
+# File-backed Engram storage for vLLM d2d649e674c75425d2d6975c87eb89fd4d55fff8.
+#
+# THIS IS A BACKING SWAP, NOT NEW PREFETCH MACHINERY. Upstream
Expand All @@ -16,22 +16,34 @@ index 0000000..b5c332c
+#
+# This module keeps the staging buffer, the stream, and the prefetch protocol.
+# It changes where the rows come from: a rank-local file whose identity is
+# verified at open, mmap'd MAP_PRIVATE without registering the whole table
+# verified at open, mmap'd read-only without registering the whole table
+# with CUDA. Token-time gathers copy requested rows into the existing bounded
+# pinned staging buffer. A missing local path refuses; there is no NFS arm.
+#
+# Prefix and offsets are gunbc.spark.v41_engram_row_store: v41_store_bytes emits
+# v41_row_store_magic then weight||scale records; v41_row_address offsets follow;
+# the store's SHA-256 is over those prefixed bytes. Open checks expected.magic
+# (supplied from that symbol, not minted here) and whole-file digest against
+# rank_digest_hex. encoding_digest / rank / format_revision are not on this
+# rank_digest_hex. expected.record_bytes is v41_row_bytes of the admitted spec,
+# supplied likewise. encoding_digest / rank / format_revision are not on this
+# type until a modeled header carries them.
+#
+# The open and gather contracts are gunbc.spark.v41_engram_file_backed:
+# v41_open_residency_verdict (verification's page-cache residue is released and
+# read back as zero over the whole file, else open refuses; its two refusal
+# branches are the two raises in _verify_and_release, and the receipt's fields
+# are V41StoreOpenReceipt's) and v41_gather_rows (one batched gather
+# at record-aligned offsets, in request order; a misaligned offset refuses).
+
+from __future__ import annotations
+
+import ctypes
+import ctypes.util
+import hashlib
+import mmap
+import os
+import time
+import warnings
+from dataclasses import dataclass
+from typing import Tuple
+
Expand All @@ -50,18 +62,76 @@ index 0000000..b5c332c
+ pass
+
+
+class EngramPageCacheResident(RuntimeError):
+ pass
+
+
+class EngramRowOffsetMisaligned(RuntimeError):
+ pass
+
+
+@dataclass(frozen=True)
+class ExpectedRowStoreIdentity:
+ rank_digest_hex: str
+ magic: bytes
+ record_bytes: int
+
+
+@dataclass(frozen=True)
+class EngramStoreOpenReceipt:
+ """What open() observed: the verified range and its page-cache residency
+ before and after the release (gunbc.spark.v41_engram_file_backed
+ V41StoreOpenReceipt)."""
+
+ store_bytes: int
+ verified_bytes: int
+ resident_bytes_after_verify: int
+ resident_bytes_after_release: int
+
+def _sha256_file(path: str) -> str:
+
+_libc = ctypes.CDLL(ctypes.util.find_library("c"), use_errno=True)
+_libc.mmap.restype = ctypes.c_void_p
+_libc.mmap.argtypes = [ctypes.c_void_p, ctypes.c_size_t, ctypes.c_int, ctypes.c_int, ctypes.c_int, ctypes.c_long]
+_libc.munmap.argtypes = [ctypes.c_void_p, ctypes.c_size_t]
+_libc.mincore.argtypes = [ctypes.c_void_p, ctypes.c_size_t, ctypes.c_void_p]
+_MAP_FAILED = ctypes.c_void_p(-1).value
+_RESIDENCY_WINDOW = 1 << 30
+
+
+def resident_bytes(fd: int, size: int) -> int:
+ """Page-cache resident bytes of [0, size) of fd, by mincore over a
+ transient mapping that is never touched (so reading does not fault)."""
+ page = os.sysconf("SC_PAGE_SIZE")
+ resident_pages = 0
+ off = 0
+ while off < size:
+ length = min(_RESIDENCY_WINDOW, size - off)
+ addr = _libc.mmap(None, length, mmap.PROT_READ, mmap.MAP_SHARED, fd, off)
+ if addr in (None, _MAP_FAILED):
+ raise OSError(ctypes.get_errno(), "mmap for residency readback failed")
+ try:
+ pages = (length + page - 1) // page
+ vec = (ctypes.c_ubyte * pages)()
+ if _libc.mincore(addr, length, vec) != 0:
+ raise OSError(ctypes.get_errno(), "mincore failed")
+ resident_pages += sum(b & 1 for b in bytes(vec))
+ finally:
+ _libc.munmap(addr, length)
+ off += length
+ return min(resident_pages * page, size)
+
+
+def _sha256_fd(fd: int, magic_len: int) -> Tuple[str, bytes, int]:
+ h = hashlib.sha256()
+ with open(path, "rb") as f:
+ head = b""
+ total = 0
+ with os.fdopen(os.dup(fd), "rb", closefd=True) as f:
+ for chunk in iter(lambda: f.read(1024 * 1024), b""):
+ if len(head) < magic_len:
+ head += chunk[: magic_len - len(head)]
+ h.update(chunk)
+ return h.hexdigest()
+ total += len(chunk)
+ return h.hexdigest(), head, total
+
+
+class FileBackedEngramStorage:
Expand All @@ -73,38 +143,123 @@ index 0000000..b5c332c
+ f"rank-local Engram row store missing at {path}; "
+ "token-time NFS fallback is refused"
+ )
+ opened = _sha256_file(path)
+ if opened != expected.rank_digest_hex:
+ raise EngramRowStoreIdentityMismatch(
+ f"opened {path} digest {opened}, expected rank digest "
+ f"{expected.rank_digest_hex}"
+ )
+ self.path = os.path.abspath(path)
+ self.expected = expected
+ self._fd = os.open(self.path, os.O_RDONLY)
+ self._mmap = mmap.mmap(self._fd, 0, access=mmap.ACCESS_READ)
+ if self._mmap[: len(expected.magic)] != expected.magic:
+ raise EngramRowStoreIdentityMismatch("row-store MAGIC mismatch")
+ try:
+ self.receipt = self._verify_and_release()
+ self._mmap = mmap.mmap(self._fd, 0, access=mmap.ACCESS_READ)
+ except BaseException:
+ os.close(self._fd)
+ raise
+ with warnings.catch_warnings():
+ # The mmap is read-only; the table is only ever read through index.
+ warnings.simplefilter("ignore", UserWarning)
+ self._table = torch.frombuffer(self._mmap, dtype=torch.uint8)
+ self._span = torch.arange(expected.record_bytes, dtype=torch.int64)
+ self._header_ok = True
+
+ def _verify_and_release(self) -> EngramStoreOpenReceipt:
+ # Verification stays fail-closed and whole-file; only its cache residue goes.
+ store = os.fstat(self._fd).st_size
+ opened, head, verified = _sha256_fd(self._fd, len(self.expected.magic))
+ if verified != store:
+ raise EngramRowStoreIdentityMismatch(
+ f"verified {verified} of {store} store bytes; verification is over the whole file"
+ )
+ if opened != self.expected.rank_digest_hex:
+ raise EngramRowStoreIdentityMismatch(
+ f"opened {self.path} digest {opened}, expected rank digest "
+ f"{self.expected.rank_digest_hex}"
+ )
+ if head != self.expected.magic:
+ raise EngramRowStoreIdentityMismatch("row-store MAGIC mismatch")
+ after_verify = resident_bytes(self._fd, verified)
+ os.posix_fadvise(self._fd, 0, verified, os.POSIX_FADV_DONTNEED)
+ after_release = resident_bytes(self._fd, verified)
+ receipt = EngramStoreOpenReceipt(
+ store_bytes=store,
+ verified_bytes=verified,
+ resident_bytes_after_verify=after_verify,
+ resident_bytes_after_release=after_release,
+ )
+ if after_release != 0:
+ raise EngramPageCacheResident(
+ f"{self.path}: {after_release} of {verified} verified bytes remain "
+ f"page-cache resident after POSIX_FADV_DONTNEED; receipt {receipt}"
+ )
+ return receipt
+
+ def close(self) -> None:
+ del self._table
+ self._mmap.close()
+ os.close(self._fd)
+
+ def fetch_rows(self, offsets: list[int], nbytes: int, dest: torch.Tensor) -> None:
+ """Copy selected rows into an already-allocated bounded staging tensor."""
+ def _check_dest(self, dest: torch.Tensor) -> None:
+ if dest.is_pinned() and dest.numel() * dest.element_size() >= self._mmap.size():
+ raise EngramFullTablePinned(
+ "staging destination covers the whole table; this runtime refuses "
+ "full-Engram pinned allocation"
+ )
+
+ def _record_index(self, offsets: list[int], nbytes: int) -> torch.Tensor:
+ record = self.expected.record_bytes
+ if nbytes != record:
+ raise EngramRowOffsetMisaligned(
+ f"fetch of {nbytes} bytes per row; the store's record is {record} bytes"
+ )
+ magic = len(self.expected.magic)
+ size = self._mmap.size()
+ for off in offsets:
+ if off < magic or (off - magic) % record != 0 or off + record > size:
+ raise EngramRowOffsetMisaligned(
+ f"offset {off} is not magic + k * {record} inside a {size}-byte store"
+ )
+ return (torch.as_tensor(offsets, dtype=torch.int64).unsqueeze(1) + self._span).view(-1)
+
+ def fetch_rows(self, offsets: list[int], nbytes: int, dest: torch.Tensor) -> None:
+ """Copy selected rows, in request order, into an already-allocated
+ bounded staging tensor with one batched gather."""
+ self._check_dest(dest)
+ index = self._record_index(offsets, nbytes)
+ rows = dest.view(torch.uint8).view(dest.shape[0], -1)[: len(offsets), :nbytes]
+ if rows.is_contiguous():
+ torch.index_select(self._table, 0, index, out=rows.view(-1))
+ else:
+ rows.copy_(torch.index_select(self._table, 0, index).view(len(offsets), nbytes))
+
+ def _fetch_rows_per_row(self, offsets: list[int], nbytes: int, dest: torch.Tensor) -> None:
+ # The replaced per-row copy, kept ONLY as the differential oracle for
+ # measure_fetch_rows. It is never a fallback for fetch_rows.
+ self._check_dest(dest)
+ self._record_index(offsets, nbytes)
+ view = memoryview(self._mmap)
+ for i, off in enumerate(offsets):
+ dest[i].view(torch.uint8).view(-1)[:nbytes].copy_(
+ torch.frombuffer(view[off : off + nbytes], dtype=torch.uint8)
+ )
+
+
+def measure_fetch_rows(storage: FileBackedEngramStorage, rows: int, repeats: int, seed: int = 0) -> dict:
+ """THE INSTRUMENT: seconds per fetch_rows call of `rows` random records,
+ batched gather versus the per-row oracle, and whether their bytes agree."""
+ record = storage.expected.record_bytes
+ count = (storage._mmap.size() - len(storage.expected.magic)) // record
+ g = torch.Generator().manual_seed(seed)
+ local = torch.randint(0, count, (rows,), generator=g)
+ offsets = (len(storage.expected.magic) + local * record).tolist()
+ a = torch.zeros((rows, record), dtype=torch.uint8)
+ b = torch.zeros((rows, record), dtype=torch.uint8)
+ timings = {}
+ for name, fn, dest in (("batched", storage.fetch_rows, a), ("per_row", storage._fetch_rows_per_row, b)):
+ fn(offsets, record, dest)
+ start = time.perf_counter()
+ for _ in range(repeats):
+ fn(offsets, record, dest)
+ timings[name] = (time.perf_counter() - start) / repeats
+ return {"rows": rows, "repeats": repeats, "seconds_per_call": timings, "bytes_agree": bool(torch.equal(a, b))}
+
+
+def allocate_file_backed(
+ path: str,
+ expected: ExpectedRowStoreIdentity,
Expand All @@ -119,3 +274,22 @@ index 0000000..b5c332c
+ "pinned staging is not smaller than the row store; refusing full-table pin"
+ )
+ return storage, staging
+
+
+if __name__ == "__main__":
+ import argparse
+ import json
+
+ p = argparse.ArgumentParser(description="measure fetch_rows over a rank-local Engram row store")
+ p.add_argument("--store", required=True)
+ p.add_argument("--rank-digest-hex", required=True)
+ p.add_argument("--magic-hex", required=True)
+ p.add_argument("--record-bytes", type=int, required=True)
+ p.add_argument("--rows", type=int, required=True)
+ p.add_argument("--repeats", type=int, default=100)
+ a = p.parse_args()
+ s = FileBackedEngramStorage(a.store, ExpectedRowStoreIdentity(a.rank_digest_hex, bytes.fromhex(a.magic_hex), a.record_bytes))
+ try:
+ print(json.dumps({"open_receipt": s.receipt.__dict__, "fetch_rows": measure_fetch_rows(s, a.rows, a.repeats)}))
+ finally:
+ s.close()
92 changes: 92 additions & 0 deletions dag/gunbc/spark/v41_engram_file_backed.dag
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
module gunbc.spark.v41_engram_file_backed

import std.types { String, Bool, List, NonEmptyStr }
import std.nat { Nat }
import std.integer { UInt8 }
import std.measure { ByteSize, byte_size_count }
import v2.std.optional { Present, Absent }
import v2.std.algebra { filter, length, skip }
import gunbc.spark.v41_engram_row_store {
V41RowStoreEncodingSpec, V41CompleteStore,
v41_row_bytes, v41_row_store_magic_size,
}

// ── THE FILE-BACKED OVERLAY'S OPEN AND GATHER CONTRACTS ─────────────────────────────────────────
//
// dag/gunbc/spark/patches/d2d649e6/engram_file_backed.py.patch FileBackedEngramStorage realizes
// what this module declares; the patch is the realization and these are its facts.
//
// OPEN. Whole-file SHA-256 verification stays, and it stays fail-closed. What goes is its residue:
// on a DGX Spark the CPU and GPU share one LPDDR5x pool, so the page cache the sequential verify
// read leaves behind is ordinary memory taken from weights and KV. After verifying, open releases
// the verified range (POSIX_FADV_DONTNEED) and reads its residency back (mincore). The receipt
// carries both readings, and open is admitted only when the readback is zero over the whole file.
//
// GATHER. fetch_rows is one batched gather: every requested offset must be a record boundary of
// this store (magic, then k whole records, with the record inside the file), and the rows come back
// in request order, each exactly the record v41_place_row put at that address. A shifted offset
// would read the tail of one record and the head of the next, a plausible row of the wrong bytes,
// so it refuses rather than being read.
//
// THE REALIZATION MIRRORS THIS, BRANCH FOR BRANCH. FileBackedEngramStorage._verify_and_release raises
// on exactly the two refusals of v41_open_residency_verdict (verified != fstat size; any residue after
// release), and EngramStoreOpenReceipt carries V41StoreOpenReceipt's four fields. _record_index raises
// on exactly the non-record offsets v41_is_record_offset rejects. The Python refuses at runtime because
// a Spark rank cannot call this fold; the fold is the declaration those raises realize.
//
// CONSUMER: A DECLARED FRONTIER, NOT YET IN THIS CLOSURE. The consumer of v41_open_residency_verdict is
// the (layer, rank) open-receipt reading the wiring patch introduces: when
// ParallelEngramEmbedding._allocate_weights constructs FileBackedEngramStorage (the trigger
// gunbc.spark.v41_engram_row_store already names), each rank's EngramStoreOpenReceipt is read back into
// gunbc.spark.v41_runtime_candidate beside v41_published_store_readback_readings and folded through this
// verdict, so a rank that verified but stayed resident refuses the candidate, not only its own open.
// v41_gather_rows' consumer is the same change's differential over boundary rows of a published store.
// Until then both are exercised by test.claim.spark.v41_engram_file_backed_witness over supplied
// receipts and fixture stores. The patch's measure_fetch_rows / __main__ is the named instrument for
// fetch_rows time; its readings are not transcribed here.

type V41StoreOpenReceipt sole_constructor {
store_bytes: ByteSize
verified_bytes: ByteSize
resident_after_verify: ByteSize
resident_after_release: ByteSize
}

fn v41_store_open_receipt(store_bytes: ByteSize, verified_bytes: ByteSize, resident_after_verify: ByteSize, resident_after_release: ByteSize) -> V41StoreOpenReceipt {
V41StoreOpenReceipt { store_bytes: store_bytes, verified_bytes: verified_bytes, resident_after_verify: resident_after_verify, resident_after_release: resident_after_release }
}

type V41OpenResidency
= V41OpenReleased { receipt: V41StoreOpenReceipt }
| V41OpenResidueRefused { defect: NonEmptyStr }

fn v41_open_residency_verdict(r: V41StoreOpenReceipt) -> V41OpenResidency {
if byte_size_count(b: r.verified_bytes) != byte_size_count(b: r.store_bytes) {
V41OpenResidueRefused { defect: join(["verified ", to_string(byte_size_count(b: r.verified_bytes)), " of ", to_string(byte_size_count(b: r.store_bytes)), " store bytes; verification is over the whole file"], "") as NonEmptyStr }
} else if byte_size_count(b: r.resident_after_release) != 0 {
V41OpenResidueRefused { defect: join([to_string(byte_size_count(b: r.resident_after_release)), " verified bytes remain page-cache resident after release (", to_string(byte_size_count(b: r.resident_after_verify)), " after verify)"], "") as NonEmptyStr }
} else {
V41OpenReleased { receipt: r }
}
}

type V41Gather
= V41Gathered { rows: List<List<UInt8>> }
| V41GatherRefused { defect: NonEmptyStr }

fn v41_is_record_offset(spec: V41RowStoreEncodingSpec, store_len: Nat, offset: ByteSize) -> Bool {
let magic = byte_size_count(b: v41_row_store_magic_size())
let record = byte_size_count(b: v41_row_bytes(g: spec.geometry))
let o = byte_size_count(b: offset)
o >= magic && (o - magic) % record == 0 && o + record <= store_len
}

fn v41_gather_rows(store: V41CompleteStore, spec: V41RowStoreEncodingSpec, offsets: List<ByteSize>) -> V41Gather {
let record = byte_size_count(b: v41_row_bytes(g: spec.geometry))
let n = length(store.bytes)
let misaligned = filter(offsets, o => !v41_is_record_offset(spec: spec, store_len: n, offset: o))
match first(misaligned) {
Present { value: o } => V41GatherRefused { defect: join(["offset ", to_string(byte_size_count(b: o)), " is not magic + k * ", to_string(record), " inside a ", to_string(n), "-byte store"], "") as NonEmptyStr }
Absent => V41Gathered { rows: map(offsets, o => store.bytes |> skip(n: byte_size_count(b: o)) |> take(n: record)) }
}
}
Loading