Download src/core/peer_experts.cpp from WineryLabs/Winery-Strata: direct link, hf CLI and curl.
- Browser
- Download file 15.2 kB
-
https://huggingface.co/WineryLabs/Winery-Strata/resolve/main/src/core/peer_experts.cpp
- Command line
-
hf download hf://WineryLabs/Winery-Strata/src/core/peer_experts.cpp
-
curl -L -o peer_experts.cpp https://huggingface.co/WineryLabs/Winery-Strata/resolve/main/src/core/peer_experts.cpp
15.2 kB
| // src/core/peer_experts.cpp - see include/strata/core/peer_experts.hpp. | |
| 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<std::pair<int32_t, int32_t>>& 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<int64_t>(ff, f.n_ff); | |
| const size_t scratch = std::max<size_t>((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<std::pair<int32_t, int32_t>> pick; | |
| std::vector<int64_t> 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<double, std::milli>(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<Swap> swaps; | |
| std::vector<std::pair<float, int32_t>> 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 | |