summaryrefslogtreecommitdiff
path: root/src/mongo/db/pipeline/document_source_change_stream_test.cpp
diff options
context:
space:
mode:
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.cpp324
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