package purge import ( "context" "errors" "log/slog" "sync" "github.com/klarkxy/nekonest-cloud/relay/internal/controlplane" "github.com/klarkxy/nekonest-cloud/relay/internal/tenantpurge" ) type Registry interface { QuiesceForPurge(tenantID string, generation int64) error } type ControlPlane interface { AdvancePurge(context.Context, controlplane.PurgeAdvance) error } type Config struct { DataRoot string BackupRoot string Registry Registry ControlPlane ControlPlane MaxConcurrent int Logger *slog.Logger } type Manager struct { config Config mu sync.Mutex inFlight map[string]struct{} semaphore chan struct{} wait sync.WaitGroup } func New(config Config) (*Manager, error) { if config.DataRoot == "" || config.BackupRoot == "" || config.Registry == nil || config.ControlPlane == nil { return nil, errors.New("purge manager requires data, backup, registry, and control-plane ports") } if config.MaxConcurrent <= 0 { config.MaxConcurrent = 1 } if config.MaxConcurrent > 4 { return nil, errors.New("purge concurrency is unreasonably high") } if config.Logger == nil { config.Logger = slog.Default() } return &Manager{ config: config, inFlight: make(map[string]struct{}), semaphore: make(chan struct{}, config.MaxConcurrent), }, nil } func (manager *Manager) Handle(ctx context.Context, assignments []controlplane.PurgeAssignment) { for _, assignment := range assignments { manager.mu.Lock() if _, exists := manager.inFlight[assignment.PurgeID]; exists { manager.mu.Unlock() continue } select { case manager.semaphore <- struct{}{}: manager.inFlight[assignment.PurgeID] = struct{}{} manager.wait.Add(1) manager.mu.Unlock() go func(assignment controlplane.PurgeAssignment) { defer func() { <-manager.semaphore manager.mu.Lock() delete(manager.inFlight, assignment.PurgeID) manager.mu.Unlock() manager.wait.Done() }() if err := manager.process(ctx, assignment); err != nil && ctx.Err() == nil { manager.config.Logger.Error("Relay tenant purge failed", "purge_id", assignment.PurgeID, "tenant_id", assignment.TenantID, "error", err) } }(assignment) default: manager.mu.Unlock() } } } func (manager *Manager) process(ctx context.Context, assignment controlplane.PurgeAssignment) error { if assignment.PurgeID == "" || assignment.TenantID == "" || assignment.PlacementGeneration < 1 { return errors.New("purge assignment is incomplete") } if err := manager.config.Registry.QuiesceForPurge( assignment.TenantID, assignment.PlacementGeneration, ); err != nil { return manager.fail(ctx, assignment, "relay_purge_quiesce_failed", err) } result, err := tenantpurge.Purge( ctx, manager.config.DataRoot, manager.config.BackupRoot, assignment.TenantID, assignment.PlacementGeneration, ) if err != nil { return manager.fail(ctx, assignment, "relay_purge_delete_failed", err) } return manager.config.ControlPlane.AdvancePurge(ctx, controlplane.PurgeAdvance{ PurgeID: assignment.PurgeID, Action: "completed", EvidenceSHA256: result.EvidenceSHA256, }) } func (manager *Manager) fail( ctx context.Context, assignment controlplane.PurgeAssignment, code string, cause error, ) error { if err := manager.config.ControlPlane.AdvancePurge(ctx, controlplane.PurgeAdvance{ PurgeID: assignment.PurgeID, Action: "failed", ErrorCode: code, }); err != nil { return errors.Join(cause, err) } return cause } func (manager *Manager) Wait() { manager.wait.Wait() }