Download src/core/expert_source.cpp from WineryLabs/Winery-Strata: direct link, hf CLI and curl.
- Browser
- Download file 131 kB
-
https://huggingface.co/WineryLabs/Winery-Strata/resolve/main/src/core/expert_source.cpp
- Command line
-
hf download hf://WineryLabs/Winery-Strata/src/core/expert_source.cpp
-
curl -L -o expert_source.cpp https://huggingface.co/WineryLabs/Winery-Strata/resolve/main/src/core/expert_source.cpp
131 kB
| // src/core/expert_source.cpp - the adapter. See the header for the three clauses of the contract. | |
| // a 64-bit seek (as in pinned.cu): the 32-bit `fseek` wraps past 4 GiB, and the spelling differs per platform | |
| namespace strata::core { | |
| namespace detail { | |
| bool cgroup_available_bytes(uint64_t limit, const CgroupMemoryStat& stat, uint64_t& bytes) { | |
| bytes = 0; | |
| if (!stat.valid) return false; | |
| // memory.stat's inactive_file can race memory.current, so bound it to charged usage first. | |
| uint64_t reclaimable = std::min(stat.inactive_file, stat.current); | |
| reclaimable = stat.file_dirty >= reclaimable ? 0 : reclaimable - stat.file_dirty; | |
| reclaimable = stat.file_writeback >= reclaimable ? 0 : reclaimable - stat.file_writeback; | |
| // Reclaiming clean file pages reduces usage; saturating subtraction also handles a transient over-limit read. | |
| const uint64_t usage_after_reclaim = stat.current - reclaimable; | |
| bytes = usage_after_reclaim < limit ? limit - usage_after_reclaim : 0; | |
| return true; | |
| } | |
| bool make_cache_complement_plan( | |
| int64_t n_layers, int64_t n_expert, const std::vector<uint64_t>& layer_blob_bytes, | |
| const std::vector<std::pair<int32_t, int32_t>>& primary_gpu_pairs, | |
| const std::vector<std::pair<int32_t, int32_t>>& additional_gpu_pairs, | |
| std::vector<uint64_t>& offsets, uint64_t& bytes, std::string& err) { | |
| offsets.clear(); | |
| bytes = 0; | |
| err.clear(); | |
| if (n_layers <= 0 || n_expert <= 0 || layer_blob_bytes.size() != (size_t) n_layers) { | |
| err = "FileExpertSource: invalid geometry for the cache complement plan"; | |
| return false; | |
| } | |
| if ((uint64_t) n_layers > (uint64_t) std::numeric_limits<size_t>::max() / (uint64_t) n_expert) { | |
| err = "FileExpertSource: cache complement index table is too large"; | |
| return false; | |
| } | |
| const size_t count = (size_t) n_layers * (size_t) n_expert; | |
| std::vector<uint8_t> omitted(count, 0); | |
| auto mark_pairs = [&](const std::vector<std::pair<int32_t, int32_t>>& pairs, uint8_t bit, | |
| const char* label) -> bool { | |
| for (const auto& pair : pairs) { | |
| if (pair.first < 0 || pair.second < 0 || pair.first >= n_layers || pair.second >= n_expert) { | |
| err = std::string("FileExpertSource: ") + label + " pair is outside the expert geometry"; | |
| return false; | |
| } | |
| const size_t index = (size_t) pair.first * (size_t) n_expert + (size_t) pair.second; | |
| if ((omitted[index] & bit) != 0) { | |
| err = std::string("FileExpertSource: duplicate ") + label + " pair in the cache complement plan"; | |
| return false; | |
| } | |
| if (bit == 2 && (omitted[index] & 1) != 0) { | |
| err = "FileExpertSource: the primary and additional GPU expert tiers overlap"; | |
| return false; | |
| } | |
| omitted[index] |= bit; | |
| } | |
| return true; | |
| }; | |
| if (!mark_pairs(primary_gpu_pairs, 1, "primary GPU") || | |
| !mark_pairs(additional_gpu_pairs, 2, "additional GPU")) return false; | |
| for (uint64_t blob_bytes : layer_blob_bytes) { | |
| if (blob_bytes == 0) { | |
| err = "FileExpertSource: cache complement layer has zero-sized expert blobs"; | |
| return false; | |
| } | |
| } | |
| offsets.assign(count, kNoCacheComplement); | |
| for (int64_t layer = 0; layer < n_layers; ++layer) { | |
| const uint64_t blob_bytes = layer_blob_bytes[(size_t) layer]; | |
| for (int64_t expert = 0; expert < n_expert; ++expert) { | |
| const size_t index = (size_t) layer * (size_t) n_expert + (size_t) expert; | |
| if (omitted[index] != 0) continue; | |
| if (bytes > std::numeric_limits<uint64_t>::max() - blob_bytes) { | |
| offsets.clear(); | |
| bytes = 0; | |
| err = "FileExpertSource: cache complement size overflows"; | |
| return false; | |
| } | |
| offsets[index] = bytes; | |
| bytes += blob_bytes; | |
| } | |
| } | |
| if (bytes > (uint64_t) std::numeric_limits<size_t>::max()) { | |
| offsets.clear(); | |
| bytes = 0; | |
| err = "FileExpertSource: cache complement exceeds the host address space"; | |
| return false; | |
| } | |
| return true; | |
| } | |
| const uint8_t* cache_complement_blob_or_fallback( | |
| size_t index, const std::vector<uint64_t>& offsets, const uint8_t* complement_host, | |
| const uint8_t* mapped_fallback) { | |
| if (complement_host != nullptr && index < offsets.size() && offsets[index] != kNoCacheComplement) | |
| return complement_host + (size_t) offsets[index]; | |
| return mapped_fallback; | |
| } | |
| int64_t choose_resident_keep_from(const std::vector<uint64_t>& slot_bytes, uint64_t base_bytes, uint64_t budget, | |
| int64_t lend_from) { | |
| if (base_bytes > budget) return -1; | |
| const int64_t slots = (int64_t) slot_bytes.size(); | |
| if (lend_from < 0 || lend_from > slots) lend_from = slots; // no lend region: only the experts no slot holds | |
| int64_t keep = slots; | |
| uint64_t bytes = base_bytes; | |
| while (keep > lend_from) { | |
| const uint64_t b = slot_bytes[(size_t) keep - 1]; | |
| if (b > budget - bytes) break; | |
| bytes += b; | |
| --keep; | |
| } | |
| return keep; | |
| } | |
| bool exchange_cache_complement(std::vector<uint64_t>& offsets, size_t in, size_t out) { | |
| if (in == out || in >= offsets.size() || out >= offsets.size() || offsets[in] == kNoCacheComplement || | |
| offsets[out] != kNoCacheComplement) return false; | |
| offsets[out] = offsets[in]; | |
| offsets[in] = kNoCacheComplement; | |
| return true; | |
| } | |
| } // namespace detail | |
| namespace { | |
| bool read_cgroup_memory_stat(const std::filesystem::path& path, uint64_t current, | |
| detail::CgroupMemoryStat& stat) { | |
| std::ifstream input(path / "memory.stat"); | |
| if (!input) return false; | |
| bool inactive_file = false, file_dirty = false, file_writeback = false; | |
| std::string line; | |
| while (std::getline(input, line)) { | |
| std::istringstream fields(line); | |
| std::string key; | |
| uint64_t value = 0; | |
| if (!(fields >> key >> value)) return false; | |
| fields >> std::ws; | |
| if (!fields.eof()) return false; | |
| if (key == "inactive_file") { | |
| if (inactive_file) return false; | |
| inactive_file = true; | |
| stat.inactive_file = value; | |
| } else if (key == "file_dirty") { | |
| if (file_dirty) return false; | |
| file_dirty = true; | |
| stat.file_dirty = value; | |
| } else if (key == "file_writeback") { | |
| if (file_writeback) return false; | |
| file_writeback = true; | |
| stat.file_writeback = value; | |
| } | |
| } | |
| if (!input.eof() || !inactive_file || !file_dirty || !file_writeback) return false; | |
| stat.current = current; | |
| stat.valid = true; | |
| return true; | |
| } | |
| bool available_memory_bytes(uint64_t& bytes) { | |
| MEMORYSTATUSEX status{}; | |
| status.dwLength = sizeof(status); | |
| if (!GlobalMemoryStatusEx(&status)) return false; | |
| bytes = (uint64_t) status.ullAvailPhys; | |
| return bytes > 0; | |
| // MemAvailable includes reclaimable page cache, unlike _SC_AVPHYS_PAGES. | |
| std::ifstream info("/proc/meminfo"); | |
| std::string line; | |
| bytes = 0; | |
| while (std::getline(info, line)) { | |
| std::istringstream fields(line); | |
| std::string key, unit; | |
| uint64_t value = 0; | |
| if (fields >> key >> value >> unit && key == "MemAvailable:" && unit == "kB" && | |
| value <= std::numeric_limits<uint64_t>::max() / 1024) bytes = value * 1024; | |
| } | |
| if (bytes == 0) return false; | |
| // Account for the tightest cgroup-v2 ancestor limit when its normal mount is visible. | |
| // This is a point-in-time guard, not a reservation against concurrent allocations. | |
| std::ifstream groups("/proc/self/cgroup"); | |
| if (!groups) return false; | |
| bool resolved_v2 = false; | |
| while (std::getline(groups, line)) { | |
| if (line.rfind("0::/", 0) != 0) continue; | |
| const std::filesystem::path root("/sys/fs/cgroup"); | |
| auto path = (root / line.substr(4)).lexically_normal(); | |
| if (path.string().rfind(root.string(), 0) != 0 || !std::filesystem::is_directory(path)) return false; | |
| resolved_v2 = true; | |
| while (path.string().rfind(root.string(), 0) == 0) { | |
| std::ifstream limit_file(path / "memory.max"), current_file(path / "memory.current"); | |
| std::string limit; | |
| uint64_t current = 0; | |
| const bool readable = bool(limit_file >> limit) && bool(current_file >> current); | |
| // The host's root cgroup has no memory.max; ordinary child groups must expose their limits. | |
| if (!readable && !(path == root && !std::filesystem::exists(path / "memory.max") && | |
| std::filesystem::exists(path / "cgroup.controllers"))) return false; | |
| if (readable && limit != "max") { | |
| try { | |
| size_t consumed = 0; | |
| const uint64_t cap = std::stoull(limit, &consumed); | |
| if (consumed != limit.size()) return false; | |
| detail::CgroupMemoryStat stat; | |
| if (!read_cgroup_memory_stat(path, current, stat)) return false; | |
| uint64_t cgroup_available = 0; | |
| if (!detail::cgroup_available_bytes(cap, stat, cgroup_available)) return false; | |
| bytes = std::min(bytes, cgroup_available); | |
| } catch (...) { return false; } | |
| } | |
| if (path == root) break; | |
| path = path.parent_path(); | |
| } | |
| } | |
| return resolved_v2; | |
| const long pages = sysconf(_SC_AVPHYS_PAGES); | |
| const long page_bytes = sysconf(_SC_PAGESIZE); | |
| if (pages <= 0 || page_bytes <= 0 || | |
| (uint64_t) pages > std::numeric_limits<uint64_t>::max() / (uint64_t) page_bytes) return false; | |
| bytes = (uint64_t) pages * (uint64_t) page_bytes; | |
| return bytes > 0; | |
| } | |
| } // namespace | |
| // ================================ THE FILE-BACKED SOURCE ================================ | |
| FileExpertSource::~FileExpertSource() { close(); } | |
| bool FileExpertSource::open(const std::string& pack_dir, int64_t n_layers, int64_t n_expert, std::string& err) { | |
| close(); | |
| if (n_layers <= 0 || n_expert <= 0) { err = "FileExpertSource: the geometry is empty"; return false; } | |
| const auto& layout = strata::kernels::cpu::expert_layout(); | |
| if (layout.n_layers != n_layers || layout.n_expert != n_expert) { | |
| err = "FileExpertSource: the requested geometry does not match the loaded expert layout"; | |
| return false; | |
| } | |
| if ((uint64_t) n_layers > (uint64_t) std::numeric_limits<int64_t>::max() / (uint64_t) n_expert) { | |
| err = "FileExpertSource: the expert count overflows"; | |
| return false; | |
| } | |
| const uint64_t blob_count = (uint64_t) n_layers * (uint64_t) n_expert; | |
| if (blob_count > (uint64_t) std::numeric_limits<int64_t>::max() || | |
| (uint64_t) n_layers > (uint64_t) std::numeric_limits<size_t>::max()) { | |
| err = "FileExpertSource: the expert count overflows"; | |
| return false; | |
| } | |
| std::vector<uint64_t> layer_offsets((size_t) n_layers), layer_blob_bytes((size_t) n_layers); | |
| const uint64_t want = layout.total; | |
| if (want == 0 || want > (uint64_t) std::numeric_limits<size_t>::max()) { | |
| err = "FileExpertSource: the loaded expert layout has an invalid size"; | |
| return false; | |
| } | |
| if (!layout.native) { | |
| if (blob_count > std::numeric_limits<uint64_t>::max() / (uint64_t) strata::kernels::cpu::BLOB) { | |
| err = "FileExpertSource: the canonical expert size overflows"; | |
| return false; | |
| } | |
| const uint64_t canonical_size = blob_count * (uint64_t) strata::kernels::cpu::BLOB; | |
| if (want != canonical_size) { | |
| err = "FileExpertSource: the canonical expert layout has an inconsistent size"; | |
| return false; | |
| } | |
| const uint64_t bytes = (uint64_t) strata::kernels::cpu::BLOB; | |
| const uint64_t layer_bytes = (uint64_t) n_expert * bytes; | |
| for (int64_t layer = 0; layer < n_layers; ++layer) { | |
| layer_offsets[(size_t) layer] = (uint64_t) layer * layer_bytes; | |
| layer_blob_bytes[(size_t) layer] = bytes; | |
| } | |
| } else { | |
| if (layout.offset.size() != (size_t) n_layers || layout.bytes.size() != (size_t) n_layers || | |
| layout.fmt.size() != (size_t) n_layers) { | |
| err = "FileExpertSource: the native expert layout is incomplete"; | |
| return false; | |
| } | |
| uint64_t at = 0; | |
| for (int64_t layer = 0; layer < n_layers; ++layer) { | |
| const size_t i = (size_t) layer; | |
| const uint64_t bytes = (uint64_t) layout.fmt[i].bytes; | |
| if (layout.offset[i] != at || bytes == 0 || layout.bytes[i] != bytes || | |
| bytes > std::numeric_limits<uint64_t>::max() / (uint64_t) n_expert) { | |
| err = "FileExpertSource: the native expert layout is invalid at layer " + std::to_string(layer); | |
| return false; | |
| } | |
| const uint64_t layer_bytes = bytes * (uint64_t) n_expert; | |
| if (at > want || layer_bytes > want - at) { | |
| err = "FileExpertSource: the native expert layout exceeds its declared size at layer " + | |
| std::to_string(layer); | |
| return false; | |
| } | |
| layer_offsets[i] = layout.offset[i]; | |
| layer_blob_bytes[i] = bytes; | |
| at += layer_bytes; | |
| } | |
| if (at != want) { | |
| err = "FileExpertSource: the native expert layout has an inconsistent size"; | |
| return false; | |
| } | |
| } | |
| const std::string path = pack_dir + "/experts.bin"; | |
| if (layout.native && !gguf_.empty() && !std::filesystem::exists(path)) { | |
| // CS-T: no experts.bin - the model's GGUF shards, read in place | |
| blobs_ = (int64_t) blob_count; | |
| n_layers_ = n_layers; | |
| n_expert_ = n_expert; | |
| layer_offsets_ = std::move(layer_offsets); | |
| layer_blob_bytes_ = std::move(layer_blob_bytes); | |
| if (!open_gguf(err)) { close(); return false; } | |
| return true; | |
| } | |
| // UTF-8 -> UTF-16: the pack may live under a path with non-ASCII characters, and `CreateFileA` would | |
| // silently mangle it into a file-not-found. | |
| const int wide = MultiByteToWideChar(CP_UTF8, 0, path.c_str(), -1, nullptr, 0); | |
| std::vector<wchar_t> wpath((size_t) (wide > 0 ? wide : 1)); | |
| if (wide > 0) MultiByteToWideChar(CP_UTF8, 0, path.c_str(), -1, wpath.data(), wide); | |
| // **`FILE_FLAG_RANDOM_ACCESS` WAS HERE AND IT COST 14x.** | |
| // | |
| // The design depends on the OS page cache holding the whole 34 GB expert set, because this machine has | |
| // 64 GB of DDR5 and `L9` measured the CPU path at 44.14 GB/s from DRAM. `FILE_FLAG_RANDOM_ACCESS` tells | |
| // the cache manager the opposite: it disables read-ahead AND it lets the manager drop the pages again | |
| // quickly, on the assumption that a large randomly-accessed file will not be re-read. Measured, on | |
| // `strata generate --max-new 24`: **1.93 GB/s** - disk speed, 344 ms/token, and it never warmed up over 25 | |
| // tokens, because the pages were being evicted as fast as they were faulted in. | |
| // | |
| // The correct flag is NO flag. The access pattern IS random (10 of 512 experts per layer, a different 10 | |
| // each layer), but every byte read is read again on the next token, so retention is the whole game. | |
| HANDLE f = CreateFileW(wpath.data(), GENERIC_READ, FILE_SHARE_READ, nullptr, OPEN_EXISTING, | |
| FILE_ATTRIBUTE_NORMAL, nullptr); | |
| if (f == INVALID_HANDLE_VALUE) { | |
| err = "FileExpertSource: cannot open " + path; | |
| return false; | |
| } | |
| LARGE_INTEGER sz{}; | |
| if (!GetFileSizeEx(f, &sz)) { | |
| CloseHandle(f); | |
| err = "FileExpertSource: cannot size " + path; | |
| return false; | |
| } | |
| if ((uint64_t) sz.QuadPart != want) { | |
| char buf[400]; | |
| std::snprintf(buf, sizeof buf, | |
| "FileExpertSource: %s is %llu B but the loaded expert layout requires %llu B - this is not " | |
| "the pack this geometry came from", | |
| path.c_str(), (unsigned long long) sz.QuadPart, (unsigned long long) want); | |
| CloseHandle(f); | |
| err = buf; | |
| return false; | |
| } | |
| HANDLE m = CreateFileMappingW(f, nullptr, PAGE_READONLY, 0, 0, nullptr); | |
| if (m == nullptr) { | |
| CloseHandle(f); | |
| err = "FileExpertSource: CreateFileMapping failed on " + path; | |
| return false; | |
| } | |
| void* view = MapViewOfFile(m, FILE_MAP_READ, 0, 0, 0); | |
| if (view == nullptr) { | |
| CloseHandle(m); | |
| CloseHandle(f); | |
| err = "FileExpertSource: MapViewOfFile failed on " + path; | |
| return false; | |
| } | |
| file_ = f; | |
| mapping_ = m; | |
| base_ = (const uint8_t*) view; | |
| paths_.assign(1, path); | |
| const int fd = ::open(path.c_str(), O_RDONLY); | |
| if (fd < 0) { err = "FileExpertSource: cannot open " + path; return false; } | |
| struct stat st{}; | |
| if (fstat(fd, &st) != 0) { ::close(fd); err = "FileExpertSource: cannot stat " + path; return false; } | |
| if (st.st_size < 0 || (uint64_t) st.st_size != want) { | |
| char buf[400]; | |
| std::snprintf(buf, sizeof buf, | |
| "FileExpertSource: %s is %llu B but the loaded expert layout requires %llu B - this is not " | |
| "the pack this geometry came from", | |
| path.c_str(), (unsigned long long) (st.st_size < 0 ? 0 : st.st_size), | |
| (unsigned long long) want); | |
| ::close(fd); | |
| err = buf; | |
| return false; | |
| } | |
| void* view = mmap(nullptr, (size_t) want, PROT_READ, MAP_SHARED, fd, 0); | |
| if (view == MAP_FAILED) { ::close(fd); err = "FileExpertSource: mmap failed on " + path; return false; } | |
| fd_ = fd; | |
| base_ = (const uint8_t*) view; | |
| blobs_ = (int64_t) blob_count; | |
| n_layers_ = n_layers; | |
| n_expert_ = n_expert; | |
| mapped_bytes_ = want; | |
| layer_offsets_ = std::move(layer_offsets); | |
| layer_blob_bytes_ = std::move(layer_blob_bytes); | |
| return true; | |
| } | |
| void FileExpertSource::close() { | |
| if (complement_arena_ != nullptr) { | |
| if (complement_pinned_ && !complement_partial_) (void) cudaFreeHost(complement_arena_); | |
| else { | |
| if (complement_partial_) (void) cudaHostUnregister(complement_arena_); | |
| if (complement_locked_ > 0) | |
| strata::platform::unlock_resident((uint8_t*) complement_arena_ + complement_lock_off_, complement_locked_); | |
| std::free(complement_arena_); | |
| } | |
| } | |
| if (xstage_ != nullptr) { | |
| if (xstage_pinned_) (void) cudaFreeHost(xstage_); | |
| else std::free(xstage_); | |
| } | |
| xstage_ = nullptr; | |
| xstage_pinned_ = false; | |
| xstage_cap_ = 0; | |
| xstage_blob_ = 0; | |
| override_.clear(); | |
| staged_.clear(); | |
| exchanges_ = 0; | |
| file_reads_.store(0); | |
| complement_arena_ = nullptr; | |
| complement_host_ = nullptr; | |
| complement_device_ = nullptr; | |
| complement_bytes_ = 0; | |
| complement_offsets_.clear(); | |
| complement_pinned_ = false; | |
| complement_partial_ = false; | |
| complement_pin_limit_ = 0; | |
| complement_lock_off_ = 0; | |
| complement_ready_ = false; | |
| complement_locked_ = 0; | |
| complement_lent_slots_ = 0; | |
| if (!maps_.empty()) { | |
| for (Map& m : maps_) { | |
| if (m.base != nullptr) UnmapViewOfFile((LPCVOID) m.base); | |
| if (m.mapping != nullptr) CloseHandle((HANDLE) m.mapping); | |
| if (m.file != nullptr) CloseHandle((HANDLE) m.file); | |
| if (m.base != nullptr) munmap((void*) m.base, (size_t) m.bytes); | |
| if (m.fd >= 0) ::close(m.fd); | |
| } | |
| maps_.clear(); | |
| base_ = nullptr; // one of the views above | |
| } | |
| role_ptr_.clear(); | |
| role_bytes_.clear(); | |
| role_file_.clear(); | |
| paths_.clear(); | |
| for (void* h : direct_) CloseHandle((HANDLE) h); | |
| direct_.clear(); | |
| { | |
| std::lock_guard<std::mutex> lk(stage_mu_); | |
| stage_buf_.clear(); | |
| stage_key_.clear(); | |
| stage_epoch_.clear(); | |
| stage_used_.clear(); | |
| stage_busy_.clear(); | |
| stage_of_.clear(); | |
| stage_blob_ = 0; | |
| stage_seq_ = 0; | |
| epoch_ = 0; | |
| last_layer_ = -1; | |
| stage_grew_ = false; | |
| } | |
| ram_reads_.store(0); | |
| warm_stamp_.reset(); | |
| warm_hits_.store(0); | |
| warm_count_.store(0); | |
| file_read_bytes_.store(0); | |
| file_blob_bytes_.store(0); | |
| file_us_.store(0); | |
| if (base_ != nullptr) UnmapViewOfFile((LPCVOID) base_); | |
| if (mapping_ != nullptr) CloseHandle((HANDLE) mapping_); | |
| if (file_ != nullptr) CloseHandle((HANDLE) file_); | |
| mapping_ = nullptr; | |
| file_ = nullptr; | |
| if (base_ != nullptr) munmap((void*) base_, (size_t) mapped_bytes_); | |
| if (fd_ >= 0) ::close(fd_); | |
| fd_ = -1; | |
| base_ = nullptr; | |
| blobs_ = 0; | |
| n_layers_ = 0; | |
| n_expert_ = 0; | |
| mapped_bytes_ = 0; | |
| layer_offsets_.clear(); | |
| layer_blob_bytes_.clear(); | |
| reads_ = 0; | |
| } | |
| bool ExpertSource::copy_blob(int64_t layer, int64_t expert, uint8_t* dst) { | |
| const uint8_t* b = blob(layer, expert); | |
| if (b == nullptr || dst == nullptr) return false; | |
| std::memcpy(dst, b, (size_t) strata::kernels::cpu::expert_layout().blob_bytes(layer)); | |
| return true; | |
| } | |
| // ================================ CS-T: THE GGUF SHARDS IN PLACE ================================ | |
| // | |
| // A native pack without experts.bin: every file native_experts.txt names is mapped (MapViewOfFile / mmap, no | |
| // flag - the same retention argument as experts.bin above), and an expert's blob [gate rows | up rows | down rows] | |
| // is three slices of three tensors, possibly in two shards (UD-Q4_K_XL's layer 11). Nothing is read at open; a | |
| // blob is assembled when it is asked for (`blob`, into a small pool of buffers) or copied where it is needed | |
| // (`copy_blob`: the RAM copy, the prompt path's pinned stager buffers). | |
| bool FileExpertSource::open_gguf(std::string& err) { | |
| const auto& lay = strata::kernels::cpu::expert_layout(); | |
| if (!check_experts_gguf(gguf_, lay, err)) { err = "FileExpertSource: " + err; return false; } | |
| const size_t cut = gguf_.find_last_of("/\\"); | |
| const std::string dir = cut == std::string::npos ? std::string() : gguf_.substr(0, cut + 1); | |
| std::map<std::string, size_t> index; | |
| role_ptr_.assign((size_t) (3 * n_layers_), nullptr); | |
| role_bytes_.assign((size_t) (3 * n_layers_), 0); | |
| role_file_.assign((size_t) (3 * n_layers_), 0); | |
| for (int64_t l = 0; l < n_layers_; ++l) { | |
| const auto& fm = lay.fmt[(size_t) l]; | |
| const uint64_t per[3] = {fm.up_off, fm.up_off, lay.bytes[(size_t) l] - fm.down_off}; | |
| for (int r = 0; r < 3; ++r) { | |
| const size_t i = (size_t) (3 * l + r); | |
| const std::string path = lay.gguf_file.size() > i && !lay.gguf_file[i].empty() ? dir + lay.gguf_file[i] | |
| : gguf_; | |
| auto it = index.find(path); | |
| if (it == index.end()) { | |
| Map m; | |
| const int wide = MultiByteToWideChar(CP_UTF8, 0, path.c_str(), -1, nullptr, 0); | |
| std::vector<wchar_t> wpath((size_t) (wide > 0 ? wide : 1)); | |
| if (wide > 0) MultiByteToWideChar(CP_UTF8, 0, path.c_str(), -1, wpath.data(), wide); | |
| HANDLE f = CreateFileW(wpath.data(), GENERIC_READ, FILE_SHARE_READ, nullptr, OPEN_EXISTING, | |
| FILE_ATTRIBUTE_NORMAL, nullptr); | |
| LARGE_INTEGER sz{}; | |
| if (f == INVALID_HANDLE_VALUE || !GetFileSizeEx(f, &sz)) { | |
| if (f != INVALID_HANDLE_VALUE) CloseHandle(f); | |
| err = "FileExpertSource: cannot open " + path; | |
| return false; | |
| } | |
| HANDLE mh = CreateFileMappingW(f, nullptr, PAGE_READONLY, 0, 0, nullptr); | |
| void* view = mh != nullptr ? MapViewOfFile(mh, FILE_MAP_READ, 0, 0, 0) : nullptr; | |
| if (view == nullptr) { | |
| if (mh != nullptr) CloseHandle(mh); | |
| CloseHandle(f); | |
| err = "FileExpertSource: cannot map " + path; | |
| return false; | |
| } | |
| m.file = f; | |
| m.mapping = mh; | |
| m.bytes = (uint64_t) sz.QuadPart; | |
| m.base = (const uint8_t*) view; | |
| const int fd = ::open(path.c_str(), O_RDONLY); | |
| struct stat st{}; | |
| if (fd < 0 || fstat(fd, &st) != 0 || st.st_size <= 0) { | |
| if (fd >= 0) ::close(fd); | |
| err = "FileExpertSource: cannot open " + path; | |
| return false; | |
| } | |
| void* view = mmap(nullptr, (size_t) st.st_size, PROT_READ, MAP_SHARED, fd, 0); | |
| if (view == MAP_FAILED) { ::close(fd); err = "FileExpertSource: cannot map " + path; return false; } | |
| m.fd = fd; | |
| m.bytes = (uint64_t) st.st_size; | |
| m.base = (const uint8_t*) view; | |
| maps_.push_back(m); | |
| paths_.push_back(path); | |
| it = index.emplace(path, maps_.size() - 1).first; | |
| } | |
| const Map& m = maps_[it->second]; | |
| role_file_[i] = (int) it->second; | |
| const uint64_t at = lay.gguf_off[i], bytes = per[r] * (uint64_t) n_expert_; | |
| if (at > m.bytes || bytes > m.bytes - at) { // check_experts_gguf proved it; the mapping must agree | |
| err = "FileExpertSource: an expert span runs past the end of " + path; | |
| return false; | |
| } | |
| role_ptr_[i] = m.base + (size_t) at; | |
| role_bytes_[i] = per[r]; | |
| } | |
| } | |
| base_ = maps_.front().base; // "opened"; mapped_blob answers nullptr in this mode | |
| warm_stamp_.reset(new std::atomic<uint32_t>[(size_t) (n_layers_ * n_expert_)]()); | |
| for (uint64_t b : layer_blob_bytes_) stage_blob_ = std::max(stage_blob_, b); | |
| return true; | |
| } | |
| bool FileExpertSource::copy_from_files(int64_t layer, int64_t expert, uint8_t* dst) const { | |
| if (dst == nullptr || layer < 0 || expert < 0 || layer >= n_layers_ || expert >= n_expert_) return false; | |
| if (!direct_.empty()) { // #286: from the drive; a failed read falls back to the mapping below | |
| const Fill f{0, layer, expert, dst}; | |
| if (read_direct(&f, 1)) return true; | |
| } | |
| if (!role_ptr_.empty()) { | |
| uint64_t at = 0; | |
| for (int r = 0; r < 3; ++r) { | |
| const size_t i = (size_t) (3 * layer + r); | |
| const uint64_t per = role_bytes_[i]; | |
| std::memcpy(dst + at, role_ptr_[i] + (size_t) ((uint64_t) expert * per), (size_t) per); | |
| at += per; | |
| } | |
| return true; | |
| } | |
| const uint8_t* b = mapped_blob(layer, expert); | |
| if (b == nullptr) return false; | |
| std::memcpy(dst, b, (size_t) layer_blob_bytes_[(size_t) layer]); | |
| return true; | |
| } | |
| // A blob assembled from the three role slices. The buffer of a (layer, expert) is reused for another only once | |
| // its blob has not been asked for during `kStageAge` layers (begin_layer) or 256 assemblies, whichever comes | |
| // first, and never while it is being filled: the pool computes a layer's misses before it starts the next, and a | |
| // fill (the GPU cache at startup, an adaptive swap, a helper GPU) copies the blob right away. | |
| // | |
| // `claim_stage` finds or reserves the buffer of `key` (stage_mu_ held): true when the blob is already there (or | |
| // being filled by another thread - the caller then waits), false when the caller must fill buffer `v`. | |
| bool FileExpertSource::claim_stage(int64_t key, size_t& v, bool& fill) { | |
| constexpr uint64_t kStageSeq = 256; | |
| const uint64_t seq = ++stage_seq_; | |
| auto it = stage_of_.find(key); | |
| fill = false; | |
| if (it != stage_of_.end()) { | |
| v = it->second; | |
| stage_epoch_[v] = epoch_; | |
| stage_used_[v] = seq; | |
| return true; | |
| } | |
| v = stage_buf_.size(); | |
| uint64_t oldest = std::numeric_limits<uint64_t>::max(); | |
| for (size_t i = 0; i < stage_buf_.size(); ++i) | |
| if (!stage_busy_[i] && (stage_epoch_[i] + kStageAge <= epoch_ || stage_used_[i] + kStageSeq <= seq) && | |
| stage_used_[i] < oldest) { | |
| oldest = stage_used_[i]; | |
| v = i; | |
| } | |
| if (v == stage_buf_.size()) { | |
| stage_buf_.emplace_back(new (std::nothrow) uint8_t[(size_t) stage_blob_]); | |
| if (!stage_buf_.back()) { stage_buf_.pop_back(); return false; } | |
| stage_key_.push_back(-1); | |
| stage_epoch_.push_back(0); | |
| stage_used_.push_back(0); | |
| stage_busy_.push_back(0); | |
| if (stage_buf_.size() == 512 && !stage_grew_) { | |
| stage_grew_ = true; | |
| std::fprintf(stderr, "FileExpertSource: %zu blobs assembled from the GGUF are in use at once (%.2f GiB)\n", | |
| stage_buf_.size(), (double) stage_buf_.size() * (double) stage_blob_ / 1073741824.0); | |
| } | |
| } else { | |
| stage_of_.erase(stage_key_[v]); | |
| } | |
| stage_key_[v] = key; | |
| stage_epoch_[v] = epoch_; | |
| stage_used_[v] = seq; | |
| stage_busy_[v] = 1; | |
| stage_of_[key] = v; | |
| fill = true; | |
| return false; | |
| } | |
| // Fills buffer `v` (reserved by claim_stage) outside the lock, then publishes it. | |
| bool FileExpertSource::fill_stage(size_t v, int64_t layer, int64_t expert, uint8_t* dst) { | |
| const auto t0 = std::chrono::steady_clock::now(); | |
| const bool ok = copy_from_files(layer, expert, dst); | |
| publish_stage(v, layer, ok, std::chrono::duration<double, std::micro>(std::chrono::steady_clock::now() - t0).count()); | |
| return ok; | |
| } | |
| void FileExpertSource::publish_stage(size_t v, int64_t layer, bool ok, double us) { | |
| file_us_.fetch_add((uint64_t) us, std::memory_order_relaxed); | |
| if (ok) { | |
| file_read_bytes_.fetch_add(layer_blob_bytes_[(size_t) layer], std::memory_order_relaxed); | |
| file_blob_bytes_.fetch_add(layer_blob_bytes_[(size_t) layer], std::memory_order_relaxed); | |
| } | |
| { | |
| std::lock_guard<std::mutex> lk(stage_mu_); | |
| stage_busy_[v] = 0; | |
| if (!ok) { | |
| stage_of_.erase(stage_key_[v]); | |
| stage_key_[v] = -1; | |
| } | |
| } | |
| stage_cv_.notify_all(); | |
| } | |
| const uint8_t* FileExpertSource::staged_blob(int64_t layer, int64_t expert) { | |
| const int64_t key = layer * n_expert_ + expert; | |
| size_t v = 0; | |
| bool fill = false; | |
| uint8_t* dst = nullptr; | |
| { | |
| std::unique_lock<std::mutex> lk(stage_mu_); | |
| const bool have = claim_stage(key, v, fill); | |
| if (!have && !fill) return nullptr; | |
| dst = stage_buf_[v].get(); | |
| if (have) { | |
| // another thread (a prefetch, the adaptive tier) is filling it: wait for that | |
| stage_cv_.wait(lk, [&] { return !stage_busy_[v] || stage_key_[v] != key; }); | |
| if (stage_key_[v] != key) return nullptr; // its fill failed | |
| return dst; | |
| } | |
| } | |
| return fill_stage(v, layer, expert, dst) ? dst : nullptr; | |
| } | |
| void FileExpertSource::prefetch(int64_t layer, const int64_t* experts, int64_t n) { | |
| if (!staged() || n <= 0 || layer < 0 || layer >= n_layers_) return; | |
| std::vector<Fill> todo; | |
| { | |
| std::lock_guard<std::mutex> lk(stage_mu_); | |
| for (int64_t i = 0; i < n; ++i) { | |
| const int64_t e = experts[i]; | |
| if (e < 0 || e >= n_expert_) continue; | |
| const size_t index = (size_t) layer * (size_t) n_expert_ + (size_t) e; | |
| if (complement_ready_ && index < complement_offsets_.size() && complement_offsets_[index] != kNoComplement) | |
| continue; // in the RAM copy | |
| if (!override_.empty() && override_[index] != nullptr) continue; | |
| size_t v = 0; | |
| bool fill = false; | |
| if (!claim_stage(layer * n_expert_ + e, v, fill) && fill) { | |
| todo.push_back({v, layer, e, stage_buf_[v].get()}); | |
| if (warm_stamp_) { | |
| const uint32_t s = warm_stamp_[index].load(std::memory_order_relaxed); | |
| if (s != 0 && (uint64_t) s + 3 >= epoch_ + 1) warm_hits_.fetch_add(1, std::memory_order_relaxed); | |
| } | |
| } | |
| } | |
| } | |
| fill_many(todo); | |
| } | |
| void FileExpertSource::prefetch_pairs(const std::pair<int32_t, int32_t>* pairs, int64_t n) { | |
| if (direct_.empty() || pairs == nullptr || n <= 0) return; // unbuffered only: the mapped fill stays as it was | |
| std::vector<Fill> todo; | |
| { | |
| std::lock_guard<std::mutex> lk(stage_mu_); | |
| for (int64_t i = 0; i < std::min<int64_t>(n, 64); ++i) { | |
| const int64_t l = pairs[i].first, e = pairs[i].second; | |
| if (l < 0 || e < 0 || l >= n_layers_ || e >= n_expert_) continue; | |
| size_t v = 0; | |
| bool fill = false; | |
| if (!claim_stage(l * n_expert_ + e, v, fill) && fill) todo.push_back({v, l, e, stage_buf_[v].get()}); | |
| } | |
| } | |
| fill_many(todo); | |
| } | |
| void FileExpertSource::fill_many(const std::vector<Fill>& todo) { | |
| if (todo.empty()) return; | |
| if (!direct_.empty()) { | |
| // #286: overlapped batches - every role window of 16 blobs in flight at once; a big batch (the profile | |
| // fill) on up to 4 threads, a decode layer's few misses on this one | |
| std::atomic<size_t> next{0}; | |
| auto work = [&] { | |
| for (size_t at; (at = next.fetch_add(16)) < todo.size();) { | |
| const size_t k = std::min<size_t>(16, todo.size() - at); | |
| const auto t0 = std::chrono::steady_clock::now(); | |
| const bool ok = read_direct(todo.data() + at, k); | |
| const double us = std::chrono::duration<double, std::micro>(std::chrono::steady_clock::now() - t0).count(); | |
| for (size_t i = at; i < at + k; ++i) { | |
| const Fill& f = todo[i]; | |
| // a failed batch: each blob again on its own (copy_from_files falls back to the mapping) | |
| if (ok) publish_stage(f.v, f.layer, true, us / (double) k); | |
| else (void) fill_stage(f.v, f.layer, f.e, f.dst); | |
| } | |
| } | |
| }; | |
| const size_t nt = std::min<size_t>(4, (todo.size() + 15) / 16); | |
| std::vector<std::thread> th; | |
| for (size_t t = 1; t < nt; ++t) th.emplace_back(work); | |
| work(); | |
| for (auto& t : th) t.join(); | |
| return; | |
| } | |
| // One PrefetchVirtualMemory call for every slice about to be copied: the memory manager reads them in large | |
| // requests, all queued at once, where the copies' page faults would read a few clusters each. The copies below | |
| // then find the pages resident (or in flight). STRATA_FETCH_PVM=0 is the A/B arm. | |
| { | |
| using Pvm = BOOL(WINAPI*)(HANDLE, ULONG_PTR, PWIN32_MEMORY_RANGE_ENTRY, ULONG); | |
| static const Pvm pvm = [] { | |
| const char* v = std::getenv("STRATA_FETCH_PVM"); | |
| if (v != nullptr && std::atoi(v) == 0) return (Pvm) nullptr; | |
| return (Pvm) (void*) GetProcAddress(GetModuleHandleW(L"kernel32.dll"), "PrefetchVirtualMemory"); | |
| }(); | |
| if (pvm != nullptr) { | |
| std::vector<WIN32_MEMORY_RANGE_ENTRY> ranges; | |
| ranges.reserve(todo.size() * 3); | |
| for (const Fill& f : todo) | |
| for (int r = 0; r < 3; ++r) { | |
| const size_t i = (size_t) (3 * f.layer + r); | |
| ranges.push_back({(PVOID) (role_ptr_[i] + (size_t) ((uint64_t) f.e * role_bytes_[i])), | |
| (SIZE_T) role_bytes_[i]}); | |
| } | |
| (void) pvm(GetCurrentProcess(), (ULONG_PTR) ranges.size(), ranges.data(), 0); | |
| } | |
| } | |
| // the page faults of a mapped read are one outstanding request each: several threads keep the SSD's queue full | |
| std::atomic<size_t> next{0}; | |
| auto work = [&] { | |
| for (size_t i; (i = next.fetch_add(1)) < todo.size();) | |
| (void) fill_stage(todo[i].v, todo[i].layer, todo[i].e, todo[i].dst); | |
| }; | |
| const size_t nt = std::min<size_t>(todo.size(), (size_t) fetch_threads_); | |
| std::vector<std::thread> th; | |
| for (size_t t = 1; t < nt; ++t) th.emplace_back(work); | |
| work(); | |
| for (auto& t : th) t.join(); | |
| } | |
| bool FileExpertSource::set_unbuffered(uint64_t ram_bytes, std::string& why) { | |
| if (base_ == nullptr || paths_.empty() || !direct_.empty()) { | |
| why = !direct_.empty() ? "already unbuffered" : "no expert files open"; | |
| return !direct_.empty(); | |
| } | |
| if (!experts_unbuffered(paths_, ram_bytes, why, /*cache_counts=*/false)) return false; | |
| for (const std::string& path : paths_) { | |
| const int wide = MultiByteToWideChar(CP_UTF8, 0, path.c_str(), -1, nullptr, 0); | |
| std::vector<wchar_t> w((size_t) (wide > 0 ? wide : 1), L'\0'); | |
| if (wide > 0) MultiByteToWideChar(CP_UTF8, 0, path.c_str(), -1, w.data(), wide); | |
| HANDLE h = CreateFileW(w.data(), GENERIC_READ, FILE_SHARE_READ, nullptr, OPEN_EXISTING, | |
| FILE_FLAG_NO_BUFFERING | FILE_FLAG_OVERLAPPED, nullptr); | |
| if (h == INVALID_HANDLE_VALUE) { | |
| why += "; cannot open " + path + " unbuffered (error " + std::to_string((unsigned long long) GetLastError()) + | |
| "), read through the file cache"; | |
| for (void* d : direct_) CloseHandle((HANDLE) d); | |
| direct_.clear(); | |
| return false; | |
| } | |
| direct_.push_back(h); | |
| } | |
| if (role_ptr_.empty()) { // experts.bin: blob() now assembles into the stage buffers, sized for the largest blob | |
| std::lock_guard<std::mutex> lk(stage_mu_); | |
| for (uint64_t b : layer_blob_bytes_) stage_blob_ = std::max(stage_blob_, b); | |
| } | |
| return true; | |
| (void) ram_bytes; | |
| why = "through the file cache (not Windows)"; | |
| return false; | |
| } | |
| bool FileExpertSource::read_direct(const Fill* fills, size_t n) const { | |
| // NTFS runs the unbuffered reads of a file one at a time while the file is mapped or cached anywhere (the engine | |
| // maps the GGUF shards for the token embedding and the PLE table): measured on a PCIe 5 drive, 512 KiB reads at | |
| // queue depth 48 make 10.4 GB/s on a file nobody maps and 3.4 GB/s on a mapped one, 8 MiB reads 8.6 GB/s. So | |
| // windows less than kGap apart (the same role of nearby experts) are merged into requests of up to kMerge bytes. | |
| constexpr uint64_t kSector = 4096, kGap = 1ull << 20, kMerge = 32ull << 20; | |
| struct Window { int file; uint64_t a0, size, skip, n, at; uint8_t* dst; size_t req; uint64_t in_req; }; | |
| struct Req { HANDLE h; uint64_t a0, size, pos; }; | |
| // this thread's aligned buffer and events, kept for its next batch | |
| struct Scratch { | |
| uint8_t* buf = nullptr; | |
| size_t cap = 0; | |
| std::vector<HANDLE> ev; | |
| std::vector<OVERLAPPED> ov; | |
| std::vector<Window> win; | |
| std::vector<size_t> order; | |
| std::vector<Req> req; | |
| std::vector<DWORD> got; | |
| ~Scratch() { | |
| if (buf != nullptr) VirtualFree(buf, 0, MEM_RELEASE); | |
| for (HANDLE e : ev) CloseHandle(e); | |
| } | |
| }; | |
| thread_local Scratch sc; | |
| sc.win.clear(); | |
| for (size_t k = 0; k < n; ++k) { | |
| const Fill& f = fills[k]; | |
| if (f.dst == nullptr || f.layer < 0 || f.e < 0 || f.layer >= n_layers_ || f.e >= n_expert_) return false; | |
| if (role_ptr_.empty()) { // experts.bin: the blob is one contiguous range | |
| const uint64_t per = layer_blob_bytes_[(size_t) f.layer]; | |
| const uint64_t off = layer_offsets_[(size_t) f.layer] + (uint64_t) f.e * per; | |
| const uint64_t a0 = off / kSector * kSector, a1 = (off + per + kSector - 1) / kSector * kSector; | |
| sc.win.push_back({0, a0, a1 - a0, off - a0, per, 0, f.dst, 0, 0}); | |
| continue; | |
| } | |
| uint64_t at = 0; | |
| for (int r = 0; r < 3; ++r) { | |
| const size_t i = (size_t) (3 * f.layer + r); | |
| const uint64_t per = role_bytes_[i]; | |
| const Map& m = maps_[(size_t) role_file_[i]]; | |
| const uint64_t off = (uint64_t) (role_ptr_[i] - m.base) + (uint64_t) f.e * per; | |
| const uint64_t a0 = off / kSector * kSector, a1 = (off + per + kSector - 1) / kSector * kSector; | |
| sc.win.push_back({role_file_[i], a0, a1 - a0, off - a0, per, at, f.dst, 0, 0}); | |
| at += per; | |
| } | |
| } | |
| sc.order.resize(sc.win.size()); | |
| for (size_t w = 0; w < sc.win.size(); ++w) sc.order[w] = w; | |
| std::sort(sc.order.begin(), sc.order.end(), [&](size_t a, size_t b) { | |
| return sc.win[a].file != sc.win[b].file ? sc.win[a].file < sc.win[b].file : sc.win[a].a0 < sc.win[b].a0; | |
| }); | |
| sc.req.clear(); | |
| uint64_t total = 0; | |
| int last_file = -1; | |
| for (size_t w : sc.order) { | |
| Window& x = sc.win[w]; | |
| if (!sc.req.empty() && x.file == last_file) { | |
| Req& q = sc.req.back(); | |
| const uint64_t end = q.a0 + q.size, xend = x.a0 + x.size; | |
| if (x.a0 <= end + kGap && std::max(end, xend) - q.a0 <= kMerge) { | |
| if (xend > end) { | |
| total += xend - end; | |
| q.size = xend - q.a0; | |
| } | |
| x.req = sc.req.size() - 1; | |
| x.in_req = x.a0 - q.a0; | |
| continue; | |
| } | |
| } | |
| sc.req.push_back({(HANDLE) direct_[(size_t) x.file], x.a0, x.size, total}); | |
| total += x.size; | |
| last_file = x.file; | |
| x.req = sc.req.size() - 1; | |
| x.in_req = 0; | |
| } | |
| if (total > sc.cap) { | |
| if (sc.buf != nullptr) VirtualFree(sc.buf, 0, MEM_RELEASE); | |
| sc.cap = (size_t) ((total + (1u << 20) - 1) >> 20 << 20); | |
| sc.buf = (uint8_t*) VirtualAlloc(nullptr, sc.cap, MEM_COMMIT | MEM_RESERVE, PAGE_READWRITE); | |
| if (sc.buf == nullptr) { sc.cap = 0; return false; } | |
| } | |
| while (sc.ev.size() < sc.req.size()) { | |
| HANDLE e = CreateEventW(nullptr, TRUE, FALSE, nullptr); | |
| if (e == nullptr) return false; | |
| sc.ev.push_back(e); | |
| } | |
| sc.ov.assign(sc.req.size(), OVERLAPPED{}); | |
| // every request in flight before the first wait: the drive sees the whole batch as one queue | |
| size_t issued = 0; | |
| bool ok = true; | |
| for (size_t q = 0; q < sc.req.size(); ++q) { | |
| const Req& r = sc.req[q]; | |
| OVERLAPPED& o = sc.ov[q]; | |
| o.Offset = (DWORD) r.a0; | |
| o.OffsetHigh = (DWORD) (r.a0 >> 32); | |
| o.hEvent = sc.ev[q]; | |
| if (!ReadFile(r.h, sc.buf + r.pos, (DWORD) r.size, nullptr, &o) && GetLastError() != ERROR_IO_PENDING) { | |
| ok = false; | |
| break; | |
| } | |
| ++issued; | |
| } | |
| std::vector<DWORD>& got = sc.got; | |
| got.assign(sc.req.size(), 0); | |
| for (size_t q = 0; q < issued; ++q) | |
| if (!GetOverlappedResult(sc.req[q].h, &sc.ov[q], &got[q], TRUE)) ok = false; | |
| if (!ok || issued < sc.req.size()) return false; | |
| for (const Window& x : sc.win) { | |
| // a request may run past the end of the file: only the role's own bytes have to arrive | |
| if ((uint64_t) got[x.req] < x.in_req + x.skip + x.n) return false; | |
| std::memcpy(x.dst + x.at, sc.buf + sc.req[x.req].pos + x.in_req + x.skip, (size_t) x.n); | |
| } | |
| return true; | |
| (void) fills; (void) n; | |
| return false; | |
| } | |
| void FileExpertSource::warm(int64_t layer, const int64_t* experts, int64_t n) { | |
| if (role_ptr_.empty() || n <= 0 || layer < 0 || layer >= n_layers_) return; | |
| uint32_t stamp; | |
| { | |
| std::lock_guard<std::mutex> lk(stage_mu_); | |
| stamp = (uint32_t) epoch_ + 1; | |
| } | |
| using Pvm = BOOL(WINAPI*)(HANDLE, ULONG_PTR, PWIN32_MEMORY_RANGE_ENTRY, ULONG); | |
| static const Pvm pvm = (Pvm) (void*) GetProcAddress(GetModuleHandleW(L"kernel32.dll"), "PrefetchVirtualMemory"); | |
| std::vector<WIN32_MEMORY_RANGE_ENTRY> ranges; | |
| for (int64_t j = 0; j < n; ++j) { | |
| const int64_t e = experts[j]; | |
| if (e < 0 || e >= n_expert_) continue; | |
| const size_t index = (size_t) layer * (size_t) n_expert_ + (size_t) e; | |
| if (complement_ready_ && index < complement_offsets_.size() && complement_offsets_[index] != kNoComplement) | |
| continue; // in the RAM copy | |
| if (warm_stamp_) warm_stamp_[index].store(stamp, std::memory_order_relaxed); | |
| warm_count_.fetch_add(1, std::memory_order_relaxed); | |
| for (int r = 0; r < 3; ++r) { | |
| const size_t i = (size_t) (3 * layer + r); | |
| const uint8_t* p = role_ptr_[i] + (size_t) ((uint64_t) e * role_bytes_[i]); | |
| ranges.push_back({(PVOID) p, (SIZE_T) role_bytes_[i]}); | |
| const uintptr_t pg = 4096, a = (uintptr_t) p & ~(pg - 1); | |
| (void) madvise((void*) a, (size_t) ((uintptr_t) p + role_bytes_[i] - a), MADV_WILLNEED); | |
| } | |
| } | |
| if (pvm != nullptr && !ranges.empty()) (void) pvm(GetCurrentProcess(), (ULONG_PTR) ranges.size(), ranges.data(), 0); | |
| } | |
| RouterLookahead::~RouterLookahead() { | |
| { | |
| std::lock_guard<std::mutex> lk(mu_); | |
| quit_ = true; | |
| } | |
| cv_.notify_all(); | |
| if (thread_.joinable()) thread_.join(); | |
| } | |
| bool RouterLookahead::start(std::vector<std::vector<uint16_t>> routers, int64_t n_embd, int64_t n_expert, int k, | |
| ExpertSource* src, std::string& err) { | |
| if (src == nullptr || !src->warms()) { err = "RouterLookahead: the expert source does not warm"; return false; } | |
| if (n_embd % 8 != 0) { err = "RouterLookahead: n_embd is not a multiple of 8"; return false; } | |
| for (const auto& r : routers) | |
| if (r.size() != (size_t) (n_embd * n_expert)) { err = "RouterLookahead: a router of another shape"; return false; } | |
| routers_ = std::move(routers); | |
| n_embd_ = n_embd; | |
| n_expert_ = n_expert; | |
| k_ = k < 1 ? 1 : k > (int) n_expert ? (int) n_expert : k; | |
| src_ = src; | |
| x_.assign((size_t) (8 * n_embd), 0.f); | |
| thread_ = std::thread([this] { run(); }); | |
| return true; | |
| } | |
| void RouterLookahead::submit(int64_t layer, const float* x, int64_t n_tok, const int32_t* host_res) { | |
| if (layer + 1 >= (int64_t) routers_.size() || n_tok <= 0 || x == nullptr) return; | |
| { | |
| std::lock_guard<std::mutex> lk(mu_); | |
| if (busy_ || pending_) { skipped_.fetch_add(1, std::memory_order_relaxed); return; } | |
| n_tok_ = std::min<int64_t>(n_tok, 8); | |
| std::memcpy(x_.data(), x, (size_t) (n_tok_ * n_embd_) * sizeof(float)); | |
| layer_ = layer + 1; | |
| host_res_ = host_res; | |
| pending_ = true; | |
| } | |
| cv_.notify_one(); | |
| } | |
| void RouterLookahead::run() { | |
| std::vector<float> logits((size_t) (8 * n_expert_)); | |
| std::vector<int32_t> order((size_t) n_expert_); | |
| std::vector<int64_t> want; | |
| for (;;) { | |
| int64_t layer, nt; | |
| const int32_t* host_res; | |
| { | |
| std::unique_lock<std::mutex> lk(mu_); | |
| cv_.wait(lk, [&] { return quit_ || pending_; }); | |
| if (quit_) return; | |
| pending_ = false; | |
| busy_ = true; | |
| layer = layer_; | |
| nt = n_tok_; | |
| host_res = host_res_; | |
| } | |
| const auto t0 = std::chrono::steady_clock::now(); | |
| want.clear(); | |
| strata::kernels::cpu::bf16_rows_dot_multi(routers_[(size_t) layer].data(), (int) n_expert_, (int) n_embd_, | |
| x_.data(), (int) nt, logits.data()); | |
| for (int64_t t = 0; t < nt; ++t) { | |
| const float* lt = logits.data() + (size_t) (t * n_expert_); | |
| for (int64_t e = 0; e < n_expert_; ++e) order[(size_t) e] = (int32_t) e; | |
| std::partial_sort(order.begin(), order.begin() + k_, order.end(), | |
| [&](int32_t a, int32_t b) { return lt[(size_t) a] > lt[(size_t) b]; }); | |
| for (int j = 0; j < k_; ++j) { | |
| const int64_t e = order[(size_t) j]; | |
| if (host_res != nullptr && host_res[(size_t) (layer * n_expert_ + e)] >= 0) continue; // on the GPU | |
| if (std::find(want.begin(), want.end(), e) == want.end()) want.push_back(e); | |
| } | |
| } | |
| src_->warm(layer, want.data(), (int64_t) want.size()); | |
| predicted_.fetch_add((int64_t) want.size(), std::memory_order_relaxed); | |
| busy_us_.fetch_add((uint64_t) std::chrono::duration<double, std::micro>(std::chrono::steady_clock::now() - t0).count(), | |
| std::memory_order_relaxed); | |
| { | |
| std::lock_guard<std::mutex> lk(mu_); | |
| busy_ = false; | |
| } | |
| } | |
| } | |
| void FileExpertSource::begin_layer(int64_t layer, const int32_t* ids, int64_t k) { | |
| (void) ids; | |
| (void) k; | |
| if (!staged()) return; | |
| std::lock_guard<std::mutex> lk(stage_mu_); | |
| if (layer != last_layer_) { | |
| ++epoch_; | |
| last_layer_ = layer; | |
| } | |
| } | |
| bool FileExpertSource::transient(int64_t layer, int64_t expert) const { | |
| if (!staged() || layer < 0 || expert < 0 || layer >= n_layers_ || expert >= n_expert_) return false; | |
| const size_t index = (size_t) layer * (size_t) n_expert_ + (size_t) expert; | |
| if (complement_ready_ && index < complement_offsets_.size() && complement_offsets_[index] != kNoComplement) | |
| return false; | |
| return override_.empty() || override_[index] == nullptr; | |
| } | |
| bool FileExpertSource::copy_blob(int64_t layer, int64_t expert, uint8_t* dst) { | |
| if (base_ == nullptr || dst == nullptr || layer < 0 || expert < 0 || layer >= n_layers_ || expert >= n_expert_) | |
| return false; | |
| const size_t index = (size_t) layer * (size_t) n_expert_ + (size_t) expert; | |
| const uint64_t bytes = layer_blob_bytes_[(size_t) layer]; | |
| if (complement_ready_) { | |
| const uint8_t* held = | |
| detail::cache_complement_blob_or_fallback(index, complement_offsets_, complement_host_, nullptr); | |
| if (held == nullptr && !override_.empty()) held = override_[index]; | |
| if (held != nullptr) { | |
| std::memcpy(dst, held, (size_t) bytes); | |
| return true; | |
| } | |
| } | |
| const auto t0 = std::chrono::steady_clock::now(); | |
| if (!copy_from_files(layer, expert, dst)) return false; | |
| file_us_.fetch_add((uint64_t) std::chrono::duration<double, std::micro>(std::chrono::steady_clock::now() - t0).count(), | |
| std::memory_order_relaxed); | |
| file_read_bytes_.fetch_add(bytes, std::memory_order_relaxed); | |
| return true; | |
| } | |
| const uint8_t* FileExpertSource::mapped_blob(int64_t layer, int64_t expert) const { | |
| if (!role_ptr_.empty()) return nullptr; // the GGUF in place: no contiguous blob in any file | |
| if (base_ == nullptr || layer < 0 || expert < 0 || layer >= n_layers_ || expert >= n_expert_) return nullptr; | |
| const size_t i = (size_t) layer; | |
| if (i >= layer_offsets_.size() || i >= layer_blob_bytes_.size()) return nullptr; | |
| const uint64_t blob_bytes = layer_blob_bytes_[i]; | |
| if (blob_bytes == 0 || (uint64_t) expert > std::numeric_limits<uint64_t>::max() / blob_bytes) return nullptr; | |
| const uint64_t expert_offset = (uint64_t) expert * blob_bytes; | |
| const uint64_t layer_offset = layer_offsets_[i]; | |
| if (layer_offset > mapped_bytes_ || expert_offset > mapped_bytes_ - layer_offset) return nullptr; | |
| const uint64_t offset = layer_offset + expert_offset; | |
| if (blob_bytes > mapped_bytes_ - offset) return nullptr; | |
| return base_ + (size_t) offset; | |
| } | |
| bool FileExpertSource::pin_cache_complement( | |
| const ExpertCache& cache, std::string& err, bool pin, | |
| const std::vector<std::pair<int32_t, int32_t>>& additional_gpu_pairs, int64_t lend_from_slot, | |
| uint64_t headroom_bytes, uint64_t budget_bytes, const std::vector<std::pair<int32_t, int32_t>>* rank) { | |
| err.clear(); | |
| if (base_ == nullptr) { err = "FileExpertSource: open the mapped experts before pinning a complement"; return false; } | |
| if (complement_ready_) { err = "FileExpertSource: the cache complement is already pinned"; return false; } | |
| if (!cache.valid()) { err = "FileExpertSource: the GPU expert cache is not open"; return false; } | |
| if (cache.fills() != cache.resident()) { | |
| err = "FileExpertSource: the GPU expert cache is not fully filled"; | |
| return false; | |
| } | |
| const cudaError_t sync = cudaDeviceSynchronize(); | |
| if (sync != cudaSuccess) { | |
| err = std::string("FileExpertSource: GPU expert cache is not ready: ") + cudaGetErrorString(sync); | |
| (void) cudaGetLastError(); | |
| return false; | |
| } | |
| // The GPU cache's experts, and the bytes each slot's expert takes here (for the lend region below). | |
| const int64_t n_slots = cache.slots(); | |
| std::vector<std::pair<int32_t, int32_t>> primary_gpu_pairs; | |
| std::vector<int32_t> pair_slot; | |
| std::vector<uint64_t> slot_bytes((size_t) std::max<int64_t>(n_slots, 0), 0); | |
| primary_gpu_pairs.reserve((size_t) cache.resident()); | |
| pair_slot.reserve((size_t) cache.resident()); | |
| for (int64_t layer = 0; layer < n_layers_; ++layer) { | |
| for (int64_t expert = 0; expert < n_expert_; ++expert) { | |
| const int32_t slot = cache.slot_of(layer, expert); | |
| if (slot == kNotResident) continue; | |
| primary_gpu_pairs.emplace_back((int32_t) layer, (int32_t) expert); | |
| pair_slot.push_back(slot); | |
| if (slot >= 0 && slot < n_slots) slot_bytes[(size_t) slot] = layer_blob_bytes_[(size_t) layer]; | |
| } | |
| } | |
| std::vector<uint64_t> offsets; | |
| uint64_t bytes = 0; | |
| if (!detail::make_cache_complement_plan(n_layers_, n_expert_, layer_blob_bytes_, primary_gpu_pairs, | |
| additional_gpu_pairs, offsets, bytes, err)) return false; | |
| // #467: the GPU cache's pre-fill touched its experts through the mapping (~19 GiB on a 24 GB card), and Windows | |
| // counts those file pages in this process's working set, not as available: a 32 GB PC read 0.44 GiB here | |
| // (20.7 GiB before the start). Trimmed, they move to the standby list (still cached, counted as available). | |
| // Resident mode only: nothing else calls this function. Locked/pinned pages stay; the rest fault back softly. | |
| { | |
| uint64_t before = 0, after = 0; | |
| const bool read_before = available_memory_bytes(before); | |
| (void) SetProcessWorkingSetSize(GetCurrentProcess(), (SIZE_T) -1, (SIZE_T) -1); | |
| if (read_before && available_memory_bytes(after)) | |
| std::fprintf(stderr, "FileExpertSource: available RAM %.2f GiB, %.2f GiB after the mapped experts left the " | |
| "process working set (#467)\n", | |
| (double) before / 1073741824.0, (double) after / 1073741824.0); | |
| } | |
| const bool what_fits = budget_bytes == kResidentWhatFits; // #467: the soft mode's second try | |
| uint64_t budget_physical = 0; // #403: the RAM reading a budget was sized from (0: no budget) | |
| if (budget_bytes > 0) { | |
| // CS-T: a RAM budget. The complement's experts in `rank` order (the expert profile, hottest first) while | |
| // they fit, the rest left on the mapped files; clamped to what the RAM has room for. | |
| uint64_t physical = 0; | |
| if (!available_memory_bytes(physical)) { | |
| err = "FileExpertSource: cannot determine available RAM for --resident-budget-gib"; | |
| return false; | |
| } | |
| budget_physical = physical; | |
| const uint64_t room = physical > headroom_bytes ? physical - headroom_bytes : 0; | |
| if (budget_bytes > room) { | |
| // #403: 256 MiB under the room, so the engine's own allocations after this reading still leave the | |
| // headroom (a budget clamped to exactly the room failed the safety check below on a reading a few MB | |
| // lower). An unclamped budget is unchanged. | |
| const uint64_t margin = 256ull << 20; | |
| const uint64_t clamped = room > margin ? room - margin : 0; | |
| if (what_fits) | |
| std::fprintf(stderr, "FileExpertSource: RAM room for the complement: %.2f GiB (%.2f GiB available " | |
| "minus %.0f GiB headroom and a 0.25 GiB margin)\n", | |
| (double) clamped / 1073741824.0, (double) physical / 1073741824.0, | |
| (double) headroom_bytes / 1073741824.0); | |
| else | |
| std::fprintf(stderr, "FileExpertSource: --resident-budget-gib %.2f is more than the RAM has room for " | |
| "(%.2f GiB available minus %.0f GiB headroom and a 0.25 GiB margin): %.2f GiB\n", | |
| (double) budget_bytes / 1073741824.0, (double) physical / 1073741824.0, | |
| (double) headroom_bytes / 1073741824.0, (double) clamped / 1073741824.0); | |
| budget_bytes = clamped; | |
| } | |
| std::vector<uint64_t> ranked(offsets.size(), kNoComplement); | |
| uint64_t at = 0; | |
| int64_t held = 0; | |
| if (rank != nullptr) | |
| for (const auto& pr : *rank) { | |
| if (pr.first < 0 || pr.second < 0 || pr.first >= n_layers_ || pr.second >= n_expert_) continue; | |
| const size_t i = (size_t) pr.first * (size_t) n_expert_ + (size_t) pr.second; | |
| if (offsets[i] == kNoComplement || ranked[i] != kNoComplement) continue; // on a GPU, or twice | |
| const uint64_t b = layer_blob_bytes_[(size_t) pr.first]; | |
| if (b > budget_bytes - at) continue; | |
| ranked[i] = at; | |
| at += b; | |
| ++held; | |
| } | |
| if (what_fits && held == 0) { // #467: nothing to keep - the caller's plain mmap fallback, not an empty copy | |
| err = "FileExpertSource: the RAM has no room for any expert of the complement"; | |
| return false; | |
| } | |
| std::fprintf(stderr, "FileExpertSource: RAM budget %.2f GiB: %lld of the %.2f GiB of experts the GPU cache does " | |
| "not hold, by profile rank; the rest are read from the files\n", | |
| (double) budget_bytes / 1073741824.0, (long long) held, (double) bytes / 1073741824.0); | |
| offsets.swap(ranked); | |
| bytes = at; | |
| lend_from_slot = -1; | |
| } | |
| const bool lend = lend_from_slot >= 0 && lend_from_slot < n_slots && additional_gpu_pairs.empty(); | |
| uint64_t budget = std::numeric_limits<uint64_t>::max(); | |
| if (bytes > 0 || lend) { | |
| // #403: with a budget, the reading it was sized from - a second reading a few MB lower (the engine's own | |
| // allocations, the file cache) failed a budget the first one had clamped. (A budget turns `lend` off.) | |
| uint64_t physical = budget_physical; | |
| if (physical == 0 && !available_memory_bytes(physical)) { | |
| err = "FileExpertSource: cannot determine available RAM for the resident-memory safety check"; | |
| return false; | |
| } | |
| budget = physical > headroom_bytes ? physical - headroom_bytes : 0; | |
| if (bytes > budget) { | |
| char message[320]; | |
| std::snprintf(message, sizeof message, | |
| "FileExpertSource: resident complement %.2f GiB exceeds available RAM (%.2f GiB) minus the " | |
| "%.0f GiB safety headroom", | |
| (double) bytes / 1073741824.0, (double) physical / 1073741824.0, | |
| (double) headroom_bytes / 1073741824.0); | |
| err = message; | |
| return false; | |
| } | |
| } | |
| // The prompt path's lend region: its slots' experts are streamed from here during a prompt and copied back into | |
| // their slots after it, so the ones that fit are kept here too (from the last slot down: a short prompt lends | |
| // only the last few). The rest keep the mapped-file fallback. | |
| int64_t keep_from = n_slots; | |
| if (lend) { | |
| keep_from = detail::choose_resident_keep_from(slot_bytes, bytes, budget, lend_from_slot); | |
| if (keep_from < 0) keep_from = n_slots; | |
| if (keep_from < n_slots) { | |
| std::vector<std::pair<int32_t, int32_t>> core; | |
| core.reserve(primary_gpu_pairs.size()); | |
| for (size_t i = 0; i < primary_gpu_pairs.size(); ++i) | |
| if (pair_slot[i] < keep_from) core.push_back(primary_gpu_pairs[i]); | |
| if (!detail::make_cache_complement_plan(n_layers_, n_expert_, layer_blob_bytes_, core, | |
| additional_gpu_pairs, offsets, bytes, err)) return false; | |
| } | |
| } | |
| void* arena = nullptr; | |
| const uint8_t* host = nullptr; | |
| const uint8_t* device = nullptr; | |
| bool pinned_ok = false; | |
| uint64_t locked = 0; | |
| uint64_t partial_pin = 0; ///< CS-T: a registered prefix of a locked arena | |
| uint64_t lock_off = 0; ///< where the working-set lock starts (after the registered prefix) | |
| std::string note; | |
| auto release = [&]() { | |
| if (arena == nullptr) return; | |
| if (pinned_ok) (void) cudaFreeHost(arena); | |
| else { | |
| if (partial_pin > 0) (void) cudaHostUnregister(arena); | |
| if (locked > 0) strata::platform::unlock_resident((uint8_t*) arena + lock_off, locked); | |
| std::free(arena); | |
| } | |
| arena = nullptr; | |
| }; | |
| if (bytes > 0) { | |
| std::fprintf(stderr, "FileExpertSource: allocating %.2f GiB %s cache complement\n", | |
| (double) bytes / 1073741824.0, pin ? "page-locked" : "pageable resident"); | |
| std::fflush(stderr); | |
| if (pin) { | |
| const cudaError_t allocated = cudaHostAlloc(&arena, (size_t) bytes, | |
| cudaHostAllocMapped | cudaHostAllocPortable); | |
| if (allocated == cudaSuccess) { | |
| void* alias = nullptr; | |
| const cudaError_t aliased = cudaHostGetDevicePointer(&alias, arena, 0); | |
| if (aliased == cudaSuccess && alias != nullptr) { | |
| device = (const uint8_t*) alias; | |
| pinned_ok = true; | |
| note = "page-locked and mapped"; | |
| } else { | |
| note = std::string("no device alias (") + cudaGetErrorString(aliased) + ")"; | |
| (void) cudaGetLastError(); | |
| (void) cudaFreeHost(arena); | |
| arena = nullptr; | |
| } | |
| } else { | |
| // Refused (the driver's page-locked limit): the same bytes in ordinary memory, locked in the working | |
| // set instead, as the arena does - resident either way, only copied by the CPU instead of by DMA. | |
| note = std::string("page-locking refused (") + cudaGetErrorString(allocated) + ")"; | |
| (void) cudaGetLastError(); | |
| arena = nullptr; | |
| } | |
| } | |
| if (arena == nullptr) { | |
| arena = std::malloc((size_t) bytes); | |
| if (arena == nullptr) { | |
| err = "FileExpertSource: pageable resident complement allocation failed"; | |
| return false; | |
| } | |
| if (pin) { | |
| // CS-T, a RAM budget: its bytes are in profile order, hottest first, so the driver is asked to | |
| // register the largest prefix it takes (from the cap down in 2 GiB steps). Those experts can be | |
| // read by the GPU over PCIe (--pcie-frac) and copied by DMA; only the rest is locked in the working | |
| // set (the registered prefix is page-locked by the driver already - locking it twice made the next | |
| // device allocation fail). | |
| // opt-in (STRATA_PARTIAL_PIN=1): on the RTX 5070 PC the GPU's PCIe share of the misses measured no | |
| // faster than the CPU computing them (7.30 / 7.44 tok/s with 24 / 16 GiB registered against 7.05-7.74 | |
| // unpinned at a 40 GiB budget), and registering adds startup time and driver memory pressure | |
| static const bool partial_on = [] { | |
| const char* v = std::getenv("STRATA_PARTIAL_PIN"); | |
| return v != nullptr && std::atoi(v) != 0; | |
| }(); | |
| // at most STRATA_PARTIAL_PIN_GIB (default 24): registering 30 GiB of a 40 GiB arena left the driver | |
| // unable to page-lock the prompt path's small buffers afterwards (RTX 5070, WDDM) | |
| static const uint64_t pin_cap = [] { | |
| const char* v = std::getenv("STRATA_PARTIAL_PIN_GIB"); | |
| return (uint64_t) ((v != nullptr && std::atof(v) > 0 ? std::atof(v) : 24.0) * 1073741824.0); | |
| }(); | |
| if (budget_bytes > 0 && partial_on) { | |
| const uint64_t step = 2ull << 30; | |
| for (uint64_t want = std::min(bytes, pin_cap); want >= step; want = want > step ? want - step : 0) { | |
| // cut at an expert boundary: a blob that started inside the registered range and ran past it | |
| // would be taken as page-locked by a cudaMemcpyAsync and refused ("adaptive refill failed") | |
| uint64_t w = want; | |
| for (size_t i = 0; i < offsets.size(); ++i) { | |
| if (offsets[i] == kNoComplement) continue; | |
| const uint64_t b = layer_blob_bytes_[i / (size_t) n_expert_]; | |
| if (offsets[i] < want && offsets[i] + b > want) { w = offsets[i]; break; } | |
| } | |
| if (w == 0) break; | |
| if (cudaHostRegister(arena, (size_t) w, cudaHostRegisterMapped | cudaHostRegisterPortable) == | |
| cudaSuccess) { | |
| void* alias = nullptr; | |
| if (cudaHostGetDevicePointer(&alias, arena, 0) == cudaSuccess && alias != nullptr) { | |
| device = (const uint8_t*) alias; | |
| partial_pin = w; | |
| } else { | |
| (void) cudaGetLastError(); | |
| (void) cudaHostUnregister(arena); | |
| } | |
| break; | |
| } | |
| (void) cudaGetLastError(); | |
| if (want <= step) break; | |
| } | |
| char msg[160]; | |
| std::snprintf(msg, sizeof msg, "%.2f GiB of it registered for the GPU (the hottest)", | |
| (double) partial_pin / 1073741824.0); | |
| note += std::string("; ") + msg; | |
| } | |
| lock_off = partial_pin; | |
| const strata::platform::LockResult lr = | |
| strata::platform::lock_resident((uint8_t*) arena + lock_off, bytes - lock_off); | |
| locked = lr.locked_bytes; | |
| note += (note.empty() ? "" : "; ") + lr.note; | |
| } | |
| } | |
| host = (const uint8_t*) arena; | |
| } | |
| // Copied layer by layer on a few threads: the page faults of the mapped file are the cost, and they overlap. | |
| const long page_size = sysconf(_SC_PAGESIZE); | |
| if (page_size <= 0) { | |
| err = "FileExpertSource: cannot determine page size for mapped-page release"; | |
| release(); | |
| return false; | |
| } | |
| std::atomic<int64_t> next_layer{0}, layers_done{0}; | |
| std::atomic<uint64_t> copied{0}; | |
| std::atomic<bool> failed{false}; | |
| std::mutex fail_mu; | |
| std::string fail_msg; | |
| auto fail = [&](const std::string& m) { | |
| std::lock_guard<std::mutex> lock(fail_mu); | |
| if (fail_msg.empty()) fail_msg = m; | |
| failed.store(true); | |
| }; | |
| auto worker = [&]() { | |
| for (;;) { | |
| const int64_t layer = next_layer.fetch_add(1); | |
| if (layer >= n_layers_ || failed.load()) return; | |
| const uint64_t blob_bytes = layer_blob_bytes_[(size_t) layer]; | |
| std::vector<Fill> batch; // #286, unbuffered: 32 blobs' reads in flight at once, merged where near | |
| auto flush = [&]() { | |
| if (batch.empty()) return true; | |
| bool ok = read_direct(batch.data(), batch.size()); | |
| for (size_t i = 0; !ok && i < batch.size(); ++i) | |
| if (!copy_from_files(batch[i].layer, batch[i].e, batch[i].dst)) return false; | |
| copied.fetch_add(blob_bytes * (uint64_t) batch.size()); | |
| batch.clear(); | |
| return true; | |
| }; | |
| for (int64_t expert = 0; expert < n_expert_; ++expert) { | |
| const size_t index = (size_t) layer * (size_t) n_expert_ + (size_t) expert; | |
| const uint64_t offset = offsets[index]; | |
| if (offset == kNoComplement) continue; | |
| if (offset > bytes || blob_bytes > bytes - offset) { | |
| fail("FileExpertSource: invalid blob bounds while building the cache complement"); | |
| return; | |
| } | |
| uint8_t* dst = (uint8_t*) host + (size_t) offset; | |
| if (!direct_.empty()) { | |
| batch.push_back({0, layer, expert, dst}); | |
| if (batch.size() == 32 && !flush()) { | |
| fail("FileExpertSource: an expert could not be read while building the cache complement"); | |
| return; | |
| } | |
| continue; | |
| } | |
| if (!copy_from_files(layer, expert, dst)) { | |
| fail("FileExpertSource: invalid blob bounds while building the cache complement"); | |
| return; | |
| } | |
| copied.fetch_add(blob_bytes); | |
| } | |
| if (!flush()) { | |
| fail("FileExpertSource: an expert could not be read while building the cache complement"); | |
| return; | |
| } | |
| if (role_ptr_.empty()) { // experts.bin; the GGUF in place leaves its pages to the OS | |
| const uint64_t layer_offset = layer_offsets_[(size_t) layer]; | |
| const uint64_t layer_bytes = blob_bytes * (uint64_t) n_expert_; | |
| const uint64_t layer_end = layer_offset + layer_bytes; | |
| const uint64_t page = (uint64_t) page_size; | |
| const uint64_t advice_start = layer_offset - layer_offset % page; | |
| const uint64_t end_remainder = layer_end % page; | |
| const uint64_t extra = end_remainder == 0 ? 0 : page - end_remainder; | |
| const uint64_t advice_end = extra > mapped_bytes_ - layer_end ? mapped_bytes_ : layer_end + extra; | |
| if (advice_end > advice_start && | |
| madvise((void*) (base_ + (size_t) advice_start), (size_t) (advice_end - advice_start), MADV_DONTNEED) != 0) { | |
| fail("FileExpertSource: madvise could not release mapped expert layer " + std::to_string(layer)); | |
| return; | |
| } | |
| if (posix_fadvise(fd_, (off_t) layer_offset, (off_t) layer_bytes, POSIX_FADV_DONTNEED) != 0) { | |
| fail("FileExpertSource: posix_fadvise could not release expert layer " + std::to_string(layer)); | |
| return; | |
| } | |
| } | |
| const int64_t done = layers_done.fetch_add(1) + 1; | |
| if (done % 8 == 0 || done == n_layers_) | |
| std::fprintf(stderr, "FileExpertSource: copied cache complement through layer %lld/%lld (%.2f GiB)\n", | |
| (long long) done, (long long) n_layers_, (double) copied.load() / 1073741824.0); | |
| } | |
| }; | |
| { | |
| const int threads = (int) std::max<int64_t>(1, std::min<int64_t>(6, n_layers_)); | |
| std::vector<std::thread> pool; | |
| for (int i = 1; i < threads; ++i) pool.emplace_back(worker); | |
| worker(); | |
| for (auto& t : pool) t.join(); | |
| } | |
| std::fflush(stderr); | |
| if (failed.load()) { | |
| err = fail_msg; | |
| release(); | |
| return false; | |
| } | |
| // The mapped pages this process touched (the GPU cache's fill and this copy) leave its working set for the | |
| // standby list: VirtualUnlock on pages that are not locked does exactly that (it then reports ERROR_NOT_LOCKED). | |
| if (maps_.empty()) (void) VirtualUnlock((LPVOID) base_, (SIZE_T) mapped_bytes_); | |
| for (const Map& m : maps_) (void) VirtualUnlock((LPVOID) m.base, (SIZE_T) m.bytes); | |
| complement_arena_ = arena; | |
| complement_host_ = host; | |
| complement_device_ = device; | |
| complement_bytes_ = bytes; | |
| complement_offsets_ = std::move(offsets); | |
| complement_pinned_ = (pinned_ok || partial_pin > 0) && bytes > 0; | |
| complement_pin_limit_ = pinned_ok ? bytes : partial_pin; | |
| complement_partial_ = !pinned_ok && partial_pin > 0; | |
| complement_locked_ = locked; | |
| complement_lock_off_ = lock_off; | |
| complement_lent_slots_ = lend ? n_slots - keep_from : 0; | |
| complement_ready_ = true; | |
| std::fprintf(stderr, "FileExpertSource: %s cache complement ready: resident %.2f GiB, pinned %.2f GiB%s%s\n", | |
| complement_pinned_ ? "mapped pinned" : pin ? "locked resident" : "pageable resident", | |
| (double) resident_bytes() / 1073741824.0, (double) pinned_bytes() / 1073741824.0, | |
| note.empty() ? "" : "; ", note.c_str()); | |
| if (lend) | |
| std::fprintf(stderr, "FileExpertSource: %lld of the prompt path's %lld lendable slots keep their experts in RAM " | |
| "too%s\n", (long long) complement_lent_slots_, (long long) (n_slots - lend_from_slot), | |
| complement_lent_slots_ < n_slots - lend_from_slot | |
| ? " (the others are read from the file when lent: not enough RAM for them)" : ""); | |
| if (!additional_gpu_pairs.empty()) { | |
| std::fprintf(stderr, "FileExpertSource: %zu verified additional-GPU experts remain on the mmap fallback\n", | |
| additional_gpu_pairs.size()); | |
| } | |
| std::fflush(stderr); | |
| return true; | |
| } | |
| bool FileExpertSource::has_resident(int64_t layer, int64_t expert) const { | |
| if (!complement_ready_ || layer < 0 || expert < 0 || layer >= n_layers_ || expert >= n_expert_) return false; | |
| const size_t index = (size_t) layer * (size_t) n_expert_ + (size_t) expert; | |
| return index < complement_offsets_.size() && complement_offsets_[index] != kNoComplement; | |
| } | |
| bool FileExpertSource::reserve_exchanges(int64_t n, std::string& err) { | |
| err.clear(); | |
| if (n <= xstage_cap_) return true; | |
| if (!staged_.empty()) { err = "FileExpertSource: exchange buffers are in use"; return false; } | |
| uint64_t blob = 0; | |
| for (const uint64_t b : layer_blob_bytes_) blob = std::max(blob, b); | |
| if (blob == 0 || n <= 0) { err = "FileExpertSource: no expert geometry for the exchange buffers"; return false; } | |
| if (xstage_ != nullptr) { | |
| if (xstage_pinned_) (void) cudaFreeHost(xstage_); | |
| else std::free(xstage_); | |
| xstage_ = nullptr; | |
| xstage_cap_ = 0; | |
| } | |
| const size_t total = (size_t) n * (size_t) blob; | |
| void* p = nullptr; | |
| if (cudaHostAlloc(&p, total, cudaHostAllocDefault) == cudaSuccess && p != nullptr) { | |
| xstage_pinned_ = true; | |
| } else { | |
| (void) cudaGetLastError(); | |
| p = std::malloc(total); | |
| xstage_pinned_ = false; | |
| if (p == nullptr) { err = "FileExpertSource: cannot allocate the exchange buffers"; return false; } | |
| } | |
| xstage_ = (uint8_t*) p; | |
| xstage_cap_ = n; | |
| xstage_blob_ = blob; | |
| return true; | |
| } | |
| uint8_t* FileExpertSource::exchange_buffer(int64_t q) const { | |
| if (xstage_ == nullptr || q < 0 || q >= xstage_cap_) return nullptr; | |
| return xstage_ + (size_t) q * (size_t) xstage_blob_; | |
| } | |
| bool FileExpertSource::stage_exchange(int64_t layer, int64_t in, int64_t out, int64_t q) { | |
| if (!has_resident(layer, in) || has_resident(layer, out) || exchange_buffer(q) == nullptr) return false; | |
| const size_t i_in = (size_t) layer * (size_t) n_expert_ + (size_t) in; | |
| const size_t i_out = (size_t) layer * (size_t) n_expert_ + (size_t) out; | |
| if (override_.empty()) override_.assign((size_t) blobs_, nullptr); | |
| if (override_[i_out] != nullptr) return false; | |
| for (const Exchange& x : staged_) | |
| if (x.in == i_in || x.q == q) return false; | |
| override_[i_out] = exchange_buffer(q); | |
| staged_.push_back({i_in, i_out, q, layer_blob_bytes_[(size_t) layer]}); | |
| return true; | |
| } | |
| int64_t FileExpertSource::commit_exchanges() { | |
| int64_t n = 0; | |
| for (const Exchange& x : staged_) { | |
| const uint8_t* src = override_[x.out]; | |
| const uint64_t at = complement_offsets_[x.in]; | |
| if (src != nullptr && at != kNoComplement && at <= complement_bytes_ && x.bytes <= complement_bytes_ - at && | |
| complement_host_ != nullptr) { | |
| std::memcpy((uint8_t*) complement_host_ + (size_t) at, src, (size_t) x.bytes); | |
| if (detail::exchange_cache_complement(complement_offsets_, x.in, x.out)) ++n; | |
| } | |
| override_[x.out] = nullptr; | |
| } | |
| staged_.clear(); | |
| exchanges_ += n; | |
| return n; | |
| } | |
| const uint8_t* FileExpertSource::blob(int64_t layer, int64_t expert) { | |
| if (base_ == nullptr || layer < 0 || expert < 0 || layer >= n_layers_ || expert >= n_expert_) return nullptr; | |
| const size_t index = (size_t) layer * (size_t) n_expert_ + (size_t) expert; | |
| const uint8_t* result = nullptr; | |
| bool from_files = true; | |
| if (complement_ready_) { | |
| result = detail::cache_complement_blob_or_fallback(index, complement_offsets_, complement_host_, nullptr); | |
| if (result != nullptr) { | |
| ram_reads_.fetch_add(1, std::memory_order_relaxed); | |
| from_files = false; | |
| } else if (!override_.empty() && override_[index] != nullptr) { | |
| result = override_[index]; | |
| from_files = false; | |
| } | |
| } | |
| if (from_files) { | |
| if (staged()) { | |
| result = staged_blob(layer, expert); // counts its bytes | |
| } else { | |
| result = mapped_blob(layer, expert); | |
| if (result != nullptr) | |
| file_read_bytes_.fetch_add(layer_blob_bytes_[(size_t) layer], std::memory_order_relaxed); | |
| } | |
| if (complement_ready_ && result != nullptr) file_reads_.fetch_add(1, std::memory_order_relaxed); | |
| } | |
| if (result != nullptr) ++reads_; | |
| return result; | |
| } | |
| bool FileExpertSource::pinned(int64_t layer, int64_t expert) const { | |
| if (!complement_ready_ || !complement_pinned_ || complement_host_ == nullptr || layer < 0 || expert < 0 || | |
| layer >= n_layers_ || expert >= n_expert_) return false; | |
| const size_t index = (size_t) layer * (size_t) n_expert_ + (size_t) expert; | |
| if (index >= complement_offsets_.size() || complement_offsets_[index] == kNoComplement) return false; | |
| // a partial pin (CS-T): only the registered prefix | |
| return !complement_partial_ || | |
| complement_offsets_[index] + layer_blob_bytes_[(size_t) layer] <= complement_pin_limit_; | |
| } | |
| const uint8_t* FileExpertSource::device_alias(int64_t layer, int64_t expert) const { | |
| if (!pinned(layer, expert) || complement_device_ == nullptr) return nullptr; | |
| const size_t index = (size_t) layer * (size_t) n_expert_ + (size_t) expert; | |
| return complement_device_ + (size_t) complement_offsets_[index]; | |
| } | |
| bool FileExpertSource::pcie_layer(int64_t layer) const { | |
| if (complement_ready_ && complement_pinned_ && complement_device_ != nullptr) | |
| return layer >= 0 && layer < n_layers_; | |
| return device_alias(layer, 0) != nullptr; | |
| } | |
| // ================================ THE ADAPTER ================================ | |
| void expert_pool_dispatch(void* user, const float* x_f, const int32_t* ids, const float* weights, int64_t n_embd, | |
| int64_t k, float* out) { | |
| (void) weights; // clause 2: `moe_combine` applies it on the device. Not an oversight. | |
| ExpertDispatch& d = *(ExpertDispatch*) user; | |
| if (d.failed) return; // a previous layer already failed; do not make it worse | |
| using namespace strata::kernels::cpu; | |
| // Clause 3: the blob's internal offsets are compile-time constants, so a mismatched geometry does not | |
| // produce a wrong answer - it produces a walk off the end of the blob into the next expert's bytes, which | |
| // is finite and plausible. Refuse, name the number, and let the driver report it. | |
| if (n_embd != H) { | |
| d.failed = true; | |
| d.fail = "the expert kernel is compiled for a 2560-wide activation"; | |
| d.fail_layer = d.layers; | |
| return; | |
| } | |
| if (expert_layout().native) { | |
| // plan v0.3 P6: a native pack runs its experts in verify windows only (the driver guarantees it) | |
| d.failed = true; | |
| d.fail = "the single-token expert path does not take a native (IQ) pack"; | |
| d.fail_layer = d.layers; | |
| return; | |
| } | |
| if (k > (int64_t) d.jobs.size()) d.jobs.resize((size_t) k); | |
| d.src->begin_layer(d.layers, ids, k); | |
| // Clause 1: rebuilt from `x_f` on EVERY call. `x_f` is mapped pinned memory whose address never changes, | |
| // so anything cached against it would be layer 0's activation reused 48 times. | |
| act_quant_q8_1(x_f, H, d.act); | |
| // ---- R4.2c: THE POOL'S HALF OF THE SPLIT. **IT DOES NOT DECIDE ANYTHING - `Launch` ALREADY DID.** | |
| // | |
| // The decision has to be made on THIS layer's ids, and `Launch` is the only callback that runs before the | |
| // pool while the ids are known (the doorbell publishes them when the ring fires). So `expert_hit_run` | |
| // decides, and this consumes `d.is_hit`. The first version decided here instead, which meant `Launch` | |
| // computed the PREVIOUS layer's experts into this layer's rows: C1 went from mean KL 9.69e-02 to 1.03e+00. | |
| // | |
| // `njobs` indexes the JOB ARRAY and `i` indexes the OUTPUT - they are the same only when nothing is a hit. | |
| const bool graph_hits = d.host_res != nullptr; | |
| const bool use_hits = graph_hits || (d.hits_ready() && d.decided); | |
| int64_t njobs = 0; | |
| if (d.remote_count > 0) { | |
| int32_t kind[32]; | |
| if (k > 32) { | |
| d.failed = true; d.fail = "remote experts: routing width exceeds 32"; return; | |
| } | |
| for (int64_t i = 0; i < k; ++i) | |
| kind[i] = use_hits && ids[i] >= 0 && ids[i] < d.n_expert && (graph_hits | |
| ? d.host_res[(size_t) d.layers * (size_t) d.n_expert + (size_t) ids[i]] >= 0 | |
| : d.is_hit[(size_t) i] != 0) ? 0 : -1; | |
| static thread_local std::string remote_error; | |
| for (int r = 0; r < d.remote_count; ++r) | |
| if (!d.remote[r]->begin(d.layers, x_f, ids, 1, k, kind, d.host_res, remote_error)) { | |
| d.failed = true; d.fail = remote_error.c_str(); d.fail_layer = d.layers; return; | |
| } | |
| } | |
| for (int64_t i = 0; i < k; ++i) { | |
| const int64_t e = ids[i]; | |
| if (e < 0 || e >= d.n_expert) { | |
| d.failed = true; | |
| d.fail = "a routed expert id is out of range"; | |
| d.fail_layer = d.layers; | |
| d.fail_expert = e; | |
| return; | |
| } | |
| const uint8_t* b = d.src->blob(d.layers, e); | |
| if (b == nullptr) { | |
| // The one failure the loop cannot see. Leaving `out` at its previous contents would feed the NEXT | |
| // layer a stale expert vector, which `moe_combine` would weight and add - the token would still be | |
| // finite and would still be wrong, 48 layers deep. | |
| d.failed = true; | |
| d.fail = "the expert source could not produce a blob"; | |
| d.fail_layer = d.layers; | |
| d.fail_expert = e; | |
| ++d.missing; | |
| return; | |
| } | |
| // A hit's row was zeroed by `Launch` and belongs to the GPU; the pool must not touch it. | |
| if (use_hits && (graph_hits ? d.host_res[(size_t) d.layers * (size_t) d.n_expert + (size_t) e] >= 0 | |
| : d.is_hit[(size_t) i] != 0)) { | |
| if (graph_hits) ++d.cache_hits; | |
| // The GPU owns this row and `hit_out` is zeroed, so the CPU's contribution is zero - but | |
| // `y_miss` is a REUSED pinned buffer, so the row must be written, not merely skipped. | |
| std::memset(out + (size_t) i * (size_t) n_embd, 0, (size_t) n_embd * sizeof(float)); | |
| continue; | |
| } | |
| if (graph_hits) ++d.cache_refused; // token graph: a miss (nothing is admitted during a token) | |
| bool remote_owns = false; | |
| for (int r = 0; r < d.remote_count; ++r) remote_owns |= d.remote[r]->owns(i); | |
| if (remote_owns) { | |
| std::memset(out + (size_t) i * (size_t) n_embd, 0, (size_t) n_embd * sizeof(float)); | |
| continue; | |
| } | |
| // `njobs` indexes the JOB ARRAY and `i` indexes the OUTPUT - they are the same only when nothing is a | |
| // hit, and using one for the other is how a hit's row would get two experts summed into it. | |
| ExpertJob& j = d.jobs[(size_t) njobs++]; | |
| j.blob = b; | |
| j.act = &d.act; // SHARED across the batch: one conversion serves all ten experts | |
| j.out = out + (size_t) i * (size_t) n_embd; | |
| j.weight = 1.0f; // clause 2: a diagnostic field, NOT the router weight | |
| j.slot = (int) i; | |
| } | |
| // Plan v0.3 P4: rows of every expert across all threads (bitwise the same as `run`). | |
| if (d.split_rows) d.pool->run_split(d.jobs.data(), (int) njobs); | |
| else d.pool->run(d.jobs.data(), (int) njobs); | |
| if (d.remote_count > 0) { | |
| static thread_local std::string remote_error; | |
| for (int r = 0; r < d.remote_count; ++r) | |
| if (!d.remote[r]->finish(out, remote_error)) { | |
| d.failed = true; d.fail = remote_error.c_str(); d.fail_layer = d.layers; return; | |
| } | |
| } | |
| ++d.layers; | |
| d.experts += k; | |
| } | |
| namespace { | |
| // the verify window's per-entry tables in `expert_pool_dispatch_multi` (`kind`, `distinct`, `first_of`) | |
| // are fixed arrays of this many entries: MAXT tokens of the model's 10 routed experts must fit, and a larger k is | |
| // refused at run time rather than written past them. | |
| constexpr int64_t kMaxWindowEntries = 128; | |
| static_assert(strata::kernels::cpu::MAXT * 10 <= kMaxWindowEntries, "a verify window's entries overflow the tables"); | |
| } // namespace | |
| void expert_pool_dispatch_multi(ExpertDispatch& d, const float* x_f, const int32_t* ids, int64_t n_tok, int64_t k, | |
| float* out) { | |
| using namespace strata::kernels::cpu; | |
| if (d.failed) return; | |
| if (n_tok < 1 || n_tok > MAXT) { | |
| d.failed = true; | |
| d.fail = "a verify window has more tokens than the multi-token expert kernel takes"; | |
| d.fail_layer = d.layers; | |
| return; | |
| } | |
| if (d.lookahead != nullptr) d.lookahead->submit(d.layers, x_f, n_tok, d.host_res); // CS-T: warm layer + 1 | |
| if (k < 1 || n_tok * k > kMaxWindowEntries) { | |
| d.failed = true; | |
| d.fail = "a verify window routes more entries than the expert pool's window tables hold"; | |
| d.fail_layer = d.layers; | |
| return; | |
| } | |
| if ((int64_t) d.act_multi.size() < n_tok) d.act_multi.resize((size_t) MAXT); | |
| const ExpertLayout& lay = expert_layout(); | |
| const bool native = lay.native; | |
| if (native && d.nact_multi.size() < (size_t) MAXT * kNativeActBytes) d.nact_multi.resize((size_t) MAXT * kNativeActBytes); | |
| if (d.job_of.size() != (size_t) d.n_expert) d.job_of.assign((size_t) d.n_expert, (int16_t) -1); | |
| if (d.jobs_multi.size() < (size_t) (n_tok * k)) d.jobs_multi.resize((size_t) (MAXT * k)); | |
| static const bool ptrace = std::getenv("STRATA_POOL_TRACE") != nullptr; | |
| auto pt = [&](const char* what, long long a = -1) { | |
| if (ptrace) { std::fprintf(stderr, "pool trace: layer %lld %s %lld\n", (long long) d.layers, what, a); std::fflush(stderr); } | |
| }; | |
| const auto c0 = std::chrono::steady_clock::now(); | |
| pt("begin"); | |
| d.src->begin_layer(d.layers, ids, n_tok * k); | |
| pt("begun"); | |
| if (!d.usage.empty()) | |
| for (int64_t i = 0; i < n_tok * k; ++i) | |
| if (ids[i] >= 0 && ids[i] < d.n_expert) d.usage[(size_t) d.layers * (size_t) d.n_expert + (size_t) ids[i]] += 1.0f; | |
| // ---- plan v0.3 P6: the GPU's share, decided and published FIRST so the GPU starts while the CPU works. | |
| // Distinct experts in routing order; resident ones and the last pcie_num/256 of the missed ones go to the GPU. | |
| const int64_t n = n_tok * k; | |
| int32_t kind[kMaxWindowEntries]; // per entry: -1 CPU, 0 VRAM, 1 PCIe | |
| if (d.plan != nullptr && n <= kMaxWindowEntries && n <= d.plan->cap) { | |
| int64_t distinct[kMaxWindowEntries], first_of[kMaxWindowEntries]; | |
| int nd = 0, nmiss = 0; | |
| for (int64_t i = 0; i < n; ++i) { | |
| first_of[i] = i; | |
| for (int64_t j = 0; j < i; ++j) | |
| if (ids[j] == ids[i]) { first_of[i] = first_of[j]; break; } | |
| if (first_of[i] == i) { | |
| distinct[nd++] = i; | |
| const int32_t e = ids[i]; | |
| if (e >= 0 && e < d.n_expert && d.host_res[(size_t) d.layers * (size_t) d.n_expert + (size_t) e] < 0 && | |
| !(d.peer != nullptr && d.peer->has(d.layers, e))) ++nmiss; | |
| } | |
| } | |
| const bool pcie_ok = d.pcie_num > 0 && d.src->pcie_layer(d.layers); | |
| const int m = pcie_ok ? (nmiss * d.pcie_num) >> 8 : 0; | |
| int miss_rank = 0, groups = 0, entries = 0, fetches = 0; | |
| GpuPlanSink& P = *d.plan; | |
| const uint8_t* dma_src[64]; | |
| int64_t pcie_i0[64]; | |
| for (int q = 0; q < nd; ++q) { | |
| const int64_t i0 = distinct[q]; | |
| const int32_t e = ids[i0]; | |
| int kd = -1; | |
| unsigned long long ptr = 0; | |
| if (e >= 0 && e < d.n_expert) { | |
| const int32_t slot = d.host_res[(size_t) d.layers * (size_t) d.n_expert + (size_t) e]; | |
| if (slot >= 0) { | |
| kd = 0; | |
| ptr = (unsigned long long) (d.cache_base + (d.cache_slot_off ? (size_t) d.cache_slot_off[slot] | |
| : (size_t) slot * (size_t) d.cache_blob)); | |
| } else if (d.peer != nullptr && d.peer->has(d.layers, e)) { | |
| kd = 2; // multi-GPU: the second GPU computes it | |
| } else { | |
| if (miss_rank >= nmiss - m && fetches < P.staging_cap && fetches < 64) { | |
| const uint8_t* src = d.src->pinned(d.layers, e) ? d.src->blob(d.layers, e) : nullptr; | |
| if (src != nullptr) { | |
| kd = 1; | |
| dma_src[fetches] = src; | |
| pcie_i0[fetches] = i0; | |
| ++fetches; | |
| } | |
| } | |
| ++miss_rank; | |
| } | |
| } | |
| for (int64_t i = i0; i < n; ++i) | |
| if (first_of[i] == i0) kind[i] = kd; | |
| if (kd != 0) continue; // the VRAM groups first; the PCIe groups below | |
| P.ptr[groups] = ptr; | |
| P.start[groups] = entries; | |
| for (int64_t i = i0; i < n; ++i) | |
| if (first_of[i] == i0) { | |
| P.dst[entries] = (int32_t) i; | |
| P.tok[entries] = (int32_t) (i / k); | |
| ++entries; | |
| } | |
| ++groups; | |
| } | |
| P.start[groups] = entries; | |
| const uint64_t bb = lay.blob_bytes(d.layers); | |
| for (int q = 0; q < fetches; ++q) { // the PCIe groups: staging slot q, entries after the VRAM ones | |
| const int64_t i0 = pcie_i0[q]; | |
| P.ptr2[q] = P.pcie_mode != 0 ? (unsigned long long) d.src->device_alias(d.layers, ids[i0]) | |
| : P.staging + (unsigned long long) q * (unsigned long long) bb; | |
| P.start2[q] = entries; | |
| for (int64_t i = i0; i < n; ++i) | |
| if (first_of[i] == i0) { | |
| P.dst[entries] = (int32_t) i; | |
| P.tok[entries] = (int32_t) (i / k); | |
| ++entries; | |
| } | |
| ++d.pcie_experts; | |
| } | |
| P.start2[fetches] = entries; | |
| P.counts[0] = groups; | |
| P.counts[1] = entries; | |
| P.counts[2] = fetches; | |
| std::atomic_thread_fence(std::memory_order_seq_cst); | |
| pt("publish", fetches); | |
| if (P.publish) P.publish(P.ctx); | |
| pt("fetch", fetches); | |
| if (P.fetch) P.fetch(P.ctx, dma_src, P.pcie_mode != 0 ? 0 : fetches, (size_t) bb); // the copy engine, beside the CPU's work | |
| } else { | |
| for (int64_t i = 0; i < n; ++i) { | |
| const int32_t e = ids[i]; | |
| kind[i] = (e >= 0 && e < d.n_expert && d.host_res != nullptr && | |
| d.host_res[(size_t) d.layers * (size_t) d.n_expert + (size_t) e] >= 0) ? 0 | |
| : (e >= 0 && e < d.n_expert && d.peer != nullptr && d.peer->has(d.layers, e)) ? 2 : -1; | |
| } | |
| } | |
| if (d.peer != nullptr) { // multi-GPU: start the second GPU's share before the CPU's own work | |
| std::string perr; | |
| if (!d.peer->launch(d.layers, x_f, ids, n_tok, k, kind, perr, out)) { | |
| std::fprintf(stderr, "strata: %s (layer %lld)\n", perr.c_str(), (long long) d.layers); | |
| d.failed = true; | |
| d.fail = "the peer GPU's experts could not be launched"; | |
| d.fail_layer = d.layers; | |
| return; | |
| } | |
| } | |
| if (d.remote_count > 0) { | |
| static thread_local std::string remote_error; | |
| for (int r = 0; r < d.remote_count; ++r) { | |
| if (!d.remote[r]->begin(d.layers, x_f, ids, n_tok, k, kind, d.host_res, remote_error)) { | |
| d.failed = true; d.fail = remote_error.c_str(); d.fail_layer = d.layers; return; | |
| } | |
| for (int64_t i = 0; i < n; ++i) if (d.remote[r]->owns(i)) kind[i] = 2; | |
| } | |
| } | |
| const auto c1 = std::chrono::steady_clock::now(); | |
| if (native && lay.fmt[(size_t) d.layers].gu_type == 42) // a native Q2_0 pack: the Q2_0 kernels' activations | |
| for (int64_t t = 0; t < n_tok; ++t) act_quant_any(x_f + (size_t) t * H, H, d.act_multi[(size_t) t]); | |
| else if (native) | |
| for (int64_t t = 0; t < n_tok; ++t) | |
| native_quant_act(lay.fmt[(size_t) d.layers], x_f + (size_t) t * H, d.nact_multi.data() + (size_t) t * kNativeActBytes); | |
| else | |
| for (int64_t t = 0; t < n_tok; ++t) act_quant_q8_1(x_f + (size_t) t * H, H, d.act_multi[(size_t) t]); | |
| const auto c2 = std::chrono::steady_clock::now(); | |
| { // CS-T: the experts the CPU computes, fetched together (the GGUF in place reads them on several threads) | |
| static thread_local std::vector<int64_t> miss; | |
| miss.clear(); | |
| for (int64_t i = 0; i < n_tok * k; ++i) | |
| if (kind[i] < 0 && ids[i] >= 0 && ids[i] < d.n_expert && | |
| std::find(miss.begin(), miss.end(), (int64_t) ids[i]) == miss.end()) | |
| miss.push_back(ids[i]); | |
| d.src->prefetch(d.layers, miss.data(), (int64_t) miss.size()); | |
| } | |
| int njobs = 0; | |
| for (int64_t t = 0; t < n_tok; ++t) | |
| for (int64_t j = 0; j < k; ++j) { | |
| const int64_t i = t * k + j; | |
| const int64_t e = ids[i]; | |
| float* row = out + (size_t) i * H; | |
| if (e < 0 || e >= d.n_expert) { | |
| d.failed = true; | |
| d.fail = "a routed expert id is out of range"; | |
| d.fail_layer = d.layers; | |
| d.fail_expert = e; | |
| return; | |
| } | |
| if (kind[i] >= 0) { // CUDA0, PCIe, or a remote/peer result staged into this row below | |
| if (kind[i] == 0) ++d.cache_hits; | |
| else if (kind[i] == 2 && d.peer != nullptr) ++d.peer_entries; | |
| // multi-GPU: a direct peer launch is writing this row right now - zeroing it would race it | |
| if (!(kind[i] == 2 && d.peer != nullptr && d.peer->launched_direct())) | |
| std::memset(row, 0, (size_t) H * sizeof(float)); | |
| continue; | |
| } | |
| ++d.cache_refused; | |
| int16_t& jo = d.job_of[(size_t) e]; | |
| if (jo < 0) { | |
| const uint8_t* b = d.src->blob(d.layers, e); | |
| if (b == nullptr) { | |
| d.failed = true; | |
| d.fail = "the expert source could not produce a blob"; | |
| d.fail_layer = d.layers; | |
| d.fail_expert = e; | |
| ++d.missing; | |
| return; | |
| } | |
| jo = (int16_t) njobs++; | |
| ExpertJobMulti& nj = d.jobs_multi[(size_t) jo]; | |
| nj.blob = b; | |
| nj.nt = 0; | |
| } | |
| ExpertJobMulti& jb = d.jobs_multi[(size_t) jo]; | |
| jb.act[jb.nt] = &d.act_multi[(size_t) t]; | |
| jb.nact[jb.nt] = native ? d.nact_multi.data() + (size_t) t * kNativeActBytes : nullptr; | |
| jb.out[jb.nt] = row; | |
| ++jb.nt; | |
| ++d.multi_entries; | |
| } | |
| const auto c3 = std::chrono::steady_clock::now(); | |
| pt("run", njobs); | |
| if (native) d.pool->run_split_multi_native(lay.fmt[(size_t) d.layers], d.jobs_multi.data(), njobs); | |
| else d.pool->run_split_multi(d.jobs_multi.data(), njobs); | |
| if (d.remote_count > 0) { | |
| static thread_local std::string remote_error; | |
| for (int r = 0; r < d.remote_count; ++r) | |
| if (!d.remote[r]->finish(out, remote_error)) { | |
| d.failed = true; d.fail = remote_error.c_str(); d.fail_layer = d.layers; return; | |
| } | |
| } | |
| if (d.peer != nullptr) { // multi-GPU: the second GPU's rows, into the same mapped rows | |
| std::string perr; | |
| if (!d.peer->finish(out, perr)) { | |
| std::fprintf(stderr, "strata: %s (layer %lld)\n", perr.c_str(), (long long) d.layers); | |
| d.failed = true; | |
| d.fail = "the peer GPU's experts failed"; | |
| d.fail_layer = d.layers; | |
| return; | |
| } | |
| } | |
| const auto c4 = std::chrono::steady_clock::now(); | |
| pt("ran"); | |
| auto ms = [](auto a, auto b) { return std::chrono::duration<double, std::milli>(b - a).count(); }; | |
| d.ms_plan += ms(c0, c1); | |
| d.ms_actq += ms(c1, c2); | |
| d.ms_jobs += ms(c2, c3); | |
| d.ms_run += ms(c3, c4); | |
| for (int64_t i = 0; i < n_tok * k; ++i) { | |
| const int64_t e = ids[i]; | |
| if (e >= 0 && e < d.n_expert) d.job_of[(size_t) e] = -1; | |
| } | |
| d.multi_misses += njobs; | |
| ++d.layers; | |
| d.experts += n_tok * k; | |
| } | |
| void expert_hit_run(void* user, void* stream, HitPhase phase, const int32_t* ids, int64_t k) { | |
| ExpertDispatch& d = *(ExpertDispatch*) user; | |
| if (d.failed) return; | |
| cudaStream_t cs = (cudaStream_t) stream; | |
| if (phase == HitPhase::Launch) { | |
| d.decided = false; | |
| d.hit_pending = false; | |
| if (!d.hits_ready() || ids == nullptr || k <= 0) return; | |
| if ((int64_t) d.is_hit.size() < k) d.is_hit.resize((size_t) k); | |
| // ================================ THE DECISION, ONCE, ON THIS LAYER'S IDS ================================ | |
| // | |
| // Every routed expert is asked of the cache. Resident -> the GPU computes it. Not resident -> it is | |
| // admitted and filled if there is room (which makes it a hit on THIS call, because the fill and the | |
| // kernel are on one stream in that order), and otherwise it stays a miss for the CPU. | |
| d.n_hits = 0; | |
| for (int64_t i = 0; i < k; ++i) { | |
| const int64_t e = ids[i]; | |
| d.is_hit[(size_t) i] = 0; | |
| if (e < 0 || e >= d.n_expert) continue; // out of range: the pool refuses it, with a message | |
| int32_t slot = d.cache->slot_of(d.layers, e); | |
| if (slot == kNotResident) { | |
| const int32_t cand = d.cache->admit(d.layers, e); | |
| if (cand == kNotResident) { | |
| ++d.cache_refused; | |
| continue; | |
| } | |
| // `blob` is asked ONLY for an expert about to be filled, so the source's read counter stays a | |
| // count of distinct experts moved rather than of looks. | |
| const uint8_t* b = d.src->blob(d.layers, e); | |
| std::string ferr; | |
| if (b == nullptr || !d.cache->fill_slot(cand, b, cs, ferr, (int64_t) strata::kernels::cpu::expert_layout().blob_bytes(d.layers))) { | |
| d.failed = true; | |
| d.fail = "the expert cache could not fill a slot"; | |
| d.fail_layer = d.layers; | |
| d.fail_expert = e; | |
| return; | |
| } | |
| ++d.cache_admitted; | |
| slot = cand; | |
| } else { | |
| ++d.cache_hits; | |
| } | |
| d.is_hit[(size_t) i] = 1; | |
| d.h_slot[(size_t) d.n_hits] = slot; | |
| d.h_dst[(size_t) d.n_hits] = (int32_t) i; | |
| ++d.n_hits; | |
| } | |
| d.decided = true; | |
| if (d.n_hits <= 0) return; // nothing resident yet: no GPU work, and nothing for `Combine` to add | |
| const size_t list_bytes = (size_t) d.n_hits * sizeof(int32_t); | |
| // `hit_out` is ZEROED rather than overwritten: the kernel writes only the rows this layer's hits own, | |
| // so a row that was a hit last layer and a miss this one would still hold last layer's expert and | |
| // `add_inplace` would sum it in. Finite, plausible, wrong. | |
| if (cudaMemsetAsync(d.hit_out, 0, (size_t) d.parts_elems * sizeof(float), cs) != cudaSuccess || | |
| cudaMemcpyAsync(d.d_slot, d.h_slot.data(), list_bytes, cudaMemcpyHostToDevice, cs) != cudaSuccess || | |
| cudaMemcpyAsync(d.d_dst, d.h_dst.data(), list_bytes, cudaMemcpyHostToDevice, cs) != cudaSuccess) { | |
| d.hit_fail = "the hit list could not be staged"; | |
| d.failed = true; | |
| d.fail = d.hit_fail; | |
| return; | |
| } | |
| // The activation is quantized HERE rather than reused from `s.moe.x_q8_0`, which `post[l-1]` wrote from | |
| // the PREVIOUS layer's `mixed`. `pre[l]` has since overwritten `mixed`, so that buffer is a layer stale | |
| // - and a stale activation produces a perfectly finite expert for the wrong input. | |
| // **R4.2h: THE SCALED QUANTIZER, SO A HIT REPRODUCES A MISS.** The CPU pool quantizes this same | |
| // activation with `act_quant_q8_1` and multiplies by the fp32 `ActQ::scale`; `quantize_q8_0` writes | |
| // an fp16 `d` instead, and `bench/micro/act_quant_parity.cu` measured **80 of 80 chunks differing by | |
| // up to 4.761e-04 relative**. `quantize_q8_0_scaled` adopts the CPU's rule and scale, and the kernel | |
| // takes the fp32 array. Falling back to the old path would silently reintroduce the divergence, so | |
| // the scales are required here rather than optional. | |
| if (d.x_q8_0_hit_scale == nullptr) { | |
| d.failed = true; | |
| d.fail = "the hit path has no fp32 activation scales (R4.2h)"; | |
| return; | |
| } | |
| strata::kernels::quantize_q8_0_scaled(d.mixed, d.x_q8_0_hit, d.x_q8_0_hit_scale, strata::kernels::cpu::H, | |
| cs); | |
| if (d.hit_cpu_order) | |
| strata::kernels::moe_hit_grouped_s2_cpu_order(d.cache_base, d.d_slot, d.d_dst, d.n_hits, | |
| d.cache_blob, d.x_q8_0_hit, d.hit_scratch, d.hit_out, cs, d.x_q8_0_hit_scale); | |
| else | |
| strata::kernels::moe_hit_grouped_s2(d.cache_base, d.d_slot, d.d_dst, d.n_hits, d.cache_blob, | |
| d.x_q8_0_hit, d.hit_scratch, d.hit_out, cs, d.x_q8_0_hit_scale); | |
| d.hit_pending = true; | |
| if (d.hit_done != nullptr) cudaEventRecord((cudaEvent_t) d.hit_done, cs); | |
| // The A/B arm: ONE driver entry here, and nothing else changes. If the work was waiting for the host | |
| // to enter the driver, this is what lets it start while the pool runs. | |
| if (d.hit_poke && d.hit_done != nullptr) (void) cudaEventQuery((cudaEvent_t) d.hit_done); | |
| return; | |
| } | |
| // Combine: `parts += hit_out`, stream-ordered after the misses were copied into `parts`. | |
| if (!d.hit_pending) return; | |
| d.hit_pending = false; | |
| // Did the GPU get the hit work done while the CPU was in the pool? This query is itself a driver entry, | |
| // so it is the LAST chance to observe a late start: a NOT-READY here means the work had not finished by the | |
| // time the pool returned, and with no poke in front of it that can only be because it began after. | |
| if (d.hit_done != nullptr) { | |
| if (cudaEventQuery((cudaEvent_t) d.hit_done) == cudaSuccess) ++d.hit_ready; | |
| else ++d.hit_late; | |
| } | |
| strata::kernels::add_inplace(d.parts_out, d.hit_out, d.parts_elems, cs); | |
| } | |
| // ================================ THE RESIDENT ARENA (R2.1) ================================ | |
| namespace { | |
| void hash_u64(uint64_t& h, uint64_t v) { | |
| h = fnv1a64((const uint8_t*) &v, sizeof v, h); | |
| } | |
| void hash_text(uint64_t& h, const std::string& s) { | |
| h = fnv1a64((const uint8_t*) s.data(), (uint64_t) s.size(), h); | |
| } | |
| bool hash_small_file(const std::filesystem::path& path, uint64_t& h, std::string& err) { | |
| hash_text(h, path.filename().string()); | |
| std::ifstream f(path, std::ios::binary); | |
| if (!f) { | |
| hash_u64(h, 0); | |
| return true; | |
| } | |
| hash_u64(h, 1); | |
| std::vector<uint8_t> buf(64u << 10); | |
| for (;;) { | |
| f.read((char*) buf.data(), (std::streamsize) buf.size()); | |
| const std::streamsize n = f.gcount(); | |
| if (n > 0) h = fnv1a64(buf.data(), (uint64_t) n, h); | |
| if (f.eof()) break; | |
| if (!f) { | |
| err = "ArenaExpertSource: cannot hash pack metadata " + path.string(); | |
| return false; | |
| } | |
| } | |
| return true; | |
| } | |
| bool hash_sampled_file(const std::filesystem::path& path, uint64_t& h, std::string& err) { | |
| std::ifstream f(path, std::ios::binary | std::ios::ate); | |
| if (!f) { | |
| err = "ArenaExpertSource: cannot sample pack source " + path.string(); | |
| return false; | |
| } | |
| const std::streamoff end = f.tellg(); | |
| if (end < 0) { | |
| err = "ArenaExpertSource: cannot size pack source " + path.string(); | |
| return false; | |
| } | |
| const uint64_t bytes = (uint64_t) end; | |
| hash_text(h, path.filename().string()); | |
| hash_u64(h, bytes); | |
| constexpr uint64_t sample = 64u << 10; | |
| const uint64_t starts[3] = {0, bytes / 2, bytes > sample ? bytes - sample : 0}; | |
| std::vector<uint8_t> buf((size_t) std::min<uint64_t>(sample, bytes)); | |
| for (uint64_t off : starts) { | |
| if (buf.empty()) break; | |
| const uint64_t at = std::min<uint64_t>(off, bytes - (uint64_t) buf.size()); | |
| f.clear(); | |
| f.seekg((std::streamoff) at); | |
| f.read((char*) buf.data(), (std::streamsize) buf.size()); | |
| if ((size_t) f.gcount() != buf.size()) { | |
| err = "ArenaExpertSource: short read while hashing pack source " + path.string(); | |
| return false; | |
| } | |
| hash_u64(h, at); | |
| h = fnv1a64(buf.data(), (uint64_t) buf.size(), h); | |
| } | |
| return true; | |
| } | |
| bool shared_arena_pack_hash(const std::string& pack_dir, const std::string& experts_path, | |
| const std::string& gguf, const strata::kernels::cpu::ExpertLayout& lay, | |
| uint64_t& out, std::string& err) { | |
| uint64_t h = 1469598103934665603ull; | |
| hash_text(h, "strata-shared-expert-arena-pack-v1"); | |
| hash_u64(h, (uint64_t) lay.n_layers); | |
| hash_u64(h, (uint64_t) lay.n_expert); | |
| hash_u64(h, lay.total); | |
| hash_u64(h, lay.max_blob); | |
| hash_u64(h, lay.native ? 1 : 0); | |
| const std::filesystem::path pack(pack_dir); | |
| for (const char* name : {"manifest.json", "index.txt", "native_experts.txt"}) { | |
| if (!hash_small_file(pack / name, h, err)) return false; | |
| } | |
| if (std::filesystem::exists(experts_path)) { | |
| if (!hash_sampled_file(experts_path, h, err)) return false; | |
| } else if (!gguf.empty()) { | |
| // Native packs may read experts straight from one or more GGUF shards. Sample every distinct source | |
| // file named by native_experts.txt; this keeps the fingerprint cheap while still tying it to the model | |
| // bytes rather than only to an equal-size layout. | |
| const std::filesystem::path first(gguf); | |
| std::vector<std::filesystem::path> sources{first}; | |
| for (const std::string& name : lay.gguf_file) { | |
| if (name.empty()) continue; | |
| const std::filesystem::path p = first.parent_path() / name; | |
| if (std::find(sources.begin(), sources.end(), p) == sources.end()) sources.push_back(p); | |
| } | |
| for (const auto& p : sources) { | |
| if (!hash_sampled_file(p, h, err)) return false; | |
| } | |
| } | |
| out = h == 0 ? 1 : h; | |
| return true; | |
| } | |
| } // namespace | |
| namespace { | |
| /// The GGUF file that holds role `r` (0 gate, 1 up, 2 down) of layer `l`: a name beside the --native shard | |
| /// (native_experts.txt v3 per layer, v4 per role), or the --native shard itself. | |
| std::string expert_gguf_file(const std::string& gguf, const strata::kernels::cpu::ExpertLayout& lay, int64_t l, int r) { | |
| const size_t i = (size_t) (3 * l + r); | |
| if (lay.gguf_file.size() <= i || lay.gguf_file[i].empty()) return gguf; | |
| const size_t cut = gguf.find_last_of("/\\"); | |
| return (cut == std::string::npos ? std::string() : gguf.substr(0, cut + 1)) + lay.gguf_file[i]; | |
| } | |
| } // namespace | |
| bool check_experts_gguf(const std::string& gguf, const strata::kernels::cpu::ExpertLayout& lay, std::string& err) { | |
| static const char* roles[3] = {"gate", "up", "down"}; | |
| if (lay.gguf_off.size() != (size_t) (3 * lay.n_layers)) { | |
| err = "native_experts.txt has no GGUF offsets (a pack older than v2): repack it with tools/iq_pack.py"; | |
| return false; | |
| } | |
| try { | |
| std::map<std::string, std::unique_ptr<strata::GgufFile>> files; | |
| for (int64_t l = 0; l < lay.n_layers; ++l) { | |
| const auto& fm = lay.fmt[(size_t) l]; | |
| const uint64_t blob = lay.bytes[(size_t) l]; | |
| const uint64_t per[3] = {fm.up_off, fm.up_off, blob - fm.down_off}; | |
| for (int r = 0; r < 3; ++r) { | |
| const std::string path = expert_gguf_file(gguf, lay, l, r); | |
| auto& f = files[path]; | |
| if (!f) f = std::make_unique<strata::GgufFile>(path); | |
| const std::string name = "blk." + std::to_string(l) + ".ffn_" + roles[r] + "_exps.weight"; | |
| const strata::TensorInfo* t = f->find(name); | |
| const uint64_t want_type = (uint64_t) (r < 2 ? fm.gu_type : fm.d_type); | |
| // GGUF order: dim 0 is the row (the input), dim 1 the rows, dim 2 the experts | |
| const uint64_t cols = (uint64_t) (r < 2 ? fm.n_embd : fm.n_ff); | |
| const uint64_t rows = (uint64_t) (r < 2 ? fm.n_ff : fm.n_embd); | |
| const uint64_t bytes = per[r] * (uint64_t) lay.n_expert; | |
| const uint64_t payload = f->file_size() - f->data_start(); | |
| std::string why; | |
| if (t == nullptr) why = "is not in it"; | |
| else if (t->type != want_type) | |
| why = std::string("is ") + t->type_name() + ", the pack says type " + std::to_string(want_type); | |
| else if (t->shape.size() != 3 || t->shape[0] != cols || t->shape[1] != rows || | |
| t->shape[2] != (uint64_t) lay.n_expert) | |
| why = "is not [" + std::to_string(cols) + ", " + std::to_string(rows) + ", " + | |
| std::to_string(lay.n_expert) + "]"; | |
| else if (strata::tensor_payload_bytes(*t) != bytes) | |
| why = "is not " + std::to_string(per[r]) + " B per expert"; | |
| else if (f->data_start() + t->offset != lay.gguf_off[(size_t) (3 * l + r)]) | |
| why = "starts at byte " + std::to_string(f->data_start() + t->offset) + ", the pack says " + | |
| std::to_string(lay.gguf_off[(size_t) (3 * l + r)]); | |
| else if (t->offset > payload || bytes > payload - t->offset) | |
| why = "runs past the end of the file (a truncated shard?)"; | |
| if (!why.empty()) { | |
| err = "the pack's native_experts.txt does not match the model: " + name + " in " + path + " " + | |
| why + " - repack with tools/iq_pack.py from this model's shards"; | |
| return false; | |
| } | |
| } | |
| } | |
| return true; | |
| } catch (const std::exception& e) { | |
| err = std::string("native experts from the GGUF: ") + e.what(); | |
| return false; | |
| } | |
| } | |
| // Plan v0.3 P6: the arena from the model's GGUF shards. Each layer's gate, up and down tensors hold the 512 | |
| // experts one after another; they are read in chunks and each expert's slice lands at its place in the blob | |
| // [gate rows | up rows | down rows] - the layout tools/iq_pack.py would have written to experts.bin. Each role | |
| // is read from its own file (native_experts.txt v4: a shard boundary can fall inside a layer; per role as in | |
| // #255, gopinath87607). The caller checks the spans first (check_experts_gguf). | |
| // `unbuffered` (Windows, experts_unbuffered): each chunk's 4 KiB-aligned window is read with FILE_FLAG_NO_BUFFERING into | |
| // an aligned buffer and scattered into the blobs - no copy through the file cache when the drive is read anyway. | |
| LoadStats load_experts_gguf(const std::string& gguf, uint8_t* dst, const strata::kernels::cpu::ExpertLayout& lay, | |
| int threads, bool unbuffered) { | |
| LoadStats st; | |
| st.layers = (uint64_t) lay.n_layers; | |
| const auto t0 = std::chrono::steady_clock::now(); | |
| std::atomic<int64_t> next{0}; | |
| std::atomic<bool> bad{false}; | |
| if (unbuffered) { | |
| uint64_t max_chunk = 0; | |
| for (int64_t l = 0; l < lay.n_layers; ++l) { | |
| const auto& fm = lay.fmt[(size_t) l]; | |
| max_chunk = std::max<uint64_t>(max_chunk, std::max<uint64_t>(fm.up_off, lay.bytes[(size_t) l] - fm.down_off) * 16); | |
| } | |
| std::mutex err_mu; | |
| std::string err; | |
| auto worker = [&]() { | |
| constexpr uint64_t kSector = 4096; | |
| const uint64_t cap = (max_chunk + 2 * kSector + kSector - 1) / kSector * kSector; | |
| uint8_t* buf = (uint8_t*) VirtualAlloc(nullptr, (size_t) cap, MEM_COMMIT | MEM_RESERVE, PAGE_READWRITE); | |
| HANDLE h = INVALID_HANDLE_VALUE; | |
| std::string open_name; | |
| auto fail = [&](const std::string& what) { | |
| std::lock_guard<std::mutex> g(err_mu); | |
| if (err.empty()) err = what; | |
| bad = true; | |
| }; | |
| if (buf == nullptr) fail("cannot allocate a read buffer"); | |
| for (;;) { | |
| const int64_t l = next.fetch_add(1); | |
| if (l >= lay.n_layers || bad) break; | |
| const auto& fm = lay.fmt[(size_t) l]; | |
| const uint64_t blob = lay.bytes[(size_t) l]; | |
| const uint64_t per[3] = {fm.up_off, fm.up_off, blob - fm.down_off}; | |
| const uint64_t at[3] = {0, fm.up_off, fm.down_off}; | |
| for (int r = 0; r < 3 && !bad; ++r) { | |
| // the handle is kept while consecutive roles share a file (every layer of a v3 pack) | |
| const std::string name = expert_gguf_file(gguf, lay, l, r); | |
| if (name != open_name) { | |
| if (h != INVALID_HANDLE_VALUE) CloseHandle(h); | |
| const int wide = MultiByteToWideChar(CP_UTF8, 0, name.c_str(), -1, nullptr, 0); | |
| std::vector<wchar_t> w((size_t) std::max(wide, 1), L'\0'); | |
| if (wide > 0) MultiByteToWideChar(CP_UTF8, 0, name.c_str(), -1, w.data(), wide); | |
| h = CreateFileW(w.data(), GENERIC_READ, FILE_SHARE_READ, nullptr, OPEN_EXISTING, | |
| FILE_FLAG_NO_BUFFERING | FILE_FLAG_SEQUENTIAL_SCAN, nullptr); | |
| if (h == INVALID_HANDLE_VALUE) { | |
| fail("cannot open " + name + " (error " + std::to_string((unsigned long long) GetLastError()) + ")"); | |
| break; | |
| } | |
| open_name = name; | |
| } | |
| const uint64_t src = lay.gguf_off[(size_t) (3 * l + r)]; | |
| const uint64_t total = per[r] * (uint64_t) lay.n_expert; | |
| const uint64_t chunk = per[r] * 16; // 16 experts per read | |
| for (uint64_t done = 0; done < total; done += chunk) { | |
| const uint64_t n = std::min<uint64_t>(chunk, total - done); | |
| const uint64_t a0 = (src + done) / kSector * kSector; | |
| const uint64_t a1 = (src + done + n + kSector - 1) / kSector * kSector; | |
| OVERLAPPED ov{}; | |
| ov.Offset = (DWORD) a0; | |
| ov.OffsetHigh = (DWORD) (a0 >> 32); | |
| DWORD got = 0; | |
| // the window may run past the end of the file: only the tensor's own bytes have to arrive | |
| if (!ReadFile(h, buf, (DWORD) (a1 - a0), &got, &ov) || (uint64_t) got < src + done - a0 + n) { | |
| fail("short unbuffered read of layer " + std::to_string(l) + " in " + name + " (error " + | |
| std::to_string((unsigned long long) GetLastError()) + ")"); | |
| break; | |
| } | |
| const uint8_t* q = buf + (src + done - a0); | |
| for (uint64_t k = 0; k < n / per[r]; ++k) { | |
| const uint64_t e = done / per[r] + k; | |
| std::memcpy(dst + lay.blob_offset(l, (int64_t) e) + at[r], q + k * per[r], (size_t) per[r]); | |
| } | |
| } | |
| } | |
| } | |
| if (h != INVALID_HANDLE_VALUE) CloseHandle(h); | |
| if (buf != nullptr) VirtualFree(buf, 0, MEM_RELEASE); | |
| }; | |
| std::vector<std::thread> pool; | |
| for (int i = 1; i < threads; ++i) pool.emplace_back(worker); | |
| worker(); | |
| for (auto& t : pool) t.join(); | |
| if (bad) { | |
| st.seconds = -1.0; | |
| st.ok = false; | |
| st.error = err.empty() ? "unreadable shard while reading the experts from the GGUF" : err; | |
| return st; | |
| } | |
| st.bytes = lay.total; | |
| st.seconds = std::chrono::duration<double>(std::chrono::steady_clock::now() - t0).count(); | |
| return st; | |
| } | |
| (void) unbuffered; | |
| auto worker = [&]() { | |
| // #230: `fread` on a `FILE*`, as load_experts_ranges (#89): MSVC's `std::ifstream::read` splits a request | |
| // into 4095-byte freads, which took this path to 0.02 GiB/s on a Windows install without experts.bin. | |
| // The guard closes the handle on every return. | |
| struct Closer { | |
| FILE* f = nullptr; | |
| ~Closer() { if (f != nullptr) std::fclose(f); } | |
| } file; | |
| std::string open_name; | |
| std::vector<uint8_t> buf; | |
| for (;;) { | |
| const int64_t l = next.fetch_add(1); | |
| if (l >= lay.n_layers || bad) break; | |
| const auto& fm = lay.fmt[(size_t) l]; | |
| const uint64_t blob = lay.bytes[(size_t) l]; | |
| const uint64_t per[3] = {fm.up_off, fm.up_off, blob - fm.down_off}; | |
| const uint64_t at[3] = {0, fm.up_off, fm.down_off}; | |
| for (int r = 0; r < 3; ++r) { | |
| // the handle is kept while consecutive roles share a file (every layer of a v3 pack) | |
| const std::string name = expert_gguf_file(gguf, lay, l, r); | |
| if (name != open_name) { | |
| if (file.f != nullptr) std::fclose(file.f); | |
| file.f = std::fopen(name.c_str(), "rb"); | |
| if (file.f == nullptr) { bad = true; return; } | |
| open_name = name; | |
| } | |
| FILE* const f = file.f; | |
| const uint64_t src = lay.gguf_off[(size_t) (3 * l + r)]; | |
| const uint64_t total = per[r] * (uint64_t) lay.n_expert; | |
| const uint64_t chunk = per[r] * 16; // 16 experts per read | |
| buf.resize((size_t) chunk); | |
| for (uint64_t done = 0; done < total; done += chunk) { | |
| const uint64_t n = std::min<uint64_t>(chunk, total - done); | |
| // 64-bit seek: a shard is tens of GB | |
| if (STRATA_FSEEK64(f, src + done) != 0) { bad = true; return; } | |
| if (std::fread(buf.data(), 1, (size_t) n, f) != (size_t) n) { bad = true; return; } | |
| for (uint64_t k = 0; k < n / per[r]; ++k) { | |
| const uint64_t e = done / per[r] + k; | |
| std::memcpy(dst + lay.blob_offset(l, (int64_t) e) + at[r], buf.data() + k * per[r], (size_t) per[r]); | |
| } | |
| } | |
| } | |
| } | |
| }; | |
| std::vector<std::thread> pool; | |
| for (int i = 1; i < threads; ++i) pool.emplace_back(worker); | |
| worker(); | |
| for (auto& t : pool) t.join(); | |
| if (bad) { | |
| st.seconds = -1.0; | |
| st.ok = false; | |
| st.error = "short read or unreadable shard while reading the experts from the GGUF"; | |
| return st; | |
| } | |
| st.bytes = lay.total; | |
| st.seconds = std::chrono::duration<double>(std::chrono::steady_clock::now() - t0).count(); | |
| return st; | |
| } | |
| ArenaExpertSource::~ArenaExpertSource() { close(); } | |
| bool ArenaExpertSource::open(const std::string& pack_dir, int64_t n_layers, int64_t n_expert, int threads, | |
| std::string& err, uint64_t max_pinned_bytes, | |
| const std::string& shared_arena_file) { | |
| close(); | |
| const std::string path = pack_dir + "/experts.bin"; | |
| // plan v0.3 P6: the layout (canonical, or a native pack's per-layer blobs) was loaded by the driver | |
| const strata::kernels::cpu::ExpertLayout& lay = strata::kernels::cpu::expert_layout(); | |
| if (lay.n_layers != n_layers || lay.n_expert != n_expert) { | |
| err = "ArenaExpertSource: the expert layout was loaded for a different geometry"; | |
| return false; | |
| } | |
| const int64_t blob = (int64_t) lay.max_blob; | |
| const uint64_t want = lay.total; | |
| // plan v0.3 P6: no experts.bin in a native pack -> the experts come straight from the GGUF | |
| const bool from_gguf = !std::ifstream(path, std::ios::binary) && lay.native && !lay.gguf_off.empty() && !gguf_.empty(); | |
| // SIZE CHECK BEFORE THE ALLOCATION, not after. A wrong pack should name the two numbers rather than spend | |
| // 34 GB and a minute of loading first. | |
| if (from_gguf) { | |
| // every (file, offset) of native_experts.txt must be the tensor it claims, of the pack's type and | |
| // dimensions and inside its file - before the allocation, so a pack of another model or a truncated | |
| // shard is a message rather than an arena of plausible wrong experts | |
| if (!check_experts_gguf(gguf_, lay, err)) { err = "ArenaExpertSource: " + err; return false; } | |
| } else { | |
| std::ifstream f(path, std::ios::binary | std::ios::ate); | |
| if (!f) { err = "ArenaExpertSource: cannot open " + path; return false; } | |
| const uint64_t got = (uint64_t) f.tellg(); | |
| if (got != want) { | |
| char buf[400]; | |
| std::snprintf(buf, sizeof buf, | |
| "ArenaExpertSource: %s is %llu B but %lld layers x %lld experts (blobs up to %lld B) " | |
| "make %llu B - this is not the pack this geometry came from", | |
| path.c_str(), (unsigned long long) got, (long long) n_layers, (long long) n_expert, | |
| (long long) blob, (unsigned long long) want); | |
| err = buf; | |
| return false; | |
| } | |
| } | |
| uint64_t pack_hash = 0; | |
| if (!shared_arena_file.empty() && | |
| !shared_arena_pack_hash(pack_dir, path, gguf_, lay, pack_hash, err)) return false; | |
| // one layer per registration slice, so no expert straddles two registrations. The arena is one blob | |
| // longer than the file: a copy of a whole VRAM slot (the largest blob) may then start at any expert. | |
| std::vector<uint64_t> bounds, loff, lbytes; | |
| for (int64_t l = 0; l < n_layers; ++l) { | |
| bounds.push_back(lay.layer_offset(l)); | |
| loff.push_back(lay.layer_offset(l)); | |
| lbytes.push_back(lay.blob_bytes(l) * (uint64_t) n_expert); | |
| } | |
| bounds.push_back(want); | |
| PinnedArena* a = new PinnedArena(want + (uint64_t) blob, bounds, max_pinned_bytes, | |
| shared_arena_file, pack_hash); | |
| if (!a->valid()) { | |
| const std::string why = a->note; | |
| delete a; | |
| err = "ArenaExpertSource: the arena could not be reserved (" + | |
| std::to_string(want + (uint64_t) blob) + " B)" + | |
| (why.empty() ? std::string{} : ": " + why); | |
| return false; | |
| } | |
| // #285: unbuffered when the drive is read anyway and the file cache could not keep the experts for the next | |
| // start either (a 64 GB PC); otherwise the buffered readers, which a warm restart serves from the cache | |
| std::vector<std::string> files; | |
| if (from_gguf) { | |
| for (int64_t l = 0; l < n_layers; ++l) | |
| for (int r = 0; r < 3; ++r) { | |
| const std::string f = expert_gguf_file(gguf_, lay, l, r); | |
| if (std::find(files.begin(), files.end(), f) == files.end()) files.push_back(f); | |
| } | |
| } else { | |
| files.push_back(path); | |
| } | |
| std::string why; | |
| const bool unbuffered = experts_unbuffered(files, want + (uint64_t) blob, why); | |
| const int readers = unbuffered ? std::max(threads, 16) : threads; // 16 keep a PCIe 5 drive's queue full | |
| LoadStats st; | |
| if (from_gguf) { | |
| st = load_experts_gguf(gguf_, a->data(), lay, readers, unbuffered); | |
| } else { | |
| if (unbuffered) st = load_experts_direct(path, a->data(), loff, lbytes, readers, /*chunk=*/8u << 20); | |
| if (!unbuffered || (!st.ok && st.error.empty())) // unaligned ranges: the buffered reader | |
| st = load_experts_ranges(path, a->data(), loff, lbytes, threads, /*chunk=*/8u << 20); | |
| } | |
| std::fprintf(stderr, "strata generate: expert arena read %s (%s)\n", unbuffered ? "unbuffered" : "through the file cache", | |
| why.c_str()); | |
| if (!st.ok) { | |
| delete a; | |
| err = "ArenaExpertSource: the expert load was refused: " + (st.error.empty() ? std::string("unknown") : st.error); | |
| return false; | |
| } | |
| if (st.bytes != want) { | |
| delete a; | |
| err = "ArenaExpertSource: the load read " + std::to_string(st.bytes) + " B of " + std::to_string(want); | |
| return false; | |
| } | |
| arena_ = a; | |
| base_ = a->data(); | |
| pinned_bytes_ = a->registered_bytes; | |
| // plan v0.3 P6: device aliases of the mapped registration, for the PCIe share of the misses | |
| dev_slice_.clear(); | |
| slice_bytes_ = a->slice_bytes; | |
| if (a->registered_bytes > 0) { | |
| std::vector<uint64_t> starts = a->slice_bytes > 0 ? a->slice_starts : std::vector<uint64_t>{0}; | |
| for (uint64_t off : starts) { | |
| void* d = nullptr; | |
| if (cudaHostGetDevicePointer(&d, (void*) (base_ + off), 0) != cudaSuccess) { | |
| (void) cudaGetLastError(); | |
| dev_slice_.clear(); | |
| break; | |
| } | |
| dev_slice_.push_back((const uint8_t*) d); | |
| } | |
| } | |
| blobs_ = n_layers * n_expert; | |
| n_expert_ = n_expert; | |
| reads_ = 0; | |
| note_ = a->note; | |
| gib_per_s_ = st.gib_per_second(); | |
| load_seconds_ = st.seconds; | |
| load_read_s_ = st.read_seconds; | |
| load_copy_s_ = st.copy_seconds; | |
| return true; | |
| } | |
| void ArenaExpertSource::close() { | |
| if (arena_ != nullptr) { | |
| delete (PinnedArena*) arena_; | |
| arena_ = nullptr; | |
| } | |
| base_ = nullptr; | |
| blobs_ = 0; | |
| n_expert_ = 0; | |
| } | |
| bool ArenaExpertSource::pinned(int64_t layer, int64_t expert) const { | |
| if (base_ == nullptr || layer < 0 || expert < 0 || expert >= n_expert_) return false; | |
| const auto& lay = strata::kernels::cpu::expert_layout(); | |
| return lay.blob_offset(layer, expert) + lay.blob_bytes(layer) <= pinned_bytes_; | |
| } | |
| const uint8_t* ArenaExpertSource::device_alias(int64_t layer, int64_t expert) const { | |
| if (dev_slice_.empty() || !pinned(layer, expert)) return nullptr; | |
| const auto& lay = strata::kernels::cpu::expert_layout(); | |
| if (slice_bytes_ == 0) return dev_slice_[0] + lay.blob_offset(layer, expert); | |
| // one registration slice per layer | |
| if ((size_t) layer >= dev_slice_.size()) return nullptr; | |
| return dev_slice_[(size_t) layer] + (uint64_t) expert * lay.blob_bytes(layer); | |
| } | |
| const uint8_t* ArenaExpertSource::blob(int64_t layer, int64_t expert) { | |
| if (base_ == nullptr) return nullptr; | |
| if (layer < 0 || expert < 0 || expert >= n_expert_) return nullptr; | |
| const int64_t idx = layer * n_expert_ + expert; | |
| if (idx < 0 || idx >= blobs_) return nullptr; | |
| ++reads_; | |
| // Pointer arithmetic into resident memory. No fault, no copy, no mapping - which is the entire point of | |
| // this class over `FileExpertSource`. | |
| return base_ + strata::kernels::cpu::expert_layout().blob_offset(layer, expert); | |
| } | |
| } // namespace strata::core | |