Files
browser/platform/repo/mysql_scum.go
T
npc0-hue a822e9250a Store SCUM facts in typed platform tables
Run pushes typed SCUM facts through POST /api/v1/run/scum/facts, but the
production router rejected that path before the signature middleware because
runServiceRequest listed individual Run channel prefixes, and MySQL kept SCUM
rows inside the whole metadata snapshot instead of platform tables.

- /api/v1/run/ is now the signed machine channel space while
  /api/v1/run/endpoints keeps normal bearer/admin authorization.
- MySQL gets real scum_user, scum_user_trajectory, scum_vehicle,
  scum_vehicle_trajectory and scum_vehicle_lock tables with parameterized
  per-row repositories instead of full snapshot rewrites. Snapshot-shaped
  tables from the unreleased interim build are replaced, and SCUM rows still
  inside a metadata snapshot are migrated once.
- Facts ingest verifies the target server plugin type and converges stale
  online users to offline after SCUMUserOfflineAfter.
- The plugin page and browser read one bounded /scum/surface response instead
  of five list calls per refresh.

Call-count budget for one server: per 5s facts batch, one SELECT plus one
INSERT/UPDATE per reported user and vehicle, one INSERT per moved trajectory
sample or new lock row, one bounded stale-user SELECT, and a trajectory
retention DELETE at most once per hour. One browser refresh issues one surface
request every 15s instead of five list requests.
2026-09-15 12:12:11 +08:00

909 lines
32 KiB
Go

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,
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)
}
}
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, 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=?, 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, 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,
&registeredAt, &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, &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
}