diff --git a/coderd/coderd.go b/coderd/coderd.go index 3b2f204bd6..ab0332f821 100644 --- a/coderd/coderd.go +++ b/coderd/coderd.go @@ -166,7 +166,7 @@ type Options struct { Pubsub pubsub.Pubsub // ReplicaSyncPubsub is used explicitly to instantiate the replicasync manager downstream if it exists. // All other consumers of pubsub should reference Options.Pubsub. - ReplicaSyncPubsub *pubsub.PGPubsub + ReplicaSyncPubsub pubsub.Pubsub RuntimeConfig *runtimeconfig.Manager // CacheDir is used for caching files served by the API. diff --git a/coderd/coderdtest/coderdtest.go b/coderd/coderdtest/coderdtest.go index c790db406c..7ab7ca8fcc 100644 --- a/coderd/coderdtest/coderdtest.go +++ b/coderd/coderdtest/coderdtest.go @@ -35,6 +35,7 @@ import ( "github.com/go-chi/chi/v5" "github.com/golang-jwt/jwt/v4" "github.com/google/uuid" + "github.com/nats-io/nats-server/v2/server" "github.com/prometheus/client_golang/prometheus" "github.com/smallstep/pkcs7" "github.com/stretchr/testify/assert" @@ -93,6 +94,7 @@ import ( "github.com/coder/coder/v2/coderd/workspacestats" "github.com/coder/coder/v2/coderd/wsbuilder" "github.com/coder/coder/v2/coderd/x/chatd/chatprovider" + natspubsub "github.com/coder/coder/v2/coderd/x/nats" "github.com/coder/coder/v2/codersdk" "github.com/coder/coder/v2/codersdk/agentsdk" "github.com/coder/coder/v2/codersdk/drpcsdk" @@ -171,7 +173,7 @@ type Options struct { // test instances are running against the same database. Database database.Store Pubsub pubsub.Pubsub - ReplicaSyncPubsub *pubsub.PGPubsub + ReplicaSyncPubsub pubsub.Pubsub // APIMiddleware inserts middleware before api.RootHandler, this can be // useful in certain tests where you want to intercept requests before @@ -290,13 +292,28 @@ func NewOptions(t testing.TB, options *Options) (func(http.Handler), context.Can usageInserter.Store(&options.UsageInserter) } if options.Database == nil { - options.Database, options.Pubsub = dbtestutil.NewDB(t) + var ps pubsub.Pubsub + options.Database, ps = dbtestutil.NewDB(t) + var ok bool + options.ReplicaSyncPubsub, ok = ps.(*pubsub.PGPubsub) + require.True(t, ok) } if options.ReplicaSyncPubsub == nil { - pgPubsub, ok := options.Pubsub.(*pubsub.PGPubsub) - require.True(t, ok, "ReplicaSyncPubsub must be a PGPubsub") - options.ReplicaSyncPubsub = pgPubsub + // To get here, the database must have been passed in, but not the ReplicSyncPubsub. We can't create a PGPubsub + // just from the database.Store since it could be anything including a mock. We need this to be independent from + // the main Pubsub in case it's NATS, since that uses the ReplicaSync to bootstrap the cluster. The in-mem + // pubsub satisfies these requirements. + options.ReplicaSyncPubsub = pubsub.NewInMemory() } + if options.Pubsub == nil { + natsCtx, natsCancel := context.WithCancel(context.Background()) + t.Cleanup(natsCancel) + natPS, err := natspubsub.New(natsCtx, *options.Logger, natspubsub.Options{ClusterPort: server.RANDOM_PORT}) + require.NoError(t, err) + t.Cleanup(func() { _ = natPS.Close() }) + options.Pubsub = natPS + } + if options.CoordinatorResumeTokenProvider == nil { options.CoordinatorResumeTokenProvider = tailnet.NewInsecureTestResumeTokenProvider() } diff --git a/enterprise/coderd/coderd.go b/enterprise/coderd/coderd.go index 764c03eda9..a38b7ee830 100644 --- a/enterprise/coderd/coderd.go +++ b/enterprise/coderd/coderd.go @@ -721,6 +721,12 @@ func New(ctx context.Context, options *Options) (_ *API, err error) { return nil, xerrors.Errorf("mount scim routes: %w", mountScimError) } + // The NATS pubsub, if enabled, used the Replica Manager for clustering. It's a layering violation if the pubsub + // for the Replica Manager *is* the NATS pubsub, because then we have a dependency loop. + if _, isNats := options.ReplicaSyncPubsub.(*nats.Pubsub); isNats { + return nil, xerrors.Errorf("replica sync pubsub cannot be the NATS pubsub") + } + // We always want to run the replica manager even if we don't have DERP // enabled, since it's used to detect other coder servers for licensing, // and NATS clustering for HA pubsub. diff --git a/enterprise/coderd/coderd_test.go b/enterprise/coderd/coderd_test.go index 9582fe2b29..e023da7b1d 100644 --- a/enterprise/coderd/coderd_test.go +++ b/enterprise/coderd/coderd_test.go @@ -212,7 +212,7 @@ func TestEntitlements(t *testing.T) { }), }) require.NoError(t, err) - err = api.Pubsub.Publish(coderd.PubsubEventLicenses, []byte{}) + err = api.ReplicaSyncPubsub.Publish(coderd.PubsubEventLicenses, []byte{}) require.NoError(t, err) require.Eventually(t, func() bool { entitlements, err := anotherClient.Entitlements(context.Background()) diff --git a/enterprise/coderd/workspaceproxy_test.go b/enterprise/coderd/workspaceproxy_test.go index 50bb7cf68c..144fbfde1a 100644 --- a/enterprise/coderd/workspaceproxy_test.go +++ b/enterprise/coderd/workspaceproxy_test.go @@ -44,13 +44,9 @@ func TestRegions(t *testing.T) { t.Run("OK", func(t *testing.T) { t.Parallel() - db, pubsub := dbtestutil.NewDB(t) - - client, _ := coderdenttest.New(t, &coderdenttest.Options{ + client, db, _ := coderdenttest.NewWithDatabase(t, &coderdenttest.Options{ Options: &coderdtest.Options{ AppHostname: appHostname, - Database: db, - Pubsub: pubsub, }, })