From e4bd1ad4439a878cda2b999c39439b184b5e0dae Mon Sep 17 00:00:00 2001 From: Pasha Sviderski Date: Thu, 1 Oct 2026 20:59:15 +1000 Subject: [PATCH] feat(store): improve waiting for store replication to skip unavailable members --- api/pb/machine.proto | 2 +- api/pb/machine_grpc.pb.go | 4 +- internal/corrosion/admin.go | 9 +++- internal/machine/machine.go | 2 +- internal/machine/store/replication.go | 76 +++++++++++++++++++++------ internal/machine/store/store.go | 7 +-- pkg/client/machine.go | 2 +- 7 files changed, 76 insertions(+), 26 deletions(-) diff --git a/api/pb/machine.proto b/api/pb/machine.proto index a4ee1791..7a86cca2 100644 --- a/api/pb/machine.proto +++ b/api/pb/machine.proto @@ -20,7 +20,7 @@ service Machine { // InspectMachine retrieves detailed information about the machine. Supports broadcasting to multiple machines. rpc InspectMachine(google.protobuf.Empty) returns (InspectMachineResponse); // WaitForStoreVersion waits until the cluster store on this machine has reached each requested actor version, with - // no known missing or pending transactions through those versions. + // no known missing or pending transactions through those versions from active members. // Corrosion may satisfy a version by applying its surviving changes or by marking it complete because its changes // have been superseded. // diff --git a/api/pb/machine_grpc.pb.go b/api/pb/machine_grpc.pb.go index 23604205..4dc01131 100644 --- a/api/pb/machine_grpc.pb.go +++ b/api/pb/machine_grpc.pb.go @@ -48,7 +48,7 @@ type MachineClient interface { // InspectMachine retrieves detailed information about the machine. Supports broadcasting to multiple machines. InspectMachine(ctx context.Context, in *emptypb.Empty, opts ...grpc.CallOption) (*InspectMachineResponse, error) // WaitForStoreVersion waits until the cluster store on this machine has reached each requested actor version, with - // no known missing or pending transactions through those versions. + // no known missing or pending transactions through those versions from active members. // Corrosion may satisfy a version by applying its surviving changes or by marking it complete because its changes // have been superseded. // @@ -224,7 +224,7 @@ type MachineServer interface { // InspectMachine retrieves detailed information about the machine. Supports broadcasting to multiple machines. InspectMachine(context.Context, *emptypb.Empty) (*InspectMachineResponse, error) // WaitForStoreVersion waits until the cluster store on this machine has reached each requested actor version, with - // no known missing or pending transactions through those versions. + // no known missing or pending transactions through those versions from active members. // Corrosion may satisfy a version by applying its surviving changes or by marking it complete because its changes // have been superseded. // diff --git a/internal/corrosion/admin.go b/internal/corrosion/admin.go index bef67b07..4b50e36c 100644 --- a/internal/corrosion/admin.go +++ b/internal/corrosion/admin.go @@ -13,6 +13,8 @@ import ( "time" ) +const adminCommandTimeout = 5 * time.Second + // AdminClient is a client for the Corrosion admin API. type AdminClient struct { sockPath string @@ -32,10 +34,15 @@ type Response struct { // The channel will be closed after sending the last or error response. The caller must read from the channel until // it is closed. func (c *AdminClient) SendCommand(cmd []byte) (<-chan Response, error) { - conn, err := net.Dial("unix", c.sockPath) + // TODO: Accept contexts in admin methods so callers can cancel requests and set their deadlines. + conn, err := net.DialTimeout("unix", c.sockPath, adminCommandTimeout) if err != nil { return nil, fmt.Errorf("connect to admin socket: %w", err) } + if err = conn.SetDeadline(time.Now().Add(adminCommandTimeout)); err != nil { + conn.Close() + return nil, fmt.Errorf("set admin connection deadline: %w", err) + } if _, err = conn.Write(encodeFrame(cmd)); err != nil { conn.Close() diff --git a/internal/machine/machine.go b/internal/machine/machine.go index 3f6f5b06..2b37998b 100644 --- a/internal/machine/machine.go +++ b/internal/machine/machine.go @@ -262,11 +262,11 @@ func NewMachine(config *Config) (*Machine, error) { if err != nil { return nil, fmt.Errorf("create corrosion API client: %w", err) } - corroStore := store.New(corro) corroAdmin, err := corrosion.NewAdminClient(config.CorrosionAdminSockPath) if err != nil { return nil, fmt.Errorf("create corrosion admin client: %w", err) } + corroStore := store.New(corro, corroAdmin) initialised := make(chan struct{}) clusterReady := make(chan struct{}) diff --git a/internal/machine/store/replication.go b/internal/machine/store/replication.go index 1d020dd8..571c721c 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/internal/corrosion" "github.com/psviderski/uncloud/pkg/api" ) @@ -15,8 +16,8 @@ import ( var ErrInvalidStoreVersion = errors.New("invalid store version") // WaitForVersion waits until the local store has reached each actor's minimum version in minVersion, with no known -// missing or pending transactions through those versions. Corrosion may satisfy a version by applying its surviving -// changes or by marking it complete because its changes have been superseded. +// missing or pending transactions through those versions from active members. Corrosion may satisfy a version +// by applying its surviving changes or by marking it complete because its changes have been superseded. // // Waiting normally makes the captured data available locally. However, if another write replaces some of that data // before it arrives, Corrosion can complete the older version without transferring the replaced data. If the @@ -53,28 +54,64 @@ func (s *Store) WaitForVersion(ctx context.Context, minVersion api.StoreVersion) return nil } - // Check all bookkeeping in one SQLite snapshot. Raw buffered rows can remain after application, - // but sequence bookkeeping is removed in the same transaction that applies or clears a version. - query := `WITH min_version(actor_id, version) AS (VALUES ` + strings.Join(placeholders, ", ") + `) + queryPrefix := `WITH min_version(actor_id, version) AS (VALUES ` + strings.Join(placeholders, ", ") + `) SELECT NOT EXISTS ( SELECT 1 FROM min_version LEFT JOIN crsql_db_versions AS current ON current.site_id = min_version.actor_id - WHERE COALESCE(current.db_version, 0) < min_version.version - OR EXISTS ( - SELECT 1 FROM __corro_bookkeeping_gaps - WHERE actor_id = min_version.actor_id AND start <= min_version.version - ) - OR EXISTS ( - SELECT 1 FROM __corro_seq_bookkeeping - WHERE site_id = min_version.actor_id AND db_version <= min_version.version - ) - )` + WHERE COALESCE(current.db_version, 0) < min_version.version` + + versionsReached := false ticker := time.NewTicker(100 * time.Millisecond) defer ticker.Stop() for { - rows, err := s.corro.QueryContext(ctx, query, args...) + query := queryPrefix + queryArgs := args + // Check for gaps and pending transactions only after the requested versions are reached. + if versionsReached { + states, err := s.corroAdmin.ClusterMembershipStates(true) + if err != nil { + return fmt.Errorf("get cluster membership for store replication: %w", err) + } + + activePlaceholders := make([]string, 0, len(states)) + queryArgs = make([]any, len(args), len(args)+len(states)) + copy(queryArgs, args) + for _, state := range states { + if state.State != corrosion.MembershipStateAlive && state.State != corrosion.MembershipStateSuspect { + continue + } + actorID, err := uuid.Parse(state.ID) + if err != nil { + return fmt.Errorf("parse cluster member actor '%s': %w", state.ID, err) + } + activePlaceholders = append(activePlaceholders, "?") + queryArgs = append(queryArgs, [16]byte(actorID)) + } + if len(activePlaceholders) > 0 { + // Check active members' gaps and pending sequences in the same query for simplicity. + // Buffered rows can remain in __corro_buffered_changes after application, but sequence bookkeeping + // is removed when a version completes. + query += ` + OR ( + min_version.actor_id IN (` + strings.Join(activePlaceholders, ", ") + `) + AND ( + EXISTS ( + SELECT 1 FROM __corro_bookkeeping_gaps + WHERE actor_id = min_version.actor_id AND start <= min_version.version + ) + OR EXISTS ( + SELECT 1 FROM __corro_seq_bookkeeping + WHERE site_id = min_version.actor_id AND db_version <= min_version.version + ) + ) + )` + } + } + query += ")" + + rows, err := s.corro.QueryContext(ctx, query, queryArgs...) if err != nil { return fmt.Errorf("check store replication: %w", err) } @@ -84,6 +121,7 @@ func (s *Store) WaitForVersion(ctx context.Context, minVersion api.StoreVersion) } return errors.New("check store replication: store version query returned no rows") } + var reached int if err = rows.Scan(&reached); err != nil { rows.Close() @@ -95,7 +133,11 @@ func (s *Store) WaitForVersion(ctx context.Context, minVersion api.StoreVersion) return fmt.Errorf("check store replication: %w", err) } if reached == 1 { - return ctx.Err() + if versionsReached { + return nil + } + versionsReached = true + continue } select { diff --git a/internal/machine/store/store.go b/internal/machine/store/store.go index 6fa45c5f..d452abc0 100644 --- a/internal/machine/store/store.go +++ b/internal/machine/store/store.go @@ -24,11 +24,12 @@ var ( // Store is a cluster store backed by a distributed Corrosion database. type Store struct { - corro *corrosion.APIClient + corro *corrosion.APIClient + corroAdmin *corrosion.AdminClient } -func New(corro *corrosion.APIClient) *Store { - return &Store{corro: corro} +func New(corro *corrosion.APIClient, corroAdmin *corrosion.AdminClient) *Store { + return &Store{corro: corro, corroAdmin: corroAdmin} } // Get retrieves an unnamespaced legacy value. diff --git a/pkg/client/machine.go b/pkg/client/machine.go index 98f781b8..45f23675 100644 --- a/pkg/client/machine.go +++ b/pkg/client/machine.go @@ -143,7 +143,7 @@ func (cli *Client) WaitClusterReady(ctx context.Context, timeout time.Duration) } // WaitForStoreVersion waits until the cluster store on the target machine has reached each requested actor version -// in minVersion, with no known missing or pending transactions through those versions. +// in minVersion, with no known missing or pending transactions through those versions from active members. // The context controls cancellation and the deadline. An empty minVersion requires no replication. // This method observes replication without initiating synchronisation. //