summaryrefslogtreecommitdiff
path: root/src/mongo/db/pipeline/pipeline.cpp
diff options
context:
space:
mode:
Diffstat (limited to 'src/mongo/db/pipeline/pipeline.cpp')
-rw-r--r--src/mongo/db/pipeline/pipeline.cpp83
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,