mirror of
https://github.com/psviderski/uncloud.git
synced 2026-10-06 13:18:58 +00:00
refactor(api): replace map[string]uint64 with api.StoreVersion for store version handling
This commit is contained in:
1 parent
e4455e936a
commit
e50c5fbe57
8 files changed
+37
-22
No files matched your search
@@ -8,6 +8,7 @@ import (
|
|||||||
"slices"
|
"slices"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/psviderski/uncloud/pkg/api"
|
||||||
"github.com/psviderski/uncloud/pkg/client"
|
"github.com/psviderski/uncloud/pkg/client"
|
||||||
"github.com/psviderski/uncloud/pkg/distlock"
|
"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.
|
// 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)
|
ctx, cancel := context.WithTimeout(ctx, distlock.DefaultMaxNodeCallTimeout)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
resp, err := s.client.MachineClient.InspectMachine(client.ProxyMachinesContext(ctx, nil), nil)
|
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)
|
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))
|
machines := make([]string, 0, len(resp.Machines))
|
||||||
for _, m := range resp.Machines {
|
for _, m := range resp.Machines {
|
||||||
if m.Metadata.Error != "" {
|
if m.Metadata.Error != "" {
|
||||||
@@ -119,9 +120,7 @@ func (s *Storage) clusterStoreVersion(ctx context.Context, log *slog.Logger) (ma
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
machines = append(machines, m.Metadata.MachineName)
|
machines = append(machines, m.Metadata.MachineName)
|
||||||
for actor, v := range m.StoreVersion {
|
maxVersion.MergeMax(m.StoreVersion)
|
||||||
maxVersion[actor] = max(maxVersion[actor], v)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
slices.Sort(machines)
|
slices.Sort(machines)
|
||||||
return maxVersion, machines, nil
|
return maxVersion, machines, nil
|
||||||
|
|||||||
+1
-1
@@ -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.
|
// 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{})
|
inspectResp, err = c.MachineClient.InspectMachine(ctx, &emptypb.Empty{})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// TODO(lhf): remove Unimplemented check when v0.17.0 is released.
|
// TODO(lhf): remove Unimplemented check when v0.17.0 is released.
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import (
|
|||||||
|
|
||||||
"github.com/psviderski/uncloud/internal/machine/network"
|
"github.com/psviderski/uncloud/internal/machine/network"
|
||||||
"github.com/psviderski/uncloud/internal/secret"
|
"github.com/psviderski/uncloud/internal/secret"
|
||||||
|
"github.com/psviderski/uncloud/pkg/api"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -31,7 +32,7 @@ type State struct {
|
|||||||
// MinStoreVersion is the cluster store version this machine must reach before participating.
|
// 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
|
// Per-actor vector (Corrosion actor UUID → max processed db_version) captured from an existing
|
||||||
// member at join time. Cleared once reached.
|
// 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 authenticates requests to the local Corrosion API.
|
||||||
CorrosionAPIToken secret.Secret `json:",omitempty"`
|
CorrosionAPIToken secret.Secret `json:",omitempty"`
|
||||||
|
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/google/uuid"
|
"github.com/google/uuid"
|
||||||
|
"github.com/psviderski/uncloud/pkg/api"
|
||||||
)
|
)
|
||||||
|
|
||||||
// ErrInvalidStoreVersion indicates an invalid actor UUID in a store version vector.
|
// 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.
|
// The method observes native replication without initiating synchronisation.
|
||||||
// An empty minVersion requires no replication. The context controls cancellation and the deadline.
|
// 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 {
|
if err := ctx.Err(); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import (
|
|||||||
"github.com/google/uuid"
|
"github.com/google/uuid"
|
||||||
"github.com/psviderski/uncloud/api/pb"
|
"github.com/psviderski/uncloud/api/pb"
|
||||||
"github.com/psviderski/uncloud/internal/corrosion"
|
"github.com/psviderski/uncloud/internal/corrosion"
|
||||||
|
"github.com/psviderski/uncloud/pkg/api"
|
||||||
"google.golang.org/protobuf/encoding/protojson"
|
"google.golang.org/protobuf/encoding/protojson"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -81,14 +82,14 @@ func (s *Store) Delete(ctx context.Context, key string) error {
|
|||||||
// limitations documented there.
|
// limitations documented there.
|
||||||
//
|
//
|
||||||
// Capturing a vector does not wait for replication or prevent further writes.
|
// 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")
|
rows, err := s.corro.QueryContext(ctx, "SELECT site_id, db_version FROM crsql_db_versions")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("query crsql_db_versions: %w", err)
|
return nil, fmt.Errorf("query crsql_db_versions: %w", err)
|
||||||
}
|
}
|
||||||
defer rows.Close()
|
defer rows.Close()
|
||||||
|
|
||||||
versions := make(map[string]uint64)
|
versions := make(api.StoreVersion)
|
||||||
for rows.Next() {
|
for rows.Next() {
|
||||||
var (
|
var (
|
||||||
siteID []byte
|
siteID []byte
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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.
|
// 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.
|
// 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})
|
_, err := cli.MachineClient.WaitForStoreVersion(ctx, &pb.WaitForStoreVersionRequest{MinVersion: minVersion})
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
+10
-11
@@ -13,6 +13,7 @@ import (
|
|||||||
"github.com/google/uuid"
|
"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/api"
|
||||||
"github.com/psviderski/uncloud/pkg/client"
|
"github.com/psviderski/uncloud/pkg/client"
|
||||||
"github.com/psviderski/uncloud/pkg/distlock"
|
"github.com/psviderski/uncloud/pkg/distlock"
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
@@ -195,25 +196,23 @@ func TestClusterLifecycle(t *testing.T) {
|
|||||||
clients[i] = cli
|
clients[i] = cli
|
||||||
}
|
}
|
||||||
|
|
||||||
storeVersion := func(clis ...*client.Client) map[string]uint64 {
|
storeVersion := func(clis ...*client.Client) api.StoreVersion {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
callCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
|
callCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
version := make(map[string]uint64)
|
version := make(api.StoreVersion)
|
||||||
for _, cli := range clis {
|
for _, cli := range clis {
|
||||||
resp, err := cli.MachineClient.InspectMachine(callCtx, &emptypb.Empty{})
|
resp, err := cli.MachineClient.InspectMachine(callCtx, nil)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
require.Len(t, resp.Machines, 1)
|
require.Len(t, resp.Machines, 1)
|
||||||
m := resp.Machines[0]
|
m := resp.Machines[0]
|
||||||
require.Len(t, m.StoreVersion, 3)
|
require.Len(t, m.StoreVersion, 3)
|
||||||
for actor, v := range m.StoreVersion {
|
version.MergeMax(m.StoreVersion)
|
||||||
version[actor] = max(version[actor], v)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
return version
|
return version
|
||||||
}
|
}
|
||||||
waitForStoreVersion := func(cli *client.Client, version map[string]uint64) {
|
waitForStoreVersion := func(cli *client.Client, version api.StoreVersion) {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
waitCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
|
waitCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
@@ -352,12 +351,12 @@ func TestClusterLifecycle(t *testing.T) {
|
|||||||
})
|
})
|
||||||
|
|
||||||
t.Run("zero version for unknown actor", func(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)
|
require.NoError(t, err)
|
||||||
})
|
})
|
||||||
|
|
||||||
t.Run("invalid actor UUID", func(t *testing.T) {
|
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))
|
require.Equal(t, codes.InvalidArgument, status.Code(err))
|
||||||
})
|
})
|
||||||
|
|
||||||
@@ -366,7 +365,7 @@ func TestClusterLifecycle(t *testing.T) {
|
|||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
// Background writes cannot satisfy a target for an actor that does not exist.
|
// 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))
|
require.Equal(t, codes.DeadlineExceeded, status.Code(err))
|
||||||
})
|
})
|
||||||
|
|
||||||
@@ -386,7 +385,7 @@ func TestClusterLifecycle(t *testing.T) {
|
|||||||
timer := time.AfterFunc(500*time.Millisecond, cancel)
|
timer := time.AfterFunc(500*time.Millisecond, cancel)
|
||||||
defer timer.Stop()
|
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))
|
require.Equal(t, codes.Canceled, status.Code(err))
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|||||||
Reference in new issue
Block a user