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 }