diff options
Diffstat (limited to 'src/mongo/db/pipeline/pipeline.cpp')
| -rw-r--r-- | src/mongo/db/pipeline/pipeline.cpp | 83 |
1 files changed, 64 insertions, 19 deletions
diff --git a/src/mongo/db/pipeline/pipeline.cpp b/src/mongo/db/pipeline/pipeline.cpp index 97a896e5898..22ca4e7a85f 100644 --- a/src/mongo/db/pipeline/pipeline.cpp +++ b/src/mongo/db/pipeline/pipeline.cpp @@ -40,6 +40,7 @@ #include "mongo/db/jsobj.h" #include "mongo/db/operation_context.h" #include "mongo/db/pipeline/accumulator.h" +#include "mongo/db/pipeline/change_stream_helpers.h" #include "mongo/db/pipeline/document_source.h" #include "mongo/db/pipeline/document_source_match.h" #include "mongo/db/pipeline/document_source_merge.h" @@ -105,14 +106,32 @@ void validateTopLevelPipeline(const Pipeline& pipeline) { // If the first stage is a $changeStream stage, then all stages in the pipeline must be // either $changeStream stages or allowlisted as being able to run in a change stream. - if (firstStageConstraints.isChangeStreamStage()) { - for (auto&& source : sources) { - uassert(ErrorCodes::IllegalOperation, - str::stream() << source->getSourceName() - << " is not permitted in a $changeStream pipeline", - source->constraints().isAllowedInChangeStream()); + const bool isChangeStream = firstStageConstraints.isChangeStreamStage(); + // Record whether any of the stages in the pipeline is a $changeStreamSplitLargeEvent. + bool hasChangeStreamSplitLargeEventStage = false; + for (auto&& source : sources) { + uassert(ErrorCodes::IllegalOperation, + str::stream() << source->getSourceName() + << " is not permitted in a $changeStream pipeline", + !(isChangeStream && !source->constraints().isAllowedInChangeStream())); + // Check whether any stages must only be run in a change stream pipeline. + uassert(ErrorCodes::IllegalOperation, + str::stream() << source->getSourceName() + << " can only be used in a $changeStream pipeline", + !(source->constraints().requiresChangeStream() && !isChangeStream)); + // Check whether this is a change stream split stage. + if ("$changeStreamSplitLargeEvent"_sd == source->getSourceName()) { + hasChangeStreamSplitLargeEventStage = true; } } + auto expCtx = pipeline.getContext(); + auto spec = isChangeStream ? expCtx->changeStreamSpec : boost::none; + auto hasSplitEventResumeToken = spec && + change_stream::resolveResumeTokenFromSpec(expCtx, *spec).fragmentNum.has_value(); + uassert(ErrorCodes::ChangeStreamFatalError, + "To resume from a split event, the $changeStream pipeline must include a " + "$changeStreamSplitLargeEvent stage", + !(hasSplitEventResumeToken && !hasChangeStreamSplitLargeEventStage)); } // Verify that usage of $searchMeta and $search is legal. @@ -231,39 +250,41 @@ std::unique_ptr<Pipeline, PipelineDeleter> Pipeline::create( } void Pipeline::validateCommon(bool alreadyOptimized) const { - size_t i = 0; - uassert(ErrorCodes::FailedToParse, str::stream() << "Pipeline length must be no longer than " << internalPipelineLengthLimit << " stages", static_cast<int>(_sources.size()) <= internalPipelineLengthLimit); - for (auto&& stage : _sources) { + // Keep track of stages which can only appear once. + std::set<StringData> singleUseStages; + + for (auto sourceIter = _sources.begin(); sourceIter != _sources.end(); ++sourceIter) { + auto& stage = *sourceIter; auto constraints = stage->constraints(_splitState); // Verify that all stages adhere to their PositionRequirement constraints. uassert(40602, str::stream() << stage->getSourceName() << " is only valid as the first stage in a pipeline", - !(constraints.requiredPosition == PositionRequirement::kFirst && i != 0)); - uassert(40603, - str::stream() << stage->getSourceName() - << " is only valid as the first stage in an optimized pipeline", - !(alreadyOptimized && - constraints.requiredPosition == PositionRequirement::kFirstAfterOptimization && - i != 0)); + !(constraints.requiredPosition == PositionRequirement::kFirst && + sourceIter != _sources.begin())); + // TODO SERVER-73790: use PositionRequirement::kCustom to validate $match. auto matchStage = dynamic_cast<DocumentSourceMatch*>(stage.get()); uassert(17313, "$match with $text is only allowed as the first pipeline stage", - !(i != 0 && matchStage && matchStage->isTextQuery())); + !(sourceIter != _sources.begin() && matchStage && matchStage->isTextQuery())); uassert(40601, str::stream() << stage->getSourceName() << " can only be the final stage in the pipeline", !(constraints.requiredPosition == PositionRequirement::kLast && - i != _sources.size() - 1)); - ++i; + std::next(sourceIter) != _sources.end())); + + // If the stage has a special requirement about its position, validate it. + if (constraints.requiredPosition == PositionRequirement::kCustom) { + stage->validatePipelinePosition(alreadyOptimized, sourceIter, _sources); + } // Verify that we are not attempting to run a mongoS-only stage on mongoD. uassert(40644, @@ -275,6 +296,12 @@ void Pipeline::validateCommon(bool alreadyOptimized) const { str::stream() << "Stage not supported inside of a multi-document transaction: " << stage->getSourceName(), !(pCtx->opCtx->inMultiDocumentTransaction() && !constraints.isAllowedInTransaction())); + + // Verify that a stage which can only appear once doesn't appear more than that. + uassert(7183900, + str::stream() << stage->getSourceName() << " can only be used once in the pipeline", + !(constraints.canAppearOnlyOnceInPipeline && + !singleUseStages.insert(stage->getSourceName()).second)); } } @@ -721,6 +748,24 @@ Pipeline::SourceContainer::iterator Pipeline::optimizeEndOfPipeline( return std::next(itr); } +Pipeline::SourceContainer::iterator Pipeline::optimizeAtEndOfPipeline( + Pipeline::SourceContainer::iterator itr, Pipeline::SourceContainer* container) { + if (itr == container->end()) { + return itr; + } + itr = std::next(itr); + try { + while (itr != container->end()) { + invariant((*itr).get()); + itr = (*itr).get()->optimizeAt(itr, container); + } + } catch (DBException& ex) { + ex.addContext("Failed to optimize pipeline"); + throw; + } + return itr; +} + std::unique_ptr<Pipeline, PipelineDeleter> Pipeline::makePipelineFromViewDefinition( const boost::intrusive_ptr<ExpressionContext>& subPipelineExpCtx, ExpressionContext::ResolvedNamespace resolvedNs, |
