File size: 16,217 Bytes
bbb6388 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 | #include "strata/core/remote_experts.hpp"
#include "strata/kernels/cpu/expert_layout.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 <cstdlib>
#include <cstring>
namespace strata::core {
namespace {
constexpr int64_t H = strata::kernels::cpu::H;
constexpr int64_t FF = strata::kernels::cpu::FF;
constexpr int64_t CAP = strata::kernels::cpu::MAXT * 10;
// Small enough to copy as one pinned buffer per layer. The kernels read the
// individual arrays through pointers into the same device allocation.
struct RemoteMeta {
unsigned long long ptr[CAP];
int32_t start[CAP + 1];
int32_t dst[CAP];
int32_t tok[CAP];
int32_t count;
};
struct DeviceScope {
int previous = -1;
bool ok = false;
cudaError_t status = cudaSuccess;
const char* failed_step = nullptr;
explicit DeviceScope(int device) {
status = cudaGetDevice(&previous);
if (status != cudaSuccess) {
previous = -1;
failed_step = "cudaGetDevice";
return;
}
status = cudaSetDevice(device);
if (status != cudaSuccess) {
failed_step = "cudaSetDevice";
return;
}
ok = true;
}
~DeviceScope() { if (previous >= 0) cudaSetDevice(previous); }
std::string error(int device) const {
return std::string("CUDA") + std::to_string(device) + " experts: " +
(failed_step ? failed_step : "device switch") +
"(" + std::to_string(device) + ") failed: " + cudaGetErrorString(status) +
" (CUDA error " + std::to_string((int) status) + ")";
}
};
bool check(cudaError_t result, const char* what, std::string& err, int device) {
if (result == cudaSuccess) return true;
err = "CUDA" + std::to_string(device) + " experts: " + what + ": " + cudaGetErrorString(result);
return false;
}
} // namespace
RemoteExperts::~RemoteExperts() { close(); }
bool RemoteExperts::preflight(int device, double& free_gib, std::string& err) {
int count = 0;
if (!check(cudaGetDeviceCount(&count), "cudaGetDeviceCount", err, device)) return false;
if (device < 1 || device >= count) {
err = "CUDA" + std::to_string(device) + " experts: CUDA device is not visible";
return false;
}
// The layer waits for this GPU on the CPU pool's critical path: spin instead of sleeping, whose wake-up
// costs more than a small expert batch takes (measured: ~0.3 ms per round trip on Windows). Only possible
// before the device's context exists, so first thing; STRATA_REMOTE_SPIN=0 keeps the driver's default.
#if !defined(STRATA_USE_HIP)
const char* spin = std::getenv("STRATA_REMOTE_SPIN");
if (!(spin && spin[0] == '0')) cudaInitDevice(device, cudaDeviceScheduleSpin | cudaDeviceMapHost, 0);
cudaGetLastError();
#endif
// HIP has no cudaInitDevice equivalent; retain its default scheduling policy.
DeviceScope scope(device);
if (!scope.ok) { err = scope.error(device); return false; }
size_t free_bytes = 0, total_bytes = 0;
if (!check(cudaMemGetInfo(&free_bytes, &total_bytes), "preflight free memory", err, device)) return false;
free_gib = (double) free_bytes / 1073741824.0;
return true;
}
void RemoteExperts::close() {
if (device_ < 0) return;
DeviceScope scope(device_);
if (scope.ok) {
if (stream_) cudaStreamSynchronize(stream_);
cache_.close();
if (d_x_) cudaFree(d_x_);
if (d_out_) cudaFree(d_out_);
if (d_q8_) cudaFree(d_q8_);
if (d_scales_) cudaFree(d_scales_);
if (d_scratch_) cudaFree(d_scratch_);
if (d_meta_) cudaFree(d_meta_);
if (h_x_) cudaFreeHost(h_x_);
if (h_out_) cudaFreeHost(h_out_);
if (h_meta_) cudaFreeHost(h_meta_);
if (stream_) cudaStreamDestroy(stream_);
}
device_ = -1;
stream_ = nullptr;
h_x_ = h_out_ = d_x_ = d_out_ = nullptr;
h_meta_ = d_meta_ = nullptr;
d_q8_ = nullptr;
d_scales_ = nullptr;
d_scratch_ = nullptr;
d_start_ = d_dst_ = d_tok_ = d_count_ = nullptr;
d_ptr_ = nullptr;
original_row_.clear();
layers_present_.clear();
}
bool RemoteExperts::open(int device, int slots, int64_t layers, int64_t experts,
const std::vector<std::pair<int32_t, int32_t>>& ranked,
const ExpertCache& primary, ExpertSource& source,
std::vector<uint8_t>& claimed, std::string& err) {
close();
int count = 0;
if (!check(cudaGetDeviceCount(&count), "cudaGetDeviceCount", err, device)) return false;
if (device < 1 || device >= count || slots <= 0 || ranked.empty() ||
layers <= 0 || experts <= 0 || claimed.size() != (size_t) layers * (size_t) experts) {
err = "CUDA" + std::to_string(device) + " experts: need the device, ranked experts and positive slot count";
return false;
}
DeviceScope scope(device);
if (!scope.ok) { err = scope.error(device); return false; }
device_ = device;
n_expert_ = experts;
const auto& lay = strata::kernels::cpu::expert_layout();
std::vector<std::pair<int32_t, int32_t>> selected;
selected.reserve((size_t) slots);
std::vector<uint8_t> picked(claimed.size(), 0);
for (const auto& pair : ranked) {
if (pair.first < 0 || pair.first >= layers || pair.second < 0 || pair.second >= experts) continue;
const size_t index = (size_t) pair.first * (size_t) experts + (size_t) pair.second;
if (primary.slot_of(pair.first, pair.second) < 0 && !claimed[index] && !picked[index]) {
selected.push_back(pair);
picked[index] = 1;
if ((int) selected.size() >= slots) break;
}
}
if (selected.empty()) { err = "CUDA" + std::to_string(device) + " experts: no unclaimed experts remain"; close(); return false; }
std::vector<int64_t> sizes;
if (lay.native) {
sizes.reserve(selected.size());
for (const auto& pair : selected) sizes.push_back((int64_t) lay.blob_bytes(pair.first));
}
size_t free_bytes = 0, total_bytes = 0;
if (!check(cudaMemGetInfo(&free_bytes, &total_bytes), "free memory", err, device)) { close(); return false; }
uint64_t needed = 0;
for (const auto& pair : selected)
needed += lay.native ? (lay.blob_bytes(pair.first) + 255) / 256 * 256 : lay.max_blob;
// Leave room for the CUDA context, staging and later driver allocations, especially under WDDM.
if (needed + (512ull << 20) > free_bytes) {
err = "CUDA" + std::to_string(device) + " experts: slots leave less than 512 MiB free; reduce --expert-cache-device" + std::to_string(device);
close(); return false;
}
const bool cache_ok = lay.native ? cache_.open_sized(sizes, layers, experts, err)
: cache_.open((int64_t) selected.size(), layers, experts,
(int64_t) lay.max_blob, err);
if (!cache_ok) { err = "CUDA" + std::to_string(device) + " experts: " + err; close(); return false; }
for (const auto& pair : selected) {
const int32_t slot = cache_.admit(pair.first, pair.second);
const uint8_t* blob = source.blob(pair.first, pair.second);
if (slot < 0 || !blob || !cache_.fill_slot_blocking(slot, blob, err, (int64_t) lay.blob_bytes(pair.first))) {
err = "CUDA" + std::to_string(device) + " experts: " +
(err.empty() ? "cache fill failed" : err);
close(); return false;
}
}
const auto& first = selected.front();
if (!cache_.verify_slot(cache_.slot_of(first.first, first.second), source.blob(first.first, first.second),
err, (int64_t) lay.blob_bytes(first.first))) {
err = "CUDA" + std::to_string(device) + " experts: " + err;
close(); return false;
}
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 allocated =
check(cudaStreamCreateWithFlags(&stream_, cudaStreamNonBlocking), "stream", err, device) &&
check(cudaHostAlloc((void**) &h_x_, (size_t) CAP * H * sizeof(float), cudaHostAllocPortable | cudaHostAllocMapped), "input staging", err, device) &&
check(cudaHostAlloc((void**) &h_out_, (size_t) CAP * H * sizeof(float), cudaHostAllocPortable | cudaHostAllocMapped), "result staging", err, device) &&
check(cudaHostAlloc(&h_meta_, sizeof(RemoteMeta), cudaHostAllocPortable), "metadata staging", err, device) &&
check(cudaMalloc((void**) &d_x_, (size_t) CAP * H * sizeof(float)), "input", err, device) &&
check(cudaMalloc((void**) &d_out_, (size_t) CAP * H * sizeof(float)), "result", err, device) &&
check(cudaMalloc((void**) &d_q8_, (size_t) CAP * (H / 32) * 36), "activation", err, device) &&
check(cudaMalloc((void**) &d_scales_, (size_t) CAP * (H / 32) * sizeof(float)), "activation scales", err, device) &&
check(cudaMalloc(&d_scratch_, scratch), "scratch", err, device) &&
check(cudaMalloc(&d_meta_, sizeof(RemoteMeta)), "group metadata", err, device);
if (!allocated) { close(); return false; }
// Zero-copy: the helper reads its input from, and writes its compact rows into, the pinned host buffers
// directly - two copies fewer per layer, each of which is a PCIe round trip. STRATA_REMOTE_ZEROCOPY=0 copies.
const char* zc = std::getenv("STRATA_REMOTE_ZEROCOPY");
zero_copy_ = !(zc && zc[0] == '0') &&
cudaHostGetDevicePointer((void**) &z_x_, h_x_, 0) == cudaSuccess &&
cudaHostGetDevicePointer((void**) &z_out_, h_out_, 0) == cudaSuccess;
cudaGetLastError();
auto* meta = (RemoteMeta*) d_meta_;
d_ptr_ = meta->ptr;
d_start_ = meta->start;
d_dst_ = meta->dst;
d_tok_ = meta->tok;
d_count_ = &meta->count;
owned_.resize(CAP);
layers_present_.assign((size_t) layers, 0);
group_of_.resize(CAP);
group_id_.reserve(CAP);
ptr_.reserve(CAP);
start_.reserve(CAP + 1);
dst_.reserve(CAP);
tok_.reserve(CAP);
original_row_.reserve(CAP);
computed_ = 0;
launched_layers_ = 0;
returned_bytes_ = full_row_bytes_ = 0;
for (const auto& pair : selected) {
layers_present_[(size_t) pair.first] = 1;
claimed[(size_t) pair.first * (size_t) experts + (size_t) pair.second] = 1;
}
return true;
}
bool RemoteExperts::begin(int64_t layer, const float* x, const int32_t* ids, int64_t n_tok,
int64_t k, const int32_t* kind, const int32_t* primary_res,
std::string& err) {
// cumulative host time in here (staging and launches), reported per request by the driver
struct Timer { double& acc; std::chrono::steady_clock::time_point t0 = std::chrono::steady_clock::now();
~Timer() { acc += std::chrono::duration<double, std::milli>(std::chrono::steady_clock::now() - t0).count(); } } timer{ms_begin_};
const int64_t n = n_tok * k;
if (n <= 0 || n > CAP || n_tok > strata::kernels::cpu::MAXT || k != 10 || layer < 0 ||
(size_t) layer >= layers_present_.size() || device_ < 0) {
err = "CUDA" + std::to_string(device_) + " experts: invalid layer, routing width or window size";
return false;
}
std::fill(owned_.begin(), owned_.begin() + n, 0);
group_id_.clear(); ptr_.clear(); start_.clear(); dst_.clear(); tok_.clear(); original_row_.clear();
if (!layers_present_[(size_t) layer]) return true;
std::fill(group_of_.begin(), group_of_.begin() + n, -1);
for (int64_t i = 0; i < n; ++i) {
const int32_t e = ids[i];
if ((kind && kind[i] != -1) || e < 0 || e >= n_expert_ ||
(primary_res && primary_res[(size_t) layer * (size_t) n_expert_ + (size_t) e] >= 0)) continue;
const int32_t slot = cache_.slot_of(layer, e);
if (slot < 0) continue;
owned_[(size_t) i] = 1;
int32_t group = -1;
for (size_t g = 0; g < group_id_.size(); ++g)
if (group_id_[g] == e) { group = (int32_t) g; break; }
if (group < 0) {
group = (int32_t) group_id_.size();
group_id_.push_back(e);
ptr_.push_back((unsigned long long) cache_.device_slot(slot));
}
group_of_[(size_t) i] = group;
++computed_;
}
if (group_id_.empty()) return true;
for (size_t g = 0; g < group_id_.size(); ++g) {
start_.push_back((int32_t) dst_.size());
for (int64_t i = 0; i < n; ++i) if (group_of_[(size_t) i] == (int32_t) g) {
// The grouped kernels read the activation from `tok`, so their
// output row can instead be packed densely for the USB4 return.
original_row_.push_back((int32_t) i);
dst_.push_back((int32_t) dst_.size());
tok_.push_back((int32_t) (i / k));
}
}
start_.push_back((int32_t) dst_.size());
// Private pinned buffers survive until this GPU has consumed them. The CPU
// pool can write other output rows without a cross-device race.
std::memcpy(h_x_, x, (size_t) n_tok * H * sizeof(float));
auto* meta = (RemoteMeta*) h_meta_;
std::memcpy(meta->ptr, ptr_.data(), ptr_.size() * sizeof(ptr_[0]));
std::memcpy(meta->start, start_.data(), start_.size() * sizeof(start_[0]));
std::memcpy(meta->dst, dst_.data(), dst_.size() * sizeof(dst_[0]));
std::memcpy(meta->tok, tok_.data(), tok_.size() * sizeof(tok_[0]));
meta->count = (int32_t) group_id_.size();
DeviceScope scope(device_);
if (!scope.ok) { err = scope.error(device_); return false; }
groups_ = (int32_t) group_id_.size();
const cudaStream_t s = stream_;
const bool staged =
(zero_copy_ || check(cudaMemcpyAsync(d_x_, h_x_, (size_t) n_tok * H * sizeof(float), cudaMemcpyHostToDevice, s), "copy input", err, device_)) &&
check(cudaMemcpyAsync(d_meta_, h_meta_, sizeof(RemoteMeta), cudaMemcpyHostToDevice, s), "copy group metadata", err, device_);
if (!staged) return false;
const auto& lay = strata::kernels::cpu::expert_layout();
if (lay.native) {
strata::kernels::quantize_q8_1_rows(zero_copy_ ? z_x_ : d_x_, n_tok, H, d_q8_, s);
const auto& fmt = lay.fmt[(size_t) layer];
auto L = strata::kernels::native_expert_layout(fmt.gu_type, fmt.d_type, fmt.n_embd, fmt.n_ff);
strata::kernels::native_expert_grouped(L, d_ptr_, d_start_, d_count_, d_dst_, d_tok_,
groups_, (int64_t) dst_.size(), d_q8_, d_scratch_, zero_copy_ ? z_out_ : d_out_, s);
} else {
strata::kernels::quantize_q8_0_scaled(zero_copy_ ? z_x_ : d_x_, d_q8_, d_scales_, n_tok * H, s);
strata::kernels::moe_grouped_s2(d_ptr_, d_start_, d_count_, d_dst_, d_tok_,
groups_, (int64_t) dst_.size(), d_q8_, d_scales_, d_scratch_, zero_copy_ ? z_out_ : d_out_, s);
}
const uint64_t compact_bytes = (uint64_t) dst_.size() * H * sizeof(float);
if (!zero_copy_ && !check(cudaMemcpyAsync(h_out_, d_out_, (size_t) compact_bytes, cudaMemcpyDeviceToHost, s),
"copy results", err, device_)) return false;
++launched_layers_;
returned_bytes_ += compact_bytes;
full_row_bytes_ += (uint64_t) n * H * sizeof(float);
return true;
}
bool RemoteExperts::finish(float* out, std::string& err) {
if (group_id_.empty()) return true;
DeviceScope scope(device_);
if (!scope.ok) { err = scope.error(device_); return false; }
const auto w0 = std::chrono::steady_clock::now();
if (!check(cudaStreamSynchronize(stream_), "finish", err, device_)) return false;
ms_wait_ += std::chrono::duration<double, std::milli>(std::chrono::steady_clock::now() - w0).count();
for (size_t i = 0; i < original_row_.size(); ++i)
std::memcpy(out + (size_t) original_row_[i] * H, h_out_ + i * H, (size_t) H * sizeof(float));
return true;
}
} // namespace strata::core
|