375 lines
11 KiB
Go
375 lines
11 KiB
Go
// 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"
|
|
ReasonCtx = "ctx"
|
|
)
|
|
|
|
// 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)
|
|
}
|
|
|
|
// StatusCount is one (route, host, status) group of live requests, for the metrics endpoint, which
|
|
// needs the per-status breakdown Usage cannot give.
|
|
type StatusCount struct {
|
|
Route, Host string
|
|
Status int
|
|
Count int64
|
|
}
|
|
|
|
// 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)
|
|
}
|
|
|
|
// StatusCounts groups the live requests at or after since by (route, host, status). It reads the
|
|
// requests table only; the rolled-up requests_daily rows are not in it (Prune has moved them out of
|
|
// requests), so counts cover only traffic still in the live table.
|
|
func (s *Store) StatusCounts(since time.Time) ([]StatusCount, error) {
|
|
sinceMs := since.UnixMilli()
|
|
rows, err := s.db.Query(`
|
|
SELECT route, host, status, COUNT(*)
|
|
FROM requests WHERE started >= ?
|
|
GROUP BY route, host, status
|
|
ORDER BY route, host, status`, sinceMs)
|
|
if err != nil {
|
|
return nil, wrap(err)
|
|
}
|
|
defer rows.Close()
|
|
return scanStatusCounts(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
|
|
}
|