refactor(store): unify logic and enhance logging for initial store sync (waitStoreSync)

This commit is contained in:
Pasha Sviderski committed 2026-09-08 12:04:26 +10:00
1 parent 0f6c03240b
commit 842f796e38
3 files changed
+98 -127

No files matched your search

+45 -97
View File
@@ -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 {
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.
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)
}
// Clear MinStoreVersion so next restart doesn't wait for sync.
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 <-time.After(500 * time.Millisecond):
}
}
// Clear MinStoreVersion so the next restart doesn't wait for sync.
cc.state.mu.Lock()
cc.state.MinStoreVersion = nil
err = cc.state.Save()
err := cc.state.Save()
if err != nil {
cc.state.MinStoreVersion = minVersion
}
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))
slog.Info("Cluster store completed the initial sync.", "min_version", minVersion)
return nil
}
if len(lagging) != lastLagging {
slog.Info("Syncing cluster store.", "lagging_actors", lagging)
lastLagging = len(lagging)
}
}
}
}
// 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}
}
}
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))
}
}
}
// runMachineSync keeps this machine's info in the cluster store in sync with the local state and the Docker engine
-30
View File
@@ -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 {
+53
View File
@@ -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)