Skip to content

Commit 2d959b1

Browse files
api: §10.20 cached aggregation endpoints + cache helper (#22)
Adds the server-side aggregation layer the dashboard needs to retire its client-side resource-list scans for the Usage panel and sidebar counts. The §13 freshness matrix calls for eventual-consistent caching on these read paths (gating paths stay real-time per their own bucket). cache.GetOrSet[T](ctx, rdb, key, ttl, fn) collapses concurrent identical requests through golang.org/x/sync/singleflight and fails open when Redis is unreachable (logs + falls through to fn). Negative caching is allowed; corrupt cache entries are treated as misses and healed on the next SET. GET /api/v1/billing/usage — Redis-cached 30s, Cache-Control: private, max-age=30, stale-while-revalidate=60. Aggregates storage bytes per service (postgres/redis/mongodb), counts of deployments/webhooks/vault/ members. Carries `as_of` ISO timestamp + `freshness_seconds` so the dashboard can render an "as of Ns ago" footnote. GET /api/v1/team/summary — Redis-cached 5m, Cache-Control: private, max-age=300. Resource breakdown by type (one GROUP BY query), deployment count, member count, vault key count. Consumed by the sidebar SidebarUpgradeCard for resource/member counts that previously didn't render because the dashboard had no live source. Cache keys are team-scoped (`billing:usage:<team_id>`, `team:summary:<team_id>`) so two teams never share an entry. Tests: - internal/cache: singleflight collapses concurrent callers to one fn invocation; Redis-down falls through to fn; corrupt entry heals; zero-value (negative) caching works; Invalidate clears the key. - internal/handlers: two requests for the same team in <30s/<5min run exactly one set of DB queries (asserted via sqlmock strict mode); different teams trigger separate aggregations. Quota check paths (POST /db/new, /cache/new, etc.) are intentionally NOT touched — they read fresh per §13 because they gate provisioning. Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
1 parent 477e3c2 commit 2d959b1

7 files changed

Lines changed: 1243 additions & 0 deletions

File tree

‎internal/cache/redis.go‎

Lines changed: 148 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,148 @@
1+
// Package cache wraps the Redis client with a typed GetOrSet helper that
2+
// collapses concurrent identical requests via singleflight and fails open
3+
// when Redis is unavailable.
4+
//
5+
// Designed for the §13 eventual-consistency surfaces (billing/usage,
6+
// team/summary) where:
7+
//
8+
// - The per-team aggregation is expensive enough that N concurrent
9+
// dashboard tabs should NOT trigger N DB scans — singleflight collapses
10+
// them to one in-process compute + one cache write.
11+
// - A Redis outage MUST NOT break the read endpoint (the underlying DB is
12+
// still authoritative). GetOrSet falls through to fn on every Redis
13+
// error so the user sees data, just without the cache amortisation.
14+
// - Hot-path callers prefer a typed result (struct, not []byte). The
15+
// generic `T any` parameter keeps callers off encoding/json directly.
16+
//
17+
// Real-time paths (POST /db/new quota checks, webhook handlers) MUST NOT
18+
// use this helper — they read fresh per the §13 freshness matrix.
19+
package cache
20+
21+
import (
22+
"context"
23+
"encoding/json"
24+
"errors"
25+
"fmt"
26+
"log/slog"
27+
"time"
28+
29+
"github.com/redis/go-redis/v9"
30+
"golang.org/x/sync/singleflight"
31+
)
32+
33+
// group is the per-process singleflight that collapses concurrent calls to
34+
// GetOrSet sharing the same key. Keys live in one global namespace so callers
35+
// must scope them (e.g. "billing:usage:" + teamID).
36+
var group singleflight.Group
37+
38+
// GetOrSet returns the cached value for key when present and fresh.
39+
//
40+
// Miss path (cache empty or returns a NOT FOUND): runs fn under singleflight,
41+
// stores the encoded result with TTL ttl, returns the result.
42+
//
43+
// Failure modes (intentional fail-open semantics):
44+
//
45+
// - Redis GET errored — log + skip cache, run fn, return its result without
46+
// attempting another SET (the cache layer is currently broken; don't
47+
// hammer it). This matches the "Redis down → fall through" cell in the
48+
// §13 freshness matrix.
49+
// - JSON unmarshal of the cached value failed — treat as miss. Most likely
50+
// cause is a serialised value shape change across deploys; the next SET
51+
// after fn runs heals the cache entry.
52+
// - fn returned an error — propagate it without touching the cache.
53+
// - Redis SET errored on the way back — log + return the freshly-computed
54+
// value anyway. The next call will re-attempt the SET.
55+
//
56+
// Negative caching (fn returned a zero-value T) is allowed and uses the same
57+
// ttl — callers that want a shorter negative TTL should branch outside.
58+
func GetOrSet[T any](
59+
ctx context.Context,
60+
rdb *redis.Client,
61+
key string,
62+
ttl time.Duration,
63+
fn func(context.Context) (T, error),
64+
) (T, error) {
65+
var zero T
66+
67+
// Fast path: try the cache. A nil client means cache is disabled — go
68+
// straight to fn without using singleflight (no point — there's nothing
69+
// to collapse on).
70+
if rdb != nil {
71+
raw, err := rdb.Get(ctx, key).Bytes()
72+
switch {
73+
case err == nil:
74+
var out T
75+
if jerr := json.Unmarshal(raw, &out); jerr == nil {
76+
return out, nil
77+
}
78+
// Corrupt cache entry — treat as miss, log so the shape skew is
79+
// visible. Don't return the unmarshal error to the caller.
80+
slog.Warn("cache.get_unmarshal_failed", "key", key, "error", "json decode")
81+
case errors.Is(err, redis.Nil):
82+
// True miss — fall through to fn under singleflight.
83+
default:
84+
// Redis is unreachable / down. Fail open: run fn without the
85+
// cache wrapper and skip the SET path entirely so we don't
86+
// hammer a flapping Redis. Bypassing singleflight here means
87+
// N concurrent callers will all hit the DB during an outage,
88+
// which is acceptable — the cache being down IS the
89+
// degradation, the DB is the source of truth.
90+
slog.Warn("cache.get_failed_fail_open", "key", key, "error", err.Error())
91+
return fn(ctx)
92+
}
93+
}
94+
95+
// Miss path: collapse concurrent callers to one fn invocation.
96+
//
97+
// singleflight returns (value, error, shared). We ignore `shared`; both
98+
// the leader and the followers see the same value+error pair. The leader
99+
// is the only one that touches Redis SET — followers piggyback on the
100+
// returned value.
101+
v, err, _ := group.Do(key, func() (interface{}, error) {
102+
out, fnErr := fn(ctx)
103+
if fnErr != nil {
104+
return out, fnErr
105+
}
106+
if rdb != nil {
107+
encoded, jerr := json.Marshal(out)
108+
if jerr != nil {
109+
// Encoding failure is a programmer error (T can't be
110+
// marshalled). Don't poison the cache; log + return the
111+
// value so the request still succeeds.
112+
slog.Warn("cache.set_marshal_failed", "key", key, "error", jerr.Error())
113+
return out, nil
114+
}
115+
if setErr := rdb.Set(ctx, key, encoded, ttl).Err(); setErr != nil {
116+
// Same fail-open as GET: log but return the value.
117+
slog.Warn("cache.set_failed", "key", key, "error", setErr.Error())
118+
}
119+
}
120+
return out, nil
121+
})
122+
123+
if err != nil {
124+
return zero, err
125+
}
126+
// singleflight returns the leader's value via interface{}. The type
127+
// parameter T is the same for every caller of this key, so the assertion
128+
// is safe under normal use; a panic here would indicate two callers
129+
// using the same cache key with different T (a bug in caller code).
130+
out, ok := v.(T)
131+
if !ok {
132+
return zero, fmt.Errorf("cache.GetOrSet: type mismatch for key %q", key)
133+
}
134+
return out, nil
135+
}
136+
137+
// Invalidate deletes a cache key. Use it from write paths that change the
138+
// underlying aggregate (e.g. a deploy completing should invalidate
139+
// billing:usage:<team>). A nil client is a no-op so callers can wire this
140+
// in without conditional checks.
141+
func Invalidate(ctx context.Context, rdb *redis.Client, key string) {
142+
if rdb == nil {
143+
return
144+
}
145+
if err := rdb.Del(ctx, key).Err(); err != nil {
146+
slog.Warn("cache.invalidate_failed", "key", key, "error", err.Error())
147+
}
148+
}

‎internal/cache/redis_test.go‎

Lines changed: 244 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,244 @@
1+
package cache_test
2+
3+
import (
4+
"context"
5+
"errors"
6+
"sync"
7+
"sync/atomic"
8+
"testing"
9+
"time"
10+
11+
"github.com/alicebob/miniredis/v2"
12+
"github.com/redis/go-redis/v9"
13+
"github.com/stretchr/testify/assert"
14+
"github.com/stretchr/testify/require"
15+
16+
"instant.dev/internal/cache"
17+
)
18+
19+
// newMiniRedis returns a *redis.Client backed by an in-memory miniredis
20+
// instance plus a cleanup func. Used everywhere we need a real-shaped
21+
// Redis without a Docker container.
22+
func newMiniRedis(t *testing.T) (*redis.Client, func()) {
23+
t.Helper()
24+
mr, err := miniredis.Run()
25+
require.NoError(t, err)
26+
rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()})
27+
return rdb, func() {
28+
rdb.Close()
29+
mr.Close()
30+
}
31+
}
32+
33+
type usagePayload struct {
34+
Postgres int64 `json:"postgres"`
35+
Redis int64 `json:"redis"`
36+
}
37+
38+
// TestGetOrSet_MissRunsFnOnceAndCaches verifies the basic Redis-miss path:
39+
// the first call runs fn, the second call short-circuits to the cache.
40+
func TestGetOrSet_MissRunsFnOnceAndCaches(t *testing.T) {
41+
rdb, cleanup := newMiniRedis(t)
42+
defer cleanup()
43+
44+
var calls atomic.Int32
45+
fn := func(_ context.Context) (usagePayload, error) {
46+
calls.Add(1)
47+
return usagePayload{Postgres: 100, Redis: 50}, nil
48+
}
49+
50+
ctx := context.Background()
51+
v1, err := cache.GetOrSet(ctx, rdb, "test:k1", 60*time.Second, fn)
52+
require.NoError(t, err)
53+
assert.Equal(t, usagePayload{Postgres: 100, Redis: 50}, v1)
54+
55+
v2, err := cache.GetOrSet(ctx, rdb, "test:k1", 60*time.Second, fn)
56+
require.NoError(t, err)
57+
assert.Equal(t, usagePayload{Postgres: 100, Redis: 50}, v2)
58+
59+
assert.Equal(t, int32(1), calls.Load(), "fn should have run exactly once across both calls")
60+
}
61+
62+
// TestGetOrSet_SingleflightCollapsesConcurrentCallers — the headline §10.20
63+
// guarantee: N concurrent identical requests collapse to 1 fn invocation.
64+
// Without singleflight, N callers would race past the empty-cache check and
65+
// all run fn before any of them got to SET. With singleflight, the leader
66+
// runs fn and the followers receive its result.
67+
func TestGetOrSet_SingleflightCollapsesConcurrentCallers(t *testing.T) {
68+
rdb, cleanup := newMiniRedis(t)
69+
defer cleanup()
70+
71+
const concurrency = 20
72+
var calls atomic.Int32
73+
// gate holds fn open until every goroutine is in flight, so they all
74+
// observe the same "cache empty" snapshot. Without it the test races —
75+
// goroutine #N might run after goroutine #1 already set the cache.
76+
gate := make(chan struct{})
77+
fn := func(_ context.Context) (usagePayload, error) {
78+
<-gate
79+
calls.Add(1)
80+
// A small sleep makes the singleflight window visible — the leader
81+
// is still inside fn when followers arrive. Without it the timing
82+
// can occasionally let a follower miss the inflight entry.
83+
time.Sleep(20 * time.Millisecond)
84+
return usagePayload{Postgres: 42}, nil
85+
}
86+
87+
ctx := context.Background()
88+
results := make(chan usagePayload, concurrency)
89+
errs := make(chan error, concurrency)
90+
91+
var wg sync.WaitGroup
92+
for i := 0; i < concurrency; i++ {
93+
wg.Add(1)
94+
go func() {
95+
defer wg.Done()
96+
v, err := cache.GetOrSet(ctx, rdb, "test:sf", 60*time.Second, fn)
97+
results <- v
98+
errs <- err
99+
}()
100+
}
101+
// Let every goroutine reach the gate before any of them runs fn.
102+
time.Sleep(50 * time.Millisecond)
103+
close(gate)
104+
wg.Wait()
105+
close(results)
106+
close(errs)
107+
108+
for err := range errs {
109+
require.NoError(t, err)
110+
}
111+
for v := range results {
112+
assert.Equal(t, usagePayload{Postgres: 42}, v)
113+
}
114+
assert.Equal(t, int32(1), calls.Load(), "singleflight should collapse %d concurrent callers to 1 fn invocation", concurrency)
115+
}
116+
117+
// TestGetOrSet_RedisDownFailsOpen verifies that when Redis errors on GET,
118+
// GetOrSet falls through to fn and returns its result. The cache being
119+
// unreachable must never break the read path.
120+
func TestGetOrSet_RedisDownFailsOpen(t *testing.T) {
121+
// Point at a closed port — the dial will fail fast.
122+
rdb := redis.NewClient(&redis.Options{
123+
Addr: "127.0.0.1:1", // reserved low port, refuses connections
124+
DialTimeout: 50 * time.Millisecond,
125+
})
126+
defer rdb.Close()
127+
128+
var calls atomic.Int32
129+
fn := func(_ context.Context) (usagePayload, error) {
130+
calls.Add(1)
131+
return usagePayload{Postgres: 7}, nil
132+
}
133+
134+
ctx := context.Background()
135+
v, err := cache.GetOrSet(ctx, rdb, "test:down", 60*time.Second, fn)
136+
require.NoError(t, err)
137+
assert.Equal(t, usagePayload{Postgres: 7}, v)
138+
assert.Equal(t, int32(1), calls.Load(), "fn must run when redis is down")
139+
140+
// A second call must also reach fn — we bypass singleflight on the
141+
// Redis-down path to avoid hammering a flapping cache, and the cache
142+
// itself can't serve the entry. (See §10.20 fail-open contract.)
143+
v2, err := cache.GetOrSet(ctx, rdb, "test:down", 60*time.Second, fn)
144+
require.NoError(t, err)
145+
assert.Equal(t, usagePayload{Postgres: 7}, v2)
146+
assert.Equal(t, int32(2), calls.Load())
147+
}
148+
149+
// TestGetOrSet_NilClientPassesThrough — a nil *redis.Client means "no cache
150+
// configured"; GetOrSet should still call fn and return its result. Useful
151+
// in tests and in dev configs where Redis isn't wired.
152+
func TestGetOrSet_NilClientPassesThrough(t *testing.T) {
153+
var calls atomic.Int32
154+
fn := func(_ context.Context) (usagePayload, error) {
155+
calls.Add(1)
156+
return usagePayload{Postgres: 1}, nil
157+
}
158+
v, err := cache.GetOrSet(context.Background(), nil, "test:nil", 60*time.Second, fn)
159+
require.NoError(t, err)
160+
assert.Equal(t, usagePayload{Postgres: 1}, v)
161+
assert.Equal(t, int32(1), calls.Load())
162+
}
163+
164+
// TestGetOrSet_FnErrorPropagates — a fn error must not be cached and must
165+
// surface to the caller verbatim.
166+
func TestGetOrSet_FnErrorPropagates(t *testing.T) {
167+
rdb, cleanup := newMiniRedis(t)
168+
defer cleanup()
169+
170+
sentinel := errors.New("aggregate failed")
171+
fn := func(_ context.Context) (usagePayload, error) {
172+
return usagePayload{}, sentinel
173+
}
174+
175+
_, err := cache.GetOrSet(context.Background(), rdb, "test:err", 60*time.Second, fn)
176+
require.Error(t, err)
177+
assert.ErrorIs(t, err, sentinel)
178+
179+
// Confirm the cache was NOT populated.
180+
_, ferr := rdb.Get(context.Background(), "test:err").Bytes()
181+
assert.ErrorIs(t, ferr, redis.Nil)
182+
}
183+
184+
// TestGetOrSet_ZeroValueCachesNegative — fn returning a zero-value T is a
185+
// valid result (e.g. a team with no resources). It must still be cached so
186+
// the next caller doesn't re-run the aggregate.
187+
func TestGetOrSet_ZeroValueCachesNegative(t *testing.T) {
188+
rdb, cleanup := newMiniRedis(t)
189+
defer cleanup()
190+
191+
var calls atomic.Int32
192+
fn := func(_ context.Context) (usagePayload, error) {
193+
calls.Add(1)
194+
return usagePayload{}, nil
195+
}
196+
ctx := context.Background()
197+
_, err := cache.GetOrSet(ctx, rdb, "test:empty", 60*time.Second, fn)
198+
require.NoError(t, err)
199+
_, err = cache.GetOrSet(ctx, rdb, "test:empty", 60*time.Second, fn)
200+
require.NoError(t, err)
201+
assert.Equal(t, int32(1), calls.Load(), "zero-value results must still be cached")
202+
}
203+
204+
// TestGetOrSet_CorruptCacheEntryFallsThrough — if a cache entry was
205+
// serialised under an older shape, json.Unmarshal returns an error and
206+
// GetOrSet treats it as a miss. The next SET heals the entry.
207+
func TestGetOrSet_CorruptCacheEntryFallsThrough(t *testing.T) {
208+
rdb, cleanup := newMiniRedis(t)
209+
defer cleanup()
210+
211+
// Plant a value that doesn't decode as usagePayload.
212+
require.NoError(t, rdb.Set(context.Background(), "test:corrupt", "not-json", time.Minute).Err())
213+
214+
var calls atomic.Int32
215+
fn := func(_ context.Context) (usagePayload, error) {
216+
calls.Add(1)
217+
return usagePayload{Postgres: 999}, nil
218+
}
219+
v, err := cache.GetOrSet(context.Background(), rdb, "test:corrupt", time.Minute, fn)
220+
require.NoError(t, err)
221+
assert.Equal(t, usagePayload{Postgres: 999}, v)
222+
assert.Equal(t, int32(1), calls.Load())
223+
}
224+
225+
// TestInvalidate_DeletesKey ensures Invalidate clears the cache and a nil
226+
// client is a no-op.
227+
func TestInvalidate_DeletesKey(t *testing.T) {
228+
rdb, cleanup := newMiniRedis(t)
229+
defer cleanup()
230+
231+
fn := func(_ context.Context) (usagePayload, error) {
232+
return usagePayload{Postgres: 5}, nil
233+
}
234+
ctx := context.Background()
235+
_, err := cache.GetOrSet(ctx, rdb, "test:inv", time.Minute, fn)
236+
require.NoError(t, err)
237+
238+
cache.Invalidate(ctx, rdb, "test:inv")
239+
_, err = rdb.Get(ctx, "test:inv").Bytes()
240+
assert.ErrorIs(t, err, redis.Nil)
241+
242+
// nil client → no panic.
243+
cache.Invalidate(ctx, nil, "test:inv")
244+
}

0 commit comments

Comments
 (0)