mirror of
https://github.com/psviderski/uncloud.git
synced 2026-10-06 13:18:58 +00:00
feat(caddystorage): integrate CaddyStorage gRPC service in machine API
This commit is contained in:
1 parent
5e01db5e20
commit
c6d0371f15
4 files changed
+194
-7
No files matched your search
@@ -2,6 +2,7 @@
|
|||||||
*
|
*
|
||||||
|
|
||||||
# Allow files and directories.
|
# Allow files and directories.
|
||||||
|
!api/
|
||||||
!cmd/
|
!cmd/
|
||||||
!internal/
|
!internal/
|
||||||
!pkg/
|
!pkg/
|
||||||
|
|||||||
@@ -29,6 +29,7 @@ import (
|
|||||||
"github.com/psviderski/uncloud/internal/journal"
|
"github.com/psviderski/uncloud/internal/journal"
|
||||||
apiproxy "github.com/psviderski/uncloud/internal/machine/api/proxy"
|
apiproxy "github.com/psviderski/uncloud/internal/machine/api/proxy"
|
||||||
"github.com/psviderski/uncloud/internal/machine/caddyconfig"
|
"github.com/psviderski/uncloud/internal/machine/caddyconfig"
|
||||||
|
"github.com/psviderski/uncloud/internal/machine/caddystorage"
|
||||||
"github.com/psviderski/uncloud/internal/machine/cluster"
|
"github.com/psviderski/uncloud/internal/machine/cluster"
|
||||||
"github.com/psviderski/uncloud/internal/machine/constants"
|
"github.com/psviderski/uncloud/internal/machine/constants"
|
||||||
"github.com/psviderski/uncloud/internal/machine/corromigrate"
|
"github.com/psviderski/uncloud/internal/machine/corromigrate"
|
||||||
@@ -307,8 +308,15 @@ func NewMachine(config *Config) (*Machine, error) {
|
|||||||
WaitForNetworkReady: m.WaitForNetworkReady,
|
WaitForNetworkReady: m.WaitForNetworkReady,
|
||||||
})
|
})
|
||||||
caddyServer := caddyconfig.NewServer(caddyconfig.NewService(config.CaddyConfigDir))
|
caddyServer := caddyconfig.NewServer(caddyconfig.NewService(config.CaddyConfigDir))
|
||||||
|
|
||||||
|
caddyStore, err := corroStore.Keyspace(caddystorage.Namespace)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("create namespaced cluster store for Caddy storage: %w", err)
|
||||||
|
}
|
||||||
|
caddyStorageServer := caddystorage.NewServer(caddyStore)
|
||||||
|
|
||||||
leaseServer := distlockgrpc.NewServer(distlock.NewMemoryStore())
|
leaseServer := distlockgrpc.NewServer(distlock.NewMemoryStore())
|
||||||
m.localMachineServer = newGRPCServer(m, c, m.dockerServer, caddyServer, leaseServer)
|
m.localMachineServer = newGRPCServer(m, c, m.dockerServer, caddyServer, caddyStorageServer, leaseServer)
|
||||||
|
|
||||||
if m.Initialised() {
|
if m.Initialised() {
|
||||||
close(m.initialised)
|
close(m.initialised)
|
||||||
@@ -322,6 +330,7 @@ func newGRPCServer(
|
|||||||
c pb.ClusterServer,
|
c pb.ClusterServer,
|
||||||
d pb.DockerServer,
|
d pb.DockerServer,
|
||||||
caddy pb.CaddyServer,
|
caddy pb.CaddyServer,
|
||||||
|
caddyStorage pb.CaddyStorageServer,
|
||||||
lease distlockgrpc.LeaseServer,
|
lease distlockgrpc.LeaseServer,
|
||||||
) *grpc.Server {
|
) *grpc.Server {
|
||||||
s := grpc.NewServer()
|
s := grpc.NewServer()
|
||||||
@@ -329,6 +338,7 @@ func newGRPCServer(
|
|||||||
pb.RegisterClusterServer(s, c)
|
pb.RegisterClusterServer(s, c)
|
||||||
pb.RegisterDockerServer(s, d)
|
pb.RegisterDockerServer(s, d)
|
||||||
pb.RegisterCaddyServer(s, caddy)
|
pb.RegisterCaddyServer(s, caddy)
|
||||||
|
pb.RegisterCaddyStorageServer(s, caddyStorage)
|
||||||
distlockgrpc.RegisterLeaseServer(s, lease)
|
distlockgrpc.RegisterLeaseServer(s, lease)
|
||||||
return s
|
return s
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -26,6 +26,7 @@ type Client struct {
|
|||||||
pb.MachineClient
|
pb.MachineClient
|
||||||
pb.ClusterClient
|
pb.ClusterClient
|
||||||
Caddy pb.CaddyClient
|
Caddy pb.CaddyClient
|
||||||
|
CaddyStorage pb.CaddyStorageClient
|
||||||
// Docker is a namespaced client for the Docker service to distinguish Uncloud-specific service container operations
|
// Docker is a namespaced client for the Docker service to distinguish Uncloud-specific service container operations
|
||||||
// from generic Docker operations.
|
// from generic Docker operations.
|
||||||
Docker *docker.Client
|
Docker *docker.Client
|
||||||
@@ -58,6 +59,7 @@ func New(ctx context.Context, connector Connector) (*Client, error) {
|
|||||||
c.MachineClient = pb.NewMachineClient(c.conn)
|
c.MachineClient = pb.NewMachineClient(c.conn)
|
||||||
c.ClusterClient = pb.NewClusterClient(c.conn)
|
c.ClusterClient = pb.NewClusterClient(c.conn)
|
||||||
c.Caddy = pb.NewCaddyClient(c.conn)
|
c.Caddy = pb.NewCaddyClient(c.conn)
|
||||||
|
c.CaddyStorage = pb.NewCaddyStorageClient(c.conn)
|
||||||
c.Docker = docker.NewClient(c.conn)
|
c.Docker = docker.NewClient(c.conn)
|
||||||
c.leases = distlockgrpc.NewLeaseClient(c.conn)
|
c.leases = distlockgrpc.NewLeaseClient(c.conn)
|
||||||
|
|
||||||
|
|||||||
+180
-6
@@ -1,6 +1,7 @@
|
|||||||
package e2e
|
package e2e
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"os"
|
"os"
|
||||||
@@ -9,6 +10,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
dockerclient "github.com/docker/docker/client"
|
dockerclient "github.com/docker/docker/client"
|
||||||
|
"github.com/google/uuid"
|
||||||
"github.com/psviderski/uncloud/api/pb"
|
"github.com/psviderski/uncloud/api/pb"
|
||||||
"github.com/psviderski/uncloud/internal/ucind"
|
"github.com/psviderski/uncloud/internal/ucind"
|
||||||
"github.com/psviderski/uncloud/pkg/client"
|
"github.com/psviderski/uncloud/pkg/client"
|
||||||
@@ -122,21 +124,21 @@ func TestClusterLifecycle(t *testing.T) {
|
|||||||
})
|
})
|
||||||
|
|
||||||
t.Run("distributed lock", func(t *testing.T) {
|
t.Run("distributed lock", func(t *testing.T) {
|
||||||
firstClient, err := c.Machines[0].Connect(ctx)
|
cli0, err := c.Machines[0].Connect(ctx)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
t.Cleanup(func() {
|
t.Cleanup(func() {
|
||||||
require.NoError(t, firstClient.Close())
|
require.NoError(t, cli0.Close())
|
||||||
})
|
})
|
||||||
|
|
||||||
secondClient, err := c.Machines[1].Connect(ctx)
|
cli1, err := c.Machines[1].Connect(ctx)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
t.Cleanup(func() {
|
t.Cleanup(func() {
|
||||||
require.NoError(t, secondClient.Close())
|
require.NoError(t, cli1.Close())
|
||||||
})
|
})
|
||||||
|
|
||||||
firstLocker, err := firstClient.NewLocker(distlock.Config{})
|
firstLocker, err := cli0.NewLocker(distlock.Config{})
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
secondLocker, err := secondClient.NewLocker(distlock.Config{})
|
secondLocker, err := cli1.NewLocker(distlock.Config{})
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
acquire := func(locker *distlock.Locker) (*distlock.Lease, error) {
|
acquire := func(locker *distlock.Locker) (*distlock.Lease, error) {
|
||||||
@@ -175,6 +177,178 @@ func TestClusterLifecycle(t *testing.T) {
|
|||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
t.Run("Caddy storage replication", func(t *testing.T) {
|
||||||
|
cli0, err := c.Machines[0].Connect(ctx)
|
||||||
|
require.NoError(t, err)
|
||||||
|
t.Cleanup(func() {
|
||||||
|
require.NoError(t, cli0.Close())
|
||||||
|
})
|
||||||
|
|
||||||
|
cli1, err := c.Machines[1].Connect(ctx)
|
||||||
|
require.NoError(t, err)
|
||||||
|
t.Cleanup(func() {
|
||||||
|
require.NoError(t, cli1.Close())
|
||||||
|
})
|
||||||
|
|
||||||
|
prefix := "e2e/caddy-storage/" + uuid.NewString()
|
||||||
|
key := prefix + "/key/path"
|
||||||
|
|
||||||
|
// Keep reused test clusters clean if an assertion stops the test before its explicit deletes.
|
||||||
|
t.Cleanup(func() {
|
||||||
|
cleanupCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
_, _ = cli1.CaddyStorage.Delete(client.ProxyMachinesContext(cleanupCtx, nil),
|
||||||
|
&pb.DeleteCaddyStorageRequest{Key: prefix})
|
||||||
|
})
|
||||||
|
|
||||||
|
// Verify both creation and overwrite, allowing each value to replicate before writing the next.
|
||||||
|
var updatedAt time.Time
|
||||||
|
for _, value := range [][]byte{[]byte("test-value"), []byte("replacement-value")} {
|
||||||
|
_, err = cli0.CaddyStorage.Store(ctx, &pb.StoreCaddyStorageRequest{Key: key, Value: value})
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// A Load through another machine must find the value on the machine that accepted the local write,
|
||||||
|
// regardless of whether Corrosion has replicated it to the other machines yet.
|
||||||
|
loadResp, err := cli1.CaddyStorage.Load(client.ProxySingleMachineContext(ctx, c.Machines[0].ID),
|
||||||
|
&pb.LoadCaddyStorageRequest{Key: key})
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Len(t, loadResp.Messages, 1)
|
||||||
|
originResult := loadResp.Messages[0]
|
||||||
|
require.Nil(t, originResult.Metadata,
|
||||||
|
"Proxy to a single machine should not inject metadata into the response")
|
||||||
|
require.Equal(t, value, originResult.Value)
|
||||||
|
require.NoError(t, originResult.UpdatedAt.CheckValid())
|
||||||
|
modified := originResult.UpdatedAt.AsTime()
|
||||||
|
require.False(t, modified.IsZero(), "Stored value should have a valid updated_at timestamp")
|
||||||
|
if !updatedAt.IsZero() {
|
||||||
|
require.True(t, modified.After(updatedAt), "Overwriting a value should advance updated_at")
|
||||||
|
}
|
||||||
|
updatedAt = modified
|
||||||
|
|
||||||
|
require.Eventually(t, func() bool {
|
||||||
|
callCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
resp, err := cli1.CaddyStorage.Load(client.ProxyMachinesContext(callCtx, nil),
|
||||||
|
&pb.LoadCaddyStorageRequest{Key: key})
|
||||||
|
if err != nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
require.Len(t, resp.Messages, 3)
|
||||||
|
for _, m := range resp.Messages {
|
||||||
|
require.NotNil(t, m.Metadata)
|
||||||
|
if m.Metadata.Error != "" || !bytes.Equal(m.Value,
|
||||||
|
value) || !m.UpdatedAt.AsTime().Equal(updatedAt) {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
}, 30*time.Second, 100*time.Millisecond, "Caddy storage value %q should replicate to every machine", value)
|
||||||
|
|
||||||
|
// Every machine should report the same Caddy storage key information.
|
||||||
|
statResp, err := cli1.CaddyStorage.Stat(client.ProxyMachinesContext(ctx, nil),
|
||||||
|
&pb.StatCaddyStorageRequest{Key: key})
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Len(t, statResp.Messages, 3)
|
||||||
|
for _, m := range statResp.Messages {
|
||||||
|
require.NotNil(t, m.Metadata)
|
||||||
|
assert.Equal(t, "", m.Metadata.Error)
|
||||||
|
assert.Equal(t, key, m.Key)
|
||||||
|
assert.True(t, m.UpdatedAt.AsTime().Equal(updatedAt))
|
||||||
|
assert.EqualValues(t, len(value), m.Size)
|
||||||
|
assert.True(t, m.IsTerminal)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A path with descendants should exist as a directory even though no value is stored at that key.
|
||||||
|
statResp, err := cli1.CaddyStorage.Stat(client.ProxyMachinesContext(ctx, nil),
|
||||||
|
&pb.StatCaddyStorageRequest{Key: prefix + "/key"})
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Len(t, statResp.Messages, 3)
|
||||||
|
for _, m := range statResp.Messages {
|
||||||
|
require.NotNil(t, m.Metadata)
|
||||||
|
assert.Equal(t, "", m.Metadata.Error)
|
||||||
|
assert.Equal(t, prefix+"/key", m.Key)
|
||||||
|
assert.Nil(t, m.UpdatedAt)
|
||||||
|
assert.EqualValues(t, 0, m.Size)
|
||||||
|
assert.False(t, m.IsTerminal)
|
||||||
|
}
|
||||||
|
|
||||||
|
listResp, err := cli1.CaddyStorage.List(client.ProxyMachinesContext(ctx, nil),
|
||||||
|
&pb.ListCaddyStorageRequest{Prefix: prefix, Recursive: true})
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Len(t, listResp.Messages, 3)
|
||||||
|
for _, m := range listResp.Messages {
|
||||||
|
require.NotNil(t, m.Metadata)
|
||||||
|
assert.Equal(t, "", m.Metadata.Error)
|
||||||
|
assert.Equal(t, []string{prefix + "/key", key}, m.Keys)
|
||||||
|
}
|
||||||
|
|
||||||
|
// A non-recursive list should only return the immediate child keys.
|
||||||
|
listResp, err = cli1.CaddyStorage.List(client.ProxyMachinesContext(ctx, nil),
|
||||||
|
&pb.ListCaddyStorageRequest{Prefix: prefix, Recursive: false})
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Len(t, listResp.Messages, 3)
|
||||||
|
for _, m := range listResp.Messages {
|
||||||
|
require.NotNil(t, m.Metadata)
|
||||||
|
assert.Equal(t, "", m.Metadata.Error)
|
||||||
|
assert.Equal(t, []string{prefix + "/key"}, m.Keys)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Delete through one machine. The deletion must reach the other replicas through Corrosion.
|
||||||
|
_, err = cli1.CaddyStorage.Delete(ctx, &pb.DeleteCaddyStorageRequest{Key: prefix})
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
require.Eventually(t, func() bool {
|
||||||
|
callCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
resp, err := cli0.CaddyStorage.Load(client.ProxyMachinesContext(callCtx, nil),
|
||||||
|
&pb.LoadCaddyStorageRequest{Key: key})
|
||||||
|
if err != nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
require.Len(t, resp.Messages, 3)
|
||||||
|
for _, m := range resp.Messages {
|
||||||
|
require.NotNil(t, m.Metadata)
|
||||||
|
if codes.Code(m.Metadata.Status.GetCode()) != codes.NotFound {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
}, 30*time.Second, 100*time.Millisecond, "Caddy storage deletion should replicate to every machine")
|
||||||
|
|
||||||
|
statResp, err = cli0.CaddyStorage.Stat(client.ProxyMachinesContext(ctx, nil),
|
||||||
|
&pb.StatCaddyStorageRequest{Key: key})
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Len(t, statResp.Messages, 3)
|
||||||
|
for _, m := range statResp.Messages {
|
||||||
|
require.NotNil(t, m.Metadata)
|
||||||
|
assert.Equal(t, codes.NotFound, codes.Code(m.Metadata.Status.GetCode()))
|
||||||
|
}
|
||||||
|
|
||||||
|
listResp, err = cli0.CaddyStorage.List(client.ProxyMachinesContext(ctx, nil),
|
||||||
|
&pb.ListCaddyStorageRequest{Prefix: prefix, Recursive: true})
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Len(t, listResp.Messages, 3)
|
||||||
|
for _, m := range listResp.Messages {
|
||||||
|
require.NotNil(t, m.Metadata)
|
||||||
|
assert.Equal(t, codes.NotFound, codes.Code(m.Metadata.Status.GetCode()))
|
||||||
|
}
|
||||||
|
|
||||||
|
// Delete is idempotent, so every machine should still return a successful response.
|
||||||
|
deleteResp, err := cli1.CaddyStorage.Delete(client.ProxyMachinesContext(ctx, nil),
|
||||||
|
&pb.DeleteCaddyStorageRequest{Key: prefix})
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Len(t, deleteResp.Messages, 3)
|
||||||
|
for _, m := range deleteResp.Messages {
|
||||||
|
require.NotNil(t, m.Metadata)
|
||||||
|
assert.Equal(t, "", m.Metadata.Error)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
t.Run("remove", func(t *testing.T) {
|
t.Run("remove", func(t *testing.T) {
|
||||||
err := p.RemoveCluster(ctx, name)
|
err := p.RemoveCluster(ctx, name)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|||||||
Reference in new issue
Block a user