From 842f796e3838091cf82393deba20e20b3f51e4cb Mon Sep 17 00:00:00 2001 From: Pasha Sviderski Date: Tue, 8 Sep 2026 12:04:26 +1000 Subject: [PATCH] refactor(store): unify logic and enhance logging for initial store sync (waitStoreSync) --- internal/machine/cluster.go | 150 +++++++++++--------------------- internal/machine/store/store.go | 30 ------- test/e2e/cluster_test.go | 53 +++++++++++ 3 files changed, 102 insertions(+), 131 deletions(-) diff --git a/internal/machine/cluster.go b/internal/machine/cluster.go index dfae1473..c4368d65 100644 --- a/internal/machine/cluster.go +++ b/internal/machine/cluster.go @@ -189,9 +189,10 @@ func (cc *clusterController) Run(ctx context.Context) error { return nil }) - // Wait for the store database to sync to the minimum version before starting store-dependent components. - // This prevents issues with using partially replicated data when the machine just joined the cluster, - // e.g., an empty machine list causing WireGuard peer misconfiguration. + // Wait for initial replication before starting store-dependent components, such as WireGuard peer reconciliation. + // This prevents issues with using partially replicated data when the machine just joined the cluster, e.g., + // an empty machine list causing WireGuard peer misconfiguration. + // This uses Store.WaitForVersion's replication-progress semantics, not an exact snapshot of the join-time data. if err = cc.waitStoreSync(ctx); err != nil { return fmt.Errorf("wait initial cluster store sync: %w", err) } @@ -366,120 +367,67 @@ func (cc *clusterController) handleEndpointChanges(ctx context.Context) { } } -// waitStoreSync blocks until the local store version >= state.MinStoreVersion and any known gaps are synced. -// No-op when MinStoreVersion is empty. Clears state.MinStoreVersion when reached. +// waitStoreSync waits for initial replication using [store.Store.WaitForVersion]. +// No-op when MinStoreVersion is empty. Retries until cancellation and clears state.MinStoreVersion after success. func (cc *clusterController) waitStoreSync(ctx context.Context) error { - target := cc.state.MinStoreVersion - if len(target) == 0 { + minVersion := cc.state.MinStoreVersion + if len(minVersion) == 0 { return nil } - slog.Info("Waiting for the initial cluster store sync.", "actors", len(target)) + slog.Info("Waiting for the initial cluster store sync.", "min_version", minVersion) - ticker := time.NewTicker(500 * time.Millisecond) - defer ticker.Stop() - // Periodic warning to surface stuck NAT/connectivity issues without aborting. + // Periodic deadlines surface stuck NAT/connectivity issues without aborting startup. warnInterval := 5 * time.Minute - warnTimer := time.NewTimer(warnInterval) - defer warnTimer.Stop() - - var ( - lastLagging int - lastErrLogTime time.Time - ) + nextWarning := time.Now().Add(warnInterval) + var lastErrLogTime time.Time for { + waitCtx, cancel := context.WithDeadline(ctx, nextWarning) + err := cc.store.WaitForVersion(waitCtx, minVersion) + cancel() + if ctx.Err() != nil { + return nil + } + if err == nil { + break + } + if errors.Is(err, store.ErrInvalidStoreVersion) { + return fmt.Errorf("wait for minimum cluster store version: %w", err) + } + + if !time.Now().Before(nextWarning) { + slog.Warn("Cluster store sync still pending. Check connectivity to peers.") + nextWarning = time.Now().Add(warnInterval) + continue + } + + // Retry store errors, throttling logs to once every 5 seconds. + if time.Since(lastErrLogTime) >= 5*time.Second { + slog.Error("Failed to check cluster store replication, retrying.", "err", err) + lastErrLogTime = time.Now() + } select { case <-ctx.Done(): return nil - case <-warnTimer.C: - local, err := cc.store.Version(ctx) - if err == nil { - slog.Error("Cluster store sync still pending. Check connectivity to peers.", - "lagging_actors", laggingActors(local, target)) - } else { - slog.Error("Cluster store sync still pending. Check connectivity to peers.", "err", err) - } - warnTimer.Reset(warnInterval) - case <-ticker.C: - local, err := cc.store.Version(ctx) - if err != nil { - // Throttle error logs to once every 5 seconds. - if time.Since(lastErrLogTime) >= 5*time.Second { - slog.Error("Failed to get the cluster store version, retrying.", "err", err) - lastErrLogTime = time.Now() - } - continue - } - - lagging := laggingActors(local, target) - if len(lagging) == 0 { - // Per-actor max doesn't imply contiguous apply: corrosion can buffer X:N before - // X:N-1 arrives and track the gap separately. Wait for any remaining gaps to be synced. - if err := cc.waitKnownMissingChanges(ctx); err != nil { - return fmt.Errorf("wait for known missing changes: %w", err) - } - // If the context was cancelled mid-gap-fill, don't persist a "synced" state. - if ctx.Err() != nil { - return nil - } - - // Clear MinStoreVersion so next restart doesn't wait for sync. - cc.state.mu.Lock() - cc.state.MinStoreVersion = nil - err = cc.state.Save() - cc.state.mu.Unlock() - if err != nil { - return fmt.Errorf("save machine state after the initial cluster store sync: %w", err) - } - - slog.Info("Cluster store completed the initial sync.", "actors", len(target)) - return nil - } - - if len(lagging) != lastLagging { - slog.Info("Syncing cluster store.", "lagging_actors", lagging) - lastLagging = len(lagging) - } + case <-time.After(500 * time.Millisecond): } } -} -// laggingActors returns target actors whose local version is below the required value, as [have, need]. -func laggingActors(local, target map[string]uint64) map[string][2]uint64 { - lagging := make(map[string][2]uint64) - for actor, need := range target { - if have := local[actor]; have < need { - lagging[actor] = [2]uint64{have, need} - } + // Clear MinStoreVersion so the next restart doesn't wait for sync. + cc.state.mu.Lock() + cc.state.MinStoreVersion = nil + err := cc.state.Save() + if err != nil { + cc.state.MinStoreVersion = minVersion } - return lagging -} - -// waitKnownMissingChanges polls the store until all known missing changes have been synced. -func (cc *clusterController) waitKnownMissingChanges(ctx context.Context) error { - ticker := time.NewTicker(1 * time.Second) - defer ticker.Stop() - - for { - select { - case <-ctx.Done(): - return nil - case <-ticker.C: - changes, err := cc.store.KnownMissingChanges(ctx) - if err != nil { - return fmt.Errorf("query known missing changes from cluster store: %w", err) - } - - if len(changes) == 0 { - slog.Debug("All known missing changes have been synced to the cluster store.") - return nil - } - - slog.Debug("Waiting for known missing changes to be synced to the cluster store.", "remaining", - len(changes)) - } + cc.state.mu.Unlock() + if err != nil { + return fmt.Errorf("save machine state after the initial cluster store sync: %w", err) } + + slog.Info("Cluster store completed the initial sync.", "min_version", minVersion) + return nil } // runMachineSync keeps this machine's info in the cluster store in sync with the local state and the Docker engine diff --git a/internal/machine/store/store.go b/internal/machine/store/store.go index dde737a3..00ad250b 100644 --- a/internal/machine/store/store.go +++ b/internal/machine/store/store.go @@ -3,7 +3,6 @@ package store import ( "context" _ "embed" - "encoding/hex" "errors" "fmt" "log/slog" @@ -107,35 +106,6 @@ func (s *Store) Version(ctx context.Context) (map[string]uint64, error) { return versions, nil } -type MissingChange struct { - ActorID string - StartVersion uint64 - EndVersion uint64 -} - -// KnownMissingChanges returns a list of currently known missing changes in the Corrosion database. -func (s *Store) KnownMissingChanges(ctx context.Context) ([]MissingChange, error) { - rows, err := s.corro.QueryContext(ctx, "SELECT actor_id, start, end FROM __corro_bookkeeping_gaps") - if err != nil { - return nil, fmt.Errorf("query missing changes: %w", err) - } - defer rows.Close() - - var changes []MissingChange - for rows.Next() { - var c MissingChange - var actorBytes []byte - if err = rows.Scan(&actorBytes, &c.StartVersion, &c.EndVersion); err != nil { - return nil, fmt.Errorf("scan missing change: %w", err) - } - - c.ActorID = hex.EncodeToString(actorBytes) - changes = append(changes, c) - } - - return changes, nil -} - func (s *Store) CreateMachine(ctx context.Context, m *pb.MachineInfo) error { mJSON, err := protojson.Marshal(m) if err != nil { diff --git a/test/e2e/cluster_test.go b/test/e2e/cluster_test.go index 05bdec55..6ee4856d 100644 --- a/test/e2e/cluster_test.go +++ b/test/e2e/cluster_test.go @@ -349,6 +349,59 @@ func TestClusterLifecycle(t *testing.T) { } }) + t.Run("wait for store version edge cases", func(t *testing.T) { + ctx, cancel := context.WithTimeout(ctx, 15*time.Second) + defer cancel() + + cli, err := c.Machines[0].Connect(ctx) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, cli.Close()) + }) + + t.Run("empty vector", func(t *testing.T) { + _, err := cli.WaitForStoreVersion(ctx, &pb.WaitForStoreVersionRequest{}) + require.NoError(t, err) + }) + + t.Run("zero version for unknown actor", func(t *testing.T) { + _, err := cli.WaitForStoreVersion(ctx, &pb.WaitForStoreVersionRequest{ + MinVersion: map[string]uint64{uuid.NewString(): 0}, + }) + require.NoError(t, err) + }) + + t.Run("invalid actor UUID", func(t *testing.T) { + _, err := cli.WaitForStoreVersion(ctx, &pb.WaitForStoreVersionRequest{ + MinVersion: map[string]uint64{"not-a-uuid": 1}, + }) + require.Equal(t, codes.InvalidArgument, status.Code(err)) + }) + + t.Run("unknown actor times out", func(t *testing.T) { + waitCtx, cancel := context.WithTimeout(ctx, 500*time.Millisecond) + defer cancel() + + // Background writes cannot satisfy a target for an actor that does not exist. + _, err := cli.WaitForStoreVersion(waitCtx, &pb.WaitForStoreVersionRequest{ + MinVersion: map[string]uint64{uuid.NewString(): 1}, + }) + require.Equal(t, codes.DeadlineExceeded, status.Code(err)) + }) + + t.Run("cancel pending wait", func(t *testing.T) { + waitCtx, cancel := context.WithCancel(ctx) + defer cancel() + timer := time.AfterFunc(500*time.Millisecond, cancel) + defer timer.Stop() + + _, err := cli.WaitForStoreVersion(waitCtx, &pb.WaitForStoreVersionRequest{ + MinVersion: map[string]uint64{uuid.NewString(): 1}, + }) + require.Equal(t, codes.Canceled, status.Code(err)) + }) + }) + t.Run("remove", func(t *testing.T) { err := p.RemoveCluster(ctx, name) require.NoError(t, err)