mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-11 00:14:00 +00:00
feat: promote accepted proposals into matches
This commit is contained in:
@@ -0,0 +1,227 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"sort"
|
||||
"time"
|
||||
|
||||
"github.com/cosmic-clash/cosmic-clash/server/domain"
|
||||
)
|
||||
|
||||
// AcceptedMatchPlan is the durable hand-off from an accepted proposal to
|
||||
// allocation. Team and slot originate from the matcher formation and are
|
||||
// persisted before allocation so later roster issuance cannot re-partition a
|
||||
// match after players have accepted it.
|
||||
type AcceptedMatchPlan struct {
|
||||
MatchID string
|
||||
ProposalID string
|
||||
Region string
|
||||
Protocol int
|
||||
Players []MatchPlayer
|
||||
}
|
||||
|
||||
type MatchPlayer struct {
|
||||
PlayerID string
|
||||
Team int
|
||||
Slot int
|
||||
}
|
||||
|
||||
const AcceptedProposalLockSQL = `SELECT playlist, state
|
||||
FROM proposals
|
||||
WHERE proposal_id = $1
|
||||
FOR UPDATE`
|
||||
|
||||
const AcceptedProposalParticipantsSQL = `SELECT player_id, ticket_id, response
|
||||
FROM proposal_participants
|
||||
WHERE proposal_id = $1
|
||||
ORDER BY player_id
|
||||
FOR UPDATE`
|
||||
|
||||
const AcceptedMatchInsertSQL = `INSERT INTO matches
|
||||
(match_id, playlist, state, region, protocol_version)
|
||||
VALUES ($1, $2, 'ALLOCATING', $3, $4)
|
||||
ON CONFLICT (match_id) DO NOTHING`
|
||||
|
||||
const AcceptedMatchSelectSQL = `SELECT playlist, state, region, protocol_version, server_id
|
||||
FROM matches
|
||||
WHERE match_id = $1
|
||||
FOR UPDATE`
|
||||
|
||||
const AcceptedMatchParticipantsSQL = `SELECT player_id, ticket_id, slot, team
|
||||
FROM match_participants
|
||||
WHERE match_id = $1
|
||||
ORDER BY player_id`
|
||||
|
||||
const AcceptedTicketSQL = `UPDATE queue_tickets
|
||||
SET state = 'ACCEPTED', revision = revision + 1
|
||||
WHERE ticket_id = $1 AND player_id = $2 AND state = 'PROPOSED'
|
||||
RETURNING protocol_version`
|
||||
|
||||
const AcceptedMatchParticipantInsertSQL = `INSERT INTO match_participants
|
||||
(match_id, player_id, ticket_id, slot, team)
|
||||
VALUES ($1, $2, $3, $4, $5)`
|
||||
|
||||
// CreateMatchFromAcceptedProposal atomically promotes the exact accepted
|
||||
// roster into an ALLOCATING match. An existing match ID is an idempotent retry
|
||||
// only if every durable field and participant assignment matches the request.
|
||||
func CreateMatchFromAcceptedProposal(ctx context.Context, db *sql.DB, plan AcceptedMatchPlan, now time.Time) error {
|
||||
if db == nil || now.IsZero() || !validAcceptedMatchPlan(plan) {
|
||||
return fmt.Errorf("invalid accepted match plan")
|
||||
}
|
||||
return RunSerializable(ctx, db, DefaultSerializableAttempts, func(ctx context.Context, tx *sql.Tx) error {
|
||||
var playlist, proposalState string
|
||||
if err := tx.QueryRowContext(ctx, AcceptedProposalLockSQL, plan.ProposalID).Scan(&playlist, &proposalState); err != nil {
|
||||
return err
|
||||
}
|
||||
if proposalState != string(domain.Accepted) {
|
||||
return fmt.Errorf("proposal is not accepted")
|
||||
}
|
||||
participants, err := acceptedProposalParticipants(ctx, tx, plan)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
inserted, err := tx.ExecContext(ctx, AcceptedMatchInsertSQL, plan.MatchID, playlist, plan.Region, plan.Protocol)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
changed, err := inserted.RowsAffected()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if changed == 0 {
|
||||
return verifyAcceptedMatchReplay(ctx, tx, plan, domain.Playlist(playlist), participants)
|
||||
}
|
||||
for _, player := range plan.Players {
|
||||
ticketID := participants[player.PlayerID]
|
||||
var protocol int
|
||||
if err := tx.QueryRowContext(ctx, AcceptedTicketSQL, ticketID, player.PlayerID).Scan(&protocol); err != nil {
|
||||
return fmt.Errorf("accepted ticket transition: %w", err)
|
||||
}
|
||||
if protocol != plan.Protocol {
|
||||
return fmt.Errorf("accepted ticket protocol mismatch")
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, AcceptedMatchParticipantInsertSQL, plan.MatchID, player.PlayerID, ticketID, player.Slot, player.Team); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
func validAcceptedMatchPlan(plan AcceptedMatchPlan) bool {
|
||||
if plan.MatchID == "" || plan.ProposalID == "" || (plan.Region != "EU" && plan.Region != "NA") || plan.Protocol < 1 || len(plan.Players) < 2 || len(plan.Players) > 6 {
|
||||
return false
|
||||
}
|
||||
players := make(map[string]struct{}, len(plan.Players))
|
||||
slots := make(map[int]struct{}, len(plan.Players))
|
||||
teams := [2]int{}
|
||||
for _, player := range plan.Players {
|
||||
if player.PlayerID == "" || player.Team < 0 || player.Team > 1 || player.Slot < 0 || player.Slot > 5 {
|
||||
return false
|
||||
}
|
||||
if _, exists := players[player.PlayerID]; exists {
|
||||
return false
|
||||
}
|
||||
if _, exists := slots[player.Slot]; exists {
|
||||
return false
|
||||
}
|
||||
players[player.PlayerID] = struct{}{}
|
||||
slots[player.Slot] = struct{}{}
|
||||
teams[player.Team]++
|
||||
}
|
||||
return teams[0] > 0 && teams[1] > 0
|
||||
}
|
||||
|
||||
func acceptedProposalParticipants(ctx context.Context, tx *sql.Tx, plan AcceptedMatchPlan) (map[string]string, error) {
|
||||
rows, err := tx.QueryContext(ctx, AcceptedProposalParticipantsSQL, plan.ProposalID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
participants := make(map[string]string, len(plan.Players))
|
||||
for rows.Next() {
|
||||
var playerID, ticketID, response string
|
||||
if err := rows.Scan(&playerID, &ticketID, &response); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if response != string(domain.AcceptedResponse) {
|
||||
return nil, fmt.Errorf("proposal participant has not accepted")
|
||||
}
|
||||
participants[playerID] = ticketID
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(participants) != len(plan.Players) {
|
||||
return nil, fmt.Errorf("proposal participants do not match accepted plan")
|
||||
}
|
||||
for _, player := range plan.Players {
|
||||
if participants[player.PlayerID] == "" {
|
||||
return nil, fmt.Errorf("accepted plan includes non-participant")
|
||||
}
|
||||
}
|
||||
return participants, nil
|
||||
}
|
||||
|
||||
func verifyAcceptedMatchReplay(ctx context.Context, tx *sql.Tx, plan AcceptedMatchPlan, playlist domain.Playlist, tickets map[string]string) error {
|
||||
var existingPlaylist, state, region string
|
||||
var protocol int
|
||||
var serverID sql.NullString
|
||||
if err := tx.QueryRowContext(ctx, AcceptedMatchSelectSQL, plan.MatchID).Scan(&existingPlaylist, &state, ®ion, &protocol, &serverID); err != nil {
|
||||
return err
|
||||
}
|
||||
if existingPlaylist != string(playlist) || state != string(domain.Allocating) || region != plan.Region || protocol != plan.Protocol || serverID.Valid {
|
||||
return domain.ErrConflict
|
||||
}
|
||||
rows, err := tx.QueryContext(ctx, AcceptedMatchParticipantsSQL, plan.MatchID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer rows.Close()
|
||||
existing := make(map[string]MatchPlayer, len(plan.Players))
|
||||
for rows.Next() {
|
||||
var player MatchPlayer
|
||||
var ticketID string
|
||||
if err := rows.Scan(&player.PlayerID, &ticketID, &player.Slot, &player.Team); err != nil {
|
||||
return err
|
||||
}
|
||||
if tickets[player.PlayerID] != ticketID {
|
||||
return domain.ErrConflict
|
||||
}
|
||||
existing[player.PlayerID] = player
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return err
|
||||
}
|
||||
if len(existing) != len(plan.Players) {
|
||||
return domain.ErrConflict
|
||||
}
|
||||
for _, player := range plan.Players {
|
||||
if existing[player.PlayerID] != player {
|
||||
return domain.ErrConflict
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// MatchPlayersFromTeams turns the deterministic matcher partition into the
|
||||
// persisted six-slot topology. Each team is sorted by player ID first, so slot
|
||||
// assignment does not depend on cache/database row order.
|
||||
func MatchPlayersFromTeams(teams domain.Teams) ([]MatchPlayer, error) {
|
||||
if len(teams.Team0) == 0 || len(teams.Team1) == 0 || len(teams.Team0)+len(teams.Team1) > 6 {
|
||||
return nil, fmt.Errorf("invalid match teams")
|
||||
}
|
||||
result := make([]MatchPlayer, 0, len(teams.Team0)+len(teams.Team1))
|
||||
add := func(team int, players []domain.Candidate) {
|
||||
ordered := append([]domain.Candidate(nil), players...)
|
||||
sort.Slice(ordered, func(i, j int) bool { return ordered[i].PlayerID < ordered[j].PlayerID })
|
||||
for index, player := range ordered {
|
||||
result = append(result, MatchPlayer{PlayerID: player.PlayerID, Team: team, Slot: team*3 + index})
|
||||
}
|
||||
}
|
||||
add(0, teams.Team0)
|
||||
add(1, teams.Team1)
|
||||
return result, nil
|
||||
}
|
||||
@@ -0,0 +1,74 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/cosmic-clash/cosmic-clash/server/domain"
|
||||
)
|
||||
|
||||
func TestAcceptedMatchSQLPreservesAtomicProposalToMatchBoundary(t *testing.T) {
|
||||
checks := map[string][]string{
|
||||
AcceptedProposalLockSQL: {"FOR UPDATE", "proposal_id = $1"},
|
||||
AcceptedProposalParticipantsSQL: {"response", "ORDER BY player_id", "FOR UPDATE"},
|
||||
AcceptedMatchInsertSQL: {"'ALLOCATING'", "ON CONFLICT (match_id) DO NOTHING"},
|
||||
AcceptedTicketSQL: {"state = 'ACCEPTED'", "state = 'PROPOSED'", "revision = revision + 1"},
|
||||
AcceptedMatchParticipantInsertSQL: {"match_participants", "slot", "team"},
|
||||
}
|
||||
for query, fragments := range checks {
|
||||
for _, fragment := range fragments {
|
||||
if !contains(query, fragment) {
|
||||
t.Fatalf("query %q missing %q", query, fragment)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAcceptedMatchPlanRejectsInvalidPlansBeforeDatabaseUse(t *testing.T) {
|
||||
valid := AcceptedMatchPlan{
|
||||
MatchID: "match-1", ProposalID: "proposal-1", Region: "EU", Protocol: 1,
|
||||
Players: []MatchPlayer{{PlayerID: "player-a", Team: 0, Slot: 0}, {PlayerID: "player-b", Team: 1, Slot: 3}},
|
||||
}
|
||||
if !validAcceptedMatchPlan(valid) {
|
||||
t.Fatal("valid accepted match plan rejected")
|
||||
}
|
||||
for name, mutate := range map[string]func(*AcceptedMatchPlan){
|
||||
"no second team": func(p *AcceptedMatchPlan) { p.Players[1].Team = 0 },
|
||||
"duplicate slot": func(p *AcceptedMatchPlan) { p.Players[1].Slot = 0 },
|
||||
"duplicate player": func(p *AcceptedMatchPlan) { p.Players[1].PlayerID = "player-a" },
|
||||
"bad region": func(p *AcceptedMatchPlan) { p.Region = "AP" },
|
||||
} {
|
||||
plan := valid
|
||||
plan.Players = append([]MatchPlayer(nil), valid.Players...)
|
||||
mutate(&plan)
|
||||
if validAcceptedMatchPlan(plan) {
|
||||
t.Fatalf("%s plan accepted", name)
|
||||
}
|
||||
}
|
||||
if err := CreateMatchFromAcceptedProposal(nil, nil, valid, time.Now()); err == nil {
|
||||
t.Fatal("nil database accepted")
|
||||
}
|
||||
}
|
||||
|
||||
func TestMatchPlayersFromTeamsUsesDeterministicTeamSlots(t *testing.T) {
|
||||
teams := domain.Teams{
|
||||
Team0: []domain.Candidate{{PlayerID: "bravo"}, {PlayerID: "alpha"}},
|
||||
Team1: []domain.Candidate{{PlayerID: "delta"}, {PlayerID: "charlie"}},
|
||||
}
|
||||
players, err := MatchPlayersFromTeams(teams)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
want := []MatchPlayer{
|
||||
{PlayerID: "alpha", Team: 0, Slot: 0}, {PlayerID: "bravo", Team: 0, Slot: 1},
|
||||
{PlayerID: "charlie", Team: 1, Slot: 3}, {PlayerID: "delta", Team: 1, Slot: 4},
|
||||
}
|
||||
if len(players) != len(want) {
|
||||
t.Fatalf("players = %+v", players)
|
||||
}
|
||||
for index := range want {
|
||||
if players[index] != want[index] {
|
||||
t.Fatalf("player %d = %+v, want %+v", index, players[index], want[index])
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -94,6 +94,55 @@ func TestPostgreSQLAllocatorClaimReplayAndCapacityFence(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestPostgreSQLAcceptedProposalPromotesOneAtomicAllocatingMatch(t *testing.T) {
|
||||
db := openIntegrationPostgres(t)
|
||||
applyIntegrationMigrations(t, db)
|
||||
now := time.Now().UTC().Truncate(time.Microsecond)
|
||||
ctx := context.Background()
|
||||
for _, player := range []string{"promote-a", "promote-b"} {
|
||||
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ($1, $1)`, player); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
for index, player := range []string{"promote-a", "promote-b"} {
|
||||
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ($1, $2, 'casual', 'PROPOSED', 'build-1', 1, $3, $4)`, fmt.Sprintf("promote-ticket-%d", index), player, now, now.Add(time.Minute)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if _, err := db.ExecContext(ctx, `INSERT INTO proposals (proposal_id, playlist, state, expires_at, revision) VALUES ('promote-proposal', 'casual', 'ACCEPTED', $1, 2)`, now.Add(time.Minute)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for index, player := range []string{"promote-a", "promote-b"} {
|
||||
if _, err := db.ExecContext(ctx, `INSERT INTO proposal_participants (proposal_id, player_id, ticket_id, response) VALUES ('promote-proposal', $1, $2, 'ACCEPTED')`, player, fmt.Sprintf("promote-ticket-%d", index)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
plan := AcceptedMatchPlan{MatchID: "promote-match", ProposalID: "promote-proposal", Region: "EU", Protocol: 1, Players: []MatchPlayer{{PlayerID: "promote-a", Team: 0, Slot: 0}, {PlayerID: "promote-b", Team: 1, Slot: 3}}}
|
||||
if err := CreateMatchFromAcceptedProposal(ctx, db, plan, now); err != nil {
|
||||
t.Fatalf("promote accepted proposal: %v", err)
|
||||
}
|
||||
var state string
|
||||
if err := db.QueryRowContext(ctx, `SELECT state FROM matches WHERE match_id = 'promote-match'`).Scan(&state); err != nil || state != "ALLOCATING" {
|
||||
t.Fatalf("match state=%q err=%v", state, err)
|
||||
}
|
||||
var acceptedTickets, participantCount int
|
||||
if err := db.QueryRowContext(ctx, `SELECT count(*) FROM queue_tickets WHERE ticket_id LIKE 'promote-ticket-%' AND state = 'ACCEPTED'`).Scan(&acceptedTickets); err != nil || acceptedTickets != 2 {
|
||||
t.Fatalf("accepted tickets=%d err=%v", acceptedTickets, err)
|
||||
}
|
||||
if err := db.QueryRowContext(ctx, `SELECT count(*) FROM match_participants WHERE match_id = 'promote-match'`).Scan(&participantCount); err != nil || participantCount != 2 {
|
||||
t.Fatalf("participants=%d err=%v", participantCount, err)
|
||||
}
|
||||
if err := CreateMatchFromAcceptedProposal(ctx, db, plan, now.Add(time.Second)); err != nil {
|
||||
t.Fatalf("identical match promotion replay: %v", err)
|
||||
}
|
||||
conflict := plan
|
||||
conflict.Players = append([]MatchPlayer(nil), plan.Players...)
|
||||
conflict.Players[1].Slot = 4
|
||||
if err := CreateMatchFromAcceptedProposal(ctx, db, conflict, now.Add(2*time.Second)); err == nil {
|
||||
t.Fatal("conflicting match promotion replay was accepted")
|
||||
}
|
||||
}
|
||||
|
||||
func TestPostgreSQLQueueAdapterAgainstRealDatabase(t *testing.T) {
|
||||
db := openIntegrationPostgres(t)
|
||||
applyIntegrationMigrations(t, db)
|
||||
|
||||
Reference in New Issue
Block a user