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;