From 1c29d4d5d03a244d3b58d1418d0b1b09f2a2a4b0 Mon Sep 17 00:00:00 2001 From: Pasha Sviderski Date: Fri, 28 Aug 2026 20:59:10 +1000 Subject: [PATCH] feat(distlock): integrate distributed lease management into Machine API --- internal/machine/machine.go | 14 +++++- pkg/client/client.go | 3 ++ pkg/client/locker.go | 95 +++++++++++++++++++++++++++++++++++++ 3 files changed, 110 insertions(+), 2 deletions(-) create mode 100644 pkg/client/locker.go diff --git a/internal/machine/machine.go b/internal/machine/machine.go index a4576d47..84dbfa34 100644 --- a/internal/machine/machine.go +++ b/internal/machine/machine.go @@ -42,6 +42,8 @@ import ( "github.com/psviderski/uncloud/internal/secret" "github.com/psviderski/uncloud/internal/version" "github.com/psviderski/uncloud/pkg/api" + "github.com/psviderski/uncloud/pkg/distlock" + distlockgrpc "github.com/psviderski/uncloud/pkg/distlock/grpc" "github.com/psviderski/unregistry" "github.com/siderolabs/grpc-proxy/proxy" "golang.org/x/sync/errgroup" @@ -305,7 +307,8 @@ func NewMachine(config *Config) (*Machine, error) { WaitForNetworkReady: m.WaitForNetworkReady, }) caddyServer := caddyconfig.NewServer(caddyconfig.NewService(config.CaddyConfigDir)) - m.localMachineServer = newGRPCServer(m, c, m.dockerServer, caddyServer) + leaseServer := distlockgrpc.NewServer(distlock.NewMemoryStore()) + m.localMachineServer = newGRPCServer(m, c, m.dockerServer, caddyServer, leaseServer) if m.Initialised() { close(m.initialised) @@ -314,12 +317,19 @@ func NewMachine(config *Config) (*Machine, error) { return m, nil } -func newGRPCServer(m pb.MachineServer, c pb.ClusterServer, d pb.DockerServer, caddy pb.CaddyServer) *grpc.Server { +func newGRPCServer( + m pb.MachineServer, + c pb.ClusterServer, + d pb.DockerServer, + caddy pb.CaddyServer, + lease distlockgrpc.LeaseServer, +) *grpc.Server { s := grpc.NewServer() pb.RegisterMachineServer(s, m) pb.RegisterClusterServer(s, c) pb.RegisterDockerServer(s, d) pb.RegisterCaddyServer(s, caddy) + distlockgrpc.RegisterLeaseServer(s, lease) return s } diff --git a/pkg/client/client.go b/pkg/client/client.go index f8926522..e4b382f8 100644 --- a/pkg/client/client.go +++ b/pkg/client/client.go @@ -10,6 +10,7 @@ import ( "github.com/psviderski/uncloud/internal/machine/api/pb" "github.com/psviderski/uncloud/internal/machine/docker" "github.com/psviderski/uncloud/pkg/api" + distlockgrpc "github.com/psviderski/uncloud/pkg/distlock/grpc" "golang.org/x/net/proxy" "google.golang.org/grpc" "google.golang.org/grpc/metadata" @@ -28,6 +29,7 @@ type Client struct { // Docker is a namespaced client for the Docker service to distinguish Uncloud-specific service container operations // from generic Docker operations. Docker *docker.Client + leases distlockgrpc.LeaseClient } var _ api.Client = (*Client)(nil) @@ -57,6 +59,7 @@ func New(ctx context.Context, connector Connector) (*Client, error) { c.ClusterClient = pb.NewClusterClient(c.conn) c.Caddy = pb.NewCaddyClient(c.conn) c.Docker = docker.NewClient(c.conn) + c.leases = distlockgrpc.NewLeaseClient(c.conn) return c, nil } diff --git a/pkg/client/locker.go b/pkg/client/locker.go new file mode 100644 index 00000000..3137a9d8 --- /dev/null +++ b/pkg/client/locker.go @@ -0,0 +1,95 @@ +package client + +import ( + "context" + "fmt" + "time" + + "github.com/psviderski/uncloud/pkg/distlock" + distlockgrpc "github.com/psviderski/uncloud/pkg/distlock/grpc" + "google.golang.org/protobuf/types/known/durationpb" +) + +// NewLocker creates a distlock.Locker that acquires automatically renewed distributed leases over the machines +// in the cluster. The Locker uses the client's connection, so callers must release its active leases before closing +// the client. +func (cli *Client) NewLocker(config distlock.Config) (*distlock.Locker, error) { + return distlock.New(&lockCluster{client: cli}, config) +} + +type lockCluster struct { + client *Client +} + +var _ distlock.Cluster = (*lockCluster)(nil) + +// Nodes returns a point-in-time snapshot of the registered machines in the cluster, including temporarily unavailable +// ones. The Locker calls Nodes at the start of each Acquire and retains the returned snapshot across acquisition +// retries and for the lifetime of any acquired lease. +// +// Adding or removing machines, combined with eventual replication of the machine list, can cause different +// acquisitions to use different snapshots while active leases continue using older ones. This adapter does not +// version snapshots or coordinate membership transitions. If old and new snapshots allow disjoint quorums, two +// clients can acquire leases for the same resource. Membership changes must preserve quorum overlap while leases from +// older snapshots may remain valid. +func (c *lockCluster) Nodes(ctx context.Context) ([]distlock.Node, error) { + machines, err := c.client.ListMachines(ctx, nil) + if err != nil { + return nil, fmt.Errorf("list machines: %w", err) + } + + nodes := make([]distlock.Node, 0, len(machines)) + for _, m := range machines { + nodes = append(nodes, &lockNode{ + id: m.Machine.Id, + leases: c.client.leases, + }) + } + return nodes, nil +} + +type lockNode struct { + id string + leases distlockgrpc.LeaseClient +} + +var _ distlock.Node = (*lockNode)(nil) + +func (n *lockNode) Acquire( + ctx context.Context, resource string, token []byte, ttl time.Duration, +) (bool, error) { + resp, err := n.leases.Acquire(ProxySingleMachineContext(ctx, n.id), &distlockgrpc.AcquireLeaseRequest{ + Resource: resource, + Token: token, + Ttl: durationpb.New(ttl), + }) + if err != nil { + return false, err + } + return resp.Acquired, nil +} + +func (n *lockNode) Renew( + ctx context.Context, resource string, token []byte, ttl time.Duration, +) (bool, error) { + resp, err := n.leases.Renew(ProxySingleMachineContext(ctx, n.id), &distlockgrpc.RenewLeaseRequest{ + Resource: resource, + Token: token, + Ttl: durationpb.New(ttl), + }) + if err != nil { + return false, err + } + return resp.Renewed, nil +} + +func (n *lockNode) Release(ctx context.Context, resource string, token []byte) (bool, error) { + resp, err := n.leases.Release(ProxySingleMachineContext(ctx, n.id), &distlockgrpc.ReleaseLeaseRequest{ + Resource: resource, + Token: token, + }) + if err != nil { + return false, err + } + return resp.Released, nil +}