Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
cffb4151a6 | ||
|
|
16dc7fe462 | ||
|
|
a74cd49fe5 | ||
|
|
cc2a2cad27 | ||
|
|
f707d07fd8 | ||
|
|
ca33a3db82 | ||
|
|
507af1dbff | ||
|
|
45bc9fe27c |
+235
-29
@@ -5,6 +5,7 @@ package main
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -17,11 +18,13 @@ import (
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
"strings"
|
||||
"sync"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"gpu-turnstile/internal/comfy"
|
||||
"gpu-turnstile/internal/config"
|
||||
"gpu-turnstile/internal/control"
|
||||
"gpu-turnstile/internal/game"
|
||||
"gpu-turnstile/internal/lock"
|
||||
"gpu-turnstile/internal/metrics"
|
||||
@@ -56,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, 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 {
|
||||
@@ -77,13 +80,17 @@ func parseFlags(args []string) (configPath string, install, remove, noCopy, help
|
||||
showVersion = true
|
||||
case args[i] == "--force-update" || args[i] == "-force-update":
|
||||
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, 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.
|
||||
@@ -98,6 +105,11 @@ Usage:
|
||||
gpu-turnstile -v | --version print just the version
|
||||
gpu-turnstile --force-update check for a signed update now,
|
||||
apply it and restart the service
|
||||
(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:
|
||||
@@ -128,12 +140,12 @@ func fatalUsage(format string, args ...any) {
|
||||
}
|
||||
|
||||
func main() {
|
||||
configPath, install, remove, noCopy, help, showVersion, forceUpdate, 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 && !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
|
||||
@@ -156,12 +168,20 @@ func main() {
|
||||
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 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:
|
||||
os.Exit(serviceCommand(configPath, false, noCopy, elevatedChild))
|
||||
case forceUpdate:
|
||||
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, " "))
|
||||
@@ -347,6 +367,13 @@ func forceUpdateCommand(configPath string, elevatedChild bool) int {
|
||||
fmt.Printf("%s: APP_VER=dev, updates disabled\n", versionLine())
|
||||
return 0
|
||||
}
|
||||
// A running service can do the privileged work itself (its account owns
|
||||
// the install dir and it knows when the GPU is idle): ask it over the
|
||||
// local control channel first, no admin rights needed. Fails fast when
|
||||
// no service is listening, in which case we do the direct check below.
|
||||
if reply, err := control.Ask(control.CmdUpdateNow); err == nil {
|
||||
return printControlReply(reply)
|
||||
}
|
||||
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
|
||||
@@ -397,6 +424,28 @@ func forceUpdateCommand(configPath string, elevatedChild bool) int {
|
||||
return 0
|
||||
}
|
||||
|
||||
// printControlReply prints the service's answer to a control-channel
|
||||
// request: "OK ..." on stdout (exit 0), "ERR ..." on stderr (exit 1).
|
||||
func printControlReply(reply string) int {
|
||||
if msg, ok := strings.CutPrefix(reply, "ERR "); ok {
|
||||
fmt.Fprintf(os.Stderr, "gpu-turnstile: %s\n", msg)
|
||||
return 1
|
||||
}
|
||||
fmt.Println("gpu-turnstile: " + strings.TrimPrefix(reply, "OK "))
|
||||
return 0
|
||||
}
|
||||
|
||||
// updateNowCommand only goes through the running service's control channel
|
||||
// (no direct check, no elevation): the unprivileged update trigger.
|
||||
func updateNowCommand() int {
|
||||
reply, err := control.Ask(control.CmdUpdateNow)
|
||||
if err != nil {
|
||||
fmt.Fprintf(os.Stderr, "%s\n\ngpu-turnstile: no running service to ask — use --force-update for a direct check\n", versionLine())
|
||||
return 1
|
||||
}
|
||||
return printControlReply(reply)
|
||||
}
|
||||
|
||||
// reportElevatedUpdate prints the parent's summary of an elevated
|
||||
// --force-update child: exitCodeStaged means the child staged a new binary,
|
||||
// 0 means it found nothing to do. to is the tag the parent's own check
|
||||
@@ -527,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
|
||||
@@ -636,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
|
||||
@@ -687,8 +740,43 @@ 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 {
|
||||
go updateLoop(ctx, cfg, log, lk, isService)
|
||||
if p, err := os.Executable(); err != nil {
|
||||
log.Warn("auto-update disabled: cannot locate executable", "err", err)
|
||||
} else {
|
||||
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) {
|
||||
if !isService {
|
||||
log.Warn("auto-update: new binary staged; restart gpu-turnstile to apply", "version", to)
|
||||
return
|
||||
}
|
||||
log.Warn("auto-update: staged; restarting once the GPU is idle", "version", to)
|
||||
once.Do(func() {
|
||||
go func() {
|
||||
if waitForIdle(ctx, lk, 24*time.Hour) {
|
||||
log.Warn("auto-update: restarting to apply update")
|
||||
os.Exit(exitCodeUpdate)
|
||||
}
|
||||
}()
|
||||
})
|
||||
}
|
||||
go updateLoop(ctx, cfg.UpdateInterval, 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 {
|
||||
@@ -790,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{}
|
||||
@@ -807,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".
|
||||
@@ -830,42 +921,157 @@ func healthLoop(ctx context.Context, interval, probeTimeout time.Duration, log *
|
||||
}
|
||||
}
|
||||
|
||||
// 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}
|
||||
// updateLoop checks for signed updates on startup and every interval;
|
||||
// applyStaged decides what a staged update means (restart when idle as a
|
||||
// service, log only interactively).
|
||||
func updateLoop(ctx context.Context, interval time.Duration, log *slog.Logger, u *update.Updater, exePath string, applyStaged func(to string)) {
|
||||
for {
|
||||
staged, to, 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", "version", to)
|
||||
return
|
||||
}
|
||||
log.Warn("auto-update: staged; restarting once the GPU is idle", "version", to)
|
||||
if waitForIdle(ctx, lk, 24*time.Hour) {
|
||||
log.Warn("auto-update: restarting to apply update")
|
||||
os.Exit(exitCodeUpdate)
|
||||
}
|
||||
applyStaged(to)
|
||||
return
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-time.After(cfg.UpdateInterval):
|
||||
case <-time.After(interval):
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// serveControl opens the local control channel (named pipe on Windows,
|
||||
// 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 {
|
||||
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()
|
||||
return fmt.Sprintf("ERR rate limited: retry in %ds", int(wait.Seconds())+1)
|
||||
}
|
||||
lastTrigger = time.Now()
|
||||
mu.Unlock()
|
||||
cctx, cancel := context.WithTimeout(ctx, 30*time.Second)
|
||||
defer cancel()
|
||||
staged, to, err := u.Check(cctx, exePath)
|
||||
if err != nil {
|
||||
return "ERR update check failed: " + err.Error()
|
||||
}
|
||||
if !staged {
|
||||
return "OK " + version + " is up to date"
|
||||
}
|
||||
applyStaged(to)
|
||||
return "OK updated from " + version + " to " + to + "; the service restarts once the GPU is idle"
|
||||
}
|
||||
if err := control.Serve(ctx, h, log); err != nil {
|
||||
log.Warn("control channel disabled", "err", err)
|
||||
}
|
||||
}
|
||||
|
||||
// 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 {
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -6,7 +6,7 @@
|
||||
# ComfyUI --listen 0.0.0.0 --port 8189).
|
||||
services:
|
||||
gpu-turnstile:
|
||||
image: git.rambossek.at/public/gpu-turnstile:v0.2.1
|
||||
image: git.rambossek.at/public/gpu-turnstile:v0.2.4
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
# Each consumer is enabled by setting its URL; leave one unset to
|
||||
|
||||
@@ -0,0 +1,64 @@
|
||||
// Package control exposes a local-only command channel into a running
|
||||
// gpu-turnstile service: a named pipe on Windows, a unix socket on Linux.
|
||||
// It lets unprivileged local users ask the service to do privileged work
|
||||
// that is safe to offer — currently triggering an update check, whose
|
||||
// payload is signature-verified regardless of who asks. The channel never
|
||||
// accepts data beyond a one-word command, and the server rate-limits
|
||||
// triggers, so the worst a local user can cause is a cheap, throttled
|
||||
// check and a GPU-idle-gated restart onto a signed binary.
|
||||
//
|
||||
// Protocol: the client writes one command line, the server answers with
|
||||
// one reply line ("OK ..." or "ERR ...") and hangs up.
|
||||
package control
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// CmdUpdateNow asks the service to check for, stage and (once the GPU is
|
||||
// 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")
|
||||
|
||||
// Handler answers one command; the returned string is sent back as one
|
||||
// line. It must start with "OK " or "ERR ".
|
||||
type Handler func(cmd string) string
|
||||
|
||||
// serveConn runs the line protocol on one accepted connection.
|
||||
func serveConn(c io.ReadWriteCloser, h Handler) {
|
||||
defer c.Close()
|
||||
line, err := bufio.NewReader(io.LimitReader(c, 4096)).ReadString('\n')
|
||||
cmd := strings.TrimSpace(line)
|
||||
if cmd == "" {
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
fmt.Fprintln(c, "ERR empty command")
|
||||
return
|
||||
}
|
||||
fmt.Fprintln(c, h(cmd))
|
||||
}
|
||||
|
||||
// readReply writes cmd and reads the server's one-line reply.
|
||||
func readReply(c io.ReadWriteCloser, cmd string) (string, error) {
|
||||
if _, err := fmt.Fprintln(c, cmd); err != nil {
|
||||
return "", err
|
||||
}
|
||||
// The server hangs up after its reply; a broken-pipe error after the
|
||||
// last byte still leaves the reply in the buffer.
|
||||
data, _ := io.ReadAll(io.LimitReader(c, 4096))
|
||||
line := strings.TrimSpace(string(data))
|
||||
if line == "" {
|
||||
return "", ErrUnavailable
|
||||
}
|
||||
return line, nil
|
||||
}
|
||||
@@ -0,0 +1,54 @@
|
||||
//go:build linux
|
||||
|
||||
package control
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log/slog"
|
||||
"net"
|
||||
"os"
|
||||
"time"
|
||||
)
|
||||
|
||||
// sockPath lives in the unit's RuntimeDirectory; mode 0666 lets every
|
||||
// local user ask, nothing can reach it from off the machine.
|
||||
const sockPath = "/run/gpu-turnstile/control.sock"
|
||||
|
||||
// Serve starts the socket listener in the background and returns; only a
|
||||
// setup failure is reported. Each client connection is answered in its own
|
||||
// goroutine.
|
||||
func Serve(ctx context.Context, h Handler, log *slog.Logger) error {
|
||||
os.Remove(sockPath) // stale socket from a previous run
|
||||
ln, err := net.Listen("unix", sockPath)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := os.Chmod(sockPath, 0o666); err != nil {
|
||||
ln.Close()
|
||||
return err
|
||||
}
|
||||
go func() {
|
||||
<-ctx.Done()
|
||||
ln.Close()
|
||||
}()
|
||||
go func() {
|
||||
for {
|
||||
c, err := ln.Accept()
|
||||
if err != nil {
|
||||
return // shutting down
|
||||
}
|
||||
go serveConn(c, h)
|
||||
}
|
||||
}()
|
||||
return nil
|
||||
}
|
||||
|
||||
// Ask sends one command to the running service and returns its reply.
|
||||
func Ask(cmd string) (string, error) {
|
||||
c, err := net.DialTimeout("unix", sockPath, 2*time.Second)
|
||||
if err != nil {
|
||||
return "", ErrUnavailable
|
||||
}
|
||||
defer c.Close()
|
||||
return readReply(c, cmd)
|
||||
}
|
||||
@@ -0,0 +1,18 @@
|
||||
//go:build !windows && !linux
|
||||
|
||||
package control
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log/slog"
|
||||
)
|
||||
|
||||
// Serve is a no-op on platforms without a control channel.
|
||||
func Serve(ctx context.Context, h Handler, log *slog.Logger) error {
|
||||
return ErrUnavailable
|
||||
}
|
||||
|
||||
// Ask always reports the channel as unavailable.
|
||||
func Ask(cmd string) (string, error) {
|
||||
return "", ErrUnavailable
|
||||
}
|
||||
@@ -0,0 +1,46 @@
|
||||
package control
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"net"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestRoundTrip(t *testing.T) {
|
||||
server, client := net.Pipe()
|
||||
go serveConn(server, func(cmd string) string {
|
||||
if cmd != CmdUpdateNow {
|
||||
return "ERR unknown command: " + cmd
|
||||
}
|
||||
return "OK v0.2.2 is up to date"
|
||||
})
|
||||
reply, err := readReply(client, CmdUpdateNow)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if reply != "OK v0.2.2 is up to date" {
|
||||
t.Fatalf("reply = %q", reply)
|
||||
}
|
||||
|
||||
server2, client2 := net.Pipe()
|
||||
go serveConn(server2, func(cmd string) string { return "ERR unknown command: " + cmd })
|
||||
reply, err = readReply(client2, "bogus")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !strings.HasPrefix(reply, "ERR ") {
|
||||
t.Fatalf("reply = %q, want ERR prefix", reply)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEmptyReplyIsUnavailable(t *testing.T) {
|
||||
server, client := net.Pipe()
|
||||
go serveConn(server, func(cmd string) string {
|
||||
server.Close() // hang up without answering
|
||||
return ""
|
||||
})
|
||||
if _, err := readReply(client, CmdUpdateNow); !errors.Is(err, ErrUnavailable) {
|
||||
t.Fatalf("err = %v, want ErrUnavailable", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,112 @@
|
||||
//go:build windows
|
||||
|
||||
package control
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"os"
|
||||
"unsafe"
|
||||
|
||||
"golang.org/x/sys/windows"
|
||||
)
|
||||
|
||||
// pipePath is a kernel-local named pipe: no TCP, no firewall prompt.
|
||||
const pipePath = `\\.\pipe\gpu-turnstile`
|
||||
|
||||
// sddlPipe grants full access to Administrators, SYSTEM and the pipe owner,
|
||||
// and read+write to authenticated users — except network logons, so the
|
||||
// pipe cannot be reached from another machine over SMB.
|
||||
const sddlPipe = "D:(D;;GRGW;;;NU)(A;;GA;;;BA)(A;;GA;;;SY)(A;;GA;;;OW)(A;;GRGW;;;AU)"
|
||||
|
||||
var (
|
||||
procConvertSDDL = windows.NewLazySystemDLL("advapi32.dll").
|
||||
NewProc("ConvertStringSecurityDescriptorToSecurityDescriptorW")
|
||||
procWaitNamedPipe = windows.NewLazySystemDLL("kernel32.dll").
|
||||
NewProc("WaitNamedPipeW")
|
||||
)
|
||||
|
||||
func waitNamedPipe(name *uint16, timeout uint32) error {
|
||||
r, _, err := procWaitNamedPipe.Call(uintptr(unsafe.Pointer(name)), uintptr(timeout))
|
||||
if r == 0 {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func securityAttributesFromSDDL(sddl string) (*windows.SecurityAttributes, error) {
|
||||
s, err := windows.UTF16PtrFromString(sddl)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var sd *uint16 // SECURITY_DESCRIPTOR*, kept for the process lifetime
|
||||
r, _, callErr := procConvertSDDL.Call(
|
||||
uintptr(unsafe.Pointer(s)), 1, /* SDDL_REVISION_1 */
|
||||
uintptr(unsafe.Pointer(&sd)), 0)
|
||||
if r == 0 {
|
||||
return nil, fmt.Errorf("invalid SDDL: %w", callErr)
|
||||
}
|
||||
sa := &windows.SecurityAttributes{
|
||||
Length: uint32(unsafe.Sizeof(windows.SecurityAttributes{})),
|
||||
SecurityDescriptor: (*windows.SECURITY_DESCRIPTOR)(unsafe.Pointer(sd)),
|
||||
}
|
||||
return sa, nil
|
||||
}
|
||||
|
||||
// Serve starts the pipe listener in the background and returns; only a
|
||||
// setup failure is reported. Each client connection is answered in its own
|
||||
// goroutine. On shutdown the process exit reaps everything.
|
||||
func Serve(ctx context.Context, h Handler, log *slog.Logger) error {
|
||||
sa, err := securityAttributesFromSDDL(sddlPipe)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
name, err := windows.UTF16PtrFromString(pipePath)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
go func() {
|
||||
for ctx.Err() == nil {
|
||||
pipe, err := windows.CreateNamedPipe(name,
|
||||
windows.PIPE_ACCESS_DUPLEX,
|
||||
windows.PIPE_TYPE_BYTE|windows.PIPE_READMODE_BYTE|windows.PIPE_WAIT,
|
||||
16, 4096, 4096, 0, sa)
|
||||
if err != nil {
|
||||
log.Warn("control channel stopped", "err", err)
|
||||
return
|
||||
}
|
||||
go func() {
|
||||
// Blocks until a client connects; on process exit the
|
||||
// handle goes away with everything else.
|
||||
if err := windows.ConnectNamedPipe(pipe, nil); err != nil {
|
||||
windows.CloseHandle(pipe)
|
||||
return
|
||||
}
|
||||
f := os.NewFile(uintptr(pipe), pipePath)
|
||||
serveConn(f, h) // closes f, and with it the pipe handle
|
||||
}()
|
||||
}
|
||||
}()
|
||||
return nil
|
||||
}
|
||||
|
||||
// Ask sends one command to the running service and returns its reply.
|
||||
func Ask(cmd string) (string, error) {
|
||||
name, err := windows.UTF16PtrFromString(pipePath)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := waitNamedPipe(name, 2000); err != nil {
|
||||
return "", ErrUnavailable
|
||||
}
|
||||
handle, err := windows.CreateFile(name,
|
||||
windows.GENERIC_READ|windows.GENERIC_WRITE, 0, nil,
|
||||
windows.OPEN_EXISTING, 0, 0)
|
||||
if err != nil {
|
||||
return "", ErrUnavailable
|
||||
}
|
||||
f := os.NewFile(uintptr(handle), pipePath)
|
||||
defer f.Close()
|
||||
return readReply(f, cmd)
|
||||
}
|
||||
+75
-2
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -15,6 +15,8 @@ import (
|
||||
"os/signal"
|
||||
"path/filepath"
|
||||
"syscall"
|
||||
|
||||
"gpu-turnstile/internal/supervise"
|
||||
)
|
||||
|
||||
// Name matches the Windows service name; the systemd unit is Name + ".service".
|
||||
@@ -66,7 +68,8 @@ func Run(run func(ctx context.Context) error) error {
|
||||
// BindPaths hole through ProtectHome/ProtectSystem: it reads its venv and
|
||||
// writes output/temp/user data under COMFY_DIR. A venv whose base
|
||||
// interpreter (pyvenv.cfg home) lives outside COMFY_DIR gets an additional
|
||||
// read-only bind.
|
||||
// read-only bind, and a Comfy-Desktop shared data dir (models, input,
|
||||
// output) a read-write one.
|
||||
func renderUnit(exePath, configPath, comfyDir string) string {
|
||||
bind := ""
|
||||
if comfyDir != "" {
|
||||
@@ -74,6 +77,9 @@ func renderUnit(exePath, configPath, comfyDir string) string {
|
||||
if home := comfyVenvHome(comfyDir); home != "" {
|
||||
bind += "BindReadOnlyPaths=" + home + "\n"
|
||||
}
|
||||
if shared := supervise.DesktopSharedDir(comfyDir); shared != "" {
|
||||
bind += "BindPaths=" + shared + "\n"
|
||||
}
|
||||
}
|
||||
return fmt.Sprintf(`[Unit]
|
||||
Description=gpu-turnstile GPU arbitration proxy for Ollama and ComfyUI
|
||||
@@ -89,6 +95,8 @@ RestartSec=5s
|
||||
|
||||
DynamicUser=yes
|
||||
StateDirectory=%s
|
||||
RuntimeDirectory=%s
|
||||
RuntimeDirectoryMode=0755
|
||||
%sProtectSystem=strict
|
||||
ProtectHome=yes
|
||||
PrivateTmp=yes
|
||||
@@ -111,7 +119,7 @@ SystemCallErrorNumber=EPERM
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
`, exePath, configPath, Name, bind)
|
||||
`, exePath, configPath, Name, Name, bind)
|
||||
}
|
||||
|
||||
// copyFile copies src to dst, creating dst with the given mode.
|
||||
|
||||
@@ -17,6 +17,7 @@ func TestRenderUnit(t *testing.T) {
|
||||
"WantedBy=multi-user.target",
|
||||
"DynamicUser=yes",
|
||||
"StateDirectory=gpu-turnstile",
|
||||
"RuntimeDirectory=gpu-turnstile",
|
||||
"ProtectSystem=strict",
|
||||
"NoNewPrivileges=yes",
|
||||
"RestrictAddressFamilies=AF_UNIX AF_INET AF_INET6",
|
||||
|
||||
@@ -24,6 +24,8 @@ import (
|
||||
"golang.org/x/sys/windows"
|
||||
"golang.org/x/sys/windows/svc"
|
||||
"golang.org/x/sys/windows/svc/mgr"
|
||||
|
||||
"gpu-turnstile/internal/supervise"
|
||||
)
|
||||
|
||||
// Name is the Windows service name.
|
||||
@@ -382,6 +384,13 @@ func grantAll(exe, configPath string) error {
|
||||
return err
|
||||
}
|
||||
}
|
||||
// Comfy-Desktop keeps models/input/output in a shared dir next
|
||||
// to the install; the managed instance writes output there.
|
||||
if shared := supervise.DesktopSharedDir(comfyDir); shared != "" {
|
||||
if err := grantAccessTree(shared, "(OI)(CI)(M)"); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
"runtime"
|
||||
"strings"
|
||||
"sync"
|
||||
@@ -43,7 +44,9 @@ func ComfyLayout(goos, dir string) (python, script string) {
|
||||
// DefaultComfyCommand builds the launch command for the standard venv
|
||||
// layout (see ComfyLayout): the script is passed relative to dir so dir
|
||||
// stays the working directory, and --port is taken from comfyURL when the
|
||||
// URL carries one.
|
||||
// URL carries one. On a Comfy-Desktop standalone install the shared data
|
||||
// directory (models, input, output) is added as --*-directory flags so the
|
||||
// managed instance sees the desktop app's models.
|
||||
func DefaultComfyCommand(goos, dir, comfyURL string) string {
|
||||
python, script := ComfyLayout(goos, dir)
|
||||
rel, err := filepath.Rel(dir, script)
|
||||
@@ -54,9 +57,34 @@ func DefaultComfyCommand(goos, dir, comfyURL string) string {
|
||||
if u, err := url.Parse(comfyURL); err == nil && u.Port() != "" {
|
||||
cmd += " --port " + u.Port()
|
||||
}
|
||||
if shared := DesktopSharedDir(dir); shared != "" {
|
||||
for _, sub := range []string{"models", "input", "output"} {
|
||||
p := filepath.Join(shared, sub)
|
||||
if st, err := os.Stat(p); err == nil && st.IsDir() {
|
||||
cmd += ` --` + sub + `-directory "` + p + `"`
|
||||
}
|
||||
}
|
||||
}
|
||||
return cmd
|
||||
}
|
||||
|
||||
// DesktopSharedDir returns the Comfy-Desktop shared data directory
|
||||
// (<root>/ComfyUI-Shared) when dir looks like a desktop standalone install
|
||||
// (<root>/ComfyUI-Installs/<name>/ComfyUI) and the shared models directory
|
||||
// exists; "" otherwise. The desktop app keeps models, input and output
|
||||
// there rather than inside the ComfyUI tree.
|
||||
func DesktopSharedDir(dir string) string {
|
||||
installs := filepath.Dir(filepath.Dir(dir))
|
||||
if filepath.Base(installs) != "ComfyUI-Installs" {
|
||||
return ""
|
||||
}
|
||||
shared := filepath.Join(filepath.Dir(installs), "ComfyUI-Shared")
|
||||
if st, err := os.Stat(filepath.Join(shared, "models")); err == nil && st.IsDir() {
|
||||
return shared
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// Process is one managed child process.
|
||||
type Process struct {
|
||||
name string
|
||||
@@ -136,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 {
|
||||
@@ -185,6 +231,9 @@ func (p *Process) EnsureRunning() error {
|
||||
p.external = false
|
||||
cmd := exec.Command(p.argv[0], p.argv[1:]...)
|
||||
cmd.Dir = p.dir
|
||||
// Ask the child not to colorize (ComfyUI ignores this and colors
|
||||
// anyway, so pipeLog also strips escape sequences).
|
||||
cmd.Env = append(os.Environ(), "NO_COLOR=1", "TERM=dumb")
|
||||
stdout, err := cmd.StdoutPipe()
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -278,7 +327,9 @@ func (p *Process) WatchIdle(ctx context.Context, idleTimeout time.Duration, gpuI
|
||||
}
|
||||
|
||||
// pipeLog forwards one child output stream to the log at INFO, line by
|
||||
// line, prefixed with the process name.
|
||||
// line, prefixed with the process name. ANSI escape sequences are
|
||||
// stripped: ComfyUI colorizes unconditionally, and the escapes only
|
||||
// render as garbage in a log file.
|
||||
func (p *Process) pipeLog(r io.Reader) {
|
||||
buf := make([]byte, 4096)
|
||||
var line string
|
||||
@@ -290,18 +341,25 @@ func (p *Process) pipeLog(r io.Reader) {
|
||||
if i < 0 {
|
||||
break
|
||||
}
|
||||
p.log.Info(p.name + ": " + strings.TrimRight(line[:i], "\r"))
|
||||
p.log.Info(p.name + ": " + stripANSI(strings.TrimRight(line[:i], "\r")))
|
||||
line = line[i+1:]
|
||||
}
|
||||
if err != nil {
|
||||
if strings.TrimSpace(line) != "" {
|
||||
p.log.Info(p.name + ": " + line)
|
||||
p.log.Info(p.name + ": " + stripANSI(line))
|
||||
}
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ansiPattern matches CSI escape sequences (colors, cursor moves, …).
|
||||
var ansiPattern = regexp.MustCompile("\x1b\\[[0-9;?]*[a-zA-Z]")
|
||||
|
||||
func stripANSI(s string) string {
|
||||
return ansiPattern.ReplaceAllString(s, "")
|
||||
}
|
||||
|
||||
// stopTree kills cmd's process, including its children on Windows (python
|
||||
// launchers tend to spawn some). The Wait goroutine reaps it.
|
||||
func stopTree(cmd *exec.Cmd) {
|
||||
|
||||
@@ -239,3 +239,38 @@ func TestComfyLayoutAndDefaultCommand(t *testing.T) {
|
||||
t.Errorf("missing-layout script = %s", script)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStripANSI(t *testing.T) {
|
||||
cases := map[string]string{
|
||||
"\x1b[32m[INFO]\x1b[0m Starting server": "[INFO] Starting server",
|
||||
"\x1b[1m\x1b[31m[ERROR]\x1b[0m boom": "[ERROR] boom",
|
||||
"plain line": "plain line",
|
||||
"\x1b[33mWARN\x1b[0m: \x1b[1mbold\x1b[0m": "WARN: bold",
|
||||
}
|
||||
for in, want := range cases {
|
||||
if got := stripANSI(in); got != want {
|
||||
t.Errorf("stripANSI(%q) = %q, want %q", in, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestDesktopSharedDir(t *testing.T) {
|
||||
root := t.TempDir()
|
||||
comfy := filepath.Join(root, "ComfyUI-Installs", "rtx5080", "ComfyUI")
|
||||
if err := os.MkdirAll(comfy, 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := DesktopSharedDir(comfy); got != "" {
|
||||
t.Fatalf("no shared dir yet: got %q, want empty", got)
|
||||
}
|
||||
shared := filepath.Join(root, "ComfyUI-Shared")
|
||||
if err := os.MkdirAll(filepath.Join(shared, "models"), 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := DesktopSharedDir(comfy); got != shared {
|
||||
t.Fatalf("got %q, want %q", got, shared)
|
||||
}
|
||||
if got := DesktopSharedDir(filepath.Join(root, "plain", "ComfyUI")); got != "" {
|
||||
t.Fatalf("non-desktop layout: got %q, want empty", got)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user