17 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
mram 19e19281de Pin compose example to v0.1.3
ci / test (push) Successful in 12s
ci / docker (push) Successful in 1m8s
ci / release (push) Failing after 25s
2026-09-20 22:16:09 +02:00
mram 08d02d8fa8 Embed release public key; rename signing secret to RELEASE_SIGNING_KEY 2026-09-20 22:15:19 +02:00
mram d1e01f9b78 Ignore signing/ directory (Ed25519 release keys) 2026-09-20 22:13:10 +02:00
mram bacb26772a Add LLM busy modes: wait (hang) or reject with Retry-After
LLM_BUSY_MODE=reject answers blocked LLM requests immediately with
LLM_BUSY_STATUS (default 503, 429 works) and Retry-After, so routers
like LiteLLM can cool down and retry instead of holding a hung
connection. The default wait mode now also sends Retry-After when
LLM_WAIT_TIMEOUT expires. Document the service account (LocalSystem
default, NT SERVICE virtual-account hardening) and the Program
Files / ProgramData install layout.
2026-09-20 21:43:03 +02:00
mram e363c9e4e4 Native Windows deployment: config file, service, signed auto-update
- internal/config: .env-style config file (gpu-turnstile.env next to the
  exe, -config flag or GPU_TURNSTILE_CONFIG); process env overrides file.
- internal/service: Windows service via golang.org/x/sys/windows/svc —
  graceful SCM stop, 'service install/remove' commands, restart-on-failure
  recovery (also applies staged updates). First external dependency,
  Windows-only; Linux/Docker build unaffected (go.mod stays at 1.23).
- internal/update: polls the Gitea releases API, verifies the Ed25519
  signature of the downloaded binary against an embedded public key
  (openssl-signed by CI), swaps it in next to the running exe, and once
  the GPU lock is idle exits with code 3 so service recovery restarts
  onto the new version. Dev builds and empty pubkey never update.
- CI: tag builds additionally produce gpu-turnstile.exe + .sig + .sha256
  attached to a Gitea release.
- LOG_FILE env var so the service has somewhere to log.
2026-09-20 21:33:29 +02:00
mram 5f6a22c2cf Pin compose example to v0.1.2
ci / test (push) Successful in 47s
ci / docker (push) Successful in 1m5s
2026-09-20 20:07:38 +02:00
mram 42e1386811 Extend retry backoff; rework logging (LOGLEVEL, request lines, colors)
Backoff now covers all dial-phase errors (refused, timeout, DNS), TLS
handshake failures, and 5xx responses with replayable bodies; streamed
POSTs are never replayed to avoid duplicate work. First retry of an
episode logs at WARN, subsequent attempts at INFO.

