mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-11 00:14:00 +00:00
feat: persist queue heartbeat and cancel mutations
This commit is contained in:
@@ -25,6 +25,16 @@ FOR UPDATE`
|
||||
protocol_version, enqueued_at, expires_at, revision
|
||||
FROM queue_tickets
|
||||
WHERE ticket_id = $1 AND player_id = $2`
|
||||
QueueTicketHeartbeatSQL = `UPDATE queue_tickets SET revision = revision + 1,
|
||||
expires_at = $4 + INTERVAL '30 seconds'
|
||||
WHERE ticket_id = $1 AND player_id = $2 AND revision = $3
|
||||
AND state IN ('QUEUED', 'PROPOSED') AND expires_at > $4
|
||||
RETURNING ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at, revision`
|
||||
QueueTicketCancelSQL = `UPDATE queue_tickets SET state = 'CANCELLED',
|
||||
revision = revision + 1, expires_at = $4
|
||||
WHERE ticket_id = $1 AND player_id = $2 AND revision = $3
|
||||
AND state NOT IN ('COMPLETED', 'CANCELLED', 'EXPIRED', 'FAILED')
|
||||
RETURNING ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at, revision`
|
||||
)
|
||||
|
||||
func CreateQueueTicket(ctx context.Context, db *sql.DB, ticketID, playerID, idempotencyKey string, spec domain.QueueSpec, now time.Time) (domain.QueueTicket, error) {
|
||||
@@ -92,6 +102,58 @@ func GetQueueTicket(ctx context.Context, db *sql.DB, playerID, ticketID string)
|
||||
return queueTicketRecordToDomain(record), nil
|
||||
}
|
||||
|
||||
func HeartbeatQueueTicket(ctx context.Context, db *sql.DB, playerID, ticketID, idempotencyKey string, expectedRevision uint64, now time.Time) (domain.QueueTicket, error) {
|
||||
return mutateQueueTicket(ctx, db, playerID, ticketID, idempotencyKey, expectedRevision, now, "heartbeat", QueueTicketHeartbeatSQL)
|
||||
}
|
||||
|
||||
func CancelQueueTicket(ctx context.Context, db *sql.DB, playerID, ticketID, idempotencyKey string, expectedRevision uint64, now time.Time) (domain.QueueTicket, error) {
|
||||
return mutateQueueTicket(ctx, db, playerID, ticketID, idempotencyKey, expectedRevision, now, "cancel", QueueTicketCancelSQL)
|
||||
}
|
||||
|
||||
func mutateQueueTicket(ctx context.Context, db *sql.DB, playerID, ticketID, idempotencyKey string, expectedRevision uint64, now time.Time, operation, mutationSQL string) (ticket domain.QueueTicket, err error) {
|
||||
if db == nil || playerID == "" || ticketID == "" || len(idempotencyKey) < 16 || len(idempotencyKey) > 128 || now.IsZero() || (operation != "heartbeat" && operation != "cancel") {
|
||||
return domain.QueueTicket{}, fmt.Errorf("invalid queue mutation arguments")
|
||||
}
|
||||
digest := sha256.Sum256([]byte(fmt.Sprintf("%s|%s|%s|%d", operation, playerID, ticketID, expectedRevision)))
|
||||
err = RunSerializable(ctx, db, DefaultSerializableAttempts, func(ctx context.Context, tx *sql.Tx) error {
|
||||
result, err := tx.ExecContext(ctx, QueueIdempotencyInsertSQL, QueueIdempotencyScope+"."+operation, 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, QueueIdempotencySelectSQL, QueueIdempotencyScope+"."+operation, idempotencyKey).Scan(&priorDigest, &priorResult); err != nil {
|
||||
return err
|
||||
}
|
||||
if !bytes.Equal(priorDigest, digest[:]) {
|
||||
return fmt.Errorf("queue mutation 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
|
||||
}
|
||||
var record queueTicketRecord
|
||||
if err := tx.QueryRowContext(ctx, mutationSQL, ticketID, playerID, expectedRevision, now).Scan(&record.TicketID, &record.PlayerID, &record.Playlist, &record.State, &record.ClientBuild, &record.ProtocolVersion, &record.EnqueuedAt, &record.ExpiresAt, &record.Revision); err != nil {
|
||||
return fmt.Errorf("queue mutation rejected: %w", err)
|
||||
}
|
||||
ticket = queueTicketRecordToDomain(record)
|
||||
stored, err := json.Marshal(record)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = tx.ExecContext(ctx, `UPDATE idempotency_keys SET result = $3 WHERE scope = $1 AND idempotency_key = $2`, QueueIdempotencyScope+"."+operation, idempotencyKey, stored)
|
||||
return err
|
||||
})
|
||||
return ticket, err
|
||||
}
|
||||
|
||||
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}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user