feat(scum): add local game player intelligence
This commit is contained in:
@@ -0,0 +1,337 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"browser.local/platform/domain"
|
||||
"browser.local/platform/repo"
|
||||
"crypto/hmac"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
const gameAccessRetention = 30 * 24 * time.Hour
|
||||
const failedAccessWindow = 15 * time.Minute
|
||||
const failedAccessThreshold = 5
|
||||
|
||||
func (svc *CoreService) ListGamePlayersForSession(sessionID string, filter domain.GamePlayerFilter) ([]domain.GamePlayer, error) {
|
||||
if err := svc.authorizeServerLifecycle(sessionID, filter.ServerInstanceID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := svc.pruneGamePlayerEvidence(filter.ServerInstanceID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
values, err := svc.store.GamePlayers().List(filter)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
sort.Slice(values, func(i, j int) bool { return values[i].LastSeenAt.After(values[j].LastSeenAt) })
|
||||
if filter.Limit > 0 && len(values) > filter.Limit {
|
||||
values = values[:filter.Limit]
|
||||
}
|
||||
return values, nil
|
||||
}
|
||||
|
||||
func (svc *CoreService) GetGamePlayerProfileForSession(sessionID, playerRecordID string) (domain.GamePlayerProfile, error) {
|
||||
player, err := svc.store.GamePlayers().Get(playerRecordID)
|
||||
if err != nil {
|
||||
return domain.GamePlayerProfile{}, err
|
||||
}
|
||||
if err := svc.authorizeServerLifecycle(sessionID, player.ServerInstanceID); err != nil {
|
||||
return domain.GamePlayerProfile{}, err
|
||||
}
|
||||
if err := svc.pruneGamePlayerEvidence(player.ServerInstanceID); err != nil {
|
||||
return domain.GamePlayerProfile{}, err
|
||||
}
|
||||
aliases, err := svc.store.GamePlayerAliases().List(domain.GamePlayerAliasFilter{GamePlayerRecordID: player.ID})
|
||||
if err != nil {
|
||||
return domain.GamePlayerProfile{}, err
|
||||
}
|
||||
sessions, err := svc.store.GamePlayerSessions().List(domain.GamePlayerSessionFilter{GamePlayerRecordID: player.ID})
|
||||
if err != nil {
|
||||
return domain.GamePlayerProfile{}, err
|
||||
}
|
||||
attempts, err := svc.store.GameAccessAttempts().List(domain.GameAccessAttemptFilter{GamePlayerRecordID: player.ID})
|
||||
if err != nil {
|
||||
return domain.GamePlayerProfile{}, err
|
||||
}
|
||||
signals, err := svc.store.GameSecuritySignals().List(domain.GameSecuritySignalFilter{GamePlayerRecordID: player.ID})
|
||||
if err != nil {
|
||||
return domain.GamePlayerProfile{}, err
|
||||
}
|
||||
return domain.GamePlayerProfile{Player: player, Aliases: aliases, Sessions: sessions, AccessAttempts: attempts, SecuritySignals: signals}, nil
|
||||
}
|
||||
|
||||
func (svc *CoreService) projectGamePlayerEvents(batch domain.LogBatchIngest) error {
|
||||
for _, entry := range batch.Entries {
|
||||
if err := svc.projectGamePlayerEvent(batch, entry); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return svc.pruneGamePlayerEvidence(batch.ServerInstanceID)
|
||||
}
|
||||
func (svc *CoreService) projectGamePlayerEvent(batch domain.LogBatchIngest, entry domain.LogEntry) error {
|
||||
fields := entry.Fields
|
||||
if fields == nil {
|
||||
return nil
|
||||
}
|
||||
eventType := strings.TrimSpace(fields["eventType"])
|
||||
if eventType != "scum.login" && eventType != "scum.logout" {
|
||||
return nil
|
||||
}
|
||||
gameID, name := strings.TrimSpace(fields["playerId"]), strings.TrimSpace(fields["playerName"])
|
||||
if gameID == "" || name == "" {
|
||||
return nil
|
||||
}
|
||||
occurred := entry.Timestamp
|
||||
if raw := strings.TrimSpace(fields["occurredAt"]); raw != "" {
|
||||
if parsed, err := time.Parse(time.RFC3339, raw); err == nil {
|
||||
occurred = parsed
|
||||
}
|
||||
}
|
||||
if occurred.IsZero() {
|
||||
occurred = svc.now()
|
||||
}
|
||||
recordID := gamePlayerRecordID(batch.ServerInstanceID, gameID)
|
||||
player, err := svc.store.GamePlayers().Get(recordID)
|
||||
if err == repo.ErrNotFound {
|
||||
player = domain.GamePlayer{ID: recordID, ServerInstanceID: batch.ServerInstanceID, GamePlayerID: gameID, DisplayName: name, FirstSeenAt: occurred, LastSeenAt: occurred, LastEventAt: occurred, CreatedAt: svc.now(), UpdatedAt: svc.now()}
|
||||
if err = svc.store.GamePlayers().Create(player); err != nil {
|
||||
return err
|
||||
}
|
||||
} else if err != nil {
|
||||
return err
|
||||
}
|
||||
if occurred.After(player.LastEventAt) || occurred.Equal(player.LastEventAt) {
|
||||
if name != player.DisplayName {
|
||||
player.DisplayName = name
|
||||
}
|
||||
player.LastSeenAt = maxTime(player.LastSeenAt, occurred)
|
||||
player.LastEventAt = occurred
|
||||
player.UpdatedAt = svc.now()
|
||||
if err := svc.store.GamePlayers().Update(player); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if err := svc.upsertGamePlayerAlias(player, name, occurred); err != nil {
|
||||
return err
|
||||
}
|
||||
sourceSession := strings.TrimSpace(fields["sessionId"])
|
||||
if sourceSession == "" {
|
||||
sourceSession = "event-" + entryID(batch.LogStreamID, entry.Seq)
|
||||
}
|
||||
if eventType == "scum.login" {
|
||||
outcome := strings.TrimSpace(fields["outcome"])
|
||||
if outcome == "accepted" {
|
||||
if err := svc.recordSuccessfulGameAccess(player, batch, entry, occurred, strings.TrimSpace(fields["networkFingerprint"])); err != nil {
|
||||
return err
|
||||
}
|
||||
return svc.openGamePlayerSession(player, sourceSession, occurred)
|
||||
}
|
||||
return svc.recordFailedGameAccess(player, batch, entry, occurred, strings.TrimSpace(fields["networkFingerprint"]))
|
||||
}
|
||||
return svc.closeGamePlayerSession(player, sourceSession, occurred, strings.TrimSpace(fields["reason"]))
|
||||
}
|
||||
|
||||
func (svc *CoreService) recordSuccessfulGameAccess(player domain.GamePlayer, batch domain.LogBatchIngest, entry domain.LogEntry, at time.Time, raw string) error {
|
||||
id := "attempt-" + entryID(batch.LogStreamID, entry.Seq)
|
||||
if _, err := svc.store.GameAccessAttempts().Get(id); err == nil {
|
||||
return nil
|
||||
} else if err != repo.ErrNotFound {
|
||||
return err
|
||||
}
|
||||
key := svc.networkCorrelation(batch.ServerInstanceID, raw)
|
||||
if err := svc.store.GameAccessAttempts().Create(domain.GameAccessAttempt{ID: id, ServerInstanceID: batch.ServerInstanceID, GamePlayerRecordID: player.ID, EventID: entryID(batch.LogStreamID, entry.Seq), OccurredAt: at, Outcome: "accepted", Reason: "login-accepted", NetworkCorrelationKey: key, ExpiresAt: at.Add(gameAccessRetention)}); err != nil {
|
||||
return err
|
||||
}
|
||||
if key != "" {
|
||||
return svc.refreshPossibleAltSignal(player, key, at)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (svc *CoreService) refreshPossibleAltSignal(player domain.GamePlayer, key string, at time.Time) error {
|
||||
attempts, err := svc.store.GameAccessAttempts().List(domain.GameAccessAttemptFilter{ServerInstanceID: player.ServerInstanceID})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
players := map[string]struct{}{}
|
||||
for _, attempt := range attempts {
|
||||
if attempt.Outcome == "accepted" && attempt.NetworkCorrelationKey == key && !attempt.OccurredAt.Before(at.Add(-gameAccessRetention)) {
|
||||
players[attempt.GamePlayerRecordID] = struct{}{}
|
||||
}
|
||||
}
|
||||
if len(players) < 2 {
|
||||
return nil
|
||||
}
|
||||
id := "signal-alt-" + fingerprintID(player.ServerInstanceID, key)
|
||||
signal, err := svc.store.GameSecuritySignals().Get(id)
|
||||
if err == repo.ErrNotFound {
|
||||
return svc.store.GameSecuritySignals().Create(domain.GameSecuritySignal{ID: id, ServerInstanceID: player.ServerInstanceID, GamePlayerRecordID: player.ID, RuleKey: "possible-alt-account", Status: domain.GameSecuritySignalReviewRequired, EvidenceCount: len(players), Summary: "Multiple game identities share server-local access evidence; manual review required", FirstObservedAt: at, LastObservedAt: at, ExpiresAt: at.Add(gameAccessRetention)})
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
signal.EvidenceCount = len(players)
|
||||
signal.LastObservedAt = at
|
||||
signal.ExpiresAt = at.Add(gameAccessRetention)
|
||||
return svc.store.GameSecuritySignals().Update(signal)
|
||||
}
|
||||
func (svc *CoreService) upsertGamePlayerAlias(player domain.GamePlayer, name string, at time.Time) error {
|
||||
id := gamePlayerAliasID(player.ID, name)
|
||||
item, err := svc.store.GamePlayerAliases().Get(id)
|
||||
if err == repo.ErrNotFound {
|
||||
return svc.store.GamePlayerAliases().Create(domain.GamePlayerAlias{ID: id, GamePlayerRecordID: player.ID, ServerInstanceID: player.ServerInstanceID, Alias: name, FirstSeenAt: at, LastSeenAt: at})
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if at.After(item.LastSeenAt) {
|
||||
item.LastSeenAt = at
|
||||
return svc.store.GamePlayerAliases().Update(item)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
func (svc *CoreService) openGamePlayerSession(player domain.GamePlayer, source string, at time.Time) error {
|
||||
id := gamePlayerSessionID(player.ID, source)
|
||||
item, err := svc.store.GamePlayerSessions().Get(id)
|
||||
if err == repo.ErrNotFound {
|
||||
return svc.store.GamePlayerSessions().Create(domain.GamePlayerSession{ID: id, GamePlayerRecordID: player.ID, ServerInstanceID: player.ServerInstanceID, SourceSessionID: source, StartedAt: at, LastEventAt: at})
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if at.After(item.LastEventAt) {
|
||||
item.LastEventAt = at
|
||||
return svc.store.GamePlayerSessions().Update(item)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
func (svc *CoreService) closeGamePlayerSession(player domain.GamePlayer, source string, at time.Time, reason string) error {
|
||||
id := gamePlayerSessionID(player.ID, source)
|
||||
item, err := svc.store.GamePlayerSessions().Get(id)
|
||||
if err == repo.ErrNotFound {
|
||||
return nil
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if item.StartedAt.After(at) || (!item.EndedAt.IsZero() && !at.After(item.EndedAt)) {
|
||||
return nil
|
||||
}
|
||||
item.EndedAt = at
|
||||
item.EndReason = bounded(reason, 40)
|
||||
item.LastEventAt = maxTime(item.LastEventAt, at)
|
||||
return svc.store.GamePlayerSessions().Update(item)
|
||||
}
|
||||
func (svc *CoreService) recordFailedGameAccess(player domain.GamePlayer, batch domain.LogBatchIngest, entry domain.LogEntry, at time.Time, raw string) error {
|
||||
id := "attempt-" + entryID(batch.LogStreamID, entry.Seq)
|
||||
if _, err := svc.store.GameAccessAttempts().Get(id); err == nil {
|
||||
return nil
|
||||
} else if err != repo.ErrNotFound {
|
||||
return err
|
||||
}
|
||||
key := svc.networkCorrelation(batch.ServerInstanceID, raw)
|
||||
attempt := domain.GameAccessAttempt{ID: id, ServerInstanceID: batch.ServerInstanceID, GamePlayerRecordID: player.ID, EventID: entryID(batch.LogStreamID, entry.Seq), OccurredAt: at, Outcome: "rejected", Reason: "login-rejected", NetworkCorrelationKey: key, ExpiresAt: at.Add(gameAccessRetention)}
|
||||
if err := svc.store.GameAccessAttempts().Create(attempt); err != nil {
|
||||
return err
|
||||
}
|
||||
if key != "" {
|
||||
return svc.refreshFailedAccessSignal(player, key, at)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
func (svc *CoreService) refreshFailedAccessSignal(player domain.GamePlayer, key string, at time.Time) error {
|
||||
attempts, err := svc.store.GameAccessAttempts().List(domain.GameAccessAttemptFilter{ServerInstanceID: player.ServerInstanceID})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
count := 0
|
||||
for _, a := range attempts {
|
||||
if a.NetworkCorrelationKey == key && !a.OccurredAt.Before(at.Add(-failedAccessWindow)) {
|
||||
count++
|
||||
}
|
||||
}
|
||||
if count < failedAccessThreshold {
|
||||
return nil
|
||||
}
|
||||
id := "signal-failed-" + fingerprintID(player.ServerInstanceID, key)
|
||||
s, err := svc.store.GameSecuritySignals().Get(id)
|
||||
if err == repo.ErrNotFound {
|
||||
return svc.store.GameSecuritySignals().Create(domain.GameSecuritySignal{ID: id, ServerInstanceID: player.ServerInstanceID, GamePlayerRecordID: player.ID, RuleKey: "excessive-failed-access", Status: domain.GameSecuritySignalReviewRequired, EvidenceCount: count, Summary: "Repeated failed access requires manual review", FirstObservedAt: at, LastObservedAt: at, ExpiresAt: at.Add(gameAccessRetention)})
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
s.EvidenceCount = count
|
||||
s.LastObservedAt = at
|
||||
s.ExpiresAt = at.Add(gameAccessRetention)
|
||||
return svc.store.GameSecuritySignals().Update(s)
|
||||
}
|
||||
func (svc *CoreService) networkCorrelation(serverID, raw string) string {
|
||||
raw = strings.TrimSpace(raw)
|
||||
if raw == "" {
|
||||
return ""
|
||||
}
|
||||
mac := hmac.New(sha256.New, svc.networkFingerprintKey)
|
||||
mac.Write([]byte(serverID))
|
||||
mac.Write([]byte{0})
|
||||
mac.Write([]byte(raw))
|
||||
return hex.EncodeToString(mac.Sum(nil))
|
||||
}
|
||||
func (svc *CoreService) pruneGamePlayerEvidence(serverID string) error {
|
||||
now := svc.now()
|
||||
attempts, err := svc.store.GameAccessAttempts().List(domain.GameAccessAttemptFilter{ServerInstanceID: serverID})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, a := range attempts {
|
||||
if !a.ExpiresAt.After(now) {
|
||||
if err := svc.store.GameAccessAttempts().Delete(a.ID); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
signals, err := svc.store.GameSecuritySignals().List(domain.GameSecuritySignalFilter{ServerInstanceID: serverID})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, s := range signals {
|
||||
if !s.ExpiresAt.After(now) {
|
||||
if err := svc.store.GameSecuritySignals().Delete(s.ID); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
func gamePlayerRecordID(server, id string) string { return "game-player-" + fingerprintID(server, id) }
|
||||
func gamePlayerAliasID(player, name string) string {
|
||||
return "game-player-alias-" + fingerprintID(player, name)
|
||||
}
|
||||
func gamePlayerSessionID(player, session string) string {
|
||||
return "game-player-session-" + fingerprintID(player, session)
|
||||
}
|
||||
func entryID(stream string, seq uint64) string { return stream + "-" + itoa(seq) }
|
||||
func fingerprintID(a, b string) string {
|
||||
sum := sha256.Sum256([]byte(a + "\x00" + b))
|
||||
return hex.EncodeToString(sum[:])[:24]
|
||||
}
|
||||
func itoa(v uint64) string {
|
||||
return strconv.FormatUint(v, 10)
|
||||
}
|
||||
func maxTime(a, b time.Time) time.Time {
|
||||
if b.After(a) {
|
||||
return b
|
||||
}
|
||||
return a
|
||||
}
|
||||
func bounded(v string, n int) string {
|
||||
v = strings.TrimSpace(v)
|
||||
if len(v) > n {
|
||||
return v[:n]
|
||||
}
|
||||
return v
|
||||
}
|
||||
@@ -0,0 +1,119 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"browser.local/platform/domain"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestSCUMGamePlayerProjectionIsIdempotentAndRedactsNetworkMaterial(t *testing.T) {
|
||||
svc, token := newRegisteredLogIngestService(t)
|
||||
createLogStreamFixture(t, svc)
|
||||
base := time.Date(2026, 7, 28, 12, 0, 0, 0, time.UTC)
|
||||
login := gamePlayerBatch(t, token, 1, []domain.LogEntry{{Seq: 1, Timestamp: base, Line: "login", Fields: map[string]string{"eventType": "scum.login", "playerId": "steam-1", "playerName": "Moon", "sessionId": "session-1", "outcome": "accepted", "networkFingerprint": "203.0.113.9"}}})
|
||||
if _, err := svc.IngestLogBatch(login); err != nil {
|
||||
t.Fatalf("login projection: %v", err)
|
||||
}
|
||||
if _, err := svc.IngestLogBatch(login); err != nil {
|
||||
t.Fatalf("duplicate projection: %v", err)
|
||||
}
|
||||
players, err := svc.store.GamePlayers().List(domain.GamePlayerFilter{ServerInstanceID: "server-1"})
|
||||
if err != nil || len(players) != 1 {
|
||||
t.Fatalf("players=%+v err=%v", players, err)
|
||||
}
|
||||
sessions, err := svc.store.GamePlayerSessions().List(domain.GamePlayerSessionFilter{GamePlayerRecordID: players[0].ID})
|
||||
if err != nil || len(sessions) != 1 {
|
||||
t.Fatalf("sessions=%+v err=%v", sessions, err)
|
||||
}
|
||||
raw, err := svc.QueryLogStream(domain.LogStreamCursorQuery{LogStreamID: "log-1", Limit: 10})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(raw.Entries) != 1 || strings.Contains(strings.Join(mapValues(raw.Entries[0].Fields), " "), "203.0.113.9") {
|
||||
t.Fatalf("raw network material leaked into log storage: %+v", raw.Entries)
|
||||
}
|
||||
logout := gamePlayerBatch(t, token, 2, []domain.LogEntry{{Seq: 2, Timestamp: base.Add(time.Minute), Line: "logout", Fields: map[string]string{"eventType": "scum.logout", "playerId": "steam-1", "playerName": "Moon Renamed", "sessionId": "session-1", "reason": "disconnect"}}})
|
||||
if _, err := svc.IngestLogBatch(logout); err != nil {
|
||||
t.Fatalf("logout projection: %v", err)
|
||||
}
|
||||
aliases, _ := svc.store.GamePlayerAliases().List(domain.GamePlayerAliasFilter{GamePlayerRecordID: players[0].ID})
|
||||
sessions, _ = svc.store.GamePlayerSessions().List(domain.GamePlayerSessionFilter{GamePlayerRecordID: players[0].ID})
|
||||
if len(aliases) != 2 || len(sessions) != 1 || sessions[0].EndedAt.IsZero() {
|
||||
t.Fatalf("expected alias history and closed session: aliases=%+v sessions=%+v", aliases, sessions)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSCUMRejectedAccessCreatesReviewOnlySignal(t *testing.T) {
|
||||
svc, token := newRegisteredLogIngestService(t)
|
||||
createLogStreamFixture(t, svc)
|
||||
entries := make([]domain.LogEntry, 0, 5)
|
||||
at := time.Date(2026, 7, 28, 12, 0, 0, 0, time.UTC)
|
||||
for seq := uint64(1); seq <= 5; seq++ {
|
||||
entries = append(entries, domain.LogEntry{Seq: seq, Timestamp: at.Add(time.Duration(seq) * time.Minute), Line: "rejected", Fields: map[string]string{"eventType": "scum.login", "playerId": "steam-rejected", "playerName": "Rejected", "sessionId": "failed", "outcome": "rejected", "networkFingerprint": "198.51.100.8"}})
|
||||
}
|
||||
if _, err := svc.IngestLogBatch(gamePlayerBatch(t, token, 1, entries)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
signals, err := svc.store.GameSecuritySignals().List(domain.GameSecuritySignalFilter{ServerInstanceID: "server-1"})
|
||||
if err != nil || len(signals) != 1 {
|
||||
t.Fatalf("signals=%+v err=%v", signals, err)
|
||||
}
|
||||
if signals[0].Status != domain.GameSecuritySignalReviewRequired || signals[0].EvidenceCount != 5 || strings.Contains(signals[0].Summary, "198.51.100.8") {
|
||||
t.Fatalf("unsafe signal=%+v", signals[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestSCUMSharedServerLocalFingerprintSignalsPossibleAltAccount(t *testing.T) {
|
||||
svc, token := newRegisteredLogIngestService(t)
|
||||
createLogStreamFixture(t, svc)
|
||||
at := time.Date(2026, 7, 28, 12, 0, 0, 0, time.UTC)
|
||||
entries := []domain.LogEntry{{Seq: 1, Timestamp: at, Line: "login", Fields: map[string]string{"eventType": "scum.login", "playerId": "steam-a", "playerName": "A", "sessionId": "a", "outcome": "accepted", "networkFingerprint": "198.51.100.9"}}, {Seq: 2, Timestamp: at.Add(time.Minute), Line: "login", Fields: map[string]string{"eventType": "scum.login", "playerId": "steam-b", "playerName": "B", "sessionId": "b", "outcome": "accepted", "networkFingerprint": "198.51.100.9"}}}
|
||||
if _, err := svc.IngestLogBatch(gamePlayerBatch(t, token, 1, entries)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
signals, err := svc.store.GameSecuritySignals().List(domain.GameSecuritySignalFilter{ServerInstanceID: "server-1"})
|
||||
if err != nil || len(signals) != 1 || signals[0].RuleKey != "possible-alt-account" || signals[0].EvidenceCount != 2 {
|
||||
t.Fatalf("signals=%+v err=%v", signals, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSCUMGamePlayerProjectionIgnoresOutOfOrderLogoutAndPrunesExpiredEvidence(t *testing.T) {
|
||||
svc, token := newRegisteredLogIngestService(t)
|
||||
createLogStreamFixture(t, svc)
|
||||
now := time.Date(2026, 7, 28, 12, 0, 0, 0, time.UTC)
|
||||
login := domain.LogEntry{Seq: 1, Timestamp: now, Line: "login", Fields: map[string]string{"eventType": "scum.login", "playerId": "steam-order", "playerName": "Order", "sessionId": "order", "outcome": "accepted"}}
|
||||
logout := domain.LogEntry{Seq: 2, Timestamp: now.Add(-time.Minute), Line: "logout", Fields: map[string]string{"eventType": "scum.logout", "playerId": "steam-order", "playerName": "Order", "sessionId": "order", "reason": "disconnect"}}
|
||||
if _, err := svc.IngestLogBatch(gamePlayerBatch(t, token, 1, []domain.LogEntry{login, logout})); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
players, _ := svc.store.GamePlayers().List(domain.GamePlayerFilter{ServerInstanceID: "server-1"})
|
||||
sessions, _ := svc.store.GamePlayerSessions().List(domain.GamePlayerSessionFilter{GamePlayerRecordID: players[0].ID})
|
||||
if len(sessions) != 1 || !sessions[0].EndedAt.IsZero() {
|
||||
t.Fatalf("stale logout closed active session: %+v", sessions)
|
||||
}
|
||||
if _, err := svc.ListGamePlayersForSession("", domain.GamePlayerFilter{ServerInstanceID: "server-1"}); err != ErrUnauthorized {
|
||||
t.Fatalf("expected unauthorized player query, got %v", err)
|
||||
}
|
||||
if err := svc.store.GameAccessAttempts().Create(domain.GameAccessAttempt{ID: "expired", ServerInstanceID: "server-1", GamePlayerRecordID: players[0].ID, ExpiresAt: now.Add(-time.Hour)}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
svc.now = func() time.Time { return now }
|
||||
if err := svc.pruneGamePlayerEvidence("server-1"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := svc.store.GameAccessAttempts().Get("expired"); err == nil {
|
||||
t.Fatal("expired access evidence was retained")
|
||||
}
|
||||
}
|
||||
func gamePlayerBatch(t *testing.T, token string, first uint64, entries []domain.LogEntry) domain.LogBatchIngest {
|
||||
t.Helper()
|
||||
return domain.LogBatchIngest{RunEndpointID: "run-local", SessionToken: token, LogStreamID: "log-1", ServerInstanceID: "server-1", StreamKey: "stdout", Source: domain.LogStreamSourceProcess, FirstSeq: first, LastSeq: entries[len(entries)-1].Seq, Compression: "none", Checksum: checksumForEntries(t, entries), Entries: entries}
|
||||
}
|
||||
func mapValues(values map[string]string) []string {
|
||||
out := make([]string, 0, len(values))
|
||||
for _, value := range values {
|
||||
out = append(out, value)
|
||||
}
|
||||
return out
|
||||
}
|
||||
@@ -9,6 +9,7 @@ const defaultLogQueryLimit = 100
|
||||
|
||||
func (svc *CoreService) IngestLogBatch(batch domain.LogBatchIngest) (domain.LogBatchIngestResult, error) {
|
||||
batch = domain.CopyLogBatchIngest(batch)
|
||||
projectionBatch := domain.CopyLogBatchIngest(batch)
|
||||
if err := validator.ValidateLogBatchIngest(batch); err != nil {
|
||||
return domain.LogBatchIngestResult{}, err
|
||||
}
|
||||
@@ -31,6 +32,9 @@ func (svc *CoreService) IngestLogBatch(batch domain.LogBatchIngest) (domain.LogB
|
||||
return domain.LogBatchIngestResult{}, err
|
||||
}
|
||||
if exists && record.LastSeq == batch.LastSeq && record.Checksum == batch.Checksum {
|
||||
if err := svc.projectGamePlayerEvents(projectionBatch); err != nil {
|
||||
return domain.LogBatchIngestResult{}, err
|
||||
}
|
||||
return domain.LogBatchIngestResult{
|
||||
Accepted: true,
|
||||
LogStreamID: batch.LogStreamID,
|
||||
@@ -47,11 +51,13 @@ func (svc *CoreService) IngestLogBatch(batch domain.LogBatchIngest) (domain.LogB
|
||||
return domain.LogBatchIngestResult{}, validationError("log batch firstSeq must follow latest acknowledged sequence")
|
||||
}
|
||||
|
||||
storedBatch := domain.CopyLogBatchIngest(batch)
|
||||
sanitizeGamePlayerNetworkFields(&storedBatch)
|
||||
record := domain.CopyLogBatchRecord(domain.LogBatchRecord{
|
||||
Checksum: batch.Checksum,
|
||||
FirstSeq: batch.FirstSeq,
|
||||
LastSeq: batch.LastSeq,
|
||||
Entries: batch.Entries,
|
||||
Entries: storedBatch.Entries,
|
||||
})
|
||||
if err := svc.logStore.AppendBatch(batch.LogStreamID, record); err != nil {
|
||||
return domain.LogBatchIngestResult{}, err
|
||||
@@ -61,6 +67,9 @@ func (svc *CoreService) IngestLogBatch(batch domain.LogBatchIngest) (domain.LogB
|
||||
if err := svc.store.LogStreams().Update(stream); err != nil {
|
||||
return domain.LogBatchIngestResult{}, err
|
||||
}
|
||||
if err := svc.projectGamePlayerEvents(projectionBatch); err != nil {
|
||||
return domain.LogBatchIngestResult{}, err
|
||||
}
|
||||
return domain.LogBatchIngestResult{
|
||||
Accepted: true,
|
||||
LogStreamID: batch.LogStreamID,
|
||||
@@ -71,6 +80,21 @@ func (svc *CoreService) IngestLogBatch(batch domain.LogBatchIngest) (domain.LogB
|
||||
}, nil
|
||||
}
|
||||
|
||||
// sanitizeGamePlayerNetworkFields removes raw network material before the durable log body is written.
|
||||
func sanitizeGamePlayerNetworkFields(batch *domain.LogBatchIngest) {
|
||||
for index := range batch.Entries {
|
||||
fields := batch.Entries[index].Fields
|
||||
if fields == nil {
|
||||
continue
|
||||
}
|
||||
if fields["eventType"] == "scum.login" {
|
||||
delete(fields, "networkFingerprint")
|
||||
delete(fields, "ip")
|
||||
delete(fields, "ipAddress")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (svc *CoreService) QueryLogStream(query domain.LogStreamCursorQuery) (domain.LogStreamCursorResult, error) {
|
||||
if err := validator.ValidateLogStreamCursorQuery(query); err != nil {
|
||||
return domain.LogStreamCursorResult{}, err
|
||||
|
||||
@@ -203,6 +203,8 @@ type Core interface {
|
||||
QueryLogStreamForSession(string, domain.LogStreamCursorQuery) (domain.LogStreamCursorResult, error)
|
||||
IngestLogBatch(domain.LogBatchIngest) (domain.LogBatchIngestResult, error)
|
||||
QueryLogStream(domain.LogStreamCursorQuery) (domain.LogStreamCursorResult, error)
|
||||
ListGamePlayersForSession(string, domain.GamePlayerFilter) ([]domain.GamePlayer, error)
|
||||
GetGamePlayerProfileForSession(string, string) (domain.GamePlayerProfile, error)
|
||||
CreateAuditEvent(domain.AuditEvent) (domain.AuditEvent, error)
|
||||
GetAuditEvent(string) (domain.AuditEvent, error)
|
||||
ListAuditEvents(domain.AuditEventFilter) ([]domain.AuditEvent, error)
|
||||
@@ -210,28 +212,29 @@ type Core interface {
|
||||
}
|
||||
|
||||
type CoreService struct {
|
||||
store repo.Store
|
||||
now func() time.Time
|
||||
authMu sync.Mutex
|
||||
authSessions map[string]string
|
||||
controlMu sync.Mutex
|
||||
runSessions map[string]domain.RunControlSession
|
||||
runSessionSeq uint64
|
||||
jobMu sync.Mutex
|
||||
bridgeMu sync.Mutex
|
||||
bridgeSeq uint64
|
||||
logStore LogBodyStore
|
||||
artifactStore ArtifactBodyStore
|
||||
artifactMu sync.Mutex
|
||||
artifactTransfers map[string]domain.ArtifactTransferSession
|
||||
artifactPayloads map[string][]byte
|
||||
artifactTransferSeq uint64
|
||||
auditMu sync.Mutex
|
||||
auditSeq uint64
|
||||
productionMu sync.Mutex
|
||||
sourceRCONCommands *sourceRCONCommandBroker
|
||||
aiProviderClient AIProviderClient
|
||||
secretEnvelope SecretEnvelope
|
||||
store repo.Store
|
||||
now func() time.Time
|
||||
authMu sync.Mutex
|
||||
authSessions map[string]string
|
||||
controlMu sync.Mutex
|
||||
runSessions map[string]domain.RunControlSession
|
||||
runSessionSeq uint64
|
||||
jobMu sync.Mutex
|
||||
bridgeMu sync.Mutex
|
||||
bridgeSeq uint64
|
||||
logStore LogBodyStore
|
||||
artifactStore ArtifactBodyStore
|
||||
artifactMu sync.Mutex
|
||||
artifactTransfers map[string]domain.ArtifactTransferSession
|
||||
artifactPayloads map[string][]byte
|
||||
artifactTransferSeq uint64
|
||||
auditMu sync.Mutex
|
||||
auditSeq uint64
|
||||
productionMu sync.Mutex
|
||||
sourceRCONCommands *sourceRCONCommandBroker
|
||||
aiProviderClient AIProviderClient
|
||||
secretEnvelope SecretEnvelope
|
||||
networkFingerprintKey []byte
|
||||
}
|
||||
|
||||
var _ Core = (*CoreService)(nil)
|
||||
@@ -254,17 +257,18 @@ func newCoreServiceWithLogStore(store repo.Store, logStore LogBodyStore, now fun
|
||||
}
|
||||
artifactStore := NewMemoryArtifactBodyStore()
|
||||
service := &CoreService{
|
||||
store: store,
|
||||
now: now,
|
||||
authSessions: map[string]string{},
|
||||
runSessions: map[string]domain.RunControlSession{},
|
||||
logStore: logStore,
|
||||
artifactStore: artifactStore,
|
||||
artifactTransfers: map[string]domain.ArtifactTransferSession{},
|
||||
artifactPayloads: map[string][]byte{},
|
||||
sourceRCONCommands: newSourceRCONCommandBroker(now),
|
||||
aiProviderClient: MockAIProviderClient{},
|
||||
secretEnvelope: newSecretEnvelope(developmentSecretEnvelopeKey),
|
||||
store: store,
|
||||
now: now,
|
||||
authSessions: map[string]string{},
|
||||
runSessions: map[string]domain.RunControlSession{},
|
||||
logStore: logStore,
|
||||
artifactStore: artifactStore,
|
||||
artifactTransfers: map[string]domain.ArtifactTransferSession{},
|
||||
artifactPayloads: map[string][]byte{},
|
||||
sourceRCONCommands: newSourceRCONCommandBroker(now),
|
||||
aiProviderClient: MockAIProviderClient{},
|
||||
secretEnvelope: newSecretEnvelope(developmentSecretEnvelopeKey),
|
||||
networkFingerprintKey: []byte(developmentSecretEnvelopeKey),
|
||||
}
|
||||
return service
|
||||
}
|
||||
|
||||
@@ -78,6 +78,7 @@ func (svc *CoreService) ConfigureSecretEnvelopeKey(secret string) error {
|
||||
return validationError("PLATFORM_SECRET_ENVELOPE_KEY must be at least 32 characters")
|
||||
}
|
||||
svc.secretEnvelope = newSecretEnvelope(secret)
|
||||
svc.networkFingerprintKey = []byte(secret)
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user