Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
379 changes: 379 additions & 0 deletions pkg/api/SUBSTRATE-FINDINGS.md

Large diffs are not rendered by default.

117 changes: 117 additions & 0 deletions pkg/api/backtest.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,117 @@
package api

import (
"fmt"
"time"

"github.com/agenticode/kilter/pkg/backtest"
)

// Replaying a cluster's own history.
//
// The time-keyed snapshot bucket in pkg/store is only half of what
// `kilter backtest --cluster` was refused for. The other half is this
// function, and it exists because a populated bucket does NOT by itself make
// the refusal safe to delete.
//
// backtest.Run over an empty or too-short history does not fail. It returns a
// Scorecard with the same shape, the same field names and the same confident
// tone as a real one — `snapshots 0`, `regret $0.00` — and an operator reading
// that has no way to tell it means "nothing was replayed" rather than "the
// policy is perfect". That is strictly worse than the refusal it replaces, so
// the precondition is checked here, once, using backtest's OWN report of how
// much it scored rather than a re-derived predicate that could drift from it.

// MinReplaySnapshots is the floor below which there is no history to replay.
// Two is the arithmetic minimum for anything to have changed; Instants below
// is the real gate.
const MinReplaySnapshots = 2

// ErrNoHistory reports a brain with no persistent store: without one there is
// no snapshot history at all, only the single in-memory latest snapshot.
type ErrNoHistory struct{ Cluster string }

func (e ErrNoHistory) Error() string {
return fmt.Sprintf("backtest %s: this brain has no persistent store, so no snapshot history is kept; "+
"start it with a --db path and let it ingest", e.Cluster)
}

// ErrHistoryTooShort reports a history that exists but cannot support a
// replay. It names what is there and what was needed, because "not enough
// history" without numbers is indistinguishable from a bug.
type ErrHistoryTooShort struct {
Cluster string
From, To time.Time
Snapshots int
Instants int
Interval time.Duration
Horizon time.Duration
}

func (e ErrHistoryTooShort) Error() string {
span := e.To.Sub(e.From)
if e.Snapshots < MinReplaySnapshots {
return fmt.Sprintf("backtest %s: refused — the retained history holds %d snapshot(s) in [%s, %s); "+
"scoring that would produce a scorecard shaped exactly like a real one, and its zeros would read "+
"as a verdict rather than as an empty replay",
e.Cluster, e.Snapshots, e.From.UTC().Format(time.RFC3339), e.To.UTC().Format(time.RFC3339))
}
return fmt.Sprintf("backtest %s: refused — %d snapshot(s) span %v, which yields no decision instant at a "+
"%v interval with a %v horizon. Nothing would be replayed, and a scorecard over nothing reads as a "+
"perfect policy. Widen the window or shorten --interval",
e.Cluster, e.Snapshots, span.Round(time.Minute), e.Interval, e.Horizon)
}

