diff --git a/caddystorage/locker.go b/caddystorage/locker.go index 7e0d99d4..1731fec2 100644 --- a/caddystorage/locker.go +++ b/caddystorage/locker.go @@ -8,6 +8,7 @@ import ( "slices" "time" + "github.com/psviderski/uncloud/pkg/api" "github.com/psviderski/uncloud/pkg/client" "github.com/psviderski/uncloud/pkg/distlock" ) @@ -102,7 +103,7 @@ func (s *Storage) Lock(ctx context.Context, name string) (err error) { } // clusterStoreVersion returns the per-actor maximum store versions from responding machines and their names. -func (s *Storage) clusterStoreVersion(ctx context.Context, log *slog.Logger) (map[string]uint64, []string, error) { +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) @@ -110,7 +111,7 @@ func (s *Storage) clusterStoreVersion(ctx context.Context, log *slog.Logger) (ma return nil, nil, fmt.Errorf("inspect machines for store versions: %w", err) } - maxVersion := make(map[string]uint64) + maxVersion := make(api.StoreVersion) machines := make([]string, 0, len(resp.Machines)) for _, m := range resp.Machines { if m.Metadata.Error != "" { @@ -119,9 +120,7 @@ func (s *Storage) clusterStoreVersion(ctx context.Context, log *slog.Logger) (ma continue } machines = append(machines, m.Metadata.MachineName) - for actor, v := range m.StoreVersion { - maxVersion[actor] = max(maxVersion[actor], v) - } + maxVersion.MergeMax(m.StoreVersion) } slices.Sort(machines) return maxVersion, machines, nil diff --git a/internal/cli/cli.go b/internal/cli/cli.go index e34720c2..8cdd38bd 100644 --- a/internal/cli/cli.go +++ b/internal/cli/cli.go @@ -458,7 +458,7 @@ func (cli *CLI) AddMachine(ctx context.Context, opts AddMachineOptions) (_ *clie } // Snapshot the cluster store version so the new machine can catch up before participating. - var storeVersion map[string]uint64 + var storeVersion api.StoreVersion inspectResp, err = c.MachineClient.InspectMachine(ctx, &emptypb.Empty{}) if err != nil { // TODO(lhf): remove Unimplemented check when v0.17.0 is released. diff --git a/internal/machine/state.go b/internal/machine/state.go index d2ad4461..c3df3544 100644 --- a/internal/machine/state.go +++ b/internal/machine/state.go @@ -10,6 +10,7 @@ import ( "github.com/psviderski/uncloud/internal/machine/network" "github.com/psviderski/uncloud/internal/secret" + "github.com/psviderski/uncloud/pkg/api" ) const ( @@ -31,7 +32,7 @@ type State struct { // MinStoreVersion is the cluster store version this machine must reach before participating. // Per-actor vector (Corrosion actor UUID → max processed db_version) captured from an existing // member at join time. Cleared once reached. - MinStoreVersion map[string]uint64 `json:",omitempty"` + MinStoreVersion api.StoreVersion `json:",omitempty"` // CorrosionAPIToken authenticates requests to the local Corrosion API. CorrosionAPIToken secret.Secret `json:",omitempty"` diff --git a/internal/machine/store/replication.go b/internal/machine/store/replication.go index 7b13f6e0..1d020dd8 100644 --- a/internal/machine/store/replication.go +++ b/internal/machine/store/replication.go @@ -8,6 +8,7 @@ import ( "time" "github.com/google/uuid" + "github.com/psviderski/uncloud/pkg/api" ) // ErrInvalidStoreVersion indicates an invalid actor UUID in a store version vector. @@ -27,7 +28,7 @@ var ErrInvalidStoreVersion = errors.New("invalid store version") // // The method observes native replication without initiating synchronisation. // An empty minVersion requires no replication. The context controls cancellation and the deadline. -func (s *Store) WaitForVersion(ctx context.Context, minVersion map[string]uint64) error { +func (s *Store) WaitForVersion(ctx context.Context, minVersion api.StoreVersion) error { if err := ctx.Err(); err != nil { return err } diff --git a/internal/machine/store/store.go b/internal/machine/store/store.go index 00ad250b..6fa45c5f 100644 --- a/internal/machine/store/store.go +++ b/internal/machine/store/store.go @@ -10,6 +10,7 @@ import ( "github.com/google/uuid" "github.com/psviderski/uncloud/api/pb" "github.com/psviderski/uncloud/internal/corrosion" + "github.com/psviderski/uncloud/pkg/api" "google.golang.org/protobuf/encoding/protojson" ) @@ -81,14 +82,14 @@ func (s *Store) Delete(ctx context.Context, key string) error { // limitations documented there. // // Capturing a vector does not wait for replication or prevent further writes. -func (s *Store) Version(ctx context.Context) (map[string]uint64, error) { +func (s *Store) Version(ctx context.Context) (api.StoreVersion, error) { rows, err := s.corro.QueryContext(ctx, "SELECT site_id, db_version FROM crsql_db_versions") if err != nil { return nil, fmt.Errorf("query crsql_db_versions: %w", err) } defer rows.Close() - versions := make(map[string]uint64) + versions := make(api.StoreVersion) for rows.Next() { var ( siteID []byte diff --git a/pkg/api/store.go b/pkg/api/store.go new file mode 100644 index 00000000..e0ba3d86 --- /dev/null +++ b/pkg/api/store.go @@ -0,0 +1,14 @@ +package api + +// StoreVersion records the highest processed database version for each Corrosion actor UUID. +// A version can include changes that were superseded, so it does not guarantee that all data +// through that version is available locally. +type StoreVersion map[string]uint64 + +// MergeMax keeps the highest observed version for each actor in other. +// The receiver must be initialised before calling MergeMax. +func (v StoreVersion) MergeMax(other StoreVersion) { + for actor, version := range other { + v[actor] = max(v[actor], version) + } +} diff --git a/pkg/client/machine.go b/pkg/client/machine.go index 8a84de8f..98f781b8 100644 --- a/pkg/client/machine.go +++ b/pkg/client/machine.go @@ -157,7 +157,7 @@ func (cli *Client) WaitClusterReady(ctx context.Context, timeout time.Duration) // // Success does not guarantee an exact snapshot or delivery of every historical value. // Callers that require a specific record or condition should verify it after waiting. -func (cli *Client) WaitForStoreVersion(ctx context.Context, minVersion map[string]uint64) error { +func (cli *Client) WaitForStoreVersion(ctx context.Context, minVersion api.StoreVersion) error { _, err := cli.MachineClient.WaitForStoreVersion(ctx, &pb.WaitForStoreVersionRequest{MinVersion: minVersion}) return err } diff --git a/test/e2e/cluster_test.go b/test/e2e/cluster_test.go index dd5c366d..5e3c559c 100644 --- a/test/e2e/cluster_test.go +++ b/test/e2e/cluster_test.go @@ -13,6 +13,7 @@ import ( "github.com/google/uuid" "github.com/psviderski/uncloud/api/pb" "github.com/psviderski/uncloud/internal/ucind" + "github.com/psviderski/uncloud/pkg/api" "github.com/psviderski/uncloud/pkg/client" "github.com/psviderski/uncloud/pkg/distlock" "github.com/stretchr/testify/assert" @@ -195,25 +196,23 @@ func TestClusterLifecycle(t *testing.T) { clients[i] = cli } - storeVersion := func(clis ...*client.Client) map[string]uint64 { + storeVersion := func(clis ...*client.Client) api.StoreVersion { t.Helper() callCtx, cancel := context.WithTimeout(ctx, 5*time.Second) defer cancel() - version := make(map[string]uint64) + version := make(api.StoreVersion) for _, cli := range clis { - resp, err := cli.MachineClient.InspectMachine(callCtx, &emptypb.Empty{}) + resp, err := cli.MachineClient.InspectMachine(callCtx, nil) require.NoError(t, err) require.Len(t, resp.Machines, 1) m := resp.Machines[0] require.Len(t, m.StoreVersion, 3) - for actor, v := range m.StoreVersion { - version[actor] = max(version[actor], v) - } + version.MergeMax(m.StoreVersion) } return version } - waitForStoreVersion := func(cli *client.Client, version map[string]uint64) { + waitForStoreVersion := func(cli *client.Client, version api.StoreVersion) { t.Helper() waitCtx, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() @@ -352,12 +351,12 @@ func TestClusterLifecycle(t *testing.T) { }) t.Run("zero version for unknown actor", func(t *testing.T) { - err := cli.WaitForStoreVersion(ctx, map[string]uint64{uuid.NewString(): 0}) + err := cli.WaitForStoreVersion(ctx, api.StoreVersion{uuid.NewString(): 0}) require.NoError(t, err) }) t.Run("invalid actor UUID", func(t *testing.T) { - err := cli.WaitForStoreVersion(ctx, map[string]uint64{"not-a-uuid": 1}) + err := cli.WaitForStoreVersion(ctx, api.StoreVersion{"not-a-uuid": 1}) require.Equal(t, codes.InvalidArgument, status.Code(err)) }) @@ -366,7 +365,7 @@ func TestClusterLifecycle(t *testing.T) { defer cancel() // Background writes cannot satisfy a target for an actor that does not exist. - err := cli.WaitForStoreVersion(waitCtx, map[string]uint64{uuid.NewString(): 1}) + err := cli.WaitForStoreVersion(waitCtx, api.StoreVersion{uuid.NewString(): 1}) require.Equal(t, codes.DeadlineExceeded, status.Code(err)) }) @@ -386,7 +385,7 @@ func TestClusterLifecycle(t *testing.T) { timer := time.AfterFunc(500*time.Millisecond, cancel) defer timer.Stop() - err := cli.WaitForStoreVersion(waitCtx, map[string]uint64{uuid.NewString(): 1}) + err := cli.WaitForStoreVersion(waitCtx, api.StoreVersion{uuid.NewString(): 1}) require.Equal(t, codes.Canceled, status.Code(err)) }) })