mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-11 00:14:00 +00:00
feat: add ranked season rollover policy
This commit is contained in:
+68
-30
@@ -11,24 +11,24 @@ import (
|
||||
|
||||
const (
|
||||
QueueHeartbeatInterval = 10 * time.Second
|
||||
QueueExpiryWindow = 30 * time.Second
|
||||
QueueExpiryWindow = 30 * time.Second
|
||||
)
|
||||
|
||||
var (
|
||||
ErrPlayerQueued = errors.New("player already owns an active queue ticket")
|
||||
ErrPlayerQueued = errors.New("player already owns an active queue ticket")
|
||||
ErrTicketNotFound = errors.New("queue ticket not found")
|
||||
ErrNotTicketOwner = errors.New("queue ticket is owned by another player")
|
||||
ErrTicketExpired = errors.New("queue ticket expired")
|
||||
ErrTicketExpired = errors.New("queue ticket expired")
|
||||
)
|
||||
|
||||
type QueueTicket struct {
|
||||
TicketID string
|
||||
PlayerID string
|
||||
Candidate Candidate
|
||||
State State
|
||||
Revision uint64
|
||||
TicketID string
|
||||
PlayerID string
|
||||
Candidate Candidate
|
||||
State State
|
||||
Revision uint64
|
||||
EnqueuedAt time.Time
|
||||
ExpiresAt time.Time
|
||||
ExpiresAt time.Time
|
||||
}
|
||||
|
||||
type queueMutation struct {
|
||||
@@ -37,8 +37,8 @@ type queueMutation struct {
|
||||
}
|
||||
|
||||
type Queue struct {
|
||||
tickets map[string]QueueTicket
|
||||
byPlayer map[string]string
|
||||
tickets map[string]QueueTicket
|
||||
byPlayer map[string]string
|
||||
mutations map[string]queueMutation
|
||||
}
|
||||
|
||||
@@ -52,14 +52,20 @@ func NewQueue() *Queue {
|
||||
func (q *Queue) Create(playerID, ticketID, idempotencyKey string, candidate Candidate, now time.Time) (QueueTicket, error) {
|
||||
digest := sha256.Sum256([]byte(createPayload(playerID, ticketID, candidate)))
|
||||
if prior, ok := q.mutations[idempotencyKey]; ok {
|
||||
if prior.digest != digest { return QueueTicket{}, fmt.Errorf("%w: create payload changed", ErrConflict) }
|
||||
if prior.digest != digest {
|
||||
return QueueTicket{}, fmt.Errorf("%w: create payload changed", ErrConflict)
|
||||
}
|
||||
return prior.ticket, nil
|
||||
}
|
||||
if idempotencyKey == "" || playerID == "" || ticketID == "" || candidate.TicketID != ticketID {
|
||||
return QueueTicket{}, fmt.Errorf("%w: invalid queue create", ErrConflict)
|
||||
}
|
||||
if _, ok := q.byPlayer[playerID]; ok { return QueueTicket{}, ErrPlayerQueued }
|
||||
if _, ok := q.tickets[ticketID]; ok { return QueueTicket{}, fmt.Errorf("%w: ticket ID already exists", ErrConflict) }
|
||||
if _, ok := q.byPlayer[playerID]; ok {
|
||||
return QueueTicket{}, ErrPlayerQueued
|
||||
}
|
||||
if _, ok := q.tickets[ticketID]; ok {
|
||||
return QueueTicket{}, fmt.Errorf("%w: ticket ID already exists", ErrConflict)
|
||||
}
|
||||
ticket := QueueTicket{TicketID: ticketID, PlayerID: playerID, Candidate: candidate, State: Queued, EnqueuedAt: now, ExpiresAt: now.Add(QueueExpiryWindow)}
|
||||
q.tickets[ticketID] = ticket
|
||||
q.byPlayer[playerID] = ticketID
|
||||
@@ -70,15 +76,27 @@ func (q *Queue) Create(playerID, ticketID, idempotencyKey string, candidate Cand
|
||||
func (q *Queue) Heartbeat(playerID, ticketID, idempotencyKey string, expectedRevision uint64, now time.Time) (QueueTicket, error) {
|
||||
digest := sha256.Sum256([]byte(fmt.Sprintf("heartbeat:%s:%d", ticketID, expectedRevision)))
|
||||
if prior, ok := q.mutations[idempotencyKey]; ok {
|
||||
if prior.digest != digest { return QueueTicket{}, fmt.Errorf("%w: heartbeat payload changed", ErrConflict) }
|
||||
if prior.digest != digest {
|
||||
return QueueTicket{}, fmt.Errorf("%w: heartbeat payload changed", ErrConflict)
|
||||
}
|
||||
return prior.ticket, nil
|
||||
}
|
||||
ticket, err := q.ownedTicket(playerID, ticketID)
|
||||
if err != nil { return QueueTicket{}, err }
|
||||
if now.After(ticket.ExpiresAt) || now.Equal(ticket.ExpiresAt) { return QueueTicket{}, ErrTicketExpired }
|
||||
if ticket.Revision != expectedRevision { return QueueTicket{}, ErrStaleRevision }
|
||||
if ticket.State != Queued && ticket.State != Proposed { return QueueTicket{}, fmt.Errorf("%w: heartbeat in %s", ErrConflict, ticket.State) }
|
||||
if idempotencyKey == "" { return QueueTicket{}, fmt.Errorf("%w: empty heartbeat key", ErrConflict) }
|
||||
if err != nil {
|
||||
return QueueTicket{}, err
|
||||
}
|
||||
if now.After(ticket.ExpiresAt) || now.Equal(ticket.ExpiresAt) {
|
||||
return QueueTicket{}, ErrTicketExpired
|
||||
}
|
||||
if ticket.Revision != expectedRevision {
|
||||
return QueueTicket{}, ErrStaleRevision
|
||||
}
|
||||
if ticket.State != Queued && ticket.State != Proposed {
|
||||
return QueueTicket{}, fmt.Errorf("%w: heartbeat in %s", ErrConflict, ticket.State)
|
||||
}
|
||||
if idempotencyKey == "" {
|
||||
return QueueTicket{}, fmt.Errorf("%w: empty heartbeat key", ErrConflict)
|
||||
}
|
||||
ticket.Revision++
|
||||
ticket.ExpiresAt = now.Add(QueueExpiryWindow)
|
||||
q.tickets[ticketID] = ticket
|
||||
@@ -89,13 +107,21 @@ func (q *Queue) Heartbeat(playerID, ticketID, idempotencyKey string, expectedRev
|
||||
func (q *Queue) Cancel(playerID, ticketID, idempotencyKey string, expectedRevision uint64, now time.Time) (QueueTicket, error) {
|
||||
digest := sha256.Sum256([]byte(fmt.Sprintf("cancel:%s:%d", ticketID, expectedRevision)))
|
||||
if prior, ok := q.mutations[idempotencyKey]; ok {
|
||||
if prior.digest != digest { return QueueTicket{}, fmt.Errorf("%w: cancel payload changed", ErrConflict) }
|
||||
if prior.digest != digest {
|
||||
return QueueTicket{}, fmt.Errorf("%w: cancel payload changed", ErrConflict)
|
||||
}
|
||||
return prior.ticket, nil
|
||||
}
|
||||
ticket, err := q.ownedTicket(playerID, ticketID)
|
||||
if err != nil { return QueueTicket{}, err }
|
||||
if ticket.Revision != expectedRevision { return QueueTicket{}, ErrStaleRevision }
|
||||
if idempotencyKey == "" { return QueueTicket{}, fmt.Errorf("%w: empty cancel key", ErrConflict) }
|
||||
if err != nil {
|
||||
return QueueTicket{}, err
|
||||
}
|
||||
if ticket.Revision != expectedRevision {
|
||||
return QueueTicket{}, ErrStaleRevision
|
||||
}
|
||||
if idempotencyKey == "" {
|
||||
return QueueTicket{}, fmt.Errorf("%w: empty cancel key", ErrConflict)
|
||||
}
|
||||
ticket.State = Cancelled
|
||||
ticket.Revision++
|
||||
ticket.ExpiresAt = now
|
||||
@@ -124,10 +150,14 @@ func (q *Queue) Candidates(now time.Time) []Candidate {
|
||||
q.Expire(now)
|
||||
result := make([]Candidate, 0)
|
||||
for _, ticket := range q.tickets {
|
||||
if ticket.State == Queued { result = append(result, ticket.Candidate) }
|
||||
if ticket.State == Queued {
|
||||
result = append(result, ticket.Candidate)
|
||||
}
|
||||
}
|
||||
sort.Slice(result, func(i, j int) bool {
|
||||
if !result[i].EnqueuedAt.Equal(result[j].EnqueuedAt) { return result[i].EnqueuedAt.Before(result[j].EnqueuedAt) }
|
||||
if !result[i].EnqueuedAt.Equal(result[j].EnqueuedAt) {
|
||||
return result[i].EnqueuedAt.Before(result[j].EnqueuedAt)
|
||||
}
|
||||
return result[i].TicketID < result[j].TicketID
|
||||
})
|
||||
return result
|
||||
@@ -135,16 +165,24 @@ func (q *Queue) Candidates(now time.Time) []Candidate {
|
||||
|
||||
func (q *Queue) ownedTicket(playerID, ticketID string) (QueueTicket, error) {
|
||||
ticket, ok := q.tickets[ticketID]
|
||||
if !ok { return QueueTicket{}, ErrTicketNotFound }
|
||||
if ticket.PlayerID != playerID { return QueueTicket{}, ErrNotTicketOwner }
|
||||
if !ok {
|
||||
return QueueTicket{}, ErrTicketNotFound
|
||||
}
|
||||
if ticket.PlayerID != playerID {
|
||||
return QueueTicket{}, ErrNotTicketOwner
|
||||
}
|
||||
return ticket, nil
|
||||
}
|
||||
|
||||
func createPayload(playerID, ticketID string, candidate Candidate) string {
|
||||
regions := make([]string, 0, len(candidate.PredictedRTT))
|
||||
for region := range candidate.PredictedRTT { regions = append(regions, region) }
|
||||
for region := range candidate.PredictedRTT {
|
||||
regions = append(regions, region)
|
||||
}
|
||||
sort.Strings(regions)
|
||||
rtts := make([]string, 0, len(regions))
|
||||
for _, region := range regions { rtts = append(rtts, fmt.Sprintf("%s=%.9f", region, candidate.PredictedRTT[region])) }
|
||||
for _, region := range regions {
|
||||
rtts = append(rtts, fmt.Sprintf("%s=%.9f", region, candidate.PredictedRTT[region]))
|
||||
}
|
||||
return strings.Join([]string{playerID, ticketID, candidate.PlayerID, candidate.TicketID, fmt.Sprintf("%.9f", candidate.Rating), candidate.EnqueuedAt.UTC().Format(time.RFC3339Nano), strings.Join(rtts, ",")}, "\x00")
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user