4 Commits
Author SHA1 Message Date
mram 9481fd8418 Embed rotated release public key; pin compose example to v0.1.4
ci / test (push) Successful in 13s
ci / docker (push) Failing after 50s
ci / release (push) Successful in 27s
2026-09-20 22:37:28 +02:00
mram 83812cebf3 Add Linux systemd support and --install-service/--remove-service flags
The service package now has a Linux implementation alongside the Windows
one: systemd unit install/remove (/etc/systemd/system), readiness
notification (READY=1 via go-systemd) sent only after the listeners are
bound, a 30s watchdog, and STOPPING on shutdown. Listeners are pre-bound
so port conflicts fail fast and the readiness signal is truthful. The
notify calls are no-ops without NOTIFY_SOCKET (containers, shells) and
on non-Linux builds. --install-service/--remove-service work on both
platforms; the 'service install|remove' subcommand remains as an alias.
2026-09-20 22:35:59 +02:00
mram 30a1f55aae Make consumers optional: a URL enables its mode, empty disables it
OLLAMA_URL and COMFY_URL no longer have defaults; each consumer
(listener, client, startup probe, lock participation) is enabled by
setting its URL and disabled by leaving it empty. At least one must be
set. Ollama-only mode is a pure pass-through; ComfyUI-only mode skips
the unload and warm-reload steps. /metrics is now served on both
listeners. This is the extension pattern for future consumers such as
local game detection.
2026-09-20 22:29:16 +02:00
mram 642cc36a39 Pass signing key via env so its value is never echoed in CI logs 2026-09-20 22:21:25 +02:00
17 changed files with 526 additions and 88 deletions
+5 -1
View File
@@ -80,8 +80,12 @@ jobs:
-o gpu-turnstile.exe ./cmd/gpu-turnstile -o gpu-turnstile.exe ./cmd/gpu-turnstile
- name: Sign and checksum - name: Sign and checksum
env:
RELEASE_SIGNING_KEY: ${{ secrets.RELEASE_SIGNING_KEY }}
run: | run: |
printf '%s\n' "${{ secrets.RELEASE_SIGNING_KEY }}" > key.pem # The key comes via the environment so its value never appears in
# the echoed command line of the run log.
printf '%s\n' "$RELEASE_SIGNING_KEY" > key.pem
chmod 600 key.pem chmod 600 key.pem
openssl pkeyutl -sign -inkey key.pem -rawin \ openssl pkeyutl -sign -inkey key.pem -rawin \
-in gpu-turnstile.exe -out gpu-turnstile.exe.sig -in gpu-turnstile.exe -out gpu-turnstile.exe.sig
+29 -6
View File
@@ -29,6 +29,12 @@ Open WebUI / n8n ────► :8188 ───┘
- Everything else (including websockets and all streaming) passes through - Everything else (including websockets and all streaming) passes through
transparently and unbuffered. transparently and unbuffered.
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.
## Configuration ## Configuration
Configuration comes from environment variables and/or an `.env`-style Configuration comes from environment variables and/or an `.env`-style
@@ -41,8 +47,8 @@ override file values. Invalid values fail at startup.
|---|---|---| |---|---|---|
| `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener | | `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener |
| `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener | | `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener |
| `OLLAMA_URL` | `http://127.0.0.1:11435` | Ollama upstream | | `OLLAMA_URL` | _(empty = disabled)_ | Ollama upstream; set to enable the Ollama consumer |
| `COMFY_URL` | `http://127.0.0.1:8189` | ComfyUI upstream | | `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 |
| `JOB_TIMEOUT` | `15m` | Wait for a ComfyUI job to finish | | `JOB_TIMEOUT` | `15m` | Wait for a ComfyUI job to finish |
| `LLM_WAIT_TIMEOUT` | `10m` | Max lock wait for an LLM request before 503 (wait mode) | | `LLM_WAIT_TIMEOUT` | `10m` | Max lock wait for an LLM request before 503 (wait mode) |
@@ -70,7 +76,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","llm_inflight":N,"image_pending":B}`
- `GET /metrics` (Ollama listener): 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`
(histogram, `kind="llm|image"`), `gpu_turnstile_unload_seconds`. (histogram, `kind="llm|image"`), `gpu_turnstile_unload_seconds`.
@@ -88,15 +94,15 @@ go build ./cmd/gpu-turnstile
./gpu-turnstile ./gpu-turnstile
``` ```
### Run natively on Windows (primary deployment) ### Run natively on Windows (current primary deployment)
Download `gpu-turnstile.exe` from a release, put a `gpu-turnstile.env` Download `gpu-turnstile.exe` from a release, put a `gpu-turnstile.env`
next to it, and run it — or install it as a Windows service from an next to it, and run it — or install it as a Windows service from an
elevated shell: elevated shell:
```sh ```sh
gpu-turnstile.exe service install # auto-start service, recovery = restart gpu-turnstile.exe --install-service # auto-start service, recovery = restart
gpu-turnstile.exe service remove gpu-turnstile.exe --remove-service
``` ```
The service uses the config file (services have no convenient The service uses the config file (services have no convenient
@@ -109,6 +115,23 @@ the install directory for self-updates. For least privilege, run it as the
virtual account `NT SERVICE\gpu-turnstile` and grant write access to just virtual account `NT SERVICE\gpu-turnstile` and grant write access to just
those two directories. those two directories.
### Run natively on Linux (systemd)
The same binary works on Linux. Install it as a systemd service as root:
```sh
gpu-turnstile --install-service # writes + enables + starts the unit
gpu-turnstile --remove-service
```
The unit (`/etc/systemd/system/gpu-turnstile.service`) is `Type=notify`:
`systemctl start` blocks until the listeners are actually bound, a 30 s
watchdog restarts the process if it wedges, and logs land in the journal
(`journalctl -u gpu-turnstile -f`) unless `LOG_FILE` is set. Put the
config in a `gpu-turnstile.env` next to the binary (or pass
`-config /path` during install). The notify integration is a no-op in
containers and interactive shells.
**Auto-update is on by default**: the binary checks the repo's latest **Auto-update is on by default**: the binary checks the repo's latest
release on startup and every `UPDATE_INTERVAL`, verifies the Ed25519 release on startup and every `UPDATE_INTERVAL`, verifies the Ed25519
signature of the download against the public key embedded at build time, signature of the download against the public key embedded at build time,
+61 -15
View File
@@ -42,6 +42,21 @@ listener is an `httputil.ReverseProxy` to its upstream. Websocket upgrades
(ComfyUI `/ws`) and streaming bodies (Ollama NDJSON / SSE) must pass through (ComfyUI `/ws`) and streaming bodies (Ollama NDJSON / SSE) must pass through
unbuffered (`FlushInterval = -1`). unbuffered (`FlushInterval = -1`).
### Modes of operation
Each GPU consumer is enabled by setting its URL and disabled by leaving it
empty — no separate flags. At least one URL must be set; a disabled
consumer gets no listener, no startup probe, and no lock participation:
- **Both set** (default deployment): full arbitration as described below.
- **Only `OLLAMA_URL`**: pure pass-through for Ollama; the LLM lock never
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.
### Lock semantics ### Lock semantics
Two-mode lock with image priority (writer-preferring RW lock, where "readers" Two-mode lock with image priority (writer-preferring RW lock, where "readers"
@@ -83,7 +98,8 @@ ComfyUI listener (`:8188` → `COMFY_URL`):
### Image job flow (`POST /prompt`) ### Image job flow (`POST /prompt`)
1. `AcquireImage()`. 1. `AcquireImage()`.
2. Unload Ollama: `GET /api/ps`; for each model `POST /api/generate 2. Unload Ollama (skipped when `OLLAMA_URL` is unset): `GET /api/ps`; for
each model `POST /api/generate
{"model":M,"keep_alive":0}`; if that returns non-2xx (embedding-only {"model":M,"keep_alive":0}`; if that returns non-2xx (embedding-only
models), `POST /api/embed {"model":M,"input":"x","keep_alive":0}`. Poll models), `POST /api/embed {"model":M,"input":"x","keep_alive":0}`. Poll
`/api/ps` every `UNLOAD_POLL_INTERVAL` (default 500 ms) until empty or `/api/ps` every `UNLOAD_POLL_INTERVAL` (default 500 ms) until empty or
@@ -117,8 +133,8 @@ override file values. A missing file is fine; a malformed one is fatal.
|---|---|---| |---|---|---|
| `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener | | `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener |
| `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener | | `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener |
| `OLLAMA_URL` | `http://127.0.0.1:11435` | upstream | | `OLLAMA_URL` | _(empty = disabled)_ | Ollama upstream; set to enable the Ollama consumer |
| `COMFY_URL` | `http://127.0.0.1:8189` | upstream | | `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 |
| `JOB_TIMEOUT` | `15m` | wait for ComfyUI job | | `JOB_TIMEOUT` | `15m` | wait for ComfyUI job |
| `LLM_WAIT_TIMEOUT` | `10m` | max time an LLM request waits for the lock before 503 (wait mode) | | `LLM_WAIT_TIMEOUT` | `10m` | max time an LLM request waits for the lock before 503 (wait mode) |
@@ -143,17 +159,25 @@ 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 |
Startup fails fast on unparsable values. Both upstreams are probed once at Startup fails fast on unparsable values and when neither consumer URL is
start (`/api/version`, `/system_stats`); failure is logged, not fatal. set. Enabled upstreams are probed once at start (`/api/version`,
`/system_stats`); failure is logged, not fatal.
## Native Windows deployment ## Native deployment (Windows and Linux)
The binary runs natively on Windows (the primary deployment) as well as in The binary runs natively on Windows (the current primary deployment) and on
Docker. Linux with systemd (the future GPU server), as well as in Docker.
- `gpu-turnstile.exe service install [-config path]` registers an Service management is the same on both platforms:
auto-start Windows service (needs an elevated shell). Recovery actions `gpu-turnstile --install-service [-config path]` registers and starts an
restart it 5 s after any failure. `service remove` uninstalls. auto-start service; `--remove-service` stops and unregisters it (both need
an elevated/root shell). The legacy form `gpu-turnstile service
install|remove` does the same thing.
### Windows
- `--install-service` registers a Windows service; recovery actions restart
it 5 s after any failure.
- **Layout**: install to `C:\Program Files\gpu-turnstile\` (exe plus - **Layout**: install to `C:\Program Files\gpu-turnstile\` (exe plus
`gpu-turnstile.env`); logs belong in `C:\ProgramData\gpu-turnstile\` via `gpu-turnstile.env`); logs belong in `C:\ProgramData\gpu-turnstile\` via
`LOG_FILE`. The service must be able to write its install directory for `LOG_FILE`. The service must be able to write its install directory for
@@ -165,6 +189,27 @@ Docker.
log directories only (no network logon, no user profile). log directories only (no network logon, no user profile).
- Use a config file (above) for the service — Windows services have no - Use a config file (above) for the service — Windows services have no
convenient environment. Logs go to `LOG_FILE` since there is no console. convenient environment. Logs go to `LOG_FILE` since there is no console.
### Linux (systemd)
- `--install-service` writes `/etc/systemd/system/gpu-turnstile.service`
with `ExecStart` pointing at the current executable and the `-config`
file, then runs `systemctl daemon-reload` and `enable --now`. The unit
runs as root (it must be able to overwrite its own binary for
self-updates); harden with `ProtectSystem=strict` plus a writable
`ReadWritePaths` if desired.
- The unit is `Type=notify`: the binary sends `READY=1` via
`github.com/coreos/go-systemd` only after the listeners are bound, so
`systemctl start` blocks until the proxy accepts connections. A 30 s
watchdog (`WatchdogSec=`) is pinged as long as the process runs; three
missed pings make systemd restart it. `STOPPING=1` is sent on shutdown.
All notify calls are no-ops when `NOTIFY_SOCKET` is unset (containers,
interactive shells), and the whole integration is Linux-only — Windows
builds carry no-op stubs.
- Logs go to the journal (`journalctl -u gpu-turnstile`) or to `LOG_FILE`.
- **Auto-update** works the same as on Windows: `Restart=on-failure` with
`RestartSec=5s` brings up the staged binary after the updater exits with
code 3.
- **Auto-update**: on startup and every `UPDATE_INTERVAL`, the binary - **Auto-update**: on startup and every `UPDATE_INTERVAL`, the binary
checks `UPDATE_REPO`'s latest release; if its tag is a newer `vX.Y.Z`, checks `UPDATE_REPO`'s latest release; if its tag is a newer `vX.Y.Z`,
it downloads `UPDATE_ASSET` plus its `.sig` (and `.sha256` when present) it downloads `UPDATE_ASSET` plus its `.sig` (and `.sha256` when present)
@@ -184,7 +229,7 @@ Docker.
- `GET /healthz` on both listeners: 200 with JSON - `GET /healthz` on both listeners: 200 with JSON
`{"state":"idle|llm|image","llm_inflight":N,"image_pending":B}`. `{"state":"idle|llm|image","llm_inflight":N,"image_pending":B}`.
- `GET /metrics` on the Ollama listener: Prometheus text format, no external - `GET /metrics` on both listeners: Prometheus text format, no external
dependency needed: dependency needed:
`gpu_turnstile_state{state="…"} 1`, `gpu_turnstile_llm_inflight`, `gpu_turnstile_state{state="…"} 1`, `gpu_turnstile_llm_inflight`,
`gpu_turnstile_image_jobs_total`, `gpu_turnstile_lock_wait_seconds` `gpu_turnstile_image_jobs_total`, `gpu_turnstile_lock_wait_seconds`
@@ -233,7 +278,7 @@ gpu-turnstile/
internal/metrics/ # Prometheus exposition internal/metrics/ # Prometheus exposition
internal/config/ # env + .env file configuration internal/config/ # env + .env file configuration
internal/update/ # signed auto-updater (public key in pubkey.go) internal/update/ # signed auto-updater (public key in pubkey.go)
internal/service/ # Windows service integration internal/service/ # Windows SCM + Linux systemd (notify/watchdog) integration
Dockerfile Dockerfile
.gitea/workflows/ci.yml .gitea/workflows/ci.yml
README.md README.md
@@ -261,8 +306,9 @@ are new.
## Build and CI ## Build and CI
- Go 1.23+, `golang.org/x/sys` is the only external dependency (Windows - Go 1.23+, two external dependencies: `golang.org/x/sys` (Windows service
service integration; not used in the Linux build). `CGO_ENABLED=0`, integration) and `github.com/coreos/go-systemd` (systemd notify/watchdog,
Linux build only). `CGO_ENABLED=0`,
`-ldflags="-s -w"`, version from `git describe` injected via `-ldflags="-s -w"`, version from `git describe` injected via
`-X main.version=`. `-X main.version=`.
- Dockerfile: multi-stage, final image `gcr.io/distroless/static` (or - Dockerfile: multi-stage, final image `gcr.io/distroless/static` (or
+84 -21
View File
@@ -8,6 +8,7 @@ import (
"fmt" "fmt"
"io" "io"
"log/slog" "log/slog"
"net"
"net/http" "net/http"
"os" "os"
"os/signal" "os/signal"
@@ -34,12 +35,22 @@ var version = "dev"
const exitCodeUpdate = 3 const exitCodeUpdate = 3
func main() { func main() {
configPath, args := splitConfigFlag(os.Args[1:]) configPath, install, remove, args := parseFlags(os.Args[1:])
switch {
case install && remove:
fmt.Fprintf(os.Stderr, "gpu-turnstile: --install-service and --remove-service are mutually exclusive\n")
os.Exit(2)
case install:
os.Exit(serviceCommand(configPath, []string{"install"}))
case remove:
os.Exit(serviceCommand(configPath, []string{"remove"}))
}
if len(args) > 0 && args[0] == "service" { if len(args) > 0 && args[0] == "service" {
os.Exit(serviceCommand(configPath, args[1:])) os.Exit(serviceCommand(configPath, args[1:]))
} }
if len(args) > 0 { if len(args) > 0 {
fmt.Fprintf(os.Stderr, "usage: gpu-turnstile [-config path] | gpu-turnstile service install|remove [-config path]\n") fmt.Fprintf(os.Stderr, "usage: gpu-turnstile [-config path] [--install-service | --remove-service]\n")
fmt.Fprintf(os.Stderr, " gpu-turnstile service install|remove [-config path]\n")
os.Exit(2) os.Exit(2)
} }
@@ -70,10 +81,10 @@ func main() {
} }
} }
// splitConfigFlag extracts -config <path> (or -config=<path>) from args. // parseFlags extracts -config <path> (or -config=<path>) and the
func splitConfigFlag(args []string) (string, []string) { // --install-service / --remove-service switches from args.
var configPath string func parseFlags(args []string) (configPath string, install, remove bool, rest []string) {
rest := args[:0] rest = args[:0]
for i := 0; i < len(args); i++ { for i := 0; i < len(args); i++ {
switch { switch {
case args[i] == "-config" && i+1 < len(args): case args[i] == "-config" && i+1 < len(args):
@@ -81,11 +92,15 @@ func splitConfigFlag(args []string) (string, []string) {
i++ i++
case strings.HasPrefix(args[i], "-config="): case strings.HasPrefix(args[i], "-config="):
configPath = strings.TrimPrefix(args[i], "-config=") configPath = strings.TrimPrefix(args[i], "-config=")
case args[i] == "--install-service" || args[i] == "-install-service":
install = true
case args[i] == "--remove-service" || args[i] == "-remove-service":
remove = true
default: default:
rest = append(rest, args[i]) rest = append(rest, args[i])
} }
} }
return configPath, rest return configPath, install, remove, rest
} }
// defaultConfigPath returns gpu-turnstile.env next to the executable. // defaultConfigPath returns gpu-turnstile.env next to the executable.
@@ -180,6 +195,14 @@ func serviceCommand(configPath string, args []string) int {
return 0 return 0
} }
// orDisabled renders an empty URL as "disabled" for the startup dump.
func orDisabled(url string) string {
if url == "" {
return "disabled"
}
return url
}
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 {
// 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.
@@ -187,8 +210,8 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
"version", version, "version", version,
"listen_ollama", cfg.ListenOllama, "listen_ollama", cfg.ListenOllama,
"listen_comfy", cfg.ListenComfy, "listen_comfy", cfg.ListenComfy,
"ollama_url", cfg.OllamaURL, "ollama_url", orDisabled(cfg.OllamaURL),
"comfy_url", cfg.ComfyURL, "comfy_url", orDisabled(cfg.ComfyURL),
"unload_timeout", cfg.UnloadTimeout, "unload_timeout", cfg.UnloadTimeout,
"job_timeout", cfg.JobTimeout, "job_timeout", cfg.JobTimeout,
"llm_wait_timeout", cfg.LLMWaitTimeout, "llm_wait_timeout", cfg.LLMWaitTimeout,
@@ -215,14 +238,21 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
) )
lk := lock.New(log) lk := lock.New(log)
ollamaClient, err := ollama.New(cfg.OllamaURL, log) // Each consumer is enabled by setting its URL; a disabled consumer gets
if err != nil { // no client, no listener and no probe.
var ollamaClient *ollama.Client
var err error
if cfg.OllamaURL != "" {
if ollamaClient, err = ollama.New(cfg.OllamaURL, log); err != nil {
return err return err
} }
comfyClient, err := comfy.New(cfg.ComfyURL, log) }
if err != nil { var comfyClient *comfy.Client
if cfg.ComfyURL != "" {
if comfyClient, err = comfy.New(cfg.ComfyURL, log); err != nil {
return err return err
} }
}
srv, err := proxy.New(proxy.Config{ srv, err := proxy.New(proxy.Config{
OllamaURL: cfg.OllamaURL, OllamaURL: cfg.OllamaURL,
@@ -253,22 +283,54 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
return err return err
} }
// Probe both upstreams once; failure is logged, not fatal. // Probe the enabled upstreams once; failure is logged, not fatal.
probeCtx, probeCancel := context.WithTimeout(ctx, cfg.ProbeTimeout) probeCtx, probeCancel := context.WithTimeout(ctx, cfg.ProbeTimeout)
if ollamaClient != nil {
if err := ollamaClient.Probe(probeCtx); err != nil { if err := ollamaClient.Probe(probeCtx); err != nil {
log.Warn("ollama probe failed", "url", cfg.OllamaURL, "err", err) log.Warn("ollama probe failed", "url", cfg.OllamaURL, "err", err)
} }
}
if comfyClient != nil {
if err := comfyClient.Probe(probeCtx); err != nil { if err := comfyClient.Probe(probeCtx); err != nil {
log.Warn("comfy probe failed", "url", cfg.ComfyURL, "err", err) log.Warn("comfy probe failed", "url", cfg.ComfyURL, "err", err)
} }
}
probeCancel() probeCancel()
ollamaSrv := &http.Server{Addr: cfg.ListenOllama, Handler: srv.OllamaHandler()} // Bind the listeners up front so a port conflict fails fast and the
comfySrv := &http.Server{Addr: cfg.ListenComfy, Handler: srv.ComfyHandler()} // readiness notification below really means "accepting connections".
var servers []*http.Server
var listeners []net.Listener
bind := func(addr string, handler http.Handler, consumer string) error {
ln, err := net.Listen("tcp", addr)
if err != nil {
return fmt.Errorf("listen %s on %s: %w", consumer, addr, err)
}
servers = append(servers, &http.Server{Addr: addr, Handler: handler})
listeners = append(listeners, ln)
log.Warn("listening", "consumer", consumer, "addr", addr)
return nil
}
if ollamaClient != nil {
if err := bind(cfg.ListenOllama, srv.OllamaHandler(), "ollama"); err != nil {
return err
}
}
if comfyClient != nil {
if err := bind(cfg.ListenComfy, srv.ComfyHandler(), "comfy"); err != nil {
return err
}
}
errCh := make(chan error, 2) errCh := make(chan error, len(servers))
go func() { errCh <- ollamaSrv.ListenAndServe() }() for i := range servers {
go func() { errCh <- comfySrv.ListenAndServe() }() go func(s *http.Server, ln net.Listener) { errCh <- s.Serve(ln) }(servers[i], listeners[i])
}
// Tell systemd we are up and start the watchdog pings; both are no-ops
// when not running under a notify/watchdog unit.
service.NotifyReady()
service.StartWatchdog(ctx)
if cfg.AutoUpdate { if cfg.AutoUpdate {
go updateLoop(ctx, cfg, log, lk, isService) go updateLoop(ctx, cfg, log, lk, isService)
@@ -285,8 +347,9 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), cfg.ShutdownTimeout) shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), cfg.ShutdownTimeout)
defer shutdownCancel() defer shutdownCancel()
ollamaSrv.Shutdown(shutdownCtx) for _, s := range servers {
comfySrv.Shutdown(shutdownCtx) s.Shutdown(shutdownCtx)
}
return nil return nil
} }
+3 -1
View File
@@ -6,9 +6,11 @@
# 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.3 image: git.rambossek.at/public/gpu-turnstile:v0.1.4
restart: unless-stopped restart: unless-stopped
environment: environment:
# Each consumer is enabled by setting its URL; leave one unset to
# disable that side (no listener, no probe, no lock participation).
# Services on the Docker host itself: # Services on the Docker host itself:
OLLAMA_URL: http://host.docker.internal:11435 OLLAMA_URL: http://host.docker.internal:11435
COMFY_URL: http://host.docker.internal:8189 COMFY_URL: http://host.docker.internal:8189
+2
View File
@@ -3,3 +3,5 @@ module gpu-turnstile
go 1.23 go 1.23
require golang.org/x/sys v0.29.0 require golang.org/x/sys v0.29.0
require github.com/coreos/go-systemd/v22 v22.7.0
+2
View File
@@ -1,2 +1,4 @@
github.com/coreos/go-systemd/v22 v22.7.0 h1:LAEzFkke61DFROc7zNLX/WA2i5J8gYqe0rSj9KI28KA=
github.com/coreos/go-systemd/v22 v22.7.0/go.mod h1:xNUYtjHu2EDXbsxz1i41wouACIwT7Ybq9o0BQhMwD0w=
golang.org/x/sys v0.29.0 h1:TPYlXGxvx1MGTn2GiZDhnjPA9wZzZeGKHHmKhHYvgaU= golang.org/x/sys v0.29.0 h1:TPYlXGxvx1MGTn2GiZDhnjPA9wZzZeGKHHmKhHYvgaU=
golang.org/x/sys v0.29.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.29.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
+5 -3
View File
@@ -51,13 +51,12 @@ type Config struct {
} }
// Defaults returns the configuration used when neither the environment nor // Defaults returns the configuration used when neither the environment nor
// a config file sets a value. // a config file sets a value. The upstream URLs default to empty: a
// consumer is enabled by setting its URL, disabled by leaving it empty.
func Defaults() Config { func Defaults() Config {
return Config{ return Config{
ListenOllama: ":11434", ListenOllama: ":11434",
ListenComfy: ":8188", ListenComfy: ":8188",
OllamaURL: "http://127.0.0.1:11435",
ComfyURL: "http://127.0.0.1:8189",
UnloadTimeout: time.Minute, UnloadTimeout: time.Minute,
JobTimeout: 15 * time.Minute, JobTimeout: 15 * time.Minute,
LLMWaitTimeout: 10 * time.Minute, LLMWaitTimeout: 10 * time.Minute,
@@ -219,5 +218,8 @@ 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.OllamaURL == "" && cfg.ComfyURL == "" {
return cfg, fmt.Errorf("at least one of OLLAMA_URL or COMFY_URL must be set (each URL enables its consumer)")
}
return cfg, nil return cfg, nil
} }
+19 -1
View File
@@ -8,13 +8,21 @@ import (
) )
func TestDefaults(t *testing.T) { func TestDefaults(t *testing.T) {
cfg, err := Load(func(string) string { return "" }) cfg, err := Load(func(k string) string {
if k == "OLLAMA_URL" {
return "http://127.0.0.1:11435"
}
return ""
})
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if cfg.ListenOllama != ":11434" || cfg.ListenComfy != ":8188" { if cfg.ListenOllama != ":11434" || cfg.ListenComfy != ":8188" {
t.Fatalf("listen addrs = %s %s", cfg.ListenOllama, cfg.ListenComfy) t.Fatalf("listen addrs = %s %s", cfg.ListenOllama, cfg.ListenComfy)
} }
if cfg.ComfyURL != "" {
t.Fatalf("ComfyURL default = %q, want empty (disabled)", cfg.ComfyURL)
}
if cfg.UnloadTimeout != time.Minute || cfg.JobTimeout != 15*time.Minute { if cfg.UnloadTimeout != time.Minute || cfg.JobTimeout != 15*time.Minute {
t.Fatalf("timeouts = %v %v", cfg.UnloadTimeout, cfg.JobTimeout) t.Fatalf("timeouts = %v %v", cfg.UnloadTimeout, cfg.JobTimeout)
} }
@@ -26,6 +34,13 @@ func TestDefaults(t *testing.T) {
} }
} }
func TestLoadRequiresConsumer(t *testing.T) {
_, err := Load(func(string) string { return "" })
if err == nil || !strings.Contains(err.Error(), "OLLAMA_URL") {
t.Fatalf("err = %v, want missing-consumer error", 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
@@ -91,6 +106,9 @@ func TestLoadErrors(t *testing.T) {
if k == tc.key { if k == tc.key {
return tc.value return tc.value
} }
if k == "OLLAMA_URL" {
return "http://127.0.0.1:11435"
}
return "" return ""
}) })
if err == nil { if err == nil {
+37 -13
View File
@@ -102,24 +102,37 @@ type Server struct {
comfyProxy *httputil.ReverseProxy comfyProxy *httputil.ReverseProxy
} }
// New builds a Server, validating the upstream URLs. // New builds a Server, validating the upstream URLs. At least one of
// OllamaURL / ComfyURL must be set; an empty URL disables that consumer —
// its handler is then never served, its client may be nil, and the image
// job flow skips the Ollama unload/warm steps.
func New(cfg Config) (*Server, error) { func New(cfg Config) (*Server, error) {
ollamaURL, err := url.Parse(cfg.OllamaURL) if cfg.OllamaURL == "" && cfg.ComfyURL == "" {
if err != nil || ollamaURL.Scheme == "" || ollamaURL.Host == "" { return nil, fmt.Errorf("at least one of OllamaURL or ComfyURL is required")
}
var ollamaURL, comfyURL *url.URL
if cfg.OllamaURL != "" {
u, err := url.Parse(cfg.OllamaURL)
if err != nil || u.Scheme == "" || u.Host == "" {
return nil, fmt.Errorf("invalid OLLAMA_URL %q", cfg.OllamaURL) return nil, fmt.Errorf("invalid OLLAMA_URL %q", cfg.OllamaURL)
} }
comfyURL, err := url.Parse(cfg.ComfyURL) ollamaURL = u
if err != nil || comfyURL.Scheme == "" || comfyURL.Host == "" { }
if cfg.ComfyURL != "" {
u, err := url.Parse(cfg.ComfyURL)
if err != nil || u.Scheme == "" || u.Host == "" {
return nil, fmt.Errorf("invalid COMFY_URL %q", cfg.ComfyURL) return nil, fmt.Errorf("invalid COMFY_URL %q", cfg.ComfyURL)
} }
comfyURL = u
}
log := cfg.Log log := cfg.Log
if log == nil { if log == nil {
log = slog.Default() log = slog.Default()
} }
if cfg.UnloadPollInterval > 0 { if cfg.UnloadPollInterval > 0 && cfg.Ollama != nil {
cfg.Ollama.PollInterval = cfg.UnloadPollInterval cfg.Ollama.PollInterval = cfg.UnloadPollInterval
} }
if cfg.HistoryPollInterval > 0 { if cfg.HistoryPollInterval > 0 && cfg.Comfy != nil {
cfg.Comfy.PollInterval = cfg.HistoryPollInterval cfg.Comfy.PollInterval = cfg.HistoryPollInterval
} }
freeTimeout := cfg.FreeTimeout freeTimeout := cfg.FreeTimeout
@@ -160,7 +173,7 @@ func New(cfg Config) (*Server, error) {
max: backoffMax, max: backoffMax,
log: log, log: log,
} }
return &Server{ s := &Server{
cfg: cfg, cfg: cfg,
log: log, log: log,
logWriter: cfg.LogWriter, logWriter: cfg.LogWriter,
@@ -172,9 +185,14 @@ func New(cfg Config) (*Server, error) {
busyMode: busyMode, busyMode: busyMode,
busyStatus: busyStatus, busyStatus: busyStatus,
busyRetryAfter: busyRetryAfter, busyRetryAfter: busyRetryAfter,
ollamaProxy: newReverseProxy(ollamaURL, retry, log.With("upstream", "ollama")), }
comfyProxy: newReverseProxy(comfyURL, retry, log.With("upstream", "comfy")), if ollamaURL != nil {
}, nil s.ollamaProxy = newReverseProxy(ollamaURL, retry, log.With("upstream", "ollama"))
}
if comfyURL != nil {
s.comfyProxy = newReverseProxy(comfyURL, retry, log.With("upstream", "comfy"))
}
return s, nil
} }
// retryTransport retries requests whose failure means the upstream never // retryTransport retries requests whose failure means the upstream never
@@ -470,9 +488,13 @@ func (s *Server) OllamaHandler() http.Handler {
// ComfyHandler serves the ComfyUI-facing listener. // ComfyHandler serves the ComfyUI-facing listener.
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) {
if r.URL.Path == "/healthz" { switch r.URL.Path {
case "/healthz":
s.writeHealthz(w) s.writeHealthz(w)
return return
case "/metrics":
s.writeMetrics(w)
return
} }
if r.Method == http.MethodPost && r.URL.Path == "/prompt" { if r.Method == http.MethodPost && r.URL.Path == "/prompt" {
s.handlePrompt(w, r) s.handlePrompt(w, r)
@@ -527,6 +549,7 @@ 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.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)
ucancel() ucancel()
@@ -541,6 +564,7 @@ func (s *Server) handlePrompt(w http.ResponseWriter, r *http.Request) {
default: default:
log.Info("ollama models unloaded", "seconds", elapsed.Seconds()) log.Info("ollama models unloaded", "seconds", elapsed.Seconds())
} }
}
cw := &captureWriter{ResponseWriter: w, status: http.StatusOK, limit: s.captureLimit} cw := &captureWriter{ResponseWriter: w, status: http.StatusOK, limit: s.captureLimit}
s.comfyProxy.ServeHTTP(cw, r) s.comfyProxy.ServeHTTP(cw, r)
@@ -585,7 +609,7 @@ 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.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 {
wctx, wcancel := context.WithTimeout(context.Background(), s.warmTimeout) wctx, wcancel := context.WithTimeout(context.Background(), s.warmTimeout)
if err := s.cfg.Ollama.Warm(wctx, s.cfg.WarmModel); err != nil { if err := s.cfg.Ollama.Warm(wctx, s.cfg.WarmModel); err != nil {
+79
View File
@@ -451,3 +451,82 @@ func TestLLMBusyWaitTimeoutRetryAfter(t *testing.T) {
t.Fatalf("Retry-After = %q", got) t.Fatalf("Retry-After = %q", got)
} }
} }
func TestComfyOnlyModeSkipsOllama(t *testing.T) {
f := newFakes(t)
comfyClient, err := comfy.New(f.comfy.URL, nil)
if err != nil {
t.Fatal(err)
}
comfyClient.PollInterval = 5 * time.Millisecond
srv, err := New(Config{
ComfyURL: f.comfy.URL,
Lock: lock.New(nil),
Comfy: comfyClient,
Metrics: metrics.New(),
JobTimeout: 2 * time.Second,
})
if err != nil {
t.Fatal(err)
}
front := httptest.NewServer(srv.ComfyHandler())
defer front.Close()
resp, err := http.Post(front.URL+"/prompt", "application/json", strings.NewReader(`{}`))
if err != nil {
t.Fatal(err)
}
resp.Body.Close()
if resp.StatusCode != 200 {
t.Fatalf("prompt status = %d", resp.StatusCode)
}
f.completeJob()
select {
case <-f.freeCh:
case <-time.After(3 * time.Second):
t.Fatal("/free never called")
}
// With Ollama disabled the unload steps must not happen.
if i := f.rec.index("unload"); i >= 0 {
f.rec.mu.Lock()
t.Fatalf("unload called with ollama disabled; events: %v", f.rec.events)
}
}
func TestOllamaOnlyMode(t *testing.T) {
f := newFakes(t)
ollamaClient, err := ollama.New(f.ollama.URL, nil)
if err != nil {
t.Fatal(err)
}
srv, err := New(Config{
OllamaURL: f.ollama.URL,
Lock: lock.New(nil),
Ollama: ollamaClient,
Metrics: metrics.New(),
LLMWaitTimeout: time.Second,
})
if err != nil {
t.Fatal(err)
}
front := httptest.NewServer(srv.OllamaHandler())
defer front.Close()
resp, err := http.Post(front.URL+"/api/chat", "application/json", strings.NewReader(`{}`))
if err != nil {
t.Fatal(err)
}
resp.Body.Close()
if resp.StatusCode != 200 {
t.Fatalf("chat status = %d", resp.StatusCode)
}
}
func TestNewRequiresConsumer(t *testing.T) {
_, err := New(Config{Lock: lock.New(nil), Metrics: metrics.New()})
if err == nil {
t.Fatal("New with no upstream URLs should fail")
}
}
+46
View File
@@ -0,0 +1,46 @@
//go:build linux
package service
import (
"context"
"os"
"strconv"
"time"
"github.com/coreos/go-systemd/v22/daemon"
)
// NotifyReady tells systemd the service is up (Type=notify). It is a no-op
// when NOTIFY_SOCKET is unset, e.g. in a container or interactive shell.
func NotifyReady() {
daemon.SdNotify(false, daemon.SdNotifyReady)
}
// NotifyStopping tells systemd the service is shutting down.
func NotifyStopping() {
daemon.SdNotify(false, daemon.SdNotifyStopping)
}
// StartWatchdog pings the systemd watchdog every half of WATCHDOG_USEC
// until ctx is cancelled. It is a no-op unless systemd started the process
// with a watchdog configured (WatchdogSec= in the unit).
func StartWatchdog(ctx context.Context) {
usec, err := strconv.Atoi(os.Getenv("WATCHDOG_USEC"))
if err != nil || usec <= 0 {
return
}
interval := time.Duration(usec) * time.Microsecond / 2
go func() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
daemon.SdNotify(false, daemon.SdNotifyWatchdog)
}
}
}()
}
+14
View File
@@ -0,0 +1,14 @@
//go:build !linux
package service
import "context"
// NotifyReady is a no-op outside Linux (no systemd notify socket).
func NotifyReady() {}
// NotifyStopping is a no-op outside Linux.
func NotifyStopping() {}
// StartWatchdog is a no-op outside Linux.
func StartWatchdog(context.Context) {}
+90
View File
@@ -0,0 +1,90 @@
//go:build linux
// Package service integrates gpu-turnstile with systemd on Linux: running
// under a unit with readiness notification and watchdog, plus
// install/remove helpers that manage a system unit.
package service
import (
"context"
"fmt"
"os"
"os/exec"
"os/signal"
"path/filepath"
"syscall"
)
// Name matches the Windows service name; the systemd unit is Name + ".service".
const Name = "gpu-turnstile"
// unitPath is where Install writes the unit file.
const unitPath = "/etc/systemd/system/" + Name + ".service"
// IsService reports whether the process was started by systemd.
func IsService() bool { return os.Getenv("INVOCATION_ID") != "" }
// Run executes run with SIGINT/SIGTERM cancellation (which is how systemctl
// stop signals the process) and tells systemd when the shutdown begins.
func Run(run func(ctx context.Context) error) error {
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
defer NotifyStopping()
return run(ctx)
}
// renderUnit builds the systemd unit: Type=notify so systemctl start blocks
// until the listeners are bound, a 30s watchdog, and restart-on-failure
// with a 5s delay — which is also what brings up a staged update after the
// updater exits with a non-zero code.
func renderUnit(exePath, configPath string) string {
return fmt.Sprintf(`[Unit]
Description=gpu-turnstile GPU arbitration proxy for Ollama and ComfyUI
After=network-online.target
Wants=network-online.target
[Service]
Type=notify
WatchdogSec=30s
ExecStart=%q -config %q
Restart=on-failure
RestartSec=5s
[Install]
WantedBy=multi-user.target
`, exePath, configPath)
}
// Install writes the unit for the current executable and the given config
// file, then enables and starts it. Needs root.
func Install(configPath string) error {
exe, err := os.Executable()
if err != nil {
return err
}
if abs, absErr := filepath.Abs(exe); absErr == nil {
exe = abs
}
if err := os.WriteFile(unitPath, []byte(renderUnit(exe, configPath)), 0o644); err != nil {
return fmt.Errorf("write %s (run as root): %w", unitPath, err)
}
if out, err := exec.Command("systemctl", "daemon-reload").CombinedOutput(); err != nil {
return fmt.Errorf("systemctl daemon-reload: %w (%s)", err, out)
}
if out, err := exec.Command("systemctl", "enable", "--now", Name+".service").CombinedOutput(); err != nil {
return fmt.Errorf("systemctl enable --now: %w (%s)", err, out)
}
return nil
}
// Remove stops and disables the service and deletes the unit file.
func Remove() error {
exec.Command("systemctl", "disable", "--now", Name+".service").Run() // ignore: may not exist
if err := os.Remove(unitPath); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("remove %s: %w", unitPath, err)
}
if out, err := exec.Command("systemctl", "daemon-reload").CombinedOutput(); err != nil {
return fmt.Errorf("systemctl daemon-reload: %w (%s)", err, out)
}
return nil
}
+23
View File
@@ -0,0 +1,23 @@
//go:build linux
package service
import (
"strings"
"testing"
)
func TestRenderUnit(t *testing.T) {
unit := renderUnit("/usr/local/bin/gpu-turnstile", "/etc/gpu-turnstile.env")
for _, want := range []string{
"Type=notify",
"WatchdogSec=30s",
`ExecStart="/usr/local/bin/gpu-turnstile" -config "/etc/gpu-turnstile.env"`,
"Restart=on-failure",
"WantedBy=multi-user.target",
} {
if !strings.Contains(unit, want) {
t.Fatalf("unit missing %q:\n%s", want, unit)
}
}
}
+5 -5
View File
@@ -1,8 +1,8 @@
//go:build !windows //go:build !windows && !linux
// Package service provides the non-Windows stubs for the Windows service // Package service provides the stubs for platforms without service
// integration. Run falls back to plain signal handling; install/remove // integration (Windows uses the SCM, Linux uses systemd). Run falls back
// are unsupported. // to plain signal handling; install/remove are unsupported.
package service package service
import ( import (
@@ -15,7 +15,7 @@ import (
// Name matches the Windows service name. // Name matches the Windows service name.
const Name = "gpu-turnstile" const Name = "gpu-turnstile"
var errUnsupported = errors.New("service management is only supported on Windows") var errUnsupported = errors.New("service management is only supported on Windows and Linux (systemd)")
// IsService is always false on non-Windows platforms. // IsService is always false on non-Windows platforms.
func IsService() bool { return false } func IsService() bool { return false }
+1 -1
View File
@@ -11,6 +11,6 @@ package update
// the RELEASE_SIGNING_KEY repository secret. When empty, the updater refuses // the RELEASE_SIGNING_KEY repository secret. When empty, the updater refuses
// to update (e.g. development builds). // to update (e.g. development builds).
var publicKeyPEM = `-----BEGIN PUBLIC KEY----- var publicKeyPEM = `-----BEGIN PUBLIC KEY-----
MCowBQYDK2VwAyEA+gJbSvgeYX58woPQGbSC8x8Zw4OTDiiQ7/19seZKfSQ= MCowBQYDK2VwAyEAqTAJ0CCeAQI7MhFlgc5xNmF/CfvLVUAY3ZAoeAS0tT8=
-----END PUBLIC KEY----- -----END PUBLIC KEY-----
` `