903 lines
36 KiB
Go
903 lines
36 KiB
Go
package service
|
|
|
|
import (
|
|
"fmt"
|
|
"math"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"browser.local/platform/domain"
|
|
"browser.local/platform/repo"
|
|
"browser.local/platform/validator"
|
|
)
|
|
|
|
func (svc *CoreService) ApplySCUMObservationResult(result domain.SCUMObservationResult) (domain.SCUMDataObservation, error) {
|
|
result = domain.CopySCUMObservationResult(result)
|
|
if strings.TrimSpace(result.ServerInstanceID) == "" {
|
|
return domain.SCUMDataObservation{}, validationError("serverInstanceId is required")
|
|
}
|
|
instance, err := svc.store.ServerInstances().Get(result.ServerInstanceID)
|
|
if err != nil {
|
|
return domain.SCUMDataObservation{}, err
|
|
}
|
|
if strings.TrimSpace(result.PluginID) == "" {
|
|
result.PluginID = instance.PluginID
|
|
}
|
|
if result.PluginID != instance.PluginID {
|
|
return domain.SCUMDataObservation{}, validationError("pluginId must match server instance")
|
|
}
|
|
if strings.TrimSpace(result.QueryKey) == "" {
|
|
return domain.SCUMDataObservation{}, validationError("queryKey is required")
|
|
}
|
|
if result.ReceivedAt.IsZero() {
|
|
result.ReceivedAt = svc.now()
|
|
}
|
|
if result.ObservedAt.IsZero() {
|
|
result.ObservedAt = result.ReceivedAt
|
|
}
|
|
if result.Status == "" {
|
|
result.Status = domain.SCUMObservationAccepted
|
|
}
|
|
latest, err := svc.latestSCUMObservation(result.ServerInstanceID, result.PluginID, result.QueryKey)
|
|
if err != nil {
|
|
return domain.SCUMDataObservation{}, err
|
|
}
|
|
if result.Status == domain.SCUMObservationAccepted && !latest.ObservedAt.IsZero() && scumObservationOlder(result, latest) {
|
|
result.Status = domain.SCUMObservationStale
|
|
result.ErrorCode = "older_observation"
|
|
result.SafeSummary = domain.SCUMSafeSummary{Title: "旧观察已忽略", Message: "Run 返回的 SCUM.db 观察早于当前本地投影,未覆盖 last-known-good 数据。"}
|
|
}
|
|
observation := domain.SCUMDataObservation{ID: scumObservationID(result), ServerInstanceID: result.ServerInstanceID, PluginID: result.PluginID, Source: result.Source, QueryKey: result.QueryKey, Sequence: result.Sequence, Checksum: result.Checksum, Status: result.Status, ErrorCode: result.ErrorCode, SafeSummary: result.SafeSummary, ObservedAt: result.ObservedAt, ReceivedAt: result.ReceivedAt}
|
|
if err := svc.upsertSCUMObservation(observation); err != nil {
|
|
return domain.SCUMDataObservation{}, err
|
|
}
|
|
if result.Status != domain.SCUMObservationAccepted {
|
|
if result.Status == domain.SCUMObservationFailed {
|
|
return observation, svc.markSCUMQueryStale(result, "observation_failed")
|
|
}
|
|
return observation, nil
|
|
}
|
|
freshness := domain.SCUMProjectionFreshnessState{Status: domain.SCUMProjectionFresh, ObservationID: observation.ID, Source: observation.Source, QueryKey: observation.QueryKey, Sequence: observation.Sequence, Checksum: observation.Checksum, ObservedAt: observation.ObservedAt, ReceivedAt: observation.ReceivedAt}
|
|
target, err := svc.scumRowTarget(result.PluginID, result.QueryKey)
|
|
if err != nil {
|
|
return domain.SCUMDataObservation{}, err
|
|
}
|
|
if err := svc.applySCUMRows(target, result.PluginID, result.QueryKey, result.ServerInstanceID, result.Rows, freshness); err != nil {
|
|
return domain.SCUMDataObservation{}, err
|
|
}
|
|
return observation, nil
|
|
}
|
|
|
|
func (svc *CoreService) ListSCUMDataRowsForSession(sessionID string, filter domain.SCUMProjectionFilter) ([]domain.SCUMDataRow, error) {
|
|
if err := svc.authorizeServerLifecycle(sessionID, filter.ServerInstanceID); err != nil {
|
|
return nil, err
|
|
}
|
|
values, err := svc.store.SCUMDataRows().List(filter)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
limitSCUMProjectionSlice(&values, filter.Limit)
|
|
return values, nil
|
|
}
|
|
|
|
func (svc *CoreService) ListSCUMPlayerLiveStatesForSession(sessionID string, filter domain.SCUMProjectionFilter) ([]domain.SCUMPlayerLiveState, error) {
|
|
if err := svc.authorizeServerLifecycle(sessionID, filter.ServerInstanceID); err != nil {
|
|
return nil, err
|
|
}
|
|
values, err := svc.store.SCUMPlayerLiveStates().List(filter)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
limitSCUMProjectionSlice(&values, filter.Limit)
|
|
return values, nil
|
|
}
|
|
|
|
func (svc *CoreService) ListSCUMSquadsForSession(sessionID string, filter domain.SCUMProjectionFilter) ([]domain.SCUMSquad, error) {
|
|
if err := svc.authorizeServerLifecycle(sessionID, filter.ServerInstanceID); err != nil {
|
|
return nil, err
|
|
}
|
|
values, err := svc.store.SCUMSquads().List(filter)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
limitSCUMProjectionSlice(&values, filter.Limit)
|
|
return values, nil
|
|
}
|
|
|
|
func (svc *CoreService) ListSCUMSquadMembersForSession(sessionID string, filter domain.SCUMProjectionFilter) ([]domain.SCUMSquadMember, error) {
|
|
if err := svc.authorizeServerLifecycle(sessionID, filter.ServerInstanceID); err != nil {
|
|
return nil, err
|
|
}
|
|
values, err := svc.store.SCUMSquadMembers().List(filter)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
limitSCUMProjectionSlice(&values, filter.Limit)
|
|
return values, nil
|
|
}
|
|
|
|
func (svc *CoreService) ListSCUMVehiclesForSession(sessionID string, filter domain.SCUMProjectionFilter) ([]domain.SCUMVehicle, error) {
|
|
if err := svc.authorizeServerLifecycle(sessionID, filter.ServerInstanceID); err != nil {
|
|
return nil, err
|
|
}
|
|
values, err := svc.store.SCUMVehicles().List(filter)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
limitSCUMProjectionSlice(&values, filter.Limit)
|
|
return values, nil
|
|
}
|
|
|
|
func (svc *CoreService) ListSCUMFlagsForSession(sessionID string, filter domain.SCUMProjectionFilter) ([]domain.SCUMFlag, error) {
|
|
if err := svc.authorizeServerLifecycle(sessionID, filter.ServerInstanceID); err != nil {
|
|
return nil, err
|
|
}
|
|
values, err := svc.store.SCUMFlags().List(filter)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
limitSCUMProjectionSlice(&values, filter.Limit)
|
|
return values, nil
|
|
}
|
|
|
|
func (svc *CoreService) ListSCUMCurrentPositionsForSession(sessionID string, filter domain.SCUMProjectionFilter) ([]domain.SCUMCurrentPosition, error) {
|
|
if err := svc.authorizeServerLifecycle(sessionID, filter.ServerInstanceID); err != nil {
|
|
return nil, err
|
|
}
|
|
values, err := svc.store.SCUMCurrentPositions().List(filter)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
limitSCUMProjectionSlice(&values, filter.Limit)
|
|
return values, nil
|
|
}
|
|
|
|
func (svc *CoreService) latestSCUMObservation(serverID, pluginID, queryKey string) (domain.SCUMDataObservation, error) {
|
|
observations, err := svc.store.SCUMDataObservations().List(domain.SCUMProjectionFilter{ServerInstanceID: serverID, QueryKey: queryKey})
|
|
if err != nil {
|
|
return domain.SCUMDataObservation{}, err
|
|
}
|
|
var latest domain.SCUMDataObservation
|
|
for _, observation := range observations {
|
|
if pluginID != "" && observation.PluginID != pluginID {
|
|
continue
|
|
}
|
|
if latest.ObservedAt.IsZero() || observation.Sequence > latest.Sequence || (observation.Sequence == latest.Sequence && observation.ObservedAt.After(latest.ObservedAt)) {
|
|
latest = observation
|
|
}
|
|
}
|
|
return latest, nil
|
|
}
|
|
|
|
func scumObservationOlder(next domain.SCUMObservationResult, latest domain.SCUMDataObservation) bool {
|
|
if next.Sequence > 0 && latest.Sequence > 0 && next.Sequence <= latest.Sequence {
|
|
return true
|
|
}
|
|
return !next.ObservedAt.IsZero() && !latest.ObservedAt.IsZero() && next.ObservedAt.Before(latest.ObservedAt)
|
|
}
|
|
|
|
func (svc *CoreService) upsertSCUMObservation(observation domain.SCUMDataObservation) error {
|
|
if existing, err := svc.store.SCUMDataObservations().Get(observation.ID); err == nil {
|
|
existing.Status = observation.Status
|
|
existing.ErrorCode = observation.ErrorCode
|
|
existing.SafeSummary = observation.SafeSummary
|
|
existing.ReceivedAt = observation.ReceivedAt
|
|
return svc.store.SCUMDataObservations().Update(existing)
|
|
} else if err != repo.ErrNotFound {
|
|
return err
|
|
}
|
|
return svc.store.SCUMDataObservations().Create(observation)
|
|
}
|
|
|
|
func (svc *CoreService) applySCUMRows(target domain.SCUMRowTargetDeclaration, pluginID, queryKey, serverID string, rows []map[string]any, freshness domain.SCUMProjectionFreshnessState) error {
|
|
if target.TargetTable != "" {
|
|
for _, row := range rows {
|
|
if err := svc.upsertSCUMDataRow(target, pluginID, queryKey, serverID, row, freshness); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
// Compatibility declarations retain the old projections without guessing from substrings.
|
|
for _, row := range rows {
|
|
switch queryKey {
|
|
case "scum.player.profile":
|
|
if err := svc.applySCUMPlayerRow(serverID, row, freshness); err != nil {
|
|
return err
|
|
}
|
|
case "scum.squads":
|
|
if err := svc.applySCUMSquadRow(serverID, row, freshness); err != nil {
|
|
return err
|
|
}
|
|
case "scum.squad-members":
|
|
if err := svc.applySCUMSquadMemberRow(serverID, row, freshness); err != nil {
|
|
return err
|
|
}
|
|
case "scum.vehicles":
|
|
if err := svc.applySCUMVehicleRow(serverID, row, freshness); err != nil {
|
|
return err
|
|
}
|
|
case "scum.flags":
|
|
if err := svc.applySCUMFlagRow(serverID, row, freshness); err != nil {
|
|
return err
|
|
}
|
|
case "scum.positions":
|
|
if err := svc.applySCUMPositionRow(serverID, row, freshness); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (svc *CoreService) scumRowTarget(pluginID, queryKey string) (domain.SCUMRowTargetDeclaration, error) {
|
|
plugin, err := svc.store.GamePlugins().Get(pluginID)
|
|
if err != nil {
|
|
return domain.SCUMRowTargetDeclaration{}, err
|
|
}
|
|
for _, template := range plugin.GameClientBridge.QueryTemplates {
|
|
if template.Key == queryKey && template.RowTarget != nil {
|
|
return domain.CopySCUMRowTargetDeclaration(*template.RowTarget), nil
|
|
}
|
|
}
|
|
if _, ok := legacySCUMQueryKeys[queryKey]; ok {
|
|
return domain.SCUMRowTargetDeclaration{}, nil
|
|
}
|
|
return domain.SCUMRowTargetDeclaration{}, validationError("queryKey does not declare a SCUM row target")
|
|
}
|
|
|
|
var legacySCUMQueryKeys = map[string]struct{}{
|
|
"scum.player.profile": {}, "scum.squads": {}, "scum.squad-members": {}, "scum.vehicles": {}, "scum.flags": {}, "scum.positions": {},
|
|
}
|
|
|
|
func (svc *CoreService) upsertSCUMDataRow(target domain.SCUMRowTargetDeclaration, pluginID, queryKey, serverID string, row map[string]any, freshness domain.SCUMProjectionFreshnessState) error {
|
|
table := domain.SCUMDataSet(strings.TrimSpace(target.TargetTable))
|
|
if !validSCUMDataSet(table) || len(target.UpsertKeys) == 0 {
|
|
return validationError("SCUM row target is invalid")
|
|
}
|
|
fields := map[string]any{}
|
|
for destination, source := range target.ColumnMappings {
|
|
if value, ok := row[source]; ok {
|
|
fields[destination] = value
|
|
}
|
|
}
|
|
keyValues := make([]string, 0, len(target.UpsertKeys))
|
|
for _, key := range target.UpsertKeys {
|
|
value, ok := fields[key]
|
|
if !ok {
|
|
value, ok = row[key]
|
|
}
|
|
text := firstString(map[string]any{"value": value}, "value")
|
|
if !ok || text == "" {
|
|
return validationError("SCUM row is missing declared upsert key " + key)
|
|
}
|
|
keyValues = append(keyValues, text)
|
|
}
|
|
upsertKey := strings.Join(keyValues, "\x00")
|
|
id := scumProjectionID(string(table), serverID, upsertKey)
|
|
value, err := svc.store.SCUMDataRows().Get(id)
|
|
if err == repo.ErrNotFound {
|
|
value = domain.SCUMDataRow{ID: id, ServerInstanceID: serverID, TargetTable: table, UpsertKey: upsertKey, CreatedAt: svc.now()}
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
if isProjectionOlder(freshness, value.Freshness) {
|
|
return nil
|
|
}
|
|
value.Fields = domain.CopyGameClientBridgePayload(fields)
|
|
value.Payload = domain.CopyGameClientBridgePayload(row)
|
|
value.PluginID, value.QueryKey, value.Freshness, value.UpdatedAt = pluginID, queryKey, freshness, svc.now()
|
|
if err == repo.ErrNotFound {
|
|
return svc.store.SCUMDataRows().Create(value)
|
|
}
|
|
return svc.store.SCUMDataRows().Update(value)
|
|
}
|
|
|
|
func validSCUMDataSet(value domain.SCUMDataSet) bool {
|
|
switch value {
|
|
case domain.SCUMDataSetUsers, domain.SCUMDataSetSquads, domain.SCUMDataSetMembers, domain.SCUMDataSetVehicles, domain.SCUMDataSetFlags, domain.SCUMDataSetActivity, domain.SCUMDataSetGiftEvents, domain.SCUMDataSetMapPoints:
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
func (svc *CoreService) applySCUMPlayerRow(serverID string, row map[string]any, freshness domain.SCUMProjectionFreshnessState) error {
|
|
gamePlayerID := firstString(row, "gamePlayerId", "playerId", "steamId", "steam_id")
|
|
profileID := firstString(row, "userProfileId", "user_profile_id", "profileId")
|
|
steamID := firstString(row, "steamId", "steam_id")
|
|
name := firstString(row, "displayName", "name", "playerName")
|
|
if gamePlayerID == "" && steamID != "" {
|
|
gamePlayerID = steamID
|
|
}
|
|
if gamePlayerID == "" && profileID == "" {
|
|
return nil
|
|
}
|
|
playerRecordID := ""
|
|
if gamePlayerID != "" {
|
|
playerRecordID = gamePlayerRecordID(serverID, gamePlayerID)
|
|
if err := svc.upsertSCUMGamePlayer(serverID, playerRecordID, gamePlayerID, name, freshness.ObservedAt); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
idSource := gamePlayerID
|
|
if idSource == "" {
|
|
idSource = "profile-" + profileID
|
|
}
|
|
id := scumProjectionID("player-live", serverID, idSource)
|
|
state, err := svc.store.SCUMPlayerLiveStates().Get(id)
|
|
if err == repo.ErrNotFound {
|
|
state = domain.SCUMPlayerLiveState{ID: id, ServerInstanceID: serverID, GamePlayerRecordID: playerRecordID, GamePlayerID: gamePlayerID, UserProfileID: profileID, SteamID: steamID, DisplayName: name, Freshness: domain.SCUMProjectionStateUnknown(), CreatedAt: svc.now()}
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
if isProjectionOlder(freshness, state.Freshness) {
|
|
return nil
|
|
}
|
|
state.GamePlayerRecordID = coalesceString(playerRecordID, state.GamePlayerRecordID)
|
|
state.GamePlayerID = coalesceString(gamePlayerID, state.GamePlayerID)
|
|
state.UserProfileID = coalesceString(profileID, state.UserProfileID)
|
|
state.SteamID = coalesceString(steamID, state.SteamID)
|
|
state.DisplayName = coalesceString(name, state.DisplayName)
|
|
state.SquadID = coalesceString(firstString(row, "squadId", "squad_id"), state.SquadID)
|
|
state.SquadName = coalesceString(firstString(row, "squadName", "squad_name"), state.SquadName)
|
|
if value, ok := firstFloat(row, "famePoints", "fame_points", "fame"); ok {
|
|
state.FamePoints = value
|
|
}
|
|
if value, ok := firstFloat(row, "normalBalance", "currencyNormal", "money", "normal_balance"); ok {
|
|
state.NormalBalance = value
|
|
}
|
|
if value, ok := firstFloat(row, "goldBalance", "currencyGold", "gold", "gold_balance"); ok {
|
|
state.GoldBalance = value
|
|
}
|
|
if value, ok := firstBool(row, "online", "isOnline"); ok {
|
|
state.Online = value
|
|
}
|
|
state.LastLoginAt = coalesceTime(firstTime(row, "lastLoginAt", "last_login_at"), state.LastLoginAt)
|
|
state.LastLogoutAt = coalesceTime(firstTime(row, "lastLogoutAt", "last_logout_at"), state.LastLogoutAt)
|
|
state.LastSaveTime = coalesceTime(firstTime(row, "lastSaveTime", "last_save_time"), state.LastSaveTime)
|
|
if position, ok := scumPositionFromRow(serverID, domain.SCUMProjectionSubjectPlayer, gamePlayerID, row, freshness); ok {
|
|
position.GamePlayerRecordID = playerRecordID
|
|
position.GamePlayerID = gamePlayerID
|
|
state.Position = position
|
|
if err := svc.upsertSCUMPosition(position); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
state.UnknownFields = unknownRowFields(row, "gamePlayerId", "playerId", "steamId", "steam_id", "userProfileId", "user_profile_id", "profileId", "displayName", "name", "playerName", "squadId", "squad_id", "squadName", "squad_name", "famePoints", "fame_points", "fame", "normalBalance", "currencyNormal", "money", "normal_balance", "goldBalance", "currencyGold", "gold", "gold_balance", "online", "isOnline", "lastLoginAt", "last_login_at", "lastLogoutAt", "last_logout_at", "lastSaveTime", "last_save_time", "x", "y", "z", "worldX", "worldY", "worldZ", "mapId", "mapVersion")
|
|
state.Freshness = freshness
|
|
state.UpdatedAt = svc.now()
|
|
if err == repo.ErrNotFound {
|
|
return svc.store.SCUMPlayerLiveStates().Create(state)
|
|
}
|
|
return svc.store.SCUMPlayerLiveStates().Update(state)
|
|
}
|
|
|
|
func (svc *CoreService) applySCUMSquadRow(serverID string, row map[string]any, freshness domain.SCUMProjectionFreshnessState) error {
|
|
squadID := firstString(row, "squadId", "squad_id", "id")
|
|
if squadID == "" {
|
|
return nil
|
|
}
|
|
id := scumProjectionID("squad", serverID, squadID)
|
|
value, err := svc.store.SCUMSquads().Get(id)
|
|
if err == repo.ErrNotFound {
|
|
value = domain.SCUMSquad{ID: id, ServerInstanceID: serverID, SquadID: squadID, Freshness: domain.SCUMProjectionStateUnknown(), CreatedAt: svc.now()}
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
if isProjectionOlder(freshness, value.Freshness) {
|
|
return nil
|
|
}
|
|
value.Name = coalesceString(firstString(row, "name", "squadName", "squad_name"), value.Name)
|
|
value.LeaderProfileID = coalesceString(firstString(row, "leaderProfileId", "leader_profile_id"), value.LeaderProfileID)
|
|
value.LeaderPlayerID = coalesceString(firstString(row, "leaderPlayerId", "leader_player_id", "leaderSteamId"), value.LeaderPlayerID)
|
|
if memberCount, ok := firstInt(row, "memberCount", "member_count"); ok {
|
|
value.MemberCount = memberCount
|
|
}
|
|
if score, ok := firstFloat(row, "score", "fame", "points"); ok {
|
|
value.Score = score
|
|
}
|
|
value.UnknownFields = unknownRowFields(row, "squadId", "squad_id", "id", "name", "squadName", "squad_name", "leaderProfileId", "leader_profile_id", "leaderPlayerId", "leader_player_id", "leaderSteamId", "memberCount", "member_count", "score", "fame", "points")
|
|
value.Freshness = freshness
|
|
value.UpdatedAt = svc.now()
|
|
if err == repo.ErrNotFound {
|
|
return svc.store.SCUMSquads().Create(value)
|
|
}
|
|
return svc.store.SCUMSquads().Update(value)
|
|
}
|
|
|
|
func (svc *CoreService) applySCUMSquadMemberRow(serverID string, row map[string]any, freshness domain.SCUMProjectionFreshnessState) error {
|
|
squadID := firstString(row, "squadId", "squad_id")
|
|
profileID := firstString(row, "userProfileId", "user_profile_id", "profileId")
|
|
gamePlayerID := firstString(row, "gamePlayerId", "playerId", "steamId", "steam_id")
|
|
if squadID == "" || (profileID == "" && gamePlayerID == "") {
|
|
return nil
|
|
}
|
|
playerRecordID := ""
|
|
if gamePlayerID != "" {
|
|
playerRecordID = gamePlayerRecordID(serverID, gamePlayerID)
|
|
if err := svc.upsertSCUMGamePlayer(serverID, playerRecordID, gamePlayerID, firstString(row, "displayName", "name", "playerName"), freshness.ObservedAt); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
id := scumProjectionID("squad-member", serverID, squadID+"/"+coalesceString(profileID, gamePlayerID))
|
|
value, err := svc.store.SCUMSquadMembers().Get(id)
|
|
if err == repo.ErrNotFound {
|
|
value = domain.SCUMSquadMember{ID: id, ServerInstanceID: serverID, SquadID: squadID, UserProfileID: profileID, GamePlayerRecordID: playerRecordID, GamePlayerID: gamePlayerID, Freshness: domain.SCUMProjectionStateUnknown(), CreatedAt: svc.now()}
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
if isProjectionOlder(freshness, value.Freshness) {
|
|
return nil
|
|
}
|
|
value.UserProfileID = coalesceString(profileID, value.UserProfileID)
|
|
value.GamePlayerRecordID = coalesceString(playerRecordID, value.GamePlayerRecordID)
|
|
value.GamePlayerID = coalesceString(gamePlayerID, value.GamePlayerID)
|
|
value.SteamID = coalesceString(firstString(row, "steamId", "steam_id"), value.SteamID)
|
|
value.DisplayName = coalesceString(firstString(row, "displayName", "name", "playerName"), value.DisplayName)
|
|
value.Rank = coalesceString(firstString(row, "rank", "role"), value.Rank)
|
|
if isLeader, ok := firstBool(row, "isLeader", "leader"); ok {
|
|
value.IsLeader = isLeader
|
|
}
|
|
value.JoinedAt = coalesceTime(firstTime(row, "joinedAt", "joined_at"), value.JoinedAt)
|
|
value.UnknownFields = unknownRowFields(row, "squadId", "squad_id", "userProfileId", "user_profile_id", "profileId", "gamePlayerId", "playerId", "steamId", "steam_id", "displayName", "name", "playerName", "rank", "role", "isLeader", "leader", "joinedAt", "joined_at")
|
|
value.Freshness = freshness
|
|
value.UpdatedAt = svc.now()
|
|
if err == repo.ErrNotFound {
|
|
return svc.store.SCUMSquadMembers().Create(value)
|
|
}
|
|
return svc.store.SCUMSquadMembers().Update(value)
|
|
}
|
|
|
|
func (svc *CoreService) applySCUMVehicleRow(serverID string, row map[string]any, freshness domain.SCUMProjectionFreshnessState) error {
|
|
vehicleID := firstString(row, "vehicleId", "vehicle_id", "id")
|
|
entityID := firstString(row, "entityId", "entity_id")
|
|
if vehicleID == "" && entityID != "" {
|
|
vehicleID = entityID
|
|
}
|
|
if vehicleID == "" {
|
|
return nil
|
|
}
|
|
id := scumProjectionID("vehicle", serverID, vehicleID)
|
|
value, err := svc.store.SCUMVehicles().Get(id)
|
|
if err == repo.ErrNotFound {
|
|
value = domain.SCUMVehicle{ID: id, ServerInstanceID: serverID, VehicleID: vehicleID, Freshness: domain.SCUMProjectionStateUnknown(), CreatedAt: svc.now()}
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
if isProjectionOlder(freshness, value.Freshness) {
|
|
return nil
|
|
}
|
|
value.EntityID = coalesceString(entityID, value.EntityID)
|
|
value.ClassName = coalesceString(firstString(row, "className", "class", "type"), value.ClassName)
|
|
value.Label = coalesceString(firstString(row, "label", "vehicleName", "name"), value.Label)
|
|
if value.Label == "" {
|
|
value.Label = coalesceString(value.ClassName, "Unknown vehicle")
|
|
}
|
|
value.OwnerProfileID = coalesceString(firstString(row, "ownerProfileId", "owner_profile_id", "userProfileId", "user_profile_id"), value.OwnerProfileID)
|
|
value.OwnerPlayerID = coalesceString(firstString(row, "ownerPlayerId", "owner_player_id", "steamId", "steam_id"), value.OwnerPlayerID)
|
|
value.SquadID = coalesceString(firstString(row, "squadId", "squad_id"), value.SquadID)
|
|
if position, ok := scumPositionFromRow(serverID, domain.SCUMProjectionSubjectVehicle, vehicleID, row, freshness); ok {
|
|
position.VehicleID = vehicleID
|
|
position.EntityID = entityID
|
|
value.Position = position
|
|
if err := svc.upsertSCUMPosition(position); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
value.UnknownFields = unknownRowFields(row, "vehicleId", "vehicle_id", "id", "entityId", "entity_id", "className", "class", "type", "label", "vehicleName", "name", "ownerProfileId", "owner_profile_id", "userProfileId", "user_profile_id", "ownerPlayerId", "owner_player_id", "steamId", "steam_id", "squadId", "squad_id", "x", "y", "z", "worldX", "worldY", "worldZ", "mapId", "mapVersion")
|
|
value.Freshness = freshness
|
|
value.UpdatedAt = svc.now()
|
|
if err == repo.ErrNotFound {
|
|
return svc.store.SCUMVehicles().Create(value)
|
|
}
|
|
return svc.store.SCUMVehicles().Update(value)
|
|
}
|
|
|
|
func (svc *CoreService) applySCUMFlagRow(serverID string, row map[string]any, freshness domain.SCUMProjectionFreshnessState) error {
|
|
flagID := firstString(row, "flagId", "flag_id", "baseElementId", "base_element_id", "id")
|
|
entityID := firstString(row, "entityId", "entity_id")
|
|
if flagID == "" && entityID != "" {
|
|
flagID = entityID
|
|
}
|
|
if flagID == "" {
|
|
return nil
|
|
}
|
|
id := scumProjectionID("flag", serverID, flagID)
|
|
value, err := svc.store.SCUMFlags().Get(id)
|
|
if err == repo.ErrNotFound {
|
|
value = domain.SCUMFlag{ID: id, ServerInstanceID: serverID, FlagID: flagID, Freshness: domain.SCUMProjectionStateUnknown(), CreatedAt: svc.now()}
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
if isProjectionOlder(freshness, value.Freshness) {
|
|
return nil
|
|
}
|
|
value.EntityID = coalesceString(entityID, value.EntityID)
|
|
value.OwnerProfileID = coalesceString(firstString(row, "ownerProfileId", "owner_profile_id", "userProfileId", "user_profile_id"), value.OwnerProfileID)
|
|
value.OwnerPlayerID = coalesceString(firstString(row, "ownerPlayerId", "owner_player_id", "steamId", "steam_id"), value.OwnerPlayerID)
|
|
value.OwnerSquadID = coalesceString(firstString(row, "ownerSquadId", "owner_squad_id", "squadId", "squad_id"), value.OwnerSquadID)
|
|
value.OwnerSquadName = coalesceString(firstString(row, "ownerSquadName", "owner_squad_name", "squadName", "squad_name"), value.OwnerSquadName)
|
|
value.OwnershipConfidence = coalesceString(firstString(row, "ownershipConfidence", "ownership_confidence"), value.OwnershipConfidence)
|
|
if value.OwnershipConfidence == "" {
|
|
value.OwnershipConfidence = "unknown"
|
|
}
|
|
if position, ok := scumPositionFromRow(serverID, domain.SCUMProjectionSubjectFlag, flagID, row, freshness); ok {
|
|
position.EntityID = entityID
|
|
value.Position = position
|
|
if err := svc.upsertSCUMPosition(position); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
value.UnknownFields = unknownRowFields(row, "flagId", "flag_id", "baseElementId", "base_element_id", "id", "entityId", "entity_id", "ownerProfileId", "owner_profile_id", "userProfileId", "user_profile_id", "ownerPlayerId", "owner_player_id", "steamId", "steam_id", "ownerSquadId", "owner_squad_id", "squadId", "squad_id", "ownerSquadName", "owner_squad_name", "squadName", "squad_name", "ownershipConfidence", "ownership_confidence", "x", "y", "z", "worldX", "worldY", "worldZ", "mapId", "mapVersion")
|
|
value.Freshness = freshness
|
|
value.UpdatedAt = svc.now()
|
|
if err == repo.ErrNotFound {
|
|
return svc.store.SCUMFlags().Create(value)
|
|
}
|
|
return svc.store.SCUMFlags().Update(value)
|
|
}
|
|
|
|
func (svc *CoreService) applySCUMPositionRow(serverID string, row map[string]any, freshness domain.SCUMProjectionFreshnessState) error {
|
|
subjectType := domain.SCUMProjectionSubject(firstString(row, "subjectType", "subject_type"))
|
|
if subjectType == "" {
|
|
if firstString(row, "vehicleId", "vehicle_id") != "" {
|
|
subjectType = domain.SCUMProjectionSubjectVehicle
|
|
} else {
|
|
subjectType = domain.SCUMProjectionSubjectPlayer
|
|
}
|
|
}
|
|
subjectID := firstString(row, "subjectId", "subject_id", "gamePlayerId", "playerId", "vehicleId", "flagId", "entityId", "id")
|
|
position, ok := scumPositionFromRow(serverID, subjectType, subjectID, row, freshness)
|
|
if !ok {
|
|
return nil
|
|
}
|
|
position.GamePlayerID = firstString(row, "gamePlayerId", "playerId", "steamId", "steam_id")
|
|
if position.GamePlayerID != "" {
|
|
position.GamePlayerRecordID = gamePlayerRecordID(serverID, position.GamePlayerID)
|
|
}
|
|
position.VehicleID = firstString(row, "vehicleId", "vehicle_id")
|
|
position.EntityID = firstString(row, "entityId", "entity_id")
|
|
return svc.upsertSCUMPosition(position)
|
|
}
|
|
|
|
func (svc *CoreService) upsertSCUMGamePlayer(serverID, recordID, gamePlayerID, displayName string, observedAt time.Time) error {
|
|
if gamePlayerID == "" || recordID == "" {
|
|
return nil
|
|
}
|
|
if observedAt.IsZero() {
|
|
observedAt = svc.now()
|
|
}
|
|
player, err := svc.store.GamePlayers().Get(recordID)
|
|
if err == repo.ErrNotFound {
|
|
return svc.store.GamePlayers().Create(domain.GamePlayer{ID: recordID, ServerInstanceID: serverID, GamePlayerID: gamePlayerID, DisplayName: displayName, FirstSeenAt: observedAt, LastSeenAt: observedAt, LastEventAt: observedAt, CreatedAt: svc.now(), UpdatedAt: svc.now()})
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if observedAt.Before(player.LastEventAt) {
|
|
return nil
|
|
}
|
|
player.DisplayName = coalesceString(displayName, player.DisplayName)
|
|
player.LastSeenAt = maxTime(player.LastSeenAt, observedAt)
|
|
player.LastEventAt = observedAt
|
|
player.UpdatedAt = svc.now()
|
|
return svc.store.GamePlayers().Update(player)
|
|
}
|
|
|
|
func (svc *CoreService) projectSCUMLoginLiveState(player domain.GamePlayer, batch domain.LogBatchIngest, entry domain.LogEntry, observedAt time.Time, online bool, reason string) error {
|
|
if player.ID == "" || player.GamePlayerID == "" {
|
|
return nil
|
|
}
|
|
freshness := domain.SCUMProjectionFreshnessState{Status: domain.SCUMProjectionFresh, ObservationID: entryID(batch.LogStreamID, entry.Seq), Source: "login-log", QueryKey: strings.TrimSpace(entry.Fields["eventType"]), Sequence: entry.Seq, Checksum: validator.LogLineChecksum(entry.Line), ObservedAt: observedAt, ReceivedAt: svc.now()}
|
|
id := scumProjectionID("player-live", player.ServerInstanceID, player.GamePlayerID)
|
|
state, err := svc.store.SCUMPlayerLiveStates().Get(id)
|
|
if err == repo.ErrNotFound {
|
|
state = domain.SCUMPlayerLiveState{ID: id, ServerInstanceID: player.ServerInstanceID, GamePlayerRecordID: player.ID, GamePlayerID: player.GamePlayerID, DisplayName: player.DisplayName, Freshness: domain.SCUMProjectionStateUnknown(), CreatedAt: svc.now()}
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
if isProjectionOlder(freshness, state.Freshness) {
|
|
return nil
|
|
}
|
|
state.GamePlayerRecordID = player.ID
|
|
state.GamePlayerID = player.GamePlayerID
|
|
state.DisplayName = player.DisplayName
|
|
state.Online = online
|
|
if online {
|
|
state.LastLoginAt = observedAt
|
|
} else {
|
|
state.LastLogoutAt = observedAt
|
|
}
|
|
state.Freshness = freshness
|
|
if reason != "" {
|
|
state.UnknownFields = domain.CopyGameClientBridgePayload(map[string]any{"lastLogoutReason": bounded(reason, 80)})
|
|
}
|
|
state.UpdatedAt = svc.now()
|
|
if err == repo.ErrNotFound {
|
|
return svc.store.SCUMPlayerLiveStates().Create(state)
|
|
}
|
|
return svc.store.SCUMPlayerLiveStates().Update(state)
|
|
}
|
|
|
|
func (svc *CoreService) upsertSCUMPosition(position domain.SCUMCurrentPosition) error {
|
|
existing, err := svc.store.SCUMCurrentPositions().Get(position.ID)
|
|
if err == repo.ErrNotFound {
|
|
position.CreatedAt = svc.now()
|
|
position.UpdatedAt = svc.now()
|
|
return svc.store.SCUMCurrentPositions().Create(position)
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if isProjectionOlder(position.Freshness, existing.Freshness) {
|
|
return nil
|
|
}
|
|
position.CreatedAt = existing.CreatedAt
|
|
position.UpdatedAt = svc.now()
|
|
return svc.store.SCUMCurrentPositions().Update(position)
|
|
}
|
|
|
|
func (svc *CoreService) markSCUMQueryStale(result domain.SCUMObservationResult, reason string) error {
|
|
freshness := domain.SCUMProjectionFreshnessState{Status: domain.SCUMProjectionStale, ObservationID: scumObservationID(result), Source: result.Source, QueryKey: result.QueryKey, Sequence: result.Sequence, Checksum: result.Checksum, StaleReason: reason, ObservedAt: result.ObservedAt, ReceivedAt: result.ReceivedAt}
|
|
target, err := svc.scumRowTarget(result.PluginID, result.QueryKey)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if target.TargetTable != "" {
|
|
values, err := svc.store.SCUMDataRows().List(domain.SCUMProjectionFilter{ServerInstanceID: result.ServerInstanceID, TargetTable: domain.SCUMDataSet(target.TargetTable)})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, value := range values {
|
|
if !isProjectionOlder(freshness, value.Freshness) {
|
|
value.Freshness, value.UpdatedAt = freshness, svc.now()
|
|
if err := svc.store.SCUMDataRows().Update(value); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
switch result.QueryKey {
|
|
case "scum.player.profile":
|
|
values, err := svc.store.SCUMPlayerLiveStates().List(domain.SCUMProjectionFilter{ServerInstanceID: result.ServerInstanceID})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, value := range values {
|
|
if !isProjectionOlder(freshness, value.Freshness) {
|
|
value.Freshness = freshness
|
|
value.UpdatedAt = svc.now()
|
|
if err := svc.store.SCUMPlayerLiveStates().Update(value); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
case "scum.squads", "scum.squad-members":
|
|
values, err := svc.store.SCUMSquads().List(domain.SCUMProjectionFilter{ServerInstanceID: result.ServerInstanceID})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, value := range values {
|
|
if !isProjectionOlder(freshness, value.Freshness) {
|
|
value.Freshness = freshness
|
|
value.UpdatedAt = svc.now()
|
|
if err := svc.store.SCUMSquads().Update(value); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
case "scum.vehicles":
|
|
values, err := svc.store.SCUMVehicles().List(domain.SCUMProjectionFilter{ServerInstanceID: result.ServerInstanceID})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, value := range values {
|
|
if !isProjectionOlder(freshness, value.Freshness) {
|
|
value.Freshness = freshness
|
|
value.UpdatedAt = svc.now()
|
|
if err := svc.store.SCUMVehicles().Update(value); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
case "scum.flags":
|
|
values, err := svc.store.SCUMFlags().List(domain.SCUMProjectionFilter{ServerInstanceID: result.ServerInstanceID})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, value := range values {
|
|
if !isProjectionOlder(freshness, value.Freshness) {
|
|
value.Freshness = freshness
|
|
value.UpdatedAt = svc.now()
|
|
if err := svc.store.SCUMFlags().Update(value); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func scumPositionFromRow(serverID string, subjectType domain.SCUMProjectionSubject, subjectID string, row map[string]any, freshness domain.SCUMProjectionFreshnessState) (domain.SCUMCurrentPosition, bool) {
|
|
x, hasX := firstFloat(row, "x", "worldX", "world_x", "locationX")
|
|
y, hasY := firstFloat(row, "y", "worldY", "world_y", "locationY")
|
|
z, hasZ := firstFloat(row, "z", "worldZ", "world_z", "locationZ")
|
|
if !hasX || !hasY {
|
|
return domain.SCUMCurrentPosition{}, false
|
|
}
|
|
if subjectID == "" {
|
|
return domain.SCUMCurrentPosition{}, false
|
|
}
|
|
position := domain.SCUMCurrentPosition{ID: scumProjectionID("position-"+string(subjectType), serverID, subjectID), ServerInstanceID: serverID, SubjectType: subjectType, SubjectID: subjectID, MapID: coalesceString(firstString(row, "mapId", "map_id"), domain.SCUMMapTrajectoryMapID), MapVersion: coalesceString(firstString(row, "mapVersion", "map_version"), "0.9"), X: x, Y: y, HasCoordinates: true, LastSaveTime: firstTime(row, "lastSaveTime", "last_save_time"), Freshness: freshness}
|
|
if hasZ && !math.IsNaN(z) {
|
|
position.Z = z
|
|
}
|
|
return position, true
|
|
}
|
|
|
|
func isProjectionOlder(next, current domain.SCUMProjectionFreshnessState) bool {
|
|
if current.Status == "" || current.Status == domain.SCUMProjectionUnknown {
|
|
return false
|
|
}
|
|
if next.Source == current.Source && next.QueryKey == current.QueryKey && next.Sequence > 0 && current.Sequence > 0 && next.Sequence < current.Sequence {
|
|
return true
|
|
}
|
|
return !next.ObservedAt.IsZero() && !current.ObservedAt.IsZero() && next.ObservedAt.Before(current.ObservedAt)
|
|
}
|
|
|
|
func scumObservationID(result domain.SCUMObservationResult) string {
|
|
seed := fmt.Sprintf("%s/%s/%s/%d/%s", result.ServerInstanceID, result.PluginID, result.QueryKey, result.Sequence, result.Checksum)
|
|
if result.Checksum == "" {
|
|
seed = fmt.Sprintf("%s/%s/%s/%d/%s", result.ServerInstanceID, result.PluginID, result.QueryKey, result.Sequence, result.ObservedAt.Format(time.RFC3339Nano))
|
|
}
|
|
return "scum-observation-" + fingerprintID(result.ServerInstanceID, seed)
|
|
}
|
|
|
|
func scumProjectionID(kind, serverID, subject string) string {
|
|
return "scum-" + kind + "-" + fingerprintID(serverID, subject)
|
|
}
|
|
|
|
func firstString(row map[string]any, keys ...string) string {
|
|
for _, key := range keys {
|
|
if value, ok := row[key]; ok {
|
|
switch typed := value.(type) {
|
|
case string:
|
|
if trimmed := strings.TrimSpace(typed); trimmed != "" {
|
|
return trimmed
|
|
}
|
|
case fmt.Stringer:
|
|
if trimmed := strings.TrimSpace(typed.String()); trimmed != "" {
|
|
return trimmed
|
|
}
|
|
case int, int64, uint64, float64:
|
|
return fmt.Sprint(typed)
|
|
}
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func firstFloat(row map[string]any, keys ...string) (float64, bool) {
|
|
for _, key := range keys {
|
|
if value, ok := row[key]; ok {
|
|
switch typed := value.(type) {
|
|
case float64:
|
|
return typed, true
|
|
case float32:
|
|
return float64(typed), true
|
|
case int:
|
|
return float64(typed), true
|
|
case int64:
|
|
return float64(typed), true
|
|
case uint64:
|
|
return float64(typed), true
|
|
case string:
|
|
parsed, err := strconv.ParseFloat(strings.TrimSpace(typed), 64)
|
|
if err == nil {
|
|
return parsed, true
|
|
}
|
|
}
|
|
}
|
|
}
|
|
return 0, false
|
|
}
|
|
|
|
func firstInt(row map[string]any, keys ...string) (int, bool) {
|
|
value, ok := firstFloat(row, keys...)
|
|
if !ok {
|
|
return 0, false
|
|
}
|
|
return int(value), true
|
|
}
|
|
|
|
func firstBool(row map[string]any, keys ...string) (bool, bool) {
|
|
for _, key := range keys {
|
|
if value, ok := row[key]; ok {
|
|
switch typed := value.(type) {
|
|
case bool:
|
|
return typed, true
|
|
case string:
|
|
parsed, err := strconv.ParseBool(strings.TrimSpace(typed))
|
|
if err == nil {
|
|
return parsed, true
|
|
}
|
|
case int:
|
|
return typed != 0, true
|
|
case int64:
|
|
return typed != 0, true
|
|
case float64:
|
|
return typed != 0, true
|
|
}
|
|
}
|
|
}
|
|
return false, false
|
|
}
|
|
|
|
func firstTime(row map[string]any, keys ...string) time.Time {
|
|
for _, key := range keys {
|
|
if value, ok := row[key]; ok {
|
|
switch typed := value.(type) {
|
|
case time.Time:
|
|
return typed
|
|
case string:
|
|
trimmed := strings.TrimSpace(typed)
|
|
if trimmed == "" {
|
|
continue
|
|
}
|
|
if parsed, err := time.Parse(time.RFC3339Nano, trimmed); err == nil {
|
|
return parsed
|
|
}
|
|
if parsed, err := time.Parse("2006-01-02 15:04:05", trimmed); err == nil {
|
|
return parsed.UTC()
|
|
}
|
|
case int64:
|
|
return time.Unix(typed, 0).UTC()
|
|
case float64:
|
|
return time.Unix(int64(typed), 0).UTC()
|
|
}
|
|
}
|
|
}
|
|
return time.Time{}
|
|
}
|
|
|
|
func unknownRowFields(row map[string]any, known ...string) map[string]any {
|
|
knownSet := map[string]struct{}{}
|
|
for _, key := range known {
|
|
knownSet[key] = struct{}{}
|
|
}
|
|
unknown := map[string]any{}
|
|
for key, value := range row {
|
|
if _, ok := knownSet[key]; ok {
|
|
continue
|
|
}
|
|
unknown[key] = value
|
|
}
|
|
if len(unknown) == 0 {
|
|
return nil
|
|
}
|
|
return domain.CopyGameClientBridgePayload(unknown)
|
|
}
|
|
|
|
func coalesceString(next, current string) string {
|
|
if strings.TrimSpace(next) != "" {
|
|
return strings.TrimSpace(next)
|
|
}
|
|
return current
|
|
}
|
|
|
|
func coalesceTime(next, current time.Time) time.Time {
|
|
if !next.IsZero() {
|
|
return next
|
|
}
|
|
return current
|
|
}
|
|
|
|
func limitSCUMProjectionSlice[T any](values *[]T, limit int) {
|
|
if limit > 0 && len(*values) > limit {
|
|
*values = (*values)[:limit]
|
|
}
|
|
}
|