BigMoeOnEdge/tools/iobench.cpp
Raffaele 47924565c1
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.
2026-09-07 20:43:02 +02:00

406 lines
17 KiB
C++

// bmoe-iobench — how far does flash read bandwidth scale with concurrent lanes?
//
// The streamer's "queue depth" is exactly its lane count: every lane issues one blocking pread
// at a time, because the async submission APIs (io_uring, kernel AIO) are not reachable from an
// Android app. So the only way to keep more reads in flight is more lanes, and the engine caps
// that at MoeStreamConfig::io_threads_max. This measures whether that cap sits below the
// device's ceiling — the question that has to be answered BEFORE raising it, since a cap that is
// already at the hardware limit costs memory (one fd + one bounce buffer per lane) for nothing.
//
// It deliberately drives bmoe::FileReader, the same read path the engine uses, rather than a
// hand-rolled pread loop: the alignment, the bounce buffer and the O_DIRECT verification/fallback
// are part of what is being measured. Nothing here links the engine or llama.cpp — the I/O layer
// stands alone, so this builds in seconds and cannot perturb the streamer.
//
// Bandwidth is judged against the ALIGNED WINDOW the drive actually served (FileReader::read's
// return), not the bytes requested — the difference is real device traffic.
//
// `--compute-load N` runs N CPU-burning threads alongside the lanes. Without it the sweep measures
// the drive on an idle CPU, which is not the condition the streamer reads under: the engine reads
// while ggml's compute threads spin. If per-lane bandwidth falls once the CPU is busy, the lanes are
// starved of scheduler time and the bottleneck is contention, not the flash.
//
// `--scatter N` splits every logical read into N preads of slice/N bytes at independent random
// offsets — the same total bytes, delivered as N small windows instead of one contiguous one. This
// is the expert-layout question: today one routed expert costs one pread per projection at offsets
// an expert-stride apart (scatter ~3), while a contiguous per-expert sidecar would serve the same
// bytes in a single window (scatter 1). If bandwidth at scatter 1 is not materially above scatter
// 3 for the same slice, the sidecar has no prize and must not be built.
//
// bmoe-iobench --model M.gguf [--lanes 1,2,4,8,16] [--slice-kb 4096] [--seconds 5] [--buffered]
// [--compute-load N] [--scatter N]
#include "file_reader.h"
#include "platform_io.h"
#if defined(_WIN32)
#include <windows.h>
#else
#include <fcntl.h>
#include <sys/mman.h>
#include <unistd.h>
#endif
#include "platform_io.h"
#include <atomic>
#include <chrono>
#include <cmath>
#include <cstdio>
#include <cstdlib>
#include <cstring>
#include <string>
#include <thread>
#include <vector>
namespace {
using clock_t_ = std::chrono::steady_clock;
// Deterministic per-lane offset stream. A real RNG would make two lane counts read different
// regions and turn a bandwidth comparison into a luck comparison; this way lane L always walks
// the same sequence whatever the lane count, so the only variable between rows is concurrency.
struct Lcg {
uint64_t s;
explicit Lcg(uint64_t seed) : s(seed * 6364136223846793005ULL + 1442695040888963407ULL) {}
uint64_t next() {
s = s * 6364136223846793005ULL + 1442695040888963407ULL;
return s >> 16;
}
};
struct LaneResult {
long long window_bytes = 0; // what the drive served (aligned)
long long reads = 0;
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
// sub-alignment tail would otherwise take FileReader's buffered fallback and quietly measure the
// page cache instead of the drive).
void lane_worker(bmoe::FileReader * r,
int lane,
size_t slice,
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);
// --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 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) {
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 = 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) {
if (fresh)
bmoe::pio::vm_release(dst, piece);
else
bmoe::pio::aligned_free(dst);
*out = acc;
return;
}
acc.window_bytes += got;
acc.reads++;
acc.busy_ns += std::chrono::duration_cast<std::chrono::nanoseconds>(t1 - t0).count();
}
}
if (fresh)
bmoe::pio::vm_release(dst, piece);
else
bmoe::pio::aligned_free(dst);
*out = acc;
}
// A stand-in for the compute threads the reads actually compete with. It only has to keep a core
// busy in userspace; the arithmetic is irrelevant, but it must not be optimised away, hence the
// volatile sink.
void load_worker(clock_t_::time_point deadline, std::atomic<bool> * stop) {
volatile double sink = 0.0;
double x = 1.000001;
while (!stop->load(std::memory_order_relaxed) && clock_t_::now() < deadline) {
for (int i = 0; i < 4096; ++i)
x = x * 1.0000001 + 1e-9;
sink = sink + x;
if (x > 1e6) x = 1.000001;
}
(void) sink;
}
// One row of the sweep: open a reader with `lanes` lanes, run them all for `seconds`, report the
// aggregate. The reader is opened and closed per row so each row pays its own warm-up and no lane
// inherits another row's fd state.
bool run_row(const std::string & path,
int lanes,
size_t slice,
int scatter,
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
// characterise, so a device with a 16 KiB page must be measured at its own page size, not ours.
const size_t align = bmoe::pio::vm_page();
const size_t bounce_cap = slice + 2 * align; // mirrors what the streamer asks for
if (!r.open(path, lanes, direct, align, bounce_cap)) {
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<LaneResult> res((size_t) lanes);
std::vector<std::thread> th;
th.reserve((size_t) lanes);
const auto t0 = clock_t_::now();
const auto deadline = t0 + std::chrono::milliseconds((long long) (seconds * 1000.0));
std::atomic<bool> stop{false};
std::vector<std::thread> loaders;
loaders.reserve((size_t) (load > 0 ? load : 0));
for (int i = 0; i < load; ++i)
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(), range_bytes, fresh, deadline,
&res[(size_t) i]);
for (auto & t : th)
t.join();
const double wall_s = std::chrono::duration<double>(clock_t_::now() - t0).count();
stop.store(true); // lanes are done; do not let the load run past the measured window
for (auto & t : loaders)
t.join();
long long bytes = 0, reads = 0, busy_ns = 0;
for (const auto & x : res) {
bytes += x.window_bytes;
reads += x.reads;
busy_ns += x.busy_ns;
}
const double mib = (double) bytes / (1024.0 * 1024.0);
const double mibs = mib / wall_s;
// Mean latency is per-read wall inside a lane; with N lanes busy, N*wall is the service budget,
// so latency rising while throughput plateaus is the saturation signature.
const double lat_ms = reads ? (double) busy_ns / (double) reads / 1e6 : 0.0;
std::printf("%6d %12.1f %10.2f %9lld %11.2f %10s\n", lanes, mibs, lat_ms, reads, mib,
r.direct() ? "direct" : "BUFFERED");
std::fflush(stdout);
r.close();
fm.close();
if (mibs_out) *mibs_out = mibs;
return true;
}
void usage(const char * a0) {
std::fprintf(stderr,
"usage: %s --model PATH [--lanes 1,2,4,8,16] [--slice-kb N] [--seconds S] [--buffered]\n"
" --slice-kb KiB per logical read, default 4096 (= 4 MiB, a typical expert slice)\n"
" --scatter preads per logical read (default 1): N splits the slice into N\n"
" 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"
" --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);
}
} // namespace
int main(int argc, char ** argv) {
std::string model;
std::string lanes_spec = "1,2,4,8,12,16,24,32";
size_t slice_kb = 4096;
int scatter = 1;
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];
auto next = [&](const char * what) -> const char * {
if (i + 1 >= argc) {
std::fprintf(stderr, "%s needs a value\n", what);
std::exit(2);
}
return argv[++i];
};
if (a == "--model" || a == "-m")
model = next("--model");
else if (a == "--lanes")
lanes_spec = next("--lanes");
else if (a == "--slice-kb")
slice_kb = (size_t) std::atoll(next("--slice-kb"));
else if (a == "--scatter")
scatter = std::atoi(next("--scatter"));
else if (a == "--seconds")
seconds = std::atof(next("--seconds"));
else if (a == "--buffered")
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;
}
}
if (model.empty()) {
usage(argv[0]);
return 2;
}
std::vector<int> lanes;
for (size_t p = 0; p < lanes_spec.size();) {
const size_t c = lanes_spec.find(',', p);
const std::string tok = lanes_spec.substr(p, c == std::string::npos ? std::string::npos : c - p);
if (!tok.empty()) lanes.push_back(std::atoi(tok.c_str()));
if (c == std::string::npos) break;
p = c + 1;
}
if (lanes.empty()) {
std::fprintf(stderr, "no lane counts parsed\n");
return 2;
}
if (scatter < 1) {
std::fprintf(stderr, "--scatter must be >= 1\n");
return 2;
}
const size_t slice = slice_kb * 1024;
std::printf("model %s\n", model.c_str());
std::printf("read size %zu KiB x scatter %d, %.1f s per row, compute load %d threads\n\n", slice_kb, scatter,
seconds, load);
std::printf("%6s %12s %10s %9s %11s %10s\n", "lanes", "MiB/s", "lat_ms", "reads", "read_MiB", "mode");
double best = 0.0;
int best_lanes = 0;
for (int L : lanes) {
if (L < 1) continue;
double mibs = 0.0;
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;
}
}
std::printf("\npeak %.1f MiB/s at %d lanes\n", best, best_lanes);
return 0;
}