15 Commits
Author SHA1 Message Date
mram da02457fd5 Pin compose example to v0.1.6
ci / test (push) Successful in 14s
ci / docker (push) Successful in 1m7s
ci / release (push) Successful in 14s
2026-09-21 08:22:42 +02:00
mram 6a6359a630 Add --force-update: immediate signed-update check, stage + service restart 2026-09-21 08:22:11 +02:00
mram e4e4348e8d Start the Windows service right after install (parity with enable --now) 2026-09-21 08:12:13 +02:00
mram b421bb7bfb Add --version and print the version atop every help and error screen 2026-09-21 08:10:46 +02:00
mram 909918657f Fix install/remove success message (installd -> installed) 2026-09-21 08:08:24 +02:00
mram a0435c858e Bare run in a terminal prints the help screen; add -h/--help
A zero-argument invocation now shows the usage text when stdout is a
console (double-clicked exe, interactive shell). Without a terminal —
Docker entrypoint, services, pipes — a bare invocation still starts the
proxy, so the container image and service behavior are unchanged.
2026-09-21 08:06:35 +02:00
mram fbab0bba33 Relaunch through UAC when (un)installing the service unprivileged
--install-service/--remove-service on Windows no longer fail with
'Access is denied' from a normal shell: the process re-runs itself via
ShellExecuteEx 'runas', waits for the elevated child and mirrors its
exit code. The child gets --elevated-child and pauses for a keypress so
its console output stays readable. Declining the prompt reports
'UAC prompt declined'.
2026-09-21 08:00:24 +02:00
mram 802a64280f Self-install into canonical layout on both platforms, --no-copy to opt out
Windows: --install-service creates %ProgramFiles%\gpu-turnstile and
%ProgramData%\gpu-turnstile, copies the exe and (if absent) the env
file in, and registers the copy. Linux: binary goes to
/var/lib/gpu-turnstile (not /usr/local/sbin: replacing a running binary
needs directory write, which must not be granted on a shared system dir
to a sandboxed service). --no-copy registers the current location
as-is on both platforms.
2026-09-21 07:52:33 +02:00
mram a88955e35c Sandbox the systemd unit: DynamicUser, read-only FS, no capabilities
The Linux install now mirrors the Windows virtual-account hardening: the
unit runs with DynamicUser=yes (transient per-service UID, no login),
ProtectSystem=strict with only StateDirectory writable (the install dir,
so self-update can rewrite the binary), NoNewPrivileges, empty
capability sets, restricted address families and a @system-service
syscall filter. Install copies the binary to /var/lib/gpu-turnstile and
the config to /etc/gpu-turnstile.env; Remove cleans up the unit and
binary but keeps the config.
2026-09-21 00:11:25 +02:00
mram 97624470eb Install the Windows service as the NT SERVICE virtual account only
--install-service now registers the service under
NT SERVICE\gpu-turnstile (low-privilege, per-service, no password) and
grants it modify access to the install dir (for self-updates) and the
LOG_FILE dir, plus read access to an external config file. Grants run
after CreateService because the virtual account's SID does not exist
before registration; a failed grant rolls back the registration.
2026-09-20 23:53:54 +02:00
mram 75f16a0229 Copy go.sum into the Docker build; pin compose example to v0.1.5
ci / test (push) Successful in 13s
ci / docker (push) Successful in 1m7s
ci / release (push) Successful in 14s
The image build broke with the first Linux-imported dependency
(go-systemd): the Dockerfile copied go.mod only, and the missing go.sum
never mattered while golang.org/x/sys was Windows-only.
2026-09-20 22:49:55 +02:00
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
20 changed files with 1240 additions and 149 deletions
+5 -1
View File
@@ -80,8 +80,12 @@ jobs:
-o gpu-turnstile.exe ./cmd/gpu-turnstile -o gpu-turnstile.exe ./cmd/gpu-turnstile
- name: Sign and checksum - name: Sign and checksum
env:
RELEASE_SIGNING_KEY: ${{ secrets.RELEASE_SIGNING_KEY }}
run: | run: |
printf '%s\n' "${{ secrets.RELEASE_SIGNING_KEY }}" > key.pem # The key comes via the environment so its value never appears in
# the echoed command line of the run log.
printf '%s\n' "$RELEASE_SIGNING_KEY" > key.pem
chmod 600 key.pem chmod 600 key.pem
openssl pkeyutl -sign -inkey key.pem -rawin \ openssl pkeyutl -sign -inkey key.pem -rawin \
-in gpu-turnstile.exe -out gpu-turnstile.exe.sig -in gpu-turnstile.exe -out gpu-turnstile.exe.sig
+1 -1
View File
@@ -1,6 +1,6 @@
FROM golang:1.23 AS build FROM golang:1.23 AS build
WORKDIR /src WORKDIR /src
COPY go.mod ./ COPY go.mod go.sum ./
COPY cmd ./cmd COPY cmd ./cmd
COPY internal ./internal COPY internal ./internal
ARG VERSION=dev ARG VERSION=dev
+51 -18
View File
@@ -29,6 +29,12 @@ Open WebUI / n8n ────► :8188 ───┘
- Everything else (including websockets and all streaming) passes through - Everything else (including websockets and all streaming) passes through
transparently and unbuffered. transparently and unbuffered.
Each consumer is enabled by setting its URL (`OLLAMA_URL`, `COMFY_URL`) and
disabled by leaving it empty — at least one is required. With only Ollama
the proxy is a pass-through (no image jobs can arrive); with only ComfyUI
the Ollama unload/warm steps are skipped. Future consumers (e.g. local game
detection) plug into the same lock the same way.
## Configuration ## Configuration
Configuration comes from environment variables and/or an `.env`-style Configuration comes from environment variables and/or an `.env`-style
@@ -41,8 +47,8 @@ override file values. Invalid values fail at startup.
|---|---|---| |---|---|---|
| `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener | | `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener |
| `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener | | `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener |
| `OLLAMA_URL` | `http://127.0.0.1:11435` | Ollama upstream | | `OLLAMA_URL` | _(empty = disabled)_ | Ollama upstream; set to enable the Ollama consumer |
| `COMFY_URL` | `http://127.0.0.1:8189` | ComfyUI upstream | | `COMFY_URL` | _(empty = disabled)_ | ComfyUI upstream; set to enable the ComfyUI consumer |
| `UNLOAD_TIMEOUT` | `60s` | Wait for Ollama to unload before an image job | | `UNLOAD_TIMEOUT` | `60s` | Wait for Ollama to unload before an image job |
| `JOB_TIMEOUT` | `15m` | Wait for a ComfyUI job to finish | | `JOB_TIMEOUT` | `15m` | Wait for a ComfyUI job to finish |
| `LLM_WAIT_TIMEOUT` | `10m` | Max lock wait for an LLM request before 503 (wait mode) | | `LLM_WAIT_TIMEOUT` | `10m` | Max lock wait for an LLM request before 503 (wait mode) |
@@ -70,7 +76,7 @@ override file values. Invalid values fail at startup.
## Observability ## Observability
- `GET /healthz` (both listeners): `{"state":"idle|llm|image","llm_inflight":N,"image_pending":B}` - `GET /healthz` (both listeners): `{"state":"idle|llm|image","llm_inflight":N,"image_pending":B}`
- `GET /metrics` (Ollama listener): Prometheus text format — `gpu_turnstile_state`, - `GET /metrics` (both listeners): Prometheus text format — `gpu_turnstile_state`,
`gpu_turnstile_llm_inflight`, `gpu_turnstile_image_pending`, `gpu_turnstile_llm_inflight`, `gpu_turnstile_image_pending`,
`gpu_turnstile_image_jobs_total`, `gpu_turnstile_lock_wait_seconds` `gpu_turnstile_image_jobs_total`, `gpu_turnstile_lock_wait_seconds`
(histogram, `kind="llm|image"`), `gpu_turnstile_unload_seconds`. (histogram, `kind="llm|image"`), `gpu_turnstile_unload_seconds`.
@@ -85,29 +91,54 @@ override file values. Invalid values fail at startup.
```sh ```sh
go build ./cmd/gpu-turnstile go build ./cmd/gpu-turnstile
./gpu-turnstile OLLAMA_URL=http://127.0.0.1:11435 COMFY_URL=http://127.0.0.1:8189 ./gpu-turnstile
``` ```
### Run natively on Windows (primary deployment) Running the binary with no arguments in a terminal prints the help screen
(same as `-h`/`--help`); without a terminal (services, containers) a bare
invocation starts the proxy.
Download `gpu-turnstile.exe` from a release, put a `gpu-turnstile.env` ### Run natively on Windows (current primary deployment)
next to it, and run it — or install it as a Windows service from an
elevated shell: Download `gpu-turnstile.exe` from a release and install it as a Windows
service — no admin shell needed, a UAC prompt appears automatically and
the elevated child does the work (its window waits for Enter so you can
read the result):
```sh ```sh
gpu-turnstile.exe service install # auto-start service, recovery = restart gpu-turnstile.exe --install-service # installs into Program Files, auto-start
gpu-turnstile.exe service remove gpu-turnstile.exe --install-service --no-copy # register in place instead
gpu-turnstile.exe --remove-service
``` ```
The service uses the config file (services have no convenient Layout: `C:\Program Files\gpu-turnstile\` holds the exe and
environment); set `LOG_FILE` in it since there is no console. `gpu-turnstile.env`, logs go to `C:\ProgramData\gpu-turnstile\` (set
`LOG_FILE` in the env file — there is no console). The service always runs
as the virtual account `NT SERVICE\gpu-turnstile` (low-privilege,
per-service, no password); the installer automatically grants it write
access to the install and data directories — nothing else to do.
Suggested layout: `C:\Program Files\gpu-turnstile\` for the exe and ### Run natively on Linux (systemd)
`gpu-turnstile.env`, logs under `C:\ProgramData\gpu-turnstile\` via
`LOG_FILE`. The service runs as `LocalSystem` by default, which can write The same binary works on Linux. Install it as a systemd service as root:
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 ```sh
those two directories. gpu-turnstile --install-service # installs into /var/lib/gpu-turnstile, enables + starts
gpu-turnstile --install-service --no-copy # register in place instead
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. Install
copies the binary to `/var/lib/gpu-turnstile/` and the config to
`/etc/gpu-turnstile.env` (edit that one after installing). The service
runs sandboxed with `DynamicUser=yes` — a transient low-privilege UID,
read-only filesystem except its install dir (so self-update keeps
working), no capabilities, syscall-filtered: same least-privilege idea as
the Windows virtual account. The notify integration is a no-op in
containers and interactive shells.
**Auto-update is on by default**: the binary checks the repo's latest **Auto-update is on by default**: the binary checks the repo's latest
release on startup and every `UPDATE_INTERVAL`, verifies the Ed25519 release on startup and every `UPDATE_INTERVAL`, verifies the Ed25519
@@ -118,6 +149,8 @@ OpenSSL; the matching public key lives in `internal/update/pubkey.go`
(one-time setup: `openssl genpkey -algorithm ed25519 -out private.pem`, (one-time setup: `openssl genpkey -algorithm ed25519 -out private.pem`,
`openssl pkey -in private.pem -pubout -out public.pem`; private key goes `openssl pkey -in private.pem -pubout -out public.pem`; private key goes
to the `RELEASE_SIGNING_KEY` repo secret, public key is committed). to the `RELEASE_SIGNING_KEY` repo secret, public key is committed).
`gpu-turnstile --force-update` checks immediately, stages the new binary
and restarts the running service (elevating via UAC only if needed).
### Docker ### Docker
+104 -26
View File
@@ -42,6 +42,21 @@ listener is an `httputil.ReverseProxy` to its upstream. Websocket upgrades
(ComfyUI `/ws`) and streaming bodies (Ollama NDJSON / SSE) must pass through (ComfyUI `/ws`) and streaming bodies (Ollama NDJSON / SSE) must pass through
unbuffered (`FlushInterval = -1`). unbuffered (`FlushInterval = -1`).
### Modes of operation
Each GPU consumer is enabled by setting its URL and disabled by leaving it
empty — no separate flags. At least one URL must be set; a disabled
consumer gets no listener, no startup probe, and no lock participation:
- **Both set** (default deployment): full arbitration as described below.
- **Only `OLLAMA_URL`**: pure pass-through for Ollama; the LLM lock never
blocks since no image jobs can arrive.
- **Only `COMFY_URL`**: image jobs are tracked and ComfyUI's VRAM is freed
afterwards, but the Ollama unload and warm-reload steps are skipped.
- Future consumers (e.g. detecting a local game holding VRAM) plug into the
same lock the same way: enabled by their config knob, excluded when
absent.
### Lock semantics ### Lock semantics
Two-mode lock with image priority (writer-preferring RW lock, where "readers" Two-mode lock with image priority (writer-preferring RW lock, where "readers"
@@ -83,7 +98,8 @@ ComfyUI listener (`:8188` → `COMFY_URL`):
### Image job flow (`POST /prompt`) ### Image job flow (`POST /prompt`)
1. `AcquireImage()`. 1. `AcquireImage()`.
2. Unload Ollama: `GET /api/ps`; for each model `POST /api/generate 2. Unload Ollama (skipped when `OLLAMA_URL` is unset): `GET /api/ps`; for
each model `POST /api/generate
{"model":M,"keep_alive":0}`; if that returns non-2xx (embedding-only {"model":M,"keep_alive":0}`; if that returns non-2xx (embedding-only
models), `POST /api/embed {"model":M,"input":"x","keep_alive":0}`. Poll models), `POST /api/embed {"model":M,"input":"x","keep_alive":0}`. Poll
`/api/ps` every `UNLOAD_POLL_INTERVAL` (default 500 ms) until empty or `/api/ps` every `UNLOAD_POLL_INTERVAL` (default 500 ms) until empty or
@@ -117,8 +133,8 @@ override file values. A missing file is fine; a malformed one is fatal.
|---|---|---| |---|---|---|
| `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener | | `LISTEN_OLLAMA` | `:11434` | Ollama-facing listener |
| `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener | | `LISTEN_COMFY` | `:8188` | ComfyUI-facing listener |
| `OLLAMA_URL` | `http://127.0.0.1:11435` | upstream | | `OLLAMA_URL` | _(empty = disabled)_ | Ollama upstream; set to enable the Ollama consumer |
| `COMFY_URL` | `http://127.0.0.1:8189` | upstream | | `COMFY_URL` | _(empty = disabled)_ | ComfyUI upstream; set to enable the ComfyUI consumer |
| `UNLOAD_TIMEOUT` | `60s` | wait for Ollama to unload | | `UNLOAD_TIMEOUT` | `60s` | wait for Ollama to unload |
| `JOB_TIMEOUT` | `15m` | wait for ComfyUI job | | `JOB_TIMEOUT` | `15m` | wait for ComfyUI job |
| `LLM_WAIT_TIMEOUT` | `10m` | max time an LLM request waits for the lock before 503 (wait mode) | | `LLM_WAIT_TIMEOUT` | `10m` | max time an LLM request waits for the lock before 503 (wait mode) |
@@ -143,28 +159,84 @@ override file values. A missing file is fine; a malformed one is fatal.
| `UPDATE_REPO` | `https://git.rambossek.at/PUBLIC/gpu-turnstile` | repository to check for releases | | `UPDATE_REPO` | `https://git.rambossek.at/PUBLIC/gpu-turnstile` | repository to check for releases |
| `UPDATE_ASSET` | `gpu-turnstile.exe` | release asset to download | | `UPDATE_ASSET` | `gpu-turnstile.exe` | release asset to download |
Startup fails fast on unparsable values. Both upstreams are probed once at Startup fails fast on unparsable values and when neither consumer URL is
start (`/api/version`, `/system_stats`); failure is logged, not fatal. set. Enabled upstreams are probed once at start (`/api/version`,
`/system_stats`); failure is logged, not fatal.
## Native Windows deployment ## Native deployment (Windows and Linux)
The binary runs natively on Windows (the primary deployment) as well as in The binary runs natively on Windows (the current primary deployment) and on
Docker. Linux with systemd (the future GPU server), as well as in Docker.
- `gpu-turnstile.exe service install [-config path]` registers an Service management is the same on both platforms:
auto-start Windows service (needs an elevated shell). Recovery actions `gpu-turnstile --install-service [-config path]` installs, registers and
restart it 5 s after any failure. `service remove` uninstalls. starts an auto-start service; `--remove-service` stops and uninstalls it.
- **Layout**: install to `C:\Program Files\gpu-turnstile\` (exe plus Both need admin/root; on Windows a non-elevated shell triggers a UAC
`gpu-turnstile.env`); logs belong in `C:\ProgramData\gpu-turnstile\` via prompt instead of failing — the command relaunches itself elevated, waits
`LOG_FILE`. The service must be able to write its install directory for for the child, and mirrors its exit code. The legacy form
self-updates — Program Files is writable by LocalSystem and admins, which `gpu-turnstile service install|remove` does the same thing.
is why running as the default `LocalSystem` account is the simple choice.
- **Account**: the default `LocalSystem` works out of the box. For least By default install creates the canonical layout and copies the binary into
privilege, create the service with the virtual account it (Windows: `%ProgramFiles%\gpu-turnstile\`, plus
`NT SERVICE\gpu-turnstile` and grant it write access to the install and `%ProgramData%\gpu-turnstile\` for logs; Linux: `/var/lib/gpu-turnstile/`
log directories only (no network logon, no user profile). with the config at `/etc/gpu-turnstile.env`). An existing config in the
- Use a config file (above) for the service — Windows services have no target location is never overwritten. `--no-copy` registers the current
convenient environment. Logs go to `LOG_FILE` since there is no console. executable location as-is instead.
### Windows
- `--install-service` creates `%ProgramFiles%\gpu-turnstile\` and
`%ProgramData%\gpu-turnstile\`, copies the exe and (if none exists there
yet) the `gpu-turnstile.env` into the Program Files directory, and
registers that copy as a Windows service; recovery actions restart it
5 s after any failure. Logs go to the ProgramData directory via
`LOG_FILE` since there is no console.
- **Account**: the service always runs as the virtual account
`NT SERVICE\gpu-turnstile` — a per-service low-privilege identity the
SCM manages (no password, automatic logon-as-a-service right, no admin
rights, gone when the service is removed). The installer grants it
modify access to the install and data directories (self-updates rewrite
the exe) and the `LOG_FILE` directory (created if missing), plus read
access to the config file when it lives elsewhere. The grants happen
after service registration because the virtual account's SID only exists
from that point on; if a grant fails the service registration is rolled
back.
### Linux (systemd)
- `--install-service` copies the binary to `/var/lib/gpu-turnstile/`,
copies the config to `/etc/gpu-turnstile.env` if none exists there yet,
writes `/etc/systemd/system/gpu-turnstile.service`, then runs `systemctl
daemon-reload` and `enable --now`. `--remove-service` removes the unit
and the installed binary; the `/etc` config stays. The binary does not
go to `/usr/local/sbin` on purpose: replacing a running binary needs
write access to its *directory*, and granting the sandboxed service
write access to a shared system directory would let a compromised
service overwrite other binaries — `/var/lib/gpu-turnstile` is
exclusively ours.
- **Sandboxing** mirrors the Windows virtual account: the unit runs with
`DynamicUser=yes` — a transient per-service UID with no login, no home
and no password, managed entirely by systemd. `ProtectSystem=strict`
makes the filesystem read-only except `StateDirectory=gpu-turnstile`
(the install dir, so self-updates can rewrite the binary), plus
`NoNewPrivileges`, `ProtectHome`, `PrivateTmp`, `ProtectKernel*`,
`ProtectControlGroups`, `RestrictNamespaces`, `RestrictSUIDSGID`,
`RestrictRealtime`, `LockPersonality`, `MemoryDenyWriteExecute`, empty
capability sets, `RestrictAddressFamilies=AF_UNIX AF_INET AF_INET6` and
`SystemCallFilter=@system-service`. The proxy needs only outbound
TCP/UDP and the notify socket, so it loses nothing.
- The unit is `Type=notify`: the binary sends `READY=1` via
`github.com/coreos/go-systemd` only after the listeners are bound, so
`systemctl start` blocks until the proxy accepts connections. A 30 s
watchdog (`WatchdogSec=`) is pinged as long as the process runs; three
missed pings make systemd restart it. `STOPPING=1` is sent on shutdown.
All notify calls are no-ops when `NOTIFY_SOCKET` is unset (containers,
interactive shells), and the whole integration is Linux-only — Windows
builds carry no-op stubs.
- Logs go to the journal (`journalctl -u gpu-turnstile`) or to `LOG_FILE`.
- **Auto-update** works the same as on Windows: `Restart=on-failure` with
`RestartSec=5s` brings up the staged binary after the updater exits with
code 3.
- **Auto-update**: on startup and every `UPDATE_INTERVAL`, the binary - **Auto-update**: on startup and every `UPDATE_INTERVAL`, the binary
checks `UPDATE_REPO`'s latest release; if its tag is a newer `vX.Y.Z`, checks `UPDATE_REPO`'s latest release; if its tag is a newer `vX.Y.Z`,
it downloads `UPDATE_ASSET` plus its `.sig` (and `.sha256` when present) it downloads `UPDATE_ASSET` plus its `.sig` (and `.sha256` when present)
@@ -174,6 +246,11 @@ Docker.
idle the process exits with code 3 so the service recovery restarts it idle the process exits with code 3 so the service recovery restarts it
on the new version. Interactive runs only log "restart to apply". on the new version. Interactive runs only log "restart to apply".
`dev` builds and builds without an embedded public key never update. `dev` builds and builds without an embedded public key never update.
- **`--force-update`** runs the same check immediately: it downloads,
verifies and stages a newer release, and if the service is running it
restarts it right away (otherwise the new version applies on next
start). On Windows it elevates via UAC only when the stage or restart
needs permissions the caller does not have.
- **Signing setup (one time)**: `openssl genpkey -algorithm ed25519 -out - **Signing setup (one time)**: `openssl genpkey -algorithm ed25519 -out
private.pem`; `openssl pkey -in private.pem -pubout -out public.pem`. private.pem`; `openssl pkey -in private.pem -pubout -out public.pem`.
Private key → repo secret `RELEASE_SIGNING_KEY`; public key → committed into Private key → repo secret `RELEASE_SIGNING_KEY`; public key → committed into
@@ -184,7 +261,7 @@ Docker.
- `GET /healthz` on both listeners: 200 with JSON - `GET /healthz` on both listeners: 200 with JSON
`{"state":"idle|llm|image","llm_inflight":N,"image_pending":B}`. `{"state":"idle|llm|image","llm_inflight":N,"image_pending":B}`.
- `GET /metrics` on the Ollama listener: Prometheus text format, no external - `GET /metrics` on both listeners: Prometheus text format, no external
dependency needed: dependency needed:
`gpu_turnstile_state{state="…"} 1`, `gpu_turnstile_llm_inflight`, `gpu_turnstile_state{state="…"} 1`, `gpu_turnstile_llm_inflight`,
`gpu_turnstile_image_jobs_total`, `gpu_turnstile_lock_wait_seconds` `gpu_turnstile_image_jobs_total`, `gpu_turnstile_lock_wait_seconds`
@@ -233,7 +310,7 @@ gpu-turnstile/
internal/metrics/ # Prometheus exposition internal/metrics/ # Prometheus exposition
internal/config/ # env + .env file configuration internal/config/ # env + .env file configuration
internal/update/ # signed auto-updater (public key in pubkey.go) internal/update/ # signed auto-updater (public key in pubkey.go)
internal/service/ # Windows service integration internal/service/ # Windows SCM + Linux systemd (notify/watchdog) integration
Dockerfile Dockerfile
.gitea/workflows/ci.yml .gitea/workflows/ci.yml
README.md README.md
@@ -261,8 +338,9 @@ are new.
## Build and CI ## Build and CI
- Go 1.23+, `golang.org/x/sys` is the only external dependency (Windows - Go 1.23+, two external dependencies: `golang.org/x/sys` (Windows service
service integration; not used in the Linux build). `CGO_ENABLED=0`, integration) and `github.com/coreos/go-systemd` (systemd notify/watchdog,
Linux build only). `CGO_ENABLED=0`,
`-ldflags="-s -w"`, version from `git describe` injected via `-ldflags="-s -w"`, version from `git describe` injected via
`-X main.version=`. `-X main.version=`.
- Dockerfile: multi-stage, final image `gcr.io/distroless/static` (or - Dockerfile: multi-stage, final image `gcr.io/distroless/static` (or
+312 -51
View File
@@ -3,11 +3,14 @@
package main package main
import ( import (
"bufio"
"context" "context"
"errors" "errors"
"fmt" "fmt"
"io" "io"
"io/fs"
"log/slog" "log/slog"
"net"
"net/http" "net/http"
"os" "os"
"os/signal" "os/signal"
@@ -33,14 +36,127 @@ var version = "dev"
// process: a signed update has been staged and the GPU lock is idle. // process: a signed update has been staged and the GPU lock is idle.
const exitCodeUpdate = 3 const exitCodeUpdate = 3
// stdoutIsTerminal reports whether stdout is a console (char device), as
// opposed to a pipe or file — which is what Docker containers and services
// see.
func stdoutIsTerminal() bool {
fi, err := os.Stdout.Stat()
return err == nil && fi.Mode()&os.ModeCharDevice != 0
}
// parseFlags extracts -config <path> (or -config=<path>), the
// --install-service / --remove-service switches, --no-copy, -h/--help,
// -v/--version, --force-update and the hidden --elevated-child marker from
// args.
func parseFlags(args []string) (configPath string, install, remove, noCopy, help, showVersion, forceUpdate, elevatedChild 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
case args[i] == "--no-copy" || args[i] == "-no-copy":
noCopy = true
case args[i] == "-h" || args[i] == "--help" || args[i] == "-help":
help = true
case args[i] == "-v" || args[i] == "--version" || args[i] == "-version":
showVersion = true
case args[i] == "--force-update" || args[i] == "-force-update":
forceUpdate = true
case args[i] == "--elevated-child":
elevatedChild = true
default:
rest = append(rest, args[i])
}
}
return configPath, install, remove, noCopy, help, showVersion, forceUpdate, elevatedChild, rest
}
// versionLine is printed at the top of every help and error screen.
func versionLine() string { return "gpu-turnstile " + version }
const usageText = `GPU arbitration proxy for Ollama + ComfyUI
Usage:
gpu-turnstile -config <path> run the proxy
gpu-turnstile --install-service [--no-copy] [-config path] install + start as a service
gpu-turnstile --remove-service stop + uninstall the service
gpu-turnstile -v | --version print just the version
gpu-turnstile --force-update check for a signed update now,
apply it and restart the service
gpu-turnstile -h | --help this help
Options:
-config <path> config file (default: gpu-turnstile.env next to the exe)
--install-service copies the binary into the canonical location
(%ProgramFiles%\gpu-turnstile or /var/lib/gpu-turnstile)
unless --no-copy; on Windows a UAC prompt appears when
the shell is not elevated
--no-copy with --install-service: register the current location as-is
All runtime settings are environment variables or KEY=VALUE lines in the
config file (OLLAMA_URL, COMFY_URL, LOGLEVEL, ...); see README.md.
`
// printHelp prints the version header plus the full help text.
func printHelp() {
fmt.Printf("%s — %s", versionLine(), usageText)
}
// fatalUsage prints the version header, an error message and the one-line
// usage summary, then exits with code 2.
func fatalUsage(format string, args ...any) {
fmt.Fprintf(os.Stderr, "%s\n\n", versionLine())
fmt.Fprintf(os.Stderr, format+"\n\n", args...)
fmt.Fprintln(os.Stderr, "usage: gpu-turnstile [-config path] [--install-service [--no-copy] | --remove-service]")
fmt.Fprintln(os.Stderr, " gpu-turnstile service install|remove [-config path]")
os.Exit(2)
}
func main() { func main() {
configPath, args := splitConfigFlag(os.Args[1:]) configPath, install, remove, noCopy, help, showVersion, forceUpdate, elevatedChild, args := parseFlags(os.Args[1:])
if showVersion {
fmt.Println(version)
return
}
bare := configPath == "" && !install && !remove && !forceUpdate && !elevatedChild && len(args) == 0
if help || (bare && stdoutIsTerminal()) {
// Bare invocation in a terminal (e.g. double-clicked on Windows)
// shows the help instead of starting a proxy window with no visible
// explanation. Without a terminal — Docker containers, services,
// pipes — a bare invocation starts the proxy as before.
printHelp()
return
}
if len(args) > 0 && args[0] == "service" { if len(args) > 0 && args[0] == "service" {
os.Exit(serviceCommand(configPath, args[1:])) // Legacy subcommand form: gpu-turnstile service install|remove.
if len(args) != 2 || (args[1] != "install" && args[1] != "remove") {
fatalUsage("error: expected 'service install' or 'service remove'")
}
install = args[1] == "install"
remove = !install
args = nil
}
switch {
case install && remove:
fatalUsage("error: --install-service and --remove-service are mutually exclusive")
case forceUpdate && (install || remove):
fatalUsage("error: --force-update cannot be combined with --install-service/--remove-service")
case install:
os.Exit(serviceCommand(configPath, true, noCopy, elevatedChild))
case remove:
os.Exit(serviceCommand(configPath, false, noCopy, elevatedChild))
case forceUpdate:
os.Exit(forceUpdateCommand(configPath, elevatedChild))
} }
if len(args) > 0 { if len(args) > 0 {
fmt.Fprintf(os.Stderr, "usage: gpu-turnstile [-config path] | gpu-turnstile service install|remove [-config path]\n") fatalUsage("error: unknown arguments: %s", strings.Join(args, " "))
os.Exit(2)
} }
if exePath, err := os.Executable(); err == nil { if exePath, err := os.Executable(); err == nil {
@@ -49,7 +165,7 @@ func main() {
cfg, err := loadMergedConfig(configPath) cfg, err := loadMergedConfig(configPath)
if err != nil { if err != nil {
fmt.Fprintf(os.Stderr, "gpu-turnstile: %v\n", err) fmt.Fprintf(os.Stderr, "%s\n\ngpu-turnstile: %v\n", versionLine(), err)
os.Exit(1) os.Exit(1)
} }
log, logOut, logCloser := newLogger(cfg) log, logOut, logCloser := newLogger(cfg)
@@ -70,24 +186,6 @@ func main() {
} }
} }
// splitConfigFlag extracts -config <path> (or -config=<path>) from args.
func splitConfigFlag(args []string) (string, []string) {
var configPath 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=")
default:
rest = append(rest, args[i])
}
}
return configPath, rest
}
// defaultConfigPath returns gpu-turnstile.env next to the executable. // defaultConfigPath returns gpu-turnstile.env next to the executable.
func defaultConfigPath() string { func defaultConfigPath() string {
exe, err := os.Executable() exe, err := os.Executable()
@@ -157,29 +255,152 @@ func newLogger(cfg config.Config) (*slog.Logger, io.Writer, io.Closer) {
return log, out, closer return log, out, closer
} }
func serviceCommand(configPath string, args []string) int { // waitForEnter keeps an elevated child's console window open until the
if len(args) != 1 || (args[0] != "install" && args[0] != "remove") { // user has read the output.
fmt.Fprintf(os.Stderr, "usage: gpu-turnstile service install|remove [-config path]\n") func waitForEnter() {
return 2 fmt.Print("\nPress Enter to close this window...")
bufio.NewReader(os.Stdin).ReadString('\n')
}
// elevateAndMirror relaunches the current command elevated (UAC) and
// mirrors the child's exit code. verb is used in messages.
func elevateAndMirror(verb string) (int, bool) {
args := append(append([]string{}, os.Args[1:]...), "--elevated-child")
code, err := service.RelaunchElevated(args)
if errors.Is(err, service.ErrUserCancelled) {
fmt.Fprintln(os.Stderr, "gpu-turnstile: UAC prompt declined")
return 1, true
}
if err != nil {
fmt.Fprintf(os.Stderr, "gpu-turnstile: could not elevate: %v\n", err)
return 1, true
}
if code != 0 {
fmt.Fprintf(os.Stderr, "gpu-turnstile %s failed in the elevated process (exit %d)\n", verb, code)
return code, true
}
return 0, true
}
// isPermission reports whether err is a permission problem (Windows
// ERROR_ACCESS_DENIED, POSIX EACCES/EPERM, possibly wrapped).
func isPermission(err error) bool {
return errors.Is(err, fs.ErrPermission) || strings.Contains(strings.ToLower(err.Error()), "access is denied")
}
// forceUpdateCommand checks for a signed update immediately, stages it if
// newer, and restarts the service when it is running so the new binary
// takes effect. Staging into a system directory and restarting a service
// need admin rights; instead of prompting unconditionally, permission
// failures trigger the UAC relaunch so a dev copy in a user-writable
// directory updates without a prompt.
func forceUpdateCommand(configPath string, elevatedChild bool) int {
if elevatedChild {
defer waitForEnter()
}
cfg, err := loadMergedConfig(configPath)
if err != nil {
fmt.Fprintf(os.Stderr, "%s\n\ngpu-turnstile: %v\n", versionLine(), err)
return 1
}
exePath, err := os.Executable()
if err != nil {
fmt.Fprintf(os.Stderr, "gpu-turnstile: cannot locate executable: %v\n", err)
return 1
}
log, _, logCloser := newLogger(cfg)
defer logCloser.Close()
u := &update.Updater{Repo: cfg.UpdateRepo, Asset: cfg.UpdateAsset, Version: version, Log: log}
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancel()
staged, err := u.Check(ctx, exePath)
if err != nil && isPermission(err) && !service.Elevated() {
code, _ := elevateAndMirror("--force-update")
if code == 0 {
fmt.Println("update applied (elevated)")
}
return code
}
if err != nil {
fmt.Fprintf(os.Stderr, "%s\n\ngpu-turnstile: update check failed: %v\n", versionLine(), err)
return 1
}
if !staged {
fmt.Printf("%s is up to date\n", versionLine())
return 0
}
fmt.Printf("%s: update staged\n", versionLine())
restarted, err := service.RestartIfRunning()
if err != nil && isPermission(err) && !service.Elevated() {
code, _ := elevateAndMirror("--force-update")
if code == 0 {
fmt.Println("update applied (elevated)")
}
return code
}
if err != nil {
fmt.Fprintf(os.Stderr, "gpu-turnstile: update staged but service restart failed: %v\n", err)
return 1
}
if restarted {
fmt.Println("service restarted on the new version")
} else {
fmt.Println("no running service; the new version applies on next start")
}
return 0
}
// serviceCommand installs (copyBin = register the canonical-layout copy)
// or removes the service and reports the result. On Windows, when the
// shell is not elevated, the command relaunches itself through a UAC
// prompt and mirrors the elevated child's exit code. An elevated child
// waits for a keypress so its console window does not flash closed before
// the output can be read.
func serviceCommand(configPath string, install, noCopy, elevatedChild bool) int {
verb, doneVerb := "remove", "removed"
if install {
verb, doneVerb = "install", "installed"
}
if elevatedChild {
defer waitForEnter()
}
if !service.Elevated() {
code, done := elevateAndMirror(verb)
if done && code != 0 {
return code
}
if done {
fmt.Printf("service %s: %s (elevated)\n", service.Name, doneVerb)
return 0
}
} }
var err error var err error
if args[0] == "install" { if install {
path := resolveConfigPath(configPath) path := resolveConfigPath(configPath)
if abs, absErr := filepath.Abs(path); absErr == nil { if abs, absErr := filepath.Abs(path); absErr == nil {
path = abs path = abs
} }
err = service.Install(path) err = service.Install(path, !noCopy)
} else { } else {
err = service.Remove() err = service.Remove()
} }
if err != nil { if err != nil {
fmt.Fprintf(os.Stderr, "gpu-turnstile service %s: %v\n", args[0], err) fmt.Fprintf(os.Stderr, "gpu-turnstile service %s: %v\n", verb, err)
return 1 return 1
} }
fmt.Printf("service %s: %sd\n", service.Name, args[0]) fmt.Printf("service %s: %s\n", service.Name, doneVerb)
return 0 return 0
} }
// orDisabled renders an empty URL as "disabled" for the startup dump.
func orDisabled(url string) string {
if url == "" {
return "disabled"
}
return url
}
func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Writer, isService bool) error { func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Writer, isService bool) error {
// The startup line carries the version and every setting and is emitted // The startup line carries the version and every setting and is emitted
// at WARN so it is visible even with the default (quiet) log level. // at WARN so it is visible even with the default (quiet) log level.
@@ -187,8 +408,8 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
"version", version, "version", version,
"listen_ollama", cfg.ListenOllama, "listen_ollama", cfg.ListenOllama,
"listen_comfy", cfg.ListenComfy, "listen_comfy", cfg.ListenComfy,
"ollama_url", cfg.OllamaURL, "ollama_url", orDisabled(cfg.OllamaURL),
"comfy_url", cfg.ComfyURL, "comfy_url", orDisabled(cfg.ComfyURL),
"unload_timeout", cfg.UnloadTimeout, "unload_timeout", cfg.UnloadTimeout,
"job_timeout", cfg.JobTimeout, "job_timeout", cfg.JobTimeout,
"llm_wait_timeout", cfg.LLMWaitTimeout, "llm_wait_timeout", cfg.LLMWaitTimeout,
@@ -215,13 +436,20 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
) )
lk := lock.New(log) lk := lock.New(log)
ollamaClient, err := ollama.New(cfg.OllamaURL, log) // Each consumer is enabled by setting its URL; a disabled consumer gets
if err != nil { // no client, no listener and no probe.
return err 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) var comfyClient *comfy.Client
if err != nil { if cfg.ComfyURL != "" {
return err if comfyClient, err = comfy.New(cfg.ComfyURL, log); err != nil {
return err
}
} }
srv, err := proxy.New(proxy.Config{ srv, err := proxy.New(proxy.Config{
@@ -253,22 +481,54 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
return err return err
} }
// Probe both upstreams once; failure is logged, not fatal. // Probe the enabled upstreams once; failure is logged, not fatal.
probeCtx, probeCancel := context.WithTimeout(ctx, cfg.ProbeTimeout) probeCtx, probeCancel := context.WithTimeout(ctx, cfg.ProbeTimeout)
if err := ollamaClient.Probe(probeCtx); err != nil { if ollamaClient != nil {
log.Warn("ollama probe failed", "url", cfg.OllamaURL, "err", err) if err := ollamaClient.Probe(probeCtx); err != nil {
log.Warn("ollama probe failed", "url", cfg.OllamaURL, "err", err)
}
} }
if err := comfyClient.Probe(probeCtx); err != nil { if comfyClient != nil {
log.Warn("comfy probe failed", "url", cfg.ComfyURL, "err", err) if err := comfyClient.Probe(probeCtx); err != nil {
log.Warn("comfy probe failed", "url", cfg.ComfyURL, "err", err)
}
} }
probeCancel() probeCancel()
ollamaSrv := &http.Server{Addr: cfg.ListenOllama, Handler: srv.OllamaHandler()} // Bind the listeners up front so a port conflict fails fast and the
comfySrv := &http.Server{Addr: cfg.ListenComfy, Handler: srv.ComfyHandler()} // readiness notification below really means "accepting connections".
var servers []*http.Server
var listeners []net.Listener
bind := func(addr string, handler http.Handler, consumer string) error {
ln, err := net.Listen("tcp", addr)
if err != nil {
return fmt.Errorf("listen %s on %s: %w", consumer, addr, err)
}
servers = append(servers, &http.Server{Addr: addr, Handler: handler})
listeners = append(listeners, ln)
log.Warn("listening", "consumer", consumer, "addr", addr)
return nil
}
if ollamaClient != nil {
if err := bind(cfg.ListenOllama, srv.OllamaHandler(), "ollama"); err != nil {
return err
}
}
if comfyClient != nil {
if err := bind(cfg.ListenComfy, srv.ComfyHandler(), "comfy"); err != nil {
return err
}
}
errCh := make(chan error, 2) errCh := make(chan error, len(servers))
go func() { errCh <- ollamaSrv.ListenAndServe() }() for i := range servers {
go func() { errCh <- comfySrv.ListenAndServe() }() go func(s *http.Server, ln net.Listener) { errCh <- s.Serve(ln) }(servers[i], listeners[i])
}
// Tell systemd we are up and start the watchdog pings; both are no-ops
// when not running under a notify/watchdog unit.
service.NotifyReady()
service.StartWatchdog(ctx)
if cfg.AutoUpdate { if cfg.AutoUpdate {
go updateLoop(ctx, cfg, log, lk, isService) go updateLoop(ctx, cfg, log, lk, isService)
@@ -285,8 +545,9 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), cfg.ShutdownTimeout) shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), cfg.ShutdownTimeout)
defer shutdownCancel() defer shutdownCancel()
ollamaSrv.Shutdown(shutdownCtx) for _, s := range servers {
comfySrv.Shutdown(shutdownCtx) s.Shutdown(shutdownCtx)
}
return nil return nil
} }
+3 -1
View File
@@ -6,9 +6,11 @@
# ComfyUI --listen 0.0.0.0 --port 8189). # ComfyUI --listen 0.0.0.0 --port 8189).
services: services:
gpu-turnstile: gpu-turnstile:
image: git.rambossek.at/public/gpu-turnstile:v0.1.3 image: git.rambossek.at/public/gpu-turnstile:v0.1.6
restart: unless-stopped restart: unless-stopped
environment: environment:
# Each consumer is enabled by setting its URL; leave one unset to
# disable that side (no listener, no probe, no lock participation).
# Services on the Docker host itself: # Services on the Docker host itself:
OLLAMA_URL: http://host.docker.internal:11435 OLLAMA_URL: http://host.docker.internal:11435
COMFY_URL: http://host.docker.internal:8189 COMFY_URL: http://host.docker.internal:8189
+2
View File
@@ -3,3 +3,5 @@ module gpu-turnstile
go 1.23 go 1.23
require golang.org/x/sys v0.29.0 require golang.org/x/sys v0.29.0
require github.com/coreos/go-systemd/v22 v22.7.0
+2
View File
@@ -1,2 +1,4 @@
github.com/coreos/go-systemd/v22 v22.7.0 h1:LAEzFkke61DFROc7zNLX/WA2i5J8gYqe0rSj9KI28KA=
github.com/coreos/go-systemd/v22 v22.7.0/go.mod h1:xNUYtjHu2EDXbsxz1i41wouACIwT7Ybq9o0BQhMwD0w=
golang.org/x/sys v0.29.0 h1:TPYlXGxvx1MGTn2GiZDhnjPA9wZzZeGKHHmKhHYvgaU= golang.org/x/sys v0.29.0 h1:TPYlXGxvx1MGTn2GiZDhnjPA9wZzZeGKHHmKhHYvgaU=
golang.org/x/sys v0.29.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.29.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
+5 -3
View File
@@ -51,13 +51,12 @@ type Config struct {
} }
// Defaults returns the configuration used when neither the environment nor // Defaults returns the configuration used when neither the environment nor
// a config file sets a value. // a config file sets a value. The upstream URLs default to empty: a
// consumer is enabled by setting its URL, disabled by leaving it empty.
func Defaults() Config { func Defaults() Config {
return Config{ return Config{
ListenOllama: ":11434", ListenOllama: ":11434",
ListenComfy: ":8188", ListenComfy: ":8188",
OllamaURL: "http://127.0.0.1:11435",
ComfyURL: "http://127.0.0.1:8189",
UnloadTimeout: time.Minute, UnloadTimeout: time.Minute,
JobTimeout: 15 * time.Minute, JobTimeout: 15 * time.Minute,
LLMWaitTimeout: 10 * time.Minute, LLMWaitTimeout: 10 * time.Minute,
@@ -219,5 +218,8 @@ func Load(getenv func(string) string) (Config, error) {
default: default:
return cfg, fmt.Errorf("LOG_FORMAT: must be \"text\" or \"json\"") return cfg, fmt.Errorf("LOG_FORMAT: must be \"text\" or \"json\"")
} }
if cfg.OllamaURL == "" && cfg.ComfyURL == "" {
return cfg, fmt.Errorf("at least one of OLLAMA_URL or COMFY_URL must be set (each URL enables its consumer)")
}
return cfg, nil return cfg, nil
} }
+19 -1
View File
@@ -8,13 +8,21 @@ import (
) )
func TestDefaults(t *testing.T) { func TestDefaults(t *testing.T) {
cfg, err := Load(func(string) string { return "" }) cfg, err := Load(func(k string) string {
if k == "OLLAMA_URL" {
return "http://127.0.0.1:11435"
}
return ""
})
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if cfg.ListenOllama != ":11434" || cfg.ListenComfy != ":8188" { if cfg.ListenOllama != ":11434" || cfg.ListenComfy != ":8188" {
t.Fatalf("listen addrs = %s %s", cfg.ListenOllama, cfg.ListenComfy) t.Fatalf("listen addrs = %s %s", cfg.ListenOllama, cfg.ListenComfy)
} }
if cfg.ComfyURL != "" {
t.Fatalf("ComfyURL default = %q, want empty (disabled)", cfg.ComfyURL)
}
if cfg.UnloadTimeout != time.Minute || cfg.JobTimeout != 15*time.Minute { if cfg.UnloadTimeout != time.Minute || cfg.JobTimeout != 15*time.Minute {
t.Fatalf("timeouts = %v %v", cfg.UnloadTimeout, cfg.JobTimeout) t.Fatalf("timeouts = %v %v", cfg.UnloadTimeout, cfg.JobTimeout)
} }
@@ -26,6 +34,13 @@ func TestDefaults(t *testing.T) {
} }
} }
func TestLoadRequiresConsumer(t *testing.T) {
_, err := Load(func(string) string { return "" })
if err == nil || !strings.Contains(err.Error(), "OLLAMA_URL") {
t.Fatalf("err = %v, want missing-consumer error", err)
}
}
func TestParseEnvFile(t *testing.T) { func TestParseEnvFile(t *testing.T) {
input := `# comment input := `# comment
OLLAMA_URL=http://host:11435 OLLAMA_URL=http://host:11435
@@ -91,6 +106,9 @@ func TestLoadErrors(t *testing.T) {
if k == tc.key { if k == tc.key {
return tc.value return tc.value
} }
if k == "OLLAMA_URL" {
return "http://127.0.0.1:11435"
}
return "" return ""
}) })
if err == nil { if err == nil {
+52 -28
View File
@@ -102,24 +102,37 @@ type Server struct {
comfyProxy *httputil.ReverseProxy comfyProxy *httputil.ReverseProxy
} }
// New builds a Server, validating the upstream URLs. // New builds a Server, validating the upstream URLs. At least one of
// OllamaURL / ComfyURL must be set; an empty URL disables that consumer —
// its handler is then never served, its client may be nil, and the image
// job flow skips the Ollama unload/warm steps.
func New(cfg Config) (*Server, error) { func New(cfg Config) (*Server, error) {
ollamaURL, err := url.Parse(cfg.OllamaURL) if cfg.OllamaURL == "" && cfg.ComfyURL == "" {
if err != nil || ollamaURL.Scheme == "" || ollamaURL.Host == "" { return nil, fmt.Errorf("at least one of OllamaURL or ComfyURL is required")
return nil, fmt.Errorf("invalid OLLAMA_URL %q", cfg.OllamaURL)
} }
comfyURL, err := url.Parse(cfg.ComfyURL) var ollamaURL, comfyURL *url.URL
if err != nil || comfyURL.Scheme == "" || comfyURL.Host == "" { if cfg.OllamaURL != "" {
return nil, fmt.Errorf("invalid COMFY_URL %q", cfg.ComfyURL) 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 log := cfg.Log
if log == nil { if log == nil {
log = slog.Default() log = slog.Default()
} }
if cfg.UnloadPollInterval > 0 { if cfg.UnloadPollInterval > 0 && cfg.Ollama != nil {
cfg.Ollama.PollInterval = cfg.UnloadPollInterval cfg.Ollama.PollInterval = cfg.UnloadPollInterval
} }
if cfg.HistoryPollInterval > 0 { if cfg.HistoryPollInterval > 0 && cfg.Comfy != nil {
cfg.Comfy.PollInterval = cfg.HistoryPollInterval cfg.Comfy.PollInterval = cfg.HistoryPollInterval
} }
freeTimeout := cfg.FreeTimeout freeTimeout := cfg.FreeTimeout
@@ -160,7 +173,7 @@ func New(cfg Config) (*Server, error) {
max: backoffMax, max: backoffMax,
log: log, log: log,
} }
return &Server{ s := &Server{
cfg: cfg, cfg: cfg,
log: log, log: log,
logWriter: cfg.LogWriter, logWriter: cfg.LogWriter,
@@ -172,9 +185,14 @@ func New(cfg Config) (*Server, error) {
busyMode: busyMode, busyMode: busyMode,
busyStatus: busyStatus, busyStatus: busyStatus,
busyRetryAfter: busyRetryAfter, busyRetryAfter: busyRetryAfter,
ollamaProxy: newReverseProxy(ollamaURL, retry, log.With("upstream", "ollama")), }
comfyProxy: newReverseProxy(comfyURL, retry, log.With("upstream", "comfy")), if ollamaURL != nil {
}, nil s.ollamaProxy = newReverseProxy(ollamaURL, retry, log.With("upstream", "ollama"))
}
if comfyURL != nil {
s.comfyProxy = newReverseProxy(comfyURL, retry, log.With("upstream", "comfy"))
}
return s, nil
} }
// retryTransport retries requests whose failure means the upstream never // retryTransport retries requests whose failure means the upstream never
@@ -470,9 +488,13 @@ func (s *Server) OllamaHandler() http.Handler {
// ComfyHandler serves the ComfyUI-facing listener. // ComfyHandler serves the ComfyUI-facing listener.
func (s *Server) ComfyHandler() http.Handler { func (s *Server) ComfyHandler() http.Handler {
return s.logRequests("comfy", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { return s.logRequests("comfy", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/healthz" { switch r.URL.Path {
case "/healthz":
s.writeHealthz(w) s.writeHealthz(w)
return return
case "/metrics":
s.writeMetrics(w)
return
} }
if r.Method == http.MethodPost && r.URL.Path == "/prompt" { if r.Method == http.MethodPost && r.URL.Path == "/prompt" {
s.handlePrompt(w, r) s.handlePrompt(w, r)
@@ -527,19 +549,21 @@ func (s *Server) handlePrompt(w http.ResponseWriter, r *http.Request) {
s.cfg.Metrics.ObserveLockWait("image", time.Since(start).Seconds()) s.cfg.Metrics.ObserveLockWait("image", time.Since(start).Seconds())
log.Info("image lock acquired") log.Info("image lock acquired")
uctx, ucancel := context.WithTimeout(r.Context(), s.cfg.UnloadTimeout) if s.cfg.Ollama != nil {
elapsed, uerr := s.cfg.Ollama.UnloadAll(uctx) uctx, ucancel := context.WithTimeout(r.Context(), s.cfg.UnloadTimeout)
ucancel() elapsed, uerr := s.cfg.Ollama.UnloadAll(uctx)
s.cfg.Metrics.ObserveUnload(elapsed.Seconds()) ucancel()
switch { s.cfg.Metrics.ObserveUnload(elapsed.Seconds())
case r.Context().Err() != nil: switch {
s.cfg.Lock.ReleaseImage() case r.Context().Err() != nil:
return s.cfg.Lock.ReleaseImage()
case uerr != nil: return
// Degrade, don't fail the user's request on a misbehaving neighbour. case uerr != nil:
log.Warn("ollama unload incomplete; continuing", "err", uerr) // Degrade, don't fail the user's request on a misbehaving neighbour.
default: log.Warn("ollama unload incomplete; continuing", "err", uerr)
log.Info("ollama models unloaded", "seconds", elapsed.Seconds()) default:
log.Info("ollama models unloaded", "seconds", elapsed.Seconds())
}
} }
cw := &captureWriter{ResponseWriter: w, status: http.StatusOK, limit: s.captureLimit} cw := &captureWriter{ResponseWriter: w, status: http.StatusOK, limit: s.captureLimit}
@@ -585,7 +609,7 @@ func (s *Server) finishImageJob(promptID string) {
s.cfg.Lock.ReleaseImage() s.cfg.Lock.ReleaseImage()
log.Info("image lock released") log.Info("image lock released")
if s.cfg.WarmModel != "" { if s.cfg.Ollama != nil && s.cfg.WarmModel != "" {
if state, _, _ := s.cfg.Lock.Snapshot(); state == lock.StateIdle { if state, _, _ := s.cfg.Lock.Snapshot(); state == lock.StateIdle {
wctx, wcancel := context.WithTimeout(context.Background(), s.warmTimeout) wctx, wcancel := context.WithTimeout(context.Background(), s.warmTimeout)
if err := s.cfg.Ollama.Warm(wctx, s.cfg.WarmModel); err != nil { if err := s.cfg.Ollama.Warm(wctx, s.cfg.WarmModel); err != nil {
+79
View File
@@ -451,3 +451,82 @@ func TestLLMBusyWaitTimeoutRetryAfter(t *testing.T) {
t.Fatalf("Retry-After = %q", got) t.Fatalf("Retry-After = %q", got)
} }
} }
func TestComfyOnlyModeSkipsOllama(t *testing.T) {
f := newFakes(t)
comfyClient, err := comfy.New(f.comfy.URL, nil)
if err != nil {
t.Fatal(err)
}
comfyClient.PollInterval = 5 * time.Millisecond
srv, err := New(Config{
ComfyURL: f.comfy.URL,
Lock: lock.New(nil),
Comfy: comfyClient,
Metrics: metrics.New(),
JobTimeout: 2 * time.Second,
})
if err != nil {
t.Fatal(err)
}
front := httptest.NewServer(srv.ComfyHandler())
defer front.Close()
resp, err := http.Post(front.URL+"/prompt", "application/json", strings.NewReader(`{}`))
if err != nil {
t.Fatal(err)
}
resp.Body.Close()
if resp.StatusCode != 200 {
t.Fatalf("prompt status = %d", resp.StatusCode)
}
f.completeJob()
select {
case <-f.freeCh:
case <-time.After(3 * time.Second):
t.Fatal("/free never called")
}
// With Ollama disabled the unload steps must not happen.
if i := f.rec.index("unload"); i >= 0 {
f.rec.mu.Lock()
t.Fatalf("unload called with ollama disabled; events: %v", f.rec.events)
}
}
func TestOllamaOnlyMode(t *testing.T) {
f := newFakes(t)
ollamaClient, err := ollama.New(f.ollama.URL, nil)
if err != nil {
t.Fatal(err)
}
srv, err := New(Config{
OllamaURL: f.ollama.URL,
Lock: lock.New(nil),
Ollama: ollamaClient,
Metrics: metrics.New(),
LLMWaitTimeout: time.Second,
})
if err != nil {
t.Fatal(err)
}
front := httptest.NewServer(srv.OllamaHandler())
defer front.Close()
resp, err := http.Post(front.URL+"/api/chat", "application/json", strings.NewReader(`{}`))
if err != nil {
t.Fatal(err)
}
resp.Body.Close()
if resp.StatusCode != 200 {
t.Fatalf("chat status = %d", resp.StatusCode)
}
}
func TestNewRequiresConsumer(t *testing.T) {
_, err := New(Config{Lock: lock.New(nil), Metrics: metrics.New()})
if err == nil {
t.Fatal("New with no upstream URLs should fail")
}
}
+8
View File
@@ -0,0 +1,8 @@
package service
import "errors"
// ErrUserCancelled is returned by RelaunchElevated when the user declines
// the UAC prompt. Windows-only in practice; defined here so cross-platform
// callers can compare against it.
var ErrUserCancelled = errors.New("UAC prompt declined")
+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) {}
+201
View File
@@ -0,0 +1,201 @@
//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 hardened system unit.
package service
import (
"context"
"fmt"
"io"
"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"
// stateDir holds the installed binary (and staged updates); the unit's
// StateDirectory= directive makes systemd own it and grant the dynamic
// user write access. etcConfig is the config file the unit loads.
const (
stateDir = "/var/lib/" + Name
etcConfig = "/etc/" + Name + ".env"
)
// IsService reports whether the process was started by systemd.
func IsService() bool { return os.Getenv("INVOCATION_ID") != "" }
// Elevated is always true on Linux: there is no UAC equivalent; privilege
// errors surface from the failing operation with a "run as root" hint.
func Elevated() bool { return true }
// RelaunchElevated is a Windows-only concept (UAC).
func RelaunchElevated([]string) (int, error) {
return 0, fmt.Errorf("elevated relaunch is only supported on Windows")
}
// 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 hardened 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.
//
// Sandboxing mirrors the Windows virtual account: DynamicUser=yes gives
// the service a transient per-service UID with no login and no home, the
// filesystem is read-only except StateDirectory (the install dir, so
// self-updates can rewrite the binary), and the usual no-privilege-escalation
// directives apply. The proxy needs nothing but outbound TCP/UDP and the
// notify socket, so it loses nothing.
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
DynamicUser=yes
StateDirectory=%s
ProtectSystem=strict
ProtectHome=yes
PrivateTmp=yes
NoNewPrivileges=yes
ProtectKernelTunables=yes
ProtectKernelModules=yes
ProtectKernelLogs=yes
ProtectControlGroups=yes
ProtectClock=yes
RestrictNamespaces=yes
RestrictSUIDSGID=yes
RestrictRealtime=yes
LockPersonality=yes
MemoryDenyWriteExecute=yes
CapabilityBoundingSet=
AmbientCapabilities=
RestrictAddressFamilies=AF_UNIX AF_INET AF_INET6
SystemCallFilter=@system-service
SystemCallErrorNumber=EPERM
[Install]
WantedBy=multi-user.target
`, exePath, configPath, Name)
}
// copyFile copies src to dst, creating dst with the given mode.
func copyFile(src, dst string, mode os.FileMode) error {
in, err := os.Open(src)
if err != nil {
return err
}
defer in.Close()
out, err := os.OpenFile(dst, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, mode)
if err != nil {
return err
}
defer out.Close()
if _, err := io.Copy(out, in); err != nil {
return err
}
return out.Close()
}
// Install copies the current executable into /var/lib/gpu-turnstile, makes
// sure /etc/gpu-turnstile.env exists (copied from the given config file if
// provided), writes the hardened unit, then enables and starts it. With
// copyBin=false the current executable location and config path are
// registered as-is instead. Needs root.
//
// The binary lives in the StateDirectory rather than /usr/local/sbin on
// purpose: replacing a running binary needs write access to its
// *directory*, and granting the sandboxed service user write access to a
// shared system directory would let a compromised service overwrite other
// binaries. /var/lib/gpu-turnstile is exclusively ours.
func Install(configPath string, copyBin bool) error {
exe, err := os.Executable()
if err != nil {
return err
}
if abs, absErr := filepath.Abs(exe); absErr == nil {
exe = abs
}
cfg := etcConfig
if copyBin {
if err := os.MkdirAll(stateDir, 0o755); err != nil {
return fmt.Errorf("create %s (run as root): %w", stateDir, err)
}
installedExe := filepath.Join(stateDir, Name)
if exe != installedExe {
if err := copyFile(exe, installedExe, 0o755); err != nil {
return fmt.Errorf("install binary to %s: %w", installedExe, err)
}
}
exe = installedExe
if _, err := os.Stat(etcConfig); os.IsNotExist(err) && configPath != "" {
// Missing config is not fatal: the service fails fast with a
// clear "no consumer URL" error until the user writes one.
copyFile(configPath, etcConfig, 0o644) //nolint:errcheck // best effort
}
} else if configPath != "" {
if abs, absErr := filepath.Abs(configPath); absErr == nil {
cfg = abs
}
}
if err := os.WriteFile(unitPath, []byte(renderUnit(exe, cfg)), 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 and the
// installed binary. The config file in /etc is left in place (user data).
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)
}
os.RemoveAll(stateDir) // installed binary + staged updates; ignore error
if out, err := exec.Command("systemctl", "daemon-reload").CombinedOutput(); err != nil {
return fmt.Errorf("systemctl daemon-reload: %w (%s)", err, out)
}
return nil
}
// RestartIfRunning restarts the systemd unit when it is active (used after
// a forced update staged a new binary). Reports whether a restart
// happened. An inactive or missing unit is not an error. Needs root.
func RestartIfRunning() (bool, error) {
if err := exec.Command("systemctl", "is-active", "--quiet", Name+".service").Run(); err != nil {
return false, nil // inactive or not installed
}
if out, err := exec.Command("systemctl", "restart", Name+".service").CombinedOutput(); err != nil {
return false, fmt.Errorf("systemctl restart (run as root): %w (%s)", err, out)
}
return true, nil
}
+29
View File
@@ -0,0 +1,29 @@
//go:build linux
package service
import (
"strings"
"testing"
)
func TestRenderUnit(t *testing.T) {
unit := renderUnit("/var/lib/gpu-turnstile/gpu-turnstile", "/etc/gpu-turnstile.env")
for _, want := range []string{
"Type=notify",
"WatchdogSec=30s",
`ExecStart="/var/lib/gpu-turnstile/gpu-turnstile" -config "/etc/gpu-turnstile.env"`,
"Restart=on-failure",
"WantedBy=multi-user.target",
"DynamicUser=yes",
"StateDirectory=gpu-turnstile",
"ProtectSystem=strict",
"NoNewPrivileges=yes",
"RestrictAddressFamilies=AF_UNIX AF_INET AF_INET6",
"SystemCallFilter=@system-service",
} {
if !strings.Contains(unit, want) {
t.Fatalf("unit missing %q:\n%s", want, unit)
}
}
}
+18 -8
View File
@@ -1,8 +1,8 @@
//go:build !windows //go:build !windows && !linux
// Package service provides the non-Windows stubs for the Windows service // Package service provides the stubs for platforms without service
// integration. Run falls back to plain signal handling; install/remove // integration (Windows uses the SCM, Linux uses systemd). Run falls back
// are unsupported. // to plain signal handling; install/remove are unsupported.
package service package service
import ( import (
@@ -15,7 +15,7 @@ import (
// Name matches the Windows service name. // Name matches the Windows service name.
const Name = "gpu-turnstile" const Name = "gpu-turnstile"
var errUnsupported = errors.New("service management is only supported on Windows") var errUnsupported = errors.New("service management is only supported on Windows and Linux (systemd)")
// IsService is always false on non-Windows platforms. // IsService is always false on non-Windows platforms.
func IsService() bool { return false } func IsService() bool { return false }
@@ -28,8 +28,18 @@ func Run(run func(ctx context.Context) error) error {
return run(ctx) return run(ctx)
} }
// Install is unsupported on non-Windows platforms. // Install is unsupported on non-Windows, non-Linux platforms.
func Install(string) error { return errUnsupported } func Install(string, bool) error { return errUnsupported }
// Remove is unsupported on non-Windows platforms. // Remove is unsupported on non-Windows, non-Linux platforms.
func Remove() error { return errUnsupported } func Remove() error { return errUnsupported }
// Elevated is always true here: there is no UAC concept, and privilege
// errors surface from the failing operation with a "run as root" hint.
func Elevated() bool { return true }
// RelaunchElevated is unsupported on non-Windows, non-Linux platforms.
func RelaunchElevated([]string) (int, error) { return 0, errUnsupported }
// RestartIfRunning is a no-op on platforms without service integration.
func RestartIfRunning() (bool, error) { return false, nil }
+288 -10
View File
@@ -2,22 +2,53 @@
// Package service integrates gpu-turnstile with the Windows Service // Package service integrates gpu-turnstile with the Windows Service
// Control Manager: running as a service with graceful stop, plus // Control Manager: running as a service with graceful stop, plus
// install/remove helpers. // install/remove helpers. Installed services always run as the virtual
// account NT SERVICE\gpu-turnstile — a per-service low-privilege identity
// managed by the SCM, with no password and no admin rights.
package service package service
import ( import (
"context" "context"
"errors"
"fmt" "fmt"
"io"
"os" "os"
"os/exec"
"path/filepath"
"strings"
"syscall"
"time" "time"
"unsafe"
"golang.org/x/sys/windows"
"golang.org/x/sys/windows/svc" "golang.org/x/sys/windows/svc"
"golang.org/x/sys/windows/svc/mgr" "golang.org/x/sys/windows/svc/mgr"
"gpu-turnstile/internal/config"
) )
// Name is the Windows service name. // Name is the Windows service name.
const Name = "gpu-turnstile" const Name = "gpu-turnstile"
// virtualAccount is the per-service identity the service runs as. The SCM
// manages it: no password, automatic "log on as a service" right, gone
// when the service is removed.
const virtualAccount = `NT SERVICE\` + Name
// installDirs returns the canonical install (Program Files) and data
// (ProgramData) directories.
func installDirs() (install, data string) {
pf := os.Getenv("ProgramFiles")
if pf == "" {
pf = `C:\Program Files`
}
pd := os.Getenv("ProgramData")
if pd == "" {
pd = `C:\ProgramData`
}
return filepath.Join(pf, Name), filepath.Join(pd, Name)
}
// IsService reports whether the process is running as a Windows service. // IsService reports whether the process is running as a Windows service.
func IsService() bool { func IsService() bool {
isSvc, err := svc.IsWindowsService() isSvc, err := svc.IsWindowsService()
@@ -64,15 +95,53 @@ func (h *handler) Execute(_ []string, requests <-chan svc.ChangeRequest, status
} }
} }
// Install registers gpu-turnstile as an auto-start Windows service whose // Install registers gpu-turnstile as an auto-start Windows service running
// binPath loads the given config file. Recovery actions restart the // as the NT SERVICE\gpu-turnstile virtual account, whose binPath loads the
// service after 5s on failure — this is also what brings up a staged // given config file. With copyBin it first creates the canonical layout —
// update after the updater exits with a non-zero code. // the binary is copied into %ProgramFiles%\gpu-turnstile and the config
func Install(configPath string) error { // next to it (an existing config there is kept), %ProgramData%\gpu-turnstile
// is created for logs — and registers that copy; with copyBin=false the
// current executable location is registered as-is. 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. After registering,
// the virtual account is granted modify access to the install and data
// directories (self-updates rewrite the exe), and read access to the
// config file if it lives elsewhere. The grants must come after
// CreateService: the virtual account's SID only exists once the service is
// registered.
func Install(configPath string, copyBin bool) error {
exe, err := os.Executable() exe, err := os.Executable()
if err != nil { if err != nil {
return err return err
} }
if abs, absErr := filepath.Abs(exe); absErr == nil {
exe = abs
}
if configPath != "" {
if abs, absErr := filepath.Abs(configPath); absErr == nil {
configPath = abs
}
}
installDir, _ := installDirs()
if copyBin && !strings.EqualFold(filepath.Dir(exe), installDir) {
if err := os.MkdirAll(installDir, 0o755); err != nil {
return fmt.Errorf("create %s: %w", installDir, err)
}
installedExe := filepath.Join(installDir, "gpu-turnstile.exe")
if err := copyFile(exe, installedExe); err != nil {
return fmt.Errorf("copy binary to %s: %w", installedExe, err)
}
exe = installedExe
targetCfg := filepath.Join(installDir, "gpu-turnstile.env")
if configPath != "" && !strings.EqualFold(configPath, targetCfg) {
if _, statErr := os.Stat(targetCfg); os.IsNotExist(statErr) {
copyFile(configPath, targetCfg) //nolint:errcheck // best effort
}
configPath = targetCfg
}
}
m, err := mgr.Connect() m, err := mgr.Connect()
if err != nil { if err != nil {
return fmt.Errorf("connect to service manager (run as administrator): %w", err) return fmt.Errorf("connect to service manager (run as administrator): %w", err)
@@ -81,9 +150,10 @@ func Install(configPath string) error {
binPath := fmt.Sprintf(`"%s" -config "%s"`, exe, configPath) binPath := fmt.Sprintf(`"%s" -config "%s"`, exe, configPath)
s, err := m.CreateService(Name, binPath, mgr.Config{ s, err := m.CreateService(Name, binPath, mgr.Config{
StartType: mgr.StartAutomatic, StartType: mgr.StartAutomatic,
DisplayName: "gpu-turnstile", DisplayName: "gpu-turnstile",
Description: "GPU arbitration proxy for Ollama and ComfyUI", Description: "GPU arbitration proxy for Ollama and ComfyUI",
ServiceStartName: virtualAccount,
}) })
if err != nil { if err != nil {
return fmt.Errorf("create service: %w", err) return fmt.Errorf("create service: %w", err)
@@ -97,10 +167,71 @@ func Install(configPath string) error {
if err := s.SetRecoveryActionsOnNonCrashFailures(true); err != nil { if err := s.SetRecoveryActionsOnNonCrashFailures(true); err != nil {
return fmt.Errorf("set failure actions flag: %w", err) return fmt.Errorf("set failure actions flag: %w", err)
} }
if err := grantAll(exe, configPath); err != nil {
s.Delete() // roll back so a retry starts clean
return err
}
// Best effort: start now instead of waiting for the next boot. A
// missing config (no consumer URLs) fails the start; the service stays
// registered and can be started once the config exists.
s.Start()
return nil return nil
} }
// Remove stops (if running) and unregisters the service. // copyFile copies src to dst (0755 on the new file).
func copyFile(src, dst string) error {
in, err := os.Open(src)
if err != nil {
return err
}
defer in.Close()
out, err := os.OpenFile(dst, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o755)
if err != nil {
return err
}
if _, err := io.Copy(out, in); err != nil {
out.Close()
return err
}
return out.Close()
}
// grantAll gives the virtual account every ACL the service needs: modify
// on the install and ProgramData directories and the LOG_FILE directory
// (if configured elsewhere), read on a config file outside the install
// directory.
func grantAll(exe, configPath string) error {
_, dataDir := installDirs()
if err := os.MkdirAll(dataDir, 0o755); err != nil {
return fmt.Errorf("create %s: %w", dataDir, err)
}
if err := grantAccess(dataDir, "(OI)(CI)(M)"); err != nil {
return err
}
exeDir := filepath.Dir(exe)
if err := grantAccess(exeDir, "(OI)(CI)(M)"); err != nil {
return err
}
if configPath != "" && !strings.HasPrefix(strings.ToLower(configPath), strings.ToLower(exeDir)+`\`) {
if err := grantAccess(configPath, "(R)"); err != nil {
return err
}
}
if logFile := configuredLogFile(configPath); logFile != "" {
dir := filepath.Dir(logFile)
if err := os.MkdirAll(dir, 0o755); err == nil {
if err := grantAccess(dir, "(OI)(CI)(M)"); err != nil {
return err
}
}
}
return nil
}
// Remove stops (if running) and unregisters the service. The virtual
// account ceases to exist with it; the ACL grants on the install and log
// directories are left in place (harmless without the account).
func Remove() error { func Remove() error {
m, err := mgr.Connect() m, err := mgr.Connect()
if err != nil { if err != nil {
@@ -118,3 +249,150 @@ func Remove() error {
} }
return nil return nil
} }
var procShellExecuteExW = windows.NewLazySystemDLL("shell32.dll").NewProc("ShellExecuteExW")
const seeMaskNoCloseProcess = 0x40
// shellExecuteInfo mirrors SHELLEXECUTEINFOW (64-bit layout).
type shellExecuteInfo struct {
cbSize uint32
fMask uint32
hwnd uintptr
lpVerb *uint16
lpFile *uint16
lpParameters *uint16
lpDirectory *uint16
nShow int32
_ int32
hInstApp uintptr
lpIDList unsafe.Pointer
lpClass *uint16
hkeyClass uintptr
dwHotKey uint32
_ uint32
hIcon uintptr
hProcess windows.Handle
}
// Elevated reports whether the current process token is UAC-elevated.
func Elevated() bool {
var token windows.Token
if err := windows.OpenProcessToken(windows.CurrentProcess(), windows.TOKEN_QUERY, &token); err != nil {
return false
}
defer token.Close()
return token.IsElevated()
}
// RelaunchElevated re-runs the current executable elevated via the UAC
// "runas" verb with the given arguments, waits for the child, and returns
// its exit code. The child gets a fresh console window for its output.
func RelaunchElevated(args []string) (int, error) {
exe, err := os.Executable()
if err != nil {
return 0, err
}
quoted := make([]string, len(args))
for i, a := range args {
quoted[i] = syscall.EscapeArg(a)
}
cwd, _ := os.Getwd()
verb, _ := windows.UTF16PtrFromString("runas")
exeP, _ := windows.UTF16PtrFromString(exe)
params, _ := windows.UTF16PtrFromString(strings.Join(quoted, " "))
dir, _ := windows.UTF16PtrFromString(cwd)
info := shellExecuteInfo{
fMask: seeMaskNoCloseProcess,
lpVerb: verb,
lpFile: exeP,
lpParameters: params,
lpDirectory: dir,
nShow: windows.SW_NORMAL,
}
info.cbSize = uint32(unsafe.Sizeof(info))
r, _, callErr := procShellExecuteExW.Call(uintptr(unsafe.Pointer(&info)))
if r == 0 {
if errors.Is(callErr, syscall.Errno(1223)) { // ERROR_CANCELLED
return 0, ErrUserCancelled
}
return 0, fmt.Errorf("ShellExecuteEx: %w", callErr)
}
defer windows.CloseHandle(windows.Handle(info.hProcess))
windows.WaitForSingleObject(windows.Handle(info.hProcess), windows.INFINITE)
var code uint32
if err := windows.GetExitCodeProcess(windows.Handle(info.hProcess), &code); err != nil {
return 0, err
}
return int(code), nil
}
// grantAccess gives the virtual account the icacls permission set (e.g.
// "(OI)(CI)(M)") on path.
func grantAccess(path, perms string) error {
out, err := exec.Command("icacls", path, "/grant", virtualAccount+":"+perms).CombinedOutput()
if err != nil {
return fmt.Errorf("grant %s access to %s: %w (%s)", virtualAccount, path, err, strings.TrimSpace(string(out)))
}
return nil
}
// configuredLogFile reads LOG_FILE from the config file so the installer
// can pre-create and ACL the log directory. "" when unset or unreadable.
func configuredLogFile(configPath string) string {
f, err := os.Open(configPath)
if err != nil {
return ""
}
defer f.Close()
values, err := config.ParseEnvFile(f)
if err != nil {
return ""
}
return values["LOG_FILE"]
}
// RestartIfRunning restarts the service when it is installed and running
// (used after a forced update staged a new binary). Reports whether a
// restart happened. A service that is not installed or not running is not
// an error. Needs elevation.
func RestartIfRunning() (bool, error) {
m, err := mgr.Connect()
if err != nil {
return false, fmt.Errorf("connect to service manager (run as administrator): %w", err)
}
defer m.Disconnect()
s, err := m.OpenService(Name)
if err != nil {
return false, nil // not installed
}
defer s.Close()
st, err := s.Query()
if err != nil {
return false, fmt.Errorf("query service: %w", err)
}
if st.State != svc.Running && st.State != svc.StartPending {
return false, nil
}
if _, err := s.Control(svc.Stop); err != nil {
return false, fmt.Errorf("stop service: %w", err)
}
deadline := time.Now().Add(30 * time.Second)
for {
st, err = s.Query()
if err != nil {
return false, fmt.Errorf("query service: %w", err)
}
if st.State == svc.Stopped {
break
}
if time.Now().After(deadline) {
return false, fmt.Errorf("service did not stop within 30s")
}
time.Sleep(300 * time.Millisecond)
}
if err := s.Start(); err != nil {
return false, fmt.Errorf("start service: %w", err)
}
return true, nil
}
+1 -1
View File
@@ -11,6 +11,6 @@ package update
// the RELEASE_SIGNING_KEY repository secret. When empty, the updater refuses // the RELEASE_SIGNING_KEY repository secret. When empty, the updater refuses
// to update (e.g. development builds). // to update (e.g. development builds).
var publicKeyPEM = `-----BEGIN PUBLIC KEY----- var publicKeyPEM = `-----BEGIN PUBLIC KEY-----
MCowBQYDK2VwAyEA+gJbSvgeYX58woPQGbSC8x8Zw4OTDiiQ7/19seZKfSQ= MCowBQYDK2VwAyEAqTAJ0CCeAQI7MhFlgc5xNmF/CfvLVUAY3ZAoeAS0tT8=
-----END PUBLIC KEY----- -----END PUBLIC KEY-----
` `