rpc : turn GGML_RPC_DEBUG into a verbosity level and add logs

GGML_RPC_DEBUG is now parsed as a number: 0/unset disables debug logs,
1-3 emit increasingly detailed output (events, per-command trace,
transport detail). Non-numeric values fall back to 1.

The duplicated env/macro blocks in ggml-rpc.cpp and transport.cpp are
replaced by a shared log.h, which becomes the single choke point for
all logging of the RPC backend: LOG_ERROR/LOG_WARN/LOG_INFO for
unconditional severity logs and LOG_DBG/LOG_DBG2/LOG_DBG3 for the
verbosity-gated ones. The transport files no longer need ggml-impl.h,
and the server banner now also goes through the ggml logger (stderr).

Missing logs are added on both the client (handshake, buffer ops,
tensor transfers, graph computes, cache decisions) and the server
(per-command dispatch, graph nodes), including the negotiated
transport via the new socket_t::transport_name().

Assisted-by: pi:llama.cpp/MiMo-V2.6-Flash-RL
This commit is contained in:
Georgi Gerganov
2026-09-27 23:00:44 +03:00
parent 136887b665
commit 1777b02be1
6 changed files with 228 additions and 86 deletions
+125 -57
View File
@@ -3,8 +3,10 @@
#include "ggml-backend-impl.h"
#include "ggml-cpp.h"
#include "transport.h"
#include "log.h"
#include <array>
#include <chrono>
#include <cinttypes>
#include <optional>
#include <string>
@@ -23,12 +25,6 @@
#include <atomic>
#include <thread>
static const char * RPC_DEBUG = std::getenv("GGML_RPC_DEBUG");
#define LOG_DBG(...) \
do { if (RPC_DEBUG) GGML_LOG_DEBUG(__VA_ARGS__); } while (0)
namespace fs = std::filesystem;
// macro for nicer error messages on server crash
@@ -83,6 +79,31 @@ enum rpc_cmd {
static_assert(RPC_CMD_HELLO == 14, "RPC_CMD_HELLO must be always 14");
static const char * rpc_cmd_name(enum rpc_cmd cmd) {
switch (cmd) {
case RPC_CMD_ALLOC_BUFFER: return "ALLOC_BUFFER";
case RPC_CMD_GET_ALIGNMENT: return "GET_ALIGNMENT";
case RPC_CMD_GET_MAX_SIZE: return "GET_MAX_SIZE";
case RPC_CMD_BUFFER_GET_BASE: return "BUFFER_GET_BASE";
case RPC_CMD_FREE_BUFFER: return "FREE_BUFFER";
case RPC_CMD_BUFFER_CLEAR: return "BUFFER_CLEAR";
case RPC_CMD_SET_TENSOR: return "SET_TENSOR";
case RPC_CMD_SET_TENSOR_HASH: return "SET_TENSOR_HASH";
case RPC_CMD_GET_TENSOR: return "GET_TENSOR";
case RPC_CMD_COPY_TENSOR: return "COPY_TENSOR";
case RPC_CMD_GRAPH_COMPUTE: return "GRAPH_COMPUTE";
case RPC_CMD_GET_DEVICE_MEMORY: return "GET_DEVICE_MEMORY";
case RPC_CMD_INIT_TENSOR: return "INIT_TENSOR";
case RPC_CMD_GET_ALLOC_SIZE: return "GET_ALLOC_SIZE";
case RPC_CMD_HELLO: return "HELLO";
case RPC_CMD_DEVICE_COUNT: return "DEVICE_COUNT";
case RPC_CMD_GRAPH_RECOMPUTE: return "GRAPH_RECOMPUTE";
case RPC_CMD_MEMSET_TENSOR: return "MEMSET_TENSOR";
case RPC_CMD_NONE: return "NONE";
default: return "UNKNOWN";
}
}
// Try RPC_CMD_SET_TENSOR_HASH first when data size is larger than this threshold
const size_t HASH_THRESHOLD = 10 * 1024 * 1024;
@@ -284,7 +305,7 @@ static bool recv_msg(socket_ptr sock, std::vector<uint8_t> & input) {
try {
input.resize(size);
} catch (const std::bad_alloc & e) {
GGML_LOG_ERROR("Failed to allocate input buffer of size %" PRIu64 "\n", size);
LOG_ERROR("Failed to allocate input buffer of size %" PRIu64 "\n", size);
return false;
}
return sock->recv_data(input.data(), size);
@@ -354,12 +375,14 @@ static bool negotiate_hello(const std::shared_ptr<socket_t> & sock) {
RPC_STATUS_ASSERT(status);
if (response.major != RPC_PROTO_MAJOR_VERSION || response.minor > RPC_PROTO_MINOR_VERSION) {
GGML_LOG_ERROR("RPC server version mismatch: %d.%d.%d\n",
LOG_ERROR("RPC server version mismatch: %d.%d.%d\n",
response.major, response.minor, response.patch);
return false;
}
sock->update_caps(response.conn_caps);
LOG_DBG("[%s] handshake ok: server version %d.%d.%d, transport: %s\n",
__func__, response.major, response.minor, response.patch, sock->transport_name());
return true;
}
@@ -457,6 +480,7 @@ void rpc_dispatcher::send(enum rpc_cmd cmd, std::shared_ptr<const void> input, s
msg->input_size = input_size;
msg->output = nullptr;
msg->output_size = 0;
LOG_DBG3("[%s] enqueue %s (in: %zu)\n", __func__, rpc_cmd_name(cmd), input_size);
GGML_ASSERT(queue.push(msg));
auto future = msg->completion.get_future();
future.wait();
@@ -469,6 +493,7 @@ void rpc_dispatcher::send_async(enum rpc_cmd cmd, std::shared_ptr<const void> in
msg->input_size = input_size;
msg->output = nullptr;
msg->output_size = 0;
LOG_DBG3("[%s] enqueue %s (in: %zu)\n", __func__, rpc_cmd_name(cmd), input_size);
GGML_ASSERT(queue.push(msg));
}
@@ -479,6 +504,7 @@ void rpc_dispatcher::send(enum rpc_cmd cmd, std::shared_ptr<const void> input, s
msg->input_size = input_size;
msg->output = output;
msg->output_size = output_size;
LOG_DBG3("[%s] enqueue %s (in: %zu, out: %zu)\n", __func__, rpc_cmd_name(cmd), input_size, output_size);
GGML_ASSERT(queue.push(msg));
auto future = msg->completion.get_future();
future.wait();
@@ -491,6 +517,7 @@ void rpc_dispatcher::send_async(enum rpc_cmd cmd, std::shared_ptr<const void> in
msg->input_size = input_size;
msg->output = output;
msg->output_size = output_size;
LOG_DBG3("[%s] enqueue %s (in: %zu, out: %zu)\n", __func__, rpc_cmd_name(cmd), input_size, output_size);
GGML_ASSERT(queue.push(msg));
}
@@ -561,6 +588,9 @@ void rpc_dispatcher::work() {
break;
}
if (msg_ptr->cmd != RPC_CMD_NONE) {
LOG_DBG2("[%s] %s (in: %zu, out: %zu)\n", __func__, rpc_cmd_name(msg_ptr->cmd), msg_ptr->input_size, msg_ptr->output_size);
const bool trace = rpc_debug_level() >= 3;
const auto t0 = std::chrono::steady_clock::now();
if (msg_ptr->output) {
bool status = send_rpc_cmd(sock, msg_ptr->cmd, msg_ptr->input.get(), msg_ptr->input_size, msg_ptr->output, msg_ptr->output_size);
RPC_STATUS_ASSERT(status);
@@ -568,6 +598,12 @@ void rpc_dispatcher::work() {
bool status = send_rpc_cmd(sock, msg_ptr->cmd, msg_ptr->input.get(), msg_ptr->input_size);
RPC_STATUS_ASSERT(status);
}
if (trace) {
const std::chrono::duration<double, std::milli> elapsed = std::chrono::steady_clock::now() - t0;
LOG_DBG3("[%s] %s done in %.3f ms\n", __func__, rpc_cmd_name(msg_ptr->cmd), elapsed.count());
}
} else {
LOG_DBG3("[%s] sync barrier\n", __func__);
}
msg_ptr->completion.set_value();
}
@@ -604,6 +640,7 @@ static void ggml_backend_rpc_buffer_free_buffer(ggml_backend_buffer_t buffer) {
ggml_backend_rpc_buffer_context * ctx = (ggml_backend_rpc_buffer_context *)buffer->context;
auto request = std::make_shared<rpc_msg_free_buffer_req>();
request->remote_ptr = ctx->remote_ptr;
LOG_DBG("[%s] remote_ptr: 0x%" PRIx64 "\n", __func__, request->remote_ptr);
ctx->dispatcher->send(RPC_CMD_FREE_BUFFER, request, sizeof(*request));
delete ctx;
}
@@ -615,6 +652,7 @@ static void * ggml_backend_rpc_buffer_get_base(ggml_backend_buffer_t buffer) {
}
auto request = std::make_shared<rpc_msg_buffer_get_base_req>();
request->remote_ptr = ctx->remote_ptr;
LOG_DBG3("[%s] remote_ptr: 0x%" PRIx64 "\n", __func__, ctx->remote_ptr);
rpc_msg_buffer_get_base_rsp response;
ctx->dispatcher->send(RPC_CMD_BUFFER_GET_BASE, request, sizeof(*request), &response, sizeof(response));
ctx->base_ptr = reinterpret_cast<void *>(response.base_ptr);
@@ -679,6 +717,7 @@ static enum ggml_status ggml_backend_rpc_buffer_init_tensor(ggml_backend_buffer_
// Due to bandwidth constraints, we only call the server init tensor functions if necessary.
// In particular, only quantized tensors need padding
if (ggml_is_quantized(tensor->type) && (tensor->ne[0] % 512 != 0) && (tensor->view_src == nullptr)) {
LOG_DBG2("[%s] tensor: %s\n", __func__, tensor->name);
auto request = std::make_shared<rpc_msg_init_tensor_req>();
request->tensor = serialize_tensor(tensor);
ctx->dispatcher->send(RPC_CMD_INIT_TENSOR, request, sizeof(*request));
@@ -694,6 +733,7 @@ static void ggml_backend_rpc_buffer_memset_tensor(
request->offset = offset;
request->size = size;
request->value = value;
LOG_DBG2("[%s] tensor: %s, offset: %zu, size: %zu, value: %u\n", __func__, tensor->name, offset, size, value);
ctx->dispatcher->send(RPC_CMD_MEMSET_TENSOR, request, sizeof(*request));
}
@@ -721,6 +761,7 @@ static void ggml_backend_rpc_buffer_set_tensor(ggml_backend_buffer_t buffer, ggm
ggml_backend_rpc_buffer_context * ctx = (ggml_backend_rpc_buffer_context *)buffer->context;
rpc_tensor rpc_tensor = serialize_tensor(tensor);
uint8_t cache_flag = 0;
LOG_DBG("[%s] tensor: %s, offset: %zu, size: %zu\n", __func__, tensor->name, offset, size);
if (rpc_use_hash_cache(tensor, size)) {
auto request = std::make_shared<rpc_msg_set_tensor_hash_req>();
request->tensor = rpc_tensor;
@@ -728,6 +769,7 @@ static void ggml_backend_rpc_buffer_set_tensor(ggml_backend_buffer_t buffer, ggm
request->hash = fnv_hash((const uint8_t*)data, size);
rpc_msg_set_tensor_hash_rsp response;
ctx->dispatcher->send(RPC_CMD_SET_TENSOR_HASH, request, sizeof(*request), &response, sizeof(response));
LOG_DBG2("[%s] tensor: %s, hash: 0x%" PRIx64 ", cache: %s\n", __func__, tensor->name, request->hash, response.result ? "hit" : "miss");
if (response.result) {
// the server has the same data, no need to send it
return;
@@ -742,6 +784,7 @@ static void ggml_backend_rpc_buffer_set_tensor(ggml_backend_buffer_t buffer, ggm
static void ggml_backend_rpc_buffer_get_tensor(ggml_backend_buffer_t buffer, const ggml_tensor * tensor, void * data, size_t offset, size_t size) {
ggml_backend_rpc_buffer_context * ctx = (ggml_backend_rpc_buffer_context *)buffer->context;
LOG_DBG("[%s] tensor: %s, offset: %zu, size: %zu\n", __func__, tensor->name, offset, size);
auto request = std::make_shared<rpc_msg_get_tensor_req>();
request->tensor = serialize_tensor(tensor);
request->offset = offset;
@@ -757,9 +800,11 @@ static bool ggml_backend_rpc_buffer_cpy_tensor(ggml_backend_buffer_t buffer, con
ggml_backend_buffer_t dst_buffer = dst->buffer;
ggml_backend_rpc_buffer_context * dst_ctx = (ggml_backend_rpc_buffer_context *)dst_buffer->context;
if (src_ctx->dispatcher != dst_ctx->dispatcher) {
LOG_DBG2("[%s] src and dst are on different servers, falling back to a host copy\n", __func__);
return false;
}
ggml_backend_rpc_buffer_context * ctx = (ggml_backend_rpc_buffer_context *)buffer->context;
LOG_DBG("[%s] src: %s -> dst: %s\n", __func__, src->name, dst->name);
auto request = std::make_shared<rpc_msg_copy_tensor_req>();
request->src = serialize_tensor(src);
request->dst = serialize_tensor(dst);
@@ -775,6 +820,7 @@ static void ggml_backend_rpc_buffer_clear(ggml_backend_buffer_t buffer, uint8_t
auto request = std::make_shared<rpc_msg_buffer_clear_req>();
request->remote_ptr = ctx->remote_ptr;
request->value = value;
LOG_DBG2("[%s] remote_ptr: 0x%" PRIx64 ", value: %u\n", __func__, request->remote_ptr, request->value);
ctx->dispatcher->send(RPC_CMD_BUFFER_CLEAR, request, sizeof(*request));
}
@@ -807,12 +853,16 @@ static ggml_backend_buffer_t ggml_backend_rpc_buffer_type_alloc_buffer(ggml_back
auto dispatcher = get_dispatcher(buft_ctx->endpoint);
dispatcher->send(RPC_CMD_ALLOC_BUFFER, request, sizeof(*request), &response, sizeof(response));
if (response.remote_ptr != 0) {
LOG_DBG("[%s] endpoint: %s, device: %u, size: %" PRIu64 " -> remote_ptr: 0x%" PRIx64 ", remote_size: %" PRIu64 "\n",
__func__, buft_ctx->endpoint.c_str(), buft_ctx->device, request->size, response.remote_ptr, response.remote_size);
ggml_backend_buffer_t buffer = ggml_backend_buffer_init(buft,
ggml_backend_rpc_buffer_interface,
new ggml_backend_rpc_buffer_context{dispatcher, nullptr, response.remote_ptr},
response.remote_size);
return buffer;
} else {
LOG_DBG("[%s] endpoint: %s, device: %u, size: %" PRIu64 " -> failed\n",
__func__, buft_ctx->endpoint.c_str(), buft_ctx->device, request->size);
return nullptr;
}
}
@@ -822,6 +872,7 @@ static size_t get_alignment(const std::shared_ptr<rpc_dispatcher> & dispatcher,
request->device = device;
rpc_msg_get_alignment_rsp response;
dispatcher->send(RPC_CMD_GET_ALIGNMENT, request, sizeof(*request), &response, sizeof(response));
LOG_DBG2("[%s] device: %u, alignment: %" PRIu64 "\n", __func__, device, response.alignment);
return response.alignment;
}
@@ -835,6 +886,7 @@ static size_t get_max_size(const std::shared_ptr<rpc_dispatcher> & dispatcher, u
request->device = device;
rpc_msg_get_max_size_rsp response;
dispatcher->send(RPC_CMD_GET_MAX_SIZE, request, sizeof(*request), &response, sizeof(response));
LOG_DBG2("[%s] device: %u, max_size: %" PRIu64 "\n", __func__, device, response.max_size);
return response.max_size;
}
@@ -900,6 +952,7 @@ static size_t ggml_backend_rpc_buffer_type_get_alloc_size(ggml_backend_buffer_ty
std::lock_guard<std::mutex> lock(cache_mutex);
auto it = cache.find(cache_hash);
if (it != cache.end()) {
LOG_DBG3("[%s] cache hit: %s [%s] -> %zu\n", __func__, ggml_op_name(tensor->op), tensor->name, it->second);
return std::max<size_t>(it->second, min_size);
}
}
@@ -913,6 +966,7 @@ static size_t ggml_backend_rpc_buffer_type_get_alloc_size(ggml_backend_buffer_ty
request->srcs[i] = serialize_tensor(tensor->src[i]);
}
LOG_DBG2("[%s] cache miss: %s [%s], querying server\n", __func__, ggml_op_name(tensor->op), tensor->name);
rpc_msg_get_alloc_size_rsp response;
auto dispatcher = get_dispatcher(buft_ctx->endpoint);
dispatcher->send(RPC_CMD_GET_ALLOC_SIZE, request, sizeof(*request), &response, sizeof(response));
@@ -922,6 +976,8 @@ static size_t ggml_backend_rpc_buffer_type_get_alloc_size(ggml_backend_buffer_ty
cache[cache_hash] = response.alloc_size;
}
LOG_DBG2("[%s] %s [%s] -> alloc_size: %" PRIu64 "\n", __func__, ggml_op_name(tensor->op), tensor->name, response.alloc_size);
return std::max<size_t>(response.alloc_size, min_size);
}
@@ -953,6 +1009,7 @@ static void ggml_backend_rpc_set_tensor_async(ggml_backend_t backend, ggml_tenso
ggml_backend_rpc_context * ctx = (ggml_backend_rpc_context *)backend->context;
rpc_tensor rpc_tensor = serialize_tensor(tensor);
uint8_t cache_flag = 0;
LOG_DBG("[%s] tensor: %s, offset: %zu, size: %zu\n", __func__, tensor->name, offset, size);
if (rpc_use_hash_cache(tensor, size)) {
auto request = std::make_shared<rpc_msg_set_tensor_hash_req>();
request->tensor = rpc_tensor;
@@ -961,6 +1018,7 @@ static void ggml_backend_rpc_set_tensor_async(ggml_backend_t backend, ggml_tenso
rpc_msg_set_tensor_hash_rsp response;
// TODO: make this async
ctx->dispatcher->send(RPC_CMD_SET_TENSOR_HASH, request, sizeof(*request), &response, sizeof(response));
LOG_DBG2("[%s] tensor: %s, hash: 0x%" PRIx64 ", cache: %s\n", __func__, tensor->name, request->hash, response.result ? "hit" : "miss");
if (response.result) {
// the server has the same data, no need to send it
return;
@@ -1042,6 +1100,7 @@ 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;
LOG_DBG("[%s] uid: %" PRIu64 ", n_nodes: %u, reuse: %s\n", __func__, cgraph->uid, cgraph->n_nodes, reuse ? "yes" : "no");
if (reuse) {
auto request = std::make_shared<rpc_msg_graph_recompute_req>();
request->device = rpc_ctx->device;
@@ -1050,6 +1109,7 @@ static enum ggml_status ggml_backend_rpc_graph_compute(ggml_backend_t backend, g
rpc_dev_ctx->last_graph_uid = cgraph->uid;
size_t input_size = 0;
uint8_t * input = serialize_graph(rpc_ctx->device, cgraph, rpc_ctx->dispatcher, &input_size);
LOG_DBG2("[%s] serialized graph: %zu bytes\n", __func__, input_size);
std::shared_ptr<uint8_t> input_ptr(input, std::default_delete<uint8_t[]>());
rpc_ctx->dispatcher->send_async(RPC_CMD_GRAPH_COMPUTE, input_ptr, input_size);
}
@@ -1113,6 +1173,7 @@ ggml_backend_buffer_type_t ggml_backend_rpc_buffer_type(const char * endpoint, u
/* .context = */ buft_ctx
};
buft_map[buft_name] = buft;
LOG_DBG2("[%s] created %s (alignment: %zu, max_size: %zu)\n", __func__, buft_name.c_str(), alignment, max_size);
return buft;
}
@@ -1125,6 +1186,7 @@ ggml_backend_t ggml_backend_rpc_init(const char * endpoint, uint32_t device) {
/* .name = */ dev_name,
};
auto reg = ggml_backend_rpc_add_server(endpoint);
LOG_DBG2("[%s] %s\n", __func__, dev_name.c_str());
ggml_backend_t backend = new ggml_backend {
/* .guid = */ ggml_backend_rpc_guid(),
/* .iface = */ ggml_backend_rpc_interface,
@@ -1144,6 +1206,7 @@ void ggml_backend_rpc_get_device_memory(const char * endpoint, uint32_t device,
request->device = device;
rpc_msg_get_device_memory_rsp response;
dispatcher->send(RPC_CMD_GET_DEVICE_MEMORY, request, sizeof(*request), &response, sizeof(response));
LOG_DBG3("[%s] device: %u, free: %" PRIu64 ", total: %" PRIu64 "\n", __func__, device, response.free_mem, response.total_mem);
*free = response.free_mem;
*total = response.total_mem;
}
@@ -1222,7 +1285,7 @@ bool rpc_server::get_alloc_size(const rpc_msg_get_alloc_size_req & request, rpc_
ggml_tensor * tensor = deserialize_tensor(ctx, &request.tensor);
if (tensor == nullptr) {
GGML_LOG_ERROR("Null tensor pointer passed to server get_alloc_size function.\n");
LOG_ERROR("Null tensor pointer passed to server get_alloc_size function.\n");
return false;
}
for (int i = 0; i < GGML_MAX_SRC; i++) {
@@ -1293,7 +1356,7 @@ bool rpc_server::buffer_get_base(const rpc_msg_buffer_get_base_req & request, rp
LOG_DBG("[%s] remote_ptr: %" PRIx64 "\n", __func__, request.remote_ptr);
ggml_backend_buffer_t buffer = reinterpret_cast<ggml_backend_buffer_t>(request.remote_ptr);
if (buffers.find(buffer) == buffers.end()) {
GGML_LOG_ERROR("[%s] buffer not found\n", __func__);
LOG_ERROR("[%s] buffer not found\n", __func__);
return false;
}
void * base = ggml_backend_buffer_get_base(buffer);
@@ -1305,7 +1368,7 @@ bool rpc_server::free_buffer(const rpc_msg_free_buffer_req & request) {
LOG_DBG("[%s] remote_ptr: %" PRIx64 "\n", __func__, request.remote_ptr);
ggml_backend_buffer_t buffer = reinterpret_cast<ggml_backend_buffer_t>(request.remote_ptr);
if (buffers.find(buffer) == buffers.end()) {
GGML_LOG_ERROR("[%s] buffer not found\n", __func__);
LOG_ERROR("[%s] buffer not found\n", __func__);
return false;
}
// Discard all cached graphs to avoid use-after-free in graph_recompute,
@@ -1322,7 +1385,7 @@ bool rpc_server::buffer_clear(const rpc_msg_buffer_clear_req & request) {
LOG_DBG("[%s] remote_ptr: %" PRIx64 ", value: %u\n", __func__, request.remote_ptr, request.value);
ggml_backend_buffer_t buffer = reinterpret_cast<ggml_backend_buffer_t>(request.remote_ptr);
if (buffers.find(buffer) == buffers.end()) {
GGML_LOG_ERROR("[%s] buffer not found\n", __func__);
LOG_ERROR("[%s] buffer not found\n", __func__);
return false;
}
ggml_backend_buffer_clear(buffer, request.value);
@@ -1340,13 +1403,13 @@ bool rpc_server::memset_tensor(const rpc_msg_memset_tensor_req & request) {
ggml_context * ctx = ctx_ptr.get();
ggml_tensor * tensor = deserialize_tensor(ctx, &request.tensor);
if (tensor == nullptr || tensor->buffer == nullptr) {
GGML_LOG_ERROR("[%s] error deserializing tensor\n", __func__);
LOG_ERROR("[%s] error deserializing tensor\n", __func__);
return false;
}
const uint64_t tensor_size = ggml_nbytes(tensor);
if (request.offset > tensor_size || request.size > tensor_size - request.offset) {
GGML_LOG_ERROR("[%s] tensor region (offset=%" PRIu64 ", size=%" PRIu64 ") out of tensor bounds [0, %" PRIu64 ")\n",
LOG_ERROR("[%s] tensor region (offset=%" PRIu64 ", size=%" PRIu64 ") out of tensor bounds [0, %" PRIu64 ")\n",
__func__, request.offset, request.size, tensor_size);
return false;
}
@@ -1354,18 +1417,18 @@ bool rpc_server::memset_tensor(const rpc_msg_memset_tensor_req & request) {
const uint64_t buffer_start = (uint64_t) ggml_backend_buffer_get_base(tensor->buffer);
const uint64_t buffer_size = ggml_backend_buffer_get_size(tensor->buffer);
if (request.tensor.data < buffer_start) {
GGML_LOG_ERROR("[%s] tensor data before buffer start\n", __func__);
LOG_ERROR("[%s] tensor data before buffer start\n", __func__);
return false;
}
const uint64_t data_offset = request.tensor.data - buffer_start;
if (data_offset > buffer_size ||
request.offset > buffer_size - data_offset ||
request.size > buffer_size - data_offset - request.offset) {
GGML_LOG_ERROR("[%s] tensor region out of buffer bounds\n", __func__);
LOG_ERROR("[%s] tensor region out of buffer bounds\n", __func__);
return false;
}
if (tensor->buffer->iface.memset_tensor == nullptr) {
GGML_LOG_ERROR("[%s] memset not implemented by backend buffer\n", __func__);
LOG_ERROR("[%s] memset not implemented by backend buffer\n", __func__);
return false;
}
@@ -1378,13 +1441,13 @@ bool rpc_server::memset_tensor(const rpc_msg_memset_tensor_req & request) {
ggml_tensor * rpc_server::deserialize_tensor(struct ggml_context * ctx, const rpc_tensor * tensor) {
// Validate tensor type before using it
if (tensor->type >= GGML_TYPE_COUNT) {
GGML_LOG_ERROR("[%s] invalid tensor type received: %u\n", __func__, tensor->type);
LOG_ERROR("[%s] invalid tensor type received: %u\n", __func__, tensor->type);
return nullptr;
}
// Fix: Prevent division by zero if blck_size is 0 (e.g., deprecated types)
if (ggml_blck_size((enum ggml_type)tensor->type) == 0) {
GGML_LOG_ERROR("[%s] invalid tensor type received (blck_size is 0): %u\n", __func__, tensor->type);
LOG_ERROR("[%s] invalid tensor type received (blck_size is 0): %u\n", __func__, tensor->type);
return nullptr;
}
@@ -1393,7 +1456,7 @@ ggml_tensor * rpc_server::deserialize_tensor(struct ggml_context * ctx, const rp
// ggml_new_tensor_4d might fail if dimensions are invalid, although less likely to crash than invalid type
if (result == nullptr) {
GGML_LOG_ERROR("[%s] ggml_new_tensor_4d failed for type %u\n", __func__, tensor->type);
LOG_ERROR("[%s] ggml_new_tensor_4d failed for type %u\n", __func__, tensor->type);
return nullptr;
}
@@ -1448,7 +1511,7 @@ bool rpc_server::set_tensor(const std::vector<uint8_t> & input) {
ggml_context * ctx = ctx_ptr.get();
ggml_tensor * tensor = deserialize_tensor(ctx, in_tensor);
if (tensor == nullptr || tensor->buffer == nullptr) {
GGML_LOG_ERROR("[%s] error deserializing tensor\n", __func__);
LOG_ERROR("[%s] error deserializing tensor\n", __func__);
return false;
}
LOG_DBG("[%s] buffer: %p, data: %p, offset: %" PRIu64 ", size: %zu\n", __func__, (void*)tensor->buffer, tensor->data, offset, size);
@@ -1459,7 +1522,7 @@ bool rpc_server::set_tensor(const std::vector<uint8_t> & input) {
const size_t p1 = p0 + ggml_backend_buffer_get_size(tensor->buffer);
if (in_tensor->data + offset < p0 || in_tensor->data + offset >= p1 || size > (p1 - in_tensor->data - offset)) {
GGML_LOG_ERROR("[%s] tensor data region (data=0x%" PRIx64 ", offset=%" PRIu64 ", size=%zu) out of buffer bounds [0x%zx, 0x%zx)\n",
LOG_ERROR("[%s] tensor data region (data=0x%" PRIx64 ", offset=%" PRIu64 ", size=%zu) out of buffer bounds [0x%zx, 0x%zx)\n",
__func__, in_tensor->data, offset, size, p0, p1);
return false;
}
@@ -1474,7 +1537,7 @@ bool rpc_server::set_tensor(const std::vector<uint8_t> & input) {
fs::path cache_file = fs::path(cache_dir) / hash_str;
std::ofstream ofs(cache_file, std::ios::binary);
ofs.write((const char *)data, size);
GGML_LOG_INFO("[%s] saved to '%s'\n", __func__, cache_file.string().c_str());
LOG_INFO("[%s] saved to '%s'\n", __func__, cache_file.string().c_str());
}
ggml_backend_tensor_set(tensor, data, offset, size);
return true;
@@ -1504,10 +1567,12 @@ bool rpc_server::set_tensor_hash(const rpc_msg_set_tensor_hash_req & request, rp
{
std::vector<uint8_t> cached_file;
if (!get_cached_file(request.hash, cached_file)) {
LOG_DBG2("[%s] hash: 0x%" PRIx64 ", cache miss\n", __func__, request.hash);
response.result = 0;
return true;
}
size_t size = cached_file.size();
LOG_DBG2("[%s] hash: 0x%" PRIx64 ", cache hit (%zu bytes)\n", __func__, request.hash, size);
struct ggml_init_params params {
/*.mem_size =*/ ggml_tensor_overhead(),
/*.mem_buffer =*/ NULL,
@@ -1518,7 +1583,7 @@ bool rpc_server::set_tensor_hash(const rpc_msg_set_tensor_hash_req & request, rp
ggml_context * ctx = ctx_ptr.get();
ggml_tensor * tensor = deserialize_tensor(ctx, &request.tensor);
if (tensor == nullptr || tensor->buffer == nullptr) {
GGML_LOG_ERROR("[%s] error deserializing tensor\n", __func__);
LOG_ERROR("[%s] error deserializing tensor\n", __func__);
return false;
}
LOG_DBG("[%s] buffer: %p, data: %p, offset: %" PRIu64 ", size: %zu, hash: %" PRIx64 "\n",
@@ -1532,7 +1597,7 @@ bool rpc_server::set_tensor_hash(const rpc_msg_set_tensor_hash_req & request, rp
if (request.tensor.data + request.offset < p0
|| request.tensor.data + request.offset >= p1
|| size > (p1 - request.tensor.data - request.offset)) {
GGML_LOG_ERROR("[%s] tensor data region (data=0x%" PRIx64 ", offset=%" PRIu64 ", size=%zu, hash=0x%" PRIx64 ") out of buffer bounds [0x%zx, 0x%zx)\n",
LOG_ERROR("[%s] tensor data region (data=0x%" PRIx64 ", offset=%" PRIu64 ", size=%zu, hash=0x%" PRIx64 ") out of buffer bounds [0x%zx, 0x%zx)\n",
__func__, request.tensor.data, request.offset, size, request.hash, p0, p1);
return false;
}
@@ -1553,7 +1618,7 @@ bool rpc_server::init_tensor(const rpc_msg_init_tensor_req & request) {
ggml_context * ctx = ctx_ptr.get();
ggml_tensor * tensor = deserialize_tensor(ctx, &request.tensor);
if (tensor == nullptr) {
GGML_LOG_ERROR("Null tensor pointer passed to server init_tensor function.\n");
LOG_ERROR("Null tensor pointer passed to server init_tensor function.\n");
return false;
}
LOG_DBG("[%s] buffer: %p, data: %p\n", __func__, (void*)tensor->buffer, tensor->data);
@@ -1563,14 +1628,14 @@ bool rpc_server::init_tensor(const rpc_msg_init_tensor_req & request) {
buffer->iface.init_tensor(buffer, tensor);
} else {
if (!buffer) {
GGML_LOG_ERROR("Tensor with null buffer passed to init_tensor function\n");
LOG_ERROR("Tensor with null buffer passed to init_tensor function\n");
}
}
if (tensor->extra != nullptr) {
// This pointer can either be passed around client/server, or probably better stored server-side and kept track of.
// Currently unimplemented.
GGML_LOG_ERROR("tensor->extra populated by the backend, this is currently unsupported.\n");
LOG_ERROR("tensor->extra populated by the backend, this is currently unsupported.\n");
return false;
}
@@ -1588,7 +1653,7 @@ bool rpc_server::get_tensor(const rpc_msg_get_tensor_req & request, std::vector<
ggml_context * ctx = ctx_ptr.get();
ggml_tensor * tensor = deserialize_tensor(ctx, &request.tensor);
if (tensor == nullptr || tensor->buffer == nullptr) {
GGML_LOG_ERROR("[%s] error deserializing tensor\n", __func__);
LOG_ERROR("[%s] error deserializing tensor\n", __func__);
return false;
}
LOG_DBG("[%s] buffer: %p, data: %p, offset: %" PRIu64 ", size: %" PRIu64 "\n", __func__, (void*)tensor->buffer, tensor->data, request.offset, request.size);
@@ -1601,7 +1666,7 @@ bool rpc_server::get_tensor(const rpc_msg_get_tensor_req & request, std::vector<
if (request.tensor.data + request.offset < p0 ||
request.tensor.data + request.offset >= p1 ||
request.size > (p1 - request.tensor.data - request.offset)) {
GGML_LOG_ERROR("[%s] requested tensor region (data=0x%" PRIx64 ", offset=%" PRIu64 ", size=%" PRIu64 ") out of buffer bounds [0x%zx, 0x%zx)\n",
LOG_ERROR("[%s] requested tensor region (data=0x%" PRIx64 ", offset=%" PRIu64 ", size=%" PRIu64 ") out of buffer bounds [0x%zx, 0x%zx)\n",
__func__, request.tensor.data, request.offset, request.size, p0, p1);
return false;
}
@@ -1625,7 +1690,7 @@ bool rpc_server::copy_tensor(const rpc_msg_copy_tensor_req & request, rpc_msg_co
ggml_tensor * src = deserialize_tensor(ctx, &request.src);
ggml_tensor * dst = deserialize_tensor(ctx, &request.dst);
if (src == nullptr || dst == nullptr || src->buffer == nullptr || dst->buffer == nullptr) {
GGML_LOG_ERROR("[%s] error deserializing tensors\n", __func__);
LOG_ERROR("[%s] error deserializing tensors\n", __func__);
return false;
}
@@ -1635,7 +1700,7 @@ bool rpc_server::copy_tensor(const rpc_msg_copy_tensor_req & request, rpc_msg_co
uint64_t dst_buf_sz = (uint64_t) ggml_backend_buffer_get_size(dst->buffer);
if (dst_data + src_size > dst_base + dst_buf_sz) {
GGML_LOG_ERROR("[%s] out-of-bounds write in rpc_server::copy_tensor:\n"
LOG_ERROR("[%s] out-of-bounds write in rpc_server::copy_tensor:\n"
" write range : [0x%" PRIx64 ", 0x%" PRIx64 "]\n"
" buffer base: [0x%" PRIx64 ", 0x%" PRIx64 "]\n",
__func__,
@@ -1672,7 +1737,7 @@ ggml_tensor * rpc_server::create_node(uint64_t id,
return nullptr;
}
if (result->buffer == nullptr && result->data != nullptr) {
GGML_LOG_ERROR("[%s] invalid data ptr", __func__);
LOG_ERROR("[%s] invalid data ptr", __func__);
return nullptr;
}
tensor_map[id] = result;
@@ -1684,7 +1749,7 @@ ggml_tensor * rpc_server::create_node(uint64_t id,
result->src[i] = create_node(tensor->src[i], ctx, tensor_ptrs, tensor_map);
// If the recursive call failed for a non-zero ID, propagate the error
if (result->src[i] == nullptr) {
GGML_LOG_ERROR("[%s] failed to create source node %d (src_id=%" PRIu64 ") for node id %" PRIu64 "\n",
LOG_ERROR("[%s] failed to create source node %d (src_id=%" PRIu64 ") for node id %" PRIu64 "\n",
__func__, i, tensor->src[i], id);
// Must return nullptr to signal failure up the call stack
return nullptr;
@@ -1699,7 +1764,7 @@ ggml_tensor * rpc_server::create_node(uint64_t id,
result->view_src = create_node(tensor->view_src, ctx, tensor_ptrs, tensor_map);
// If the recursive call failed for a non-zero ID, propagate the error
if (result->view_src == nullptr) {
GGML_LOG_ERROR("[%s] failed to create view_src node (view_src_id=%" PRIu64 ") for node id %" PRIu64 "\n",
LOG_ERROR("[%s] failed to create view_src node (view_src_id=%" PRIu64 ") for node id %" PRIu64 "\n",
__func__, tensor->view_src, id);
// Must return nullptr to signal failure up the call stack
return nullptr;
@@ -1768,10 +1833,11 @@ bool rpc_server::graph_compute(const std::vector<uint8_t> & input) {
// If id was 0, create_node returning nullptr is expected.
// If id was non-zero and create_node returned nullptr, it indicates a deserialization error.
if (graph->nodes[i] == nullptr && id != 0) {
GGML_LOG_ERROR("[%s] failed to create graph node %d (id=%" PRId64 ")\n", __func__, i, id);
LOG_ERROR("[%s] failed to create graph node %d (id=%" PRId64 ")\n", __func__, i, id);
return false;
}
if (graph->nodes[i] != nullptr) {
LOG_DBG3("[%s] node %u: id=%" PRId64 ", op=%s, name=%s\n", __func__, i, id, ggml_op_name(graph->nodes[i]->op), graph->nodes[i]->name);
const size_t hash_pos = ggml_hash_insert(&graph->visited_hash_set, graph->nodes[i]);
graph->use_counts[hash_pos] = tensor_ptrs.at(id)->use_count;
}
@@ -1825,7 +1891,7 @@ static void rpc_serve_client(const std::vector<ggml_backend_t> & backends, const
return;
}
if (cmd != RPC_CMD_HELLO) {
GGML_LOG_ERROR("Expected HELLO command, update client\n");
LOG_ERROR("Expected HELLO command, update client\n");
return;
}
@@ -1836,7 +1902,7 @@ static void rpc_serve_client(const std::vector<ggml_backend_t> & backends, const
}
if (hello_input_size != sizeof(rpc_msg_hello_req)) {
GGML_LOG_ERROR("HELLO request size mismatch (%zu vs %zu) — client needs upgrade to protocol v%d.x\n",
LOG_ERROR("HELLO request size mismatch (%zu vs %zu) — client needs upgrade to protocol v%d.x\n",
(size_t)hello_input_size, sizeof(rpc_msg_hello_req), RPC_PROTO_MAJOR_VERSION);
return;
}
@@ -1856,15 +1922,17 @@ static void rpc_serve_client(const std::vector<ggml_backend_t> & backends, const
// Activate transport upgrade using client's caps
sock->update_caps(req.conn_caps);
LOG_DBG("[%s] client connected, transport: %s\n", __func__, sock->transport_name());
while (true) {
if (!sock->recv_data(&cmd, 1)) {
break;
}
if (cmd >= RPC_CMD_COUNT) {
// fail fast if the command is invalid
GGML_LOG_ERROR("Unknown command: %d\n", cmd);
LOG_ERROR("Unknown command: %d\n", cmd);
break;
}
LOG_DBG2("[%s] %s\n", __func__, rpc_cmd_name((enum rpc_cmd)cmd));
switch (cmd) {
case RPC_CMD_HELLO: {
// HELLO command is handled above
@@ -2078,7 +2146,7 @@ static void rpc_serve_client(const std::vector<ggml_backend_t> & backends, const
break;
}
default: {
GGML_LOG_ERROR("Unknown command: %d\n", cmd);
LOG_ERROR("Unknown command: %d\n", cmd);
return;
}
}
@@ -2088,26 +2156,26 @@ static void rpc_serve_client(const std::vector<ggml_backend_t> & 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) {
if (n_devices == 0 || devices == nullptr) {
fprintf(stderr, "Invalid arguments to ggml_backend_rpc_start_server\n");
LOG_ERROR("Invalid arguments to ggml_backend_rpc_start_server\n");
return;
}
std::vector<ggml_backend_t> backends;
printf("Starting RPC server v%d.%d.%d\n",
LOG_INFO("Starting RPC server v%d.%d.%d\n",
RPC_PROTO_MAJOR_VERSION,
RPC_PROTO_MINOR_VERSION,
RPC_PROTO_PATCH_VERSION);
printf(" endpoint : %s\n", endpoint);
printf(" local cache : %s\n", cache_dir ? cache_dir : "n/a");
printf("Devices:\n");
LOG_INFO(" endpoint : %s\n", endpoint);
LOG_INFO(" local cache : %s\n", cache_dir ? cache_dir : "n/a");
LOG_INFO("Devices:\n");
for (size_t i = 0; i < n_devices; i++) {
auto dev = devices[i];
size_t free, total;
ggml_backend_dev_memory(dev, &free, &total);
printf(" %s: %s (%zu MiB, %zu MiB free)\n", ggml_backend_dev_name(dev), ggml_backend_dev_description(dev),
total / 1024 / 1024, free / 1024 / 1024);
LOG_INFO(" %s: %s (%zu MiB, %zu MiB free)\n", ggml_backend_dev_name(dev), ggml_backend_dev_description(dev),
total / 1024 / 1024, free / 1024 / 1024);
auto backend = ggml_backend_dev_init(dev, nullptr);
if (!backend) {
fprintf(stderr, "Failed to create backend for device %s\n", dev->iface.get_name(dev));
LOG_ERROR("Failed to create backend for device %s\n", dev->iface.get_name(dev));
return;
}
backends.push_back(backend);
@@ -2127,30 +2195,28 @@ void ggml_backend_rpc_start_server(const char * endpoint, const char * cache_dir
}
#ifdef GGML_RPC_RDMA
printf(" transport : TCP (RDMA auto-negotiate enabled)\n");
LOG_INFO(" transport : TCP (RDMA auto-negotiate enabled)\n");
#else
printf(" transport : TCP\n");
LOG_INFO(" transport : TCP\n");
#endif // GGML_RPC_RDMA
if (!rpc_transport_init()) {
fprintf(stderr, "Failed to initialize RPC transport\n");
LOG_ERROR("Failed to initialize RPC transport\n");
return;
}
auto server_socket = socket_t::create_server(host.c_str(), port);
if (server_socket == nullptr) {
fprintf(stderr, "Failed to create server socket\n");
LOG_ERROR("Failed to create server socket\n");
return;
}
while (true) {
auto client_socket = server_socket->accept();
if (client_socket == nullptr) {
fprintf(stderr, "Failed to accept client connection\n");
LOG_ERROR("Failed to accept client connection\n");
return;
}
printf("Accepted client connection\n");
fflush(stdout);
LOG_INFO("Accepted client connection\n");
rpc_serve_client(backends, cache_dir, client_socket);
printf("Client connection closed\n");
fflush(stdout);
LOG_INFO("Client connection closed\n");
}
rpc_transport_shutdown();
for (auto backend : backends) {
@@ -2341,9 +2407,11 @@ ggml_backend_reg_t ggml_backend_rpc_add_server(const char * endpoint) {
static uint32_t dev_id = 0;
std::lock_guard<std::mutex> lock(mutex);
if (reg_map.find(endpoint) != reg_map.end()) {
LOG_DBG3("[%s] endpoint %s already registered\n", __func__, endpoint);
return reg_map[endpoint];
}
uint32_t dev_count = ggml_backend_rpc_get_device_count(endpoint);
LOG_DBG2("[%s] endpoint: %s, device count: %u\n", __func__, endpoint, dev_count);
if (dev_count == 0) {
return nullptr;
}
+47
View File
@@ -0,0 +1,47 @@
#pragma once
// Log infrastructure for the RPC backend: all RPC logging goes through this header.
//
// LOG_ERROR / LOG_WARN / LOG_INFO are unconditional severity logs.
// LOG_DBG* are debug logs gated by the GGML_RPC_DEBUG verbosity variable:
// unset / 0 - disabled (only warnings and errors are printed)
// 1 - high-level events: connections, handshake, buffer operations, tensor transfers, graph computes
// 2 - per-command trace: every RPC message sent/received, cache hits/misses, queue events
// 3 - transport detail: byte counts, per-command timings, graph node details
// non-numeric values fall back to 1
#include "ggml-impl.h"
#include <cstdlib>
static inline int rpc_debug_level() {
static const int level = [] {
const char * env = std::getenv("GGML_RPC_DEBUG");
if (env == nullptr) {
return 0;
}
if (env[0] == '\0') {
return 1;
}
int value = 0;
for (const char * p = env; *p != '\0'; p++) {
if (*p < '0' || *p > '9') {
return 1;
}
value = value * 10 + (*p - '0');
if (value > 3) {
value = 3;
}
}
return value;
}();
return level;
}
#define LOG_INFO(...) GGML_LOG_INFO(__VA_ARGS__)
#define LOG_WARN(...) GGML_LOG_WARN(__VA_ARGS__)
#define LOG_ERROR(...) GGML_LOG_ERROR(__VA_ARGS__)
#define LOG_DBG(...) do { if (rpc_debug_level() >= 1) GGML_LOG_DEBUG(__VA_ARGS__); } while (0)
#define LOG_DBG2(...) do { if (rpc_debug_level() >= 2) GGML_LOG_DEBUG(__VA_ARGS__); } while (0)
#define LOG_DBG3(...) do { if (rpc_debug_level() >= 3) GGML_LOG_DEBUG(__VA_ARGS__); } while (0)
+15 -11
View File
@@ -1,6 +1,6 @@
#include "transport-apple.h"
#include "transport.h"
#include "ggml-impl.h"
#include "log.h"
#include <infiniband/verbs.h>
@@ -199,6 +199,7 @@ static bool rdma_library_present() {
// address, i.e. the one cabled to the peer.
std::unique_ptr<apple_rdma> apple_rdma::probe(int fd, const uint8_t * target_gid, uint8_t * caps) {
if (!rdma_library_present()) {
LOG_DBG2("[RDMA(Apple)] librdma.dylib not present, continuing with TCP\n");
return nullptr;
}
int ndev = 0;
@@ -224,7 +225,10 @@ std::unique_ptr<apple_rdma> apple_rdma::probe(int fd, const uint8_t * target_gid
break;
}
ibv_free_device_list(devs);
if (!ctx) return nullptr;
if (!ctx) {
LOG_DBG2("[RDMA(Apple)] no RDMA device matched the connection address, continuing with TCP\n");
return nullptr;
}
std::unique_ptr<impl> c(new impl());
c->fd = fd;
@@ -286,7 +290,7 @@ std::unique_ptr<apple_rdma> apple_rdma::probe(int fd, const uint8_t * target_gid
memcpy(rc.gid, gid.raw, RDMA_GID_SIZE);
memcpy(caps, &rc, sizeof(rc));
GGML_LOG_INFO("RDMA(Apple/UC) probed: dev=%s port=%u gid=%d qpn=%u lid=%u mtu=%d ring=%d x %zu KiB\n",
LOG_INFO("RDMA(Apple/UC) probed: dev=%s port=%u gid=%d qpn=%u lid=%u mtu=%d ring=%d x %zu KiB\n",
matched.c_str(), port, gid_idx, c->qpn, (unsigned)pa.lid, 128 << c->path_mtu,
RDMA_NBUF, RDMA_STRIDE / 1024);
return std::unique_ptr<apple_rdma>(new apple_rdma(std::move(c)));
@@ -317,7 +321,7 @@ bool apple_rdma::activate(const uint8_t * caps) {
memcpy(&a.ah_attr.grh.dgid, rc.gid, RDMA_GID_SIZE);
if (ibv_modify_qp(c->qp, &a,
IBV_QP_STATE | IBV_QP_AV | IBV_QP_PATH_MTU | IBV_QP_DEST_QPN | IBV_QP_RQ_PSN) != 0) {
GGML_LOG_ERROR("RDMA(Apple/UC) RTR failed: %s\n", strerror(errno));
LOG_ERROR("RDMA(Apple/UC) RTR failed: %s\n", strerror(errno));
ok = false;
}
}
@@ -326,7 +330,7 @@ bool apple_rdma::activate(const uint8_t * caps) {
a.qp_state = IBV_QPS_RTS;
a.sq_psn = RDMA_PSN;
if (ibv_modify_qp(c->qp, &a, IBV_QP_STATE | IBV_QP_SQ_PSN) != 0) {
GGML_LOG_ERROR("RDMA(Apple/UC) RTS failed: %s\n", strerror(errno));
LOG_ERROR("RDMA(Apple/UC) RTS failed: %s\n", strerror(errno));
ok = false;
}
}
@@ -334,7 +338,7 @@ bool apple_rdma::activate(const uint8_t * caps) {
// Recvs are posted only now: the controller starts processing them at RTR.
for (int i = 0; ok && i < RDMA_NBUF; i++) {
if (!c->post_recv(i)) {
GGML_LOG_ERROR("RDMA(Apple/UC) post_recv %d/%d failed\n", i, RDMA_NBUF);
LOG_ERROR("RDMA(Apple/UC) post_recv %d/%d failed\n", i, RDMA_NBUF);
ok = false;
}
}
@@ -350,7 +354,7 @@ bool apple_rdma::activate(const uint8_t * caps) {
return false;
}
GGML_LOG_INFO("RDMA(Apple/UC) activated: qpn=%u->%u mtu=%d rx_depth=%d\n",
LOG_INFO("RDMA(Apple/UC) activated: qpn=%u->%u mtu=%d rx_depth=%d\n",
c->qpn, rc.qpn, 128 << c->path_mtu, RDMA_NBUF);
return true;
}
@@ -360,20 +364,20 @@ bool apple_rdma::activate(const uint8_t * caps) {
int apple_rdma::impl::progress() {
struct ibv_wc wc[RDMA_NBUF * 2];
int n = ibv_poll_cq(cq, RDMA_NBUF * 2, wc);
if (n < 0) { GGML_LOG_ERROR("RDMA(Apple/UC) poll_cq failed\n"); broken = true; return -1; }
if (n < 0) { LOG_ERROR("RDMA(Apple/UC) poll_cq failed\n"); broken = true; return -1; }
for (int j = 0; j < n; j++) {
uint64_t id = wc[j].wr_id;
bool is_recv = (id & RDMA_RECV_WR) != 0;
if (wc[j].status != IBV_WC_SUCCESS) {
GGML_LOG_ERROR("RDMA(Apple/UC) %s wc error: status=%d\n", is_recv ? "recv" : "send", wc[j].status);
LOG_ERROR("RDMA(Apple/UC) %s wc error: status=%d\n", is_recv ? "recv" : "send", wc[j].status);
broken = true;
return -1;
}
if (is_recv) {
int b = (int)(id & RDMA_WR_IDX_MASK);
const rdma_seg_hdr * h = (const rdma_seg_hdr *)(recv_mem + (size_t)b * RDMA_STRIDE);
if (h->magic != RDMA_SEG_MAGIC) { GGML_LOG_ERROR("RDMA(Apple/UC) bad frame magic\n"); broken = true; return -1; }
if (h->len > RDMA_PAYLOAD) { GGML_LOG_ERROR("RDMA(Apple/UC) frame len %u exceeds payload\n", h->len); broken = true; return -1; }
if (h->magic != RDMA_SEG_MAGIC) { LOG_ERROR("RDMA(Apple/UC) bad frame magic\n"); broken = true; return -1; }
if (h->len > RDMA_PAYLOAD) { LOG_ERROR("RDMA(Apple/UC) frame len %u exceeds payload\n", h->len); broken = true; return -1; }
int slot = (inq_head + inq_count) % RDMA_NBUF;
inq[slot].buf = b;
inq[slot].off = 0;
+28 -17
View File
@@ -1,5 +1,5 @@
#include "transport.h"
#include "ggml-impl.h"
#include "log.h"
#ifdef _WIN32
# define WIN32_LEAN_AND_MEAN
@@ -43,11 +43,6 @@ using ssize_t = __int64;
typedef int sockfd_t;
#endif
static const char * RPC_DEBUG = std::getenv("GGML_RPC_DEBUG");
#define LOG_DBG(...) \
do { if (RPC_DEBUG) GGML_LOG_DEBUG(__VA_ARGS__); } while (0)
#ifdef GGML_RPC_RDMA
static constexpr size_t RDMA_GID_SIZE = 16; // RoCE GID / IB GID is always 16 bytes
using rdma_gid_t = std::array<uint8_t, RDMA_GID_SIZE>;
@@ -342,7 +337,7 @@ bool socket_t::impl::rdma_probe() {
} else if (gid_version == IBV_GID_TYPE_ROCE_V1) {
ver_str = " RoCEv1";
}
GGML_LOG_INFO("RDMA probed: dev=%s gid=%d%s qpn=%u inline=%u\n",
LOG_INFO("RDMA probed: dev=%s gid=%d%s qpn=%u inline=%u\n",
matched_dev, gid_idx, ver_str, rdma_local.qpn, rdma->max_inline);
return true;
}
@@ -407,7 +402,7 @@ bool socket_t::impl::rdma_activate(uint32_t remote_qpn, uint32_t remote_psn, con
rdma->last_active = std::chrono::steady_clock::now();
GGML_LOG_INFO("RDMA activated: qpn=%u->%u mtu=%d rx_depth=%d\n",
LOG_INFO("RDMA activated: qpn=%u->%u mtu=%d rx_depth=%d\n",
rdma_local.qpn, remote_qpn, 128 << rdma_local.path_mtu, RDMA_RX_DEPTH);
return true;
}
@@ -445,7 +440,7 @@ bool socket_t::impl::rdma_poll(struct ibv_cq * cq, struct ibv_wc * wc) {
if (n > 0) {
c->last_active = std::chrono::steady_clock::now();
if (wc->status != IBV_WC_SUCCESS) {
GGML_LOG_ERROR("RDMA CQ wc error: status=%d (%s) vendor_err=0x%x\n",
LOG_ERROR("RDMA CQ wc error: status=%d (%s) vendor_err=0x%x\n",
wc->status, ibv_wc_status_str(wc->status), wc->vendor_err);
}
return wc->status == IBV_WC_SUCCESS;
@@ -537,6 +532,7 @@ bool socket_t::impl::rdma_recv(void * data, size_t size) {
#endif // GGML_RPC_RDMA
bool socket_t::impl::send_data(const void * data, size_t size) {
LOG_DBG3("[%s] transport: %s, size: %zu\n", __func__, use_rdma ? "RDMA" : "TCP", size);
#ifdef GGML_RPC_RDMA_APPLE
if (use_rdma) {
return rdma->send(data, size);
@@ -551,7 +547,7 @@ bool socket_t::impl::send_data(const void * data, size_t size) {
size_t size_to_send = std::min(size - bytes_sent, MAX_CHUNK_SIZE);
ssize_t n = send(fd, (const char *)data + bytes_sent, size_to_send, 0);
if (n < 0) {
GGML_LOG_ERROR("send failed (bytes_sent=%zu, size_to_send=%zu)\n",
LOG_ERROR("send failed (bytes_sent=%zu, size_to_send=%zu)\n",
bytes_sent, size_to_send);
return false;
}
@@ -561,6 +557,7 @@ bool socket_t::impl::send_data(const void * data, size_t size) {
}
bool socket_t::impl::recv_data(void * data, size_t size) {
LOG_DBG3("[%s] transport: %s, size: %zu\n", __func__, use_rdma ? "RDMA" : "TCP", size);
#ifdef GGML_RPC_RDMA_APPLE
if (use_rdma) {
return rdma->recv(data, size);
@@ -575,7 +572,7 @@ bool socket_t::impl::recv_data(void * data, size_t size) {
size_t size_to_recv = std::min(size - bytes_recv, MAX_CHUNK_SIZE);
ssize_t n = recv(fd, (char *)data + bytes_recv, size_to_recv, 0);
if (n < 0) {
GGML_LOG_ERROR("recv failed (bytes_recv=%zu, size_to_recv=%zu)\n",
LOG_ERROR("recv failed (bytes_recv=%zu, size_to_recv=%zu)\n",
bytes_recv, size_to_recv);
return false;
}
@@ -599,6 +596,9 @@ void socket_t::impl::get_caps(uint8_t * local_caps) {
if (target_gid) {
rdma = apple_rdma::probe(fd, target_gid->data(), local_caps);
}
if (!rdma) {
LOG_DBG2("[%s] RDMA probe failed, continuing with TCP\n", __func__);
}
# else
rdma_local = {};
if (rdma_probe()) {
@@ -609,6 +609,7 @@ void socket_t::impl::get_caps(uint8_t * local_caps) {
memcpy(local_caps, &rc, sizeof(rc));
} else {
rdma.reset();
LOG_DBG2("[%s] RDMA probe failed, continuing with TCP\n", __func__);
}
# endif
#endif // GGML_RPC_RDMA
@@ -623,6 +624,9 @@ void socket_t::impl::update_caps(const uint8_t * remote_caps) {
remote_rdma |= remote_caps[i] != 0;
}
if (!rdma || !remote_rdma) {
if (rdma && !remote_rdma) {
LOG_DBG("[%s] peer does not support RDMA, continuing with TCP\n", __func__);
}
rdma.reset();
return;
}
@@ -636,7 +640,7 @@ void socket_t::impl::update_caps(const uint8_t * remote_caps) {
if (activated) {
use_rdma = true;
} else {
GGML_LOG_ERROR("RDMA activate failed, staying on TCP\n");
LOG_ERROR("RDMA activate failed, staying on TCP\n");
rdma.reset();
}
#else
@@ -679,6 +683,10 @@ void socket_t::update_caps(const uint8_t * remote_caps) {
return pimpl->update_caps(remote_caps);
}
const char * socket_t::transport_name() const {
return pimpl->use_rdma ? "RDMA" : "TCP";
}
static bool is_valid_fd(sockfd_t sockfd) {
#ifdef _WIN32
return sockfd != INVALID_SOCKET;
@@ -706,9 +714,10 @@ socket_ptr socket_t::accept() {
return nullptr;
}
if (!set_no_delay(client_socket_fd)) {
GGML_LOG_ERROR("Failed to set TCP_NODELAY\n");
LOG_ERROR("Failed to set TCP_NODELAY\n");
return nullptr;
}
LOG_DBG("[%s] accepted client connection\n", __func__);
return socket_ptr(new socket_t(std::make_unique<impl>(client_socket_fd)));
}
@@ -718,11 +727,11 @@ socket_ptr socket_t::create_server(const char * host, int port) {
return nullptr;
}
if (!set_reuse_addr(sockfd)) {
GGML_LOG_ERROR("Failed to set SO_REUSEADDR\n");
LOG_ERROR("Failed to set SO_REUSEADDR\n");
return nullptr;
}
if (inet_addr(host) == INADDR_NONE) {
GGML_LOG_ERROR("Invalid host address: %s\n", host);
LOG_ERROR("Invalid host address: %s\n", host);
return nullptr;
}
struct sockaddr_in serv_addr;
@@ -736,6 +745,7 @@ socket_ptr socket_t::create_server(const char * host, int port) {
if (listen(sockfd, 1) < 0) {
return nullptr;
}
LOG_DBG("[%s] listening on %s:%d\n", __func__, host, port);
return socket_ptr(new socket_t(std::make_unique<impl>(sockfd)));
}
@@ -745,7 +755,7 @@ socket_ptr socket_t::connect(const char * host, int port) {
return nullptr;
}
if (!set_no_delay(sockfd)) {
GGML_LOG_ERROR("Failed to set TCP_NODELAY\n");
LOG_ERROR("Failed to set TCP_NODELAY\n");
return nullptr;
}
struct sockaddr_in addr;
@@ -753,13 +763,14 @@ socket_ptr socket_t::connect(const char * host, int port) {
addr.sin_port = htons(port);
struct hostent * server = gethostbyname(host);
if (server == NULL) {
GGML_LOG_ERROR("Cannot resolve host '%s'\n", host);
LOG_ERROR("Cannot resolve host '%s'\n", host);
return nullptr;
}
memcpy(&addr.sin_addr.s_addr, server->h_addr, server->h_length);
if (::connect(sockfd, (struct sockaddr *)&addr, sizeof(addr)) < 0) {
return nullptr;
}
LOG_DBG("[%s] connected to %s:%d\n", __func__, host, port);
return socket_ptr(new socket_t(std::make_unique<impl>(sockfd)));
}
+2
View File
@@ -24,6 +24,8 @@ struct socket_t {
void get_caps(uint8_t * local_caps);
void update_caps(const uint8_t * remote_caps);
// transport in use after the HELLO negotiation: "TCP" or "RDMA"
const char * transport_name() const;
static socket_ptr create_server(const char * host, int port);
static socket_ptr connect(const char * host, int port);
+11 -1
View File
@@ -113,8 +113,18 @@ $ GGML_RPC_NO_RDMA=1 bin/ggml-rpc-server
### Troubleshooting
Use the `GGML_RPC_DEBUG` environment variable to enable debug messages from `ggml-rpc-server`:
The `GGML_RPC_DEBUG` environment variable controls the verbosity of the logs emitted by the RPC backend.
It can be set on the server, on the client (e.g. `llama-cli`), or both. Larger values produce more detailed output:
- unset / `0` - disabled (only warnings and errors are printed)
- `1` - high-level events: connections, handshake, buffer operations, tensor transfers, graph computes
- `2` - per-command trace: every RPC message sent/received, cache hits/misses, queue events
- `3` - transport detail: byte counts, per-command timings, graph node details
Non-numeric values are treated as `1`.
```bash
$ GGML_RPC_DEBUG=1 bin/ggml-rpc-server
$ GGML_RPC_DEBUG=2 bin/llama-cli -hf ggml-org/gemma-3-1b-it-GGUF -ngl 99 --rpc 192.168.88.10:50052
```