mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
feat: add nats experiment (#25703)
This commit is contained in:
@@ -45,6 +45,7 @@ import (
|
||||
agplschedule "github.com/coder/coder/v2/coderd/schedule"
|
||||
agplusage "github.com/coder/coder/v2/coderd/usage"
|
||||
"github.com/coder/coder/v2/coderd/wsbuilder"
|
||||
"github.com/coder/coder/v2/coderd/x/nats"
|
||||
"github.com/coder/coder/v2/codersdk"
|
||||
"github.com/coder/coder/v2/enterprise/aiseats"
|
||||
"github.com/coder/coder/v2/enterprise/coderd/connectionlog"
|
||||
@@ -655,7 +656,7 @@ func New(ctx context.Context, options *Options) (_ *API, err error) {
|
||||
|
||||
// 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.
|
||||
api.replicaManager, err = replicasync.New(ctx, options.Logger, options.Database, options.Pubsub, &replicasync.Options{
|
||||
api.replicaManager, err = replicasync.New(ctx, options.Logger, options.Database, options.ReplicaSyncPubsub, &replicasync.Options{
|
||||
ID: api.AGPL.ID,
|
||||
RelayAddress: options.DERPServerRelayAddress,
|
||||
// #nosec G115 - DERP region IDs are small and fit in int32
|
||||
@@ -757,6 +758,10 @@ type Options struct {
|
||||
|
||||
ExternalTokenEncryption []dbcrypt.Cipher
|
||||
|
||||
// ReplicaManager detects and syncs multiple Coder replicas. When provided,
|
||||
// the API owns and closes it.
|
||||
ReplicaManager *replicasync.Manager
|
||||
|
||||
// Used for high availability.
|
||||
ReplicaSyncUpdateInterval time.Duration
|
||||
ReplicaErrorGracePeriod time.Duration
|
||||
@@ -965,7 +970,12 @@ func (api *API) updateEntitlements(ctx context.Context) error {
|
||||
coordinator = haCoordinator
|
||||
}
|
||||
|
||||
api.replicaManager.SetCallback(func() {
|
||||
if natsPubsub, ok := api.Pubsub.(*nats.Pubsub); ok {
|
||||
natsPubsub.SetPeerFetcher(api.replicaManager)
|
||||
api.replicaManager.SetCallback("nats", natsPubsub.RefreshPeers)
|
||||
}
|
||||
|
||||
api.replicaManager.SetCallback("derp", func() {
|
||||
// Only update DERP mesh if the built-in server is enabled.
|
||||
if api.Options.DeploymentValues.DERP.Server.Enable {
|
||||
addresses := make([]string, 0)
|
||||
@@ -985,11 +995,16 @@ func (api *API) updateEntitlements(ctx context.Context) error {
|
||||
if api.Options.DeploymentValues.DERP.Server.Enable {
|
||||
api.derpMesh.SetAddresses([]string{}, false)
|
||||
}
|
||||
api.replicaManager.SetCallback(func() {
|
||||
api.replicaManager.SetCallback("derp", func() {
|
||||
// If the amount of replicas change, so should our entitlements.
|
||||
// This is to display a warning in the UI if the user is unlicensed.
|
||||
_ = api.updateEntitlements(api.ctx)
|
||||
})
|
||||
|
||||
if natsPubsub, ok := api.Pubsub.(*nats.Pubsub); ok {
|
||||
natsPubsub.SetPeerFetcher(nats.NopPeerFetcher{})
|
||||
api.replicaManager.SetCallback("nats", nil)
|
||||
}
|
||||
}
|
||||
|
||||
// Recheck changed in case the HA coordinator failed to set up.
|
||||
|
||||
@@ -18,6 +18,7 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
natsserver "github.com/nats-io/nats-server/v2/server"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"go.uber.org/goleak"
|
||||
@@ -36,6 +37,7 @@ import (
|
||||
"github.com/coder/coder/v2/coderd/database/dbmock"
|
||||
"github.com/coder/coder/v2/coderd/database/dbtestutil"
|
||||
"github.com/coder/coder/v2/coderd/database/dbtime"
|
||||
"github.com/coder/coder/v2/coderd/database/pubsub"
|
||||
"github.com/coder/coder/v2/coderd/entitlements"
|
||||
"github.com/coder/coder/v2/coderd/httpapi"
|
||||
agplprebuilds "github.com/coder/coder/v2/coderd/prebuilds"
|
||||
@@ -43,6 +45,7 @@ import (
|
||||
"github.com/coder/coder/v2/coderd/rbac/policy"
|
||||
"github.com/coder/coder/v2/coderd/util/namesgenerator"
|
||||
"github.com/coder/coder/v2/coderd/util/ptr"
|
||||
natspubsub "github.com/coder/coder/v2/coderd/x/nats"
|
||||
"github.com/coder/coder/v2/codersdk"
|
||||
"github.com/coder/coder/v2/codersdk/workspacesdk"
|
||||
"github.com/coder/coder/v2/enterprise/audit"
|
||||
@@ -624,6 +627,95 @@ func TestMultiReplica_EmptyRelayAddress_DisabledDERP(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultiReplica_NATSPubsubPeers(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
logger := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true})
|
||||
db, pgPubsub := dbtestutil.NewDB(t)
|
||||
clusterToken := "shared-token"
|
||||
|
||||
natsA, err := natspubsub.New(ctx, logger.Named("nats-a"), natspubsub.Options{
|
||||
ClusterHost: "127.0.0.1",
|
||||
ClusterPort: natsserver.RANDOM_PORT,
|
||||
ClusterAuthToken: clusterToken,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() { _ = natsA.Close() })
|
||||
|
||||
dv := coderdtest.DeploymentValues(t)
|
||||
dv.Experiments = []string{string(codersdk.ExperimentNATSPubsub)}
|
||||
_, _ = coderdenttest.New(t, &coderdenttest.Options{
|
||||
EntitlementsUpdateInterval: 25 * time.Millisecond,
|
||||
ReplicaSyncUpdateInterval: 25 * time.Millisecond,
|
||||
Options: &coderdtest.Options{
|
||||
Logger: &logger,
|
||||
Database: db,
|
||||
Pubsub: natsA,
|
||||
ReplicaSyncPubsub: pgPubsub.(*pubsub.PGPubsub),
|
||||
DeploymentValues: dv,
|
||||
},
|
||||
LicenseOptions: &coderdenttest.LicenseOptions{
|
||||
Features: license.Features{
|
||||
codersdk.FeatureHighAvailability: 1,
|
||||
},
|
||||
},
|
||||
})
|
||||
|
||||
natsB, err := natspubsub.New(ctx, logger.Named("nats-b"), natspubsub.Options{
|
||||
ClusterHost: "127.0.0.1",
|
||||
ClusterPort: natsserver.RANDOM_PORT,
|
||||
ClusterAuthToken: clusterToken,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() { _ = natsB.Close() })
|
||||
|
||||
mgr, err := replicasync.New(ctx, logger.Named("replica-b"), db, pgPubsub, &replicasync.Options{
|
||||
ID: uuid.New(),
|
||||
RelayAddress: fmt.Sprintf("nats://127.0.0.1:%d", natsB.Server.ClusterAddr().Port),
|
||||
RegionID: 12345,
|
||||
UpdateInterval: testutil.IntervalFast,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() { _ = mgr.Close() })
|
||||
|
||||
subject := "nats.replica"
|
||||
messages := make(chan []byte, 1)
|
||||
cancel, err := natsB.Subscribe(subject, func(_ context.Context, msg []byte) {
|
||||
messages <- msg
|
||||
})
|
||||
require.NoError(t, err)
|
||||
defer cancel()
|
||||
|
||||
payload := []byte("from-replicasync-peers")
|
||||
var publishErr error
|
||||
var flushErr error
|
||||
var updateErr error
|
||||
require.Eventually(t, func() bool {
|
||||
updateErr = mgr.PublishUpdate()
|
||||
if updateErr != nil {
|
||||
return false
|
||||
}
|
||||
publishErr = natsA.Publish(subject, payload)
|
||||
if publishErr != nil {
|
||||
return false
|
||||
}
|
||||
flushErr = natsA.Flush()
|
||||
if flushErr != nil {
|
||||
return false
|
||||
}
|
||||
select {
|
||||
case got := <-messages:
|
||||
return string(got) == string(payload)
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}, testutil.WaitShort, testutil.IntervalFast)
|
||||
require.NoError(t, updateErr)
|
||||
require.NoError(t, publishErr)
|
||||
require.NoError(t, flushErr)
|
||||
}
|
||||
|
||||
func TestSCIMDisabled(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user