// 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 }