diff --git a/cmd/gpu-turnstile/main.go b/cmd/gpu-turnstile/main.go index ec40f0a..c38b5b3 100644 --- a/cmd/gpu-turnstile/main.go +++ b/cmd/gpu-turnstile/main.go @@ -5,6 +5,7 @@ package main import ( "bufio" "context" + "encoding/json" "errors" "fmt" "io" @@ -58,7 +59,7 @@ func stdoutIsTerminal() bool { // --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, updateNow, elevatedChild bool, rest []string) { +func parseFlags(args []string) (configPath string, install, remove, noCopy, help, showVersion, forceUpdate, updateNow, monitor, elevatedChild bool, rest []string) { rest = args[:0] for i := 0; i < len(args); i++ { switch { @@ -81,13 +82,15 @@ func parseFlags(args []string) (configPath string, install, remove, noCopy, help forceUpdate = true case args[i] == "--update-now" || args[i] == "-update-now": updateNow = true + case args[i] == "--monitor" || args[i] == "-monitor": + monitor = true case args[i] == "--elevated-child": elevatedChild = true default: rest = append(rest, args[i]) } } - return configPath, install, remove, noCopy, help, showVersion, forceUpdate, updateNow, elevatedChild, rest + return configPath, install, remove, noCopy, help, showVersion, forceUpdate, updateNow, monitor, elevatedChild, rest } // versionLine is printed at the top of every help and error screen. @@ -105,6 +108,8 @@ Usage: (no admin needed when the service runs) gpu-turnstile --update-now like --force-update, but only through the running service + gpu-turnstile --monitor live status view (downstreams, + GPU lock, queue); Ctrl+C quits gpu-turnstile -h | --help this help Options: @@ -135,12 +140,12 @@ func fatalUsage(format string, args ...any) { } func main() { - configPath, install, remove, noCopy, help, showVersion, forceUpdate, updateNow, elevatedChild, args := parseFlags(os.Args[1:]) + configPath, install, remove, noCopy, help, showVersion, forceUpdate, updateNow, monitor, elevatedChild, args := parseFlags(os.Args[1:]) if showVersion { fmt.Println(version) return } - bare := configPath == "" && !install && !remove && !forceUpdate && !updateNow && !elevatedChild && len(args) == 0 + bare := configPath == "" && !install && !remove && !forceUpdate && !updateNow && !monitor && !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 @@ -165,6 +170,8 @@ func main() { fatalUsage("error: --force-update cannot be combined with --install-service/--remove-service") case updateNow && (install || remove || forceUpdate): fatalUsage("error: --update-now cannot be combined with other commands") + case monitor && (install || remove || forceUpdate || updateNow): + fatalUsage("error: --monitor cannot be combined with other commands") case install: os.Exit(serviceCommand(configPath, true, noCopy, elevatedChild)) case remove: @@ -173,6 +180,8 @@ func main() { os.Exit(forceUpdateCommand(configPath, elevatedChild)) case updateNow: os.Exit(updateNowCommand()) + case monitor: + os.Exit(monitorCommand()) } if len(args) > 0 { fatalUsage("error: unknown arguments: %s", strings.Join(args, " ")) @@ -567,6 +576,8 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri ) lk := lock.New(log) + health := newHealthTracker() + started := time.Now() // Each consumer is enabled by setting its URL; a disabled consumer gets // no client, no listener and no probe. var ollamaClient *ollama.Client @@ -676,13 +687,15 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri } probeCtx, probeCancel := context.WithTimeout(ctx, cfg.ProbeTimeout) for name, probe := range probes { - if err := probe(probeCtx); err != nil && !errors.Is(err, errManagedDown) { + err := probe(probeCtx) + health.set(name, err == nil) + if 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) + go healthLoop(ctx, cfg.HealthInterval, cfg.ProbeTimeout, log, probes, health) } // Foreign GPU holders (games, other ML jobs) — enabled by GAME_PROCS @@ -727,17 +740,20 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri service.NotifyReady() service.StartWatchdog(ctx) + var u *update.Updater + var exePath string + applyStaged := func(to string) {} if cfg.AutoUpdate { - exePath, err := os.Executable() - if err != nil { + if p, err := os.Executable(); err != nil { log.Warn("auto-update disabled: cannot locate executable", "err", err) } else { - u := &update.Updater{Repo: cfg.UpdateRepo, Asset: cfg.UpdateAsset, Version: version, Desired: cfg.AppVersion, Log: log} + exePath = p + u = &update.Updater{Repo: cfg.UpdateRepo, Asset: cfg.UpdateAsset, Version: version, Desired: cfg.AppVersion, Log: log} // applyStaged is shared by the hourly loop and the control // channel; the once guard keeps a second trigger from // double-waiting on the GPU lock. var once sync.Once - applyStaged := func(to string) { + applyStaged = func(to string) { if !isService { log.Warn("auto-update: new binary staged; restart gpu-turnstile to apply", "version", to) return @@ -753,12 +769,16 @@ func run(ctx context.Context, cfg config.Config, log *slog.Logger, logOut io.Wri }) } go updateLoop(ctx, cfg.UpdateInterval, log, u, exePath, applyStaged) - if isService { - serveControl(ctx, log, u, exePath, applyStaged) - } } } + // The control channel (status for --monitor, update-now trigger) is + // served whenever running as a service, independent of AUTO_UPDATE. + if isService { + serveControl(ctx, log, u, exePath, applyStaged, + statusProvider(cfg, lk, comfySup, health, started)) + } + select { case err := <-errCh: if err != nil && !errors.Is(err, http.ErrServerClosed) { @@ -858,7 +878,8 @@ func freeVRAM(ctx context.Context, unloadTimeout time.Duration, log *slog.Logger // 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) { +// Every result goes into the tracker for the status channel. +func healthLoop(ctx context.Context, interval, probeTimeout time.Duration, log *slog.Logger, probes map[string]func(context.Context) error, tracker *healthTracker) { ticker := time.NewTicker(interval) defer ticker.Stop() up := map[string]bool{} @@ -875,9 +896,11 @@ func healthLoop(ctx context.Context, interval, probeTimeout time.Duration, log * cancel() if errors.Is(err, errManagedDown) { managed[name] = true // intentionally stopped; not an outage + tracker.set(name, false) continue } now := err == nil + tracker.set(name, now) if managed[name] { // First real probe after an idle stop only re-baselines — // an on-demand start is not a "recovery". @@ -920,17 +943,25 @@ func updateLoop(ctx context.Context, interval time.Duration, log *slog.Logger, u } // serveControl opens the local control channel (named pipe on Windows, -// unix socket on Linux) so unprivileged local users can trigger an update -// check via --force-update without admin rights. The check payload is -// signature-verified regardless of who asks; triggers are rate-limited to -// one per minute so the channel cannot be used to spam restarts. -func serveControl(ctx context.Context, log *slog.Logger, u *update.Updater, exePath string, applyStaged func(to string)) { +// unix socket on Linux) so unprivileged local users can query status +// (--monitor) and trigger an update check (--force-update/--update-now) +// without admin rights. The update payload is signature-verified regardless +// of who asks; triggers are rate-limited to one per minute so the channel +// cannot be used to spam restarts. u is nil when AUTO_UPDATE=false. +func serveControl(ctx context.Context, log *slog.Logger, u *update.Updater, exePath string, applyStaged func(to string), status func() string) { var mu sync.Mutex var lastTrigger time.Time h := func(cmd string) string { - if cmd != control.CmdUpdateNow { + switch cmd { + case control.CmdStatus: + return "OK " + status() + case control.CmdUpdateNow: + default: return "ERR unknown command: " + cmd } + if u == nil { + return "ERR auto-update is disabled on this instance" + } mu.Lock() if wait := time.Minute - time.Since(lastTrigger); wait > 0 { mu.Unlock() @@ -955,6 +986,92 @@ func serveControl(ctx context.Context, log *slog.Logger, u *update.Updater, exeP } } +// healthTracker records the latest probe result per upstream for the +// status channel. +type healthTracker struct { + mu sync.Mutex + up map[string]bool +} + +func newHealthTracker() *healthTracker { + return &healthTracker{up: map[string]bool{}} +} + +func (h *healthTracker) set(name string, up bool) { + h.mu.Lock() + h.up[name] = up + h.mu.Unlock() +} + +func (h *healthTracker) get(name string) bool { + h.mu.Lock() + defer h.mu.Unlock() + return h.up[name] +} + +// statusDownstream/statusLock/statusSnapshot are the JSON the control +// channel serves on CmdStatus; the monitor mode renders them. +type statusDownstream struct { + Name string `json:"name"` + URL string `json:"url"` + Up bool `json:"up"` + Managed string `json:"managed,omitempty"` +} + +type statusLock struct { + State string `json:"state"` + Detail string `json:"detail,omitempty"` + LLMInflight int `json:"llm_inflight"` + LLMWaiting int `json:"llm_waiting"` + ImageQueue int `json:"image_queue"` + External string `json:"external,omitempty"` + SinceS int64 `json:"since_s"` +} + +type statusSnapshot struct { + Version string `json:"version"` + UptimeS int64 `json:"uptime_s"` + Downstreams []statusDownstream `json:"downstreams"` + Lock statusLock `json:"lock"` +} + +// statusProvider assembles the one-line JSON snapshot for CmdStatus. +func statusProvider(cfg config.Config, lk *lock.Lock, comfySup *supervise.Process, health *healthTracker, started time.Time) func() string { + return func() string { + snap := statusSnapshot{ + Version: version, + UptimeS: int64(time.Since(started).Seconds()), + } + if cfg.OllamaURL != "" { + snap.Downstreams = append(snap.Downstreams, statusDownstream{ + Name: "ollama", URL: cfg.OllamaURL, Up: health.get("ollama"), + }) + } + if cfg.ComfyURL != "" { + d := statusDownstream{Name: "comfy", URL: cfg.ComfyURL, Up: health.get("comfy")} + if comfySup != nil { + d.Managed = comfySup.Status() + } + snap.Downstreams = append(snap.Downstreams, d) + } + st := lk.Status() + snap.Lock = statusLock{ + State: string(st.State), + Detail: st.Detail, + LLMInflight: st.LLMInflight, + LLMWaiting: st.LLMWaiting, + ImageQueue: st.ImageQueue, + External: st.External, + SinceS: int64(time.Since(st.Since).Seconds()), + } + b, err := json.Marshal(snap) + if err != nil { + return `{"version":"` + version + `"}` + } + return string(b) + } +} + // 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 { diff --git a/cmd/gpu-turnstile/monitor.go b/cmd/gpu-turnstile/monitor.go new file mode 100644 index 0000000..7de306e --- /dev/null +++ b/cmd/gpu-turnstile/monitor.go @@ -0,0 +1,146 @@ +package main + +import ( + "encoding/json" + "fmt" + "os" + "strings" + "time" + + "gpu-turnstile/internal/control" +) + +// ANSI colors for the monitor frame. +const ( + cReset = "\x1b[0m" + cDim = "\x1b[2m" + cRed = "\x1b[31m" + cGreen = "\x1b[32m" + cYellow = "\x1b[33m" + cCyan = "\x1b[36m" +) + +// monitorCommand renders a live status view of the running service, +// refreshed every second from the control channel. Ctrl+C quits. +func monitorCommand() int { + if !stdoutIsTerminal() { + fmt.Fprintln(os.Stderr, "gpu-turnstile: --monitor needs an interactive terminal") + return 1 + } + enableVirtualTerminal() + fmt.Print("\x1b[2J") // clear once; frames then redraw in place + defer fmt.Print(cReset + "\n") + for { + frame := renderWaiting() + if reply, err := control.Ask(control.CmdStatus); err == nil { + if msg, ok := strings.CutPrefix(reply, "OK "); ok { + var snap statusSnapshot + if json.Unmarshal([]byte(msg), &snap) == nil { + frame = renderMonitor(snap, termWidth()) + } + } + } + fmt.Print("\x1b[H" + frame + "\x1b[J") // home, frame, clear below + time.Sleep(time.Second) + } +} + +func renderWaiting() string { + return cDim + " gpu-turnstile — waiting for a running service…" + cReset + "\x1b[K\n" +} + +// renderMonitor draws one full frame. Each line ends with \x1b[K (clear to +// end of line) so shrinking content leaves no residue. +func renderMonitor(snap statusSnapshot, width int) string { + if width < 40 { + width = 80 + } + var b strings.Builder + + left := " gpu-turnstile" + if snap.UptimeS > 0 { + left += " " + cDim + "up " + fmtDur(snap.UptimeS) + cReset + } + right := snap.Version + pad := width - printableLen(" gpu-turnstile up "+fmtDur(snap.UptimeS)) - len(right) - 1 + if snap.UptimeS == 0 { + pad = width - len(" gpu-turnstile") - len(right) - 1 + } + if pad < 1 { + pad = 1 + } + b.WriteString(cDim + left + strings.Repeat(" ", pad) + right + cReset + "\x1b[K\n") + b.WriteString(cDim + " " + strings.Repeat("─", width-2) + cReset + "\x1b[K\n") + + for _, d := range snap.Downstreams { + b.WriteString(renderDownstream(d) + "\x1b[K\n") + } + b.WriteString("\x1b[K\n") + b.WriteString(renderLock(snap.Lock) + "\x1b[K\n") + if snap.Lock.ImageQueue > 0 { + b.WriteString(fmt.Sprintf(" Queue: %s%d image job(s) waiting%s\x1b[K\n", + cYellow, snap.Lock.ImageQueue, cReset)) + } + return b.String() +} + +// printableLen counts characters without ANSI escapes (ASCII-only content). +func printableLen(s string) int { return len(s) } + +func renderDownstream(d statusDownstream) string { + url := cDim + d.URL + cReset + switch d.Managed { + case "stopped": + return fmt.Sprintf(" %s○%s %-8s %sstopped (managed — starts on demand)%s %s", + cDim, cReset, d.Name, cDim, cReset, url) + case "starting": + return fmt.Sprintf(" %s◌%s %-8s %sstarting…%s %s", + cYellow, cReset, d.Name, cYellow, cReset, url) + } + suffix := "" + if d.Managed == "external" { + suffix = " (external)" + } + if d.Up { + return fmt.Sprintf(" %s●%s %-8s %sUP%s%s %s", cGreen, cReset, d.Name, cGreen, cReset, suffix, url) + } + return fmt.Sprintf(" %s●%s %-8s %sDOWN%s %s", cRed, cReset, d.Name, cRed, cReset, url) +} + +func renderLock(l statusLock) string { + dur := cDim + "(" + fmtDur(l.SinceS) + ")" + cReset + switch l.State { + case "idle": + return fmt.Sprintf(" Lock: %sidle%s %s", cGreen, cReset, dur) + case "llm": + s := fmt.Sprintf(" Lock: %sLLM%s — %d in flight", cCyan, cReset, l.LLMInflight) + if l.LLMWaiting > 0 { + s += fmt.Sprintf(", %d waiting", l.LLMWaiting) + } + if l.Detail != "" { + s += " — " + l.Detail + } + return s + " " + dur + case "image": + s := fmt.Sprintf(" Lock: %sIMAGE%s", cYellow, cReset) + if l.Detail != "" { + s += " — " + l.Detail + } + return s + " " + dur + case "external": + return fmt.Sprintf(" Lock: %sEXTERNAL%s — %s %s", cRed, cReset, l.External, dur) + } + return " Lock: unknown" +} + +// fmtDur renders seconds as a compact duration ("1m32s", "2h07m"). +func fmtDur(s int64) string { + if s < 0 { + s = 0 + } + d := time.Duration(s) * time.Second + if d >= time.Hour { + return fmt.Sprintf("%dh%02dm", int(d.Hours()), int(d.Minutes())%60) + } + return d.Round(time.Second).String() +} diff --git a/cmd/gpu-turnstile/monitor_test.go b/cmd/gpu-turnstile/monitor_test.go new file mode 100644 index 0000000..b6bbecf --- /dev/null +++ b/cmd/gpu-turnstile/monitor_test.go @@ -0,0 +1,41 @@ +package main + +import ( + "strings" + "testing" +) + +func TestRenderMonitor(t *testing.T) { + snap := statusSnapshot{ + Version: "v0.2.3", + UptimeS: 3725, + Downstreams: []statusDownstream{ + {Name: "ollama", URL: "http://127.0.0.1:11435", Up: true}, + {Name: "comfy", URL: "http://127.0.0.1:8189", Managed: "stopped"}, + }, + Lock: statusLock{State: "llm", LLMInflight: 2, LLMWaiting: 1, Detail: "ollama: POST /api/generate", SinceS: 95}, + } + frame := renderMonitor(snap, 80) + for _, want := range []string{"v0.2.3", "1h02m", "ollama", "UP", "comfy", "stopped", "LLM", "2 in flight", "1 waiting", "1m35s"} { + if !strings.Contains(frame, want) { + t.Errorf("frame missing %q:\n%s", want, frame) + } + } + + snap.Lock = statusLock{State: "external", External: "cyberpunk2077.exe (pid 1234)", SinceS: 3, ImageQueue: 2} + frame = renderMonitor(snap, 40) // narrow: falls back to 80 + for _, want := range []string{"EXTERNAL", "cyberpunk2077.exe", "2 image job(s) waiting"} { + if !strings.Contains(frame, want) { + t.Errorf("frame missing %q:\n%s", want, frame) + } + } +} + +func TestFmtDur(t *testing.T) { + cases := map[int64]string{0: "0s", 5: "5s", 95: "1m35s", 3725: "1h02m", -3: "0s"} + for in, want := range cases { + if got := fmtDur(in); got != want { + t.Errorf("fmtDur(%d) = %q, want %q", in, got, want) + } + } +} diff --git a/cmd/gpu-turnstile/monitor_unix.go b/cmd/gpu-turnstile/monitor_unix.go new file mode 100644 index 0000000..4685106 --- /dev/null +++ b/cmd/gpu-turnstile/monitor_unix.go @@ -0,0 +1,21 @@ +//go:build !windows + +package main + +import ( + "os" + + "golang.org/x/sys/unix" +) + +// enableVirtualTerminal is a no-op: unix terminals speak ANSI natively. +func enableVirtualTerminal() {} + +// termWidth is the terminal width in columns, 80 when unknown. +func termWidth() int { + ws, err := unix.IoctlGetWinsize(int(os.Stdout.Fd()), unix.TIOCGWINSZ) + if err != nil || ws.Col == 0 { + return 80 + } + return int(ws.Col) +} diff --git a/cmd/gpu-turnstile/monitor_windows.go b/cmd/gpu-turnstile/monitor_windows.go new file mode 100644 index 0000000..f577d8d --- /dev/null +++ b/cmd/gpu-turnstile/monitor_windows.go @@ -0,0 +1,29 @@ +//go:build windows + +package main + +import ( + "os" + + "golang.org/x/sys/windows" +) + +// enableVirtualTerminal asks the console to honor ANSI escapes (Windows 10+); +// mintty/Git Bash already does, so errors are ignored. +func enableVirtualTerminal() { + h := windows.Handle(os.Stdout.Fd()) + var mode uint32 + if err := windows.GetConsoleMode(h, &mode); err != nil { + return + } + windows.SetConsoleMode(h, mode|windows.ENABLE_VIRTUAL_TERMINAL_PROCESSING) //nolint:errcheck +} + +// termWidth is the console window width in columns, 80 when unknown. +func termWidth() int { + var info windows.ConsoleScreenBufferInfo + if err := windows.GetConsoleScreenBufferInfo(windows.Handle(os.Stdout.Fd()), &info); err != nil { + return 80 + } + return int(info.Window.Right-info.Window.Left) + 1 +} diff --git a/internal/control/control.go b/internal/control/control.go index eea29c8..02226c6 100644 --- a/internal/control/control.go +++ b/internal/control/control.go @@ -23,6 +23,9 @@ import ( // idle) restart onto a signed update immediately. const CmdUpdateNow = "update-now" +// CmdStatus asks for a one-line JSON status snapshot (monitor mode). +const CmdStatus = "status" + // ErrUnavailable means no running service offers the control channel. var ErrUnavailable = errors.New("control channel unavailable") diff --git a/internal/lock/lock.go b/internal/lock/lock.go index 792c22e..5f0a2ea 100644 --- a/internal/lock/lock.go +++ b/internal/lock/lock.go @@ -9,6 +9,7 @@ import ( "context" "log/slog" "sync" + "time" ) // State is the current GPU occupancy state. @@ -30,10 +31,13 @@ type Lock struct { change chan struct{} // closed and replaced on every state change n int // LLM requests in flight + llmWaiting int // LLM requests blocked waiting for the GPU imageActive bool // an image job holds the GPU imageQ []imageWaiter nextID uint64 - external string // non-empty: a foreign process (e.g. a game) holds the GPU + external string // non-empty: a foreign process (e.g. a game) holds the GPU + detail string // what the current holder is doing (best effort) + since time.Time // when the current state began log *slog.Logger } @@ -41,7 +45,7 @@ type Lock struct { // New returns a ready-to-use Lock. log may be nil; if set, every state // transition is logged at debug level. func New(log *slog.Logger) *Lock { - return &Lock{change: make(chan struct{}), log: log} + return &Lock{change: make(chan struct{}), log: log, since: time.Now()} } // broadcast wakes all waiters. Call with mu held. @@ -63,6 +67,7 @@ func (l *Lock) logTransition(msg string, args ...any) { func (l *Lock) SetExternal(holder string) { l.mu.Lock() l.external = holder + l.since = time.Now() l.broadcast() l.mu.Unlock() l.logTransition("lock transition", "state", StateExternal, "holder", holder) @@ -73,6 +78,7 @@ func (l *Lock) SetExternal(holder string) { func (l *Lock) ClearExternal() { l.mu.Lock() l.external = "" + l.since = time.Now() l.broadcast() l.mu.Unlock() l.logTransition("lock transition", "state", StateIdle) @@ -90,17 +96,31 @@ func (l *Lock) External() string { // while waiting; no state is changed in that case. func (l *Lock) AcquireLLM(ctx context.Context) error { l.mu.Lock() + waiting := false for l.imageActive || len(l.imageQ) > 0 || l.external != "" { + if !waiting { + l.llmWaiting++ + waiting = true + } ch := l.change l.mu.Unlock() select { case <-ctx.Done(): + l.mu.Lock() + l.llmWaiting-- + l.mu.Unlock() return ctx.Err() case <-ch: } l.mu.Lock() } + if waiting { + l.llmWaiting-- + } l.n++ + if l.n == 1 { + l.since = time.Now() + } n := l.n l.mu.Unlock() l.logTransition("lock transition", "state", StateLLM, "llm_inflight", n) @@ -129,6 +149,8 @@ func (l *Lock) ReleaseLLM() { l.n-- n := l.n if l.n == 0 { + l.since = time.Now() + l.detail = "" l.broadcast() } l.mu.Unlock() @@ -156,6 +178,7 @@ func (l *Lock) AcquireImage(ctx context.Context) error { if l.imageQ[0].id == w.id && l.n == 0 && !l.imageActive && l.external == "" { l.imageQ = l.imageQ[1:] l.imageActive = true + l.since = time.Now() l.mu.Unlock() l.logTransition("lock transition", "state", StateImage) return nil @@ -184,6 +207,8 @@ func (l *Lock) AcquireImage(ctx context.Context) error { func (l *Lock) ReleaseImage() { l.mu.Lock() l.imageActive = false + l.since = time.Now() + l.detail = "" l.broadcast() l.mu.Unlock() l.logTransition("lock transition", "state", StateIdle) @@ -206,3 +231,51 @@ func (l *Lock) Snapshot() (state State, llmInflight int, imagePending bool) { } return state, l.n, l.imageActive || len(l.imageQ) > 0 } + +// SetDetail records what the current holder is doing (e.g. the request +// path), for status displays. Best effort: overwritten by each new holder, +// cleared when the GPU goes idle. +func (l *Lock) SetDetail(detail string) { + l.mu.Lock() + l.detail = detail + l.mu.Unlock() +} + +// Status is a point-in-time view of the lock for monitoring. +type Status struct { + State State + Detail string + LLMInflight int + LLMWaiting int + ImageActive bool + ImageQueue int + External string + Since time.Time +} + +// Status reports the full lock state, including waiters and how long the +// current state has held. +func (l *Lock) Status() Status { + l.mu.Lock() + defer l.mu.Unlock() + s := Status{ + Detail: l.detail, + LLMInflight: l.n, + LLMWaiting: l.llmWaiting, + ImageActive: l.imageActive, + ImageQueue: len(l.imageQ), + External: l.external, + Since: l.since, + } + switch { + case l.imageActive: + s.State = StateImage + case l.n > 0: + s.State = StateLLM + case l.external != "": + s.State = StateExternal + default: + s.State = StateIdle + } + return s +} diff --git a/internal/proxy/proxy.go b/internal/proxy/proxy.go index f481921..8dcdc79 100644 --- a/internal/proxy/proxy.go +++ b/internal/proxy/proxy.go @@ -491,6 +491,7 @@ func (s *Server) OllamaHandler() http.Handler { return } defer s.cfg.Lock.ReleaseLLM() + s.cfg.Lock.SetDetail("ollama: " + r.Method + " " + r.URL.Path) s.ollamaProxy.ServeHTTP(w, r) })) } @@ -590,6 +591,7 @@ func (s *Server) handlePrompt(w http.ResponseWriter, r *http.Request) { } s.cfg.Metrics.ObserveLockWait("image", time.Since(start).Seconds()) log.Info("image lock acquired") + s.cfg.Lock.SetDetail("comfy: POST /prompt") if s.cfg.ComfySup != nil && !comfyFirst { if err := s.cfg.ComfySup.EnsureRunning(); err != nil { diff --git a/internal/supervise/supervise.go b/internal/supervise/supervise.go index 1928ac4..a19fb7b 100644 --- a/internal/supervise/supervise.go +++ b/internal/supervise/supervise.go @@ -164,6 +164,24 @@ func (p *Process) Running() bool { return p.cmd != nil } +// Status describes the child for status displays: "external" when something +// else serves the port, "running" once ready, "starting" while the child +// boots, "stopped" otherwise. +func (p *Process) Status() string { + p.mu.Lock() + defer p.mu.Unlock() + switch { + case p.external: + return "external" + case p.cmd == nil: + return "stopped" + case p.ready: + return "running" + default: + return "starting" + } +} + // Ready reports whether the server has answered a probe since its last // (re)start. Health checks use it to tell "starting up" from "outage". func (p *Process) Ready() bool {