diff options
Diffstat (limited to 'src/mongo/db/pipeline/document_source_change_stream_test.cpp')
| -rw-r--r-- | src/mongo/db/pipeline/document_source_change_stream_test.cpp | 324 |
1 files changed, 321 insertions, 3 deletions
diff --git a/src/mongo/db/pipeline/document_source_change_stream_test.cpp b/src/mongo/db/pipeline/document_source_change_stream_test.cpp index 50e0f0cdcac..8bcdb7e7f9d 100644 --- a/src/mongo/db/pipeline/document_source_change_stream_test.cpp +++ b/src/mongo/db/pipeline/document_source_change_stream_test.cpp @@ -27,6 +27,7 @@ * it in the license file. */ +#include "mongo/bson/bsontypes.h" #include "mongo/platform/basic.h" #include <boost/intrusive_ptr.hpp> @@ -52,8 +53,11 @@ #include "mongo/db/pipeline/document_source_change_stream_add_pre_image.h" #include "mongo/db/pipeline/document_source_change_stream_check_invalidate.h" #include "mongo/db/pipeline/document_source_change_stream_check_resumability.h" +#include "mongo/db/pipeline/document_source_change_stream_check_topology_change.h" #include "mongo/db/pipeline/document_source_change_stream_ensure_resume_token_present.h" +#include "mongo/db/pipeline/document_source_change_stream_handle_topology_change.h" #include "mongo/db/pipeline/document_source_change_stream_oplog_match.h" +#include "mongo/db/pipeline/document_source_change_stream_split_large_event.h" #include "mongo/db/pipeline/document_source_change_stream_transform.h" #include "mongo/db/pipeline/document_source_change_stream_unwind_transaction.h" #include "mongo/db/pipeline/document_source_limit.h" @@ -88,15 +92,26 @@ using V = Value; using DSChangeStream = DocumentSourceChangeStream; +// Deterministic values used for testing +const UUID testConstUuid = UUID::parse("6948DF80-14BD-4E04-8842-7668D9C001F5").getValue(); + +class ExecutableStubMongoProcessInterface : public StubMongoProcessInterface { + bool isExpectedToExecuteQueries() override { + return true; + } +}; + class ChangeStreamStageTestNoSetup : public AggregationContextFixture { public: ChangeStreamStageTestNoSetup() : ChangeStreamStageTestNoSetup(nss) {} explicit ChangeStreamStageTestNoSetup(NamespaceString nsString) - : AggregationContextFixture(nsString) {} + : AggregationContextFixture(nsString) { + getExpCtx()->mongoProcessInterface = + std::make_unique<ExecutableStubMongoProcessInterface>(); + }; }; -struct MockMongoInterface final : public StubMongoProcessInterface { - +struct MockMongoInterface final : public ExecutableStubMongoProcessInterface { // Used by operations which need to obtain the oplog's UUID. static const UUID& oplogUuid() { static const UUID* oplog_uuid = new UUID(UUID::gen()); @@ -4648,5 +4663,308 @@ TEST_F(MultiTokenFormatVersionTest, CanResumeFromV2HighWaterMark) { next = lastStage->getNext(); ASSERT_FALSE(next.isAdvanced()); } + +TEST_F(ChangeStreamStageTestNoSetup, DocumentSourceChangeStreamAddPostImageEmptyForQueryStats) { + auto spec = DocumentSourceChangeStreamSpec(); + spec.setFullDocument(FullDocumentModeEnum::kUpdateLookup); + + auto docSource = DocumentSourceChangeStreamAddPostImage::create(getExpCtx(), spec); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({"$_internalChangeStreamAddPostImage":{"fullDocument":"updateLookup"}})", + docSource->serialize().getDocument().toBson()); + + auto opts = SerializationOptions{LiteralSerializationPolicy::kToRepresentativeParseableValue}; + ASSERT(docSource->serialize(opts).missing()); +} + +TEST_F(ChangeStreamStageTestNoSetup, DocumentSourceChangeStreamAddPreImageEmptyForQueryStats) { + auto docSource = DocumentSourceChangeStreamAddPreImage{ + getExpCtx(), FullDocumentBeforeChangeModeEnum::kWhenAvailable}; + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({ + "$_internalChangeStreamAddPreImage": { + "fullDocumentBeforeChange": "whenAvailable" + } + })", + docSource.serialize().getDocument().toBson()); + + + auto opts = SerializationOptions{LiteralSerializationPolicy::kToRepresentativeParseableValue}; + ASSERT(docSource.serialize(opts).missing()); +} + +TEST_F(ChangeStreamStageTestNoSetup, DocumentSourceChangeStreamCheckInvalidateEmptyForQueryStats) { + DocumentSourceChangeStreamSpec spec; + spec.setResumeAfter(ResumeToken::parse(makeResumeToken(Timestamp(), + testConstUuid, + BSON("_id" << 1 << "x" << 2), + ResumeTokenData::kFromInvalidate))); + + auto docSource = DocumentSourceChangeStreamCheckInvalidate::create(getExpCtx(), spec); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({ + "$_internalChangeStreamCheckInvalidate": { + "startAfterInvalidate": { + "_data": "8200000000000000002B022C0100296F5A10046948DF8014BD4E0488427668D9C001F5461E5F6964002B021E78002B040004" + } + } + })", + docSource->serialize().getDocument().toBson()); + + auto opts = SerializationOptions{LiteralSerializationPolicy::kToRepresentativeParseableValue}; + ASSERT(docSource->serialize(opts).missing()); +} + +TEST_F(ChangeStreamStageTestNoSetup, + DocumentSourceChangeStreamCheckResumabilityEmptyForQueryStats) { + DocumentSourceChangeStreamSpec spec; + spec.setResumeAfter(ResumeToken::parse( + makeResumeToken(Timestamp(), testConstUuid, BSON("_id" << 1 << "x" << 2)))); + + auto docSource = DocumentSourceChangeStreamCheckResumability::create(getExpCtx(), spec); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({ + "$_internalChangeStreamCheckResumability": { + "resumeToken": { + "_data": "8200000000000000002B022C0100296E5A10046948DF8014BD4E0488427668D9C001F5461E5F6964002B021E78002B040004" + } + } + })", + docSource->serialize().getDocument().toBson()); + + auto opts = SerializationOptions{LiteralSerializationPolicy::kToRepresentativeParseableValue}; + ASSERT(docSource->serialize(opts).missing()); +} + +TEST_F(ChangeStreamStageTestNoSetup, + DocumentSourceChangeStreamCheckTopologyChangeEmptyForQueryStats) { + auto docSource = DocumentSourceChangeStreamCheckTopologyChange::create(getExpCtx()); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({"$_internalChangeStreamCheckTopologyChange":{}})", + docSource->serialize().getDocument().toBson()); + + auto opts = SerializationOptions{LiteralSerializationPolicy::kToRepresentativeParseableValue}; + ASSERT(docSource->serialize(opts).missing()); +} + +TEST_F(ChangeStreamStageTestNoSetup, + DocumentSourceChangeStreamEnsureResumeTokenPresentEmptyForQueryStats) { + DocumentSourceChangeStreamSpec spec; + spec.setResumeAfter(ResumeToken::parse( + makeResumeToken(Timestamp(), testConstUuid, BSON("_id" << 1 << "x" << 2)))); + + auto docSource = DocumentSourceChangeStreamEnsureResumeTokenPresent::create(getExpCtx(), spec); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({ + "$_internalChangeStreamEnsureResumeTokenPresent": { + "resumeToken": { + "_data": "8200000000000000002B022C0100296E5A10046948DF8014BD4E0488427668D9C001F5461E5F6964002B021E78002B040004" + } + } + })", + docSource->serialize().getDocument().toBson()); + + auto opts = SerializationOptions{LiteralSerializationPolicy::kToRepresentativeParseableValue}; + ASSERT(docSource->serialize(opts).missing()); +} + +TEST_F(ChangeStreamStageTestNoSetup, + DocumentSourceChangeStreamHandleTopologyChangeEmptyForQueryStats) { + auto docSource = DocumentSourceChangeStreamHandleTopologyChange::create(getExpCtx()); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({"$_internalChangeStreamHandleTopologyChange":{}})", + docSource->serialize().getDocument().toBson()); + + auto opts = SerializationOptions{LiteralSerializationPolicy::kToRepresentativeParseableValue}; + ASSERT(docSource->serialize(opts).missing()); +} + +TEST_F(ChangeStreamStageTestNoSetup, RedactDocumentSourceChangeStreamSplitLargeEvent) { + DocumentSourceChangeStreamSpec spec; + spec.setResumeAfter(ResumeToken::parse( + makeResumeToken(Timestamp(), testConstUuid, BSON("_id" << 1 << "x" << 2)))); + + auto docSource = DocumentSourceChangeStreamSplitLargeEvent::create(getExpCtx(), spec); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({"$changeStreamSplitLargeEvent":{}})", + docSource->serialize().getDocument().toBson()); + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({"$changeStreamSplitLargeEvent":{}})", + redact(*docSource)); +} + +TEST_F(ChangeStreamStageTestNoSetup, RedactDocumentSourceChangeStreamTransform) { + DocumentSourceChangeStreamSpec spec; + spec.setResumeAfter(ResumeToken::parse( + makeResumeToken(Timestamp(), testConstUuid, BSON("_id" << 1 << "x" << 2)))); + + auto docSource = DocumentSourceChangeStreamTransform::create(getExpCtx(), spec); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({ + "$_internalChangeStreamTransform": { + "resumeAfter": { + "_data": "8200000000000000002B022C0100296E5A10046948DF8014BD4E0488427668D9C001F5461E5F6964002B021E78002B040004" + }, + "fullDocument": "default", + "fullDocumentBeforeChange": "off" + } + })", + docSource->serialize().getDocument().toBson()); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({ + "$changeStream": { + "resumeAfter": { + "_data": "?string" + }, + "fullDocument": "default", + "fullDocumentBeforeChange": "off" + } + })", + redact(*docSource)); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({ + "$changeStream": { + "resumeAfter": { + "_data": "8200000000000000002B0229296E04" + }, + "fullDocument": "default", + "fullDocumentBeforeChange": "off" + } + })", + docSource + ->serialize( + SerializationOptions{LiteralSerializationPolicy::kToRepresentativeParseableValue}) + .getDocument() + .toBson()); +} + + +TEST_F(ChangeStreamStageTestNoSetup, RedactDocumentSourceChangeStreamTransformMoreFields) { + DocumentSourceChangeStreamSpec spec; + spec.setStartAfter(ResumeToken::parse( + makeResumeToken(Timestamp(), testConstUuid, BSON("_id" << 1 << "x" << 2)))); + spec.setFullDocument(FullDocumentModeEnum::kRequired); + spec.setFullDocumentBeforeChange(FullDocumentBeforeChangeModeEnum::kWhenAvailable); + spec.setShowExpandedEvents(true); + + auto docSource = DocumentSourceChangeStreamTransform::create(getExpCtx(), spec); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({ + "$_internalChangeStreamTransform": { + "startAfter": { + "_data": "8200000000000000002B022C0100296E5A10046948DF8014BD4E0488427668D9C001F5461E5F6964002B021E78002B040004" + }, + "fullDocument": "required", + "fullDocumentBeforeChange": "whenAvailable", + "showExpandedEvents": true + } + })", + docSource->serialize().getDocument().toBson()); + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({ + "$changeStream": { + "startAfter": { + "_data": "?string" + }, + "fullDocument": "required", + "fullDocumentBeforeChange": "whenAvailable", + "showExpandedEvents": true + } + })", + redact(*docSource)); + + ASSERT_BSONOBJ_EQ_AUTO( // NOLINT + R"({ + "$changeStream": { + "startAfter": { + "_data": "8200000000000000002B0229296E04" + }, + "fullDocument": "required", + "fullDocumentBeforeChange": "whenAvailable", + "showExpandedEvents": true + } + })", + docSource + ->serialize( + SerializationOptions{LiteralSerializationPolicy::kToRepresentativeParseableValue}) + .getDocument() + .toBson()); +} + +// For DocumentSource types which contain an arbitrarily internal +// MatchExpression, we don't want match the entire structure. This +// assertion allows us to check some basic structure. +void assertRedactedMatchExpressionContainsOperatorsAndRedactedFieldPaths(BSONElement el) { + // Walk the redacted BSON and assert that we have some ops and + // redacted field paths. + auto opCount = 0; + auto redactedFieldPaths = 0; + while (true) { + if (el.type() == mongo::Array) { + auto array = el.Array(); + if (array.empty()) { + break; + } + el = array[0]; + } else if (el.type() == mongo::Object) { + auto obj = el.Obj(); + if (obj.begin() == obj.end()) { + break; + } + el = obj.firstElement(); + + // Field name should be an operator or a redacted field path. + if (el.fieldName()[0] == '$') { + opCount++; + } else if (!strcmp(el.fieldName(), "$regularExpression")) { + opCount++; + // Skip $regularExpression. + continue; + } else { + if (strstr(el.fieldName(), "HASH<") != el.fieldName()) { + FAIL(std::string("Expected redacted field path: ") + el.fieldName()); + } + redactedFieldPaths++; + } + } else { + break; + } + } + + ASSERT(opCount > 0); + ASSERT(redactedFieldPaths > 0); +} + +TEST_F(ChangeStreamStageTestNoSetup, + DocumentSourceChangeStreamUnwindTransactionEmptyForQueryStats) { + auto docSource = DocumentSourceChangeStreamUnwindTransaction::create(getExpCtx()); + + auto opts = SerializationOptions{LiteralSerializationPolicy::kToRepresentativeParseableValue}; + ASSERT(docSource->serialize(opts).missing()); +} + +TEST_F(ChangeStreamStageTestNoSetup, DocumentSourceChangeStreamOplogMatchEmptyForQueryStats) { + DocumentSourceChangeStreamSpec spec; + spec.setResumeAfter(ResumeToken::parse( + makeResumeToken(Timestamp(), testConstUuid, BSON("_id" << 1 << "x" << 2)))); + + auto docSource = DocumentSourceChangeStreamOplogMatch::create(getExpCtx(), spec); + + auto opts = SerializationOptions{LiteralSerializationPolicy::kToRepresentativeParseableValue}; + ASSERT(docSource->serialize(opts).missing()); +} + } // namespace } // namespace mongo |
