Files

466 lines
20 KiB
Go

package repo
import (
"encoding/json"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"sync"
"browser.local/platform/domain"
)
type StoreSnapshot struct {
Users []domain.User `json:"users"`
AuthSessions []domain.AuthSessionRecord `json:"authSessions"`
RunControlSessions []domain.RunControlSession `json:"runControlSessions"`
AIProviders []domain.AIProvider `json:"aiProviders"`
GamePlugins []domain.GamePlugin `json:"gamePlugins"`
ServerInstances []domain.ServerInstance `json:"serverInstances"`
RunEndpoints []domain.RunEndpoint `json:"runEndpoints"`
Jobs []domain.Job `json:"jobs"`
Artifacts []domain.Artifact `json:"artifacts"`
RuntimeBindings []domain.RuntimeBinding `json:"runtimeBindings"`
EncryptedComponentKeys []domain.EncryptedComponentKey `json:"encryptedComponentKeys"`
RunDistributions []domain.RunDistribution `json:"runDistributions"`
DependencyStatuses []domain.DependencyStatus `json:"dependencyStatuses"`
RunUpdateJobs []domain.RunUpdateJob `json:"runUpdateJobs"`
LogStreams []domain.LogStream `json:"logStreams"`
MetricSamples []domain.MetricSample `json:"metricSamples"`
Backups []domain.BackupRecord `json:"backups"`
PluginLifecycles []domain.PluginLifecycleInstallation `json:"pluginLifecycles"`
GameClientBridgeCommands []domain.GameClientBridgeCommand `json:"gameClientBridgeCommands"`
GameClientBridgeSnapshots []domain.GameClientBridgeSnapshot `json:"gameClientBridgeSnapshots"`
GameClientBridgeStreams []domain.GameClientBridgeSnapshotStream `json:"gameClientBridgeStreams"`
PluginDataRecords []domain.PluginDataRecord `json:"pluginDataRecords"`
SCUMUsers []domain.SCUMUser `json:"scumUsers"`
SCUMUserTrajectories []domain.SCUMUserTrajectory `json:"scumUserTrajectories"`
SCUMVehicles []domain.SCUMVehicle `json:"scumVehicles"`
SCUMVehicleTrajectories []domain.SCUMVehicleTrajectory `json:"scumVehicleTrajectories"`
SCUMVehicleLocks []domain.SCUMVehicleLock `json:"scumVehicleLocks"`
}
type runtimeSnapshot struct {
RunControlSessions []domain.RunControlSession `json:"runControlSessions"`
RunEndpoints []domain.RunEndpoint `json:"runEndpoints"`
LogStreams []domain.LogStream `json:"logStreams"`
}
type FileStore struct {
*MemoryStore
path string
runtimePath string
persistMu sync.Mutex
runtimePersistMu sync.Mutex
}
func NewFileStore(path string) (*FileStore, error) {
path = strings.TrimSpace(path)
if path == "" {
return nil, fmt.Errorf("metadata path is required")
}
store := &FileStore{
MemoryStore: NewMemoryStore(),
path: path,
runtimePath: runtimeMetadataPath(path),
}
if err := store.load(); err != nil {
return nil, err
}
return store, nil
}
func (store *FileStore) MetadataPath() string {
return store.path
}
func (store *FileStore) Users() UserRepository {
return &persistentRepository[domain.User, domain.UserFilter]{repository: store.MemoryStore.users, persist: store.persist}
}
func (store *FileStore) AuthSessions() AuthSessionRepository {
return &persistentRepository[domain.AuthSessionRecord, domain.AuthSessionFilter]{repository: store.MemoryStore.authSessions, persist: store.persist}
}
func (store *FileStore) RunControlSessions() RunControlSessionRepository {
return &persistentRepository[domain.RunControlSession, struct{}]{repository: store.MemoryStore.runSessions, persist: store.persistRuntime}
}
func (store *FileStore) AIProviders() AIProviderRepository {
return &persistentRepository[domain.AIProvider, domain.AIProviderFilter]{repository: store.MemoryStore.aiProviders, persist: store.persist}
}
func (store *FileStore) GamePlugins() GamePluginRepository {
return &persistentRepository[domain.GamePlugin, domain.GamePluginFilter]{repository: store.MemoryStore.gamePlugins, persist: store.persist}
}
func (store *FileStore) ServerInstances() ServerInstanceRepository {
return &persistentRepository[domain.ServerInstance, domain.ServerInstanceFilter]{repository: store.MemoryStore.serverInstances, persist: store.persist}
}
func (store *FileStore) RunEndpoints() RunEndpointRepository {
return &persistentRepository[domain.RunEndpoint, domain.RunEndpointFilter]{repository: store.MemoryStore.runEndpoints, persist: store.persistRuntime}
}
func (store *FileStore) Jobs() JobRepository {
return &persistentJobRepository{
persistentRepository: &persistentRepository[domain.Job, domain.JobFilter]{repository: store.MemoryStore.jobs, persist: store.persist},
repository: store.MemoryStore.jobs,
}
}
func (store *FileStore) Artifacts() ArtifactRepository {
return &persistentRepository[domain.Artifact, domain.ArtifactFilter]{repository: store.MemoryStore.artifacts, persist: store.persist}
}
func (store *FileStore) RuntimeBindings() RuntimeBindingRepository {
return &persistentRepository[domain.RuntimeBinding, domain.RuntimeBindingFilter]{repository: store.MemoryStore.runtimeBindings, persist: store.persist}
}
func (store *FileStore) EncryptedComponentKeys() EncryptedComponentKeyRepository {
return &persistentRepository[domain.EncryptedComponentKey, domain.EncryptedComponentKeyFilter]{repository: store.MemoryStore.componentKeys, persist: store.persist}
}
func (store *FileStore) RunDistributions() RunDistributionRepository {
return &persistentRepository[domain.RunDistribution, domain.RunDistributionFilter]{repository: store.MemoryStore.runDists, persist: store.persist}
}
func (store *FileStore) DependencyStatuses() DependencyStatusRepository {
return &persistentRepository[domain.DependencyStatus, domain.DependencyStatusFilter]{repository: store.MemoryStore.dependencies, persist: store.persist}
}
func (store *FileStore) RunUpdateJobs() RunUpdateJobRepository {
return &persistentRepository[domain.RunUpdateJob, domain.RunUpdateJobFilter]{repository: store.MemoryStore.updateJobs, persist: store.persist}
}
func (store *FileStore) LogStreams() LogStreamRepository {
return &persistentRepository[domain.LogStream, domain.LogStreamFilter]{repository: store.MemoryStore.logStreams, persist: store.persistRuntime}
}
func (store *FileStore) MetricSamples() MetricSampleRepository {
return &persistentRepository[domain.MetricSample, domain.MetricSampleFilter]{repository: store.MemoryStore.metricSamples, persist: store.persist}
}
func (store *FileStore) Backups() BackupRepository {
return &persistentRepository[domain.BackupRecord, domain.BackupFilter]{repository: store.MemoryStore.backups, persist: store.persist}
}
func (store *FileStore) PluginLifecycles() PluginLifecycleRepository {
return &persistentRepository[domain.PluginLifecycleInstallation, domain.PluginLifecycleFilter]{repository: store.MemoryStore.pluginLifecycle, persist: store.persist}
}
func (store *FileStore) GameClientBridgeCommands() GameClientBridgeCommandRepository {
return &persistentGameClientBridgeCommandRepository{
persistentRepository: &persistentRepository[domain.GameClientBridgeCommand, domain.GameClientBridgeCommandFilter]{repository: store.MemoryStore.bridgeCommands, persist: store.persist},
repository: store.MemoryStore.bridgeCommands,
}
}
func (store *FileStore) GameClientBridgeSnapshots() GameClientBridgeSnapshotRepository {
return &persistentRepository[domain.GameClientBridgeSnapshot, domain.GameClientBridgeSnapshotFilter]{repository: store.MemoryStore.bridgeSnapshots, persist: store.persist}
}
func (store *FileStore) GameClientBridgeSnapshotStreams() GameClientBridgeSnapshotStreamRepository {
return &persistentRepository[domain.GameClientBridgeSnapshotStream, domain.GameClientBridgeSnapshotStreamFilter]{repository: store.MemoryStore.bridgeStreams, persist: store.persist}
}
func (store *FileStore) PluginDataRecords() PluginDataRecordRepository {
return &persistentRepository[domain.PluginDataRecord, domain.PluginDataFilter]{repository: store.MemoryStore.pluginDataRecords, persist: store.persist}
}
func (store *FileStore) SCUMUsers() SCUMUserRepository {
return &persistentRepository[domain.SCUMUser, domain.SCUMUserFilter]{repository: store.MemoryStore.scumUsers, persist: store.persist}
}
func (store *FileStore) SCUMUserTrajectories() SCUMUserTrajectoryRepository {
return &persistentRepository[domain.SCUMUserTrajectory, domain.SCUMUserTrajectoryFilter]{repository: store.MemoryStore.scumUserTracks, persist: store.persist}
}
func (store *FileStore) SCUMVehicles() SCUMVehicleRepository {
return &persistentRepository[domain.SCUMVehicle, domain.SCUMVehicleFilter]{repository: store.MemoryStore.scumVehicles, persist: store.persist}
}
func (store *FileStore) SCUMVehicleTrajectories() SCUMVehicleTrajectoryRepository {
return &persistentRepository[domain.SCUMVehicleTrajectory, domain.SCUMVehicleTrajectoryFilter]{repository: store.MemoryStore.scumVehicleTracks, persist: store.persist}
}
func (store *FileStore) SCUMVehicleLocks() SCUMVehicleLockRepository {
return &persistentRepository[domain.SCUMVehicleLock, domain.SCUMVehicleLockFilter]{repository: store.MemoryStore.scumVehicleLocks, persist: store.persist}
}
func (store *FileStore) load() error {
data, err := os.ReadFile(store.path)
if err != nil {
if os.IsNotExist(err) {
return store.loadRuntime()
}
return fmt.Errorf("read metadata snapshot: %w", err)
}
if len(strings.TrimSpace(string(data))) != 0 {
var snapshot StoreSnapshot
if err := json.Unmarshal(data, &snapshot); err != nil {
return fmt.Errorf("decode metadata snapshot: %w", err)
}
store.loadSnapshot(snapshot)
}
return store.loadRuntime()
}
func (store *FileStore) persist() error {
store.persistMu.Lock()
defer store.persistMu.Unlock()
snapshot := store.snapshot()
data, err := json.Marshal(snapshot)
if err != nil {
return fmt.Errorf("encode metadata snapshot: %w", err)
}
if err := os.MkdirAll(filepath.Dir(store.path), 0o755); err != nil {
return fmt.Errorf("create metadata directory: %w", err)
}
tmpPath := store.path + ".tmp"
if err := os.WriteFile(tmpPath, data, 0o600); err != nil {
return fmt.Errorf("write metadata snapshot: %w", err)
}
if err := os.Rename(tmpPath, store.path); err != nil {
return fmt.Errorf("replace metadata snapshot: %w", err)
}
return nil
}
func (store *FileStore) loadRuntime() error {
data, err := os.ReadFile(store.runtimePath)
if err != nil {
if os.IsNotExist(err) {
return nil
}
return fmt.Errorf("read runtime metadata snapshot: %w", err)
}
if len(strings.TrimSpace(string(data))) == 0 {
return nil
}
var snapshot runtimeSnapshot
if err := json.Unmarshal(data, &snapshot); err != nil {
return fmt.Errorf("decode runtime metadata snapshot: %w", err)
}
store.loadRuntimeSnapshot(snapshot)
return nil
}
func (store *FileStore) persistRuntime() error {
store.runtimePersistMu.Lock()
defer store.runtimePersistMu.Unlock()
snapshot := store.runtimeSnapshot()
data, err := json.Marshal(snapshot)
if err != nil {
return fmt.Errorf("encode runtime metadata snapshot: %w", err)
}
if err := os.MkdirAll(filepath.Dir(store.runtimePath), 0o755); err != nil {
return fmt.Errorf("create runtime metadata directory: %w", err)
}
if err := store.ensureMetadataFile(); err != nil {
return err
}
tmpPath := store.runtimePath + ".tmp"
if err := os.WriteFile(tmpPath, data, 0o600); err != nil {
return fmt.Errorf("write runtime metadata snapshot: %w", err)
}
if err := os.Rename(tmpPath, store.runtimePath); err != nil {
return fmt.Errorf("replace runtime metadata snapshot: %w", err)
}
return nil
}
func (store *FileStore) ensureMetadataFile() error {
if _, err := os.Stat(store.path); err == nil {
return nil
} else if !os.IsNotExist(err) {
return fmt.Errorf("stat metadata snapshot: %w", err)
}
store.persistMu.Lock()
defer store.persistMu.Unlock()
if _, err := os.Stat(store.path); err == nil {
return nil
} else if !os.IsNotExist(err) {
return fmt.Errorf("stat metadata snapshot: %w", err)
}
if err := os.MkdirAll(filepath.Dir(store.path), 0o755); err != nil {
return fmt.Errorf("create metadata directory: %w", err)
}
tmpPath := store.path + ".empty.tmp"
if err := os.WriteFile(tmpPath, []byte("{}"), 0o600); err != nil {
return fmt.Errorf("write empty metadata snapshot: %w", err)
}
if err := os.Rename(tmpPath, store.path); err != nil {
return fmt.Errorf("replace empty metadata snapshot: %w", err)
}
return nil
}
func runtimeMetadataPath(path string) string {
ext := filepath.Ext(path)
if ext == "" {
return path + ".runtime"
}
return strings.TrimSuffix(path, ext) + ".runtime" + ext
}
func (store *FileStore) snapshot() StoreSnapshot {
return normalizeStoreSnapshot(StoreSnapshot{
Users: snapshotRepository(store.MemoryStore.users),
AuthSessions: snapshotRepository(store.MemoryStore.authSessions),
RunControlSessions: snapshotRepository(store.MemoryStore.runSessions),
AIProviders: snapshotRepository(store.MemoryStore.aiProviders),
GamePlugins: snapshotRepository(store.MemoryStore.gamePlugins),
ServerInstances: snapshotRepository(store.MemoryStore.serverInstances),
RunEndpoints: snapshotRepository(store.MemoryStore.runEndpoints),
Jobs: snapshotRepository(store.MemoryStore.jobs.memoryRepository),
Artifacts: snapshotRepository(store.MemoryStore.artifacts),
RuntimeBindings: snapshotRepository(store.MemoryStore.runtimeBindings),
EncryptedComponentKeys: snapshotRepository(store.MemoryStore.componentKeys),
RunDistributions: snapshotRepository(store.MemoryStore.runDists),
DependencyStatuses: snapshotRepository(store.MemoryStore.dependencies),
RunUpdateJobs: snapshotRepository(store.MemoryStore.updateJobs),
LogStreams: snapshotRepository(store.MemoryStore.logStreams),
MetricSamples: snapshotRepository(store.MemoryStore.metricSamples),
Backups: snapshotRepository(store.MemoryStore.backups),
PluginLifecycles: snapshotRepository(store.MemoryStore.pluginLifecycle),
GameClientBridgeCommands: snapshotRepository(store.MemoryStore.bridgeCommands.memoryRepository),
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),
})
}
func (store *FileStore) runtimeSnapshot() runtimeSnapshot {
return runtimeSnapshot{
RunControlSessions: snapshotRepository(store.MemoryStore.runSessions),
RunEndpoints: snapshotRepository(store.MemoryStore.runEndpoints),
LogStreams: snapshotRepository(store.MemoryStore.logStreams),
}
}
func (store *FileStore) loadSnapshot(snapshot StoreSnapshot) {
snapshot = normalizeStoreSnapshot(snapshot)
loadRepository(store.MemoryStore.users, snapshot.Users)
loadRepository(store.MemoryStore.authSessions, snapshot.AuthSessions)
loadRepository(store.MemoryStore.runSessions, snapshot.RunControlSessions)
loadRepository(store.MemoryStore.aiProviders, snapshot.AIProviders)
loadRepository(store.MemoryStore.gamePlugins, snapshot.GamePlugins)
loadRepository(store.MemoryStore.serverInstances, snapshot.ServerInstances)
loadRepository(store.MemoryStore.runEndpoints, snapshot.RunEndpoints)
loadRepository(store.MemoryStore.jobs.memoryRepository, snapshot.Jobs)
loadRepository(store.MemoryStore.artifacts, snapshot.Artifacts)
loadRepository(store.MemoryStore.runtimeBindings, snapshot.RuntimeBindings)
loadRepository(store.MemoryStore.componentKeys, snapshot.EncryptedComponentKeys)
loadRepository(store.MemoryStore.runDists, snapshot.RunDistributions)
loadRepository(store.MemoryStore.dependencies, snapshot.DependencyStatuses)
loadRepository(store.MemoryStore.updateJobs, snapshot.RunUpdateJobs)
loadRepository(store.MemoryStore.logStreams, snapshot.LogStreams)
loadRepository(store.MemoryStore.metricSamples, snapshot.MetricSamples)
loadRepository(store.MemoryStore.backups, snapshot.Backups)
loadRepository(store.MemoryStore.pluginLifecycle, snapshot.PluginLifecycles)
loadRepository(store.MemoryStore.bridgeCommands.memoryRepository, snapshot.GameClientBridgeCommands)
loadRepository(store.MemoryStore.bridgeSnapshots.memoryRepository, snapshot.GameClientBridgeSnapshots)
loadRepository(store.MemoryStore.bridgeStreams, snapshot.GameClientBridgeStreams)
loadRepository(store.MemoryStore.pluginDataRecords, snapshot.PluginDataRecords)
loadRepository(store.MemoryStore.scumUsers, snapshot.SCUMUsers)
loadRepository(store.MemoryStore.scumUserTracks, snapshot.SCUMUserTrajectories)
loadRepository(store.MemoryStore.scumVehicles, snapshot.SCUMVehicles)
loadRepository(store.MemoryStore.scumVehicleTracks, snapshot.SCUMVehicleTrajectories)
loadRepository(store.MemoryStore.scumVehicleLocks, snapshot.SCUMVehicleLocks)
}
func (store *FileStore) loadRuntimeSnapshot(snapshot runtimeSnapshot) {
loadRepository(store.MemoryStore.runSessions, snapshot.RunControlSessions)
loadRepository(store.MemoryStore.runEndpoints, snapshot.RunEndpoints)
loadRepository(store.MemoryStore.logStreams, snapshot.LogStreams)
}
type mutableRepository[T any, F any] interface {
Create(T) error
Get(string) (T, error)
List(F) ([]T, error)
Update(T) error
Delete(string) error
Apply([]T, []string) error
}
type persistentRepository[T any, F any] struct {
repository mutableRepository[T, F]
persist func() error
}
func (repository *persistentRepository[T, F]) Create(value T) error {
if err := repository.repository.Create(value); err != nil {
return err
}
return repository.persist()
}
func (repository *persistentRepository[T, F]) Get(id string) (T, error) {
return repository.repository.Get(id)
}
func (repository *persistentRepository[T, F]) List(filter F) ([]T, error) {
return repository.repository.List(filter)
}
func (repository *persistentRepository[T, F]) Update(value T) error {
if err := repository.repository.Update(value); err != nil {
return err
}
return repository.persist()
}
func (repository *persistentRepository[T, F]) Delete(id string) error {
if err := repository.repository.Delete(id); err != nil {
return err
}
return repository.persist()
}
func (repository *persistentRepository[T, F]) Apply(upserts []T, deleteIDs []string) error {
if err := repository.repository.Apply(upserts, deleteIDs); err != nil {
return err
}
return repository.persist()
}
type persistentJobRepository struct {
*persistentRepository[domain.Job, domain.JobFilter]
repository JobRepository
}
func (repository *persistentJobRepository) GetByIdempotency(runEndpointID string, idempotencyKey string) (domain.Job, error) {
return repository.repository.GetByIdempotency(runEndpointID, idempotencyKey)
}
func snapshotRepository[T any, F any](repository *memoryRepository[T, F]) []T {
repository.mu.RLock()
defer repository.mu.RUnlock()
ids := make([]string, 0, len(repository.byID))
for id := range repository.byID {
ids = append(ids, id)
}
sort.Strings(ids)
values := make([]T, 0, len(ids))
for _, id := range ids {
values = append(values, repository.copyOf(repository.byID[id]))
}
return values
}
func loadRepository[T any, F any](repository *memoryRepository[T, F], values []T) {
repository.mu.Lock()
defer repository.mu.Unlock()
repository.byID = map[string]T{}
for _, value := range values {
repository.byID[repository.idOf(value)] = repository.copyOf(value)
}
}