From a58f7d4c0deccd73ed04e642f6d2218b7a5100c9 Mon Sep 17 00:00:00 2001 From: Tor Colvin Date: Fri, 18 Sep 2026 15:21:24 -0400 Subject: [PATCH 1/5] CBG-5838: report starting until a replicator reports running /_replicationStatus fell back to the cfg's target state while no replicator had published status, so a replication reported running before anything started. - If this is assigned to this node, return starting if it hasn't get made it to running yet. - add starting to activeOnly=true filter - Add unassigned to the openapi docs --- db/sg_replicate_cfg.go | 9 ++++- db/sg_replicate_cfg_test.go | 56 +++++++++++++++++++++++++++++ docs/api/components/parameters.yaml | 2 +- docs/api/components/schemas.yaml | 2 ++ 4 files changed, 67 insertions(+), 2 deletions(-) diff --git a/db/sg_replicate_cfg.go b/db/sg_replicate_cfg.go index cdb3b750d0..fd603f191f 100644 --- a/db/sg_replicate_cfg.go +++ b/db/sg_replicate_cfg.go @@ -1628,6 +1628,12 @@ func (m *sgReplicateManager) GetReplicationStatus(ctx context.Context, replicati ID: replicationID, Status: ReplicationStateUnassigned, } + } else if remoteCfg.TargetState == ReplicationStateRunning { + // Nothing is replicating until a replicator says so. + status = &ReplicationStatus{ + ID: replicationID, + Status: ReplicationStateStarting, + } } else { status = &ReplicationStatus{ ID: replicationID, @@ -1652,7 +1658,8 @@ func (m *sgReplicateManager) GetReplicationStatus(ctx context.Context, replicati if !options.IncludeError && status.Status == ReplicationStateError { return nil, nil } - if options.ActiveOnly && status.Status != ReplicationStateRunning { + // A replication on its way to running is still active. + if options.ActiveOnly && status.Status != ReplicationStateRunning && status.Status != ReplicationStateStarting { return nil, nil } diff --git a/db/sg_replicate_cfg_test.go b/db/sg_replicate_cfg_test.go index 52b737b8cd..ee19d97da0 100644 --- a/db/sg_replicate_cfg_test.go +++ b/db/sg_replicate_cfg_test.go @@ -1373,3 +1373,59 @@ func TestSGReplicateManagerStopDrainsClusterUpdates(t *testing.T) { })) require.False(t, ran, "cluster update ran after Stop returned") } + +// TestReplicationStatusBeforeReplicatorStarts asserts that a replication assigned to a node that has not +// initialized its replicator yet is reported as starting rather than running. Nothing has connected to the +// remote at this point, and callers waiting for running would otherwise proceed while a start is still to come. +func TestReplicationStatusBeforeReplicatorStarts(t *testing.T) { + testDB, ctx := SetupTestDB(t) + defer testDB.Close(ctx) + + const localNodeUUID = "localNode" + + mgr, err := NewSGReplicateManager(ctx, testDB.DatabaseContext, testDB.CfgSG) + require.NoError(t, err) + require.NoError(t, mgr.StartLocalNode(localNodeUUID, nil)) + + for _, targetState := range []string{ReplicationStateRunning, ReplicationStateStopped} { + t.Run(targetState, func(t *testing.T) { + replicationID := "rep_" + targetState + require.NoError(t, mgr.AddReplication(&ReplicationCfg{ + ReplicationConfig: ReplicationConfig{ + ID: replicationID, + Direction: ActiveReplicatorTypePush, + Remote: "http://localhost:4984/remotedb", + }, + AssignedNode: localNodeUUID, + TargetState: targetState, + })) + + // No RefreshReplicationCfg, so no replicator exists and no status has been published. + require.Nil(t, mgr.GetActiveReplicator(replicationID)) + + status, err := mgr.GetReplicationStatus(ctx, replicationID, DefaultReplicationStatusOptions()) + require.NoError(t, err) + expectedState := targetState + if targetState == ReplicationStateRunning { + expectedState = ReplicationStateStarting + } + require.Equal(t, expectedState, status.Status) + + // activeOnly covers a replication that is starting, so a caller does not lose sight of it + // between the config write and the first connection. + options := DefaultReplicationStatusOptions() + options.ActiveOnly = true + activeStatuses, err := mgr.GetReplicationStatusAll(ctx, options) + require.NoError(t, err) + activeIDs := make([]string, 0, len(activeStatuses)) + for _, activeStatus := range activeStatuses { + activeIDs = append(activeIDs, activeStatus.ID) + } + if targetState == ReplicationStateRunning { + require.Contains(t, activeIDs, replicationID) + } else { + require.NotContains(t, activeIDs, replicationID) + } + }) + } +} diff --git a/docs/api/components/parameters.yaml b/docs/api/components/parameters.yaml index b49828d5e1..2ada68bddb 100644 --- a/docs/api/components/parameters.yaml +++ b/docs/api/components/parameters.yaml @@ -300,7 +300,7 @@ replication-active-only: schema: type: boolean default: false - description: Only return replications that are actively running (`state=running`). + description: Only return replications that are actively running or starting up (`state=running` or `state=starting`). replication-include-config: name: includeConfig in: query diff --git a/docs/api/components/schemas.yaml b/docs/api/components/schemas.yaml index b5f00da3df..1b8162fdfe 100644 --- a/docs/api/components/schemas.yaml +++ b/docs/api/components/schemas.yaml @@ -668,6 +668,7 @@ ISGRReplicationState: - error - starting - reconnecting + - unassigned x-enumDescriptions: running: Currently running replication. stopped: Not running replication. @@ -675,6 +676,7 @@ ISGRReplicationState: error: The replication is stopped due to an error. starting: The replication is starting up. reconnecting: The replication is reconnecting to the remote database. + unassigned: The replication has not been assigned to a node. title: Replication Status Retrieved-replication: description: Properties of a replication From 98fbd45ecc964669904562fc68505403d32a3143 Mon Sep 17 00:00:00 2001 From: Tor Colvin Date: Mon, 28 Sep 2026 13:55:22 -0400 Subject: [PATCH 2/5] CBG-5838: treat unset target state and reconnecting as active - An unset target state starts like running, so report it as starting too. - activeOnly now includes reconnecting replications. - Cover replications assigned to another node, live replicator states, the activeOnly REST query, and upserts while starting. Co-Authored-By: Claude Opus 5.5 (1M context) --- db/sg_replicate_cfg.go | 10 ++- db/sg_replicate_cfg_test.go | 108 ++++++++++++++++++++++--- docs/api/components/parameters.yaml | 2 +- rest/replicatortest/replicator_test.go | 41 ++++++++++ 4 files changed, 144 insertions(+), 17 deletions(-) diff --git a/db/sg_replicate_cfg.go b/db/sg_replicate_cfg.go index fd603f191f..05a2c23012 100644 --- a/db/sg_replicate_cfg.go +++ b/db/sg_replicate_cfg.go @@ -1628,7 +1628,7 @@ func (m *sgReplicateManager) GetReplicationStatus(ctx context.Context, replicati ID: replicationID, Status: ReplicationStateUnassigned, } - } else if remoteCfg.TargetState == ReplicationStateRunning { + } else if remoteCfg.TargetState == "" || remoteCfg.TargetState == ReplicationStateRunning { // Nothing is replicating until a replicator says so. status = &ReplicationStatus{ ID: replicationID, @@ -1658,14 +1658,18 @@ func (m *sgReplicateManager) GetReplicationStatus(ctx context.Context, replicati if !options.IncludeError && status.Status == ReplicationStateError { return nil, nil } - // A replication on its way to running is still active. - if options.ActiveOnly && status.Status != ReplicationStateRunning && status.Status != ReplicationStateStarting { + if options.ActiveOnly && !isActiveReplicationState(status.Status) { return nil, nil } return status, nil } +// isActiveReplicationState returns true for a replication that is running or on its way to running. +func isActiveReplicationState(state string) bool { + return state == ReplicationStateRunning || state == ReplicationStateStarting || state == ReplicationStateReconnecting +} + // PutReplicationStatus updates the state of a replication. func (m *sgReplicateManager) PutReplicationStatus(ctx context.Context, replicationID, action string) (status *ReplicationStatus, auditEvent base.AuditID, err error) { diff --git a/db/sg_replicate_cfg_test.go b/db/sg_replicate_cfg_test.go index ee19d97da0..3305b80397 100644 --- a/db/sg_replicate_cfg_test.go +++ b/db/sg_replicate_cfg_test.go @@ -15,6 +15,7 @@ import ( "maps" "net/http" "net/http/httptest" + "strings" "sync" "sync/atomic" "testing" @@ -1381,23 +1382,41 @@ func TestReplicationStatusBeforeReplicatorStarts(t *testing.T) { testDB, ctx := SetupTestDB(t) defer testDB.Close(ctx) - const localNodeUUID = "localNode" + const ( + localNodeUUID = "localNode" + otherNodeUUID = "otherNode" + ) mgr, err := NewSGReplicateManager(ctx, testDB.DatabaseContext, testDB.CfgSG) require.NoError(t, err) require.NoError(t, mgr.StartLocalNode(localNodeUUID, nil)) - - for _, targetState := range []string{ReplicationStateRunning, ReplicationStateStopped} { - t.Run(targetState, func(t *testing.T) { - replicationID := "rep_" + targetState + require.NoError(t, mgr.RegisterNode(otherNodeUUID)) + + testCases := []struct { + name string + assignedNode string + targetState string + expectedStatus string + active bool + }{ + {name: "running", assignedNode: localNodeUUID, targetState: ReplicationStateRunning, expectedStatus: ReplicationStateStarting, active: true}, + {name: "stopped", assignedNode: localNodeUUID, targetState: ReplicationStateStopped, expectedStatus: ReplicationStateStopped, active: false}, + // startup treats an unset target state as running + {name: "unset", assignedNode: localNodeUUID, targetState: "", expectedStatus: ReplicationStateStarting, active: true}, + {name: "other node running", assignedNode: otherNodeUUID, targetState: ReplicationStateRunning, expectedStatus: ReplicationStateStarting, active: true}, + {name: "other node stopped", assignedNode: otherNodeUUID, targetState: ReplicationStateStopped, expectedStatus: ReplicationStateStopped, active: false}, + } + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + replicationID := "rep_" + strings.ReplaceAll(tc.name, " ", "_") require.NoError(t, mgr.AddReplication(&ReplicationCfg{ ReplicationConfig: ReplicationConfig{ ID: replicationID, Direction: ActiveReplicatorTypePush, Remote: "http://localhost:4984/remotedb", }, - AssignedNode: localNodeUUID, - TargetState: targetState, + AssignedNode: tc.assignedNode, + TargetState: tc.targetState, })) // No RefreshReplicationCfg, so no replicator exists and no status has been published. @@ -1405,11 +1424,7 @@ func TestReplicationStatusBeforeReplicatorStarts(t *testing.T) { status, err := mgr.GetReplicationStatus(ctx, replicationID, DefaultReplicationStatusOptions()) require.NoError(t, err) - expectedState := targetState - if targetState == ReplicationStateRunning { - expectedState = ReplicationStateStarting - } - require.Equal(t, expectedState, status.Status) + require.Equal(t, tc.expectedStatus, status.Status) // activeOnly covers a replication that is starting, so a caller does not lose sight of it // between the config write and the first connection. @@ -1421,11 +1436,78 @@ func TestReplicationStatusBeforeReplicatorStarts(t *testing.T) { for _, activeStatus := range activeStatuses { activeIDs = append(activeIDs, activeStatus.ID) } - if targetState == ReplicationStateRunning { + if tc.active { require.Contains(t, activeIDs, replicationID) } else { require.NotContains(t, activeIDs, replicationID) } + + // An upsert needs a stopped replication, and a replication that is starting is not stopped. + created, err := mgr.UpsertReplication(ctx, &ReplicationUpsertConfig{ID: replicationID, Remote: base.Ptr("http://localhost:4984/otherdb")}) + if tc.active { + require.Error(t, err) + status, _ := base.ErrorAsHTTPStatus(err) + require.Equal(t, http.StatusBadRequest, status) + } else { + require.NoError(t, err) + require.False(t, created) + } + }) + } +} + +// TestReplicationStatusActiveOnlyForReplicatorStates asserts which states of a local replicator activeOnly includes. +func TestReplicationStatusActiveOnlyForReplicatorStates(t *testing.T) { + testDB, ctx := SetupTestDB(t) + defer testDB.Close(ctx) + + const ( + localNodeUUID = "localNode" + replicationID = "rep" + ) + + mgr, err := NewSGReplicateManager(ctx, testDB.DatabaseContext, testDB.CfgSG) + require.NoError(t, err) + require.NoError(t, mgr.StartLocalNode(localNodeUUID, nil)) + require.NoError(t, mgr.AddReplication(&ReplicationCfg{ + ReplicationConfig: ReplicationConfig{ + ID: replicationID, + Direction: ActiveReplicatorTypePush, + Remote: "http://localhost:4984/remotedb", + CollectionsEnabled: !testDB.OnlyDefaultCollection(), + }, + AssignedNode: localNodeUUID, + TargetState: ReplicationStateStopped, + })) + // Initializes the replicator without starting it, so the test controls its state. + require.NoError(t, mgr.RefreshReplicationCfg(ctx)) + replicator := mgr.GetActiveReplicator(replicationID) + require.NotNil(t, replicator) + + testCases := []struct { + state string + active bool + }{ + {state: ReplicationStateRunning, active: true}, + {state: ReplicationStateStarting, active: true}, + {state: ReplicationStateReconnecting, active: true}, + {state: ReplicationStateStopped, active: false}, + {state: ReplicationStateError, active: false}, + } + for _, tc := range testCases { + t.Run(tc.state, func(t *testing.T) { + replicator.Push.setState(tc.state) + + options := DefaultReplicationStatusOptions() + options.ActiveOnly = true + status, err := mgr.GetReplicationStatus(ctx, replicationID, options) + require.NoError(t, err) + if tc.active { + require.NotNil(t, status) + require.Equal(t, tc.state, status.Status) + } else { + require.Nil(t, status) + } }) } } diff --git a/docs/api/components/parameters.yaml b/docs/api/components/parameters.yaml index 2ada68bddb..d850b7897c 100644 --- a/docs/api/components/parameters.yaml +++ b/docs/api/components/parameters.yaml @@ -300,7 +300,7 @@ replication-active-only: schema: type: boolean default: false - description: Only return replications that are actively running or starting up (`state=running` or `state=starting`). + description: Only return replications that are running or on their way to running (`state=running`, `state=starting` or `state=reconnecting`). replication-include-config: name: includeConfig in: query diff --git a/rest/replicatortest/replicator_test.go b/rest/replicatortest/replicator_test.go index 931109359f..f4a7eb7093 100644 --- a/rest/replicatortest/replicator_test.go +++ b/rest/replicatortest/replicator_test.go @@ -1453,6 +1453,47 @@ func TestGetStatusWithReplication(t *testing.T) { require.Len(t, status.Databases["db"].ReplicationStatus, 0) } +// TestReplicationStatusActiveOnly asserts that activeOnly=true returns a replication that is reconnecting +// and leaves out one that is stopped. +func TestReplicationStatusActiveOnly(t *testing.T) { + rt := rest.NewRestTester(t, &rest.RestTesterConfig{SgReplicateEnabled: true}) + defer rt.Close() + + // A closed server refuses connections, which keeps the replicator reconnecting. + unreachable := httptest.NewServer(http.NotFoundHandler()) + unreachable.Close() + remoteURL := unreachable.URL + "/db" + + const ( + reconnectingID = "reconnecting" + stoppedID = "stopped" + ) + for id, initialState := range map[string]string{reconnectingID: db.ReplicationStateRunning, stoppedID: db.ReplicationStateStopped} { + replicationConfig := db.ReplicationConfig{ + ID: id, + Remote: remoteURL, + Direction: db.ActiveReplicatorTypePull, + Continuous: true, + InitialState: initialState, + CollectionsEnabled: !rt.GetDatabase().OnlyDefaultCollection(), + } + response := rt.SendAdminRequest(http.MethodPut, "/{{.db}}/_replication/"+id, rest.MarshalConfig(t, replicationConfig)) + rest.RequireStatus(t, response, http.StatusCreated) + } + rt.WaitForReplicationStatus(reconnectingID, db.ReplicationStateReconnecting) + rt.WaitForReplicationStatus(stoppedID, db.ReplicationStateStopped) + + statusIDs := func(queryString string) []string { + var ids []string + for _, status := range rt.GetReplicationStatuses(queryString) { + ids = append(ids, status.ID) + } + return ids + } + require.ElementsMatch(t, []string{reconnectingID, stoppedID}, statusIDs("")) + require.ElementsMatch(t, []string{reconnectingID}, statusIDs("?activeOnly=true")) +} + func TestRequireReplicatorStoppedBeforeUpsert(t *testing.T) { base.SetUpTestLogging(t, base.LevelInfo, base.KeyHTTP, base.KeyHTTPResp) From 24d6f7a3173f93c1aeb8808e09ee76a7bd33ad40 Mon Sep 17 00:00:00 2001 From: Tor Colvin Date: Mon, 28 Sep 2026 14:36:40 -0400 Subject: [PATCH 3/5] CBG-5838: exercise replication status through the production paths - Create replications with UpsertReplication and let rebalance assign them, with the other node joining first where it owns the replication. - Move the unset target state case into its own test. - Run activeOnly against a real passive node, so running and reconnecting both come from a replicator, and drop the test that forced states. Co-Authored-By: Claude Opus 5.5 (1M context) --- db/sg_replicate_cfg_test.go | 119 ++++++++++++------------- rest/replicatortest/replicator_test.go | 78 ++++++++-------- 2 files changed, 98 insertions(+), 99 deletions(-) diff --git a/db/sg_replicate_cfg_test.go b/db/sg_replicate_cfg_test.go index 3305b80397..a2d84fef53 100644 --- a/db/sg_replicate_cfg_test.go +++ b/db/sg_replicate_cfg_test.go @@ -15,7 +15,6 @@ import ( "maps" "net/http" "net/http/httptest" - "strings" "sync" "sync/atomic" "testing" @@ -1379,45 +1378,59 @@ func TestSGReplicateManagerStopDrainsClusterUpdates(t *testing.T) { // initialized its replicator yet is reported as starting rather than running. Nothing has connected to the // remote at this point, and callers waiting for running would otherwise proceed while a start is still to come. func TestReplicationStatusBeforeReplicatorStarts(t *testing.T) { - testDB, ctx := SetupTestDB(t) - defer testDB.Close(ctx) - const ( localNodeUUID = "localNode" otherNodeUUID = "otherNode" ) - mgr, err := NewSGReplicateManager(ctx, testDB.DatabaseContext, testDB.CfgSG) - require.NoError(t, err) - require.NoError(t, mgr.StartLocalNode(localNodeUUID, nil)) - require.NoError(t, mgr.RegisterNode(otherNodeUUID)) - testCases := []struct { name string - assignedNode string - targetState string + otherNodeOwns bool + initialState string expectedStatus string active bool }{ - {name: "running", assignedNode: localNodeUUID, targetState: ReplicationStateRunning, expectedStatus: ReplicationStateStarting, active: true}, - {name: "stopped", assignedNode: localNodeUUID, targetState: ReplicationStateStopped, expectedStatus: ReplicationStateStopped, active: false}, - // startup treats an unset target state as running - {name: "unset", assignedNode: localNodeUUID, targetState: "", expectedStatus: ReplicationStateStarting, active: true}, - {name: "other node running", assignedNode: otherNodeUUID, targetState: ReplicationStateRunning, expectedStatus: ReplicationStateStarting, active: true}, - {name: "other node stopped", assignedNode: otherNodeUUID, targetState: ReplicationStateStopped, expectedStatus: ReplicationStateStopped, active: false}, + {name: "running", initialState: ReplicationStateRunning, expectedStatus: ReplicationStateStarting, active: true}, + {name: "stopped", initialState: ReplicationStateStopped, expectedStatus: ReplicationStateStopped, active: false}, + {name: "other node running", otherNodeOwns: true, initialState: ReplicationStateRunning, expectedStatus: ReplicationStateStarting, active: true}, + {name: "other node stopped", otherNodeOwns: true, initialState: ReplicationStateStopped, expectedStatus: ReplicationStateStopped, active: false}, } for _, tc := range testCases { t.Run(tc.name, func(t *testing.T) { - replicationID := "rep_" + strings.ReplaceAll(tc.name, " ", "_") - require.NoError(t, mgr.AddReplication(&ReplicationCfg{ - ReplicationConfig: ReplicationConfig{ - ID: replicationID, - Direction: ActiveReplicatorTypePush, - Remote: "http://localhost:4984/remotedb", - }, - AssignedNode: tc.assignedNode, - TargetState: tc.targetState, - })) + testDB, ctx := SetupTestDB(t) + defer testDB.Close(ctx) + + mgr, err := NewSGReplicateManager(ctx, testDB.DatabaseContext, testDB.CfgSG) + require.NoError(t, err) + + // The first node to join owns the replication, and a second node joining does not move it. + expectedNode := localNodeUUID + if tc.otherNodeOwns { + expectedNode = otherNodeUUID + require.NoError(t, mgr.RegisterNode(otherNodeUUID)) + } else { + require.NoError(t, mgr.StartLocalNode(localNodeUUID, nil)) + } + + const replicationID = "rep" + created, err := mgr.UpsertReplication(ctx, &ReplicationUpsertConfig{ + ID: replicationID, + Remote: base.Ptr("http://localhost:4984/remotedb"), + Direction: base.Ptr(string(ActiveReplicatorTypePush)), + CollectionsEnabled: base.Ptr(!testDB.OnlyDefaultCollection()), + InitialState: base.Ptr(tc.initialState), + }) + require.NoError(t, err) + require.True(t, created) + + if tc.otherNodeOwns { + require.NoError(t, mgr.StartLocalNode(localNodeUUID, nil)) + } else { + require.NoError(t, mgr.RegisterNode(otherNodeUUID)) + } + cluster, err := mgr.GetSGRCluster() + require.NoError(t, err) + require.Equal(t, expectedNode, cluster.Replications[replicationID].AssignedNode) // No RefreshReplicationCfg, so no replicator exists and no status has been published. require.Nil(t, mgr.GetActiveReplicator(replicationID)) @@ -1443,7 +1456,7 @@ func TestReplicationStatusBeforeReplicatorStarts(t *testing.T) { } // An upsert needs a stopped replication, and a replication that is starting is not stopped. - created, err := mgr.UpsertReplication(ctx, &ReplicationUpsertConfig{ID: replicationID, Remote: base.Ptr("http://localhost:4984/otherdb")}) + created, err = mgr.UpsertReplication(ctx, &ReplicationUpsertConfig{ID: replicationID, Remote: base.Ptr("http://localhost:4984/otherdb")}) if tc.active { require.Error(t, err) status, _ := base.ErrorAsHTTPStatus(err) @@ -1456,8 +1469,9 @@ func TestReplicationStatusBeforeReplicatorStarts(t *testing.T) { } } -// TestReplicationStatusActiveOnlyForReplicatorStates asserts which states of a local replicator activeOnly includes. -func TestReplicationStatusActiveOnlyForReplicatorStates(t *testing.T) { +// TestReplicationStatusUnsetTargetState asserts that a cfg without a target state reports like a running one, +// since startup treats it as running. No API writes this, so the cfg is written directly. +func TestReplicationStatusUnsetTargetState(t *testing.T) { testDB, ctx := SetupTestDB(t) defer testDB.Close(ctx) @@ -1471,43 +1485,20 @@ func TestReplicationStatusActiveOnlyForReplicatorStates(t *testing.T) { require.NoError(t, mgr.StartLocalNode(localNodeUUID, nil)) require.NoError(t, mgr.AddReplication(&ReplicationCfg{ ReplicationConfig: ReplicationConfig{ - ID: replicationID, - Direction: ActiveReplicatorTypePush, - Remote: "http://localhost:4984/remotedb", - CollectionsEnabled: !testDB.OnlyDefaultCollection(), + ID: replicationID, + Direction: ActiveReplicatorTypePush, + Remote: "http://localhost:4984/remotedb", }, AssignedNode: localNodeUUID, - TargetState: ReplicationStateStopped, })) - // Initializes the replicator without starting it, so the test controls its state. - require.NoError(t, mgr.RefreshReplicationCfg(ctx)) - replicator := mgr.GetActiveReplicator(replicationID) - require.NotNil(t, replicator) - testCases := []struct { - state string - active bool - }{ - {state: ReplicationStateRunning, active: true}, - {state: ReplicationStateStarting, active: true}, - {state: ReplicationStateReconnecting, active: true}, - {state: ReplicationStateStopped, active: false}, - {state: ReplicationStateError, active: false}, - } - for _, tc := range testCases { - t.Run(tc.state, func(t *testing.T) { - replicator.Push.setState(tc.state) + status, err := mgr.GetReplicationStatus(ctx, replicationID, DefaultReplicationStatusOptions()) + require.NoError(t, err) + require.Equal(t, ReplicationStateStarting, status.Status) - options := DefaultReplicationStatusOptions() - options.ActiveOnly = true - status, err := mgr.GetReplicationStatus(ctx, replicationID, options) - require.NoError(t, err) - if tc.active { - require.NotNil(t, status) - require.Equal(t, tc.state, status.Status) - } else { - require.Nil(t, status) - } - }) - } + options := DefaultReplicationStatusOptions() + options.ActiveOnly = true + status, err = mgr.GetReplicationStatus(ctx, replicationID, options) + require.NoError(t, err) + require.NotNil(t, status) } diff --git a/rest/replicatortest/replicator_test.go b/rest/replicatortest/replicator_test.go index f4a7eb7093..078c280e3d 100644 --- a/rest/replicatortest/replicator_test.go +++ b/rest/replicatortest/replicator_test.go @@ -1453,45 +1453,53 @@ func TestGetStatusWithReplication(t *testing.T) { require.Len(t, status.Databases["db"].ReplicationStatus, 0) } -// TestReplicationStatusActiveOnly asserts that activeOnly=true returns a replication that is reconnecting -// and leaves out one that is stopped. +// TestReplicationStatusActiveOnly asserts that activeOnly=true returns replications that are running or +// reconnecting, and leaves out one that is stopped. func TestReplicationStatusActiveOnly(t *testing.T) { - rt := rest.NewRestTester(t, &rest.RestTesterConfig{SgReplicateEnabled: true}) - defer rt.Close() - - // A closed server refuses connections, which keeps the replicator reconnecting. - unreachable := httptest.NewServer(http.NotFoundHandler()) - unreachable.Close() - remoteURL := unreachable.URL + "/db" + base.RequireNumTestBuckets(t, 2) - const ( - reconnectingID = "reconnecting" - stoppedID = "stopped" - ) - for id, initialState := range map[string]string{reconnectingID: db.ReplicationStateRunning, stoppedID: db.ReplicationStateStopped} { - replicationConfig := db.ReplicationConfig{ - ID: id, - Remote: remoteURL, - Direction: db.ActiveReplicatorTypePull, - Continuous: true, - InitialState: initialState, - CollectionsEnabled: !rt.GetDatabase().OnlyDefaultCollection(), + sgrRunner := rest.NewSGRTestRunner(t) + sgrRunner.Run(func(t *testing.T) { + peers := sgrRunner.SetupSGRPeers(t) + activeRT := peers.ActiveRT + + // A closed server refuses connections, which keeps the replicator reconnecting. + unreachable := httptest.NewServer(http.NotFoundHandler()) + unreachable.Close() + unreachableURL := unreachable.URL + "/db" + + const ( + runningID = "running" + reconnectingID = "reconnecting" + stoppedID = "stopped" + ) + activeRT.CreateReplication(runningID, peers.PassiveDBURL, db.ActiveReplicatorTypePull, nil, true, db.ConflictResolverDefault, "") + for id, initialState := range map[string]string{reconnectingID: db.ReplicationStateRunning, stoppedID: db.ReplicationStateStopped} { + replicationConfig := db.ReplicationConfig{ + ID: id, + Remote: unreachableURL, + Direction: db.ActiveReplicatorTypePull, + Continuous: true, + InitialState: initialState, + CollectionsEnabled: base.TestsUseNamedCollections(), + } + response := activeRT.SendAdminRequest(http.MethodPut, "/{{.db}}/_replication/"+id, rest.MarshalConfig(t, replicationConfig)) + rest.RequireStatus(t, response, http.StatusCreated) } - response := rt.SendAdminRequest(http.MethodPut, "/{{.db}}/_replication/"+id, rest.MarshalConfig(t, replicationConfig)) - rest.RequireStatus(t, response, http.StatusCreated) - } - rt.WaitForReplicationStatus(reconnectingID, db.ReplicationStateReconnecting) - rt.WaitForReplicationStatus(stoppedID, db.ReplicationStateStopped) - - statusIDs := func(queryString string) []string { - var ids []string - for _, status := range rt.GetReplicationStatuses(queryString) { - ids = append(ids, status.ID) + activeRT.WaitForReplicationStatus(runningID, db.ReplicationStateRunning) + activeRT.WaitForReplicationStatus(reconnectingID, db.ReplicationStateReconnecting) + activeRT.WaitForReplicationStatus(stoppedID, db.ReplicationStateStopped) + + statusIDs := func(queryString string) []string { + var ids []string + for _, status := range activeRT.GetReplicationStatuses(queryString) { + ids = append(ids, status.ID) + } + return ids } - return ids - } - require.ElementsMatch(t, []string{reconnectingID, stoppedID}, statusIDs("")) - require.ElementsMatch(t, []string{reconnectingID}, statusIDs("?activeOnly=true")) + require.ElementsMatch(t, []string{runningID, reconnectingID, stoppedID}, statusIDs("")) + require.ElementsMatch(t, []string{runningID, reconnectingID}, statusIDs("?activeOnly=true")) + }) } func TestRequireReplicatorStoppedBeforeUpsert(t *testing.T) { From 776158f9dd8b16ec3bc23b6e74348d4da62f9e58 Mon Sep 17 00:00:00 2001 From: Tor Colvin Date: Mon, 28 Sep 2026 16:46:29 -0400 Subject: [PATCH 4/5] CBG-5838: replace base.Ptr with new in replication status tests Co-Authored-By: Claude Opus 5.5 (1M context) --- db/sg_replicate_cfg_test.go | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/db/sg_replicate_cfg_test.go b/db/sg_replicate_cfg_test.go index fe6eea8555..2f229da0df 100644 --- a/db/sg_replicate_cfg_test.go +++ b/db/sg_replicate_cfg_test.go @@ -1415,10 +1415,10 @@ func TestReplicationStatusBeforeReplicatorStarts(t *testing.T) { const replicationID = "rep" created, err := mgr.UpsertReplication(ctx, &ReplicationUpsertConfig{ ID: replicationID, - Remote: base.Ptr("http://localhost:4984/remotedb"), - Direction: base.Ptr(string(ActiveReplicatorTypePush)), - CollectionsEnabled: base.Ptr(!testDB.OnlyDefaultCollection()), - InitialState: base.Ptr(tc.initialState), + Remote: new("http://localhost:4984/remotedb"), + Direction: new(string(ActiveReplicatorTypePush)), + CollectionsEnabled: new(!testDB.OnlyDefaultCollection()), + InitialState: new(tc.initialState), }) require.NoError(t, err) require.True(t, created) @@ -1456,7 +1456,7 @@ func TestReplicationStatusBeforeReplicatorStarts(t *testing.T) { } // An upsert needs a stopped replication, and a replication that is starting is not stopped. - created, err = mgr.UpsertReplication(ctx, &ReplicationUpsertConfig{ID: replicationID, Remote: base.Ptr("http://localhost:4984/otherdb")}) + created, err = mgr.UpsertReplication(ctx, &ReplicationUpsertConfig{ID: replicationID, Remote: new("http://localhost:4984/otherdb")}) if tc.active { require.Error(t, err) status, _ := base.ErrorAsHTTPStatus(err) From 5f7a1e885a36461cb3ab0de9f037c23aedc5ce47 Mon Sep 17 00:00:00 2001 From: Tor Colvin Date: Fri, 2 Oct 2026 11:33:36 -0400 Subject: [PATCH 5/5] Update docs/api/components/parameters.yaml Co-authored-by: factory-droid[bot] <138933559+factory-droid[bot]@users.noreply.github.com> --- docs/api/components/parameters.yaml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/api/components/parameters.yaml b/docs/api/components/parameters.yaml index d850b7897c..4698ca9a8b 100644 --- a/docs/api/components/parameters.yaml +++ b/docs/api/components/parameters.yaml @@ -300,7 +300,7 @@ replication-active-only: schema: type: boolean default: false - description: Only return replications that are running or on their way to running (`state=running`, `state=starting` or `state=reconnecting`). + description: Only return replications that are running or on their way to running (`status=running`, `status=starting` or `status=reconnecting`). replication-include-config: name: includeConfig in: query