mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-11 19:03:43 +00:00
6237a25a69
Found by reading the code, not a failing test: no path anywhere transitioned a queue ticket from PROPOSED back to QUEUED after a proposal was declined. A stranded PROPOSED ticket is invisible to the matcher (ListQueuedCandidates only ever reads state='QUEUED'), still counts as that player's one active ticket (blocking a fresh queue_create), and is renewable forever by an ordinary heartbeat -- a player proposed a match with someone who then declines had no way back into matchmaking without realising, on their own, that they needed to manually cancel first. This affects every participant, not just the decliner: an uninvolved player who never even responded was left stuck by someone else's decision. ProposalDeclineRequeueSQL requeues every participant's ticket, including the decliner's own -- nothing yet enforces the decline cooldown task 8.17 documents as a separate, not-yet-built feature, so leaving anyone behind at PROPOSED today isn't "cooldown behaviour", it's just broken. Once that cooldown exists it can exempt the decliner from this immediate requeue; today nothing does. Covered by a real PostgreSQL integration test: after one player declines, both the decliner's and an uninvolved participant's tickets land back at QUEUED with a refreshed expiry, and -- the actual end-to-end regression -- both are visible again to ListQueuedCandidates, the same query the matcher itself uses. Clean across 5 runs, plus the full integration and unit suites.
277 lines
10 KiB
Go
277 lines
10 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)`
|
|
|
|
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
|
|
}
|
|
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 !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
|
|
}
|