diff --git a/ggml/src/ggml-rpc/ggml-rpc.cpp b/ggml/src/ggml-rpc/ggml-rpc.cpp index 353b79b072..e9f88c22f1 100644 --- a/ggml/src/ggml-rpc/ggml-rpc.cpp +++ b/ggml/src/ggml-rpc/ggml-rpc.cpp @@ -3,8 +3,10 @@ #include "ggml-backend-impl.h" #include "ggml-cpp.h" #include "transport.h" +#include "log.h" #include +#include #include #include #include @@ -23,12 +25,6 @@ #include #include -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 & 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 & 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 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 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 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 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 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(); 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(); 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(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(); 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(); 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(); 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(); 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(); 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 & 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 & 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 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(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(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(); 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(); 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 input_ptr(input, std::default_delete()); 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(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(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(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 & 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 & 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 & 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 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 & 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 & 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 & 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 & 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 & 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 & 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 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 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; } diff --git a/ggml/src/ggml-rpc/log.h b/ggml/src/ggml-rpc/log.h new file mode 100644 index 0000000000..1624b75f49 --- /dev/null +++ b/ggml/src/ggml-rpc/log.h @@ -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 + +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) diff --git a/ggml/src/ggml-rpc/transport-apple.cpp b/ggml/src/ggml-rpc/transport-apple.cpp index bb24a5d4c1..18c4e7efc7 100644 --- a/ggml/src/ggml-rpc/transport-apple.cpp +++ b/ggml/src/ggml-rpc/transport-apple.cpp @@ -1,6 +1,6 @@ #include "transport-apple.h" #include "transport.h" -#include "ggml-impl.h" +#include "log.h" #include @@ -199,6 +199,7 @@ static bool rdma_library_present() { // address, i.e. the one cabled to the peer. std::unique_ptr 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::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 c(new impl()); c->fd = fd; @@ -286,7 +290,7 @@ std::unique_ptr 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(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; diff --git a/ggml/src/ggml-rpc/transport.cpp b/ggml/src/ggml-rpc/transport.cpp index b28d16605c..339eff40bb 100644 --- a/ggml/src/ggml-rpc/transport.cpp +++ b/ggml/src/ggml-rpc/transport.cpp @@ -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; @@ -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(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(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(sockfd))); } diff --git a/ggml/src/ggml-rpc/transport.h b/ggml/src/ggml-rpc/transport.h index 3f747ecffd..ea8c94d020 100644 --- a/ggml/src/ggml-rpc/transport.h +++ b/ggml/src/ggml-rpc/transport.h @@ -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); diff --git a/tools/rpc/README.md b/tools/rpc/README.md index fc51568947..c38ce84d1d 100644 --- a/tools/rpc/README.md +++ b/tools/rpc/README.md @@ -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 ```