Files
gpu-turnstile/cmd/gpu-turnstile/main.go
T

851 lines
28 KiB
Go

// gpu-turnstile is a GPU arbitration proxy that sits in front of Ollama and
// ComfyUI and guarantees only one of them uses the GPU at a time.
package main
import (
"bufio"
"context"
"errors"
"fmt"
"io"
"io/fs"
"log/slog"
"net"
"net/http"
"os"
"os/signal"
"path/filepath"
"runtime"
"strings"
"syscall"
"time"
"gpu-turnstile/internal/comfy"
"gpu-turnstile/internal/config"
"gpu-turnstile/internal/game"
"gpu-turnstile/internal/lock"
"gpu-turnstile/internal/metrics"
"gpu-turnstile/internal/ollama"
"gpu-turnstile/internal/proxy"
"gpu-turnstile/internal/service"
"gpu-turnstile/internal/supervise"
"gpu-turnstile/internal/update"
)
// version is injected at build time via -ldflags "-X main.version=...".
var version = "dev"
// exitCodeUpdate tells the service recovery configuration to restart the
// process: a signed update has been staged and the GPU lock is idle.
const exitCodeUpdate = 3
// 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() {
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" {
// 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 {
fatalUsage("error: unknown arguments: %s", strings.Join(args, " "))
}
if exePath, err := os.Executable(); err == nil {
update.CleanupOld(exePath)
}
cfg, err := loadMergedConfig(configPath)
if err != nil {
fmt.Fprintf(os.Stderr, "%s\n\ngpu-turnstile: %v\n", versionLine(), err)
os.Exit(1)
}
log, logOut, logCloser := newLogger(cfg)
defer logCloser.Close()
syncEnvFile(resolveConfigPath(configPath), cfg.LogFile, log)
if service.IsService() {
if err := service.Run(func(ctx context.Context) error { return run(ctx, cfg, log, logOut, true) }); err != nil {
log.Error("service failed", "err", err)
os.Exit(1)
}
return
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
if err := run(ctx, cfg, log, logOut, false); err != nil {
log.Error("listener failed", "err", err)
os.Exit(1)
}
}
// syncEnvFile upgrades an installer-written config file after an update:
// settings added since its CFG_VER are appended (commented out) and CFG_VER
// is bumped. Files not written by the installer (no CFG_VER), up-to-date
// files and dev builds are left untouched; a write failure is logged, not
// fatal.
func syncEnvFile(path, logFile string, log *slog.Logger) {
data, err := os.ReadFile(path)
if err != nil {
return // no config file; nothing to upgrade
}
synced, changed := config.SyncSample(string(data), version, logFile)
if !changed {
return
}
if err := os.WriteFile(path, []byte(synced), 0o644); err != nil {
log.Warn("could not append new settings to the config file", "path", path, "err", err)
return
}
log.Warn("config file updated: new settings appended", "path", path, "version", version)
}
// defaultConfigPath returns gpu-turnstile.env next to the executable.
func defaultConfigPath() string {
exe, err := os.Executable()
if err != nil {
return "gpu-turnstile.env"
}
return filepath.Join(filepath.Dir(exe), "gpu-turnstile.env")
}
// resolveConfigPath applies the precedence: -config flag, then
// GPU_TURNSTILE_CONFIG, then the default next to the executable.
func resolveConfigPath(flagValue string) string {
if flagValue != "" {
return flagValue
}
if v := os.Getenv("GPU_TURNSTILE_CONFIG"); v != "" {
return v
}
return defaultConfigPath()
}
// loadMergedConfig reads the config file (if present) and overlays process
// environment variables on top. A missing file is fine; an unreadable or
// malformed file is fatal.
func loadMergedConfig(flagValue string) (config.Config, error) {
path := resolveConfigPath(flagValue)
values := map[string]string{}
if f, err := os.Open(path); err == nil {
defer f.Close()
parsed, err := config.ParseEnvFile(f)
if err != nil {
return config.Config{}, fmt.Errorf("%s: %w", path, err)
}
values = parsed
} else if !errors.Is(err, os.ErrNotExist) {
return config.Config{}, fmt.Errorf("read config file: %w", err)
}
getenv := func(key string) string {
if v := os.Getenv(key); v != "" {
return v
}
return values[key]
}
return config.Load(getenv)
}
// newLogger builds the slog logger and returns the output writer (stderr or
// the opened LOG_FILE) plus a closer for it.
func newLogger(cfg config.Config) (*slog.Logger, io.Writer, io.Closer) {
out := io.Writer(os.Stderr)
closer := io.NopCloser(nil)
if cfg.LogFile != "" {
if f, err := os.OpenFile(cfg.LogFile, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o644); err == nil {
out = f
closer = f
} else {
fmt.Fprintf(os.Stderr, "gpu-turnstile: cannot open LOG_FILE %s: %v (logging to stderr)\n", cfg.LogFile, err)
}
}
opts := &slog.HandlerOptions{Level: cfg.LogLevel}
var handler slog.Handler = slog.NewTextHandler(out, opts)
if cfg.LogJSON {
handler = slog.NewJSONHandler(out, opts)
}
log := slog.New(handler)
slog.SetDefault(log)
return log, out, closer
}
// waitForEnter keeps an elevated child's console window open until the
// user has read the output.
func waitForEnter() {
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 errors.Is(err, config.ErrNoConsumer) {
err = nil // update settings do not depend on a consumer URL
}
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()
if cfg.AppVersion == "dev" {
fmt.Printf("%s: APP_VER=dev, updates disabled\n", versionLine())
return 0
}
u := &update.Updater{Repo: cfg.UpdateRepo, Asset: cfg.UpdateAsset, Version: version, Desired: cfg.AppVersion, Log: log}
// Single-shot: one attempt, fail fast when the server is unreachable
// instead of hanging in a TCP connect for minutes.
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
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
if install {
path := resolveConfigPath(configPath)
if abs, absErr := filepath.Abs(path); absErr == nil {
path = abs
}
err = service.Install(path, !noCopy, version)
} else {
err = service.Remove()
}
if err != nil {
fmt.Fprintf(os.Stderr, "gpu-turnstile service %s: %v\n", verb, err)
return 1
}
fmt.Printf("service %s: %s\n", service.Name, doneVerb)
return 0
}
// orDisabled renders an empty URL as "disabled" for the startup dump.
func orDisabled(url string) string {
if url == "" {
return "disabled"
}
return url
}
// managedComfyCommand resolves how ComfyUI is launched when it is managed:
// COMFY_CMD verbatim, or the standard venv layout under COMFY_DIR. Empty
// when neither is set (unmanaged).
func managedComfyCommand(cfg config.Config) string {
if cfg.ComfyCmd != "" {
return cfg.ComfyCmd
}
if cfg.ComfyDir != "" {
return supervise.DefaultComfyCommand(runtime.GOOS, cfg.ComfyDir, cfg.ComfyURL)
}
return ""
}
func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Writer, isService bool) error {
// ComfyUI can run as a managed child — COMFY_CMD verbatim, or the
// standard venv layout derived from COMFY_DIR alone: started on demand
// by the proxy, stopped after COMFY_IDLE_TIMEOUT idle (and on shutdown)
// so its VRAM is freed.
comfyCmdLine := managedComfyCommand(cfg)
// The startup line carries the version and every setting and is emitted
// at WARN so it is visible even with the default (quiet) log level.
log.Log(ctx, slog.LevelWarn, "starting gpu-turnstile",
"version", version,
"listen_ollama", cfg.ListenOllama,
"listen_comfy", cfg.ListenComfy,
"ollama_url", orDisabled(cfg.OllamaURL),
"comfy_url", orDisabled(cfg.ComfyURL),
"unload_timeout", cfg.UnloadTimeout,
"job_timeout", cfg.JobTimeout,
"llm_wait_timeout", cfg.LLMWaitTimeout,
"llm_busy_mode", cfg.LLMBusyMode,
"llm_busy_status", cfg.LLMBusyStatus,
"busy_retry_after", cfg.BusyRetryAfter,
"unload_poll_interval", cfg.UnloadPollInterval,
"history_poll_interval", cfg.HistoryPollInterval,
"probe_timeout", cfg.ProbeTimeout,
"health_interval", cfg.HealthInterval,
"free_timeout", cfg.FreeTimeout,
"warm_timeout", cfg.WarmTimeout,
"shutdown_timeout", cfg.ShutdownTimeout,
"backoff_initial", cfg.BackoffInitial,
"backoff_max", cfg.BackoffMax,
"prompt_capture_limit", cfg.PromptCaptureLimit,
"warm_model", cfg.WarmModel,
"comfy_cmd", orDisabled(comfyCmdLine),
"comfy_dir", cfg.ComfyDir,
"comfy_idle_timeout", cfg.ComfyIdleTimeout,
"comfy_start_timeout", cfg.ComfyStartTimeout,
"game_procs", cfg.GameProcs,
"gpu_foreign_vram_mb", cfg.GPUForeignVRAMMB,
"gpu_ignore_procs", cfg.GPUIgnoreProcs,
"game_poll_interval", cfg.GamePollInterval,
"auto_update", cfg.AutoUpdate,
"update_interval", cfg.UpdateInterval,
"update_repo", cfg.UpdateRepo,
"update_asset", cfg.UpdateAsset,
"log_level", cfg.LogLevel,
"log_format", map[bool]string{true: "json", false: "text"}[cfg.LogJSON],
"log_file", cfg.LogFile,
)
lk := lock.New(log)
// Each consumer is enabled by setting its URL; a disabled consumer gets
// no client, no listener and no probe.
var ollamaClient *ollama.Client
var err error
if cfg.OllamaURL != "" {
if ollamaClient, err = ollama.New(cfg.OllamaURL, log); err != nil {
return err
}
}
var comfyClient *comfy.Client
if cfg.ComfyURL != "" {
if comfyClient, err = comfy.New(cfg.ComfyURL, log); err != nil {
return err
}
}
var comfySup *supervise.Process
if comfyCmdLine != "" {
if cfg.ComfyCmd == "" {
// Derived from COMFY_DIR: flag a wrong-looking layout early,
// while the operator is still watching the startup log.
python, script := supervise.ComfyLayout(runtime.GOOS, cfg.ComfyDir)
for _, p := range []string{python, script} {
if _, err := os.Stat(p); err != nil {
log.Warn("COMFY_DIR: file not found; ComfyUI requests will fail until it exists", "path", p)
}
}
}
var err error
comfySup, err = supervise.New("comfy", comfyCmdLine, cfg.ComfyDir, comfyClient.Probe, cfg.ComfyStartTimeout, log)
if err != nil {
return err
}
defer comfySup.Stop()
gpuIdle := func() bool {
state, _, pending := lk.Snapshot()
return state == lock.StateIdle && !pending
}
go comfySup.WatchIdle(ctx, cfg.ComfyIdleTimeout, gpuIdle)
}
srv, err := proxy.New(proxy.Config{
OllamaURL: cfg.OllamaURL,
ComfyURL: cfg.ComfyURL,
Lock: lk,
Ollama: ollamaClient,
Comfy: comfyClient,
ComfySup: comfySup,
Metrics: metrics.New(),
Log: log,
LogColor: !cfg.LogJSON && cfg.LogFile == "" && os.Getenv("NO_COLOR") == "",
LogWriter: logOut,
LLMWaitTimeout: cfg.LLMWaitTimeout,
UnloadTimeout: cfg.UnloadTimeout,
JobTimeout: cfg.JobTimeout,
UnloadPollInterval: cfg.UnloadPollInterval,
HistoryPollInterval: cfg.HistoryPollInterval,
FreeTimeout: cfg.FreeTimeout,
WarmTimeout: cfg.WarmTimeout,
LLMBusyMode: cfg.LLMBusyMode,
LLMBusyStatus: cfg.LLMBusyStatus,
BusyRetryAfter: cfg.BusyRetryAfter,
BackoffInitial: cfg.BackoffInitial,
BackoffMax: cfg.BackoffMax,
PromptCaptureLimit: cfg.PromptCaptureLimit,
WarmModel: cfg.WarmModel,
})
if err != nil {
return err
}
// Probe the enabled upstreams once; failure is logged, not fatal. A
// managed ComfyUI is intentionally down at startup — the first request
// starts it — so neither the startup probe nor the health check treats
// that as an outage.
probes := map[string]func(context.Context) error{}
if ollamaClient != nil {
probes["ollama"] = ollamaClient.Probe
}
if comfyClient != nil {
if comfySup == nil {
probes["comfy"] = comfyClient.Probe
} else {
// Managed upstream: an idle-stopped or still-starting server is
// not an outage, so it is skipped until it has answered once
// (Ready resets on every spawn/stop). After that, a failed
// probe while the process lives is a real "DOWN".
probes["comfy"] = func(ctx context.Context) error {
if !comfySup.Ready() {
if comfySup.Running() {
if err := comfyClient.Probe(ctx); err == nil {
comfySup.MarkReady()
}
}
return errManagedDown
}
return comfyClient.Probe(ctx)
}
}
}
probeCtx, probeCancel := context.WithTimeout(ctx, cfg.ProbeTimeout)
for name, probe := range probes {
if err := probe(probeCtx); err != nil && !errors.Is(err, errManagedDown) {
log.Warn(name+" probe failed", "err", err)
}
}
probeCancel()
if cfg.HealthInterval > 0 {
go healthLoop(ctx, cfg.HealthInterval, cfg.ProbeTimeout, log, probes)
}
// Foreign GPU holders (games, other ML jobs) — enabled by GAME_PROCS
// and/or GPU_FOREIGN_VRAM_MB — hold the lock externally while they run.
if len(cfg.GameProcs) > 0 || cfg.GPUForeignVRAMMB > 0 {
det := game.New(cfg.GameProcs, cfg.GPUForeignVRAMMB, cfg.GPUIgnoreProcs, log)
go gameLoop(ctx, cfg, log, det, lk, ollamaClient, comfySup)
}
// Bind the listeners up front so a port conflict fails fast and the
// readiness notification below really means "accepting connections".
var servers []*http.Server
var listeners []net.Listener
bind := func(addr string, handler http.Handler, consumer string) error {
ln, err := net.Listen("tcp", addr)
if err != nil {
return fmt.Errorf("listen %s on %s: %w", consumer, addr, err)
}
servers = append(servers, &http.Server{Addr: addr, Handler: handler})
listeners = append(listeners, ln)
log.Warn("listening", "consumer", consumer, "addr", addr)
return nil
}
if ollamaClient != nil {
if err := bind(cfg.ListenOllama, srv.OllamaHandler(), "ollama"); err != nil {
return err
}
}
if comfyClient != nil {
if err := bind(cfg.ListenComfy, srv.ComfyHandler(), "comfy"); err != nil {
return err
}
}
errCh := make(chan error, len(servers))
for i := range servers {
go func(s *http.Server, ln net.Listener) { errCh <- s.Serve(ln) }(servers[i], listeners[i])
}
// Tell systemd we are up and start the watchdog pings; both are no-ops
// when not running under a notify/watchdog unit.
service.NotifyReady()
service.StartWatchdog(ctx)
if cfg.AutoUpdate {
go updateLoop(ctx, cfg, log, lk, isService)
}
select {
case err := <-errCh:
if err != nil && !errors.Is(err, http.ErrServerClosed) {
return err
}
case <-ctx.Done():
log.Info("shutting down")
}
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), cfg.ShutdownTimeout)
defer shutdownCancel()
for _, s := range servers {
s.Shutdown(shutdownCtx)
}
return nil
}
// errManagedDown marks a managed upstream that is intentionally stopped
// (idle); health checks skip it instead of logging an outage.
var errManagedDown = errors.New("managed upstream intentionally stopped")
// gameLoop polls for foreign GPU holders (a game, another ML job). While one
// is detected it holds the lock externally so new LLM and image requests
// wait (or are rejected per LLM_BUSY_MODE), and — once in-flight work has
// drained — frees VRAM for it: the managed ComfyUI is stopped and Ollama's
// resident models are unloaded.
func gameLoop(ctx context.Context, cfg config.Config, log *slog.Logger, det *game.Detector, lk *lock.Lock, ollamaClient *ollama.Client, comfySup *supervise.Process) {
ticker := time.NewTicker(cfg.GamePollInterval)
defer ticker.Stop()
held, freed := false, false
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
}
holders, err := det.Check(ctx)
if err != nil && ctx.Err() == nil {
log.Warn("game detection failed", "err", err)
}
switch {
case len(holders) > 0 && !held:
held = true
lk.SetExternal(summarizeHolders(holders))
log.Warn("GPU held by an external process; new LLM/image requests wait",
"holders", summarizeHolders(holders))
case len(holders) == 0 && held:
held, freed = false, false
lk.ClearExternal()
log.Warn("external process released the GPU; resuming")
}
if held && !freed {
if state, _, _ := lk.Snapshot(); state != lock.StateLLM && state != lock.StateImage {
freed = true
freeVRAM(ctx, cfg.UnloadTimeout, log, ollamaClient, comfySup)
}
}
}
}
// summarizeHolders joins holder descriptions for logs and busy responses,
// capping the list so a process name matching dozens of PIDs (system
// services) does not flood the log.
func summarizeHolders(holders []string) string {
const max = 5
if len(holders) > max {
return strings.Join(holders[:max], "; ") + fmt.Sprintf("; +%d more", len(holders)-max)
}
return strings.Join(holders, "; ")
}
// freeVRAM stops the managed ComfyUI (never an external server on its port)
// and unloads Ollama's resident models so the foreign process gets the GPU
// memory.
func freeVRAM(ctx context.Context, unloadTimeout time.Duration, log *slog.Logger, ollamaClient *ollama.Client, comfySup *supervise.Process) {
if comfySup != nil && comfySup.Running() {
log.Warn("stopping the managed ComfyUI to free VRAM")
comfySup.Stop()
}
if ollamaClient == nil {
return
}
uctx, cancel := context.WithTimeout(ctx, unloadTimeout)
defer cancel()
if models, err := ollamaClient.LoadedModels(uctx); err != nil || len(models) == 0 {
return // nothing resident (or ollama unreachable); nothing to free
}
elapsed, err := ollamaClient.UnloadAll(uctx)
if err != nil {
log.Warn("ollama unload incomplete; continuing", "err", err)
return
}
log.Warn("ollama models unloaded to free VRAM", "seconds", elapsed.Seconds())
}
// healthLoop probes the enabled upstreams every interval and logs status
// transitions — "is DOWN" when a previously healthy upstream stops
// answering, "recovered" when it comes back. The first round only
// establishes the baseline; the startup probe already reported that state.
func healthLoop(ctx context.Context, interval, probeTimeout time.Duration, log *slog.Logger, probes map[string]func(context.Context) error) {
ticker := time.NewTicker(interval)
defer ticker.Stop()
up := map[string]bool{}
managed := map[string]bool{}
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
}
for name, probe := range probes {
pctx, cancel := context.WithTimeout(ctx, probeTimeout)
err := probe(pctx)
cancel()
if errors.Is(err, errManagedDown) {
managed[name] = true // intentionally stopped; not an outage
continue
}
now := err == nil
if managed[name] {
// First real probe after an idle stop only re-baselines —
// an on-demand start is not a "recovery".
managed[name] = false
up[name] = now
continue
}
was, seen := up[name]
if seen && now != was {
if now {
log.Warn(name + " upstream recovered")
} else {
log.Warn(name+" upstream is DOWN", "err", err)
}
}
up[name] = now
}
}
}
// updateLoop checks for signed updates on startup and every UPDATE_INTERVAL.
// In service mode a staged update is applied by exiting with exitCodeUpdate
// once the GPU lock is idle; the service recovery configuration restarts the
// process with the new binary. Interactively it only logs.
func updateLoop(ctx context.Context, cfg config.Config, log *slog.Logger, lk *lock.Lock, isService bool) {
exePath, err := os.Executable()
if err != nil {
log.Warn("auto-update disabled: cannot locate executable", "err", err)
return
}
u := &update.Updater{Repo: cfg.UpdateRepo, Asset: cfg.UpdateAsset, Version: version, Desired: cfg.AppVersion, Log: log}
for {
staged, err := u.Check(ctx, exePath)
if err != nil && ctx.Err() == nil {
log.Warn("auto-update check failed", "err", err)
}
if staged {
if !isService {
log.Warn("auto-update: new binary staged; restart gpu-turnstile to apply")
return
}
log.Warn("auto-update: staged; restarting once the GPU is idle")
if waitForIdle(ctx, lk, 24*time.Hour) {
log.Warn("auto-update: restarting to apply update")
os.Exit(exitCodeUpdate)
}
return
}
select {
case <-ctx.Done():
return
case <-time.After(cfg.UpdateInterval):
}
}
}
// waitForIdle polls the lock until no LLM or image work is active or
// pending, max at most. Returns false on timeout or cancellation.
func waitForIdle(ctx context.Context, lk *lock.Lock, max time.Duration) bool {
deadline := time.Now().Add(max)
for {
if state, _, pending := lk.Snapshot(); state == lock.StateIdle && !pending {
return true
}
if time.Now().After(deadline) {
return false
}
select {
case <-ctx.Done():
return false
case <-time.After(5 * time.Second):
}
}
}