mirror of
https://github.com/psviderski/uncloud.git
synced 2026-10-08 22:24:54 +00:00
chore(caddystorage): move storage module code to a separate unlabs-dev/caddy-uncloud repo
This commit is contained in:
1 parent
930e7ea637
commit
4dae187aa5
3 files changed
-458
No files matched your search
@@ -1,164 +0,0 @@
|
|||||||
package caddystorage
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"errors"
|
|
||||||
"fmt"
|
|
||||||
"log/slog"
|
|
||||||
"slices"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/psviderski/uncloud/pkg/api"
|
|
||||||
"github.com/psviderski/uncloud/pkg/client"
|
|
||||||
"github.com/psviderski/uncloud/pkg/distlock"
|
|
||||||
)
|
|
||||||
|
|
||||||
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
|
|
||||||
// observed on responding machines. Writers using the same lock can then read locally. Unavailable machines may have
|
|
||||||
// writes that this wait does not cover, and reads outside a lock remain eventually consistent.
|
|
||||||
func (s *Storage) Lock(ctx context.Context, name string) (err error) {
|
|
||||||
if name == "" {
|
|
||||||
return errors.New("lock name is empty")
|
|
||||||
}
|
|
||||||
|
|
||||||
// Count the call first. If Cleanup has already observed zero, Caddy has cancelled s.ctx and the check below rejects
|
|
||||||
// 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")
|
|
||||||
}
|
|
||||||
|
|
||||||
// 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)
|
|
||||||
defer cancel(nil)
|
|
||||||
stopOnCleanup := context.AfterFunc(s.ctx, func() {
|
|
||||||
cancel(context.Cause(s.ctx))
|
|
||||||
})
|
|
||||||
defer stopOnCleanup()
|
|
||||||
|
|
||||||
log := s.log.With("lock", name)
|
|
||||||
started := time.Now()
|
|
||||||
log.Debug("acquiring lock", "ttl", time.Duration(s.LockTTL))
|
|
||||||
lease, err := s.locker.Acquire(ctx, lockPrefix+name)
|
|
||||||
if err != nil {
|
|
||||||
log.Debug("failed to acquire lock", "duration", time.Since(started), "error", err)
|
|
||||||
return fmt.Errorf("acquire lock '%s': %w", name, err)
|
|
||||||
}
|
|
||||||
log.Debug("lock lease acquired", "duration", time.Since(started))
|
|
||||||
// Release the lease if the lock acquisition fails after this point.
|
|
||||||
// Unlock will take care of releasing the lease on success.
|
|
||||||
defer func() {
|
|
||||||
if err == nil {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
if releaseErr := s.releaseLock(ctx, name, lease, "cancelled acquisition"); releaseErr != nil {
|
|
||||||
err = errors.Join(err, releaseErr)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
// Catch up with the latest store versions observed on responding machines to increase the chance of reading
|
|
||||||
// the latest writes on them locally.
|
|
||||||
version, machines, err := s.clusterStoreVersion(ctx, log)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
waitStarted := time.Now()
|
|
||||||
log.Debug("waiting for local store replication", "machine_names", machines, "store_version", version)
|
|
||||||
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)
|
|
||||||
}
|
|
||||||
log.Debug("local store replication complete", "duration", time.Since(waitStarted))
|
|
||||||
|
|
||||||
s.locksMu.Lock()
|
|
||||||
defer s.locksMu.Unlock()
|
|
||||||
if s.ctx.Err() != nil {
|
|
||||||
return errors.New("storage is closed")
|
|
||||||
}
|
|
||||||
if lost := context.Cause(lease.Context()); lost != nil {
|
|
||||||
return lost
|
|
||||||
}
|
|
||||||
if _, exists := s.locks[name]; exists {
|
|
||||||
return errors.New("lock is already tracked by this storage instance")
|
|
||||||
}
|
|
||||||
s.locks[name] = lease
|
|
||||||
|
|
||||||
log.Debug("lock acquired", "duration", time.Since(started))
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// clusterStoreVersion returns the per-actor maximum store versions from responding machines and their names.
|
|
||||||
func (s *Storage) clusterStoreVersion(ctx context.Context, log *slog.Logger) (api.StoreVersion, []string, error) {
|
|
||||||
ctx, cancel := context.WithTimeout(ctx, distlock.DefaultMaxNodeCallTimeout)
|
|
||||||
defer cancel()
|
|
||||||
resp, err := s.client.MachineClient.InspectMachine(client.ProxyMachinesContext(ctx, nil), nil)
|
|
||||||
if err != nil {
|
|
||||||
return nil, nil, fmt.Errorf("inspect machines for store versions: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
maxVersion := make(api.StoreVersion)
|
|
||||||
machines := make([]string, 0, len(resp.Machines))
|
|
||||||
for _, m := range resp.Machines {
|
|
||||||
if m.Metadata.Error != "" {
|
|
||||||
log.Warn("skipping machine when collecting store versions",
|
|
||||||
"id", m.Metadata.MachineId, "name", m.Metadata.MachineName, "error", m.Metadata.Error)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
machines = append(machines, m.Metadata.MachineName)
|
|
||||||
maxVersion.MergeMax(m.StoreVersion)
|
|
||||||
}
|
|
||||||
slices.Sort(machines)
|
|
||||||
return maxVersion, machines, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Unlock releases a previously acquired distributed lock.
|
|
||||||
func (s *Storage) Unlock(ctx context.Context, name string) error {
|
|
||||||
s.locksMu.Lock()
|
|
||||||
lease, exists := s.locks[name]
|
|
||||||
if exists {
|
|
||||||
delete(s.locks, name)
|
|
||||||
}
|
|
||||||
s.locksMu.Unlock()
|
|
||||||
if !exists {
|
|
||||||
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.
|
|
||||||
if err := s.releaseLock(ctx, name, lease, "unlock"); err != nil {
|
|
||||||
return fmt.Errorf("release lock '%s': %w", name, err)
|
|
||||||
}
|
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// 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 {
|
|
||||||
// Unlock and rollback must attempt node cleanup even if the caller or module has already been cancelled.
|
|
||||||
ctx = context.WithoutCancel(ctx)
|
|
||||||
|
|
||||||
log := s.log.With("lock", name, "reason", reason)
|
|
||||||
started := time.Now()
|
|
||||||
log.Debug("releasing lock")
|
|
||||||
if err := lease.Release(ctx); err != nil {
|
|
||||||
log.Debug("failed to release lock", "duration", time.Since(started), "error", err)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
log.Debug("lock released", "duration", time.Since(started))
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
@@ -1,211 +0,0 @@
|
|||||||
// Package caddystorage provides Caddy storage backed by an Uncloud cluster.
|
|
||||||
package caddystorage
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"errors"
|
|
||||||
"fmt"
|
|
||||||
"log/slog"
|
|
||||||
"sync"
|
|
||||||
"sync/atomic"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/caddyserver/caddy/v2"
|
|
||||||
"github.com/caddyserver/caddy/v2/caddyconfig/caddyfile"
|
|
||||||
"github.com/caddyserver/certmagic"
|
|
||||||
"github.com/psviderski/uncloud/pkg/client"
|
|
||||||
"github.com/psviderski/uncloud/pkg/client/connector"
|
|
||||||
"github.com/psviderski/uncloud/pkg/distlock"
|
|
||||||
)
|
|
||||||
|
|
||||||
const (
|
|
||||||
// ModuleID is the Caddy module ID for Uncloud storage.
|
|
||||||
ModuleID = "caddy.storage.uncloud"
|
|
||||||
// DefaultSocketPath is the default path to the Uncloud API socket.
|
|
||||||
DefaultSocketPath = "/run/uncloud/api/uncloud.sock"
|
|
||||||
// DefaultLockTTL is the default duration of a distributed lock lease.
|
|
||||||
DefaultLockTTL = 20 * time.Second
|
|
||||||
// lockCleanupTimeout bounds how long an unloaded module waits for active lock operations when cleaning up.
|
|
||||||
lockCleanupTimeout = 5 * time.Minute
|
|
||||||
)
|
|
||||||
|
|
||||||
//nolint:gochecknoinits // Caddy modules must register during package initialisation.
|
|
||||||
func init() {
|
|
||||||
caddy.RegisterModule(new(Storage))
|
|
||||||
}
|
|
||||||
|
|
||||||
// Storage implements a Caddy storage backend that uses an Uncloud cluster to store assets such as TLS certificates.
|
|
||||||
type Storage struct {
|
|
||||||
// Socket is the path to the Uncloud API socket.
|
|
||||||
// Defaults to /run/uncloud/api/uncloud.sock when not set.
|
|
||||||
Socket string `json:"socket,omitempty"`
|
|
||||||
// LockTTL is the duration of a distributed lock after which it expires if not renewed. Locks renew automatically
|
|
||||||
// until unlocked. If an instance crashes or cannot renew, expiry allows another instance to acquire the stale lock.
|
|
||||||
// Longer durations tolerate longer interruptions but delay recovery after a crash. Normal unlocks release the lock
|
|
||||||
// immediately.
|
|
||||||
// Defaults to 20 seconds when not set.
|
|
||||||
LockTTL caddy.Duration `json:"lock_ttl,omitempty"`
|
|
||||||
|
|
||||||
client *client.Client
|
|
||||||
locker *distlock.Locker
|
|
||||||
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
|
|
||||||
// locks maps successfully acquired lock names to held leases.
|
|
||||||
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.
|
|
||||||
func (*Storage) CaddyModule() caddy.ModuleInfo {
|
|
||||||
return caddy.ModuleInfo{
|
|
||||||
ID: ModuleID,
|
|
||||||
New: func() caddy.Module { return new(Storage) },
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Provision connects the storage to the local Uncloud API and initialises the distributed locker.
|
|
||||||
func (s *Storage) Provision(ctx caddy.Context) error {
|
|
||||||
s.ctx = ctx
|
|
||||||
s.log = ctx.Slogger()
|
|
||||||
|
|
||||||
if s.Socket == "" {
|
|
||||||
s.Socket = DefaultSocketPath
|
|
||||||
}
|
|
||||||
if s.LockTTL == 0 {
|
|
||||||
s.LockTTL = caddy.Duration(DefaultLockTTL)
|
|
||||||
}
|
|
||||||
if s.LockTTL < 0 {
|
|
||||||
return errors.New("lock_ttl must be positive")
|
|
||||||
}
|
|
||||||
|
|
||||||
cli, err := client.New(ctx, connector.NewUnixConnector(s.Socket))
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("connect to Uncloud API: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
locker, err := cli.NewLocker(distlock.Config{
|
|
||||||
LeaseDuration: time.Duration(s.LockTTL),
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
_ = cli.Close()
|
|
||||||
return fmt.Errorf("create distributed locker: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
s.client = cli
|
|
||||||
s.locker = locker
|
|
||||||
s.locks = make(map[string]*distlock.Lease)
|
|
||||||
|
|
||||||
s.log.Info("module provisioned", "socket", s.Socket, "lock_ttl", time.Duration(s.LockTTL))
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Cleanup keeps the Uncloud API connection open so acquired lock leases can renew until Caddy releases them. It closes
|
|
||||||
// the connection after all lock operations finish or the cleanup timeout expires.
|
|
||||||
func (s *Storage) Cleanup() error {
|
|
||||||
if s.client == nil {
|
|
||||||
// Provision failed before opening the storage.
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
started := time.Now()
|
|
||||||
s.log.Debug("cleaning up module", "locks", s.lockOps.Load())
|
|
||||||
|
|
||||||
go func() {
|
|
||||||
timer := time.NewTimer(lockCleanupTimeout)
|
|
||||||
defer timer.Stop()
|
|
||||||
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.Warn("timed out waiting for active locks to be unlocked",
|
|
||||||
"locks", remaining, "timeout", lockCleanupTimeout)
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := s.client.Close(); err != nil {
|
|
||||||
s.log.Warn("failed to clean up module", "duration", time.Since(started), "error", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
s.log.Debug("module cleanup complete", "duration", time.Since(started))
|
|
||||||
}()
|
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// CertMagicStorage returns the provisioned CertMagic storage implementation.
|
|
||||||
func (s *Storage) CertMagicStorage() (certmagic.Storage, error) {
|
|
||||||
return s, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// UnmarshalCaddyfile configures Uncloud storage from the Caddyfile global storage block.
|
|
||||||
//
|
|
||||||
// {
|
|
||||||
// storage uncloud {
|
|
||||||
// socket /run/uncloud/api/uncloud.sock
|
|
||||||
// lock_ttl 20s
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
func (s *Storage) UnmarshalCaddyfile(d *caddyfile.Dispenser) error {
|
|
||||||
d.Next() // Skip the module name 'uncloud'.
|
|
||||||
// Reject inline arguments. NextArg leaves an opening brace for NextBlock.
|
|
||||||
if d.NextArg() {
|
|
||||||
return d.ArgErr()
|
|
||||||
}
|
|
||||||
|
|
||||||
// Read the optional options block, skipping its surrounding braces.
|
|
||||||
for d.NextBlock(0) {
|
|
||||||
switch d.Val() {
|
|
||||||
case "socket":
|
|
||||||
// Require a socket path on the same line as 'socket' option.
|
|
||||||
if !d.NextArg() {
|
|
||||||
return d.ArgErr()
|
|
||||||
}
|
|
||||||
s.Socket = d.Val()
|
|
||||||
// Reject extra arguments after the socket path.
|
|
||||||
if d.NextArg() {
|
|
||||||
return d.ArgErr()
|
|
||||||
}
|
|
||||||
case "lock_ttl":
|
|
||||||
if !d.NextArg() {
|
|
||||||
return d.ArgErr()
|
|
||||||
}
|
|
||||||
ttl, err := caddy.ParseDuration(d.Val())
|
|
||||||
if err != nil {
|
|
||||||
return d.Errf("invalid lock_ttl '%s': %w", d.Val(), err)
|
|
||||||
}
|
|
||||||
if ttl <= 0 {
|
|
||||||
return d.Err("lock_ttl must be positive")
|
|
||||||
}
|
|
||||||
if d.NextArg() {
|
|
||||||
return d.ArgErr()
|
|
||||||
}
|
|
||||||
s.LockTTL = caddy.Duration(ttl)
|
|
||||||
default:
|
|
||||||
return d.Errf("unknown uncloud storage option: '%s'", d.Val())
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
var (
|
|
||||||
_ caddy.Module = (*Storage)(nil)
|
|
||||||
_ caddy.Provisioner = (*Storage)(nil)
|
|
||||||
_ caddy.CleanerUpper = (*Storage)(nil)
|
|
||||||
_ caddy.StorageConverter = (*Storage)(nil)
|
|
||||||
_ caddyfile.Unmarshaler = (*Storage)(nil)
|
|
||||||
_ certmagic.Storage = (*Storage)(nil)
|
|
||||||
)
|
|
||||||
@@ -1,83 +0,0 @@
|
|||||||
package caddystorage
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"fmt"
|
|
||||||
"io/fs"
|
|
||||||
|
|
||||||
"github.com/caddyserver/certmagic"
|
|
||||||
"github.com/psviderski/uncloud/api/pb"
|
|
||||||
"google.golang.org/grpc/codes"
|
|
||||||
"google.golang.org/grpc/status"
|
|
||||||
)
|
|
||||||
|
|
||||||
// Store writes a value to the cluster store.
|
|
||||||
func (s *Storage) Store(ctx context.Context, key string, value []byte) error {
|
|
||||||
if _, err := s.client.CaddyStorage.Store(ctx, &pb.StoreCaddyStorageRequest{Key: key, Value: value}); err != nil {
|
|
||||||
return storageError("store", key, err)
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Load reads a value from the cluster store.
|
|
||||||
func (s *Storage) Load(ctx context.Context, key string) ([]byte, error) {
|
|
||||||
resp, err := s.client.CaddyStorage.Load(ctx, &pb.LoadCaddyStorageRequest{Key: key})
|
|
||||||
if err != nil {
|
|
||||||
return nil, storageError("load", key, err)
|
|
||||||
}
|
|
||||||
return resp.Value, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Delete removes a key and its descendants from the cluster store.
|
|
||||||
func (s *Storage) Delete(ctx context.Context, key string) error {
|
|
||||||
if _, err := s.client.CaddyStorage.Delete(ctx, &pb.DeleteCaddyStorageRequest{Key: key}); err != nil {
|
|
||||||
return storageError("delete", key, err)
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Exists reports whether a key exists in the cluster store.
|
|
||||||
func (s *Storage) Exists(ctx context.Context, key string) bool {
|
|
||||||
_, err := s.Stat(ctx, key)
|
|
||||||
return err == nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// List returns keys under prefix from the cluster store.
|
|
||||||
func (s *Storage) List(ctx context.Context, prefix string, recursive bool) ([]string, error) {
|
|
||||||
resp, err := s.client.CaddyStorage.List(ctx, &pb.ListCaddyStorageRequest{
|
|
||||||
Prefix: prefix,
|
|
||||||
Recursive: recursive,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
return nil, storageError("list", prefix, err)
|
|
||||||
}
|
|
||||||
return resp.Keys, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Stat returns information about a key in the cluster store.
|
|
||||||
func (s *Storage) Stat(ctx context.Context, key string) (certmagic.KeyInfo, error) {
|
|
||||||
resp, err := s.client.CaddyStorage.Stat(ctx, &pb.StatCaddyStorageRequest{Key: key})
|
|
||||||
if err != nil {
|
|
||||||
return certmagic.KeyInfo{}, storageError("stat", key, err)
|
|
||||||
}
|
|
||||||
|
|
||||||
info := certmagic.KeyInfo{
|
|
||||||
Key: resp.Key,
|
|
||||||
Size: resp.Size,
|
|
||||||
IsTerminal: resp.IsTerminal,
|
|
||||||
}
|
|
||||||
if resp.UpdatedAt != nil {
|
|
||||||
if err := resp.UpdatedAt.CheckValid(); err != nil {
|
|
||||||
return certmagic.KeyInfo{}, storageError("stat", key, fmt.Errorf("invalid updated_at timestamp: %w", err))
|
|
||||||
}
|
|
||||||
info.Modified = resp.UpdatedAt.AsTime()
|
|
||||||
}
|
|
||||||
return info, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func storageError(operation, key string, err error) error {
|
|
||||||
if status.Code(err) == codes.NotFound {
|
|
||||||
err = fs.ErrNotExist
|
|
||||||
}
|
|
||||||
return fmt.Errorf("uncloud storage: %s key %q: %w", operation, key, err)
|
|
||||||
}
|
|
||||||
Reference in new issue
Block a user