summaryrefslogtreecommitdiff
path: root/src/mongo/db/pipeline/document_source_change_stream_transform.cpp
diff options
context:
space:
mode:
authorLucas de Castro Borges <lucas@gnuabordo.com.br>2025-02-18 17:02:53 -0300
committerLucas de Castro Borges <lucas@gnuabordo.com.br>2025-02-18 17:02:53 -0300
commit959575a5ca598bf5f37fb5cebe7ed1d80d3d71f7 (patch)
treeacc8d60aedb12b70048e676e8a7349deb0010db8 /src/mongo/db/pipeline/document_source_change_stream_transform.cpp
parent76588293975fc059cf076779e4283e6ffaf8afff (diff)
New upstream version 6.0.20upstream
Diffstat (limited to 'src/mongo/db/pipeline/document_source_change_stream_transform.cpp')
-rw-r--r--src/mongo/db/pipeline/document_source_change_stream_transform.cpp114
1 files changed, 107 insertions, 7 deletions
diff --git a/src/mongo/db/pipeline/document_source_change_stream_transform.cpp b/src/mongo/db/pipeline/document_source_change_stream_transform.cpp
index 7434190a609..3ea1312f51c 100644
--- a/src/mongo/db/pipeline/document_source_change_stream_transform.cpp
+++ b/src/mongo/db/pipeline/document_source_change_stream_transform.cpp
@@ -73,7 +73,8 @@ DocumentSourceChangeStreamTransform::createFromBson(
DocumentSourceChangeStreamTransform::DocumentSourceChangeStreamTransform(
const boost::intrusive_ptr<ExpressionContext>& expCtx, DocumentSourceChangeStreamSpec spec)
- : DocumentSource(DocumentSourceChangeStreamTransform::kStageName, expCtx),
+ : DocumentSourceInternalChangeStreamStage(DocumentSourceChangeStreamTransform::kStageName,
+ expCtx),
_changeStreamSpec(std::move(spec)),
_transformer(expCtx, _changeStreamSpec),
_isIndependentOfAnyCollection(expCtx->ns.isCollectionlessAggregateNS()) {
@@ -103,16 +104,115 @@ StageConstraints DocumentSourceChangeStreamTransform::constraints(
return constraints;
}
-Value DocumentSourceChangeStreamTransform::serialize(
- boost::optional<ExplainOptions::Verbosity> explain) const {
- if (explain) {
+namespace {
+
+template <typename T>
+void serializeSpecField(BSONObjBuilder* builder,
+ const SerializationOptions& opts,
+ const StringData& fieldName,
+ const boost::optional<T>& value) {
+ if (value) {
+ opts.serializeLiteral((*value).toBSON()).addToBsonObj(builder, fieldName);
+ }
+}
+
+template <>
+void serializeSpecField(BSONObjBuilder* builder,
+ const SerializationOptions& opts,
+ const StringData& fieldName,
+ const boost::optional<Timestamp>& value) {
+ if (value) {
+ opts.serializeLiteral(*value).addToBsonObj(builder, fieldName);
+ }
+}
+
+template <typename T>
+void serializeSpecField(BSONObjBuilder* builder,
+ const SerializationOptions& opts,
+ const StringData& fieldName,
+ const T& value) {
+ opts.appendLiteral(builder, fieldName, value);
+}
+
+template <>
+void serializeSpecField(BSONObjBuilder* builder,
+ const SerializationOptions& opts,
+ const StringData& fieldName,
+ const mongo::OptionalBool& value) {
+ if (value.has_value()) {
+ opts.appendLiteral(builder, fieldName, value.value_or(true));
+ }
+}
+
+void serializeSpec(const DocumentSourceChangeStreamSpec& spec,
+ const SerializationOptions& opts,
+ BSONObjBuilder* builder) {
+ serializeSpecField(builder,
+ opts,
+ DocumentSourceChangeStreamSpec::kResumeAfterFieldName,
+ spec.getResumeAfter());
+ serializeSpecField(
+ builder, opts, DocumentSourceChangeStreamSpec::kStartAfterFieldName, spec.getStartAfter());
+ serializeSpecField(builder,
+ opts,
+ DocumentSourceChangeStreamSpec::kStartAtOperationTimeFieldName,
+ spec.getStartAtOperationTime());
+ serializeSpecField(builder,
+ opts,
+ DocumentSourceChangeStreamSpec::kFullDocumentFieldName,
+ ::mongo::FullDocumentMode_serializer(spec.getFullDocument()));
+ serializeSpecField(
+ builder,
+ opts,
+ DocumentSourceChangeStreamSpec::kFullDocumentBeforeChangeFieldName,
+ ::mongo::FullDocumentBeforeChangeMode_serializer(spec.getFullDocumentBeforeChange()));
+ serializeSpecField(builder,
+ opts,
+ DocumentSourceChangeStreamSpec::kAllChangesForClusterFieldName,
+ spec.getAllChangesForCluster());
+ serializeSpecField(builder,
+ opts,
+ DocumentSourceChangeStreamSpec::kShowMigrationEventsFieldName,
+ spec.getShowMigrationEvents());
+ serializeSpecField(builder,
+ opts,
+ DocumentSourceChangeStreamSpec::kShowSystemEventsFieldName,
+ spec.getShowSystemEvents());
+ serializeSpecField(builder,
+ opts,
+ DocumentSourceChangeStreamSpec::kAllowToRunOnConfigDBFieldName,
+ spec.getAllowToRunOnConfigDB());
+ serializeSpecField(builder,
+ opts,
+ DocumentSourceChangeStreamSpec::kAllowToRunOnSystemNSFieldName,
+ spec.getAllowToRunOnSystemNS());
+ serializeSpecField(builder,
+ opts,
+ DocumentSourceChangeStreamSpec::kShowExpandedEventsFieldName,
+ spec.getShowExpandedEvents());
+ serializeSpecField(builder,
+ opts,
+ DocumentSourceChangeStreamSpec::kShowRawUpdateDescriptionFieldName,
+ spec.getShowRawUpdateDescription());
+}
+
+} // namespace
+
+Value DocumentSourceChangeStreamTransform::serialize(const SerializationOptions& opts) const {
+ if (opts.verbosity) {
return Value(Document{{DocumentSourceChangeStream::kStageName,
Document{{"stage"_sd, "internalTransform"_sd},
- {"options"_sd, _changeStreamSpec.toBSON()}}}});
+ {"options"_sd, _changeStreamSpec.toBSON(opts)}}}});
}
- return Value(
- Document{{DocumentSourceChangeStreamTransform::kStageName, _changeStreamSpec.toBSON()}});
+ // Internal change stream stages are not serialized for query stats. Query stats uses this stage
+ // to serialize the user specified stage, and therefore if serializing for query stats, we
+ // should use the '$changeStream' stage name.
+ auto stageName =
+ (opts.literalPolicy != LiteralSerializationPolicy::kUnchanged || opts.transformIdentifiers)
+ ? DocumentSourceChangeStream::kStageName
+ : DocumentSourceChangeStreamTransform::kStageName;
+ return Value(Document{{stageName, _changeStreamSpec.toBSON(opts)}});
}
DepsTracker::State DocumentSourceChangeStreamTransform::getDependencies(DepsTracker* deps) const {