diff --git a/docs/longhaul.md b/docs/longhaul.md index 021b93c..534c97c 100644 --- a/docs/longhaul.md +++ b/docs/longhaul.md @@ -16,6 +16,41 @@ llama-cli \ `--longhaul-cache` is the expert cache budget in GiB. It is required when `--longhaul` is used. The same options are accepted by `llama-server`. +## RPC mode + +RPC longhaul keeps the paging loop on the accelerator host so expert payloads do +not cross the network during generation. Build both hosts with `-DGGML_RPC=ON`. +Place the same GGUF shard files on both hosts, preserving their basenames, then +start the Metal host with the directory containing those shards: + +```sh +ggml-rpc-server \ + --device MTL0 \ + --longhaul-root /path/to/model-directory \ + --host 192.168.1.20 +``` + +Start the client normally: + +```sh +llama-server \ + --model /local/path/model.gguf \ + --rpc 192.168.1.20:50052 \ + --longhaul \ + --longhaul-cache 2 \ + --n-gpu-layers 99 +``` + +RPC longhaul currently requires all repeating layers on one remote Metal device. +The server validates each shard basename, size, GGUF metadata layout digest, and +every registered expert range before decoding. An older server, a server without +`--longhaul-root`, or a shard mismatch is a fatal model-load error; the client +does not fall back to sending expert slices over RPC. + +`--longhaul-root` grants connected RPC clients read access to registered byte +ranges in files in that directory. The RPC server remains experimental and +insecure and must not be exposed to an untrusted network. + Longhaul may reduce `--ubatch-size` so that every expert selected by one graph segment can be present in the cache at the same time. The effective value is logged during context creation. If one token selects more experts than the cache has slots, the routed MoE computation is split into multiple stages and the partial results are summed. This permits smaller caches at the cost of additional graph work. The normal startup warmup is skipped automatically in longhaul mode. Routed expert weights are not read until the first real decode. @@ -24,9 +59,9 @@ The normal startup warmup is skipped automatically in longhaul mode. Routed expe Longhaul currently requires: -- macOS with the Metal backend +- the Metal backend, either local on macOS or on one RPC host - Qwen3.5 MoE or Laguna architecture -- all repeating model layers assigned to Metal +- all repeating model layers assigned to local Metal or one RPC Metal device - text generation without embeddings or LoRA adapters Longhaul does not restrict the GGUF quantization type. Individual tensor types must still be supported by Metal. diff --git a/ggml/include/ggml-rpc.h b/ggml/include/ggml-rpc.h index 16ca339..49daeda 100644 --- a/ggml/include/ggml-rpc.h +++ b/ggml/include/ggml-rpc.h @@ -8,7 +8,7 @@ extern "C" { #define RPC_PROTO_MAJOR_VERSION 4 #define RPC_PROTO_MINOR_VERSION 0 -#define RPC_PROTO_PATCH_VERSION 3 +#define RPC_PROTO_PATCH_VERSION 4 #ifdef __cplusplus static_assert(GGML_OP_COUNT == 101, "GGML_OP_COUNT has changed - update RPC_PROTO_PATCH_VERSION"); @@ -27,6 +27,49 @@ GGML_BACKEND_API void ggml_backend_rpc_get_device_memory(const char * endpoint, GGML_BACKEND_API void ggml_backend_rpc_start_server(const char * endpoint, const char * cache_dir, size_t n_threads, size_t n_devices, ggml_backend_dev_t * devices); +struct ggml_backend_rpc_longhaul_shard { + const char * name; + uint64_t size; + uint64_t metadata_size; + uint8_t metadata_sha256[32]; +}; + +struct ggml_backend_rpc_longhaul_source { + struct ggml_tensor * tensor; + uint32_t shard; + uint64_t offset; + uint64_t expert_size; + int32_t layer; +}; + +struct ggml_backend_rpc_longhaul_params { + uint32_t n_layers; + uint32_t n_experts; + uint32_t n_slots; + size_t n_shards; + const struct ggml_backend_rpc_longhaul_shard * shards; + size_t n_sources; + const struct ggml_backend_rpc_longhaul_source * sources; +}; + +GGML_BACKEND_API bool ggml_backend_rpc_longhaul_register( + const struct ggml_backend_rpc_longhaul_params * params, + uint64_t * registration_id, + char * error, + size_t error_size); + +GGML_BACKEND_API void ggml_backend_rpc_longhaul_unregister( + struct ggml_tensor * tensor, + uint64_t registration_id); + +GGML_BACKEND_API void ggml_backend_rpc_start_server_with_options( + const char * endpoint, + const char * cache_dir, + const char * longhaul_root, + size_t n_threads, + size_t n_devices, + ggml_backend_dev_t * devices); + GGML_BACKEND_API ggml_backend_reg_t ggml_backend_rpc_reg(void); GGML_BACKEND_API ggml_backend_reg_t ggml_backend_rpc_add_server(const char * endpoint); diff --git a/ggml/src/ggml-rpc/ggml-rpc.cpp b/ggml/src/ggml-rpc/ggml-rpc.cpp index d380577..8bc567d 100644 --- a/ggml/src/ggml-rpc/ggml-rpc.cpp +++ b/ggml/src/ggml-rpc/ggml-rpc.cpp @@ -5,12 +5,18 @@ #include "transport.h" #include +#include +#include #include +#include +#include #include +#include #include #include #include #include +#include #include #include #include @@ -18,6 +24,11 @@ #include #include +#if !defined(_WIN32) +#include +#include +#endif + static const char * RPC_DEBUG = std::getenv("GGML_RPC_DEBUG"); #define LOG_DBG(...) \ @@ -71,6 +82,10 @@ enum rpc_cmd { RPC_CMD_HELLO, RPC_CMD_DEVICE_COUNT, RPC_CMD_GRAPH_RECOMPUTE, + RPC_CMD_LONGHAUL_REGISTER, + RPC_CMD_LONGHAUL_UNREGISTER, + RPC_CMD_LONGHAUL_GRAPH_COMPUTE, + RPC_CMD_LONGHAUL_GRAPH_RECOMPUTE, RPC_CMD_COUNT, }; @@ -190,6 +205,44 @@ struct rpc_msg_graph_recompute_req { uint32_t device; }; +struct rpc_msg_longhaul_register_header { + uint32_t n_layers; + uint32_t n_experts; + uint32_t n_slots; + uint32_t n_shards; + uint32_t n_sources; +}; + +struct rpc_msg_longhaul_shard { + uint32_t name_size; + uint64_t size; + uint64_t metadata_size; + uint8_t metadata_sha256[32]; +}; + +struct rpc_msg_longhaul_source { + rpc_tensor tensor; + uint32_t shard; + uint64_t offset; + uint64_t expert_size; + int32_t layer; +}; + +struct rpc_msg_longhaul_register_rsp { + uint8_t result; + uint64_t registration_id; + char error[255]; +}; + +struct rpc_msg_longhaul_unregister_req { + uint64_t registration_id; +}; + +struct rpc_msg_longhaul_compute_rsp { + uint8_t result; + char error[255]; +}; + #pragma pack(pop) // RPC data structures @@ -225,8 +278,19 @@ struct ggml_backend_rpc_buffer_context { std::shared_ptr sock; void * base_ptr; uint64_t remote_ptr; + uint8_t peer_patch; }; +static std::mutex g_peer_patch_mutex; +static std::unordered_map g_peer_patches; +static std::atomic g_longhaul_registration_count {0}; + +static uint8_t get_peer_patch(const std::shared_ptr & sock) { + std::lock_guard lock(g_peer_patch_mutex); + const auto it = g_peer_patches.find(sock.get()); + return it == g_peer_patches.end() ? 0 : it->second; +} + // RPC helper functions // Computes FNV-1a hash of the data @@ -342,6 +406,10 @@ static bool negotiate_hello(const std::shared_ptr & sock) { return false; } + { + std::lock_guard lock(g_peer_patch_mutex); + g_peer_patches[sock.get()] = response.patch; + } sock->update_caps(response.conn_caps); return true; } @@ -556,7 +624,7 @@ static ggml_backend_buffer_t ggml_backend_rpc_buffer_type_alloc_buffer(ggml_back if (response.remote_ptr != 0) { ggml_backend_buffer_t buffer = ggml_backend_buffer_init(buft, ggml_backend_rpc_buffer_interface, - new ggml_backend_rpc_buffer_context{sock, nullptr, response.remote_ptr}, + new ggml_backend_rpc_buffer_context{sock, nullptr, response.remote_ptr, get_peer_patch(sock)}, response.remote_size); return buffer; } else { @@ -696,6 +764,127 @@ static void serialize_graph(uint32_t device, const ggml_cgraph * cgraph, std::ve memcpy(out_tensors, tensors.data(), n_tensors * sizeof(rpc_tensor)); } +static bool graph_has_longhaul_marker(const ggml_cgraph * cgraph) { + for (int i = 0; i < cgraph->n_nodes; ++i) { + if (std::strncmp(cgraph->nodes[i]->name, "longhaul.", 9) == 0) { + return true; + } + } + return false; +} + +static void set_longhaul_error(char * error, size_t error_size, const char * message) { + if (error && error_size > 0) { + snprintf(error, error_size, "%s", message ? message : "unknown RPC longhaul error"); + } +} + +bool ggml_backend_rpc_longhaul_register( + const ggml_backend_rpc_longhaul_params * params, + uint64_t * registration_id, + char * error, + size_t error_size) { + if (!params || !params->shards || !params->sources || + params->n_layers == 0 || params->n_experts == 0 || params->n_slots == 0 || + params->n_shards == 0 || params->n_sources == 0 || + params->n_shards > UINT32_MAX || params->n_sources > UINT32_MAX) { + set_longhaul_error(error, error_size, "invalid RPC longhaul registration"); + return false; + } + + std::shared_ptr sock; + for (size_t i = 0; i < params->n_sources; ++i) { + ggml_tensor * tensor = params->sources[i].tensor; + if (!tensor || !tensor->buffer || !ggml_backend_buffer_is_rpc(tensor->buffer)) { + set_longhaul_error(error, error_size, "longhaul source tensor is not on an RPC buffer"); + return false; + } + auto * ctx = (ggml_backend_rpc_buffer_context *) tensor->buffer->context; + if (!ctx || (sock && sock != ctx->sock)) { + set_longhaul_error(error, error_size, "longhaul sources span multiple RPC endpoints"); + return false; + } + if (ctx->peer_patch < 4) { + set_longhaul_error(error, error_size, "RPC server does not support server-resident longhaul"); + return false; + } + sock = ctx->sock; + } + + rpc_msg_longhaul_register_header header = { + params->n_layers, + params->n_experts, + params->n_slots, + (uint32_t) params->n_shards, + (uint32_t) params->n_sources, + }; + std::vector input; + input.reserve(sizeof(header) + + params->n_shards * sizeof(rpc_msg_longhaul_shard) + + params->n_sources * sizeof(rpc_msg_longhaul_source) + 256 * params->n_shards); + auto append = [&](const void * data, size_t size) { + const size_t old_size = input.size(); + input.resize(old_size + size); + memcpy(input.data() + old_size, data, size); + }; + append(&header, sizeof(header)); + for (size_t i = 0; i < params->n_shards; ++i) { + const auto & shard = params->shards[i]; + if (!shard.name || shard.name[0] == '\0' || strlen(shard.name) > UINT32_MAX) { + set_longhaul_error(error, error_size, "invalid RPC longhaul shard name"); + return false; + } + rpc_msg_longhaul_shard wire = { + (uint32_t) strlen(shard.name), + shard.size, + shard.metadata_size, + {}, + }; + memcpy(wire.metadata_sha256, shard.metadata_sha256, sizeof(wire.metadata_sha256)); + append(&wire, sizeof(wire)); + append(shard.name, wire.name_size); + } + for (size_t i = 0; i < params->n_sources; ++i) { + const auto & source = params->sources[i]; + rpc_msg_longhaul_source wire = { + serialize_tensor(source.tensor), + source.shard, + source.offset, + source.expert_size, + source.layer, + }; + append(&wire, sizeof(wire)); + } + + rpc_msg_longhaul_register_rsp response = {}; + if (!send_rpc_cmd(sock, RPC_CMD_LONGHAUL_REGISTER, input.data(), input.size(), &response, sizeof(response))) { + set_longhaul_error(error, error_size, "RPC longhaul registration transport failure"); + return false; + } + if (!response.result) { + set_longhaul_error(error, error_size, response.error); + return false; + } + if (registration_id) { + *registration_id = response.registration_id; + } + g_longhaul_registration_count.fetch_add(1, std::memory_order_relaxed); + return true; +} + +void ggml_backend_rpc_longhaul_unregister(ggml_tensor * tensor, uint64_t registration_id) { + if (!tensor || !tensor->buffer || !ggml_backend_buffer_is_rpc(tensor->buffer) || registration_id == 0) { + return; + } + auto * ctx = (ggml_backend_rpc_buffer_context *) tensor->buffer->context; + rpc_msg_longhaul_unregister_req request = { registration_id }; + rpc_msg_longhaul_compute_rsp response = {}; + if (send_rpc_cmd(ctx->sock, RPC_CMD_LONGHAUL_UNREGISTER, + &request, sizeof(request), &response, sizeof(response)) && response.result) { + g_longhaul_registration_count.fetch_sub(1, std::memory_order_relaxed); + } +} + static enum ggml_status ggml_backend_rpc_graph_compute(ggml_backend_t backend, ggml_cgraph * cgraph) { ggml_backend_rpc_context * rpc_ctx = (ggml_backend_rpc_context *)backend->context; ggml_backend_dev_t rpc_dev = ggml_backend_get_device(backend); @@ -703,18 +892,43 @@ static enum ggml_status ggml_backend_rpc_graph_compute(ggml_backend_t backend, g GGML_ASSERT(cgraph->n_nodes > 0); bool reuse = cgraph->uid != 0 && rpc_dev_ctx->last_graph_uid == cgraph->uid; + const bool longhaul = + g_longhaul_registration_count.load(std::memory_order_relaxed) != 0 && + graph_has_longhaul_marker(cgraph); if (reuse) { rpc_msg_graph_recompute_req request; request.device = rpc_ctx->device; auto sock = get_socket(rpc_ctx->endpoint); - bool status = send_rpc_cmd(sock, RPC_CMD_GRAPH_RECOMPUTE, &request, sizeof(request)); + bool status; + if (longhaul) { + rpc_msg_longhaul_compute_rsp response = {}; + status = send_rpc_cmd(sock, RPC_CMD_LONGHAUL_GRAPH_RECOMPUTE, + &request, sizeof(request), &response, sizeof(response)); + if (status && !response.result) { + GGML_LOG_ERROR("RPC longhaul graph recompute failed: %s\n", response.error); + return GGML_STATUS_FAILED; + } + } else { + status = send_rpc_cmd(sock, RPC_CMD_GRAPH_RECOMPUTE, &request, sizeof(request)); + } RPC_STATUS_ASSERT(status); } else { rpc_dev_ctx->last_graph_uid = cgraph->uid; std::vector input; serialize_graph(rpc_ctx->device, cgraph, input); auto sock = get_socket(rpc_ctx->endpoint); - bool status = send_rpc_cmd(sock, RPC_CMD_GRAPH_COMPUTE, input.data(), input.size()); + bool status; + if (longhaul) { + rpc_msg_longhaul_compute_rsp response = {}; + status = send_rpc_cmd(sock, RPC_CMD_LONGHAUL_GRAPH_COMPUTE, + input.data(), input.size(), &response, sizeof(response)); + if (status && !response.result) { + GGML_LOG_ERROR("RPC longhaul graph compute failed: %s\n", response.error); + return GGML_STATUS_FAILED; + } + } else { + status = send_rpc_cmd(sock, RPC_CMD_GRAPH_COMPUTE, input.data(), input.size()); + } RPC_STATUS_ASSERT(status); } return GGML_STATUS_SUCCESS; @@ -816,10 +1030,442 @@ void ggml_backend_rpc_get_device_memory(const char * endpoint, uint32_t device, // RPC server-side implementation +namespace { + +struct rpc_sha256 { + uint32_t h[8] = { + 0x6a09e667, 0xbb67ae85, 0x3c6ef372, 0xa54ff53a, + 0x510e527f, 0x9b05688c, 0x1f83d9ab, 0x5be0cd19, + }; + uint8_t block[64] = {}; + uint64_t bits = 0; + size_t used = 0; + + static uint32_t rotr(uint32_t x, uint32_t n) { return (x >> n) | (x << (32 - n)); } + + void compress(const uint8_t * p) { + static const uint32_t k[64] = { + 0x428a2f98,0x71374491,0xb5c0fbcf,0xe9b5dba5,0x3956c25b,0x59f111f1,0x923f82a4,0xab1c5ed5, + 0xd807aa98,0x12835b01,0x243185be,0x550c7dc3,0x72be5d74,0x80deb1fe,0x9bdc06a7,0xc19bf174, + 0xe49b69c1,0xefbe4786,0x0fc19dc6,0x240ca1cc,0x2de92c6f,0x4a7484aa,0x5cb0a9dc,0x76f988da, + 0x983e5152,0xa831c66d,0xb00327c8,0xbf597fc7,0xc6e00bf3,0xd5a79147,0x06ca6351,0x14292967, + 0x27b70a85,0x2e1b2138,0x4d2c6dfc,0x53380d13,0x650a7354,0x766a0abb,0x81c2c92e,0x92722c85, + 0xa2bfe8a1,0xa81a664b,0xc24b8b70,0xc76c51a3,0xd192e819,0xd6990624,0xf40e3585,0x106aa070, + 0x19a4c116,0x1e376c08,0x2748774c,0x34b0bcb5,0x391c0cb3,0x4ed8aa4a,0x5b9cca4f,0x682e6ff3, + 0x748f82ee,0x78a5636f,0x84c87814,0x8cc70208,0x90befffa,0xa4506ceb,0xbef9a3f7,0xc67178f2, + }; + uint32_t w[64]; + for (int i = 0; i < 16; ++i) { + w[i] = uint32_t(p[4*i]) << 24 | uint32_t(p[4*i+1]) << 16 | + uint32_t(p[4*i+2]) << 8 | p[4*i+3]; + } + for (int i = 16; i < 64; ++i) { + const uint32_t s0 = rotr(w[i-15], 7) ^ rotr(w[i-15], 18) ^ (w[i-15] >> 3); + const uint32_t s1 = rotr(w[i-2], 17) ^ rotr(w[i-2], 19) ^ (w[i-2] >> 10); + w[i] = w[i-16] + s0 + w[i-7] + s1; + } + uint32_t a=h[0],b=h[1],c=h[2],d=h[3],e=h[4],f=h[5],g=h[6],hh=h[7]; + for (int i = 0; i < 64; ++i) { + const uint32_t s1 = rotr(e,6)^rotr(e,11)^rotr(e,25); + const uint32_t ch = (e&f)^((~e)&g); + const uint32_t t1 = hh+s1+ch+k[i]+w[i]; + const uint32_t s0 = rotr(a,2)^rotr(a,13)^rotr(a,22); + const uint32_t maj = (a&b)^(a&c)^(b&c); + const uint32_t t2 = s0+maj; + hh=g;g=f;f=e;e=d+t1;d=c;c=b;b=a;a=t1+t2; + } + h[0]+=a;h[1]+=b;h[2]+=c;h[3]+=d;h[4]+=e;h[5]+=f;h[6]+=g;h[7]+=hh; + } + + void update(const void * data, size_t n) { + const auto * p = static_cast(data); + bits += uint64_t(n) * 8; + while (n) { + const size_t take = std::min(n, sizeof(block) - used); + memcpy(block + used, p, take); + used += take; p += take; n -= take; + if (used == sizeof(block)) { compress(block); used = 0; } + } + } + + std::array finish() { + const uint64_t original_bits = bits; + const uint8_t one = 0x80; + update(&one, 1); + const uint8_t zero = 0; + while (used != 56) update(&zero, 1); + uint8_t len[8]; + for (int i = 0; i < 8; ++i) len[7-i] = uint8_t(original_bits >> (8*i)); + update(len, sizeof(len)); + std::array out; + for (int i = 0; i < 8; ++i) { + out[4*i]=uint8_t(h[i]>>24); out[4*i+1]=uint8_t(h[i]>>16); + out[4*i+2]=uint8_t(h[i]>>8); out[4*i+3]=uint8_t(h[i]); + } + return out; + } +}; + +static bool sha256_file_prefix(const fs::path & path, uint64_t size, std::array & digest) { + std::ifstream file(path, std::ios::binary); + if (!file) { + return false; + } + rpc_sha256 hash; + std::array buffer; + uint64_t remaining = size; + while (remaining > 0) { + const size_t count = (size_t) std::min(remaining, buffer.size()); + file.read((char *) buffer.data(), count); + if ((size_t) file.gcount() != count) { + return false; + } + hash.update(buffer.data(), count); + remaining -= count; + } + digest = hash.finish(); + return true; +} + +struct rpc_longhaul_source { + ggml_tensor * tensor; + uint32_t shard; + uint64_t offset; + uint64_t expert_size; + int32_t layer; +}; + +class rpc_longhaul_pager { +public: + rpc_longhaul_pager( + uint64_t id, + uint32_t n_layers, + uint32_t n_experts, + uint32_t n_slots, + std::vector shard_paths, + std::vector tensor_memory, + ggml_context_ptr tensor_ctx, + std::vector sources) : + id(id), + n_slots(n_slots), + n_experts(n_experts), + shard_paths(std::move(shard_paths)), + tensor_memory(std::move(tensor_memory)), + tensor_ctx(std::move(tensor_ctx)), + sources(std::move(sources)), + sources_by_layer(n_layers), + layers(n_layers) { + for (const auto & source : this->sources) { + sources_by_layer.at(source.layer).push_back(&source); + } + for (auto & layer : layers) { + layer.expert_ids.assign(n_slots, -1); + layer.expert_slots.assign(n_experts, -1); + layer.last_used.assign(n_slots, 0); + } + requested.resize(n_experts); + const unsigned int n_threads = std::max(1u, std::min(4u, std::thread::hardware_concurrency())); + for (unsigned int i = 0; i < n_threads; ++i) { + workers.emplace_back(&rpc_longhaul_pager::io_worker, this); + } + } + + ~rpc_longhaul_pager() { + { + std::lock_guard lock(io_mutex); + stopping = true; + } + io_ready.notify_all(); + for (auto & worker : workers) { + worker.join(); + } + GGML_LOG_INFO( + "rpc longhaul: registration=%" PRIu64 ", hits=%" PRIu64 ", misses=%" PRIu64 + ", read=%.2f MiB, io=%.2f ms\n", + id, n_hits, n_misses, bytes_read / 1024.0 / 1024.0, io_wall_us / 1000.0); + } + + bool remap(int layer, ggml_tensor * ids) { + last_error.clear(); + if (layer < 0 || layer >= (int) sources_by_layer.size() || sources_by_layer[layer].empty()) { + last_error = "RPC longhaul graph requested an unregistered expert layer"; + return false; + } + mutex.lock(); + locked = true; + locked_layer = layer; + const size_t n_values = ggml_nbytes(ids) / sizeof(int32_t); + id_buffer.resize(n_values); + ggml_backend_tensor_get(ids, id_buffer.data(), 0, ggml_nbytes(ids)); + + std::fill(requested.begin(), requested.end(), 0); + requested_experts.clear(); + missing_experts.clear(); + available_slots.clear(); + load_plan.clear(); + auto & state = layers.at(layer); + for (const int32_t value : id_buffer) { + if (value < 0 || value >= (int32_t) n_experts) { + last_error = "router produced an out-of-range expert ID"; + release(-1); + return false; + } + if (!requested[value]) { + requested[value] = 1; + requested_experts.push_back(value); + if (state.expert_slots[value] < 0) { + missing_experts.push_back(value); + } + } + } + if (requested_experts.size() > n_slots) { + last_error = "router requested more experts than the RPC longhaul cache can hold"; + release(-1); + return false; + } + for (uint32_t slot = 0; slot < n_slots; ++slot) { + const int32_t expert = state.expert_ids[slot]; + if (expert < 0 || !requested[expert]) { + available_slots.push_back(slot); + } + } + std::stable_sort(available_slots.begin(), available_slots.end(), [&](int32_t a, int32_t b) { + const bool a_empty = state.expert_ids[a] < 0; + const bool b_empty = state.expert_ids[b] < 0; + return a_empty != b_empty ? a_empty : state.last_used[a] < state.last_used[b]; + }); + if (missing_experts.size() > available_slots.size()) { + last_error = "RPC longhaul could not reserve enough cache slots"; + release(-1); + return false; + } + for (size_t i = 0; i < missing_experts.size(); ++i) { + load_plan.emplace_back(missing_experts[i], available_slots[i]); + } + + const int64_t io_start = ggml_time_us(); + if (!load_sources()) { + invalidate_plan(layer); + release(-1); + return false; + } + io_wall_us += ggml_time_us() - io_start; + for (const auto & item : load_plan) { + const int32_t old = state.expert_ids[item.second]; + if (old >= 0) state.expert_slots[old] = -1; + state.expert_ids[item.second] = item.first; + state.expert_slots[item.first] = item.second; + } + for (int32_t & value : id_buffer) { + const int32_t slot = state.expert_slots[value]; + GGML_ASSERT(slot >= 0); + state.last_used[slot] = ++tick; + value = slot; + } + ggml_backend_tensor_set(ids, id_buffer.data(), 0, ggml_nbytes(ids)); + n_hits += n_values - missing_experts.size(); + n_misses += missing_experts.size(); + return true; + } + + void release(int layer) { + if (locked && (layer < 0 || layer == locked_layer)) { + locked = false; + locked_layer = -1; + mutex.unlock(); + } + } + + const std::string & error() const { return last_error; } + uint64_t registration_id() const { return id; } + +private: + struct layer_state { + std::vector expert_ids; + std::vector expert_slots; + std::vector last_used; + }; + struct io_job { + const rpc_longhaul_source * source; + int32_t expert; + int32_t slot; + bool ok = false; + std::string error; + }; + + static bool is_direct(const rpc_longhaul_source & source) { + if (ggml_backend_buffer_is_host(source.tensor->buffer)) return true; + const char * name = ggml_backend_buft_name(ggml_backend_buffer_get_type(source.tensor->buffer)); + return name && std::strncmp(name, "MTL", 3) == 0 && std::strstr(name, "_Private") == nullptr; + } + + void read_source(const rpc_longhaul_source & source, int32_t expert, void * output) { +#if !defined(_WIN32) + const std::string path = shard_paths.at(source.shard).string(); + const int fd = open(path.c_str(), O_RDONLY); + if (fd < 0) throw std::runtime_error("failed to open RPC longhaul shard"); + struct fd_guard { + int fd; + ~fd_guard() { close(fd); } + } guard { fd }; +#if defined(__APPLE__) && defined(F_NOCACHE) + fcntl(fd, F_NOCACHE, 1); +#endif + uint64_t offset = source.offset + uint64_t(expert) * source.expert_size; + size_t remaining = source.expert_size; + auto * dest = static_cast(output); + while (remaining > 0) { + const ssize_t count = pread(fd, dest, remaining, (off_t) offset); + if (count < 0 && errno == EINTR) continue; + if (count <= 0) throw std::runtime_error("short read from RPC longhaul shard"); + dest += count; + offset += count; + remaining -= count; + } +#else + std::ifstream file(shard_paths.at(source.shard), std::ios::binary); + if (!file) throw std::runtime_error("failed to open RPC longhaul shard"); + const uint64_t offset = source.offset + uint64_t(expert) * source.expert_size; + file.seekg((std::streamoff) offset); + file.read((char *) output, (std::streamsize) source.expert_size); + if ((uint64_t) file.gcount() != source.expert_size) { + throw std::runtime_error("short read from RPC longhaul shard"); + } +#endif + } + + void io_worker() { + while (true) { + io_job * job = nullptr; + { + std::unique_lock lock(io_mutex); + io_ready.wait(lock, [&] { return stopping || !io_queue.empty(); }); + if (stopping && io_queue.empty()) return; + job = io_queue.front(); + io_queue.pop_front(); + } + try { + void * output = (uint8_t *) job->source->tensor->data + + size_t(job->slot) * job->source->tensor->nb[2]; + read_source(*job->source, job->expert, output); + job->ok = true; + } catch (const std::exception & e) { + job->error = e.what(); + } + { + std::lock_guard lock(io_mutex); + --io_pending; + } + io_done.notify_one(); + } + } + + bool load_sources() { + io_jobs.clear(); + for (const auto & item : load_plan) { + for (const auto * source : sources_by_layer.at(locked_layer)) { + if (is_direct(*source)) { + io_jobs.push_back({source, item.first, item.second, false, {}}); + } + } + } + if (!io_jobs.empty()) { + std::lock_guard lock(io_mutex); + io_pending = io_jobs.size(); + for (auto & job : io_jobs) io_queue.push_back(&job); + io_ready.notify_all(); + } + + bool staged_ok = true; + for (const auto & item : load_plan) { + for (const auto * source : sources_by_layer.at(locked_layer)) { + if (is_direct(*source)) continue; + staged_buffer.resize(source->expert_size); + try { + read_source(*source, item.first, staged_buffer.data()); + ggml_backend_tensor_set(source->tensor, staged_buffer.data(), + size_t(item.second) * source->tensor->nb[2], staged_buffer.size()); + bytes_read += staged_buffer.size(); + } catch (const std::exception & e) { + last_error = e.what(); + staged_ok = false; + break; + } + } + if (!staged_ok) break; + } + if (!io_jobs.empty()) { + std::unique_lock lock(io_mutex); + io_done.wait(lock, [&] { return io_pending == 0; }); + } + for (const auto & job : io_jobs) { + if (!job.ok) { + last_error = job.error; + return false; + } + bytes_read += job.source->expert_size; + } + return staged_ok; + } + + void invalidate_plan(int layer) { + auto & state = layers.at(layer); + for (const auto & item : load_plan) { + const int32_t old = state.expert_ids[item.second]; + if (old >= 0) state.expert_slots[old] = -1; + state.expert_ids[item.second] = -1; + state.last_used[item.second] = 0; + } + } + + uint64_t id; + uint32_t n_slots; + uint32_t n_experts; + std::vector shard_paths; + std::vector tensor_memory; + ggml_context_ptr tensor_ctx; + std::vector sources; + std::vector> sources_by_layer; + std::vector layers; + std::vector requested; + std::vector requested_experts; + std::vector missing_experts; + std::vector available_slots; + std::vector> load_plan; + std::vector id_buffer; + std::vector staged_buffer; + uint64_t tick = 0; + uint64_t n_hits = 0; + uint64_t n_misses = 0; + uint64_t bytes_read = 0; + uint64_t io_wall_us = 0; + std::mutex mutex; + bool locked = false; + int locked_layer = -1; + std::string last_error; + std::vector workers; + std::vector io_jobs; + std::deque io_queue; + std::mutex io_mutex; + std::condition_variable io_ready; + std::condition_variable io_done; + size_t io_pending = 0; + bool stopping = false; +}; + +struct rpc_longhaul_marker { + int index; + int layer; + bool release; +}; + +} // namespace + class rpc_server { public: - rpc_server(std::vector all_backends, const char * cache_dir) - : backends(std::move(all_backends)), cache_dir(cache_dir) { + rpc_server(std::vector all_backends, const char * cache_dir, const char * longhaul_root) + : backends(std::move(all_backends)), cache_dir(cache_dir), + longhaul_root(longhaul_root ? longhaul_root : "") { stored_graphs.resize(backends.size()); } ~rpc_server(); @@ -837,6 +1483,10 @@ public: bool copy_tensor(const rpc_msg_copy_tensor_req & request, rpc_msg_copy_tensor_rsp & response); bool graph_compute(const std::vector & input); bool graph_recompute(const rpc_msg_graph_recompute_req & request); + bool longhaul_register(const std::vector & input, rpc_msg_longhaul_register_rsp & response); + bool longhaul_unregister(const rpc_msg_longhaul_unregister_req & request, rpc_msg_longhaul_compute_rsp & response); + bool longhaul_graph_compute(const std::vector & input, rpc_msg_longhaul_compute_rsp & response); + bool longhaul_graph_recompute(const rpc_msg_graph_recompute_req & request, rpc_msg_longhaul_compute_rsp & response); bool init_tensor(const rpc_msg_init_tensor_req & request); bool get_alloc_size(const rpc_msg_get_alloc_size_req & request, rpc_msg_get_alloc_size_rsp & response); bool get_device_memory(const rpc_msg_get_device_memory_req & request, rpc_msg_get_device_memory_rsp & response); @@ -844,10 +1494,14 @@ public: struct stored_graph { std::vector buffer; ggml_cgraph * graph; + rpc_longhaul_pager * pager = nullptr; + std::vector longhaul_markers; }; private: bool get_cached_file(uint64_t hash, std::vector & data); + bool graph_compute_impl(const std::vector & input, bool longhaul, std::string & error); + bool execute_longhaul_graph(uint32_t device, stored_graph & stored, std::string & error); ggml_tensor * deserialize_tensor(struct ggml_context * ctx, const rpc_tensor * tensor); ggml_tensor * create_node(uint64_t id, struct ggml_context * ctx, @@ -857,9 +1511,13 @@ private: std::vector backends; const char * cache_dir; + std::string longhaul_root; std::unordered_set buffers; // store the last computed graph for each backend std::vector stored_graphs; + std::vector> longhaul_pagers; + std::unordered_map longhaul_by_tensor; + uint64_t next_longhaul_id = 1; }; void rpc_server::hello(rpc_msg_hello_rsp & response) { @@ -1188,6 +1846,158 @@ bool rpc_server::init_tensor(const rpc_msg_init_tensor_req & request) { return true; } +bool rpc_server::longhaul_register( + const std::vector & input, + rpc_msg_longhaul_register_rsp & response) { + auto fail = [&](const std::string & message) { + response.result = 0; + snprintf(response.error, sizeof(response.error), "%s", message.c_str()); + return true; + }; + if (longhaul_root.empty()) { + return fail("RPC server was not started with --longhaul-root"); + } + if (input.size() < sizeof(rpc_msg_longhaul_register_header)) { + return false; + } + size_t cursor = 0; + auto take = [&](void * output, size_t size) { + if (size > input.size() - cursor) return false; + memcpy(output, input.data() + cursor, size); + cursor += size; + return true; + }; + + rpc_msg_longhaul_register_header header; + if (!take(&header, sizeof(header)) || header.n_layers == 0 || header.n_experts == 0 || + header.n_slots == 0 || header.n_shards == 0 || header.n_sources == 0) { + return fail("invalid RPC longhaul registration header"); + } + + std::vector shard_paths; + shard_paths.reserve(header.n_shards); + for (uint32_t i = 0; i < header.n_shards; ++i) { + rpc_msg_longhaul_shard shard; + if (!take(&shard, sizeof(shard)) || shard.name_size == 0 || + shard.name_size > input.size() - cursor || shard.metadata_size > shard.size) { + return fail("invalid RPC longhaul shard descriptor"); + } + std::string name((const char *) input.data() + cursor, shard.name_size); + cursor += shard.name_size; + const fs::path relative(name); + if (relative.is_absolute() || relative.has_parent_path() || relative.filename().string() != name) { + return fail("RPC longhaul shard name must be a basename"); + } + std::error_code ec; + fs::path path = fs::canonical(fs::path(longhaul_root) / relative, ec); + if (ec || path.parent_path() != fs::path(longhaul_root) || !fs::is_regular_file(path, ec)) { + return fail("RPC longhaul shard is outside the configured root or is not a file: " + name); + } + const uint64_t actual_size = fs::file_size(path, ec); + if (ec || actual_size != shard.size) { + return fail("RPC longhaul shard size mismatch: " + name); + } + std::array digest; + if (!sha256_file_prefix(path, shard.metadata_size, digest) || + memcmp(digest.data(), shard.metadata_sha256, digest.size()) != 0) { + return fail("RPC longhaul shard metadata digest mismatch: " + name); + } + shard_paths.push_back(std::move(path)); + } + + if (header.n_sources > (input.size() - cursor) / sizeof(rpc_msg_longhaul_source) || + cursor + size_t(header.n_sources) * sizeof(rpc_msg_longhaul_source) != input.size()) { + return fail("invalid RPC longhaul source payload"); + } + std::vector wire_sources(header.n_sources); + if (!take(wire_sources.data(), wire_sources.size() * sizeof(wire_sources[0]))) { + return false; + } + + std::vector tensor_memory(ggml_tensor_overhead() * header.n_sources); + ggml_init_params tensor_params = { + /*.mem_size =*/ tensor_memory.size(), + /*.mem_buffer =*/ tensor_memory.data(), + /*.no_alloc =*/ true, + }; + ggml_context_ptr tensor_ctx { ggml_init(tensor_params) }; + if (!tensor_ctx) { + return fail("failed to allocate RPC longhaul tensor metadata"); + } + std::vector sources; + sources.reserve(header.n_sources); + for (const auto & wire : wire_sources) { + if (wire.shard >= shard_paths.size() || wire.layer < 0 || + wire.layer >= (int32_t) header.n_layers || wire.expert_size == 0) { + return fail("invalid RPC longhaul source descriptor"); + } + std::error_code ec; + const uint64_t shard_size = fs::file_size(shard_paths[wire.shard], ec); + if (ec || wire.expert_size > UINT64_MAX / header.n_experts || + wire.offset > shard_size || + wire.expert_size * header.n_experts > shard_size - wire.offset) { + return fail("RPC longhaul expert source is outside its shard"); + } + ggml_tensor * tensor = deserialize_tensor(tensor_ctx.get(), &wire.tensor); + if (!tensor || !tensor->buffer || tensor->ne[2] != (int64_t) header.n_slots || + tensor->nb[2] != wire.expert_size) { + return fail("RPC longhaul destination tensor does not match its source"); + } + const char * buft_name = ggml_backend_buft_name(ggml_backend_buffer_get_type(tensor->buffer)); + if (!buft_name || std::strncmp(buft_name, "MTL", 3) != 0) { + return fail("RPC longhaul destination is not a Metal buffer"); + } + if (longhaul_by_tensor.count(wire.tensor.data)) { + return fail("RPC longhaul tensor is already registered"); + } + sources.push_back({tensor, wire.shard, wire.offset, wire.expert_size, wire.layer}); + } + + const uint64_t id = next_longhaul_id++; + auto pager = std::make_unique( + id, header.n_layers, header.n_experts, header.n_slots, + std::move(shard_paths), std::move(tensor_memory), std::move(tensor_ctx), std::move(sources)); + for (const auto & wire : wire_sources) { + longhaul_by_tensor.emplace(wire.tensor.data, pager.get()); + } + longhaul_pagers.push_back(std::move(pager)); + response.result = 1; + response.registration_id = id; + response.error[0] = '\0'; + GGML_LOG_INFO("registered RPC longhaul pager %" PRIu64 ": layers=%u experts=%u slots=%u sources=%u\n", + id, header.n_layers, header.n_experts, header.n_slots, header.n_sources); + return true; +} + +bool rpc_server::longhaul_unregister( + const rpc_msg_longhaul_unregister_req & request, + rpc_msg_longhaul_compute_rsp & response) { + auto it = std::find_if(longhaul_pagers.begin(), longhaul_pagers.end(), [&](const auto & pager) { + return pager->registration_id() == request.registration_id; + }); + if (it == longhaul_pagers.end()) { + response.result = 0; + snprintf(response.error, sizeof(response.error), "%s", "RPC longhaul registration not found"); + return true; + } + rpc_longhaul_pager * pager = it->get(); + pager->release(-1); + for (auto map_it = longhaul_by_tensor.begin(); map_it != longhaul_by_tensor.end();) { + if (map_it->second == pager) map_it = longhaul_by_tensor.erase(map_it); + else ++map_it; + } + for (auto & stored : stored_graphs) { + if (stored.pager == pager) { + stored.pager = nullptr; + stored.longhaul_markers.clear(); + } + } + longhaul_pagers.erase(it); + response.result = 1; + response.error[0] = '\0'; + return true; +} + bool rpc_server::get_tensor(const rpc_msg_get_tensor_req & request, std::vector & response) { struct ggml_init_params params { /*.mem_size =*/ ggml_tensor_overhead(), @@ -1320,7 +2130,7 @@ ggml_tensor * rpc_server::create_node(uint64_t id, return result; } -bool rpc_server::graph_compute(const std::vector & input) { +bool rpc_server::graph_compute_impl(const std::vector & input, bool longhaul, std::string & error) { // serialization format: // | device (4 bytes) | n_nodes (4 bytes) | nodes (n_nodes * sizeof(uint64_t) | n_tensors (4 bytes) | tensors (n_tensors * sizeof(rpc_tensor)) | if (input.size() < 2*sizeof(uint32_t)) { @@ -1384,10 +2194,48 @@ bool rpc_server::graph_compute(const std::vector & input) { return false; } } - ggml_status status = ggml_backend_graph_compute(backends[device], graph); - GGML_ASSERT(status == GGML_STATUS_SUCCESS && "Unsuccessful graph computations are not supported with RPC"); - stored_graphs[device].graph = graph; - return true; + auto & stored = stored_graphs[device]; + stored.graph = graph; + stored.pager = nullptr; + stored.longhaul_markers.clear(); + + if (!longhaul) { + ggml_status status = ggml_backend_graph_compute(backends[device], graph); + GGML_ASSERT(status == GGML_STATUS_SUCCESS && "Unsuccessful graph computations are not supported with RPC"); + return true; + } + + for (uint32_t i = 0; i < n_tensors; ++i) { + auto it = longhaul_by_tensor.find(tensors[i].data); + if (it == longhaul_by_tensor.end()) continue; + if (stored.pager && stored.pager != it->second) { + error = "longhaul graph references more than one registered pager"; + return false; + } + stored.pager = it->second; + } + if (!stored.pager) { + error = "longhaul graph does not reference a registered expert tensor"; + return false; + } + for (int i = 0; i < graph->n_nodes; ++i) { + int layer = -1; + if (sscanf(graph->nodes[i]->name, "longhaul.remap.%d", &layer) == 1) { + stored.longhaul_markers.push_back({i, layer, false}); + } else if (sscanf(graph->nodes[i]->name, "longhaul.release.%d", &layer) == 1) { + stored.longhaul_markers.push_back({i, layer, true}); + } + } + if (stored.longhaul_markers.empty()) { + error = "longhaul graph contains no paging markers"; + return false; + } + return execute_longhaul_graph(device, stored, error); +} + +bool rpc_server::graph_compute(const std::vector & input) { + std::string error; + return graph_compute_impl(input, false, error); } bool rpc_server::graph_recompute(const rpc_msg_graph_recompute_req & request) { @@ -1405,6 +2253,70 @@ bool rpc_server::graph_recompute(const rpc_msg_graph_recompute_req & request) { return true; } +bool rpc_server::execute_longhaul_graph(uint32_t device, stored_graph & stored, std::string & error) { + int first = 0; + for (const auto & marker : stored.longhaul_markers) { + if (marker.index < first || marker.index >= stored.graph->n_nodes) { + error = "invalid stored longhaul graph marker"; + stored.pager->release(-1); + return false; + } + ggml_cgraph view = ggml_graph_view(stored.graph, first, marker.index + 1); + ggml_status status = ggml_backend_graph_compute(backends[device], &view); + if (status != GGML_STATUS_SUCCESS) { + error = "Metal graph segment failed"; + stored.pager->release(-1); + return false; + } + ggml_backend_synchronize(backends[device]); + if (marker.release) { + stored.pager->release(marker.layer); + } else if (!stored.pager->remap(marker.layer, stored.graph->nodes[marker.index])) { + error = stored.pager->error(); + stored.pager->release(-1); + return false; + } + first = marker.index + 1; + } + if (first < stored.graph->n_nodes) { + ggml_cgraph view = ggml_graph_view(stored.graph, first, stored.graph->n_nodes); + ggml_status status = ggml_backend_graph_compute(backends[device], &view); + if (status != GGML_STATUS_SUCCESS) { + error = "Metal graph tail failed"; + stored.pager->release(-1); + return false; + } + } + return true; +} + +bool rpc_server::longhaul_graph_compute( + const std::vector & input, + rpc_msg_longhaul_compute_rsp & response) { + std::string error; + response.result = graph_compute_impl(input, true, error); + snprintf(response.error, sizeof(response.error), "%s", error.c_str()); + return true; +} + +bool rpc_server::longhaul_graph_recompute( + const rpc_msg_graph_recompute_req & request, + rpc_msg_longhaul_compute_rsp & response) { + if (request.device >= stored_graphs.size()) { + return false; + } + auto & stored = stored_graphs[request.device]; + std::string error; + if (!stored.graph || !stored.pager || stored.longhaul_markers.empty()) { + error = "no stored RPC longhaul graph"; + response.result = 0; + } else { + response.result = execute_longhaul_graph(request.device, stored, error); + } + snprintf(response.error, sizeof(response.error), "%s", error.c_str()); + return true; +} + bool rpc_server::get_device_memory(const rpc_msg_get_device_memory_req & request, rpc_msg_get_device_memory_rsp & response) { uint32_t dev_id = request.device; if (dev_id >= backends.size()) { @@ -1425,9 +2337,12 @@ rpc_server::~rpc_server() { } } -static void rpc_serve_client(const std::vector & backends, const char * cache_dir, - socket_ptr sock) { - rpc_server server(backends, cache_dir); +static void rpc_serve_client( + const std::vector & backends, + const char * cache_dir, + const char * longhaul_root, + socket_ptr sock) { + rpc_server server(backends, cache_dir, longhaul_root); uint8_t cmd; if (!sock->recv_data(&cmd, 1)) { return; @@ -1670,6 +2585,54 @@ static void rpc_serve_client(const std::vector & backends, const } break; } + case RPC_CMD_LONGHAUL_REGISTER: { + std::vector input; + if (!recv_msg(sock, input)) { + return; + } + rpc_msg_longhaul_register_rsp response = {}; + if (!server.longhaul_register(input, response) || + !send_msg(sock, &response, sizeof(response))) { + return; + } + break; + } + case RPC_CMD_LONGHAUL_UNREGISTER: { + rpc_msg_longhaul_unregister_req request; + if (!recv_msg(sock, &request, sizeof(request))) { + return; + } + rpc_msg_longhaul_compute_rsp response = {}; + if (!server.longhaul_unregister(request, response) || + !send_msg(sock, &response, sizeof(response))) { + return; + } + break; + } + case RPC_CMD_LONGHAUL_GRAPH_COMPUTE: { + std::vector input; + if (!recv_msg(sock, input)) { + return; + } + rpc_msg_longhaul_compute_rsp response = {}; + if (!server.longhaul_graph_compute(input, response) || + !send_msg(sock, &response, sizeof(response))) { + return; + } + break; + } + case RPC_CMD_LONGHAUL_GRAPH_RECOMPUTE: { + rpc_msg_graph_recompute_req request; + if (!recv_msg(sock, &request, sizeof(request))) { + return; + } + rpc_msg_longhaul_compute_rsp response = {}; + if (!server.longhaul_graph_recompute(request, response) || + !send_msg(sock, &response, sizeof(response))) { + return; + } + break; + } case RPC_CMD_GET_DEVICE_MEMORY: { rpc_msg_get_device_memory_req request; if (!recv_msg(sock, &request, sizeof(request))) { @@ -1694,6 +2657,17 @@ static void rpc_serve_client(const std::vector & backends, const void ggml_backend_rpc_start_server(const char * endpoint, const char * cache_dir, size_t n_threads, size_t n_devices, ggml_backend_dev_t * devices) { + ggml_backend_rpc_start_server_with_options( + endpoint, cache_dir, nullptr, n_threads, n_devices, devices); +} + +void ggml_backend_rpc_start_server_with_options( + const char * endpoint, + const char * cache_dir, + const char * longhaul_root, + size_t n_threads, + size_t n_devices, + ggml_backend_dev_t * devices) { if (n_devices == 0 || devices == nullptr) { fprintf(stderr, "Invalid arguments to ggml_backend_rpc_start_server\n"); return; @@ -1705,6 +2679,7 @@ void ggml_backend_rpc_start_server(const char * endpoint, const char * cache_dir RPC_PROTO_PATCH_VERSION); printf(" endpoint : %s\n", endpoint); printf(" local cache : %s\n", cache_dir ? cache_dir : "n/a"); + printf(" longhaul root : %s\n", longhaul_root ? longhaul_root : "n/a"); printf("Devices:\n"); for (size_t i = 0; i < n_devices; i++) { auto dev = devices[i]; @@ -1755,7 +2730,7 @@ void ggml_backend_rpc_start_server(const char * endpoint, const char * cache_dir } printf("Accepted client connection\n"); fflush(stdout); - rpc_serve_client(backends, cache_dir, client_socket); + rpc_serve_client(backends, cache_dir, longhaul_root, client_socket); printf("Client connection closed\n"); fflush(stdout); } @@ -1887,6 +2862,15 @@ static void * ggml_backend_rpc_get_proc_address(ggml_backend_reg_t reg, const ch if (std::strcmp(name, "ggml_backend_rpc_start_server") == 0) { return (void *)ggml_backend_rpc_start_server; } + if (std::strcmp(name, "ggml_backend_rpc_start_server_with_options") == 0) { + return (void *)ggml_backend_rpc_start_server_with_options; + } + if (std::strcmp(name, "ggml_backend_rpc_longhaul_register") == 0) { + return (void *)ggml_backend_rpc_longhaul_register; + } + if (std::strcmp(name, "ggml_backend_rpc_longhaul_unregister") == 0) { + return (void *)ggml_backend_rpc_longhaul_unregister; + } return NULL; GGML_UNUSED(reg); diff --git a/src/llama-context.cpp b/src/llama-context.cpp index e78aa3f..6506b5f 100644 --- a/src/llama-context.cpp +++ b/src/llama-context.cpp @@ -1359,10 +1359,12 @@ llm_graph_result * llama_context::process_ubatch(const llama_ubatch & ubatch, ll res->reset(); ggml_backend_sched_reset(sched.get()); + auto * longhaul = model.longhaul_cache(); + const bool local_longhaul = longhaul && !longhaul->is_remote(); ggml_backend_sched_set_eval_callback( sched.get(), - model.longhaul_cache() ? longhaul_eval_callback : cparams.cb_eval, - model.longhaul_cache() ? this : cparams.cb_eval_user_data); + local_longhaul ? longhaul_eval_callback : cparams.cb_eval, + local_longhaul ? this : cparams.cb_eval_user_data); //const auto t_start_us = ggml_time_us(); @@ -2461,7 +2463,7 @@ llm_graph_params llama_context::graph_params( bool llama_context::longhaul_eval_callback(ggml_tensor * tensor, bool ask, void * user_data) { auto * ctx = static_cast(user_data); auto * cache = ctx->model.longhaul_cache(); - GGML_ASSERT(cache != nullptr); + GGML_ASSERT(cache != nullptr && !cache->is_remote()); int layer = -1; const bool remap = sscanf(tensor->name, "longhaul.remap.%d", &layer) == 1; diff --git a/src/llama-longhaul.cpp b/src/llama-longhaul.cpp index 1b9ea98..7f17843 100644 --- a/src/llama-longhaul.cpp +++ b/src/llama-longhaul.cpp @@ -1,23 +1,108 @@ #include "llama-longhaul.h" #include "llama-impl.h" +#include "llama-token-cache.h" + +#include "ggml-rpc.h" #include +#include #include #include +#include llama_longhaul_cache::llama_longhaul_cache( llama_files files, std::vector sources, + std::vector shards, size_t n_slots, uint32_t n_experts, - uint32_t n_layers) : + uint32_t n_layers, + bool remote) : files(std::move(files)), sources(std::move(sources)), + shards(std::move(shards)), sources_by_layer(n_layers), layers(n_layers), n_slots(n_slots), - n_experts(n_experts) { + n_experts(n_experts), + remote(remote) { + if (remote) { + if (this->shards.size() != this->files.size() || this->sources.empty()) { + throw std::runtime_error("RPC longhaul requires named GGUF shards"); + } + std::vector> digests(this->shards.size()); + std::vector rpc_shards(this->shards.size()); + for (size_t i = 0; i < this->shards.size(); ++i) { + const auto & shard = this->shards[i]; + if (shard.metadata_size > shard.size) { + throw std::runtime_error("invalid RPC longhaul shard metadata range"); + } + std::vector metadata(shard.metadata_size); + if (!metadata.empty()) { + this->files[i]->read_at(metadata.data(), metadata.size(), 0); + } + const std::string hex = llama_token_cache_hash(metadata.data(), metadata.size()); + if (hex.size() != 64) { + throw std::runtime_error("failed to fingerprint RPC longhaul shard"); + } + for (size_t j = 0; j < digests[i].size(); ++j) { + auto nibble = [](char c) -> uint8_t { + if (c >= '0' && c <= '9') return c - '0'; + if (c >= 'a' && c <= 'f') return c - 'a' + 10; + if (c >= 'A' && c <= 'F') return c - 'A' + 10; + return 0xff; + }; + const uint8_t hi = nibble(hex[2*j]); + const uint8_t lo = nibble(hex[2*j + 1]); + if (hi > 0x0f || lo > 0x0f) { + throw std::runtime_error("invalid RPC longhaul shard fingerprint"); + } + digests[i][j] = (hi << 4) | lo; + } + rpc_shards[i] = { + shard.name.c_str(), + shard.size, + shard.metadata_size, + {}, + }; + memcpy(rpc_shards[i].metadata_sha256, digests[i].data(), digests[i].size()); + } + + std::vector rpc_sources; + rpc_sources.reserve(this->sources.size()); + for (const auto & source : this->sources) { + rpc_sources.push_back({ + source.tensor, + source.file_idx, + source.offset, + source.expert_size, + source.layer, + }); + } + ggml_backend_rpc_longhaul_params params = { + n_layers, + n_experts, + (uint32_t) n_slots, + rpc_shards.size(), + rpc_shards.data(), + rpc_sources.size(), + rpc_sources.data(), + }; + ggml_backend_reg_t reg = ggml_backend_reg_by_name("RPC"); + auto register_fn = reg ? (decltype(ggml_backend_rpc_longhaul_register) *) + ggml_backend_reg_get_proc_address(reg, "ggml_backend_rpc_longhaul_register") : nullptr; + if (!register_fn) { + throw std::runtime_error("RPC backend does not expose longhaul registration"); + } + char error[256] = {}; + if (!register_fn(¶ms, &remote_registration_id, error, sizeof(error))) { + throw std::runtime_error(error[0] ? error : "RPC longhaul registration failed"); + } + this->files.clear(); + return; + } + for (auto & file : this->files) { file->set_no_cache(); } @@ -45,6 +130,17 @@ llama_longhaul_cache::llama_longhaul_cache( } llama_longhaul_cache::~llama_longhaul_cache() { + if (remote) { + if (!sources.empty() && remote_registration_id != 0) { + ggml_backend_reg_t reg = ggml_backend_reg_by_name("RPC"); + auto unregister_fn = reg ? (decltype(ggml_backend_rpc_longhaul_unregister) *) + ggml_backend_reg_get_proc_address(reg, "ggml_backend_rpc_longhaul_unregister") : nullptr; + if (unregister_fn) { + unregister_fn(sources.front().tensor, remote_registration_id); + } + } + return; + } { std::lock_guard lock(io_mutex); io_stopping = true; @@ -321,3 +417,7 @@ uint64_t llama_longhaul_cache::misses() const { uint64_t llama_longhaul_cache::bytes_read_count() const { return bytes_read; } + +bool llama_longhaul_cache::is_remote() const { + return remote; +} diff --git a/src/llama-longhaul.h b/src/llama-longhaul.h index cf33ec4..b4702df 100644 --- a/src/llama-longhaul.h +++ b/src/llama-longhaul.h @@ -17,9 +17,11 @@ struct llama_longhaul_cache { llama_longhaul_cache( llama_files files, std::vector sources, + std::vector shards, size_t n_slots, uint32_t n_experts, - uint32_t n_layers); + uint32_t n_layers, + bool remote); ~llama_longhaul_cache(); bool remap(int layer, ggml_tensor * ids); @@ -31,6 +33,7 @@ struct llama_longhaul_cache { bool failed() const; uint64_t misses() const; uint64_t bytes_read_count() const; + bool is_remote() const; private: struct layer_state { @@ -49,6 +52,7 @@ private: llama_files files; std::vector sources; + std::vector shards; std::vector> sources_by_layer; std::vector layers; size_t n_slots; @@ -84,6 +88,8 @@ private: std::condition_variable io_done; size_t io_pending = 0; bool io_stopping = false; + bool remote = false; + uint64_t remote_registration_id = 0; bool source_is_direct(const llama_model_loader::longhaul_source & source) const; bool load_plan_sources(); diff --git a/src/llama-model-loader.cpp b/src/llama-model-loader.cpp index ffbd649..c1fc160 100644 --- a/src/llama-model-loader.cpp +++ b/src/llama-model-loader.cpp @@ -594,6 +594,11 @@ llama_model_loader::llama_model_loader( files.emplace_back(new llama_file(fname.c_str(), "rb", use_direct_io)); contexts.emplace_back(ctx); + longhaul_shards.push_back({ + std::filesystem::path(fname).filename().string(), + files.back()->size(), + gguf_get_data_offset(metadata), + }); // Save tensors data offset of the main file. // For subsidiary files, `meta` tensor data offset must not be used, @@ -662,6 +667,11 @@ llama_model_loader::llama_model_loader( files.emplace_back(new llama_file(fname_split, "rb", use_direct_io)); contexts.emplace_back(ctx); + longhaul_shards.push_back({ + std::filesystem::path(fname_split).filename().string(), + files.back()->size(), + gguf_get_data_offset(ctx_gguf.get()), + }); // Save tensors data offset info of the shard. for (ggml_tensor * cur = ggml_get_first_tensor(ctx); cur; cur = ggml_get_next_tensor(ctx, cur)) { diff --git a/src/llama-model-loader.h b/src/llama-model-loader.h index 6224c31..38aa8bf 100644 --- a/src/llama-model-loader.h +++ b/src/llama-model-loader.h @@ -9,6 +9,7 @@ #include "ggml-cpp.h" +#include #include #include #include @@ -77,6 +78,12 @@ struct llama_model_loader { int layer; }; + struct longhaul_shard { + std::string name; + uint64_t size; + uint64_t metadata_size; + }; + int n_kv = 0; int n_tensors = 0; int n_created = 0; @@ -115,6 +122,7 @@ struct llama_model_loader { std::vector> mmaps_used; size_t longhaul_slots = 0; std::vector longhaul_sources; + std::vector longhaul_shards; // define a comparator for the buft -> ctx map to ensure that the order is well-defined: struct ggml_backend_buft_comparator { diff --git a/src/llama-model.cpp b/src/llama-model.cpp index 67ee21e..88f7352 100644 --- a/src/llama-model.cpp +++ b/src/llama-model.cpp @@ -1039,6 +1039,7 @@ struct llama_model::impl { std::vector tensor_split_owned; std::unique_ptr longhaul; + bool longhaul_remote = false; }; llama_model::llama_model(const llama_model_params & params) : params(params), pimpl(std::make_unique()) { @@ -1356,9 +1357,6 @@ bool llama_model_base::load_tensors(llama_model_loader & ml) { } if (params.load_mode == LLAMA_LOAD_MODE_LONGHAUL) { -#if !defined(__APPLE__) - throw std::runtime_error("longhaul is only supported on macOS"); -#endif if (arch != LLM_ARCH_QWEN35MOE && arch != LLM_ARCH_LAGUNA) { throw std::runtime_error("longhaul currently requires a qwen35moe or laguna model"); } @@ -1368,12 +1366,35 @@ bool llama_model_base::load_tensors(llama_model_loader & ml) { if (params.check_tensors || params.no_alloc || params.vocab_only) { throw std::runtime_error("longhaul does not support check-tensors, no-alloc, or vocab-only loading"); } + bool all_metal = true; + bool all_rpc = true; for (int il = 0; il < n_layer_all; ++il) { ggml_backend_dev_t dev = pimpl->dev_layer[il].dev; - if (ggml_backend_dev_type(dev) != GGML_BACKEND_DEVICE_TYPE_GPU || - strcmp(ggml_backend_reg_name(ggml_backend_dev_backend_reg(dev)), "MTL") != 0) { - throw std::runtime_error("longhaul requires every model layer on Metal"); + const char * reg_name = ggml_backend_reg_name(ggml_backend_dev_backend_reg(dev)); + all_metal = all_metal && + ggml_backend_dev_type(dev) == GGML_BACKEND_DEVICE_TYPE_GPU && + strcmp(reg_name, "MTL") == 0; + all_rpc = all_rpc && + ggml_backend_dev_type(dev) == GGML_BACKEND_DEVICE_TYPE_GPU && + strcmp(reg_name, "RPC") == 0; + } + if (!all_metal && !all_rpc) { + throw std::runtime_error( + "longhaul requires all model layers on local Metal or one RPC device"); + } + if (all_metal) { +#if !defined(__APPLE__) + throw std::runtime_error("local longhaul is only supported on macOS"); +#endif + } else { + ggml_backend_dev_t first = pimpl->dev_layer.front().dev; + for (int il = 1; il < n_layer_all; ++il) { + if (pimpl->dev_layer[il].dev != first) { + throw std::runtime_error( + "RPC longhaul v1 requires every model layer on one RPC device"); + } } + pimpl->longhaul_remote = true; } ml.configure_longhaul(params.longhaul_cache_bytes, n_expert, n_layer_all); if (ml.longhaul_slots < (size_t) n_expert_used) { @@ -1694,7 +1715,8 @@ bool llama_model_base::load_tensors(llama_model_loader & ml) { if (params.load_mode == LLAMA_LOAD_MODE_LONGHAUL) { pimpl->longhaul = std::make_unique( - std::move(ml.files), std::move(ml.longhaul_sources), ml.longhaul_slots, hparams.n_expert, hparams.n_layer_all); + std::move(ml.files), std::move(ml.longhaul_sources), std::move(ml.longhaul_shards), + ml.longhaul_slots, hparams.n_expert, hparams.n_layer_all, pimpl->longhaul_remote); } return true; diff --git a/src/llama.cpp b/src/llama.cpp index 8288b05..a6b4088 100644 --- a/src/llama.cpp +++ b/src/llama.cpp @@ -263,16 +263,21 @@ static bool llama_prepare_model_devices(const llama_model_params & params, llama } } - // add RPC servers at the front of the list to minimize network transfers - model->devices.insert(model->devices.begin(), rpc_servers.begin(), rpc_servers.end()); + if (params.load_mode == LLAMA_LOAD_MODE_LONGHAUL && !rpc_servers.empty()) { + // RPC longhaul v1 keeps the entire paging domain on one remote device. + model->devices.push_back(rpc_servers.front()); + } else { + // add RPC servers at the front of the list to minimize network transfers + model->devices.insert(model->devices.begin(), rpc_servers.begin(), rpc_servers.end()); - // add GPUs - model->devices.insert(model->devices.end(), gpus.begin(), gpus.end()); + // add GPUs + model->devices.insert(model->devices.end(), gpus.begin(), gpus.end()); - // add integrated GPUs only if no discrete GPUs were found - // (RPC servers do not count, otherwise the local iGPU would be dropped on iGPU+RPC setups) - if (gpus.empty()) { - model->devices.insert(model->devices.end(), igpus.begin(), igpus.end()); + // add integrated GPUs only if no discrete GPUs were found + // (RPC servers do not count, otherwise the local iGPU would be dropped on iGPU+RPC setups) + if (gpus.empty()) { + model->devices.insert(model->devices.end(), igpus.begin(), igpus.end()); + } } } diff --git a/tests/test-longhaul.cpp b/tests/test-longhaul.cpp index 39a2ebe..41dfa74 100644 --- a/tests/test-longhaul.cpp +++ b/tests/test-longhaul.cpp @@ -58,7 +58,8 @@ struct longhaul_fixture { sources.push_back({ weights_b, 0, n_experts * expert_size, expert_size, 0 }); } cache = std::make_unique( - std::move(files), std::move(sources), n_slots, n_experts, 1); + std::move(files), std::move(sources), std::vector{}, + n_slots, n_experts, 1, false); } ~longhaul_fixture() { diff --git a/tools/rpc/README.md b/tools/rpc/README.md index 655b653..d502595 100644 --- a/tools/rpc/README.md +++ b/tools/rpc/README.md @@ -36,6 +36,17 @@ flowchart TD By default, `ggml-rpc-server` exposes all available accelerator devices on the host. If there are no accelerators, it exposes a single `CPU` device. +For server-resident longhaul MoE paging, place the same GGUF shards on the RPC +host and point the server at their directory: + +```sh +$ bin/ggml-rpc-server --device MTL0 --longhaul-root /path/to/model-directory +``` + +The client can then use `--rpc`, `--longhaul`, and `--longhaul-cache` together. +RPC longhaul v1 requires one remote Metal device. Expert cache misses are read +from the RPC host's local files rather than uploaded by the client. + ## Usage ### Remote hosts @@ -107,4 +118,3 @@ Use the `GGML_RPC_DEBUG` environment variable to enable debug messages from `ggm ```bash $ GGML_RPC_DEBUG=1 bin/ggml-rpc-server ``` - diff --git a/tools/rpc/rpc-server.cpp b/tools/rpc/rpc-server.cpp index 08e6803..c6406b3 100644 --- a/tools/rpc/rpc-server.cpp +++ b/tools/rpc/rpc-server.cpp @@ -13,6 +13,7 @@ #include #include #include +#include #include #include #include @@ -173,6 +174,7 @@ struct rpc_server_params { std::string host = "127.0.0.1"; int port = 50052; bool use_cache = false; + std::string longhaul_root; int n_threads = std::max(1U, std::thread::hardware_concurrency()/2); std::vector devices; }; @@ -186,6 +188,7 @@ static void print_usage(int /*argc*/, char ** argv, rpc_server_params params) { fprintf(stderr, " -H, --host HOST host to bind to (default: %s)\n", params.host.c_str()); fprintf(stderr, " -p, --port PORT port to bind to (default: %d)\n", params.port); fprintf(stderr, " -c, --cache enable local file cache\n"); + fprintf(stderr, " --longhaul-root PATH serve longhaul expert data from this model directory\n"); fprintf(stderr, "\n"); } @@ -233,6 +236,11 @@ static bool rpc_server_params_parse(int argc, char ** argv, rpc_server_params & } } else if (arg == "-c" || arg == "--cache") { params.use_cache = true; + } else if (arg == "--longhaul-root") { + if (++i >= argc) { + return false; + } + params.longhaul_root = argv[i]; } else if (arg == "-h" || arg == "--help") { print_usage(argc, argv, params); exit(0); @@ -313,6 +321,23 @@ int main(int argc, char * argv[]) { fprintf(stderr, "No devices found\n"); return 1; } + + if (!params.longhaul_root.empty()) { + std::error_code ec; + const auto root = std::filesystem::canonical(params.longhaul_root, ec); + if (ec || !std::filesystem::is_directory(root, ec)) { + fprintf(stderr, "Invalid longhaul root: %s\n", params.longhaul_root.c_str()); + return 1; + } + params.longhaul_root = root.string(); + for (auto * dev : devices) { + auto * dev_reg = ggml_backend_dev_backend_reg(dev); + if (!dev_reg || strcmp(ggml_backend_reg_name(dev_reg), "MTL") != 0) { + fprintf(stderr, "--longhaul-root currently requires Metal devices\n"); + return 1; + } + } + } std::string endpoint = params.host + ":" + std::to_string(params.port); const char * cache_dir = nullptr; std::string cache_dir_str; @@ -331,12 +356,19 @@ int main(int argc, char * argv[]) { return 1; } - auto start_server_fn = (decltype(ggml_backend_rpc_start_server)*) ggml_backend_reg_get_proc_address(reg, "ggml_backend_rpc_start_server"); + auto start_server_fn = (decltype(ggml_backend_rpc_start_server_with_options)*) + ggml_backend_reg_get_proc_address(reg, "ggml_backend_rpc_start_server_with_options"); if (!start_server_fn) { - fprintf(stderr, "Failed to obtain RPC backend start server function\n"); + fprintf(stderr, "Failed to obtain RPC backend start server function with options\n"); return 1; } - start_server_fn(endpoint.c_str(), cache_dir, params.n_threads, devices.size(), devices.data()); + start_server_fn( + endpoint.c_str(), + cache_dir, + params.longhaul_root.empty() ? nullptr : params.longhaul_root.c_str(), + params.n_threads, + devices.size(), + devices.data()); return 0; }