// Package manager coordinates encrypted file mutations and provider reconciliation. package manager import ( "context" "crypto/sha256" "errors" "fmt" "os" "sort" "sync" "time" "cloud.campbellwireless.net/git/barkstack/treatvault/internal/provider" "cloud.campbellwireless.net/git/barkstack/treatvault/internal/vault" ) type Provider interface { Sync(context.Context, vault.Snapshot) ([]provider.SecretState, error) InUse(context.Context, string) (bool, error) } type Secret struct { Name string `json:"name"` InUse bool `json:"inUse"` Synced bool `json:"synced"` } type Status struct { State string `json:"state"` Secrets []Secret `json:"secrets"` LastSync time.Time `json:"lastSync,omitempty"` Error string `json:"error,omitempty"` } type Manager struct { store *vault.Store provider Provider pollInterval time.Duration operationMu sync.Mutex statusMu sync.RWMutex status Status } func New(store *vault.Store, secretProvider Provider, pollInterval time.Duration) *Manager { if pollInterval <= 0 { pollInterval = 500 * time.Millisecond } return &Manager{ store: store, provider: secretProvider, pollInterval: pollInterval, status: Status{State: "starting", Secrets: []Secret{}}, } } func (m *Manager) Status() Status { m.statusMu.RLock() defer m.statusMu.RUnlock() status := m.status status.Secrets = append([]Secret(nil), m.status.Secrets...) return status } func (m *Manager) Sync(ctx context.Context) error { m.operationMu.Lock() defer m.operationMu.Unlock() snapshot, err := m.store.Load() if err != nil { m.setFailure(err) return err } return m.syncSnapshot(ctx, snapshot) } func (m *Manager) Set(ctx context.Context, name string, value []byte) error { m.operationMu.Lock() defer m.operationMu.Unlock() snapshot, err := m.store.Set(name, value) if err != nil { return err } return m.syncSnapshot(ctx, snapshot) } func (m *Manager) Delete(ctx context.Context, name string) error { m.operationMu.Lock() defer m.operationMu.Unlock() inUse, err := m.provider.InUse(ctx, name) if err != nil { return fmt.Errorf("check secret consumers: %w", err) } if inUse { return fmt.Errorf("%w: %s", provider.ErrSecretInUse, name) } snapshot, err := m.store.Delete(name) if err != nil { return err } return m.syncSnapshot(ctx, snapshot) } func (m *Manager) Run(ctx context.Context) { _ = m.Sync(ctx) fingerprint, _ := fileFingerprint(m.store.Path()) ticker := time.NewTicker(m.pollInterval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: next, err := fileFingerprint(m.store.Path()) status := m.Status() if err != nil { m.setFailure(err) continue } if next == fingerprint && status.Error == "" { continue } fingerprint = next _ = m.Sync(ctx) } } } func (m *Manager) syncSnapshot(ctx context.Context, snapshot vault.Snapshot) error { states, err := m.provider.Sync(ctx, snapshot) secrets := make([]Secret, 0, len(snapshot.Secrets)) stateByName := make(map[string]provider.SecretState, len(states)) for _, state := range states { stateByName[state.Name] = state } for name := range snapshot.Secrets { state, synced := stateByName[name] secrets = append(secrets, Secret{Name: name, InUse: state.InUse, Synced: synced && err == nil}) } sort.Slice(secrets, func(i, j int) bool { return secrets[i].Name < secrets[j].Name }) status := Status{State: "ready", Secrets: secrets, LastSync: time.Now().UTC()} if err != nil { status.State = "degraded" status.Error = err.Error() } m.statusMu.Lock() m.status = status m.statusMu.Unlock() return err } func (m *Manager) setFailure(err error) { if err == nil { return } m.statusMu.Lock() m.status.State = "degraded" m.status.Error = err.Error() m.statusMu.Unlock() } func fileFingerprint(path string) ([sha256.Size]byte, error) { contents, err := os.ReadFile(path) if err != nil { return [sha256.Size]byte{}, fmt.Errorf("read encrypted file: %w", err) } return sha256.Sum256(contents), nil } func IsSecretInUse(err error) bool { return errors.Is(err, provider.ErrSecretInUse) }