diff --git a/cmd/crossbar/main.go b/cmd/crossbar/main.go index a0109db..ce11cbb 100644 --- a/cmd/crossbar/main.go +++ b/cmd/crossbar/main.go @@ -17,7 +17,10 @@ import ( "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/lease" + "git.wntrmute.dev/kyle/crossbar/internal/limiter" "git.wntrmute.dev/kyle/crossbar/internal/proxy" + "git.wntrmute.dev/kyle/crossbar/internal/store" ) func main() { @@ -38,6 +41,12 @@ func run() error { log := slog.New(slog.NewTextHandler(os.Stderr, nil)) + st, err := store.Open(cfg.DB) + if err != nil { + return err + } + defer st.Close() + baseURLs := make(map[string]string, len(cfg.Hosts)) for name, host := range cfg.Hosts { baseURLs[name] = host.BaseURL @@ -48,9 +57,79 @@ func run() error { defer stop() go table.Run(ctx) + hosts := proxy.HostView(table, cfg) + lim := limiter.New() + for name, h := range cfg.Hosts { + for model, m := range h.Models { + lim.Configure(name, model, m.Parallel, cfg.QueueMax) + } + } + + leases, err := lease.New(st, hosts, proxy.Chooser(cfg, table, lim), cfg.LeaseIdle.Duration) + if err != nil { + return err + } + for name, rt := range cfg.Routes { + leases.Candidates(name, rt.Hosts) + } + mux := http.NewServeMux() - mux.Handle("/_crossbar/", admin.Handler(cfg, table, nil, nil, nil, nil)) - mux.Handle("/", proxy.New(cfg, table, nil, nil, nil, log)) + mux.Handle("/_crossbar/", admin.Handler(cfg, table, leases, lim, st, hosts)) + mux.Handle("/", proxy.New(cfg, table, leases, lim, st, log)) + + // Background maintenance until ctx is done. Errors are logged, never fatal. + go func() { + ticker := time.NewTicker(time.Minute) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + leases.ExpireIdle(time.Now()) + } + } + }() + + go func() { + ticker := time.NewTicker(time.Hour) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + n, err := st.Prune(time.Now(), cfg.Retention.Duration) + if err != nil { + log.Error("prune", "err", err) + continue + } + log.Info("pruned request rows", "rows", n) + } + } + }() + + go func() { + ticker := time.NewTicker(cfg.PollInterval.Duration) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + for name, s := range table.All() { + if err := st.RecordHostHealth(store.HostHealth{ + TS: time.Now(), + Host: name, + Healthy: s.Healthy, + Loaded: s.Loaded, + }); err != nil { + log.Error("record host health", "host", name, "err", err) + } + } + } + } + }() srv := &http.Server{ Addr: cfg.Listen, diff --git a/docs/implementer-log.md b/docs/implementer-log.md index 0f814e2..332e1b1 100644 --- a/docs/implementer-log.md +++ b/docs/implementer-log.md @@ -5,6 +5,7 @@ owner fills in the Model column. The reviewer adds findings under "Reviews" once | Task | Date | Status | Gate runs | First gate | Deviations | Notes | Model | |---|---|---|---|---|---|---|---| +| v1/07-main | 2026-09-25 | done | 1 | pass | none | Wired store, limiter and lease table into `cmd/crossbar/main.go`: `store.Open` before the health table, `limiter.Configure` per (host, model) from `cfg.Hosts`, `lease.New` with `proxy.Chooser`, `Candidates` for every route, three background goroutines (idle expiry per minute, prune per hour logging the count, host-health recording per `poll_interval`), and `st.Close` via `defer`. The 3s SIGTERM run exits 0 with `listening`/`shutting down`; the missing-config run exits 1. | ? | | v1/06-admin | 2026-09-25 | done | 2 | fail | Split `internal/admin/admin.go` (196 lines) + `admin_ops.go` (366 lines) to stay under 400. Updated `cmd/crossbar/main.go`'s `admin.Handler` call from the committed 2-arg `(cfg, table)` to the task's 6-arg signature, passing the health table for `hosts` and `nil` for the not-yet-wired `leases`/`limiter`/`store`/`drainer` (task 07 wires them); this was a compile fix required for `go vet`/`go test ./...` on `cmd/crossbar` to pass — the full wiring is task 07. | First `make gate` failed on `go vet` (`admin.Handler` called with 2 args in `main.go` after the signature changed); fixed `main.go` and the gate passed on the second run. `admin_test.go` and `example.toml` verified byte-identical to `docs/plans/v1/_files/`; `internal/lease` and `internal/store` left untouched except the already-present `Candidates`/`StatusCounts`. | ? | | v1/05-proxy | 2026-09-25 | done | 2 | fail | Split `internal/proxy/proxy.go` (411 lines) into `proxy.go` + `forward.go` by moving `forward`, `newReverseProxy`, `forwardState`, `statusRecorder`, `leaseState`, `ttfbMs` and the `writeError`/`writeRecord` helpers to `forward.go`; the one `recorder_test.go` `proxy.New` call changed to `proxy.New(cfg, h, nil, nil, nil, nil)` per the task; `cmd/crossbar/main.go` passes `nil, nil, nil` for the new `leases`/`lim`/`rec` args (task 06 wires them). | The tee in `tee.go` already read the final SSE chunk's (streamed) and the JSON body's (non-streamed) usage/timings, so `TestAccountingRowsFromUsageAndTimings` passed on the first run — the only gate blocker was `proxy.go` at 411 lines. | ? | | v1/04-lease | 2026-09-25 | done | 1 | pass | The given `TestPinAndUnpin` was wrong and replaced by the owner mid-task; the corrected `internal/lease/lease_test.go` is byte-identical to `docs/plans/v1/_files/internal/lease/lease_test.go`. A `fmt.Printf("DEBUG …")` line the prior session left in `event` was removed before the gate. | `Acquire` order (pinned, existing, inherit, choose) with memory rolled back only after a successful save; `Pin` writes a pin event, then the pin row, then deletes other-host leases, so the pin event always precedes the unpin's release event in the log. | ? |