mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-11 18:03:43 +00:00
feat: persist queue probe metadata
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
-- Persist only server-derived placement metadata alongside queue ownership.
|
||||
-- Redis remains a rebuildable index; this JSON projection is durable source
|
||||
-- data and may be empty until the authenticated probe completes.
|
||||
ALTER TABLE queue_tickets
|
||||
ADD COLUMN predicted_rtt JSONB NOT NULL DEFAULT '{}'::jsonb;
|
||||
+36
-19
@@ -22,23 +22,23 @@ FROM idempotency_keys
|
||||
WHERE scope = $1 AND idempotency_key = $2
|
||||
FOR UPDATE`
|
||||
QueueTicketSelectSQL = `SELECT ticket_id, player_id, playlist, state, client_build,
|
||||
protocol_version, enqueued_at, expires_at, revision
|
||||
protocol_version, enqueued_at, expires_at, revision, predicted_rtt
|
||||
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`
|
||||
RETURNING ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at, revision, predicted_rtt`
|
||||
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`
|
||||
RETURNING ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at, revision, predicted_rtt`
|
||||
)
|
||||
|
||||
const QueueCandidateProjectionSQL = `SELECT ticket_id, player_id, playlist, client_build,
|
||||
protocol_version, enqueued_at
|
||||
protocol_version, enqueued_at, predicted_rtt
|
||||
FROM queue_tickets
|
||||
WHERE state = 'QUEUED' AND playlist = $1 AND expires_at > $2
|
||||
ORDER BY enqueued_at, ticket_id
|
||||
@@ -60,9 +60,13 @@ func ListQueuedCandidates(ctx context.Context, db *sql.DB, playlist domain.Playl
|
||||
for rows.Next() {
|
||||
var candidate domain.Candidate
|
||||
var playlist string
|
||||
if err := rows.Scan(&candidate.TicketID, &candidate.PlayerID, &playlist, &candidate.ClientBuild, &candidate.ProtocolVersion, &candidate.EnqueuedAt); err != nil {
|
||||
var predictedRTT []byte
|
||||
if err := rows.Scan(&candidate.TicketID, &candidate.PlayerID, &playlist, &candidate.ClientBuild, &candidate.ProtocolVersion, &candidate.EnqueuedAt, &predictedRTT); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := json.Unmarshal(predictedRTT, &candidate.PredictedRTT); err != nil {
|
||||
return nil, fmt.Errorf("decode candidate RTT: %w", err)
|
||||
}
|
||||
candidate.Playlist = domain.Playlist(playlist)
|
||||
candidates = append(candidates, candidate)
|
||||
}
|
||||
@@ -108,22 +112,27 @@ func CreateQueueTicket(ctx context.Context, db *sql.DB, ticketID, playerID, idem
|
||||
ticket = queueTicketRecordToDomain(prior)
|
||||
return nil
|
||||
}
|
||||
_, err = tx.ExecContext(ctx, QueueTicketInsertSQL, ticketID, playerID, string(spec.Playlist), string(domain.Queued), spec.ClientBuild, spec.ProtocolVersion, now, ticket.ExpiresAt)
|
||||
predictedRTT, err := json.Marshal(candidate.PredictedRTT)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = tx.ExecContext(ctx, QueueTicketInsertSQL, ticketID, playerID, string(spec.Playlist), string(domain.Queued), spec.ClientBuild, spec.ProtocolVersion, now, ticket.ExpiresAt, predictedRTT)
|
||||
return err
|
||||
})
|
||||
return ticket, err
|
||||
}
|
||||
|
||||
type queueTicketRecord struct {
|
||||
TicketID string `json:"ticket_id"`
|
||||
PlayerID string `json:"player_id"`
|
||||
Playlist string `json:"playlist"`
|
||||
State string `json:"state"`
|
||||
ClientBuild string `json:"client_build"`
|
||||
ProtocolVersion int `json:"protocol_version"`
|
||||
EnqueuedAt time.Time `json:"enqueued_at"`
|
||||
ExpiresAt time.Time `json:"expires_at"`
|
||||
Revision uint64 `json:"revision"`
|
||||
TicketID string `json:"ticket_id"`
|
||||
PlayerID string `json:"player_id"`
|
||||
Playlist string `json:"playlist"`
|
||||
State string `json:"state"`
|
||||
ClientBuild string `json:"client_build"`
|
||||
ProtocolVersion int `json:"protocol_version"`
|
||||
EnqueuedAt time.Time `json:"enqueued_at"`
|
||||
ExpiresAt time.Time `json:"expires_at"`
|
||||
Revision uint64 `json:"revision"`
|
||||
PredictedRTT map[string]float64 `json:"predicted_rtt"`
|
||||
}
|
||||
|
||||
type PostgresQueue struct{ DB *sql.DB }
|
||||
@@ -146,9 +155,13 @@ func GetQueueTicket(ctx context.Context, db *sql.DB, playerID, ticketID string,
|
||||
return domain.QueueTicket{}, fmt.Errorf("invalid queue recovery arguments")
|
||||
}
|
||||
var record queueTicketRecord
|
||||
if err := db.QueryRowContext(ctx, QueueTicketSelectSQL, ticketID, playerID).Scan(&record.TicketID, &record.PlayerID, &record.Playlist, &record.State, &record.ClientBuild, &record.ProtocolVersion, &record.EnqueuedAt, &record.ExpiresAt, &record.Revision); err != nil {
|
||||
var predictedRTT []byte
|
||||
if err := db.QueryRowContext(ctx, QueueTicketSelectSQL, ticketID, playerID).Scan(&record.TicketID, &record.PlayerID, &record.Playlist, &record.State, &record.ClientBuild, &record.ProtocolVersion, &record.EnqueuedAt, &record.ExpiresAt, &record.Revision, &predictedRTT); err != nil {
|
||||
return domain.QueueTicket{}, err
|
||||
}
|
||||
if err := json.Unmarshal(predictedRTT, &record.PredictedRTT); err != nil {
|
||||
return domain.QueueTicket{}, fmt.Errorf("decode queue RTT: %w", err)
|
||||
}
|
||||
ticket := queueTicketRecordToDomain(record)
|
||||
if (ticket.State == domain.Queued || ticket.State == domain.Proposed) && !now.Before(ticket.ExpiresAt) {
|
||||
return domain.QueueTicket{}, domain.ErrTicketExpired
|
||||
@@ -194,9 +207,13 @@ func mutateQueueTicket(ctx context.Context, db *sql.DB, playerID, ticketID, idem
|
||||
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 {
|
||||
var predictedRTT []byte
|
||||
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, &predictedRTT); err != nil {
|
||||
return fmt.Errorf("queue mutation rejected: %w", err)
|
||||
}
|
||||
if err := json.Unmarshal(predictedRTT, &record.PredictedRTT); err != nil {
|
||||
return fmt.Errorf("decode queue RTT: %w", err)
|
||||
}
|
||||
ticket = queueTicketRecordToDomain(record)
|
||||
stored, err := json.Marshal(record)
|
||||
if err != nil {
|
||||
@@ -209,9 +226,9 @@ func mutateQueueTicket(ctx context.Context, db *sql.DB, playerID, ticketID, idem
|
||||
}
|
||||
|
||||
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}
|
||||
return queueTicketRecord{TicketID: ticket.TicketID, PlayerID: ticket.PlayerID, Playlist: string(ticket.Playlist), State: string(ticket.State), ClientBuild: ticket.Candidate.ClientBuild, ProtocolVersion: ticket.Candidate.ProtocolVersion, EnqueuedAt: ticket.EnqueuedAt, ExpiresAt: ticket.ExpiresAt, Revision: ticket.Revision, PredictedRTT: ticket.Candidate.PredictedRTT}
|
||||
}
|
||||
func queueTicketRecordToDomain(record queueTicketRecord) domain.QueueTicket {
|
||||
candidate := domain.Candidate{TicketID: record.TicketID, PlayerID: record.PlayerID, Playlist: domain.Playlist(record.Playlist), ClientBuild: record.ClientBuild, ProtocolVersion: record.ProtocolVersion, EnqueuedAt: record.EnqueuedAt}
|
||||
candidate := domain.Candidate{TicketID: record.TicketID, PlayerID: record.PlayerID, Playlist: domain.Playlist(record.Playlist), ClientBuild: record.ClientBuild, ProtocolVersion: record.ProtocolVersion, EnqueuedAt: record.EnqueuedAt, PredictedRTT: record.PredictedRTT}
|
||||
return domain.QueueTicket{TicketID: record.TicketID, PlayerID: record.PlayerID, Candidate: candidate, Playlist: domain.Playlist(record.Playlist), State: domain.State(record.State), Revision: record.Revision, EnqueuedAt: record.EnqueuedAt, ExpiresAt: record.ExpiresAt}
|
||||
}
|
||||
|
||||
@@ -14,6 +14,7 @@ func TestQueueSQLUsesDurableIdempotencyAndOwnerScopedRecovery(t *testing.T) {
|
||||
QueueTicketInsertSQL: {"player_id", "playlist", "client_build", "protocol_version"},
|
||||
QueueTicketHeartbeatSQL: {"player_id = $2", "revision = $3", "expires_at > $4", "RETURNING"},
|
||||
QueueTicketCancelSQL: {"player_id = $2", "revision = $3", "state NOT IN", "RETURNING"},
|
||||
QueueCandidateProjectionSQL: {"playlist = $1", "predicted_rtt", "expires_at > $2", "LIMIT $3"},
|
||||
} {
|
||||
for _, fragment := range fragments {
|
||||
if !contains(query, fragment) {
|
||||
@@ -23,6 +24,21 @@ func TestQueueSQLUsesDurableIdempotencyAndOwnerScopedRecovery(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestListQueuedCandidatesRejectsUnscopedOrUnboundedReads(t *testing.T) {
|
||||
now := time.Unix(1000, 0)
|
||||
for _, playlist := range []domain.Playlist{"", "invalid"} {
|
||||
if _, err := ListQueuedCandidates(nil, nil, playlist, now, 4); err == nil {
|
||||
t.Fatalf("playlist %q accepted", playlist)
|
||||
}
|
||||
}
|
||||
if _, err := ListQueuedCandidates(nil, nil, domain.Casual, now, 0); err == nil {
|
||||
t.Fatal("zero limit accepted")
|
||||
}
|
||||
if _, err := ListQueuedCandidates(nil, nil, domain.Casual, now, 1001); err == nil {
|
||||
t.Fatal("unbounded limit accepted")
|
||||
}
|
||||
}
|
||||
|
||||
func TestQueueMutationAdaptersRejectInvalidArgumentsWithoutDatabase(t *testing.T) {
|
||||
now := time.Unix(1000, 0)
|
||||
if _, err := HeartbeatQueueTicket(nil, nil, "player-1", "ticket-1", "short", 0, now); err == nil {
|
||||
|
||||
@@ -62,8 +62,8 @@ var (
|
||||
// QueueTicketInsertSQL relies on the partial unique index in migration 0001
|
||||
// as the cross-replica one-active-ticket fence.
|
||||
QueueTicketInsertSQL = `INSERT INTO queue_tickets
|
||||
(ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at)
|
||||
VALUES ($1, $2, $3, 'QUEUED', $4, $5, $6, $7)`
|
||||
(ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at, predicted_rtt)
|
||||
VALUES ($1, $2, $3, 'QUEUED', $4, $5, $6, $7, $8)`
|
||||
|
||||
// CandidateClaimSQL must run in the same serializable transaction as
|
||||
// ProposalParticipantInsertSQL. SKIP LOCKED lets another matcher continue,
|
||||
|
||||
Reference in New Issue
Block a user