From 4bc5115f24605d871ef394c8d82a82f10b9e5712 Mon Sep 17 00:00:00 2001 From: Kyle Isom Date: Fri, 25 Sep 2026 02:32:50 -0700 Subject: [PATCH] Add the admin endpoints, the crossbar binary and the fake upstream Implemented-By: OpenCode session (model recorded in docs/implementer-log.md) --- cmd/crossbar/main.go | 79 +++++++++++++++++++++++++++ cmd/fakeupstream/main.go | 101 ++++++++++++++++++++++++++++++++++ docs/implementer-log.md | 1 + example.toml | 22 ++++++++ internal/admin/admin.go | 102 +++++++++++++++++++++++++++++++++++ internal/admin/admin_test.go | 99 ++++++++++++++++++++++++++++++++++ 6 files changed, 404 insertions(+) create mode 100644 cmd/crossbar/main.go create mode 100644 cmd/fakeupstream/main.go create mode 100644 example.toml create mode 100644 internal/admin/admin.go create mode 100644 internal/admin/admin_test.go diff --git a/cmd/crossbar/main.go b/cmd/crossbar/main.go new file mode 100644 index 0000000..7c34d9c --- /dev/null +++ b/cmd/crossbar/main.go @@ -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 + } +} diff --git a/cmd/fakeupstream/main.go b/cmd/fakeupstream/main.go new file mode 100644 index 0000000..dbb7801 --- /dev/null +++ b/cmd/fakeupstream/main.go @@ -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: . +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) +} diff --git a/docs/implementer-log.md b/docs/implementer-log.md index 202db05..a38f06a 100644 --- a/docs/implementer-log.md +++ b/docs/implementer-log.md @@ -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 diff --git a/example.toml b/example.toml new file mode 100644 index 0000000..7b105cc --- /dev/null +++ b/example.toml @@ -0,0 +1,22 @@ +# crossbar example configuration. Replace 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.: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.: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"] diff --git a/internal/admin/admin.go b/internal/admin/admin.go new file mode 100644 index 0000000..8bafc89 --- /dev/null +++ b/internal/admin/admin.go @@ -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) +} diff --git a/internal/admin/admin_test.go b/internal/admin/admin_test.go new file mode 100644 index 0000000..80a0dfb --- /dev/null +++ b/internal/admin/admin_test.go @@ -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) + } + } +}