diff options
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.cpp | 114 |
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 { |
