From 7a88865f64a47ab5f2450964c11a1adddb1953c8 Mon Sep 17 00:00:00 2001 From: Pasha Sviderski Date: Wed, 2 Sep 2026 17:29:30 +1000 Subject: [PATCH] feat(caddystorage): implement the machine-local CaddyStorage gRPC service --- internal/machine/api/pb/caddy_storage.proto | 1 + .../machine/api/pb/caddy_storage_grpc.pb.go | 2 + internal/machine/caddystorage/server.go | 191 ++++++++++++++++++ internal/machine/caddystorage/server_test.go | 106 ++++++++++ 4 files changed, 300 insertions(+) create mode 100644 internal/machine/caddystorage/server.go create mode 100644 internal/machine/caddystorage/server_test.go diff --git a/internal/machine/api/pb/caddy_storage.proto b/internal/machine/api/pb/caddy_storage.proto index 03f9c85e..7485b4f5 100644 --- a/internal/machine/api/pb/caddy_storage.proto +++ b/internal/machine/api/pb/caddy_storage.proto @@ -9,6 +9,7 @@ import "google/protobuf/timestamp.proto"; import "internal/machine/api/pb/common.proto"; // CaddyStorage exposes the CertMagic storage operations backed by the distributed cluster store. +// See Storage interface in https://github.com/caddyserver/certmagic/blob/master/storage.go. service CaddyStorage { rpc Store(StoreCaddyStorageRequest) returns (google.protobuf.Empty); rpc Load(LoadCaddyStorageRequest) returns (LoadCaddyStorageResponse); diff --git a/internal/machine/api/pb/caddy_storage_grpc.pb.go b/internal/machine/api/pb/caddy_storage_grpc.pb.go index 56fd4812..f5693006 100644 --- a/internal/machine/api/pb/caddy_storage_grpc.pb.go +++ b/internal/machine/api/pb/caddy_storage_grpc.pb.go @@ -32,6 +32,7 @@ const ( // For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream. // // CaddyStorage exposes the CertMagic storage operations backed by the distributed cluster store. +// See Storage interface in https://github.com/caddyserver/certmagic/blob/master/storage.go. type CaddyStorageClient interface { Store(ctx context.Context, in *StoreCaddyStorageRequest, opts ...grpc.CallOption) (*emptypb.Empty, error) Load(ctx context.Context, in *LoadCaddyStorageRequest, opts ...grpc.CallOption) (*LoadCaddyStorageResponse, error) @@ -103,6 +104,7 @@ func (c *caddyStorageClient) Stat(ctx context.Context, in *StatCaddyStorageReque // for forward compatibility. // // CaddyStorage exposes the CertMagic storage operations backed by the distributed cluster store. +// See Storage interface in https://github.com/caddyserver/certmagic/blob/master/storage.go. type CaddyStorageServer interface { Store(context.Context, *StoreCaddyStorageRequest) (*emptypb.Empty, error) Load(context.Context, *LoadCaddyStorageRequest) (*LoadCaddyStorageResponse, error) diff --git a/internal/machine/caddystorage/server.go b/internal/machine/caddystorage/server.go new file mode 100644 index 00000000..20bd7827 --- /dev/null +++ b/internal/machine/caddystorage/server.go @@ -0,0 +1,191 @@ +// Package caddystorage implements the machine-local Caddy storage API. +// See internal/machine/api/pb/caddy_storage.proto for the gRPC service definition. +package caddystorage + +import ( + "context" + "errors" + "fmt" + "maps" + "path" + "slices" + "strings" + + "github.com/psviderski/uncloud/internal/machine/api/pb" + "github.com/psviderski/uncloud/internal/machine/store" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + "google.golang.org/protobuf/types/known/emptypb" + "google.golang.org/protobuf/types/known/timestamppb" +) + +// Namespace identifies Caddy storage records in the cluster key-value store. +const Namespace = "caddy_storage" + +// Server implements the machine-local CaddyStorage gRPC service. +type Server struct { + pb.UnimplementedCaddyStorageServer + store *store.Keyspace +} + +func NewServer(store *store.Keyspace) *Server { + return &Server{store: store} +} + +func (s *Server) Store(ctx context.Context, req *pb.StoreCaddyStorageRequest) (*emptypb.Empty, error) { + if err := validateKey(req.Key); err != nil { + return nil, status.Error(codes.InvalidArgument, err.Error()) + } + if err := s.store.Put(ctx, req.Key, req.Value); err != nil { + return nil, status.Errorf(codes.Internal, "put value: %v", err) + } + + return &emptypb.Empty{}, nil +} + +func (s *Server) Load(ctx context.Context, req *pb.LoadCaddyStorageRequest) (*pb.LoadCaddyStorageResponse, error) { + if err := validateKey(req.Key); err != nil { + return nil, status.Error(codes.InvalidArgument, err.Error()) + } + + record, err := s.store.Get(ctx, req.Key) + if errors.Is(err, store.ErrKeyNotFound) { + return nil, status.Errorf(codes.NotFound, "Caddy storage key %q not found", req.Key) + } + if err != nil { + return nil, status.Errorf(codes.Internal, "get value: %v", err) + } + + return &pb.LoadCaddyStorageResponse{ + Messages: []*pb.MachineCaddyStorageValue{{ + Value: record.Value, + UpdatedAt: timestamppb.New(record.UpdatedAt), + }}, + }, nil +} + +func (s *Server) Delete(ctx context.Context, req *pb.DeleteCaddyStorageRequest) (*pb.EmptyResponse, error) { + if err := validateKey(req.Key); err != nil { + return nil, status.Error(codes.InvalidArgument, err.Error()) + } + + // Delete the key and all keys prefixed by it in case it's a directory. + if err := errors.Join( + s.store.Delete(ctx, req.Key, store.KeyspaceDeleteOptions{}), + s.store.Delete(ctx, req.Key+"/", store.KeyspaceDeleteOptions{Prefix: true}), + ); err != nil { + return nil, status.Errorf(codes.Internal, "delete key: %v", err) + } + + return &pb.EmptyResponse{ + Messages: []*pb.Empty{{}}, + }, nil +} + +func (s *Server) List(ctx context.Context, req *pb.ListCaddyStorageRequest) (*pb.ListCaddyStorageResponse, error) { + if req.Prefix != "" { + if err := validateKey(req.Prefix); err != nil { + return nil, status.Error(codes.InvalidArgument, err.Error()) + } + } + + storagePrefix := req.Prefix + if storagePrefix != "" { + storagePrefix += "/" + } + records, err := s.store.List(ctx, storagePrefix, store.KeyspaceListOptions{KeysOnly: true}) + if err != nil { + return nil, status.Errorf(codes.Internal, "list keys: %v", err) + } + // No descendants can mean an existing terminal key or a missing path. + // Check the exact key for non-root prefixes. An empty prefix lists the storage root. + if len(records) == 0 && req.Prefix != "" { + if _, err = s.store.Get(ctx, req.Prefix); errors.Is(err, store.ErrKeyNotFound) { + return nil, status.Errorf(codes.NotFound, "Caddy storage prefix %q not found", req.Prefix) + } else if err != nil { + return nil, status.Errorf(codes.Internal, "get prefix value: %v", err) + } + } + + return &pb.ListCaddyStorageResponse{ + Messages: []*pb.MachineCaddyStorageKeys{{ + Keys: listKeys(records, req.Prefix, req.Recursive), + }}, + }, nil +} + +func listKeys(records []store.Record, prefix string, recursive bool) []string { + keys := make(map[string]struct{}) + for _, record := range records { + relative := record.Key + if prefix != "" { + relative = strings.TrimPrefix(record.Key, prefix+"/") + } + components := strings.Split(relative, "/") + if !recursive { + keys[path.Join(prefix, components[0])] = struct{}{} + continue + } + + for i := range components { + keys[path.Join(prefix, path.Join(components[:i+1]...))] = struct{}{} + } + } + + return slices.Sorted(maps.Keys(keys)) +} + +func (s *Server) Stat(ctx context.Context, req *pb.StatCaddyStorageRequest) (*pb.StatCaddyStorageResponse, error) { + if err := validateKey(req.Key); err != nil { + return nil, status.Error(codes.InvalidArgument, err.Error()) + } + + record, err := s.store.Get(ctx, req.Key) + if err == nil { + return &pb.StatCaddyStorageResponse{ + Messages: []*pb.MachineCaddyStorageKeyInfo{{ + Key: req.Key, + UpdatedAt: timestamppb.New(record.UpdatedAt), + Size: int64(len(record.Value)), + IsTerminal: true, + }}, + }, nil + } + if !errors.Is(err, store.ErrKeyNotFound) { + return nil, status.Errorf(codes.Internal, "get value: %v", err) + } + + // A terminal key is not found. Check if there are any keys prefixed by the key to determine if it's a directory. + records, err := s.store.List(ctx, req.Key+"/", store.KeyspaceListOptions{KeysOnly: true}) + if err != nil { + return nil, status.Errorf(codes.Internal, "list keys: %v", err) + } + if len(records) == 0 { + return nil, status.Errorf(codes.NotFound, "Caddy storage key %q not found", req.Key) + } + + return &pb.StatCaddyStorageResponse{ + Messages: []*pb.MachineCaddyStorageKeyInfo{{ + Key: req.Key, + IsTerminal: false, + }}, + }, nil +} + +func validateKey(key string) error { + if key == "" { + return fmt.Errorf("key is empty") + } + if strings.HasPrefix(key, "/") || strings.HasSuffix(key, "/") { + return fmt.Errorf("key %q must not have a leading or trailing slash", key) + } + if strings.Contains(key, "\\") { + return fmt.Errorf("key %q must use forward slashes", key) + } + for component := range strings.SplitSeq(key, "/") { + if component == "" || component == "." || component == ".." { + return fmt.Errorf("key %q contains an invalid path component", key) + } + } + return nil +} diff --git a/internal/machine/caddystorage/server_test.go b/internal/machine/caddystorage/server_test.go new file mode 100644 index 00000000..eedead20 --- /dev/null +++ b/internal/machine/caddystorage/server_test.go @@ -0,0 +1,106 @@ +package caddystorage + +import ( + "testing" + + "github.com/psviderski/uncloud/internal/machine/store" + "github.com/stretchr/testify/require" +) + +func TestListKeys(t *testing.T) { + tests := []struct { + name string + records []store.Record + prefix string + recursive bool + want []string + }{ + { + name: "empty root", + recursive: true, + }, + { + name: "root immediate children", + records: []store.Record{ + {Key: "ocsp/example.com"}, + {Key: "certificates/issuer/example.com/example.com.key"}, + {Key: "certificates/issuer/example.com/example.com.crt"}, + }, + want: []string{ + "certificates", + "ocsp", + }, + }, + { + name: "root recursive", + records: []store.Record{ + {Key: "ocsp/example.com"}, + {Key: "certificates/issuer/example.com/example.com.key"}, + {Key: "certificates/issuer/example.com/example.com.crt"}, + }, + recursive: true, + want: []string{ + "certificates", + "certificates/issuer", + "certificates/issuer/example.com", + "certificates/issuer/example.com/example.com.crt", + "certificates/issuer/example.com/example.com.key", + "ocsp", + "ocsp/example.com", + }, + }, + { + name: "nested immediate children", + records: []store.Record{ + {Key: "certificates/issuer/example.net/example.net.crt"}, + {Key: "certificates/issuer/example.com/example.com.key"}, + {Key: "certificates/issuer/example.com/example.com.crt"}, + }, + prefix: "certificates/issuer", + want: []string{ + "certificates/issuer/example.com", + "certificates/issuer/example.net", + }, + }, + { + name: "nested terminal children", + records: []store.Record{ + {Key: "ocsp/example.net"}, + {Key: "ocsp/example.com"}, + }, + prefix: "ocsp", + want: []string{ + "ocsp/example.com", + "ocsp/example.net", + }, + }, + { + name: "nested recursive", + records: []store.Record{ + {Key: "certificates/issuer/example.net/example.net.crt"}, + {Key: "certificates/issuer/example.com/example.com.key"}, + {Key: "certificates/issuer/example.com/example.com.crt"}, + }, + prefix: "certificates/issuer", + recursive: true, + want: []string{ + "certificates/issuer/example.com", + "certificates/issuer/example.com/example.com.crt", + "certificates/issuer/example.com/example.com.key", + "certificates/issuer/example.net", + "certificates/issuer/example.net/example.net.crt", + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got := listKeys(tt.records, tt.prefix, tt.recursive) + if len(tt.want) == 0 { + require.Empty(t, got) + return + } + require.Equal(t, tt.want, got) + }) + } +}