Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
da02457fd5 | ||
|
|
6a6359a630 | ||
|
|
e4e4348e8d | ||
|
|
b421bb7bfb | ||
|
|
909918657f | ||
|
|
a0435c858e | ||
|
|
fbab0bba33 | ||
|
|
802a64280f | ||
|
|
a88955e35c | ||
|
|
97624470eb | ||
|
|
75f16a0229 | ||
|
|
9481fd8418 | ||
|
|
83812cebf3 | ||
|
|
30a1f55aae | ||
|
|
642cc36a39 |
@@ -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
@@ -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
|
||||||
|
|||||||
@@ -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
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
+306
-45
@@ -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,14 +436,21 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri
|
|||||||
)
|
)
|
||||||
|
|
||||||
lk := lock.New(log)
|
lk := lock.New(log)
|
||||||
ollamaClient, err := ollama.New(cfg.OllamaURL, log)
|
// Each consumer is enabled by setting its URL; a disabled consumer gets
|
||||||
if err != nil {
|
// no client, no listener and no probe.
|
||||||
|
var ollamaClient *ollama.Client
|
||||||
|
var err error
|
||||||
|
if cfg.OllamaURL != "" {
|
||||||
|
if ollamaClient, err = ollama.New(cfg.OllamaURL, log); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
comfyClient, err := comfy.New(cfg.ComfyURL, log)
|
}
|
||||||
if err != nil {
|
var comfyClient *comfy.Client
|
||||||
|
if cfg.ComfyURL != "" {
|
||||||
|
if comfyClient, err = comfy.New(cfg.ComfyURL, log); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
srv, err := proxy.New(proxy.Config{
|
srv, err := proxy.New(proxy.Config{
|
||||||
OllamaURL: cfg.OllamaURL,
|
OllamaURL: cfg.OllamaURL,
|
||||||
@@ -253,22 +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 ollamaClient != nil {
|
||||||
if err := ollamaClient.Probe(probeCtx); err != nil {
|
if err := ollamaClient.Probe(probeCtx); err != nil {
|
||||||
log.Warn("ollama probe failed", "url", cfg.OllamaURL, "err", err)
|
log.Warn("ollama probe failed", "url", cfg.OllamaURL, "err", err)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
if comfyClient != nil {
|
||||||
if err := comfyClient.Probe(probeCtx); err != nil {
|
if err := comfyClient.Probe(probeCtx); err != nil {
|
||||||
log.Warn("comfy probe failed", "url", cfg.ComfyURL, "err", err)
|
log.Warn("comfy probe failed", "url", cfg.ComfyURL, "err", err)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
probeCancel()
|
probeCancel()
|
||||||
|
|
||||||
ollamaSrv := &http.Server{Addr: cfg.ListenOllama, Handler: srv.OllamaHandler()}
|
// Bind the listeners up front so a port conflict fails fast and the
|
||||||
comfySrv := &http.Server{Addr: cfg.ListenComfy, Handler: srv.ComfyHandler()}
|
// readiness notification below really means "accepting connections".
|
||||||
|
var servers []*http.Server
|
||||||
|
var listeners []net.Listener
|
||||||
|
bind := func(addr string, handler http.Handler, consumer string) error {
|
||||||
|
ln, err := net.Listen("tcp", addr)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("listen %s on %s: %w", consumer, addr, err)
|
||||||
|
}
|
||||||
|
servers = append(servers, &http.Server{Addr: addr, Handler: handler})
|
||||||
|
listeners = append(listeners, ln)
|
||||||
|
log.Warn("listening", "consumer", consumer, "addr", addr)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if ollamaClient != nil {
|
||||||
|
if err := bind(cfg.ListenOllama, srv.OllamaHandler(), "ollama"); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if comfyClient != nil {
|
||||||
|
if err := bind(cfg.ListenComfy, srv.ComfyHandler(), "comfy"); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
errCh := make(chan error, 2)
|
errCh := make(chan error, len(servers))
|
||||||
go func() { errCh <- ollamaSrv.ListenAndServe() }()
|
for i := range servers {
|
||||||
go func() { errCh <- comfySrv.ListenAndServe() }()
|
go func(s *http.Server, ln net.Listener) { errCh <- s.Serve(ln) }(servers[i], listeners[i])
|
||||||
|
}
|
||||||
|
|
||||||
|
// Tell systemd we are up and start the watchdog pings; both are no-ops
|
||||||
|
// when not running under a notify/watchdog unit.
|
||||||
|
service.NotifyReady()
|
||||||
|
service.StartWatchdog(ctx)
|
||||||
|
|
||||||
if cfg.AutoUpdate {
|
if cfg.AutoUpdate {
|
||||||
go updateLoop(ctx, cfg, log, lk, isService)
|
go updateLoop(ctx, cfg, log, lk, isService)
|
||||||
@@ -285,8 +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
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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=
|
||||||
|
|||||||
@@ -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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,13 +8,21 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
func TestDefaults(t *testing.T) {
|
func TestDefaults(t *testing.T) {
|
||||||
cfg, err := Load(func(string) string { return "" })
|
cfg, err := Load(func(k string) string {
|
||||||
|
if k == "OLLAMA_URL" {
|
||||||
|
return "http://127.0.0.1:11435"
|
||||||
|
}
|
||||||
|
return ""
|
||||||
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if cfg.ListenOllama != ":11434" || cfg.ListenComfy != ":8188" {
|
if cfg.ListenOllama != ":11434" || cfg.ListenComfy != ":8188" {
|
||||||
t.Fatalf("listen addrs = %s %s", cfg.ListenOllama, cfg.ListenComfy)
|
t.Fatalf("listen addrs = %s %s", cfg.ListenOllama, cfg.ListenComfy)
|
||||||
}
|
}
|
||||||
|
if cfg.ComfyURL != "" {
|
||||||
|
t.Fatalf("ComfyURL default = %q, want empty (disabled)", cfg.ComfyURL)
|
||||||
|
}
|
||||||
if cfg.UnloadTimeout != time.Minute || cfg.JobTimeout != 15*time.Minute {
|
if cfg.UnloadTimeout != time.Minute || cfg.JobTimeout != 15*time.Minute {
|
||||||
t.Fatalf("timeouts = %v %v", cfg.UnloadTimeout, cfg.JobTimeout)
|
t.Fatalf("timeouts = %v %v", cfg.UnloadTimeout, cfg.JobTimeout)
|
||||||
}
|
}
|
||||||
@@ -26,6 +34,13 @@ func TestDefaults(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestLoadRequiresConsumer(t *testing.T) {
|
||||||
|
_, err := Load(func(string) string { return "" })
|
||||||
|
if err == nil || !strings.Contains(err.Error(), "OLLAMA_URL") {
|
||||||
|
t.Fatalf("err = %v, want missing-consumer error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestParseEnvFile(t *testing.T) {
|
func TestParseEnvFile(t *testing.T) {
|
||||||
input := `# comment
|
input := `# comment
|
||||||
OLLAMA_URL=http://host:11435
|
OLLAMA_URL=http://host:11435
|
||||||
@@ -91,6 +106,9 @@ func TestLoadErrors(t *testing.T) {
|
|||||||
if k == tc.key {
|
if k == tc.key {
|
||||||
return tc.value
|
return tc.value
|
||||||
}
|
}
|
||||||
|
if k == "OLLAMA_URL" {
|
||||||
|
return "http://127.0.0.1:11435"
|
||||||
|
}
|
||||||
return ""
|
return ""
|
||||||
})
|
})
|
||||||
if err == nil {
|
if err == nil {
|
||||||
|
|||||||
+37
-13
@@ -102,24 +102,37 @@ type Server struct {
|
|||||||
comfyProxy *httputil.ReverseProxy
|
comfyProxy *httputil.ReverseProxy
|
||||||
}
|
}
|
||||||
|
|
||||||
// New builds a Server, validating the upstream URLs.
|
// New builds a Server, validating the upstream URLs. At least one of
|
||||||
|
// OllamaURL / ComfyURL must be set; an empty URL disables that consumer —
|
||||||
|
// its handler is then never served, its client may be nil, and the image
|
||||||
|
// job flow skips the Ollama unload/warm steps.
|
||||||
func New(cfg Config) (*Server, error) {
|
func New(cfg Config) (*Server, error) {
|
||||||
ollamaURL, err := url.Parse(cfg.OllamaURL)
|
if cfg.OllamaURL == "" && cfg.ComfyURL == "" {
|
||||||
if err != nil || ollamaURL.Scheme == "" || ollamaURL.Host == "" {
|
return nil, fmt.Errorf("at least one of OllamaURL or ComfyURL is required")
|
||||||
|
}
|
||||||
|
var ollamaURL, comfyURL *url.URL
|
||||||
|
if cfg.OllamaURL != "" {
|
||||||
|
u, err := url.Parse(cfg.OllamaURL)
|
||||||
|
if err != nil || u.Scheme == "" || u.Host == "" {
|
||||||
return nil, fmt.Errorf("invalid OLLAMA_URL %q", cfg.OllamaURL)
|
return nil, fmt.Errorf("invalid OLLAMA_URL %q", cfg.OllamaURL)
|
||||||
}
|
}
|
||||||
comfyURL, err := url.Parse(cfg.ComfyURL)
|
ollamaURL = u
|
||||||
if err != nil || comfyURL.Scheme == "" || comfyURL.Host == "" {
|
}
|
||||||
|
if cfg.ComfyURL != "" {
|
||||||
|
u, err := url.Parse(cfg.ComfyURL)
|
||||||
|
if err != nil || u.Scheme == "" || u.Host == "" {
|
||||||
return nil, fmt.Errorf("invalid COMFY_URL %q", cfg.ComfyURL)
|
return nil, fmt.Errorf("invalid COMFY_URL %q", cfg.ComfyURL)
|
||||||
}
|
}
|
||||||
|
comfyURL = u
|
||||||
|
}
|
||||||
log := cfg.Log
|
log := cfg.Log
|
||||||
if log == nil {
|
if log == nil {
|
||||||
log = slog.Default()
|
log = slog.Default()
|
||||||
}
|
}
|
||||||
if cfg.UnloadPollInterval > 0 {
|
if cfg.UnloadPollInterval > 0 && cfg.Ollama != nil {
|
||||||
cfg.Ollama.PollInterval = cfg.UnloadPollInterval
|
cfg.Ollama.PollInterval = cfg.UnloadPollInterval
|
||||||
}
|
}
|
||||||
if cfg.HistoryPollInterval > 0 {
|
if cfg.HistoryPollInterval > 0 && cfg.Comfy != nil {
|
||||||
cfg.Comfy.PollInterval = cfg.HistoryPollInterval
|
cfg.Comfy.PollInterval = cfg.HistoryPollInterval
|
||||||
}
|
}
|
||||||
freeTimeout := cfg.FreeTimeout
|
freeTimeout := cfg.FreeTimeout
|
||||||
@@ -160,7 +173,7 @@ func New(cfg Config) (*Server, error) {
|
|||||||
max: backoffMax,
|
max: backoffMax,
|
||||||
log: log,
|
log: log,
|
||||||
}
|
}
|
||||||
return &Server{
|
s := &Server{
|
||||||
cfg: cfg,
|
cfg: cfg,
|
||||||
log: log,
|
log: log,
|
||||||
logWriter: cfg.LogWriter,
|
logWriter: cfg.LogWriter,
|
||||||
@@ -172,9 +185,14 @@ func New(cfg Config) (*Server, error) {
|
|||||||
busyMode: busyMode,
|
busyMode: busyMode,
|
||||||
busyStatus: busyStatus,
|
busyStatus: busyStatus,
|
||||||
busyRetryAfter: busyRetryAfter,
|
busyRetryAfter: busyRetryAfter,
|
||||||
ollamaProxy: newReverseProxy(ollamaURL, retry, log.With("upstream", "ollama")),
|
}
|
||||||
comfyProxy: newReverseProxy(comfyURL, retry, log.With("upstream", "comfy")),
|
if ollamaURL != nil {
|
||||||
}, nil
|
s.ollamaProxy = newReverseProxy(ollamaURL, retry, log.With("upstream", "ollama"))
|
||||||
|
}
|
||||||
|
if comfyURL != nil {
|
||||||
|
s.comfyProxy = newReverseProxy(comfyURL, retry, log.With("upstream", "comfy"))
|
||||||
|
}
|
||||||
|
return s, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// retryTransport retries requests whose failure means the upstream never
|
// retryTransport retries requests whose failure means the upstream never
|
||||||
@@ -470,9 +488,13 @@ func (s *Server) OllamaHandler() http.Handler {
|
|||||||
// ComfyHandler serves the ComfyUI-facing listener.
|
// ComfyHandler serves the ComfyUI-facing listener.
|
||||||
func (s *Server) ComfyHandler() http.Handler {
|
func (s *Server) ComfyHandler() http.Handler {
|
||||||
return s.logRequests("comfy", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
return s.logRequests("comfy", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
if r.URL.Path == "/healthz" {
|
switch r.URL.Path {
|
||||||
|
case "/healthz":
|
||||||
s.writeHealthz(w)
|
s.writeHealthz(w)
|
||||||
return
|
return
|
||||||
|
case "/metrics":
|
||||||
|
s.writeMetrics(w)
|
||||||
|
return
|
||||||
}
|
}
|
||||||
if r.Method == http.MethodPost && r.URL.Path == "/prompt" {
|
if r.Method == http.MethodPost && r.URL.Path == "/prompt" {
|
||||||
s.handlePrompt(w, r)
|
s.handlePrompt(w, r)
|
||||||
@@ -527,6 +549,7 @@ func (s *Server) handlePrompt(w http.ResponseWriter, r *http.Request) {
|
|||||||
s.cfg.Metrics.ObserveLockWait("image", time.Since(start).Seconds())
|
s.cfg.Metrics.ObserveLockWait("image", time.Since(start).Seconds())
|
||||||
log.Info("image lock acquired")
|
log.Info("image lock acquired")
|
||||||
|
|
||||||
|
if s.cfg.Ollama != nil {
|
||||||
uctx, ucancel := context.WithTimeout(r.Context(), s.cfg.UnloadTimeout)
|
uctx, ucancel := context.WithTimeout(r.Context(), s.cfg.UnloadTimeout)
|
||||||
elapsed, uerr := s.cfg.Ollama.UnloadAll(uctx)
|
elapsed, uerr := s.cfg.Ollama.UnloadAll(uctx)
|
||||||
ucancel()
|
ucancel()
|
||||||
@@ -541,6 +564,7 @@ func (s *Server) handlePrompt(w http.ResponseWriter, r *http.Request) {
|
|||||||
default:
|
default:
|
||||||
log.Info("ollama models unloaded", "seconds", elapsed.Seconds())
|
log.Info("ollama models unloaded", "seconds", elapsed.Seconds())
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
cw := &captureWriter{ResponseWriter: w, status: http.StatusOK, limit: s.captureLimit}
|
cw := &captureWriter{ResponseWriter: w, status: http.StatusOK, limit: s.captureLimit}
|
||||||
s.comfyProxy.ServeHTTP(cw, r)
|
s.comfyProxy.ServeHTTP(cw, r)
|
||||||
@@ -585,7 +609,7 @@ func (s *Server) finishImageJob(promptID string) {
|
|||||||
s.cfg.Lock.ReleaseImage()
|
s.cfg.Lock.ReleaseImage()
|
||||||
log.Info("image lock released")
|
log.Info("image lock released")
|
||||||
|
|
||||||
if s.cfg.WarmModel != "" {
|
if s.cfg.Ollama != nil && s.cfg.WarmModel != "" {
|
||||||
if state, _, _ := s.cfg.Lock.Snapshot(); state == lock.StateIdle {
|
if state, _, _ := s.cfg.Lock.Snapshot(); state == lock.StateIdle {
|
||||||
wctx, wcancel := context.WithTimeout(context.Background(), s.warmTimeout)
|
wctx, wcancel := context.WithTimeout(context.Background(), s.warmTimeout)
|
||||||
if err := s.cfg.Ollama.Warm(wctx, s.cfg.WarmModel); err != nil {
|
if err := s.cfg.Ollama.Warm(wctx, s.cfg.WarmModel); err != nil {
|
||||||
|
|||||||
@@ -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")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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")
|
||||||
@@ -0,0 +1,46 @@
|
|||||||
|
//go:build linux
|
||||||
|
|
||||||
|
package service
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"os"
|
||||||
|
"strconv"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/coreos/go-systemd/v22/daemon"
|
||||||
|
)
|
||||||
|
|
||||||
|
// NotifyReady tells systemd the service is up (Type=notify). It is a no-op
|
||||||
|
// when NOTIFY_SOCKET is unset, e.g. in a container or interactive shell.
|
||||||
|
func NotifyReady() {
|
||||||
|
daemon.SdNotify(false, daemon.SdNotifyReady)
|
||||||
|
}
|
||||||
|
|
||||||
|
// NotifyStopping tells systemd the service is shutting down.
|
||||||
|
func NotifyStopping() {
|
||||||
|
daemon.SdNotify(false, daemon.SdNotifyStopping)
|
||||||
|
}
|
||||||
|
|
||||||
|
// StartWatchdog pings the systemd watchdog every half of WATCHDOG_USEC
|
||||||
|
// until ctx is cancelled. It is a no-op unless systemd started the process
|
||||||
|
// with a watchdog configured (WatchdogSec= in the unit).
|
||||||
|
func StartWatchdog(ctx context.Context) {
|
||||||
|
usec, err := strconv.Atoi(os.Getenv("WATCHDOG_USEC"))
|
||||||
|
if err != nil || usec <= 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
interval := time.Duration(usec) * time.Microsecond / 2
|
||||||
|
go func() {
|
||||||
|
ticker := time.NewTicker(interval)
|
||||||
|
defer ticker.Stop()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return
|
||||||
|
case <-ticker.C:
|
||||||
|
daemon.SdNotify(false, daemon.SdNotifyWatchdog)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
@@ -0,0 +1,14 @@
|
|||||||
|
//go:build !linux
|
||||||
|
|
||||||
|
package service
|
||||||
|
|
||||||
|
import "context"
|
||||||
|
|
||||||
|
// NotifyReady is a no-op outside Linux (no systemd notify socket).
|
||||||
|
func NotifyReady() {}
|
||||||
|
|
||||||
|
// NotifyStopping is a no-op outside Linux.
|
||||||
|
func NotifyStopping() {}
|
||||||
|
|
||||||
|
// StartWatchdog is a no-op outside Linux.
|
||||||
|
func StartWatchdog(context.Context) {}
|
||||||
@@ -0,0 +1,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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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 }
|
||||||
|
|||||||
@@ -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)
|
||||||
@@ -84,6 +153,7 @@ func Install(configPath string) error {
|
|||||||
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
|
||||||
|
}
|
||||||
|
|||||||
@@ -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-----
|
||||||
`
|
`
|
||||||
|
|||||||
Reference in New Issue
Block a user