Files
CosmicClash/server/store/proposal_recovery_sql.go
T
Josh Creek 4627dd58fb fix(multiplayer): requeue every participant after a proposal times out
The timeout sibling of the previous commit's decline fix: a proposal
that simply times out (the 10s window elapses with no unanimous
response) hits ProposalExpireSQL/ProposalParticipantExpireSQL, and
neither of those -- same as the decline path -- ever touched
queue_tickets. Same severe consequence: every participant still
holding a PROPOSED ticket, response pending or already accepted, is
left stranded (invisible to the matcher, blocking a fresh
queue_create, renewable forever by heartbeat) with no automatic way
back into matchmaking. This path is reached from both GetProposal
(the recovery/read boundary -- a client that missed the expiry event
entirely) and RespondToProposal (a response arriving after the
window), so both needed the fix.

ProposalExpireRequeueSQL mirrors ProposalDeclineRequeueSQL, guarded on
state = 'EXPIRED' so it's safe to call unconditionally right after
ProposalExpireSQL: a no-op on a proposal that's still OPEN, and a
no-op on a proposal that was already EXPIRED on a prior pass (nothing
left at PROPOSED to requeue a second time).

Covered by a real PostgreSQL integration test via GetProposal (nobody
ever responds; recovering the proposal well after its window expires
it and must requeue both participants), confirming both tickets land
back at QUEUED with a refreshed expiry and are visible again to
ListQueuedCandidates. Clean across 5 runs, plus the full integration
and unit suites.
2026-09-01 14:36:34 +01:00

299 lines
12 KiB
Go

