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
5 changes: 5 additions & 0 deletions internal/util/clusterinfo.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package util

import (
"cmp"
"context"

"github.com/10gen/migration-verifier/internal/logger"
Expand All @@ -23,6 +24,10 @@ const (
TopologyReplset ClusterTopology = "replset"
)

func CmpMinorVersions(a, b [2]int) int {
return cmp.Or(cmp.Compare(a[0], b[0]), cmp.Compare(a[1], b[1]))
}

func GetClusterInfo(ctx context.Context, logger *logger.Logger, client *mongo.Client) (ClusterInfo, error) {
va, err := getVersionArray(ctx, client)
if err != nil {
Expand Down
47 changes: 24 additions & 23 deletions internal/verifier/change_stream.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ func (uee UnknownEventError) Error() string {

type changeEventBatch struct {
events []ParsedEvent
resumeToken bson.Raw
clusterTime bson.Timestamp
}

Expand Down Expand Up @@ -127,6 +128,21 @@ func (verifier *Verifier) initializeChangeStreamReaders() {
func (verifier *Verifier) RunChangeEventHandler(ctx context.Context, reader *ChangeStreamReader) error {
var err error

var lastPersistedTime time.Time
persistResumeTokenIfNeeded := func(ctx context.Context, token bson.Raw) {
if time.Since(lastPersistedTime) >= minChangeStreamPersistInterval {
persistErr := reader.persistChangeStreamResumeToken(ctx, token)
if persistErr != nil {
verifier.logger.Warn().
Stringer("changeReader", reader).
Err(persistErr).
Msg("Failed to persist resume token. Because of this, if the verifier restarts, it will have to re-process already-handled change events. This error may be transient, but if it recurs, investigate.")
} else {
lastPersistedTime = time.Now()
}
}
}

HandlerLoop:
for err == nil {
select {
Expand Down Expand Up @@ -156,6 +172,10 @@ HandlerLoop:
verifier.HandleChangeStreamEvents(ctx, batch, reader.readerType),
"failed to handle change stream events",
)

if err == nil && batch.resumeToken != nil {
persistResumeTokenIfNeeded(ctx, batch.resumeToken)
}
}
}

Expand Down Expand Up @@ -470,6 +490,8 @@ func (csr *ChangeStreamReader) readAndHandleOneChangeEventBatch(
case csr.changeEventBatchChan <- changeEventBatch{
events: changeEvents,

resumeToken: cs.ResumeToken(),

// NB: We know by now that OperationTime is non-nil.
clusterTime: *sess.OperationTime(),
}:
Expand All @@ -493,21 +515,6 @@ func (csr *ChangeStreamReader) iterateChangeStream(
cs *mongo.ChangeStream,
sess *mongo.Session,
) error {
var lastPersistedTime time.Time

persistResumeTokenIfNeeded := func() error {
if time.Since(lastPersistedTime) <= minChangeStreamPersistInterval {
return nil
}

err := csr.persistChangeStreamResumeToken(ctx, cs)
if err == nil {
lastPersistedTime = time.Now()
}

return err
}

for {
var err error
var gotwritesOffTimestamp bool
Expand Down Expand Up @@ -571,10 +578,6 @@ func (csr *ChangeStreamReader) iterateChangeStream(
default:
err = csr.readAndHandleOneChangeEventBatch(ctx, ri, cs, sess)

if err == nil {
err = persistResumeTokenIfNeeded()
}

if err != nil {
return err
}
Expand Down Expand Up @@ -659,7 +662,7 @@ func (csr *ChangeStreamReader) createChangeStream(
return nil, nil, bson.Timestamp{}, errors.Wrap(err, "failed to open change stream")
}

err = csr.persistChangeStreamResumeToken(ctx, changeStream)
err = csr.persistChangeStreamResumeToken(ctx, changeStream.ResumeToken())
if err != nil {
return nil, nil, bson.Timestamp{}, err
}
Expand Down Expand Up @@ -852,9 +855,7 @@ func (csr *ChangeStreamReader) resumeTokenDocID() string {
}
}

func (csr *ChangeStreamReader) persistChangeStreamResumeToken(ctx context.Context, cs *mongo.ChangeStream) error {
token := cs.ResumeToken()

func (csr *ChangeStreamReader) persistChangeStreamResumeToken(ctx context.Context, token bson.Raw) error {
coll := csr.getChangeStreamMetadataCollection()
_, err := coll.ReplaceOne(
ctx,
Expand Down
148 changes: 148 additions & 0 deletions internal/verifier/change_stream_test.go
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
package verifier

import (
"bytes"
"context"
"fmt"
"io"
"strings"
"time"
Expand All @@ -17,10 +19,12 @@ import (
"github.com/pkg/errors"
"github.com/rs/zerolog"
"github.com/samber/lo"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.mongodb.org/mongo-driver/v2/bson"
"go.mongodb.org/mongo-driver/v2/mongo"
"go.mongodb.org/mongo-driver/v2/mongo/options"
"go.mongodb.org/mongo-driver/v2/mongo/readconcern"
)

func (suite *IntegrationTestSuite) TestChangeStreamFilter_NoNamespaces() {
Expand Down Expand Up @@ -245,6 +249,150 @@ func (suite *IntegrationTestSuite) startSrcChangeStreamReaderAndHandler(ctx cont
}()
}

func (suite *IntegrationTestSuite) TestChangeStream_Resume_NoSkip() {
ctx := suite.T().Context()

verifier1 := suite.BuildVerifier()

// Use of linearizable read concern below seems to freeze pre-4.4 servers.
srcVersion := verifier1.srcClusterInfo.VersionArray
if util.CmpMinorVersions([2]int(srcVersion), [2]int{4, 4}) == -1 {
suite.T().Skipf("Source version (%v) is too old for this test.", srcVersion)
}

srcDB := verifier1.srcClient.Database(suite.DBNameForTest())
srcColl := srcDB.Collection("coll")

require.NoError(
suite.T(),
srcDB.CreateCollection(ctx, srcColl.Name()),
)

v1Ctx, v1Cancel := contextplus.WithCancelCause(ctx)
defer v1Cancel(ctx.Err())
suite.startSrcChangeStreamReaderAndHandler(v1Ctx, verifier1)

changeStreamMetaColl := verifier1.srcChangeStreamReader.getChangeStreamMetadataCollection()

var originalResumeToken bson.Raw

assert.Eventually(
suite.T(),
func() bool {
var err error
originalResumeToken, err = changeStreamMetaColl.FindOne(ctx, bson.D{}).Raw()
return err == nil
},
time.Minute,
50*time.Millisecond,
"should see a change stream resume token persisted",
)

insertCtx, cancelInserts := contextplus.WithCancelCause(ctx)
defer cancelInserts(ctx.Err())
insertsDone := make(chan struct{})
go func() {
defer func() {
close(insertsDone)
}()

sess, err := verifier1.srcClient.StartSession(
options.Session().SetCausalConsistency(true),
)
require.NoError(suite.T(), err)

sessCtx := mongo.NewSessionContext(insertCtx, sess)

docID := int32(1)
for {
_, err := srcColl.InsertOne(
sessCtx,
bson.D{{"_id", docID}},
)

if err != nil {
require.ErrorIs(suite.T(), err, context.Canceled)
return
}

docID++
}
}()

assert.Eventually(
suite.T(),
func() bool {
rt, err := changeStreamMetaColl.FindOne(ctx, bson.D{}).Raw()
require.NoError(suite.T(), err)

suite.T().Logf("found rt: %v\n", rt)

return !bytes.Equal(rt, originalResumeToken)
},
time.Minute,
50*time.Millisecond,
"should see a new change stream resume token persisted",
)

v1Cancel(fmt.Errorf("killing verifier1"))

verifier2 := suite.BuildVerifier()
suite.startSrcChangeStreamReaderAndHandler(ctx, verifier2)

cancelInserts(fmt.Errorf("verifier2 started"))
<-insertsDone

lastIDRes := srcColl.Database().Collection(
srcColl.Name(),
options.Collection().SetReadConcern(readconcern.Linearizable()),
).FindOne(
ctx,
bson.D{},
options.FindOne().
SetSort(bson.D{{"_id", -1}}),
)
require.NoError(suite.T(), lastIDRes.Err())

lastDocID := lo.Must(lo.Must(lastIDRes.Raw()).LookupErr("_id")).Int32()

assert.Eventually(
suite.T(),
func() bool {
rechecks := suite.fetchVerifierRechecks(ctx, verifier2)

return lo.SomeBy(
rechecks,
func(cur bson.M) bool {
id := cur["_id"].(bson.D)

for _, el := range id {
if el.Key != "docID" {
continue
}

return el.Value.(int32) == lastDocID
}

panic(fmt.Sprintf("no docID in _id: %+v", id))
},
)
},
time.Minute,
100*time.Millisecond,
"last-inserted doc shows as recheck",
)

sess := lo.Must(verifier2.verificationDatabase().Client().StartSession())
sctx := mongo.NewSessionContext(ctx, sess)

rechecks := suite.fetchVerifierRechecks(sctx, verifier2)
if !assert.EqualValues(suite.T(), lastDocID, len(rechecks), "all source docs should be rechecked") {
for _, recheck := range rechecks {
suite.T().Logf("found recheck: %v", recheck)
}
}
}

// TestChangeStreamResumability creates a verifier, starts its change stream,
// terminates that verifier, updates the source cluster, starts a new
// verifier with change stream, and confirms that things look as they should.
Expand Down
2 changes: 0 additions & 2 deletions vendor/modules.txt
Original file line number Diff line number Diff line change
Expand Up @@ -150,8 +150,6 @@ github.com/modern-go/concurrent
# github.com/modern-go/reflect2 v1.0.2
## explicit; go 1.12
github.com/modern-go/reflect2
# github.com/montanaflynn/stats v0.7.1
## explicit; go 1.13
# github.com/olekukonko/tablewriter v0.0.5
## explicit; go 1.12
github.com/olekukonko/tablewriter
Expand Down