121 lines
3.5 KiB
Go
121 lines
3.5 KiB
Go
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() }
|