2b97bac51e
The lifecycle rule from Vikunja #488. Not on demand, because a 7-14B takes tens of seconds to load and a world question would meet a gap every time the card had been quiet. Not always on, because that is what holds the card. /health is answered locally and always, so Maven's prober costs nothing and works while the model is down. Everything else is reverse-proxied to llama-server, which is what makes the idle window measurable at all. Yielding is checked before starting, and both transitions are damped by a poll streak so a short-lived rocm process cannot evict the model.
248 lines
7.2 KiB
Go
248 lines
7.2 KiB
Go
// mavgpud — the workstation's GPU supervisor.
|
|
//
|
|
// It runs on the workstation (an AMD 7900 GRE, 16GB), not on homesrv, and it is
|
|
// deployed separately from the Maven daemons. Maven does not participate in any
|
|
// of this and never asks for a start: it reads /health through internal/llm.Pair
|
|
// and either gets the big model or falls back to the resident 1.7B.
|
|
//
|
|
// The rule, from Vikunja #488: keep llama-server loaded whenever the card is
|
|
// free, unload it when it has been idle too long or when another process needs
|
|
// the card. Not on demand, because a 7-14B takes tens of seconds to load and a
|
|
// world question would be answered by a gap every time the card had been quiet.
|
|
// Not always on, because that holds 16GB against the owner's own jobs.
|
|
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"flag"
|
|
"log"
|
|
"net/http"
|
|
"net/http/httputil"
|
|
"net/url"
|
|
"os"
|
|
"os/signal"
|
|
"sync/atomic"
|
|
"syscall"
|
|
"time"
|
|
)
|
|
|
|
type config struct {
|
|
Listen string `json:"listen"` // what Maven talks to
|
|
LlamaAddr string `json:"llama_addr"` // where llama-server binds
|
|
LlamaBin string `json:"llama_bin"`
|
|
// LlamaArgs must include the flags that bind LlamaAddr. They are passed
|
|
// through untouched so the model, context size and layer count stay the
|
|
// owner's business and not this daemon's schema.
|
|
LlamaArgs []string `json:"llama_args"`
|
|
|
|
KFDRoot string `json:"kfd_root"`
|
|
DRMDevice string `json:"drm_device"`
|
|
|
|
Poll duration `json:"poll"`
|
|
IdleTimeout duration `json:"idle_timeout"`
|
|
StopGrace duration `json:"stop_grace"`
|
|
MinFreeVRAM int64 `json:"min_free_vram_bytes"`
|
|
// EvictAfter and StartAfter are counted in polls, not seconds. Both exist
|
|
// to damp flapping: a one-tick blip from a short-lived rocm process must
|
|
// not evict the model, and a card that has just been released must not be
|
|
// grabbed before the previous job has finished unmapping.
|
|
EvictAfter int `json:"evict_after_polls"`
|
|
StartAfter int `json:"start_after_polls"`
|
|
}
|
|
|
|
func defaults() config {
|
|
return config{
|
|
Listen: ":8080",
|
|
LlamaAddr: "127.0.0.1:8081",
|
|
KFDRoot: "/sys/class/kfd/kfd/proc",
|
|
DRMDevice: "/sys/class/drm/card1/device",
|
|
Poll: duration(time.Second),
|
|
IdleTimeout: duration(15 * time.Minute),
|
|
StopGrace: duration(20 * time.Second),
|
|
MinFreeVRAM: 15 << 30,
|
|
EvictAfter: 2,
|
|
StartAfter: 5,
|
|
}
|
|
}
|
|
|
|
// duration lets the config file say "15m" instead of counting nanoseconds.
|
|
type duration time.Duration
|
|
|
|
func (d *duration) UnmarshalJSON(b []byte) error {
|
|
var s string
|
|
if err := json.Unmarshal(b, &s); err != nil {
|
|
return err
|
|
}
|
|
v, err := time.ParseDuration(s)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
*d = duration(v)
|
|
return nil
|
|
}
|
|
|
|
func main() {
|
|
path := flag.String("config", "/etc/mavgpud.json", "config file")
|
|
flag.Parse()
|
|
|
|
cfg := defaults()
|
|
b, err := os.ReadFile(*path)
|
|
if err != nil {
|
|
log.Fatalf("mavgpud: read config: %v", err)
|
|
}
|
|
if err := json.Unmarshal(b, &cfg); err != nil {
|
|
log.Fatalf("mavgpud: parse config: %v", err)
|
|
}
|
|
if cfg.LlamaBin == "" {
|
|
log.Fatal("mavgpud: llama_bin is required")
|
|
}
|
|
|
|
base := "http://" + cfg.LlamaAddr
|
|
run := newRunner(cfg.LlamaBin, cfg.LlamaArgs, base+"/health")
|
|
sup := &supervisor{
|
|
cfg: cfg,
|
|
probe: probe{kfdRoot: cfg.KFDRoot, drmDev: cfg.DRMDevice},
|
|
run: run,
|
|
}
|
|
sup.touch()
|
|
|
|
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
|
defer cancel()
|
|
|
|
target, err := url.Parse(base)
|
|
if err != nil {
|
|
log.Fatalf("mavgpud: llama_addr: %v", err)
|
|
}
|
|
srv := &http.Server{Addr: cfg.Listen, Handler: sup.handler(target)}
|
|
go func() {
|
|
log.Printf("mavgpud: listening on %s, model %s", cfg.Listen, cfg.LlamaBin)
|
|
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
|
log.Fatalf("mavgpud: listen: %v", err)
|
|
}
|
|
}()
|
|
|
|
sup.loop(ctx)
|
|
|
|
// The card must come back before we do. A supervisor that exits leaving
|
|
// llama-server holding 14GB is worse than one that never ran.
|
|
shut, done := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer done()
|
|
_ = srv.Shutdown(shut)
|
|
run.stop(time.Duration(cfg.StopGrace))
|
|
}
|
|
|
|
type supervisor struct {
|
|
cfg config
|
|
probe probe
|
|
run *runner
|
|
|
|
lastReq atomic.Int64 // unix nanos of the last request Maven sent
|
|
|
|
foreignStreak int
|
|
clearStreak int
|
|
}
|
|
|
|
func (s *supervisor) touch() { s.lastReq.Store(time.Now().UnixNano()) }
|
|
|
|
func (s *supervisor) idle() time.Duration {
|
|
return time.Since(time.Unix(0, s.lastReq.Load()))
|
|
}
|
|
|
|
// handler serves the two things the workstation exposes.
|
|
//
|
|
// /health is answered locally and always, with no GPU cost and no round trip,
|
|
// because it is the only thing Maven reads and Maven reads it on a timer
|
|
// forever. Everything else is llama-server's API, reverse-proxied. Proxying
|
|
// rather than pointing Maven straight at llama-server is what makes the idle
|
|
// window measurable: the supervisor cannot otherwise know when the model was
|
|
// last used.
|
|
func (s *supervisor) handler(target *url.URL) http.Handler {
|
|
proxy := httputil.NewSingleHostReverseProxy(target)
|
|
mux := http.NewServeMux()
|
|
mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) {
|
|
if !s.run.isReady() {
|
|
http.Error(w, "model not loaded", http.StatusServiceUnavailable)
|
|
return
|
|
}
|
|
w.Header().Set("Content-Type", "application/json")
|
|
_, _ = w.Write([]byte(`{"status":"ok"}`))
|
|
})
|
|
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
|
|
if !s.run.isReady() {
|
|
http.Error(w, "model not loaded", http.StatusServiceUnavailable)
|
|
return
|
|
}
|
|
s.touch()
|
|
proxy.ServeHTTP(w, r)
|
|
})
|
|
return mux
|
|
}
|
|
|
|
func (s *supervisor) loop(ctx context.Context) {
|
|
t := time.NewTicker(time.Duration(s.cfg.Poll))
|
|
defer t.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-t.C:
|
|
s.tick(ctx)
|
|
}
|
|
}
|
|
}
|
|
|
|
// tick is the whole decision. Yielding is checked before starting, and presence
|
|
// on the KFD is what triggers it — not a VRAM threshold. A ROCm process
|
|
// registers under /sys/class/kfd/kfd/proc when it initialises HIP, before it
|
|
// allocates, so we see a contender during its startup rather than after it has
|
|
// already failed to get the memory it wanted.
|
|
func (s *supervisor) tick(ctx context.Context) {
|
|
others := s.probe.foreign(s.run.pid())
|
|
if len(others) > 0 {
|
|
s.foreignStreak++
|
|
s.clearStreak = 0
|
|
} else {
|
|
s.foreignStreak = 0
|
|
s.clearStreak++
|
|
}
|
|
|
|
if s.run.running() {
|
|
s.run.refreshReady(ctx)
|
|
switch {
|
|
case s.foreignStreak >= s.cfg.EvictAfter:
|
|
log.Printf("mavgpud: yielding the card to %s", describe(others))
|
|
s.run.stop(time.Duration(s.cfg.StopGrace))
|
|
case s.idle() > time.Duration(s.cfg.IdleTimeout):
|
|
log.Printf("mavgpud: idle for %s, unloading", s.idle().Round(time.Second))
|
|
s.run.stop(time.Duration(s.cfg.StopGrace))
|
|
}
|
|
return
|
|
}
|
|
|
|
if s.clearStreak < s.cfg.StartAfter {
|
|
return
|
|
}
|
|
if free := s.probe.freeVRAM(); free < s.cfg.MinFreeVRAM {
|
|
return
|
|
}
|
|
s.touch() // the idle clock starts at load, not at the last request before it
|
|
if err := s.run.start(); err != nil {
|
|
log.Printf("mavgpud: start llama-server: %v", err)
|
|
}
|
|
}
|
|
|
|
// describe names the contenders in the log. This log is the instrument for the
|
|
// open question in #488: whether polling the KFD misses a job that wants the
|
|
// card without registering there.
|
|
func describe(procs []gpuProc) string {
|
|
out := ""
|
|
for i, p := range procs {
|
|
if i > 0 {
|
|
out += ", "
|
|
}
|
|
out += p.Comm
|
|
}
|
|
return out
|
|
}
|