diff --git a/.github/smoke-rpc-fatjar.sh b/.github/smoke-rpc-fatjar.sh new file mode 100755 index 000000000..bd322b376 --- /dev/null +++ b/.github/smoke-rpc-fatjar.sh @@ -0,0 +1,119 @@ +#!/usr/bin/env bash + +# SPDX-FileCopyrightText: 2026 Bernard Ladenthin +# +# SPDX-License-Identifier: MIT + +# RPC smoke test over two JVMs, on the real release asset: +# +# JVM A java -cp net.ladenthin.llama.RpcServer serves this runner's devices +# JVM B java -jar -m --rpc 127.0.0.1: the default NativeServer, which +# offloads its layers to A +# +# and checks that B answers a chat completion, that its load log shows a model buffer on A's +# endpoint (the layers really went over RPC, not silently to the CPU), and that A accepted a +# client. A third launch names a server nobody runs and must fail with a message naming it and a +# normal exit -- not a SIGABRT, which is what ggml-rpc did before patches/0015. +# +# Usage: smoke-rpc-fatjar.sh +# Output lands in rpc-server.log, rpc-client-out.log, rpc-client-err.log, rpc-unreachable.log +# (uploaded by the CI job on failure). +set -euo pipefail + +JAR_DIR="${1:?usage: smoke-rpc-fatjar.sh }" +JAR_GLOB="${2:?usage: smoke-rpc-fatjar.sh }" +MODEL="${3:?usage: smoke-rpc-fatjar.sh }" +RPC_PORT="${RPC_PORT:-50152}" +HTTP_PORT="${HTTP_PORT:-18181}" +UNUSED_PORT="${UNUSED_PORT:-50153}" + +fail() { + echo "::error::$*" >&2 + for f in rpc-server.log rpc-client-out.log rpc-client-err.log rpc-unreachable.log; do + [ -f "$f" ] && { echo "--- $f (tail) ---"; tail -60 "$f"; } + done + exit 1 +} + +mapfile -t JARS < <(find "$JAR_DIR" -maxdepth 1 -name "$JAR_GLOB" | sort) +[ "${#JARS[@]}" -eq 1 ] || fail "expected exactly 1 jar matching $JAR_GLOB in $JAR_DIR, got ${#JARS[@]}: ${JARS[*]:-none}" +JAR="${JARS[0]}" +[ -f "$MODEL" ] || fail "model file missing: $MODEL" + +PIDS=() +cleanup() { + for pid in "${PIDS[@]}"; do + kill "$pid" 2> /dev/null || true + done +} +trap cleanup EXIT + +# --- JVM A: the RPC server ---------------------------------------------------------------------- +java -cp "$JAR" net.ladenthin.llama.RpcServer --port "$RPC_PORT" --threads 2 --device CPU > rpc-server.log 2>&1 & +PIDS+=($!) +SERVER_PID=$! +# Up to 300 s, like the NativeServer smoke: an all-backends jar extracts every GPU backend's +# library (CUDA, ROCm and SYCL are hundreds of MB) and fails to load each before it reaches one +# this GPU-less runner can load. 60 s was too short for that on the first CI run. +for _ in $(seq 1 100); do + kill -0 "$SERVER_PID" 2> /dev/null || fail "RpcServer exited before listening" + grep -q "RpcServer listening on 127.0.0.1:$RPC_PORT" rpc-server.log && break + sleep 3 +done +grep -q "RpcServer listening on 127.0.0.1:$RPC_PORT" rpc-server.log || fail "RpcServer never reported listening" +grep -q "serving \[CPU\]" rpc-server.log || fail "RpcServer --device CPU did not serve exactly the CPU" +# The backend is chosen once. RpcServer is the one entry point that reaches the loader before +# LlamaModel, and JNI_OnLoad initializes LlamaModel, whose static block re-entered the loader and +# ran a second complete load over the library being loaded (LlamaLoader.runOnceOnThisThread). +# A manifest-less jar prints the line zero times. +[ "$(grep -c '\[jllama\] using native backend' rpc-server.log)" -le 1 ] \ + || fail "the native library was loaded more than once: $(grep -c '\[jllama\] using native backend' rpc-server.log) backend selections" +echo "RPC server up: $(grep 'RpcServer listening' rpc-server.log)" + +# --- JVM B: the model, offloaded over RPC -------------------------------------------------------- +java -jar "$JAR" -m "$MODEL" --host 127.0.0.1 --port "$HTTP_PORT" --chat-template chatml \ + --rpc "127.0.0.1:$RPC_PORT" -ngl 99 -lv 4 > rpc-client-out.log 2> rpc-client-err.log & +PIDS+=($!) +CLIENT_PID=$! +CODE="" +for _ in $(seq 1 100); do + kill -0 "$CLIENT_PID" 2> /dev/null || fail "the RPC client server exited before becoming healthy" + CODE="$(curl -s -o /dev/null -w '%{http_code}' "http://127.0.0.1:$HTTP_PORT/health" || true)" + [ "$CODE" = "200" ] && break + sleep 3 +done +[ "$CODE" = "200" ] || fail "/health never returned 200 (last code: ${CODE:-none})" + +RESPONSE="$(curl -sS --fail -X POST "http://127.0.0.1:$HTTP_PORT/v1/chat/completions" \ + -H 'Content-Type: application/json' \ + -d '{"messages":[{"role":"user","content":"Say hello."}],"max_tokens":8,"temperature":0}')" \ + || fail "chat completion over RPC failed" +echo "$RESPONSE" | python3 -c ' +import json, sys +message = json.load(sys.stdin)["choices"][0]["message"] +assert message is not None, "choices[0].message missing" +print("chat completion over RPC OK:", json.dumps(message)[:200]) +' || fail "malformed chat completion response: $RESPONSE" + +grep -h "model buffer size" rpc-client-out.log rpc-client-err.log | grep -q "127.0.0.1:$RPC_PORT" \ + || fail "no model buffer on the RPC server in the load log -- the layers did not go over RPC" +grep -q "Accepted client connection" rpc-server.log || fail "the RPC server never accepted a client" +echo "layers offloaded over RPC: $(grep -h 'model buffer size' rpc-client-out.log rpc-client-err.log | grep "127.0.0.1:$RPC_PORT" | head -1)" + +kill "$CLIENT_PID" 2> /dev/null || true +wait "$CLIENT_PID" 2> /dev/null || true + +# --- an unreachable server fails the start cleanly ----------------------------------------------- +set +e +timeout 120 java -jar "$JAR" -m "$MODEL" --host 127.0.0.1 --port "$((HTTP_PORT + 1))" \ + --rpc "127.0.0.1:$UNUSED_PORT" > rpc-unreachable.log 2>&1 +status=$? +set -e +[ "$status" -ne 0 ] || fail "a server naming an unreachable RPC endpoint started anyway" +[ "$status" -ne 124 ] || fail "a server naming an unreachable RPC endpoint hung instead of failing" +# 134 = SIGABRT: the GGML_ABORT patches/0015 removed from the registration path +[ "$status" -ne 134 ] || fail "an unreachable RPC endpoint aborted the JVM (exit 134)" +grep -q "127.0.0.1:$UNUSED_PORT" rpc-unreachable.log || fail "the failure does not name the unreachable endpoint" +echo "unreachable endpoint rejected (exit $status): $(grep -m1 "127.0.0.1:$UNUSED_PORT" rpc-unreachable.log)" + +echo "RPC smoke test PASSED" diff --git a/.github/verify-native-deps.py b/.github/verify-native-deps.py new file mode 100755 index 000000000..ab26f3812 --- /dev/null +++ b/.github/verify-native-deps.py @@ -0,0 +1,211 @@ +#!/usr/bin/env python3 +# SPDX-FileCopyrightText: 2026 Bernard Ladenthin +# +# SPDX-License-Identifier: MIT +"""Fail when a shipped native library needs a runtime library it did not need before. + +Every jllama library is ONE file with llama.cpp and ggml linked in statically, so its dynamic +dependencies are exactly what a consumer's machine must provide. A new one is a silent break on +every machine that lacks it -- the case this guards against is ggml-rpc's RDMA transport, which +upstream switches on whenever the build host has libibverbs/librdma and which would make the +library unloadable without rdma-core. It reads the dependency list straight from the file (ELF +DT_NEEDED, PE import table, Mach-O LC_LOAD_DYLIB) with the standard library only, so it runs on +any runner and checks every architecture, including the ones binutils cannot read (Windows arm64, +Mach-O). + +Usage: + verify-native-deps.py --default # exact allowlist per / + verify-native-deps.py --deny ... # classifier trees: only the denylist + +--default checks net/ladenthin/llama/// under the root against ALLOWED below: a +dependency outside the list fails, and so does an / without a list (a new platform +must be listed consciously). --deny scans every native library under the directories for DENIED +names only, because GPU classifiers legitimately need their vendor runtime. + +Exit codes: 0 clean, 1 violation, 2 nothing found to check. +""" + +import os +import struct +import sys + +# What each default-JAR library needed when this check was introduced (5.1.0 plus the RPC backend, +# which adds nothing: its sockets are libc/libSystem/WS2_32, all already present). +ALLOWED = { + "Linux/x86_64": {"libdl.so.2", "libgomp.so.1", "libpthread.so.0", "librt.so.1", "libstdc++.so.6", + "libm.so.6", "libgcc_s.so.1", "libc.so.6", "ld-linux-x86-64.so.2"}, + "Linux/aarch64": {"libgomp.so.1", "libstdc++.so.6", "libm.so.6", "libgcc_s.so.1", "libc.so.6", + "ld-linux-aarch64.so.1"}, + "Linux/s390x": {"libstdc++.so.6", "libm.so.6", "libgcc_s.so.1", "libc.so.6", "ld64.so.1"}, + "Linux-Android/aarch64": {"liblog.so", "libm.so", "libdl.so", "libc.so", "libandroid.so"}, + "Linux-Android/x86_64": {"liblog.so", "libm.so", "libdl.so", "libc.so", "libandroid.so"}, + "Windows/x86_64": {"ws2_32.dll", "kernel32.dll", "shell32.dll", "advapi32.dll", "vcomp140.dll"}, + "Windows/x86": {"ws2_32.dll", "kernel32.dll", "shell32.dll", "advapi32.dll", "vcomp140.dll"}, + "Windows/aarch64": {"ws2_32.dll", "kernel32.dll", "shell32.dll", "advapi32.dll"}, + "Mac/aarch64": {"/usr/lib/libc++.1.dylib", "/usr/lib/libSystem.B.dylib", + "/System/Library/Frameworks/Foundation.framework/Versions/C/Foundation", + "/System/Library/Frameworks/Metal.framework/Versions/A/Metal", + "/System/Library/Frameworks/MetalKit.framework/Versions/A/MetalKit", + "/System/Library/Frameworks/Accelerate.framework/Versions/A/Accelerate", + "/usr/lib/libobjc.A.dylib", + "/System/Library/Frameworks/CoreFoundation.framework/Versions/A/CoreFoundation", + "/System/Library/Frameworks/Security.framework/Versions/A/Security", + # KNOWN DEFECT, allowed only so this check reports NEW dependencies: the macOS + # build picks up the runner's Homebrew OpenSSL, so the shipped dylib does not load + # on a Mac without `brew install openssl@3`. See TODO.md ("macOS dylib links + # Homebrew OpenSSL"); remove these two lines with the fix. + "/opt/homebrew/opt/openssl@3/lib/libssl.3.dylib", + "/opt/homebrew/opt/openssl@3/lib/libcrypto.3.dylib"}, +} + +# Never acceptable in any artifact: libraries a consumer cannot be expected to have. +DENIED = ("libibverbs", "librdma", "rdma.dylib", "libmlx") + +LIB_NAMES = ("libjllama.so", "jllama.dll", "libjllama.dylib") + + +def elf_needed(data): + if data[:4] != b"\x7fELF": + raise ValueError("not an ELF file") + is64 = data[4] == 2 + end = "<" if data[5] == 1 else ">" + if is64: + shoff = struct.unpack_from(end + "Q", data, 0x28)[0] + shentsize, shnum = struct.unpack_from(end + "HH", data, 0x3A) + else: + shoff = struct.unpack_from(end + "I", data, 0x20)[0] + shentsize, shnum = struct.unpack_from(end + "HH", data, 0x2E) + sections = [] + for i in range(shnum): + off = shoff + i * shentsize + if is64: + _, sh_type, _, _, sh_offset, sh_size, sh_link = struct.unpack_from(end + "IIQQQQI", data, off) + else: + _, sh_type, _, _, sh_offset, sh_size, sh_link = struct.unpack_from(end + "IIIIIII", data, off) + sections.append((sh_type, sh_offset, sh_size, sh_link)) + out = [] + for sh_type, sh_offset, sh_size, sh_link in sections: + if sh_type != 6: # SHT_DYNAMIC + continue + strtab = sections[sh_link] + entry = 16 if is64 else 8 + for off in range(sh_offset, sh_offset + sh_size, entry): + tag, val = struct.unpack_from(end + ("qQ" if is64 else "iI"), data, off) + if tag == 0: + break + if tag == 1: # DT_NEEDED + start = strtab[1] + val + out.append(data[start:data.index(b"\0", start)].decode()) + return out + + +def pe_imports(data): + if data[:2] != b"MZ": + raise ValueError("not a PE file") + pe = struct.unpack_from("= 3 else "" + allowed = ALLOWED.get(key) + if allowed is None: + failures.append(f"{rel}: no dependency allowlist for '{key}' -- add one to ALLOWED") + continue + for d in deps: + if d.lower() not in {a.lower() for a in allowed}: + failures.append(f"{rel} needs {d}, which it did not need before (allowed: {sorted(allowed)})") + if checked == 0: + print(f"no native library found under {roots}", file=sys.stderr) + return 2 + for f in failures: + print(f"::error::{f}", file=sys.stderr) + print(f"{checked} native libraries checked, {len(failures)} violations") + return 1 if failures else 0 + + +if __name__ == "__main__": + sys.exit(main(sys.argv)) diff --git a/.github/workflows/publish.yml b/.github/workflows/publish.yml index be2197271..9342bbe1a 100644 --- a/.github/workflows/publish.yml +++ b/.github/workflows/publish.yml @@ -3544,6 +3544,17 @@ jobs: with: name: Windows-x86_64-openvino path: ${{ github.workspace }}/llama/src/main/resources_windows_openvino/net/ladenthin/llama/ + # Runtime dependencies of every shipped native library, read from the file itself (ELF + # DT_NEEDED, PE imports incl. Windows arm64, Mach-O load commands). The default tree must + # match an exact per-{OS}/{ARCH} allowlist -- a new dependency is a load failure on every + # machine that lacks it. The classifier trees legitimately need their vendor runtime, so + # there only the denylist applies. The case this was written for: ggml-rpc's RDMA + # transport, which upstream enables whenever the build host has libibverbs/librdma + # (llama/CMakeLists.txt forces it off). + - name: Verify native runtime dependencies + run: | + python3 .github/verify-native-deps.py --default llama/src/main/resources + python3 .github/verify-native-deps.py --deny llama/src/main/resources_* - uses: actions/setup-java@v6 with: distribution: 'temurin' @@ -3679,6 +3690,12 @@ jobs: run: .github/verify-bytecode-version.sh --max-major 52 fatjars - name: Run fat-jar server smoke test run: .github/smoke-test-fatjar.sh fatjars 'llama-*-all-linux-x86-64-jar-with-dependencies.jar' "models/${DRAFT_MODEL_NAME}" + # RPC over two JVMs from the same release asset: RpcServer in one, the default NativeServer + # with --rpc in the other. Requires a chat completion, the model buffer on the RPC endpoint in + # the load log (the layers really went over RPC), an accepted client on the server, and a clean + # non-SIGABRT failure naming the endpoint when --rpc names a server nobody runs. + - name: Run fat-jar RPC smoke test (two JVMs) + run: .github/smoke-rpc-fatjar.sh fatjars 'llama-*-all-linux-x86-64-jar-with-dependencies.jar' "models/${DRAFT_MODEL_NAME}" - name: Upload server logs if: failure() uses: actions/upload-artifact@v7 @@ -3687,6 +3704,10 @@ jobs: path: | server-out.log server-err.log + rpc-server.log + rpc-client-out.log + rpc-client-err.log + rpc-unreachable.log if-no-files-found: warn # The agent release asset, launched the way the README tells a user to: `java -jar` on the agent diff --git a/CHANGELOG.md b/CHANGELOG.md index 837535706..826ab3926 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -39,6 +39,28 @@ from version 5.0.0 onward. Pre-fork releases (`1.x`–`4.2.0`) were authored by documented as the no-ops they are (`common_init()` forces both on), `setLogFile` as additive. ### Added +- **Distributed inference over llama.cpp RPC**, client and server, in every artifact with no new + runtime dependency. `ModelParameters.setRpcServers(RpcEndpoint...)` (and `--rpc host:port[,…]` on + both HTTP servers) offloads layers to RPC servers on other machines; `RpcServer` serves this + machine's devices (every GPU found, else the CPU), loopback-only unless `startOnNetwork` is used, + also runnable as `java -cp net.ladenthin.llama.RpcServer`. `value.RpcEndpoint` validates + endpoints (IPv4 or host name; the transport has no IPv6). Upstream aborts the whole process on + every client-side connection problem; the new `patches/0015` turns an unreachable, malformed or + non-RPC endpoint into a `LlamaException` naming it, makes the server stoppable, and keeps a + stopped server's registered device from aborting later loads. Because llama.cpp never forgets a + registered RPC server, a load that does not ask for one gets an explicit device list without it — + including the multimodal projector's device, and `TextToSpeech` / `LlamaTrainer` loads. + A server lost in the middle of inference still terminates the process (upstream limitation). + The served devices can be chosen by name (`RpcServer.startLocal(…, List devices)`, + `--device CPU`), because a served device that cannot run an operation terminates the server + (llama.cpp's RPC client reports every operation as supported). +- **`LlamaLoader` no longer loads the native library twice** when the first class to load it is not + `LlamaModel` (e.g. `TextToSpeech`, `LlamaQuantizer`, `RpcServer`): `JNI_OnLoad` initializes + `LlamaModel`, whose static block re-entered the loader and ran a second full load — with an + all-backends fat jar that meant re-extracting every GPU backend over the library being loaded. +- **`.github/verify-native-deps.py`** holds every shipped native library to its known runtime + dependencies (ELF, PE incl. Windows arm64, Mach-O). It found that the macOS dylib has always + linked Homebrew's `openssl@3` (see `TODO.md`). - **The agent is a release asset: `llama-atmosphere-agent--jar-with-dependencies.jar`**, with `.sha256` and a GPG `.asc`, on every GitHub release and the rolling `snapshot` pre-release — never on Maven Central. It carries **no core** (~7 MB instead of hundreds, natives not in the release twice): diff --git a/CLAUDE.md b/CLAUDE.md index e934ca4f2..204a9b2ec 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -784,6 +784,7 @@ Current patches: | `0006-server-embed-native-server-jni.patch` | **Makes `server.cpp`'s `llama_server` embeddable in the JVM** so the `NativeServer` JNI bridge can run the full upstream HTTP server (WebUI included) inside `libjllama` — see "Two server modes" below. b9870 already exposes `int llama_server(int, char**)` (non-static; no `main` in the file), so the patch only adds embedded-mode support: (1) a `g_llama_server_embedded` flag + `llama_server_set_embedded()` / `llama_server_request_shutdown()` (declared in the committed `src/main/cpp/native_server_bridge.h`); (2) skips installing the process-wide SIGINT/SIGTERM handlers when embedded (they would hijack the JVM's); (3) in embedded mode parses the **forwarded** argv via `common_params_parse` instead of `common_params_parse_main` (whose `GetCommandLineW` recovery would pick up `java.exe`'s command line — the same Windows class of bug `0001` fixes). `llama_server_request_shutdown()` mirrors the SIGTERM path (invokes the installed `shutdown_handler` → `ctx_server.terminate()` unblocks `start_loop()`), giving JNI an out-of-band stop since `ctx_server` is loop-local. Applies **after `0001`** (which flips this call site to `common_params_parse_main`), so its context is the post-`0001` tree; regenerate against `0001`+source on a bump. Only touches `tools/server/server.cpp`. | | `0012-model-guard-zero-split-sum-and-name-the-device-index.patch` | **A GPU that reports zero free memory makes every model load fail with the unactionable `error loading model: vector`.** `llama_model_base::load_tensors` (`src/llama-model.cpp`) weights the per-device layer split by `ggml_backend_dev_memory()`'s `free`, then normalises: `splits[i] /= split_sum`. With a single device reporting `free == 0` that is `0/0` → **NaN** in every split point; NaN compares false against everything, so the `std::upper_bound` below returns the end iterator, `layer_gpu == n_devices()`, and `devices.at(layer_gpu)` throws `std::out_of_range` — whose libc++ `what()` is the bare string `"vector"`, which `llama.cpp`'s `catch (const std::exception &)` prints verbatim. Upstream's `free == 0 && total == 0` host-memory fallback does **not** fire, because `total` is `recommendedMaxWorkingSetSize` and is non-zero. **Reachable since b10618..b10797**: upstream `8c0b9cd04` ("metal : fix memory query under low-memory conditions", [#27701](https://github.com/ggml-org/llama.cpp/pull/27701)) changed `ggml-metal-device.m` to `*free = *total > cur ? *total - cur : 0`; before that clamp an over-committed device (`currentAllocatedSize > recommendedMaxWorkingSetSize`) *underflowed* to a huge `size_t`, which normalised fine, so the same precondition was harmless. That is why the `Java Tests macOS …` jobs went red at the b10792→b10797 step while every Linux/Windows job stayed green — **and why only a GPU build can fail this way at all**: `act_gpu_layers` is `devices.empty() ? 0 : …`, so with no GPU backend `devices` is empty, every layer returns early on `cpu_dev`, and the `.at()` line is unreachable. **Shape:** the two blocks are lifted out of `load_tensors` into free functions declared in `src/llama-model.h`, purely so they can be driven by a test — the failing state needs a real over-committed GPU and cannot be arranged through any public API. `llama_model_splits_normalize()` carries **the fix**: on `split_sum == 0` it `LLAMA_LOG_WARN`s and falls back to an even split (`splits[i] = float(i+1)/splits.size()`), the only neutral choice when no device can be preferred and exactly right for a single device. `llama_model_splits_select_device()` carries **the diagnostic**: it bounds-checks the index and throws a `std::runtime_error` naming the function, the offloaded layer, the device index, the split-point count **and the split points themselves** — with NaN splits that message prints `nan` and names the cause outright, which is precisely what was missing when this had to be diagnosed by reading source. **A second, backend-independent trigger reaches the same line**, found while writing this up and verified against the unfixed library: `--tensor-split` values are parsed with `std::stof` and never range-checked (`common/arg.cpp`), so `-ts 1,-1` cancels out, `split_sum` is 0 again, the split points become `[inf, -nan]`, and every layer maps one past the last device — on CUDA, Vulkan or ROCm just as much as on Metal, with no memory pressure involved. That is what makes this an ordinary upstream defect rather than a Metal edge case, and the warning names both causes rather than only the memory one. Also adds upstream `tests/test-model-split.cpp` (5 cases in upstream's `testing.h` style) + its `llama_build_and_test` registration. Touches `src/llama-model.{cpp,h}`, `tests/test-model-split.cpp` and `tests/CMakeLists.txt` — **none** of which any other patch touches, so it is independent of all of them. Upstream-submittable ("model: fall back to an even split when no device reports free memory"); **not yet filed upstream**. **Runnable guard: `src/test/cpp/test_model_split.cpp`** — a FetchContent subproject builds with `LLAMA_BUILD_TESTS=OFF`, so the upstream test above is applied-but-never-compiled here (same as `0001`'s test). That file drives the same two functions from `jllama_test`, which runs on **every** platform in `C++ Tests`, so a bump that drops this patch fails the build at link time everywhere instead of surfacing as one red macOS Java job. **Verification limit — read before assuming this can be dropped:** the *failing path* still cannot be reached without a GPU backend, so the guard pins the arithmetic (what actually broke), not the end-to-end load; the end-to-end proof is the macOS CI job. On a bump, re-check whether upstream added its own `split_sum == 0` guard (grep `split_sum` in `src/llama-model.cpp`) and **drop this patch rather than refreshing it** if they did — the fail-loud applier detects "does not apply", never "upstream already fixed this". | | `0014-common-log-callback-sink.patch` | **Gives `common_log` a callback sink: `common_log_set_callback(log, cb, user_data)` (`common/log.{h,cpp}`).** This is what `LlamaModel.setLogger` hooks. Before it, the Java logger was a `llama_log_set()` callback, which has two holes, both found while chasing `slot print_timing` lines interleaving with the Atmosphere agent's streamed answer: **(a)** every model load runs `common_init()`, which re-points `llama_log_set()` at `common_log_default_callback` (`common.cpp:394`), so `setLogger(…)` *before* `new LlamaModel(…)` silently lost the callback; **(b)** the server's own `SRV_*`/`SLT_*` macros are `LOG_INF` and write straight into `common_log`, which `llama_log_set()` never carried, so the per-request `slot …`/`srv …` lines could not be routed to Java at all (the reason `LlamaModelTest#testLogText/JSON` sat `@Disabled` for years). `common_log` upstream offers file, colors, prefix, timestamps, verbosity and JSONL but no hook. The patch adds one: while a callback is set, the worker thread hands every entry to it **instead of** printing to stdout/stderr (a `--log-file` still receives them); the callback gets the bare formatted message (no prefix/timestamp/colors) with the `ggml_log_callback` signature; swapping the callback pauses the worker first, so queued entries reach the *previous* sink (which is what makes `setLogger(format, null)` a synchronous drain). With the sink, `common_init()`'s `llama_log_set()` reset is harmless — it points at the default callback that feeds `common_log`, i.e. exactly the path into the sink — so the ordering problem (a) disappears without any re-install logic, and the `srv`/`slot` lines (b) arrive because they are `common_log` entries. `jllama.cpp`'s `setLogger` therefore sets `common_log_set_callback(common_log_main(), trampoline)` **plus** `llama_log_set(common_log_default_callback)` (so llama/ggml lines feed `common_log` even before the first load). Two consequences to know: the callback runs on `common_log`'s **worker thread**, a plain `std::thread` llama.cpp re-creates on every pause/resume and never attaches to the JVM — the trampoline attaches per call and detaches again (`get_jni_env_attaching`; a thread that exits while attached leaks a `JavaThread`, and this thread is not ours — the leak-free choice, not the cheapest: each attach creates a `java.lang.Thread` object, so a `thread_local` guard that detaches once at thread exit is the optimisation on file in `TODO.md`), `setLogger` must call `common_log_set_callback` **outside** `g_log_mutex`, because the pause joins the worker, which needs that mutex to read the callback, and `setLogger` callers are serialized by a **separate** `g_set_logger_mutex`: two unserialized swaps race on the worker's `std::thread` (one joins it while the other assigns a fresh thread over the still-joinable object = `std::terminate`, the whole JVM), which `LlamaLoggerTest#concurrentSetLoggerCallsDoNotRaceOnTheLogWorker` reproduced before the mutex existed. Two caveats the Javadoc carries: a caller must not hold a lock the *previous* callback needs (the drain runs it on the worker while the caller waits), and the verbosity threshold is process-wide and reset by **every** load (`common_params_parse` ends with `common_log_set_verbosity_thold(params.verbosity)`, default 3), so a load without `-lv` puts it back to 3 — a review assumed the opposite, and `LlamaLoggerTest#verbosityThresholdIsProcessWideAndEveryLoadSetsIt` now pins the measured behaviour. And the verbosity threshold applies *before* the sink: at the default (`3`) the Java logger sees errors, warnings and the server's INFO lines, while llama/ggml INFO lines (`common_log_get_verbosity` maps them to TRACE = 4) arrive only from `setLogVerbosity(4)` on — the same filtering the console gets, and a behaviour change for consumers who captured the unfiltered `llama_log_set()` stream before. **Runnable guards:** `src/test/cpp/test_common_log_callback.cpp` (6 tests over a private `common_log_init()` instance: delivery, bare text under prefix+timestamps, clear, swap-drains-to-old-sink, file kept, levels pass through) links the function on every platform, so a bump that drops the patch reds `C++ Tests` at link time; `LlamaLoggerTest` (model-free, needs only `libjllama`: a logger set before a deliberately failing load on a non-GGUF file sees the `srv … loading model` INFO line and llama's ERROR line, in TEXT and JSON) and the re-enabled `LlamaModelTest#testLogText/testLogJSON` plus `#testLoggerSetBeforeLoadSurvivesTheLoad` (vocab-only load) cover the Java side. Upstream-submittable ("common : add a callback sink to common_log for embedding hosts"); **not yet filed upstream**. Touches only `common/log.{h,cpp}`, which no other patch touches. **On a bump, check whether upstream added a hook of its own (grep `callback` in `common/log.h`) and, if so, DROP this patch and port `setLogger` to theirs rather than refreshing it.** | +| `0015-rpc-embeddable-client-and-stoppable-server.patch` | **Makes ggml-rpc usable inside a JVM** (see "RPC backend" below). Upstream treats every client-side RPC problem as `GGML_ABORT` — a malformed endpoint, a server that is not running, a failed handshake — which in a JVM kills the application over a typo in `--rpc`. **(1)** The *registration* path reports failure instead: `rpc_dispatcher::start()` returns `false`, `try_get_dispatcher()` `nullptr`, `ggml_backend_rpc_get_device_count()` `0`, `ggml_backend_rpc_add_server()` `nullptr`, and `add_rpc_devices()` (`common/arg.cpp`) throws `std::invalid_argument` naming the server where it used to register `nullptr` and silently run without it. Every path *after* registration keeps the original contract through `get_dispatcher()`, which still aborts — a server that vanishes mid-inference stays fatal (TODO). **(2)** `add_server()`'s cache hit is re-checked (a stopped server used to pass registration and abort on the first tensor upload), and `ggml_backend_rpc_get_device_memory()` of a gone server returns 0/0 instead of aborting — RPC devices stay in ggml's process-wide registry forever, and `common_init()` queries the memory of **every** registered device at `-lv 4`, so a stopped server would otherwise take down an unrelated later load. **(3)** The server becomes stoppable: `ggml_backend_rpc_stop_server()` (disconnects the client being served, wakes `accept()` with a loopback connection — the one portable way) and `ggml_backend_rpc_server_listening()`; the loop now reaches its cleanup and frees its backends on every early return, but deliberately **not** `rpc_transport_shutdown()`, whose `WSACleanup()` would invalidate every other RPC socket in the process. **(4)** transport: `MSG_NOSIGNAL`/`SO_NOSIGPIPE`, `socket_t::shutdown()`, and the fds `connect()`/`create_server()`/`accept()` leaked on their error paths are closed. Touches `ggml/include/ggml-rpc.h`, `ggml/src/ggml-rpc/{ggml-rpc.cpp,transport.cpp,transport.h}` and one function in `common/arg.cpp` (`add_rpc_devices`, away from `0001`'s `common_params_parse` hunks). **Runnable guard: `src/test/cpp/test_rpc.cpp`** — links both new functions (a bump that drops the patch fails `C++ Tests` at link time on every platform) and drives registration, stop and the memory query over loopback. Upstream-submittable; **not yet filed upstream**. On a bump, check whether upstream grew its own stop function or non-aborting registration (grep `stop_server` / `GGML_ABORT("Failed to connect` in `ggml-rpc.cpp`). | **`0010` was dropped at the b11080 bump.** Upstream merged [ggml-org/llama.cpp#28518](https://github.com/ggml-org/llama.cpp/pull/28518) @@ -901,6 +902,116 @@ prefill behavior shows up, re-check `tools/server/server-context.cpp`'s `create_ (`test_reasoning_budget_tokens_per_request` / `test_reasoning_budget_message_per_request`, byte-identical body). No local patch needed — the tree already matches what `0004` used to add. +## RPC backend: `--rpc` client and the in-JVM `RpcServer` + +llama.cpp's RPC backend (`ggml-rpc`) is compiled into **every** artifact — the default JAR and all +GPU classifiers — with `GGML_RPC=ON` forced in `llama/CMakeLists.txt`. It stays **one** `jllama` +library: `ggml_add_backend_library` makes `ggml-rpc` a static library linked into `ggml` (only +`GGML_BACKEND_DL`, which needs shared libs, would split it out), and its registration is compiled in +(`GGML_USE_RPC`). Client and server live in the same file, so the flag brings both. + +**No new runtime dependency — and this is enforced, not assumed.** The transport is plain TCP: +BSD sockets from libc/libSystem/bionic, Winsock on Windows (`jllama.dll` already imported +`WS2_32.dll` through cpp-httplib). `GGML_RPC_RDMA` is forced **OFF**: upstream turns it on whenever +the *build host* has `libibverbs` (Linux) or `librdma` (Apple), which would make `libjllama.so` need +`libibverbs.so.1` at load time and fail on every machine without rdma-core; it cannot be linked +statically in a useful way either (it `dlopen`s its hardware providers), and upstream's own CMake +comment says the Apple weak link does not survive a static `ggml-rpc`. **`.github/verify-native-deps.py`** +(the `package` job) reads each shipped library's dependency list straight from the file — ELF +`DT_NEEDED`, PE import table including Windows arm64 (which binutils cannot read), Mach-O load +commands — and holds the default tree to an exact per-`{OS}/{ARCH}` allowlist (the classifier trees +only to a denylist, since they need their vendor runtime). A new `{OS}/{ARCH}` without an allowlist +fails too. **It found a pre-existing defect on its first run**: the macOS dylib links Homebrew's +`openssl@3` (`/opt/homebrew/opt/openssl@3/lib/libssl.3.dylib` + `libcrypto`), so it does not load on a +Mac without that formula. Those two paths are in the allowlist marked as a known defect so the +check reports only *new* dependencies; the fix is on file in `TODO.md`. + +**Three places every argv goes through (`src/main/cpp/rpc_support.hpp`).** ggml's backend registry +is process-wide and has no unregister, so an RPC server registered by one `--rpc` load stays a +device for the life of the JVM — and llama.cpp's default device selection puts RPC devices *first*. +`jllama::rpc::prepare_argv()` runs before `common_params_parse` in `LlamaModel`'s load and in +`NativeServer`'s start (attach mode parses only HTTP args and is exempt): +1. every `--rpc` endpoint is registered up front, and an unreachable one throws + `std::invalid_argument` naming it → `LlamaException` (instead of upstream's generic parse failure); +2. if the registry holds an RPC device the argv did **not** ask for and the caller gave no + `--device`/`-dev`, it appends `--device ` — the local part + chosen exactly as llama.cpp's default (`src/llama.cpp`, `llama_model_load`) chooses: duplicates by + `device_id` dropped, integrated GPUs only when there is no discrete one, `none` on a CPU-only host; +3. in the same situation the **multimodal projector's** device is pinned too (`--mmproj-device `, or `--no-mmproj-offload` without one), unless the caller named it: clip does not use the + device list but takes the first registered device of type GPU, and RPC devices are type GPU and + registered after every local backend — so on a GPU-less host a stale RPC device is exactly what + it would pick; +4. otherwise the argv is untouched, so a process that never uses RPC sees upstream's behaviour. + +**Entry points that build `common_params` themselves need the same guard, and it is easy to +forget.** `TextToSpeech` (`tts_engine.cpp`) and `LlamaTrainer` (`train_engine.cpp`) never parse an +argv, so `prepare_argv` never sees them; they call `jllama::rpc::exclude_stale_devices(params.devices)` +(same selection, applied to the null-terminated device vector; TTS also applies the mmproj answer +to its `mtmd_context_params`). Found in CI, not by reading: `RpcIntegrationTest` followed by +`TtsIntegrationTest` in one fork aborted the JVM in `get_dispatcher()` while `common_fit_params` +built a context over the gone server — `RpcIntegrationTest` now ends with a `TextToSpeech` load +that pins it (verified red with the call removed). **A new native entry point that calls +`common_init_from_params` must call it too.** +The selection rules are pure functions over a descriptor list (unit-tested with literals in +`test_rpc.cpp`); only `registered_devices()` / `register_server()` touch ggml. `RpcIntegrationTest` +pins the end-to-end reason: a model over RPC, the server stopped, then a load **without** `--rpc` in +the same JVM — which would offload to the stopped server and abort without step 2. + +**`RpcServer` (root package) + `RpcServerNative` (`rpc_bridge.cpp`).** The native calls live in the +package-private `RpcServerNative`, behind the `RpcServer.Backend` seam, so `RpcServer` itself loads +without `libjllama` and its whole Java-side lifecycle is tested with a fake backend +(`RpcServerLifecycleTest`) — the analysis build (SonarCloud, no native library) sees it covered, and a +new lifecycle branch cannot hide behind "needs natives". The in-JVM `rpc-server`: ggml's blocking +`ggml_backend_rpc_start_server()` runs on a daemon Java thread, `close()` calls +`ggml_backend_rpc_stop_server()` (patch `0015`), and `start` waits for +`ggml_backend_rpc_server_listening()` so a bind failure is a `LlamaException` rather than a silently +dead thread. It serves every accelerator, else the CPU — **never an RPC device** +(`jllama::rpc::server_devices()`), which in a JVM that is also a client would forward to itself. +The overloads taking a device list (`--device` on the command line, upstream `rpc-server -d`) serve +named devices instead, resolved **before** the server thread starts so an unknown name is a +`LlamaException` listing the available ones. That is not a convenience: ggml-rpc's client answers +every `supports_op` with `true` (upstream TODO), so a served device that cannot run an operation +**aborts the process** — the macOS CI runners' paravirtual Metal GPU has no `MUL_MAT` and took the +first CI run down that way. **Every test that computes over RPC serves `CPU`** (C++ `loopback_server`, +`RpcIntegrationTest`); do not switch one back to the default choice. +**`RpcServer` is the first entry point that loads the library before `LlamaModel`, which exposed a +loader defect** (fixed in `LlamaLoader.runOnceOnThisThread`): `JNI_OnLoad`'s `GetFieldID` on +`LlamaModel` initializes that class, whose static block re-entered `LlamaLoader.initialize()` on the +loading thread — the lock is reentrant, so a second complete load ran while the first was inside +`System.load`, clearing the temp files and, with an all-backends jar, probing and extracting every +backend again over the library being loaded. The fat-jar RPC smoke timed out on it (the doubled +extraction of the CUDA/ROCm/SYCL libraries); it now waits 300 s like the other smoke and fails if the +backend is selected more than once. Only the *nested* call is skipped — later calls still run the body, +which `BackendManifestLoadTest` relies on. +**Single instance per process** (ggml keeps the server state in globals). `startLocal` binds +loopback only; `startOnNetwork` is the explicit, warned opt-in, and binding needs an IPv4 literal +(the server uses `inet_addr`). Endpoints on the client side are `value.RpcEndpoint` +(`host:port`, IPv4 or host name — the transport has no IPv6, and upstream's parser would cut an IPv6 +literal at its first colon). `main()` is the command-line form +(`java -cp net.ladenthin.llama.RpcServer --port 50052`). + +**Tests, by layer.** C++ `test_rpc.cpp` (28, every platform in `C++ Tests`): selection rules plus the +real client/server over loopback. Java: `RpcEndpointTest`, `RpcServerOptionsTest` and `RpcServerLifecycleTest` (pure; the last one +drives start/stop, bind failure, start timeout, interrupt and the single-instance slot over a fake backend), +`RpcServerTest` (native, model-free: lifecycle, restart on the same port, single instance, bind +failure, unreachable `--rpc` load), `RpcIntegrationTest` (draft model: layers on the RPC server +proven from the load log's `model buffer size` line naming the endpoint, then the stale-server +load). CI: `.github/smoke-rpc-fatjar.sh` in `smoke-fatjar-linux` runs **two JVMs from the release +asset** — `RpcServer` in one, the default `NativeServer` with `--rpc` in the other — and requires a +chat completion, the RPC model buffer in the log, an accepted client on the server, and a clean +non-SIGABRT failure naming the endpoint for a server nobody runs. + +**One visible side effect, deliberate:** `llama_supports_gpu_offload()` is `true` on a CPU-only +build now (upstream ORs in `llama_supports_rpc()`), so `-ngl` no longer prints "no usable GPU found" +there and the load log shows "offloading N layers to GPU" lines even when every layer stays on the +CPU. Nothing else reads it (checked at b11222: `common/arg.cpp` warnings + one `llama-model.cpp` log +block), and upstream's own release binaries, which all build with `GGML_RPC=ON`, behave the same. + +**Known limits (in `TODO.md`):** a server lost *mid-inference* still aborts the process (upstream has +no error path for it); Android needs the app's `INTERNET` permission even on loopback, which the AAR +deliberately does not request; the server serves one client at a time; no authentication or TLS. + ## Qwen3-TTS via `mtmd_helper::gen_audio` (was: OuteTTS build-time extraction) The `TextToSpeech` native pipeline (`tts_engine.{h,cpp}`) drives llama.cpp's upstream Qwen3-TTS @@ -1370,7 +1481,7 @@ If the local check passes (`BUILD SUCCESS`), the `mvn package` job in The library exposes **two** ways to serve a model over HTTP, on two different transports. The fat jar's `Main-Class` is `server.ServerLauncher`, a tiny dispatcher: it runs `OpenAiCompatServer` when `--jllama-openai-compat` is present (that marker is stripped, the rest forwarded) and the default `NativeServer` otherwise. Both mains are also runnable directly by class name via `java -cp`. The two modes: -1. **`server.OpenAiCompatServer` (Java transport).** OpenAI/Ollama/Anthropic-compatible JSON API on the JDK's `com.sun.net.httpserver`, driving the compiled server *core* over JNI. Embeddable, no extra dependency, and it can share/reuse a `LlamaModel`. It serves **no** static assets — its `/` route is a 404, so **no WebUI**. It has its own `main` (run via `java -cp net.ladenthin.llama.server.OpenAiCompatServer …`); its CLI (`OpenAiServerCli`) maps a curated flag subset (`-m/-c/-b/-ub/-ngl/-t/-tb/-ctk/-ctv/--jinja/--chat-template-kwargs/--host/--port/--parallel/--mmproj/--api-key/--embedding/--reranking`). +1. **`server.OpenAiCompatServer` (Java transport).** OpenAI/Ollama/Anthropic-compatible JSON API on the JDK's `com.sun.net.httpserver`, driving the compiled server *core* over JNI. Embeddable, no extra dependency, and it can share/reuse a `LlamaModel`. It serves **no** static assets — its `/` route is a 404, so **no WebUI**. It has its own `main` (run via `java -cp net.ladenthin.llama.server.OpenAiCompatServer …`); its CLI (`OpenAiServerCli`) maps a curated flag subset (`-m/-c/-b/-ub/-ngl/-t/-tb/-ctk/-ctv/--jinja/--chat-template-kwargs/--host/--port/--parallel/--mmproj/--api-key/--embedding/--reranking/--rpc`). 2. **`server.NativeServer` (native transport) — the default fat-jar server (when `--jllama-openai-compat` is absent).** Runs the **full upstream `llama_server`** (via `patches/0006` + `native_server.cpp`) inside `libjllama`, forwarding the raw llama-server argv verbatim — so **every** llama-server flag works and the **embedded WebUI is served** (when the assets are compiled in; CI's released jars have them, local `cmake` builds use the empty-asset stub). With the classic constructor it is an **independent lifecycle** (loads its own model from the argv, like `llama-server.exe`; owns the process's llama backend + stderr logging while running); the **attach constructor** (`NativeServer(LlamaModel, String...)`, via `patches/0007`'s `llama_server_attach`) instead serves an **already-loaded `LlamaModel`** — one copy of the weights, the model's worker keeps driving inference, the HTTP routes post to its queue; caller closes the server before the model. **Router mode** (start without a model argument: `--models-dir`, `GET/POST /models`, per-request model selection) works in-JVM after `NativeServer.setWorkerCommand(...)` redirects the worker spawn to a fresh JVM (`patches/0008` — upstream re-execs its own binary, which in a JVM is `java`); the typed `server.RouterClient` (+ `value.RouterModel`, `json.RouterModelsResponseParser`) wraps the model-management endpoints (list/load/unload/await-loaded with fail-fast on failed workers) so callers don't hand-roll HTTP+JSON, and its `apiKey` constructors send `Authorization: Bearer ` — required for **every** one of those calls against a router started with `--api-key` since b10519 (#26347 dropped `/models` + `/v1/models` from the public-endpoint set; `/models/load` and `/models/unload` were always gated). `awaitModelLoaded` cannot observe a model hidden by a preset with `dedup-cache-models` (b10505/#27346 omits it from `GET /models` although it still loads and serves by name), so its "not listed" message names that cause explicitly; such a model is reached by issuing the request directly instead. Either way it is **single-instance per process** (upstream keeps shutdown state in file-scope globals) and **not available on Android** (the `subprocess.h` guard). `libjllama` loading anywhere a JVM runs is what makes this "no separate `llama-server.exe`" possible. ### `getMetrics()` — one object rebuilt from two upstream tasks @@ -1476,8 +1587,8 @@ Functions with `_impl` suffix are called directly from `jllama.cpp`. An exception that escapes a native method and unwinds across the JNI boundary is **undefined behaviour and aborts the JVM** on most implementations. **Every `Java_*` entry point must therefore -convert anything that escapes into a Java exception**, and there are 40 of them across three TUs — -`jllama.cpp` (34), `native_server.cpp` (5), `train_engine.cpp` (1). +convert anything that escapes into a Java exception**, and there are 44 of them across four TUs — +`jllama.cpp` (34), `native_server.cpp` (5), `rpc_bridge.cpp` (4), `train_engine.cpp` (1). The mechanism is `jni_guard_impl(env, exception_class, [&]() -> Ret { … })` (`jni_helpers.hpp`, Layer A). It is **additive**: an entry point that already converts `std::exception` itself keeps @@ -1678,8 +1789,9 @@ ctest --test-dir build --output-on-failure -R "ResultsToJson" | `src/test/cpp/test_model_split.cpp` | 7 | The two `load_tensors()` split helpers that `patches/0012` extracts out of llama.cpp's `src/llama-model.cpp` — `llama_model_splits_normalize` (proportional split, single device, and the zero-sum case that used to produce NaN, **and the cancelling `--tensor-split` case** — `-ts 1,-1` reaches the identical line on any backend with no GPU memory pressure at all) and `llama_model_splits_select_device` (every layer maps to a real device index; malformed split points throw a message that names the function, the layer, the index and the split values instead of libc++'s bare `"vector"`). **This is the runnable guard for `0012`**: the patch also ships an upstream `tests/test-model-split.cpp`, but a FetchContent subproject builds with `LLAMA_BUILD_TESTS=OFF`, so that one is applied-but-never-compiled here. This file is the only place the two functions are linked in CI, on every platform — so a bump that drops the patch fails the `C++ Tests` build outright rather than resurfacing as one red macOS Java job. It is the one test file that includes an **internal** upstream header (`llama-model.h`, via the `${llama.cpp_SOURCE_DIR}/src` include dir added for it), which is deliberate: a signature drift should fail loudly at compile time. | | `src/test/cpp/test_model_flags.cpp` | 4 | **The contract between the Java CLI-flag registries and llama.cpp's server argument parser.** CMake reads `ModelFlag.java` + `ModelOption.java` (`cmake/extract-java-wire-names.cmake` → a generated header of `{name, contract}` pairs), and this file asserts every `SERVER_PARSER` name is in `common_params_parser_init(params, LLAMA_EXAMPLE_SERVER).options`. It exists because **no Java test can catch this class**: `ModelFlagTest`/`ModelParametersExtendedTest` pin the *string mapping* (`hasKey("--mlock")`), never that llama.cpp still accepts the string, so they stay green forever while the flag is dead — and `common_params_parse` treats an unregistered option as a hard error, so the affected builder method makes the model **unloadable**, not merely ineffective. **A grep over `arg.cpp` is not a substitute**: `--grp-attn-n`/`-w` are present there at every pinned tag but `set_examples()`-scoped to `LLAMA_EXAMPLE_COMPLETION`/`PASSKEY`, so the server parser rejects them exactly like a deleted flag — only the real option table sees that. `--vocab-only` is the one exemption, and it declares itself `CliContract.PROJECT_PSEUDO` on its own constant rather than appearing in a list inside this file; the test asserts such a name is **still unknown** to the parser (an exemption upstream later registers would be hiding a real check) and that the exempt set is non-empty. | | `src/test/cpp/test_wire_contracts.cpp` | 6 | **The same contract for the two quieter surfaces.** `RequestField` against `server_schema::make_llama_cmpl_schema(...)` (5 tests) and `TrainingField` against `jllama_train::config_keys()` (1 test). Both receivers *silently ignore* an unknown key — the schema skips it, `train_engine.cpp` reads with `j.value(key, default)` and falls back — so a dead field produces no error anywhere and every string-mapping test keeps passing. `OAI_LAYER`-declared keys (consumed by `oaicompat_*_params_parse` before the schema) are exempt from the schema check, and are checked **both** ways: still unknown to the schema (the inverted check), and read by at least one upstream reader-shaped site (the configure-time sweep — this is what caught `chat_template`, a key a public builder wrote and nothing read). See [`docs/history/parameter-wire-surface.md`](docs/history/parameter-wire-surface.md). | +| `src/test/cpp/test_rpc.cpp` | 28 | **The runnable guard for `patches/0015` and for `rpc_support.hpp`.** The device-selection rules (pure, literal inputs: `--rpc` endpoints accumulate across repeated options, an explicit `--device`/`-dev` is never overridden, stale RPC devices are replaced by the local GPUs chosen the way llama.cpp's default does — duplicates by device id dropped, iGPUs only without a discrete GPU, `none` on a CPU-only host). Then the **real** ggml-rpc client and server over loopback, no model: a `mul_mat` graph computed on the RPC backend equals the local CPU result; `ggml_backend_rpc_stop_server()` ends the server, frees the port and disconnects a client that is still connected; an unreachable or malformed endpoint is `nullptr`/`std::invalid_argument` instead of an abort; a registered server that went away is re-checked on the next registration; its device reports 0/0 memory instead of aborting; and a stale server is kept out of a later load that did not ask for it. The loopback server picks a free port from a range so parallel jobs do not collide, and serves the **CPU** by name — a served GPU that lacks an op aborts the server (see the RPC section). Also: server devices chosen by name (case-insensitive, deduplicated), an unknown name listing the available devices, and a registered RPC device refused; the mmproj device pinned the way clip would choose it minus the stale server; and `exclude_stale_devices()`, the params-level guard for TTS and the trainer. | -**Current total: 559 tests (all passing).** +**Current total: 587 tests (all passing).** #### Upstream source location (in CMake build tree) diff --git a/README.md b/README.md index 4065c43bd..555e0cf57 100644 --- a/README.md +++ b/README.md @@ -117,6 +117,7 @@ Inference of Meta's LLaMA model (and others) in pure C/C++. - **Model metadata** access (`getModelMeta()`) and **server management** (metrics, slot save/restore, runtime thread reconfiguration). - **Conversation checkpoints** — `Session.checkpoint(...)` / `rewind(...)` / `fork(...)` branch and roll back a chat (KV-cache slot save/restore + transcript snapshot) without re-prefilling. - **GGUF metadata inspection** without loading the model (`GgufInspector` — pure Java, reads header + key/value table only, big-endian aware). +- **Distributed inference over RPC** — offload a model's layers to llama.cpp RPC servers on other machines (`ModelParameters.setRpcServers(...)` / `--rpc host:port`), and serve this machine's devices to them with `RpcServer` (the in-JVM `rpc-server`). See [Distributed inference over RPC](#distributed-inference-over-rpc). - **Multi-model router mode** (`--models-dir` + per-request model selection, managed via the typed `RouterClient`) and **attach mode** (`NativeServer(LlamaModel, ...)` serves an already-loaded model over the full upstream HTTP frontend — one copy of the weights). - Pre-built native binaries in the default JAR for Linux (x86-64, aarch64, s390x), macOS (x86-64, arm64 — Metal included), Windows (x86-64, x86, arm64) and Android (arm64, x86-64); GPU backends (CUDA, Vulkan, OpenCL, ROCm/HIP, SYCL, OpenVINO) ship as Maven classifiers — see [Choosing the right classifier](#choosing-the-right-classifier). Android additionally ships as the [`llama-android` AAR](#importing-in-android) with the optional `llama-kotlin` coroutines façade. @@ -979,6 +980,73 @@ RouterClient client = new RouterClient(8080, System.getenv("LLAMA_API_KEY")); > `dedup-cache-models` still loads and still serves by name, but never appears. For those, skip the > await and issue the request directly; with autoload the router waits for the worker itself. +### Distributed inference over RPC + +llama.cpp's RPC backend spreads one model over the devices of several machines: every machine +that contributes runs an **RPC server**, and the machine that loads the model names them with +`--rpc`. Layers are then distributed over local and remote devices exactly as over several local +GPUs (`setGpuLayers`, `setTensorSplit`). Both halves are in every artifact — the default JAR and +every GPU classifier — with **no additional runtime dependency** (plain TCP over the system socket +library the library already links). + +Serve this machine's devices (every GPU this library found, else the CPU): + +```java +try (RpcServer server = RpcServer.startLocal(RpcEndpoint.DEFAULT_PORT)) { // 127.0.0.1:50052 + server.awaitTermination(); +} +``` + +or from the command line, with the fat jar: + +```bash +java -cp llama--jar-with-dependencies.jar net.ladenthin.llama.RpcServer --port 50052 +``` + +Use the servers from a model: + +```java +ModelParameters params = new ModelParameters() + .setModel("models/big-model.gguf") + .setGpuLayers(99) + .setRpcServers(RpcEndpoint.parse("10.0.0.2:50052"), RpcEndpoint.parse("10.0.0.3:50052")); +``` + +The same works for both HTTP servers: `--rpc 10.0.0.2:50052,10.0.0.3:50052` is forwarded to the +native server as-is, and `OpenAiCompatServer` accepts it too. Any upstream `rpc-server` works as a +server, and this library's `RpcServer` works for any llama.cpp client. + +> [!WARNING] +> The RPC protocol has **no authentication and no encryption**: whoever reaches the port can use the +> devices and read or write the tensors on them. `RpcServer.startLocal` therefore binds to loopback +> only; `RpcServer.startOnNetwork(address, …)` (or `--host` on the command line) is the explicit +> opt-in for another interface and logs a warning. Across machines, use a trusted network or a +> tunnel (SSH, WireGuard). + +What to know: + +- **An unreachable server fails the load** with a `LlamaException` naming it, instead of reaching + llama.cpp. A server that disappears **after** the model loaded still terminates the process — + llama.cpp has no error path for a device lost mid-inference. +- **One `RpcServer` per process.** A second `start` while one runs throws `IllegalStateException`. +- **Endpoints are IPv4 addresses or host names** (`host:port`); llama.cpp's RPC transport has no + IPv6. `RpcServer` binds to an IPv4 literal (`127.0.0.1`, `0.0.0.0`, an interface address). +- **Registered servers stay registered.** llama.cpp keeps RPC devices in a process-wide registry + with no way to remove them; this library therefore gives every later load that does not ask for a + server an explicit device list without it, so a model loaded without `--rpc` never offloads to a + server an earlier model used — the same holds for the multimodal projector, `TextToSpeech` and + `LlamaTrainer`. An explicit `setDevices(...)` / `--device` is never overridden. +- **Android** needs the `android.permission.INTERNET` permission for RPC, even over loopback — the + `llama-android` AAR does not request it, so an app that wants RPC must declare it itself. +- `RpcServer.startLocal(port, threads, cacheDir)` enables upstream's tensor cache: a client that + loads the same model again sends the large tensors only once. +- **Choose the served devices when the default is wrong.** `RpcServer.startLocal(port, threads, + cacheDir, Arrays.asList("CPU"))` (or `--device CPU` on the command line; names as llama.cpp + prints them, e.g. `CUDA0`, `Vulkan1`, `MTL0`) replaces the default of every accelerator. It + matters because llama.cpp's RPC client treats every operation as supported by the remote device: + a served GPU that cannot run one terminates the server process on the first graph that needs it. + The paravirtual GPU of a macOS virtual machine is such a device — serve `CPU` there. + ### LangChain4j integration A separate artifact, **`net.ladenthin:llama-langchain4j`**, adapts a `LlamaModel` to diff --git a/TODO.md b/TODO.md index c16871b5e..903c1ff01 100644 --- a/TODO.md +++ b/TODO.md @@ -17,6 +17,49 @@ so everything below is genuinely still open. ## Open — jllama-specific +### macOS dylib links Homebrew OpenSSL (found by `verify-native-deps.py`) + +- **The shipped `Mac/aarch64/libjllama.dylib` needs `/opt/homebrew/opt/openssl@3/lib/libssl.3.dylib` + and `libcrypto.3.dylib`** (verified on the published 5.1.0 jar and the current snapshot). The macOS + build finds the runner's Homebrew OpenSSL and links it dynamically, so the default JAR does not load + on a Mac without `brew install openssl@3` — and the macOS smoke cannot see it, because the runner + has it. Likely fix: build BoringSSL statically on macOS as on Windows + (`LLAMA_BUILD_BORINGSSL`, `llama/CMakeLists.txt`), or turn HTTPS off there (`-DLLAMA_OPENSSL=OFF`; + the library only needs it for URL model downloads). Then delete the two allowlist lines marked + KNOWN DEFECT in `.github/verify-native-deps.py`. Needs a macOS CI run to verify, which is why it is + not folded into the RPC PR that surfaced it. + +### RPC backend — follow-ups + +- **A server lost mid-inference still aborts the JVM.** Patch `0015` makes *registration* fail + softly; every call after it (`get_dispatcher()`, `RPC_STATUS_ASSERT` in the dispatcher's `work()`) + still ends in `GGML_ABORT`, because the ggml backend interface has no error return for a lost + device (`graph_compute` returns a status, but buffer `set/get_tensor` are `void`). A real fix is an + upstream change: carry a failed-state flag through the dispatcher, fail the pending futures, and + surface it as a `GGML_STATUS_FAILED` at the next `graph_compute`, which llama.cpp already turns into + a decode error. File upstream first; do not carry it downstream. +- **A served device that cannot run an operation aborts the server process.** ggml-rpc's client + answers every `supports_op` with `true` (upstream `//TODO: call the remote backend and cache the + results` in `ggml_backend_rpc_device_supports_op`), so the scheduler never falls back and the + server hits `GGML_ABORT("unsupported op")` in the device's graph compute. Seen on the macOS CI + runners, whose paravirtual Metal GPU has no `MUL_MAT`: an in-JVM `RpcServer` serving it takes the + JVM down on the first inference. Mitigated, not fixed: `RpcServer`'s device list (`--device CPU`) + and the tests serve the CPU. The fix is upstream — forward `supports_op` over the protocol (a new + command, with a per-op cache on the client) — and belongs there, not in `0015`. +- **File patch `0015` upstream** (non-aborting registration, `ggml_backend_rpc_stop_server()`, + `ggml_backend_rpc_server_listening()`, the transport fd/SIGPIPE fixes) and drop it once merged. +- **Android RPC is untested on a device.** Bionic sockets build (upstream ships RPC in its Android + release too), but the app needs `android.permission.INTERNET` even for loopback, which the AAR + deliberately does not declare and the emulator fixture does not have. A loopback test on the + emulator needs that permission in the fixture's manifest only. +- **RPC smoke on the other fat-jar platforms.** `smoke-rpc-fatjar.sh` runs in `smoke-fatjar-linux` + only; the Java `RpcServerTest`/`RpcIntegrationTest` already run on every `test-java-*` job + (Windows and macOS included), so this is about the packaged asset, not the code path. +- **Authentication / TLS.** Upstream has none; the documented answer is a trusted network or a + tunnel. Only worth doing if it lands upstream. +- **RDMA transport** (`GGML_RPC_RDMA`) as its own classifier, since it needs `libibverbs` at runtime. +- **Several clients at once.** Upstream's server serves one connection at a time. + ### Logging sink (`patches/0014`) — follow-ups - **Keep the log worker attached instead of attaching per line.** `LlamaModel.setLogger`'s trampoline diff --git a/llama/CMakeLists.txt b/llama/CMakeLists.txt index 9c6e87470..21a191b5a 100644 --- a/llama/CMakeLists.txt +++ b/llama/CMakeLists.txt @@ -159,6 +159,21 @@ endif() set(GGML_FMA ON CACHE BOOL "" FORCE) set(GGML_F16C ON CACHE BOOL "" FORCE) set(GGML_AVX512 OFF CACHE BOOL "" FORCE) +# RPC backend (ggml-rpc): lets a model offload layers to rpc-servers on other machines +# (--rpc host:port) and lets this library serve its own devices (RpcServer). It is a static +# library linked into jllama like every other ggml backend, so the shipped artifact stays ONE +# jllama library. TCP only: plain BSD sockets / Winsock, which jllama already links (cpp-httplib +# pulls in ws2_32), so no new runtime dependency on any platform. +# +# GGML_RPC_RDMA must be forced OFF, not left to upstream's default: ggml-rpc's CMakeLists +# switches it ON whenever the BUILD host has libibverbs (Linux) or librdma (Apple), which would +# make libjllama.so NEED libibverbs.so.1 at load time -- it would then fail to load on every +# machine without rdma-core. libibverbs cannot be linked statically in any useful way either +# (it dlopens its hardware providers), and upstream notes the Apple weak link does not survive a +# static ggml-rpc. CI asserts the resulting DT_NEEDED set, so a toolchain image that grows +# rdma-core cannot turn this back on unnoticed. +set(GGML_RPC ON CACHE BOOL "" FORCE) +set(GGML_RPC_RDMA OFF CACHE BOOL "" FORCE) # b9305 removed the top-level LLAMA_BUILD_WEBUI -> LLAMA_BUILD_UI shim; set the # new name directly. (The old name no longer forwards at top level; the shim # survives in tools/ui/CMakeLists.txt but that subdir is not configured in @@ -367,6 +382,7 @@ add_library(jllama SHARED src/main/cpp/jllama.cpp src/main/cpp/tts_engine.cpp src/main/cpp/train_engine.cpp + src/main/cpp/rpc_bridge.cpp src/main/cpp/utils.hpp ${llama.cpp_SOURCE_DIR}/tools/server/server-common.cpp ${llama.cpp_SOURCE_DIR}/tools/server/server-chat.cpp) @@ -612,6 +628,7 @@ if(BUILD_TESTING) src/test/cpp/test_model_split.cpp src/test/cpp/test_model_flags.cpp src/test/cpp/test_wire_contracts.cpp + src/test/cpp/test_rpc.cpp ${llama.cpp_SOURCE_DIR}/tools/server/server-common.cpp ${llama.cpp_SOURCE_DIR}/tools/server/server-chat.cpp ${llama.cpp_SOURCE_DIR}/tools/server/server-context.cpp diff --git a/llama/patches/0015-rpc-embeddable-client-and-stoppable-server.patch b/llama/patches/0015-rpc-embeddable-client-and-stoppable-server.patch new file mode 100644 index 000000000..f7ecfde63 --- /dev/null +++ b/llama/patches/0015-rpc-embeddable-client-and-stoppable-server.patch @@ -0,0 +1,506 @@ +ggml-rpc: let an embedding host register servers and stop its own server safely + +An embedding host (here: the java-llama.cpp JNI binding, which runs llama.cpp inside a +JVM) cannot use ggml-rpc as it is: + +1. Every client-side connection problem is a GGML_ABORT: a malformed endpoint, a server + that is not running, a failed handshake. For a standalone tool that is an error exit; + inside a JVM it kills the whole application because of a typo in "--rpc". The + registration path now reports these as failures instead: + - rpc_dispatcher::start() returns false; + - try_get_dispatcher() returns nullptr; + - ggml_backend_rpc_get_device_count() returns 0; + - ggml_backend_rpc_add_server() returns nullptr. + add_rpc_devices() (common/arg.cpp) now throws std::invalid_argument naming the server, + where it used to register nullptr and silently run without the requested server. + + Every path that runs AFTER registration (buffer/graph/tensor calls) keeps the original + contract through get_dispatcher(), which still aborts. A server that vanishes mid-run + stays fatal; that needs a real error path through the backend interface. + +2. add_server() returned a cached registration without checking that the server is still + there, so re-using an endpoint whose server had stopped passed registration and aborted + on the first tensor upload. The cache hit is now re-checked. + + ggml_backend_rpc_get_device_memory() of a server that went away reports 0/0 instead of + aborting. Registered devices stay in ggml's process-wide registry forever, and code that + never uses them still queries their memory: common_init() lists every device at -lv 4, + and --list-devices does too. 0/0 is what --fit already reads as an unknown budget. + ggml_backend_rpc_init() keeps aborting if the registration fails. + +3. ggml_backend_rpc_start_server() cannot be stopped: while(true) around accept(). Two new + entry points, both also exposed through get_proc_address: + - ggml_backend_rpc_stop_server() disconnects the client being served and wakes accept() + with a loopback connection (the one portable way to interrupt it). The loop then + releases its backends; it used to never reach that code. + - ggml_backend_rpc_server_listening() tells the host whether the socket was bound. + The transport stays initialised on the way out, because WSACleanup() would invalidate + every other RPC socket in the process. + +4. transport: no SIGPIPE when the peer went away (MSG_NOSIGNAL / SO_NOSIGPIPE), a + socket_t::shutdown(), and the file descriptors that connect()/create_server()/accept() + leaked on their error paths are closed. + +Upstream-submittable; not filed upstream yet. The runnable guard is +src/test/cpp/test_rpc.cpp: it links both new functions and exercises the registration and +stop paths over loopback. + +diff --git a/common/arg.cpp b/common/arg.cpp +index 42bbc5601..7fe6b0ce2 100644 +--- a/common/arg.cpp ++++ b/common/arg.cpp +@@ -1177,6 +1177,11 @@ static void add_rpc_devices(const std::string & servers) { + } + for (const auto & server : rpc_servers) { + auto reg = ggml_backend_rpc_add_server_fn(server.c_str()); ++ if (!reg) { ++ // was: registering nullptr, i.e. silently running without the requested server ++ throw std::invalid_argument("failed to connect to RPC server " + server + ++ " (expected host:port of a running rpc-server)"); ++ } + ggml_backend_register(reg); + } + } +diff --git a/ggml/include/ggml-rpc.h b/ggml/include/ggml-rpc.h +index 1f8cb7906..a8d50ee01 100644 +--- a/ggml/include/ggml-rpc.h ++++ b/ggml/include/ggml-rpc.h +@@ -26,6 +26,11 @@ GGML_BACKEND_API void ggml_backend_rpc_get_device_memory(const char * endpoint, + + GGML_BACKEND_API void ggml_backend_rpc_start_server(const char * endpoint, const char * cache_dir, + size_t n_threads, size_t n_devices, ggml_backend_dev_t * devices); ++// Makes a ggml_backend_rpc_start_server() running on another thread return: the client being ++// served is disconnected and the listening socket closed. One server per process. ++GGML_BACKEND_API void ggml_backend_rpc_stop_server(void); ++// true while ggml_backend_rpc_start_server() is accepting connections ++GGML_BACKEND_API bool ggml_backend_rpc_server_listening(void); + + GGML_BACKEND_API ggml_backend_reg_t ggml_backend_rpc_reg(void); + GGML_BACKEND_API ggml_backend_reg_t ggml_backend_rpc_add_server(const char * endpoint); +diff --git a/ggml/src/ggml-rpc/ggml-rpc.cpp b/ggml/src/ggml-rpc/ggml-rpc.cpp +index 353b79b07..abfe8b37d 100644 +--- a/ggml/src/ggml-rpc/ggml-rpc.cpp ++++ b/ggml/src/ggml-rpc/ggml-rpc.cpp +@@ -351,7 +351,12 @@ static bool negotiate_hello(const std::shared_ptr & sock) { + sock->get_caps(request.conn_caps); + + bool status = send_rpc_cmd(sock, RPC_CMD_HELLO, &request, sizeof(request), &response, sizeof(response)); +- RPC_STATUS_ASSERT(status); ++ if (!status) { ++ // a peer that is not an RPC server (or went away) fails the handshake, it does not ++ // take the calling process down: the caller decides what an unreachable server means ++ GGML_LOG_ERROR("RPC handshake: no valid HELLO response\n"); ++ return false; ++ } + + if (response.major != RPC_PROTO_MAJOR_VERSION || response.minor > RPC_PROTO_MINOR_VERSION) { + GGML_LOG_ERROR("RPC server version mismatch: %d.%d.%d\n", +@@ -419,7 +424,7 @@ public: + void event_record(ggml_backend_event_t event); + void synchronize(); + +- void start(const std::string & endpoint); ++ bool start(const std::string & endpoint); + void work(); + + ~rpc_dispatcher(); +@@ -532,26 +537,32 @@ void rpc_dispatcher::synchronize() { + msg->completion.get_future().wait(); + } + +-void rpc_dispatcher::start(const std::string & endpoint) { ++bool rpc_dispatcher::start(const std::string & endpoint) { + std::string host; + int port; + if (!parse_endpoint(endpoint, host, port)) { +- GGML_ABORT("Failed to parse endpoint: %s\n", endpoint.c_str()); ++ GGML_LOG_ERROR("Failed to parse endpoint: %s\n", endpoint.c_str()); ++ return false; + } + if (!rpc_transport_init()) { +- GGML_ABORT("RPC transport initialization failed\n"); ++ GGML_LOG_ERROR("RPC transport initialization failed\n"); ++ return false; + } + + sock = socket_t::connect(host.c_str(), port); + if (sock == nullptr) { +- GGML_ABORT("Failed to connect to %s\n", endpoint.c_str()); ++ GGML_LOG_ERROR("Failed to connect to %s\n", endpoint.c_str()); ++ return false; + } + if (!negotiate_hello(sock)) { +- GGML_ABORT("RPC handshake failed for %s\n", endpoint.c_str()); ++ GGML_LOG_ERROR("RPC handshake failed for %s\n", endpoint.c_str()); ++ sock = nullptr; ++ return false; + } + LOG_DBG("[%s] connected to %s\n", __func__, endpoint.c_str()); + running = true; + thread = std::thread(rpc_dispatcher_trampoline, this); ++ return true; + } + + void rpc_dispatcher::work() { +@@ -582,7 +593,8 @@ rpc_dispatcher::~rpc_dispatcher() { + } + } + +-static std::shared_ptr get_dispatcher(const std::string & endpoint) { ++// Returns nullptr when the endpoint cannot be parsed, reached or handshaken with. ++static std::shared_ptr try_get_dispatcher(const std::string & endpoint) { + static std::mutex mutex; + std::lock_guard lock(mutex); + static std::unordered_map> dispatchers; +@@ -595,11 +607,23 @@ static std::shared_ptr get_dispatcher(const std::string & endpoi + } + + auto dispatcher = std::make_shared(); +- dispatcher->start(endpoint); ++ if (!dispatcher->start(endpoint)) { ++ return nullptr; ++ } + dispatchers[endpoint] = dispatcher; + return dispatcher; + } + ++// Every path below ggml_backend_rpc_add_server() keeps the original contract: a server that ++// disappears after it was registered is still fatal (see the RPC_STATUS_ASSERT in work()). ++static std::shared_ptr get_dispatcher(const std::string & endpoint) { ++ auto dispatcher = try_get_dispatcher(endpoint); ++ if (dispatcher == nullptr) { ++ GGML_ABORT("Failed to connect to RPC server %s\n", endpoint.c_str()); ++ } ++ return dispatcher; ++} ++ + 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(); +@@ -1125,6 +1149,10 @@ ggml_backend_t ggml_backend_rpc_init(const char * endpoint, uint32_t device) { + /* .name = */ dev_name, + }; + auto reg = ggml_backend_rpc_add_server(endpoint); ++ if (reg == nullptr) { ++ // add_server() reports an unreachable server as nullptr now; this entry point keeps aborting ++ GGML_ABORT("Failed to register RPC server %s\n", endpoint); ++ } + ggml_backend_t backend = new ggml_backend { + /* .guid = */ ggml_backend_rpc_guid(), + /* .iface = */ ggml_backend_rpc_interface, +@@ -1139,7 +1167,15 @@ bool ggml_backend_is_rpc(ggml_backend_t backend) { + } + + void ggml_backend_rpc_get_device_memory(const char * endpoint, uint32_t device, size_t * free, size_t * total) { +- auto dispatcher = get_dispatcher(endpoint); ++ auto dispatcher = try_get_dispatcher(endpoint); ++ if (dispatcher == nullptr) { ++ // A registered server that is gone reports no memory instead of aborting: every device in ++ // the process-wide registry is queried by code that never uses it (common_init's device ++ // listing, --list-devices), and an unknown budget is exactly what 0/0 means to --fit. ++ *free = 0; ++ *total = 0; ++ return; ++ } + auto request = std::make_shared(); + request->device = device; + rpc_msg_get_device_memory_rsp response; +@@ -2085,6 +2121,41 @@ static void rpc_serve_client(const std::vector & backends, const + } + } + ++// State of the one server a process can run, so that another thread can stop it. ++// ggml_backend_rpc_start_server() blocks for the lifetime of the server; an embedding host ++// (e.g. a JVM) runs it on a thread of its own and needs a way to end it without ending the process. ++static std::mutex g_rpc_server_mutex; ++static std::atomic g_rpc_server_stop{false}; ++static std::atomic g_rpc_server_listening{false}; ++static socket_ptr g_rpc_server_client; // the client being served, if any ++static std::string g_rpc_server_wake; // endpoint that reaches the listening socket ++ ++void ggml_backend_rpc_stop_server(void) { ++ std::string wake; ++ { ++ std::lock_guard lock(g_rpc_server_mutex); ++ g_rpc_server_stop = true; ++ if (g_rpc_server_client) { ++ // ends rpc_serve_client(): its blocking recv returns an error ++ g_rpc_server_client->shutdown(); ++ } ++ wake = g_rpc_server_wake; ++ } ++ if (!wake.empty()) { ++ // accept() is not interruptible from another thread in a portable way; a connection ++ // is. The loop sees the stop flag as soon as accept() returns and closes it unserved. ++ std::string host; ++ int port; ++ if (parse_endpoint(wake, host, port)) { ++ socket_t::connect(host.c_str(), port); ++ } ++ } ++} ++ ++bool ggml_backend_rpc_server_listening(void) { ++ return g_rpc_server_listening; ++} ++ + 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) { +@@ -2092,6 +2163,12 @@ void ggml_backend_rpc_start_server(const char * endpoint, const char * cache_dir + return; + } + std::vector backends; ++ auto free_backends = [&backends]() { ++ for (auto backend : backends) { ++ ggml_backend_free(backend); ++ } ++ backends.clear(); ++ }; + printf("Starting RPC server v%d.%d.%d\n", + RPC_PROTO_MAJOR_VERSION, + RPC_PROTO_MINOR_VERSION, +@@ -2108,6 +2185,7 @@ void ggml_backend_rpc_start_server(const char * endpoint, const char * cache_dir + 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)); ++ free_backends(); + return; + } + backends.push_back(backend); +@@ -2123,6 +2201,7 @@ void ggml_backend_rpc_start_server(const char * endpoint, const char * cache_dir + std::string host; + int port; + if (!parse_endpoint(endpoint, host, port)) { ++ free_backends(); + return; + } + +@@ -2133,29 +2212,57 @@ void ggml_backend_rpc_start_server(const char * endpoint, const char * cache_dir + #endif // GGML_RPC_RDMA + if (!rpc_transport_init()) { + fprintf(stderr, "Failed to initialize RPC transport\n"); ++ free_backends(); + return; + } + auto server_socket = socket_t::create_server(host.c_str(), port); + if (server_socket == nullptr) { + fprintf(stderr, "Failed to create server socket\n"); ++ free_backends(); + return; + } +- while (true) { ++ { ++ std::lock_guard lock(g_rpc_server_mutex); ++ g_rpc_server_stop = false; ++ // a server bound to every interface is reached for the wake-up on loopback ++ g_rpc_server_wake = (host == "0.0.0.0" ? std::string("127.0.0.1") : host) + ":" + std::to_string(port); ++ } ++ g_rpc_server_listening = true; ++ while (!g_rpc_server_stop) { + auto client_socket = server_socket->accept(); + if (client_socket == nullptr) { +- fprintf(stderr, "Failed to accept client connection\n"); +- return; ++ if (!g_rpc_server_stop) { ++ fprintf(stderr, "Failed to accept client connection\n"); ++ } ++ break; ++ } ++ { ++ std::lock_guard lock(g_rpc_server_mutex); ++ if (g_rpc_server_stop) { ++ break; // the wake-up connection of ggml_backend_rpc_stop_server() ++ } ++ g_rpc_server_client = client_socket; + } + printf("Accepted client connection\n"); + fflush(stdout); + rpc_serve_client(backends, cache_dir, client_socket); ++ { ++ std::lock_guard lock(g_rpc_server_mutex); ++ g_rpc_server_client = nullptr; ++ } + printf("Client connection closed\n"); + fflush(stdout); + } +- rpc_transport_shutdown(); +- for (auto backend : backends) { +- ggml_backend_free(backend); ++ { ++ std::lock_guard lock(g_rpc_server_mutex); ++ g_rpc_server_wake.clear(); + } ++ g_rpc_server_listening = false; ++ server_socket = nullptr; ++ // Deliberately no rpc_transport_shutdown(): the transport (Winsock) is shared with every RPC ++ // client dispatcher in this process, and WSACleanup() would invalidate their sockets. The ++ // original loop never got here, so it never cleaned up either. ++ free_backends(); + } + + static const char * ggml_backend_rpc_device_get_name(ggml_backend_dev_t dev) { +@@ -2299,6 +2406,12 @@ static void * ggml_backend_rpc_get_proc_address(ggml_backend_reg_t reg, const ch + if (std::strcmp(name, "ggml_backend_rpc_start_server") == 0) { + return (void *)ggml_backend_rpc_start_server; + } ++ if (std::strcmp(name, "ggml_backend_rpc_stop_server") == 0) { ++ return (void *)ggml_backend_rpc_stop_server; ++ } ++ if (std::strcmp(name, "ggml_backend_rpc_server_listening") == 0) { ++ return (void *)ggml_backend_rpc_server_listening; ++ } + return NULL; + + GGML_UNUSED(reg); +@@ -2321,8 +2434,13 @@ ggml_backend_reg_t ggml_backend_rpc_reg(void) { + return &ggml_backend_rpc_reg; + } + ++// 0 when the server cannot be reached, so that registering an unreachable server is an ++// ordinary failure for the caller rather than an abort of the whole process. + static uint32_t ggml_backend_rpc_get_device_count(const char * endpoint) { +- auto dispatcher = get_dispatcher(endpoint); ++ auto dispatcher = try_get_dispatcher(endpoint); ++ if (dispatcher == nullptr) { ++ return 0; ++ } + rpc_msg_device_count_rsp response; + dispatcher->send(RPC_CMD_DEVICE_COUNT, nullptr, 0, &response, sizeof(response)); + return response.device_count; +@@ -2341,6 +2459,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()) { ++ // an endpoint registered earlier may have gone away since: re-check, so the caller ++ // is told now instead of the first tensor upload aborting the process ++ if (ggml_backend_rpc_get_device_count(endpoint) == 0) { ++ return nullptr; ++ } + return reg_map[endpoint]; + } + uint32_t dev_count = ggml_backend_rpc_get_device_count(endpoint); +diff --git a/ggml/src/ggml-rpc/transport.cpp b/ggml/src/ggml-rpc/transport.cpp +index b28d16605..5609b733a 100644 +--- a/ggml/src/ggml-rpc/transport.cpp ++++ b/ggml/src/ggml-rpc/transport.cpp +@@ -549,7 +549,12 @@ bool socket_t::impl::send_data(const void * data, size_t size) { + size_t bytes_sent = 0; + while (bytes_sent < size) { + size_t size_to_send = std::min(size - bytes_sent, MAX_CHUNK_SIZE); ++#ifdef MSG_NOSIGNAL ++ // a peer that went away must fail this call, not raise SIGPIPE in the host process ++ ssize_t n = send(fd, (const char *)data + bytes_sent, size_to_send, MSG_NOSIGNAL); ++#else + ssize_t n = send(fd, (const char *)data + bytes_sent, size_to_send, 0); ++#endif + if (n < 0) { + GGML_LOG_ERROR("send failed (bytes_sent=%zu, size_to_send=%zu)\n", + bytes_sent, size_to_send); +@@ -694,6 +699,32 @@ static bool set_no_delay(sockfd_t sockfd) { + return ret == 0; + } + ++static void close_fd(sockfd_t sockfd) { ++#ifdef _WIN32 ++ closesocket(sockfd); ++#else ++ close(sockfd); ++#endif ++} ++ ++// macOS has no MSG_NOSIGNAL; the per-socket option does the same job there. ++static void set_no_sigpipe(sockfd_t sockfd) { ++#ifdef SO_NOSIGPIPE ++ int flag = 1; ++ setsockopt(sockfd, SOL_SOCKET, SO_NOSIGPIPE, (char *)&flag, sizeof(int)); ++#else ++ (void)sockfd; ++#endif ++} ++ ++void socket_t::shutdown() { ++#ifdef _WIN32 ++ ::shutdown(pimpl->fd, SD_BOTH); ++#else ++ ::shutdown(pimpl->fd, SHUT_RDWR); ++#endif ++} ++ + static bool set_reuse_addr(sockfd_t sockfd) { + int flag = 1; + int ret = setsockopt(sockfd, SOL_SOCKET, SO_REUSEADDR, (char *)&flag, sizeof(int)); +@@ -707,8 +738,10 @@ socket_ptr socket_t::accept() { + } + if (!set_no_delay(client_socket_fd)) { + GGML_LOG_ERROR("Failed to set TCP_NODELAY\n"); ++ close_fd(client_socket_fd); + return nullptr; + } ++ set_no_sigpipe(client_socket_fd); + return socket_ptr(new socket_t(std::make_unique(client_socket_fd))); + } + +@@ -719,10 +752,12 @@ socket_ptr socket_t::create_server(const char * host, int port) { + } + if (!set_reuse_addr(sockfd)) { + GGML_LOG_ERROR("Failed to set SO_REUSEADDR\n"); ++ close_fd(sockfd); + return nullptr; + } + if (inet_addr(host) == INADDR_NONE) { + GGML_LOG_ERROR("Invalid host address: %s\n", host); ++ close_fd(sockfd); + return nullptr; + } + struct sockaddr_in serv_addr; +@@ -731,9 +766,11 @@ socket_ptr socket_t::create_server(const char * host, int port) { + serv_addr.sin_port = htons(port); + + if (bind(sockfd, (struct sockaddr *) &serv_addr, sizeof(serv_addr)) < 0) { ++ close_fd(sockfd); + return nullptr; + } + if (listen(sockfd, 1) < 0) { ++ close_fd(sockfd); + return nullptr; + } + return socket_ptr(new socket_t(std::make_unique(sockfd))); +@@ -746,18 +783,22 @@ socket_ptr socket_t::connect(const char * host, int port) { + } + if (!set_no_delay(sockfd)) { + GGML_LOG_ERROR("Failed to set TCP_NODELAY\n"); ++ close_fd(sockfd); + return nullptr; + } ++ set_no_sigpipe(sockfd); + struct sockaddr_in addr; + addr.sin_family = AF_INET; + addr.sin_port = htons(port); + struct hostent * server = gethostbyname(host); + if (server == NULL) { + GGML_LOG_ERROR("Cannot resolve host '%s'\n", host); ++ close_fd(sockfd); + return nullptr; + } + memcpy(&addr.sin_addr.s_addr, server->h_addr, server->h_length); + if (::connect(sockfd, (struct sockaddr *)&addr, sizeof(addr)) < 0) { ++ close_fd(sockfd); + return nullptr; + } + 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 3f747ecff..09616d983 100644 +--- a/ggml/src/ggml-rpc/transport.h ++++ b/ggml/src/ggml-rpc/transport.h +@@ -22,6 +22,10 @@ struct socket_t { + + socket_ptr accept(); + ++ // Ends any blocking send/recv/accept on this socket from another thread, so ++ // an embedding host can stop a server that is waiting for or serving a client. ++ void shutdown(); ++ + void get_caps(uint8_t * local_caps); + void update_caps(const uint8_t * remote_caps); + diff --git a/llama/spotbugs-exclude.xml b/llama/spotbugs-exclude.xml index 0858fe282..32809ffc7 100644 --- a/llama/spotbugs-exclude.xml +++ b/llama/spotbugs-exclude.xml @@ -666,6 +666,32 @@ SPDX-License-Identifier: MIT + + + + + + + + + + + +