191 lines
7.1 KiB
Go
191 lines
7.1 KiB
Go
package service
|
|
|
|
import (
|
|
"errors"
|
|
"sort"
|
|
"strings"
|
|
|
|
"browser.local/platform/domain"
|
|
"browser.local/platform/repo"
|
|
)
|
|
|
|
const (
|
|
defaultPluginDataListLimit = 500
|
|
maxPluginDataListLimit = 2000
|
|
)
|
|
|
|
func (svc *CoreService) ListPluginDataForSession(sessionID string, filter domain.PluginDataFilter) ([]domain.PluginDataRecord, error) {
|
|
if err := svc.authorizePluginData(sessionID, filter.PluginID, filter.ServerInstanceID, filter.Collection); err != nil {
|
|
return nil, err
|
|
}
|
|
values, err := svc.store.PluginDataRecords().List(filter)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
sort.SliceStable(values, func(left, right int) bool {
|
|
if !values[left].UpdatedAt.Equal(values[right].UpdatedAt) {
|
|
return values[left].UpdatedAt.After(values[right].UpdatedAt)
|
|
}
|
|
if !values[left].CreatedAt.Equal(values[right].CreatedAt) {
|
|
return values[left].CreatedAt.After(values[right].CreatedAt)
|
|
}
|
|
return values[left].Key < values[right].Key
|
|
})
|
|
limit := boundedPluginDataLimit(filter.Limit)
|
|
if len(values) > limit {
|
|
values = values[:limit]
|
|
}
|
|
return values, nil
|
|
}
|
|
|
|
func (svc *CoreService) DeletePluginDataForSession(sessionID, pluginID, serverInstanceID, collection, key string) error {
|
|
transaction := domain.PluginDataTransaction{PluginID: pluginID, ServerInstanceID: serverInstanceID, Collection: collection, Mutations: []domain.PluginDataMutation{{Operation: domain.PluginDataMutationDelete, Key: key}}}
|
|
_, err := svc.ApplyPluginDataTransactionForSession(sessionID, transaction)
|
|
return err
|
|
}
|
|
|
|
func (svc *CoreService) ApplyPluginDataTransactionForSession(sessionID string, transaction domain.PluginDataTransaction) ([]domain.PluginDataRecord, error) {
|
|
if err := svc.authorizePluginData(sessionID, transaction.PluginID, transaction.ServerInstanceID, transaction.Collection); err != nil {
|
|
return nil, err
|
|
}
|
|
return svc.applyPluginDataTransaction(transaction)
|
|
}
|
|
|
|
func (svc *CoreService) applyPluginDataTransaction(transaction domain.PluginDataTransaction) ([]domain.PluginDataRecord, error) {
|
|
if len(transaction.Mutations) == 0 {
|
|
return nil, validationError("plugin data mutations are required")
|
|
}
|
|
if err := svc.validatePluginDataMutations(transaction); err != nil {
|
|
return nil, err
|
|
}
|
|
stamp := svc.now()
|
|
upserts := make([]domain.PluginDataRecord, 0, len(transaction.Mutations))
|
|
deleteIDs := make([]string, 0, len(transaction.Mutations))
|
|
seen := make(map[string]struct{}, len(transaction.Mutations))
|
|
for _, mutation := range transaction.Mutations {
|
|
key := strings.TrimSpace(mutation.Key)
|
|
if key == "" {
|
|
return nil, validationError("plugin data mutation key is required")
|
|
}
|
|
if _, exists := seen[key]; exists {
|
|
return nil, validationError("plugin data mutation keys must be unique")
|
|
}
|
|
seen[key] = struct{}{}
|
|
id := pluginDataID(transaction.ServerInstanceID, transaction.PluginID, transaction.Collection, key)
|
|
switch mutation.Operation {
|
|
case domain.PluginDataMutationPut:
|
|
if mutation.Value == nil {
|
|
return nil, validationError("plugin data mutation value is required")
|
|
}
|
|
createdAt := stamp
|
|
if existing, err := svc.store.PluginDataRecords().Get(id); err == nil {
|
|
createdAt = existing.CreatedAt
|
|
} else if !errors.Is(err, repo.ErrNotFound) {
|
|
return nil, err
|
|
}
|
|
upserts = append(upserts, domain.PluginDataRecord{ID: id, PluginID: transaction.PluginID, ServerInstanceID: transaction.ServerInstanceID, Collection: transaction.Collection, Key: key, Value: domain.CopyGameClientBridgePayload(mutation.Value), CreatedAt: createdAt, UpdatedAt: stamp})
|
|
case domain.PluginDataMutationDelete:
|
|
deleteIDs = append(deleteIDs, id)
|
|
default:
|
|
return nil, validationError("plugin data mutation operation is invalid")
|
|
}
|
|
}
|
|
if err := svc.store.PluginDataRecords().Apply(upserts, deleteIDs); err != nil {
|
|
return nil, err
|
|
}
|
|
result := make([]domain.PluginDataRecord, len(upserts))
|
|
for index, value := range upserts {
|
|
result[index] = domain.CopyPluginDataRecord(value)
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
func (svc *CoreService) PutPluginDataForSession(sessionID string, value domain.PluginDataRecord) (domain.PluginDataRecord, error) {
|
|
if err := svc.authorizePluginData(sessionID, value.PluginID, value.ServerInstanceID, value.Collection); err != nil {
|
|
return domain.PluginDataRecord{}, err
|
|
}
|
|
if strings.TrimSpace(value.Key) == "" {
|
|
return domain.PluginDataRecord{}, validationError("plugin data key is required")
|
|
}
|
|
if value.Value == nil {
|
|
return domain.PluginDataRecord{}, validationError("plugin data value is required")
|
|
}
|
|
if legacySCUMPluginDataCollection(value.PluginID, value.Collection) {
|
|
return domain.PluginDataRecord{}, validationError("SCUM user, vehicle, and trajectory data must be written to platform SCUM tables")
|
|
}
|
|
value.ID = pluginDataID(value.ServerInstanceID, value.PluginID, value.Collection, value.Key)
|
|
stamp := svc.now()
|
|
existing, err := svc.store.PluginDataRecords().Get(value.ID)
|
|
if err == repo.ErrNotFound {
|
|
value.CreatedAt, value.UpdatedAt = stamp, stamp
|
|
if err := svc.store.PluginDataRecords().Create(value); err != nil {
|
|
return domain.PluginDataRecord{}, err
|
|
}
|
|
return domain.CopyPluginDataRecord(value), nil
|
|
}
|
|
if err != nil {
|
|
return domain.PluginDataRecord{}, err
|
|
}
|
|
existing.Value, existing.UpdatedAt = domain.CopyGameClientBridgePayload(value.Value), stamp
|
|
if err := svc.store.PluginDataRecords().Update(existing); err != nil {
|
|
return domain.PluginDataRecord{}, err
|
|
}
|
|
return domain.CopyPluginDataRecord(existing), nil
|
|
}
|
|
|
|
func (svc *CoreService) authorizePluginData(sessionID, pluginID, serverInstanceID, collection string) error {
|
|
if strings.TrimSpace(pluginID) == "" || strings.TrimSpace(serverInstanceID) == "" || strings.TrimSpace(collection) == "" {
|
|
return validationError("pluginId, serverInstanceId, and collection are required")
|
|
}
|
|
if err := svc.authorizeServerLifecycle(sessionID, serverInstanceID); err != nil {
|
|
return err
|
|
}
|
|
instance, err := svc.store.ServerInstances().Get(serverInstanceID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if instance.PluginID != pluginID {
|
|
return ErrForbidden
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (svc *CoreService) validatePluginDataMutations(transaction domain.PluginDataTransaction) error {
|
|
if !legacySCUMPluginDataCollection(transaction.PluginID, transaction.Collection) {
|
|
return nil
|
|
}
|
|
for _, mutation := range transaction.Mutations {
|
|
if mutation.Operation == domain.PluginDataMutationPut {
|
|
return validationError("SCUM user, vehicle, and trajectory data must be written to platform SCUM tables")
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func boundedPluginDataLimit(requested int) int {
|
|
if requested <= 0 {
|
|
return defaultPluginDataListLimit
|
|
}
|
|
if requested > maxPluginDataListLimit {
|
|
return maxPluginDataListLimit
|
|
}
|
|
return requested
|
|
}
|
|
|
|
func legacySCUMPluginDataCollection(pluginID string, collection string) bool {
|
|
canonical := domain.CanonicalGamePluginID(pluginID)
|
|
if canonical != "game.scum" && canonical != "server.scum" {
|
|
return false
|
|
}
|
|
switch strings.ToLower(strings.TrimSpace(collection)) {
|
|
case "scum_users", "scum_players", "scum_trajectories", "scum_user_trajectories", "scum_vehicles", "scum_vehicle_trajectories", "scum_vehicle_locks":
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
func pluginDataID(serverID, pluginID, collection, key string) string {
|
|
return "plugin-data-" + fingerprintID(serverID, pluginID+"\x00"+collection+"\x00"+key)
|
|
}
|