fix: fence concurrent queue mutations

This commit is contained in:
Josh Creek
2026-08-31 21:09:58 +01:00
parent 58e8a5c523
commit caa875f15c
3 changed files with 48 additions and 2 deletions
+17 -1
View File
@@ -6,6 +6,7 @@ import (
"fmt"
"sort"
"strings"
"sync"
"time"
)
@@ -37,6 +38,7 @@ type queueMutation struct {
}
type Queue struct {
mu sync.Mutex
tickets map[string]QueueTicket
byPlayer map[string]string
mutations map[string]queueMutation
@@ -50,6 +52,8 @@ func NewQueue() *Queue {
// production adapter must perform the same check in one transaction and use
// the same idempotency semantics.
func (q *Queue) Create(playerID, ticketID, idempotencyKey string, candidate Candidate, now time.Time) (QueueTicket, error) {
q.mu.Lock()
defer q.mu.Unlock()
digest := sha256.Sum256([]byte(createPayload(playerID, ticketID, candidate)))
if prior, ok := q.mutations[idempotencyKey]; ok {
if prior.digest != digest {
@@ -74,6 +78,8 @@ 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) {
q.mu.Lock()
defer q.mu.Unlock()
digest := sha256.Sum256([]byte(fmt.Sprintf("heartbeat:%s:%d", ticketID, expectedRevision)))
if prior, ok := q.mutations[idempotencyKey]; ok {
if prior.digest != digest {
@@ -105,6 +111,8 @@ func (q *Queue) Heartbeat(playerID, ticketID, idempotencyKey string, expectedRev
}
func (q *Queue) Cancel(playerID, ticketID, idempotencyKey string, expectedRevision uint64, now time.Time) (QueueTicket, error) {
q.mu.Lock()
defer q.mu.Unlock()
digest := sha256.Sum256([]byte(fmt.Sprintf("cancel:%s:%d", ticketID, expectedRevision)))
if prior, ok := q.mutations[idempotencyKey]; ok {
if prior.digest != digest {
@@ -132,6 +140,12 @@ func (q *Queue) Cancel(playerID, ticketID, idempotencyKey string, expectedRevisi
}
func (q *Queue) Expire(now time.Time) []QueueTicket {
q.mu.Lock()
defer q.mu.Unlock()
return q.expireLocked(now)
}
func (q *Queue) expireLocked(now time.Time) []QueueTicket {
var expired []QueueTicket
for id, ticket := range q.tickets {
if (ticket.State == Queued || ticket.State == Proposed) && !now.Before(ticket.ExpiresAt) {
@@ -147,7 +161,9 @@ func (q *Queue) Expire(now time.Time) []QueueTicket {
}
func (q *Queue) Candidates(now time.Time) []Candidate {
q.Expire(now)
q.mu.Lock()
defer q.mu.Unlock()
q.expireLocked(now)
result := make([]Candidate, 0)
for _, ticket := range q.tickets {
if ticket.State == Queued {
+30
View File
@@ -3,6 +3,7 @@ package domain
import (
"errors"
"reflect"
"sync"
"testing"
"time"
)
@@ -77,3 +78,32 @@ func TestQueueCreateIdempotencyIncludesCandidatePayload(t *testing.T) {
t.Fatalf("changed create payload error = %v", err)
}
}
func TestQueueConcurrentCreateKeepsOneActiveTicketPerPlayer(t *testing.T) {
q := NewQueue()
now := time.Unix(1000, 0)
var wg sync.WaitGroup
results := make(chan error, 2)
for i := 0; i < 2; i++ {
wg.Add(1)
go func(i int) {
defer wg.Done()
id := string(rune('a' + i))
_, err := q.Create("same-player", "ticket-"+id, "create-"+id+"-123456", Candidate{TicketID: "ticket-" + id, PlayerID: "same-player", EnqueuedAt: now}, now)
results <- err
}(i)
}
wg.Wait()
close(results)
succeeded := 0
for err := range results {
if err == nil {
succeeded++
} else if !errors.Is(err, ErrPlayerQueued) {
t.Fatalf("unexpected concurrent create error: %v", err)
}
}
if succeeded != 1 {
t.Fatalf("concurrent creates succeeded %d times", succeeded)
}
}