mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-10 16:04:04 +00:00
320ec46ba2
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.
268 lines
9.8 KiB
Go
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
|
|
}
|