package repo import ( "context" "database/sql" "errors" "fmt" "strings" "sync" "time" "browser.local/platform/domain" "github.com/go-sql-driver/mysql" ) // The SCUM tables are platform-owned storage for the typed SCUM facts Run // pushes for one server instance. Columns mirror platform/model/scum.go and // carry the indexes the platform queries below rely on. const ( scumMaxListLimit = 1000 scumTrajectoryRetention = 30 * 24 * time.Hour scumTrajectoryPruneEvery = time.Hour ) var scumTableDDL = []struct { table string statement string }{ {"scum_user", `CREATE TABLE IF NOT EXISTS scum_user ( id VARCHAR(191) PRIMARY KEY, server_instance_id VARCHAR(191) NOT NULL, steam_id VARCHAR(64) NOT NULL, display_name VARCHAR(191) NOT NULL DEFAULT '', squad_id BIGINT NOT NULL DEFAULT 0, registered_at DATETIME(6) NULL, last_login_at DATETIME(6) NULL, last_activity_at DATETIME(6) NULL, last_login_ip VARCHAR(64) NOT NULL DEFAULT '', online BOOLEAN NOT NULL DEFAULT FALSE, x DOUBLE NULL, y DOUBLE NULL, z DOUBLE NULL, bank_balance BIGINT NULL, gold_bars BIGINT NULL, ridden_vehicle_id VARCHAR(191) NOT NULL DEFAULT '', game_vehicle_id VARCHAR(191) NOT NULL DEFAULT '', welcome_queued_at DATETIME(6) NULL, welcome_last_error VARCHAR(512) NOT NULL DEFAULT '', created_at DATETIME(6) NOT NULL, updated_at DATETIME(6) NOT NULL, UNIQUE KEY uq_scum_user_instance_steam (server_instance_id, steam_id), KEY ix_scum_user_presence (server_instance_id, online, last_activity_at), KEY ix_scum_user_updated (server_instance_id, updated_at) )`}, {"scum_user_trajectory", `CREATE TABLE IF NOT EXISTS scum_user_trajectory ( id VARCHAR(191) PRIMARY KEY, server_instance_id VARCHAR(191) NOT NULL, scum_user_id VARCHAR(191) NOT NULL, steam_id VARCHAR(64) NOT NULL, display_name VARCHAR(191) NOT NULL DEFAULT '', x DOUBLE NOT NULL, y DOUBLE NOT NULL, z DOUBLE NOT NULL, ridden_vehicle_id VARCHAR(191) NOT NULL DEFAULT '', game_vehicle_id VARCHAR(191) NOT NULL DEFAULT '', observed_at DATETIME(6) NOT NULL, created_at DATETIME(6) NOT NULL, KEY ix_scum_user_track_subject (server_instance_id, scum_user_id, observed_at), KEY ix_scum_user_track_steam (server_instance_id, steam_id, observed_at), KEY ix_scum_user_track_retention (observed_at) )`}, {"scum_vehicle", `CREATE TABLE IF NOT EXISTS scum_vehicle ( id VARCHAR(191) PRIMARY KEY, server_instance_id VARCHAR(191) NOT NULL, game_vehicle_id VARCHAR(191) NOT NULL, vehicle_class VARCHAR(64) NOT NULL DEFAULT '', display_name VARCHAR(191) NOT NULL DEFAULT '', exists_in_game BOOLEAN NOT NULL DEFAULT FALSE, locked BOOLEAN NOT NULL DEFAULT FALSE, functional BOOLEAN NULL, x DOUBLE NULL, y DOUBLE NULL, z DOUBLE NULL, last_observed_at DATETIME(6) NULL, created_at DATETIME(6) NOT NULL, updated_at DATETIME(6) NOT NULL, UNIQUE KEY uq_scum_vehicle_instance_game (server_instance_id, game_vehicle_id), KEY ix_scum_vehicle_presence (server_instance_id, exists_in_game, last_observed_at), KEY ix_scum_vehicle_updated (server_instance_id, updated_at) )`}, {"scum_vehicle_trajectory", `CREATE TABLE IF NOT EXISTS scum_vehicle_trajectory ( id VARCHAR(191) PRIMARY KEY, server_instance_id VARCHAR(191) NOT NULL, scum_vehicle_id VARCHAR(191) NOT NULL, game_vehicle_id VARCHAR(191) NOT NULL, vehicle_class VARCHAR(64) NOT NULL DEFAULT '', x DOUBLE NOT NULL, y DOUBLE NOT NULL, z DOUBLE NOT NULL, observed_at DATETIME(6) NOT NULL, created_at DATETIME(6) NOT NULL, KEY ix_scum_vehicle_track_vehicle (server_instance_id, scum_vehicle_id, observed_at), KEY ix_scum_vehicle_track_game (server_instance_id, game_vehicle_id, observed_at), KEY ix_scum_vehicle_track_retention (observed_at) )`}, {"scum_vehicle_lock", `CREATE TABLE IF NOT EXISTS scum_vehicle_lock ( id VARCHAR(191) PRIMARY KEY, server_instance_id VARCHAR(191) NOT NULL, scum_vehicle_id VARCHAR(191) NOT NULL, game_vehicle_id VARCHAR(191) NOT NULL, scum_user_id VARCHAR(191) NOT NULL, steam_id VARCHAR(64) NOT NULL, locked_at DATETIME(6) NOT NULL, created_at DATETIME(6) NOT NULL, KEY ix_scum_vehicle_lock_vehicle (server_instance_id, scum_vehicle_id, locked_at), KEY ix_scum_vehicle_lock_user (server_instance_id, scum_user_id, locked_at) )`}, } // initializeSCUMTables creates the typed SCUM tables. An earlier unreleased // build created a throwaway snapshot-shaped table with a payload JSON column // that no code path ever read; that shape is replaced here instead of being // kept around with missing typed columns. func (store *MySQLStore) initializeSCUMTables(ctx context.Context) error { for _, entry := range scumTableDDL { legacy, err := store.scumTableUsesLegacyPayload(ctx, entry.table) if err != nil { return err } if legacy { if _, err := store.db.ExecContext(ctx, "DROP TABLE IF EXISTS "+entry.table); err != nil { return fmt.Errorf("drop legacy SCUM table %s: %w", entry.table, err) } } if _, err := store.db.ExecContext(ctx, entry.statement); err != nil { return fmt.Errorf("create SCUM table %s: %w", entry.table, err) } } if err := store.ensureSCUMVehicleFunctionalColumn(ctx); err != nil { return err } return nil } func (store *MySQLStore) ensureSCUMVehicleFunctionalColumn(ctx context.Context) error { var found int if err := store.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'scum_vehicle' AND COLUMN_NAME = 'functional'`).Scan(&found); err != nil { return fmt.Errorf("inspect SCUM vehicle functional column: %w", err) } if found > 0 { return nil } if _, err := store.db.ExecContext(ctx, `ALTER TABLE scum_vehicle ADD COLUMN functional BOOLEAN NULL AFTER locked`); err != nil { return fmt.Errorf("add SCUM vehicle functional column: %w", err) } return nil } func (store *MySQLStore) scumTableUsesLegacyPayload(ctx context.Context, table string) (bool, error) { var found int err := store.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND COLUMN_NAME = 'payload' AND EXISTS (SELECT 1 FROM information_schema.COLUMNS c WHERE c.TABLE_SCHEMA = DATABASE() AND c.TABLE_NAME = ? AND c.COLUMN_NAME = 'server_instance_id')`, table, table).Scan(&found) if err != nil { return false, fmt.Errorf("inspect SCUM table %s: %w", table, err) } return found > 0, nil } func (store *MySQLStore) SCUMUsers() SCUMUserRepository { return &mysqlSCUMUserRepository{db: store.db} } func (store *MySQLStore) SCUMUserTrajectories() SCUMUserTrajectoryRepository { return &mysqlSCUMUserTrajectoryRepository{db: store.db, pruner: &store.scumUserTrackPruner} } func (store *MySQLStore) SCUMVehicles() SCUMVehicleRepository { return &mysqlSCUMVehicleRepository{db: store.db} } func (store *MySQLStore) SCUMVehicleTrajectories() SCUMVehicleTrajectoryRepository { return &mysqlSCUMVehicleTrajectoryRepository{db: store.db, pruner: &store.scumVehicleTrackPruner} } func (store *MySQLStore) SCUMVehicleLocks() SCUMVehicleLockRepository { return &mysqlSCUMVehicleLockRepository{db: store.db} } type mysqlSCUMUserRepository struct{ db *sql.DB } const scumUserColumns = `id, server_instance_id, steam_id, display_name, squad_id, registered_at, last_login_at, last_activity_at, last_login_ip, online, x, y, z, bank_balance, gold_bars, ridden_vehicle_id, game_vehicle_id, welcome_queued_at, welcome_last_error, created_at, updated_at` func (repository *mysqlSCUMUserRepository) Create(value domain.SCUMUser) error { ctx, cancel := scumQueryContext() defer cancel() _, err := repository.db.ExecContext(ctx, `INSERT INTO scum_user (`+scumUserColumns+`) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`, scumUserArgs(value)...) if err != nil { if isMySQLDuplicateKey(err) { return ErrDuplicate } return fmt.Errorf("create scum_user: %w", err) } return nil } func (repository *mysqlSCUMUserRepository) Get(id string) (domain.SCUMUser, error) { ctx, cancel := scumQueryContext() defer cancel() value, err := scanSCUMUser(repository.db.QueryRowContext(ctx, `SELECT `+scumUserColumns+` FROM scum_user WHERE id = ?`, id)) if errors.Is(err, sql.ErrNoRows) { return domain.SCUMUser{}, ErrNotFound } if err != nil { return domain.SCUMUser{}, fmt.Errorf("read scum_user: %w", err) } return value, nil } func (repository *mysqlSCUMUserRepository) List(filter domain.SCUMUserFilter) ([]domain.SCUMUser, error) { query := `SELECT ` + scumUserColumns + ` FROM scum_user WHERE server_instance_id = ?` args := []any{filter.ServerInstanceID} if filter.SteamID != "" { query += " AND steam_id = ?" args = append(args, filter.SteamID) } if filter.Online != nil { query += " AND online = ?" args = append(args, *filter.Online) } if !filter.ChangedAfter.IsZero() { query += " AND updated_at > ?" args = append(args, filter.ChangedAfter.UTC()) } if !filter.StaleBefore.IsZero() { query += " AND COALESCE(last_activity_at, updated_at) < ?" args = append(args, filter.StaleBefore.UTC()) } query += " ORDER BY last_activity_at DESC, steam_id ASC LIMIT ?" args = append(args, scumListLimit(filter.Limit)) ctx, cancel := scumQueryContext() defer cancel() rows, err := repository.db.QueryContext(ctx, query, args...) if err != nil { return nil, fmt.Errorf("list scum_user: %w", err) } defer rows.Close() values := make([]domain.SCUMUser, 0, 16) for rows.Next() { value, err := scanSCUMUser(rows) if err != nil { return nil, fmt.Errorf("scan scum_user: %w", err) } values = append(values, value) } if err := rows.Err(); err != nil { return nil, fmt.Errorf("list scum_user: %w", err) } return values, nil } func (repository *mysqlSCUMUserRepository) Update(value domain.SCUMUser) error { ctx, cancel := scumQueryContext() defer cancel() result, err := repository.db.ExecContext(ctx, `UPDATE scum_user SET server_instance_id=?, steam_id=?, display_name=?, squad_id=?, registered_at=?, last_login_at=?, last_activity_at=?, last_login_ip=?, online=?, x=?, y=?, z=?, bank_balance=?, gold_bars=?, ridden_vehicle_id=?, game_vehicle_id=?, welcome_queued_at=?, welcome_last_error=?, created_at=?, updated_at=? WHERE id=?`, scumUserUpdateArgs(value)...) if err != nil { return fmt.Errorf("update scum_user: %w", err) } return scumUpdateResult(result, func() error { _, err := repository.Get(value.ID) return err }) } func (repository *mysqlSCUMUserRepository) Delete(id string) error { ctx, cancel := scumQueryContext() defer cancel() result, err := repository.db.ExecContext(ctx, `DELETE FROM scum_user WHERE id = ?`, id) if err != nil { return fmt.Errorf("delete scum_user: %w", err) } return scumDeleteResult(result) } type mysqlSCUMVehicleRepository struct{ db *sql.DB } const scumVehicleColumns = `id, server_instance_id, game_vehicle_id, vehicle_class, display_name, exists_in_game, locked, functional, x, y, z, last_observed_at, created_at, updated_at` func (repository *mysqlSCUMVehicleRepository) Create(value domain.SCUMVehicle) error { ctx, cancel := scumQueryContext() defer cancel() _, err := repository.db.ExecContext(ctx, `INSERT INTO scum_vehicle (`+scumVehicleColumns+`) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)`, scumVehicleArgs(value)...) if err != nil { if isMySQLDuplicateKey(err) { return ErrDuplicate } return fmt.Errorf("create scum_vehicle: %w", err) } return nil } func (repository *mysqlSCUMVehicleRepository) Get(id string) (domain.SCUMVehicle, error) { ctx, cancel := scumQueryContext() defer cancel() value, err := scanSCUMVehicle(repository.db.QueryRowContext(ctx, `SELECT `+scumVehicleColumns+` FROM scum_vehicle WHERE id = ?`, id)) if errors.Is(err, sql.ErrNoRows) { return domain.SCUMVehicle{}, ErrNotFound } if err != nil { return domain.SCUMVehicle{}, fmt.Errorf("read scum_vehicle: %w", err) } return value, nil } func (repository *mysqlSCUMVehicleRepository) List(filter domain.SCUMVehicleFilter) ([]domain.SCUMVehicle, error) { query := `SELECT ` + scumVehicleColumns + ` FROM scum_vehicle WHERE server_instance_id = ?` args := []any{filter.ServerInstanceID} if filter.GameVehicleID != "" { query += " AND game_vehicle_id = ?" args = append(args, filter.GameVehicleID) } if filter.Exists != nil { query += " AND exists_in_game = ?" args = append(args, *filter.Exists) } if !filter.ChangedAfter.IsZero() { query += " AND updated_at > ?" args = append(args, filter.ChangedAfter.UTC()) } query += " ORDER BY last_observed_at DESC, game_vehicle_id ASC LIMIT ?" args = append(args, scumListLimit(filter.Limit)) ctx, cancel := scumQueryContext() defer cancel() rows, err := repository.db.QueryContext(ctx, query, args...) if err != nil { return nil, fmt.Errorf("list scum_vehicle: %w", err) } defer rows.Close() values := make([]domain.SCUMVehicle, 0, 16) for rows.Next() { value, err := scanSCUMVehicle(rows) if err != nil { return nil, fmt.Errorf("scan scum_vehicle: %w", err) } values = append(values, value) } if err := rows.Err(); err != nil { return nil, fmt.Errorf("list scum_vehicle: %w", err) } return values, nil } func (repository *mysqlSCUMVehicleRepository) Update(value domain.SCUMVehicle) error { ctx, cancel := scumQueryContext() defer cancel() result, err := repository.db.ExecContext(ctx, `UPDATE scum_vehicle SET server_instance_id=?, game_vehicle_id=?, vehicle_class=?, display_name=?, exists_in_game=?, locked=?, functional=?, x=?, y=?, z=?, last_observed_at=?, created_at=?, updated_at=? WHERE id=?`, scumVehicleUpdateArgs(value)...) if err != nil { return fmt.Errorf("update scum_vehicle: %w", err) } return scumUpdateResult(result, func() error { _, err := repository.Get(value.ID) return err }) } func (repository *mysqlSCUMVehicleRepository) Delete(id string) error { ctx, cancel := scumQueryContext() defer cancel() result, err := repository.db.ExecContext(ctx, `DELETE FROM scum_vehicle WHERE id = ?`, id) if err != nil { return fmt.Errorf("delete scum_vehicle: %w", err) } return scumDeleteResult(result) } type mysqlSCUMUserTrajectoryRepository struct { db *sql.DB pruner *scumTrajectoryPruner } const scumUserTrajectoryColumns = `id, server_instance_id, scum_user_id, steam_id, display_name, x, y, z, ridden_vehicle_id, game_vehicle_id, observed_at, created_at` func (repository *mysqlSCUMUserTrajectoryRepository) Create(value domain.SCUMUserTrajectory) error { if err := repository.pruner.prune(repository.db, "scum_user_trajectory", "observed_at"); err != nil { return err } ctx, cancel := scumQueryContext() defer cancel() _, err := repository.db.ExecContext(ctx, `INSERT INTO scum_user_trajectory (`+scumUserTrajectoryColumns+`) VALUES (?,?,?,?,?,?,?,?,?,?,?,?)`, value.ID, value.ServerInstanceID, value.SCUMUserID, value.SteamID, value.DisplayName, value.X, value.Y, value.Z, value.RiddenVehicleID, value.GameVehicleID, value.ObservedAt.UTC(), value.CreatedAt.UTC()) if err != nil { if isMySQLDuplicateKey(err) { return ErrDuplicate } return fmt.Errorf("create scum_user_trajectory: %w", err) } return nil } func (repository *mysqlSCUMUserTrajectoryRepository) Get(id string) (domain.SCUMUserTrajectory, error) { ctx, cancel := scumQueryContext() defer cancel() value, err := scanSCUMUserTrajectory(repository.db.QueryRowContext(ctx, `SELECT `+scumUserTrajectoryColumns+` FROM scum_user_trajectory WHERE id = ?`, id)) if errors.Is(err, sql.ErrNoRows) { return domain.SCUMUserTrajectory{}, ErrNotFound } if err != nil { return domain.SCUMUserTrajectory{}, fmt.Errorf("read scum_user_trajectory: %w", err) } return value, nil } func (repository *mysqlSCUMUserTrajectoryRepository) List(filter domain.SCUMUserTrajectoryFilter) ([]domain.SCUMUserTrajectory, error) { query := `SELECT ` + scumUserTrajectoryColumns + ` FROM scum_user_trajectory WHERE server_instance_id = ?` args := []any{filter.ServerInstanceID} if filter.SCUMUserID != "" { query += " AND scum_user_id = ?" args = append(args, filter.SCUMUserID) } if filter.SteamID != "" { query += " AND steam_id = ?" args = append(args, filter.SteamID) } if !filter.After.IsZero() { query += " AND observed_at > ?" args = append(args, filter.After.UTC()) } query += " ORDER BY observed_at DESC, id DESC LIMIT ?" args = append(args, scumListLimit(filter.Limit)) return querySCUMUserTrajectories(repository.db, query, args) } func (repository *mysqlSCUMUserTrajectoryRepository) Update(value domain.SCUMUserTrajectory) error { ctx, cancel := scumQueryContext() defer cancel() result, err := repository.db.ExecContext(ctx, `UPDATE scum_user_trajectory SET server_instance_id=?, scum_user_id=?, steam_id=?, display_name=?, x=?, y=?, z=?, ridden_vehicle_id=?, game_vehicle_id=?, observed_at=?, created_at=? WHERE id=?`, value.ServerInstanceID, value.SCUMUserID, value.SteamID, value.DisplayName, value.X, value.Y, value.Z, value.RiddenVehicleID, value.GameVehicleID, value.ObservedAt.UTC(), value.CreatedAt.UTC(), value.ID) if err != nil { return fmt.Errorf("update scum_user_trajectory: %w", err) } return scumUpdateResult(result, func() error { _, err := repository.Get(value.ID) return err }) } func (repository *mysqlSCUMUserTrajectoryRepository) Delete(id string) error { ctx, cancel := scumQueryContext() defer cancel() result, err := repository.db.ExecContext(ctx, `DELETE FROM scum_user_trajectory WHERE id = ?`, id) if err != nil { return fmt.Errorf("delete scum_user_trajectory: %w", err) } return scumDeleteResult(result) } type mysqlSCUMVehicleTrajectoryRepository struct { db *sql.DB pruner *scumTrajectoryPruner } const scumVehicleTrajectoryColumns = `id, server_instance_id, scum_vehicle_id, game_vehicle_id, vehicle_class, x, y, z, observed_at, created_at` func (repository *mysqlSCUMVehicleTrajectoryRepository) Create(value domain.SCUMVehicleTrajectory) error { if err := repository.pruner.prune(repository.db, "scum_vehicle_trajectory", "observed_at"); err != nil { return err } ctx, cancel := scumQueryContext() defer cancel() _, err := repository.db.ExecContext(ctx, `INSERT INTO scum_vehicle_trajectory (`+scumVehicleTrajectoryColumns+`) VALUES (?,?,?,?,?,?,?,?,?,?)`, value.ID, value.ServerInstanceID, value.SCUMVehicleID, value.GameVehicleID, value.VehicleClass, value.X, value.Y, value.Z, value.ObservedAt.UTC(), value.CreatedAt.UTC()) if err != nil { if isMySQLDuplicateKey(err) { return ErrDuplicate } return fmt.Errorf("create scum_vehicle_trajectory: %w", err) } return nil } func (repository *mysqlSCUMVehicleTrajectoryRepository) Get(id string) (domain.SCUMVehicleTrajectory, error) { ctx, cancel := scumQueryContext() defer cancel() value, err := scanSCUMVehicleTrajectory(repository.db.QueryRowContext(ctx, `SELECT `+scumVehicleTrajectoryColumns+` FROM scum_vehicle_trajectory WHERE id = ?`, id)) if errors.Is(err, sql.ErrNoRows) { return domain.SCUMVehicleTrajectory{}, ErrNotFound } if err != nil { return domain.SCUMVehicleTrajectory{}, fmt.Errorf("read scum_vehicle_trajectory: %w", err) } return value, nil } func (repository *mysqlSCUMVehicleTrajectoryRepository) List(filter domain.SCUMVehicleTrajectoryFilter) ([]domain.SCUMVehicleTrajectory, error) { query := `SELECT ` + scumVehicleTrajectoryColumns + ` FROM scum_vehicle_trajectory WHERE server_instance_id = ?` args := []any{filter.ServerInstanceID} if filter.SCUMVehicleID != "" { query += " AND scum_vehicle_id = ?" args = append(args, filter.SCUMVehicleID) } if filter.GameVehicleID != "" { query += " AND game_vehicle_id = ?" args = append(args, filter.GameVehicleID) } if !filter.After.IsZero() { query += " AND observed_at > ?" args = append(args, filter.After.UTC()) } query += " ORDER BY observed_at DESC, id DESC LIMIT ?" args = append(args, scumListLimit(filter.Limit)) return querySCUMVehicleTrajectories(repository.db, query, args) } func (repository *mysqlSCUMVehicleTrajectoryRepository) Update(value domain.SCUMVehicleTrajectory) error { ctx, cancel := scumQueryContext() defer cancel() result, err := repository.db.ExecContext(ctx, `UPDATE scum_vehicle_trajectory SET server_instance_id=?, scum_vehicle_id=?, game_vehicle_id=?, vehicle_class=?, x=?, y=?, z=?, observed_at=?, created_at=? WHERE id=?`, value.ServerInstanceID, value.SCUMVehicleID, value.GameVehicleID, value.VehicleClass, value.X, value.Y, value.Z, value.ObservedAt.UTC(), value.CreatedAt.UTC(), value.ID) if err != nil { return fmt.Errorf("update scum_vehicle_trajectory: %w", err) } return scumUpdateResult(result, func() error { _, err := repository.Get(value.ID) return err }) } func (repository *mysqlSCUMVehicleTrajectoryRepository) Delete(id string) error { ctx, cancel := scumQueryContext() defer cancel() result, err := repository.db.ExecContext(ctx, `DELETE FROM scum_vehicle_trajectory WHERE id = ?`, id) if err != nil { return fmt.Errorf("delete scum_vehicle_trajectory: %w", err) } return scumDeleteResult(result) } type mysqlSCUMVehicleLockRepository struct{ db *sql.DB } const scumVehicleLockColumns = `id, server_instance_id, scum_vehicle_id, game_vehicle_id, scum_user_id, steam_id, locked_at, created_at` func (repository *mysqlSCUMVehicleLockRepository) Create(value domain.SCUMVehicleLock) error { ctx, cancel := scumQueryContext() defer cancel() _, err := repository.db.ExecContext(ctx, `INSERT INTO scum_vehicle_lock (`+scumVehicleLockColumns+`) VALUES (?,?,?,?,?,?,?,?)`, value.ID, value.ServerInstanceID, value.SCUMVehicleID, value.GameVehicleID, value.SCUMUserID, value.SteamID, value.LockedAt.UTC(), value.CreatedAt.UTC()) if err != nil { if isMySQLDuplicateKey(err) { return ErrDuplicate } return fmt.Errorf("create scum_vehicle_lock: %w", err) } return nil } func (repository *mysqlSCUMVehicleLockRepository) Get(id string) (domain.SCUMVehicleLock, error) { ctx, cancel := scumQueryContext() defer cancel() value, err := scanSCUMVehicleLock(repository.db.QueryRowContext(ctx, `SELECT `+scumVehicleLockColumns+` FROM scum_vehicle_lock WHERE id = ?`, id)) if errors.Is(err, sql.ErrNoRows) { return domain.SCUMVehicleLock{}, ErrNotFound } if err != nil { return domain.SCUMVehicleLock{}, fmt.Errorf("read scum_vehicle_lock: %w", err) } return value, nil } func (repository *mysqlSCUMVehicleLockRepository) List(filter domain.SCUMVehicleLockFilter) ([]domain.SCUMVehicleLock, error) { query := `SELECT ` + scumVehicleLockColumns + ` FROM scum_vehicle_lock WHERE server_instance_id = ?` args := []any{filter.ServerInstanceID} if filter.SCUMVehicleID != "" { query += " AND scum_vehicle_id = ?" args = append(args, filter.SCUMVehicleID) } if filter.GameVehicleID != "" { query += " AND game_vehicle_id = ?" args = append(args, filter.GameVehicleID) } if filter.SCUMUserID != "" { query += " AND scum_user_id = ?" args = append(args, filter.SCUMUserID) } if filter.SteamID != "" { query += " AND steam_id = ?" args = append(args, filter.SteamID) } if !filter.After.IsZero() { query += " AND locked_at > ?" args = append(args, filter.After.UTC()) } query += " ORDER BY locked_at DESC, id DESC LIMIT ?" args = append(args, scumListLimit(filter.Limit)) ctx, cancel := scumQueryContext() defer cancel() rows, err := repository.db.QueryContext(ctx, query, args...) if err != nil { return nil, fmt.Errorf("list scum_vehicle_lock: %w", err) } defer rows.Close() values := make([]domain.SCUMVehicleLock, 0, 16) for rows.Next() { value, err := scanSCUMVehicleLock(rows) if err != nil { return nil, fmt.Errorf("scan scum_vehicle_lock: %w", err) } values = append(values, value) } if err := rows.Err(); err != nil { return nil, fmt.Errorf("list scum_vehicle_lock: %w", err) } return values, nil } func (repository *mysqlSCUMVehicleLockRepository) Update(value domain.SCUMVehicleLock) error { ctx, cancel := scumQueryContext() defer cancel() result, err := repository.db.ExecContext(ctx, `UPDATE scum_vehicle_lock SET server_instance_id=?, scum_vehicle_id=?, game_vehicle_id=?, scum_user_id=?, steam_id=?, locked_at=?, created_at=? WHERE id=?`, value.ServerInstanceID, value.SCUMVehicleID, value.GameVehicleID, value.SCUMUserID, value.SteamID, value.LockedAt.UTC(), value.CreatedAt.UTC(), value.ID) if err != nil { return fmt.Errorf("update scum_vehicle_lock: %w", err) } return scumUpdateResult(result, func() error { _, err := repository.Get(value.ID) return err }) } func (repository *mysqlSCUMVehicleLockRepository) Delete(id string) error { ctx, cancel := scumQueryContext() defer cancel() result, err := repository.db.ExecContext(ctx, `DELETE FROM scum_vehicle_lock WHERE id = ?`, id) if err != nil { return fmt.Errorf("delete scum_vehicle_lock: %w", err) } return scumDeleteResult(result) } type scumTrajectoryPruner struct { mu sync.Mutex lastPrune time.Time } // prune drops trajectory samples past the retention window. It runs at most // once per hour from the trajectory write path so the delete cannot dominate // ingest work. func (pruner *scumTrajectoryPruner) prune(db *sql.DB, table string, timeColumn string) error { pruner.mu.Lock() defer pruner.mu.Unlock() now := time.Now().UTC() if !pruner.lastPrune.IsZero() && now.Sub(pruner.lastPrune) < scumTrajectoryPruneEvery { return nil } ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) defer cancel() if _, err := db.ExecContext(ctx, "DELETE FROM "+table+" WHERE "+timeColumn+" < ?", now.Add(-scumTrajectoryRetention)); err != nil { return fmt.Errorf("prune %s: %w", table, err) } pruner.lastPrune = now return nil } // mysqlDSNWithParseTime keeps DATETIME columns readable as time.Time even when // the operator-provided DSN omits the driver flag. func mysqlDSNWithParseTime(dsn string) string { if strings.Contains(dsn, "parseTime=") { return dsn } separator := "?" if strings.Contains(dsn, "?") { separator = "&" } return dsn + separator + "parseTime=true" } func scumQueryContext() (context.Context, context.CancelFunc) { return context.WithTimeout(context.Background(), 10*time.Second) } func scumListLimit(requested int) int { if requested <= 0 { return domain.SCUMDefaultListLimit } if requested > scumMaxListLimit { return scumMaxListLimit } return requested } func isMySQLDuplicateKey(err error) bool { var mysqlError *mysql.MySQLError return errors.As(err, &mysqlError) && mysqlError.Number == 1062 } func scumUpdateResult(result sql.Result, confirm func() error) error { affected, err := result.RowsAffected() if err != nil { return err } if affected == 0 { return confirm() } return nil } func scumDeleteResult(result sql.Result) error { affected, err := result.RowsAffected() if err != nil { return err } if affected == 0 { return ErrNotFound } return nil } func scumTimeValue(value time.Time) any { if value.IsZero() { return nil } return value.UTC() } func scumTimeFromNull(value sql.NullTime) time.Time { if !value.Valid { return time.Time{} } return value.Time.UTC() } func scumFloatValue(value *float64) any { if value == nil { return nil } return *value } func scumFloatFromNull(value sql.NullFloat64) *float64 { if !value.Valid { return nil } copied := value.Float64 return &copied } func scumIntValue(value *int64) any { if value == nil { return nil } return *value } func scumIntFromNull(value sql.NullInt64) *int64 { if !value.Valid { return nil } copied := value.Int64 return &copied } func scumUserArgs(value domain.SCUMUser) []any { return []any{ value.ID, value.ServerInstanceID, value.SteamID, value.DisplayName, value.SquadID, scumTimeValue(value.RegisteredAt), scumTimeValue(value.LastLoginAt), scumTimeValue(value.LastActivityAt), value.LastLoginIP, value.Online, scumFloatValue(value.X), scumFloatValue(value.Y), scumFloatValue(value.Z), scumIntValue(value.BankBalance), scumIntValue(value.GoldBars), value.RiddenVehicleID, value.GameVehicleID, scumTimeValue(value.WelcomeQueuedAt), value.WelcomeLastError, value.CreatedAt.UTC(), value.UpdatedAt.UTC(), } } func scumUserUpdateArgs(value domain.SCUMUser) []any { args := scumUserArgs(value) return append(args[1:], args[0]) } func scumVehicleArgs(value domain.SCUMVehicle) []any { return []any{ value.ID, value.ServerInstanceID, value.GameVehicleID, value.VehicleClass, value.DisplayName, value.Exists, value.Locked, value.Functional, scumFloatValue(value.X), scumFloatValue(value.Y), scumFloatValue(value.Z), scumTimeValue(value.LastObservedAt), value.CreatedAt.UTC(), value.UpdatedAt.UTC(), } } func scumVehicleUpdateArgs(value domain.SCUMVehicle) []any { args := scumVehicleArgs(value) return append(args[1:], args[0]) } type scumRowScanner interface { Scan(dest ...any) error } func scanSCUMUser(scanner scumRowScanner) (domain.SCUMUser, error) { var ( value domain.SCUMUser registeredAt, lastLoginAt, lastActivityAt sql.NullTime welcomeQueuedAt sql.NullTime x, y, z sql.NullFloat64 bankBalance, goldBars sql.NullInt64 ) err := scanner.Scan( &value.ID, &value.ServerInstanceID, &value.SteamID, &value.DisplayName, &value.SquadID, ®isteredAt, &lastLoginAt, &lastActivityAt, &value.LastLoginIP, &value.Online, &x, &y, &z, &bankBalance, &goldBars, &value.RiddenVehicleID, &value.GameVehicleID, &welcomeQueuedAt, &value.WelcomeLastError, &value.CreatedAt, &value.UpdatedAt, ) if err != nil { return domain.SCUMUser{}, err } value.RegisteredAt = scumTimeFromNull(registeredAt) value.LastLoginAt = scumTimeFromNull(lastLoginAt) value.LastActivityAt = scumTimeFromNull(lastActivityAt) value.WelcomeQueuedAt = scumTimeFromNull(welcomeQueuedAt) value.X = scumFloatFromNull(x) value.Y = scumFloatFromNull(y) value.Z = scumFloatFromNull(z) value.BankBalance = scumIntFromNull(bankBalance) value.GoldBars = scumIntFromNull(goldBars) value.CreatedAt = value.CreatedAt.UTC() value.UpdatedAt = value.UpdatedAt.UTC() return value, nil } func scanSCUMVehicle(scanner scumRowScanner) (domain.SCUMVehicle, error) { var ( value domain.SCUMVehicle x, y, z sql.NullFloat64 lastObservedAt sql.NullTime ) err := scanner.Scan( &value.ID, &value.ServerInstanceID, &value.GameVehicleID, &value.VehicleClass, &value.DisplayName, &value.Exists, &value.Locked, &value.Functional, &x, &y, &z, &lastObservedAt, &value.CreatedAt, &value.UpdatedAt, ) if err != nil { return domain.SCUMVehicle{}, err } value.X = scumFloatFromNull(x) value.Y = scumFloatFromNull(y) value.Z = scumFloatFromNull(z) value.LastObservedAt = scumTimeFromNull(lastObservedAt) value.CreatedAt = value.CreatedAt.UTC() value.UpdatedAt = value.UpdatedAt.UTC() return value, nil } func scanSCUMUserTrajectory(scanner scumRowScanner) (domain.SCUMUserTrajectory, error) { var value domain.SCUMUserTrajectory err := scanner.Scan( &value.ID, &value.ServerInstanceID, &value.SCUMUserID, &value.SteamID, &value.DisplayName, &value.X, &value.Y, &value.Z, &value.RiddenVehicleID, &value.GameVehicleID, &value.ObservedAt, &value.CreatedAt, ) if err != nil { return domain.SCUMUserTrajectory{}, err } value.ObservedAt = value.ObservedAt.UTC() value.CreatedAt = value.CreatedAt.UTC() return value, nil } func scanSCUMVehicleTrajectory(scanner scumRowScanner) (domain.SCUMVehicleTrajectory, error) { var value domain.SCUMVehicleTrajectory err := scanner.Scan( &value.ID, &value.ServerInstanceID, &value.SCUMVehicleID, &value.GameVehicleID, &value.VehicleClass, &value.X, &value.Y, &value.Z, &value.ObservedAt, &value.CreatedAt, ) if err != nil { return domain.SCUMVehicleTrajectory{}, err } value.ObservedAt = value.ObservedAt.UTC() value.CreatedAt = value.CreatedAt.UTC() return value, nil } func scanSCUMVehicleLock(scanner scumRowScanner) (domain.SCUMVehicleLock, error) { var value domain.SCUMVehicleLock err := scanner.Scan( &value.ID, &value.ServerInstanceID, &value.SCUMVehicleID, &value.GameVehicleID, &value.SCUMUserID, &value.SteamID, &value.LockedAt, &value.CreatedAt, ) if err != nil { return domain.SCUMVehicleLock{}, err } value.LockedAt = value.LockedAt.UTC() value.CreatedAt = value.CreatedAt.UTC() return value, nil } func querySCUMUserTrajectories(db *sql.DB, query string, args []any) ([]domain.SCUMUserTrajectory, error) { ctx, cancel := scumQueryContext() defer cancel() rows, err := db.QueryContext(ctx, query, args...) if err != nil { return nil, fmt.Errorf("list scum_user_trajectory: %w", err) } defer rows.Close() values := make([]domain.SCUMUserTrajectory, 0, 16) for rows.Next() { value, err := scanSCUMUserTrajectory(rows) if err != nil { return nil, fmt.Errorf("scan scum_user_trajectory: %w", err) } values = append(values, value) } if err := rows.Err(); err != nil { return nil, fmt.Errorf("list scum_user_trajectory: %w", err) } return values, nil } func querySCUMVehicleTrajectories(db *sql.DB, query string, args []any) ([]domain.SCUMVehicleTrajectory, error) { ctx, cancel := scumQueryContext() defer cancel() rows, err := db.QueryContext(ctx, query, args...) if err != nil { return nil, fmt.Errorf("list scum_vehicle_trajectory: %w", err) } defer rows.Close() values := make([]domain.SCUMVehicleTrajectory, 0, 16) for rows.Next() { value, err := scanSCUMVehicleTrajectory(rows) if err != nil { return nil, fmt.Errorf("scan scum_vehicle_trajectory: %w", err) } values = append(values, value) } if err := rows.Err(); err != nil { return nil, fmt.Errorf("list scum_vehicle_trajectory: %w", err) } return values, nil }