Winery-Strata / src /core /peer_experts.cpp
Zandy-Wandy's picture
🍷 Winery Strata: fork of Niko1221/Strata (Winery Cuvée family, wine-red UI)
bbb6388 verified
Raw History Blame Contribute Delete
15.2 kB
// 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 <algorithm>
#include <chrono>
#include <cstdio>
#include <cstring>
#include <cstdlib>
#include <stdexcept>
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