// SDQ expert pager 實作:見 sdq_pager.h 的設計說明 #include "sdq_pager.h" #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include namespace sdq { // ============================================================ 小工具 uint64_t now_us() { struct timespec ts; clock_gettime(CLOCK_MONOTONIC, &ts); return (uint64_t) ts.tv_sec * 1000000ull + (uint64_t) ts.tv_nsec / 1000ull; } size_t process_rss_bytes() { std::ifstream f("/proc/self/statm"); size_t total = 0, resident = 0; if (f >> total >> resident) { return resident * (size_t) sysconf(_SC_PAGESIZE); } return 0; } size_t physical_ram_bytes() { std::ifstream f("/proc/meminfo"); std::string key; size_t value = 0; std::string unit; while (f >> key >> value >> unit) { if (key == "MemTotal:") { return value * 1024; } } return 0; } static size_t env_size(const char * name, size_t dflt) { const char * v = getenv(name); if (!v || !*v) { return dflt; } return (size_t) strtoull(v, nullptr, 10); } static bool env_flag(const char * name, bool dflt) { const char * v = getenv(name); if (!v || !*v) { return dflt; } return !(v[0] == '0' && v[1] == '\0'); } // ============================================================ GGUF 檔頭解析 namespace { struct gguf_reader { const uint8_t * p = nullptr; size_t size = 0; size_t o = 0; uint8_t u8() { return p[o++]; } uint32_t u32() { uint32_t v; memcpy(&v, p + o, 4); o += 4; return v; } uint64_t u64() { uint64_t v; memcpy(&v, p + o, 8); o += 8; return v; } int64_t i64() { int64_t v; memcpy(&v, p + o, 8); o += 8; return v; } float f32() { float v; memcpy(&v, p + o, 4); o += 4; return v; } std::string str() { uint64_t n = u64(); std::string s((const char *) p + o, (size_t) n); o += n; return s; } void skip(uint64_t n) { o += (size_t) n; } }; struct kv_value { enum kind { NONE, INT, REAL, STR, ARR } k = NONE; int64_t i = 0; double f = 0; std::string s; std::vector arr; }; struct tensor_info { std::string name; std::vector dims; ggml_type type = GGML_TYPE_COUNT; uint64_t offset = 0; }; size_t ggml_type_size_of(ggml_type t) { return ggml_row_size(t, 1); } } // namespace // ============================================================ pager pager & pager::instance() { static pager p; return p; } void pager::load_gguf_header(const std::string & path) { const int fd = open(path.c_str(), O_RDONLY); if (fd < 0) { throw std::runtime_error("無法開啟 GGUF:" + path); } struct stat st; fstat(fd, &st); // 只讀檔頭:metadata + tensor 資訊大約 11 MB const size_t want = (size_t) std::min(st.st_size, 64ull << 20); std::vector buf(want); ssize_t got = 0; while (got < (ssize_t) want) { ssize_t n = pread(fd, buf.data() + got, want - got, got); if (n <= 0) { break; } got += n; } buf.resize((size_t) std::max(got, 0)); gguf_reader r{buf.data(), buf.size(), 0}; if (memcmp(buf.data(), "GGUF", 4) != 0) { close(fd); throw std::runtime_error("不是 GGUF 檔:" + path); } r.o = 4; const uint32_t version = r.u32(); if (version < 2 || version > 3) { close(fd); throw std::runtime_error("不支援的 GGUF 版本"); } const uint64_t n_tensors = r.u64(); const uint64_t n_kv = r.u64(); std::map kv; for (uint64_t i = 0; i < n_kv; i++) { std::string key = r.str(); const uint32_t t = r.u32(); kv_value v; switch (t) { case 0: v.k = kv_value::INT; v.i = (int64_t) r.u32(); v.f = (double) v.i; break; case 1: v.k = kv_value::INT; v.i = (int64_t) (int8_t) r.u32(); v.f = (double) v.i; break; case 2: v.k = kv_value::INT; v.i = r.u32(); v.f = (double) v.i; break; case 3: v.k = kv_value::INT; v.i = (int16_t) r.u32(); v.f = (double) v.i; break; case 4: v.k = kv_value::INT; v.i = (int64_t) (uint32_t) r.u32(); v.f = (double) v.i; break; case 5: v.k = kv_value::INT; v.i = (int32_t) r.u32(); v.f = (double) v.i; break; case 6: v.k = kv_value::REAL; v.f = r.f32(); v.i = (int64_t) v.f; break; case 7: v.k = kv_value::INT; v.i = r.u8(); v.f = (double) v.i; break; // BOOL = 1 byte case 8: v.k = kv_value::STR; v.s = r.str(); break; case 9: { v.k = kv_value::ARR; const uint32_t et = r.u32(); const uint64_t n = r.u64(); for (uint64_t j = 0; j < n; j++) { switch (et) { case 2: v.arr.push_back(r.u32()); break; case 3: v.arr.push_back((int16_t) r.u32()); break; case 4: v.arr.push_back((uint32_t) r.u32()); break; case 5: v.arr.push_back((int32_t) r.u32()); break; case 6: v.arr.push_back((int64_t) r.f32()); break; case 7: v.arr.push_back(r.u32()); break; case 10: v.arr.push_back((int64_t) r.u64()); break; case 11: v.arr.push_back(r.i64()); break; case 12: v.arr.push_back((int64_t) (double) 0); break; // 極少見 case 8: r.skip(r.u64()); v.arr.push_back(0); break; default: break; } } } break; case 10: v.k = kv_value::INT; v.i = (int64_t) r.u64(); v.f = (double) v.i; break; case 11: v.k = kv_value::INT; v.i = r.i64(); v.f = (double) v.i; break; case 12: v.k = kv_value::REAL; { double d; memcpy(&d, buf.data() + r.o, 8); r.o += 8; v.f = d; v.i = (int64_t) d; } break; default: close(fd); throw std::runtime_error("GGUF kv 型別不支援: " + std::to_string(t)); } kv[key] = v; } std::vector tensors; tensors.reserve(n_tensors); for (uint64_t i = 0; i < n_tensors; i++) { tensor_info ti; ti.name = r.str(); const uint32_t nd = r.u32(); for (uint32_t d = 0; d < nd; d++) { ti.dims.push_back((int64_t) r.u64()); } ti.type = (ggml_type) r.u32(); ti.offset = r.u64(); tensors.push_back(std::move(ti)); } const uint64_t align = kv.count("general.alignment") ? (uint64_t) kv["general.alignment"].i : 32; const uint64_t data_start = (r.o + align - 1) / align * align; close(fd); // ---- hparams auto gi = [&](const std::string & k, int64_t dflt) -> int64_t { auto it = kv.find(k); return it == kv.end() ? dflt : it->second.i; }; info_.arch = kv.count("general.architecture") ? kv["general.architecture"].s : ""; info_.name = kv.count("general.name") ? kv["general.name"].s : ""; const std::string pfx = info_.arch + "."; info_.n_layer = gi(pfx + "block_count", 0); info_.n_expert = gi(pfx + "expert_count", 0); info_.n_expert_used= gi(pfx + "expert_used_count", 0); info_.n_embd = gi(pfx + "embedding_length", 0); info_.n_ff_exp = gi(pfx + "expert_feed_forward_length", 0); info_.n_ff_shexp = gi(pfx + "expert_shared_feed_forward_length", 0); info_.n_head = gi(pfx + "attention.head_count", 0); info_.n_head_kv = gi(pfx + "attention.head_count_kv", 0); info_.n_embd_head_k= gi(pfx + "attention.key_length", 0); info_.n_embd_head_v= gi(pfx + "attention.value_length", info_.n_embd_head_k); info_.context_length = gi(pfx + "context_length", 0); info_.full_attention_interval = gi(pfx + "full_attention_interval", 4); if (info_.n_expert <= 0 || info_.n_layer <= 0) { throw std::runtime_error("這不是 MoE 模型(expert_count=" + std::to_string(info_.n_expert) + ")"); } // ---- 頁面表:每個 (layer, expert) 的三段 layers_.assign((size_t) info_.n_layer, {}); std::map by_name; for (const auto & ti : tensors) { by_name[ti.name] = &ti; } size_t max_expert_bytes = 0; for (int64_t il = 0; il < info_.n_layer; il++) { const std::string b = "blk." + std::to_string(il) + "."; layer_desc & L = layers_[(size_t) il]; L.experts.assign((size_t) info_.n_expert, {}); auto get = [&](const std::string & suffix) -> const tensor_info * { auto it = by_name.find(b + suffix); return it == by_name.end() ? nullptr : it->second; }; const tensor_info * t_gate = get("ffn_gate_exps.weight"); const tensor_info * t_up = get("ffn_up_exps.weight"); const tensor_info * t_gu = get("ffn_gate_up_exps.weight"); const tensor_info * t_down = get("ffn_down_exps.weight"); if (!t_down || (t_gate && !t_up && !t_gu)) { continue; // 這層不是 MoE(理論上不會發生) } for (int64_t e = 0; e < info_.n_expert; e++) { expert_desc & ed = L.experts[(size_t) e]; if (t_gate || t_gu) { const tensor_info * t = t_gate ? t_gate : t_gu; const uint64_t nb = (uint64_t) ggml_row_size(t->type, t->dims[0]) * (uint64_t) t->dims[1]; ed.seg[part::GATE].off = data_start + t->offset + (uint64_t) e * nb; ed.seg[part::GATE].size = nb; ed.ne0 = t->dims[0]; ed.ne1 = t->dims[1]; ed.type[part::GATE] = t->type; } if (t_up) { const uint64_t nb = (uint64_t) ggml_row_size(t_up->type, t_up->dims[0]) * (uint64_t) t_up->dims[1]; ed.seg[part::UP].off = data_start + t_up->offset + (uint64_t) e * nb; ed.seg[part::UP].size = nb; ed.type[part::UP] = t_up->type; // gate/up 都是 Q4_K,型別一定要各自填 if (!t_gate) { ed.ne0 = t_up->dims[0]; ed.ne1 = t_up->dims[1]; ed.type[part::GATE] = t_up->type; } } { const uint64_t nb = (uint64_t) ggml_row_size(t_down->type, t_down->dims[0]) * (uint64_t) t_down->dims[1]; ed.seg[part::DOWN].off = data_start + t_down->offset + (uint64_t) e * nb; ed.seg[part::DOWN].size = nb; ed.type[part::DOWN] = t_down->type; ed.down_ne0 = t_down->dims[0]; ed.down_ne1 = t_down->dims[1]; } ed.fused_gate_up = (t_gate == nullptr && t_gu != nullptr); ed.bytes = 0; int64_t acc = 0; for (int p = 0; p < 3; p++) { if (ed.seg[p].size == 0) { continue; } acc = (acc + 63) & ~(int64_t) 63; // 每段 64B 對齊 ed.slot_off[p] = acc; acc += (int64_t) ed.seg[p].size; } ed.bytes = acc; max_expert_bytes = std::max(max_expert_bytes, (size_t) ed.bytes); } L.bytes_per_expert = L.experts.empty() ? 0 : L.experts[0].bytes; } slot_bytes_ = (max_expert_bytes + 4095) & ~(size_t) 4095; if (slot_bytes_ < (256ull << 10)) { slot_bytes_ = 256ull << 10; } // ---- 大小兩種槽位 ---------------------------------------------------- // 這個模型有 37 層的 expert 是 1900544 B(Q4_K + Q5_K),但層 34/38/39 是 // 2039808 B(down 換成 Q6_K)。若所有槽位都照最大的配,arena 有 6.8% 被浪費 // ——而 arena 就是 8 GB 預算裡最寶貴的資源(實測 +5% 槽位 ≈ +9% tok/s)。 // // 解法:槽位分成兩區。「小槽位」給 37 層用,「大槽位」只給那 3 層用, // 大槽位的數量按「那 3 層佔請求的比例」分配(3/40 = 7.5%)。 // 選「大多數層」當小的:用眾數而不是最大值,這樣即使某個模型只有一種大小, // 也不會誤判。 { std::map hist; for (const auto & L : layers_) { hist[(size_t) L.bytes_per_expert]++; } size_t small = 0; int best_n = -1; for (const auto & kv : hist) { if (kv.second > best_n) { best_n = kv.second; small = kv.first; } } const size_t mx = (size_t) (layers_.empty() ? 0 : layers_[0].bytes_per_expert); size_t big = 0; for (const auto & kv : hist) { (void) kv; } big = 0; for (const auto & L : layers_) { if ((size_t) L.bytes_per_expert > small) { big = std::max(big, (size_t) L.bytes_per_expert); } } small_bytes_ = (small + 4095) & ~(size_t) 4095; big_bytes_ = big ? ((big + 4095) & ~(size_t) 4095) : 0; if (big_bytes_ && big_bytes_ <= small_bytes_) { big_bytes_ = 0; // 其實一樣大,別分兩區 } slot_bytes_ = big_bytes_ ? big_bytes_ : small_bytes_; int n_big_layers = 0; for (auto & L : layers_) { L.needs_big_slot = big_bytes_ && (size_t) L.bytes_per_expert > small_bytes_; if (L.needs_big_slot) { n_big_layers++; } } (void) mx; fprintf(stderr, "[sdq] 槽位大小:小的 %zu B(%zu 層)/大的 %zu B(%d 層)\n", small_bytes_, (int) layers_.size() - n_big_layers, big_bytes_, n_big_layers); } } bool pager::need_big_slot(int il, int ie) const { if (!big_bytes_) { return false; } const layer_desc * L = layer(il); if (!L) { return false; } if (L->needs_big_slot) { return true; } return ie >= 0 && ie < (int) L->experts.size() && (size_t) L->experts[(size_t) ie].bytes > small_bytes_; } void pager::alloc_arena() { const size_t slot = cfg_.slot_bytes ? cfg_.slot_bytes : slot_bytes_; slot_bytes_ = slot; // arena 大小決定:RAM 是「總預算」,不是「還有多少可用」 // 總預算 = 非專家權重(2.4 GB) + KV cache + ggml 計算緩衝 + 詞表 + arena // 所以有指定 ram_budget 時,用「目前 RSS(不含 arena)」當作已用額度扣掉。 // 預留量要涵蓋 arena 之外、但「在 pager 算完之後才長出來」的記憶體: // llama.cpp 的 CPU compute buffer(125 MiB @ubatch=128、502 MiB @512) // + KV cache(ctx=4096 約 670 MiB)+ 圖與執行緒堆疊 // 這些在 sizing 的當下還沒配置,所以不會出現在 rss_now 裡,必須手動扣。 // 設太小 = 實際用量會超過 RAM 預算。實測(RLIMIT_DATA 8 GB、ctx 4096、681 token 提示詞): // 預留 700 MB → 峰值 RSS 8231 MB(超出 8 GB 預算 39 MB) // 預留 900 MB → 峰值 RSS 8035 MB(留 157 MB 餘量) // 代價是槽位數:2039 → 1936 槽,decode 5.85 → 5.37 tok/s(約 -8%)。 // 寧可少 8% 速度,也不要在目標機器的 8 GB 上 OOM —— 這是本專案的硬約束。 const size_t reserve = (size_t) env_size("SDQ_RESERVE_MB", 900) << 20; const size_t rss_now = process_rss_bytes(); size_t arena; if (cfg_.ram_budget_bytes) { const size_t used = rss_now; // 這時候 arena 還沒配置 size_t avail = cfg_.ram_budget_bytes > used + reserve ? cfg_.ram_budget_bytes - used - reserve : 512ull << 20; arena = avail; if (cfg_.arena_bytes) { arena = std::min(arena, cfg_.arena_bytes); } fprintf(stderr, "[sdq] RAM 預算 %.2f GiB:已用(RSS) %.2f GiB、預留 %.2f GiB → arena %.2f GiB\n", cfg_.ram_budget_bytes / 1073741824.0, used / 1073741824.0, reserve / 1073741824.0, arena / 1073741824.0); } else { const size_t phys = physical_ram_bytes(); const size_t budget = (size_t) (phys * 0.70); arena = cfg_.arena_bytes ? cfg_.arena_bytes : std::min(budget, (size_t) (4ull << 30)); } if (arena < slot_bytes_ * 16) { throw std::runtime_error("expert arena 太小"); } // 切成兩區:大槽位只分給「需要大槽位的層」,比例 = 那些層佔全部層的比例。 const bool force_uniform = env_flag("SDQ_FORCE_UNIFORM_SLOT", false); if (big_bytes_ && small_bytes_ && !layers_.empty() && !force_uniform) { int n_big_layers = 0; for (const auto & L : layers_) { if (L.needs_big_slot) { n_big_layers++; } } const double f = (double) n_big_layers / (double) layers_.size(); const double avg = f * (double) big_bytes_ + (1.0 - f) * (double) small_bytes_; size_t n_total = (size_t)((double) arena / avg); size_t nb = (size_t)(n_total * f); // 大槽位的下限不是「有點就好」,而是**工作集的下限**: // 一個 token 在每個大槽位層要 8 個大 expert,這個模型有 3 個大槽位層 // → 至少要 3 × 8 = 24 個大槽位,而且還要留給預取的空間。 // 只給 16 個時量到過「找不到能放 il=34 ie=114 的槽位」而中止 // (這是在 verify 新增的預取壓力測試 [2b] 抓到的)。 const size_t big_need = (size_t) (n_big_layers * info_.n_expert_used); if (nb < big_need) { nb = big_need; } if (nb > n_total / 4) { nb = n_total / 4; } // 連工作集下限都放不進去時,寧可讓小槽位少一點(n_total 保持不變), // 否則會算出一個「必定會中止」的配置。 n_big_ = nb; n_small_ = n_total > nb ? n_total - nb : 0; arena = n_small_ * small_bytes_ + n_big_ * big_bytes_; } else { // 單一尺寸:全部槽位都是「大的那一種」(slot_bytes_)。 // 這是 SDQ_FORCE_UNIFORM_SLOT=1 與「模型只有一種 expert 大小」時的路徑。 // // 這裡必須記成 n_small_=0 / n_big_=全部,不能記成「n_small_=全部、單位用大的」: // slot_ptr()/slot_size_of() 是用「索引 < n_small_ → 小槽位」來分支的, // 把它記成全部都是小槽位的話,大 expert 會拿到小槽位的 stride // (寫出槽位邊界、算錯),而且大 expert 永遠配不到槽位 → 推論卡死。 // (第一版就是這樣,實驗跑到一半 hang 才找出來。) const size_t unit = slot_bytes_ ? slot_bytes_ : small_bytes_; n_small_ = 0; n_big_ = arena / unit; arena = n_big_ * unit; } if (n_small_ == 0 && n_big_ == 0) { throw std::runtime_error("expert arena 太小(算不出槽位)"); } n_slots_ = n_small_ + n_big_; if (n_slots_ < 16) { throw std::runtime_error("expert arena 太小(算不出 16 個槽位)"); } void * p = nullptr; const size_t alloc = arena + (2ull << 20); if (posix_memalign(&p, 2ull << 20, alloc) != 0 || !p) { throw std::runtime_error("無法配置 expert arena(" + std::to_string(arena >> 20) + " MB)"); } // 讓 arena 全部變成實體 RAM,並提示核心使用 huge page(減少 TLB 壓力) memset(p, 0, alloc); madvise(p, alloc, MADV_HUGEPAGE); arena_ = (uint8_t *) p; arena_bytes_ = arena; // OPT-1:剖面計數表(n_layer × n_expert,40×256 = 10240 筆 ≈ 40 KB,可忽略) prof_.assign((size_t) info_.n_layer * (size_t) info_.n_expert, 0); if (!cfg_.profile_in.empty()) { load_profile(cfg_.profile_in); } profile_out_ = cfg_.profile_out; slots_.assign(n_slots_, {}); free_slots_.reserve(n_small_); for (size_t i = n_small_; i-- > 0;) { free_slots_.push_back((int32_t) i); } free_big_slots_.reserve(n_big_); for (size_t i = n_slots_; i-- > n_small_;) { free_big_slots_.push_back((int32_t) i); } used_.store(0); hist_.assign((size_t) info_.n_layer, {}); for (auto & per_layer : hist_) { per_layer.assign((size_t) info_.n_expert, {}); } prev_used_.assign((size_t) info_.n_layer, {}); cur_used_.assign((size_t) info_.n_layer, {}); } void pager::init(const config & cfg, ggml_context * ctx) { (void) ctx; if (ready_) { return; } cfg_ = cfg; if (cfg_.model_path.empty()) { const char * v = getenv("SDQ_MODEL_PATH"); cfg_.model_path = v ? v : ""; } if (cfg_.model_path.empty()) { throw std::runtime_error("SDQ 缺少模型路徑(設定 SDQ_MODEL_PATH)"); } open_model_fd(); if (fd_ < 0) { throw std::runtime_error("無法開啟模型檔:" + cfg_.model_path); } load_gguf_header(cfg_.model_path); // 順序有意義:I/O 池的對齊緩衝必須在 alloc_arena() 之前配置, // 否則 arena 會把它當成「不存在的 RAM」,實際上就偷用了 8 GB 預算。 io_pool_start(); alloc_arena(); if (cfg_.prefetch) { pf_ = std::thread([this] { pf_worker(); }); } ready_ = true; print_summary(); } const layer_desc * pager::layer(int il) const { if (il < 0 || il >= (int) layers_.size()) { return nullptr; } const layer_desc & L = layers_[(size_t) il]; if (L.experts.empty()) { return nullptr; } return &L; } void pager::open_model_fd() { fd_ = open(cfg_.model_path.c_str(), O_RDONLY); if (fd_ < 0) { return; } // O_DIRECT:讀資料完全不進 kernel page cache。 // 這是「RAM 帳目誠實」的關鍵 —— 否則核心預讀會讓 mmap 檔案的頁算進 RSS // (實測 4K prefill 時檔案對映 RSS 可以多到 6 GB 以上),而且量到的 // 「SSD 讀取量」會被上一輪的 page cache 灌水。 const char * no_direct = getenv("SDQ_NO_DIRECT_IO"); if (no_direct && *no_direct && *no_direct != '0') { return; } #ifdef O_DIRECT fd_direct_ = open(cfg_.model_path.c_str(), O_RDONLY | O_DIRECT); if (fd_direct_ >= 0) { use_direct_ = true; fprintf(stderr, "[sdq] 使用 O_DIRECT 讀取模型檔(不佔用 kernel page cache)\n"); } else { fprintf(stderr, "[sdq] O_DIRECT 不可用,改用 buffered + fadvise(DONTNEED)\n"); } #endif } // 有對齊的緩衝區(O_DIRECT 要求),thread_local 避免競爭 static thread_local uint8_t * g_direct_buf = nullptr; static thread_local size_t g_direct_sz = 0; // 與 read_range_direct 相同,但用呼叫端提供的緩衝區(已配置好、已算進 RSS)。 // I/O 執行緒池走這條:每條 worker 在啟動時就配置好對齊緩衝,於是 arena sizing 時 // 這筆記憶體已經在 RSS 裡,不會變成「額外偷用的 RAM」。 bool pager::read_range_staged(uint64_t off, uint8_t * dst, size_t size, uint8_t * stage, size_t stage_sz) { if (!use_direct_) { return read_range(off, dst, size); } const size_t blk = 4096; const uint64_t start = off & ~(uint64_t)(blk - 1); const uint64_t end = (off + size + blk - 1) & ~(uint64_t)(blk - 1); const size_t bytes = (size_t)(end - start); if (stage_sz < bytes) { // 預先配置不足(不該發生):退回 thread_local 的通用路徑 return read_range_direct(off, dst, size); } size_t done = 0; while (done < bytes) { const ssize_t n = pread(fd_direct_, stage + done, bytes - done, (off_t)(start + done)); if (n <= 0) { return false; } done += (size_t) n; } memcpy(dst, stage + (off - start), size); return true; } bool pager::read_range_direct(uint64_t off, uint8_t * dst, size_t size) { if (!use_direct_) { return read_range(off, dst, size); } const size_t blk = 4096; const uint64_t start = off & ~(uint64_t)(blk - 1); const uint64_t end = (off + size + blk - 1) & ~(uint64_t)(blk - 1); const size_t bytes = (size_t)(end - start); if (g_direct_sz < bytes) { if (g_direct_buf) { free(g_direct_buf); } if (posix_memalign((void **) &g_direct_buf, blk, bytes) != 0) { g_direct_buf = nullptr; g_direct_sz = 0; use_direct_ = false; return read_range(off, dst, size); } g_direct_sz = bytes; } size_t done = 0; while (done < bytes) { const ssize_t n = pread(fd_direct_, g_direct_buf + done, bytes - done, (off_t)(start + done)); if (n <= 0) { return false; } done += (size_t) n; } memcpy(dst, g_direct_buf + (off - start), size); return true; } bool pager::read_range(uint64_t off, uint8_t * dst, size_t size) { size_t done = 0; while (done < size) { const ssize_t n = pread(fd_, dst + done, size - done, (off_t) (off + done)); if (n <= 0) { return false; } done += (size_t) n; } // 讓 RAM 帳目誠實:把 kernel 對這個檔案的 page cache 丟掉,避免「藏在核心裡的 RAM」 if (cfg_.fadvise_dontneed) { posix_fadvise(fd_, (off_t) off, (off_t) size, POSIX_FADV_DONTNEED); } return true; } bool pager::read_expert(int il, int ie, uint8_t * dst) { const layer_desc * L = layer(il); if (!L || ie < 0 || ie >= (int) L->experts.size()) { return false; } const expert_desc & ed = L->experts[(size_t) ie]; for (int p = 0; p < 3; p++) { if (ed.seg[p].size == 0) { continue; } if (!read_range_direct(ed.seg[p].off, dst + ed.slot_off[p], (size_t) ed.seg[p].size)) { return false; } } return true; } // ---------------------------------------------------------------- I/O 執行緒池 // // 問題:一個 expert = gate / up / down 三段,檔案上是三段不相鄰的資料。 // 循序讀就是 3 次 pread 排在同一條執行緒上,而單次 1.9 MB 的 pread 在本機量到 // ~7 ms(≈275 MB/s 的單流頻寬)。因此「一個 expert 要多久」= 3 × 單次延遲, // 這條路徑的長度跟「有多少條執行緒」無關 —— mode C 再怎麼排程都救不了。 // // 做法:把三段交給三條執行緒同時 pread,臨界路徑降到約「單段」的時間。 // 量測(8 個 expert 全部 miss 的一層,各 9 次取中位數): // 8 執行緒 × 三段序列讀 : 12.60 ms // 16 執行緒 × 三段平行讀 : 7.57 ms // 24 執行緒 × 三段平行讀 : 7.58 ms ← 16 就飽和,所以預設 3 × n_threads // // 這個池子對 prefill 更有用:一層要 256 個 expert = 768 個讀取工作, // 序列讀時 16 執行緒只能拿到 16 路的頻寬 utilization。 void pager::io_pool_start() { if (!cfg_.io_parallel) { fprintf(stderr, "[sdq] I/O 段平行:關(每個 expert 的 gate/up/down 循序讀)\n"); return; } int n = cfg_.io_threads > 0 ? cfg_.io_threads : cfg_.n_threads * 3; if (n < 4) { n = 4; } if (n > 48) { n = 48; } // 先算「一段權重最大有多少位元組」,每條 worker 配一塊對齊緩衝。 // 這必须在 alloc_arena() 之前做完,否則這筆 RAM 會變成帳外支出。 size_t max_seg = 0; for (const auto & L : layers_) { for (const auto & ed : L.experts) { for (int p = 0; p < 3; p++) { max_seg = std::max(max_seg, (size_t) ed.seg[p].size); } } } const size_t stage_bytes = ((max_seg + 2 * 4096) + 4095) & ~(size_t) 4095; io_stage_.assign((size_t) n, io_stage{}); for (int i = 0; i < n; i++) { void * p = nullptr; if (posix_memalign(&p, 4096, stage_bytes) == 0 && p) { memset(p, 0, stage_bytes); // 實際碰頁,讓它出現在 RSS 裡 io_stage_[(size_t) i].buf = (uint8_t *) p; io_stage_[(size_t) i].sz = stage_bytes; } } io_split_ = cfg_.io_split > 0 ? cfg_.io_split : 1; io_stop_ = false; io_pool_.reserve((size_t) n); for (int i = 0; i < n; i++) { const size_t wid = (size_t) i; try { io_pool_.emplace_back([this, wid] { io_worker(wid); }); } catch (...) { break; } } io_threads_ = (int) io_pool_.size(); io_pool_on_ = io_threads_ > 0; if (io_pool_on_) { fprintf(stderr, "[sdq] I/O 段平行:%d 條讀取執行緒 × %zu KB 對齊緩衝" "(gate/up/down 同時讀,每段再切 %d 塊)\n", io_threads_, stage_bytes >> 10, io_split_); } } void pager::io_pool_stop() { if (io_pool_.empty()) { for (auto & s : io_stage_) { free(s.buf); s.buf = nullptr; } io_stage_.clear(); return; } { std::lock_guard lk(io_mu_); io_stop_ = true; } io_cv_.notify_all(); for (auto & t : io_pool_) { if (t.joinable()) { t.join(); } } io_pool_.clear(); io_pool_on_ = false; for (auto & s : io_stage_) { free(s.buf); s.buf = nullptr; s.sz = 0; } io_stage_.clear(); } void pager::io_worker(size_t wid) { for (;;) { io_job job{}; { std::unique_lock lk(io_mu_); io_cv_.wait(lk, [this] { return io_stop_ || !io_q_.empty(); }); if (io_q_.empty()) { if (io_stop_) { return; } continue; } job = io_q_.front(); io_q_.pop_front(); } io_jobs_.fetch_add(1, std::memory_order_relaxed); // 讀取本身不持鎖:io_mu_ 只保護佇列與完成計數 const bool ok = (wid < io_stage_.size() && io_stage_[wid].buf) ? read_range_staged(job.off, job.dst, job.size, io_stage_[wid].buf, io_stage_[wid].sz) : read_range_direct(job.off, job.dst, job.size); { std::lock_guard lk(io_mu_); if (!ok) { job.ctx->ok = false; } if (job.ctx->remaining.fetch_sub(1) == 1) { io_done_cv_.notify_all(); } } } } // 把 (il,ie) 的三段交給執行緒池並行讀,呼叫端等到整批完成才返回。 // ctx 放在呼叫端 stack:worker 只會在 remaining 歸零之後才再碰它, // 而我們正是等到 remaining==0 才離開函式,所以沒有 use-after-return。 bool pager::read_expert_par(int il, int ie, uint8_t * dst) { const layer_desc * L = layer(il); if (!L || ie < 0 || ie >= (int) L->experts.size()) { return false; } if (!io_pool_on_) { return read_expert(il, ie, dst); } const expert_desc & ed = L->experts[(size_t) ie]; io_ctx ctx; int n = 0; // 把每一段再切成一樣大的好幾塊,交給不同的執行緒同時讀。 // // 為什麼:單一 pread 的時間 ≈ 固定延遲 + 位元組數 / 單流頻寬。 // 本機實測單流約 275 MB/s、4K 隨機讀延遲約 0.57 ms,所以 // 590 KB 讀一次 ≈ 0.57 + 2.1 = 2.7 ms // 590 KB 切成 2 塊 ≈ 0.57 + 1.05 = 1.6 ms // 而「讀完一個 expert」的時間正是每層的關鍵路徑,切塊就是直接砍關鍵路徑。 // 設定檔 SDQ_IO_SPLIT。**實測預設是 1(不切塊)**:切塊在 24 條讀取執行緒 // 之下沒有可重現的改善(split=1..6 的中位數在 4.9~5.9 之間跳動,純噪音), // 因為 8 個 expert × 3 段 = 24 個並行請求已經把裝置的佇列灌飽, // 這時瓶頸是裝置總頻寬而不是單一請求的延遲。保留這個開關是為了其他裝置。 const int split = io_split_; const size_t blk = 4096; { std::lock_guard lk(io_mu_); for (int p = 0; p < 3; p++) { const size_t sz = (size_t) ed.seg[p].size; if (sz == 0) { continue; } uint8_t * base = dst + ed.slot_off[p]; if (split <= 1 || sz <= 2 * blk) { io_q_.push_back(io_job{ed.seg[p].off, base, sz, &ctx}); n++; continue; } size_t chunk = ((sz + (size_t) split - 1) / (size_t) split + blk - 1) & ~(blk - 1); for (size_t off = 0; off < sz; off += chunk) { const size_t take = sz - off < chunk ? sz - off : chunk; io_q_.push_back(io_job{ed.seg[p].off + off, base + off, take, &ctx}); n++; } } ctx.remaining.store(n); } if (n == 0) { return true; } io_cv_.notify_all(); const uint64_t t0 = now_us(); { std::unique_lock lk(io_mu_); io_done_cv_.wait(lk, [&ctx] { return ctx.remaining.load() == 0; }); } io_wait_ns_.fetch_add((now_us() - t0) * 1000ull, std::memory_order_relaxed); return ctx.ok; } uint64_t pager::read_expert_span(int il, const int32_t * experts, int n, const int32_t * slots) { const layer_desc * L = layer(il); if (!L || n <= 0) { return 0; } const uint64_t t_start = now_us(); // 每個執行緒一份暫存區(多執行緒會同時呼叫,不能共用一個 buffer) static thread_local std::vector scratch_buf; uint64_t total = 0; // 把 experts 切成「檔案編號連續」的區段 int i = 0; while (i < n) { int j = i + 1; while (j < n && experts[j] == experts[j - 1] + 1) { j++; } const int cnt = j - i; const expert_desc & first = L->experts[(size_t) experts[i]]; const expert_desc & last = L->experts[(size_t) experts[j - 1]]; for (int p = 0; p < 3; p++) { if (first.seg[p].size == 0 || last.seg[p].size == 0) { continue; } const uint64_t off = first.seg[p].off; const uint64_t end = last.seg[p].off + last.seg[p].size; const size_t bytes = (size_t) (end - off); if (scratch_buf.size() < bytes) { scratch_buf.assign(bytes, 0); } if (!read_range_direct(off, scratch_buf.data(), bytes)) { return total; } total += bytes; // 散到各自的槽位:scratch 從 off 開始,所以偏移是「相對於 off」 for (int k = i; k < j; k++) { const expert_desc & ed = L->experts[(size_t) experts[k]]; const size_t seg_bytes = (size_t) ed.seg[p].size; const size_t delta = (size_t) (ed.seg[p].off - off); if (slots[k] < 0) { continue; } if ((size_t) slots[k] >= n_slots_ || (uint64_t) ed.slot_off[p] + seg_bytes > slot_bytes_ || delta + seg_bytes > scratch_buf.size()) { fprintf(stderr, "[sdq-pager] span 越界:il=%d slot=%d/%zu slot_off=%lld seg=%zu delta=%zu scratch=%zu p=%d e=%d\n", il, slots[k], n_slots_, (long long) ed.slot_off[p], seg_bytes, delta, scratch_buf.size(), p, experts[k]); return total; } memcpy(slot_ptr(slots[k]) + ed.slot_off[p], scratch_buf.data() + delta, seg_bytes); } } i = j; } stats_.demand_wait_us.fetch_add(now_us() - t_start, std::memory_order_relaxed); return total; } int pager::reserve_slot(int il, int ie) { if (!ready_) { return -1; } count_request(il, ie); std::lock_guard lk(mu_); const int32_t exist = find_slot_locked(il, ie); if (exist >= 0) { slots_[(size_t) exist].last_use = ++clock_; slots_[(size_t) exist].pin++; touch_freq_locked(slots_[(size_t) exist]); return exist; } const int32_t s = victim_locked(); if (s < 0) { return -1; } slot_state & st = slots_[(size_t) s]; if (st.layer >= 0) { page_.erase(((int64_t) st.layer << 20) | (int64_t) st.expert); stats_.evictions.fetch_add(1, std::memory_order_relaxed); used_.fetch_sub(1, std::memory_order_relaxed); } st.layer = il; st.expert = ie; st.last_use = ++clock_; st.freq = 1; st.pin = 1; st.is_hot = hot_lookup(il, ie); page_[((int64_t) il << 20) | (int64_t) ie] = s; used_.fetch_add(1, std::memory_order_relaxed); stats_.expert_misses.fetch_add(1, std::memory_order_relaxed); return s; } int pager::probe(int il, int ie) { if (!ready_) { return -1; } stats_.expert_requests.fetch_add(1, std::memory_order_relaxed); count_request(il, ie); std::lock_guard lk(mu_); const int32_t s = find_slot_locked(il, ie); if (s < 0) { stats_.expert_misses.fetch_add(1, std::memory_order_relaxed); return -1; } slots_[(size_t) s].last_use = ++clock_; slots_[(size_t) s].pin++; touch_freq_locked(slots_[(size_t) s]); stats_.expert_hits.fetch_add(1, std::memory_order_relaxed); return s; } void pager::note_span_read(uint64_t bytes) { stats_.ssb_read_bytes_demand.fetch_add(bytes, std::memory_order_relaxed); stats_.demand_wait_us.fetch_add(0, std::memory_order_relaxed); } // 衰減式頻率:每 8192 次使用就把所有計數減半,讓「過去熱門、現在變冷」的 expert 能被淘汰 void pager::touch_freq_locked(slot_state & st) { if (++freq_tick_ >= 8192) { freq_tick_ = 0; for (auto & s : slots_) { s.freq /= 2; } } if (st.freq < 0xFFFFFFFFu) { st.freq++; } } // ============================== OPT-1:熱門 expert 剖面 ============================== // 問題實測(docs/03):arena 3.88 GiB 只放得下 2044 個 expert,總共卻有 10240 個, // prefill 的幾百個「只用一次」的專家會把熱區整個洗掉,命中率卡在 ~70%。 // 解法借用 Strata 的做法(src/core/expert_cache.cpp):**離線剖面 + 靜態排行榜 + // 熱頁永不淘汰**。這裡不動檔案格式,用純文字存排行榜。 void pager::count_request(int il, int ie) { if (prof_.empty() || il < 0 || ie < 0) { return; } const size_t idx = (size_t) il * (size_t) info_.n_expert + (size_t) ie; if (idx < prof_.size()) { // 單執行緒 ++ 即可:請求次數是統計用途,不參與任何正確性判斷 if (prof_[idx] < 0xFFFFFFFFu) { prof_[idx]++; } } } bool pager::hot_lookup(int il, int ie) const { if (hot_.empty() || il < 0 || ie < 0) { return false; } const size_t idx = (size_t) il * (size_t) info_.n_expert + (size_t) ie; return idx < hot_.size() && hot_[idx] != 0; } // 挑一個容得下 (il,ie) 的槽。與 victim_locked 的差別:**先只在冷頁裡找**, // 找不到才退讓去動熱頁。這就是「熱區不會被 prefill 洗掉」的機制。 // 這個槽位放得下「需要 big 的 expert」嗎? // 小槽位:只能放小 expert // 大槽位:小expert 也放得下(只是浪費一點空間) static inline bool slot_accepts(bool need_big, bool slot_is_big) { return !need_big || slot_is_big; } int pager::admit_slot_locked(int il, int ie) { // il/ie 都 < 0 代表「不指定大小」→ 兩種槽位都可以(舊的呼叫路徑)。 const bool need_big = (il >= 0 && ie >= 0) ? need_big_slot(il, ie) : false; // 完全沒用過的空槽永遠優先(它不是熱頁,佔用等於浪費) // 大 expert 先看大槽位;小 expert 先看小槽位。 std::vector & first = need_big ? free_big_slots_ : free_slots_; std::vector & second = need_big ? free_slots_ : free_big_slots_; while (!first.empty()) { const int32_t s = first.back(); if (!slots_[(size_t) s].busy) { first.pop_back(); return s; } } // 第一輪:只在「冷頁」裡挑最久未用的。 // want_hot=true → 熱門 expert 進來時,把冷頁換掉(不動其他熱頁) // want_hot=false → 冷門 expert 進來時,優先吃冷頁 // 兩種情況的規則相同,所以 hot_lookup 的結果不影響選擇順序,只影響標記。 for (int pass = 0; pass < 2; pass++) { int64_t best = -1; uint64_t best_clock = UINT64_MAX; for (size_t s = 0; s < n_slots_; s++) { const slot_state & st = slots_[s]; if (st.pin > 0 || st.busy) { continue; } if (pass == 0 && st.is_hot) { continue; // 第一輪不動熱頁 } if (!slot_accepts(need_big, is_big_slot(s))) { continue; } if (st.last_use < best_clock) { best_clock = st.last_use; best = (int64_t) s; } } if (best >= 0) { return (int) best; } } // 退讓:另一種大小的空槽。**一定要檢查放不放得下** —— 小槽位放不下大 expert, // 硬塞會讓 read_expert 寫出槽位邊界(靜默把別的 expert 寫壞)。 while (!second.empty()) { const int32_t s = second.back(); if (!slots_[(size_t) s].busy && slot_accepts(need_big, is_big_slot(s))) { second.pop_back(); return s; } if (!slots_[(size_t) s].busy) { second.pop_back(); } else { break; } } return -1; } bool pager::load_profile(const std::string & path) { std::FILE * f = std::fopen(path.c_str(), "rb"); if (!f) { fprintf(stderr, "[sdq] 找不到 expert 剖面 %s(將以動態 LRU/LFU 運作)\n", path.c_str()); return false; } char magic[8] = {0}; if (std::fread(magic, 1, 4, f) != 4 || std::memcmp(magic, "SDQP", 4) != 0) { std::fclose(f); fprintf(stderr, "[sdq] %s 不是 SDQ expert 剖面(缺少 SDQP 標頭)\n", path.c_str()); return false; } uint32_t hdr[3] = {0, 0, 0}; if (std::fread(hdr, sizeof(uint32_t), 3, f) != 3) { std::fclose(f); return false; } const uint32_t n = hdr[2]; std::vector>> ranked; ranked.reserve(n); for (uint32_t i = 0; i < n; i++) { int32_t l = 0, e = 0; uint32_t c = 0; if (std::fread(&l, sizeof(int32_t), 1, f) != 1) break; if (std::fread(&e, sizeof(int32_t), 1, f) != 1) break; if (std::fread(&c, sizeof(uint32_t), 1, f) != 1) break; ranked.push_back({c, {l, e}}); } std::fclose(f); std::sort(ranked.begin(), ranked.end(), [](const auto & a, const auto & b) { return a.first > b.first; }); hot_.assign((size_t) info_.n_layer * (size_t) info_.n_expert, 0); const int want = cfg_.hot_pin > 0 ? cfg_.hot_pin : (int) ranked.size(); n_hot_ = 0; for (const auto & r : ranked) { if (n_hot_ >= want) { break; } const int l = r.second.first, e = r.second.second; if (l < 0 || l >= info_.n_layer || e < 0 || e >= info_.n_expert) { continue; } const size_t idx = (size_t) l * (size_t) info_.n_expert + (size_t) e; if (!hot_[idx]) { hot_[idx] = 1; n_hot_++; } } fprintf(stderr, "[sdq] expert 剖面:%s → 釘選 %d / %zu 個熱門 expert(永不淘汰)\n", path.c_str(), n_hot_, ranked.size()); return n_hot_ > 0; } void pager::dump_profile(const std::string & path) { if (prof_.empty()) { return; } std::vector>> ranked; ranked.reserve(prof_.size()); uint64_t total = 0; for (int l = 0; l < info_.n_layer; l++) { for (int e = 0; e < info_.n_expert; e++) { const uint32_t c = prof_[(size_t) l * (size_t) info_.n_expert + (size_t) e]; if (c > 0) { ranked.push_back({c, {l, e}}); total += c; } } } std::sort(ranked.begin(), ranked.end(), [](const auto & a, const auto & b) { return a.first > b.first; }); std::FILE * f = std::fopen(path.c_str(), "wb"); if (!f) { fprintf(stderr, "[sdq] 無法寫出 expert 剖面 %s\n", path.c_str()); return; } std::fwrite("SDQP", 1, 4, f); const uint32_t hdr[3] = {1, (uint32_t) info_.n_layer, (uint32_t) ranked.size()}; std::fwrite(hdr, sizeof(uint32_t), 3, f); for (const auto & r : ranked) { const int32_t l = r.second.first, e = r.second.second; const uint32_t c = r.first; std::fwrite(&l, sizeof(int32_t), 1, f); std::fwrite(&e, sizeof(int32_t), 1, f); std::fwrite(&c, sizeof(uint32_t), 1, f); } std::fclose(f); // 覆盖率報告:前 K 名可以蓋掉多少請求 → 直接決定 hot_pin 該設多少 uint64_t acc = 0; const size_t ks[6] = {512, 1024, 1536, 2048, 3072, 4096}; fprintf(stderr, "[sdq] expert 剖面已寫到 %s(%zu 個有被用到,總請求 %llu)\n", path.c_str(), ranked.size(), (unsigned long long) total); size_t ki = 0; for (size_t i = 0; i < ranked.size(); i++) { acc += ranked[i].first; if (ki < 6 && i + 1 == ks[ki]) { fprintf(stderr, "[sdq] 前 %5zu 名涵蓋 %6.2f%% 的請求\n", ks[ki], 100.0 * (double) acc / (double) (total ? total : 1)); ki++; } } } int pager::find_slot_locked(int il, int ie) { const int64_t key = ((int64_t) il << 20) | (int64_t) ie; auto it = page_.find(key); return it == page_.end() ? -1 : it->second; } // 找一個「可以放 (il,ie) 而且 last_use <= stale_before」的槽位。 // // 為什麼不能直接用 victim_stale_locked():大小兩種槽位之後,那個函式會掃過**所有** // 槽位(包含小槽位),於是一個需要大槽位的 expert 可能拿到小槽位, // 預取時 read_expert 會寫出槽位邊界 → 靜默把鄰居槽位的資料寫壞。 // 這是引進大小兩種槽位時真正會咬人的地方。 // // 順序:先看對應尺寸的空槽(從沒用過,最該拿來放預取), // 再掃全部槽位挑最久未用的(放得下的那種尺寸)。 int pager::victim_stale_for(int il, int ie, uint64_t stale_before) { const bool need_big = need_big_slot(il, ie); std::vector & first = need_big ? free_big_slots_ : free_slots_; std::vector & second = need_big ? free_slots_ : free_big_slots_; for (std::vector * lst : { &first, &second }) { while (!lst->empty()) { const int32_t s = lst->back(); if (slots_[(size_t) s].busy) { break; } if (slots_[(size_t) s].last_use <= stale_before && slot_accepts(need_big, is_big_slot(s))) { lst->pop_back(); return s; } lst->pop_back(); } } int64_t best = -1; uint64_t best_clock = UINT64_MAX; for (size_t i = 0; i < n_slots_; i++) { const slot_state & st = slots_[i]; if (st.pin > 0 || st.busy || st.last_use > stale_before) { continue; } if (!slot_accepts(need_big, is_big_slot(i))) { continue; } if (st.last_use < best_clock) { best_clock = st.last_use; best = (int64_t) i; } } return (int) best; } int pager::victim_stale_locked(uint64_t stale_before) { // 只挑「last_use 很舊」而且沒被 pin 的槽:保護最近的工作集 while (!free_slots_.empty()) { const int32_t s = free_slots_.back(); if (!slots_[(size_t) s].busy && slots_[(size_t) s].last_use <= stale_before) { free_slots_.pop_back(); return s; } if (free_slots_.size() * 4 < n_slots_) { break; } free_slots_.pop_back(); } int64_t best = -1; uint64_t best_clock = UINT64_MAX; for (size_t i = 0; i < n_slots_; i++) { const slot_state & st = slots_[i]; if (st.pin > 0 || st.busy || st.last_use > stale_before) { continue; } if (st.last_use < best_clock) { best_clock = st.last_use; best = (int64_t) i; } } return (int) best; } int pager::acquire(int il, int ie, bool blocking) { if (!ready_) { return -1; } std::unique_lock lk(mu_); int32_t s = find_slot_locked(il, ie); if (s >= 0 && slots_[s].busy) { // 這個槽位**正在被填入**:另一條執行緒剛 admit 它、資料還沒讀完。 // // 這裡一定要等,不能當成命中直接用。否則讀到的是「寫到一半的權重」, // 而且那個錯誤是**非決定性**的 —— 每次執行結果可能不同(實測在 // prefill 時會讓同一個提示詞跑出不同的文字)。 // busy 的槽位在 admit/victim 那邊本來就會被跳過,所以這裡是唯一 // 還能拿到它的路徑,也因此必須在這裡等。 const uint64_t t0 = now_us(); while (slots_[s].busy) { lk.unlock(); std::this_thread::sleep_for(std::chrono::microseconds(200)); lk.lock(); if (now_us() - t0 > (uint64_t) env_size("SDQ_BUSY_WAIT_MS", 30000) * 1000) { fprintf(stderr, "[sdq] 致命:等待槽位 il=%d ie=%d 讀入超過 %u ms(busy 一直沒解除)\n", il, ie, (unsigned) env_size("SDQ_BUSY_WAIT_MS", 30000)); std::fflush(stderr); abort(); } } // 讀取可能失敗(那時槽位已被清空)→ 重新走一次正常的 miss 路徑 if (find_slot_locked(il, ie) != s) { s = -1; } } if (s >= 0) { slots_[s].last_use = ++clock_; slots_[s].pin++; touch_freq_locked(slots_[(size_t) s]); stats_.expert_requests.fetch_add(1, std::memory_order_relaxed); stats_.expert_hits.fetch_add(1, std::memory_order_relaxed); return s; } stats_.expert_requests.fetch_add(1, std::memory_order_relaxed); count_request(il, ie); stats_.expert_misses.fetch_add(1, std::memory_order_relaxed); // 先找一個已經在背景讀入這頁的槽 for (size_t i = 0; i < n_slots_; i++) { slot_state & st = slots_[i]; if (st.busy && st.layer == il && st.expert == ie) { // 等背景執行緒讀完 lk.unlock(); const uint64_t t0 = now_us(); while (true) { lk.lock(); if (!slots_[i].busy) { break; } lk.unlock(); std::this_thread::sleep_for(std::chrono::microseconds(200)); } s = (int32_t) i; lk.unlock(); stats_.prefetch_wait_us.fetch_add(now_us() - t0, std::memory_order_relaxed); std::unique_lock lk2(mu_); slots_[s].last_use = ++clock_; slots_[s].pin++; stats_.expert_hits.fetch_add(1, std::memory_order_relaxed); stats_.prefetch_used.fetch_add(1, std::memory_order_relaxed); return s; } } if (!blocking) { return -1; } // 選槽位時必須知道「這個 expert 需要哪一種大小」,否則大 expert 可能被塞進小槽位 // (寫出邊界 → 靜默算錯),或者小 expert 佔掉大槽位(浪費 6.8% 的 arena)。 s = admit_slot_locked(il, ie); if (s < 0) { // 拿不到槽位:短暫等待別的執行緒解除 pin(每層開頭會 unpin_all)。 // 這裡「重試」而不是直接回 -1,因為回 -1 會讓呼叫端**跳過這個 expert**, // 結果是靜默少算一條專家分支 —— 那是最嚴重的錯誤。 // 2000 次 × 250 µs = 最多等 0.5 秒;真的撐不到那麼久,代表 arena 被 // 預取的 busy 槽位佔住了(實測:跨層預取開啟時會發生)。 for (int attempt = 0; attempt < 2000 && s < 0; attempt++) { lk.unlock(); std::this_thread::sleep_for(std::chrono::microseconds(250)); lk.lock(); s = admit_slot_locked(il, ie); } if (s < 0) { fprintf(stderr, "[sdq] 嚴重:找不到能放 il=%d ie=%d(need_big=%d)的槽位," "大槽 %zu/小槽 %zu\n", il, ie, (int) need_big_slot(il, ie), n_big_, n_small_); // 這裡「回 -1」會讓呼叫端**跳過這個 expert**,結果是靜默少算一條專家分支 // —— 那是整個引擎最嚴重的錯誤型態(模型會安靜地給出錯誤答案,而且 // 每次執行還不一定一樣)。寧可整個行程死掉,也不要產出「看起來正常」 // 的錯誤結果。 fprintf(stderr, "[sdq] 致命:無法取得 expert 槽位。為了不給出靜默錯誤的結果," "直接中止。請回報這個情況(多半是預取把槽位佔成 busy 太久)。\n"); std::fflush(stderr); abort(); } } slot_state & st = slots_[s]; if (st.layer >= 0) { page_.erase(((int64_t) st.layer << 20) | (int64_t) st.expert); stats_.evictions.fetch_add(1, std::memory_order_relaxed); used_.fetch_sub(1, std::memory_order_relaxed); } st.layer = il; st.expert = ie; st.last_use = ++clock_; st.pin = 1; st.is_hot = hot_lookup(il, ie); // **先標 busy,再 publish 到 page_。** // 順序反過來的話(先 publish 再標 busy 中間有一個窗口)另一條執行緒會把它 // 當成命中,然後在資料還沒讀完時就拿去做矩陣乘法 → 非決定性的錯誤輸出。 // acquire() 看到 busy 會等;admit/victim 本來就會跳過 busy 的槽位。 st.busy = true; page_[((int64_t) il << 20) | (int64_t) ie] = s; lk.unlock(); const uint64_t t0 = now_us(); const bool ok = read_expert_par(il, ie, slot_ptr(s)); const uint64_t dt = now_us() - t0; std::unique_lock lk2(mu_); st.busy = false; if (!ok) { st.layer = -1; st.expert = -1; st.pin = 0; page_.erase(((int64_t) il << 20) | (int64_t) ie); return -1; } used_.fetch_add(1, std::memory_order_relaxed); stats_.demand_wait_us.fetch_add(dt, std::memory_order_relaxed); const layer_desc * L = layer(il); if (L) { stats_.ssb_read_bytes_demand.fetch_add((uint64_t) L->experts[(size_t) ie].bytes, std::memory_order_relaxed); } return s; } void pager::prefetch_expert(int il, int ie) { if (!ready_ || !cfg_.prefetch) { return; } { std::unique_lock lk(mu_); if (find_slot_locked(il, ie) >= 0) { stats_.prefetch_duplicate.fetch_add(1, std::memory_order_relaxed); return; } // 預取只能占用「空槽」或「已經很久沒被用到的頁」。 // (實測:無差別淘汰會讓命中率從 90% 掉到 27%,而預取使用率是 0%,純粹有害。) // // **槽位一定要選對尺寸**:need_big_slot() 為真的 expert 拿到小槽位會寫出邊界。 const uint64_t stale_before = clock_ - (uint64_t) env_size("SDQ_PREFETCH_STALE_TOKENS", 2) * 320; // 預取佇列也不能太深:SSD 讀不夠快時,排在後面的預取早就過期了 if (pf_q_.size() > (size_t) env_size("SDQ_PREFETCH_QUEUE", 64)) { return; } // **必須留夠空的槽位給 demand。** // 預取會把槽位標成 busy 並佔著它讀取;大槽位本來就少(工作集下限是 // n_big_layers × n_expert_used),如果預取把最後幾個空槽都吃掉, // 下一層的 demand 就會「找不到能放 il=34 ie=114 的槽位」而中止。 // (這是 verify 新增的預取壓力測試 [2b] 抓到的實際情況。) // 所以:對應尺寸的空槽少於 n_expert_used 個時,這一筆預取直接不做。 const size_t reserve = (size_t) info_.n_expert_used; const size_t n_free_right = need_big_slot(il, ie) ? free_big_slots_.size() : free_slots_.size(); if (n_free_right < reserve) { stats_.prefetch_duplicate.fetch_add(1, std::memory_order_relaxed); return; } const int32_t s = victim_stale_for(il, ie, stale_before); if (s < 0) { return; } slot_state & st = slots_[s]; if (st.layer >= 0) { page_.erase(((int64_t) st.layer << 20) | (int64_t) st.expert); stats_.evictions.fetch_add(1, std::memory_order_relaxed); used_.fetch_sub(1, std::memory_order_relaxed); } st.layer = il; st.expert = ie; st.busy = true; st.pin = 0; pf_q_.push_back({s, il, ie, now_us()}); stats_.prefetch_issued.fetch_add(1, std::memory_order_relaxed); } pf_cv_.notify_one(); } void pager::pf_worker() { std::unique_lock lk(pf_mu_); for (;;) { pf_cv_.wait(lk, [this] { return pf_stop_ || !pf_q_.empty(); }); if (pf_stop_ && pf_q_.empty()) { return; } const pf_item item = pf_q_.front(); pf_q_.pop_front(); lk.unlock(); // 這個槽可能已經被淘汰/重用 → 檢查 generation { std::unique_lock g(mu_); const slot_state & st = slots_[item.slot]; const bool still = st.busy && st.layer == item.il && st.expert == item.ie; if (!still) { lk.lock(); continue; } } // 預取也走 I/O 段平行:單獨一條執行緒循序讀三段,是這條路徑唯一的瓶頸。 const bool ok = cfg_.io_parallel ? read_expert_par(item.il, item.ie, slot_ptr(item.slot)) : read_expert(item.il, item.ie, slot_ptr(item.slot)); const layer_desc * L = layer(item.il); { std::unique_lock g(mu_); slot_state & st = slots_[item.slot]; if (st.busy && st.layer == item.il && st.expert == item.ie) { st.busy = false; st.last_use = ++clock_; if (ok) { page_[((int64_t) item.il << 20) | (int64_t) item.ie] = item.slot; used_.fetch_add(1, std::memory_order_relaxed); stats_.prefetch_completed.fetch_add(1, std::memory_order_relaxed); if (L) { stats_.ssb_read_bytes_prefetch.fetch_add( (uint64_t) L->experts[(size_t) item.ie].bytes, std::memory_order_relaxed); } } else { st.layer = -1; st.expert = -1; } } } lk.lock(); } } void pager::release(int s) { if (s < 0 || (size_t) s >= n_slots_) { return; } std::lock_guard lk(mu_); slot_state & st = slots_[(size_t) s]; if (st.pin > 0) { st.pin--; } } void pager::unpin_all() { std::unique_lock lk(mu_); for (auto & st : slots_) { st.pin = 0; } } void pager::note_used(int il, int ie) { if (il < 0 || ie < 0) { return; } std::lock_guard lk(hist_mu_); if (il >= (int) cur_used_.size()) { return; } auto & v = cur_used_[(size_t) il]; if (std::find(v.begin(), v.end(), ie) == v.end()) { v.push_back(ie); } } // 「下一層」預取:只在 layer L 的 op 開頭呼叫一次。 // // 為什麼需要這個函式:end_token() 是**整個 token 的 40 層都算完之後**才被呼叫的 // (sdq_cli.cpp 在 llama_decode 之後呼叫),所以它排的預取完全沒有和任何計算重疊 // —— 量測上 prefetch 使用率只有 0~8 %,等於整條預取路徑是關的。 // 在 layer L 開始時就替 layer L+1 排預取,才有一整層的時間(讀+算,約 4 ms) // 可以讓 SSD 讀取躲在後面。 // // 預測用的是 prev_used_[il+1](上一個 token 在該層選的專家)—— MoE 的路由在 // 鄰近 token 之間非常穩定,這是最便宜且不用訓練的預測器。 void pager::prefetch_layer_ahead(int il) { if (!ready_ || !cfg_.prefetch || il < 0 || il >= (int) layers_.size()) { return; } int32_t ids[64]; const int k = (int) info_.n_expert_used; if (k <= 0 || k > 64) { return; } predict(il, ids, k); for (int i = 0; i < k; i++) { prefetch_expert(il, ids[i]); } } void pager::end_token() { if (!ready_) { return; } // 用「上個 token 的選擇」更新一階 Markov 轉移表 { std::lock_guard lk(hist_mu_); for (size_t il = 0; il < cur_used_.size(); il++) { auto & cur = cur_used_[il]; if (cur.empty()) { continue; } auto & layer_hist = hist_[il]; // 以每層最近一次(或平均)的前一次選擇為前驅 if (!prev_used_.empty() && il < prev_used_.size() && !prev_used_[il].empty()) { for (int32_t from : prev_used_[il]) { trans & t = layer_hist[(size_t) from]; for (int32_t to : cur) { int found = -1; for (int32_t i = 0; i < t.n; i++) { if (t.to[i] == to) { t.cnt[i]++; found = i; break; } } if (found < 0 && t.n < 64) { t.to[t.n] = to; t.cnt[t.n] = 1; t.n++; } } } } prev_used_[il] = cur; cur.clear(); } } // 依預測排預取 if (!cfg_.prefetch) { return; } int32_t ids[64]; for (size_t il = 0; il < layers_.size(); il++) { const int k = (int) info_.n_expert_used; predict((int) il, ids, k); for (int i = 0; i < k; i++) { prefetch_expert((int) il, ids[i]); } } } void pager::predict(int il, int32_t * out, int k) { if (il < 0 || il >= (int) prev_used_.size()) { for (int i = 0; i < k; i++) { out[i] = i; } return; } // 以最近一次選擇為前驅;沒有歷史就直接沿用上一次的選擇 const std::vector & prev = prev_used_[(size_t) il]; std::vector> scored; if (prev.empty()) { for (int i = 0; i < k; i++) { out[i] = i; } return; } for (int32_t from : prev) { const trans & t = hist_[(size_t) il][(size_t) from]; for (int32_t i = 0; i < t.n; i++) { scored.push_back({t.cnt[i], t.to[i]}); } } if (scored.empty()) { for (size_t i = 0; i < prev.size() && i < (size_t) k; i++) { out[i] = prev[i]; } for (int i = (int) prev.size(); i < k; i++) { out[i] = i; } return; } std::sort(scored.begin(), scored.end(), [](auto & a, auto & b) { return a.first > b.first; }); for (int i = 0; i < k && i < (int) scored.size(); i++) { out[i] = scored[(size_t) i].second; } } double pager::hit_rate() const { const uint64_t req = stats_.expert_requests.load(); const uint64_t hit = stats_.expert_hits.load(); return req ? (double) hit / (double) req : 0.0; } void pager::shutdown() { if (!ready_) { return; } if (!profile_out_.empty()) { std::lock_guard lk(mu_); dump_profile(profile_out_); } if (pf_.joinable()) { { std::unique_lock lk(pf_mu_); pf_stop_ = true; } pf_cv_.notify_all(); pf_.join(); } io_pool_stop(); if (fd_ >= 0) { close(fd_); fd_ = -1; } if (fd_direct_ >= 0) { close(fd_direct_); fd_direct_ = -1; } ready_ = false; } pager::~pager() { shutdown(); if (arena_) { // 注意:arena 是 posix_memalign 配置的,必須用 free() 釋放, // 用 munmap() 會破壞 glibc 的 heap(之後任何 free 都會報 double free)。 free(arena_); arena_ = nullptr; } } static double pct(uint64_t a, uint64_t b) { return b ? 100.0 * (double) a / (double) b : 0.0; } void pager::print_summary() const { if (!ready_) { return; } const layer_desc * L0 = layer(0); if (L0 && !L0->experts.empty()) { const expert_desc & e = L0->experts[0]; fprintf(stderr, "[sdq] page-map: n_layer=%lld n_expert=%lld n_expert_used=%lld n_embd=%lld n_ff_exp=%lld" " | L0E0 gate=%llu/%llu up=%llu/%llu down=%llu/%llu (off/size) ne0=%lld ne1=%lld d_ne0=%lld d_ne1=%lld\n", (long long) info_.n_layer, (long long) info_.n_expert, (long long) info_.n_expert_used, (long long) info_.n_embd, (long long) info_.n_ff_exp, (unsigned long long) e.seg[part::GATE].off, (unsigned long long) e.seg[part::GATE].size, (unsigned long long) e.seg[part::UP].off, (unsigned long long) e.seg[part::UP].size, (unsigned long long) e.seg[part::DOWN].off, (unsigned long long) e.seg[part::DOWN].size, (long long) e.ne0, (long long) e.ne1, (long long) e.down_ne0, (long long) e.down_ne1); } fprintf(stderr, "[sdq] expert pager ready: %s\n" "[sdq] 熱門 expert 釘選:%d 個%s\n" "[sdq] slot=%.2f MiB (%zu B) slots=%zu arena=%.2f GiB (RAM 實測 RSS=%.2f GiB)\n", cfg_.model_path.c_str(), n_hot_, n_hot_ > 0 ? "(永不淘汰)" : "(無剖面,用動態 LRU/LFU)", slot_bytes_ / 1048576.0, slot_bytes_, n_slots_, arena_bytes_ / 1073741824.0, process_rss_bytes() / 1073741824.0); if (n_big_) { // 大小兩種槽位:把配置結果講清楚,否則只看「slots=N」會誤以為每個槽一樣大。 fprintf(stderr, "[sdq] 槽位分兩種:小 %.2f MiB × %zu + 大 %.2f MiB × %zu(平均 %.2f MiB)\n", small_bytes_ / 1048576.0, n_small_, big_bytes_ / 1048576.0, n_big_, (double) arena_bytes_ / (double) (n_slots_ ? n_slots_ : 1) / 1048576.0); } fprintf(stderr, "[sdq] prefetch=%s fadvise_dontneed=%d\n", cfg_.prefetch ? "on" : "off", (int) cfg_.fadvise_dontneed); } void pager::dump_stats(const char * phase) { if (cfg_.stats_path.empty()) { return; } const uint64_t tokens = stats_.tokens.load(); const uint64_t req = stats_.expert_requests.load(); const uint64_t hit = stats_.expert_hits.load(); const uint64_t issued = stats_.prefetch_issued.load(); const uint64_t used = stats_.prefetch_used.load(); const uint64_t b_dem = stats_.ssb_read_bytes_demand.load(); const uint64_t b_pre = stats_.ssb_read_bytes_prefetch.load(); std::ostringstream o; o << "{\"phase\":\"" << phase << "\"" << ",\"tokens\":" << tokens << ",\"expert_requests\":" << req << ",\"expert_hits\":" << hit << ",\"expert_misses\":" << (req - hit) << ",\"hit_rate\":" << pct(hit, req) << ",\"ssb_demand_bytes\":" << b_dem << ",\"ssb_prefetch_bytes\":" << b_pre << ",\"ssb_total_bytes\":" << (b_dem + b_pre) << ",\"ssd_bytes_per_token\":" << (tokens ? (double) (b_dem + b_pre) / (double) tokens : 0.0) << ",\"prefetch_issued\":" << issued << ",\"prefetch_used\":" << used << ",\"prefetch_used_rate\":" << pct(used, issued) << ",\"prefetch_duplicate\":" << stats_.prefetch_duplicate.load() << ",\"evictions\":" << stats_.evictions.load() << ",\"demand_wait_us\":" << stats_.demand_wait_us.load() << ",\"prefetch_wait_us\":" << stats_.prefetch_wait_us.load() << ",\"arena_bytes\":" << arena_bytes_ // 大小兩種槽位之後,「用掉幾個槽位」不能直接乘 slot_bytes_(那是大的那種)。 // dump_stats 只在階段轉換時呼叫一次,遍歷兩千幾個槽位的成本可以忽略。 << ",\"arena_used_bytes\":" << used_bytes_snapshot() << ",\"rss_bytes\":" << process_rss_bytes() << ",\"n_slots\":" << n_slots_ << ",\"n_small_slots\":" << n_small_ << ",\"n_big_slots\":" << n_big_ << ",\"small_slot_bytes\":" << small_bytes_ << ",\"big_slot_bytes\":" << big_bytes_ << ",\"slot_bytes\":" << slot_bytes_ << "}"; std::ofstream f(cfg_.stats_path, std::ios::app); if (f) { f << o.str() << "\n"; } } } // namespace sdq namespace sdq { // 由 CLI 呼叫:依環境變數/參數初始化 pager bool sdq_pager_active_for_layer(int il) { pager & P = pager::instance(); if (!P.ready() || P.layer(il) == nullptr) { return false; } // 除錯用:SDQ_MOE_LO / SDQ_MOE_HI 可以只讓部分層走分頁路徑, // 用來二分搜尋第一個發生數值差異的層。 static const int lo = [] { const char * e = getenv("SDQ_MOE_LO"); return e ? atoi(e) : 0; }(); static const int hi = [] { const char * e = getenv("SDQ_MOE_HI"); return e ? atoi(e) : 1 << 30; }(); return il >= lo && il < hi; } bool sdq_init(const std::string & model_path) { config c; c.model_path = model_path; c.arena_bytes = env_size("SDQ_ARENA_MB", 0) << 20; c.ram_budget_bytes = env_size("SDQ_RAM_BUDGET_MB", 0) << 20; c.slot_bytes = env_size("SDQ_SLOT_KB", 0) << 10; c.prefetch = env_flag("SDQ_PREFETCH", true); // ---- OPT-1:熱門 expert 剖面(Strata 式)---- if (const char * pi = getenv("SDQ_PROFILE_IN")) { c.profile_in = pi; } if (const char * po = getenv("SDQ_PROFILE_OUT")) { c.profile_out = po; } c.hot_pin = (int) env_size("SDQ_HOT_PIN", 0); c.admit_min_hits = (int) env_size("SDQ_ADMIT_MIN_HITS", 1); c.fadvise_dontneed = env_flag("SDQ_FADV_DONTNEED", true); c.io_parallel = env_flag("SDQ_IO_PARALLEL", true); c.io_threads = (int) env_size("SDQ_IO_THREADS", 0); c.io_split = (int) env_size("SDQ_IO_SPLIT", 1); c.verbose = env_flag("SDQ_VERBOSE", false); const char * sp = getenv("SDQ_STATS_FILE"); if (sp) { c.stats_path = sp; } sdq_set_config(c); try { pager::instance().init(c, nullptr); } catch (const std::exception & e) { fprintf(stderr, "[sdq] pager 初始化失敗:%s\n", e.what()); return false; } return pager::instance().ready(); } // ---------------------------------------------------------------- 除錯:dump 任意 F32 張量 // 目的:讓「原生 mul_mat_id 路徑」與「SDQ 分頁路徑」輸出同一層的 MoE 結果, // 兩邊各自 dump 一次就能直接 A/B 對照,判定到底是哪一邊算錯。 namespace { struct dump_arg { std::string path; }; void dump_op(ggml_tensor * dst, int ith, int /*nth*/, void * userdata) { if (ith != 0) { return; } const ggml_tensor * src = dst->src[0]; dump_arg * a = (dump_arg *) userdata; FILE * f = fopen(a->path.c_str(), "wb"); if (!f) { fprintf(stderr, "[sdq-dump] 無法開啟 %s\n", a->path.c_str()); return; } const int64_t meta[6] = { src->ne[0], src->ne[1], src->ne[2], src->ne[3], (int64_t) ggml_nbytes(src), 0 }; fwrite(meta, sizeof(int64_t), 6, f); // 只支援 F32 contiguous;其他情況就明確報錯,不要假設 if (src->type != GGML_TYPE_F32 || !ggml_is_contiguous(src)) { fprintf(stderr, "[sdq-dump] %s 不是 contiguous F32(type=%d)\n", a->path.c_str(), (int) src->type); } else { fwrite(src->data, sizeof(float), (size_t) (ggml_nelements(src)), f); } fclose(f); fprintf(stderr, "[sdq-dump] %s (%lld 個元素)\n", a->path.c_str(), (long long) ggml_nelements(src)); } } // namespace ggml_tensor * sdq_maybe_dump_tensor(ggml_context * ctx, ggml_tensor * t, int il) { const char * prefix = getenv("SDQ_MOE_OUT_DUMP"); if (!prefix || !t) { return t; } int layer = getenv("SDQ_MOE_OUT_DUMP_LAYER") ? atoi(getenv("SDQ_MOE_OUT_DUMP_LAYER")) : -1; if (layer >= 0 && il != layer) { return t; } dump_arg * a = new dump_arg{ std::string(prefix) + "." + std::to_string(il) }; return ggml_custom_4d(ctx, GGML_TYPE_F32, t->ne[0], t->ne[1], t->ne[2], t->ne[3], &t, 1, dump_op, 1, a); } void sdq_shutdown() { pager & P = pager::instance(); sdq_print_time_profile(); P.dump_stats("shutdown"); P.print_summary(); P.shutdown(); } } // namespace sdq