Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
75f16a0229 | ||
|
|
9481fd8418 | ||
|
|
83812cebf3 | ||
|
|
30a1f55aae | ||
|
|
642cc36a39 |
@@ -80,8 +80,12 @@ jobs:
|
||||
-o gpu-turnstile.exe ./cmd/gpu-turnstile
|
||||
|
||||
- name: Sign and checksum
|
||||
env:
|
||||
RELEASE_SIGNING_KEY: ${{ secrets.RELEASE_SIGNING_KEY }}
|
||||
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
|
||||
openssl pkeyutl -sign -inkey key.pem -rawin \
|
||||
-in gpu-turnstile.exe -out gpu-turnstile.exe.sig
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
FROM golang:1.23 AS build
|
||||
WORKDIR /src
|
||||
COPY go.mod ./
|
||||
COPY go.mod go.sum ./
|
||||
COPY cmd ./cmd
|
||||
COPY internal ./internal
|
||||
ARG VERSION=dev
|
||||
|
||||
@@ -29,6 +29,12 @@ Open WebUI / n8n ────► :8188 ───┘
|
||||
- Everything else (including websockets and all streaming) passes through
|
||||
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 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_COMFY` | `:8188` | ComfyUI-facing listener |
|
||||
| `OLLAMA_URL` | `http://127.0.0.1:11435` | Ollama upstream |
|
||||
| `COMFY_URL` | `http://127.0.0.1:8189` | ComfyUI upstream |
|
||||
| `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 |
|
||||
| `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) |
|
||||
@@ -70,7 +76,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 /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_image_jobs_total`, `gpu_turnstile_lock_wait_seconds`
|
||||
(histogram, `kind="llm|image"`), `gpu_turnstile_unload_seconds`.
|
||||
@@ -88,15 +94,15 @@ go build ./cmd/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`
|
||||
next to it, and run it — or install it as a Windows service from an
|
||||
elevated shell:
|
||||
|
||||
```sh
|
||||
gpu-turnstile.exe service install # auto-start service, recovery = restart
|
||||
gpu-turnstile.exe service remove
|
||||
gpu-turnstile.exe --install-service # auto-start service, recovery = restart
|
||||
gpu-turnstile.exe --remove-service
|
||||
```
|
||||
|
||||
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
|
||||
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
|
||||
release on startup and every `UPDATE_INTERVAL`, verifies the Ed25519
|
||||
signature of the download against the public key embedded at build time,
|
||||
|
||||
@@ -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
|
||||
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
|
||||
|
||||
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`)
|
||||
|
||||
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
|
||||
models), `POST /api/embed {"model":M,"input":"x","keep_alive":0}`. Poll
|
||||
`/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_COMFY` | `:8188` | ComfyUI-facing listener |
|
||||
| `OLLAMA_URL` | `http://127.0.0.1:11435` | upstream |
|
||||
| `COMFY_URL` | `http://127.0.0.1:8189` | upstream |
|
||||
| `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 |
|
||||
| `JOB_TIMEOUT` | `15m` | wait for ComfyUI job |
|
||||
| `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_ASSET` | `gpu-turnstile.exe` | release asset to download |
|
||||
|
||||
Startup fails fast on unparsable values. Both upstreams are probed once at
|
||||
start (`/api/version`, `/system_stats`); failure is logged, not fatal.
|
||||
Startup fails fast on unparsable values and when neither consumer URL is
|
||||
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
|
||||
Docker.
|
||||
The binary runs natively on Windows (the current primary deployment) and on
|
||||
Linux with systemd (the future GPU server), as well as in Docker.
|
||||
|
||||
- `gpu-turnstile.exe service install [-config path]` registers an
|
||||
auto-start Windows service (needs an elevated shell). Recovery actions
|
||||
restart it 5 s after any failure. `service remove` uninstalls.
|
||||
Service management is the same on both platforms:
|
||||
`gpu-turnstile --install-service [-config path]` registers and starts an
|
||||
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
|
||||
`gpu-turnstile.env`); logs belong in `C:\ProgramData\gpu-turnstile\` via
|
||||
`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).
|
||||
- Use a config file (above) for the service — Windows services have no
|
||||
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
|
||||
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)
|
||||
@@ -184,7 +229,7 @@ Docker.
|
||||
|
||||
- `GET /healthz` on both listeners: 200 with JSON
|
||||
`{"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:
|
||||
`gpu_turnstile_state{state="…"} 1`, `gpu_turnstile_llm_inflight`,
|
||||
`gpu_turnstile_image_jobs_total`, `gpu_turnstile_lock_wait_seconds`
|
||||
@@ -233,7 +278,7 @@ gpu-turnstile/
|
||||
internal/metrics/ # Prometheus exposition
|
||||
internal/config/ # env + .env file configuration
|
||||
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
|
||||
.gitea/workflows/ci.yml
|
||||
README.md
|
||||
@@ -261,8 +306,9 @@ are new.
|
||||
|
||||
## Build and CI
|
||||
|
||||
- Go 1.23+, `golang.org/x/sys` is the only external dependency (Windows
|
||||
service integration; not used in the Linux build). `CGO_ENABLED=0`,
|
||||
- Go 1.23+, two external dependencies: `golang.org/x/sys` (Windows service
|
||||
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
|
||||
`-X main.version=`.
|
||||
- Dockerfile: multi-stage, final image `gcr.io/distroless/static` (or
|
||||
|
||||
+84
-21
@@ -8,6 +8,7 @@ import (
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net"
|
||||
"net/http"
|
||||
"os"
|
||||
"os/signal"
|
||||
@@ -34,12 +35,22 @@ var version = "dev"
|
||||
const exitCodeUpdate = 3
|
||||
|
||||
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" {
|
||||
os.Exit(serviceCommand(configPath, args[1:]))
|
||||
}
|
||||
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)
|
||||
}
|
||||
|
||||
@@ -70,10 +81,10 @@ func main() {
|
||||
}
|
||||
}
|
||||
|
||||
// splitConfigFlag extracts -config <path> (or -config=<path>) from args.
|
||||
func splitConfigFlag(args []string) (string, []string) {
|
||||
var configPath string
|
||||
rest := args[:0]
|
||||
// parseFlags extracts -config <path> (or -config=<path>) and the
|
||||
// --install-service / --remove-service switches from args.
|
||||
func parseFlags(args []string) (configPath string, install, remove bool, rest []string) {
|
||||
rest = args[:0]
|
||||
for i := 0; i < len(args); i++ {
|
||||
switch {
|
||||
case args[i] == "-config" && i+1 < len(args):
|
||||
@@ -81,11 +92,15 @@ func splitConfigFlag(args []string) (string, []string) {
|
||||
i++
|
||||
case strings.HasPrefix(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:
|
||||
rest = append(rest, args[i])
|
||||
}
|
||||
}
|
||||
return configPath, rest
|
||||
return configPath, install, remove, rest
|
||||
}
|
||||
|
||||
// defaultConfigPath returns gpu-turnstile.env next to the executable.
|
||||
@@ -180,6 +195,14 @@ func serviceCommand(configPath string, args []string) int {
|
||||
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 {
|
||||
// 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.
|
||||
@@ -187,8 +210,8 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
|
||||
"version", version,
|
||||
"listen_ollama", cfg.ListenOllama,
|
||||
"listen_comfy", cfg.ListenComfy,
|
||||
"ollama_url", cfg.OllamaURL,
|
||||
"comfy_url", cfg.ComfyURL,
|
||||
"ollama_url", orDisabled(cfg.OllamaURL),
|
||||
"comfy_url", orDisabled(cfg.ComfyURL),
|
||||
"unload_timeout", cfg.UnloadTimeout,
|
||||
"job_timeout", cfg.JobTimeout,
|
||||
"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)
|
||||
ollamaClient, err := ollama.New(cfg.OllamaURL, log)
|
||||
if err != nil {
|
||||
// Each consumer is enabled by setting its URL; a disabled consumer gets
|
||||
// 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
|
||||
}
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
srv, err := proxy.New(proxy.Config{
|
||||
OllamaURL: cfg.OllamaURL,
|
||||
@@ -253,22 +283,54 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
|
||||
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)
|
||||
if ollamaClient != nil {
|
||||
if err := ollamaClient.Probe(probeCtx); err != nil {
|
||||
log.Warn("ollama probe failed", "url", cfg.OllamaURL, "err", err)
|
||||
}
|
||||
}
|
||||
if comfyClient != nil {
|
||||
if err := comfyClient.Probe(probeCtx); err != nil {
|
||||
log.Warn("comfy probe failed", "url", cfg.ComfyURL, "err", err)
|
||||
}
|
||||
}
|
||||
probeCancel()
|
||||
|
||||
ollamaSrv := &http.Server{Addr: cfg.ListenOllama, Handler: srv.OllamaHandler()}
|
||||
comfySrv := &http.Server{Addr: cfg.ListenComfy, Handler: srv.ComfyHandler()}
|
||||
// Bind the listeners up front so a port conflict fails fast and the
|
||||
// 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)
|
||||
go func() { errCh <- ollamaSrv.ListenAndServe() }()
|
||||
go func() { errCh <- comfySrv.ListenAndServe() }()
|
||||
errCh := make(chan error, len(servers))
|
||||
for i := range servers {
|
||||
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 {
|
||||
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)
|
||||
defer shutdownCancel()
|
||||
ollamaSrv.Shutdown(shutdownCtx)
|
||||
comfySrv.Shutdown(shutdownCtx)
|
||||
for _, s := range servers {
|
||||
s.Shutdown(shutdownCtx)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -6,9 +6,11 @@
|
||||
# ComfyUI --listen 0.0.0.0 --port 8189).
|
||||
services:
|
||||
gpu-turnstile:
|
||||
image: git.rambossek.at/public/gpu-turnstile:v0.1.3
|
||||
image: git.rambossek.at/public/gpu-turnstile:v0.1.5
|
||||
restart: unless-stopped
|
||||
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:
|
||||
OLLAMA_URL: http://host.docker.internal:11435
|
||||
COMFY_URL: http://host.docker.internal:8189
|
||||
|
||||
@@ -3,3 +3,5 @@ module gpu-turnstile
|
||||
go 1.23
|
||||
|
||||
require golang.org/x/sys v0.29.0
|
||||
|
||||
require github.com/coreos/go-systemd/v22 v22.7.0
|
||||
|
||||
@@ -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/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
|
||||
@@ -51,13 +51,12 @@ type Config struct {
|
||||
}
|
||||
|
||||
// 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 {
|
||||
return Config{
|
||||
ListenOllama: ":11434",
|
||||
ListenComfy: ":8188",
|
||||
OllamaURL: "http://127.0.0.1:11435",
|
||||
ComfyURL: "http://127.0.0.1:8189",
|
||||
UnloadTimeout: time.Minute,
|
||||
JobTimeout: 15 * time.Minute,
|
||||
LLMWaitTimeout: 10 * time.Minute,
|
||||
@@ -219,5 +218,8 @@ func Load(getenv func(string) string) (Config, error) {
|
||||
default:
|
||||
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
|
||||
}
|
||||
|
||||
@@ -8,13 +8,21 @@ import (
|
||||
)
|
||||
|
||||
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 {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if cfg.ListenOllama != ":11434" || cfg.ListenComfy != ":8188" {
|
||||
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 {
|
||||
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) {
|
||||
input := `# comment
|
||||
OLLAMA_URL=http://host:11435
|
||||
@@ -91,6 +106,9 @@ func TestLoadErrors(t *testing.T) {
|
||||
if k == tc.key {
|
||||
return tc.value
|
||||
}
|
||||
if k == "OLLAMA_URL" {
|
||||
return "http://127.0.0.1:11435"
|
||||
}
|
||||
return ""
|
||||
})
|
||||
if err == nil {
|
||||
|
||||
+37
-13
@@ -102,24 +102,37 @@ type Server struct {
|
||||
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) {
|
||||
ollamaURL, err := url.Parse(cfg.OllamaURL)
|
||||
if err != nil || ollamaURL.Scheme == "" || ollamaURL.Host == "" {
|
||||
if cfg.OllamaURL == "" && cfg.ComfyURL == "" {
|
||||
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)
|
||||
}
|
||||
comfyURL, err := url.Parse(cfg.ComfyURL)
|
||||
if err != nil || comfyURL.Scheme == "" || comfyURL.Host == "" {
|
||||
ollamaURL = u
|
||||
}
|
||||
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)
|
||||
}
|
||||
comfyURL = u
|
||||
}
|
||||
log := cfg.Log
|
||||
if log == nil {
|
||||
log = slog.Default()
|
||||
}
|
||||
if cfg.UnloadPollInterval > 0 {
|
||||
if cfg.UnloadPollInterval > 0 && cfg.Ollama != nil {
|
||||
cfg.Ollama.PollInterval = cfg.UnloadPollInterval
|
||||
}
|
||||
if cfg.HistoryPollInterval > 0 {
|
||||
if cfg.HistoryPollInterval > 0 && cfg.Comfy != nil {
|
||||
cfg.Comfy.PollInterval = cfg.HistoryPollInterval
|
||||
}
|
||||
freeTimeout := cfg.FreeTimeout
|
||||
@@ -160,7 +173,7 @@ func New(cfg Config) (*Server, error) {
|
||||
max: backoffMax,
|
||||
log: log,
|
||||
}
|
||||
return &Server{
|
||||
s := &Server{
|
||||
cfg: cfg,
|
||||
log: log,
|
||||
logWriter: cfg.LogWriter,
|
||||
@@ -172,9 +185,14 @@ func New(cfg Config) (*Server, error) {
|
||||
busyMode: busyMode,
|
||||
busyStatus: busyStatus,
|
||||
busyRetryAfter: busyRetryAfter,
|
||||
ollamaProxy: newReverseProxy(ollamaURL, retry, log.With("upstream", "ollama")),
|
||||
comfyProxy: newReverseProxy(comfyURL, retry, log.With("upstream", "comfy")),
|
||||
}, nil
|
||||
}
|
||||
if ollamaURL != 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
|
||||
@@ -470,9 +488,13 @@ func (s *Server) OllamaHandler() http.Handler {
|
||||
// ComfyHandler serves the ComfyUI-facing listener.
|
||||
func (s *Server) ComfyHandler() http.Handler {
|
||||
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)
|
||||
return
|
||||
case "/metrics":
|
||||
s.writeMetrics(w)
|
||||
return
|
||||
}
|
||||
if r.Method == http.MethodPost && r.URL.Path == "/prompt" {
|
||||
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())
|
||||
log.Info("image lock acquired")
|
||||
|
||||
if s.cfg.Ollama != nil {
|
||||
uctx, ucancel := context.WithTimeout(r.Context(), s.cfg.UnloadTimeout)
|
||||
elapsed, uerr := s.cfg.Ollama.UnloadAll(uctx)
|
||||
ucancel()
|
||||
@@ -541,6 +564,7 @@ func (s *Server) handlePrompt(w http.ResponseWriter, r *http.Request) {
|
||||
default:
|
||||
log.Info("ollama models unloaded", "seconds", elapsed.Seconds())
|
||||
}
|
||||
}
|
||||
|
||||
cw := &captureWriter{ResponseWriter: w, status: http.StatusOK, limit: s.captureLimit}
|
||||
s.comfyProxy.ServeHTTP(cw, r)
|
||||
@@ -585,7 +609,7 @@ func (s *Server) finishImageJob(promptID string) {
|
||||
s.cfg.Lock.ReleaseImage()
|
||||
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 {
|
||||
wctx, wcancel := context.WithTimeout(context.Background(), s.warmTimeout)
|
||||
if err := s.cfg.Ollama.Warm(wctx, s.cfg.WarmModel); err != nil {
|
||||
|
||||
@@ -451,3 +451,82 @@ func TestLLMBusyWaitTimeoutRetryAfter(t *testing.T) {
|
||||
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")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
@@ -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) {}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,8 +1,8 @@
|
||||
//go:build !windows
|
||||
//go:build !windows && !linux
|
||||
|
||||
// Package service provides the non-Windows stubs for the Windows service
|
||||
// integration. Run falls back to plain signal handling; install/remove
|
||||
// are unsupported.
|
||||
// Package service provides the stubs for platforms without service
|
||||
// integration (Windows uses the SCM, Linux uses systemd). Run falls back
|
||||
// to plain signal handling; install/remove are unsupported.
|
||||
package service
|
||||
|
||||
import (
|
||||
@@ -15,7 +15,7 @@ import (
|
||||
// Name matches the Windows service name.
|
||||
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.
|
||||
func IsService() bool { return false }
|
||||
|
||||
@@ -11,6 +11,6 @@ package update
|
||||
// the RELEASE_SIGNING_KEY repository secret. When empty, the updater refuses
|
||||
// to update (e.g. development builds).
|
||||
var publicKeyPEM = `-----BEGIN PUBLIC KEY-----
|
||||
MCowBQYDK2VwAyEA+gJbSvgeYX58woPQGbSC8x8Zw4OTDiiQ7/19seZKfSQ=
|
||||
MCowBQYDK2VwAyEAqTAJ0CCeAQI7MhFlgc5xNmF/CfvLVUAY3ZAoeAS0tT8=
|
||||
-----END PUBLIC KEY-----
|
||||
`
|
||||
|
||||
Reference in New Issue
Block a user