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
7 changes: 7 additions & 0 deletions framework/configstore/tables/webhooks.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package tables

import (
"database/sql/driver"
"encoding/json"
"fmt"
"net/url"
Expand Down Expand Up @@ -37,6 +38,12 @@ func (e WebhookEvent) IsValid() bool {
return slices.Contains(WebhookEvents, e)
}

// Value implements driver.Valuer so database drivers that append typed column
// values (e.g. clickhouse-go batch inserts) can serialize the type.
func (e WebhookEvent) Value() (driver.Value, error) {
return string(e), nil
}

// TableWebhookEndpoint represents a registered webhook endpoint in the database.
type TableWebhookEndpoint struct {
ID string `gorm:"type:varchar(36);primaryKey" json:"id"`
Expand Down
12 changes: 7 additions & 5 deletions framework/logstore/logstoreparity_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -252,11 +252,13 @@ func canonicalizeOrder(key string, v any) any {
tv[i] = canonicalizeOrder("", tv[i])
}
if key == "rankings" {
keys := make([]string, len(tv))
for i, e := range tv {
keys[i] = canonicalEntryKey(e)
}
sort.SliceStable(tv, func(i, j int) bool { return keys[i] < keys[j] })
// Compute the key from the live element: a precomputed key slice
// would not be permuted alongside tv by SliceStable's swapper, so
// after the first swap the comparator would read stale keys and
// the final order would depend on the raw SQL row order.
sort.SliceStable(tv, func(i, j int) bool {
return canonicalEntryKey(tv[i]) < canonicalEntryKey(tv[j])
})
}
return tv
default:
Expand Down
103 changes: 97 additions & 6 deletions framework/logstore/matview_count_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -109,12 +109,12 @@ func TestSearchLogsMatViewCountCanonicalModelFilter(t *testing.T) {
end := day.Add(35*time.Hour + 45*time.Minute) // 25h15m window, matview-eligible

canonical := "gpt-4o-mini"
interior := day.Add(18 * time.Hour) // deep inside the window, never a boundary sliver
insertCountTestModelLog(t, db, interior, "success", canonical, nil) // wire-model match
insertCountTestModelLog(t, db, interior.Add(time.Minute), "success", "prod-alias", &canonical) // canonical-only match
insertCountTestModelLog(t, db, interior.Add(2*time.Minute), "success", "other-model", nil) // no match
insertCountTestModelLog(t, db, start.Add(5*time.Minute), "success", "prod-alias", &canonical) // canonical match in the head boundary sliver
insertCountTestModelLog(t, db, day.Add(9*time.Hour), "success", "prod-alias", &canonical) // before start -> excluded
interior := day.Add(18 * time.Hour) // deep inside the window, never a boundary sliver
insertCountTestModelLog(t, db, interior, "success", canonical, nil) // wire-model match
insertCountTestModelLog(t, db, interior.Add(time.Minute), "success", "prod-alias", &canonical) // canonical-only match
insertCountTestModelLog(t, db, interior.Add(2*time.Minute), "success", "other-model", nil) // no match
insertCountTestModelLog(t, db, start.Add(5*time.Minute), "success", "prod-alias", &canonical) // canonical match in the head boundary sliver
insertCountTestModelLog(t, db, day.Add(9*time.Hour), "success", "prod-alias", &canonical) // before start -> excluded

refreshTestMatViews(t, db)
store.matViewsReady.Store(true)
Expand Down Expand Up @@ -253,3 +253,94 @@ func TestGetHistogramMatViewTrimsBoundaryBuckets(t *testing.T) {
assert.Equal(t, int64(1), byBucket[day.Add(35*time.Hour).Unix()],
"last bar holds only the in-window sliver row, not the post-end row from the same hour")
}

// insertCacheTestLog inserts a terminal log with an explicit cache_debug
// payload for the hybrid cache-hit tests. Pass nil for a row with no
// cache_debug at all.
func insertCacheTestLog(t *testing.T, db *gorm.DB, ts time.Time, cacheDebug *string) {
t.Helper()
err := db.Exec(`
INSERT INTO logs (id, timestamp, object_type, provider, model, status, cache_debug,
created_at, latency, cost, prompt_tokens, completion_tokens, total_tokens)
VALUES (?, ?, 'chat_completion', 'openai', 'gpt-4', 'success', ?, ?, 100, 0.01, 10, 5, 15)
`, uuid.New().String(), ts, cacheDebug, ts).Error
require.NoError(t, err, "failed to insert cache test log")
}

// TestGetStatsMatViewCacheHitsHybrid verifies cache-hit stats come from the
// interior+boundary hybrid (materialized in mv_logs_hourly, classified raw in
// the slivers) instead of a full-window raw scan, and that the matview path
// agrees exactly with the raw path, including the nil contract: fields stay
// nil when no row in the window carried valid cache_debug JSON, and are
// explicit zeros when cache rows exist but none were direct/semantic.
func TestGetStatsMatViewCacheHitsHybrid(t *testing.T) {
store, db := setupPerfTestDB(t)
ctx := context.Background()

day := time.Date(2026, 7, 20, 0, 0, 0, 0, time.UTC)
direct := `{"hit_type":"direct"}`
semantic := `{"hit_type":"semantic"}`
other := `{"hit_type":"external"}`
malformed := `not json`

// Main window: 10:30 -> next day 11:45.
start := day.Add(10*time.Hour + 30*time.Minute)
end := day.Add(35*time.Hour + 45*time.Minute)
insertCacheTestLog(t, db, day.Add(10*time.Hour+15*time.Minute), &direct) // before start -> excluded
insertCacheTestLog(t, db, day.Add(10*time.Hour+45*time.Minute), &direct) // head sliver -> raw classifier
insertCacheTestLog(t, db, day.Add(18*time.Hour), &semantic) // interior -> matview column
insertCacheTestLog(t, db, day.Add(18*time.Hour+5*time.Minute), &other) // counts as cache row, neither type
insertCacheTestLog(t, db, day.Add(19*time.Hour), &malformed) // guard rejects -> not a cache row
insertCacheTestLog(t, db, day.Add(35*time.Hour+30*time.Minute), &direct) // tail sliver -> raw classifier
insertCacheTestLog(t, db, day.Add(35*time.Hour+50*time.Minute), &semantic) // after end -> excluded

// Nil-contract window: 40h -> 70h holds one row with no cache_debug.
nilStart, nilEnd := day.Add(40*time.Hour), day.Add(70*time.Hour)
insertCacheTestLog(t, db, day.Add(50*time.Hour), nil)

// Zero-contract window: 80h -> 110h holds only an other-hit_type cache row.
zeroStart, zeroEnd := day.Add(80*time.Hour), day.Add(110*time.Hour)
insertCacheTestLog(t, db, day.Add(90*time.Hour), &other)

refreshTestMatViews(t, db)
store.matViewsReady.Store(true)

filters := SearchFilters{StartTime: &start, EndTime: &end}
require.True(t, store.canUseMatViewForFreshAggregate(filters))
stats, err := store.GetStats(ctx, filters)
require.NoError(t, err)
require.NotNil(t, stats.DirectCacheHits)
require.NotNil(t, stats.SemanticCacheHits)
assert.Equal(t, int64(2), *stats.DirectCacheHits, "head + tail sliver direct hits")
assert.Equal(t, int64(1), *stats.SemanticCacheHits, "interior semantic hit")

// The raw path over the identical filters must produce identical values.
store.matViewsReady.Store(false)
rawStats, err := store.GetStats(ctx, filters)
require.NoError(t, err)
require.NotNil(t, rawStats.DirectCacheHits)
require.NotNil(t, rawStats.SemanticCacheHits)
assert.Equal(t, *rawStats.DirectCacheHits, *stats.DirectCacheHits)
assert.Equal(t, *rawStats.SemanticCacheHits, *stats.SemanticCacheHits)
store.matViewsReady.Store(true)

// Nil contract: no valid cache_debug rows in the window -> fields omitted
// on both paths.
filters = SearchFilters{StartTime: &nilStart, EndTime: &nilEnd}
require.True(t, store.canUseMatViewForFreshAggregate(filters))
stats, err = store.GetStats(ctx, filters)
require.NoError(t, err)
assert.Nil(t, stats.DirectCacheHits, "no cache rows -> nil, matching the raw path contract")
assert.Nil(t, stats.SemanticCacheHits)

// Zero contract: cache rows exist but none direct/semantic -> explicit
// zeros on both paths.
filters = SearchFilters{StartTime: &zeroStart, EndTime: &zeroEnd}
require.True(t, store.canUseMatViewForFreshAggregate(filters))
stats, err = store.GetStats(ctx, filters)
require.NoError(t, err)
require.NotNil(t, stats.DirectCacheHits, "cache rows present -> explicit zeros, not omission")
require.NotNil(t, stats.SemanticCacheHits)
assert.Equal(t, int64(0), *stats.DirectCacheHits)
assert.Equal(t, int64(0), *stats.SemanticCacheHits)
}
119 changes: 119 additions & 0 deletions framework/logstore/matviewheal.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
package logstore

import (
"context"
"errors"
"fmt"
"time"

"github.com/jackc/pgx/v5/pgconn"
)

// This file implements runtime resilience for the materialized-view read
// path. Startup already repairs shape drift (repairMatViewShapes diffs the
// live catalog against matviewRequiredColumns and drops stale views), but a
// query can still hit a missing or stale-shaped view at runtime: a replica
// that lost the ensureMatViews advisory-lock race during a rolling deploy
// reads while the lock holder is mid-rebuild, or an operator drops a view.
// Rather than gating reads on catalog checks (the approach of 8354ca76,
// reverted in 9c3e4c7f3), the failing query itself signals staleness: the
// dispatch site serves that request from the raw logs table, the matview
// read path is disabled process-wide, and a single-flight background repair
// recreates and refreshes the views before re-enabling it. A premature
// ready=true from any source is therefore harmless - the next shape error
// restarts the cycle, and the system converges once shapes are current.

// isMatViewShapeError reports whether err indicates a materialized view that
// is missing or has a stale shape, so the caller should fall back to the raw
// logs table. Deliberately excluded: 42883 (undefined_function - our readers
// use only built-ins, so that is a code bug we want loud) and 0A000
// (cached-plan drift is structurally prevented by the two-pool connection
// lifecycle; masking it would hide a regression there).
func isMatViewShapeError(err error) bool {
var pgErr *pgconn.PgError
if !errors.As(err, &pgErr) {
return false
}
switch pgErr.Code {
case "42P01", // undefined_table: view dropped (mid-repair window, operator action)
"42703", // undefined_column: old-shape view still present (the #5384 failure)
"55000": // object_not_in_prerequisite_state: matview exists but is not populated
return true
}
return false
}

// fallBackToRaw classifies err; on a matview shape error it disables the
// matview read path for subsequent requests, kicks the background self-heal,
// and returns true so the dispatch site falls through to its raw-table
// implementation for the current request. Any other error (including nil)
// returns false and the caller returns normally.
func (s *RDBLogStore) fallBackToRaw(err error) bool {
if err == nil || !isMatViewShapeError(err) {
return false
}
s.matViewsReady.Store(false)
if s.logger != nil {
s.logger.Warn(fmt.Sprintf("logstore: matview query failed with shape error, serving from raw tables until repaired: %s", err))
}
s.triggerMatViewSelfHeal()
return true
}

// matViewHealCooldown bounds how often a process attempts a background
// repair. No dedicated retry loop exists here because recovery after a failed
// heal is owned by the periodic refresher: startMatViewRefresher re-arms
// matViewsReady on its next successful tick, and it is guaranteed to be
// running whenever self-heal can trigger (a shape error requires the matview
// path to have been enabled, which requires the boot-time ensureMatViews that
// also starts the refresher). Future shape errors additionally re-trigger the
// heal, cooldown-limited - while broken, every request keeps succeeding via
// the raw fallback.
const matViewHealCooldown = 30 * time.Second

// triggerMatViewSelfHeal starts a single-flight background repair: recreate
// any missing or stale matviews (ensureMatViews serializes cross-replica on
// the refresh advisory lock and handles shape diffing, drops, creates, and
// index builds), refresh them, and re-enable the matview read path. Repair
// failures are logged only; the raw fallback keeps serving in the meantime.
//
// A replica that loses the advisory-lock race re-enables the read path while
// the lock holder may still be mid-rebuild; that premature enable is accepted
// - the next query against a still-stale view falls back raw and re-triggers
// the heal, converging once the lock holder finishes.
func (s *RDBLogStore) triggerMatViewSelfHeal() {
if !s.matViewHealInFlight.CompareAndSwap(false, true) {
return // a heal is already running
}
go func() {
defer s.matViewHealInFlight.Store(false)
if time.Since(time.Unix(0, s.matViewHealLastAttempt.Load())) < matViewHealCooldown {
return
}
s.matViewHealLastAttempt.Store(time.Now().UnixNano())

ctx := context.Background()
if err := ensureMatViews(ctx, s.db); err != nil {
if s.logger != nil {
s.logger.Warn(fmt.Sprintf("logstore: matview self-heal creation failed: %s (still serving from raw tables)", err))
}
return
}
if err := refreshMatViews(ctx, s.db); err != nil {
if s.logger != nil {
s.logger.Warn(fmt.Sprintf("logstore: matview self-heal refresh failed: %s (still serving from raw tables)", err))
}
return
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
s.matViewsReady.Store(true)
if s.logger != nil {
s.logger.Info("logstore: materialized views self-healed after shape error")
}
}()
}

// resetMatViewHeal clears the single-flight and cooldown state. Test helper.
func (s *RDBLogStore) resetMatViewHeal() {
s.matViewHealInFlight.Store(false)
s.matViewHealLastAttempt.Store(0)
}
Loading
Loading