Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 12 additions & 1 deletion db/sg_replicate_cfg.go
Original file line number Diff line number Diff line change
Expand Up @@ -1628,6 +1628,12 @@ func (m *sgReplicateManager) GetReplicationStatus(ctx context.Context, replicati
ID: replicationID,
Status: ReplicationStateUnassigned,
}
} else if remoteCfg.TargetState == "" || remoteCfg.TargetState == ReplicationStateRunning {
// Nothing is replicating until a replicator says so.
status = &ReplicationStatus{
ID: replicationID,
Status: ReplicationStateStarting,
}
} else {
status = &ReplicationStatus{
ID: replicationID,
Expand All @@ -1652,13 +1658,18 @@ func (m *sgReplicateManager) GetReplicationStatus(ctx context.Context, replicati
if !options.IncludeError && status.Status == ReplicationStateError {
return nil, nil
}
if options.ActiveOnly && status.Status != ReplicationStateRunning {
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) {

Expand Down
129 changes: 129 additions & 0 deletions db/sg_replicate_cfg_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1373,3 +1373,132 @@ 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) {
const (
localNodeUUID = "localNode"
otherNodeUUID = "otherNode"
)

testCases := []struct {
name string
otherNodeOwns bool
initialState string
expectedStatus string
active bool
}{
{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) {
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: 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)

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))

status, err := mgr.GetReplicationStatus(ctx, replicationID, DefaultReplicationStatusOptions())
require.NoError(t, err)
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.
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 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: new("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)
}
})
}
}

// 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)

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",
},
AssignedNode: localNodeUUID,
}))

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)
require.NotNil(t, status)
}
2 changes: 1 addition & 1 deletion docs/api/components/parameters.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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 running or on their way to running (`status=running`, `status=starting` or `status=reconnecting`).
replication-include-config:
name: includeConfig
in: query
Expand Down
2 changes: 2 additions & 0 deletions docs/api/components/schemas.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -668,13 +668,15 @@ ISGRReplicationState:
- error
- starting
- reconnecting
- unassigned
x-enumDescriptions:
running: Currently running replication.
stopped: Not running replication.
resetting: The replication is resetting its state.
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
Expand Down
49 changes: 49 additions & 0 deletions rest/replicatortest/replicator_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1488,6 +1488,55 @@ func TestGetStatusWithReplication(t *testing.T) {
require.Len(t, status.Databases["db"].ReplicationStatus, 0)
}

// TestReplicationStatusActiveOnly asserts that activeOnly=true returns replications that are running or
// reconnecting, and leaves out one that is stopped.
func TestReplicationStatusActiveOnly(t *testing.T) {
base.RequireNumTestBuckets(t, 2)

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)
}
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
}
require.ElementsMatch(t, []string{runningID, reconnectingID, stoppedID}, statusIDs(""))
require.ElementsMatch(t, []string{runningID, reconnectingID}, statusIDs("?activeOnly=true"))
})
}

func TestRequireReplicatorStoppedBeforeUpsert(t *testing.T) {
base.SetUpTestLogging(t, base.LevelInfo, base.KeyHTTP, base.KeyHTTPResp)

Expand Down
Loading