Add the SQLite store for leases and accounting
Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)
This commit is contained in:
@@ -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 |
|
||||
|
||||
@@ -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
|
||||
)
|
||||
|
||||
@@ -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=
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user