From 47924565c1907bf3c34fe680ac1c8fef145e447d Mon Sep 17 00:00:00 2001 From: Raffaele <74499579+Helldez@users.noreply.github.com> Date: Mon, 7 Sep 2026 20:43:02 +0200 Subject: [PATCH] feat(io): release the model file's mapping after load (--release-mmap) (#185) llama.cpp maps every gguf it loads and keeps the mapping for the model's lifetime. On Windows that is expensive in a way nothing had attributed: while a section of a file is alive, NTFS serialises concurrent unbuffered reads on that file, and a lane opened while the section existed keeps serialising against it after the section is gone. Four I/O lanes therefore delivered exactly one lane's throughput, which is why lanes and threads have always measured dead on the desktop host and why the engine read at about a third of what the drive can serve. --release-mmap hands the mapping back after load: unmap the file, close its section, reopen the reader lanes. Both halves are needed. Whether it is safe is decided by looking rather than by reasoning, the engine asks the OS whether any weight the capture pass observed still points inside a mapping of the model files, and declines if any does. Off by default, because the check answers for the pointers the capture saw and for no others. Host A/B on Qwen3.6-35B-A3B Q4_K_M: 3.16 to 4.63 tok/s (+46%), flash stall per token 0.182 to 0.074, with bytes read, hit rate, evictions and re-reads identical to the digit and the generated text byte-identical. On the phone the read path is flat (f2fs does not serialise) but CPU per token falls about 9%; that cell is two short runs per variant and is recorded as a direction, not a number. Also fixes a bug this uncovered, independent of the flag: when a gguf carries no output.weight, llama.cpp builds the output head from the token embedding table and the model holds two identically named tensors over the same bytes. The capture pass keyed its map by name, so --dense-weights anon and ahwb rebound one and left the twin reading the mmap for the whole run, which on a model past RAM means the output projection served by page faults from flash. The capture now records every distinct leaf object by address and the dense policy rebinds every tensor over one file range onto the same buffer. Adds bmoe-iobench --mmap / --reopen-lanes / --range-mb / --fresh, the cells that isolate the mechanism, a mapping_release unit test on both platforms, the app switch "Release the model mapping", and the bench findings. README, architecture, AGENTS and roadmap updated, the last correcting a diagnosis this refutes. --- AGENTS.md | 2 +- CHANGELOG.md | 71 ++++ README.md | 2 + cli/main.cpp | 5 + core/CMakeLists.txt | 1 + core/include/bmoe/config.h | 17 + core/src/engine/session.cpp | 91 ++++- core/src/io/file_reader.cpp | 16 + core/src/io/file_reader.h | 14 +- core/src/io/mapping_release.cpp | 367 ++++++++++++++++++ core/src/io/mapping_release.h | 76 ++++ core/src/moe/dense_weights.cpp | 11 + core/src/moe/dense_weights.h | 16 + core/src/moe/expert_stream_source.cpp | 12 + core/src/moe/expert_stream_source.h | 5 + core/src/moe/router_hook.cpp | 3 + core/src/moe/router_hook.h | 10 + core/src/moe/row_stream.cpp | 7 + core/src/moe/row_stream.h | 2 + docs/architecture.md | 1 + .../2026-08-29-mmap-serialisation/findings.md | 168 ++++++++ docs/moe-streaming.md | 39 ++ docs/roadmap.md | 12 +- examples/android/README.md | 7 + .../io/bigmoeonedge/example/AppSettings.kt | 20 + .../io/bigmoeonedge/example/SettingsScreen.kt | 11 + tests/CMakeLists.txt | 7 + tests/mapping_release_test.cpp | 183 +++++++++ tools/iobench.cpp | 145 ++++++- 29 files changed, 1303 insertions(+), 18 deletions(-) create mode 100644 core/src/io/mapping_release.cpp create mode 100644 core/src/io/mapping_release.h create mode 100644 docs/bench-data/2026-08-29-mmap-serialisation/findings.md create mode 100644 tests/mapping_release_test.cpp diff --git a/AGENTS.md b/AGENTS.md index 4a7da4a..5b831dc 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -15,7 +15,7 @@ we do not fork llama.cpp. See `docs/architecture.md` and `docs/seam.md`. - `core/include/bmoe/` — ports (interfaces) + config. Pure policy, no llama.cpp include. - `core/src/io/` — `platform_io` (cross-platform O_DIRECT reads + reserve/commit/evict VM); `file_reader` (pooled positioned reader, per-consumer O_DIRECT — used by both the expert stream - and the dense loader). + and the dense loader); `mapping_release` (hands the model file's mapping back after load). - `core/src/moe/` — `gguf_offsets`, `arch_registry`, `expert_stream_source`, `router_hook`; `dense_weights` (the non-expert weight policy: mmap / warm / anon, plus the residency sensor). - `core/src/engine/runtime.cpp` — composition + greedy generation loop. diff --git a/CHANGELOG.md b/CHANGELOG.md index 341bc4f..f8c9f20 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,7 +6,78 @@ Semantic Versioning. ## [0.24.0] - unreleased +### Fixed +- **A tied output head left half the biggest dense weight reading the model's mmap.** When a gguf + carries no `output.weight`, llama.cpp builds the output head from the token embedding table, and + the model then holds **two** `ggml_tensor` objects, identically named, over the same file bytes. + The capture pass recorded weights in a map keyed by name, so it kept one of them, and + `--dense-weights anon` / `ahwb` rebound that one: the twin went on reading the mmap for the whole + run. On a model past RAM that is the output projection — a tensor whose per-token cost rivals all + the routed experts together — served by page faults from flash, which is precisely what those two + policies exist to prevent. Gemma is in this case, and so is any architecture that ties its + embeddings. + + The capture now records every distinct leaf object, deduplicated by address, and the dense policy + folds tensors over one file range into a single entry with an alias list, rebinding them all onto + the same buffer. Bytes are still read once and memory is still allocated once — on the gemma4 + gate, the same 70 anon buffers as before — and every byte-identity gate passes unchanged. The + same fix is what makes the mapping releasable on those models rather than a crash. + + What this is worth in throughput is **not measured**. The mechanism is certain, the size of it is + not: it depends on how much of the mapped table the kernel was still holding, which on a phone + under reclaim is a different story from a desktop with room to spare. Gemma-4-26B on the test + phone is the cell that would price it, and it is owed. + ### Added +- **`--release-mmap`: hand back the model file's mapping after load, which on Windows is worth + +46% decode.** llama.cpp maps every gguf it loads and keeps the mapping for the model's + lifetime. That is load-bearing for the streamer, which rebinds expert tensors onto the file's + native layout, but it is not free: on Windows, while a section of a file is alive, NTFS + serialises concurrent unbuffered reads on that file, and a lane opened while the section + existed keeps serialising against it even after the section is gone. Four I/O lanes therefore + deliver exactly one lane's throughput, which is why lanes and threads have always measured dead + on the desktop host and why the engine read at roughly a third of what the drive can serve. + + Measured with `bmoe-iobench` at the engine's own request shape (576 KiB, 4 lanes, O_DIRECT) on + an NVMe host: **2400-2660 MiB/s with no mapping, 895-930 MiB/s with one held open, 2400 again + once the mapping is dropped and the lanes are reopened.** Confining the offsets (`--range-mb`), + committing the destination afresh (`--fresh`) and running the lanes against a busy CPU + (`--compute-load`) all move the number by a few per cent; the mapping moves it by 2.6x. The + same three cells on the 12 GB Android phone are flat: f2fs does not serialise, so there is no + read-path gain to collect there. The flag is implemented on POSIX anyway (munmap of the file's + VMAs, inert placeholders left in their ranges), and the phone then turned out to gain for a + different reason — see below. + + With the flag, once nothing reads through the mapping any more the engine unmaps the file, + closes its section and reopens the reader lanes. Whether that is safe is decided by looking rather + than by reasoning: the engine asks the OS whether any weight the capture pass observed still points + inside a mapping of the model files, and declines if any does. One check covers every residency + policy at once — a dense set deliberately left mmap'd, a table held back as oversized, or a tensor + no name-based accounting could have found. + + It stays opt-in because the check answers for the pointers the capture pass saw and for no others: + a graph shape this session never built could hold another one. That residue, plus the gguf tensors + no policy owns once the MTP draft adds a second graph, is the remaining unknown. llama.cpp is not patched and + still believes it owns its mapping, so what it will unmap and close at teardown is left in place + as an inert placeholder. Off by default. + + Host A/B, Qwen3.6-35B-A3B Q4_K_M, cache 3000 MiB, 4 lanes, overlap, dense anon, 96 tokens, + interleaved: **3.16 to 4.63 tok/s (+46%)**, flash stall per token **0.182 to 0.074 s**, per-read + latency p50 **2.49 to 1.20 ms**, with bytes read, cache hit rate, evictions and re-reads + identical to the digit and the generated text byte-identical. + + On the phone the flag pays too, by a different mechanism and by less. Flash stall is unchanged in + every cell (0.087-0.093 s/token with the flag and without), exactly as the microbenchmark + predicts, but CPU time per token falls about 9% and decode comes out 5-9% ahead: holding a 20 GB + mapping registered is not free for a kernel already under memory pressure, and handing it back + removes that. Prefill, model load and TTFT are unaffected. Preliminary — two 48-token cells per + variant in both orders, against a 20% cell-to-cell spread on this device — so it is recorded as a + direction, not a number, pending a long run. + + The demo app exposes it as **"Release the model mapping"**, enabled only under a dense policy that + rebinds every weight into the app's own memory (Anon or Pinned) — under Mmap or Warm the engine + looks at its own pointers and declines, so the switch would be one that silently does nothing. Off + by default, for the same reason the flag is: the phone's number is a direction. - **`--row-stream`: dense tables the graph only gathers rows from, served from flash.** The dense policy has one shape for every non-expert weight, and it is the right shape for a weight that is multiplied: read it whole, keep it resident. A token embedding table is not that. The graph diff --git a/README.md b/README.md index c2e6a81..d57155a 100644 --- a/README.md +++ b/README.md @@ -135,6 +135,8 @@ flash, at the moment they are needed. Everything below tunes that. | Direct I/O | `--no-odirect` disables it  (on by default) | Bypasses the OS page cache, so the system holds no second copy of what the expert cache already has. Falls back where unsupported. | | I/O and compute overlap | `--overlap`  (off in the CLI, on in the app) | Issues the next reads while the current layer computes, hiding flash latency behind work. Byte-identical; needs a small optional add-on to llama.cpp ([seam](docs/seam.md)). | | Dense weights | `--dense-weights`  `mmap`, `warm`, `anon`, `ahwb`  (default `anon`) | How the always-needed non-expert weights are held. Decisive far past RAM. `ahwb` is Android-only and puts them where the kernel cannot reclaim them at all ([data](docs/bench-data/2026-07-21-pinned-dense-ab/findings.md)). | +| Row-gathered tables | `--row-stream`, with `--row-stream-mb`  (default `64`) | Serves a dense table the graph only gathers rows from, typically the token embedding, out of flash instead of RAM: only the rows about to be read are pulled in, inside a bounded window. Which tables qualify is read off the model's own graph. Lossless ([detail](docs/row-gathered-tables.md)). | +| Release the model mapping | `--release-mmap`  (off by default) | Unmaps the gguf once every weight has been rebound into the engine's own memory, and reopens the read lanes. On Windows a live mapping serialises the streamer's unbuffered reads, so this is worth +46% decode there; on Android it saves CPU instead. Needs `--dense-weights anon` or `ahwb`, and the engine declines if any weight still points into the mapping. Lossless ([data](docs/bench-data/2026-08-29-mmap-serialisation/findings.md)). | | Temporal prefetch *(experimental)* | `--prefetch`  `0` (off), `1`, `2`, `4` layers | Bets a layer reuses the previous token's experts and fetches them on idle lanes. Needs the cache. | | Predictive prefetch *(experimental)* | `--predict-prefetch`, with `--predict-spec-max`  `0` (retention only), `1`, `2`, `4` | Runs the next layer's own router early and fetches what it names. More accurate than the bet above; reading ahead on it still lost its on-device A/B ([why](docs/expert-prediction.md)). | diff --git a/cli/main.cpp b/cli/main.cpp index ed1f4ff..ea3c2b7 100644 --- a/cli/main.cpp +++ b/cli/main.cpp @@ -462,6 +462,9 @@ static void print_usage(const char * argv0) { " name list, so a model that also multiplies by its embedding\n" " table is left alone, on any architecture\n" " --row-stream-mb N resident window for those tables in MiB (default 64)\n" + " --release-mmap unmap the model file after load once nothing reads through it\n" + " (Windows: a live mapping serialises the streamer's concurrent reads;\n" + " needs --dense-weights anon|ahwb; measured neutral on Android)\n" " --load-all debug: read ALL experts each token (A/B baseline)\n" " --force-cache allow a cache-mb in the pathological band\n" " --overlap overlap async expert reads with FFN compute (needs the fork)\n" @@ -682,6 +685,8 @@ int main(int argc, char ** argv) { cfg.moe.io_threads = std::atoi(next("--io-threads")); else if (a == "--no-odirect") cfg.moe.o_direct = false; + else if (a == "--release-mmap") + cfg.moe.release_mmap = true; else if (a == "--row-stream") cfg.moe.row_stream = true; else if (a == "--row-stream-mb") diff --git a/core/CMakeLists.txt b/core/CMakeLists.txt index f370cb9..8d37f72 100644 --- a/core/CMakeLists.txt +++ b/core/CMakeLists.txt @@ -23,6 +23,7 @@ if(BMOE_HAVE_LLAMA) target_sources(bmoe_core PRIVATE src/io/platform_io.cpp src/io/file_reader.cpp + src/io/mapping_release.cpp src/moe/dense_weights.cpp src/moe/row_stream.cpp src/moe/gguf_offsets.cpp diff --git a/core/include/bmoe/config.h b/core/include/bmoe/config.h index 913e4f2..fbed9b6 100644 --- a/core/include/bmoe/config.h +++ b/core/include/bmoe/config.h @@ -120,6 +120,23 @@ struct MoeStreamConfig { bool row_stream = false; int row_stream_mb = 64; // resident window across all row-streamed tables, in MiB + // ── release the model file's mapping after load (Windows) ───────────────────────── + // llama.cpp keeps the gguf mapped for the model's lifetime, and on Windows a live section of + // the file serialises the streamer's concurrent unbuffered reads: four lanes deliver one + // lane's throughput (measured; see core/src/io/mapping_release.h). With this on, once nothing + // reads through the mapping any more (dense weights copied out under anon/ahwb, every file + // tensor either streamed, copied or row-served) the engine unmaps the file and closes its + // section. On POSIX the mapping is munmap'd the same way; Android was measured not to serialise + // reads, so there the gain is smaller and of a different kind (see docs/moe-streaming.md). + // + // Opt-in, and it must stay opt-in until the safety question has a positive answer rather than a + // blacklist. The release is correct only while nothing still dereferences the mapping, and + // llama.cpp exposes no way to enumerate a loaded model's tensors and prove that. The session + // declines on every shape where a surviving pointer is KNOWN — a dense policy that keeps the + // weights mmap'd, a tied output head, an unowned tensor under the MTP draft — but a future + // architecture can invent another one. See Session::open. + bool release_mmap = false; + // ── cache-aware expert dropping (lossy; opt-in) ────────────────────────────────── // Skip a routed expert when it is a cache MISS *and* the router weighted it below // drop_cold_frac × (1 / n_expert_used) — i.e. below that fraction of the uniform share a diff --git a/core/src/engine/session.cpp b/core/src/engine/session.cpp index 04c4b6f..f484ca9 100644 --- a/core/src/engine/session.cpp +++ b/core/src/engine/session.cpp @@ -10,6 +10,7 @@ #include "../moe/expert_stream_source.h" #include "../moe/gguf_offsets.h" #include "../io/platform_io.h" +#include "../io/mapping_release.h" #include "llama.h" #include "ggml.h" @@ -321,6 +322,9 @@ struct Session::Impl { // (the fail-open default) means the flag alone does the job and generate() adds nothing. ThinkControl think_ctl = ThinkControl::Template; bool backend_inited = false; + // What release_file_mappings left behind on purpose; released only after `model` is gone + // (its teardown still unmaps/closes what it believes is its mapping; see mapping_release.h). + pio::MappingPlaceholders mapping_placeholders; // Sampling chain, built once at open() only when sampling is requested (temp > 0); null on the // greedy default, where the decode loop stays on the argmax fast path. See open()/generate(). @@ -371,6 +375,7 @@ struct Session::Impl { ctx.reset(); hook.reset(); model.reset(); + mapping_placeholders.release(); if (backend_inited) llama_backend_free(); } }; @@ -862,6 +867,7 @@ std::unique_ptr Session::open(const SessionConfig & cfg, // their own buffers; every mode needs the list to hold back a tensor too large to be // resident at all (qwen4exp's n-gram table), which must also leave the warm sweep and the // residency sensor — see DenseWeights::hold_back_oversized. + std::vector unaccounted; // file tensors no policy owns; decides release-mmap below { const std::unordered_set expert_names = expert_tensor_names(layers); // The subset the graph only row-gathers, when the run asked for the row policy. Their @@ -869,21 +875,55 @@ std::unique_ptr Session::open(const SessionConfig & cfg, // dense list rather than a second list, so a table cannot be both resident and streamed. const std::unordered_set row_names = cfg.moe.row_stream ? im.hook->row_gathered_weights() : std::unordered_set(); + // Walk the captured OBJECTS, not the name map: a tied output head gives two tensors + // one name over one file range, and a name-keyed walk would rebind only one of them + // and leave the other reading the mmap for the life of the run. Twins are folded into + // one entry's alias list so the bytes are still read and allocated exactly once. std::vector dense, rows; - for (const auto & kv : im.hook->captured_weights()) { - const std::string & name = kv.first; + std::unordered_map dense_at; // (shard, file offset) -> index in `dense` + for (ggml_tensor * t : im.hook->captured_weight_objects()) { + if (!t) continue; + const std::string name = t->name; if (expert_names.count(name)) continue; auto off = offs.off_by_name.find(name); auto sz = offs.size_by_name.find(name); if (off == offs.off_by_name.end() || sz == offs.size_by_name.end()) continue; // not a file tensor + const int file_idx = offs.file_by_name.at(name); + const uint64_t key = ((uint64_t) (uint32_t) file_idx << 48) ^ off->second; + auto seen = dense_at.find(key); + if (seen != dense_at.end()) { + dense[seen->second].aliases.push_back(t); + continue; + } DenseTensorRef d; - d.tensor = kv.second; + d.tensor = t; d.file_off = off->second; d.size = sz->second; - d.file_idx = offs.file_by_name.at(name); + d.file_idx = file_idx; + dense_at.emplace(key, dense.size()); dense.push_back(d); + // A table with a twin is by definition also read some other way, so it cannot have + // qualified as row-gathered; the guard costs nothing and keeps that coupling explicit. if (row_names.count(name)) rows.push_back(d); } + for (const DenseTensorRef & d : dense) + if (!d.aliases.empty()) + std::fprintf(stderr, "bmoe: '%s' is bound %zu times over one range (tied head) — all rebound\n", + d.tensor->name, d.aliases.size() + 1); + // Every gguf tensor that is neither streamed nor in the dense list is one the capture + // decode never touched. With speculation off that is a tensor llama.cpp did not even load + // (load_mtp is off, so the MTP block stays in the file): nothing in this session can + // reach it. With the MTP draft on, a second context builds a second graph, and a tensor + // the capture missed may be one that graph reads, so the list then blocks release-mmap. + // Only release-mmap reads this list, so only a run that asked for it pays to build it. + if (cfg.moe.release_mmap) { + std::unordered_set owned; + owned.reserve(dense.size()); + for (const DenseTensorRef & d : dense) + if (d.tensor) owned.insert(d.tensor->name); + for (const auto & kv : offs.off_by_name) + if (!expert_names.count(kv.first) && !owned.count(kv.first)) unaccounted.push_back(kv.first); + } const uint64_t row_budget = (uint64_t) std::max(0, cfg.moe.row_stream_mb) * 1024ull * 1024ull; im.source.set_dense_tensors(std::move(dense)); im.source.set_row_tensors(std::move(rows), row_budget); @@ -896,6 +936,49 @@ std::unique_ptr Session::open(const SessionConfig & cfg, // qualified or the takeover declined, and the hook then costs exactly nothing per node. im.hook->set_row_source(im.source.row_source()); + // Drop the model file's mapping once nothing reads through it (core/src/io/mapping_release.h). + // Safety is decided here, not in the module, and it is decided by looking rather than by + // reasoning: the release is correct only while nothing still dereferences the mapping, so the + // engine asks the OS whether any weight the graph read still points inside it. That covers + // every residency policy at once — a dense set left mmap'd, a table held back as oversized, + // and the twin a tied output head creates, which carries the same name as the embedding table + // and which no name-based accounting can see. + // + // It is a check over the pointers the capture pass observed, not a proof about the ones it + // did not: a graph shape this session never built could hold another. That residue, plus the + // gguf tensors no policy owns once the MTP draft adds a second graph, is why the flag is + // opt-in rather than on by default. + std::vector weight_addrs; + if (cfg.moe.release_mmap) { + weight_addrs.reserve(im.hook->captured_weight_objects().size()); + for (const ggml_tensor * t : im.hook->captured_weight_objects()) + if (t && t->data) weight_addrs.push_back(t->data); + } + const size_t still_mapped = + cfg.moe.release_mmap ? pio::addresses_in_file_mappings(offs.shard_paths, weight_addrs) : 0; + if (cfg.moe.release_mmap) { + if (still_mapped) { + std::fprintf(stderr, "bmoe: release-mmap skipped: %zu weight(s) still read the model's mapping\n", + still_mapped); + } else if (!unaccounted.empty() && cfg.spec.is_mtp()) { + std::fprintf(stderr, "bmoe: release-mmap skipped: %zu file tensor(s) no policy owns (first: %s)\n", + unaccounted.size(), unaccounted.front().c_str()); + } else { + const pio::MappingReleaseReport r = + pio::release_file_mappings(offs.shard_paths, &im.mapping_placeholders); + // Lanes opened while the section was alive stay serialised against it even once it + // is gone (measured: iobench S3 vs S4), so the readers are reopened after the release. + if (r.supported && (r.views_unmapped || r.sections_closed) && !im.source.reopen_readers()) + std::fprintf(stderr, "bmoe: release-mmap: reader reopen failed; reads stay serialised\n"); + if (r.supported) // POSIX has nothing to undo and says nothing + std::fprintf(stderr, + "bmoe: release-mmap: %d view(s) unmapped (%llu MiB), %d section(s) closed%s%s%s\n", + r.views_unmapped, (unsigned long long) (r.bytes >> 20), r.sections_closed, + r.plugs_missed ? ", handle slot not reclaimed" : "", r.error.empty() ? "" : "; ", + r.error.c_str()); + } + } + if (route_trace) { im.route_trace = route_trace; im.hook->set_trace(true); diff --git a/core/src/io/file_reader.cpp b/core/src/io/file_reader.cpp index a8b541d..c86f12e 100644 --- a/core/src/io/file_reader.cpp +++ b/core/src/io/file_reader.cpp @@ -15,6 +15,11 @@ FileReader::~FileReader() { bool FileReader::open(const std::string & path, int lanes, bool direct, size_t align, size_t bounce_cap) { if (is_open()) return false; + // Recorded after the guard: a refused call must not overwrite what reopen() replays. + path_ = path; + lanes_ = lanes; + direct_req_ = direct; + bounce_cap_ = bounce_cap; align_ = align ? align : 4096; direct_ = direct; const int n = lanes < 1 ? 1 : lanes; @@ -179,6 +184,17 @@ long long FileReader::read(int lane, void * dst, uint64_t off, uint64_t nbytes) return window; // the aligned window pulled — what the effective bandwidth is judged against } +bool FileReader::reopen() { + if (!is_open()) return false; + const std::string path = path_; + const int lanes = lanes_; + const bool direct = direct_req_; + const size_t align = align_; + const size_t cap = bounce_cap_; + close(); + return open(path, lanes, direct, align, cap); +} + void FileReader::close() { for (pio::fd_t fd : fds_) if (pio::fd_ok(fd)) pio::close_fd(fd); diff --git a/core/src/io/file_reader.h b/core/src/io/file_reader.h index 842c6fb..712eb07 100644 --- a/core/src/io/file_reader.h +++ b/core/src/io/file_reader.h @@ -40,9 +40,15 @@ public: // bounce per lane. `direct` requests cache bypass; it is verified and silently downgraded where // the platform or the storage refuses or mis-serves it — direct() then reports the effective // mode. Reads align to `align` where the platform's direct mode demands it. Returns false on any - // open/alloc failure. A reader is opened once and not reused. + // open/alloc failure. bool open(const std::string & path, int lanes, bool direct, size_t align, size_t bounce_cap); void close(); + // Close every fd and open the same file again with the same lanes, alignment and bounce size. + // Exists for one reason: on Windows an unbuffered handle opened while the file had a live + // section keeps serialising against it even after the section is gone, so once the model's + // mapping is released the lanes must be reopened to read at the drive's rate (see + // mapping_release.h). Not thread-safe: no read may be in flight on any lane. Accounting is kept. + bool reopen(); bool is_open() const { return !fds_.empty(); } bool direct() const { return direct_; } // the effective cache-bypass mode, not the request @@ -72,6 +78,12 @@ private: // after the open's own report and every fallback bool aligned_reads_ = false; // direct_ AND the platform's direct mode rejects unaligned reads — // gates the read mechanics: window rounding, bounce, buffered tail fd + // What open() was asked for, so reopen() can ask for the same thing again. direct_ above is the + // outcome, which a reopen must be free to reach on its own rather than inherit. + std::string path_; + int lanes_ = 0; + bool direct_req_ = false; + size_t bounce_cap_ = 0; std::atomic read_bytes_{0}; std::atomic syscall_ns_{0}; }; diff --git a/core/src/io/mapping_release.cpp b/core/src/io/mapping_release.cpp new file mode 100644 index 0000000..049f9d9 --- /dev/null +++ b/core/src/io/mapping_release.cpp @@ -0,0 +1,367 @@ +#include "mapping_release.h" + +// System headers at global scope, as in platform_io.cpp. +#if defined(_WIN32) +#ifndef NOMINMAX +#define NOMINMAX // windows.h's max() macro would shadow std::max below +#endif +#include +#include +#include +#include +#include +#else +#include +#include "platform_io.h" +#endif + +namespace bmoe::pio { + +MappingPlaceholders::~MappingPlaceholders() { + release(); +} + +MappingPlaceholders::MappingPlaceholders(MappingPlaceholders && o) noexcept + : reserved_(std::move(o.reserved_)), plugs_(std::move(o.plugs_)) { + o.reserved_.clear(); + o.plugs_.clear(); +} + +MappingPlaceholders & MappingPlaceholders::operator=(MappingPlaceholders && o) noexcept { + if (this != &o) { + release(); + reserved_ = std::move(o.reserved_); + plugs_ = std::move(o.plugs_); + o.reserved_.clear(); + o.plugs_.clear(); + } + return *this; +} + +void MappingPlaceholders::release() { +#if defined(_WIN32) + for (const auto & r : reserved_) + VirtualFree(r.first, 0, MEM_RELEASE); +#else + // POSIX: the model's own munmap of "its" range has already removed the placeholder by the + // time this runs (it is called after the model is freed), and unmapping the range again could + // take out whatever was mapped there since. If the model never unmaps, an inaccessible + // reservation leaks until process exit. Either way: forget, do not touch. +#endif +#if defined(_WIN32) + // The plugs are deliberately NOT closed here. Each occupies a handle value the model still + // believes is its section, and the model's own teardown closes that value; closing it here as + // well would be a double close of whatever occupies the slot by then. If the model never + // closes it (it always does today), one inert event per section leaks until process exit. +#endif + reserved_.clear(); + plugs_.clear(); +} + +#if defined(_WIN32) + +namespace { + +std::wstring widen(const std::string & s) { + if (s.empty()) return {}; + const int n = MultiByteToWideChar(CP_UTF8, 0, s.data(), (int) s.size(), nullptr, 0); + if (n <= 0) return {}; + std::wstring w((size_t) n, L'\0'); + MultiByteToWideChar(CP_UTF8, 0, s.data(), (int) s.size(), &w[0], n); + return w; +} + +std::wstring lower(std::wstring w) { + for (wchar_t & c : w) + c = (wchar_t) std::towlower(c); + return w; +} + +// The file's NT device path ("\Device\HarddiskVolumeN\..."), which is the form the memory manager +// reports for mapped views, so the two can be compared directly. Opened with no access rights: +// identification only, and closed before returning. Empty on failure. +std::wstring nt_path_of(const std::wstring & path) { + HANDLE h = CreateFileW(path.c_str(), 0, FILE_SHARE_READ | FILE_SHARE_WRITE | FILE_SHARE_DELETE, nullptr, + OPEN_EXISTING, FILE_ATTRIBUTE_NORMAL, nullptr); + if (h == INVALID_HANDLE_VALUE) return {}; + wchar_t buf[32768]; + const DWORD n = GetFinalPathNameByHandleW(h, buf, (DWORD) (sizeof(buf) / sizeof(buf[0])), VOLUME_NAME_NT); + CloseHandle(h); + if (n == 0 || n >= sizeof(buf) / sizeof(buf[0])) return {}; + return lower(std::wstring(buf, n)); +} + +using GetMappedFileNameW_t = DWORD(WINAPI *)(HANDLE, LPVOID, LPWSTR, DWORD); +using NtQueryInformationProcess_t = NTSTATUS(NTAPI *)(HANDLE, PROCESSINFOCLASS, PVOID, ULONG, PULONG); +using NtQueryObject_t = NTSTATUS(NTAPI *)(HANDLE, OBJECT_INFORMATION_CLASS, PVOID, ULONG, PULONG); + +// Which file a view (or any address inside one) maps, as an NT device path; empty if none. +std::wstring mapped_name(GetMappedFileNameW_t fn, void * addr) { + wchar_t buf[32768]; + const DWORD n = fn(GetCurrentProcess(), addr, buf, (DWORD) (sizeof(buf) / sizeof(buf[0]))); + if (n == 0) return {}; + return lower(std::wstring(buf, n)); +} + +bool is_target(const std::vector & targets, const std::wstring & name) { + return !name.empty() && std::find(targets.begin(), targets.end(), name) != targets.end(); +} + +// Every view of the target files: [AllocationBase, end) per view, assembled from the regions +// VirtualQuery walks (a view can be reported as several regions with one AllocationBase). +std::vector> find_views(GetMappedFileNameW_t fn, const std::vector & targets) { + std::map ends; // AllocationBase -> highest end seen + MEMORY_BASIC_INFORMATION mbi; + const char * p = nullptr; + while (VirtualQuery(p, &mbi, sizeof(mbi)) == sizeof(mbi)) { + if (mbi.Type == MEM_MAPPED && mbi.State == MEM_COMMIT && is_target(targets, mapped_name(fn, mbi.BaseAddress))) { + const uintptr_t end = (uintptr_t) mbi.BaseAddress + mbi.RegionSize; + uintptr_t & e = ends[mbi.AllocationBase]; + e = std::max(e, end); + } + const char * next = (const char *) mbi.BaseAddress + mbi.RegionSize; + if (next <= p) break; // wrapped: end of the address space + p = next; + } + std::vector> out; + for (const auto & kv : ends) + out.emplace_back(kv.first, (size_t) (kv.second - (uintptr_t) kv.first)); + return out; +} + +// Hold the released range so nothing else can be allocated there: the model will still call +// UnmapViewOfFile on this base at teardown, which must find our reservation (a harmless failure), +// never a live allocation. A reservation is in 64 KiB units, so if the exact size collides with a +// neighbour the tail is given up rather than the whole placeholder. +void * reserve_placeholder(void * base, size_t size) { + if (void * r = VirtualAlloc(base, size, MEM_RESERVE, PAGE_NOACCESS)) return r; + SYSTEM_INFO si; + GetSystemInfo(&si); + const size_t g = si.dwAllocationGranularity; + const size_t down = size & ~(g - 1); + if (down && down < size) return VirtualAlloc(base, down, MEM_RESERVE, PAGE_NOACCESS); + return nullptr; +} + +struct HandleEntry { + HANDLE HandleValue; + ULONG_PTR HandleCount; + ULONG_PTR PointerCount; + ULONG GrantedAccess; + ULONG ObjectTypeIndex; + ULONG HandleAttributes; + ULONG Reserved; +}; +struct HandleSnapshot { + ULONG_PTR NumberOfHandles; + ULONG_PTR Reserved; + HandleEntry Handles[1]; +}; +constexpr int kProcessHandleInformation = 51; +constexpr NTSTATUS kStatusInfoLengthMismatch = (NTSTATUS) 0xC0000004L; +constexpr DWORD kSectionMapRead = 0x0004; // SECTION_MAP_READ +constexpr ULONG kProtectFromClose = 0x2; // HANDLE_FLAG_PROTECT_FROM_CLOSE + +std::vector handle_snapshot(NtQueryInformationProcess_t q) { + std::vector buf(1 << 16); + for (int attempt = 0; attempt < 8; ++attempt) { + ULONG need = 0; + const NTSTATUS st = + q(GetCurrentProcess(), (PROCESSINFOCLASS) kProcessHandleInformation, buf.data(), (ULONG) buf.size(), &need); + if (st == kStatusInfoLengthMismatch) { + buf.resize((size_t) need + (1 << 14)); + continue; + } + if (st < 0) return {}; + return buf; + } + return {}; +} + +bool is_section(NtQueryObject_t q, HANDLE h) { + alignas(PUBLIC_OBJECT_TYPE_INFORMATION) char buf[sizeof(PUBLIC_OBJECT_TYPE_INFORMATION) + 512]; + ULONG len = 0; + if (q(h, ObjectTypeInformation, buf, (ULONG) sizeof(buf), &len) < 0) return false; + const auto * ti = (const PUBLIC_OBJECT_TYPE_INFORMATION *) buf; + const size_t n = ti->TypeName.Length / sizeof(wchar_t); + return n == 7 && std::wstring(ti->TypeName.Buffer, n) == L"Section"; +} + +// Re-occupy a just-closed handle value with an inert event so no later handle can inherit it and +// be closed by the model's teardown in its place. Handle values are recycled lowest-free-first, so +// the first attempt normally lands; a few more cover a concurrent allocation. Returns the plug, or +// null if the slot could not be reclaimed (then the model's late CloseHandle hits whatever lives +// there, which is why the caller reports it). +HANDLE plug_handle_slot(HANDLE value) { + std::vector misses; + HANDLE plug = nullptr; + for (int i = 0; i < 16 && !plug; ++i) { + HANDLE e = CreateEventW(nullptr, TRUE, FALSE, nullptr); + if (!e) break; + if (e == value) + plug = e; + else + misses.push_back(e); + } + for (HANDLE m : misses) + CloseHandle(m); + return plug; +} + +} // namespace + +size_t addresses_in_file_mappings(const std::vector & paths, const std::vector & addresses) { + HMODULE k32 = GetModuleHandleW(L"kernel32.dll"); + auto mapped_fn = k32 ? (GetMappedFileNameW_t) (void *) GetProcAddress(k32, "K32GetMappedFileNameW") : nullptr; + if (!mapped_fn) return 0; // cannot tell; the caller's other conditions still apply + std::vector targets; + for (const std::string & p : paths) { + std::wstring nt = nt_path_of(widen(p)); + if (!nt.empty()) targets.push_back(std::move(nt)); + } + size_t hits = 0; + for (const void * a : addresses) { + if (!a) continue; + MEMORY_BASIC_INFORMATION mbi; + if (VirtualQuery(a, &mbi, sizeof(mbi)) != sizeof(mbi)) continue; + if (mbi.Type != MEM_MAPPED) continue; + if (is_target(targets, mapped_name(mapped_fn, (void *) a))) ++hits; + } + return hits; +} + +MappingReleaseReport release_file_mappings(const std::vector & paths, MappingPlaceholders * out) { + MappingReleaseReport rep; + rep.supported = true; + + HMODULE k32 = GetModuleHandleW(L"kernel32.dll"); + HMODULE ntdll = GetModuleHandleW(L"ntdll.dll"); + auto mapped_fn = k32 ? (GetMappedFileNameW_t) (void *) GetProcAddress(k32, "K32GetMappedFileNameW") : nullptr; + auto qip = + ntdll ? (NtQueryInformationProcess_t) (void *) GetProcAddress(ntdll, "NtQueryInformationProcess") : nullptr; + auto qobj = ntdll ? (NtQueryObject_t) (void *) GetProcAddress(ntdll, "NtQueryObject") : nullptr; + if (!mapped_fn || !qip || !qobj) { + rep.error = "required kernel32/ntdll entry points unavailable"; + return rep; + } + + std::vector targets; + for (const std::string & p : paths) { + std::wstring nt = nt_path_of(widen(p)); + if (nt.empty()) { + if (rep.error.empty()) rep.error = "cannot identify " + p; + continue; + } + targets.push_back(std::move(nt)); + } + if (targets.empty()) return rep; + + // 1. Views. Unmapping first means the section-closing pass below cannot pull a view out from + // under a live mapping; each view's range is held by a placeholder afterwards. + for (const auto & v : find_views(mapped_fn, targets)) { + if (!UnmapViewOfFile(v.first)) { + if (rep.error.empty()) rep.error = "UnmapViewOfFile failed"; + continue; + } + ++rep.views_unmapped; + rep.bytes += v.second; + if (void * r = reserve_placeholder(v.first, v.second)) { + if (out) out->reserved_.emplace_back(r, v.second); + } else if (rep.error.empty()) { + rep.error = "placeholder reservation failed"; + } + } + + // 2. Section handles. A section is matched by mapping one page of it and asking which file + // that page belongs to, never by guessing from handle values or types alone. + const std::vector snap = handle_snapshot(qip); + if (snap.empty()) { + if (rep.error.empty()) rep.error = "process handle snapshot unavailable"; + return rep; + } + const auto * hs = (const HandleSnapshot *) snap.data(); + for (ULONG_PTR i = 0; i < hs->NumberOfHandles; ++i) { + const HandleEntry & e = hs->Handles[i]; + if ((e.HandleAttributes & kProtectFromClose) || !(e.GrantedAccess & kSectionMapRead)) continue; + if (!is_section(qobj, e.HandleValue)) continue; + void * probe = MapViewOfFile(e.HandleValue, FILE_MAP_READ, 0, 0, 1); + if (!probe) continue; + const bool match = is_target(targets, mapped_name(mapped_fn, probe)); + UnmapViewOfFile(probe); + if (!match) continue; + if (!CloseHandle(e.HandleValue)) { + if (rep.error.empty()) rep.error = "CloseHandle on section failed"; + continue; + } + ++rep.sections_closed; + if (HANDLE plug = plug_handle_slot(e.HandleValue)) { + if (out) out->plugs_.push_back(plug); + } else { + ++rep.plugs_missed; + } + } + return rep; +} + +#else // POSIX: the file's VMAs come from /proc/self/maps; there are no handles to deal with. + +// Each region is unmapped and immediately replaced by an inaccessible anonymous placeholder at the +// same address, so the model's later munmap of "its" range removes the placeholder and nothing +// else, and no allocation can land there in between. Offered for measurement: Android (f2fs) was +// measured NOT to serialise direct reads against a live mapping, so there is no throughput to +// gain there; what a release buys on a given POSIX system is a bench question. +size_t addresses_in_file_mappings(const std::vector & paths, const std::vector & addresses) { + std::vector regs; + std::vector all; + for (const std::string & p : paths) { + const size_t slash = p.find_last_of('/'); + const std::string base = slash == std::string::npos ? p : p.substr(slash + 1); + regs.clear(); + if (file_mapped_regions(base.c_str(), regs)) all.insert(all.end(), regs.begin(), regs.end()); + } + size_t hits = 0; + for (const void * a : addresses) { + const uintptr_t x = (uintptr_t) a; + for (const MappedRegion & r : all) + if (x >= r.start && x < r.end) { + ++hits; + break; + } + } + return hits; +} + +MappingReleaseReport release_file_mappings(const std::vector & paths, MappingPlaceholders * out) { + MappingReleaseReport rep; + rep.supported = true; + for (const std::string & p : paths) { + const size_t slash = p.find_last_of('/'); + const std::string base = slash == std::string::npos ? p : p.substr(slash + 1); + std::vector regs; + if (!file_mapped_regions(base.c_str(), regs)) { + if (rep.error.empty()) rep.error = "cannot read /proc/self/maps"; + continue; + } + for (const MappedRegion & r : regs) { + void * addr = (void *) r.start; + const size_t size = (size_t) (r.end - r.start); + if (munmap(addr, size) != 0) { + if (rep.error.empty()) rep.error = "munmap failed"; + continue; + } + ++rep.views_unmapped; + rep.bytes += size; + void * ph = mmap(addr, size, PROT_NONE, MAP_PRIVATE | MAP_ANONYMOUS | MAP_FIXED | MAP_NORESERVE, -1, 0); + if (ph == MAP_FAILED) { + if (rep.error.empty()) rep.error = "placeholder mapping failed"; + } else if (out) { + out->reserved_.emplace_back(ph, size); + } + } + } + return rep; +} + +#endif + +} // namespace bmoe::pio diff --git a/core/src/io/mapping_release.h b/core/src/io/mapping_release.h new file mode 100644 index 0000000..308ab5a --- /dev/null +++ b/core/src/io/mapping_release.h @@ -0,0 +1,76 @@ +#pragma once +// Release a model file's memory mapping once nothing reads through it any more. +// +// Why this exists. llama.cpp maps every gguf it loads (use_mmap is load-bearing for the streamer: +// the native layout is what the expert rebind points into) and keeps the mapping for the model's +// lifetime. On Windows that mapping has a cost the streamer pays on every read: while a section of +// the file is alive, NTFS serialises concurrent unbuffered (FILE_FLAG_NO_BUFFERING) reads on that +// file. Measured with tools/bmoe-iobench on the same drive and request shape (576 KiB, 4 lanes): +// 2400-2660 MiB/s with no mapping, 895-930 with one, i.e. four lanes deliver exactly one lane's +// throughput. Unmapping the view alone does not help; the section must be closed too. +// +// Under `--dense-weights anon`/`ahwb` with every file tensor accounted for (streamed experts, +// copied dense weights, row-streamed tables), nothing dereferences the mapping after load, so it +// can go. llama.cpp offers no public API to drop it (its Windows unmap_fragment is a no-op), so +// this module finds the process's own views and section handles of the file and releases them. +// +// What it does NOT do: touch llama.cpp. The model still believes it owns the mapping and will +// UnmapViewOfFile/CloseHandle it at free time. Both are made harmless here rather than left to +// luck: the view's address range is re-reserved as a placeholder (so no later allocation can land +// there and be unmapped by mistake), and the closed handle's slot is re-occupied by an inert event +// (so a handle allocated later cannot inherit the value and be closed by mistake). The placeholders +// are returned to the caller, who releases them AFTER the model is freed. +// +// On POSIX the same release is done with munmap over the file's VMAs from /proc/self/maps, with +// anonymous placeholders in their place. Android (f2fs) was measured not to serialise direct reads +// against a live mapping, so there it buys nothing; it exists so that claim can be re-measured +// rather than assumed on the next platform. +#include +#include +#include + +namespace bmoe::pio { + +struct MappingReleaseReport { + bool supported = false; // false: nothing to do on this platform (POSIX); the rest is zero + int views_unmapped = 0; + int sections_closed = 0; + uint64_t bytes = 0; // total size of the views released + int plugs_missed = 0; // closed sections whose handle slot could not be re-occupied + std::string error; // first hard failure, empty on success; partial work is still reported +}; + +// Everything a release leaves behind on purpose, to be released only after the model owning the +// original mapping has been freed (see the header comment). Movable, not copyable; releasing twice +// is a no-op. +class MappingPlaceholders { +public: + MappingPlaceholders() = default; + ~MappingPlaceholders(); + MappingPlaceholders(MappingPlaceholders &&) noexcept; + MappingPlaceholders & operator=(MappingPlaceholders &&) noexcept; + MappingPlaceholders(const MappingPlaceholders &) = delete; + MappingPlaceholders & operator=(const MappingPlaceholders &) = delete; + + void release(); + bool empty() const { return reserved_.empty() && plugs_.empty(); } + +private: + friend MappingReleaseReport release_file_mappings(const std::vector &, MappingPlaceholders *); + std::vector> reserved_; // address ranges held in place of the views + std::vector plugs_; // inert handles occupying closed sections' slots +}; + +// Release every view and section this process holds on the files in `paths` (UTF-8). Files that +// cannot be opened for identification are skipped and named in `report.error`. `out` receives the +// placeholders; it must outlive the model. +MappingReleaseReport release_file_mappings(const std::vector & paths, MappingPlaceholders * out); + +// How many of `addresses` still point inside a mapping of one of `paths`. This is what makes a +// release decidable rather than guessed: releasing is correct only while nothing dereferences the +// mapping, and the caller can ask about every pointer it knows instead of reasoning about which +// tensor names ought to have been rebound. It answers for the pointers it is given and for no +// others, so a caller that cannot enumerate everything still has to keep the flag opt-in. +size_t addresses_in_file_mappings(const std::vector & paths, const std::vector & addresses); + +} // namespace bmoe::pio diff --git a/core/src/moe/dense_weights.cpp b/core/src/moe/dense_weights.cpp index bd9466d..e43e2b9 100644 --- a/core/src/moe/dense_weights.cpp +++ b/core/src/moe/dense_weights.cpp @@ -156,6 +156,15 @@ void DenseWeights::set_row_gathered(std::vector tables, uint64_t row_budget_ = budget_bytes; } +bool DenseWeights::reopen_readers() { + return rows_ ? rows_->reopen_readers() : true; +} + +bool DenseWeights::file_mapping_in_use() const { + if (!mapped_.empty()) return true; + return (mode_ == DenseWeightsMode::Mmap || mode_ == DenseWeightsMode::Warmed) && !tensors_.empty(); +} + IRowSource * DenseWeights::row_source() const { return rows_ && !rows_->empty() ? rows_.get() : nullptr; } @@ -282,6 +291,8 @@ bool DenseWeights::read_anonymous(size_t align) { done += n; } d.tensor->data = buf; // rebind the model weight onto its private copy + for (ggml_tensor * alias : d.aliases) + alias->data = buf; // and every twin over the same bytes, onto the same copy total += d.size; } std::fprintf(stderr, "bmoe: dense-weights=%s — %llu MiB in %zu %s buffers\n", pinned ? "ahwb" : "anon", diff --git a/core/src/moe/dense_weights.h b/core/src/moe/dense_weights.h index 5fd03d2..4a72ec5 100644 --- a/core/src/moe/dense_weights.h +++ b/core/src/moe/dense_weights.h @@ -40,6 +40,13 @@ struct DenseTensorRef { uint64_t file_off = 0; uint64_t size = 0; int file_idx = 0; // which shard file holds the bytes (0 for a single-file model) + // Other tensor objects over the SAME file bytes, which a residency policy must rebind together + // with `tensor` or leave together where they are. A tied output head is the case that produces + // them: llama.cpp builds the head from the token embedding table, so the model holds two + // ggml_tensors, identically named, over one range. Rebinding one and not the other leaves the + // graph reading the model's mmap every token, and makes the range unsafe to release. Empty for + // every ordinary weight, so the common path allocates nothing. + std::vector aliases; }; class DenseWeights { @@ -81,6 +88,15 @@ public: // The row policy, once init has run: null when no table qualified or the takeover failed. The // engine's graph adapter needs it to make rows present before a gather node runs. IRowSource * row_source() const; + + // Whether any tensor this module was given still reads through the model file's mapping once + // init has run: every tensor under Mmap/Warmed, and under Anonymous/Pinned the ones held back + // as oversized. Row-gathered tables are served from our own reader, not the mapping. + bool file_mapping_in_use() const; + + // Reopen the readers still in service after init: the row policy's, if any (the dense copy's + // own readers have done their one job). See FileReader::reopen for why this exists. + bool reopen_readers(); RowSourceStats row_stats() const; // Sample how much of the dense set the kernel still has in RAM (mincore), setting resident_frac(). diff --git a/core/src/moe/expert_stream_source.cpp b/core/src/moe/expert_stream_source.cpp index 1a35350..860acb8 100644 --- a/core/src/moe/expert_stream_source.cpp +++ b/core/src/moe/expert_stream_source.cpp @@ -292,6 +292,18 @@ bool ExpertStreamSource::init(const std::vector & shard_paths, return true; } +bool ExpertStreamSource::reopen_readers() { + std::lock_guard lk(io_mtx_); + if (!jobs_.empty() || !spec_jobs_.empty()) { + std::fprintf(stderr, "bmoe: reopen_readers refused: reads are queued\n"); + return false; + } + bool ok = true; + for (auto & r : readers_) + ok = r->reopen() && ok; + return dense_.reopen_readers() && ok; +} + // ── one aligned slice read on a lane ──────────────────────────────────────────────── // The bytes come from the reader; this wraps it with the domain the reader must not know about — the // per-read I/O trace that attributes a read to its (layer, expert, projection). Latency is timed here diff --git a/core/src/moe/expert_stream_source.h b/core/src/moe/expert_stream_source.h index b5d824e..5995b7f 100644 --- a/core/src/moe/expert_stream_source.h +++ b/core/src/moe/expert_stream_source.h @@ -86,6 +86,11 @@ public: // The row policy after init, for the graph adapter that must make rows present before a gather // node runs; null when nothing qualified. Its accounting, for telemetry, is row_stats(). IRowSource * row_source() const { return dense_.row_source(); } + bool dense_file_mapping_in_use() const { return dense_.file_mapping_in_use(); } + // Reopen every reader that will serve decode (expert lanes and row tables). Only legal between + // init and the first decode, when the workers are idle and nothing is queued; checked under the + // queue lock. See FileReader::reopen. + bool reopen_readers(); RowSourceStats row_stats() const { return dense_.row_stats(); } // IExpertSource diff --git a/core/src/moe/router_hook.cpp b/core/src/moe/router_hook.cpp index 68f02fe..f0ec178 100644 --- a/core/src/moe/router_hook.cpp +++ b/core/src/moe/router_hook.cpp @@ -168,6 +168,8 @@ void RouterHook::begin_capture() { for (auto & L : captured_) L = LayerExperts{}; captured_weights_.clear(); + captured_weight_objects_.clear(); + captured_weight_seen_.clear(); row_gathered_.clear(); row_disqualified_.clear(); } @@ -1379,6 +1381,7 @@ bool RouterHook::on_eval(ggml_tensor * t, bool ask) { // filtered out downstream by the gguf tensor set, so recording them here is harmless. if (src->op != GGML_OP_NONE) continue; captured_weights_[src->name] = src; + if (captured_weight_seen_.insert(src).second) captured_weight_objects_.push_back(src); // …and HOW this node used it, which is what decides whether its residency can be // reduced to the rows the graph asks for. Only the table position of a row gather // counts, and only when the index is already in memory when the node runs. diff --git a/core/src/moe/router_hook.h b/core/src/moe/router_hook.h index c64c80d..7d20cf6 100644 --- a/core/src/moe/router_hook.h +++ b/core/src/moe/router_hook.h @@ -73,6 +73,14 @@ public: // against the gguf tensor set, which those do not belong to. Only .tensor is meaningful here. const std::unordered_map & captured_weights() const { return captured_weights_; } + // Every distinct weight leaf OBJECT the graph read, deduplicated by address rather than by + // name. The map above cannot answer this: a model whose output head is tied to its token + // embedding table has two ggml_tensor objects carrying the same name over the same file + // bytes, and a name-keyed map keeps one of them. Rebinding only that one leaves the twin + // reading the model's mmap for the life of the run, which is what --dense-weights anon and + // ahwb exist to prevent. Order is first-seen, so a run is reproducible. + const std::vector & captured_weight_objects() const { return captured_weight_objects_; } + // After capture, the subset of those weights the graph only ever GATHERS ROWS from — the shape a // token embedding table has, and the one residency policy can exploit (see IRowSource). A name is // in this set only if EVERY node that referenced the tensor was a row gather taking it as the @@ -285,6 +293,8 @@ private: IRowSource * row_source_ = nullptr; // row-gathered dense tables, when the policy is on std::vector captured_; std::unordered_map captured_weights_; + std::vector captured_weight_objects_; // same leaves, deduplicated by address + std::unordered_set captured_weight_seen_; // Capture-time evidence for row_gathered_weights(): every weight seen as the TABLE of a row // gather, and every weight seen in any way that rules that out. The verdict is the difference. std::unordered_set row_gathered_; diff --git a/core/src/moe/row_stream.cpp b/core/src/moe/row_stream.cpp index 2ea0f8d..468a4a1 100644 --- a/core/src/moe/row_stream.cpp +++ b/core/src/moe/row_stream.cpp @@ -260,6 +260,13 @@ void RowStream::release() { table_bytes_ = 0; } +bool RowStream::reopen_readers() { + bool ok = true; + for (auto & r : readers_) + ok = r->reopen() && ok; + return ok; +} + void RowStream::shutdown() { release(); } diff --git a/core/src/moe/row_stream.h b/core/src/moe/row_stream.h index 491f9eb..ee9aa22 100644 --- a/core/src/moe/row_stream.h +++ b/core/src/moe/row_stream.h @@ -70,6 +70,8 @@ public: RowSourceStats stats() const override; void shutdown(); + // Reopen the per-shard readers (see FileReader::reopen). No gather may be in flight. + bool reopen_readers(); // A row is far below any device's request floor, so a slab is the unit actually read: large // enough that the read costs what the floor costs anyway, small enough that a scattered diff --git a/docs/architecture.md b/docs/architecture.md index 68a72a4..912b373 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -22,6 +22,7 @@ core/ src/ io/ platform_io — O_DIRECT reads + reserve/commit/evict VM, cross-platform file_reader — pooled positioned reader, per-consumer O_DIRECT + mapping_release - releases the model file's mapping after load (--release-mmap) moe/ gguf_offsets (tensor → (shard, offset), split ggufs included), arch_registry, expert_stream_source (one reader per shard), router_hook dense_weights — non-expert weight policy + the residency sensor diff --git a/docs/bench-data/2026-08-29-mmap-serialisation/findings.md b/docs/bench-data/2026-08-29-mmap-serialisation/findings.md new file mode 100644 index 0000000..3aa229d --- /dev/null +++ b/docs/bench-data/2026-08-29-mmap-serialisation/findings.md @@ -0,0 +1,168 @@ +# The model file's mapping serialises the streamer's reads (Windows) + +2026-08-29. Why the desktop host has always read at about a third of what its drive can serve, why +more I/O lanes and more compute threads measured dead there, and what `--release-mmap` recovers. + +Hosts: a Windows x86 laptop (8 cores, 15.5 GB RAM, NVMe SSD) and, for the control, the 12 GB / +UFS 4.x Android test phone. Model in every cell: Qwen3.6-35B-A3B-Q4_K_M (22.3 GB). Read shape in +every iobench cell: 576 KiB per read, O_DIRECT, random offsets — the shape one routed expert +projection actually has. + +## The drive, and what changes it + +`bmoe-iobench --model M.gguf --lanes 1,4 --slice-kb 576 --seconds 3`, host, one variable per row: + +| variable | 1 lane | 4 lanes | +|---|---:|---:| +| nothing (the drive) | 883-967 MiB/s | **2295-2665 MiB/s** | +| `--range-mb 440` (offsets confined to one region) | 917 | 2651 | +| `--compute-load 8` (lanes read against a busy CPU) | — | 2463 | +| `--fresh` (destination pages committed afresh per read) | 745 | 1952 | +| **`--mmap` (a read-only mapping of the file held open)** | 909 | **895-930** | +| `--mmap --reopen-lanes` (mapping dropped, lanes reopened) | — | **2402** | + +Locality, CPU contention and page-commit cost move the number by a few per cent. The mapping moves +it by **2.6x**, and four lanes under it deliver exactly one lane's throughput (4 / 2.42 ms = 1.65 +reads/ms against 1 / 0.62 = 1.61). Nothing is ever read through the mapping in these cells: only +its existence is the variable. + +Two halves, separable and both necessary: + +1. While a section of the file is alive, concurrent unbuffered reads on that file are serialised. +2. A lane opened *while* the section existed keeps serialising against it after the section is + gone — dropping the mapping alone leaves the rate at 927 MiB/s; dropping it and reopening the + lanes restores 2402. + +Sub-cells that placed the blame precisely, all at 4 lanes: mapping opened before the lanes vs after +them (927 either way, so it is not open order); reading through a hard link while the original is +mapped (927, so it is the file, not the handle); mapping the file and unmapping the view but +keeping the section (926) vs closing the section too (2050-2207, recovered); a plain unbuffered +handle held open with no mapping at all (916 — that handle had been opened under a live section +earlier in the process, which is finding 2 again). + +## What it costs the engine, and what the flag returns + +Host, `--moe-stream --cache-mb 3000 --io-threads 4 --overlap --dense-weights anon -t 8 +--ubatch 256 -c 1024`, 96 generated tokens, cells interleaved: + +| | baseline | `--release-mmap` | +|---|---:|---:| +| decode | 3.16 tok/s | **4.63 tok/s (+46%)** | +| s/token | 0.316 | 0.216 | +| compute | 0.107 | 0.114 | +| cache mgmt | 0.027 | 0.028 | +| **flash stall** | **0.182** | **0.074** | +| per-read latency p50 (`--io-trace`) | 2.49 ms | 1.20 ms | +| flash read per token | 192.5 MiB | 192.5 MiB | +| cache hit / evictions / re-reads | 60.9% / 11635 / 7848 | identical | +| generated text | — | byte-identical | + +Not one byte fewer is read and no quality knob is touched: the same reads are simply served in +parallel. The remaining decode budget is 0.114 compute + 0.028 mgmt + 0.074 stall, so the ceiling +if flash were free is 7.0 tok/s. + +Two earlier verdicts are corrected by this. "Lanes and threads are dead on the desktop host" +(2026-07-24) and "the serial path reads at exactly the one-lane rate with four lanes open" +(2026-08-15) were both this mechanism, not a read-path defect. The "at least 1.7x of headroom" from +the same session is now accounted for and collected. + +## The phone: flat, so this is a Windows defect + +Same tool, same model, same read shape, on the Android test phone via adb: + +| | 1 lane | 2 lanes | 4 lanes | +|---|---:|---:|---:| +| no mapping | 967 MiB/s | 1785 | 2180 | +| mapping held open | 1122 | 1890 | 2200 | + +No serialisation: f2fs does not take an exclusive lock against a live mapping, which is also why +lanes have always scaled in the engine there (2 to 4 lanes was worth +15-20%). The engine's own +per-read latency on the phone is 1.24 ms against the drive's 1.03 at 4 lanes, i.e. it already reads +at about 80% of the drive's rate with the compute running alongside. + +So the read-path prediction for the phone is "no gain", and the engine confirms it: flash stall is +the same with the flag and without. The flag is still faster there, for another reason. Same model +and device, `--cache-mb 2000 --io-threads 4 --overlap --dense-weights anon -t 4 --ubatch 512 +-c 1024`, 48 tokens, two cells per variant run in both orders with a thermal gate at 45 C: + +| | baseline | `--release-mmap` | +|---|---:|---:| +| decode | 3.92, 3.56 tok/s | **4.14, 4.03 tok/s** | +| **flash stall** | 0.088, 0.087 | **0.089, 0.093** | +| CPU per token | 0.634, 0.688 cpu-s | **0.592, 0.604 cpu-s** | +| major faults per token | 70.1, 0.85 | 0.00, 0.21 | +| prefill | 7.1, 8.7 tok/s | 8.4, 7.3 tok/s | +| model load | 34.7, 24.6 s | 25.8, 26.8 s | +| bytes, hit rate, evictions, re-reads | — | identical to the digit | + +Stall is flat, so this is not the Windows effect; CPU per token falls ~9% and decode follows. +Keeping a 20 GB mapping registered costs a kernel already under memory pressure, and handing it +back removes that cost from the decode's own CPU time. The 70 major faults of the first cell are a +cold-start artefact rather than the explanation — the second baseline had 0.85 and was still the +slowest cell. Prefill and model load are unaffected, which is expected: prefill on this device is +compute-bound and never waits on flash, and unmapping 20 GB costs nothing measurable. + +Read this as a direction, not a number. Two cells per variant, 48 tokens each, against a device +whose prefill alone spread 20% (7.1 to 8.7 tok/s) across cells that differ in nothing. What holds +it up is that the ordering is consistent in both directions and that CPU per token moves with it. +A 256-token protocol is owed before any of this becomes a published figure. + +Owed: the same cells on macOS, which #182 has just made worth running — a direct request there now +applies `F_NOCACHE` to the descriptor instead of falling back to a buffered read. `F_NOCACHE` is a +caching hint rather than an I/O mode, so whether a live mapping of the same file interacts with it +at all is an open question, and the answer is one `bmoe-iobench --mmap` cell on a Mac. + +## What the gates caught: a tied output head keeps a pointer nobody can see + +Turning the flag on by default failed the `moe_gates_gemma4` byte-identity gate with a segfault, and +the reason is worth recording because it bounds what this feature can ever promise. + +A bisect of the release (unmap the view but keep the section; close the section but keep the view) +put the fault on **unmapping**, so something still dereferenced the mapping. A vectored exception +handler reported the faulting addresses as file offsets 12192 and 79264, both inside +`token_embd.weight`. But the accounting said every gguf tensor was owned by a policy, and it was +right: the tiny gemma4 model has **no `output.weight`**, so llama.cpp builds the output head from +the embedding table and the model carries a *second* tensor over the same mapped bytes. The engine +rebinds the tensor the gguf names; the twin is not a gguf tensor, so no name-based accounting can +ever see it. The qwen3moe model, which has a separate `output.weight`, passes. + +The first fix was a blacklist — decline when the gguf has no `output.weight` — and it was the wrong +shape of answer. The twin is not a special case to route around, it is a weight the engine failed to +rebind: under `--dense-weights anon` the output projection was being served by page faults from the +mmap on every tied model, which is exactly what that policy exists to prevent. So the capture now +records every distinct leaf object rather than one per name, the dense policy rebinds every tensor +over a file range onto the same buffer, and the release decision became a question asked of the OS: +does any weight the capture observed still point inside the mapping? With that, gemma4 releases +cleanly, `--dense-weights mmap` reports 39 weights still mapped and stands down, and every gate +passes. + +What is left is a residue, and it is why the flag is opt-in: the check answers for pointers the +capture pass observed, and a graph shape the session never builds could hold another. + +Two alternatives were measured and refuted while looking for a design that could not fail this way: + +- **Close the section handle and keep the view mapped.** Would be safe by construction — no pointer + can dangle — but it does not lift the serialisation: 792 MiB/s against 790 with the mapping fully + alive. It is the view, not the handle. +- **Read buffered instead of unbuffered while the mapping is alive.** Buffered reads do escape the + serialisation, but the microbenchmark cannot price it honestly: `--buffered --mmap` measures + 3406 MiB/s precisely because it is being served from the page cache the mapping populated, which + is the pollution O_DIRECT exists to avoid on a model many times larger than RAM. Answering it + properly needs an engine-level A/B on Windows, not an iobench cell. + +## Two things that are NOT this, checked because they looked like it + +**The expert-cache stall accounting is not a hot-path cost.** `StallUnion` (0.23.0) takes a +process-wide mutex around each stall interval, which reads like a serialization point on the +compute path. It is not: the readiness fast path returns before `enter()`, so the mutex is taken +only when a thread is about to spin and possibly sleep anyway. Measured on the host with a twin +binary whose `StallUnion::enter/exit` are no-ops, in the configuration with the *most* stall +(baseline, no release), interleaved 96-token cells: **3.281 / 3.272 tok/s with, 3.286 / 3.289 / +3.290 without** — 0.4%, inside the run-to-run spread, and the direction is not even consistent +with the cold first cell (3.103) discarded. No regression. + +**A second engine on the same drive invalidates everything.** Two of these cells were first +measured while another process was streaming the same model on the same disk; the drive's own +curve read half its real value and the lane counts looked flat for the wrong reason. Every number +above was re-taken with the machine otherwise idle. The rule stands: one engine at a time, and +check before trusting a rate. diff --git a/docs/moe-streaming.md b/docs/moe-streaming.md index eddf6de..85127fc 100644 --- a/docs/moe-streaming.md +++ b/docs/moe-streaming.md @@ -63,6 +63,45 @@ calling thread participates as lane 0. On UFS 4.x, 4 lanes roughly triples effec bandwidth over serial. Compute threads (`-t`) show a U-shape — 4 is the measured optimum; 8 regresses badly because ggml's spin-wait contends with the synchronous reads. +## The model file's mapping (`--release-mmap`) + +llama.cpp maps the gguf and keeps it mapped for the model's lifetime. On Windows that mapping +serialises the lanes above: while a section of the file is alive, concurrent unbuffered reads on +it are taken one at a time, so `--io-threads 4` reads at one lane's rate. It is not the drive and +it is not the engine — `bmoe-iobench --model M.gguf --lanes 4 --slice-kb 576` measures 2400-2660 +MiB/s, the same command with `--mmap` measures 895-930, and adding `--reopen-lanes` recovers the +full rate. A lane opened while the section existed stays serialised after it is gone, which is why +the recovery needs both halves. + +`--release-mmap` does exactly that inside the engine: after load, once nothing reads through the +mapping any more, the file is unmapped, its section closed and the reader lanes reopened. + +Whether releasing is safe is decided by looking, not by reasoning about which tensors ought to have +been rebound: the engine asks the OS whether any weight the capture pass observed still points inside +a mapping of the model files, and declines if any does. That one question covers every residency +policy — a dense set left mmap'd under `mmap` or `warm`, a table held back as oversized, or a tensor +no name-based accounting could have found. Run with `--dense-weights mmap` and the engine reports the +count and stands down. + +On Windows the run ends with `warning: UnmapViewOfFile failed`, printed by llama.cpp rather than +by the engine. That is the designed outcome, not a defect: llama.cpp is not patched and still +believes it owns the mapping, so at teardown it unmaps a base the engine has already released. An +inert reservation is left in that range precisely so the call finds a placeholder and fails +harmlessly, instead of finding whatever was allocated there next. + +It is still opt-in, because the check answers for the pointers the capture pass saw and for no +others. A graph shape this session never builds could hold another one, and llama.cpp exposes no way +to enumerate a loaded model's tensors and settle it. The MTP draft is the concrete case: it builds a +second graph, so a gguf tensor no policy owns blocks the release there as well. Worth +46% decode on +the desktop host, byte-identical. + +Android is a different story with the same conclusion. The iobench cells above are flat there — f2fs +does not serialise, so there is no read bandwidth to recover, and the engine's flash stall is +unchanged with the flag and without. What changes is CPU: a 20 GB mapping the kernel still has to +account for costs about 9% of the decode's CPU time on a device under memory pressure, and dropping +it is worth 5-9% of throughput. That measurement is two short cells per variant and is a direction, +not a number. The flag stays off by default on both platforms. + ## Why repack must stay off The streamer rebinds `tensor->data` to a buffer it fills from the file's native byte diff --git a/docs/roadmap.md b/docs/roadmap.md index 56d62c2..951f061 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -55,9 +55,15 @@ What follows from that: base by `file_off % align` so the destination inherits the file's misalignment) does not remove the spill, only its size. Not built: the win is a memcpy, the risk is the memory accounting the whole engine's budget rests on, and it cannot be judged without a device. -- The remaining gap between the engine's effective rate and the drive's is **duty cycle, not - bandwidth**, and is not yet honestly sized: the ceiling itself falls by a third once the device - is hot, so engine and microbench must be measured interleaved at matched entry state. Owed. +- **The gap between the engine's effective rate and the drive's was the model file's own mapping, + on Windows (2026-08-29).** This entry used to call it duty cycle and leave it unsized. It is + not duty cycle: while llama.cpp's section of the gguf is alive, NTFS serialises the streamer's + concurrent unbuffered reads, so four lanes deliver one lane's throughput — which is also the + real reason lanes and threads measured dead on this host, and why the serial path once read at + exactly the one-lane rate with four lanes open. `bmoe-iobench --mmap` reproduces it in one + cell and `--mmap --reopen-lanes` recovers it; `--release-mmap` is the engine's version and is + worth +46% decode on the desktop host, lossless. The same cells are flat on the phone, so this + is a platform defect, not a read-path one, and the phone's remaining gap is still unsized. ## Warm-up diff --git a/examples/android/README.md b/examples/android/README.md index 7614f2f..b6a8162 100644 --- a/examples/android/README.md +++ b/examples/android/README.md @@ -170,6 +170,13 @@ Two worth knowing before you turn them on: - **"Stream row-gathered tables"** (`--row-stream`) serves the token embedding table out of flash instead of RAM. Lossless, and which tables it applies to is read off the model's own graph, so on a model where none qualify it does nothing. See `../../docs/row-gathered-tables.md`. +- **"Release the model mapping"** (`--release-mmap`) hands the model file's mapping back to the + kernel once every weight has been copied into the app's own memory, so it needs **Dense weights** + on Anon or Pinned and is disabled otherwise. Lossless. The mechanism that makes it worth +46% on + a Windows desktop does not exist here — on f2fs the read lanes measure the same either way — but + keeping a 20 GB mapping registered costs a kernel under memory pressure, and dropping it took + ~9% off CPU per token. Two 48-token cells on a device that spreads 20%: a direction, not a + number, which is why it is off by default. - **"Prefer cached experts"** (`--expert-substitute`) steers each routing toward experts already in RAM, so the same number of experts runs but fewer are read from flash. It changes the reply, and past 20% the reply keeps reading well while the model behind it is much worse: judge it on diff --git a/examples/android/app/src/main/java/io/bigmoeonedge/example/AppSettings.kt b/examples/android/app/src/main/java/io/bigmoeonedge/example/AppSettings.kt index 4295ce0..2e4394e 100644 --- a/examples/android/app/src/main/java/io/bigmoeonedge/example/AppSettings.kt +++ b/examples/android/app/src/main/java/io/bigmoeonedge/example/AppSettings.kt @@ -65,6 +65,18 @@ data class AppSettings( // so the only question it raises is whether the reads cost more than the RAM is worth - which // is why it is off until the on-device A/B says otherwise. val rowStream: Boolean = false, + // Hand the model file's mapping back to the kernel once every weight has been rebound onto the + // app's own memory. Needs a dense policy that does that rebinding (Anon or Pinned), which is why + // the switch is disabled under Mmap and Warm — under those the engine looks at its own pointers, + // sees weights still reading the mapping, and stands down anyway. + // + // The mechanism that makes this worth +46% on a Windows desktop (a live mapping serialises the + // streamer's unbuffered reads) does NOT exist here: on f2fs the read lanes measure the same with + // the mapping and without. What it buys on device is CPU — keeping a 20 GB mapping registered + // costs a kernel under memory pressure, and dropping it took ~9% off CPU per token. That is two + // 48-token cells against a device whose cells spread 20%, so it is a direction and not a number, + // and the switch stays off until a 256-token A/B earns it. + val releaseMmap: Boolean = false, // Cache-aware substitution, as a PERCENTAGE of the router's score range (0 = off). Before a // routing is committed, every expert already resident gets its score raised by this fraction of // the range and the top-k is taken again, so a resident expert wins a slot only when it was @@ -187,6 +199,12 @@ data class AppSettings( // discovered by the streamer's capture pass; independent of the cache and of the // dense-weight mode, since what it changes is which tensors that mode applies to. if (rowStream) a += "--row-stream" + // Only the policies that rebind every weight into the app's own memory can leave the + // mapping unreferenced. Sending it under Mmap or Warm is not unsafe — the engine checks + // its own pointers and declines — but it would be a switch that silently does nothing. + if (releaseMmap && (denseWeights == DenseWeights.ANON || denseWeights == DenseWeights.AHWB)) { + a += "--release-mmap" + } // Same cacheOn guard and for the same reason: with no cache there is nothing resident // to substitute toward, so the policy would re-rank against an all-miss mask. if (substitutePct > 0 && cacheOn) a += listOf("--expert-substitute", (substitutePct / 100.0).toString()) @@ -236,6 +254,7 @@ data class AppSettings( .putInt("routeAhead", routeAhead) .putInt("dropColdPct", dropColdPct) .putBoolean("rowStream", rowStream) + .putBoolean("releaseMmap", releaseMmap) .putInt("substitutePct", substitutePct) .putInt("sessionCtx", sessionCtx) .putString("spec", spec).putInt("mtpDraft", mtpDraft).putInt("mtpPMinPct", mtpPMinPct) @@ -379,6 +398,7 @@ data class AppSettings( routeAhead = p.getInt("routeAhead", d.routeAhead), dropColdPct = p.getInt("dropColdPct", d.dropColdPct), rowStream = p.getBoolean("rowStream", d.rowStream), + releaseMmap = p.getBoolean("releaseMmap", d.releaseMmap), substitutePct = p.getInt("substitutePct", d.substitutePct), sessionCtx = p.getInt("sessionCtx", d.sessionCtx), spec = run { diff --git a/examples/android/app/src/main/java/io/bigmoeonedge/example/SettingsScreen.kt b/examples/android/app/src/main/java/io/bigmoeonedge/example/SettingsScreen.kt index b3dc188..bb33f21 100644 --- a/examples/android/app/src/main/java/io/bigmoeonedge/example/SettingsScreen.kt +++ b/examples/android/app/src/main/java/io/bigmoeonedge/example/SettingsScreen.kt @@ -127,6 +127,17 @@ fun SettingsScreen(current: AppSettings, onChange: (AppSettings) -> Unit, onBack enabled = stream, ) { onChange(current.copy(rowStream = it)) } + SwitchRow( + "Release the model mapping", + "Once every weight has been copied into the app's own memory, the model file " + + "does not need to stay mapped. Handing the mapping back frees the kernel " + + "from tracking it, which showed up as less CPU per token. Lossless - the " + + "output is identical. Needs Dense weights on Anon or Pinned.", + current.releaseMmap, + enabled = stream && (current.denseWeights == DenseWeights.ANON || + current.denseWeights == DenseWeights.AHWB), + ) { onChange(current.copy(releaseMmap = it)) } + ExperimentalGroup { IntSetting( "Temporal prefetch (layers)", AppSettings.PREFETCH_CHOICES, current.prefetchLayers, diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 27ab2e9..8e6a6bd 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -36,6 +36,13 @@ add_executable(bmoe_config_test config_test.cpp) target_link_libraries(bmoe_config_test PRIVATE bmoe_core) add_test(NAME config_validate COMMAND bmoe_config_test) +# release_file_mappings: on Windows it must find and drop this process's own view + section of a +# file and leave the placeholders in place; on POSIX it must report unsupported and touch nothing. +add_executable(bmoe_mapping_release_test mapping_release_test.cpp) +target_link_libraries(bmoe_mapping_release_test PRIVATE bmoe_core) +target_include_directories(bmoe_mapping_release_test PRIVATE ${CMAKE_SOURCE_DIR}/core/src/io) +add_test(NAME mapping_release COMMAND bmoe_mapping_release_test) + # The n-gram draft matcher. Pure policy over token ids (no llama.cpp, no model), so the drafting # decision is covered exhaustively here; what it costs is a bench question, not a test. add_executable(bmoe_ngram_test ngram_test.cpp) diff --git a/tests/mapping_release_test.cpp b/tests/mapping_release_test.cpp new file mode 100644 index 0000000..f2aac8f --- /dev/null +++ b/tests/mapping_release_test.cpp @@ -0,0 +1,183 @@ +// release_file_mappings (core/src/io/mapping_release.h): the engine's answer to Windows serialising +// concurrent unbuffered reads on a file that has a live section. This test builds exactly the state +// llama.cpp leaves after load — a view and a section handle on a file, the file handle itself already +// closed — asks the module to release it, and checks both what it must do (drop the view, close the +// section) and what it must leave behind on purpose (a reservation where the view was, an inert +// handle in the section's slot) so that the model's later teardown of "its" mapping is harmless. +// No model, no llama.cpp: a scratch file is enough. On POSIX the same release is a munmap with an +// anonymous placeholder left in the range, checked the same way. +#include "mapping_release.h" + +#include +#include +#include +#include + +#if defined(_WIN32) +#include +#endif + +static int failures = 0; +#define CHECK(cond) \ + do { \ + if (!(cond)) { \ + std::fprintf(stderr, "FAIL %s:%d: %s\n", __FILE__, __LINE__, #cond); \ + ++failures; \ + } \ + } while (0) + +#if defined(_WIN32) + +namespace { + +std::string scratch_path() { + char dir[MAX_PATH]; + const DWORD n = GetTempPathA(MAX_PATH, dir); + if (n == 0 || n >= MAX_PATH) return {}; + char path[MAX_PATH]; + if (!GetTempFileNameA(dir, "bmr", 0, path)) return {}; + return path; +} + +bool fill(const std::string & path, size_t bytes) { + FILE * f = std::fopen(path.c_str(), "wb"); + if (!f) return false; + std::vector chunk(1 << 16, 'x'); + for (size_t done = 0; done < bytes; done += chunk.size()) + if (std::fwrite(chunk.data(), 1, chunk.size(), f) != chunk.size()) { + std::fclose(f); + return false; + } + std::fclose(f); + return true; +} + +// The post-load state of llama_mmap on Windows: view + section handle alive, file handle closed. +struct LlamaLikeMapping { + HANDLE section = nullptr; + void * view = nullptr; + bool open(const std::string & path) { + HANDLE h = CreateFileA(path.c_str(), GENERIC_READ, FILE_SHARE_READ, nullptr, OPEN_EXISTING, + FILE_ATTRIBUTE_NORMAL, nullptr); + if (h == INVALID_HANDLE_VALUE) return false; + section = CreateFileMappingA(h, nullptr, PAGE_READONLY, 0, 0, nullptr); + CloseHandle(h); + if (!section) return false; + view = MapViewOfFile(section, FILE_MAP_READ, 0, 0, 0); + return view != nullptr; + } +}; + +DWORD state_at(void * p) { + MEMORY_BASIC_INFORMATION mbi; + return VirtualQuery(p, &mbi, sizeof(mbi)) == sizeof(mbi) ? mbi.State : 0; +} + +} // namespace + +int main() { + const size_t size = 8u << 20; + const std::string path = scratch_path(); + CHECK(!path.empty()); + CHECK(fill(path, size)); + + LlamaLikeMapping m; + CHECK(m.open(path)); + CHECK(state_at(m.view) == MEM_COMMIT); + CHECK(((const volatile char *) m.view)[0] == 'x'); // the view is real before we start + + bmoe::pio::MappingPlaceholders ph; + bmoe::pio::MappingReleaseReport r = bmoe::pio::release_file_mappings({path}, &ph); + CHECK(r.supported); + CHECK(r.error.empty()); + CHECK(r.views_unmapped == 1); + CHECK(r.bytes == size); + CHECK(r.sections_closed == 1); + CHECK(r.plugs_missed == 0); + CHECK(!ph.empty()); + + // The view's range is now a reservation, not free and not mapped: nothing can be allocated + // there, so the model's late UnmapViewOfFile on this base finds no live allocation. + CHECK(state_at(m.view) == MEM_RESERVE); + + // The section's handle value is now an inert event: waiting on it times out instead of failing + // as it would on a section, and the model's late CloseHandle will close this plug, nothing else. + CHECK(WaitForSingleObject(m.section, 0) == WAIT_TIMEOUT); + + // What llama_mmap's destructor does, in its order: both must be harmless. + CHECK(UnmapViewOfFile(m.view) == FALSE); + CHECK(CloseHandle(m.section) == TRUE); + + // Idempotent: nothing of this file is left to find. + bmoe::pio::MappingPlaceholders ph2; + r = bmoe::pio::release_file_mappings({path}, &ph2); + CHECK(r.supported); + CHECK(r.views_unmapped == 0); + CHECK(r.sections_closed == 0); + CHECK(ph2.empty()); + + // Releasing the placeholders frees the range. + ph.release(); + CHECK(ph.empty()); + CHECK(state_at(m.view) == MEM_FREE); + + DeleteFileA(path.c_str()); + if (failures) { + std::fprintf(stderr, "%d check(s) failed\n", failures); + return 1; + } + std::printf("mapping_release: ok\n"); + return 0; +} + +#else + +#include +#include +#include + +int main() { + // A file this process has mapped, the way llama.cpp maps a gguf: MAP_PRIVATE, read-only. + char path[] = "/tmp/bmoe-mapping-release-XXXXXX"; + const int fd = mkstemp(path); + CHECK(fd >= 0); + const size_t size = 8u << 20; + CHECK(ftruncate(fd, (off_t) size) == 0); + void * view = mmap(nullptr, size, PROT_READ, MAP_PRIVATE, fd, 0); + CHECK(view != MAP_FAILED); + close(fd); + + bmoe::pio::MappingPlaceholders ph; + bmoe::pio::MappingReleaseReport r = bmoe::pio::release_file_mappings({path}, &ph); + CHECK(r.supported); + CHECK(r.error.empty()); + CHECK(r.views_unmapped == 1); + CHECK(r.bytes == size); + CHECK(r.sections_closed == 0); + CHECK(!ph.empty()); + + // The range is held by an inaccessible placeholder: mapping something else there with + // MAP_FIXED_NOREPLACE must fail, and the model's own late munmap of the range is harmless. +#ifdef MAP_FIXED_NOREPLACE + void * clash = mmap(view, size, PROT_READ, MAP_PRIVATE | MAP_ANONYMOUS | MAP_FIXED_NOREPLACE, -1, 0); + CHECK(clash == MAP_FAILED); +#endif + CHECK(munmap(view, size) == 0); + + // Idempotent, and forgetting the placeholders touches nothing. + bmoe::pio::MappingPlaceholders ph2; + r = bmoe::pio::release_file_mappings({path}, &ph2); + CHECK(r.supported && r.views_unmapped == 0); + ph.release(); + CHECK(ph.empty()); + + unlink(path); + if (failures) { + std::fprintf(stderr, "%d check(s) failed\n", failures); + return 1; + } + std::printf("mapping_release: ok\n"); + return 0; +} + +#endif diff --git a/tools/iobench.cpp b/tools/iobench.cpp index f646b86..28279a1 100644 --- a/tools/iobench.cpp +++ b/tools/iobench.cpp @@ -31,6 +31,14 @@ // [--compute-load N] [--scatter N] #include "file_reader.h" #include "platform_io.h" +#if defined(_WIN32) +#include +#else +#include +#include +#include +#endif +#include "platform_io.h" #include #include @@ -64,6 +72,58 @@ struct LaneResult { long long busy_ns = 0; }; +// A read-only mapping of the whole file, held open while the lanes read, which is the state +// llama.cpp leaves a gguf in for the model's lifetime. It exists because that state is not free: +// on Windows, while a section of a file is alive, NTFS serialises concurrent unbuffered reads on +// it, and four lanes deliver one lane's throughput (2400-2660 MiB/s without, 895-930 with, at +// 576 KiB). Nothing is ever read through the mapping here; only its existence is the variable. +struct FileMapping { + void * view = nullptr; + uint64_t len = 0; +#if defined(_WIN32) + HANDLE file = INVALID_HANDLE_VALUE, section = nullptr; +#else + int fd = -1; +#endif + + bool open(const std::string & path) { +#if defined(_WIN32) + file = CreateFileA(path.c_str(), GENERIC_READ, FILE_SHARE_READ, nullptr, OPEN_EXISTING, FILE_ATTRIBUTE_NORMAL, + nullptr); + if (file == INVALID_HANDLE_VALUE) return false; + LARGE_INTEGER sz; + if (!GetFileSizeEx(file, &sz)) return false; + len = (uint64_t) sz.QuadPart; + section = CreateFileMappingA(file, nullptr, PAGE_READONLY, 0, 0, nullptr); + if (!section) return false; + view = MapViewOfFile(section, FILE_MAP_READ, 0, 0, 0); + return view != nullptr; +#else + fd = ::open(path.c_str(), O_RDONLY); + if (fd < 0) return false; + len = (uint64_t) lseek(fd, 0, SEEK_END); + view = mmap(nullptr, (size_t) len, PROT_READ, MAP_PRIVATE, fd, 0); + return view != MAP_FAILED; +#endif + } + + void close() { +#if defined(_WIN32) + if (view) UnmapViewOfFile(view); + if (section) CloseHandle(section); + if (file != INVALID_HANDLE_VALUE) CloseHandle(file); + view = nullptr; + section = nullptr; + file = INVALID_HANDLE_VALUE; +#else + if (view && view != MAP_FAILED) munmap(view, (size_t) len); + if (fd >= 0) ::close(fd); + view = nullptr; + fd = -1; +#endif + } +}; + // One lane: random-offset logical reads of `slice` bytes until the deadline, each issued as // `scatter` preads of slice/scatter bytes at independent offsets (see the header comment). Offsets // are block-aligned and kept a slice away from EOF so every read is a full window (the @@ -75,28 +135,48 @@ void lane_worker(bmoe::FileReader * r, int scatter, size_t align, uint64_t fsize, + uint64_t range_bytes, + bool fresh, clock_t_::time_point deadline, LaneResult * out) { // Equal-total-bytes split, each piece its own aligned window — otherwise scatter rows would // compare different traffic volumes, not different layouts. const size_t piece = ((slice / (size_t) scatter) + align - 1) & ~(align - 1); - void * dst = bmoe::pio::alloc_aligned(align, piece); + // --fresh: reserved address space that is decommitted and recommitted before every read, so + // each copy out of the bounce lands on pages the process has never touched -- what the expert + // cache's reserve/commit/evict does to every slice it admits, and a cost a plain reused + // destination buffer hides entirely. + void * dst = fresh ? bmoe::pio::vm_reserve(piece) : bmoe::pio::alloc_aligned(align, piece); if (!dst) return; - const uint64_t span = (fsize > piece * 2) ? (fsize - piece * 2) : 0; + const uint64_t whole = (fsize > piece * 2) ? (fsize - piece * 2) : 0; + // --range-mb: confine the offsets to one region in the middle of the file, the way one layer's + // routed experts sit inside one tensor rather than spread over the whole file. + const uint64_t span = (range_bytes && range_bytes < whole) ? range_bytes : whole; + const uint64_t base = (span < whole) ? ((whole - span) / 2) & ~(uint64_t) (align - 1) : 0; if (span == 0) { - bmoe::pio::aligned_free(dst); + if (fresh) + bmoe::pio::vm_release(dst, piece); + else + bmoe::pio::aligned_free(dst); return; } Lcg rng((uint64_t) lane + 1); LaneResult acc; while (clock_t_::now() < deadline) { for (int s = 0; s < scatter; ++s) { - const uint64_t off = (rng.next() % span) & ~(uint64_t) (align - 1); + const uint64_t off = base + ((rng.next() % span) & ~(uint64_t) (align - 1)); + if (fresh) { + bmoe::pio::vm_evict(dst, piece); + if (!bmoe::pio::vm_commit(dst, piece)) break; + } const auto t0 = clock_t_::now(); const long long got = r->read(lane, dst, off, piece); const auto t1 = clock_t_::now(); if (got < 0) { - bmoe::pio::aligned_free(dst); + if (fresh) + bmoe::pio::vm_release(dst, piece); + else + bmoe::pio::aligned_free(dst); *out = acc; return; } @@ -105,7 +185,10 @@ void lane_worker(bmoe::FileReader * r, acc.busy_ns += std::chrono::duration_cast(t1 - t0).count(); } } - bmoe::pio::aligned_free(dst); + if (fresh) + bmoe::pio::vm_release(dst, piece); + else + bmoe::pio::aligned_free(dst); *out = acc; } @@ -134,6 +217,10 @@ bool run_row(const std::string & path, bool direct, double seconds, int load, + uint64_t range_bytes, + bool fresh, + bool with_mapping, + bool reopen_lanes, double * mibs_out) { bmoe::FileReader r; // Ask the OS rather than assuming 4096: alignment is exactly the variable this tool exists to @@ -144,6 +231,22 @@ bool run_row(const std::string & path, std::fprintf(stderr, "open failed (lanes=%d)\n", lanes); return false; } + // The mapping is opened AFTER the lanes and, with --reopen-lanes, dropped before they are + // reopened: on Windows a lane opened while a section was alive keeps serialising against it + // even once the section is gone, so the two halves of that behaviour are separable here. + FileMapping fm; + if (with_mapping && !fm.open(path)) { + std::fprintf(stderr, "mapping the model failed\n"); + r.close(); + return false; + } + if (reopen_lanes) { + fm.close(); + if (!r.reopen()) { + std::fprintf(stderr, "reopening the lanes failed\n"); + return false; + } + } std::vector res((size_t) lanes); std::vector th; th.reserve((size_t) lanes); @@ -157,7 +260,8 @@ bool run_row(const std::string & path, loaders.emplace_back(load_worker, deadline, &stop); for (int i = 0; i < lanes; ++i) - th.emplace_back(lane_worker, &r, i, slice, scatter, align, r.file_size(), deadline, &res[(size_t) i]); + th.emplace_back(lane_worker, &r, i, slice, scatter, align, r.file_size(), range_bytes, fresh, deadline, + &res[(size_t) i]); for (auto & t : th) t.join(); const double wall_s = std::chrono::duration(clock_t_::now() - t0).count(); @@ -180,6 +284,7 @@ bool run_row(const std::string & path, r.direct() ? "direct" : "BUFFERED"); std::fflush(stdout); r.close(); + fm.close(); if (mibs_out) *mibs_out = mibs; return true; } @@ -192,7 +297,15 @@ void usage(const char * a0) { " equal pieces at independent offsets — same bytes, scattered layout\n" " --buffered drop O_DIRECT, to see what the page cache contributes\n" " --compute-load N CPU-burning threads alongside the lanes (default 0), to read\n" - " under the contention the streamer actually faces\n", + " under the contention the streamer actually faces\n" + " --range-mb confine the random offsets to one N MiB region of the file\n" + " --fresh commit the destination pages afresh before every read, as the\n" + " expert cache does for every slice it admits\n" + " --mmap hold a read-only mapping of the file open while the lanes read,\n" + " as llama.cpp does for a loaded model (Windows: this alone\n" + " serialises unbuffered reads -- see mapping_release.h)\n" + " --reopen-lanes with --mmap: drop the mapping and reopen the lanes before\n" + " reading, which is what the engine does under --release-mmap\n", a0); } @@ -206,6 +319,10 @@ int main(int argc, char ** argv) { double seconds = 5.0; bool direct = true; int load = 0; + uint64_t range_bytes = 0; + bool fresh = false; + bool with_mapping = false; + bool reopen_lanes = false; for (int i = 1; i < argc; ++i) { const std::string a = argv[i]; @@ -230,6 +347,14 @@ int main(int argc, char ** argv) { direct = false; else if (a == "--compute-load") load = std::atoi(next("--compute-load")); + else if (a == "--range-mb") + range_bytes = (uint64_t) std::atoll(next("--range-mb")) << 20; + else if (a == "--fresh") + fresh = true; + else if (a == "--mmap") + with_mapping = true; + else if (a == "--reopen-lanes") + reopen_lanes = true; else { usage(argv[0]); return 2; @@ -268,7 +393,9 @@ int main(int argc, char ** argv) { for (int L : lanes) { if (L < 1) continue; double mibs = 0.0; - if (!run_row(model, L, slice, scatter, direct, seconds, load, &mibs)) return 1; + if (!run_row(model, L, slice, scatter, direct, seconds, load, range_bytes, fresh, with_mapping, reopen_lanes, + &mibs)) + return 1; if (mibs > best) { best = mibs; best_lanes = L;