feat(store): improve waiting for store replication to skip unavailable members

This commit is contained in:
Pasha Sviderski committed 2026-10-01 20:59:15 +10:00
1 parent 0e53f505c0
commit e4bd1ad443
7 files changed
+71 -21

No files matched your search

+1 -1
View File
@@ -20,7 +20,7 @@ service Machine {
// InspectMachine retrieves detailed information about the machine. Supports broadcasting to multiple machines. // InspectMachine retrieves detailed information about the machine. Supports broadcasting to multiple machines.
rpc InspectMachine(google.protobuf.Empty) returns (InspectMachineResponse); rpc InspectMachine(google.protobuf.Empty) returns (InspectMachineResponse);
// WaitForStoreVersion waits until the cluster store on this machine has reached each requested actor version, with // 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 // Corrosion may satisfy a version by applying its surviving changes or by marking it complete because its changes
// have been superseded. // have been superseded.
// //
+2 -2
View File
@@ -48,7 +48,7 @@ type MachineClient interface {
// InspectMachine retrieves detailed information about the machine. Supports broadcasting to multiple machines. // InspectMachine retrieves detailed information about the machine. Supports broadcasting to multiple machines.
InspectMachine(ctx context.Context, in *emptypb.Empty, opts ...grpc.CallOption) (*InspectMachineResponse, error) 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 // 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 // Corrosion may satisfy a version by applying its surviving changes or by marking it complete because its changes
// have been superseded. // have been superseded.
// //
@@ -224,7 +224,7 @@ type MachineServer interface {
// InspectMachine retrieves detailed information about the machine. Supports broadcasting to multiple machines. // InspectMachine retrieves detailed information about the machine. Supports broadcasting to multiple machines.
InspectMachine(context.Context, *emptypb.Empty) (*InspectMachineResponse, error) InspectMachine(context.Context, *emptypb.Empty) (*InspectMachineResponse, error)
// WaitForStoreVersion waits until the cluster store on this machine has reached each requested actor version, with // 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 // Corrosion may satisfy a version by applying its surviving changes or by marking it complete because its changes
// have been superseded. // have been superseded.
// //
+8 -1
View File
@@ -13,6 +13,8 @@ import (
"time" "time"
) )
const adminCommandTimeout = 5 * time.Second
// AdminClient is a client for the Corrosion admin API. // AdminClient is a client for the Corrosion admin API.
type AdminClient struct { type AdminClient struct {
sockPath string 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 // The channel will be closed after sending the last or error response. The caller must read from the channel until
// it is closed. // it is closed.
func (c *AdminClient) SendCommand(cmd []byte) (<-chan Response, error) { 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 { if err != nil {
return nil, fmt.Errorf("connect to admin socket: %w", err) 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 { if _, err = conn.Write(encodeFrame(cmd)); err != nil {
conn.Close() conn.Close()
+1 -1
View File
@@ -262,11 +262,11 @@ func NewMachine(config *Config) (*Machine, error) {
if err != nil { if err != nil {
return nil, fmt.Errorf("create corrosion API client: %w", err) return nil, fmt.Errorf("create corrosion API client: %w", err)
} }
corroStore := store.New(corro)
corroAdmin, err := corrosion.NewAdminClient(config.CorrosionAdminSockPath) corroAdmin, err := corrosion.NewAdminClient(config.CorrosionAdminSockPath)
if err != nil { if err != nil {
return nil, fmt.Errorf("create corrosion admin client: %w", err) return nil, fmt.Errorf("create corrosion admin client: %w", err)
} }
corroStore := store.New(corro, corroAdmin)
initialised := make(chan struct{}) initialised := make(chan struct{})
clusterReady := make(chan struct{}) clusterReady := make(chan struct{})
+55 -13
View File
@@ -8,6 +8,7 @@ import (
"time" "time"
"github.com/google/uuid" "github.com/google/uuid"
"github.com/psviderski/uncloud/internal/corrosion"
"github.com/psviderski/uncloud/pkg/api" "github.com/psviderski/uncloud/pkg/api"
) )
@@ -15,8 +16,8 @@ import (
var ErrInvalidStoreVersion = errors.New("invalid store version") 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 // 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 // missing or pending transactions through those versions from active members. Corrosion may satisfy a version
// changes or by marking it complete because its changes have been superseded. // 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 // 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 // before it arrives, Corrosion can complete the older version without transferring the replaced data. If the
@@ -53,14 +54,50 @@ func (s *Store) WaitForVersion(ctx context.Context, minVersion api.StoreVersion)
return nil return nil
} }
// Check all bookkeeping in one SQLite snapshot. Raw buffered rows can remain after application, queryPrefix := `WITH min_version(actor_id, version) AS (VALUES ` + strings.Join(placeholders, ", ") + `)
// 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, ", ") + `)
SELECT NOT EXISTS ( SELECT NOT EXISTS (
SELECT 1 FROM min_version SELECT 1 FROM min_version
LEFT JOIN crsql_db_versions AS current ON current.site_id = min_version.actor_id LEFT JOIN crsql_db_versions AS current ON current.site_id = min_version.actor_id
WHERE COALESCE(current.db_version, 0) < min_version.version WHERE COALESCE(current.db_version, 0) < min_version.version`
OR EXISTS (
versionsReached := false
ticker := time.NewTicker(100 * time.Millisecond)
defer ticker.Stop()
for {
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 SELECT 1 FROM __corro_bookkeeping_gaps
WHERE actor_id = min_version.actor_id AND start <= min_version.version WHERE actor_id = min_version.actor_id AND start <= min_version.version
) )
@@ -68,13 +105,13 @@ func (s *Store) WaitForVersion(ctx context.Context, minVersion api.StoreVersion)
SELECT 1 FROM __corro_seq_bookkeeping SELECT 1 FROM __corro_seq_bookkeeping
WHERE site_id = min_version.actor_id AND db_version <= min_version.version WHERE site_id = min_version.actor_id AND db_version <= min_version.version
) )
)
)` )`
}
}
query += ")"
ticker := time.NewTicker(100 * time.Millisecond) rows, err := s.corro.QueryContext(ctx, query, queryArgs...)
defer ticker.Stop()
for {
rows, err := s.corro.QueryContext(ctx, query, args...)
if err != nil { if err != nil {
return fmt.Errorf("check store replication: %w", err) 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") return errors.New("check store replication: store version query returned no rows")
} }
var reached int var reached int
if err = rows.Scan(&reached); err != nil { if err = rows.Scan(&reached); err != nil {
rows.Close() rows.Close()
@@ -95,7 +133,11 @@ func (s *Store) WaitForVersion(ctx context.Context, minVersion api.StoreVersion)
return fmt.Errorf("check store replication: %w", err) return fmt.Errorf("check store replication: %w", err)
} }
if reached == 1 { if reached == 1 {
return ctx.Err() if versionsReached {
return nil
}
versionsReached = true
continue
} }
select { select {
+3 -2
View File
@@ -25,10 +25,11 @@ var (
// Store is a cluster store backed by a distributed Corrosion database. // Store is a cluster store backed by a distributed Corrosion database.
type Store struct { type Store struct {
corro *corrosion.APIClient corro *corrosion.APIClient
corroAdmin *corrosion.AdminClient
} }
func New(corro *corrosion.APIClient) *Store { func New(corro *corrosion.APIClient, corroAdmin *corrosion.AdminClient) *Store {
return &Store{corro: corro} return &Store{corro: corro, corroAdmin: corroAdmin}
} }
// Get retrieves an unnamespaced legacy value. // Get retrieves an unnamespaced legacy value.
+1 -1
View File
@@ -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 // 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. // The context controls cancellation and the deadline. An empty minVersion requires no replication.
// This method observes replication without initiating synchronisation. // This method observes replication without initiating synchronisation.
// //