Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
482731ed9e | ||
|
|
e4bdc92ece | ||
|
|
14120bf4a4 | ||
|
|
d7566329ae | ||
|
|
e43ad02fc4 | ||
|
|
232f5b61f2 | ||
|
|
5e7a042cad | ||
|
|
14c2e30478 |
@@ -4,3 +4,6 @@
|
||||
/compose.yml
|
||||
/signing/
|
||||
|
||||
/gpu-turnstile.exe.old
|
||||
/gpu-turnstile.exe.new
|
||||
/gpu-turnstile.exe~
|
||||
|
||||
@@ -2,9 +2,10 @@
|
||||
|
||||
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
|
||||
front of both and guarantees the GPU is always in exactly one of three states:
|
||||
`idle`, `llm` (N ≥ 1 Ollama requests in flight), or `image` (exactly one
|
||||
ComfyUI job, Ollama models unloaded). See [SPEC.md](SPEC.md) for the full
|
||||
front of both and guarantees the GPU is always in exactly one of four states:
|
||||
`idle`, `llm` (N ≥ 1 Ollama requests in flight), `image` (exactly one
|
||||
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.
|
||||
|
||||
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
|
||||
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 Ollama unload/warm steps are skipped. Future consumers (e.g. local game
|
||||
detection) plug into the same lock the same way.
|
||||
the Ollama unload/warm steps are skipped. A third, URL-less consumer —
|
||||
detection of foreign GPU holders such as games — is enabled by `GAME_PROCS`
|
||||
and/or `GPU_FOREIGN_VRAM_MB` (see below).
|
||||
|
||||
## Configuration
|
||||
|
||||
@@ -45,8 +47,8 @@ override file values. Invalid values fail at startup.
|
||||
|
||||
| Var | Default | Meaning |
|
||||
|---|---|---|
|
||||
| `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener |
|
||||
| `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener |
|
||||
| `LISTEN_OLLAMA` | `:11434` | Listener for Ollama-compatible clients |
|
||||
| `LISTEN_COMFY` | `:8188` | Listener for ComfyUI clients |
|
||||
| `OLLAMA_URL` | _(empty = disabled)_ | Ollama upstream; set to enable the Ollama 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 |
|
||||
@@ -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 (400–599, e.g. 429) |
|
||||
| `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) |
|
||||
| `COMFY_CMD` | _(empty = unmanaged)_ | Supervise ComfyUI: start on demand, stop when idle to free VRAM. Requires `COMFY_URL` |
|
||||
| `COMFY_DIR` | _(empty)_ | 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` | `5s` | How often game/VRAM detection runs |
|
||||
| `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_FILE` | _(empty)_ | Append logs to this file instead of stderr |
|
||||
| `UNLOAD_POLL_INTERVAL` | `500ms` | `/api/ps` poll interval while unloading |
|
||||
| `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 |
|
||||
| `WARM_TIMEOUT` | `2m` | Warm-model reload after an image job |
|
||||
| `SHUTDOWN_TIMEOUT` | `10s` | Graceful shutdown on SIGINT/SIGTERM |
|
||||
@@ -76,7 +87,7 @@ override file values. Invalid values fail at startup.
|
||||
|
||||
## 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`,
|
||||
`gpu_turnstile_llm_inflight`, `gpu_turnstile_image_pending`,
|
||||
`gpu_turnstile_image_jobs_total`, `gpu_turnstile_lock_wait_seconds`
|
||||
@@ -88,6 +99,52 @@ override file values. Invalid values fail at startup.
|
||||
class) in text mode, which renders in `docker compose logs` on Windows
|
||||
Terminal. Set `NO_COLOR` to disable colors.
|
||||
|
||||
## Managed ComfyUI (`COMFY_CMD`)
|
||||
|
||||
Don't want ComfyUI running 24/7 (it holds VRAM even when idle — and the
|
||||
Desktop app kills its server when you close it)? Point `COMFY_CMD` at a
|
||||
standalone launch command and gpu-turnstile supervises 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. Example for a Desktop
|
||||
install (run it once manually to confirm it works):
|
||||
|
||||
```
|
||||
COMFY_URL=http://127.0.0.1:8188
|
||||
COMFY_CMD="C:\ComfyUI\.venv\Scripts\python.exe ComfyUI\main.py --port 8188"
|
||||
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` (5s):
|
||||
|
||||
```
|
||||
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
|
||||
|
||||
```sh
|
||||
|
||||
@@ -10,11 +10,12 @@ becomes very slow; on Linux it would OOM instead.
|
||||
## Goal
|
||||
|
||||
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
|
||||
- `llm` — N ≥ 1 Ollama inference requests in flight (concurrency allowed)
|
||||
- `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
|
||||
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.
|
||||
- **Only `COMFY_URL`**: image jobs are tracked and ComfyUI's VRAM is freed
|
||||
afterwards, but the Ollama unload and warm-reload steps are skipped.
|
||||
- Future consumers (e.g. detecting a local game holding VRAM) plug into the
|
||||
same lock the same way: enabled by their config knob, excluded when
|
||||
absent.
|
||||
- **Game detection** is a third, optional consumer without a URL: enabled by
|
||||
`GAME_PROCS` and/or `GPU_FOREIGN_VRAM_MB` it watches for foreign processes
|
||||
holding the GPU (see below) and plugs into the same lock the same way —
|
||||
excluded when both knobs are unset.
|
||||
|
||||
### 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
|
||||
the ComfyUI job finished and models were freed.
|
||||
- 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
|
||||
removed from the queue.
|
||||
|
||||
@@ -121,6 +128,68 @@ 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
|
||||
load time. Off by default.
|
||||
|
||||
### Managed ComfyUI (`COMFY_CMD`)
|
||||
|
||||
When `COMFY_CMD` is set, gpu-turnstile runs ComfyUI as a supervised child
|
||||
process instead of expecting an always-on server:
|
||||
|
||||
- **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 5 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 comes from environment variables and/or an `.env`-style
|
||||
@@ -131,8 +200,8 @@ override file values. A missing file is fine; a malformed one is fatal.
|
||||
|
||||
| Var | Default | Meaning |
|
||||
|---|---|---|
|
||||
| `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener |
|
||||
| `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener |
|
||||
| `LISTEN_OLLAMA` | `:11434` | listener for Ollama-compatible clients |
|
||||
| `LISTEN_COMFY` | `:8188` | listener for ComfyUI clients |
|
||||
| `OLLAMA_URL` | _(empty = disabled)_ | Ollama upstream; set to enable the Ollama consumer |
|
||||
| `COMFY_URL` | _(empty = disabled)_ | ComfyUI upstream; set to enable the ComfyUI consumer |
|
||||
| `UNLOAD_TIMEOUT` | `60s` | wait for Ollama to unload |
|
||||
@@ -142,12 +211,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 (400–599, e.g. 429) |
|
||||
| `BUSY_RETRY_AFTER` | `30` | seconds sent as `Retry-After` on busy responses (both modes) |
|
||||
| `WARM_MODEL` | `` | optional model to reload after an image job |
|
||||
| `COMFY_CMD` | _(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` | `` | 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` | `5s` | how often game/VRAM detection runs |
|
||||
| `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_FILE` | `` | append logs to this file instead of stderr (useful as a service) |
|
||||
| `UNLOAD_POLL_INTERVAL` | `500ms` | `/api/ps` poll interval while unloading |
|
||||
| `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 |
|
||||
| `WARM_TIMEOUT` | `2m` | warm-model reload after an image job |
|
||||
| `SHUTDOWN_TIMEOUT` | `10s` | graceful shutdown on SIGINT/SIGTERM |
|
||||
@@ -159,7 +237,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_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) |
|
||||
| `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`) |
|
||||
|
||||
Startup fails fast on unparsable values and when neither consumer URL is
|
||||
set. Enabled upstreams are probed once at start (`/api/version`,
|
||||
|
||||
+192
-7
@@ -21,11 +21,13 @@ import (
|
||||
|
||||
"gpu-turnstile/internal/comfy"
|
||||
"gpu-turnstile/internal/config"
|
||||
"gpu-turnstile/internal/game"
|
||||
"gpu-turnstile/internal/lock"
|
||||
"gpu-turnstile/internal/metrics"
|
||||
"gpu-turnstile/internal/ollama"
|
||||
"gpu-turnstile/internal/proxy"
|
||||
"gpu-turnstile/internal/service"
|
||||
"gpu-turnstile/internal/supervise"
|
||||
"gpu-turnstile/internal/update"
|
||||
)
|
||||
|
||||
@@ -428,6 +430,7 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
|
||||
"unload_poll_interval", cfg.UnloadPollInterval,
|
||||
"history_poll_interval", cfg.HistoryPollInterval,
|
||||
"probe_timeout", cfg.ProbeTimeout,
|
||||
"health_interval", cfg.HealthInterval,
|
||||
"free_timeout", cfg.FreeTimeout,
|
||||
"warm_timeout", cfg.WarmTimeout,
|
||||
"shutdown_timeout", cfg.ShutdownTimeout,
|
||||
@@ -435,6 +438,14 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
|
||||
"backoff_max", cfg.BackoffMax,
|
||||
"prompt_capture_limit", cfg.PromptCaptureLimit,
|
||||
"warm_model", cfg.WarmModel,
|
||||
"comfy_cmd", cfg.ComfyCmd,
|
||||
"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,
|
||||
"update_interval", cfg.UpdateInterval,
|
||||
"update_repo", cfg.UpdateRepo,
|
||||
@@ -461,12 +472,31 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
|
||||
}
|
||||
}
|
||||
|
||||
// With COMFY_CMD set, ComfyUI runs as a managed child: started on
|
||||
// demand by the proxy, stopped after COMFY_IDLE_TIMEOUT idle (and on
|
||||
// shutdown) so its VRAM is freed.
|
||||
var comfySup *supervise.Process
|
||||
if cfg.ComfyCmd != "" {
|
||||
var err error
|
||||
comfySup, err = supervise.New("comfy", cfg.ComfyCmd, 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{
|
||||
OllamaURL: cfg.OllamaURL,
|
||||
ComfyURL: cfg.ComfyURL,
|
||||
Lock: lk,
|
||||
Ollama: ollamaClient,
|
||||
Comfy: comfyClient,
|
||||
ComfySup: comfySup,
|
||||
Metrics: metrics.New(),
|
||||
Log: log,
|
||||
LogColor: !cfg.LogJSON && cfg.LogFile == "" && os.Getenv("NO_COLOR") == "",
|
||||
@@ -490,19 +520,52 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
|
||||
return err
|
||||
}
|
||||
|
||||
// Probe the enabled upstreams once; failure is logged, not fatal.
|
||||
probeCtx, probeCancel := context.WithTimeout(ctx, cfg.ProbeTimeout)
|
||||
// Probe the enabled upstreams once; failure is logged, not fatal. A
|
||||
// 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 err := ollamaClient.Probe(probeCtx); err != nil {
|
||||
log.Warn("ollama probe failed", "url", cfg.OllamaURL, "err", err)
|
||||
}
|
||||
probes["ollama"] = ollamaClient.Probe
|
||||
}
|
||||
if comfyClient != nil {
|
||||
if err := comfyClient.Probe(probeCtx); err != nil {
|
||||
log.Warn("comfy probe failed", "url", cfg.ComfyURL, "err", err)
|
||||
if comfySup == nil {
|
||||
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()
|
||||
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
|
||||
// readiness notification below really means "accepting connections".
|
||||
@@ -560,6 +623,128 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
|
||||
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.
|
||||
// In service mode a staged update is applied by exiting with exitCodeUpdate
|
||||
// once the GPU lock is idle; the service recovery configuration restarts the
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
# ComfyUI --listen 0.0.0.0 --port 8189).
|
||||
services:
|
||||
gpu-turnstile:
|
||||
image: git.rambossek.at/public/gpu-turnstile:v0.1.7
|
||||
image: git.rambossek.at/public/gpu-turnstile:v0.1.8
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
# Each consumer is enabled by setting its URL; leave one unset to
|
||||
|
||||
@@ -26,6 +26,7 @@ type Config struct {
|
||||
UnloadPollInterval time.Duration
|
||||
HistoryPollInterval time.Duration
|
||||
ProbeTimeout time.Duration
|
||||
HealthInterval time.Duration
|
||||
FreeTimeout time.Duration
|
||||
WarmTimeout time.Duration
|
||||
ShutdownTimeout time.Duration
|
||||
@@ -50,6 +51,27 @@ type Config struct {
|
||||
LLMBusyStatus int
|
||||
BusyRetryAfter int
|
||||
|
||||
// ComfyCmd spawns and supervises a ComfyUI server on demand (empty =
|
||||
// unmanaged, the current behavior). ComfyDir is its working directory.
|
||||
// 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
|
||||
LogLevel slog.Level
|
||||
LogJSON bool
|
||||
@@ -70,6 +92,7 @@ func Defaults() Config {
|
||||
UnloadPollInterval: 500 * time.Millisecond,
|
||||
HistoryPollInterval: time.Second,
|
||||
ProbeTimeout: 5 * time.Second,
|
||||
HealthInterval: 30 * time.Second,
|
||||
FreeTimeout: 30 * time.Second,
|
||||
WarmTimeout: 2 * time.Minute,
|
||||
ShutdownTimeout: 10 * time.Second,
|
||||
@@ -87,6 +110,14 @@ func Defaults() Config {
|
||||
LLMBusyStatus: 503,
|
||||
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: 5 * time.Second,
|
||||
|
||||
LogLevel: slog.LevelWarn,
|
||||
}
|
||||
}
|
||||
@@ -122,6 +153,17 @@ func ParseEnvFile(r io.Reader) (map[string]string, error) {
|
||||
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 {
|
||||
v := getenv(name)
|
||||
if v == "" {
|
||||
@@ -148,6 +190,8 @@ func Load(getenv func(string) string) (Config, error) {
|
||||
{"OLLAMA_URL", &cfg.OllamaURL},
|
||||
{"COMFY_URL", &cfg.ComfyURL},
|
||||
{"WARM_MODEL", &cfg.WarmModel},
|
||||
{"COMFY_CMD", &cfg.ComfyCmd},
|
||||
{"COMFY_DIR", &cfg.ComfyDir},
|
||||
{"UPDATE_REPO", &cfg.UpdateRepo},
|
||||
{"UPDATE_ASSET", &cfg.UpdateAsset},
|
||||
{"LOG_FILE", &cfg.LogFile},
|
||||
@@ -166,17 +210,34 @@ func Load(getenv func(string) string) (Config, error) {
|
||||
{"UNLOAD_POLL_INTERVAL", &cfg.UnloadPollInterval},
|
||||
{"HISTORY_POLL_INTERVAL", &cfg.HistoryPollInterval},
|
||||
{"PROBE_TIMEOUT", &cfg.ProbeTimeout},
|
||||
{"HEALTH_INTERVAL", &cfg.HealthInterval},
|
||||
{"FREE_TIMEOUT", &cfg.FreeTimeout},
|
||||
{"WARM_TIMEOUT", &cfg.WarmTimeout},
|
||||
{"SHUTDOWN_TIMEOUT", &cfg.ShutdownTimeout},
|
||||
{"BACKOFF_INITIAL", &cfg.BackoffInitial},
|
||||
{"BACKOFF_MAX", &cfg.BackoffMax},
|
||||
{"COMFY_IDLE_TIMEOUT", &cfg.ComfyIdleTimeout},
|
||||
{"COMFY_START_TIMEOUT", &cfg.ComfyStartTimeout},
|
||||
{"UPDATE_INTERVAL", &cfg.UpdateInterval},
|
||||
{"GAME_POLL_INTERVAL", &cfg.GamePollInterval},
|
||||
} {
|
||||
if err := envDuration(getenv, e.name, e.dst); err != nil {
|
||||
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 != "" {
|
||||
n, err := strconv.ParseInt(v, 10, 64)
|
||||
if err != nil || n < 0 {
|
||||
@@ -241,6 +302,9 @@ func Load(getenv func(string) string) (Config, error) {
|
||||
default:
|
||||
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.OllamaURL == "" && cfg.ComfyURL == "" {
|
||||
return cfg, ErrNoConsumer
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
input := `# comment
|
||||
OLLAMA_URL=http://host:11435
|
||||
@@ -134,6 +157,8 @@ func TestLoadErrors(t *testing.T) {
|
||||
{"LLM_BUSY_MODE", "bogus"},
|
||||
{"LLM_BUSY_STATUS", "200"},
|
||||
{"BUSY_RETRY_AFTER", "0"},
|
||||
{"GPU_FOREIGN_VRAM_MB", "-1"},
|
||||
{"GAME_POLL_INTERVAL", "bogus"},
|
||||
} {
|
||||
_, err := Load(func(k string) string {
|
||||
if k == tc.key {
|
||||
@@ -149,3 +174,51 @@ 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 != 5*time.Second {
|
||||
t.Fatalf("GamePollInterval = %v, want 5s 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")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
func sampleEntries(logFile string) []sampleEntry {
|
||||
return []sampleEntry{
|
||||
{"LISTEN_OLLAMA", ":11434", "Ollama-facing listener address", false},
|
||||
{"LISTEN_COMFY", ":8188", "ComfyUI-facing listener address", false},
|
||||
{"LISTEN_OLLAMA", ":11434", "Listen address for Ollama-compatible clients (gpu-turnstile poses as Ollama here)", 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},
|
||||
{"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},
|
||||
{"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: empty = unmanaged)", false},
|
||||
{"COMFY_DIR", `C:\ComfyUI`, "Working directory for COMFY_CMD (default: empty = inherit)", 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", "5s", "How often game/VRAM detection runs", 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},
|
||||
{"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 != ""},
|
||||
{"UNLOAD_POLL_INTERVAL", "500ms", "/api/ps poll interval while unloading", 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},
|
||||
{"WARM_TIMEOUT", "2m", "Timeout for the warm-model reload after an image job", false},
|
||||
{"SHUTDOWN_TIMEOUT", "10s", "Graceful shutdown timeout on SIGINT/SIGTERM", false},
|
||||
@@ -59,8 +68,15 @@ func sampleEntries(logFile string) []sampleEntry {
|
||||
// (a Windows service has no console). CFG_VER records the version that
|
||||
// wrote the file so later installs can upgrade it.
|
||||
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
|
||||
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("# The installer uses it to append newly added settings on updates.\n\n")
|
||||
b.WriteString("# gpu-turnstile configuration\n")
|
||||
|
||||
@@ -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) {
|
||||
cases := []struct {
|
||||
a, b string
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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 }
|
||||
@@ -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
@@ -1,6 +1,8 @@
|
||||
// Package lock implements the two-mode GPU arbitration lock: any number of
|
||||
// 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
|
||||
|
||||
import (
|
||||
@@ -13,9 +15,10 @@ import (
|
||||
type State string
|
||||
|
||||
const (
|
||||
StateIdle State = "idle"
|
||||
StateLLM State = "llm"
|
||||
StateImage State = "image"
|
||||
StateIdle State = "idle"
|
||||
StateLLM State = "llm"
|
||||
StateImage State = "image"
|
||||
StateExternal State = "external"
|
||||
)
|
||||
|
||||
type imageWaiter struct{ id uint64 }
|
||||
@@ -30,6 +33,7 @@ type Lock struct {
|
||||
imageActive bool // an image job holds the GPU
|
||||
imageQ []imageWaiter
|
||||
nextID uint64
|
||||
external string // non-empty: a foreign process (e.g. a game) holds the GPU
|
||||
|
||||
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
|
||||
// one in-flight LLM request. Returns ctx.Err() if the context is cancelled
|
||||
// while waiting; no state is changed in that case.
|
||||
func (l *Lock) AcquireLLM(ctx context.Context) error {
|
||||
l.mu.Lock()
|
||||
for l.imageActive || len(l.imageQ) > 0 {
|
||||
for l.imageActive || len(l.imageQ) > 0 || l.external != "" {
|
||||
ch := l.change
|
||||
l.mu.Unlock()
|
||||
select {
|
||||
@@ -76,10 +109,10 @@ func (l *Lock) AcquireLLM(ctx context.Context) error {
|
||||
|
||||
// TryAcquireLLM acquires one in-flight LLM slot without waiting and
|
||||
// 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 {
|
||||
l.mu.Lock()
|
||||
if l.imageActive || len(l.imageQ) > 0 {
|
||||
if l.imageActive || len(l.imageQ) > 0 || l.external != "" {
|
||||
l.mu.Unlock()
|
||||
return false
|
||||
}
|
||||
@@ -120,7 +153,7 @@ func (l *Lock) AcquireImage(ctx context.Context) error {
|
||||
|
||||
for {
|
||||
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.imageActive = true
|
||||
l.mu.Unlock()
|
||||
@@ -166,6 +199,8 @@ func (l *Lock) Snapshot() (state State, llmInflight int, imagePending bool) {
|
||||
state = StateImage
|
||||
case l.n > 0:
|
||||
state = StateLLM
|
||||
case l.external != "":
|
||||
state = StateExternal
|
||||
default:
|
||||
state = StateIdle
|
||||
}
|
||||
|
||||
@@ -226,3 +226,70 @@ func TestRace(t *testing.T) {
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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).
|
||||
# TYPE gpu_turnstile_state gauge
|
||||
`)
|
||||
for _, s := range []string{"idle", "llm", "image"} {
|
||||
for _, s := range []string{"idle", "llm", "image", "external"} {
|
||||
v := 0
|
||||
if s == state {
|
||||
v = 1
|
||||
|
||||
+62
-3
@@ -26,6 +26,7 @@ import (
|
||||
"gpu-turnstile/internal/lock"
|
||||
"gpu-turnstile/internal/metrics"
|
||||
"gpu-turnstile/internal/ollama"
|
||||
"gpu-turnstile/internal/supervise"
|
||||
)
|
||||
|
||||
// defaultCaptureLimit bounds how much of a /prompt response body is
|
||||
@@ -44,6 +45,11 @@ type Config struct {
|
||||
Metrics *metrics.Metrics
|
||||
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
|
||||
// the log level is above INFO (request lines are not emitted at all).
|
||||
LogColor bool
|
||||
@@ -438,7 +444,7 @@ func isLLMRequest(r *http.Request) bool {
|
||||
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 {
|
||||
return s.logRequests("ollama", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
switch r.URL.Path {
|
||||
@@ -458,10 +464,14 @@ func (s *Server) OllamaHandler() http.Handler {
|
||||
if s.busyMode == "reject" {
|
||||
if !s.cfg.Lock.TryAcquireLLM() {
|
||||
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",
|
||||
"path", r.URL.Path, "status", s.busyStatus)
|
||||
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
|
||||
}
|
||||
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 {
|
||||
return s.logRequests("comfy", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
switch r.URL.Path {
|
||||
@@ -500,6 +510,21 @@ func (s *Server) ComfyHandler() http.Handler {
|
||||
s.handlePrompt(w, r)
|
||||
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)
|
||||
}))
|
||||
}
|
||||
@@ -539,6 +564,23 @@ func (w *captureWriter) Unwrap() http.ResponseWriter { return w.ResponseWriter }
|
||||
func (s *Server) handlePrompt(w http.ResponseWriter, r *http.Request) {
|
||||
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()
|
||||
if err := s.cfg.Lock.AcquireImage(r.Context()); err != nil {
|
||||
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())
|
||||
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 {
|
||||
uctx, ucancel := context.WithTimeout(r.Context(), s.cfg.UnloadTimeout)
|
||||
elapsed, uerr := s.cfg.Ollama.UnloadAll(uctx)
|
||||
@@ -608,6 +663,10 @@ func (s *Server) finishImageJob(promptID string) {
|
||||
|
||||
s.cfg.Lock.ReleaseImage()
|
||||
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 state, _, _ := s.cfg.Lock.Snapshot(); state == lock.StateIdle {
|
||||
|
||||
@@ -46,8 +46,8 @@ type fakes struct {
|
||||
rec *recorder
|
||||
ollama *httptest.Server
|
||||
comfy *httptest.Server
|
||||
server *httptest.Server // Ollama-facing gpu-turnstile listener
|
||||
comfySrv *httptest.Server // ComfyUI-facing gpu-turnstile listener
|
||||
server *httptest.Server // gpu-turnstile listener for Ollama-compatible clients
|
||||
comfySrv *httptest.Server // gpu-turnstile listener for ComfyUI clients
|
||||
freeCh chan struct{}
|
||||
chatCh chan struct{}
|
||||
historyMu sync.Mutex
|
||||
|
||||
@@ -0,0 +1,273 @@
|
||||
// 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"
|
||||
"os/exec"
|
||||
"runtime"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// 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
|
||||
}
|
||||
@@ -0,0 +1,192 @@
|
||||
package supervise
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"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")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user