// Backtest replays a cluster's own retained history through the production
// decision path and scores it.
//
// The policy under test is THIS BRAIN'S policy: the recommender and planner
// configs it actually runs with, so the scorecard answers "how good is what
// is running here" rather than "how good is the default". The evidence
// substrate is this brain's too, which is what lets the harness score real
// OOMKills rather than only counterfactual memory violations.
//
// It refuses rather than returns an empty scorecard. See the file comment.
func (b *Brain) Backtest(cluster string, from, to time.Time, horizon time.Duration, scoring backtest.Config) (*backtest.Scorecard, error) {
if b.st == nil {
return nil, ErrNoHistory{Cluster: cluster}
}
if err := checkExplainWindow(from, to); err != nil {
return nil, err
}
snaps, err := b.st.Snapshots(cluster, from, to)
if err != nil {
return nil, err
}
if len(snaps) < MinReplaySnapshots {
return nil, ErrHistoryTooShort{
Cluster: cluster, From: from, To: to, Snapshots: len(snaps),
Interval: scoring.DecisionInterval, Horizon: horizon,
}
}
h := &backtest.Harness{
Evidence: b.mem,
// The store is the SnapshotSource: satisfied structurally, so the
// harness replays the same rows the bucket retained, thinning and all.
History: b.st,
Rec: b.cfg.Recommend,
Plan: b.cfg.Plan,
Catalog: b.catalog,
Scoring: scoring,
}
sc, err := h.Run(cluster, from, to, horizon)
if err != nil {
return nil, err
}
// backtest's own coverage report, not a predicate reimplemented here: if
// it scored no instant, there is no scorecard to show.
if sc.Instants == 0 {
return nil, ErrHistoryTooShort{
Cluster: cluster, From: from, To: to,
Snapshots: sc.Snapshots, Instants: sc.Instants,
Interval: time.Duration(sc.DecisionIntervalHours * float64(time.Hour)),
Horizon: horizon,
}
}
return sc, nil
}
54 changes: 48 additions & 6 deletions pkg/api/brain.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import (
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promhttp"

"github.com/agenticode/kilter/pkg/evidence"
"github.com/agenticode/kilter/pkg/forecast"
"github.com/agenticode/kilter/pkg/model"
"github.com/agenticode/kilter/pkg/plan"
Expand Down Expand Up @@ -52,7 +53,10 @@ type BrainConfig struct {
ForecasterURL string
Recommend recommend.Config
Plan plan.Config
Logger *slog.Logger
// Evidence bounds the L0 evidence substrate the explanation routes are
// served from. The zero value takes evidence.DefaultConfig().
Evidence evidence.Config
Logger *slog.Logger
}

func (c BrainConfig) withDefaults() BrainConfig {
Expand Down Expand Up @@ -91,6 +95,12 @@ type Brain struct {

forecaster *forecast.RemoteForecaster // nil = built-in models only

// mem is the evidence substrate every explanation is grounded in and
// every citation is re-resolved against. Never nil — see substrate.go.
mem *evidence.Memory
// subs holds per-cluster change detection for the substrate's events.
subs map[string]*substrateState

m brainMetrics
}

Expand Down Expand Up @@ -140,8 +150,24 @@ func NewBrain(cfg BrainConfig, catalog *pricing.Catalog, st *store.Store) (*Brai
demand: map[string]*demandTracker{},
ledgers: map[string]*ledgerState{},
approvals: map[string]*approvalState{},
subs: map[string]*substrateState{},
m: newBrainMetrics(),
}
// The substrate comes first: everything below may write into it, and a
// brain that cannot build one cannot verify an answer it serves.
if st != nil {
mem, err := restoreEvidence(b.cfg.Evidence, st.EvidenceCheckpoints())
if err != nil {
return nil, err
}
b.mem = mem
} else {
mem, err := evidence.NewMemory(b.cfg.Evidence)
if err != nil {
return nil, fmt.Errorf("api: evidence config: %w", err)
}
b.mem = mem
}
if b.cfg.ForecasterURL != "" {
rf, err := forecast.NewRemoteForecaster(b.cfg.ForecasterURL)
if err != nil {
Expand All @@ -166,7 +192,9 @@ func NewBrain(cfg BrainConfig, catalog *pricing.Catalog, st *store.Store) (*Brai
b.lastSnap[c] = snap
}
}
b.cfg.Logger.Info("brain restored", "clusters", len(clusters))
mstats := b.mem.Stats()
b.cfg.Logger.Info("brain restored", "clusters", len(clusters),
"evidenceEvents", mstats.Events, "evidenceSubjects", mstats.SeriesSubjects)
}
return b, nil
}
Expand Down Expand Up @@ -221,19 +249,31 @@ func (b *Brain) Ingest(snap *model.ClusterSnapshot) error {
if err := b.st.SaveSnapshot(snap); err != nil {
b.cfg.Logger.Error("persist snapshot", "err", err)
}
if count%b.cfg.CheckpointEvery == 0 {
if err := b.st.SaveRecommenderState(snap.ClusterID, r.Checkpoint()); err != nil {
b.cfg.Logger.Error("persist recommender", "err", err)
}
// The time-keyed history `kilter backtest --cluster` replays. It is
// separate from SaveSnapshot because that one is keyed by cluster and
// keeps exactly one row; this one is keyed by cluster AND time, and
// thins itself to its retention cadence (pkg/store/history.go).
if err := b.st.SaveSnapshotAt(snap); err != nil {
b.cfg.Logger.Error("persist snapshot history", "err", err)
}
}

cost := b.catalog.SnapshotCost(snap)
b.observeIntoSubstrate(snap, cost.HourlyUSD)
b.ledgerFor(snap.ClusterID).addCost(snap.Timestamp, cost.HourlyUSD)
b.m.snapshots.WithLabelValues(snap.ClusterID).Inc()
b.m.containers.WithLabelValues(snap.ClusterID).Set(float64(r.StateCount()))
b.m.costHourly.WithLabelValues(snap.ClusterID).Set(cost.HourlyUSD)
b.m.ingestSec.Observe(time.Since(start).Seconds())

// Checkpointing is last: it is the most expensive thing in this function
// and the only part whose failure costs nothing already observed.
if b.st != nil && count%b.cfg.CheckpointEvery == 0 {
if err := b.st.SaveRecommenderState(snap.ClusterID, r.Checkpoint()); err != nil {
b.cfg.Logger.Error("persist recommender", "err", err)
}
b.saveEvidence()
}
return nil
}

Expand Down Expand Up @@ -390,6 +430,7 @@ func (b *Brain) Handler() http.Handler {
writeJSON(w, http.StatusOK, map[string]any{"insights": ins})
}))
b.registerTrustRoutes(mux)
b.registerExplainRoutes(mux)
mux.HandleFunc("GET /api/v1/clusters/{id}/cost", b.auth(func(w http.ResponseWriter, r *http.Request) {
snap := b.snapshotFor(r.PathValue("id"))
if snap == nil {
Expand Down Expand Up @@ -518,6 +559,7 @@ func (b *Brain) Serve(ctx context.Context, addr string) error {
_ = b.st.SaveRecommenderState(cluster, r.Checkpoint())
}
b.mu.RUnlock()
b.saveEvidence()
}
return nil
case err := <-errCh:
Expand Down
Loading
Loading