// src/core/peer_experts.cpp - see include/strata/core/peer_experts.hpp. #include "strata/core/peer_experts.hpp" #include "strata/kernels/cpu/expert_layout.hpp" #include "strata/kernels/elementwise.hpp" #include "strata/kernels/iq_kernels.hpp" #include "strata/kernels/quantize_act.hpp" #include "strata/kernels/s2_expert_grouped.hpp" #include #include #include #include #include #include namespace strata::core { static bool g_peer_portable = false; void set_peer_portable(bool on) { g_peer_portable = on; } bool peer_portable() { return g_peer_portable; } namespace { constexpr int64_t H = strata::kernels::cpu::H; constexpr int64_t CAP = strata::kernels::cpu::MAXT * 10; // entries per layer: MAXT tokens x top-10 struct Meta { // one block, uploaded with one copy per layer unsigned long long ptr[CAP]; int32_t start[CAP + 1]; int32_t dst[CAP]; int32_t tok[CAP]; int32_t count[4]; // [0] groups, [1] entries int32_t out_row[CAP]; // multi-GPU: compact row -> entry row in the host `out` }; /// Switches the calling thread to `dev` and back. The engine's own thread lives on device 0. struct On { int prev = 0; explicit On(int dev) { cudaGetDevice(&prev); if (prev != dev) cudaSetDevice(dev); } ~On() { int cur = 0; cudaGetDevice(&cur); if (cur != prev) cudaSetDevice(prev); } }; bool ck(cudaError_t e, const char* what, std::string& err) { if (e == cudaSuccess) return true; err = std::string("peer experts: ") + what + ": " + cudaGetErrorString(e); return false; } } // namespace PeerExperts::~PeerExperts() { close(); } void PeerExperts::close() { if (device_ < 0) return; { On on(device_); if (stream_) cudaStreamSynchronize(stream_); if (refill_) cudaStreamSynchronize(refill_); cache_.close(); if (d_x_) cudaFree(d_x_); if (d_out_) cudaFree(d_out_); if (d_meta_) cudaFree(d_meta_); if (d_q8_) cudaFree(d_q8_); if (d_scratch_) cudaFree(d_scratch_); if (h_x_) cudaFreeHost(h_x_); if (h_out_) cudaFreeHost(h_out_); if (h_meta_) cudaFreeHost(h_meta_); if (refill_ev_) cudaEventDestroy(refill_ev_); if (stream_) cudaStreamDestroy(stream_); if (refill_) cudaStreamDestroy(refill_); } d_x_ = d_out_ = h_x_ = h_out_ = nullptr; d_meta_ = h_meta_ = d_scratch_ = nullptr; d_q8_ = nullptr; stream_ = refill_ = nullptr; refill_ev_ = nullptr; device_ = -1; } bool PeerExperts::open(int device, const std::vector>& ranked, const ExpertCache& primary, ExpertSource& src, int64_t n_layers, int64_t n_expert, int reserve_mib, int64_t max_slots, std::string& err) { close(); int count = 0; if (!ck(cudaGetDeviceCount(&count), "cudaGetDeviceCount", err)) return false; if (device < 1 || device >= count) { err = "peer experts: CUDA device " + std::to_string(device) + " is not visible (" + std::to_string(count) + " devices; set CUDA_VISIBLE_DEVICES)"; return false; } const auto& lay = strata::kernels::cpu::expert_layout(); n_layers_ = n_layers; n_expert_ = n_expert; src_ = &src; res_.assign((size_t) (n_layers * n_expert), kNotResident); { // direct access both ways (NVLink or another P2P path): the prompt path copies activations and results over it int a = 0, b = 0; cudaDeviceCanAccessPeer(&a, 0, device); cudaDeviceCanAccessPeer(&b, device, 0); p2p_ = a && b; if (p2p_) { On on0(0); cudaError_t e0 = cudaDeviceEnablePeerAccess(device, 0); if (e0 == cudaErrorPeerAccessAlreadyEnabled) { cudaGetLastError(); e0 = cudaSuccess; } On on1(device); cudaError_t e1 = cudaDeviceEnablePeerAccess(0, 0); if (e1 == cudaErrorPeerAccessAlreadyEnabled) { cudaGetLastError(); e1 = cudaSuccess; } p2p_ = e0 == cudaSuccess && e1 == cudaSuccess; } } On on(device); device_ = device; // scratch first, so the slots take what is really left int64_t ff = strata::kernels::cpu::FF; if (lay.native) for (const auto& f : lay.fmt) ff = std::max(ff, f.n_ff); const size_t scratch = std::max((size_t) strata::kernels::moe_hit_grouped_scratch_bytes(CAP, H, ff), strata::kernels::native_expert_scratch_bytes(CAP, ff)); const bool alloc_ok = ck(cudaStreamCreateWithFlags(&stream_, cudaStreamNonBlocking), "stream", err) && ck(cudaStreamCreateWithFlags(&refill_, cudaStreamNonBlocking), "refill stream", err) && ck(cudaEventCreateWithFlags(&refill_ev_, cudaEventDisableTiming), "event", err) && ck(cudaHostAlloc((void**) &h_x_, (size_t) CAP * H * sizeof(float), cudaHostAllocPortable | cudaHostAllocMapped), "input staging", err) && ck(cudaHostAlloc((void**) &h_out_, (size_t) CAP * H * sizeof(float), cudaHostAllocPortable | cudaHostAllocMapped), "result staging", err) && ck(cudaHostAlloc(&h_meta_, sizeof(Meta), cudaHostAllocPortable | cudaHostAllocMapped), "plan staging", err) && ck(cudaMalloc((void**) &d_x_, (size_t) CAP * H * sizeof(float)), "input", err) && ck(cudaMalloc((void**) &d_out_, (size_t) CAP * H * sizeof(float)), "result", err) && ck(cudaMalloc(&d_meta_, sizeof(Meta)), "plan", err) && ck(cudaMalloc((void**) &d_q8_, (size_t) CAP * (H / 32) * 36), "activations", err) && ck(cudaMalloc(&d_scratch_, scratch), "scratch", err); if (!alloc_ok) { close(); return false; } // the pairs the primary does not hold, in rank order, as many as fit size_t free_b = 0, total_b = 0; if (!ck(cudaMemGetInfo(&free_b, &total_b), "cudaMemGetInfo", err)) { close(); return false; } const uint64_t reserve = (uint64_t) std::max(reserve_mib, 128) << 20; const uint64_t budget = free_b > reserve ? free_b - reserve : 0; std::vector> pick; std::vector sizes; uint64_t used = 0; for (const auto& pr : ranked) { if (primary.slot_of(pr.first, pr.second) >= 0) continue; const uint64_t b = lay.blob_bytes(pr.first); const uint64_t b256 = lay.native ? (b + 255) / 256 * 256 : lay.max_blob; if (used + b256 > budget) break; if (max_slots > 0 && (int64_t) pick.size() >= max_slots) break; used += b256; pick.push_back(pr); sizes.push_back((int64_t) b); } if (pick.empty()) { err = "peer experts: no room or no expert left for the peer"; close(); return false; } const bool opened = lay.native ? cache_.open_sized(sizes, n_layers, n_expert, err) : cache_.open((int64_t) pick.size(), n_layers, n_expert, (int64_t) lay.max_blob, err); if (!opened) { err = "peer experts: " + err; close(); return false; } // open() zeroes the arena with cudaMemset on the legacy stream, which is not ordered against the non-blocking // refill stream: wait for it, or the zeroing can land on top of the fills if (!ck(cudaDeviceSynchronize(), "arena zeroing", err)) { close(); return false; } for (const auto& pr : pick) { const int32_t slot = cache_.admit(pr.first, pr.second); const uint8_t* b = src.blob(pr.first, pr.second); if (slot < 0 || b == nullptr || !cache_.fill_slot(slot, b, refill_, err, (int64_t) lay.blob_bytes(pr.first))) { err = "peer experts: fill failed: " + err; close(); return false; } res_[(size_t) (pr.first * n_expert + pr.second)] = slot; } if (!ck(cudaStreamSynchronize(refill_), "fill", err)) { close(); return false; } const auto& f0 = pick.front(); if (!cache_.verify_slot(res_[(size_t) (f0.first * n_expert + f0.second)], src.blob(f0.first, f0.second), err, (int64_t) lay.blob_bytes(f0.first))) { err = "peer experts: " + err; close(); return false; } resident_ = (int64_t) pick.size(); return true; } bool PeerExperts::launch(int64_t layer, const float* x, const int32_t* ids, int64_t n_tok, int64_t k, const int32_t* kind, std::string& err, float* out) { static const bool direct_env = [] { const char* v = std::getenv("STRATA_PEER_DIRECT"); return v == nullptr || std::atoi(v) != 0; }(); const bool direct = direct_env && out != nullptr; launched_direct_ = false; launched_rows_ = 0; row_of_.clear(); const int64_t n = n_tok * k; if (n > CAP) { err = "peer experts: window too large"; return false; } Meta& m = *(Meta*) h_meta_; int groups = 0, rows = 0; // distinct experts in routing order; each one's entries become compact rows for (int64_t i = 0; i < n; ++i) { if (kind[i] != 2) continue; bool seen = false; for (int64_t j = 0; j < i; ++j) if (kind[j] == 2 && ids[j] == ids[i]) { seen = true; break; } if (seen) continue; const int32_t slot = res_[(size_t) (layer * n_expert_ + ids[i])]; if (slot < 0) { err = "peer experts: an entry was planned for the peer but is not resident"; return false; } m.ptr[groups] = (unsigned long long) cache_.device_slot(slot); m.start[groups] = rows; for (int64_t j = i; j < n; ++j) if (kind[j] == 2 && ids[j] == ids[i]) { m.dst[rows] = rows; m.tok[rows] = (int32_t) (j / k); m.out_row[rows] = (int32_t) j; row_of_.push_back((int32_t) j); ++rows; } ++groups; ++experts_; } if (groups == 0) return true; m.start[groups] = rows; m.count[0] = groups; m.count[1] = rows; entries_ += rows; launched_rows_ = rows; // direct: x is the verifier's pinned (portable) doorbell row block - copied from where it is if (!direct) std::memcpy(h_x_, x, (size_t) (n_tok * H) * sizeof(float)); On on(device_); const cudaStream_t s = stream_; if (!ck(cudaMemcpyAsync(d_x_, direct ? x : h_x_, (size_t) (n_tok * H) * sizeof(float), cudaMemcpyHostToDevice, s), "input", err) || !ck(cudaMemcpyAsync(d_meta_, h_meta_, sizeof(Meta), cudaMemcpyHostToDevice, s), "plan", err)) return false; Meta* dm = (Meta*) d_meta_; const auto& lay = strata::kernels::cpu::expert_layout(); if (lay.native) { strata::kernels::quantize_q8_1_rows(d_x_, n_tok, H, d_q8_, s); const auto& f = lay.fmt[(size_t) layer]; const auto L = strata::kernels::native_expert_layout(f.gu_type, f.d_type, f.n_embd, f.n_ff); strata::kernels::native_expert_grouped(L, dm->ptr, dm->start, dm->count, dm->dst, dm->tok, groups, rows, d_q8_, d_scratch_, d_out_, s); } else { err = "peer experts: only native packs are supported"; return false; } if (direct) { // the rows straight into the host `out` (zero-copy, coalesced float4 writes) try { strata::kernels::scatter_rows_f32(d_out_, out, dm->out_row, rows, H, s); } catch (const std::exception& e) { err = std::string("peer experts: scatter: ") + e.what(); return false; } launched_direct_ = true; return true; } return ck(cudaMemcpyAsync(h_out_, d_out_, (size_t) rows * H * sizeof(float), cudaMemcpyDeviceToHost, s), "results", err); } bool PeerExperts::finish(float* out, std::string& err) { if (launched_rows_ == 0) return true; const auto t0 = std::chrono::steady_clock::now(); { On on(device_); // spin on the stream (a blocking sync would sleep the pool thread and wake it late) cudaError_t e; while ((e = cudaStreamQuery(stream_)) == cudaErrorNotReady) { } if (!ck(e, "finish", err)) return false; } ms_wait += std::chrono::duration(std::chrono::steady_clock::now() - t0).count(); if (launched_direct_) { launched_rows_ = 0; return true; } for (size_t r = 0; r < row_of_.size(); ++r) std::memcpy(out + (size_t) row_of_[r] * H, h_out_ + r * H, (size_t) H * sizeof(float)); launched_rows_ = 0; return true; } bool PeerExperts::adapt(const float* usage, const int32_t* res0, int max_swaps, std::string& err) { if (!pending_.empty() || max_swaps <= 0) return true; const auto& lay = strata::kernels::cpu::expert_layout(); struct Swap { float gain; int32_t layer, in, out; }; std::vector swaps; std::vector> cand, vict; for (int64_t l = 0; l < n_layers_; ++l) { cand.clear(); vict.clear(); const float* u = usage + l * n_expert_; const int32_t* r0 = res0 + l * n_expert_; const int32_t* r1 = res_.data() + l * n_expert_; for (int32_t e = 0; e < (int32_t) n_expert_; ++e) { if (r1[e] >= 0) vict.emplace_back(u[e], e); else if (r0[e] < 0 && u[e] >= 2.0f) cand.emplace_back(u[e], e); } if (cand.empty() || vict.empty()) continue; std::sort(cand.begin(), cand.end(), [](auto& a, auto& b) { return a.first > b.first; }); const size_t nc = std::min(cand.size(), vict.size()); std::partial_sort(vict.begin(), vict.begin() + (ptrdiff_t) nc, vict.end(), [](auto& a, auto& b) { return a.first < b.first; }); for (size_t i = 0; i < nc; ++i) { if (cand[i].first < vict[i].first + 1.5f) break; swaps.push_back({cand[i].first - vict[i].first, (int32_t) l, cand[i].second, vict[i].second}); } } if (swaps.empty()) return true; std::sort(swaps.begin(), swaps.end(), [](const Swap& a, const Swap& b) { return a.gain > b.gain; }); if ((int) swaps.size() > max_swaps) swaps.resize((size_t) max_swaps); On on(device_); for (const Swap& s : swaps) { const size_t in = (size_t) (s.layer * n_expert_ + s.in), out = (size_t) (s.layer * n_expert_ + s.out); const int32_t slot = res_[out]; const uint8_t* b = src_->blob(s.layer, s.in); if (slot < 0 || b == nullptr) continue; if (!ck(cudaMemcpyAsync(cache_.device_slot(slot), b, (size_t) lay.blob_bytes(s.layer), cudaMemcpyHostToDevice, refill_), "refill", err)) return false; res_[out] = kNotResident; // evicted now: the CPU computes it meanwhile pending_.emplace_back((int32_t) in, slot); ++swaps_; } return ck(cudaEventRecord(refill_ev_, refill_), "refill event", err); } void PeerExperts::apply_pending(bool wait) { if (pending_.empty()) return; { On on(device_); if (wait) cudaEventSynchronize(refill_ev_); else if (cudaEventQuery(refill_ev_) != cudaSuccess) return; } for (const auto& [i, slot] : pending_) res_[(size_t) i] = slot; pending_.clear(); } } // namespace strata::core