From 927b2cfc4d85e1f9d886dea0ecb195bf4e9addf6 Mon Sep 17 00:00:00 2001 From: Kyle Isom Date: Fri, 25 Sep 2026 09:11:54 -0700 Subject: [PATCH] v2 plan: fake upstream answers /props without counting it (task 01 replacement helper); note the finding Co-Authored-By: Claude Fable 5.1 --- docs/plans/v2/01-props.md | 6 +- docs/plans/v2/README.md | 9 + .../v2/_files/internal/proxy/helpers_test.go | 216 ++++++++++++++++++ 3 files changed, 229 insertions(+), 2 deletions(-) create mode 100644 docs/plans/v2/_files/internal/proxy/helpers_test.go diff --git a/docs/plans/v2/01-props.md b/docs/plans/v2/01-props.md index 72a7f93..05d0507 100644 --- a/docs/plans/v2/01-props.md +++ b/docs/plans/v2/01-props.md @@ -20,6 +20,7 @@ bytes as for the other endpoints. ## Files - Copy: `internal/health/props_test.go` +- Copy (**replaces** v1's): `internal/proxy/helpers_test.go` — the fake upstream now answers `/props` without counting it as a hit, so the v1 proxy tests' exact hit counts still hold once the poller asks for it - Modify: `internal/health/health.go`, `internal/admin/admin.go` (or wherever `HostView` is built), `docs/implementer-log.md` ## Interfaces @@ -49,13 +50,14 @@ Rules the tests check: ## Steps -- [ ] **1.** `git switch master && git switch -c v2`; `cp docs/plans/v2/_files/internal/health/props_test.go internal/health/`. +- [ ] **1.** `git switch master && git switch -c v2`; `cp docs/plans/v2/_files/internal/health/props_test.go internal/health/`; + `cp docs/plans/v2/_files/internal/proxy/helpers_test.go internal/proxy/`. - [ ] **2. See it fail** (compile: `NCtx` undefined). **3. Write the code.** `gofmt -w internal/`. - [ ] **4.** `go test -race -count=1 ./internal/health/ ./internal/admin/` → both `ok`. - [ ] **5.** `make gate` → `gate: ok`. **6.** Row `v2/01-props`; commit. ```sh -git add internal/health internal/admin docs/implementer-log.md +git add internal/health internal/admin internal/proxy/helpers_test.go docs/implementer-log.md git commit ``` diff --git a/docs/plans/v2/README.md b/docs/plans/v2/README.md index 71f8834..7d6dbe9 100644 --- a/docs/plans/v2/README.md +++ b/docs/plans/v2/README.md @@ -52,3 +52,12 @@ reason. the same peer; a wake target whose broadcast address is unroutable (503 within `wait`, no hang); the guard with a body of exactly `MaxBody`. 4. Findings under "Reviews" in `docs/implementer-log.md`, by fault. + +## Changes during the run + +- 2026-09-25, task 01: the new `/props` poll lands on the v1 fake upstream's `/` catch-all, which + counts hits, so two v1 proxy tests with exact hit counts failed. Ornith implemented the task + correctly, did not touch the protected file, and stopped with a `stopped` row — exactly the + procedure. Owner's fault (T19 once more: a new task changed what an earlier given file + measures, and the pre-handover walk missed it). `helpers_test.go` is now a v2 given file that + answers `/props` without counting it; resumed. diff --git a/docs/plans/v2/_files/internal/proxy/helpers_test.go b/docs/plans/v2/_files/internal/proxy/helpers_test.go new file mode 100644 index 0000000..da689db --- /dev/null +++ b/docs/plans/v2/_files/internal/proxy/helpers_test.go @@ -0,0 +1,216 @@ +package proxy_test + +// Test scaffolding shared by proxy_test.go and recorder_test.go: the fake health table, the fake +// llama-server upstream, and the rig that builds a whole crossbar over real HTTP. + +import ( + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/httptest" + "path/filepath" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/config" + "git.wntrmute.dev/kyle/crossbar/internal/health" + "git.wntrmute.dev/kyle/crossbar/internal/lease" + "git.wntrmute.dev/kyle/crossbar/internal/limiter" + "git.wntrmute.dev/kyle/crossbar/internal/proxy" + "git.wntrmute.dev/kyle/crossbar/internal/store" +) + +// fakeHealth is a hand-set health table that also records MarkDown calls. It lived in the v0 +// proxy_test.go; the v1 given test replaces that file, so recorder_test.go (which still exercises +// the nil-lease path through proxy.New) needs it here. +type fakeHealth struct { + mu sync.Mutex + st map[string]health.Status + marked []string +} + +func (f *fakeHealth) Get(name string) (health.Status, bool) { + f.mu.Lock() + defer f.mu.Unlock() + s, ok := f.st[name] + return s, ok +} + +func (f *fakeHealth) MarkDown(name, reason string) { + f.mu.Lock() + defer f.mu.Unlock() + f.marked = append(f.marked, name) + s := f.st[name] + s.Healthy = false + s.LastErr = reason + f.st[name] = s +} + +func (f *fakeHealth) markedHosts() []string { + f.mu.Lock() + defer f.mu.Unlock() + return append([]string{}, f.marked...) +} + +// upstream is a llama-server stand-in: streams N chunks with a delay, reports usage/timings in +// the final chunk, counts requests, and can be slowed down or killed. +type upstream struct { + name string + srv *httptest.Server + hits atomic.Int32 + delay time.Duration + mu sync.Mutex + last recorded +} + +type recorded struct{ method, path, host, xff, body string } + +func newUpstream(t *testing.T, name string) *upstream { + u := &upstream{name: name} + mux := http.NewServeMux() + mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"status":"ok"}`) }) + mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) { + u.mu.Lock() + u.last = recorded{r.Method, r.URL.RequestURI(), r.Host, r.Header.Get("X-Forwarded-For"), ""} + u.mu.Unlock() + fmt.Fprint(w, `{"object":"list","data":[{"id":"shared"},{"id":"`+name+`-only"}]}`) + }) + // The v2 poller also asks /props; it is a health request, not a hit, so it is not counted. + // No n_ctx here: "unknown context" is what the v1 tests and TestUnknownContextNeverBlocks want. + mux.HandleFunc("/props", func(w http.ResponseWriter, r *http.Request) { + fmt.Fprint(w, `{"model_path":"`+name+`"}`) + }) + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + u.hits.Add(1) + b, _ := io.ReadAll(r.Body) + u.mu.Lock() + u.last = recorded{r.Method, r.URL.RequestURI(), r.Host, r.Header.Get("X-Forwarded-For"), string(b)} + u.mu.Unlock() + var req struct { + Stream bool `json:"stream"` + } + _ = json.Unmarshal(b, &req) + w.Header().Set("X-Upstream", name) + time.Sleep(u.delay) + if !req.Stream { + w.Header().Set("Content-Type", "application/json") + fmt.Fprintf(w, `{"choices":[{"message":{"role":"assistant","content":"hi from %s"}}],"usage":{"prompt_tokens":100,"completion_tokens":10,"total_tokens":110},"timings":{"prompt_n":100,"cache_n":90,"predicted_n":10,"predicted_ms":50.0}}`, name) + return + } + w.Header().Set("Content-Type", "text/event-stream") + w.WriteHeader(200) + fl := w.(http.Flusher) + for i := 0; i < 3; i++ { + fmt.Fprintf(w, "data: {\"choices\":[{\"delta\":{\"content\":\"%s %d \"}}]}\n\n", name, i) + fl.Flush() + time.Sleep(10 * time.Millisecond) + } + fmt.Fprint(w, `data: {"choices":[],"usage":{"prompt_tokens":200,"completion_tokens":20,"total_tokens":220},"timings":{"prompt_n":200,"cache_n":150,"predicted_n":20,"predicted_ms":80.0}}`+"\n\n") + fl.Flush() + fmt.Fprint(w, "data: [DONE]\n\n") + }) + u.srv = httptest.NewServer(mux) + t.Cleanup(u.srv.Close) + return u +} + +func (u *upstream) lastReq() recorded { u.mu.Lock(); defer u.mu.Unlock(); return u.last } + +// rig is one crossbar: config, real health table (polled once), real lease table over a real +// SQLite store, real limiter, the proxy handler served by httptest. +type rig struct { + t *testing.T + cfg *config.Config + health *health.Table + store *store.Store + leases *lease.Table + lim *limiter.Limiter + front *httptest.Server +} + +// newRig builds crossbar from a config text where %s placeholders are the upstream base URLs. +func newRig(t *testing.T, cfgText string, ups ...*upstream) *rig { + urls := make([]any, len(ups)) + for i, u := range ups { + urls[i] = u.srv.URL + } + cfg, err := config.Parse(strings.NewReader(fmt.Sprintf(cfgText, urls...))) + if err != nil { + t.Fatal(err) + } + bases := map[string]string{} + for name, h := range cfg.Hosts { + bases[name] = h.BaseURL + } + ht := health.New(bases, time.Hour, nil) + ht.PollOnce(t.Context()) + st, err := store.Open(filepath.Join(t.TempDir(), "crossbar.db")) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = st.Close() }) + lim := limiter.New() + for name, h := range cfg.Hosts { + for model, m := range h.Models { + lim.Configure(name, model, m.Parallel, cfg.QueueMax) + } + } + lt, err := lease.New(st, proxy.HostView(ht, cfg), proxy.Chooser(cfg, ht, lim), cfg.LeaseIdle.Duration) + if err != nil { + t.Fatal(err) + } + p := proxy.New(cfg, ht, lt, lim, st, nil) + front := httptest.NewServer(p) + t.Cleanup(front.Close) + return &rig{t: t, cfg: cfg, health: ht, store: st, leases: lt, lim: lim, front: front} +} + +const twoHosts = ` +listen = "127.0.0.1:1" +queue_max = 1 +lease_idle = "30m" +[hosts.alpha] +base_url = %q +weight = 1.0 +models = { "shared" = { parallel = 2 }, "alpha-only" = { } } +[hosts.beta] +base_url = %q +weight = 2.0 +models = { "shared" = { parallel = 2 }, "beta-only" = { } } +[routes.r] +hosts = ["alpha", "beta"] +default_model = "shared" +[routes.other] +hosts = ["alpha"] +` + +func conversation(id, turn int) string { + msgs := fmt.Sprintf(`{"role":"system","content":"project"},{"role":"user","content":"conversation %d opening"}`, id) + for i := 1; i < turn; i++ { + msgs += fmt.Sprintf(`,{"role":"assistant","content":"ok"},{"role":"user","content":"turn %d"}`, i) + } + return `{"model":"shared","stream":false,"messages":[` + msgs + `]}` +} + +func (r *rig) post(path, body string, hdr ...string) *http.Response { + req, _ := http.NewRequest(http.MethodPost, r.front.URL+path, strings.NewReader(body)) + req.Header.Set("Content-Type", "application/json") + for i := 0; i+1 < len(hdr); i += 2 { + req.Header.Set(hdr[i], hdr[i+1]) + } + resp, err := http.DefaultClient.Do(req) + if err != nil { + r.t.Fatal(err) + } + return resp +} + +func drain(resp *http.Response) string { + b, _ := io.ReadAll(resp.Body) + resp.Body.Close() + return string(b) +}