157 lines
4.1 KiB
Go
157 lines
4.1 KiB
Go
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()
|
|
}
|