13 Commits
Author SHA1 Message Date
mram be0bb36317 Pin compose example to v0.1.10
ci / test (push) Successful in 28s
ci / docker (push) Successful in 1m15s
ci / release (push) Successful in 17s
2026-09-21 21:20:49 +02:00
mram 9b267533c9 COMFY_DIR alone manages ComfyUI: derive the launch command from the standard venv layout 2026-09-21 21:18:33 +02:00
mram 6f092ddc12 Default GAME_POLL_INTERVAL to 15s: nvidia-smi polls keep the GPU awake 2026-09-21 21:10:21 +02:00
mram 0228ccc296 Pin compose example to v0.1.9
ci / test (push) Successful in 17s
ci / docker (push) Successful in 1m17s
ci / release (push) Successful in 19s
2026-09-21 20:57:34 +02:00
mram 9997913929 Sync the env file at startup, not just at install: updates append new settings 2026-09-21 20:56:32 +02:00
mram 482731ed9e Pin compose example to v0.1.8
ci / test (push) Successful in 14s
ci / docker (push) Successful in 1m11s
ci / release (push) Successful in 15s
2026-09-21 17:19:16 +02:00
mram e4bdc92ece Game detection: foreign GPU holders take an external lock hold (GAME_PROCS, GPU_FOREIGN_VRAM_MB) 2026-09-21 17:16:56 +02:00
mram 14120bf4a4 Supervisor probes before spawning: an external server on the port is used, never fought or killed 2026-09-21 16:06:28 +02:00
mram d7566329ae Managed ComfyUI: COMFY_CMD starts it on demand, idle stop frees VRAM (internal/supervise) 2026-09-21 13:48:07 +02:00
mram e43ad02fc4 Ignore staged-update leftovers next to the dev exe 2026-09-21 13:14:33 +02:00
mram 232f5b61f2 Periodic upstream health checks (HEALTH_INTERVAL, default 30s); log down/recovered transitions 2026-09-21 13:14:18 +02:00
mram 5e7a042cad CFG_VER is always a concrete version: dev builds stamp v0.0.0, never "dev" 2026-09-21 12:42:32 +02:00
mram 14c2e30478 Clarify listener wording: Ollama-compatible clients, not Ollama-facing 2026-09-21 12:38:03 +02:00
21 changed files with 1672 additions and 44 deletions
+3
View File
@@ -4,3 +4,6 @@
/compose.yml /compose.yml
/signing/ /signing/
/gpu-turnstile.exe.old
/gpu-turnstile.exe.new
/gpu-turnstile.exe~
+76 -9
View File
@@ -2,9 +2,10 @@
GPU arbitration proxy for Ollama + ComfyUI. One consumer GPU is shared by an GPU arbitration proxy for Ollama + ComfyUI. One consumer GPU is shared by an
LLM server (Ollama) and an image generator (ComfyUI); gpu-turnstile sits in LLM server (Ollama) and an image generator (ComfyUI); gpu-turnstile sits in
front of both and guarantees the GPU is always in exactly one of three states: front of both and guarantees the GPU is always in exactly one of four states:
`idle`, `llm` (N ≥ 1 Ollama requests in flight), or `image` (exactly one `idle`, `llm` (N ≥ 1 Ollama requests in flight), `image` (exactly one
ComfyUI job, Ollama models unloaded). See [SPEC.md](SPEC.md) for the full ComfyUI job, Ollama models unloaded), or `external` (a foreign process such
as a game holds the GPU). See [SPEC.md](SPEC.md) for the full
design. design.
gpu-turnstile listens on the ports the services normally use; the actual gpu-turnstile listens on the ports the services normally use; the actual
@@ -32,8 +33,9 @@ Open WebUI / n8n ────► :8188 ───┘
Each consumer is enabled by setting its URL (`OLLAMA_URL`, `COMFY_URL`) and Each consumer is enabled by setting its URL (`OLLAMA_URL`, `COMFY_URL`) and
disabled by leaving it empty — at least one is required. With only Ollama disabled by leaving it empty — at least one is required. With only Ollama
the proxy is a pass-through (no image jobs can arrive); with only ComfyUI the proxy is a pass-through (no image jobs can arrive); with only ComfyUI
the Ollama unload/warm steps are skipped. Future consumers (e.g. local game the Ollama unload/warm steps are skipped. A third, URL-less consumer —
detection) plug into the same lock the same way. detection of foreign GPU holders such as games — is enabled by `GAME_PROCS`
and/or `GPU_FOREIGN_VRAM_MB` (see below).
## Configuration ## Configuration
@@ -45,8 +47,8 @@ override file values. Invalid values fail at startup.
| Var | Default | Meaning | | Var | Default | Meaning |
|---|---|---| |---|---|---|
| `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener | | `LISTEN_OLLAMA` | `:11434` | Listener for Ollama-compatible clients |
| `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener | | `LISTEN_COMFY` | `:8188` | Listener for ComfyUI clients |
| `OLLAMA_URL` | _(empty = disabled)_ | Ollama upstream; set to enable the Ollama consumer | | `OLLAMA_URL` | _(empty = disabled)_ | Ollama upstream; set to enable the Ollama consumer |
| `COMFY_URL` | _(empty = disabled)_ | ComfyUI upstream; set to enable the ComfyUI consumer | | `COMFY_URL` | _(empty = disabled)_ | ComfyUI upstream; set to enable the ComfyUI consumer |
| `UNLOAD_TIMEOUT` | `60s` | Wait for Ollama to unload before an image job | | `UNLOAD_TIMEOUT` | `60s` | Wait for Ollama to unload before an image job |
@@ -56,12 +58,21 @@ override file values. Invalid values fail at startup.
| `LLM_BUSY_STATUS` | `503` | HTTP status for rejected LLM requests in reject mode (400599, e.g. 429) | | `LLM_BUSY_STATUS` | `503` | HTTP status for rejected LLM requests in reject mode (400599, e.g. 429) |
| `BUSY_RETRY_AFTER` | `30` | Seconds sent as `Retry-After` on busy responses (both modes) | | `BUSY_RETRY_AFTER` | `30` | Seconds sent as `Retry-After` on busy responses (both modes) |
| `WARM_MODEL` | _(empty)_ | Model to reload after an image job (off by default) | | `WARM_MODEL` | _(empty)_ | Model to reload after an image job (off by default) |
| `COMFY_CMD` | _(derived from `COMFY_DIR`; both empty = unmanaged)_ | Supervise ComfyUI: start on demand, stop when idle to free VRAM. Requires `COMFY_URL` |
| `COMFY_DIR` | _(empty)_ | Standard venv install root: set alone to supervise ComfyUI with the derived command (`.venv` + `main.py`); also the working directory for `COMFY_CMD` |
| `COMFY_IDLE_TIMEOUT` | `5m` | Stop the managed ComfyUI after this long idle |
| `COMFY_START_TIMEOUT` | `2m` | Max wait for the managed ComfyUI to come up |
| `GAME_PROCS` | _(empty = disabled)_ | Process names (comma-separated); while any runs, the GPU counts as held: requests wait, Ollama unloads, managed ComfyUI stops |
| `GPU_FOREIGN_VRAM_MB` | `0` (disabled) | Also treat the GPU as held when a non-ignored process uses more VRAM than this (needs nvidia-smi) |
| `GPU_IGNORE_PROCS` | `ollama,ollama app,ollama_llama_server,python,pythonw` | Process names never counted as foreign GPU users |
| `GAME_POLL_INTERVAL` | `15s` | How often game/VRAM detection runs (don't go below ~10s — nvidia-smi polls keep the GPU awake) |
| `LOGLEVEL` | `warn` | `info` logs every request (colored arrows in text mode), `debug` adds lock transitions. `LOG_LEVEL` works as an alias | | `LOGLEVEL` | `warn` | `info` logs every request (colored arrows in text mode), `debug` adds lock transitions. `LOG_LEVEL` works as an alias |
| `LOG_FORMAT` | `text` | `json` for structured JSON logs | | `LOG_FORMAT` | `text` | `json` for structured JSON logs |
| `LOG_FILE` | _(empty)_ | Append logs to this file instead of stderr | | `LOG_FILE` | _(empty)_ | Append logs to this file instead of stderr |
| `UNLOAD_POLL_INTERVAL` | `500ms` | `/api/ps` poll interval while unloading | | `UNLOAD_POLL_INTERVAL` | `500ms` | `/api/ps` poll interval while unloading |
| `HISTORY_POLL_INTERVAL` | `1s` | `/history/<id>` poll interval while a job runs | | `HISTORY_POLL_INTERVAL` | `1s` | `/history/<id>` poll interval while a job runs |
| `PROBE_TIMEOUT` | `5s` | Startup probe of both upstreams | | `PROBE_TIMEOUT` | `5s` | Startup probe of both upstreams (also per-probe health check timeout) |
| `HEALTH_INTERVAL` | `30s` | Periodic upstream probe; down/recovered changes are logged |
| `FREE_TIMEOUT` | `30s` | `POST /free` call after an image job | | `FREE_TIMEOUT` | `30s` | `POST /free` call after an image job |
| `WARM_TIMEOUT` | `2m` | Warm-model reload after an image job | | `WARM_TIMEOUT` | `2m` | Warm-model reload after an image job |
| `SHUTDOWN_TIMEOUT` | `10s` | Graceful shutdown on SIGINT/SIGTERM | | `SHUTDOWN_TIMEOUT` | `10s` | Graceful shutdown on SIGINT/SIGTERM |
@@ -76,7 +87,7 @@ override file values. Invalid values fail at startup.
## Observability ## Observability
- `GET /healthz` (both listeners): `{"state":"idle|llm|image","llm_inflight":N,"image_pending":B}` - `GET /healthz` (both listeners): `{"state":"idle|llm|image|external","llm_inflight":N,"image_pending":B}`
- `GET /metrics` (both listeners): Prometheus text format — `gpu_turnstile_state`, - `GET /metrics` (both listeners): Prometheus text format — `gpu_turnstile_state`,
`gpu_turnstile_llm_inflight`, `gpu_turnstile_image_pending`, `gpu_turnstile_llm_inflight`, `gpu_turnstile_image_pending`,
`gpu_turnstile_image_jobs_total`, `gpu_turnstile_lock_wait_seconds` `gpu_turnstile_image_jobs_total`, `gpu_turnstile_lock_wait_seconds`
@@ -88,6 +99,62 @@ override file values. Invalid values fail at startup.
class) in text mode, which renders in `docker compose logs` on Windows class) in text mode, which renders in `docker compose logs` on Windows
Terminal. Set `NO_COLOR` to disable colors. Terminal. Set `NO_COLOR` to disable colors.
## Managed ComfyUI (`COMFY_CMD` / `COMFY_DIR`)
Don't want ComfyUI running 24/7 (it holds VRAM even when idle — and the
Desktop app kills its server when you close it)? gpu-turnstile can supervise
it: the first request starts it, it stops again after `COMFY_IDLE_TIMEOUT`
(default 5m) without work, freeing the GPU for games or the LLM.
The easy way — point `COMFY_DIR` at a standard venv install (a folder with
`.venv` and `main.py`, or `.venv` and `ComfyUI\main.py`) and the launch
command is derived from it, including `--port` from `COMFY_URL`:
```
COMFY_URL=http://127.0.0.1:8189
COMFY_DIR=C:\ComfyUI
```
For other layouts, spell the command out yourself (run it once manually to
confirm it works):
```
COMFY_URL=http://127.0.0.1:8189
COMFY_CMD="C:\ComfyUI\.venv\Scripts\python.exe" main.py --port 8189
COMFY_DIR=C:\ComfyUI
```
`POST /prompt` waits for the server to answer before taking the GPU lock
(LLM traffic keeps flowing while torch loads); crashes are logged and the
next request respawns. Shutting gpu-turnstile down stops the child too.
Running the ComfyUI Desktop app alongside is safe: if something already
answers on the port, gpu-turnstile just uses it instead of spawning
(and never kills it — it only ever stops its own child). If the managed
instance already holds the port when you open the desktop app, the
desktop's server is the one that fails to bind.
## Game detection
Want to game on the same GPU without Ollama/ComfyUI squatting on the VRAM?
gpu-turnstile can watch for foreign GPU holders and, while one is active,
make LLM/image requests wait (or 503, per `LLM_BUSY_MODE`), unload Ollama's
models and stop the managed ComfyUI so the game gets the memory. Two
detection paths, each optional, polled every `GAME_POLL_INTERVAL` (15s):
```
GAME_PROCS=cyberpunk2077.exe,bg3.exe # the reliable way on Windows
GPU_FOREIGN_VRAM_MB=1024 # catch-all via nvidia-smi
```
`GAME_PROCS` matches running process names (case-insensitive, `.exe`
optional). `GPU_FOREIGN_VRAM_MB` asks nvidia-smi which processes hold GPU
memory and treats anything not in `GPU_IGNORE_PROCS` above the threshold as
foreign — handy as a catch-all, but note that under Windows' WDDM driver
graphics-only games may not show up in nvidia-smi's per-process list, so
name your games in `GAME_PROCS` there; on Linux both paths work. When the
game exits, requests resume automatically.
## Build and run ## Build and run
```sh ```sh
+91 -8
View File
@@ -10,11 +10,12 @@ becomes very slow; on Linux it would OOM instead.
## Goal ## Goal
A single Go binary that sits in front of **both** services and guarantees that A single Go binary that sits in front of **both** services and guarantees that
at any moment the GPU is in exactly one of three states: at any moment the GPU is in exactly one of four states:
- `idle` — nothing in flight - `idle` — nothing in flight
- `llm` — N ≥ 1 Ollama inference requests in flight (concurrency allowed) - `llm` — N ≥ 1 Ollama inference requests in flight (concurrency allowed)
- `image` — exactly one ComfyUI job in flight, Ollama models unloaded - `image` — exactly one ComfyUI job in flight, Ollama models unloaded
- `external` — a foreign process (e.g. a game) holds the GPU; new work waits
Clients (LiteLLM, Open WebUI, n8n) point at gpu-turnstile instead of at the Clients (LiteLLM, Open WebUI, n8n) point at gpu-turnstile instead of at the
services. gpu-turnstile is transparent for everything that does not touch the services. gpu-turnstile is transparent for everything that does not touch the
@@ -53,9 +54,10 @@ consumer gets no listener, no startup probe, and no lock participation:
blocks since no image jobs can arrive. blocks since no image jobs can arrive.
- **Only `COMFY_URL`**: image jobs are tracked and ComfyUI's VRAM is freed - **Only `COMFY_URL`**: image jobs are tracked and ComfyUI's VRAM is freed
afterwards, but the Ollama unload and warm-reload steps are skipped. afterwards, but the Ollama unload and warm-reload steps are skipped.
- Future consumers (e.g. detecting a local game holding VRAM) plug into the - **Game detection** is a third, optional consumer without a URL: enabled by
same lock the same way: enabled by their config knob, excluded when `GAME_PROCS` and/or `GPU_FOREIGN_VRAM_MB` it watches for foreign processes
absent. holding the GPU (see below) and plugs into the same lock the same way —
excluded when both knobs are unset.
### Lock semantics ### Lock semantics
@@ -75,6 +77,11 @@ are LLM requests and the single "writer" is an image job):
requests start), waits until n == 0, sets state := `image`. Released after requests start), waits until n == 0, sets state := `image`. Released after
the ComfyUI job finished and models were freed. the ComfyUI job finished and models were freed.
- Concurrent image jobs queue FIFO behind each other. - Concurrent image jobs queue FIFO behind each other.
- **External hold**: game detection calls `SetExternal(holder)` while a
foreign process holds the GPU. New LLM and image grants block (same busy
handling as above) until `ClearExternal()`; in-flight work is not
preempted, it drains. Snapshot reports state `external` once nothing else
is in flight.
- All waits are context-aware: a client that disconnects while waiting is - All waits are context-aware: a client that disconnects while waiting is
removed from the queue. removed from the queue.
@@ -121,6 +128,73 @@ state is `idle`, send `POST /api/generate {"model":WARM_MODEL,"keep_alive":-1}`
with empty prompt to reload the chat model so the next chat doesn't pay the with empty prompt to reload the chat model so the next chat doesn't pay the
load time. Off by default. load time. Off by default.
### Managed ComfyUI (`COMFY_CMD` / `COMFY_DIR`)
When `COMFY_CMD` is set, gpu-turnstile runs ComfyUI as a supervised child
process instead of expecting an always-on server. Setting only `COMFY_DIR`
enables the same management with the launch command derived from the
standard venv layout under it (`.venv\Scripts\python.exe` on Windows,
`.venv/bin/python` on Linux; `ComfyUI\main.py`, or a flat `main.py` when
that is what exists; `--port` from the `COMFY_URL` port). Missing layout
files are flagged in the startup log.
- **Start on demand**: any ComfyUI request spawns it (double quotes in the
command line group arguments with spaces; `COMFY_DIR` sets the working
directory). `POST /prompt` additionally waits for readiness
(`/system_stats`) for up to `COMFY_START_TIMEOUT` *before* taking the GPU
lock, so LLM traffic flows while torch loads. Other requests are bridged
by the normal retry backoff. Spawn failure or a readiness timeout
answers 502.
- **Idle stop**: after `COMFY_IDLE_TIMEOUT` without requests or finished
jobs — and only while the GPU lock is idle — the process tree is killed,
freeing the VRAM ComfyUI holds. The next request restarts it.
- **Crash**: an unexpected exit is logged; the next request respawns.
gpu-turnstile's own shutdown stops the child too.
- **Coexistence**: before spawning, the URL is probed — if another server
already answers (e.g. the ComfyUI desktop app), it is used as-is and no
child is spawned; the idle watcher and shutdown only ever stop the
supervisor's own process, never the external one. The other direction —
starting the desktop app while the managed instance holds the port —
makes the *desktop* server fail to bind; gpu-turnstile is unaffected.
- Its stdout/stderr is forwarded to the log at INFO. The health check
skips the intentionally-stopped/starting states; a failed probe while
the process is alive and was previously ready is logged as DOWN.
## Game detection (foreign GPU holders)
Games and other foreign GPU users sit outside the URL-based consumer model —
nothing proxies through gpu-turnstile for them. Two independent detection
paths, polled every `GAME_POLL_INTERVAL` (default 15 s); either one being
configured enables the feature:
- **Process watch list** (`GAME_PROCS`, comma-separated, case-insensitive,
`.exe` optional): while any listed process runs, the GPU counts as held.
This is the reliable path on Windows.
- **Foreign VRAM threshold** (`GPU_FOREIGN_VRAM_MB`): `nvidia-smi
--query-compute-apps` lists per-process GPU memory; any process not in
`GPU_IGNORE_PROCS` (default: Ollama and python — ComfyUI runs under python)
holding more than the threshold counts as a foreign holder. Needs
nvidia-smi on the PATH (absent: logged once, path disabled) and works best
on Linux — under Windows' WDDM driver, graphics-only games may not appear
in the per-process list.
While a holder is detected, gpu-turnstile:
1. takes the external hold on the lock (`SetExternal`), so new LLM/image
requests wait or are rejected per `LLM_BUSY_MODE`; busy responses name the
holder,
2. once in-flight work has drained, **frees VRAM for the foreign process**:
the managed ComfyUI is stopped (an external server on its port is never
touched) and Ollama's resident models are unloaded,
3. refuses to spawn the managed ComfyUI; ComfyUI requests that would need a
spawn are answered 503 + `Retry-After` (a running desktop instance keeps
being proxied),
4. logs the transitions at WARN ("GPU held by an external process …" /
"released the GPU; resuming") and reports state `external` in `/healthz`
and `/metrics`.
When the holder disappears, the hold is lifted and queued requests proceed.
## Configuration (env) ## Configuration (env)
Configuration comes from environment variables and/or an `.env`-style Configuration comes from environment variables and/or an `.env`-style
@@ -131,8 +205,8 @@ override file values. A missing file is fine; a malformed one is fatal.
| Var | Default | Meaning | | Var | Default | Meaning |
|---|---|---| |---|---|---|
| `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener | | `LISTEN_OLLAMA` | `:11434` | listener for Ollama-compatible clients |
| `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener | | `LISTEN_COMFY` | `:8188` | listener for ComfyUI clients |
| `OLLAMA_URL` | _(empty = disabled)_ | Ollama upstream; set to enable the Ollama consumer | | `OLLAMA_URL` | _(empty = disabled)_ | Ollama upstream; set to enable the Ollama consumer |
| `COMFY_URL` | _(empty = disabled)_ | ComfyUI upstream; set to enable the ComfyUI consumer | | `COMFY_URL` | _(empty = disabled)_ | ComfyUI upstream; set to enable the ComfyUI consumer |
| `UNLOAD_TIMEOUT` | `60s` | wait for Ollama to unload | | `UNLOAD_TIMEOUT` | `60s` | wait for Ollama to unload |
@@ -142,12 +216,21 @@ override file values. A missing file is fine; a malformed one is fatal.
| `LLM_BUSY_STATUS` | `503` | HTTP status for rejected LLM requests in reject mode (400599, e.g. 429) | | `LLM_BUSY_STATUS` | `503` | HTTP status for rejected LLM requests in reject mode (400599, e.g. 429) |
| `BUSY_RETRY_AFTER` | `30` | seconds sent as `Retry-After` on busy responses (both modes) | | `BUSY_RETRY_AFTER` | `30` | seconds sent as `Retry-After` on busy responses (both modes) |
| `WARM_MODEL` | `` | optional model to reload after an image job | | `WARM_MODEL` | `` | optional model to reload after an image job |
| `COMFY_CMD` | _(derived from `COMFY_DIR`; both empty = unmanaged)_ | spawn and supervise ComfyUI on demand: first request starts it, idle stop after `COMFY_IDLE_TIMEOUT` frees its VRAM. Requires `COMFY_URL` |
| `COMFY_DIR` | `` | standard venv install root: set alone to supervise ComfyUI with the derived launch command (`.venv` + `main.py`, `--port` from `COMFY_URL`); also the working directory for `COMFY_CMD` |
| `COMFY_IDLE_TIMEOUT` | `5m` | stop the managed ComfyUI after this long without requests or jobs |
| `COMFY_START_TIMEOUT` | `2m` | how long a request waits for the managed ComfyUI to come up |
| `GAME_PROCS` | _(empty = disabled)_ | comma-separated process names (case-insensitive, `.exe` optional); while any runs, the GPU counts as held by it: requests wait, Ollama unloads, the managed ComfyUI stops |
| `GPU_FOREIGN_VRAM_MB` | `0` (disabled) | also treat the GPU as held when a process not in `GPU_IGNORE_PROCS` uses more VRAM than this; needs nvidia-smi |
| `GPU_IGNORE_PROCS` | `ollama,ollama app,ollama_llama_server,python,pythonw` | process names never counted as foreign GPU users |
| `GAME_POLL_INTERVAL` | `15s` | how often game/VRAM detection runs (nvidia-smi polls keep the GPU awake; don't go below ~10s) |
| `LOGLEVEL` | `warn` | `info` logs every request (colored arrows in text mode), `debug` adds lock transitions. `LOG_LEVEL` is accepted as an alias | | `LOGLEVEL` | `warn` | `info` logs every request (colored arrows in text mode), `debug` adds lock transitions. `LOG_LEVEL` is accepted as an alias |
| `LOG_FORMAT` | `text` | `json` for structured JSON logs | | `LOG_FORMAT` | `text` | `json` for structured JSON logs |
| `LOG_FILE` | `` | append logs to this file instead of stderr (useful as a service) | | `LOG_FILE` | `` | append logs to this file instead of stderr (useful as a service) |
| `UNLOAD_POLL_INTERVAL` | `500ms` | `/api/ps` poll interval while unloading | | `UNLOAD_POLL_INTERVAL` | `500ms` | `/api/ps` poll interval while unloading |
| `HISTORY_POLL_INTERVAL` | `1s` | `/history/<id>` poll interval while a job runs | | `HISTORY_POLL_INTERVAL` | `1s` | `/history/<id>` poll interval while a job runs |
| `PROBE_TIMEOUT` | `5s` | startup probe of both upstreams | | `PROBE_TIMEOUT` | `5s` | startup probe of both upstreams (also the per-probe health check timeout) |
| `HEALTH_INTERVAL` | `30s` | periodic probe of enabled upstreams; status changes (down/recovered) are logged |
| `FREE_TIMEOUT` | `30s` | `POST /free` call after an image job | | `FREE_TIMEOUT` | `30s` | `POST /free` call after an image job |
| `WARM_TIMEOUT` | `2m` | warm-model reload after an image job | | `WARM_TIMEOUT` | `2m` | warm-model reload after an image job |
| `SHUTDOWN_TIMEOUT` | `10s` | graceful shutdown on SIGINT/SIGTERM | | `SHUTDOWN_TIMEOUT` | `10s` | graceful shutdown on SIGINT/SIGTERM |
@@ -159,7 +242,7 @@ override file values. A missing file is fine; a malformed one is fatal.
| `UPDATE_REPO` | `https://git.rambossek.at/PUBLIC/gpu-turnstile` | repository to check for releases | | `UPDATE_REPO` | `https://git.rambossek.at/PUBLIC/gpu-turnstile` | repository to check for releases |
| `UPDATE_ASSET` | `gpu-turnstile.exe` | release asset to download | | `UPDATE_ASSET` | `gpu-turnstile.exe` | release asset to download |
| `APP_VER` | `stable` | version to run: `dev` disables updates, `stable` tracks the latest release, or an exact `vX.Y.Z` pin (up- or downgraded to) | | `APP_VER` | `stable` | version to run: `dev` disables updates, `stable` tracks the latest release, or an exact `vX.Y.Z` pin (up- or downgraded to) |
| `CFG_VER` | _(installer-managed)_ | config format reference written by `--install-service`; missing = the file is replaced with a fresh sample (backup `.bak`) | | `CFG_VER` | _(installer-managed)_ | config format reference written by `--install-service` (always a concrete `vX.Y.Z`; a dev build stamps `v0.0.0`); missing = the file is replaced with a fresh sample (backup `.bak`). New settings are appended (commented out) at install and at every startup after an update changed the version |
Startup fails fast on unparsable values and when neither consumer URL is Startup fails fast on unparsable values and when neither consumer URL is
set. Enabled upstreams are probed once at start (`/api/version`, set. Enabled upstreams are probed once at start (`/api/version`,
+241 -7
View File
@@ -15,17 +15,20 @@ import (
"os" "os"
"os/signal" "os/signal"
"path/filepath" "path/filepath"
"runtime"
"strings" "strings"
"syscall" "syscall"
"time" "time"
"gpu-turnstile/internal/comfy" "gpu-turnstile/internal/comfy"
"gpu-turnstile/internal/config" "gpu-turnstile/internal/config"
"gpu-turnstile/internal/game"
"gpu-turnstile/internal/lock" "gpu-turnstile/internal/lock"
"gpu-turnstile/internal/metrics" "gpu-turnstile/internal/metrics"
"gpu-turnstile/internal/ollama" "gpu-turnstile/internal/ollama"
"gpu-turnstile/internal/proxy" "gpu-turnstile/internal/proxy"
"gpu-turnstile/internal/service" "gpu-turnstile/internal/service"
"gpu-turnstile/internal/supervise"
"gpu-turnstile/internal/update" "gpu-turnstile/internal/update"
) )
@@ -170,6 +173,7 @@ func main() {
} }
log, logOut, logCloser := newLogger(cfg) log, logOut, logCloser := newLogger(cfg)
defer logCloser.Close() defer logCloser.Close()
syncEnvFile(resolveConfigPath(configPath), cfg.LogFile, log)
if service.IsService() { if service.IsService() {
if err := service.Run(func(ctx context.Context) error { return run(ctx, cfg, log, logOut, true) }); err != nil { if err := service.Run(func(ctx context.Context) error { return run(ctx, cfg, log, logOut, true) }); err != nil {
@@ -186,6 +190,27 @@ func main() {
} }
} }
// syncEnvFile upgrades an installer-written config file after an update:
// settings added since its CFG_VER are appended (commented out) and CFG_VER
// is bumped. Files not written by the installer (no CFG_VER), up-to-date
// files and dev builds are left untouched; a write failure is logged, not
// fatal.
func syncEnvFile(path, logFile string, log *slog.Logger) {
data, err := os.ReadFile(path)
if err != nil {
return // no config file; nothing to upgrade
}
synced, changed := config.SyncSample(string(data), version, logFile)
if !changed {
return
}
if err := os.WriteFile(path, []byte(synced), 0o644); err != nil {
log.Warn("could not append new settings to the config file", "path", path, "err", err)
return
}
log.Warn("config file updated: new settings appended", "path", path, "version", version)
}
// defaultConfigPath returns gpu-turnstile.env next to the executable. // defaultConfigPath returns gpu-turnstile.env next to the executable.
func defaultConfigPath() string { func defaultConfigPath() string {
exe, err := os.Executable() exe, err := os.Executable()
@@ -410,7 +435,26 @@ func orDisabled(url string) string {
return url return url
} }
// managedComfyCommand resolves how ComfyUI is launched when it is managed:
// COMFY_CMD verbatim, or the standard venv layout under COMFY_DIR. Empty
// when neither is set (unmanaged).
func managedComfyCommand(cfg config.Config) string {
if cfg.ComfyCmd != "" {
return cfg.ComfyCmd
}
if cfg.ComfyDir != "" {
return supervise.DefaultComfyCommand(runtime.GOOS, cfg.ComfyDir, cfg.ComfyURL)
}
return ""
}
func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Writer, isService bool) error { func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Writer, isService bool) error {
// ComfyUI can run as a managed child — COMFY_CMD verbatim, or the
// standard venv layout derived from COMFY_DIR alone: started on demand
// by the proxy, stopped after COMFY_IDLE_TIMEOUT idle (and on shutdown)
// so its VRAM is freed.
comfyCmdLine := managedComfyCommand(cfg)
// The startup line carries the version and every setting and is emitted // The startup line carries the version and every setting and is emitted
// at WARN so it is visible even with the default (quiet) log level. // at WARN so it is visible even with the default (quiet) log level.
log.Log(ctx, slog.LevelWarn, "starting gpu-turnstile", log.Log(ctx, slog.LevelWarn, "starting gpu-turnstile",
@@ -428,6 +472,7 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
"unload_poll_interval", cfg.UnloadPollInterval, "unload_poll_interval", cfg.UnloadPollInterval,
"history_poll_interval", cfg.HistoryPollInterval, "history_poll_interval", cfg.HistoryPollInterval,
"probe_timeout", cfg.ProbeTimeout, "probe_timeout", cfg.ProbeTimeout,
"health_interval", cfg.HealthInterval,
"free_timeout", cfg.FreeTimeout, "free_timeout", cfg.FreeTimeout,
"warm_timeout", cfg.WarmTimeout, "warm_timeout", cfg.WarmTimeout,
"shutdown_timeout", cfg.ShutdownTimeout, "shutdown_timeout", cfg.ShutdownTimeout,
@@ -435,6 +480,14 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
"backoff_max", cfg.BackoffMax, "backoff_max", cfg.BackoffMax,
"prompt_capture_limit", cfg.PromptCaptureLimit, "prompt_capture_limit", cfg.PromptCaptureLimit,
"warm_model", cfg.WarmModel, "warm_model", cfg.WarmModel,
"comfy_cmd", orDisabled(comfyCmdLine),
"comfy_dir", cfg.ComfyDir,
"comfy_idle_timeout", cfg.ComfyIdleTimeout,
"comfy_start_timeout", cfg.ComfyStartTimeout,
"game_procs", cfg.GameProcs,
"gpu_foreign_vram_mb", cfg.GPUForeignVRAMMB,
"gpu_ignore_procs", cfg.GPUIgnoreProcs,
"game_poll_interval", cfg.GamePollInterval,
"auto_update", cfg.AutoUpdate, "auto_update", cfg.AutoUpdate,
"update_interval", cfg.UpdateInterval, "update_interval", cfg.UpdateInterval,
"update_repo", cfg.UpdateRepo, "update_repo", cfg.UpdateRepo,
@@ -461,12 +514,38 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
} }
} }
var comfySup *supervise.Process
if comfyCmdLine != "" {
if cfg.ComfyCmd == "" {
// Derived from COMFY_DIR: flag a wrong-looking layout early,
// while the operator is still watching the startup log.
python, script := supervise.ComfyLayout(runtime.GOOS, cfg.ComfyDir)
for _, p := range []string{python, script} {
if _, err := os.Stat(p); err != nil {
log.Warn("COMFY_DIR: file not found; ComfyUI requests will fail until it exists", "path", p)
}
}
}
var err error
comfySup, err = supervise.New("comfy", comfyCmdLine, cfg.ComfyDir, comfyClient.Probe, cfg.ComfyStartTimeout, log)
if err != nil {
return err
}
defer comfySup.Stop()
gpuIdle := func() bool {
state, _, pending := lk.Snapshot()
return state == lock.StateIdle && !pending
}
go comfySup.WatchIdle(ctx, cfg.ComfyIdleTimeout, gpuIdle)
}
srv, err := proxy.New(proxy.Config{ srv, err := proxy.New(proxy.Config{
OllamaURL: cfg.OllamaURL, OllamaURL: cfg.OllamaURL,
ComfyURL: cfg.ComfyURL, ComfyURL: cfg.ComfyURL,
Lock: lk, Lock: lk,
Ollama: ollamaClient, Ollama: ollamaClient,
Comfy: comfyClient, Comfy: comfyClient,
ComfySup: comfySup,
Metrics: metrics.New(), Metrics: metrics.New(),
Log: log, Log: log,
LogColor: !cfg.LogJSON && cfg.LogFile == "" && os.Getenv("NO_COLOR") == "", LogColor: !cfg.LogJSON && cfg.LogFile == "" && os.Getenv("NO_COLOR") == "",
@@ -490,19 +569,52 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
return err return err
} }
// Probe the enabled upstreams once; failure is logged, not fatal. // Probe the enabled upstreams once; failure is logged, not fatal. A
probeCtx, probeCancel := context.WithTimeout(ctx, cfg.ProbeTimeout) // managed ComfyUI is intentionally down at startup — the first request
// starts it — so neither the startup probe nor the health check treats
// that as an outage.
probes := map[string]func(context.Context) error{}
if ollamaClient != nil { if ollamaClient != nil {
if err := ollamaClient.Probe(probeCtx); err != nil { probes["ollama"] = ollamaClient.Probe
log.Warn("ollama probe failed", "url", cfg.OllamaURL, "err", err)
}
} }
if comfyClient != nil { if comfyClient != nil {
if err := comfyClient.Probe(probeCtx); err != nil { if comfySup == nil {
log.Warn("comfy probe failed", "url", cfg.ComfyURL, "err", err) probes["comfy"] = comfyClient.Probe
} else {
// Managed upstream: an idle-stopped or still-starting server is
// not an outage, so it is skipped until it has answered once
// (Ready resets on every spawn/stop). After that, a failed
// probe while the process lives is a real "DOWN".
probes["comfy"] = func(ctx context.Context) error {
if !comfySup.Ready() {
if comfySup.Running() {
if err := comfyClient.Probe(ctx); err == nil {
comfySup.MarkReady()
}
}
return errManagedDown
}
return comfyClient.Probe(ctx)
}
}
}
probeCtx, probeCancel := context.WithTimeout(ctx, cfg.ProbeTimeout)
for name, probe := range probes {
if err := probe(probeCtx); err != nil && !errors.Is(err, errManagedDown) {
log.Warn(name+" probe failed", "err", err)
} }
} }
probeCancel() probeCancel()
if cfg.HealthInterval > 0 {
go healthLoop(ctx, cfg.HealthInterval, cfg.ProbeTimeout, log, probes)
}
// Foreign GPU holders (games, other ML jobs) — enabled by GAME_PROCS
// and/or GPU_FOREIGN_VRAM_MB — hold the lock externally while they run.
if len(cfg.GameProcs) > 0 || cfg.GPUForeignVRAMMB > 0 {
det := game.New(cfg.GameProcs, cfg.GPUForeignVRAMMB, cfg.GPUIgnoreProcs, log)
go gameLoop(ctx, cfg, log, det, lk, ollamaClient, comfySup)
}
// Bind the listeners up front so a port conflict fails fast and the // Bind the listeners up front so a port conflict fails fast and the
// readiness notification below really means "accepting connections". // readiness notification below really means "accepting connections".
@@ -560,6 +672,128 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
return nil return nil
} }
// errManagedDown marks a managed upstream that is intentionally stopped
// (idle); health checks skip it instead of logging an outage.
var errManagedDown = errors.New("managed upstream intentionally stopped")
// gameLoop polls for foreign GPU holders (a game, another ML job). While one
// is detected it holds the lock externally so new LLM and image requests
// wait (or are rejected per LLM_BUSY_MODE), and — once in-flight work has
// drained — frees VRAM for it: the managed ComfyUI is stopped and Ollama's
// resident models are unloaded.
func gameLoop(ctx context.Context, cfg config.Config, log *slog.Logger, det *game.Detector, lk *lock.Lock, ollamaClient *ollama.Client, comfySup *supervise.Process) {
ticker := time.NewTicker(cfg.GamePollInterval)
defer ticker.Stop()
held, freed := false, false
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
}
holders, err := det.Check(ctx)
if err != nil && ctx.Err() == nil {
log.Warn("game detection failed", "err", err)
}
switch {
case len(holders) > 0 && !held:
held = true
lk.SetExternal(summarizeHolders(holders))
log.Warn("GPU held by an external process; new LLM/image requests wait",
"holders", summarizeHolders(holders))
case len(holders) == 0 && held:
held, freed = false, false
lk.ClearExternal()
log.Warn("external process released the GPU; resuming")
}
if held && !freed {
if state, _, _ := lk.Snapshot(); state != lock.StateLLM && state != lock.StateImage {
freed = true
freeVRAM(ctx, cfg.UnloadTimeout, log, ollamaClient, comfySup)
}
}
}
}
// summarizeHolders joins holder descriptions for logs and busy responses,
// capping the list so a process name matching dozens of PIDs (system
// services) does not flood the log.
func summarizeHolders(holders []string) string {
const max = 5
if len(holders) > max {
return strings.Join(holders[:max], "; ") + fmt.Sprintf("; +%d more", len(holders)-max)
}
return strings.Join(holders, "; ")
}
// freeVRAM stops the managed ComfyUI (never an external server on its port)
// and unloads Ollama's resident models so the foreign process gets the GPU
// memory.
func freeVRAM(ctx context.Context, unloadTimeout time.Duration, log *slog.Logger, ollamaClient *ollama.Client, comfySup *supervise.Process) {
if comfySup != nil && comfySup.Running() {
log.Warn("stopping the managed ComfyUI to free VRAM")
comfySup.Stop()
}
if ollamaClient == nil {
return
}
uctx, cancel := context.WithTimeout(ctx, unloadTimeout)
defer cancel()
if models, err := ollamaClient.LoadedModels(uctx); err != nil || len(models) == 0 {
return // nothing resident (or ollama unreachable); nothing to free
}
elapsed, err := ollamaClient.UnloadAll(uctx)
if err != nil {
log.Warn("ollama unload incomplete; continuing", "err", err)
return
}
log.Warn("ollama models unloaded to free VRAM", "seconds", elapsed.Seconds())
}
// healthLoop probes the enabled upstreams every interval and logs status
// transitions — "is DOWN" when a previously healthy upstream stops
// answering, "recovered" when it comes back. The first round only
// establishes the baseline; the startup probe already reported that state.
func healthLoop(ctx context.Context, interval, probeTimeout time.Duration, log *slog.Logger, probes map[string]func(context.Context) error) {
ticker := time.NewTicker(interval)
defer ticker.Stop()
up := map[string]bool{}
managed := map[string]bool{}
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
}
for name, probe := range probes {
pctx, cancel := context.WithTimeout(ctx, probeTimeout)
err := probe(pctx)
cancel()
if errors.Is(err, errManagedDown) {
managed[name] = true // intentionally stopped; not an outage
continue
}
now := err == nil
if managed[name] {
// First real probe after an idle stop only re-baselines —
// an on-demand start is not a "recovery".
managed[name] = false
up[name] = now
continue
}
was, seen := up[name]
if seen && now != was {
if now {
log.Warn(name + " upstream recovered")
} else {
log.Warn(name+" upstream is DOWN", "err", err)
}
}
up[name] = now
}
}
}
// updateLoop checks for signed updates on startup and every UPDATE_INTERVAL. // updateLoop checks for signed updates on startup and every UPDATE_INTERVAL.
// In service mode a staged update is applied by exiting with exitCodeUpdate // In service mode a staged update is applied by exiting with exitCodeUpdate
// once the GPU lock is idle; the service recovery configuration restarts the // once the GPU lock is idle; the service recovery configuration restarts the
+1 -1
View File
@@ -6,7 +6,7 @@
# ComfyUI --listen 0.0.0.0 --port 8189). # ComfyUI --listen 0.0.0.0 --port 8189).
services: services:
gpu-turnstile: gpu-turnstile:
image: git.rambossek.at/public/gpu-turnstile:v0.1.7 image: git.rambossek.at/public/gpu-turnstile:v0.1.10
restart: unless-stopped restart: unless-stopped
environment: environment:
# Each consumer is enabled by setting its URL; leave one unset to # Each consumer is enabled by setting its URL; leave one unset to
+70
View File
@@ -26,6 +26,7 @@ type Config struct {
UnloadPollInterval time.Duration UnloadPollInterval time.Duration
HistoryPollInterval time.Duration HistoryPollInterval time.Duration
ProbeTimeout time.Duration ProbeTimeout time.Duration
HealthInterval time.Duration
FreeTimeout time.Duration FreeTimeout time.Duration
WarmTimeout time.Duration WarmTimeout time.Duration
ShutdownTimeout time.Duration ShutdownTimeout time.Duration
@@ -50,6 +51,30 @@ type Config struct {
LLMBusyStatus int LLMBusyStatus int
BusyRetryAfter int BusyRetryAfter int
// ComfyCmd spawns and supervises a ComfyUI server on demand. When
// ComfyCmd is empty but ComfyDir is set, management is enabled with the
// standard venv layout under ComfyDir (.venv + main.py or
// ComfyUI/main.py; --port from the COMFY_URL port) — ComfyCmd is the
// override for other layouts and doubles as the working directory when
// set explicitly. The managed server is stopped after ComfyIdleTimeout
// without requests, freeing its VRAM; ComfyStartTimeout bounds how long
// a request waits for it to come up.
ComfyCmd string
ComfyDir string
ComfyIdleTimeout time.Duration
ComfyStartTimeout time.Duration
// GameProcs (GAME_PROCS) is a watch list of process names; while any of
// them runs, the GPU is treated as held by a foreign process. The
// nvidia-smi path (GPUForeignVRAMMB, GPU_FOREIGN_VRAM_MB) does the same
// when a process not in GPUIgnoreProcs (GPU_IGNORE_PROCS) holds more than
// that many MiB of VRAM. GamePollInterval (GAME_POLL_INTERVAL) is how
// often both checks run.
GameProcs []string
GPUForeignVRAMMB int
GPUIgnoreProcs []string
GamePollInterval time.Duration
WarmModel string WarmModel string
LogLevel slog.Level LogLevel slog.Level
LogJSON bool LogJSON bool
@@ -70,6 +95,7 @@ func Defaults() Config {
UnloadPollInterval: 500 * time.Millisecond, UnloadPollInterval: 500 * time.Millisecond,
HistoryPollInterval: time.Second, HistoryPollInterval: time.Second,
ProbeTimeout: 5 * time.Second, ProbeTimeout: 5 * time.Second,
HealthInterval: 30 * time.Second,
FreeTimeout: 30 * time.Second, FreeTimeout: 30 * time.Second,
WarmTimeout: 2 * time.Minute, WarmTimeout: 2 * time.Minute,
ShutdownTimeout: 10 * time.Second, ShutdownTimeout: 10 * time.Second,
@@ -87,6 +113,14 @@ func Defaults() Config {
LLMBusyStatus: 503, LLMBusyStatus: 503,
BusyRetryAfter: 30, BusyRetryAfter: 30,
ComfyIdleTimeout: 5 * time.Minute,
ComfyStartTimeout: 2 * time.Minute,
// ComfyUI runs under python; excluding it (and Ollama) by name keeps
// our own consumers from tripping the foreign-VRAM check.
GPUIgnoreProcs: []string{"ollama", "ollama app", "ollama_llama_server", "python", "pythonw"},
GamePollInterval: 15 * time.Second,
LogLevel: slog.LevelWarn, LogLevel: slog.LevelWarn,
} }
} }
@@ -122,6 +156,17 @@ func ParseEnvFile(r io.Reader) (map[string]string, error) {
return values, scanner.Err() return values, scanner.Err()
} }
// splitList parses a comma-separated setting into trimmed, non-empty items.
func splitList(v string) []string {
var out []string
for _, item := range strings.Split(v, ",") {
if item = strings.TrimSpace(item); item != "" {
out = append(out, item)
}
}
return out
}
func envDuration(getenv func(string) string, name string, dst *time.Duration) error { func envDuration(getenv func(string) string, name string, dst *time.Duration) error {
v := getenv(name) v := getenv(name)
if v == "" { if v == "" {
@@ -148,6 +193,8 @@ func Load(getenv func(string) string) (Config, error) {
{"OLLAMA_URL", &cfg.OllamaURL}, {"OLLAMA_URL", &cfg.OllamaURL},
{"COMFY_URL", &cfg.ComfyURL}, {"COMFY_URL", &cfg.ComfyURL},
{"WARM_MODEL", &cfg.WarmModel}, {"WARM_MODEL", &cfg.WarmModel},
{"COMFY_CMD", &cfg.ComfyCmd},
{"COMFY_DIR", &cfg.ComfyDir},
{"UPDATE_REPO", &cfg.UpdateRepo}, {"UPDATE_REPO", &cfg.UpdateRepo},
{"UPDATE_ASSET", &cfg.UpdateAsset}, {"UPDATE_ASSET", &cfg.UpdateAsset},
{"LOG_FILE", &cfg.LogFile}, {"LOG_FILE", &cfg.LogFile},
@@ -166,17 +213,34 @@ func Load(getenv func(string) string) (Config, error) {
{"UNLOAD_POLL_INTERVAL", &cfg.UnloadPollInterval}, {"UNLOAD_POLL_INTERVAL", &cfg.UnloadPollInterval},
{"HISTORY_POLL_INTERVAL", &cfg.HistoryPollInterval}, {"HISTORY_POLL_INTERVAL", &cfg.HistoryPollInterval},
{"PROBE_TIMEOUT", &cfg.ProbeTimeout}, {"PROBE_TIMEOUT", &cfg.ProbeTimeout},
{"HEALTH_INTERVAL", &cfg.HealthInterval},
{"FREE_TIMEOUT", &cfg.FreeTimeout}, {"FREE_TIMEOUT", &cfg.FreeTimeout},
{"WARM_TIMEOUT", &cfg.WarmTimeout}, {"WARM_TIMEOUT", &cfg.WarmTimeout},
{"SHUTDOWN_TIMEOUT", &cfg.ShutdownTimeout}, {"SHUTDOWN_TIMEOUT", &cfg.ShutdownTimeout},
{"BACKOFF_INITIAL", &cfg.BackoffInitial}, {"BACKOFF_INITIAL", &cfg.BackoffInitial},
{"BACKOFF_MAX", &cfg.BackoffMax}, {"BACKOFF_MAX", &cfg.BackoffMax},
{"COMFY_IDLE_TIMEOUT", &cfg.ComfyIdleTimeout},
{"COMFY_START_TIMEOUT", &cfg.ComfyStartTimeout},
{"UPDATE_INTERVAL", &cfg.UpdateInterval}, {"UPDATE_INTERVAL", &cfg.UpdateInterval},
{"GAME_POLL_INTERVAL", &cfg.GamePollInterval},
} { } {
if err := envDuration(getenv, e.name, e.dst); err != nil { if err := envDuration(getenv, e.name, e.dst); err != nil {
return cfg, err return cfg, err
} }
} }
if v := getenv("GAME_PROCS"); v != "" {
cfg.GameProcs = splitList(v)
}
if v := getenv("GPU_IGNORE_PROCS"); v != "" {
cfg.GPUIgnoreProcs = splitList(v)
}
if v := getenv("GPU_FOREIGN_VRAM_MB"); v != "" {
n, err := strconv.Atoi(v)
if err != nil || n < 0 {
return cfg, fmt.Errorf("GPU_FOREIGN_VRAM_MB: must be a non-negative integer (MiB, 0 = disabled)")
}
cfg.GPUForeignVRAMMB = n
}
if v := getenv("PROMPT_CAPTURE_LIMIT"); v != "" { if v := getenv("PROMPT_CAPTURE_LIMIT"); v != "" {
n, err := strconv.ParseInt(v, 10, 64) n, err := strconv.ParseInt(v, 10, 64)
if err != nil || n < 0 { if err != nil || n < 0 {
@@ -241,6 +305,12 @@ func Load(getenv func(string) string) (Config, error) {
default: default:
return cfg, fmt.Errorf("LOG_FORMAT: must be \"text\" or \"json\"") return cfg, fmt.Errorf("LOG_FORMAT: must be \"text\" or \"json\"")
} }
if cfg.ComfyCmd != "" && cfg.ComfyURL == "" {
return cfg, fmt.Errorf("COMFY_CMD requires COMFY_URL to be set (the proxy needs somewhere to forward)")
}
if cfg.ComfyCmd == "" && cfg.ComfyDir != "" && cfg.ComfyURL == "" {
return cfg, fmt.Errorf("COMFY_DIR without COMFY_CMD requires COMFY_URL to be set (it enables the managed ComfyUI)")
}
if cfg.OllamaURL == "" && cfg.ComfyURL == "" { if cfg.OllamaURL == "" && cfg.ComfyURL == "" {
return cfg, ErrNoConsumer return cfg, ErrNoConsumer
} }
+98
View File
@@ -74,6 +74,29 @@ func TestAppVersion(t *testing.T) {
} }
} }
func TestComfyCmdRequiresURL(t *testing.T) {
_, err := Load(func(k string) string {
if k == "COMFY_CMD" {
return "python main.py"
}
return ""
})
if err == nil || !strings.Contains(err.Error(), "COMFY_CMD requires COMFY_URL") {
t.Fatalf("err = %v, want COMFY_CMD/COMFY_URL validation error", err)
}
if _, err := Load(func(k string) string {
switch k {
case "COMFY_CMD":
return "python main.py"
case "COMFY_URL":
return "http://127.0.0.1:8188"
}
return ""
}); err != nil {
t.Fatalf("COMFY_CMD with COMFY_URL must load: %v", err)
}
}
func TestParseEnvFile(t *testing.T) { func TestParseEnvFile(t *testing.T) {
input := `# comment input := `# comment
OLLAMA_URL=http://host:11435 OLLAMA_URL=http://host:11435
@@ -134,6 +157,8 @@ func TestLoadErrors(t *testing.T) {
{"LLM_BUSY_MODE", "bogus"}, {"LLM_BUSY_MODE", "bogus"},
{"LLM_BUSY_STATUS", "200"}, {"LLM_BUSY_STATUS", "200"},
{"BUSY_RETRY_AFTER", "0"}, {"BUSY_RETRY_AFTER", "0"},
{"GPU_FOREIGN_VRAM_MB", "-1"},
{"GAME_POLL_INTERVAL", "bogus"},
} { } {
_, err := Load(func(k string) string { _, err := Load(func(k string) string {
if k == tc.key { if k == tc.key {
@@ -149,3 +174,76 @@ func TestLoadErrors(t *testing.T) {
} }
} }
} }
func TestGameDetectionSettings(t *testing.T) {
cfg, err := Load(func(k string) string {
switch k {
case "OLLAMA_URL":
return "http://127.0.0.1:11435"
case "GAME_PROCS":
return " cyberpunk2077.exe, hl2.exe ,, "
case "GPU_FOREIGN_VRAM_MB":
return "1024"
case "GPU_IGNORE_PROCS":
return "ollama, my-trainer"
}
return ""
})
if err != nil {
t.Fatal(err)
}
if len(cfg.GameProcs) != 2 || cfg.GameProcs[0] != "cyberpunk2077.exe" || cfg.GameProcs[1] != "hl2.exe" {
t.Fatalf("GameProcs = %v", cfg.GameProcs)
}
if cfg.GPUForeignVRAMMB != 1024 {
t.Fatalf("GPUForeignVRAMMB = %d", cfg.GPUForeignVRAMMB)
}
if len(cfg.GPUIgnoreProcs) != 2 || cfg.GPUIgnoreProcs[1] != "my-trainer" {
t.Fatalf("GPUIgnoreProcs = %v", cfg.GPUIgnoreProcs)
}
if cfg.GamePollInterval != 15*time.Second {
t.Fatalf("GamePollInterval = %v, want 15s default", cfg.GamePollInterval)
}
// Defaults: both detection paths off, ignore list covers our consumers.
cfg, err = Load(func(k string) string {
if k == "OLLAMA_URL" {
return "http://127.0.0.1:11435"
}
return ""
})
if err != nil {
t.Fatal(err)
}
if len(cfg.GameProcs) != 0 || cfg.GPUForeignVRAMMB != 0 {
t.Fatalf("detection must be off by default: %v %d", cfg.GameProcs, cfg.GPUForeignVRAMMB)
}
if len(cfg.GPUIgnoreProcs) == 0 {
t.Fatal("GPUIgnoreProcs default must not be empty")
}
}
func TestComfyDirOnlyEnablesManaged(t *testing.T) {
// COMFY_DIR without COMFY_CMD and without COMFY_URL is a mistake.
_, err := Load(func(k string) string {
if k == "COMFY_DIR" {
return `C:\ComfyUI`
}
return ""
})
if err == nil || !strings.Contains(err.Error(), "COMFY_DIR") {
t.Fatalf("err = %v, want COMFY_DIR/COMFY_URL validation error", err)
}
// With COMFY_URL it loads — the launch command is derived from the dir.
if _, err := Load(func(k string) string {
switch k {
case "COMFY_DIR":
return `C:\ComfyUI`
case "COMFY_URL":
return "http://127.0.0.1:8189"
}
return ""
}); err != nil {
t.Fatalf("COMFY_DIR with COMFY_URL must load: %v", err)
}
}
+22 -5
View File
@@ -23,11 +23,19 @@ const appVerComment = `Version to run: "dev" disables updates, "stable" tracks t
// CFG_VER and APP_VER are not entries — they head the file, always active. // CFG_VER and APP_VER are not entries — they head the file, always active.
func sampleEntries(logFile string) []sampleEntry { func sampleEntries(logFile string) []sampleEntry {
return []sampleEntry{ return []sampleEntry{
{"LISTEN_OLLAMA", ":11434", "Ollama-facing listener address", false}, {"LISTEN_OLLAMA", ":11434", "Listen address for Ollama-compatible clients (gpu-turnstile poses as Ollama here)", false},
{"LISTEN_COMFY", ":8188", "ComfyUI-facing listener address", false}, {"LISTEN_COMFY", ":8188", "Listen address for ComfyUI clients (gpu-turnstile poses as ComfyUI here)", false},
{"OLLAMA_URL", "http://127.0.0.1:11434", "Ollama upstream URL; setting it enables the Ollama consumer (default: empty = disabled)", false}, {"OLLAMA_URL", "http://127.0.0.1:11434", "Ollama upstream URL; setting it enables the Ollama consumer (default: empty = disabled)", false},
{"COMFY_URL", "http://127.0.0.1:8188", "ComfyUI upstream URL; setting it enables the ComfyUI consumer (default: empty = disabled)", false}, {"COMFY_URL", "http://127.0.0.1:8188", "ComfyUI upstream URL; setting it enables the ComfyUI consumer (default: empty = disabled)", false},
{"WARM_MODEL", "", "Optional model to reload after an image job (default: empty = none)", false}, {"WARM_MODEL", "", "Optional model to reload after an image job (default: empty = none)", false},
{"COMFY_CMD", `"C:\ComfyUI\.venv\Scripts\python.exe" ComfyUI\main.py --port 8188`, "Spawn and supervise ComfyUI on demand: the first request starts it, it stops after COMFY_IDLE_TIMEOUT to free VRAM (default: derived from COMFY_DIR; both empty = unmanaged)", false},
{"COMFY_DIR", `C:\ComfyUI`, "Root of a standard ComfyUI venv install (.venv + main.py): setting it alone supervises ComfyUI with the derived launch command; also the working directory for COMFY_CMD", false},
{"COMFY_IDLE_TIMEOUT", "5m", "Stop the managed ComfyUI after this long without requests or jobs (frees VRAM)", false},
{"COMFY_START_TIMEOUT", "2m", "How long a request waits for the managed ComfyUI to come up", false},
{"GAME_PROCS", "cyberpunk2077.exe,hl2.exe", "While a listed process runs, the GPU counts as held by it: requests wait, Ollama unloads, managed ComfyUI stops (default: empty = disabled)", false},
{"GPU_FOREIGN_VRAM_MB", "1024", "Also treat the GPU as held when a process not in GPU_IGNORE_PROCS uses more VRAM than this (needs nvidia-smi; 0/empty = disabled)", false},
{"GPU_IGNORE_PROCS", "ollama,ollama app,ollama_llama_server,python,pythonw", "Process names never counted as foreign GPU users (ComfyUI runs under python)", false},
{"GAME_POLL_INTERVAL", "15s", "How often game/VRAM detection runs (nvidia-smi polls keep the GPU awake; don't go below ~10s)", false},
{"UNLOAD_TIMEOUT", "60s", "How long to wait for Ollama to unload a model", false}, {"UNLOAD_TIMEOUT", "60s", "How long to wait for Ollama to unload a model", false},
{"JOB_TIMEOUT", "15m", "Maximum time to wait for a ComfyUI job", false}, {"JOB_TIMEOUT", "15m", "Maximum time to wait for a ComfyUI job", false},
{"LLM_WAIT_TIMEOUT", "10m", "Max time an LLM request waits for the GPU before being answered 503 (wait mode)", false}, {"LLM_WAIT_TIMEOUT", "10m", "Max time an LLM request waits for the GPU before being answered 503 (wait mode)", false},
@@ -39,7 +47,8 @@ func sampleEntries(logFile string) []sampleEntry {
{"LOG_FILE", logFile, "Append logs to this file instead of stderr (a Windows service has no console)", logFile != ""}, {"LOG_FILE", logFile, "Append logs to this file instead of stderr (a Windows service has no console)", logFile != ""},
{"UNLOAD_POLL_INTERVAL", "500ms", "/api/ps poll interval while unloading", false}, {"UNLOAD_POLL_INTERVAL", "500ms", "/api/ps poll interval while unloading", false},
{"HISTORY_POLL_INTERVAL", "1s", "/history/<id> poll interval while a job runs", false}, {"HISTORY_POLL_INTERVAL", "1s", "/history/<id> poll interval while a job runs", false},
{"PROBE_TIMEOUT", "5s", "Startup probe timeout for the enabled upstreams", false}, {"PROBE_TIMEOUT", "5s", "Probe timeout for the startup probe and the periodic health check", false},
{"HEALTH_INTERVAL", "30s", "How often enabled upstreams are probed; status changes are logged", false},
{"FREE_TIMEOUT", "30s", "Timeout for the POST /free call after an image job", false}, {"FREE_TIMEOUT", "30s", "Timeout for the POST /free call after an image job", false},
{"WARM_TIMEOUT", "2m", "Timeout for the warm-model reload after an image job", false}, {"WARM_TIMEOUT", "2m", "Timeout for the warm-model reload after an image job", false},
{"SHUTDOWN_TIMEOUT", "10s", "Graceful shutdown timeout on SIGINT/SIGTERM", false}, {"SHUTDOWN_TIMEOUT", "10s", "Graceful shutdown timeout on SIGINT/SIGTERM", false},
@@ -57,10 +66,18 @@ func sampleEntries(logFile string) []sampleEntry {
// comment line. Everything is commented out — so all defaults apply — // comment line. Everything is commented out — so all defaults apply —
// except the CFG_VER/APP_VER header and LOG_FILE when logFile is non-empty // except the CFG_VER/APP_VER header and LOG_FILE when logFile is non-empty
// (a Windows service has no console). CFG_VER records the version that // (a Windows service has no console). CFG_VER records the version that
// wrote the file so later installs can upgrade it. // wrote the file so installs — and startups after an update — can upgrade
// it.
func SampleEnv(version, logFile string) string { func SampleEnv(version, logFile string) string {
// CFG_VER is always a concrete vX.Y.Z — never "dev". A dev build
// stamps v0.0.0, which sorts older than any release, so the next
// release install upgrades the file and stamps a proper version.
cfgVer := version
if cfgVer == "" || cfgVer == "dev" {
cfgVer = "v0.0.0"
}
var b strings.Builder var b strings.Builder
fmt.Fprintf(&b, "CFG_VER=%s\n", version) fmt.Fprintf(&b, "CFG_VER=%s\n", cfgVer)
b.WriteString("# Config format reference, written by the installer — do not edit.\n") b.WriteString("# Config format reference, written by the installer — do not edit.\n")
b.WriteString("# The installer uses it to append newly added settings on updates.\n\n") b.WriteString("# The installer uses it to append newly added settings on updates.\n\n")
b.WriteString("# gpu-turnstile configuration\n") b.WriteString("# gpu-turnstile configuration\n")
+13
View File
@@ -122,6 +122,19 @@ func TestSyncSampleAppendsMissingAppVer(t *testing.T) {
} }
} }
func TestSampleEnvDevStampsConcreteVersion(t *testing.T) {
// A dev build must never write CFG_VER=dev: it stamps v0.0.0, and the
// next release install upgrades the file to a proper version.
dev := SampleEnv("dev", sampleLogPath)
if !strings.HasPrefix(dev, "CFG_VER=v0.0.0\n") {
t.Errorf("dev sample first line: %q", strings.SplitN(dev, "\n", 2)[0])
}
out, changed := SyncSample(dev, "v0.1.7", sampleLogPath)
if !changed || !strings.HasPrefix(out, "CFG_VER=v0.1.7\n") {
t.Errorf("dev-stamped file was not upgraded to v0.1.7 (changed=%v)", changed)
}
}
func TestCompareVersions(t *testing.T) { func TestCompareVersions(t *testing.T) {
cases := []struct { cases := []struct {
a, b string a, b string
+159
View File
@@ -0,0 +1,159 @@
// Package game detects processes outside gpu-turnstile's control that hold
// the GPU — typically a game — so the proxy can block new GPU work and free
// VRAM while they run. Two detection paths: an explicit process watch list
// (GAME_PROCS) and a foreign-VRAM threshold via nvidia-smi
// (GPU_FOREIGN_VRAM_MB) that catches anything not on the ignore list.
package game
import (
"context"
"errors"
"fmt"
"log/slog"
"os/exec"
"strconv"
"strings"
)
// Process is one running OS process.
type Process struct {
PID int
Name string
}
// computeApp is one process holding GPU memory, as reported by nvidia-smi.
type computeApp struct {
PID int
UsedMB int
}
// Detector checks whether a foreign process holds the GPU. The zero value
// (no watch list, no threshold) never detects anything; main only starts the
// poll loop when at least one path is configured.
type Detector struct {
procs map[string]bool // normalized names from GAME_PROCS
vramMB int // foreign VRAM threshold; 0 = disabled
ignore map[string]bool // normalized names never counted as foreign
log *slog.Logger
noNvidia bool // nvidia-smi was not found; VRAM path disabled for good
}
// New builds a Detector from the configured watch list, VRAM threshold in
// MiB (0 disables the nvidia-smi path) and ignore list. Names are matched
// case-insensitively, with or without a trailing ".exe".
func New(procs []string, vramMB int, ignore []string, log *slog.Logger) *Detector {
if log == nil {
log = slog.Default()
}
return &Detector{
procs: nameSet(procs),
vramMB: vramMB,
ignore: nameSet(ignore),
log: log,
}
}
// normName lowercases a process name and strips a trailing ".exe" so the
// watch and ignore lists match on Windows and Linux spellings alike.
func normName(s string) string {
return strings.TrimSuffix(strings.ToLower(strings.TrimSpace(s)), ".exe")
}
func nameSet(names []string) map[string]bool {
set := make(map[string]bool, len(names))
for _, n := range names {
if n = normName(n); n != "" {
set[n] = true
}
}
return set
}
// Check looks once for foreign GPU holders and returns a human-readable
// description of each (empty when the GPU is free for gpu-turnstile's
// consumers). A failing nvidia-smi call is returned as an error only when
// the process list found nothing; a missing nvidia-smi binary disables the
// VRAM path permanently (logged once).
func (d *Detector) Check(ctx context.Context) ([]string, error) {
ps, psErr := processes()
if d.vramMB <= 0 || d.noNvidia {
return d.detect(ps, nil), psErr
}
apps, err := queryComputeApps(ctx)
if errors.Is(err, exec.ErrNotFound) {
d.noNvidia = true
d.log.Warn("GPU_FOREIGN_VRAM_MB is set but nvidia-smi was not found; VRAM detection disabled")
return d.detect(ps, nil), nil
}
if err != nil {
return d.detect(ps, nil), err
}
return d.detect(ps, apps), nil
}
// detect is the pure core of Check: given the process table and (optionally)
// the nvidia-smi compute-apps list, it returns the foreign holders.
func (d *Detector) detect(ps []Process, apps []computeApp) []string {
var holders []string
for _, p := range ps {
if d.procs[normName(p.Name)] {
holders = append(holders, fmt.Sprintf("%s (pid %d)", p.Name, p.PID))
}
}
if d.vramMB > 0 && apps != nil {
names := make(map[int]string, len(ps))
for _, p := range ps {
names[p.PID] = p.Name
}
for _, a := range apps {
name := names[a.PID]
if d.ignore[normName(name)] || a.UsedMB < d.vramMB {
continue
}
if name == "" {
name = "unknown process"
}
holders = append(holders, fmt.Sprintf("%s (pid %d) using %d MiB VRAM", name, a.PID, a.UsedMB))
}
}
return holders
}
// queryComputeApps runs nvidia-smi and parses the per-process VRAM list.
// Note: under Windows' WDDM driver, nvidia-smi only sees compute
// allocations, so graphics-only games may not appear there — GAME_PROCS is
// the reliable path on Windows; on Linux both work.
func queryComputeApps(ctx context.Context) ([]computeApp, error) {
out, err := exec.CommandContext(ctx, "nvidia-smi",
"--query-compute-apps=pid,used_memory", "--format=csv,noheader,nounits").Output()
if err != nil {
return nil, err
}
return parseComputeApps(string(out))
}
// parseComputeApps parses "pid, used_memory" CSV lines (no header, MiB
// units). Unsupported rows ("N/A" on WDDM) are skipped.
func parseComputeApps(out string) ([]computeApp, error) {
var apps []computeApp
for _, line := range strings.Split(out, "\n") {
line = strings.TrimSpace(line)
if line == "" {
continue
}
pidStr, memStr, ok := strings.Cut(line, ",")
if !ok {
return nil, fmt.Errorf("nvidia-smi: unexpected line %q", line)
}
pid, err := strconv.Atoi(strings.TrimSpace(pidStr))
if err != nil {
return nil, fmt.Errorf("nvidia-smi: unexpected pid in %q", line)
}
mem, err := strconv.Atoi(strings.TrimSpace(memStr))
if err != nil {
continue // "N/A" and friends: unsupported under WDDM
}
apps = append(apps, computeApp{PID: pid, UsedMB: mem})
}
return apps, nil
}
+96
View File
@@ -0,0 +1,96 @@
package game
import (
"runtime"
"slices"
"testing"
)
func TestNormName(t *testing.T) {
for in, want := range map[string]string{
"Cyberpunk2077.exe": "cyberpunk2077",
"ollama": "ollama",
"OLLAMA APP.EXE": "ollama app",
" python.exe ": "python",
"hl2": "hl2",
} {
if got := normName(in); got != want {
t.Errorf("normName(%q) = %q, want %q", in, got, want)
}
}
}
func TestParseComputeApps(t *testing.T) {
apps, err := parseComputeApps("1234, 512\n 42 , 8192 \n\n")
if err != nil {
t.Fatal(err)
}
want := []computeApp{{PID: 1234, UsedMB: 512}, {PID: 42, UsedMB: 8192}}
if !slices.Equal(apps, want) {
t.Errorf("got %+v, want %+v", apps, want)
}
// WDDM reports "N/A" for memory; those rows are skipped, not fatal.
apps, err = parseComputeApps("1234, N/A\n42, 1024\n")
if err != nil {
t.Fatal(err)
}
if !slices.Equal(apps, []computeApp{{PID: 42, UsedMB: 1024}}) {
t.Errorf("got %+v", apps)
}
if _, err := parseComputeApps("garbage\n"); err == nil {
t.Error("expected an error for a malformed line")
}
}
func TestDetect(t *testing.T) {
d := New([]string{"Cyberpunk2077.exe", "hl2"}, 1024,
[]string{"ollama", "python", "pythonw"}, nil)
ps := []Process{
{PID: 10, Name: "ollama.exe"},
{PID: 20, Name: "python.exe"},
{PID: 30, Name: "cyberpunk2077.exe"},
}
apps := []computeApp{
{PID: 10, UsedMB: 8192}, // ignored: ollama
{PID: 20, UsedMB: 4096}, // ignored: python (ComfyUI)
{PID: 40, UsedMB: 2048}, // foreign, above threshold
{PID: 50, UsedMB: 100}, // foreign but below threshold
}
holders := d.detect(ps, apps)
if len(holders) != 2 {
t.Fatalf("got %v, want 2 holders", holders)
}
if holders[0] != "cyberpunk2077.exe (pid 30)" {
t.Errorf("holders[0] = %q", holders[0])
}
if holders[1] != "unknown process (pid 40) using 2048 MiB VRAM" {
t.Errorf("holders[1] = %q", holders[1])
}
}
func TestDetectNothingConfigured(t *testing.T) {
d := New(nil, 0, nil, nil)
if got := d.detect([]Process{{PID: 1, Name: "game.exe"}}, nil); len(got) != 0 {
t.Errorf("got %v, want none", got)
}
}
func TestProcessesLive(t *testing.T) {
if runtime.GOOS != "windows" && runtime.GOOS != "linux" {
t.Skip("no process listing on this platform")
}
ps, err := processes()
if err != nil {
t.Fatal(err)
}
if len(ps) == 0 {
t.Fatal("no processes listed")
}
for _, p := range ps {
if p.Name == "" {
t.Errorf("pid %d has an empty name", p.PID)
}
}
}
+30
View File
@@ -0,0 +1,30 @@
//go:build linux
package game
import (
"os"
"strconv"
"strings"
)
// processes lists the running processes from /proc/<pid>/comm.
func processes() ([]Process, error) {
entries, err := os.ReadDir("/proc")
if err != nil {
return nil, err
}
var ps []Process
for _, e := range entries {
pid, err := strconv.Atoi(e.Name())
if err != nil {
continue
}
comm, err := os.ReadFile("/proc/" + e.Name() + "/comm")
if err != nil {
continue // process vanished mid-walk
}
ps = append(ps, Process{PID: pid, Name: strings.TrimSpace(string(comm))})
}
return ps, nil
}
+7
View File
@@ -0,0 +1,7 @@
//go:build !windows && !linux
package game
// processes is unsupported on this platform; the process watch list never
// matches and the nvidia-smi path reports PIDs without names.
func processes() ([]Process, error) { return nil, nil }
+35
View File
@@ -0,0 +1,35 @@
//go:build windows
package game
import (
"unsafe"
"golang.org/x/sys/windows"
)
// processes lists the running processes via the Toolhelp32 snapshot API.
func processes() ([]Process, error) {
h, err := windows.CreateToolhelp32Snapshot(windows.TH32CS_SNAPPROCESS, 0)
if err != nil {
return nil, err
}
defer windows.CloseHandle(h) //nolint:errcheck // best effort
var entry windows.ProcessEntry32
entry.Size = uint32(unsafe.Sizeof(entry))
if err := windows.Process32First(h, &entry); err != nil {
return nil, err
}
var ps []Process
for {
ps = append(ps, Process{
PID: int(entry.ProcessID),
Name: windows.UTF16ToString(entry.ExeFile[:]),
})
if err := windows.Process32Next(h, &entry); err != nil {
break // ERROR_NO_MORE_FILES ends the walk
}
}
return ps, nil
}
+43 -8
View File
@@ -1,6 +1,8 @@
// Package lock implements the two-mode GPU arbitration lock: any number of // Package lock implements the two-mode GPU arbitration lock: any number of
// concurrent LLM requests ("readers") or exactly one image job ("writer"), // concurrent LLM requests ("readers") or exactly one image job ("writer"),
// with image jobs taking priority over newly arriving LLM requests. // with image jobs taking priority over newly arriving LLM requests. An
// external hold (SetExternal) blocks new grants of both kinds while a
// foreign process — e.g. a game — holds the GPU; in-flight work drains.
package lock package lock
import ( import (
@@ -13,9 +15,10 @@ import (
type State string type State string
const ( const (
StateIdle State = "idle" StateIdle State = "idle"
StateLLM State = "llm" StateLLM State = "llm"
StateImage State = "image" StateImage State = "image"
StateExternal State = "external"
) )
type imageWaiter struct{ id uint64 } type imageWaiter struct{ id uint64 }
@@ -30,6 +33,7 @@ type Lock struct {
imageActive bool // an image job holds the GPU imageActive bool // an image job holds the GPU
imageQ []imageWaiter imageQ []imageWaiter
nextID uint64 nextID uint64
external string // non-empty: a foreign process (e.g. a game) holds the GPU
log *slog.Logger log *slog.Logger
} }
@@ -52,12 +56,41 @@ func (l *Lock) logTransition(msg string, args ...any) {
} }
} }
// SetExternal records that a process outside gpu-turnstile's control (a
// game, another ML job) holds the GPU: new LLM and image grants block until
// ClearExternal. In-flight work is not preempted. holder describes the
// process for logs and busy responses.
func (l *Lock) SetExternal(holder string) {
l.mu.Lock()
l.external = holder
l.broadcast()
l.mu.Unlock()
l.logTransition("lock transition", "state", StateExternal, "holder", holder)
}
// ClearExternal lifts the external hold; waiting LLM and image requests
// proceed.
func (l *Lock) ClearExternal() {
l.mu.Lock()
l.external = ""
l.broadcast()
l.mu.Unlock()
l.logTransition("lock transition", "state", StateIdle)
}
// External returns the current external holder, or "" when none.
func (l *Lock) External() string {
l.mu.Lock()
defer l.mu.Unlock()
return l.external
}
// AcquireLLM blocks until no image job is active or pending, then registers // AcquireLLM blocks until no image job is active or pending, then registers
// one in-flight LLM request. Returns ctx.Err() if the context is cancelled // one in-flight LLM request. Returns ctx.Err() if the context is cancelled
// while waiting; no state is changed in that case. // while waiting; no state is changed in that case.
func (l *Lock) AcquireLLM(ctx context.Context) error { func (l *Lock) AcquireLLM(ctx context.Context) error {
l.mu.Lock() l.mu.Lock()
for l.imageActive || len(l.imageQ) > 0 { for l.imageActive || len(l.imageQ) > 0 || l.external != "" {
ch := l.change ch := l.change
l.mu.Unlock() l.mu.Unlock()
select { select {
@@ -76,10 +109,10 @@ func (l *Lock) AcquireLLM(ctx context.Context) error {
// TryAcquireLLM acquires one in-flight LLM slot without waiting and // TryAcquireLLM acquires one in-flight LLM slot without waiting and
// reports whether it succeeded. It fails when an image job is active or // reports whether it succeeded. It fails when an image job is active or
// pending. // pending or an external hold is set.
func (l *Lock) TryAcquireLLM() bool { func (l *Lock) TryAcquireLLM() bool {
l.mu.Lock() l.mu.Lock()
if l.imageActive || len(l.imageQ) > 0 { if l.imageActive || len(l.imageQ) > 0 || l.external != "" {
l.mu.Unlock() l.mu.Unlock()
return false return false
} }
@@ -120,7 +153,7 @@ func (l *Lock) AcquireImage(ctx context.Context) error {
for { for {
l.mu.Lock() l.mu.Lock()
if l.imageQ[0].id == w.id && l.n == 0 && !l.imageActive { if l.imageQ[0].id == w.id && l.n == 0 && !l.imageActive && l.external == "" {
l.imageQ = l.imageQ[1:] l.imageQ = l.imageQ[1:]
l.imageActive = true l.imageActive = true
l.mu.Unlock() l.mu.Unlock()
@@ -166,6 +199,8 @@ func (l *Lock) Snapshot() (state State, llmInflight int, imagePending bool) {
state = StateImage state = StateImage
case l.n > 0: case l.n > 0:
state = StateLLM state = StateLLM
case l.external != "":
state = StateExternal
default: default:
state = StateIdle state = StateIdle
} }
+67
View File
@@ -226,3 +226,70 @@ func TestRace(t *testing.T) {
t.Fatalf("leaked lock state: state=%s n=%d pending=%v", state, n, pending) t.Fatalf("leaked lock state: state=%s n=%d pending=%v", state, n, pending)
} }
} }
func TestExternalHoldBlocksBoth(t *testing.T) {
lk := New(nil)
ctx := context.Background()
lk.SetExternal("game.exe (pid 42)")
if lk.TryAcquireLLM() {
t.Fatal("TryAcquireLLM succeeded during external hold")
}
if got := lk.External(); got != "game.exe (pid 42)" {
t.Fatalf("External() = %q", got)
}
if state, _, _ := lk.Snapshot(); state != StateExternal {
t.Fatalf("state=%s, want external", state)
}
llmAcquired := make(chan struct{})
go func() {
if err := lk.AcquireLLM(ctx); err != nil {
t.Error(err)
}
close(llmAcquired)
}()
assertBlocked(t, llmAcquired, "LLM acquire during external hold")
imageAcquired := make(chan struct{})
go func() {
if err := lk.AcquireImage(ctx); err != nil {
t.Error(err)
}
close(imageAcquired)
}()
assertBlocked(t, imageAcquired, "image acquire during external hold")
lk.ClearExternal()
// The queued image job wins over the LLM waiter (image priority).
waitFor(t, imageAcquired, "image acquire after ClearExternal")
assertBlocked(t, llmAcquired, "LLM acquire while image active")
lk.ReleaseImage()
waitFor(t, llmAcquired, "LLM acquire after image release")
lk.ReleaseLLM()
}
func TestExternalHoldDoesNotPreempt(t *testing.T) {
lk := New(nil)
ctx := context.Background()
if err := lk.AcquireLLM(ctx); err != nil {
t.Fatal(err)
}
lk.SetExternal("game.exe")
// In-flight LLM work keeps the llm state; the hold blocks new grants.
if state, _, _ := lk.Snapshot(); state != StateLLM {
t.Fatalf("state=%s, want llm while work in flight", state)
}
if lk.TryAcquireLLM() {
t.Fatal("TryAcquireLLM succeeded during external hold")
}
lk.ReleaseLLM()
if state, _, _ := lk.Snapshot(); state != StateExternal {
t.Fatalf("state=%s, want external after drain", state)
}
lk.ClearExternal()
if state, _, _ := lk.Snapshot(); state != StateIdle {
t.Fatalf("state=%s, want idle", state)
}
}
+1 -1
View File
@@ -89,7 +89,7 @@ func (m *Metrics) Render(w io.Writer, state string, llmInflight int, imagePendin
fmt.Fprint(w, `# HELP gpu_turnstile_state Current GPU state (1 for the active state). fmt.Fprint(w, `# HELP gpu_turnstile_state Current GPU state (1 for the active state).
# TYPE gpu_turnstile_state gauge # TYPE gpu_turnstile_state gauge
`) `)
for _, s := range []string{"idle", "llm", "image"} { for _, s := range []string{"idle", "llm", "image", "external"} {
v := 0 v := 0
if s == state { if s == state {
v = 1 v = 1
+62 -3
View File
@@ -26,6 +26,7 @@ import (
"gpu-turnstile/internal/lock" "gpu-turnstile/internal/lock"
"gpu-turnstile/internal/metrics" "gpu-turnstile/internal/metrics"
"gpu-turnstile/internal/ollama" "gpu-turnstile/internal/ollama"
"gpu-turnstile/internal/supervise"
) )
// defaultCaptureLimit bounds how much of a /prompt response body is // defaultCaptureLimit bounds how much of a /prompt response body is
@@ -44,6 +45,11 @@ type Config struct {
Metrics *metrics.Metrics Metrics *metrics.Metrics
Log *slog.Logger Log *slog.Logger
// ComfySup, when non-nil, is the managed ComfyUI process: any ComfyUI
// request starts it on demand, /prompt additionally waits for
// readiness before taking the GPU lock. nil = unmanaged upstream.
ComfySup *supervise.Process
// LogColor enables ANSI colors in per-request log lines. Ignored when // LogColor enables ANSI colors in per-request log lines. Ignored when
// the log level is above INFO (request lines are not emitted at all). // the log level is above INFO (request lines are not emitted at all).
LogColor bool LogColor bool
@@ -438,7 +444,7 @@ func isLLMRequest(r *http.Request) bool {
return r.Method == http.MethodPost && llmPaths[r.URL.Path] return r.Method == http.MethodPost && llmPaths[r.URL.Path]
} }
// OllamaHandler serves the Ollama-facing listener. // OllamaHandler serves the listener for Ollama-compatible clients.
func (s *Server) OllamaHandler() http.Handler { func (s *Server) OllamaHandler() http.Handler {
return s.logRequests("ollama", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { return s.logRequests("ollama", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path { switch r.URL.Path {
@@ -458,10 +464,14 @@ func (s *Server) OllamaHandler() http.Handler {
if s.busyMode == "reject" { if s.busyMode == "reject" {
if !s.cfg.Lock.TryAcquireLLM() { if !s.cfg.Lock.TryAcquireLLM() {
s.cfg.Metrics.ObserveLockWait("llm", time.Since(start).Seconds()) s.cfg.Metrics.ObserveLockWait("llm", time.Since(start).Seconds())
msg := "GPU busy: image job active or queued"
if holder := s.cfg.Lock.External(); holder != "" {
msg = "GPU busy: " + holder
}
s.log.Info("llm request rejected; GPU busy", s.log.Info("llm request rejected; GPU busy",
"path", r.URL.Path, "status", s.busyStatus) "path", r.URL.Path, "status", s.busyStatus)
w.Header().Set("Retry-After", strconv.Itoa(s.busyRetryAfter)) w.Header().Set("Retry-After", strconv.Itoa(s.busyRetryAfter))
http.Error(w, "GPU busy: image job active or queued", s.busyStatus) http.Error(w, msg, s.busyStatus)
return return
} }
s.cfg.Metrics.ObserveLockWait("llm", time.Since(start).Seconds()) s.cfg.Metrics.ObserveLockWait("llm", time.Since(start).Seconds())
@@ -485,7 +495,7 @@ func (s *Server) OllamaHandler() http.Handler {
})) }))
} }
// ComfyHandler serves the ComfyUI-facing listener. // ComfyHandler serves the listener for ComfyUI clients.
func (s *Server) ComfyHandler() http.Handler { func (s *Server) ComfyHandler() http.Handler {
return s.logRequests("comfy", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { return s.logRequests("comfy", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path { switch r.URL.Path {
@@ -500,6 +510,21 @@ func (s *Server) ComfyHandler() http.Handler {
s.handlePrompt(w, r) s.handlePrompt(w, r)
return return
} }
if s.cfg.ComfySup != nil {
// Any other ComfyUI request also wakes the managed server; the
// retry backoff bridges the time it needs to come up. While a
// foreign process holds the GPU we refuse to spawn it — the
// request gets the busy answer instead of fighting for VRAM.
if holder := s.cfg.Lock.External(); holder != "" && !s.cfg.ComfySup.Running() {
w.Header().Set("Retry-After", strconv.Itoa(s.busyRetryAfter))
http.Error(w, "GPU busy: "+holder, http.StatusServiceUnavailable)
return
}
if err := s.cfg.ComfySup.EnsureRunning(); err != nil {
http.Error(w, fmt.Sprintf("cannot start ComfyUI: %v", err), http.StatusBadGateway)
return
}
}
s.comfyProxy.ServeHTTP(w, r) s.comfyProxy.ServeHTTP(w, r)
})) }))
} }
@@ -539,6 +564,23 @@ func (w *captureWriter) Unwrap() http.ResponseWriter { return w.ResponseWriter }
func (s *Server) handlePrompt(w http.ResponseWriter, r *http.Request) { func (s *Server) handlePrompt(w http.ResponseWriter, r *http.Request) {
log := s.log.With("op", "image") log := s.log.With("op", "image")
// Normally the managed server is brought up *before* taking the GPU
// lock: torch can take a minute to load, and LLM traffic should keep
// flowing in the meantime. While a foreign process holds the GPU, LLM
// traffic is blocked anyway and a fresh ComfyUI would fight it for
// VRAM — so the lock comes first in that case.
comfyFirst := s.cfg.ComfySup != nil && s.cfg.Lock.External() == ""
if comfyFirst {
if err := s.cfg.ComfySup.EnsureRunning(); err != nil {
http.Error(w, fmt.Sprintf("cannot start ComfyUI: %v", err), http.StatusBadGateway)
return
}
if err := s.cfg.ComfySup.WaitReady(r.Context()); err != nil {
http.Error(w, fmt.Sprintf("ComfyUI did not become ready: %v", err), http.StatusBadGateway)
return
}
}
start := time.Now() start := time.Now()
if err := s.cfg.Lock.AcquireImage(r.Context()); err != nil { if err := s.cfg.Lock.AcquireImage(r.Context()); err != nil {
if errors.Is(err, context.DeadlineExceeded) { if errors.Is(err, context.DeadlineExceeded) {
@@ -549,6 +591,19 @@ func (s *Server) handlePrompt(w http.ResponseWriter, r *http.Request) {
s.cfg.Metrics.ObserveLockWait("image", time.Since(start).Seconds()) s.cfg.Metrics.ObserveLockWait("image", time.Since(start).Seconds())
log.Info("image lock acquired") log.Info("image lock acquired")
if s.cfg.ComfySup != nil && !comfyFirst {
if err := s.cfg.ComfySup.EnsureRunning(); err != nil {
s.cfg.Lock.ReleaseImage()
http.Error(w, fmt.Sprintf("cannot start ComfyUI: %v", err), http.StatusBadGateway)
return
}
if err := s.cfg.ComfySup.WaitReady(r.Context()); err != nil {
s.cfg.Lock.ReleaseImage()
http.Error(w, fmt.Sprintf("ComfyUI did not become ready: %v", err), http.StatusBadGateway)
return
}
}
if s.cfg.Ollama != nil { if s.cfg.Ollama != nil {
uctx, ucancel := context.WithTimeout(r.Context(), s.cfg.UnloadTimeout) uctx, ucancel := context.WithTimeout(r.Context(), s.cfg.UnloadTimeout)
elapsed, uerr := s.cfg.Ollama.UnloadAll(uctx) elapsed, uerr := s.cfg.Ollama.UnloadAll(uctx)
@@ -608,6 +663,10 @@ func (s *Server) finishImageJob(promptID string) {
s.cfg.Lock.ReleaseImage() s.cfg.Lock.ReleaseImage()
log.Info("image lock released") log.Info("image lock released")
if s.cfg.ComfySup != nil {
// The idle clock starts when the job ends, not when it began.
s.cfg.ComfySup.NoteActivity()
}
if s.cfg.Ollama != nil && s.cfg.WarmModel != "" { if s.cfg.Ollama != nil && s.cfg.WarmModel != "" {
if state, _, _ := s.cfg.Lock.Snapshot(); state == lock.StateIdle { if state, _, _ := s.cfg.Lock.Snapshot(); state == lock.StateIdle {
+2 -2
View File
@@ -46,8 +46,8 @@ type fakes struct {
rec *recorder rec *recorder
ollama *httptest.Server ollama *httptest.Server
comfy *httptest.Server comfy *httptest.Server
server *httptest.Server // Ollama-facing gpu-turnstile listener server *httptest.Server // gpu-turnstile listener for Ollama-compatible clients
comfySrv *httptest.Server // ComfyUI-facing gpu-turnstile listener comfySrv *httptest.Server // gpu-turnstile listener for ComfyUI clients
freeCh chan struct{} freeCh chan struct{}
chatCh chan struct{} chatCh chan struct{}
historyMu sync.Mutex historyMu sync.Mutex
+314
View File
@@ -0,0 +1,314 @@
// Package supervise runs an upstream server (ComfyUI) as a managed child
// process: started on demand when a request needs it, stopped after an
// idle timeout so the GPU memory it holds is freed, and stopped with the
// parent. Crashes are logged; the next request respawns it.
package supervise
import (
"context"
"fmt"
"io"
"log/slog"
"net/url"
"os"
"os/exec"
"path/filepath"
"runtime"
"strings"
"sync"
"time"
)
// ComfyLayout returns the interpreter and script path of a standard ComfyUI
// venv install rooted at dir for the given GOOS: .venv\Scripts\python.exe
// on Windows, .venv/bin/python elsewhere. The script is main.py — either in
// a ComfyUI subdirectory or directly under dir, whichever exists (the
// subdirectory form wins ties and is the default when neither exists yet,
// so the caller's missing-file warning points at the documented layout).
func ComfyLayout(goos, dir string) (python, script string) {
if goos == "windows" {
python = filepath.Join(dir, ".venv", "Scripts", "python.exe")
} else {
python = filepath.Join(dir, ".venv", "bin", "python")
}
script = filepath.Join(dir, "ComfyUI", "main.py")
if _, err := os.Stat(script); err != nil {
if _, err := os.Stat(filepath.Join(dir, "main.py")); err == nil {
script = filepath.Join(dir, "main.py")
}
}
return python, script
}
// DefaultComfyCommand builds the launch command for the standard venv
// layout (see ComfyLayout): the script is passed relative to dir so dir
// stays the working directory, and --port is taken from comfyURL when the
// URL carries one.
func DefaultComfyCommand(goos, dir, comfyURL string) string {
python, script := ComfyLayout(goos, dir)
rel, err := filepath.Rel(dir, script)
if err != nil {
rel = script
}
cmd := `"` + python + `" ` + rel
if u, err := url.Parse(comfyURL); err == nil && u.Port() != "" {
cmd += " --port " + u.Port()
}
return cmd
}
// Process is one managed child process.
type Process struct {
name string
argv []string
dir string
probe func(context.Context) error
startTimeout time.Duration
log *slog.Logger
mu sync.Mutex
cmd *exec.Cmd
stopping bool
ready bool
external bool // someone else serves the port; not our process
lastActivity time.Time
}
// New parses cmdLine (double quotes group arguments containing spaces) and
// prepares a managed process. probe reports whether the server answers
// (e.g. the comfy client's Probe); startTimeout bounds WaitReady. dir is
// the child's working directory; empty inherits ours.
func New(name, cmdLine, dir string, probe func(context.Context) error, startTimeout time.Duration, log *slog.Logger) (*Process, error) {
argv, err := splitCommandLine(cmdLine)
if err != nil {
return nil, fmt.Errorf("%s command: %w", name, err)
}
if len(argv) == 0 {
return nil, fmt.Errorf("%s command is empty", name)
}
if log == nil {
log = slog.Default()
}
if startTimeout <= 0 {
startTimeout = 2 * time.Minute
}
return &Process{name: name, argv: argv, dir: dir, probe: probe, startTimeout: startTimeout, log: log}, nil
}
// splitCommandLine splits a command line on whitespace, treating
// double-quoted sections as one argument (quotes removed). Backslashes are
// literal — this matches Windows paths.
func splitCommandLine(s string) ([]string, error) {
var argv []string
var cur strings.Builder
inQuote := false
have := false
flush := func() {
if have {
argv = append(argv, cur.String())
cur.Reset()
have = false
}
}
for _, r := range s {
switch {
case r == '"':
inQuote = !inQuote
have = true
case (r == ' ' || r == '\t') && !inQuote:
flush()
default:
cur.WriteRune(r)
have = true
}
}
if inQuote {
return nil, fmt.Errorf("unterminated quote in %q", s)
}
flush()
return argv, nil
}
// Running reports whether the child process is currently alive.
func (p *Process) Running() bool {
p.mu.Lock()
defer p.mu.Unlock()
return p.cmd != nil
}
// Ready reports whether the server has answered a probe since its last
// (re)start. Health checks use it to tell "starting up" from "outage".
func (p *Process) Ready() bool {
p.mu.Lock()
defer p.mu.Unlock()
return p.ready
}
// MarkReady records that the server answered.
func (p *Process) MarkReady() {
p.mu.Lock()
p.ready = true
p.mu.Unlock()
}
// NoteActivity resets the idle clock; called for every request served.
func (p *Process) NoteActivity() {
p.mu.Lock()
p.lastActivity = time.Now()
p.mu.Unlock()
}
// EnsureRunning starts the child if it is not running. It returns as soon
// as the process is spawned; readiness is WaitReady's job (and the proxy's
// retry backoff bridges the gap for plain proxied requests). When the URL
// already answers — e.g. the ComfyUI desktop app grabbed the port — no
// child is spawned: the external server is used as-is, and the idle
// watcher never touches it (it only kills its own child).
func (p *Process) EnsureRunning() error {
p.mu.Lock()
defer p.mu.Unlock()
p.lastActivity = time.Now()
if p.cmd != nil {
return nil
}
pctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
err := p.probe(pctx)
cancel()
if err == nil {
p.ready = true
if !p.external {
p.external = true
p.log.Info(p.name + " is already served externally; not spawning a managed instance")
}
return nil
}
p.external = false
cmd := exec.Command(p.argv[0], p.argv[1:]...)
cmd.Dir = p.dir
stdout, err := cmd.StdoutPipe()
if err != nil {
return err
}
stderr, err := cmd.StderrPipe()
if err != nil {
return err
}
if err := cmd.Start(); err != nil {
return fmt.Errorf("start %s: %w", p.name, err)
}
p.cmd = cmd
p.stopping = false
p.ready = false
go p.pipeLog(stdout)
go p.pipeLog(stderr)
go func() {
err := cmd.Wait()
p.mu.Lock()
p.cmd = nil
p.ready = false
intentional := p.stopping
p.mu.Unlock()
if intentional {
p.log.Info(p.name + " stopped")
} else {
p.log.Warn(p.name+" exited unexpectedly; the next request restarts it", "err", err)
}
}()
p.log.Info(p.name+" starting", "pid", cmd.Process.Pid, "cmd", strings.Join(p.argv, " "))
return nil
}
// WaitReady blocks until the probe succeeds, ctx ends, or the start
// timeout passes.
func (p *Process) WaitReady(ctx context.Context) error {
ctx, cancel := context.WithTimeout(ctx, p.startTimeout)
defer cancel()
for {
pctx, pcancel := context.WithTimeout(ctx, 5*time.Second)
err := p.probe(pctx)
pcancel()
if err == nil {
p.NoteActivity()
p.MarkReady()
return nil
}
select {
case <-ctx.Done():
return fmt.Errorf("%s did not become ready: %w", p.name, ctx.Err())
case <-time.After(500 * time.Millisecond):
}
}
}
// Stop kills the child process (the whole tree on Windows) if running.
func (p *Process) Stop() {
p.mu.Lock()
cmd := p.cmd
if cmd == nil {
p.mu.Unlock()
return
}
p.stopping = true
p.ready = false
p.mu.Unlock()
stopTree(cmd)
}
// WatchIdle stops the child after idleTimeout without activity, but only
// when gpuIdle reports the GPU lock is free (no active or pending work).
// Returns when ctx ends.
func (p *Process) WatchIdle(ctx context.Context, idleTimeout time.Duration, gpuIdle func() bool) {
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
}
p.mu.Lock()
idleFor := time.Since(p.lastActivity)
running := p.cmd != nil
p.mu.Unlock()
if running && idleFor > idleTimeout && gpuIdle() {
p.log.Info(p.name+" idle; stopping to free the GPU", "idle_for", idleFor.Round(time.Second))
p.Stop()
}
}
}
// pipeLog forwards one child output stream to the log at INFO, line by
// line, prefixed with the process name.
func (p *Process) pipeLog(r io.Reader) {
buf := make([]byte, 4096)
var line string
for {
n, err := r.Read(buf)
line += string(buf[:n])
for {
i := strings.IndexByte(line, '\n')
if i < 0 {
break
}
p.log.Info(p.name + ": " + strings.TrimRight(line[:i], "\r"))
line = line[i+1:]
}
if err != nil {
if strings.TrimSpace(line) != "" {
p.log.Info(p.name + ": " + line)
}
return
}
}
}
// stopTree kills cmd's process, including its children on Windows (python
// launchers tend to spawn some). The Wait goroutine reaps it.
func stopTree(cmd *exec.Cmd) {
if runtime.GOOS == "windows" {
exec.Command("taskkill", "/T", "/F", "/PID",
fmt.Sprint(cmd.Process.Pid)).Run() //nolint:errcheck // best effort
return
}
cmd.Process.Kill() //nolint:errcheck // best effort
}
+241
View File
@@ -0,0 +1,241 @@
package supervise
import (
"context"
"errors"
"log/slog"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"testing"
"time"
)
func TestSplitCommandLine(t *testing.T) {
cases := []struct {
in string
want []string
}{
{`python main.py --port 8188`, []string{"python", "main.py", "--port", "8188"}},
{`"C:\Program Files\py\python.exe" main.py`, []string{`C:\Program Files\py\python.exe`, "main.py"}},
{` spaced out `, []string{"spaced", "out"}},
{`a "b c" d`, []string{"a", "b c", "d"}},
{"", nil},
}
for _, c := range cases {
got, err := splitCommandLine(c.in)
if err != nil {
t.Errorf("splitCommandLine(%q): %v", c.in, err)
continue
}
if len(got) != len(c.want) {
t.Errorf("splitCommandLine(%q) = %v, want %v", c.in, got, c.want)
continue
}
for i := range got {
if got[i] != c.want[i] {
t.Errorf("splitCommandLine(%q) = %v, want %v", c.in, got, c.want)
break
}
}
}
if _, err := splitCommandLine(`"unterminated`); err == nil {
t.Error("expected error for unterminated quote")
}
}
// TestHelperProcess is the child executed by the lifecycle tests: it just
// sleeps. The env marker is set only around the spawn, so in the normal
// test run this returns immediately.
func TestHelperProcess(t *testing.T) {
if os.Getenv("GO_HELPER_PROCESS") != "1" {
return
}
time.Sleep(30 * time.Second)
os.Exit(0)
}
func newHelper(t *testing.T, name string) *Process {
t.Helper()
// The probe always fails: nothing external serves the port, so
// EnsureRunning spawns the helper child.
p, err := New(name, `"`+os.Args[0]+`" -test.run=TestHelperProcess`, "",
func(context.Context) error { return errors.New("nothing there") }, 5*time.Second, slog.Default())
if err != nil {
t.Fatal(err)
}
return p
}
// startHelper spawns the child with the marker set; exec.Command inherits
// the environment at spawn time, so it can be unset right after.
func startHelper(t *testing.T, p *Process) {
t.Helper()
os.Setenv("GO_HELPER_PROCESS", "1")
defer os.Unsetenv("GO_HELPER_PROCESS")
if err := p.EnsureRunning(); err != nil {
t.Fatal(err)
}
}
func waitStopped(t *testing.T, p *Process, timeout time.Duration) {
t.Helper()
deadline := time.Now().Add(timeout)
for p.Running() && time.Now().Before(deadline) {
time.Sleep(50 * time.Millisecond)
}
if p.Running() {
t.Fatal("process still running")
}
}
func TestEnsureRunningAndStop(t *testing.T) {
p := newHelper(t, "helper")
if p.Running() {
t.Fatal("Running before start")
}
startHelper(t, p)
if !p.Running() {
t.Fatal("not Running after EnsureRunning")
}
if err := p.EnsureRunning(); err != nil {
t.Fatal("second EnsureRunning must be a no-op")
}
p.Stop()
waitStopped(t, p, 5*time.Second)
}
func TestEnsureRunningPrefersExternalServer(t *testing.T) {
// The port is already served (e.g. the ComfyUI desktop app): no child
// is spawned, the supervisor reports ready, and Stop is a no-op.
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
}))
defer srv.Close()
probe := func(ctx context.Context) error {
req, _ := http.NewRequestWithContext(ctx, http.MethodGet, srv.URL, nil)
resp, err := http.DefaultClient.Do(req)
if err != nil {
return err
}
resp.Body.Close()
return nil
}
p, err := New("external", `"`+os.Args[0]+`"`, "", probe, 5*time.Second, slog.Default())
if err != nil {
t.Fatal(err)
}
if err := p.EnsureRunning(); err != nil {
t.Fatal(err)
}
if p.Running() {
t.Fatal("spawned a child even though the port is already served")
}
if !p.Ready() {
t.Fatal("external server should count as ready")
}
p.Stop() // must not touch the external server
}
func TestWaitReady(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
}))
defer srv.Close()
probe := func(ctx context.Context) error {
req, _ := http.NewRequestWithContext(ctx, http.MethodGet, srv.URL, nil)
resp, err := http.DefaultClient.Do(req)
if err != nil {
return err
}
resp.Body.Close()
return nil
}
p, err := New("ready", `"`+os.Args[0]+`"`, "", probe, 5*time.Second, slog.Default())
if err != nil {
t.Fatal(err)
}
if err := p.WaitReady(context.Background()); err != nil {
t.Fatal(err)
}
failing, err := New("failing", `"`+os.Args[0]+`"`, "",
func(context.Context) error { return errors.New("no") }, 500*time.Millisecond, slog.Default())
if err != nil {
t.Fatal(err)
}
if err := failing.WaitReady(context.Background()); err == nil {
t.Fatal("expected timeout error from WaitReady")
}
}
func TestWatchIdleStops(t *testing.T) {
p := newHelper(t, "idle")
startHelper(t, p)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go p.WatchIdle(ctx, 200*time.Millisecond, func() bool { return true })
waitStopped(t, p, 10*time.Second)
}
func TestWatchIdleRespectsBusyGPU(t *testing.T) {
p := newHelper(t, "busy")
startHelper(t, p)
defer p.Stop()
ctx, cancel := context.WithCancel(context.Background())
go p.WatchIdle(ctx, 100*time.Millisecond, func() bool { return false })
time.Sleep(600 * time.Millisecond)
cancel()
if !p.Running() {
t.Fatal("process was stopped while the GPU was busy")
}
}
func TestComfyLayoutAndDefaultCommand(t *testing.T) {
// Nested layout (ComfyUI/main.py under dir) wins.
dir := t.TempDir()
nested := filepath.Join(dir, "ComfyUI", "main.py")
if err := os.MkdirAll(filepath.Dir(nested), 0o755); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(nested, []byte("x"), 0o644); err != nil {
t.Fatal(err)
}
python, script := ComfyLayout("windows", dir)
if want := filepath.Join(dir, ".venv", "Scripts", "python.exe"); python != want {
t.Errorf("python = %s, want %s", python, want)
}
if script != nested {
t.Errorf("script = %s, want %s", script, nested)
}
cmd := DefaultComfyCommand("windows", dir, "http://127.0.0.1:8189")
want := `"` + filepath.Join(dir, ".venv", "Scripts", "python.exe") + `" ` + filepath.Join("ComfyUI", "main.py") + " --port 8189"
if cmd != want {
t.Errorf("cmd = %q, want %q", cmd, want)
}
// Flat layout (main.py directly under dir) is found too.
flat := t.TempDir()
if err := os.WriteFile(filepath.Join(flat, "main.py"), []byte("x"), 0o644); err != nil {
t.Fatal(err)
}
if _, script := ComfyLayout("linux", flat); script != filepath.Join(flat, "main.py") {
t.Errorf("flat script = %s", script)
}
cmd = DefaultComfyCommand("linux", flat, "http://comfy.internal")
if strings.Contains(cmd, "--port") {
t.Errorf("cmd = %q, want no --port for a port-less URL", cmd)
}
if !strings.HasSuffix(cmd, `" main.py`) {
t.Errorf("cmd = %q, want quoted python + relative main.py", cmd)
}
// Neither exists yet: default to the documented nested form so the
// startup warning points there.
empty := t.TempDir()
if _, script := ComfyLayout("windows", empty); script != filepath.Join(empty, "ComfyUI", "main.py") {
t.Errorf("missing-layout script = %s", script)
}
}