package store
import (
"bytes"
"context"
"crypto/sha256"
"database/sql"
"encoding/json"
"fmt"
"time"
"github.com/cosmic-clash/cosmic-clash/server/domain"
)
const ProposalExpireSQL = `UPDATE proposals
SET state = 'EXPIRED', revision = revision + 1
WHERE proposal_id = $1 AND state = 'OPEN' AND expires_at <= $2`
const ProposalParticipantExpireSQL = `UPDATE proposal_participants
SET response = 'TIMED_OUT', responded_at = $2
WHERE proposal_id = $1 AND response = 'PENDING'
AND EXISTS (SELECT 1 FROM proposals WHERE proposals.proposal_id = proposal_participants.proposal_id AND proposals.expires_at <= $2)`
// ProposalExpireRequeueSQL is the timeout sibling of
// ProposalDeclineRequeueSQL: a proposal that simply times out (no unanimous
// response inside the 10s window) leaves any participant still holding a
// PROPOSED ticket exactly as stranded as an explicit decline does, and for
// the identical reason -- nothing else ever moves a PROPOSED ticket back to
// QUEUED. The `state = 'EXPIRED'` guard makes this safe to call
// unconditionally right after ProposalExpireSQL: it's a no-op on a proposal
// that was already OPEN and stays OPEN (nothing to requeue) or one that was
// already EXPIRED on a prior pass (its participants' tickets, if any were
// still PROPOSED, were already requeued then).
const ProposalExpireRequeueSQL = `UPDATE queue_tickets q
SET state = 'QUEUED', expires_at = $2, revision = revision + 1
FROM proposal_participants pp
WHERE pp.proposal_id = $1 AND q.ticket_id = pp.ticket_id AND q.player_id = pp.player_id AND q.state = 'PROPOSED'
AND EXISTS (SELECT 1 FROM proposals WHERE proposals.proposal_id = $1 AND proposals.state = 'EXPIRED')`
const ProposalRecoverySelectSQL = `SELECT proposal_id, playlist, state, revision, expires_at
FROM proposals
WHERE proposal_id = $1
AND EXISTS (SELECT 1 FROM proposal_participants WHERE proposal_id = proposals.proposal_id AND player_id = $2)`
const ProposalParticipantsSelectSQL = `SELECT player_id, response
FROM proposal_participants
WHERE proposal_id = $1
ORDER BY player_id`
const ProposalResponseIdempotencyScope = "proposal.respond"
const ProposalResponseIdempotencyInsertSQL = `INSERT INTO idempotency_keys
(scope, idempotency_key, payload_digest, result)
VALUES ($1, $2, $3, $4)
ON CONFLICT (scope, idempotency_key) DO NOTHING`
const ProposalResponseIdempotencySelectSQL = `SELECT payload_digest, result
FROM idempotency_keys
WHERE scope = $1 AND idempotency_key = $2
FOR UPDATE`
const ProposalLockSQL = `SELECT playlist, state, revision, expires_at
FROM proposals
WHERE proposal_id = $1
FOR UPDATE`
const ProposalParticipantLockSQL = `SELECT response
FROM proposal_participants
WHERE proposal_id = $1 AND player_id = $2
FOR UPDATE`
const ProposalParticipantRespondSQL = `UPDATE proposal_participants
SET response = $3, responded_at = $4
WHERE proposal_id = $1 AND player_id = $2 AND response = 'PENDING'`
const ProposalCountPendingSQL = `SELECT COUNT(*)
FROM proposal_participants
WHERE proposal_id = $1 AND response = 'PENDING'`
const ProposalAcceptSQL = `UPDATE proposals
SET state = 'ACCEPTED', revision = revision + 1
WHERE proposal_id = $1 AND state = 'OPEN'`
const ProposalDeclineSQL = `UPDATE proposals
SET state = 'DECLINED', revision = revision + 1
WHERE proposal_id = $1 AND state = 'OPEN'`
// ProposalDeclineRequeueSQL requeues every participant's ticket, including
// the decliner's own: nothing yet enforces the decline cooldown §8.17
// documents as a separate, not-yet-built feature, so leaving any ticket
// behind at PROPOSED here isn't "cooldown behaviour", it's just a stranded
// ticket -- invisible to the matcher (which only ever reads state='QUEUED'),
// still counted as this player's one active ticket (blocking a fresh
// queue_create), and renewable forever by an ordinary heartbeat, so a player
// left in this state has no path back into matchmaking without realising
// they need to cancel and start over. Once §8.17's cooldown exists, it can
// exempt the decliner from this immediate requeue; today nothing does.
const ProposalDeclineRequeueSQL = `UPDATE queue_tickets q
SET state = 'QUEUED', expires_at = $2, revision = revision + 1
FROM proposal_participants pp
WHERE pp.proposal_id = $1 AND q.ticket_id = pp.ticket_id AND q.player_id = pp.player_id AND q.state = 'PROPOSED'`
const ProposalRevisionBumpSQL = `UPDATE proposals
SET revision = revision + 1
WHERE proposal_id = $1 AND state = 'OPEN'`
var ErrProposalResponseConflict = fmt.Errorf("proposal response conflict")
// GetProposal recovers the full proposal only after proving the caller is a
// participant. Expiry is advanced in the same transaction as the read so a
// missed event cannot leave a durable proposal indefinitely OPEN.
func GetProposal(ctx context.Context, db *sql.DB, playerID, proposalID string, now time.Time) (domain.Proposal, error) {
if db == nil || playerID == "" || proposalID == "" || now.IsZero() {
return domain.Proposal{}, fmt.Errorf("invalid proposal recovery arguments")
}
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return domain.Proposal{}, err
}
defer tx.Rollback()
if _, err := tx.ExecContext(ctx, ProposalExpireSQL, proposalID, now); err != nil {
return domain.Proposal{}, err
}
if _, err := tx.ExecContext(ctx, ProposalParticipantExpireSQL, proposalID, now); err != nil {
return domain.Proposal{}, err
}
if _, err := tx.ExecContext(ctx, ProposalExpireRequeueSQL, proposalID, now.Add(domain.QueueExpiryWindow)); err != nil {
return domain.Proposal{}, err
}
var proposal domain.Proposal
var playlist, state string
if err := tx.QueryRowContext(ctx, ProposalRecoverySelectSQL, proposalID, playerID).Scan(&proposal.ProposalID, &playlist, &state, &proposal.Revision, &proposal.ExpiresAt); err != nil {
return domain.Proposal{}, err
}
proposal.Playlist = domain.Playlist(playlist)
proposal.State = domain.State(state)
rows, err := tx.QueryContext(ctx, ProposalParticipantsSelectSQL, proposalID)
if err != nil {
return domain.Proposal{}, err
}
defer rows.Close()
for rows.Next() {
var participant domain.ProposalParticipant
if err := rows.Scan(&participant.PlayerID, &participant.Response); err != nil {
return domain.Proposal{}, err
}
proposal.Participants = append(proposal.Participants, participant)
}
if err := rows.Err(); err != nil {
return domain.Proposal{}, err
}
if len(proposal.Participants) == 0 {
return domain.Proposal{}, fmt.Errorf("proposal has no participants")
}
if err := tx.Commit(); err != nil {
return domain.Proposal{}, err
}
return proposal, nil
}
// RespondToProposal is the durable mutation counterpart to GetProposal. The
// proposal row and participant row are locked in one transaction; the result
// is stored under the idempotency key before the transaction commits.
func RespondToProposal(ctx context.Context, db *sql.DB, playerID, proposalID, idempotencyKey string, accept bool, expectedRevision uint64, now time.Time) (domain.Proposal, error) {
if db == nil || playerID == "" || proposalID == "" || len(idempotencyKey) < 16 || len(idempotencyKey) > 128 || now.IsZero() {
return domain.Proposal{}, fmt.Errorf("invalid proposal response arguments")
}
digest := sha256.Sum256([]byte(fmt.Sprintf("%s|%s|%t|%d", playerID, proposalID, accept, expectedRevision)))
var proposal domain.Proposal
err := RunSerializable(ctx, db, DefaultSerializableAttempts, func(ctx context.Context, tx *sql.Tx) error {
result, err := tx.ExecContext(ctx, ProposalResponseIdempotencyInsertSQL, ProposalResponseIdempotencyScope, idempotencyKey, digest[:], []byte("{}"))
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, ProposalResponseIdempotencySelectSQL, ProposalResponseIdempotencyScope, idempotencyKey).Scan(&priorDigest, &priorResult); err != nil {
return err
}
if !bytes.Equal(priorDigest, digest[:]) {
return ErrProposalResponseConflict
}
if err := json.Unmarshal(priorResult, &proposal); err != nil {
return fmt.Errorf("invalid stored proposal response: %w", err)
}
return nil
}
var playlist, state string
var revision uint64
var expiresAt time.Time
if err := tx.QueryRowContext(ctx, ProposalLockSQL, proposalID).Scan(&playlist, &state, &revision, &expiresAt); err != nil {
return err
}
// A mutation is also a recovery boundary. If the response arrives after
// the window, advance both the proposal and its pending participants in
// this same transaction before returning the closed error. Otherwise a
// client that missed the expiry event could observe OPEN/PENDING forever
// when its first durable interaction is an accept/decline.
if _, err := tx.ExecContext(ctx, ProposalExpireSQL, proposalID, now); err != nil {
return err
}
if _, err := tx.ExecContext(ctx, ProposalParticipantExpireSQL, proposalID, now); err != nil {
return err
}
if _, err := tx.ExecContext(ctx, ProposalExpireRequeueSQL, proposalID, now.Add(domain.QueueExpiryWindow)); err != nil {
return err
}
if !now.Before(expiresAt) {
return domain.ErrProposalClosed
}
if state != string(domain.Open) || !now.Before(expiresAt) {
return domain.ErrProposalClosed
}
if revision != expectedRevision {
return domain.ErrStaleRevision
}
var response string
if err := tx.QueryRowContext(ctx, ProposalParticipantLockSQL, proposalID, playerID).Scan(&response); err != nil {
return domain.ErrNotParticipant
}
if response != string(domain.Pending) {
return domain.ErrConflict
}
response = string(domain.DeclinedResponse)
if accept {
response = string(domain.AcceptedResponse)
}
changed, err := tx.ExecContext(ctx, ProposalParticipantRespondSQL, proposalID, playerID, response, now)
if err != nil {
return err
}
if count, err := changed.RowsAffected(); err != nil || count != 1 {
return ErrProposalResponseConflict
}
targetState := string(domain.Declined)
if accept {
var pending int
if err := tx.QueryRowContext(ctx, ProposalCountPendingSQL, proposalID).Scan(&pending); err != nil {
return err
}
if pending == 0 {
targetState = string(domain.Accepted)
} else {
targetState = state
}
}
if targetState != state {
if targetState == string(domain.Accepted) {
_, err = tx.ExecContext(ctx, ProposalAcceptSQL, proposalID)
} else {
_, err = tx.ExecContext(ctx, ProposalDeclineSQL, proposalID)
if err != nil {
return err
}
_, err = tx.ExecContext(ctx, ProposalDeclineRequeueSQL, proposalID, now.Add(domain.QueueExpiryWindow))
}
if err != nil {
return err
}
revision++
} else {
if _, err := tx.ExecContext(ctx, ProposalRevisionBumpSQL, proposalID); err != nil {
return err
}
revision++
}
proposal = domain.Proposal{ProposalID: proposalID, Playlist: domain.Playlist(playlist), State: domain.State(targetState), Revision: revision, ExpiresAt: expiresAt}
rows, err := tx.QueryContext(ctx, ProposalParticipantsSelectSQL, proposalID)
if err != nil {
return err
}
for rows.Next() {
var participant domain.ProposalParticipant
if err := rows.Scan(&participant.PlayerID, &participant.Response); err != nil {
rows.Close()
return err
}
proposal.Participants = append(proposal.Participants, participant)
}
if err := rows.Err(); err != nil {
rows.Close()
return err
}
rows.Close()
stored, err := json.Marshal(proposal)
if err != nil {
return err
}
_, err = tx.ExecContext(ctx, `UPDATE idempotency_keys SET result = $3 WHERE scope = $1 AND idempotency_key = $2`, ProposalResponseIdempotencyScope, idempotencyKey, stored)
return err
})
return proposal, err
}