mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-11 00:14:00 +00:00
feat: add durable queue ticket repository
This commit is contained in:
@@ -0,0 +1,101 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/cosmic-clash/cosmic-clash/server/domain"
|
||||
)
|
||||
|
||||
const (
|
||||
QueueIdempotencyScope = "queue.create"
|
||||
QueueIdempotencyInsertSQL = `INSERT INTO idempotency_keys (scope, idempotency_key, payload_digest, result)
|
||||
VALUES ($1, $2, $3, $4)
|
||||
ON CONFLICT (scope, idempotency_key) DO NOTHING`
|
||||
QueueIdempotencySelectSQL = `SELECT payload_digest, result
|
||||
FROM idempotency_keys
|
||||
WHERE scope = $1 AND idempotency_key = $2
|
||||
FOR UPDATE`
|
||||
QueueTicketSelectSQL = `SELECT ticket_id, player_id, playlist, state, client_build,
|
||||
protocol_version, enqueued_at, expires_at, revision
|
||||
FROM queue_tickets
|
||||
WHERE ticket_id = $1 AND player_id = $2`
|
||||
)
|
||||
|
||||
func CreateQueueTicket(ctx context.Context, db *sql.DB, ticketID, playerID, idempotencyKey string, spec domain.QueueSpec, now time.Time) (domain.QueueTicket, error) {
|
||||
if db == nil || ticketID == "" || playerID == "" || len(idempotencyKey) < 16 || len(idempotencyKey) > 128 || (spec.Playlist != domain.Casual && spec.Playlist != domain.Ranked) || spec.ClientBuild == "" || len(spec.ClientBuild) > 128 || spec.ProtocolVersion < 1 || now.IsZero() {
|
||||
return domain.QueueTicket{}, fmt.Errorf("invalid queue transaction arguments")
|
||||
}
|
||||
digest := sha256.Sum256([]byte(fmt.Sprintf("%s|%s|%s|%s|%d", ticketID, playerID, spec.Playlist, spec.ClientBuild, spec.ProtocolVersion)))
|
||||
var ticket domain.QueueTicket
|
||||
err := RunSerializable(ctx, db, DefaultSerializableAttempts, func(ctx context.Context, tx *sql.Tx) error {
|
||||
candidate := domain.Candidate{TicketID: ticketID, PlayerID: playerID, Playlist: spec.Playlist, ClientBuild: spec.ClientBuild, ProtocolVersion: spec.ProtocolVersion, EnqueuedAt: now}
|
||||
ticket = domain.QueueTicket{TicketID: ticketID, PlayerID: playerID, Candidate: candidate, Playlist: spec.Playlist, State: domain.Queued, EnqueuedAt: now, ExpiresAt: now.Add(domain.QueueExpiryWindow)}
|
||||
stored, err := json.Marshal(queueTicketRecordFromDomain(ticket))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
result, err := tx.ExecContext(ctx, QueueIdempotencyInsertSQL, QueueIdempotencyScope, idempotencyKey, digest[:], stored)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
inserted, err := result.RowsAffected()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if inserted == 0 {
|
||||
var priorDigest, priorResult []byte
|
||||
if err := tx.QueryRowContext(ctx, QueueIdempotencySelectSQL, QueueIdempotencyScope, idempotencyKey).Scan(&priorDigest, &priorResult); err != nil {
|
||||
return err
|
||||
}
|
||||
if !bytes.Equal(priorDigest, digest[:]) {
|
||||
return fmt.Errorf("queue create idempotency conflict")
|
||||
}
|
||||
var prior queueTicketRecord
|
||||
if err := json.Unmarshal(priorResult, &prior); err != nil {
|
||||
return fmt.Errorf("invalid stored queue result: %w", err)
|
||||
}
|
||||
ticket = queueTicketRecordToDomain(prior)
|
||||
return nil
|
||||
}
|
||||
_, err = tx.ExecContext(ctx, QueueTicketInsertSQL, ticketID, playerID, string(spec.Playlist), string(domain.Queued), spec.ClientBuild, spec.ProtocolVersion, now, ticket.ExpiresAt)
|
||||
return err
|
||||
})
|
||||
return ticket, err
|
||||
}
|
||||
|
||||
type queueTicketRecord struct {
|
||||
TicketID string `json:"ticket_id"`
|
||||
PlayerID string `json:"player_id"`
|
||||
Playlist string `json:"playlist"`
|
||||
State string `json:"state"`
|
||||
ClientBuild string `json:"client_build"`
|
||||
ProtocolVersion int `json:"protocol_version"`
|
||||
EnqueuedAt time.Time `json:"enqueued_at"`
|
||||
ExpiresAt time.Time `json:"expires_at"`
|
||||
Revision uint64 `json:"revision"`
|
||||
}
|
||||
|
||||
func GetQueueTicket(ctx context.Context, db *sql.DB, playerID, ticketID string) (domain.QueueTicket, error) {
|
||||
if db == nil || playerID == "" || ticketID == "" {
|
||||
return domain.QueueTicket{}, fmt.Errorf("invalid queue recovery arguments")
|
||||
}
|
||||
var record queueTicketRecord
|
||||
if err := db.QueryRowContext(ctx, QueueTicketSelectSQL, ticketID, playerID).Scan(&record.TicketID, &record.PlayerID, &record.Playlist, &record.State, &record.ClientBuild, &record.ProtocolVersion, &record.EnqueuedAt, &record.ExpiresAt, &record.Revision); err != nil {
|
||||
return domain.QueueTicket{}, err
|
||||
}
|
||||
return queueTicketRecordToDomain(record), nil
|
||||
}
|
||||
|
||||
func queueTicketRecordFromDomain(ticket domain.QueueTicket) queueTicketRecord {
|
||||
return queueTicketRecord{ticket.TicketID, ticket.PlayerID, string(ticket.Playlist), string(ticket.State), ticket.Candidate.ClientBuild, ticket.Candidate.ProtocolVersion, ticket.EnqueuedAt, ticket.ExpiresAt, ticket.Revision}
|
||||
}
|
||||
func queueTicketRecordToDomain(record queueTicketRecord) domain.QueueTicket {
|
||||
candidate := domain.Candidate{TicketID: record.TicketID, PlayerID: record.PlayerID, Playlist: domain.Playlist(record.Playlist), ClientBuild: record.ClientBuild, ProtocolVersion: record.ProtocolVersion, EnqueuedAt: record.EnqueuedAt}
|
||||
return domain.QueueTicket{TicketID: record.TicketID, PlayerID: record.PlayerID, Candidate: candidate, Playlist: domain.Playlist(record.Playlist), State: domain.State(record.State), Revision: record.Revision, EnqueuedAt: record.EnqueuedAt, ExpiresAt: record.ExpiresAt}
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"github.com/cosmic-clash/cosmic-clash/server/domain"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestQueueSQLUsesDurableIdempotencyAndOwnerScopedRecovery(t *testing.T) {
|
||||
for query, fragments := range map[string][]string{
|
||||
QueueIdempotencyInsertSQL: {"idempotency_keys", "ON CONFLICT (scope, idempotency_key) DO NOTHING", "payload_digest"},
|
||||
QueueIdempotencySelectSQL: {"scope = $1", "idempotency_key = $2", "FOR UPDATE"},
|
||||
QueueTicketSelectSQL: {"ticket_id = $1", "player_id = $2"}, QueueTicketInsertSQL: {"player_id", "playlist", "client_build", "protocol_version"},
|
||||
} {
|
||||
for _, fragment := range fragments {
|
||||
if !contains(query, fragment) {
|
||||
t.Fatalf("query %q missing %q", query, fragment)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestCreateQueueTicketRejectsInvalidArgumentsWithoutDatabase(t *testing.T) {
|
||||
if _, err := CreateQueueTicket(nil, nil, "ticket-1", "player-1", "short", domain.QueueSpec{Playlist: domain.Casual, ClientBuild: "build-1", ProtocolVersion: 1}, time.Unix(1000, 0)); err == nil {
|
||||
t.Fatal("invalid arguments accepted")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user