diff --git a/docs/implementer-log.md b/docs/implementer-log.md index 474983b..86a9a19 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/01-store | 2026-09-25 | done | 1 | pass | none | Gate passed on the first run once the owner gofmt'd the three previously-un-clean _files plan-tests under docs/plans/v1/_files/; the blocker in the stopped row no longer applies. | ? | | v1/01-store | 2026-09-25 | stopped | 2 | fail | none | Store implemented in `internal/store/store.go` + `schema.go`; `go test -race -count=1 ./internal/store/` is ok and `go vet`/`check-lines` pass. `make gate` cannot print `gate: ok` here: its `gofmt -l .` step flags three committed plan-tests under `docs/plans/v1/_files/` (admin, choose, proxy) that are not gofmt-clean under Go 1.26.7 (formatted by a gofmt that aligns one-line function bodies two columns wider; same diff on a pristine master). They live under `docs/plans/` (must not edit) and the gate covers them; the check cannot be scoped down without weakening it. Code left uncommitted for review. | llama.cpp/ornith-1.5-35b-a3b | | 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. | llama.cpp/ornith-1.5-35b-a3b | | 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`. | llama.cpp/ornith-1.5-35b-a3b | diff --git a/go.mod b/go.mod index ab77dc4..d26e30d 100644 --- a/go.mod +++ b/go.mod @@ -2,4 +2,19 @@ module git.wntrmute.dev/kyle/crossbar go 1.26 -require github.com/BurntSushi/toml v1.6.0 +require ( + github.com/BurntSushi/toml v1.6.0 + modernc.org/sqlite v1.59.0 +) + +require ( + github.com/dustin/go-humanize v1.0.1 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/mattn/go-isatty v0.0.24 // indirect + github.com/ncruces/go-strftime v1.0.0 // indirect + github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect + golang.org/x/sys v0.47.0 // indirect + modernc.org/libc v1.75.7 // indirect + modernc.org/mathutil v1.7.1 // indirect + modernc.org/memory v1.12.1 // indirect +) diff --git a/go.sum b/go.sum index f74b269..e3728f9 100644 --- a/go.sum +++ b/go.sum @@ -1,2 +1,52 @@ -github.com/BurntSushi/toml v1.6.0 h1:dRaEfpa2VI55EwlIW72hMRHdWouJeRF7TPYhI+AUQjk= github.com/BurntSushi/toml v1.6.0/go.mod h1:ukJfTF/6rtPPRCnwkur4qwRxa8vTRFBF0uk2lLoLwho= +github.com/BurntSushi/toml v1.6.0 h1:dRaEfpa2VI55EwlIW72hMRHdWouJeRF7TPYhI+AUQjk= +github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= +github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= +github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3/go.mod h1:jl5iWTm0/hd5PjEYEOuwAJ57L/CibdZfrqZ5XA5GrCk= +github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3 h1:LMLX+LgTNWpfvCBdFebv6EsYotImrt/Ppc5cXIriCSo= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= +github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= +github.com/mattn/go-isatty v0.0.24/go.mod h1:nMCL3Zebbrt45jsMDgnfIwz6ydEQApk5oEI3HqDio6A= +github.com/mattn/go-isatty v0.0.24 h1:tGZZoVgT/KiqK1c8ocVLeDS8BSWMRd47J3Lbz7vsReI= +github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls= +github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE= +golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= +golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= +golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek= +golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= +golang.org/x/tools v0.48.0/go.mod h1:08xX0orndb/F7jJxGDicx061tyd5pcMto75YMAXr6lk= +golang.org/x/tools v0.48.0 h1:3+hClM1aLL5mjMKm5ovokw9epgRXPuu2tILgismM6RE= +modernc.org/ccgo/v4 v4.35.0/go.mod h1:qrVGs9S3Sr2Ztcg9ve+kTAYMp5a3YvWjo+SoN06kJ5I= +modernc.org/ccgo/v4 v4.35.0 h1:F+TUsmw09QxLzmi3aeYYGxjAXarmZaKgj3mKQHNaA8w= +modernc.org/cc/v4 v4.29.2/go.mod h1:OnovgIhbbMXMu1aISnJ0wvVD1KnW+cAUJkIrAWh+kVI= +modernc.org/cc/v4 v4.29.2 h1:h6+9ciCnPKutf4I03CvheAvDLX7+IHlqR6Iy6J+cgd8= +modernc.org/fileutil v1.4.0/go.mod h1:EqdKFDxiByqxLk8ozOxObDSfcVOv/54xDs/DUHdvCUU= +modernc.org/fileutil v1.4.0 h1:j6ZzNTftVS054gi281TyLjHPp6CPHr2KCxEXjEbD6SM= +modernc.org/gc/v2 v2.6.5/go.mod h1:YgIahr1ypgfe7chRuJi2gD7DBQiKSLMPgBQe9oIiito= +modernc.org/gc/v2 v2.6.5 h1:nyqdV8q46KvTpZlsw66kWqwXRHdjIlJOhG6kxiV/9xI= +modernc.org/gc/v3 v3.1.5/go.mod h1:HFK/6AGESC7Ex+EZJhJ2Gni6cTaYpSMmU/cT9RmlfYY= +modernc.org/gc/v3 v3.1.5 h1:21ldfPfRYE31Tb7B3mwAK8gy1AxP4+dKjrOQPfqakoc= +modernc.org/goabi0 v0.2.0/go.mod h1:CEFRnnJhKvWT1c1JTI3Avm+tgOWbkOu5oPA8eH8LnMI= +modernc.org/goabi0 v0.2.0 h1:HvEowk7LxcPd0eq6mVOAEMai46V+i7Jrj13t4AzuNks= +modernc.org/libc v1.75.7/go.mod h1:bO5o2ztHxBb2rjz0PgdHN0sSMw57CgxGFLZ3Qd/QpVQ= +modernc.org/libc v1.75.7 h1:o3DTP9/0p9pKmY2WCKQaySW6wIiZhNM7wc2lUoyhfew= +modernc.org/mathutil v1.7.1/go.mod h1:4p5IwJITfppl0G4sUEDtCr4DthTaT47/N3aT6MhfgJg= +modernc.org/mathutil v1.7.1 h1:GCZVGXdaN8gTqB1Mf/usp1Y/hSqgI2vAGGP4jZMCxOU= +modernc.org/memory v1.12.1/go.mod h1:/JP4VbVC+K5sU2wZi9bHoq2MAkCnrt2r98UGeSK7Mjw= +modernc.org/memory v1.12.1 h1:nFMiWrpStgZczNl6XI9GnIk/rWhYIyHGUaR04pGbp9g= +modernc.org/opt v0.2.0/go.mod h1:03fq9lsNfvkYSfxrfUhZCWPk1lm4cq4N+Bh//bEtgns= +modernc.org/opt v0.2.0 h1:tGyef5ApycA7FSEOMraay9SaTk5zmbx7Tu+cJs4QKZg= +modernc.org/sortutil v1.2.1/go.mod h1:7ZI3a3REbai7gzCLcotuw9AC4VZVpYMjDzETGsSMqJE= +modernc.org/sortutil v1.2.1 h1:+xyoGf15mM3NMlPDnFqrteY07klSFxLElE2PVuWIJ7w= +modernc.org/sqlite v1.59.0/go.mod h1:+paeT2A3iPRHkQDwG7oA6Tk0zQd5woMEI8q7orfry8k= +modernc.org/sqlite v1.59.0 h1:X1es1GpqBlS/5T+vbM4HLUdaa8OtQx468DF2vrx+38A= +modernc.org/strutil v1.2.1/go.mod h1:EHkiggD70koQxjVdSBM3JKM7k6L0FbGE5eymy9i3B9A= +modernc.org/strutil v1.2.1 h1:UneZBkQA+DX2Rp35KcM69cSsNES9ly8mQWD71HKlOA0= +modernc.org/token v1.1.0/go.mod h1:UGzOrNV1mAFSEB63lOFHIpNRUVMvYTc6yu1SMY/XTDM= +modernc.org/token v1.1.0 h1:Xl7Ap9dKaEs5kLoOQeQmPWevfnk/DM5qcLcYlA8ys6Y= diff --git a/internal/store/schema.go b/internal/store/schema.go new file mode 100644 index 0000000..b05d3e7 --- /dev/null +++ b/internal/store/schema.go @@ -0,0 +1,156 @@ +package store + +import ( + "database/sql" + "encoding/json" + "fmt" + "time" +) + +// schema is the durable layout: the lease table, the lease-event log, the +// per-request accounting rows, the host-health poll log, and the daily +// rollup that Prune writes into. Times are Unix milliseconds. +const schema = ` +CREATE TABLE IF NOT EXISTS leases ( + route TEXT NOT NULL, + fp TEXT NOT NULL, + model TEXT NOT NULL, + host TEXT NOT NULL, + state TEXT NOT NULL, + created INTEGER NOT NULL, + last_used INTEGER NOT NULL, + PRIMARY KEY (route, fp, model) +); +CREATE TABLE IF NOT EXISTS lease_events ( + ts INTEGER NOT NULL, + route TEXT NOT NULL, + model TEXT NOT NULL, + from_host TEXT NOT NULL, + to_host TEXT NOT NULL, + reason TEXT NOT NULL +); +CREATE TABLE IF NOT EXISTS requests ( + id INTEGER PRIMARY KEY, + route TEXT NOT NULL, + fp TEXT NOT NULL, + model TEXT NOT NULL, + host TEXT NOT NULL, + started INTEGER NOT NULL, + queued_ms INTEGER NOT NULL, + ttfb_ms INTEGER NOT NULL, + total_ms INTEGER NOT NULL, + status INTEGER NOT NULL, + streamed INTEGER NOT NULL, + prompt_tokens INTEGER NOT NULL, + cached_tokens INTEGER NOT NULL, + completion_tokens INTEGER NOT NULL, + err TEXT NOT NULL +); +CREATE TABLE IF NOT EXISTS host_health ( + ts INTEGER NOT NULL, + host TEXT NOT NULL, + healthy INTEGER NOT NULL, + loaded_models TEXT NOT NULL +); +CREATE TABLE IF NOT EXISTS requests_daily ( + day INTEGER NOT NULL, + route TEXT NOT NULL, + model TEXT NOT NULL, + host TEXT NOT NULL, + requests INTEGER NOT NULL, + errors INTEGER NOT NULL, + busy_ms INTEGER NOT NULL, + queued_ms INTEGER NOT NULL, + prompt_tokens INTEGER NOT NULL, + cached_tokens INTEGER NOT NULL, + completion_tokens INTEGER NOT NULL, + PRIMARY KEY (day, route, model, host) +); +` + +// byColumn maps a grouping to its table column, rejecting anything else. +func byColumn(b By) (string, error) { + switch b { + case ByRoute: + return "route", nil + case ByModel: + return "model", nil + case ByHost: + return "host", nil + default: + return "", fmt.Errorf("store: unknown grouping %q", b) + } +} + +// midnight returns UTC midnight of the day holding ms. +func midnight(ms int64) int64 { + return time.UnixMilli(ms).Truncate(24 * time.Hour).UnixMilli() +} + +// btoi converts a bool to 0/1 for storage. +func btoi(b bool) int64 { + if b { + return 1 + } + return 0 +} + +// encodeLoaded serialises a Loaded slice as JSON, writing nil as []. +func encodeLoaded(loaded []string) ([]byte, error) { + if loaded == nil { + loaded = []string{} + } + return json.Marshal(loaded) +} + +// wrap prefixes a driver error with the package name. +func wrap(err error) error { + if err == nil { + return nil + } + return fmt.Errorf("store: %w", err) +} + +func scanLeases(rows *sql.Rows) ([]Lease, error) { + var out []Lease + for rows.Next() { + var ( + l Lease + created, last int64 + ) + if err := rows.Scan(&l.Route, &l.FP, &l.Model, &l.Host, (*string)(&l.State), &created, &last); err != nil { + return nil, wrap(err) + } + l.Created = time.UnixMilli(created).UTC() + l.LastUsed = time.UnixMilli(last).UTC() + out = append(out, l) + } + return out, rows.Err() +} + +func scanUsage(rows *sql.Rows) ([]UsageRow, error) { + var out []UsageRow + for rows.Next() { + var u UsageRow + if err := rows.Scan(&u.Key, &u.Requests, &u.Errors, &u.BusyMs, &u.QueuedMs, + &u.PromptTokens, &u.CachedTokens, &u.CompletionTokens); err != nil { + return nil, wrap(err) + } + out = append(out, u) + } + return out, rows.Err() +} + +func scanEvents(rows *sql.Rows) ([]LeaseEvent, error) { + var out []LeaseEvent + for rows.Next() { + var e LeaseEvent + var ts int64 + if err := rows.Scan(&ts, &e.Route, &e.Model, &e.FromHost, &e.ToHost, (*string)(&e.Reason)); err != nil { + return nil, wrap(err) + } + e.TS = time.UnixMilli(ts).UTC() + out = append(out, e) + } + return out, rows.Err() +} diff --git a/internal/store/store.go b/internal/store/store.go new file mode 100644 index 0000000..daaf39a --- /dev/null +++ b/internal/store/store.go @@ -0,0 +1,348 @@ +// Package store is crossbar's durable state: the lease table and the +// accounting log (requests, lease events, host-health polls) with rollup +// queries that answer per-route / per-model / per-host usage. All times are +// stored as Unix milliseconds (INTEGER) and returned as time.Time in UTC. +package store + +import ( + "database/sql" + "fmt" + "os" + "path/filepath" + "time" + + _ "modernc.org/sqlite" +) + +// State is the lease's placement state. +type State string + +const ( + Active State = "active" + Pinned State = "pinned" +) + +// Reasons a lease event carries. +const ( + ReasonNew = "new" + ReasonUnhealthy = "unhealthy" + ReasonIdle = "idle" + ReasonPin = "pin" + ReasonRelease = "release" + ReasonDrain = "drain" +) + +// By selects the grouping column of a Usage query. +type By string + +const ( + ByRoute By = "route" + ByModel By = "model" + ByHost By = "host" +) + +// Lease is one routed model on one host. +type Lease struct { + Route, FP, Model, Host string + State State + Created, LastUsed time.Time +} + +// LeaseEvent records a change to a lease. +type LeaseEvent struct { + TS time.Time + Route, Model, FromHost, ToHost string + Reason string +} + +// Request is one proxied request, for the accounting log. +type Request struct { + Route, FP, Model, Host string + Started time.Time + QueuedMs, TTFBMs, TotalMs int64 + Status int + Streamed bool + PromptTokens, CachedTokens, CompletionTokens int64 + Err string +} + +// HostHealth is one poller observation of a host. +type HostHealth struct { + TS time.Time + Host string + Healthy bool + Loaded []string +} + +// UsageRow is one group of a Usage query. +type UsageRow struct { + Key string `json:"key"` + Requests int64 `json:"requests"` + Errors int64 `json:"errors"` + BusyMs int64 `json:"busy_ms"` + QueuedMs int64 `json:"queued_ms"` + PromptTokens int64 `json:"prompt_tokens"` + CachedTokens int64 `json:"cached_tokens"` + CompletionTokens int64 `json:"completion_tokens"` +} + +// CacheHitRatio is CachedTokens / PromptTokens, or 0 when there were no prompt +// tokens to spend. +func (u UsageRow) CacheHitRatio() float64 { + if u.PromptTokens == 0 { + return 0 + } + return float64(u.CachedTokens) / float64(u.PromptTokens) +} + +// Store holds the SQLite connection to crossbar's durable state. +type Store struct { + db *sql.DB +} + +// Open connects to path in WAL mode and creates the tables if they are missing. +// It fails if the directory that holds the file does not exist. +func Open(path string) (*Store, error) { + if dir := filepath.Dir(path); dir != "" { + if _, err := os.Stat(dir); err != nil { + return nil, fmt.Errorf("store: %s: %w", dir, err) + } + } + db, err := sql.Open("sqlite", "file:"+path+"?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)") + if err != nil { + return nil, wrap(err) + } + if _, err := db.Exec(schema); err != nil { + _ = db.Close() + return nil, wrap(err) + } + return &Store{db: db}, nil +} + +// Close releases the connection. +func (s *Store) Close() error { + return wrap(s.db.Close()) +} + +// JournalMode reports the active journal mode ("wal"). +func (s *Store) JournalMode() string { + var mode string + _ = s.db.QueryRow(`PRAGMA journal_mode`).Scan(&mode) + return mode +} + +// SaveLease inserts or replaces a lease on (route, fp, model). +func (s *Store) SaveLease(l Lease) error { + _, err := s.db.Exec(` + INSERT OR REPLACE INTO leases (route, fp, model, host, state, created, last_used) + VALUES (?, ?, ?, ?, ?, ?, ?)`, + l.Route, l.FP, l.Model, l.Host, string(l.State), + l.Created.UnixMilli(), l.LastUsed.UnixMilli()) + return wrap(err) +} + +// DeleteLease removes the lease identified by (route, fp, model). +func (s *Store) DeleteLease(route, fp, model string) error { + _, err := s.db.Exec(`DELETE FROM leases WHERE route = ? AND fp = ? AND model = ?`, route, fp, model) + return wrap(err) +} + +// ListLeases returns every lease, ordered by (route, fp, model). +func (s *Store) ListLeases() ([]Lease, error) { + rows, err := s.db.Query(` + SELECT route, fp, model, host, state, created, last_used + FROM leases ORDER BY route, fp, model`) + if err != nil { + return nil, wrap(err) + } + defer rows.Close() + return scanLeases(rows) +} + +// RecordEvent appends a lease event. +func (s *Store) RecordEvent(e LeaseEvent) error { + _, err := s.db.Exec(` + INSERT INTO lease_events (ts, route, model, from_host, to_host, reason) + VALUES (?, ?, ?, ?, ?, ?)`, + e.TS.UnixMilli(), e.Route, e.Model, e.FromHost, e.ToHost, e.Reason) + return wrap(err) +} + +// RecordRequest appends one request to the accounting log. +func (s *Store) RecordRequest(r Request) error { + _, err := s.db.Exec(` + INSERT INTO requests (route, fp, model, host, started, queued_ms, ttfb_ms, total_ms, status, streamed, prompt_tokens, cached_tokens, completion_tokens, err) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + r.Route, r.FP, r.Model, r.Host, + r.Started.UnixMilli(), r.QueuedMs, r.TTFBMs, r.TotalMs, + r.Status, btoi(r.Streamed), + r.PromptTokens, r.CachedTokens, r.CompletionTokens, r.Err) + return wrap(err) +} + +// RecordHostHealth appends one host-health observation. +func (s *Store) RecordHostHealth(h HostHealth) error { + b, err := encodeLoaded(h.Loaded) + if err != nil { + return wrap(err) + } + _, err = s.db.Exec(` + INSERT INTO host_health (ts, host, healthy, loaded_models) + VALUES (?, ?, ?, ?)`, + h.TS.UnixMilli(), h.Host, btoi(h.Healthy), string(b)) + return wrap(err) +} + +// Usage sums requests at or after since, grouped by by. A zero since means all +// time. It counts the live requests table plus the requests_daily rollup whose +// day is at or after since. Results are ordered by Key. +func (s *Store) Usage(since time.Time, by By) ([]UsageRow, error) { + col, err := byColumn(by) + if err != nil { + return nil, err + } + sinceMs := since.UnixMilli() + query := ` + SELECT key, + SUM(requests) AS requests, + SUM(errors) AS errors, + SUM(busy_ms) AS busy_ms, + SUM(queued_ms) AS queued_ms, + SUM(prompt_tokens) AS prompt_tokens, + SUM(cached_tokens) AS cached_tokens, + SUM(completion_tokens) AS completion_tokens + FROM ( + SELECT ` + col + ` AS key, 1 AS requests, + CASE WHEN status >= 400 THEN 1 ELSE 0 END AS errors, + total_ms AS busy_ms, queued_ms, + prompt_tokens, cached_tokens, completion_tokens + FROM requests WHERE started >= ? + UNION ALL + SELECT ` + col + ` AS key, requests, errors, busy_ms, queued_ms, + prompt_tokens, cached_tokens, completion_tokens + FROM requests_daily WHERE day >= ? + ) + GROUP BY key + ORDER BY key` + rows, err := s.db.Query(query, sinceMs, sinceMs) + if err != nil { + return nil, wrap(err) + } + defer rows.Close() + return scanUsage(rows) +} + +// Events returns lease events at or after since, oldest first, at most limit. +func (s *Store) Events(since time.Time, limit int) ([]LeaseEvent, error) { + rows, err := s.db.Query(` + SELECT ts, route, model, from_host, to_host, reason + FROM lease_events WHERE ts >= ? + ORDER BY ts ASC LIMIT ?`, since.UnixMilli(), limit) + if err != nil { + return nil, wrap(err) + } + defer rows.Close() + return scanEvents(rows) +} + +// Prune moves every request older than now-retention into requests_daily (adding +// into the existing daily row for that day/route/model/host), deletes the live +// rows, and returns how many were removed. All in one transaction. +func (s *Store) Prune(now time.Time, retention time.Duration) (int64, error) { + threshold := now.Add(-retention).UnixMilli() + tx, err := s.db.Begin() + if err != nil { + return 0, wrap(err) + } + defer func() { _ = tx.Rollback() }() + + rows, err := tx.Query(` + SELECT route, model, host, started, total_ms, queued_ms, + status, prompt_tokens, cached_tokens, completion_tokens + FROM requests WHERE started < ?`, threshold) + if err != nil { + _ = tx.Rollback() + return 0, wrap(err) + } + + type bucket struct { + day int64 + route string + model string + host string + } + type agg struct { + requests int64 + errors int64 + busyMs int64 + queuedMs int64 + prompt int64 + cached int64 + done int64 + } + aggs := map[bucket]*agg{} + count := int64(0) + for rows.Next() { + var ( + route, model, host string + started int64 + totalMs, queuedMs int64 + status int + prompt, cached, done int64 + ) + if err := rows.Scan(&route, &model, &host, &started, &totalMs, &queuedMs, &status, &prompt, &cached, &done); err != nil { + _ = rows.Close() + _ = tx.Rollback() + return 0, wrap(err) + } + count++ + k := bucket{day: midnight(started), route: route, model: model, host: host} + a := aggs[k] + if a == nil { + a = &agg{} + aggs[k] = a + } + a.requests++ + if status >= 400 { + a.errors++ + } + a.busyMs += totalMs + a.queuedMs += queuedMs + a.prompt += prompt + a.cached += cached + a.done += done + } + if err := rows.Err(); err != nil { + _ = rows.Close() + _ = tx.Rollback() + return 0, wrap(err) + } + _ = rows.Close() + + upsert := ` + INSERT INTO requests_daily (day, route, model, host, requests, errors, busy_ms, queued_ms, prompt_tokens, cached_tokens, completion_tokens) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (day, route, model, host) DO UPDATE SET + requests = requests + excluded.requests, + errors = errors + excluded.errors, + busy_ms = busy_ms + excluded.busy_ms, + queued_ms = queued_ms + excluded.queued_ms, + prompt_tokens = prompt_tokens + excluded.prompt_tokens, + cached_tokens = cached_tokens + excluded.cached_tokens, + completion_tokens = completion_tokens + excluded.completion_tokens` + for k, a := range aggs { + if _, err := tx.Exec(upsert, k.day, k.route, k.model, k.host, + a.requests, a.errors, a.busyMs, a.queuedMs, a.prompt, a.cached, a.done); err != nil { + _ = tx.Rollback() + return 0, wrap(err) + } + } + if _, err := tx.Exec(`DELETE FROM requests WHERE started < ?`, threshold); err != nil { + _ = tx.Rollback() + return 0, wrap(err) + } + if err := tx.Commit(); err != nil { + return 0, wrap(err) + } + return count, nil +} diff --git a/internal/store/store_test.go b/internal/store/store_test.go new file mode 100644 index 0000000..e7f2c28 --- /dev/null +++ b/internal/store/store_test.go @@ -0,0 +1,185 @@ +package store_test + +import ( + "path/filepath" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/store" +) + +func open(t *testing.T, dir string) *store.Store { + s, err := store.Open(filepath.Join(dir, "crossbar.db")) + if err != nil { + t.Fatalf("Open: %v", err) + } + t.Cleanup(func() { _ = s.Close() }) + return s +} + +func TestOpenIsIdempotentAndWAL(t *testing.T) { + dir := t.TempDir() + s := open(t, dir) + if got := s.JournalMode(); got != "wal" { + t.Errorf("journal_mode = %q, want wal", got) + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + open(t, dir) // second open on the same file must not fail on existing tables +} + +func TestLeasesSurviveReopen(t *testing.T) { + dir := t.TempDir() + s := open(t, dir) + now := time.Date(2026, 9, 25, 10, 0, 0, 0, time.UTC) + l := store.Lease{Route: "opencode-a", FP: "abc", Model: "m", Host: "alpha", State: store.Active, Created: now, LastUsed: now} + if err := s.SaveLease(l); err != nil { + t.Fatal(err) + } + l2 := l + l2.FP = "def" + l2.Host = "beta" + l2.State = store.Pinned + if err := s.SaveLease(l2); err != nil { + t.Fatal(err) + } + // Saving the same key again replaces, not duplicates. + l.Host = "beta" + l.LastUsed = now.Add(time.Minute) + if err := s.SaveLease(l); err != nil { + t.Fatal(err) + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + s = open(t, dir) + got, err := s.ListLeases() + if err != nil { + t.Fatal(err) + } + if len(got) != 2 { + t.Fatalf("ListLeases = %d rows, want 2: %+v", len(got), got) + } + byFP := map[string]store.Lease{} + for _, x := range got { + byFP[x.FP] = x + } + if a := byFP["abc"]; a.Host != "beta" || !a.LastUsed.Equal(now.Add(time.Minute)) || a.State != store.Active { + t.Errorf("abc = %+v", a) + } + if d := byFP["def"]; d.State != store.Pinned || d.Host != "beta" { + t.Errorf("def = %+v", d) + } + if err := s.DeleteLease("opencode-a", "abc", "m"); err != nil { + t.Fatal(err) + } + got, _ = s.ListLeases() + if len(got) != 1 || got[0].FP != "def" { + t.Errorf("after delete: %+v", got) + } +} + +func TestEventsAndRequestsAndUsage(t *testing.T) { + s := open(t, t.TempDir()) + t0 := time.Date(2026, 9, 25, 10, 0, 0, 0, time.UTC) + must := func(err error) { + if err != nil { + t.Fatal(err) + } + } + must(s.RecordEvent(store.LeaseEvent{TS: t0, Route: "r1", Model: "m", FromHost: "", ToHost: "alpha", Reason: store.ReasonNew})) + must(s.RecordEvent(store.LeaseEvent{TS: t0.Add(time.Hour), Route: "r1", Model: "m", FromHost: "alpha", ToHost: "beta", Reason: store.ReasonUnhealthy})) + reqs := []store.Request{ + {Route: "r1", FP: "a", Model: "m", Host: "alpha", Started: t0, QueuedMs: 0, TTFBMs: 100, TotalMs: 1000, Status: 200, Streamed: true, PromptTokens: 1000, CachedTokens: 900, CompletionTokens: 50}, + {Route: "r1", FP: "a", Model: "m", Host: "alpha", Started: t0.Add(time.Minute), QueuedMs: 40, TTFBMs: 120, TotalMs: 2000, Status: 200, Streamed: true, PromptTokens: 1100, CachedTokens: 1000, CompletionTokens: 60}, + {Route: "r2", FP: "b", Model: "m", Host: "beta", Started: t0.Add(2 * time.Minute), TotalMs: 500, Status: 502, Err: "upstream failed"}, + {Route: "r2", FP: "b", Model: "m", Host: "beta", Started: t0.Add(-48 * time.Hour), TotalMs: 300, Status: 200, PromptTokens: 10, CompletionTokens: 5}, + } + for _, r := range reqs { + must(s.RecordRequest(r)) + } + must(s.RecordHostHealth(store.HostHealth{TS: t0, Host: "alpha", Healthy: true, Loaded: []string{"m"}})) + + rows, err := s.Usage(t0.Add(-time.Hour), store.ByRoute) + must(err) + if len(rows) != 2 { + t.Fatalf("Usage by route since t0-1h: %d rows, want 2 (r1, r2): %+v", len(rows), rows) + } + byKey := map[string]store.UsageRow{} + for _, r := range rows { + byKey[r.Key] = r + } + r1 := byKey["r1"] + if r1.Requests != 2 || r1.Errors != 0 || r1.BusyMs != 3000 || r1.QueuedMs != 40 { + t.Errorf("r1 = %+v", r1) + } + if r1.PromptTokens != 2100 || r1.CachedTokens != 1900 || r1.CompletionTokens != 110 { + t.Errorf("r1 tokens = %+v", r1) + } + if got := r1.CacheHitRatio(); got < 0.904 || got > 0.905 { + t.Errorf("r1 cache hit ratio = %v, want 1900/2100", got) + } + r2 := byKey["r2"] + if r2.Requests != 1 || r2.Errors != 1 || r2.BusyMs != 500 { + t.Errorf("r2 = %+v (the 48h-old request is outside since)", r2) + } + if r2.CacheHitRatio() != 0 { + t.Errorf("no prompt tokens: ratio must be 0, got %v", r2.CacheHitRatio()) + } + byHost, err := s.Usage(time.Time{}, store.ByHost) + must(err) + if len(byHost) != 2 { + t.Errorf("by host, all time: %+v", byHost) + } + for _, r := range byHost { + if r.Key == "beta" && r.Requests != 2 { + t.Errorf("beta all-time requests = %d, want 2", r.Requests) + } + } + byModel, err := s.Usage(time.Time{}, store.ByModel) + must(err) + if len(byModel) != 1 || byModel[0].Key != "m" || byModel[0].Requests != 4 { + t.Errorf("by model: %+v", byModel) + } + ev, err := s.Events(t0.Add(-time.Minute), 10) + must(err) + if len(ev) != 2 || ev[0].Reason != store.ReasonNew || ev[1].ToHost != "beta" { + t.Errorf("events = %+v", ev) + } +} + +func TestPruneRollsUpOldRequests(t *testing.T) { + s := open(t, t.TempDir()) + t0 := time.Date(2026, 9, 25, 10, 0, 0, 0, time.UTC) + old := t0.Add(-200 * 24 * time.Hour) + for i := 0; i < 3; i++ { + if err := s.RecordRequest(store.Request{Route: "r", Model: "m", Host: "h", Started: old.Add(time.Duration(i) * time.Minute), TotalMs: 100, Status: 200, PromptTokens: 10, CachedTokens: 5, CompletionTokens: 1}); err != nil { + t.Fatal(err) + } + } + if err := s.RecordRequest(store.Request{Route: "r", Model: "m", Host: "h", Started: t0, TotalMs: 100, Status: 200}); err != nil { + t.Fatal(err) + } + n, err := s.Prune(t0, 180*24*time.Hour) + if err != nil { + t.Fatal(err) + } + if n != 3 { + t.Errorf("Prune removed %d rows, want 3", n) + } + rows, _ := s.Usage(time.Time{}, store.ByRoute) + if len(rows) != 1 || rows[0].Requests != 4 || rows[0].PromptTokens != 30 { + t.Errorf("usage must still include pruned traffic through the daily rollup: %+v", rows) + } + live, _ := s.Usage(old.Add(24*time.Hour), store.ByRoute) + if len(live) != 1 || live[0].Requests != 1 { + t.Errorf("recent-only usage = %+v", live) + } +} + +func TestBadPath(t *testing.T) { + if _, err := store.Open(filepath.Join(t.TempDir(), "no", "such", "dir", "x.db")); err == nil { + t.Fatal("Open must fail when the directory does not exist") + } +}