feat: add ranked season maintenance role

This commit is contained in:
Josh Creek
2026-09-01 09:46:57 +01:00
parent e729570010
commit 4bbaf0976f
5 changed files with 167 additions and 3 deletions
+67
View File
@@ -0,0 +1,67 @@
package main
import (
"context"
"database/sql"
"flag"
"fmt"
"log"
"os"
"os/signal"
"syscall"
"time"
"github.com/cosmic-clash/cosmic-clash/server/migrations"
"github.com/cosmic-clash/cosmic-clash/server/store"
_ "github.com/jackc/pgx/v5/stdlib"
)
func main() {
dsn := flag.String("dsn", os.Getenv("COSMIC_CLASH_POSTGRES_DSN"), "PostgreSQL connection string")
migrationDir := flag.String("migrations", "migrations", "directory containing numbered SQL migrations")
interval := flag.Duration("interval", time.Minute, "maintenance poll interval")
batch := flag.Int("batch", 100, "maximum player rollovers per pass")
flag.Parse()
if *dsn == "" {
fatalf("--dsn or COSMIC_CLASH_POSTGRES_DSN is required")
}
if *interval <= 0 || *batch < 1 || *batch > 1000 {
fatalf("invalid interval or batch")
}
db, err := sql.Open("pgx", *dsn)
if err != nil {
fatalf("open PostgreSQL: %v", err)
}
defer db.Close()
startupCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
if err := db.PingContext(startupCtx); err != nil {
fatalf("ping PostgreSQL: %v", err)
}
if err := migrations.Apply(startupCtx, db, *migrationDir); err != nil {
fatalf("apply migrations: %v", err)
}
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
for {
count, err := store.RolloverDueSeasons(ctx, db, time.Now().UTC(), *batch)
if err != nil {
fatalf("season maintenance: %v", err)
}
if count > 0 {
log.Printf("applied %d ranked season rollovers", count)
}
timer := time.NewTimer(*interval)
select {
case <-ctx.Done():
timer.Stop()
return
case <-timer.C:
}
}
}
func fatalf(format string, args ...any) {
fmt.Fprintf(os.Stderr, "maintenance: "+format+"\n", args...)
os.Exit(1)
}
+67
View File
@@ -0,0 +1,67 @@
package store
import (
"context"
"database/sql"
"fmt"
"time"
"github.com/cosmic-clash/cosmic-clash/server/domain"
)
const DueSeasonRolloversSQL = `SELECT s.season_id, r.player_id, r.rating, r.deviation,
r.volatility, r.ranked_games
FROM seasons s
CROSS JOIN ratings r
LEFT JOIN ranked_season_rollovers rr ON rr.season_id = s.season_id AND rr.player_id = r.player_id
WHERE s.playlist = 'ranked' AND s.ends_at <= $1 AND rr.player_id IS NULL
ORDER BY s.ends_at, s.season_id, r.player_id
LIMIT $2`
const MarkSeasonRolledOverSQL = `UPDATE seasons SET rolled_over_at = $2
WHERE season_id = $1 AND rolled_over_at IS NULL
AND NOT EXISTS (SELECT 1 FROM ratings r
LEFT JOIN ranked_season_rollovers rr ON rr.season_id = $1 AND rr.player_id = r.player_id
WHERE rr.player_id IS NULL)`
type dueSeasonRollover struct {
seasonID string
playerID string
profile domain.RankedProfile
}
// RolloverDueSeasons processes a bounded batch. Each player update is its own
// exactly-once SERIALIZABLE transaction, so a worker crash can safely resume.
func RolloverDueSeasons(ctx context.Context, db *sql.DB, now time.Time, limit int) (int, error) {
if db == nil || now.IsZero() || limit < 1 || limit > 1000 {
return 0, fmt.Errorf("invalid season maintenance arguments")
}
rows, err := db.QueryContext(ctx, DueSeasonRolloversSQL, now, limit)
if err != nil {
return 0, err
}
defer rows.Close()
var due []dueSeasonRollover
for rows.Next() {
var item dueSeasonRollover
if err := rows.Scan(&item.seasonID, &item.playerID, &item.profile.Value, &item.profile.RD, &item.profile.Volatility, &item.profile.RankedGames); err != nil {
return 0, err
}
due = append(due, item)
}
if err := rows.Err(); err != nil {
return 0, err
}
count := 0
for _, item := range due {
if _, applied, err := ApplyRankedSeasonRollover(ctx, db, item.playerID, item.seasonID, item.profile, now); err != nil {
return count, err
} else if applied {
count++
}
if _, err := db.ExecContext(ctx, MarkSeasonRolledOverSQL, item.seasonID, now); err != nil {
return count, err
}
}
return count, nil
}
+28
View File
@@ -0,0 +1,28 @@
package store
import (
"testing"
"time"
)
func TestMaintenanceSQLEnumeratesOnlyUnrolledRankedPlayers(t *testing.T) {
for _, fragment := range []string{"s.playlist = 'ranked'", "ends_at <= $1", "rr.player_id IS NULL", "ORDER BY s.ends_at", "LIMIT $2"} {
if !contains(DueSeasonRolloversSQL, fragment) {
t.Fatalf("due query missing %q", fragment)
}
}
for _, fragment := range []string{"rolled_over_at IS NULL", "NOT EXISTS", "ranked_season_rollovers"} {
if !contains(MarkSeasonRolledOverSQL, fragment) {
t.Fatalf("mark query missing %q", fragment)
}
}
}
func TestRolloverDueSeasonsRejectsUnboundedMaintenance(t *testing.T) {
if _, err := RolloverDueSeasons(nil, nil, time.Unix(1000, 0), 0); err == nil {
t.Fatal("zero batch accepted")
}
if _, err := RolloverDueSeasons(nil, nil, time.Unix(1000, 0), 1001); err == nil {
t.Fatal("oversized batch accepted")
}
}