mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
feat: Add high availability for multiple replicas (#4555)
* feat: HA tailnet coordinator * fixup! feat: HA tailnet coordinator * fixup! feat: HA tailnet coordinator * remove printlns * close all connections on coordinator * impelement high availability feature * fixup! impelement high availability feature * fixup! impelement high availability feature * fixup! impelement high availability feature * fixup! impelement high availability feature * Add replicas * Add DERP meshing to arbitrary addresses * Move packages to highavailability folder * Move coordinator to high availability package * Add flags for HA * Rename to replicasync * Denest packages for replicas * Add test for multiple replicas * Fix coordination test * Add HA to the helm chart * Rename function pointer * Add warnings for HA * Add the ability to block endpoints * Add flag to disable P2P connections * Wow, I made the tests pass * Add replicas endpoint * Ensure close kills replica * Update sql * Add database latency to high availability * Pipe TLS to DERP mesh * Fix DERP mesh with TLS * Add tests for TLS * Fix replica sync TLS * Fix RootCA for replica meshing * Remove ID from replicasync * Fix getting certificates for meshing * Remove excessive locking * Fix linting * Store mesh key in the database * Fix replica key for tests * Fix types gen * Fix unlocking unlocked * Fix race in tests * Update enterprise/derpmesh/derpmesh.go Co-authored-by: Colin Adler <colin1adler@gmail.com> * Rename to syncReplicas * Reuse http client * Delete old replicas on a CRON * Fix race condition in connection tests * Fix linting * Fix nil type * Move pubsub to in-memory for twenty test * Add comment for configuration tweaking * Fix leak with transport * Fix close leak in derpmesh * Fix race when creating server * Remove handler update * Skip test on Windows * Fix DERP mesh test * Wrap HTTP handler replacement in mutex * Fix error message for relay * Fix API handler for normal tests * Fix speedtest * Fix replica resend * Fix derpmesh send * Ping async * Increase wait time of template version jobd * Fix race when closing replica sync * Add name to client * Log the derpmap being used * Don't connect if DERP is empty * Improve agent coordinator logging * Fix lock in coordinator * Fix relay addr * Fix race when updating durations * Fix client publish race * Run pubsub loop in a queue * Store agent nodes in order * Fix coordinator locking * Check for closed pipe Co-authored-by: Colin Adler <colin1adler@gmail.com>
This commit is contained in:
co-authored by
Colin Adler
parent
dc3519e973
commit
2ba4a62a0d
@@ -0,0 +1,391 @@
|
||||
package replicasync
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"database/sql"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"golang.org/x/xerrors"
|
||||
|
||||
"cdr.dev/slog"
|
||||
|
||||
"github.com/coder/coder/buildinfo"
|
||||
"github.com/coder/coder/coderd/database"
|
||||
)
|
||||
|
||||
var (
|
||||
PubsubEvent = "replica"
|
||||
)
|
||||
|
||||
type Options struct {
|
||||
CleanupInterval time.Duration
|
||||
UpdateInterval time.Duration
|
||||
PeerTimeout time.Duration
|
||||
RelayAddress string
|
||||
RegionID int32
|
||||
TLSConfig *tls.Config
|
||||
}
|
||||
|
||||
// New registers the replica with the database and periodically updates to ensure
|
||||
// it's healthy. It contacts all other alive replicas to ensure they are reachable.
|
||||
func New(ctx context.Context, logger slog.Logger, db database.Store, pubsub database.Pubsub, options *Options) (*Manager, error) {
|
||||
if options == nil {
|
||||
options = &Options{}
|
||||
}
|
||||
if options.PeerTimeout == 0 {
|
||||
options.PeerTimeout = 3 * time.Second
|
||||
}
|
||||
if options.UpdateInterval == 0 {
|
||||
options.UpdateInterval = 5 * time.Second
|
||||
}
|
||||
if options.CleanupInterval == 0 {
|
||||
// The cleanup interval can be quite long, because it's
|
||||
// primary purpose is to clean up dead replicas.
|
||||
options.CleanupInterval = 30 * time.Minute
|
||||
}
|
||||
hostname, err := os.Hostname()
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("get hostname: %w", err)
|
||||
}
|
||||
databaseLatency, err := db.Ping(ctx)
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("ping database: %w", err)
|
||||
}
|
||||
id := uuid.New()
|
||||
replica, err := db.InsertReplica(ctx, database.InsertReplicaParams{
|
||||
ID: id,
|
||||
CreatedAt: database.Now(),
|
||||
StartedAt: database.Now(),
|
||||
UpdatedAt: database.Now(),
|
||||
Hostname: hostname,
|
||||
RegionID: options.RegionID,
|
||||
RelayAddress: options.RelayAddress,
|
||||
Version: buildinfo.Version(),
|
||||
DatabaseLatency: int32(databaseLatency.Microseconds()),
|
||||
})
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("insert replica: %w", err)
|
||||
}
|
||||
err = pubsub.Publish(PubsubEvent, []byte(id.String()))
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("publish new replica: %w", err)
|
||||
}
|
||||
ctx, cancelFunc := context.WithCancel(ctx)
|
||||
manager := &Manager{
|
||||
id: id,
|
||||
options: options,
|
||||
db: db,
|
||||
pubsub: pubsub,
|
||||
self: replica,
|
||||
logger: logger,
|
||||
closed: make(chan struct{}),
|
||||
closeCancel: cancelFunc,
|
||||
}
|
||||
err = manager.syncReplicas(ctx)
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("run replica: %w", err)
|
||||
}
|
||||
peers := manager.Regional()
|
||||
if len(peers) > 0 {
|
||||
self := manager.Self()
|
||||
if self.RelayAddress == "" {
|
||||
return nil, xerrors.Errorf("a relay address must be specified when running multiple replicas in the same region")
|
||||
}
|
||||
}
|
||||
|
||||
err = manager.subscribe(ctx)
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("subscribe: %w", err)
|
||||
}
|
||||
manager.closeWait.Add(1)
|
||||
go manager.loop(ctx)
|
||||
return manager, nil
|
||||
}
|
||||
|
||||
// Manager keeps the replica up to date and in sync with other replicas.
|
||||
type Manager struct {
|
||||
id uuid.UUID
|
||||
options *Options
|
||||
db database.Store
|
||||
pubsub database.Pubsub
|
||||
logger slog.Logger
|
||||
|
||||
closeWait sync.WaitGroup
|
||||
closeMutex sync.Mutex
|
||||
closed chan (struct{})
|
||||
closeCancel context.CancelFunc
|
||||
|
||||
self database.Replica
|
||||
mutex sync.Mutex
|
||||
peers []database.Replica
|
||||
callback func()
|
||||
}
|
||||
|
||||
// updateInterval is used to determine a replicas state.
|
||||
// If the replica was updated > the time, it's considered healthy.
|
||||
// If the replica was updated < the time, it's considered stale.
|
||||
func (m *Manager) updateInterval() time.Time {
|
||||
return database.Now().Add(-3 * m.options.UpdateInterval)
|
||||
}
|
||||
|
||||
// loop runs the replica update sequence on an update interval.
|
||||
func (m *Manager) loop(ctx context.Context) {
|
||||
defer m.closeWait.Done()
|
||||
updateTicker := time.NewTicker(m.options.UpdateInterval)
|
||||
defer updateTicker.Stop()
|
||||
deleteTicker := time.NewTicker(m.options.CleanupInterval)
|
||||
defer deleteTicker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-deleteTicker.C:
|
||||
err := m.db.DeleteReplicasUpdatedBefore(ctx, m.updateInterval())
|
||||
if err != nil {
|
||||
m.logger.Warn(ctx, "delete old replicas", slog.Error(err))
|
||||
}
|
||||
continue
|
||||
case <-updateTicker.C:
|
||||
}
|
||||
err := m.syncReplicas(ctx)
|
||||
if err != nil && !errors.Is(err, context.Canceled) {
|
||||
m.logger.Warn(ctx, "run replica update loop", slog.Error(err))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// subscribe listens for new replica information!
|
||||
func (m *Manager) subscribe(ctx context.Context) error {
|
||||
var (
|
||||
needsUpdate = false
|
||||
updating = false
|
||||
updateMutex = sync.Mutex{}
|
||||
)
|
||||
|
||||
// This loop will continually update nodes as updates are processed.
|
||||
// The intent is to always be up to date without spamming the run
|
||||
// function, so if a new update comes in while one is being processed,
|
||||
// it will reprocess afterwards.
|
||||
var update func()
|
||||
update = func() {
|
||||
err := m.syncReplicas(ctx)
|
||||
if err != nil && !errors.Is(err, context.Canceled) {
|
||||
m.logger.Warn(ctx, "run replica from subscribe", slog.Error(err))
|
||||
}
|
||||
updateMutex.Lock()
|
||||
if needsUpdate {
|
||||
needsUpdate = false
|
||||
updateMutex.Unlock()
|
||||
update()
|
||||
return
|
||||
}
|
||||
updating = false
|
||||
updateMutex.Unlock()
|
||||
}
|
||||
cancelFunc, err := m.pubsub.Subscribe(PubsubEvent, func(ctx context.Context, message []byte) {
|
||||
updateMutex.Lock()
|
||||
defer updateMutex.Unlock()
|
||||
id, err := uuid.Parse(string(message))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
// Don't process updates for ourself!
|
||||
if id == m.id {
|
||||
return
|
||||
}
|
||||
if updating {
|
||||
needsUpdate = true
|
||||
return
|
||||
}
|
||||
updating = true
|
||||
go update()
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
go func() {
|
||||
<-ctx.Done()
|
||||
cancelFunc()
|
||||
}()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *Manager) syncReplicas(ctx context.Context) error {
|
||||
m.closeMutex.Lock()
|
||||
m.closeWait.Add(1)
|
||||
m.closeMutex.Unlock()
|
||||
defer m.closeWait.Done()
|
||||
// Expect replicas to update once every three times the interval...
|
||||
// If they don't, assume death!
|
||||
replicas, err := m.db.GetReplicasUpdatedAfter(ctx, m.updateInterval())
|
||||
if err != nil {
|
||||
return xerrors.Errorf("get replicas: %w", err)
|
||||
}
|
||||
|
||||
m.mutex.Lock()
|
||||
m.peers = make([]database.Replica, 0, len(replicas))
|
||||
for _, replica := range replicas {
|
||||
if replica.ID == m.id {
|
||||
continue
|
||||
}
|
||||
m.peers = append(m.peers, replica)
|
||||
}
|
||||
m.mutex.Unlock()
|
||||
|
||||
client := http.Client{
|
||||
Timeout: m.options.PeerTimeout,
|
||||
Transport: &http.Transport{
|
||||
TLSClientConfig: m.options.TLSConfig,
|
||||
},
|
||||
}
|
||||
defer client.CloseIdleConnections()
|
||||
var wg sync.WaitGroup
|
||||
var mu sync.Mutex
|
||||
failed := make([]string, 0)
|
||||
for _, peer := range m.Regional() {
|
||||
wg.Add(1)
|
||||
go func(peer database.Replica) {
|
||||
defer wg.Done()
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, peer.RelayAddress, nil)
|
||||
if err != nil {
|
||||
m.logger.Warn(ctx, "create http request for relay probe",
|
||||
slog.F("relay_address", peer.RelayAddress), slog.Error(err))
|
||||
return
|
||||
}
|
||||
res, err := client.Do(req)
|
||||
if err != nil {
|
||||
mu.Lock()
|
||||
failed = append(failed, fmt.Sprintf("relay %s (%s): %s", peer.Hostname, peer.RelayAddress, err))
|
||||
mu.Unlock()
|
||||
return
|
||||
}
|
||||
_ = res.Body.Close()
|
||||
}(peer)
|
||||
}
|
||||
wg.Wait()
|
||||
replicaError := ""
|
||||
if len(failed) > 0 {
|
||||
replicaError = fmt.Sprintf("Failed to dial peers: %s", strings.Join(failed, ", "))
|
||||
}
|
||||
|
||||
databaseLatency, err := m.db.Ping(ctx)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("ping database: %w", err)
|
||||
}
|
||||
|
||||
replica, err := m.db.UpdateReplica(ctx, database.UpdateReplicaParams{
|
||||
ID: m.self.ID,
|
||||
UpdatedAt: database.Now(),
|
||||
StartedAt: m.self.StartedAt,
|
||||
StoppedAt: m.self.StoppedAt,
|
||||
RelayAddress: m.self.RelayAddress,
|
||||
RegionID: m.self.RegionID,
|
||||
Hostname: m.self.Hostname,
|
||||
Version: m.self.Version,
|
||||
Error: replicaError,
|
||||
DatabaseLatency: int32(databaseLatency.Microseconds()),
|
||||
})
|
||||
if err != nil {
|
||||
return xerrors.Errorf("update replica: %w", err)
|
||||
}
|
||||
m.mutex.Lock()
|
||||
defer m.mutex.Unlock()
|
||||
if m.self.Error != replica.Error {
|
||||
// Publish an update occurred!
|
||||
err = m.pubsub.Publish(PubsubEvent, []byte(m.self.ID.String()))
|
||||
if err != nil {
|
||||
return xerrors.Errorf("publish replica update: %w", err)
|
||||
}
|
||||
}
|
||||
m.self = replica
|
||||
if m.callback != nil {
|
||||
go m.callback()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Self represents the current replica.
|
||||
func (m *Manager) Self() database.Replica {
|
||||
m.mutex.Lock()
|
||||
defer m.mutex.Unlock()
|
||||
return m.self
|
||||
}
|
||||
|
||||
// All returns every replica, including itself.
|
||||
func (m *Manager) All() []database.Replica {
|
||||
m.mutex.Lock()
|
||||
defer m.mutex.Unlock()
|
||||
return append(m.peers[:], m.self)
|
||||
}
|
||||
|
||||
// Regional returns all replicas in the same region excluding itself.
|
||||
func (m *Manager) Regional() []database.Replica {
|
||||
m.mutex.Lock()
|
||||
defer m.mutex.Unlock()
|
||||
replicas := make([]database.Replica, 0)
|
||||
for _, replica := range m.peers {
|
||||
if replica.RegionID != m.self.RegionID {
|
||||
continue
|
||||
}
|
||||
replicas = append(replicas, replica)
|
||||
}
|
||||
return replicas
|
||||
}
|
||||
|
||||
// SetCallback sets a function to execute whenever new peers
|
||||
// are refreshed or updated.
|
||||
func (m *Manager) SetCallback(callback func()) {
|
||||
m.mutex.Lock()
|
||||
defer m.mutex.Unlock()
|
||||
m.callback = callback
|
||||
// Instantly call the callback to inform replicas!
|
||||
go callback()
|
||||
}
|
||||
|
||||
func (m *Manager) Close() error {
|
||||
m.closeMutex.Lock()
|
||||
select {
|
||||
case <-m.closed:
|
||||
m.closeMutex.Unlock()
|
||||
return nil
|
||||
default:
|
||||
}
|
||||
close(m.closed)
|
||||
m.closeCancel()
|
||||
m.closeWait.Wait()
|
||||
m.closeMutex.Unlock()
|
||||
m.mutex.Lock()
|
||||
defer m.mutex.Unlock()
|
||||
ctx, cancelFunc := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancelFunc()
|
||||
_, err := m.db.UpdateReplica(ctx, database.UpdateReplicaParams{
|
||||
ID: m.self.ID,
|
||||
UpdatedAt: database.Now(),
|
||||
StartedAt: m.self.StartedAt,
|
||||
StoppedAt: sql.NullTime{
|
||||
Time: database.Now(),
|
||||
Valid: true,
|
||||
},
|
||||
RelayAddress: m.self.RelayAddress,
|
||||
RegionID: m.self.RegionID,
|
||||
Hostname: m.self.Hostname,
|
||||
Version: m.self.Version,
|
||||
Error: m.self.Error,
|
||||
})
|
||||
if err != nil {
|
||||
return xerrors.Errorf("update replica: %w", err)
|
||||
}
|
||||
err = m.pubsub.Publish(PubsubEvent, []byte(m.self.ID.String()))
|
||||
if err != nil {
|
||||
return xerrors.Errorf("publish replica update: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,239 @@
|
||||
package replicasync_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/require"
|
||||
"go.uber.org/goleak"
|
||||
|
||||
"cdr.dev/slog/sloggers/slogtest"
|
||||
"github.com/coder/coder/coderd/database"
|
||||
"github.com/coder/coder/coderd/database/databasefake"
|
||||
"github.com/coder/coder/coderd/database/dbtestutil"
|
||||
"github.com/coder/coder/enterprise/replicasync"
|
||||
"github.com/coder/coder/testutil"
|
||||
)
|
||||
|
||||
func TestMain(m *testing.M) {
|
||||
goleak.VerifyTestMain(m)
|
||||
}
|
||||
|
||||
func TestReplica(t *testing.T) {
|
||||
t.Parallel()
|
||||
t.Run("CreateOnNew", func(t *testing.T) {
|
||||
// This ensures that a new replica is created on New.
|
||||
t.Parallel()
|
||||
db, pubsub := dbtestutil.NewDB(t)
|
||||
closeChan := make(chan struct{}, 1)
|
||||
cancel, err := pubsub.Subscribe(replicasync.PubsubEvent, func(ctx context.Context, message []byte) {
|
||||
closeChan <- struct{}{}
|
||||
})
|
||||
require.NoError(t, err)
|
||||
defer cancel()
|
||||
server, err := replicasync.New(context.Background(), slogtest.Make(t, nil), db, pubsub, nil)
|
||||
require.NoError(t, err)
|
||||
<-closeChan
|
||||
_ = server.Close()
|
||||
require.NoError(t, err)
|
||||
})
|
||||
t.Run("ErrorsWithoutRelayAddress", func(t *testing.T) {
|
||||
// Ensures that the replica reports a successful status for
|
||||
// accessing all of its peers.
|
||||
t.Parallel()
|
||||
db, pubsub := dbtestutil.NewDB(t)
|
||||
_, err := db.InsertReplica(context.Background(), database.InsertReplicaParams{
|
||||
ID: uuid.New(),
|
||||
CreatedAt: database.Now(),
|
||||
StartedAt: database.Now(),
|
||||
UpdatedAt: database.Now(),
|
||||
Hostname: "something",
|
||||
})
|
||||
require.NoError(t, err)
|
||||
_, err = replicasync.New(context.Background(), slogtest.Make(t, nil), db, pubsub, nil)
|
||||
require.Error(t, err)
|
||||
require.Equal(t, "a relay address must be specified when running multiple replicas in the same region", err.Error())
|
||||
})
|
||||
t.Run("ConnectsToPeerReplica", func(t *testing.T) {
|
||||
// Ensures that the replica reports a successful status for
|
||||
// accessing all of its peers.
|
||||
t.Parallel()
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}))
|
||||
defer srv.Close()
|
||||
db, pubsub := dbtestutil.NewDB(t)
|
||||
peer, err := db.InsertReplica(context.Background(), database.InsertReplicaParams{
|
||||
ID: uuid.New(),
|
||||
CreatedAt: database.Now(),
|
||||
StartedAt: database.Now(),
|
||||
UpdatedAt: database.Now(),
|
||||
Hostname: "something",
|
||||
RelayAddress: srv.URL,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
server, err := replicasync.New(context.Background(), slogtest.Make(t, nil), db, pubsub, &replicasync.Options{
|
||||
RelayAddress: "http://169.254.169.254",
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Len(t, server.Regional(), 1)
|
||||
require.Equal(t, peer.ID, server.Regional()[0].ID)
|
||||
require.Empty(t, server.Self().Error)
|
||||
_ = server.Close()
|
||||
})
|
||||
t.Run("ConnectsToPeerReplicaTLS", func(t *testing.T) {
|
||||
// Ensures that the replica reports a successful status for
|
||||
// accessing all of its peers.
|
||||
t.Parallel()
|
||||
rawCert := testutil.GenerateTLSCertificate(t, "hello.org")
|
||||
certificate, err := x509.ParseCertificate(rawCert.Certificate[0])
|
||||
require.NoError(t, err)
|
||||
pool := x509.NewCertPool()
|
||||
pool.AddCert(certificate)
|
||||
// nolint:gosec
|
||||
tlsConfig := &tls.Config{
|
||||
Certificates: []tls.Certificate{rawCert},
|
||||
ServerName: "hello.org",
|
||||
RootCAs: pool,
|
||||
}
|
||||
srv := httptest.NewUnstartedServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}))
|
||||
srv.TLS = tlsConfig
|
||||
srv.StartTLS()
|
||||
defer srv.Close()
|
||||
db, pubsub := dbtestutil.NewDB(t)
|
||||
peer, err := db.InsertReplica(context.Background(), database.InsertReplicaParams{
|
||||
ID: uuid.New(),
|
||||
CreatedAt: database.Now(),
|
||||
StartedAt: database.Now(),
|
||||
UpdatedAt: database.Now(),
|
||||
Hostname: "something",
|
||||
RelayAddress: srv.URL,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
server, err := replicasync.New(context.Background(), slogtest.Make(t, nil), db, pubsub, &replicasync.Options{
|
||||
RelayAddress: "http://169.254.169.254",
|
||||
TLSConfig: tlsConfig,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Len(t, server.Regional(), 1)
|
||||
require.Equal(t, peer.ID, server.Regional()[0].ID)
|
||||
require.Empty(t, server.Self().Error)
|
||||
_ = server.Close()
|
||||
})
|
||||
t.Run("ConnectsToFakePeerWithError", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
db, pubsub := dbtestutil.NewDB(t)
|
||||
peer, err := db.InsertReplica(context.Background(), database.InsertReplicaParams{
|
||||
ID: uuid.New(),
|
||||
CreatedAt: database.Now().Add(time.Minute),
|
||||
StartedAt: database.Now().Add(time.Minute),
|
||||
UpdatedAt: database.Now().Add(time.Minute),
|
||||
Hostname: "something",
|
||||
// Fake address to dial!
|
||||
RelayAddress: "http://127.0.0.1:1",
|
||||
})
|
||||
require.NoError(t, err)
|
||||
server, err := replicasync.New(context.Background(), slogtest.Make(t, nil), db, pubsub, &replicasync.Options{
|
||||
PeerTimeout: 1 * time.Millisecond,
|
||||
RelayAddress: "http://127.0.0.1:1",
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Len(t, server.Regional(), 1)
|
||||
require.Equal(t, peer.ID, server.Regional()[0].ID)
|
||||
require.NotEmpty(t, server.Self().Error)
|
||||
require.Contains(t, server.Self().Error, "Failed to dial peers")
|
||||
_ = server.Close()
|
||||
})
|
||||
t.Run("RefreshOnPublish", func(t *testing.T) {
|
||||
// Refresh when a new replica appears!
|
||||
t.Parallel()
|
||||
db, pubsub := dbtestutil.NewDB(t)
|
||||
server, err := replicasync.New(context.Background(), slogtest.Make(t, nil), db, pubsub, nil)
|
||||
require.NoError(t, err)
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}))
|
||||
defer srv.Close()
|
||||
peer, err := db.InsertReplica(context.Background(), database.InsertReplicaParams{
|
||||
ID: uuid.New(),
|
||||
RelayAddress: srv.URL,
|
||||
UpdatedAt: database.Now(),
|
||||
})
|
||||
require.NoError(t, err)
|
||||
// Publish multiple times to ensure it can handle that case.
|
||||
err = pubsub.Publish(replicasync.PubsubEvent, []byte(peer.ID.String()))
|
||||
require.NoError(t, err)
|
||||
err = pubsub.Publish(replicasync.PubsubEvent, []byte(peer.ID.String()))
|
||||
require.NoError(t, err)
|
||||
require.Eventually(t, func() bool {
|
||||
return len(server.Regional()) == 1
|
||||
}, testutil.WaitShort, testutil.IntervalFast)
|
||||
_ = server.Close()
|
||||
})
|
||||
t.Run("DeletesOld", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
db, pubsub := dbtestutil.NewDB(t)
|
||||
_, err := db.InsertReplica(context.Background(), database.InsertReplicaParams{
|
||||
ID: uuid.New(),
|
||||
UpdatedAt: database.Now().Add(-time.Hour),
|
||||
})
|
||||
require.NoError(t, err)
|
||||
server, err := replicasync.New(context.Background(), slogtest.Make(t, nil), db, pubsub, &replicasync.Options{
|
||||
RelayAddress: "google.com",
|
||||
CleanupInterval: time.Millisecond,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
defer server.Close()
|
||||
require.Eventually(t, func() bool {
|
||||
return len(server.Regional()) == 0
|
||||
}, testutil.WaitShort, testutil.IntervalFast)
|
||||
})
|
||||
t.Run("TwentyConcurrent", func(t *testing.T) {
|
||||
// Ensures that twenty concurrent replicas can spawn and all
|
||||
// discover each other in parallel!
|
||||
t.Parallel()
|
||||
// This doesn't use the database fake because creating
|
||||
// this many PostgreSQL connections takes some
|
||||
// configuration tweaking.
|
||||
db := databasefake.New()
|
||||
pubsub := database.NewPubsubInMemory()
|
||||
logger := slogtest.Make(t, nil)
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}))
|
||||
defer srv.Close()
|
||||
var wg sync.WaitGroup
|
||||
count := 20
|
||||
wg.Add(count)
|
||||
for i := 0; i < count; i++ {
|
||||
server, err := replicasync.New(context.Background(), logger, db, pubsub, &replicasync.Options{
|
||||
RelayAddress: srv.URL,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() {
|
||||
_ = server.Close()
|
||||
})
|
||||
done := false
|
||||
server.SetCallback(func() {
|
||||
if len(server.All()) != count {
|
||||
return
|
||||
}
|
||||
if done {
|
||||
return
|
||||
}
|
||||
done = true
|
||||
wg.Done()
|
||||
})
|
||||
}
|
||||
wg.Wait()
|
||||
})
|
||||
}
|
||||
Reference in New Issue
Block a user