Add the admin endpoints, the crossbar binary and the fake upstream
Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)
This commit is contained in:
@@ -0,0 +1,79 @@
|
||||
// Command crossbar is the affinity router for the fleet's llama-server instances. It wires config,
|
||||
// health, proxy and admin into one HTTP server with graceful shutdown on SIGINT/SIGTERM.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"flag"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/admin"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/config"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/health"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/proxy"
|
||||
)
|
||||
|
||||
func main() {
|
||||
if err := run(); err != nil {
|
||||
fmt.Fprintln(os.Stderr, "crossbar: "+err.Error())
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
|
||||
func run() error {
|
||||
configPath := flag.String("config", "crossbar.toml", "path to the crossbar config file")
|
||||
flag.Parse()
|
||||
|
||||
cfg, err := config.Load(*configPath)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
log := slog.New(slog.NewTextHandler(os.Stderr, nil))
|
||||
|
||||
baseURLs := make(map[string]string, len(cfg.Hosts))
|
||||
for name, host := range cfg.Hosts {
|
||||
baseURLs[name] = host.BaseURL
|
||||
}
|
||||
|
||||
table := health.New(baseURLs, cfg.PollInterval.Duration, nil)
|
||||
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
||||
defer stop()
|
||||
go table.Run(ctx)
|
||||
|
||||
mux := http.NewServeMux()
|
||||
mux.Handle("/_crossbar/", admin.Handler(cfg, table))
|
||||
mux.Handle("/", proxy.New(cfg, table, log))
|
||||
|
||||
srv := &http.Server{
|
||||
Addr: cfg.Listen,
|
||||
Handler: mux,
|
||||
ReadHeaderTimeout: 10 * time.Second,
|
||||
}
|
||||
|
||||
serverErr := make(chan error, 1)
|
||||
go func() {
|
||||
log.Info("listening", "addr", srv.Addr)
|
||||
serverErr <- srv.ListenAndServe()
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
log.Info("shutting down")
|
||||
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
return srv.Shutdown(shutdownCtx)
|
||||
case err := <-serverErr:
|
||||
if errors.Is(err, http.ErrServerClosed) {
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,101 @@
|
||||
// fakeupstream stands in for a llama-server router in tests and the smoke run. Do not edit.
|
||||
//
|
||||
// fakeupstream -listen 127.0.0.1:18081 -name alpha -models a,b -down-file /tmp/alpha.down
|
||||
//
|
||||
// /health answers 503 while the down file exists, 200 otherwise. /v1/models lists -models.
|
||||
// /props answers a small JSON object. /v1/chat/completions echoes: a streamed answer of five
|
||||
// SSE chunks 200 ms apart when the body has "stream": true, one JSON answer otherwise. Every
|
||||
// response carries X-Upstream: <name>.
|
||||
package main
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"flag"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
func main() {
|
||||
listen := flag.String("listen", "127.0.0.1:18081", "address to listen on")
|
||||
name := flag.String("name", "fake", "name reported in X-Upstream and answers")
|
||||
models := flag.String("models", "m", "comma-separated model ids for /v1/models")
|
||||
downFile := flag.String("down-file", "", "while this file exists, /health answers 503")
|
||||
flag.Parse()
|
||||
|
||||
ids := strings.Split(*models, ",")
|
||||
mux := http.NewServeMux()
|
||||
stamp := func(w http.ResponseWriter) { w.Header().Set("X-Upstream", *name) }
|
||||
|
||||
mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) {
|
||||
stamp(w)
|
||||
if *downFile != "" {
|
||||
if _, err := os.Stat(*downFile); err == nil {
|
||||
http.Error(w, `{"error":{"message":"Loading model"}}`, http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
}
|
||||
writeJSON(w, map[string]string{"status": "ok"})
|
||||
})
|
||||
mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) {
|
||||
stamp(w)
|
||||
data := []map[string]any{}
|
||||
for _, id := range ids {
|
||||
data = append(data, map[string]any{"id": id, "object": "model", "owned_by": *name})
|
||||
}
|
||||
writeJSON(w, map[string]any{"object": "list", "data": data})
|
||||
})
|
||||
mux.HandleFunc("/props", func(w http.ResponseWriter, r *http.Request) {
|
||||
stamp(w)
|
||||
writeJSON(w, map[string]any{"default_generation_settings": map[string]any{"n_ctx": 8192}, "total_slots": 2, "model_path": *name})
|
||||
})
|
||||
mux.HandleFunc("/v1/chat/completions", func(w http.ResponseWriter, r *http.Request) {
|
||||
stamp(w)
|
||||
body, _ := io.ReadAll(io.LimitReader(r.Body, 1<<20))
|
||||
var req struct {
|
||||
Model string `json:"model"`
|
||||
Stream bool `json:"stream"`
|
||||
}
|
||||
_ = json.Unmarshal(body, &req)
|
||||
if !req.Stream {
|
||||
writeJSON(w, map[string]any{
|
||||
"id": "chatcmpl-fake", "object": "chat.completion", "model": req.Model,
|
||||
"choices": []map[string]any{{"index": 0, "message": map[string]string{"role": "assistant", "content": "hello from " + *name}, "finish_reason": "stop"}},
|
||||
"usage": map[string]int{"prompt_tokens": 3, "completion_tokens": 3, "total_tokens": 6},
|
||||
})
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
w.Header().Set("Cache-Control", "no-cache")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
fl, _ := w.(http.Flusher)
|
||||
for i := 1; i <= 5; i++ {
|
||||
chunk := map[string]any{"id": "chatcmpl-fake", "object": "chat.completion.chunk", "model": req.Model,
|
||||
"choices": []map[string]any{{"index": 0, "delta": map[string]string{"content": fmt.Sprintf("%s chunk %d ", *name, i)}}}}
|
||||
b, _ := json.Marshal(chunk)
|
||||
fmt.Fprintf(w, "data: %s\n\n", b)
|
||||
if fl != nil {
|
||||
fl.Flush()
|
||||
}
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
}
|
||||
fmt.Fprint(w, "data: [DONE]\n\n")
|
||||
})
|
||||
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
|
||||
stamp(w)
|
||||
http.Error(w, `{"error":"not found"}`, http.StatusNotFound)
|
||||
})
|
||||
|
||||
log.Printf("fakeupstream %s listening on %s models=%v", *name, *listen, ids)
|
||||
srv := &http.Server{Addr: *listen, Handler: mux, ReadHeaderTimeout: 5 * time.Second}
|
||||
log.Fatal(srv.ListenAndServe())
|
||||
}
|
||||
|
||||
func writeJSON(w http.ResponseWriter, v any) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_ = json.NewEncoder(w).Encode(v)
|
||||
}
|
||||
@@ -8,5 +8,6 @@ owner fills in the Model column. The reviewer adds findings under "Reviews" once
|
||||
| v0/01-module-gate-config | 2026-09-25 | done | 1 | pass | none | `go mod download` fetched the module (network available); gate passed on the first run. | ? |
|
||||
| v0/02-health | 2026-09-25 | done | 1 | pass | none | First gate run passed. `MarkDown` initially forgot to write the entry back; caught by `TestMarkDown`. | ? |
|
||||
| v0/03-proxy | 2026-09-25 | done | 1 | pass | none | `SplitRoute` must reject an empty first segment (`/`, `//x`) as `ok=false`; the model peek restores the body and leaves non-JSON/empty as `""`. | ? |
|
||||
| v0/04-admin-main | 2026-09-25 | done | 1 | pass | none | `timeout --signal=TERM 3` exits 124 on a timed-out child on this GNU system, so the task's `exit=0` is not observable through it; sent SIGTERM directly and confirmed crossbar's own exit code is 0 with both log lines. | ? |
|
||||
|
||||
## Reviews
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
# crossbar example configuration. Replace <tailnet> and the addresses with your own.
|
||||
listen = "127.0.0.1:17777" # never 0.0.0.0 — bind the tailnet address in production
|
||||
poll_interval = "1s" # 60s in production; 1s makes the smoke run quick
|
||||
queue_max = 8
|
||||
|
||||
[hosts.alpha]
|
||||
base_url = "http://127.0.0.1:18081" # e.g. http://straylight.<tailnet>:11434
|
||||
weight = 1.0
|
||||
models = { "ornith-1.5-35b-a3b" = { parallel = 4 }, "small-9b" = { parallel = 6 } }
|
||||
|
||||
[hosts.beta]
|
||||
base_url = "http://127.0.0.1:18082" # e.g. http://titan.<tailnet>:8081
|
||||
weight = 2.0
|
||||
models = { "ornith-1.5-35b-a3b" = { parallel = 4 } }
|
||||
|
||||
# v0: a route is a preference list; the first healthy host that has the model wins.
|
||||
[routes.opencode-a]
|
||||
hosts = ["alpha", "beta"]
|
||||
default_model = "ornith-1.5-35b-a3b"
|
||||
|
||||
[routes.hermes-x]
|
||||
hosts = ["beta", "alpha"]
|
||||
@@ -0,0 +1,102 @@
|
||||
// Package admin serves the operator's view of crossbar: the health table and the routes as JSON,
|
||||
// mounted at /_crossbar/ on the same listener as the proxy. The shape of /_crossbar/hosts is
|
||||
// fixed so operators can read why a request went where it went.
|
||||
package admin
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/config"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/health"
|
||||
)
|
||||
|
||||
// Hosts is what the admin handler needs from the health table.
|
||||
type Hosts interface {
|
||||
All() map[string]health.Status
|
||||
}
|
||||
|
||||
type HostView struct {
|
||||
Healthy bool `json:"healthy"`
|
||||
Loaded []string `json:"loaded"` // never null: an empty slice when nothing is loaded
|
||||
LastOK string `json:"last_ok"` // time.RFC3339 in UTC, or "" if never
|
||||
LastErr string `json:"last_err"`
|
||||
}
|
||||
|
||||
type RouteView struct {
|
||||
Hosts []string `json:"hosts"`
|
||||
DefaultModel string `json:"default_model"`
|
||||
}
|
||||
|
||||
// Handler serves GET /_crossbar/hosts and GET /_crossbar/routes. Any other method on those paths is
|
||||
// a 405 with an Allow: GET header; anything else under the handler is a 404.
|
||||
func Handler(cfg *config.Config, h Hosts) http.Handler {
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/_crossbar/hosts", hostsHandler(h))
|
||||
mux.HandleFunc("/_crossbar/routes", routesHandler(cfg))
|
||||
mux.HandleFunc("/", notFound)
|
||||
return mux
|
||||
}
|
||||
|
||||
func hostsHandler(h Hosts) http.HandlerFunc {
|
||||
return func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodGet {
|
||||
wrongMethod(w)
|
||||
return
|
||||
}
|
||||
views := make(map[string]HostView, len(h.All()))
|
||||
for name, s := range h.All() {
|
||||
views[name] = hostView(s)
|
||||
}
|
||||
writeJSON(w, http.StatusOK, views)
|
||||
}
|
||||
}
|
||||
|
||||
func routesHandler(cfg *config.Config) http.HandlerFunc {
|
||||
return func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodGet {
|
||||
wrongMethod(w)
|
||||
return
|
||||
}
|
||||
views := make(map[string]RouteView, len(cfg.Routes))
|
||||
for name, route := range cfg.Routes {
|
||||
hosts := make([]string, len(route.Hosts))
|
||||
copy(hosts, route.Hosts)
|
||||
views[name] = RouteView{Hosts: hosts, DefaultModel: route.DefaultModel}
|
||||
}
|
||||
writeJSON(w, http.StatusOK, views)
|
||||
}
|
||||
}
|
||||
|
||||
func hostView(s health.Status) HostView {
|
||||
loaded := s.Loaded
|
||||
if loaded == nil {
|
||||
loaded = []string{}
|
||||
}
|
||||
lastOK := ""
|
||||
if !s.LastOK.IsZero() {
|
||||
lastOK = s.LastOK.UTC().Format(time.RFC3339)
|
||||
}
|
||||
return HostView{
|
||||
Healthy: s.Healthy,
|
||||
Loaded: loaded,
|
||||
LastOK: lastOK,
|
||||
LastErr: s.LastErr,
|
||||
}
|
||||
}
|
||||
|
||||
func wrongMethod(w http.ResponseWriter) {
|
||||
w.Header().Set("Allow", "GET")
|
||||
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
|
||||
}
|
||||
|
||||
func notFound(w http.ResponseWriter, r *http.Request) {
|
||||
writeJSON(w, http.StatusNotFound, map[string]string{"error": "not found"})
|
||||
}
|
||||
|
||||
func writeJSON(w http.ResponseWriter, status int, v any) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(status)
|
||||
_ = json.NewEncoder(w).Encode(v)
|
||||
}
|
||||
@@ -0,0 +1,99 @@
|
||||
package admin_test
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/admin"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/config"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/health"
|
||||
)
|
||||
|
||||
type fakeHosts map[string]health.Status
|
||||
|
||||
func (f fakeHosts) All() map[string]health.Status { return f }
|
||||
|
||||
func testConfig(t *testing.T) *config.Config {
|
||||
c, err := config.Parse(strings.NewReader(`
|
||||
listen = "127.0.0.1:1"
|
||||
[hosts.alpha]
|
||||
base_url = "http://alpha:1"
|
||||
models = { "m" = { } }
|
||||
[hosts.beta]
|
||||
base_url = "http://beta:1"
|
||||
models = { "m" = { } }
|
||||
[routes.r]
|
||||
hosts = ["alpha", "beta"]
|
||||
default_model = "m"
|
||||
`))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
func TestHosts(t *testing.T) {
|
||||
when := time.Date(2026, 9, 25, 8, 0, 0, 0, time.UTC)
|
||||
h := admin.Handler(testConfig(t), fakeHosts{
|
||||
"alpha": {Healthy: true, Loaded: []string{"m"}, LastOK: when},
|
||||
"beta": {Healthy: false, LastErr: "HTTP 503"},
|
||||
})
|
||||
rec := httptest.NewRecorder()
|
||||
h.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/_crossbar/hosts", nil))
|
||||
if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") {
|
||||
t.Fatalf("status %d, content-type %q", rec.Code, rec.Header().Get("Content-Type"))
|
||||
}
|
||||
var out map[string]admin.HostView
|
||||
if err := json.Unmarshal(rec.Body.Bytes(), &out); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if a := out["alpha"]; !a.Healthy || len(a.Loaded) != 1 || a.LastOK != "2026-09-25T08:00:00Z" || a.LastErr != "" {
|
||||
t.Errorf("alpha = %+v", a)
|
||||
}
|
||||
if b := out["beta"]; b.Healthy || b.LastOK != "" || b.LastErr != "HTTP 503" || b.Loaded == nil {
|
||||
t.Errorf("beta = %+v (loaded must be [] not null)", b)
|
||||
}
|
||||
if !strings.Contains(rec.Body.String(), `"loaded":[]`) {
|
||||
t.Errorf("beta.loaded must encode as []: %s", rec.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestRoutes(t *testing.T) {
|
||||
h := admin.Handler(testConfig(t), fakeHosts{})
|
||||
rec := httptest.NewRecorder()
|
||||
h.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/_crossbar/routes", nil))
|
||||
var out map[string]admin.RouteView
|
||||
if err := json.Unmarshal(rec.Body.Bytes(), &out); err != nil {
|
||||
t.Fatalf("%v: %s", err, rec.Body.String())
|
||||
}
|
||||
r := out["r"]
|
||||
if len(r.Hosts) != 2 || r.Hosts[0] != "alpha" || r.DefaultModel != "m" {
|
||||
t.Errorf("routes = %+v", out)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMethodsAndUnknown(t *testing.T) {
|
||||
h := admin.Handler(testConfig(t), fakeHosts{})
|
||||
for _, tc := range []struct {
|
||||
method, path string
|
||||
want int
|
||||
}{
|
||||
{http.MethodPost, "/_crossbar/hosts", 405},
|
||||
{http.MethodDelete, "/_crossbar/routes", 405},
|
||||
{http.MethodGet, "/_crossbar/nope", 404},
|
||||
{http.MethodGet, "/_crossbar/", 404},
|
||||
} {
|
||||
rec := httptest.NewRecorder()
|
||||
h.ServeHTTP(rec, httptest.NewRequest(tc.method, tc.path, nil))
|
||||
if rec.Code != tc.want {
|
||||
t.Errorf("%s %s = %d, want %d", tc.method, tc.path, rec.Code, tc.want)
|
||||
}
|
||||
if !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") {
|
||||
t.Errorf("%s %s: errors are JSON too", tc.method, tc.path)
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user