feat: TreatVault encrypted secret manager with Docker Swarm sync and console plugin
This commit is contained in:
170
internal/manager/manager.go
Normal file
170
internal/manager/manager.go
Normal file
@@ -0,0 +1,170 @@
|
||||
// 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)
|
||||
}
|
||||
99
internal/manager/manager_test.go
Normal file
99
internal/manager/manager_test.go
Normal file
@@ -0,0 +1,99 @@
|
||||
package manager
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"cloud.campbellwireless.net/git/barkstack/treatvault/internal/provider"
|
||||
"cloud.campbellwireless.net/git/barkstack/treatvault/internal/vault"
|
||||
"filippo.io/age"
|
||||
)
|
||||
|
||||
func TestRunSynchronizesExternalEncryptedFileChanges(t *testing.T) {
|
||||
identity, err := age.GenerateX25519Identity()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
path := t.TempDir() + "/secrets.age"
|
||||
if err := vault.Initialize(path, identity); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
store := vault.New(path, identity)
|
||||
provider := &recordingProvider{synced: make(chan vault.Snapshot, 4)}
|
||||
manager := New(store, provider, 10*time.Millisecond)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
go manager.Run(ctx)
|
||||
|
||||
initial := waitSnapshot(t, provider.synced)
|
||||
if len(initial.Secrets) != 0 {
|
||||
t.Fatalf("initial secrets = %#v", initial.Secrets)
|
||||
}
|
||||
if _, err := store.Set("database_password", []byte("changed outside manager")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
changed := waitSnapshot(t, provider.synced)
|
||||
if string(changed.Secrets["database_password"].Value) != "changed outside manager" {
|
||||
t.Fatalf("changed snapshot = %#v", changed)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteRefusesSecretInUse(t *testing.T) {
|
||||
identity, err := age.GenerateX25519Identity()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
path := t.TempDir() + "/secrets.age"
|
||||
if err := vault.Initialize(path, identity); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
store := vault.New(path, identity)
|
||||
if _, err := store.Set("database_password", []byte("value")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
manager := New(store, &recordingProvider{inUse: true}, time.Second)
|
||||
if err := manager.Delete(context.Background(), "database_password"); !IsSecretInUse(err) {
|
||||
t.Fatalf("Delete() error = %v", err)
|
||||
}
|
||||
loaded, err := store.Load()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, exists := loaded.Secrets["database_password"]; !exists {
|
||||
t.Fatal("in-use secret was removed from encrypted file")
|
||||
}
|
||||
}
|
||||
|
||||
type recordingProvider struct {
|
||||
mu sync.Mutex
|
||||
synced chan vault.Snapshot
|
||||
inUse bool
|
||||
}
|
||||
|
||||
func (p *recordingProvider) Sync(_ context.Context, snapshot vault.Snapshot) ([]provider.SecretState, error) {
|
||||
if p.synced != nil {
|
||||
p.synced <- snapshot
|
||||
}
|
||||
states := make([]provider.SecretState, 0, len(snapshot.Secrets))
|
||||
for name, record := range snapshot.Secrets {
|
||||
states = append(states, provider.SecretState{Name: name, Revision: record.Revision, InUse: p.inUse})
|
||||
}
|
||||
return states, nil
|
||||
}
|
||||
|
||||
func (p *recordingProvider) InUse(context.Context, string) (bool, error) {
|
||||
return p.inUse, nil
|
||||
}
|
||||
|
||||
func waitSnapshot(t *testing.T, snapshots <-chan vault.Snapshot) vault.Snapshot {
|
||||
t.Helper()
|
||||
select {
|
||||
case snapshot := <-snapshots:
|
||||
return snapshot
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("timed out waiting for synchronization")
|
||||
return vault.Snapshot{}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user