Implement SCUM direct data plane
This commit is contained in:
@@ -231,6 +231,7 @@ type Core interface {
|
||||
ListSCUMVehiclesForSession(string, domain.SCUMProjectionFilter) ([]domain.SCUMVehicle, error)
|
||||
ListSCUMFlagsForSession(string, domain.SCUMProjectionFilter) ([]domain.SCUMFlag, error)
|
||||
ListSCUMCurrentPositionsForSession(string, domain.SCUMProjectionFilter) ([]domain.SCUMCurrentPosition, error)
|
||||
ListSCUMDataRowsForSession(string, domain.SCUMProjectionFilter) ([]domain.SCUMDataRow, error)
|
||||
RequestSCUMOperationForSession(string, string, domain.SCUMOperationRequest) (domain.SCUMOperationRequest, error)
|
||||
ListSCUMOperationsForSession(string, domain.SCUMOperationRequestFilter) ([]domain.SCUMOperationRequest, error)
|
||||
ApproveSCUMOperationForSession(string, string) (domain.SCUMOperationRequest, error)
|
||||
|
||||
@@ -59,12 +59,28 @@ func (svc *CoreService) ApplySCUMObservationResult(result domain.SCUMObservation
|
||||
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}
|
||||
if err := svc.applySCUMRows(result.QueryKey, result.ServerInstanceID, result.Rows, freshness); err != nil {
|
||||
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
|
||||
@@ -174,44 +190,39 @@ func (svc *CoreService) upsertSCUMObservation(observation domain.SCUMDataObserva
|
||||
return svc.store.SCUMDataObservations().Create(observation)
|
||||
}
|
||||
|
||||
func (svc *CoreService) applySCUMRows(queryKey, serverID string, rows []map[string]any, freshness domain.SCUMProjectionFreshnessState) error {
|
||||
lower := strings.ToLower(queryKey)
|
||||
if strings.Contains(lower, "player") || strings.Contains(lower, "profile") || strings.Contains(lower, "economy") {
|
||||
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
|
||||
}
|
||||
}
|
||||
}
|
||||
if strings.Contains(lower, "squad-member") || strings.Contains(lower, "squad.member") || strings.Contains(lower, "member") {
|
||||
for _, row := range rows {
|
||||
if err := svc.applySCUMSquadMemberRow(serverID, row, freshness); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
} else if strings.Contains(lower, "squad") {
|
||||
for _, row := range rows {
|
||||
case "scum.squads":
|
||||
if err := svc.applySCUMSquadRow(serverID, row, freshness); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
if strings.Contains(lower, "vehicle") {
|
||||
for _, row := range rows {
|
||||
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
|
||||
}
|
||||
}
|
||||
}
|
||||
if strings.Contains(lower, "flag") {
|
||||
for _, row := range rows {
|
||||
case "scum.flags":
|
||||
if err := svc.applySCUMFlagRow(serverID, row, freshness); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
if strings.Contains(lower, "position") || strings.Contains(lower, "coordinate") {
|
||||
for _, row := range rows {
|
||||
case "scum.positions":
|
||||
if err := svc.applySCUMPositionRow(serverID, row, freshness); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -220,6 +231,78 @@ func (svc *CoreService) applySCUMRows(queryKey, serverID string, rows []map[stri
|
||||
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.SCUMDataSetGifts, 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")
|
||||
@@ -558,8 +641,27 @@ func (svc *CoreService) upsertSCUMPosition(position domain.SCUMCurrentPosition)
|
||||
|
||||
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}
|
||||
lower := strings.ToLower(result.QueryKey)
|
||||
if strings.Contains(lower, "player") || strings.Contains(lower, "profile") || strings.Contains(lower, "economy") {
|
||||
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
|
||||
@@ -573,8 +675,7 @@ func (svc *CoreService) markSCUMQueryStale(result domain.SCUMObservationResult,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if strings.Contains(lower, "squad") {
|
||||
case "scum.squads", "scum.squad-members":
|
||||
values, err := svc.store.SCUMSquads().List(domain.SCUMProjectionFilter{ServerInstanceID: result.ServerInstanceID})
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -588,8 +689,7 @@ func (svc *CoreService) markSCUMQueryStale(result domain.SCUMObservationResult,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if strings.Contains(lower, "vehicle") {
|
||||
case "scum.vehicles":
|
||||
values, err := svc.store.SCUMVehicles().List(domain.SCUMProjectionFilter{ServerInstanceID: result.ServerInstanceID})
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -603,8 +703,7 @@ func (svc *CoreService) markSCUMQueryStale(result domain.SCUMObservationResult,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if strings.Contains(lower, "flag") {
|
||||
case "scum.flags":
|
||||
values, err := svc.store.SCUMFlags().List(domain.SCUMProjectionFilter{ServerInstanceID: result.ServerInstanceID})
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
@@ -108,3 +108,29 @@ func TestSCUMLoginLogsProjectLiveStateAndDatabaseSaveTimeDoesNotProveOnline(t *t
|
||||
t.Fatalf("last_save_time was incorrectly treated as online proof: %+v", states[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestSCUMObservationUsesDeclaredRowTargetInsteadOfQueryKeyName(t *testing.T) {
|
||||
svc, _ := newRegisteredLogIngestService(t)
|
||||
plugin, err := svc.store.GamePlugins().Get("server.scum")
|
||||
if err != nil {
|
||||
t.Fatalf("get plugin: %v", err)
|
||||
}
|
||||
plugin.GameClientBridge.QueryTemplates = append(plugin.GameClientBridge.QueryTemplates, domain.GameClientBridgeQueryTemplateDeclaration{Key: "v57.catalog.people", RowTarget: &domain.SCUMRowTargetDeclaration{TargetTable: string(domain.SCUMDataSetUsers), UpsertKeys: []string{"profileId"}, ColumnMappings: map[string]string{"profileId": "user_profile_id", "name": "display_name"}}})
|
||||
if err := svc.store.GamePlugins().Update(plugin); err != nil {
|
||||
t.Fatalf("update plugin: %v", err)
|
||||
}
|
||||
_, err = svc.ApplySCUMObservationResult(domain.SCUMObservationResult{ServerInstanceID: "server-1", PluginID: "server.scum", Source: "run.sqlite.read", QueryKey: "v57.catalog.people", Sequence: 1, ObservedAt: time.Now().UTC(), Rows: []map[string]any{{"user_profile_id": "profile-1", "display_name": "Moon", "unmapped": "kept"}}})
|
||||
if err != nil {
|
||||
t.Fatalf("apply declared data row: %v", err)
|
||||
}
|
||||
rows, err := svc.store.SCUMDataRows().List(domain.SCUMProjectionFilter{ServerInstanceID: "server-1", TargetTable: domain.SCUMDataSetUsers})
|
||||
if err != nil || len(rows) != 1 {
|
||||
t.Fatalf("rows=%+v err=%v", rows, err)
|
||||
}
|
||||
if rows[0].Fields["profileId"] != "profile-1" || rows[0].Payload["unmapped"] != "kept" {
|
||||
t.Fatalf("unexpected declared row: %+v", rows[0])
|
||||
}
|
||||
if states, err := svc.store.SCUMPlayerLiveStates().List(domain.SCUMProjectionFilter{ServerInstanceID: "server-1"}); err != nil || len(states) != 0 {
|
||||
t.Fatalf("query key leaked into legacy projection dispatch: states=%+v err=%v", states, err)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user