a6c6e49f112567d4c61b8ed83282e9f53303d120

Author
TheEdgeOfRage <git@theedgeofrage.com>
Committer
TheEdgeOfRage <git@theedgeofrage.com>
Date

Message

Add internal/services package for self-hosted model services

Process manager that spawns llama-server + audiocpp_server as subprocesses:
AMD free-VRAM gate (16 GiB), all-or-nothing loopback spawn (9931/9932),
concurrent TTS/STT/LLM health probes, and SIGTERM->grace->SIGKILL teardown.
Exported NewManager/Start/Stop; no callers yet.

Diff

This diff is truncated to protect this page.

  1diff --git a/internal/services/manager.go b/internal/services/manager.go
  2new file mode 100644
  3index 0000000000000000000000000000000000000000..5bcb9b99d9b50228075099d7d57da4ef6c7fa695
  4--- /dev/null
  5+++ b/internal/services/manager.go
  6@@ -0,0 +1,172 @@
  7+// Package services spawns, health-checks, and tears down the two self-hosted
  8+// model services (llama.cpp router and audio.cpp) as subprocesses. It owns the
  9+// process lifecycle only; it does not serve game or LLM logic.
 10+package services
 11+
 12+import (
 13+	"context"
 14+	"fmt"
 15+	"os/exec"
 16+	"syscall"
 17+	"time"
 18+)
 19+
 20+const (
 21+	LLMListen   = "127.0.0.1:9931"
 22+	AudioListen = "127.0.0.1:9932"
 23+
 24+	minFreeVRAMBytes  = 16 << 30         // 16 GiB free required to start
 25+	modelReadyTimeout = 5 * time.Minute  // total budget for both services to become healthy
 26+	killGrace         = 10 * time.Second // SIGTERM -> wait -> SIGKILL
 27+)
 28+
 29+const (
 30+	llamaBinary = "llama-server"
 31+	audioBinary = "audiocpp_server"
 32+	audioConfig = "audio.cpp.json"
 33+)
 34+
 35+// llamaServerArgs mirrors the Makefile `llama` target but bound to loopback.
 36+// Keep in sync with the Makefile.
 37+var llamaServerArgs = []string{
 38+	"--host", "127.0.0.1",
 39+	"--port", "9931",
 40+	"--hf-repo", "unsloth/gemma-4-12B-it-qat-GGUF:UD-Q4_K_XL",
 41+	"--n-gpu-layers", "99",
 42+	"--n-gpu-layers-draft", "99",
 43+	"--fit", "off",
 44+	"--flash-attn", "on",
 45+	"--parallel", "5",
 46+	"--kv-unified",
 47+	"--ctx-size", "32768",
 48+	"--batch-size", "1024",
 49+	"--ubatch-size", "1024",
 50+	"--temp", "1.0",
 51+	"--top-p", "0.95",
 52+	"--top-k", "64",
 53+	"--spec-type", "draft-mtp",
 54+	"--spec-draft-n-max", "4",
 55+	"--reasoning", "off",
 56+}
 57+
 58+type child struct {
 59+	cmd  *exec.Cmd
 60+	pgid int
 61+	done chan struct{}
 62+}
 63+
 64+// Manager owns the spawned model-service subprocesses. Create one with
 65+// NewManager, Start it, and Stop it on exit. Stop is safe to call multiple times.
 66+type Manager struct {
 67+	llama   *child
 68+	audio   *child
 69+	stopped bool
 70+}
 71+
 72+func NewManager() *Manager { return &Manager{} }
 73+
 74+// Start runs the VRAM gate, spawns both services (all-or-nothing), and blocks
 75+// until all three probes are healthy or ctx / the readiness deadline expires.
 76+// On any failure it tears down whatever was started and returns an error.
 77+func (m *Manager) Start(ctx context.Context) error {
 78+	if err := checkVRAM(); err != nil {
 79+		return err
 80+	}
 81+
 82+	llama, err := startChild("llama-server", llamaBinary, llamaServerArgs)
 83+	if err != nil {
 84+		return err
 85+	}
 86+	m.llama = llama
 87+
 88+	audio, err := startChild("audiocpp_server", audioBinary, []string{"--config", audioConfig})
 89+	if err != nil {
 90+		m.Stop()
 91+		return err
 92+	}
 93+	m.audio = audio
 94+
 95+	if err := m.waitHealthy(ctx); err != nil {
 96+		m.Stop()
 97+		return err
 98+	}
 99+	return nil
100+}
101+
102+// Stop SIGTERMs each child process group, waits up to killGrace, SIGKILLs any
103+// survivor, and reaps both. Safe to call multiple times; a no-op if nothing was
104+// started.
105+func (m *Manager) Stop() {
106diff --git a/internal/services/probe.go b/internal/services/probe.go
107new file mode 100644
108index 0000000000000000000000000000000000000000..871ebe1692b59cb4243ac7c520d1a0dd55852b1a
109--- /dev/null
110+++ b/internal/services/probe.go
111@@ -0,0 +1,147 @@
112+package services
113+
114+import (
115+	"bytes"
116+	"context"
117+	"encoding/binary"
118+	"fmt"
119+	"io"
120+	"math"
121+	"net/http"
122+	"strings"
123+	"sync"
124+	"time"
125+
126+	"git.theedgeofrage.com/TheEdgeOfRage/kaiwari/internal/stt"
127+	"git.theedgeofrage.com/TheEdgeOfRage/kaiwari/internal/tts"
128+)
129+
130+const (
131+	probeTimeout = 10 * time.Second
132+	pollInterval = 2 * time.Second
133+)
134+
135+// waitHealthy polls the TTS, STT, and LLM probes concurrently until all three
136+// succeed or the total readiness deadline (or ctx) expires. Each probe is a real
137+// request round-trip through the same client methods the app uses.
138+func (m *Manager) waitHealthy(ctx context.Context) error {
139+	ctx, cancel := context.WithDeadline(ctx, time.Now().Add(modelReadyTimeout))
140+	defer cancel()
141+
142+	hc := newProbeHTTP()
143+	ttsClient := tts.NewClient(httpURL(AudioListen), hc)
144+	sttClient := stt.NewASRClient(httpURL(AudioListen), hc)
145+	llmURL := httpURL(LLMListen) + "/v1/chat/completions"
146+	audio := probeWAV()
147+
148+	var failed []string
149+	for {
150+		failed = probeAll(ctx, ttsClient, sttClient, llmURL, hc, audio)
151+		if len(failed) == 0 {
152+			return nil
153+		}
154+		if !sleepCtx(ctx, pollInterval) {
155+			break
156+		}
157+	}
158+	return fmt.Errorf("services: model services not healthy after %s (unhealthy: %s)", modelReadyTimeout, strings.Join(failed, ", "))
159+}
160+
161+func probeAll(ctx context.Context, ttsClient *tts.Client, sttClient *stt.ASRClient, llmURL string, hc *http.Client, audio []byte) []string {
162+	names := []string{"tts", "stt", "llm"}
163+	probes := []func(context.Context) error{
164+		func(ctx context.Context) error { _, e := ttsClient.Speech(ctx, "あ"); return e },
165+		func(ctx context.Context) error { _, e := sttClient.TranscribeBytes(ctx, audio); return e },
166+		func(ctx context.Context) error { return probeLLMCompletion(ctx, llmURL, hc) },
167+	}
168+
169+	errs := make([]error, len(names))
170+	var wg sync.WaitGroup
171+	for i := range probes {
172+		wg.Add(1)
173+		go func(i int) {
174+			defer wg.Done()
175+			errs[i] = probes[i](ctx)
176+		}(i)
177+	}
178+	wg.Wait()
179+
180+	var failed []string
181+	for i, e := range errs {
182+		if e != nil {
183+			failed = append(failed, names[i])
184+		}
185+	}
186+	return failed
187+}
188+
189+func probeLLMCompletion(ctx context.Context, url string, hc *http.Client) error {
190+	body := `{"messages":[{"role":"user","content":"Ready."}],"max_tokens":1,"stream":false}`
191+	req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, strings.NewReader(body))
192+	if err != nil {
193+		return fmt.Errorf("services: llm probe build request: %w", err)
194+	}
195+	req.Header.Set("Content-Type", "application/json")
196+	resp, err := hc.Do(req)
197+	if err != nil {
198+		return fmt.Errorf("services: llm probe request to %s: %w", url, err)
199+	}
200+	defer func() { _ = resp.Body.Close() }()
201+	detail, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
202+	if resp.StatusCode < 200 || resp.StatusCode >= 300 {
203+		return fmt.Errorf("services: llm probe %s returned HTTP %d: %s", url, resp.StatusCode, strings.TrimSpace(string(detail)))
204+	}
205+	return nil
206+}
207+
208+func newProbeHTTP() *http.Client { return &http.Client{Timeout: probeTimeout} }
209+
210+func httpURL(hostPort string) string { return "http://" + hostPort }
211diff --git a/internal/services/vram.go b/internal/services/vram.go
212new file mode 100644
213index 0000000000000000000000000000000000000000..d5df437b0cf79b438a942fe4f087f7a10df8ee6b
214--- /dev/null
215+++ b/internal/services/vram.go
216@@ -0,0 +1,67 @@
217+package services
218+
219+import (
220+	"fmt"
221+	"os"
222+	"path/filepath"
223+	"strconv"
224+	"strings"
225+)
226+
227+const vramTotalGlob = "/sys/class/drm/card*/device/mem_info_vram_total"
228+
229+// checkVRAM sums free VRAM across all AMD cards from sysfs and requires at
230+// least minFreeVRAMBytes free. It fails closed: no readable AMD VRAM sysfs is
231+// an error (only AMD GPUs are supported).
232+func checkVRAM() error {
233+	total, used, err := readVRAM()
234+	if err != nil {
235+		return fmt.Errorf("services: %w", err)
236+	}
237+	free := total - used
238+	if free < minFreeVRAMBytes {
239+		return fmt.Errorf("services: insufficient free VRAM to start models: need %d bytes, have %d bytes", minFreeVRAMBytes, free)
240+	}
241+	return nil
242+}
243+
244+func readVRAM() (int64, int64, error) {
245+	totals, err := filepath.Glob(vramTotalGlob)
246+	if err != nil {
247+		return 0, 0, fmt.Errorf("glob %s: %w", vramTotalGlob, err)
248+	}
249+	if len(totals) == 0 {
250+		return 0, 0, fmt.Errorf("no AMD GPU VRAM sysfs found at %s (AMD only)", vramTotalGlob)
251+	}
252+	var total, used int64
253+	for _, tf := range totals {
254+		t, err := readMemInfoBytes(tf)
255+		if err != nil {
256+			return 0, 0, err
257+		}
258+		uf := strings.Replace(tf, "mem_info_vram_total", "mem_info_vram_used", 1)
259+		u, err := readMemInfoBytes(uf)
260+		if err != nil {
261+			return 0, 0, err
262+		}
263+		total += t
264+		used += u
265+	}
266+	return total, used, nil
267+}
268+
269+func readMemInfoBytes(path string) (int64, error) {
270+	raw, err := os.ReadFile(path)
271+	if err != nil {
272+		return 0, fmt.Errorf("read %s: %w", path, err)
273+	}
274+	fields := strings.Fields(string(raw))
275+	if len(fields) == 0 {
276+		return 0, fmt.Errorf("%s is empty", path)
277+	}
278+	n, err := strconv.ParseInt(fields[0], 10, 64)
279+	if err != nil {
280+		return 0, fmt.Errorf("parse %s: %w", path, err)
281+	}
282+	return n, nil
283+}