From a822e9250a776deea913fcf1c178370ded180370 Mon Sep 17 00:00:00 2001 From: npc0-hue Date: Tue, 15 Sep 2026 12:12:11 +0800 Subject: [PATCH] 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. --- platform/api/authorization.go | 15 +- platform/api/authorization_test.go | 44 + platform/api/resource_handlers.go | 1 + platform/api/routes.md | 2 +- platform/api/scum_handlers.go | 13 + platform/domain/scum.go | 13 +- platform/dto/scum.go | 17 + platform/repo/mysql_scum.go | 908 ++++++++++++++++++ platform/repo/mysql_store.go | 94 +- platform/repo/resources.go | 17 +- platform/service/resources.go | 1 + platform/service/scum.go | 111 +++ platform/service/scum_test.go | 38 + platform_web/api/client.ts | 5 + platform_web/api/types.ts | 1 + platform_web/contracts/pluginPageHost.ts | 3 +- platform_web/pages/PluginPageHostPage.tsx | 1 + .../scum-server-plugin/features/page-data.ts | 17 +- .../scum-server-plugin/features/page.ts | 35 +- plugins/tests/scum-feature-module.test.ts | 24 +- 20 files changed, 1288 insertions(+), 72 deletions(-) create mode 100644 platform/repo/mysql_scum.go diff --git a/platform/api/authorization.go b/platform/api/authorization.go index bf14b17..2f057c9 100644 --- a/platform/api/authorization.go +++ b/platform/api/authorization.go @@ -40,12 +40,15 @@ func publicAPIRequest(r *http.Request) bool { } func runServiceRequest(r *http.Request) bool { - return strings.HasPrefix(r.URL.Path, "/api/v1/run/control/") || - strings.HasPrefix(r.URL.Path, "/api/v1/run/lifecycle/") || - strings.HasPrefix(r.URL.Path, "/api/v1/run/jobs/") || - strings.HasPrefix(r.URL.Path, "/api/v1/run/logs/") || - strings.HasPrefix(r.URL.Path, "/api/v1/run/artifacts/") || - strings.HasPrefix(r.URL.Path, "/api/v1/run/metrics/") + // Run owns the authenticated machine channels under /run/. Keep the + // browser's endpoint-management API on the normal bearer/admin path, but do + // not maintain a list of individual Run channel prefixes here. New typed + // plugin channels (for example /run/scum/facts) must reach their own + // signature/session middleware instead of being rejected by this router. + path := strings.TrimSuffix(r.URL.Path, "/") + return strings.HasPrefix(path, "/api/v1/run/") && + path != "/api/v1/run/endpoints" && + !strings.HasPrefix(path, "/api/v1/run/endpoints/") } func platformAdminRequest(r *http.Request) bool { diff --git a/platform/api/authorization_test.go b/platform/api/authorization_test.go index 6975574..d9bea94 100644 --- a/platform/api/authorization_test.go +++ b/platform/api/authorization_test.go @@ -263,6 +263,50 @@ func TestAuthorizedRouterAllowsRunLifecycleReportWithoutBearer(t *testing.T) { } } +func TestAuthorizedRouterAcceptsSignedRunSCUMFactsWithoutBearer(t *testing.T) { + store := repo.NewMemoryStore() + core := service.NewCoreService(store) + if _, err := core.CreateGamePlugin(validGamePluginRequest().ToDomain()); err != nil { + t.Fatalf("create plugin: %v", err) + } + if _, err := core.CreateRunEndpoint(validRunEndpointRequest().ToDomain()); err != nil { + t.Fatalf("create endpoint: %v", err) + } + if _, err := core.CreateServerInstance(domain.ServerInstance{ID: "server-scum-facts", PluginID: "server.scum", RunEndpointID: "run-local", Name: "SCUM Facts", State: domain.ServerInstanceStateReady, ConfigVersion: 1}); err != nil { + t.Fatalf("create server: %v", err) + } + router := NewAuthorizedRouterWithCore(core) + // The browser endpoint-management API under /run/ keeps the normal bearer path. + assertErrorResponse(t, performRaw(t, router, http.MethodGet, "/api/v1/run/endpoints", ""), http.StatusUnauthorized, errorCodeUnauthorized) + + hello := decodeBody[dto.RunControlHelloResponse](t, performRunControlHello(t, router, validRunControlHelloRequest())) + observedAt := time.Now().UTC() + body, err := json.Marshal(dto.SCUMFactIngestRequest{ + RunEndpointID: "run-local", SessionToken: hello.SessionToken, ServerInstanceID: "server-scum-facts", + Users: []dto.SCUMUserFactBody{{SteamID: "76561198000000009", DisplayName: "Signed Run", Online: true, Login: true, ObservedAt: observedAt, LoginObservedAt: observedAt, Position: &dto.SCUMPositionBody{X: 1, Y: 2, Z: 3}}}, + }) + if err != nil { + t.Fatalf("marshal scum facts: %v", err) + } + signed := signedRunRequest(t, router, "/api/v1/run/scum/facts", body, hello.SessionToken, "nonce-api-scum-facts", time.Now().UTC()) + assertStatus(t, signed, http.StatusAccepted) + + users, err := store.SCUMUsers().List(domain.SCUMUserFilter{ServerInstanceID: "server-scum-facts"}) + if err != nil { + t.Fatalf("list scum users: %v", err) + } + if len(users) != 1 || users[0].SteamID != "76561198000000009" || !users[0].Online { + t.Fatalf("signed Run SCUM facts did not reach the platform tables: %+v", users) + } + tracks, err := store.SCUMUserTrajectories().List(domain.SCUMUserTrajectoryFilter{ServerInstanceID: "server-scum-facts"}) + if err != nil { + t.Fatalf("list scum trajectories: %v", err) + } + if len(tracks) != 1 { + t.Fatalf("expected one platform trajectory for the reported position, got %+v", tracks) + } +} + func signedRunRequest(t *testing.T, router http.Handler, path string, body []byte, token string, nonce string, stamp time.Time) *httptest.ResponseRecorder { t.Helper() timestamp := strconv.FormatInt(stamp.Unix(), 10) diff --git a/platform/api/resource_handlers.go b/platform/api/resource_handlers.go index b2f37fa..b6e7de6 100644 --- a/platform/api/resource_handlers.go +++ b/platform/api/resource_handlers.go @@ -84,6 +84,7 @@ func (h *coreHandlers) register(mux *http.ServeMux) { mux.HandleFunc("/api/v1/server-instances/{id}/scum/vehicles", h.serverSCUMVehicles) mux.HandleFunc("/api/v1/server-instances/{id}/scum/vehicle-trajectories", h.serverSCUMVehicleTrajectories) mux.HandleFunc("/api/v1/server-instances/{id}/scum/vehicle-locks", h.serverSCUMVehicleLocks) + mux.HandleFunc("/api/v1/server-instances/{id}/scum/surface", h.serverSCUMSurface) mux.HandleFunc("/api/v1/server-instances/{id}/plugin-data/{collection}", h.serverPluginDataCollection) mux.HandleFunc("/api/v1/server-instances/{id}/plugin-data/{collection}/transaction", h.serverPluginDataTransaction) mux.HandleFunc("/api/v1/server-instances/{id}/dependencies/check", h.serverDependenciesCheck) diff --git a/platform/api/routes.md b/platform/api/routes.md index 9dc2245..0ba352e 100644 --- a/platform/api/routes.md +++ b/platform/api/routes.md @@ -154,7 +154,7 @@ Server-scoped terminal log streaming (`GET /api/v1/server-instances/{id}/logs/ev Runtime distribution APIs require the current bearer session, server visibility, plugin-declared permissions, complete runtime bindings only for actions that truly depend on external logical bindings, and platform-builder readiness. Run-side lifecycle commands separately require run endpoint capability support and use plugin-declared lifecycle actions without making manual runtime-profile binding a user prerequisite. Responses and summaries expose artifact IDs, job IDs, checksums, key generations, fingerprints, status, and redacted `secret://runtime-keys/.../current` refs only. They do not expose raw run keys, FTP passwords, database DSNs, RCON passwords, host paths, direct sockets, run endpoint private addresses, build workspace paths, or large inline logs. -SCUM product APIs expose platform-owned `scum_user`, trajectory, vehicle, and lock table records plus typed operation/workflow requests, approval status, confirmation status, blocker reasons, and bounded summaries. They never expose direct game database SQL text, DB paths, DSNs, RCON command text, raw request payloads, run sockets, host paths, or credentials. +SCUM product APIs expose platform-owned `scum_user`, trajectory, vehicle, and lock table records plus typed operation/workflow requests, approval status, confirmation status, blocker reasons, and bounded summaries. The browser reads one bounded server surface (`GET /api/v1/server-instances/{id}/scum/surface`) instead of issuing one list call per table. Signed Run facts arrive through `POST /api/v1/run/scum/facts`; browser Run endpoint management under `/api/v1/run/endpoints` keeps the normal bearer/admin authorization while the remaining `/api/v1/run/` paths belong to the signed machine channels. These APIs never expose direct game database SQL text, DB paths, DSNs, RCON command text, raw request payloads, run sockets, host paths, or credentials. `POST /api/v1/server-instances/workflows/create` requires only the plugin type and server name. A runtime binding may still be maintained internally for advanced logical transports, but browser lifecycle controls must not force operators to choose a runtime profile before start/stop or run-package generation when the plugin deployment/lifecycle declaration is sufficient. Platform builds distributions itself and never needs a registered Run endpoint with `distribution.build` to do so. diff --git a/platform/api/scum_handlers.go b/platform/api/scum_handlers.go index 073944f..14eaf56 100644 --- a/platform/api/scum_handlers.go +++ b/platform/api/scum_handlers.go @@ -94,6 +94,19 @@ func (h *coreHandlers) serverSCUMVehicleLocks(w http.ResponseWriter, r *http.Req writeJSON(w, http.StatusOK, dto.SCUMVehicleLockListFromDomain(items)) } +func (h *coreHandlers) serverSCUMSurface(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + writeMethodNotAllowed(w, http.MethodGet) + return + } + value, err := h.core.GetSCUMSurfaceForSession(bearerToken(r), r.PathValue("id")) + if err != nil { + writeServiceError(w, err) + return + } + writeJSON(w, http.StatusOK, dto.SCUMSurfaceFromDomain(value)) +} + func (h *coreHandlers) runSCUMFacts(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { writeMethodNotAllowed(w, http.MethodPost) diff --git a/platform/domain/scum.go b/platform/domain/scum.go index d4fcf5d..537f4d3 100644 --- a/platform/domain/scum.go +++ b/platform/domain/scum.go @@ -45,7 +45,10 @@ type SCUMUserFilter struct { SteamID string Online *bool ChangedAfter time.Time - Limit int + // StaleBefore matches users whose latest activity is older than the given + // instant. It backs the offline convergence query for stale online users. + StaleBefore time.Time + Limit int } type SCUMUserTrajectory struct { @@ -137,6 +140,14 @@ type SCUMVehicleLockFilter struct { Limit int } +type SCUMSurface struct { + Users []SCUMUser + Vehicles []SCUMVehicle + UserTrajectories []SCUMUserTrajectory + VehicleTrajectories []SCUMVehicleTrajectory + VehicleLocks []SCUMVehicleLock +} + type SCUMFactIngest struct { RunEndpointID string SessionToken string diff --git a/platform/dto/scum.go b/platform/dto/scum.go index 7f4581d..79e9a43 100644 --- a/platform/dto/scum.go +++ b/platform/dto/scum.go @@ -141,6 +141,23 @@ type SCUMVehicleLockListResponse struct { Count int `json:"count"` } +type SCUMSurfaceResponse struct { + Users SCUMUserListResponse `json:"users"` + Vehicles SCUMVehicleListResponse `json:"vehicles"` + UserTrajectories SCUMUserTrajectoryListResponse `json:"userTrajectories"` + VehicleTrajectories SCUMVehicleTrajectoryListResponse `json:"vehicleTrajectories"` + VehicleLocks SCUMVehicleLockListResponse `json:"vehicleLocks"` +} + +func SCUMSurfaceFromDomain(value domain.SCUMSurface) SCUMSurfaceResponse { + return SCUMSurfaceResponse{ + Users: SCUMUserListFromDomain(value.Users), Vehicles: SCUMVehicleListFromDomain(value.Vehicles), + UserTrajectories: SCUMUserTrajectoryListFromDomain(value.UserTrajectories), + VehicleTrajectories: SCUMVehicleTrajectoryListFromDomain(value.VehicleTrajectories), + VehicleLocks: SCUMVehicleLockListFromDomain(value.VehicleLocks), + } +} + type SCUMFactIngestRequest struct { RunEndpointID string `json:"runEndpointId"` SessionToken string `json:"sessionToken"` diff --git a/platform/repo/mysql_scum.go b/platform/repo/mysql_scum.go new file mode 100644 index 0000000..743bebf --- /dev/null +++ b/platform/repo/mysql_scum.go @@ -0,0 +1,908 @@ +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, + ®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, &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 +} diff --git a/platform/repo/mysql_store.go b/platform/repo/mysql_store.go index 6570844..adfc52a 100644 --- a/platform/repo/mysql_store.go +++ b/platform/repo/mysql_store.go @@ -4,6 +4,7 @@ import ( "context" "database/sql" "encoding/json" + "errors" "fmt" "strings" "sync" @@ -21,9 +22,11 @@ const ( type MySQLStore struct { *MemoryStore - db *sql.DB - persistMu sync.Mutex - runtimePersistMu sync.Mutex + db *sql.DB + persistMu sync.Mutex + runtimePersistMu sync.Mutex + scumUserTrackPruner scumTrajectoryPruner + scumVehicleTrackPruner scumTrajectoryPruner } func NewMySQLStore(dsn string) (*MySQLStore, error) { @@ -31,7 +34,7 @@ func NewMySQLStore(dsn string) (*MySQLStore, error) { if dsn == "" { return nil, fmt.Errorf("PLATFORM_MYSQL_DSN is required when PLATFORM_STORAGE_BACKEND=mysql") } - db, err := sql.Open("mysql", dsn) + db, err := sql.Open("mysql", mysqlDSNWithParseTime(dsn)) if err != nil { return nil, fmt.Errorf("open mysql metadata store: %w", err) } @@ -150,21 +153,6 @@ func (store *MySQLStore) GameClientBridgeSnapshotStreams() GameClientBridgeSnaps func (store *MySQLStore) PluginDataRecords() PluginDataRecordRepository { return &persistentRepository[domain.PluginDataRecord, domain.PluginDataFilter]{repository: store.MemoryStore.pluginDataRecords, persist: store.persist} } -func (store *MySQLStore) SCUMUsers() SCUMUserRepository { - return &persistentRepository[domain.SCUMUser, domain.SCUMUserFilter]{repository: store.MemoryStore.scumUsers, persist: store.persist} -} -func (store *MySQLStore) SCUMUserTrajectories() SCUMUserTrajectoryRepository { - return &persistentRepository[domain.SCUMUserTrajectory, domain.SCUMUserTrajectoryFilter]{repository: store.MemoryStore.scumUserTracks, persist: store.persist} -} -func (store *MySQLStore) SCUMVehicles() SCUMVehicleRepository { - return &persistentRepository[domain.SCUMVehicle, domain.SCUMVehicleFilter]{repository: store.MemoryStore.scumVehicles, persist: store.persist} -} -func (store *MySQLStore) SCUMVehicleTrajectories() SCUMVehicleTrajectoryRepository { - return &persistentRepository[domain.SCUMVehicleTrajectory, domain.SCUMVehicleTrajectoryFilter]{repository: store.MemoryStore.scumVehicleTracks, persist: store.persist} -} -func (store *MySQLStore) SCUMVehicleLocks() SCUMVehicleLockRepository { - return &persistentRepository[domain.SCUMVehicleLock, domain.SCUMVehicleLockFilter]{repository: store.MemoryStore.scumVehicleLocks, persist: store.persist} -} func (store *MySQLStore) initialize() error { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() @@ -180,6 +168,9 @@ CREATE TABLE IF NOT EXISTS platform_metadata_snapshots ( if err != nil { return fmt.Errorf("create mysql metadata snapshot table: %w", err) } + if err := store.initializeSCUMTables(ctx); err != nil { + return err + } return nil } @@ -199,9 +190,63 @@ func (store *MySQLStore) load() error { return fmt.Errorf("decode mysql metadata snapshot: %w", err) } store.loadSnapshot(snapshot) + if err := store.migrateSnapshotSCUM(ctx, snapshot); err != nil { + return err + } return store.loadRuntime() } +// migrateSnapshotSCUM copies SCUM rows that an older build kept inside the +// metadata snapshot into the typed SCUM tables. SCUM is no longer stored in +// the snapshot, so this runs at most once for a snapshot written before the +// typed tables existed. +func (store *MySQLStore) migrateSnapshotSCUM(ctx context.Context, snapshot StoreSnapshot) error { + if len(snapshot.SCUMUsers)+len(snapshot.SCUMVehicles)+len(snapshot.SCUMUserTrajectories)+len(snapshot.SCUMVehicleTrajectories)+len(snapshot.SCUMVehicleLocks) == 0 { + return nil + } + users := store.SCUMUsers() + for _, value := range snapshot.SCUMUsers { + if _, err := users.Get(value.ID); err == nil { + continue + } else if !errors.Is(err, ErrNotFound) { + return err + } + if err := users.Create(value); err != nil && !errors.Is(err, ErrDuplicate) { + return fmt.Errorf("migrate scum_user %s: %w", value.ID, err) + } + } + vehicles := store.SCUMVehicles() + for _, value := range snapshot.SCUMVehicles { + if _, err := vehicles.Get(value.ID); err == nil { + continue + } else if !errors.Is(err, ErrNotFound) { + return err + } + if err := vehicles.Create(value); err != nil && !errors.Is(err, ErrDuplicate) { + return fmt.Errorf("migrate scum_vehicle %s: %w", value.ID, err) + } + } + userTracks := store.SCUMUserTrajectories() + for _, value := range snapshot.SCUMUserTrajectories { + if err := userTracks.Create(value); err != nil && !errors.Is(err, ErrDuplicate) { + return fmt.Errorf("migrate scum_user_trajectory %s: %w", value.ID, err) + } + } + vehicleTracks := store.SCUMVehicleTrajectories() + for _, value := range snapshot.SCUMVehicleTrajectories { + if err := vehicleTracks.Create(value); err != nil && !errors.Is(err, ErrDuplicate) { + return fmt.Errorf("migrate scum_vehicle_trajectory %s: %w", value.ID, err) + } + } + locks := store.SCUMVehicleLocks() + for _, value := range snapshot.SCUMVehicleLocks { + if err := locks.Create(value); err != nil && !errors.Is(err, ErrDuplicate) { + return fmt.Errorf("migrate scum_vehicle_lock %s: %w", value.ID, err) + } + } + return nil +} + func (store *MySQLStore) loadRuntime() error { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() @@ -288,11 +333,12 @@ func (store *MySQLStore) snapshot() StoreSnapshot { GameClientBridgeSnapshots: snapshotRepository(store.MemoryStore.bridgeSnapshots.memoryRepository), GameClientBridgeStreams: snapshotRepository(store.MemoryStore.bridgeStreams), PluginDataRecords: snapshotRepository(store.MemoryStore.pluginDataRecords), - SCUMUsers: snapshotRepository(store.MemoryStore.scumUsers), - SCUMUserTrajectories: snapshotRepository(store.MemoryStore.scumUserTracks), - SCUMVehicles: snapshotRepository(store.MemoryStore.scumVehicles), - SCUMVehicleTrajectories: snapshotRepository(store.MemoryStore.scumVehicleTracks), - SCUMVehicleLocks: snapshotRepository(store.MemoryStore.scumVehicleLocks), + // SCUM rows live in the typed scum_* tables, never in the metadata snapshot. + SCUMUsers: nil, + SCUMUserTrajectories: nil, + SCUMVehicles: nil, + SCUMVehicleTrajectories: nil, + SCUMVehicleLocks: nil, }) } diff --git a/platform/repo/resources.go b/platform/repo/resources.go index d097e89..ef8daaa 100644 --- a/platform/repo/resources.go +++ b/platform/repo/resources.go @@ -4,6 +4,7 @@ import ( "errors" "sort" "sync" + "time" "browser.local/platform/domain" ) @@ -190,6 +191,7 @@ type SCUMUserRepository interface { Get(string) (domain.SCUMUser, error) List(domain.SCUMUserFilter) ([]domain.SCUMUser, error) Update(domain.SCUMUser) error + Delete(string) error } type SCUMUserTrajectoryRepository interface { @@ -197,6 +199,7 @@ type SCUMUserTrajectoryRepository interface { Get(string) (domain.SCUMUserTrajectory, error) List(domain.SCUMUserTrajectoryFilter) ([]domain.SCUMUserTrajectory, error) Update(domain.SCUMUserTrajectory) error + Delete(string) error } type SCUMVehicleRepository interface { @@ -204,6 +207,7 @@ type SCUMVehicleRepository interface { Get(string) (domain.SCUMVehicle, error) List(domain.SCUMVehicleFilter) ([]domain.SCUMVehicle, error) Update(domain.SCUMVehicle) error + Delete(string) error } type SCUMVehicleTrajectoryRepository interface { @@ -211,6 +215,7 @@ type SCUMVehicleTrajectoryRepository interface { Get(string) (domain.SCUMVehicleTrajectory, error) List(domain.SCUMVehicleTrajectoryFilter) ([]domain.SCUMVehicleTrajectory, error) Update(domain.SCUMVehicleTrajectory) error + Delete(string) error } type SCUMVehicleLockRepository interface { @@ -218,6 +223,7 @@ type SCUMVehicleLockRepository interface { Get(string) (domain.SCUMVehicleLock, error) List(domain.SCUMVehicleLockFilter) ([]domain.SCUMVehicleLock, error) Update(domain.SCUMVehicleLock) error + Delete(string) error } type Store interface { @@ -758,7 +764,16 @@ func matchSCUMUser(value domain.SCUMUser, filter domain.SCUMUserFilter) bool { return (filter.ServerInstanceID == "" || value.ServerInstanceID == filter.ServerInstanceID) && (filter.SteamID == "" || value.SteamID == filter.SteamID) && (filter.Online == nil || value.Online == *filter.Online) && - (filter.ChangedAfter.IsZero() || value.UpdatedAt.After(filter.ChangedAfter) || value.LastActivityAt.After(filter.ChangedAfter)) + (filter.ChangedAfter.IsZero() || value.UpdatedAt.After(filter.ChangedAfter) || value.LastActivityAt.After(filter.ChangedAfter)) && + (filter.StaleBefore.IsZero() || scumUserActivityAt(value).Before(filter.StaleBefore)) +} + +// scumUserActivityAt reports the timestamp used for offline convergence. +func scumUserActivityAt(value domain.SCUMUser) time.Time { + if value.LastActivityAt.IsZero() { + return value.UpdatedAt + } + return value.LastActivityAt } func matchSCUMUserTrajectory(value domain.SCUMUserTrajectory, filter domain.SCUMUserTrajectoryFilter) bool { diff --git a/platform/service/resources.go b/platform/service/resources.go index b330fa2..96f2fd6 100644 --- a/platform/service/resources.go +++ b/platform/service/resources.go @@ -186,6 +186,7 @@ type Core interface { ListSCUMVehiclesForSession(string, domain.SCUMVehicleFilter) ([]domain.SCUMVehicle, error) ListSCUMVehicleTrajectoriesForSession(string, domain.SCUMVehicleTrajectoryFilter) ([]domain.SCUMVehicleTrajectory, error) ListSCUMVehicleLocksForSession(string, domain.SCUMVehicleLockFilter) ([]domain.SCUMVehicleLock, error) + GetSCUMSurfaceForSession(string, string) (domain.SCUMSurface, error) IngestSCUMFacts(domain.SCUMFactIngest) (domain.SCUMFactIngestResult, error) ListPluginDataForSession(string, domain.PluginDataFilter) ([]domain.PluginDataRecord, error) PutPluginDataForSession(string, domain.PluginDataRecord) (domain.PluginDataRecord, error) diff --git a/platform/service/scum.go b/platform/service/scum.go index f56313a..9494a06 100644 --- a/platform/service/scum.go +++ b/platform/service/scum.go @@ -19,6 +19,9 @@ func (svc *CoreService) ListSCUMUsersForSession(sessionID string, filter domain. if _, err := svc.GetServerInstanceForSession(sessionID, filter.ServerInstanceID); err != nil { return nil, err } + if _, err := svc.reconcileSCUMUserPresence(filter.ServerInstanceID, svc.now()); err != nil { + return nil, err + } users, err := svc.store.SCUMUsers().List(filter) if err != nil { return nil, err @@ -138,6 +141,79 @@ func (svc *CoreService) ListSCUMVehicleLocksForSession(sessionID string, filter return domain.CopySCUMVehicleLocks(items), nil } +func (svc *CoreService) GetSCUMSurfaceForSession(sessionID, serverInstanceID string) (domain.SCUMSurface, error) { + if strings.TrimSpace(serverInstanceID) == "" { + return domain.SCUMSurface{}, validationError("serverInstanceId is required") + } + if _, err := svc.GetServerInstanceForSession(sessionID, serverInstanceID); err != nil { + return domain.SCUMSurface{}, err + } + if _, err := svc.reconcileSCUMUserPresence(serverInstanceID, svc.now()); err != nil { + return domain.SCUMSurface{}, err + } + users, err := svc.store.SCUMUsers().List(domain.SCUMUserFilter{ServerInstanceID: serverInstanceID, Limit: domain.SCUMDefaultListLimit}) + if err != nil { + return domain.SCUMSurface{}, err + } + vehicles, err := svc.store.SCUMVehicles().List(domain.SCUMVehicleFilter{ServerInstanceID: serverInstanceID, Limit: domain.SCUMDefaultListLimit}) + if err != nil { + return domain.SCUMSurface{}, err + } + userTracks, err := svc.store.SCUMUserTrajectories().List(domain.SCUMUserTrajectoryFilter{ServerInstanceID: serverInstanceID, Limit: domain.SCUMDefaultTrajectoryLimit}) + if err != nil { + return domain.SCUMSurface{}, err + } + vehicleTracks, err := svc.store.SCUMVehicleTrajectories().List(domain.SCUMVehicleTrajectoryFilter{ServerInstanceID: serverInstanceID, Limit: domain.SCUMDefaultTrajectoryLimit}) + if err != nil { + return domain.SCUMSurface{}, err + } + locks, err := svc.store.SCUMVehicleLocks().List(domain.SCUMVehicleLockFilter{ServerInstanceID: serverInstanceID, Limit: domain.SCUMDefaultTrajectoryLimit}) + if err != nil { + return domain.SCUMSurface{}, err + } + return domain.SCUMSurface{Users: limitSCUMUsers(users), Vehicles: limitSCUMVehicles(vehicles), UserTrajectories: limitSCUMUserTrajectories(userTracks), VehicleTrajectories: limitSCUMVehicleTrajectories(vehicleTracks), VehicleLocks: limitSCUMVehicleLocks(locks)}, nil +} + +func limitSCUMUsers(values []domain.SCUMUser) []domain.SCUMUser { + sort.SliceStable(values, func(i, j int) bool { return values[i].LastActivityAt.After(values[j].LastActivityAt) }) + if len(values) > domain.SCUMDefaultListLimit { + values = values[:domain.SCUMDefaultListLimit] + } + return domain.CopySCUMUsers(values) +} + +func limitSCUMVehicles(values []domain.SCUMVehicle) []domain.SCUMVehicle { + sort.SliceStable(values, func(i, j int) bool { return values[i].LastObservedAt.After(values[j].LastObservedAt) }) + if len(values) > domain.SCUMDefaultListLimit { + values = values[:domain.SCUMDefaultListLimit] + } + return domain.CopySCUMVehicles(values) +} + +func limitSCUMUserTrajectories(values []domain.SCUMUserTrajectory) []domain.SCUMUserTrajectory { + sort.SliceStable(values, func(i, j int) bool { return values[i].ObservedAt.After(values[j].ObservedAt) }) + if len(values) > domain.SCUMDefaultTrajectoryLimit { + values = values[:domain.SCUMDefaultTrajectoryLimit] + } + return domain.CopySCUMUserTrajectories(values) +} + +func limitSCUMVehicleTrajectories(values []domain.SCUMVehicleTrajectory) []domain.SCUMVehicleTrajectory { + sort.SliceStable(values, func(i, j int) bool { return values[i].ObservedAt.After(values[j].ObservedAt) }) + if len(values) > domain.SCUMDefaultTrajectoryLimit { + values = values[:domain.SCUMDefaultTrajectoryLimit] + } + return domain.CopySCUMVehicleTrajectories(values) +} + +func limitSCUMVehicleLocks(values []domain.SCUMVehicleLock) []domain.SCUMVehicleLock { + sort.SliceStable(values, func(i, j int) bool { return values[i].LockedAt.After(values[j].LockedAt) }) + if len(values) > domain.SCUMDefaultTrajectoryLimit { + values = values[:domain.SCUMDefaultTrajectoryLimit] + } + return domain.CopySCUMVehicleLocks(values) +} + func (svc *CoreService) IngestSCUMFacts(batch domain.SCUMFactIngest) (domain.SCUMFactIngestResult, error) { batch = domain.CopySCUMFactIngest(batch) if strings.TrimSpace(batch.ServerInstanceID) == "" { @@ -159,7 +235,17 @@ func (svc *CoreService) IngestSCUMFacts(batch domain.SCUMFactIngest) (domain.SCU if instance.RunEndpointID != batch.RunEndpointID { return domain.SCUMFactIngestResult{}, ErrForbidden } + plugin, err := svc.store.GamePlugins().Get(instance.PluginID) + if err != nil { + return domain.SCUMFactIngestResult{}, err + } + if !strings.EqualFold(strings.TrimSpace(plugin.ServerType), "scum") { + return domain.SCUMFactIngestResult{}, ErrForbidden + } stamp := svc.now() + if _, err := svc.reconcileSCUMUserPresence(instance.ID, stamp); err != nil { + return domain.SCUMFactIngestResult{}, err + } welcomeCount := 0 for _, fact := range batch.Users { if err := svc.ingestSCUMUserFact(instance, fact, stamp, &welcomeCount); err != nil { @@ -174,6 +260,31 @@ func (svc *CoreService) IngestSCUMFacts(batch domain.SCUMFactIngest) (domain.SCU return domain.SCUMFactIngestResult{Accepted: true, AcceptedUserCount: len(batch.Users), AcceptedVehicleCount: len(batch.Vehicles), WelcomeQueuedCount: welcomeCount, ServerTime: stamp}, nil } +// reconcileSCUMUserPresence converges stale online flags from the last typed +// Run fact. A missing offline fact is treated as a presence timeout only; no +// trajectory row is written and the last observed coordinates stay untouched. +func (svc *CoreService) reconcileSCUMUserPresence(serverInstanceID string, stamp time.Time) (int, error) { + online := true + stale, err := svc.store.SCUMUsers().List(domain.SCUMUserFilter{ + ServerInstanceID: serverInstanceID, + Online: &online, + StaleBefore: stamp.Add(-domain.SCUMUserOfflineAfter), + }) + if err != nil { + return 0, err + } + updated := 0 + for _, user := range stale { + user.Online = false + user.UpdatedAt = stamp + if err := svc.store.SCUMUsers().Update(user); err != nil { + return updated, err + } + updated++ + } + return updated, nil +} + func (svc *CoreService) ingestSCUMUserFact(instance domain.ServerInstance, fact domain.SCUMUserFact, stamp time.Time, welcomeCount *int) error { steamID := strings.TrimSpace(fact.SteamID) if steamID == "" { diff --git a/platform/service/scum_test.go b/platform/service/scum_test.go index d98754c..964f98e 100644 --- a/platform/service/scum_test.go +++ b/platform/service/scum_test.go @@ -1,6 +1,7 @@ package service import ( + "errors" "testing" "time" @@ -126,3 +127,40 @@ func TestSCUMFactIngestQueuesFreshWelcomeAfterActivityGap(t *testing.T) { t.Fatalf("expected two welcome jobs after activity gap, jobs=%+v err=%v", jobs, err) } } + +func TestSCUMStaleOnlineUserConvergesToOffline(t *testing.T) { + svc, session, runSession, instance := newSourceRCONFixture(t) + staleAt := fixedTime.Add(-16 * time.Minute) + if _, err := svc.IngestSCUMFacts(domain.SCUMFactIngest{ + RunEndpointID: instance.RunEndpointID, SessionToken: runSession, ServerInstanceID: instance.ID, + Users: []domain.SCUMUserFact{{SteamID: "76561198000000003", DisplayName: "Stale", Online: true, Login: true, ObservedAt: staleAt, LoginObservedAt: staleAt}}, + }); err != nil { + t.Fatalf("ingest stale SCUM user fact: %v", err) + } + + users, err := svc.ListSCUMUsersForSession(session, domain.SCUMUserFilter{ServerInstanceID: instance.ID}) + if err != nil || len(users) != 1 { + t.Fatalf("list stale SCUM users: users=%+v err=%v", users, err) + } + if users[0].Online { + t.Fatalf("expected a user without activity for %s to converge offline: %+v", domain.SCUMUserOfflineAfter, users[0]) + } +} + +func TestSCUMFactsRejectedForNonSCUMServerInstance(t *testing.T) { + svc, _, runSession, instance := newSourceRCONFixture(t) + if err := svc.store.GamePlugins().Create(domain.GamePlugin{ID: "server.other", Name: "Other", Version: "1.0.0", ServerType: "other"}); err != nil { + t.Fatalf("create non-SCUM plugin: %v", err) + } + if err := svc.store.ServerInstances().Create(domain.ServerInstance{ID: "server-other", PluginID: "server.other", RunEndpointID: instance.RunEndpointID, Name: "Other Server", State: domain.ServerInstanceStateRunning, ConfigVersion: 1}); err != nil { + t.Fatalf("create non-SCUM server: %v", err) + } + + _, err := svc.IngestSCUMFacts(domain.SCUMFactIngest{ + RunEndpointID: instance.RunEndpointID, SessionToken: runSession, ServerInstanceID: "server-other", + Users: []domain.SCUMUserFact{{SteamID: "76561198000000004", Online: true}}, + }) + if !errors.Is(err, ErrForbidden) { + t.Fatalf("expected SCUM facts for a non-SCUM server to be rejected, got %v", err) + } +} diff --git a/platform_web/api/client.ts b/platform_web/api/client.ts index 9867a90..c12867c 100644 --- a/platform_web/api/client.ts +++ b/platform_web/api/client.ts @@ -89,6 +89,7 @@ import type { ServerMetricsListResponse, ScumListRequest, ScumUserListResponse, + ScumSurfaceResponse, ScumUserTrajectoryListResponse, ScumVehicleLockListResponse, ScumVehicleListResponse, @@ -537,6 +538,10 @@ export class PlatformApiClient { return this.request(`/server-instances/${encodeURIComponent(serverInstanceId)}/scum/vehicle-locks${scumListQuery(request)}`); } + async getScumSurface(serverInstanceId: string): Promise { + return this.request(`/server-instances/${encodeURIComponent(serverInstanceId)}/scum/surface`); + } + async listPluginData(serverInstanceId: string, collection: string, key?: string): Promise<{ items: Array<{ key: string; value: Record }>; count: number }> { const params = new URLSearchParams(); if (key) params.set("key", key); diff --git a/platform_web/api/types.ts b/platform_web/api/types.ts index 2f63f50..d80824f 100644 --- a/platform_web/api/types.ts +++ b/platform_web/api/types.ts @@ -192,6 +192,7 @@ export interface ScumUserResponse { } export interface ScumUserListResponse { items: ScumUserResponse[]; count: number; } +export interface ScumSurfaceResponse { users: ScumUserListResponse; vehicles: ScumVehicleListResponse; userTrajectories: ScumUserTrajectoryListResponse; vehicleTrajectories: ScumVehicleTrajectoryListResponse; vehicleLocks: ScumVehicleLockListResponse; } export interface ScumUserTrajectoryResponse { id: string; diff --git a/platform_web/contracts/pluginPageHost.ts b/platform_web/contracts/pluginPageHost.ts index 56f595d..c18dad8 100644 --- a/platform_web/contracts/pluginPageHost.ts +++ b/platform_web/contracts/pluginPageHost.ts @@ -1,5 +1,5 @@ import type { PluginBridgeExecuteEnvelope, PluginBridgeExecutionResult } from "./pluginBridge"; -import type { GameClientBridgeCommandFilterRequest, GameClientBridgeQueueRequest, GameClientBridgeSnapshotQuery, LogStreamCursorRequest, LogStreamCursorResponse, LogStreamListResponse, ScumListRequest, ScumUserListResponse, ScumUserTrajectoryListResponse, ScumVehicleLockListResponse, ScumVehicleListResponse, ScumVehicleTrajectoryListResponse } from "../api/types"; +import type { GameClientBridgeCommandFilterRequest, GameClientBridgeQueueRequest, GameClientBridgeSnapshotQuery, LogStreamCursorRequest, LogStreamCursorResponse, LogStreamListResponse, ScumListRequest, ScumSurfaceResponse, ScumUserListResponse, ScumUserTrajectoryListResponse, ScumVehicleLockListResponse, ScumVehicleListResponse, ScumVehicleTrajectoryListResponse } from "../api/types"; export interface PluginDataMutation { operation: "put" | "delete"; @@ -25,6 +25,7 @@ export interface PluginPageWorkspaceActions { query: (request: LogStreamCursorRequest) => Promise; }; scum?: { + surface: () => Promise; users: (request?: ScumListRequest) => Promise; userTrajectories: (request?: ScumListRequest) => Promise; vehicles: (request?: ScumListRequest) => Promise; diff --git a/platform_web/pages/PluginPageHostPage.tsx b/platform_web/pages/PluginPageHostPage.tsx index 7482b20..5d5d291 100644 --- a/platform_web/pages/PluginPageHostPage.tsx +++ b/platform_web/pages/PluginPageHostPage.tsx @@ -94,6 +94,7 @@ export function PluginPageHostPage({ params, onNavigate, initialPlugin, embedded } }, scum: { + surface: () => platformApiClient.getScumSurface(serverId), users: (request) => platformApiClient.listScumUsers(serverId, request), userTrajectories: (request) => platformApiClient.listScumUserTrajectories(serverId, request), vehicles: (request) => platformApiClient.listScumVehicles(serverId, request), diff --git a/plugins/examples/scum-server-plugin/features/page-data.ts b/plugins/examples/scum-server-plugin/features/page-data.ts index 7b12cbf..286d698 100644 --- a/plugins/examples/scum-server-plugin/features/page-data.ts +++ b/plugins/examples/scum-server-plugin/features/page-data.ts @@ -16,11 +16,7 @@ export type PluginBridgeExecuteEnvelope = { requestId: string; action: string; p export type PluginBridgeExecutionResult = { status?: string; result?: Record; error?: { message?: string } }; export type SCUMPlatformActions = { - users: (query?: { limit?: number; changedAfter?: string; online?: boolean; steamId?: string }) => Promise; - vehicles: (query?: { limit?: number; changedAfter?: string; exists?: boolean; gameVehicleId?: string }) => Promise; - userTrajectories: (query?: { limit?: number; after?: string; steamId?: string; scumUserId?: string }) => Promise; - vehicleTrajectories: (query?: { limit?: number; after?: string; gameVehicleId?: string; scumVehicleId?: string }) => Promise; - vehicleLocks: (query?: { limit?: number; after?: string; gameVehicleId?: string; scumVehicleId?: string; scumUserId?: string; steamId?: string }) => Promise; + surface: () => Promise<{ users: unknown; vehicles: unknown; userTrajectories: unknown; vehicleTrajectories: unknown; vehicleLocks: unknown }>; }; export type SCUMWorkspaceActions = { @@ -119,12 +115,11 @@ export async function loadSCUMSurface(actions: SCUMWorkspaceActions, pageKey: st async function loadPlatformSCUMTables(actions: SCUMWorkspaceActions, data: SCUMSurfaceData, keys: SCUMPlatformSurfaceKey[]): Promise { if (!keys.length) return; if (!actions.scum) throw new Error("平台 SCUM 数据能力不可用。"); - const reads: Array> = []; - if (keys.includes("players")) reads.push(actions.scum.users({ limit: 500 }).then((response) => { data.players = collectionRecords(response); })); - if (keys.includes("vehicles")) reads.push(actions.scum.vehicles({ limit: 500 }).then((response) => { data.vehicles = collectionRecords(response); })); - if (keys.includes("vehicleLocks")) reads.push(actions.scum.vehicleLocks({ limit: 500 }).then((response) => { data.vehicleLocks = collectionRecords(response); })); - if (keys.includes("trajectories")) reads.push(Promise.all([actions.scum.userTrajectories({ limit: 500 }), actions.scum.vehicleTrajectories({ limit: 500 })]).then(([players, vehicles]) => { data.trajectories = [...collectionRecords(players), ...collectionRecords(vehicles)]; })); - await Promise.all(reads); + const surface = await actions.scum.surface(); + if (keys.includes("players")) data.players = collectionRecords(surface.users); + if (keys.includes("vehicles")) data.vehicles = collectionRecords(surface.vehicles); + if (keys.includes("vehicleLocks")) data.vehicleLocks = collectionRecords(surface.vehicleLocks); + if (keys.includes("trajectories")) data.trajectories = [...collectionRecords(surface.userTrajectories), ...collectionRecords(surface.vehicleTrajectories)]; } export async function saveGiftDefinition(actions: SCUMWorkspaceActions, gift: RecordMap): Promise { diff --git a/plugins/examples/scum-server-plugin/features/page.ts b/plugins/examples/scum-server-plugin/features/page.ts index f1456e3..160688f 100644 --- a/plugins/examples/scum-server-plugin/features/page.ts +++ b/plugins/examples/scum-server-plugin/features/page.ts @@ -121,24 +121,31 @@ export function renderSCUMFeaturePage(react: ReactLike, input: SCUMPageContext) const [produceEditorOpen, setProduceEditorOpen] = usePluginState(react, false); const [giftEditorOpen, setGiftEditorOpen] = usePluginState(react, false); const [mapSettingsOpen, setMapSettingsOpen] = usePluginState(react, false); + const [refreshSignal, setRefreshSignal] = usePluginState(react, 0); const pageKey = input.pageKey ?? "players"; - - const refresh = () => { - if (!input.serverInstanceId || !input.workspaceActions?.scum || !input.workspaceActions?.pluginData) { - setState({ status: "error", reason: "插件页面没有绑定服务器、平台 SCUM 数据能力或通用 pluginData 能力。" }); - return; - } - void loadSCUMSurface(input.workspaceActions, pageKey) - .then((data) => setState({ status: "ready", data })) - .catch((error) => setState({ status: "error", reason: errorMessage(error, "SCUM 插件数据读取失败。") })); - }; - if (react.useEffect) react.useEffect(() => { + let inFlight = false; + let generation = 0; + let disposed = false; + const refresh = () => { + if (inFlight || disposed) return; + inFlight = true; + const currentGeneration = ++generation; + if (!input.serverInstanceId || !input.workspaceActions?.scum || !input.workspaceActions?.pluginData) { + setState({ status: "error", reason: "插件页面没有绑定服务器、平台 SCUM 数据能力或通用 pluginData 能力。" }); + inFlight = false; + return; + } + void loadSCUMSurface(input.workspaceActions, pageKey) + .then((data) => { if (!disposed && currentGeneration === generation) setState({ status: "ready", data }); }) + .catch((error) => { if (!disposed && currentGeneration === generation) setState({ status: "error", reason: errorMessage(error, "SCUM 插件数据读取失败。") }); }) + .finally(() => { inFlight = false; }); + }; if (playerPanel.kind === "closed") refresh(); if (playerPanel.kind !== "closed") return; const interval = setInterval(refresh, scumSurfaceRefreshMs); - return () => clearInterval(interval); - }, [input.serverInstanceId, pageKey, input.workspaceActions, playerPanel.kind]); + return () => { disposed = true; generation += 1; clearInterval(interval); }; + }, [input.serverInstanceId, pageKey, input.workspaceActions, playerPanel.kind, refreshSignal]); const data = state.status === "ready" ? state.data : emptySCUMSurfaceData; return e("section", { className: "console-panel scum-workbench", "aria-label": input.pageTitle ?? surfaceTitle(pageKey) }, @@ -154,7 +161,7 @@ export function renderSCUMFeaturePage(react: ReactLike, input: SCUMPageContext) deliveryGift, setDeliveryGift, deliveryPlayer, setDeliveryPlayer, mapSearch, setMapSearch, mapLayers, setMapLayers, selectedMapPoint, setSelectedMapPoint, mapCustomEnabled, setMapCustomEnabled, mapCenterX, setMapCenterX, mapCenterY, setMapCenterY, mapWidthKm, setMapWidthKm, mapHeightKm, setMapHeightKm, eventEditorOpen, setEventEditorOpen, produceEditorOpen, setProduceEditorOpen, giftEditorOpen, setGiftEditorOpen, mapSettingsOpen, setMapSettingsOpen, - setAction, refresh + setAction, refresh: () => setRefreshSignal((value) => value + 1) }) : null ); } diff --git a/plugins/tests/scum-feature-module.test.ts b/plugins/tests/scum-feature-module.test.ts index 83f06f7..e56aafc 100644 --- a/plugins/tests/scum-feature-module.test.ts +++ b/plugins/tests/scum-feature-module.test.ts @@ -92,16 +92,16 @@ describe("SCUM plugin feature module", () => { const scum = scumActions(); const data = await loadSCUMSurface({ pluginData: pluginDataActions({ list }), scum }, "gifts"); expect(list.mock.calls.map(([collection]) => collection)).toEqual([scumCollections.gifts, scumCollections.giftClaims, scumCollections.pendingGifts, scumCollections.giftDeliveries, scumCollections.timedGiftEvents, scumCollections.tradeGoods]); - expect(scum.users).toHaveBeenCalledWith({ limit: 500 }); + expect(scum.surface).toHaveBeenCalledTimes(1); expect(data.gifts[0]).toMatchObject({ collection: scumCollections.gifts, _recordKey: `${scumCollections.gifts}-1` }); }); it("reads player records from platform scum_user table only", async () => { const pluginData = pluginDataActions(); - const scum = scumActions({ users: async () => ({ items: [{ gamePlayerId: "steam-1", steamId: "steam-1", displayName: "Mira", online: false, source: "platform.scum_user" }], count: 1 }) }); + const scum = scumActions({ surface: async () => ({ users: { items: [{ gamePlayerId: "steam-1", steamId: "steam-1", displayName: "Mira", online: false, source: "platform.scum_user" }], count: 1 }, vehicles: { items: [], count: 0 }, userTrajectories: { items: [], count: 0 }, vehicleTrajectories: { items: [], count: 0 }, vehicleLocks: { items: [], count: 0 } }) }); const data = await loadSCUMSurface({ pluginData, scum }, "players"); expect(pluginData.list).not.toHaveBeenCalledWith("scum_users"); - expect(scum.users).toHaveBeenCalledWith({ limit: 500 }); + expect(scum.surface).toHaveBeenCalledTimes(1); expect(data.players[0]).toMatchObject({ gamePlayerId: "steam-1", displayName: "Mira", online: false }); expect(dataClientSource).not.toContain("scum-client-manager"); }); @@ -111,11 +111,7 @@ describe("SCUM plugin feature module", () => { const scum = scumActions(); await expect(loadSCUMSurface({ pluginData: pluginDataActions(), scum, dispatch }, "live-map")).resolves.toMatchObject({ players: expect.any(Array), vehicles: expect.any(Array), trajectories: expect.any(Array) }); expect(dispatch).not.toHaveBeenCalled(); - expect(scum.users).toHaveBeenCalledWith({ limit: 500 }); - expect(scum.vehicles).toHaveBeenCalledWith({ limit: 500 }); - expect(scum.userTrajectories).toHaveBeenCalledWith({ limit: 500 }); - expect(scum.vehicleTrajectories).toHaveBeenCalledWith({ limit: 500 }); - expect(scum.vehicleLocks).toHaveBeenCalledWith({ limit: 500 }); + expect(scum.surface).toHaveBeenCalledTimes(1); expect(dataClientSource).not.toContain("remote.run.db.sqlite"); expect(dataClientSource).not.toContain("input.templateKey"); expect(dataClientSource).not.toContain("projectSCUMLoginLogs"); @@ -341,11 +337,13 @@ function pluginDataActions(overrides: Partial<{ list: (collection: string, key?: function scumActions(overrides: Partial> = {}): NonNullable { return { - users: vi.fn(overrides.users ?? (async () => ({ items: surfaceData.players, count: surfaceData.players.length }))), - vehicles: vi.fn(overrides.vehicles ?? (async () => ({ items: surfaceData.vehicles, count: surfaceData.vehicles.length }))), - userTrajectories: vi.fn(overrides.userTrajectories ?? (async () => ({ items: surfaceData.trajectories.filter((row) => row.subjectType === "player"), count: 1 }))), - vehicleTrajectories: vi.fn(overrides.vehicleTrajectories ?? (async () => ({ items: surfaceData.trajectories.filter((row) => row.subjectType === "vehicle"), count: 1 }))), - vehicleLocks: vi.fn(overrides.vehicleLocks ?? (async () => ({ items: surfaceData.vehicleLocks, count: surfaceData.vehicleLocks.length }))) + surface: vi.fn(overrides.surface ?? (async () => ({ + users: { items: surfaceData.players, count: surfaceData.players.length }, + vehicles: { items: surfaceData.vehicles, count: surfaceData.vehicles.length }, + userTrajectories: { items: surfaceData.trajectories.filter((row) => row.subjectType === "player"), count: 1 }, + vehicleTrajectories: { items: surfaceData.trajectories.filter((row) => row.subjectType === "vehicle"), count: 1 }, + vehicleLocks: { items: surfaceData.vehicleLocks, count: surfaceData.vehicleLocks.length } + }))) }; }