LOGLEVEL (LOG_LEVEL kept as alias) now defaults to warn: startup logs
version plus every setting; INFO adds one line per incoming request and
per response with status/duration, ANSI-colored in text mode (bypasses
slog's escaping so colors render in docker compose logs); NO_COLOR or
LOG_FORMAT=json disables colors.
2026-09-20 20:00:57 +02:00
mram 20d5439b5f Retry upstream connection-refused with exponential backoff
A refused dial (service down/restarting) is retried with a wait that
doubles from BACKOFF_INITIAL (1s) up to BACKOFF_MAX (60s) until the
upstream answers or the client disconnects. Handles the Windows WSA
errno (10061) as well as POSIX ECONNREFUSED. compose.yaml.example now
uses host.docker.internal like the working local deployment.
2026-09-20 19:45:48 +02:00
mram 3aaa5d80a9 CI: lowercase image repository name; compose example too
ci / test (push) Successful in 49s
ci / docker (push) Successful in 1m7s
Docker registry names must be lowercase, so PUBLIC/gpu-turnstile was
rejected by buildx. The workflow now lowercases gitea.repository.
2026-09-20 19:25:37 +02:00
mram cef845c2e3 Add compose.yaml.example; ignore local compose files 2026-09-20 19:24:26 +02:00
mram cf9be55a6c Make poll intervals and operational timeouts configurable
ci / test (push) Successful in 47s
ci / docker (push) Failing after 30s
New env vars: UNLOAD_POLL_INTERVAL, HISTORY_POLL_INTERVAL, PROBE_TIMEOUT,
FREE_TIMEOUT, WARM_TIMEOUT, SHUTDOWN_TIMEOUT, PROMPT_CAPTURE_LIMIT.
Defaults unchanged; invalid values fail fast at startup.
2026-09-20 19:21:25 +02:00
mram 0f950f3134 CI: use gitea.actor as registry login username
ci / test (push) Successful in 48s
ci / docker (push) Skipped
Gitea authenticates registry pushes by token alone; a single
REGISTRY_TOKEN secret suffices.
2026-09-20 18:55:46 +02:00
mram 73fb8ac3ad CI: log in to registry with REGISTRY_USERNAME/REGISTRY_TOKEN secrets
ci / test (push) Successful in 47s
ci / docker (push) Skipped
The automatic GITEA_TOKEN has no write:package scope, so docker login to
git.rambossek.at failed with unauthorized.
2026-09-20 18:24:54 +02:00
24 changed files with 2690 additions and 209 deletions
+60 -3
View File
@@ -37,13 +37,19 @@ jobs:
exit 1
fi
- name: Compute lowercase image name
id: meta
run: |
REPO=git.rambossek.at/$(echo "${{ gitea.repository }}" | tr '[:upper:]' '[:lower:]')
echo "image=$REPO" >> "$GITHUB_OUTPUT"
- uses: docker/setup-buildx-action@v3
- uses: docker/login-action@v3
with:
registry: git.rambossek.at
username: ${{ gitea.actor }}
password: ${{ secrets.GITEA_TOKEN }}
password: ${{ secrets.REGISTRY_TOKEN }}
- uses: docker/build-push-action@v6
with:
@@ -52,5 +58,56 @@ jobs:
build-args: |
VERSION=${{ gitea.ref_name }}
tags: |
git.rambossek.at/${{ gitea.repository }}:${{ gitea.ref_name }}
git.rambossek.at/${{ gitea.repository }}:latest
${{ steps.meta.outputs.image }}:${{ gitea.ref_name }}
${{ steps.meta.outputs.image }}:latest
# On version tags: build the signed Windows binary and attach it (plus
# signature and checksum) to a Gitea release for the auto-updater.
release:
if: gitea.ref_type == 'tag'
needs: test
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version: "1.23"
- name: Build Windows binary
run: |
GOOS=windows GOARCH=amd64 CGO_ENABLED=0 go build \
-ldflags="-s -w -X main.version=${{ gitea.ref_name }}" \
-o gpu-turnstile.exe ./cmd/gpu-turnstile
- name: Sign and checksum
env:
RELEASE_SIGNING_KEY: ${{ secrets.RELEASE_SIGNING_KEY }}
run: |
# 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
rm -f key.pem
sha256sum gpu-turnstile.exe > gpu-turnstile.exe.sha256
- name: Create release and upload assets
env:
TOKEN: ${{ secrets.GITEA_TOKEN }}
API: https://git.rambossek.at/api/v1/repos/${{ gitea.repository }}
TAG: ${{ gitea.ref_name }}
run: |
set -e
ID=$(curl -sf -H "Authorization: token $TOKEN" "$API/releases/tags/$TAG" | jq -r .id || true)
if [ -z "$ID" ] || [ "$ID" = "null" ]; then
ID=$(curl -sf -X POST -H "Authorization: token $TOKEN" \
-H "Content-Type: application/json" \
-d "{\"tag_name\":\"$TAG\",\"name\":\"$TAG\"}" \
"$API/releases" | jq -r .id)
fi
for f in gpu-turnstile.exe gpu-turnstile.exe.sig gpu-turnstile.exe.sha256; do
curl -sf -X POST -H "Authorization: token $TOKEN" \
-F "attachment=@$f" "$API/releases/$ID/assets?name=$f" > /dev/null
echo "uploaded $f"
done
+4
View File
@@ -1,2 +1,6 @@
/gpu-turnstile
/gpu-turnstile.exe
/compose.yaml
/compose.yml
/signing/
+107 -13
View File
@@ -18,7 +18,10 @@ Open WebUI / n8n ────► :8188 ───┘
- LLM endpoints (`/api/generate`, `/api/chat`, `/api/embed`, `/v1/*`) take the
LLM lock: concurrent requests allowed, but blocked while an image job is
active or waiting (image priority).
active or waiting (image priority). Blocked requests either hang until the
lock is free (`LLM_BUSY_MODE=wait`, default) or fail immediately with 503
(or 429) + `Retry-After` (`LLM_BUSY_MODE=reject`) — the latter lets routers
like LiteLLM cool down and retry instead of holding a hung connection.
- `POST /prompt` on the ComfyUI listener takes the image lock: new LLM
requests block, in-flight LLMs drain, Ollama models are unloaded, the prompt
is forwarded, and the lock is held until the job finishes and ComfyUI frees
@@ -26,31 +29,63 @@ 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
All configuration is via environment variables; invalid values fail at
startup.
Configuration comes from environment variables and/or an `.env`-style
config file (`KEY=VALUE` lines, `#` comments). File lookup order:
`-config <path>` flag, then `GPU_TURNSTILE_CONFIG`, then
`gpu-turnstile.env` next to the executable. Process environment variables
override file values. Invalid values fail at startup.
| Var | Default | Meaning |
|---|---|---|
| `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 |
| `LLM_WAIT_TIMEOUT` | `10m` | Max lock wait for an LLM request before 503 (wait mode) |
| `LLM_BUSY_MODE` | `wait` | `wait` = hold blocked LLM requests; `reject` = fail them immediately |
| `LLM_BUSY_STATUS` | `503` | HTTP status for rejected LLM requests in reject mode (400599, e.g. 429) |
| `BUSY_RETRY_AFTER` | `30` | Seconds sent as `Retry-After` on busy responses (both modes) |
| `WARM_MODEL` | _(empty)_ | Model to reload after an image job (off by default) |
| `LOG_LEVEL` | `info` | `debug` logs every lock transition |
| `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 |
| `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 |
| `BACKOFF_INITIAL` | `1s` | First retry wait when an upstream refuses a connection |
| `BACKOFF_MAX` | `60s` | Cap for the exponential retry backoff |
| `PROMPT_CAPTURE_LIMIT` | `65536` | Bytes of the `/prompt` response buffered to find `prompt_id` (pass-through is unaffected) |
| `AUTO_UPDATE` | `true` | Poll the Gitea releases API for signed updates |
| `UPDATE_INTERVAL` | `6h` | Auto-update check interval |
| `UPDATE_REPO` | `https://git.rambossek.at/PUBLIC/gpu-turnstile` | Repository checked for releases |
| `UPDATE_ASSET` | `gpu-turnstile.exe` | Release asset to download |
## 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`.
- Logs: startup logs the version and every setting (visible even at the
default `warn` level). With `LOGLEVEL=info` or `debug`, every request
logs a `-->` incoming line and a `<--` response line with status and
duration — ANSI-colored (cyan incoming; green/yellow/red by status
class) in text mode, which renders in `docker compose logs` on Windows
Terminal. Set `NO_COLOR` to disable colors.
## Build and run
@@ -59,6 +94,56 @@ go build ./cmd/gpu-turnstile
./gpu-turnstile
```
### 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 --install-service # auto-start service, recovery = restart
gpu-turnstile.exe --remove-service
```
The service uses the config file (services have no convenient
environment); set `LOG_FILE` in it since there is no console.
Suggested layout: `C:\Program Files\gpu-turnstile\` for the exe and
`gpu-turnstile.env`, logs under `C:\ProgramData\gpu-turnstile\` via
`LOG_FILE`. The service runs as `LocalSystem` by default, which can write
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,
and — once the GPU lock is idle — restarts the service onto the new
version. Disable with `AUTO_UPDATE=false`. Releases are signed by CI with
OpenSSL; the matching public key lives in `internal/update/pubkey.go`
(one-time setup: `openssl genpkey -algorithm ed25519 -out private.pem`,
`openssl pkey -in private.pem -pubout -out public.pem`; private key goes
to the `RELEASE_SIGNING_KEY` repo secret, public key is committed).
### Docker
```sh
docker build -t gpu-turnstile .
docker run --rm -p 11434:11434 -p 8188:8188 \
@@ -69,9 +154,14 @@ docker run --rm -p 11434:11434 -p 8188:8188 \
Releases are built by Gitea Actions (`.gitea/workflows/ci.yml`): every push
runs `go vet` and `go test -race`, and pushing a semantic-version tag
`vX.Y.Z` builds and publishes
`git.rambossek.at/<owner>/gpu-turnstile:vX.Y.Z` (and updates `:latest`).
No images are built from branches.
`vX.Y.Z` publishes the container image
(`git.rambossek.at/<owner>/gpu-turnstile:vX.Y.Z` plus `:latest`) and a
signed Windows binary attached to a Gitea release. Nothing is built from
branches.
The registry login needs one repository secret (Settings → Actions →
Secrets): `REGISTRY_TOKEN` — an access token with `write:package` scope.
The automatic `GITEA_TOKEN` cannot push packages.
## Development
@@ -80,13 +170,17 @@ go vet ./...
go test -race ./...
```
Stdlib only, Go 1.23+. Layout:
Go 1.23+; the only external dependency is `golang.org/x/sys` (Windows
service integration, unused in the Linux build). Layout:
```
cmd/gpu-turnstile/main.go wiring, config, listeners
cmd/gpu-turnstile/main.go wiring, config, listeners, service + updater
internal/lock/ two-mode lock (LLM readers / image writer, FIFO)
internal/ollama/ ps / unload / warm client
internal/comfy/ history poll / free client
internal/proxy/ handlers for both listeners
internal/metrics/ Prometheus exposition, no dependencies
internal/config/ env + .env file configuration
internal/update/ signed auto-updater
internal/service/ Windows service integration
```
+161 -19
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
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"
@@ -51,6 +66,11 @@ are LLM requests and the single "writer" is an image job):
`image` **or while an image job is waiting**. Then state := `llm`, n++.
On completion (response fully written, including streamed bodies, or client
disconnect) n--; if n == 0 state := `idle`.
`LLM_BUSY_MODE` selects what a blocked LLM request sees: `wait` (default)
hangs until the lock is free or `LLM_WAIT_TIMEOUT` expires (then 503 +
`Retry-After`); `reject` answers immediately with `LLM_BUSY_STATUS`
(default 503; 429 works too) + `Retry-After: BUSY_RETRY_AFTER`, which
routers like LiteLLM honor for cooldowns/retries.
- **Image job**: `AcquireImage()` marks "image pending" (so no new LLM
requests start), waits until n == 0, sets state := `image`. Released after
the ComfyUI job finished and models were freed.
@@ -78,15 +98,18 @@ 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 500 ms until empty or `UNLOAD_TIMEOUT`. On timeout: log and
`/api/ps` every `UNLOAD_POLL_INTERVAL` (default 500 ms) until empty or
`UNLOAD_TIMEOUT`. On timeout: log and
continue (degrade, don't fail the user's request).
3. Forward the original request body to ComfyUI `/prompt`, return status,
headers and body to the caller unchanged, flush.
4. If the response is 200 and contains `prompt_id`: in a goroutine, poll
`GET /history/<prompt_id>` every 1 s until the entry has
`GET /history/<prompt_id>` every `HISTORY_POLL_INTERVAL` (default 1 s)
until the entry has
`status.completed == true`, `status.status_str == "error"`, or
`JOB_TIMEOUT`. Then `POST /free {"unload_models":true,"free_memory":true}`.
Then release the image lock.
@@ -100,32 +123,125 @@ load time. Off by default.
## Configuration (env)
Configuration comes from environment variables and/or an `.env`-style
config file (`KEY=VALUE` lines, `#` comments). File lookup order:
`-config <path>` flag, then `GPU_TURNSTILE_CONFIG`, then
`gpu-turnstile.env` next to the executable. Process environment variables
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 |
| `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 |
| `LLM_WAIT_TIMEOUT` | `10m` | max time an LLM request waits for the lock before 503 (wait mode) |
| `LLM_BUSY_MODE` | `wait` | `wait` = hold blocked LLM requests; `reject` = fail them immediately |
| `LLM_BUSY_STATUS` | `503` | HTTP status for rejected LLM requests in reject mode (400599, e.g. 429) |
| `BUSY_RETRY_AFTER` | `30` | seconds sent as `Retry-After` on busy responses (both modes) |
| `WARM_MODEL` | `` | optional model to reload after an image job |
| `LOG_LEVEL` | `info` | `debug` logs every lock transition |
| `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 |
| `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 |
| `BACKOFF_INITIAL` | `1s` | first retry wait when an upstream refuses a connection |
| `BACKOFF_MAX` | `60s` | cap for the exponential retry backoff |
| `PROMPT_CAPTURE_LIMIT` | `65536` | bytes of the `/prompt` response buffered to find `prompt_id` (pass-through is unaffected) |
| `AUTO_UPDATE` | `true` | poll the Gitea releases API for signed updates |
| `UPDATE_INTERVAL` | `6h` | auto-update check interval |
| `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 deployment (Windows and Linux)
The binary runs natively on Windows (the current primary deployment) and on
Linux with systemd (the future GPU server), as well as in Docker.
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
self-updates — Program Files is writable by LocalSystem and admins, which
is why running as the default `LocalSystem` account is the simple choice.
- **Account**: the default `LocalSystem` works out of the box. For least
privilege, create the service with the virtual account
`NT SERVICE\gpu-turnstile` and grant it write access to the install and
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)
and verifies an Ed25519 signature against the public key embedded in
`internal/update/pubkey.go`. A verified binary is swapped in next to the
running exe (rename-aside, allowed on Windows), and once the GPU lock is
idle the process exits with code 3 so the service recovery restarts it
on the new version. Interactive runs only log "restart to apply".
`dev` builds and builds without an embedded public key never update.
- **Signing setup (one time)**: `openssl genpkey -algorithm ed25519 -out
private.pem`; `openssl pkey -in private.pem -pubout -out public.pem`.
Private key → repo secret `RELEASE_SIGNING_KEY`; public key → committed into
`internal/update/pubkey.go`. CI signs release binaries with
`openssl pkeyutl -sign -rawin`.
## Observability
- `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`
(histogram, label `kind="llm|image"`), `gpu_turnstile_unload_seconds`.
- Structured logs (`log/slog`, JSON when `LOG_FORMAT=json`), one line per
state transition and per image job phase with `prompt_id`.
state transition and per image job phase with `prompt_id`. Startup logs
the version and every setting (visible even at the default `warn`
level). With `LOGLEVEL=info` or `debug`, every request logs a `-->`
incoming line and a `<--` response line with status and duration —
ANSI-colored (cyan incoming; green/yellow/red by status class) in text
mode, which renders in `docker compose logs` on Windows Terminal. Set
`NO_COLOR` to disable colors.
## Edge cases to handle
@@ -137,6 +253,14 @@ start (`/api/version`, `/system_stats`); failure is logged, not fatal.
`JOB_TIMEOUT` releases the lock; log at warn.
- Ollama unreachable during unload: continue with the image job; the whole
point is not to block users on a misbehaving neighbour.
- Upstream unreachable while proxying (connection refused, dial timeout,
DNS failure, TLS handshake error): retry with exponential backoff —
`BACKOFF_INITIAL`, doubling per attempt, capped at `BACKOFF_MAX` — until
the upstream answers or the client disconnects. These are safe to retry:
the request never reached the upstream application. 5xx responses are
retried the same way, but only when the request body can be replayed
(GETs, or bodies with `GetBody`); streamed POSTs are never replayed to
avoid duplicate work such as a double-enqueued ComfyUI prompt.
- `POST /prompt` with a body that ComfyUI rejects (400): lock released
immediately, body passed back.
- Websocket `/ws` connections are long-lived and never take the lock.
@@ -146,11 +270,15 @@ start (`/api/version`, `/system_stats`); failure is logged, not fatal.
```
gpu-turnstile/
cmd/gpu-turnstile/main.go # wiring, config, listeners
cmd/gpu-turnstile/main.go # wiring, config, listeners, service + updater
internal/lock/lock.go # two-mode lock + tests
internal/ollama/client.go # ps / unload / warm
internal/comfy/client.go # history poll / free
internal/proxy/ # handlers for both listeners
internal/metrics/ # Prometheus exposition
internal/config/ # env + .env file configuration
internal/update/ # signed auto-updater (public key in pubkey.go)
internal/service/ # Windows SCM + Linux systemd (notify/watchdog) integration
Dockerfile
.gitea/workflows/ci.yml
README.md
@@ -172,11 +300,17 @@ are new.
held until `/free` was called.
- Streaming test: fake Ollama emits chunks with delays; assert the client
receives the first chunk before the last is sent (no buffering).
- `internal/config`: env-file parsing, precedence, fail-fast values.
- `internal/update`: fake Gitea releases API; staged update happy path,
tampered signature rejected, older versions and dev builds skipped.
## Build and CI
- Go 1.23+, stdlib only. `CGO_ENABLED=0`, `-ldflags="-s -w"`, version from
`git describe` injected via `-X main.version=`.
- 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
`scratch`), non-root user, `EXPOSE 8188 11434`,
`ENTRYPOINT ["/gpu-turnstile"]`.
@@ -185,16 +319,24 @@ are new.
available in the runner image
2. on a version tag only (`vX.Y.Z`, enforced): build the image with buildx
and push it to the Gitea registry
`git.rambossek.at/<owner>/gpu-turnstile` tagged `:<tag>` and `:latest`,
using the workflow token (`${{ secrets.GITEA_TOKEN }}` / `gitea.actor`)
- Release: a git tag `vX.Y.Z` produces the versioned image; the Open WebUI
compose pins that tag. No images are built from branches.
`git.rambossek.at/<owner>/gpu-turnstile` tagged `:<tag>` and `:latest`
(the repository path is lowercased in the workflow; Docker registry
names must be lowercase).
Login uses the repo secret `REGISTRY_TOKEN` (an access token with
`write:package` scope) because the automatic `GITEA_TOKEN` cannot push
packages; the username is just `gitea.actor`.
3. on a version tag: also build the Windows binary, sign it with OpenSSL
(`RELEASE_SIGNING_KEY` secret), and attach `gpu-turnstile.exe`, `.sig` and
`.sha256` to a Gitea release for the auto-updater.
- Release: a git tag `vX.Y.Z` produces the versioned image and the signed
Windows binary; the Open WebUI compose pins that tag. No images or
binaries are built from branches.
## Deployment (target)
```yaml
gpu-turnstile:
image: git.rambossek.at/<owner>/gpu-turnstile:v0.1.0
image: git.rambossek.at/<owner>/gpu-turnstile:v0.1.0 # owner lowercased, e.g. "public"
environment:
OLLAMA_URL: http://<workstation-ip>:11435
COMFY_URL: http://<workstation-ip>:8189
+345 -129
View File
@@ -6,188 +6,404 @@ import (
"context"
"errors"
"fmt"
"io"
"log/slog"
"net"
"net/http"
"os"
"os/signal"
"path/filepath"
"strings"
"syscall"
"time"
"gpu-turnstile/internal/comfy"
"gpu-turnstile/internal/config"
"gpu-turnstile/internal/lock"
"gpu-turnstile/internal/metrics"
"gpu-turnstile/internal/ollama"
"gpu-turnstile/internal/proxy"
"gpu-turnstile/internal/service"
"gpu-turnstile/internal/update"
)
// version is injected at build time via -ldflags "-X main.version=...".
var version = "dev"
type config struct {
listenOllama string
listenComfy string
ollamaURL string
comfyURL string
unloadTimeout time.Duration
jobTimeout time.Duration
llmWaitTimeout time.Duration
warmModel string
logLevel slog.Level
logJSON bool
}
func envDuration(getenv func(string) string, name string, dst *time.Duration) error {
v := getenv(name)
if v == "" {
return nil
}
d, err := time.ParseDuration(v)
if err != nil {
return fmt.Errorf("%s: %w", name, err)
}
*dst = d
return nil
}
func loadConfig(getenv func(string) string) (config, error) {
cfg := 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,
logLevel: slog.LevelInfo,
}
for _, e := range []struct {
name string
dst *string
}{
{"LISTEN_OLLAMA", &cfg.listenOllama},
{"LISTEN_COMFY", &cfg.listenComfy},
{"OLLAMA_URL", &cfg.ollamaURL},
{"COMFY_URL", &cfg.comfyURL},
{"WARM_MODEL", &cfg.warmModel},
} {
if v := getenv(e.name); v != "" {
*e.dst = v
}
}
for _, e := range []struct {
name string
dst *time.Duration
}{
{"UNLOAD_TIMEOUT", &cfg.unloadTimeout},
{"JOB_TIMEOUT", &cfg.jobTimeout},
{"LLM_WAIT_TIMEOUT", &cfg.llmWaitTimeout},
} {
if err := envDuration(getenv, e.name, e.dst); err != nil {
return cfg, err
}
}
if v := getenv("LOG_LEVEL"); v != "" {
var level slog.Level
if err := level.UnmarshalText([]byte(v)); err != nil {
return cfg, fmt.Errorf("LOG_LEVEL: %w", err)
}
cfg.logLevel = level
}
switch strings.ToLower(getenv("LOG_FORMAT")) {
case "", "text":
case "json":
cfg.logJSON = true
default:
return cfg, fmt.Errorf("LOG_FORMAT: must be \"text\" or \"json\"")
}
return cfg, nil
}
// exitCodeUpdate tells the service recovery configuration to restart the
// process: a signed update has been staged and the GPU lock is idle.
const exitCodeUpdate = 3
func main() {
cfg, err := loadConfig(os.Getenv)
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] [--install-service | --remove-service]\n")
fmt.Fprintf(os.Stderr, " gpu-turnstile service install|remove [-config path]\n")
os.Exit(2)
}
if exePath, err := os.Executable(); err == nil {
update.CleanupOld(exePath)
}
cfg, err := loadMergedConfig(configPath)
if err != nil {
fmt.Fprintf(os.Stderr, "gpu-turnstile: %v\n", err)
os.Exit(1)
}
log, logOut, logCloser := newLogger(cfg)
defer logCloser.Close()
opts := &slog.HandlerOptions{Level: cfg.logLevel}
var handler slog.Handler = slog.NewTextHandler(os.Stderr, opts)
if cfg.logJSON {
handler = slog.NewJSONHandler(os.Stderr, opts)
if service.IsService() {
if err := service.Run(func(ctx context.Context) error { return run(ctx, cfg, log, logOut, true) }); err != nil {
log.Error("service failed", "err", err)
os.Exit(1)
}
return
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
if err := run(ctx, cfg, log, logOut, false); err != nil {
log.Error("listener failed", "err", err)
os.Exit(1)
}
}
// 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):
configPath = args[i+1]
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, install, remove, rest
}
// defaultConfigPath returns gpu-turnstile.env next to the executable.
func defaultConfigPath() string {
exe, err := os.Executable()
if err != nil {
return "gpu-turnstile.env"
}
return filepath.Join(filepath.Dir(exe), "gpu-turnstile.env")
}
// resolveConfigPath applies the precedence: -config flag, then
// GPU_TURNSTILE_CONFIG, then the default next to the executable.
func resolveConfigPath(flagValue string) string {
if flagValue != "" {
return flagValue
}
if v := os.Getenv("GPU_TURNSTILE_CONFIG"); v != "" {
return v
}
return defaultConfigPath()
}
// loadMergedConfig reads the config file (if present) and overlays process
// environment variables on top. A missing file is fine; an unreadable or
// malformed file is fatal.
func loadMergedConfig(flagValue string) (config.Config, error) {
path := resolveConfigPath(flagValue)
values := map[string]string{}
if f, err := os.Open(path); err == nil {
defer f.Close()
parsed, err := config.ParseEnvFile(f)
if err != nil {
return config.Config{}, fmt.Errorf("%s: %w", path, err)
}
values = parsed
} else if !errors.Is(err, os.ErrNotExist) {
return config.Config{}, fmt.Errorf("read config file: %w", err)
}
getenv := func(key string) string {
if v := os.Getenv(key); v != "" {
return v
}
return values[key]
}
return config.Load(getenv)
}
// newLogger builds the slog logger and returns the output writer (stderr or
// the opened LOG_FILE) plus a closer for it.
func newLogger(cfg config.Config) (*slog.Logger, io.Writer, io.Closer) {
out := io.Writer(os.Stderr)
closer := io.NopCloser(nil)
if cfg.LogFile != "" {
if f, err := os.OpenFile(cfg.LogFile, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o644); err == nil {
out = f
closer = f
} else {
fmt.Fprintf(os.Stderr, "gpu-turnstile: cannot open LOG_FILE %s: %v (logging to stderr)\n", cfg.LogFile, err)
}
}
opts := &slog.HandlerOptions{Level: cfg.LogLevel}
var handler slog.Handler = slog.NewTextHandler(out, opts)
if cfg.LogJSON {
handler = slog.NewJSONHandler(out, opts)
}
log := slog.New(handler)
slog.SetDefault(log)
return log, out, closer
}
log.Info("starting gpu-turnstile",
func serviceCommand(configPath string, args []string) int {
if len(args) != 1 || (args[0] != "install" && args[0] != "remove") {
fmt.Fprintf(os.Stderr, "usage: gpu-turnstile service install|remove [-config path]\n")
return 2
}
var err error
if args[0] == "install" {
path := resolveConfigPath(configPath)
if abs, absErr := filepath.Abs(path); absErr == nil {
path = abs
}
err = service.Install(path)
} else {
err = service.Remove()
}
if err != nil {
fmt.Fprintf(os.Stderr, "gpu-turnstile service %s: %v\n", args[0], err)
return 1
}
fmt.Printf("service %s: %sd\n", service.Name, args[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 {
// 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.
log.Log(ctx, slog.LevelWarn, "starting gpu-turnstile",
"version", version,
"listen_ollama", cfg.listenOllama,
"listen_comfy", cfg.listenComfy,
"ollama_url", cfg.ollamaURL,
"comfy_url", cfg.comfyURL,
"listen_ollama", cfg.ListenOllama,
"listen_comfy", cfg.ListenComfy,
"ollama_url", orDisabled(cfg.OllamaURL),
"comfy_url", orDisabled(cfg.ComfyURL),
"unload_timeout", cfg.UnloadTimeout,
"job_timeout", cfg.JobTimeout,
"llm_wait_timeout", cfg.LLMWaitTimeout,
"llm_busy_mode", cfg.LLMBusyMode,
"llm_busy_status", cfg.LLMBusyStatus,
"busy_retry_after", cfg.BusyRetryAfter,
"unload_poll_interval", cfg.UnloadPollInterval,
"history_poll_interval", cfg.HistoryPollInterval,
"probe_timeout", cfg.ProbeTimeout,
"free_timeout", cfg.FreeTimeout,
"warm_timeout", cfg.WarmTimeout,
"shutdown_timeout", cfg.ShutdownTimeout,
"backoff_initial", cfg.BackoffInitial,
"backoff_max", cfg.BackoffMax,
"prompt_capture_limit", cfg.PromptCaptureLimit,
"warm_model", cfg.WarmModel,
"auto_update", cfg.AutoUpdate,
"update_interval", cfg.UpdateInterval,
"update_repo", cfg.UpdateRepo,
"update_asset", cfg.UpdateAsset,
"log_level", cfg.LogLevel,
"log_format", map[bool]string{true: "json", false: "text"}[cfg.LogJSON],
"log_file", cfg.LogFile,
)
ollamaClient, err := ollama.New(cfg.ollamaURL, log)
if err != nil {
log.Error("invalid configuration", "err", err)
os.Exit(1)
lk := lock.New(log)
// 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 {
log.Error("invalid configuration", "err", err)
os.Exit(1)
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,
ComfyURL: cfg.comfyURL,
Lock: lock.New(log),
Ollama: ollamaClient,
Comfy: comfyClient,
Metrics: metrics.New(),
Log: log,
LLMWaitTimeout: cfg.llmWaitTimeout,
UnloadTimeout: cfg.unloadTimeout,
JobTimeout: cfg.jobTimeout,
WarmModel: cfg.warmModel,
OllamaURL: cfg.OllamaURL,
ComfyURL: cfg.ComfyURL,
Lock: lk,
Ollama: ollamaClient,
Comfy: comfyClient,
Metrics: metrics.New(),
Log: log,
LogColor: !cfg.LogJSON && cfg.LogFile == "" && os.Getenv("NO_COLOR") == "",
LogWriter: logOut,
LLMWaitTimeout: cfg.LLMWaitTimeout,
UnloadTimeout: cfg.UnloadTimeout,
JobTimeout: cfg.JobTimeout,
UnloadPollInterval: cfg.UnloadPollInterval,
HistoryPollInterval: cfg.HistoryPollInterval,
FreeTimeout: cfg.FreeTimeout,
WarmTimeout: cfg.WarmTimeout,
LLMBusyMode: cfg.LLMBusyMode,
LLMBusyStatus: cfg.LLMBusyStatus,
BusyRetryAfter: cfg.BusyRetryAfter,
BackoffInitial: cfg.BackoffInitial,
BackoffMax: cfg.BackoffMax,
PromptCaptureLimit: cfg.PromptCaptureLimit,
WarmModel: cfg.WarmModel,
})
if err != nil {
log.Error("invalid configuration", "err", err)
os.Exit(1)
return err
}
// Probe both upstreams once; failure is logged, not fatal.
probeCtx, probeCancel := context.WithTimeout(context.Background(), 5*time.Second)
if err := ollamaClient.Probe(probeCtx); err != nil {
log.Warn("ollama probe failed", "url", cfg.ollamaURL, "err", err)
// 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 err := comfyClient.Probe(probeCtx); err != nil {
log.Warn("comfy probe failed", "url", cfg.comfyURL, "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])
}
sigCtx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
// 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)
}
select {
case err := <-errCh:
if err != nil && !errors.Is(err, http.ErrServerClosed) {
log.Error("listener failed", "err", err)
os.Exit(1)
return err
}
case <-sigCtx.Done():
case <-ctx.Done():
log.Info("shutting down")
}
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 10*time.Second)
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
}
// 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
// process with the new binary. Interactively it only logs.
func updateLoop(ctx context.Context, cfg config.Config, log *slog.Logger, lk *lock.Lock, isService bool) {
exePath, err := os.Executable()
if err != nil {
log.Warn("auto-update disabled: cannot locate executable", "err", err)
return
}
u := &update.Updater{Repo: cfg.UpdateRepo, Asset: cfg.UpdateAsset, Version: version, Log: log}
for {
staged, err := u.Check(ctx, exePath)
if err != nil && ctx.Err() == nil {
log.Warn("auto-update check failed", "err", err)
}
if staged {
if !isService {
log.Warn("auto-update: new binary staged; restart gpu-turnstile to apply")
return
}
log.Warn("auto-update: staged; restarting once the GPU is idle")
if waitForIdle(ctx, lk, 24*time.Hour) {
log.Warn("auto-update: restarting to apply update")
os.Exit(exitCodeUpdate)
}
return
}
select {
case <-ctx.Done():
return
case <-time.After(cfg.UpdateInterval):
}
}
}
// waitForIdle polls the lock until no LLM or image work is active or
// pending, max at most. Returns false on timeout or cancellation.
func waitForIdle(ctx context.Context, lk *lock.Lock, max time.Duration) bool {
deadline := time.Now().Add(max)
for {
if state, _, pending := lk.Snapshot(); state == lock.StateIdle && !pending {
return true
}
if time.Now().After(deadline) {
return false
}
select {
case <-ctx.Done():
return false
case <-time.After(5 * time.Second):
}
}
}
+36
View File
@@ -0,0 +1,36 @@
# Example deployment for gpu-turnstile. Copy to compose.yaml and adjust.
#
# gpu-turnstile listens on the ports the services normally use; the actual
# Ollama and ComfyUI instances run one port higher (11435 / 8189) and must
# bind 0.0.0.0 so the container can reach them (OLLAMA_HOST=0.0.0.0:11435,
# ComfyUI --listen 0.0.0.0 --port 8189).
services:
gpu-turnstile:
image: git.rambossek.at/public/gpu-turnstile:v0.1.4
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
# Services on another machine: use its LAN IP instead, e.g.
# OLLAMA_URL: http://192.168.1.10:11435
# COMFY_URL: http://192.168.1.10:8189
# UNLOAD_TIMEOUT: 60s
# JOB_TIMEOUT: 15m
# LLM_WAIT_TIMEOUT: 10m
# LLM_BUSY_MODE: reject # wait (default) hangs; reject fails fast
# LLM_BUSY_STATUS: 429 # status in reject mode (default 503)
# BUSY_RETRY_AFTER: 30 # Retry-After seconds on busy responses
# WARM_MODEL: qwen3:14b
# LOG_LEVEL: info
# LOG_FORMAT: json
ports:
- "11434:11434" # LiteLLM api_base -> http://gpu-turnstile:11434
- "8188:8188" # Open WebUI COMFYUI_BASE_URL -> http://gpu-turnstile:8188
networks: [internal]
networks:
internal:
external: true
+4
View File
@@ -1,3 +1,7 @@
module gpu-turnstile
go 1.23
require golang.org/x/sys v0.29.0
require github.com/coreos/go-systemd/v22 v22.7.0
+4
View File
@@ -0,0 +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=
+225
View File
@@ -0,0 +1,225 @@
// Package config loads gpu-turnstile's configuration from environment
// variables and an optional .env-style config file.
package config
import (
"bufio"
"fmt"
"io"
"log/slog"
"strconv"
"strings"
"time"
)
// Config holds every gpu-turnstile setting.
type Config struct {
ListenOllama string
ListenComfy string
OllamaURL string
ComfyURL string
UnloadTimeout time.Duration
JobTimeout time.Duration
LLMWaitTimeout time.Duration
UnloadPollInterval time.Duration
HistoryPollInterval time.Duration
ProbeTimeout time.Duration
FreeTimeout time.Duration
WarmTimeout time.Duration
ShutdownTimeout time.Duration
BackoffInitial time.Duration
BackoffMax time.Duration
PromptCaptureLimit int64
AutoUpdate bool
UpdateInterval time.Duration
UpdateRepo string
UpdateAsset string
// LLMBusyMode is "wait" (hold requests until the lock is free or
// LLMWaitTimeout expires) or "reject" (immediately answer with
// LLMBusyStatus + Retry-After when an image job is active or pending).
LLMBusyMode string
LLMBusyStatus int
BusyRetryAfter int
WarmModel string
LogLevel slog.Level
LogJSON bool
LogFile string
}
// Defaults returns the configuration used when neither the environment nor
// 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",
UnloadTimeout: time.Minute,
JobTimeout: 15 * time.Minute,
LLMWaitTimeout: 10 * time.Minute,
UnloadPollInterval: 500 * time.Millisecond,
HistoryPollInterval: time.Second,
ProbeTimeout: 5 * time.Second,
FreeTimeout: 30 * time.Second,
WarmTimeout: 2 * time.Minute,
ShutdownTimeout: 10 * time.Second,
BackoffInitial: time.Second,
BackoffMax: time.Minute,
PromptCaptureLimit: 64 * 1024,
AutoUpdate: true,
UpdateInterval: 6 * time.Hour,
UpdateRepo: "https://git.rambossek.at/PUBLIC/gpu-turnstile",
UpdateAsset: "gpu-turnstile.exe",
LLMBusyMode: "wait",
LLMBusyStatus: 503,
BusyRetryAfter: 30,
LogLevel: slog.LevelWarn,
}
}
// ParseEnvFile parses a .env-style file: KEY=VALUE lines, blank lines and
// #-comments are ignored, no quoting. A line without '=' is an error.
func ParseEnvFile(r io.Reader) (map[string]string, error) {
values := make(map[string]string)
scanner := bufio.NewScanner(r)
scanner.Buffer(make([]byte, 64*1024), 1024*1024)
lineNo := 0
for scanner.Scan() {
lineNo++
line := strings.TrimSpace(scanner.Text())
if line == "" || strings.HasPrefix(line, "#") {
continue
}
key, value, ok := strings.Cut(line, "=")
if !ok {
return nil, fmt.Errorf("line %d: expected KEY=VALUE", lineNo)
}
key = strings.TrimSpace(key)
if key == "" {
return nil, fmt.Errorf("line %d: empty key", lineNo)
}
values[key] = strings.TrimSpace(value)
}
return values, scanner.Err()
}
func envDuration(getenv func(string) string, name string, dst *time.Duration) error {
v := getenv(name)
if v == "" {
return nil
}
d, err := time.ParseDuration(v)
if err != nil {
return fmt.Errorf("%s: %w", name, err)
}
*dst = d
return nil
}
// Load overlays values from getenv onto the Defaults. Unknown keys are
// ignored. Invalid values are fatal.
func Load(getenv func(string) string) (Config, error) {
cfg := Defaults()
for _, e := range []struct {
name string
dst *string
}{
{"LISTEN_OLLAMA", &cfg.ListenOllama},
{"LISTEN_COMFY", &cfg.ListenComfy},
{"OLLAMA_URL", &cfg.OllamaURL},
{"COMFY_URL", &cfg.ComfyURL},
{"WARM_MODEL", &cfg.WarmModel},
{"UPDATE_REPO", &cfg.UpdateRepo},
{"UPDATE_ASSET", &cfg.UpdateAsset},
{"LOG_FILE", &cfg.LogFile},
} {
if v := getenv(e.name); v != "" {
*e.dst = v
}
}
for _, e := range []struct {
name string
dst *time.Duration
}{
{"UNLOAD_TIMEOUT", &cfg.UnloadTimeout},
{"JOB_TIMEOUT", &cfg.JobTimeout},
{"LLM_WAIT_TIMEOUT", &cfg.LLMWaitTimeout},
{"UNLOAD_POLL_INTERVAL", &cfg.UnloadPollInterval},
{"HISTORY_POLL_INTERVAL", &cfg.HistoryPollInterval},
{"PROBE_TIMEOUT", &cfg.ProbeTimeout},
{"FREE_TIMEOUT", &cfg.FreeTimeout},
{"WARM_TIMEOUT", &cfg.WarmTimeout},
{"SHUTDOWN_TIMEOUT", &cfg.ShutdownTimeout},
{"BACKOFF_INITIAL", &cfg.BackoffInitial},
{"BACKOFF_MAX", &cfg.BackoffMax},
{"UPDATE_INTERVAL", &cfg.UpdateInterval},
} {
if err := envDuration(getenv, e.name, e.dst); err != nil {
return cfg, err
}
}
if v := getenv("PROMPT_CAPTURE_LIMIT"); v != "" {
n, err := strconv.ParseInt(v, 10, 64)
if err != nil || n < 0 {
return cfg, fmt.Errorf("PROMPT_CAPTURE_LIMIT: must be a non-negative integer (bytes)")
}
cfg.PromptCaptureLimit = n
}
if v := getenv("AUTO_UPDATE"); v != "" {
b, err := strconv.ParseBool(v)
if err != nil {
return cfg, fmt.Errorf("AUTO_UPDATE: must be a boolean (true/false)")
}
cfg.AutoUpdate = b
}
if v := getenv("LLM_BUSY_MODE"); v != "" {
if v != "wait" && v != "reject" {
return cfg, fmt.Errorf("LLM_BUSY_MODE: must be \"wait\" or \"reject\"")
}
cfg.LLMBusyMode = v
}
if v := getenv("LLM_BUSY_STATUS"); v != "" {
n, err := strconv.Atoi(v)
if err != nil || n < 400 || n > 599 {
return cfg, fmt.Errorf("LLM_BUSY_STATUS: must be an HTTP status in 400-599")
}
cfg.LLMBusyStatus = n
}
if v := getenv("BUSY_RETRY_AFTER"); v != "" {
n, err := strconv.Atoi(v)
if err != nil || n <= 0 {
return cfg, fmt.Errorf("BUSY_RETRY_AFTER: must be a positive integer (seconds)")
}
cfg.BusyRetryAfter = n
}
// LOGLEVEL is the canonical spelling; LOG_LEVEL is kept as an alias.
logLevelValue := getenv("LOGLEVEL")
if logLevelValue == "" {
logLevelValue = getenv("LOG_LEVEL")
}
if logLevelValue != "" {
var level slog.Level
if err := level.UnmarshalText([]byte(logLevelValue)); err != nil {
return cfg, fmt.Errorf("LOGLEVEL: %w", err)
}
cfg.LogLevel = level
}
switch strings.ToLower(getenv("LOG_FORMAT")) {
case "", "text":
case "json":
cfg.LogJSON = true
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
}
+118
View File
@@ -0,0 +1,118 @@
package config
import (
"log/slog"
"strings"
"testing"
"time"
)
func TestDefaults(t *testing.T) {
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)
}
if !cfg.AutoUpdate || cfg.UpdateInterval != 6*time.Hour {
t.Fatalf("update = %v %v", cfg.AutoUpdate, cfg.UpdateInterval)
}
if cfg.LogLevel != slog.LevelWarn {
t.Fatalf("log level = %v", cfg.LogLevel)
}
}
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
LOGLEVEL=debug
SPACED = value with spaces
`
values, err := ParseEnvFile(strings.NewReader(input))
if err != nil {
t.Fatal(err)
}
if values["OLLAMA_URL"] != "http://host:11435" {
t.Fatalf("OLLAMA_URL = %q", values["OLLAMA_URL"])
}
if values["LOGLEVEL"] != "debug" {
t.Fatalf("LOGLEVEL = %q", values["LOGLEVEL"])
}
if values["SPACED"] != "value with spaces" {
t.Fatalf("SPACED = %q", values["SPACED"])
}
}
func TestParseEnvFileMalformed(t *testing.T) {
_, err := ParseEnvFile(strings.NewReader("OK=1\nNOT_A_PAIR\n"))
if err == nil || !strings.Contains(err.Error(), "line 2") {
t.Fatalf("err = %v, want line 2 error", err)
}
}
func TestEnvOverridesFile(t *testing.T) {
file := map[string]string{"OLLAMA_URL": "http://file:1", "UNLOAD_TIMEOUT": "42s"}
env := map[string]string{"OLLAMA_URL": "http://env:2"}
getenv := func(k string) string {
if v := env[k]; v != "" {
return v
}
return file[k]
}
cfg, err := Load(getenv)
if err != nil {
t.Fatal(err)
}
if cfg.OllamaURL != "http://env:2" {
t.Fatalf("OllamaURL = %q, want env value", cfg.OllamaURL)
}
if cfg.UnloadTimeout != 42*time.Second {
t.Fatalf("UnloadTimeout = %v, want file value", cfg.UnloadTimeout)
}
}
func TestLoadErrors(t *testing.T) {
for _, tc := range []struct{ key, value string }{
{"UNLOAD_TIMEOUT", "bogus"},
{"PROMPT_CAPTURE_LIMIT", "-5"},
{"AUTO_UPDATE", "maybe"},
{"LOGLEVEL", "shouty"},
{"LOG_FORMAT", "yaml"},
{"LLM_BUSY_MODE", "bogus"},
{"LLM_BUSY_STATUS", "200"},
{"BUSY_RETRY_AFTER", "0"},
} {
_, err := Load(func(k string) string {
if k == tc.key {
return tc.value
}
if k == "OLLAMA_URL" {
return "http://127.0.0.1:11435"
}
return ""
})
if err == nil {
t.Errorf("%s=%s: expected error", tc.key, tc.value)
}
}
}
+16
View File
@@ -74,6 +74,22 @@ func (l *Lock) AcquireLLM(ctx context.Context) error {
return nil
}
// TryAcquireLLM acquires one in-flight LLM slot without waiting and
// reports whether it succeeded. It fails when an image job is active or
// pending.
func (l *Lock) TryAcquireLLM() bool {
l.mu.Lock()
if l.imageActive || len(l.imageQ) > 0 {
l.mu.Unlock()
return false
}
l.n++
n := l.n
l.mu.Unlock()
l.logTransition("lock transition", "state", StateLLM, "llm_inflight", n)
return true
}
// ReleaseLLM marks one LLM request as finished.
func (l *Lock) ReleaseLLM() {
l.mu.Lock()
+51
View File
@@ -0,0 +1,51 @@
package lock
import (
"context"
"testing"
"time"
)
func TestTryAcquireLLM(t *testing.T) {
lk := New(nil)
if !lk.TryAcquireLLM() {
t.Fatal("TryAcquireLLM on idle lock should succeed")
}
if _, n, _ := lk.Snapshot(); n != 1 {
t.Fatalf("n = %d, want 1", n)
}
// While an image job is pending, TryAcquireLLM must fail.
imageWaiting := make(chan struct{})
go func() {
lk.AcquireImage(context.Background())
close(imageWaiting)
}()
deadline := time.Now().Add(2 * time.Second)
for {
lk.mu.Lock()
queued := len(lk.imageQ)
lk.mu.Unlock()
if queued == 1 {
break
}
if time.Now().After(deadline) {
t.Fatal("image waiter never queued")
}
time.Sleep(time.Millisecond)
}
if lk.TryAcquireLLM() {
t.Fatal("TryAcquireLLM with image pending should fail")
}
lk.ReleaseLLM()
<-imageWaiting
if lk.TryAcquireLLM() {
t.Fatal("TryAcquireLLM with image active should fail")
}
lk.ReleaseImage()
if !lk.TryAcquireLLM() {
t.Fatal("TryAcquireLLM after image release should succeed")
}
lk.ReleaseLLM()
}
+383 -45
View File
@@ -4,15 +4,22 @@
package proxy
import (
"bufio"
"bytes"
"context"
"crypto/tls"
"encoding/json"
"errors"
"fmt"
"io"
"log/slog"
"net"
"net/http"
"net/http/httputil"
"net/url"
"os"
"strconv"
"strings"
"time"
"gpu-turnstile/internal/comfy"
@@ -21,10 +28,10 @@ import (
"gpu-turnstile/internal/ollama"
)
// captureLimit bounds how much of a /prompt response body is buffered while
// looking for prompt_id. The body still passes through to the client
// unchanged regardless of size.
const captureLimit = 64 * 1024
// defaultCaptureLimit bounds how much of a /prompt response body is
// buffered while looking for prompt_id. The body still passes through to
// the client unchanged regardless of size.
const defaultCaptureLimit = 64 * 1024
// Config wires a Server.
type Config struct {
@@ -37,45 +44,266 @@ type Config struct {
Metrics *metrics.Metrics
Log *slog.Logger
// 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
// LogWriter receives the colored per-request lines; nil means stderr.
LogWriter io.Writer
LLMWaitTimeout time.Duration
UnloadTimeout time.Duration
JobTimeout time.Duration
WarmModel string
// LLMBusyMode is "wait" (default) or "reject". In reject mode an LLM
// request that arrives while an image job is active or pending is
// answered immediately with LLMBusyStatus and a Retry-After header
// (BusyRetryAfter seconds) instead of waiting for the lock. In wait
// mode the Retry-After header is sent when LLMWaitTimeout expires.
LLMBusyMode string
LLMBusyStatus int
BusyRetryAfter int
// BackoffInitial and BackoffMax control the exponential retry backoff
// when an upstream refuses a connection: the wait doubles from
// BackoffInitial up to BackoffMax between attempts. Zero selects the
// defaults (1s / 60s).
BackoffInitial time.Duration
BackoffMax time.Duration
// UnloadPollInterval and HistoryPollInterval override the clients'
// /api/ps and /history poll intervals when > 0.
UnloadPollInterval time.Duration
HistoryPollInterval time.Duration
// FreeTimeout and WarmTimeout bound the /free call and the warm-model
// reload; zero selects the defaults.
FreeTimeout time.Duration
WarmTimeout time.Duration
// PromptCaptureLimit overrides defaultCaptureLimit when > 0.
PromptCaptureLimit int64
WarmModel string
}
// Server serves both gpu-turnstile listeners.
type Server struct {
cfg Config
log *slog.Logger
cfg Config
log *slog.Logger
logWriter io.Writer
freeTimeout time.Duration
warmTimeout time.Duration
captureLimit int64
backoffInitial time.Duration
backoffMax time.Duration
busyMode string
busyStatus int
busyRetryAfter int
ollamaProxy *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) {
ollamaURL, err := url.Parse(cfg.OllamaURL)
if err != nil || ollamaURL.Scheme == "" || ollamaURL.Host == "" {
return nil, fmt.Errorf("invalid OLLAMA_URL %q", cfg.OllamaURL)
if cfg.OllamaURL == "" && cfg.ComfyURL == "" {
return nil, fmt.Errorf("at least one of OllamaURL or ComfyURL is required")
}
comfyURL, err := url.Parse(cfg.ComfyURL)
if err != nil || comfyURL.Scheme == "" || comfyURL.Host == "" {
return nil, fmt.Errorf("invalid COMFY_URL %q", cfg.ComfyURL)
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)
}
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()
}
return &Server{
cfg: cfg,
log: log,
ollamaProxy: newReverseProxy(ollamaURL, log.With("upstream", "ollama")),
comfyProxy: newReverseProxy(comfyURL, log.With("upstream", "comfy")),
}, nil
if cfg.UnloadPollInterval > 0 && cfg.Ollama != nil {
cfg.Ollama.PollInterval = cfg.UnloadPollInterval
}
if cfg.HistoryPollInterval > 0 && cfg.Comfy != nil {
cfg.Comfy.PollInterval = cfg.HistoryPollInterval
}
freeTimeout := cfg.FreeTimeout
if freeTimeout <= 0 {
freeTimeout = 30 * time.Second
}
warmTimeout := cfg.WarmTimeout
if warmTimeout <= 0 {
warmTimeout = 2 * time.Minute
}
captureLimit := cfg.PromptCaptureLimit
if captureLimit <= 0 {
captureLimit = defaultCaptureLimit
}
backoffInitial := cfg.BackoffInitial
if backoffInitial <= 0 {
backoffInitial = time.Second
}
backoffMax := cfg.BackoffMax
if backoffMax <= 0 {
backoffMax = time.Minute
}
busyMode := "wait"
if cfg.LLMBusyMode == "reject" {
busyMode = "reject"
}
busyStatus := cfg.LLMBusyStatus
if busyStatus == 0 {
busyStatus = http.StatusServiceUnavailable
}
busyRetryAfter := cfg.BusyRetryAfter
if busyRetryAfter <= 0 {
busyRetryAfter = 30
}
retry := &retryTransport{
base: http.DefaultTransport,
initial: backoffInitial,
max: backoffMax,
log: log,
}
s := &Server{
cfg: cfg,
log: log,
logWriter: cfg.LogWriter,
freeTimeout: freeTimeout,
warmTimeout: warmTimeout,
captureLimit: captureLimit,
backoffInitial: backoffInitial,
backoffMax: backoffMax,
busyMode: busyMode,
busyStatus: busyStatus,
busyRetryAfter: busyRetryAfter,
}
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
}
func newReverseProxy(target *url.URL, log *slog.Logger) *httputil.ReverseProxy {
// retryTransport retries requests whose failure means the upstream never
// saw them — any dial-phase error (connection refused, dial timeout, DNS
// failure), TLS handshake errors — plus 5xx responses when the request
// body can be replayed (GETs and requests with GetBody set). The wait
// doubles from initial up to max between attempts. The loop runs until
// the request succeeds, fails in a non-retryable way, or the client's
// context is cancelled.
type retryTransport struct {
base http.RoundTripper
initial time.Duration
max time.Duration
log *slog.Logger
}
// shouldRetry reports whether a RoundTrip error means the request never
// reached the upstream application and is therefore safe to send again.
func shouldRetry(err error) bool {
// Dial-phase failures: refused, timeout, unreachable, DNS (wrapped).
var opErr *net.OpError
if errors.As(err, &opErr) && opErr.Op == "dial" {
return true
}
var dnsErr *net.DNSError
if errors.As(err, &dnsErr) {
return true
}
// TLS handshake failures: the HTTP request was never written.
var recordErr tls.RecordHeaderError
if errors.As(err, &recordErr) {
return true
}
var certErr *tls.CertificateVerificationError
if errors.As(err, &certErr) {
return true
}
var alertErr tls.AlertError
if errors.As(err, &alertErr) {
return true
}
return false
}
// replayable reports whether the request body can be sent again. Bodies
// streamed from the client (GetBody == nil) cannot, so 5xx responses to
// POSTs are not retried: the upstream may have partially processed them,
// and re-sending could duplicate work (e.g. a second ComfyUI prompt).
func replayable(req *http.Request) bool {
return req.Body == nil || req.Body == http.NoBody || req.GetBody != nil
}
func (t *retryTransport) logAt(level slog.Level, msg string, args ...any) {
if t.log != nil && t.log.Enabled(context.Background(), level) {
t.log.Log(context.Background(), level, msg, args...)
}
}
func (t *retryTransport) RoundTrip(req *http.Request) (*http.Response, error) {
wait := t.initial
attempt := 0
for {
resp, err := t.base.RoundTrip(req)
switch {
case err != nil && !shouldRetry(err):
return nil, err
case err == nil && (resp.StatusCode < 500 || !replayable(req)):
return resp, nil
}
// Retryable failure: a transport error, or a 5xx response.
var reason string
if err != nil {
reason = err.Error()
} else {
reason = resp.Status
io.Copy(io.Discard, resp.Body)
resp.Body.Close()
if req.GetBody != nil {
if body, berr := req.GetBody(); berr == nil {
req.Body = body
}
}
}
attempt++
// One WARN per outage episode; subsequent attempts at INFO.
level := slog.LevelInfo
if attempt == 1 {
level = slog.LevelWarn
}
t.logAt(level, "upstream unavailable; retrying with backoff",
"path", req.URL.Path, "reason", reason, "attempt", attempt, "retry_in", wait)
select {
case <-req.Context().Done():
if err != nil {
return nil, err
}
return nil, req.Context().Err()
case <-time.After(wait):
}
wait *= 2
if wait > t.max {
wait = t.max
}
}
}
func newReverseProxy(target *url.URL, transport http.RoundTripper, log *slog.Logger) *httputil.ReverseProxy {
return &httputil.ReverseProxy{
Transport: transport,
Rewrite: func(pr *httputil.ProxyRequest) {
pr.SetURL(target)
pr.SetXForwarded()
@@ -106,6 +334,94 @@ func (s *Server) writeMetrics(w http.ResponseWriter) {
s.cfg.Metrics.Render(w, string(state), n, pending)
}
// ANSI colors for per-request log lines.
const (
ansiReset = "\x1b[0m"
ansiCyan = "\x1b[36m"
ansiGreen = "\x1b[32m"
ansiYellow = "\x1b[33m"
ansiRed = "\x1b[31m"
)
// statusRecorder remembers the response status while passing everything
// through, including streaming flushes and websocket hijacks.
type statusRecorder struct {
http.ResponseWriter
status int
}
func (r *statusRecorder) WriteHeader(code int) {
r.status = code
r.ResponseWriter.WriteHeader(code)
}
func (r *statusRecorder) Flush() {
if f, ok := r.ResponseWriter.(http.Flusher); ok {
f.Flush()
}
}
func (r *statusRecorder) Hijack() (net.Conn, *bufio.ReadWriter, error) {
h, ok := r.ResponseWriter.(http.Hijacker)
if !ok {
return nil, nil, errors.New("response writer does not support hijacking")
}
return h.Hijack()
}
func (r *statusRecorder) Unwrap() http.ResponseWriter { return r.ResponseWriter }
// reqLine emits one request log line. With color enabled slog cannot be
// used: its text handler escapes the ANSI sequences, so the line is
// written to stderr directly in the same key=value shape. Without color
// it is a plain slog INFO line.
func (s *Server) reqLine(log *slog.Logger, code, line string, attrs ...any) {
if !s.cfg.LogColor {
log.Info(line, attrs...)
return
}
var sb strings.Builder
sb.WriteString("time=" + time.Now().Format("2006-01-02T15:04:05.000Z07:00") + " level=INFO ")
sb.WriteString(code + line + ansiReset)
for i := 0; i+1 < len(attrs); i += 2 {
fmt.Fprintf(&sb, " %v=%v", attrs[i], attrs[i+1])
}
sb.WriteByte('\n')
w := s.logWriter
if w == nil {
w = os.Stderr
}
io.WriteString(w, sb.String())
}
// logRequests logs one line per incoming request and one per completed
// response at INFO level, colored when enabled: cyan "-->" for incoming,
// green/yellow/red "<--" for responses by status class. At log levels
// above INFO it is a pass-through.
func (s *Server) logRequests(listener string, next http.Handler) http.Handler {
log := s.log.With("listener", listener)
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if !log.Enabled(r.Context(), slog.LevelInfo) {
next.ServeHTTP(w, r)
return
}
start := time.Now()
s.reqLine(log, ansiCyan, "--> "+r.Method+" "+r.URL.RequestURI(),
"listener", listener, "remote", r.RemoteAddr)
rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK}
next.ServeHTTP(rec, r)
code := ansiGreen
switch {
case rec.status >= 500:
code = ansiRed
case rec.status >= 400:
code = ansiYellow
}
s.reqLine(log, code, "<-- "+strconv.Itoa(rec.status)+" "+r.Method+" "+r.URL.RequestURI(),
"listener", listener, "ms", time.Since(start).Milliseconds())
})
}
// llmPaths are the Ollama endpoints that load models into VRAM and therefore
// take the LLM lock. Everything else passes through unlocked.
var llmPaths = map[string]bool{
@@ -124,7 +440,7 @@ func isLLMRequest(r *http.Request) bool {
// OllamaHandler serves the Ollama-facing listener.
func (s *Server) OllamaHandler() http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
return s.logRequests("ollama", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/healthz":
s.writeHealthz(w)
@@ -139,42 +455,62 @@ func (s *Server) OllamaHandler() http.Handler {
}
start := time.Now()
if s.busyMode == "reject" {
if !s.cfg.Lock.TryAcquireLLM() {
s.cfg.Metrics.ObserveLockWait("llm", time.Since(start).Seconds())
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)
return
}
s.cfg.Metrics.ObserveLockWait("llm", time.Since(start).Seconds())
defer s.cfg.Lock.ReleaseLLM()
s.ollamaProxy.ServeHTTP(w, r)
return
}
wctx, cancel := context.WithTimeout(r.Context(), s.cfg.LLMWaitTimeout)
err := s.cfg.Lock.AcquireLLM(wctx)
cancel()
s.cfg.Metrics.ObserveLockWait("llm", time.Since(start).Seconds())
if err != nil {
if errors.Is(err, context.DeadlineExceeded) && r.Context().Err() == nil {
w.Header().Set("Retry-After", strconv.Itoa(s.busyRetryAfter))
http.Error(w, "GPU busy: timed out waiting for the lock", http.StatusServiceUnavailable)
}
return
}
defer s.cfg.Lock.ReleaseLLM()
s.ollamaProxy.ServeHTTP(w, r)
})
}))
}
// ComfyHandler serves the ComfyUI-facing listener.
func (s *Server) ComfyHandler() http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/healthz" {
return s.logRequests("comfy", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
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)
return
}
s.comfyProxy.ServeHTTP(w, r)
})
}))
}
// captureWriter passes the response through unchanged while recording the
// status code and the first captureLimit bytes of the body.
// status code and the first limit bytes of the body.
type captureWriter struct {
http.ResponseWriter
status int
buf bytes.Buffer
limit int64
}
func (w *captureWriter) WriteHeader(code int) {
@@ -183,7 +519,7 @@ func (w *captureWriter) WriteHeader(code int) {
}
func (w *captureWriter) Write(p []byte) (int, error) {
if w.buf.Len() < captureLimit {
if int64(w.buf.Len()) < w.limit {
w.buf.Write(p)
}
return w.ResponseWriter.Write(p)
@@ -213,22 +549,24 @@ func (s *Server) handlePrompt(w http.ResponseWriter, r *http.Request) {
s.cfg.Metrics.ObserveLockWait("image", time.Since(start).Seconds())
log.Info("image lock acquired")
uctx, ucancel := context.WithTimeout(r.Context(), s.cfg.UnloadTimeout)
elapsed, uerr := s.cfg.Ollama.UnloadAll(uctx)
ucancel()
s.cfg.Metrics.ObserveUnload(elapsed.Seconds())
switch {
case r.Context().Err() != nil:
s.cfg.Lock.ReleaseImage()
return
case uerr != nil:
// Degrade, don't fail the user's request on a misbehaving neighbour.
log.Warn("ollama unload incomplete; continuing", "err", uerr)
default:
log.Info("ollama models unloaded", "seconds", elapsed.Seconds())
if s.cfg.Ollama != nil {
uctx, ucancel := context.WithTimeout(r.Context(), s.cfg.UnloadTimeout)
elapsed, uerr := s.cfg.Ollama.UnloadAll(uctx)
ucancel()
s.cfg.Metrics.ObserveUnload(elapsed.Seconds())
switch {
case r.Context().Err() != nil:
s.cfg.Lock.ReleaseImage()
return
case uerr != nil:
// Degrade, don't fail the user's request on a misbehaving neighbour.
log.Warn("ollama unload incomplete; continuing", "err", uerr)
default:
log.Info("ollama models unloaded", "seconds", elapsed.Seconds())
}
}
cw := &captureWriter{ResponseWriter: w, status: http.StatusOK}
cw := &captureWriter{ResponseWriter: w, status: http.StatusOK, limit: s.captureLimit}
s.comfyProxy.ServeHTTP(cw, r)
var accepted struct {
@@ -262,7 +600,7 @@ func (s *Server) finishImageJob(promptID string) {
log.Info("image job completed")
}
freeCtx, freeCancel := context.WithTimeout(context.Background(), 30*time.Second)
freeCtx, freeCancel := context.WithTimeout(context.Background(), s.freeTimeout)
if err := s.cfg.Comfy.Free(freeCtx); err != nil {
log.Warn("failed to free ComfyUI models", "err", err)
}
@@ -271,9 +609,9 @@ 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(), 2*time.Minute)
wctx, wcancel := context.WithTimeout(context.Background(), s.warmTimeout)
if err := s.cfg.Ollama.Warm(wctx, s.cfg.WarmModel); err != nil {
log.Warn("warm model reload failed", "model", s.cfg.WarmModel, "err", err)
} else {
+185
View File
@@ -2,6 +2,7 @@ package proxy
import (
"bufio"
"context"
"fmt"
"io"
"net/http"
@@ -345,3 +346,187 @@ func TestPassThroughNoLock(t *testing.T) {
t.Fatalf("pass-through = %d %s", resp.StatusCode, body)
}
}
func TestLLMBusyReject(t *testing.T) {
f := newFakes(t)
lk := lock.New(nil)
if err := lk.AcquireImage(context.Background()); err != nil {
t.Fatal(err)
}
ollamaClient, err := ollama.New(f.ollama.URL, nil)
if err != nil {
t.Fatal(err)
}
comfyClient, err := comfy.New(f.comfy.URL, nil)
if err != nil {
t.Fatal(err)
}
srv, err := New(Config{
OllamaURL: f.ollama.URL,
ComfyURL: f.comfy.URL,
Lock: lk,
Ollama: ollamaClient,
Comfy: comfyClient,
Metrics: metrics.New(),
LLMWaitTimeout: 2 * time.Second,
LLMBusyMode: "reject",
BusyRetryAfter: 17,
})
if err != nil {
t.Fatal(err)
}
front := httptest.NewServer(srv.OllamaHandler())
defer front.Close()
start := time.Now()
resp, err := http.Post(front.URL+"/api/chat", "application/json", strings.NewReader(`{}`))
if err != nil {
t.Fatal(err)
}
body, _ := io.ReadAll(resp.Body)
resp.Body.Close()
if resp.StatusCode != http.StatusServiceUnavailable {
t.Fatalf("busy chat status = %d %s", resp.StatusCode, body)
}
if got := resp.Header.Get("Retry-After"); got != "17" {
t.Fatalf("Retry-After = %q", got)
}
if elapsed := time.Since(start); elapsed > time.Second {
t.Fatalf("reject was not immediate: %v", elapsed)
}
// After the image lock is released the next LLM request goes through.
lk.ReleaseImage()
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 after release status = %d", resp.StatusCode)
}
}
func TestLLMBusyWaitTimeoutRetryAfter(t *testing.T) {
f := newFakes(t)
lk := lock.New(nil)
if err := lk.AcquireImage(context.Background()); err != nil {
t.Fatal(err)
}
defer lk.ReleaseImage()
ollamaClient, err := ollama.New(f.ollama.URL, nil)
if err != nil {
t.Fatal(err)
}
comfyClient, err := comfy.New(f.comfy.URL, nil)
if err != nil {
t.Fatal(err)
}
srv, err := New(Config{
OllamaURL: f.ollama.URL,
ComfyURL: f.comfy.URL,
Lock: lk,
Ollama: ollamaClient,
Comfy: comfyClient,
Metrics: metrics.New(),
LLMWaitTimeout: 50 * time.Millisecond,
})
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 != http.StatusServiceUnavailable {
t.Fatalf("timed-out chat status = %d", resp.StatusCode)
}
if got := resp.Header.Get("Retry-After"); got != "30" {
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")
}
}
+207
View File
@@ -0,0 +1,207 @@
package proxy
import (
"context"
"crypto/tls"
"io"
"net"
"net/http"
"strings"
"syscall"
"testing"
"time"
)
// stubTransport fails with ECONNREFUSED for the first fails requests, then
// returns a 200 response.
type stubTransport struct {
fails int
calls int
}
func (s *stubTransport) RoundTrip(req *http.Request) (*http.Response, error) {
s.calls++
if s.calls <= s.fails {
return nil, &net.OpError{Op: "dial", Net: "tcp", Err: syscall.ECONNREFUSED}
}
return &http.Response{
StatusCode: 200,
Body: io.NopCloser(strings.NewReader("ok")),
Header: make(http.Header),
}, nil
}
func TestRetryTransportBackoff(t *testing.T) {
st := &stubTransport{fails: 3}
rt := &retryTransport{
base: st,
initial: 10 * time.Millisecond,
max: 25 * time.Millisecond,
}
start := time.Now()
req, _ := http.NewRequest(http.MethodGet, "http://upstream/api/version", nil)
resp, err := rt.RoundTrip(req)
if err != nil {
t.Fatal(err)
}
resp.Body.Close()
if st.calls != 4 {
t.Fatalf("calls = %d, want 4", st.calls)
}
// Waits: 10ms + 20ms + 25ms (capped) = 55ms minimum.
elapsed := time.Since(start)
if elapsed < 50*time.Millisecond {
t.Fatalf("elapsed = %v, want >= ~55ms of backoff", elapsed)
}
if elapsed > 5*time.Second {
t.Fatalf("elapsed = %v, suspiciously long", elapsed)
}
}
func TestRetryTransportNonRefusedErrorNotRetried(t *testing.T) {
rt := &retryTransport{
base: &stubTransport{fails: 0},
initial: time.Millisecond,
max: time.Millisecond,
}
req, _ := http.NewRequest(http.MethodGet, "http://upstream/", nil)
resp, err := rt.RoundTrip(req)
if err != nil || resp.StatusCode != 200 {
t.Fatalf("resp=%v err=%v", resp, err)
}
}
type alwaysRefused struct{ calls int }
func (a *alwaysRefused) RoundTrip(*http.Request) (*http.Response, error) {
a.calls++
return nil, &net.OpError{Op: "dial", Net: "tcp", Err: syscall.ECONNREFUSED}
}
func TestRetryTransportContextCancel(t *testing.T) {
ar := &alwaysRefused{}
rt := &retryTransport{base: ar, initial: time.Second, max: time.Second}
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, http.MethodGet, "http://upstream/", nil)
start := time.Now()
_, err := rt.RoundTrip(req)
if err == nil {
t.Fatal("expected error after context cancellation")
}
if time.Since(start) > 2*time.Second {
t.Fatal("retry loop did not stop on context cancellation")
}
}
func TestRetryViaReverseProxy(t *testing.T) {
// Point the proxy at a port nothing listens on; the request should be
// retried (not instantly 502) until the client context ends.
srv, err := New(Config{
OllamaURL: "http://127.0.0.1:1",
ComfyURL: "http://127.0.0.1:1",
Lock: nil,
Metrics: nil,
BackoffInitial: 10 * time.Millisecond,
BackoffMax: 20 * time.Millisecond,
})
if err != nil {
t.Fatal(err)
}
_ = srv // construction must not panic with minimal config
rt := &retryTransport{base: http.DefaultTransport, initial: 10 * time.Millisecond, max: 20 * time.Millisecond}
ctx, cancel := context.WithTimeout(context.Background(), 200*time.Millisecond)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, http.MethodGet, "http://127.0.0.1:1/", nil)
_, err = rt.RoundTrip(req)
if err == nil || !shouldRetry(err) {
t.Fatalf("err = %v, want connection refused", err)
}
}
// flakyStatus returns 500 for the first fails requests, then 200.
type flakyStatus struct {
fails int
calls int
}
func (s *flakyStatus) RoundTrip(req *http.Request) (*http.Response, error) {
s.calls++
code := 200
if s.calls <= s.fails {
code = 500
}
return &http.Response{
StatusCode: code,
Status: http.StatusText(code),
Body: io.NopCloser(strings.NewReader("")),
Header: make(http.Header),
}, nil
}
func TestRetryTransport5xxGet(t *testing.T) {
fs := &flakyStatus{fails: 2}
rt := &retryTransport{base: fs, initial: time.Millisecond, max: 2 * time.Millisecond}
req, _ := http.NewRequest(http.MethodGet, "http://upstream/api/version", nil)
resp, err := rt.RoundTrip(req)
if err != nil {
t.Fatal(err)
}
resp.Body.Close()
if resp.StatusCode != 200 {
t.Fatalf("status = %d, want 200", resp.StatusCode)
}
if fs.calls != 3 {
t.Fatalf("calls = %d, want 3", fs.calls)
}
}
func TestRetryTransport5xxPostNotRetried(t *testing.T) {
fs := &flakyStatus{fails: 10}
rt := &retryTransport{base: fs, initial: time.Millisecond, max: time.Millisecond}
// A streamed body (no GetBody) must not be replayed after a 500.
req, _ := http.NewRequest(http.MethodPost, "http://upstream/prompt", io.NopCloser(strings.NewReader("{}")))
resp, err := rt.RoundTrip(req)
if err != nil {
t.Fatal(err)
}
resp.Body.Close()
if resp.StatusCode != 500 {
t.Fatalf("status = %d, want 500", resp.StatusCode)
}
if fs.calls != 1 {
t.Fatalf("calls = %d, want 1 (no retry for streamed POST)", fs.calls)
}
}
func TestShouldRetryClassification(t *testing.T) {
cases := []struct {
name string
err error
want bool
}{
{"dial refused", &net.OpError{Op: "dial", Net: "tcp", Err: syscall.ECONNREFUSED}, true},
{"dial timeout", &net.OpError{Op: "dial", Net: "tcp", Err: timeoutErr{}}, true},
{"dns", &net.DNSError{Err: "no such host", IsNotFound: true}, true},
{"tls record", tls.RecordHeaderError{Msg: "bad"}, true},
{"tls alert", tls.AlertError(42), true},
{"read error mid-request", &net.OpError{Op: "read", Net: "tcp", Err: syscall.ECONNRESET}, false},
{"plain error", io.EOF, false},
}
for _, c := range cases {
if got := shouldRetry(c.err); got != c.want {
t.Errorf("%s: shouldRetry = %v, want %v", c.name, got, c.want)
}
}
}
type timeoutErr struct{}
func (timeoutErr) Error() string { return "i/o timeout" }
func (timeoutErr) Timeout() bool { return true }
func (timeoutErr) Temporary() bool { return true }
+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)
}
}
}
+35
View File
@@ -0,0 +1,35 @@
//go:build !windows && !linux
// 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 (
"context"
"errors"
"os/signal"
"syscall"
)
// Name matches the Windows service name.
const Name = "gpu-turnstile"
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 }
// Run executes run with SIGINT/SIGTERM cancellation, mirroring the
// interactive behavior.
func Run(run func(ctx context.Context) error) error {
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
return run(ctx)
}
// Install is unsupported on non-Windows platforms.
func Install(string) error { return errUnsupported }
// Remove is unsupported on non-Windows platforms.
func Remove() error { return errUnsupported }
+120
View File
@@ -0,0 +1,120 @@
//go:build windows
// Package service integrates gpu-turnstile with the Windows Service
// Control Manager: running as a service with graceful stop, plus
// install/remove helpers.
package service
import (
"context"
"fmt"
"os"
"time"
"golang.org/x/sys/windows/svc"
"golang.org/x/sys/windows/svc/mgr"
)
// Name is the Windows service name.
const Name = "gpu-turnstile"
// IsService reports whether the process is running as a Windows service.
func IsService() bool {
isSvc, err := svc.IsWindowsService()
return err == nil && isSvc
}
// Run executes run as a Windows service. SCM Stop and Shutdown cancel the
// context passed to run, triggering the same graceful shutdown as SIGTERM
// in interactive mode.
func Run(run func(ctx context.Context) error) error {
return svc.Run(Name, &handler{run: run})
}
type handler struct {
run func(ctx context.Context) error
}
func (h *handler) Execute(_ []string, requests <-chan svc.ChangeRequest, status chan<- svc.Status) (bool, uint32) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
status <- svc.Status{State: svc.StartPending}
errCh := make(chan error, 1)
go func() { errCh <- h.run(ctx) }()
status <- svc.Status{State: svc.Running, Accepts: svc.AcceptStop | svc.AcceptShutdown}
for {
select {
case err := <-errCh:
status <- svc.Status{State: svc.Stopped}
if err != nil {
return true, 1
}
return false, 0
case c := <-requests:
switch c.Cmd {
case svc.Interrogate:
status <- c.CurrentStatus
case svc.Stop, svc.Shutdown:
status <- svc.Status{State: svc.StopPending}
cancel()
}
}
}
}
// Install registers gpu-turnstile as an auto-start Windows service whose
// binPath loads the given config file. Recovery actions restart the
// service after 5s on failure — this is also what brings up a staged
// update after the updater exits with a non-zero code.
func Install(configPath string) error {
exe, err := os.Executable()
if err != nil {
return err
}
m, err := mgr.Connect()
if err != nil {
return fmt.Errorf("connect to service manager (run as administrator): %w", err)
}
defer m.Disconnect()
binPath := fmt.Sprintf(`"%s" -config "%s"`, exe, configPath)
s, err := m.CreateService(Name, binPath, mgr.Config{
StartType: mgr.StartAutomatic,
DisplayName: "gpu-turnstile",
Description: "GPU arbitration proxy for Ollama and ComfyUI",
})
if err != nil {
return fmt.Errorf("create service: %w", err)
}
defer s.Close()
restart := mgr.RecoveryAction{Type: mgr.ServiceRestart, Delay: 5 * time.Second}
if err := s.SetRecoveryActions([]mgr.RecoveryAction{restart, restart, restart}, 24*60*60); err != nil {
return fmt.Errorf("set recovery actions: %w", err)
}
if err := s.SetRecoveryActionsOnNonCrashFailures(true); err != nil {
return fmt.Errorf("set failure actions flag: %w", err)
}
return nil
}
// Remove stops (if running) and unregisters the service.
func Remove() error {
m, err := mgr.Connect()
if err != nil {
return fmt.Errorf("connect to service manager (run as administrator): %w", err)
}
defer m.Disconnect()
s, err := m.OpenService(Name)
if err != nil {
return fmt.Errorf("open service: %w", err)
}
defer s.Close()
s.Control(svc.Stop) // ignore error: may already be stopped
if err := s.Delete(); err != nil {
return fmt.Errorf("delete service: %w", err)
}
return nil
}
+16
View File
@@ -0,0 +1,16 @@
package update
// publicKeyPEM is the PEM-encoded Ed25519 public key that matches the
// RELEASE_SIGNING_KEY secret used by CI to sign release binaries. Generate a
// keypair once with:
//
// openssl genpkey -algorithm ed25519 -out private.pem
// openssl pkey -in private.pem -pubout -out public.pem
//
// Paste the contents of public.pem here and commit; store private.pem as
// the RELEASE_SIGNING_KEY repository secret. When empty, the updater refuses
// to update (e.g. development builds).
var publicKeyPEM = `-----BEGIN PUBLIC KEY-----
MCowBQYDK2VwAyEAqTAJ0CCeAQI7MhFlgc5xNmF/CfvLVUAY3ZAoeAS0tT8=
-----END PUBLIC KEY-----
`
+254
View File
@@ -0,0 +1,254 @@
// Package update implements gpu-turnstile's self-updater: it polls the
// Gitea releases API, downloads the Windows binary of newer releases, and
// verifies its Ed25519 signature (produced by CI with OpenSSL) before
// swapping it in next to the running executable.
package update
import (
"context"
"crypto/ed25519"
"crypto/sha256"
"crypto/x509"
"encoding/hex"
"encoding/json"
"encoding/pem"
"fmt"
"io"
"log/slog"
"net/http"
"net/url"
"os"
"strconv"
"strings"
"time"
)
// maxAssetSize bounds release asset downloads.
const maxAssetSize = 512 << 20
// Updater checks one Gitea repository for newer releases.
type Updater struct {
Repo string // e.g. https://git.rambossek.at/PUBLIC/gpu-turnstile
Asset string // e.g. gpu-turnstile.exe
Version string // current version, e.g. v0.1.2 ("dev" disables updates)
Log *slog.Logger
Client *http.Client
}
type release struct {
TagName string `json:"tag_name"`
Assets []struct {
Name string `json:"name"`
BrowserDownloadURL string `json:"browser_download_url"`
} `json:"assets"`
}
func (u *Updater) logger() *slog.Logger {
if u.Log != nil {
return u.Log
}
return slog.Default()
}
func (u *Updater) httpClient() *http.Client {
if u.Client != nil {
return u.Client
}
return &http.Client{Timeout: 5 * time.Minute}
}
// apiURL derives <scheme>://<host>/api/v1/repos/<owner>/<name> from Repo.
func (u *Updater) apiURL() (string, error) {
repoURL, err := url.Parse(u.Repo)
if err != nil || repoURL.Scheme == "" || repoURL.Host == "" {
return "", fmt.Errorf("invalid UPDATE_REPO %q", u.Repo)
}
ownerName := strings.Trim(repoURL.Path, "/")
if len(strings.Split(ownerName, "/")) != 2 {
return "", fmt.Errorf("UPDATE_REPO %q: expected path /<owner>/<name>", u.Repo)
}
return fmt.Sprintf("%s://%s/api/v1/repos/%s", repoURL.Scheme, repoURL.Host, ownerName), nil
}
// newerVersion reports whether latest is a higher vX.Y.Z version than
// current. Both may carry a leading "v".
func newerVersion(current, latest string) (bool, error) {
parse := func(s string) ([3]int, error) {
var v [3]int
parts := strings.Split(strings.TrimPrefix(s, "v"), ".")
if len(parts) != 3 {
return v, fmt.Errorf("not a vX.Y.Z version: %q", s)
}
for i, p := range parts {
n, err := strconv.Atoi(p)
if err != nil {
return v, fmt.Errorf("not a vX.Y.Z version: %q", s)
}
v[i] = n
}
return v, nil
}
cur, err := parse(current)
if err != nil {
return false, err
}
lat, err := parse(latest)
if err != nil {
return false, err
}
for i := 0; i < 3; i++ {
if lat[i] != cur[i] {
return lat[i] > cur[i], nil
}
}
return false, nil
}
func (u *Updater) get(ctx context.Context, url string) ([]byte, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
if err != nil {
return nil, err
}
resp, err := u.httpClient().Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
io.Copy(io.Discard, resp.Body)
return nil, fmt.Errorf("GET %s: %s", url, resp.Status)
}
return io.ReadAll(io.LimitReader(resp.Body, maxAssetSize))
}
func verifySignature(pubKeyPEM string, data, sig []byte) error {
block, _ := pem.Decode([]byte(pubKeyPEM))
if block == nil {
return fmt.Errorf("invalid embedded public key PEM")
}
key, err := x509.ParsePKIXPublicKey(block.Bytes)
if err != nil {
return fmt.Errorf("parse public key: %w", err)
}
pub, ok := key.(ed25519.PublicKey)
if !ok {
return fmt.Errorf("public key is not Ed25519")
}
if !ed25519.Verify(pub, data, sig) {
return fmt.Errorf("signature verification failed")
}
return nil
}
// stage swaps data into place at exePath: the running executable is
// renamed aside (allowed on Windows) and the new file takes its name.
func stage(exePath string, data []byte) error {
newPath := exePath + ".new"
oldPath := exePath + ".old"
os.Remove(oldPath) // leftover from a previous update
if err := os.WriteFile(newPath, data, 0o755); err != nil {
return err
}
if err := os.Rename(exePath, oldPath); err != nil {
os.Remove(newPath)
return err
}
if err := os.Rename(newPath, exePath); err != nil {
os.Rename(oldPath, exePath) // roll back
return err
}
return nil
}
// CleanupOld removes the .old binary left behind by a staged update.
// Call once at startup.
func CleanupOld(exePath string) {
os.Remove(exePath + ".old")
os.Remove(exePath + ".new")
}
// Check performs a single update check. staged is true when a newer,
// signature-verified binary has been swapped into place at exePath; the
// caller should then restart the process. A nil error with staged=false
// means "no action" (up to date, disabled, or dev build); a non-nil error
// means the check failed and the running binary is untouched.
func (u *Updater) Check(ctx context.Context, exePath string) (staged bool, err error) {
log := u.logger()
if u.Version == "" || u.Version == "dev" {
log.Debug("auto-update: dev build, skipping")
return false, nil
}
if publicKeyPEM == "" {
log.Debug("auto-update: no public key embedded, skipping")
return false, nil
}
api, err := u.apiURL()
if err != nil {
return false, err
}
body, err := u.get(ctx, api+"/releases/latest")
if err != nil {
return false, fmt.Errorf("fetch latest release: %w", err)
}
var rel release
if err := json.Unmarshal(body, &rel); err != nil {
return false, fmt.Errorf("parse release: %w", err)
}
newer, err := newerVersion(u.Version, rel.TagName)
if err != nil {
return false, err
}
if !newer {
log.Debug("auto-update: up to date", "version", u.Version, "latest", rel.TagName)
return false, nil
}
urls := make(map[string]string, len(rel.Assets))
for _, a := range rel.Assets {
urls[a.Name] = a.BrowserDownloadURL
}
assetURL, ok := urls[u.Asset]
if !ok {
return false, fmt.Errorf("release %s has no asset %q", rel.TagName, u.Asset)
}
sigURL, ok := urls[u.Asset+".sig"]
if !ok {
return false, fmt.Errorf("release %s has no signature asset %q", rel.TagName, u.Asset+".sig")
}
data, err := u.get(ctx, assetURL)
if err != nil {
return false, fmt.Errorf("download %s: %w", u.Asset, err)
}
sig, err := u.get(ctx, sigURL)
if err != nil {
return false, fmt.Errorf("download signature: %w", err)
}
if sumURL, ok := urls[u.Asset+".sha256"]; ok {
sumText, err := u.get(ctx, sumURL)
if err != nil {
return false, fmt.Errorf("download checksum: %w", err)
}
want := strings.Fields(string(sumText))[0]
got := hex.EncodeToString(sha256Bytes(data))
if !strings.EqualFold(want, got) {
return false, fmt.Errorf("sha256 mismatch: got %s, want %s", got, want)
}
}
if err := verifySignature(publicKeyPEM, data, sig); err != nil {
return false, err
}
if err := stage(exePath, data); err != nil {
return false, fmt.Errorf("stage update: %w", err)
}
log.Info("auto-update: new version staged", "from", u.Version, "to", rel.TagName)
return true, nil
}
func sha256Bytes(data []byte) []byte {
sum := sha256.Sum256(data)
return sum[:]
}
+186
View File
@@ -0,0 +1,186 @@
package update
import (
"context"
"crypto/ed25519"
"crypto/rand"
"crypto/x509"
"encoding/json"
"encoding/pem"
"fmt"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"testing"
)
// fakeGitea serves a Gitea-flavored releases API with one release.
type fakeGitea struct {
srv *httptest.Server
pubPEM string
asset []byte
tag string
tamper bool
noSig bool
}
func newFakeGitea(t *testing.T, tag string, assetContent []byte) *fakeGitea {
t.Helper()
pub, priv, err := ed25519.GenerateKey(rand.Reader)
if err != nil {
t.Fatal(err)
}
der, err := x509.MarshalPKIXPublicKey(pub)
if err != nil {
t.Fatal(err)
}
f := &fakeGitea{
pubPEM: string(pem.EncodeToMemory(&pem.Block{Type: "PUBLIC KEY", Bytes: der})),
asset: assetContent,
tag: tag,
}
sign := func() []byte { return ed25519.Sign(priv, f.asset) }
mux := http.NewServeMux()
mux.HandleFunc("/api/v1/repos/o/r/releases/latest", func(w http.ResponseWriter, r *http.Request) {
assets := []map[string]string{
{"name": "gpu-turnstile.exe", "browser_download_url": f.srv.URL + "/dl/exe"},
{"name": "gpu-turnstile.exe.sha256", "browser_download_url": f.srv.URL + "/dl/sha"},
}
if !f.noSig {
assets = append(assets, map[string]string{"name": "gpu-turnstile.exe.sig", "browser_download_url": f.srv.URL + "/dl/sig"})
}
json.NewEncoder(w).Encode(map[string]any{"tag_name": f.tag, "assets": assets})
})
mux.HandleFunc("/dl/exe", func(w http.ResponseWriter, r *http.Request) { w.Write(f.asset) })
mux.HandleFunc("/dl/sig", func(w http.ResponseWriter, r *http.Request) {
sig := sign()
if f.tamper {
sig[0] ^= 0xff
}
w.Write(sig)
})
mux.HandleFunc("/dl/sha", func(w http.ResponseWriter, r *http.Request) {
fmt.Fprintf(w, "%x gpu-turnstile.exe\n", sha256Bytes(f.asset))
})
f.srv = httptest.NewServer(mux)
t.Cleanup(f.srv.Close)
return f
}
func (f *fakeGitea) updater(version string) *Updater {
return &Updater{Repo: f.srv.URL + "/o/r", Asset: "gpu-turnstile.exe", Version: version}
}
func fakeExe(t *testing.T) string {
t.Helper()
exe := filepath.Join(t.TempDir(), "gpu-turnstile.exe")
if err := os.WriteFile(exe, []byte("old-binary"), 0o755); err != nil {
t.Fatal(err)
}
return exe
}
func withPublicKey(t *testing.T, pem string) {
t.Helper()
old := publicKeyPEM
publicKeyPEM = pem
t.Cleanup(func() { publicKeyPEM = old })
}
func TestCheckStagesUpdate(t *testing.T) {
f := newFakeGitea(t, "v9.9.9", []byte("new-binary"))
withPublicKey(t, f.pubPEM)
exe := fakeExe(t)
staged, err := f.updater("v0.1.2").Check(context.Background(), exe)
if err != nil {
t.Fatal(err)
}
if !staged {
t.Fatal("expected staged update")
}
content, _ := os.ReadFile(exe)
if string(content) != "new-binary" {
t.Fatalf("exe content = %q", content)
}
old, _ := os.ReadFile(exe + ".old")
if string(old) != "old-binary" {
t.Fatalf(".old content = %q", old)
}
}
func TestCheckRejectsTamperedSignature(t *testing.T) {
f := newFakeGitea(t, "v9.9.9", []byte("new-binary"))
f.tamper = true
withPublicKey(t, f.pubPEM)
exe := fakeExe(t)
staged, err := f.updater("v0.1.2").Check(context.Background(), exe)
if err == nil {
t.Fatal("expected signature error")
}
if staged {
t.Fatal("must not stage on bad signature")
}
content, _ := os.ReadFile(exe)
if string(content) != "old-binary" {
t.Fatal("exe was modified despite bad signature")
}
}
func TestCheckSkipsOlderOrEqual(t *testing.T) {
for _, tag := range []string{"v0.1.2", "v0.1.1", "v0.0.9"} {
f := newFakeGitea(t, tag, []byte("new-binary"))
withPublicKey(t, f.pubPEM)
staged, err := f.updater("v0.1.2").Check(context.Background(), fakeExe(t))
if err != nil {
t.Fatal(err)
}
if staged {
t.Fatalf("tag %s must not stage over v0.1.2", tag)
}
}
}
func TestCheckSkipsWithoutPublicKey(t *testing.T) {
f := newFakeGitea(t, "v9.9.9", []byte("new-binary"))
withPublicKey(t, "")
staged, err := f.updater("v0.1.2").Check(context.Background(), fakeExe(t))
if err != nil || staged {
t.Fatalf("staged=%v err=%v, want no action without key", staged, err)
}
}
func TestCheckSkipsDevBuild(t *testing.T) {
f := newFakeGitea(t, "v9.9.9", []byte("new-binary"))
withPublicKey(t, f.pubPEM)
staged, err := f.updater("dev").Check(context.Background(), fakeExe(t))
if err != nil || staged {
t.Fatalf("staged=%v err=%v, want no action for dev build", staged, err)
}
}
func TestNewerVersion(t *testing.T) {
cases := []struct {
cur, lat string
want bool
}{
{"v0.1.2", "v0.1.3", true},
{"0.1.2", "0.2.0", true},
{"v0.1.2", "v1.0.0", true},
{"v0.1.2", "v0.1.2", false},
{"v1.2.3", "v1.2.10", true},
{"v1.2.10", "v1.2.3", false},
}
for _, c := range cases {
got, err := newerVersion(c.cur, c.lat)
if err != nil || got != c.want {
t.Errorf("newerVersion(%s, %s) = %v, %v; want %v", c.cur, c.lat, got, err, c.want)
}
}
if _, err := newerVersion("v0.1", "v0.1.2"); err == nil {
t.Error("expected error for malformed version")
}
}