Files
Josh Creek 320ec46ba2 fix(server): partition and bound the Redis candidate projection
Both playlists shared one hash and sorted set, causing two independent
failures.

Starvation: Snapshot performed an unbounded ZRANGEBYSCORE and HMGET,
decoded the whole queue, and the matcher then truncated to its candidate
limit *before* filtering by playlist. A large casual prefix could
therefore leave the ranked worker with zero candidates indefinitely even
while ranked tickets were queued further down the set.

Mutual erasure: each matcher captured only its own playlist as the
durable source, but Rebuild replaced the shared keys, so a casual repair
wiped ranked projections and vice versa.

Namespace the keys per playlist, push the limit into Redis (LIMIT 0 N)
so reads no longer scale with total queue depth, and scope Rebuild to
one namespace. Rebuild now rejects a candidate whose playlist does not
match the namespace, which would reintroduce the starvation. Upsert
derives the namespace from the candidate; Remove takes the playlist,
since a ticket ID alone no longer identifies its namespace.

Add tests for a 300-deep casual backlog not starving ranked, for neither
playlist's rebuild erasing the other, and for the limit being applied
without losing enqueue ordering.
2026-09-05 10:23:52 +01:00

268 lines
9.8 KiB
Go

package store
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/cosmic-clash/cosmic-clash/server/domain"
"github.com/redis/go-redis/v9"
)
// RedisCandidateIndex is a rebuildable acceleration index. It never decides
// ownership or claims a match; callers must source candidates from the
// durable queue projection before rebuilding it.
type RedisCandidateIndex struct {
Client *redis.Client
Prefix string
TTL time.Duration
}
// DurableCandidateSource is the authoritative queue projection used to
// repair Redis. Implementations must apply queue state and expiry rules before
// returning candidates. It is playlist- and limit-scoped so a repair reads
// only the namespace it is about to rebuild, and never an unbounded queue.
type DurableCandidateSource func(context.Context, domain.Playlist, time.Time, int) ([]domain.Candidate, error)
// CandidateProjection couples the transient index to its durable repair
// source. A cache miss, partial write, malformed payload, or Redis restart is
// repaired before candidates are returned to a matcher.
type CandidateProjection struct {
Index RedisCandidateIndex
Source DurableCandidateSource
}
func (p CandidateProjection) Repair(ctx context.Context, playlist domain.Playlist, now time.Time, limit int) error {
if p.Source == nil || now.IsZero() {
return fmt.Errorf("invalid candidate repair source")
}
candidates, err := p.Source(ctx, playlist, now, limit)
if err != nil {
return err
}
return p.Index.Rebuild(ctx, playlist, candidates)
}
// Snapshot never fails just because Redis specifically is unreachable.
// RedisCandidateIndex is documented everywhere (this type's own comment,
// cmd/matcher, cmd/control-plane's --redis-addr help text) as an optional,
// rebuildable acceleration layer over PostgreSQL authority -- but until this
// fix, a genuine Redis outage (not merely an empty or partial cache, an
// actual connection failure) made Snapshot fail outright: the old code
// treated "the index errored" and "the index came back empty" identically,
// funnelling both into Repair, which itself calls Index.Rebuild -- a second
// Redis round-trip that fails for exactly the same reason the first one did.
// A Redis failover or restart would have taken matchmaking down completely
// even though the authoritative Source (PostgreSQL) was perfectly healthy.
// Now: an index error or an empty read both fall back to serving Source
// directly, and only attempt to repopulate Redis on a best-effort basis --
// its outcome is deliberately ignored, since a caller must never be denied
// service just because the rebuild's own Redis write also failed.
func (p CandidateProjection) Snapshot(ctx context.Context, playlist domain.Playlist, now time.Time, limit int) ([]domain.Candidate, error) {
if p.Source == nil {
return nil, fmt.Errorf("invalid candidate repair source")
}
candidates, err := p.Index.Snapshot(ctx, playlist, now, limit)
if err == nil && len(candidates) > 0 {
return candidates, nil
}
// Either the index errored outright, or came back empty -- indistinguishable
// from a Redis restart or a lost keyspace. Consult PostgreSQL, the
// authoritative source, either way.
source, sourceErr := p.Source(ctx, playlist, now, limit)
if sourceErr != nil {
return nil, sourceErr
}
// Rebuild only this playlist's namespace. When the keys were shared, a
// casual repair replaced the keys ranked candidates lived in and vice
// versa, so each worker could erase the other's projection.
_ = p.Index.Rebuild(ctx, playlist, source)
return source, nil
}
// keys are namespaced per playlist. They used to be shared, which caused two
// independent failures: the matcher truncated a mixed snapshot to its
// candidate limit before filtering by playlist, so a large casual backlog
// could starve the ranked worker indefinitely; and Rebuild replaced the shared
// keys, so one playlist's repair erased the other's projection.
func (r RedisCandidateIndex) keys(playlist domain.Playlist) (string, string) {
prefix := r.Prefix
if prefix == "" {
prefix = "cosmic-clash"
}
base := prefix + ":queue:candidates:" + string(playlist)
return base + ":data", base + ":order"
}
func validRedisPlaylist(playlist domain.Playlist) error {
if playlist != domain.Casual && playlist != domain.Ranked {
return fmt.Errorf("invalid candidate playlist %q", playlist)
}
return nil
}
func (r RedisCandidateIndex) validate() error {
if r.Client == nil || r.TTL <= 0 {
return fmt.Errorf("invalid Redis candidate index")
}
return nil
}
func validateRedisCandidate(candidate domain.Candidate) error {
if candidate.TicketID == "" || candidate.PlayerID == "" || candidate.EnqueuedAt.IsZero() {
return fmt.Errorf("invalid candidate")
}
return nil
}
// Upsert stores the candidate payload and its deterministic enqueue ordering.
// Both keys receive a TTL so a Redis restart or abandoned index cannot become
// a permanent source of stale presence.
func (r RedisCandidateIndex) Upsert(ctx context.Context, candidate domain.Candidate) error {
if err := r.validate(); err != nil {
return err
}
if err := validateRedisCandidate(candidate); err != nil {
return err
}
if err := validRedisPlaylist(candidate.Playlist); err != nil {
return err
}
payload, err := json.Marshal(candidate)
if err != nil {
return err
}
dataKey, orderKey := r.keys(candidate.Playlist)
pipe := r.Client.TxPipeline()
pipe.HSet(ctx, dataKey, candidate.TicketID, payload)
pipe.ZAdd(ctx, orderKey, redis.Z{Score: float64(candidate.EnqueuedAt.UnixNano()), Member: candidate.TicketID})
pipe.Expire(ctx, dataKey, r.TTL)
pipe.Expire(ctx, orderKey, r.TTL)
_, err = pipe.Exec(ctx)
return err
}
func (r RedisCandidateIndex) Remove(ctx context.Context, playlist domain.Playlist, ticketID string) error {
if err := r.validate(); err != nil {
return err
}
if err := validRedisPlaylist(playlist); err != nil {
return err
}
if ticketID == "" {
return fmt.Errorf("ticket ID is required")
}
dataKey, orderKey := r.keys(playlist)
pipe := r.Client.TxPipeline()
pipe.HDel(ctx, dataKey, ticketID)
pipe.ZRem(ctx, orderKey, ticketID)
_, err := pipe.Exec(ctx)
return err
}
// Snapshot reads only candidates whose enqueue timestamp is not in the
// future. Missing payloads are ignored; the durable rebuild path repairs such
// partial cache state without allowing it to affect ownership.
func (r RedisCandidateIndex) Snapshot(ctx context.Context, playlist domain.Playlist, now time.Time, limit int) ([]domain.Candidate, error) {
if err := r.validate(); err != nil {
return nil, err
}
if err := validRedisPlaylist(playlist); err != nil {
return nil, err
}
if now.IsZero() {
return nil, fmt.Errorf("authoritative time is required")
}
if limit < 1 || limit > 1000 {
return nil, fmt.Errorf("invalid candidate snapshot limit")
}
dataKey, orderKey := r.keys(playlist)
// The limit is applied by Redis (LIMIT 0 N), not after transfer. The
// unbounded range and HMGET decoded the entire queue on every one-second
// poll, allocating and transferring in proportion to total queue depth.
tickets, err := r.Client.ZRangeByScore(ctx, orderKey, &redis.ZRangeBy{
Min: "-inf", Max: fmt.Sprint(now.UnixNano()), Offset: 0, Count: int64(limit),
}).Result()
if err != nil {
return nil, err
}
if len(tickets) == 0 {
return []domain.Candidate{}, nil
}
payloads, err := r.Client.HMGet(ctx, dataKey, tickets...).Result()
if err != nil {
return nil, err
}
result := make([]domain.Candidate, 0, len(payloads))
for i, raw := range payloads {
var encoded []byte
switch value := raw.(type) {
case string:
encoded = []byte(value)
case []byte:
encoded = value
default:
return nil, fmt.Errorf("candidate payload missing for %s", tickets[i])
}
var candidate domain.Candidate
if err := json.Unmarshal(encoded, &candidate); err != nil {
return nil, fmt.Errorf("invalid candidate payload for %s: %w", tickets[i], err)
}
if err := validateRedisCandidate(candidate); err != nil {
return nil, fmt.Errorf("invalid candidate payload for %s: %w", tickets[i], err)
}
if candidate.EnqueuedAt.After(now) {
return nil, fmt.Errorf("candidate payload is newer than its index for %s", tickets[i])
}
result = append(result, candidate)
}
return result, nil
}
// Rebuild atomically replaces both Redis keys from the authoritative queue
// projection. It is the required path after Redis restart/failover or cache
// loss, and rejects duplicate ticket IDs before touching Redis.
func (r RedisCandidateIndex) Rebuild(ctx context.Context, playlist domain.Playlist, candidates []domain.Candidate) error {
if err := r.validate(); err != nil {
return err
}
if err := validRedisPlaylist(playlist); err != nil {
return err
}
seen := make(map[string]struct{}, len(candidates))
values := make([]interface{}, 0, len(candidates)*2)
scores := make([]redis.Z, 0, len(candidates))
for _, candidate := range candidates {
if err := validateRedisCandidate(candidate); err != nil {
return err
}
if candidate.Playlist != playlist {
// A rebuild that mixed playlists would write foreign candidates
// into this namespace, reintroducing the starvation it fixes.
return fmt.Errorf("rebuild candidate %s is %q, not %q", candidate.TicketID, candidate.Playlist, playlist)
}
if _, exists := seen[candidate.TicketID]; exists {
return fmt.Errorf("duplicate candidate in rebuild")
}
seen[candidate.TicketID] = struct{}{}
payload, err := json.Marshal(candidate)
if err != nil {
return err
}
values = append(values, candidate.TicketID, payload)
scores = append(scores, redis.Z{Score: float64(candidate.EnqueuedAt.UnixNano()), Member: candidate.TicketID})
}
dataKey, orderKey := r.keys(playlist)
pipe := r.Client.TxPipeline()
pipe.Del(ctx, dataKey, orderKey)
if len(values) > 0 {
pipe.HSet(ctx, dataKey, values...)
pipe.ZAdd(ctx, orderKey, scores...)
}
pipe.Expire(ctx, dataKey, r.TTL)
pipe.Expire(ctx, orderKey, r.TTL)
_, err := pipe.Exec(ctx)
return err
}