mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-10 16:04:04 +00:00
79318b56bd
The last two commits fixed the severe stranding bug in decline and timeout, but left a real responsiveness gap: cancelling a ticket directly while it's part of an OPEN proposal used to leave the OTHER participant waiting out the full response window for something the system already knew couldn't happen -- their proposal partner just abandoned the queue. ProposalExpireRequeueSQL eventually rescues them, but only after the full window elapses, not immediately. CascadeCancelToOpenProposal runs inside the same transaction as the cancel itself: if the cancelled ticket belonged to a currently-OPEN proposal, decline that proposal right now and requeue every other participant immediately via the same ProposalDeclineRequeueSQL the decline path already uses. The cancelling player's own ticket correctly stays CANCELLED, not swept back into the requeue meant for everyone else (ProposalDeclineRequeueSQL only touches tickets still at PROPOSED). Covered by a real PostgreSQL integration test: cancelling one participant's ticket mid-proposal immediately declines the proposal and requeues the other participant with a refreshed expiry, while the cancelling player's own ticket stays CANCELLED. First draft used a stale expected revision (0) for the cancel call -- CreateProposal's own QueueTicketProposeSQL already bumps a ticket's revision to 1 when forming the proposal, caught immediately by actually running the test against real Postgres rather than assuming. Clean across 5 runs after the fix, plus the full integration and unit suites.
328 lines
13 KiB
Go
328 lines
13 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 OpenProposalForCancelledTicketSQL = `SELECT pp.proposal_id
|
|
FROM proposal_participants pp
|
|
JOIN proposals p ON p.proposal_id = pp.proposal_id
|
|
WHERE pp.ticket_id = $1 AND pp.player_id = $2 AND p.state = 'OPEN'`
|
|
|
|
// CascadeCancelToOpenProposal declines and requeues an OPEN proposal
|
|
// immediately when one of its participants cancels their own queue ticket
|
|
// directly, rather than leaving every other participant to wait out the
|
|
// full response window for something the system already knows can't happen
|
|
// -- ProposalExpireRequeueSQL would eventually rescue them anyway, but not
|
|
// for up to ProposalWindow's full duration for no reason. Must run inside
|
|
// the same transaction as the ticket cancel itself; a no-op if the ticket
|
|
// wasn't part of any currently-OPEN proposal.
|
|
func CascadeCancelToOpenProposal(ctx context.Context, tx *sql.Tx, ticketID, playerID string, now time.Time) error {
|
|
var proposalID string
|
|
err := tx.QueryRowContext(ctx, OpenProposalForCancelledTicketSQL, ticketID, playerID).Scan(&proposalID)
|
|
if err == sql.ErrNoRows {
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if _, err := tx.ExecContext(ctx, ProposalDeclineSQL, proposalID); err != nil {
|
|
return err
|
|
}
|
|
_, err = tx.ExecContext(ctx, ProposalDeclineRequeueSQL, proposalID, now.Add(domain.QueueExpiryWindow))
|
|
return err
|
|
}
|
|
|
|
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
|
|
}
|