mirror of
https://github.com/psviderski/uncloud.git
synced 2026-10-06 13:18:58 +00:00
feat(caddystorage): simplify active locks tracking and cleanup
This commit is contained in:
1 parent
9a4098694c
commit
6ab02f95fe
2 files changed
+85
-87
No files matched your search
+41
-48
@@ -12,7 +12,10 @@ import (
|
|||||||
"github.com/psviderski/uncloud/pkg/distlock"
|
"github.com/psviderski/uncloud/pkg/distlock"
|
||||||
)
|
)
|
||||||
|
|
||||||
const lockResourcePrefix = "caddy_storage:"
|
const (
|
||||||
|
lockPrefix = "caddy_storage:"
|
||||||
|
storeReplicationTimeout = 10 * time.Second
|
||||||
|
)
|
||||||
|
|
||||||
// Lock acquires an automatically renewed distributed lock and waits for the local store to catch up with versions
|
// Lock acquires an automatically renewed distributed lock and waits for the local store to catch up with versions
|
||||||
// observed on responding machines. Writers using the same lock can then read locally. Unavailable machines may have
|
// observed on responding machines. Writers using the same lock can then read locally. Unavailable machines may have
|
||||||
@@ -21,14 +24,23 @@ func (s *Storage) Lock(ctx context.Context, name string) (err error) {
|
|||||||
if name == "" {
|
if name == "" {
|
||||||
return errors.New("lock name is empty")
|
return errors.New("lock name is empty")
|
||||||
}
|
}
|
||||||
s.locksMu.Lock()
|
|
||||||
if s.locks == nil {
|
// Count the call first. If Cleanup has already observed zero, Caddy has cancelled s.ctx and the check below rejects
|
||||||
s.locksMu.Unlock()
|
// this call before it uses the client.
|
||||||
|
s.lockOps.Add(1)
|
||||||
|
// A failed Lock owns its lifecycle through any lease rollback. A successful Lock transfers that responsibility to
|
||||||
|
// Unlock, which keeps the client open until its release attempt finishes.
|
||||||
|
defer func() {
|
||||||
|
if err != nil {
|
||||||
|
s.lockOps.Add(-1)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
if s.ctx.Err() != nil {
|
||||||
return errors.New("storage is closed")
|
return errors.New("storage is closed")
|
||||||
}
|
}
|
||||||
s.locksMu.Unlock()
|
|
||||||
|
|
||||||
// Stop acquisition retries when Caddy unloads the module, even if the caller's context is still active.
|
// Stop acquisition retries when Caddy unloads the module (cancels s.ctx), even if the caller's context
|
||||||
|
// is still active.
|
||||||
ctx, cancel := context.WithCancelCause(ctx)
|
ctx, cancel := context.WithCancelCause(ctx)
|
||||||
defer cancel(nil)
|
defer cancel(nil)
|
||||||
stopOnCleanup := context.AfterFunc(s.ctx, func() {
|
stopOnCleanup := context.AfterFunc(s.ctx, func() {
|
||||||
@@ -38,70 +50,50 @@ func (s *Storage) Lock(ctx context.Context, name string) (err error) {
|
|||||||
|
|
||||||
log := s.log.With("lock", name)
|
log := s.log.With("lock", name)
|
||||||
started := time.Now()
|
started := time.Now()
|
||||||
stage := "acquire_lease"
|
log.Debug("acquiring lock", "ttl", time.Duration(s.LockTTL))
|
||||||
log.Debug("acquiring lock", "lock_ttl", time.Duration(s.LockTTL))
|
lease, err := s.locker.Acquire(ctx, lockPrefix+name)
|
||||||
defer func() {
|
|
||||||
if err != nil {
|
|
||||||
log.Debug("failed to acquire lock",
|
|
||||||
"stage", stage, "duration", time.Since(started), "error", err)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
lease, err := s.locker.Acquire(ctx, lockResourcePrefix+name)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
log.Debug("failed to acquire lock", "duration", time.Since(started), "error", err)
|
||||||
return fmt.Errorf("acquire lock '%s': %w", name, err)
|
return fmt.Errorf("acquire lock '%s': %w", name, err)
|
||||||
}
|
}
|
||||||
log.Debug("lock lease acquired", "duration", time.Since(started))
|
log.Debug("lock lease acquired", "duration", time.Since(started))
|
||||||
// Keep observing after Lock returns so lease loss during protected work remains visible.
|
// Release the lease if the lock acquisition fails after this point.
|
||||||
context.AfterFunc(lease.Context(), func() {
|
// Unlock will take care of releasing the lease on success.
|
||||||
if cause := context.Cause(lease.Context()); errors.Is(cause, distlock.ErrLeaseLost) {
|
|
||||||
log.Error("lock lease lost", "error", cause)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
|
|
||||||
defer func() {
|
defer func() {
|
||||||
if err == nil {
|
if err == nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// Capture lease loss before Release cancels the lease context itself.
|
if releaseErr := s.releaseLock(ctx, name, lease, "cancelled acquisition"); releaseErr != nil {
|
||||||
err = errors.Join(err, context.Cause(lease.Context()))
|
err = errors.Join(err, releaseErr)
|
||||||
err = fmt.Errorf("acquire lock '%s': %w", name,
|
}
|
||||||
errors.Join(err, s.releaseLock(ctx, name, lease, "failed acquisition")))
|
|
||||||
}()
|
}()
|
||||||
|
|
||||||
stopOnLeaseLoss := context.AfterFunc(lease.Context(), func() {
|
// Catch up with the latest store versions observed on responding machines to increase the chance of reading
|
||||||
cancel(context.Cause(lease.Context()))
|
// the latest writes on them locally.
|
||||||
})
|
|
||||||
defer stopOnLeaseLoss()
|
|
||||||
|
|
||||||
stage = "collect_store_versions"
|
|
||||||
version, machines, err := s.clusterStoreVersion(ctx, log)
|
version, machines, err := s.clusterStoreVersion(ctx, log)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
stage = "wait_for_replication"
|
|
||||||
waitStarted := time.Now()
|
waitStarted := time.Now()
|
||||||
log.Debug("waiting for local store replication", "machine_names", machines, "store_version", version)
|
log.Debug("waiting for local store replication", "machine_names", machines, "store_version", version)
|
||||||
if err := s.client.WaitForStoreVersion(ctx, version); err != nil {
|
waitCtx, cancelWait := context.WithTimeout(ctx, storeReplicationTimeout)
|
||||||
|
err = s.client.WaitForStoreVersion(waitCtx, version)
|
||||||
|
cancelWait()
|
||||||
|
if err != nil {
|
||||||
return fmt.Errorf("wait for local store replication: %w", err)
|
return fmt.Errorf("wait for local store replication: %w", err)
|
||||||
}
|
}
|
||||||
log.Debug("local store replication complete", "duration", time.Since(waitStarted))
|
log.Debug("local store replication complete", "duration", time.Since(waitStarted))
|
||||||
|
|
||||||
stage = "register_lock"
|
|
||||||
s.locksMu.Lock()
|
s.locksMu.Lock()
|
||||||
defer s.locksMu.Unlock()
|
defer s.locksMu.Unlock()
|
||||||
if err := context.Cause(ctx); err != nil {
|
if s.ctx.Err() != nil {
|
||||||
return err
|
|
||||||
}
|
|
||||||
if err := context.Cause(lease.Context()); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if s.locks == nil {
|
|
||||||
// A lease obtained during cleanup must be released instead of reopening the module's lock map.
|
|
||||||
return errors.New("storage is closed")
|
return errors.New("storage is closed")
|
||||||
}
|
}
|
||||||
|
if lost := context.Cause(lease.Context()); lost != nil {
|
||||||
|
return lost
|
||||||
|
}
|
||||||
if _, exists := s.locks[name]; exists {
|
if _, exists := s.locks[name]; exists {
|
||||||
return errors.New("lock is already held by this storage instance")
|
return errors.New("lock is already tracked by this storage instance")
|
||||||
}
|
}
|
||||||
s.locks[name] = lease
|
s.locks[name] = lease
|
||||||
|
|
||||||
@@ -146,6 +138,8 @@ func (s *Storage) Unlock(ctx context.Context, name string) error {
|
|||||||
if !exists {
|
if !exists {
|
||||||
return fmt.Errorf("lock '%s' is not held by this storage instance", name)
|
return fmt.Errorf("lock '%s' is not held by this storage instance", name)
|
||||||
}
|
}
|
||||||
|
defer s.lockOps.Add(-1)
|
||||||
|
|
||||||
// Release stops renewal even on error. Any nodes that cannot be reached will let the lease expire.
|
// Release stops renewal even on error. Any nodes that cannot be reached will let the lease expire.
|
||||||
if err := s.releaseLock(ctx, name, lease, "unlock"); err != nil {
|
if err := s.releaseLock(ctx, name, lease, "unlock"); err != nil {
|
||||||
return fmt.Errorf("release lock '%s': %w", name, err)
|
return fmt.Errorf("release lock '%s': %w", name, err)
|
||||||
@@ -154,11 +148,10 @@ func (s *Storage) Unlock(ctx context.Context, name string) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// releaseLock logs releases consistently across unlock, failed acquisition, and module cleanup.
|
// releaseLock releases the lease with a context that ignores cancellation and logs the release attempt.
|
||||||
func (s *Storage) releaseLock(ctx context.Context, name string, lease *distlock.Lease, reason string) error {
|
func (s *Storage) releaseLock(ctx context.Context, name string, lease *distlock.Lease, reason string) error {
|
||||||
// Unlock and rollback must attempt node cleanup even if the caller or module has already been cancelled.
|
// Unlock and rollback must attempt node cleanup even if the caller or module has already been cancelled.
|
||||||
ctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), distlock.DefaultMaxNodeCallTimeout)
|
ctx = context.WithoutCancel(ctx)
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
log := s.log.With("lock", name, "reason", reason)
|
log := s.log.With("lock", name, "reason", reason)
|
||||||
started := time.Now()
|
started := time.Now()
|
||||||
|
|||||||
+44
-39
@@ -7,6 +7,7 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"sync"
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/caddyserver/caddy/v2"
|
"github.com/caddyserver/caddy/v2"
|
||||||
@@ -24,6 +25,8 @@ const (
|
|||||||
DefaultSocketPath = "/run/uncloud/uncloud.sock"
|
DefaultSocketPath = "/run/uncloud/uncloud.sock"
|
||||||
// DefaultLockTTL is the default duration of a distributed lock lease.
|
// DefaultLockTTL is the default duration of a distributed lock lease.
|
||||||
DefaultLockTTL = 20 * time.Second
|
DefaultLockTTL = 20 * time.Second
|
||||||
|
// lockCleanupTimeout bounds how long an unloaded module waits for active lock operations when cleaning up.
|
||||||
|
lockCleanupTimeout = 5 * time.Minute
|
||||||
)
|
)
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
@@ -46,8 +49,15 @@ type Storage struct {
|
|||||||
locker *distlock.Locker
|
locker *distlock.Locker
|
||||||
log *slog.Logger
|
log *slog.Logger
|
||||||
|
|
||||||
|
// ctx is the module context from Provision. Caddy cancels it when it unloads the module and calls Cleanup.
|
||||||
|
ctx context.Context
|
||||||
|
|
||||||
|
// locksMu protects locks.
|
||||||
locksMu sync.Mutex
|
locksMu sync.Mutex
|
||||||
|
// locks maps successfully acquired lock names to held leases.
|
||||||
locks map[string]*distlock.Lease
|
locks map[string]*distlock.Lease
|
||||||
|
// lockOps counts Lock calls until they fail or their locks are released.
|
||||||
|
lockOps atomic.Int64
|
||||||
}
|
}
|
||||||
|
|
||||||
// CaddyModule returns the Caddy module information.
|
// CaddyModule returns the Caddy module information.
|
||||||
@@ -60,6 +70,7 @@ func (*Storage) CaddyModule() caddy.ModuleInfo {
|
|||||||
|
|
||||||
// Provision connects the storage to the local Uncloud API and initialises the distributed locker.
|
// Provision connects the storage to the local Uncloud API and initialises the distributed locker.
|
||||||
func (s *Storage) Provision(ctx caddy.Context) error {
|
func (s *Storage) Provision(ctx caddy.Context) error {
|
||||||
|
s.ctx = ctx
|
||||||
s.log = ctx.Slogger()
|
s.log = ctx.Slogger()
|
||||||
|
|
||||||
if s.Socket == "" {
|
if s.Socket == "" {
|
||||||
@@ -93,50 +104,44 @@ func (s *Storage) Provision(ctx caddy.Context) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Cleanup releases active locks and closes the Uncloud API connection.
|
// Cleanup keeps the Uncloud API connection open so acquired lock leases can renew until Caddy releases them. It closes
|
||||||
func (s *Storage) Cleanup() (err error) {
|
// the connection after all lock operations finish or the cleanup timeout expires.
|
||||||
s.locksMu.Lock()
|
func (s *Storage) Cleanup() error {
|
||||||
locks := s.locks
|
if s.client == nil {
|
||||||
s.locks = nil
|
// Provision failed before opening the storage.
|
||||||
s.locksMu.Unlock()
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
started := time.Now()
|
started := time.Now()
|
||||||
s.log.Debug("cleaning up module", "locks", len(locks))
|
s.log.Debug("cleaning up module", "locks", s.lockOps.Load())
|
||||||
defer func() {
|
|
||||||
if err != nil {
|
go func() {
|
||||||
s.log.Debug("failed to clean up module", "duration", time.Since(started), "error", err)
|
timer := time.NewTimer(lockCleanupTimeout)
|
||||||
} else {
|
defer timer.Stop()
|
||||||
s.log.Debug("module cleanup complete", "duration", time.Since(started))
|
ticker := time.NewTicker(1 * time.Second)
|
||||||
|
defer ticker.Stop()
|
||||||
|
|
||||||
|
timedOut := false
|
||||||
|
for !timedOut && s.lockOps.Load() > 0 {
|
||||||
|
select {
|
||||||
|
case <-ticker.C:
|
||||||
|
case <-timer.C:
|
||||||
|
timedOut = true
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
if remaining := s.lockOps.Load(); timedOut && remaining > 0 {
|
||||||
|
s.log.Debug("timed out waiting for active locks to be unlocked",
|
||||||
|
"locks", remaining, "timeout", lockCleanupTimeout)
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := s.client.Close(); err != nil {
|
||||||
|
s.log.Debug("failed to clean up module", "duration", time.Since(started), "error", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
s.log.Debug("module cleanup complete", "duration", time.Since(started))
|
||||||
}()
|
}()
|
||||||
|
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), distlock.DefaultMaxNodeCallTimeout)
|
return nil
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
errCh := make(chan error, len(locks))
|
|
||||||
var wg sync.WaitGroup
|
|
||||||
for name, lease := range locks {
|
|
||||||
wg.Go(func() {
|
|
||||||
if err := s.releaseLock(ctx, name, lease, "cleanup"); err != nil {
|
|
||||||
errCh <- fmt.Errorf("release lock '%s': %w", name, err)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
}
|
|
||||||
wg.Wait()
|
|
||||||
close(errCh)
|
|
||||||
|
|
||||||
errs := make([]error, 0, len(errCh)+1)
|
|
||||||
for err := range errCh {
|
|
||||||
errs = append(errs, err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if s.client != nil {
|
|
||||||
errs = append(errs, s.client.Close())
|
|
||||||
s.client = nil
|
|
||||||
s.locker = nil
|
|
||||||
}
|
|
||||||
|
|
||||||
return errors.Join(errs...)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// CertMagicStorage returns the provisioned CertMagic storage implementation.
|
// CertMagicStorage returns the provisioned CertMagic storage implementation.
|
||||||
|
|||||||
Reference in new issue
Block a user