summaryrefslogtreecommitdiff
path: root/src/mongo/db/pipeline/dispatch_shard_pipeline_test.cpp
diff options
context:
space:
mode:
Diffstat (limited to 'src/mongo/db/pipeline/dispatch_shard_pipeline_test.cpp')
-rw-r--r--src/mongo/db/pipeline/dispatch_shard_pipeline_test.cpp42
1 files changed, 32 insertions, 10 deletions
diff --git a/src/mongo/db/pipeline/dispatch_shard_pipeline_test.cpp b/src/mongo/db/pipeline/dispatch_shard_pipeline_test.cpp
index 069a7e2f0b2..ac8924a13ed 100644
--- a/src/mongo/db/pipeline/dispatch_shard_pipeline_test.cpp
+++ b/src/mongo/db/pipeline/dispatch_shard_pipeline_test.cpp
@@ -30,7 +30,7 @@
#include "mongo/db/pipeline/aggregation_request_helper.h"
#include "mongo/db/pipeline/sharded_agg_helpers.h"
#include "mongo/s/query/sharded_agg_test_fixture.h"
-#include "mongo/s/router.h"
+#include "mongo/s/router_role.h"
namespace mongo {
namespace {
@@ -53,10 +53,14 @@ TEST_F(DispatchShardPipelineTest, DoesNotSplitPipelineIfTargetingOneShard) {
const Document serializedCommand = aggregation_request_helper::serializeToCommandDoc(
AggregateCommandRequest(expCtx()->ns, stages));
const bool hasChangeStream = false;
+ const bool startsWithDocuments = false;
auto future = launchAsync([&] {
- auto results = sharded_agg_helpers::dispatchShardPipeline(
- serializedCommand, hasChangeStream, std::move(pipeline));
+ auto results = sharded_agg_helpers::dispatchShardPipeline(serializedCommand,
+ hasChangeStream,
+ startsWithDocuments,
+ std::move(pipeline),
+ boost::none /*explain*/);
ASSERT_EQ(results.remoteCursors.size(), 1UL);
ASSERT(!results.splitPipeline);
});
@@ -84,10 +88,14 @@ TEST_F(DispatchShardPipelineTest, DoesSplitPipelineIfMatchSpansTwoShards) {
const Document serializedCommand = aggregation_request_helper::serializeToCommandDoc(
AggregateCommandRequest(expCtx()->ns, stages));
const bool hasChangeStream = false;
+ const bool startsWithDocuments = false;
auto future = launchAsync([&] {
- auto results = sharded_agg_helpers::dispatchShardPipeline(
- serializedCommand, hasChangeStream, std::move(pipeline));
+ auto results = sharded_agg_helpers::dispatchShardPipeline(serializedCommand,
+ hasChangeStream,
+ startsWithDocuments,
+ std::move(pipeline),
+ boost::none /*explain*/);
ASSERT_EQ(results.remoteCursors.size(), 2UL);
ASSERT(bool(results.splitPipeline));
});
@@ -118,10 +126,14 @@ TEST_F(DispatchShardPipelineTest, DispatchShardPipelineRetriesOnNetworkError) {
const Document serializedCommand = aggregation_request_helper::serializeToCommandDoc(
AggregateCommandRequest(expCtx()->ns, stages));
const bool hasChangeStream = false;
+ const bool startsWithDocuments = false;
auto future = launchAsync([&] {
// Shouldn't throw.
- auto results = sharded_agg_helpers::dispatchShardPipeline(
- serializedCommand, hasChangeStream, std::move(pipeline));
+ auto results = sharded_agg_helpers::dispatchShardPipeline(serializedCommand,
+ hasChangeStream,
+ startsWithDocuments,
+ std::move(pipeline),
+ boost::none /*explain*/);
ASSERT_EQ(results.remoteCursors.size(), 2UL);
ASSERT(bool(results.splitPipeline));
});
@@ -163,9 +175,14 @@ TEST_F(DispatchShardPipelineTest, DispatchShardPipelineDoesNotRetryOnStaleConfig
const Document serializedCommand = aggregation_request_helper::serializeToCommandDoc(
AggregateCommandRequest(expCtx()->ns, stages));
const bool hasChangeStream = false;
+ const bool startsWithDocuments = false;
+
auto future = launchAsync([&] {
- ASSERT_THROWS_CODE(sharded_agg_helpers::dispatchShardPipeline(
- serializedCommand, hasChangeStream, std::move(pipeline)),
+ ASSERT_THROWS_CODE(sharded_agg_helpers::dispatchShardPipeline(serializedCommand,
+ hasChangeStream,
+ startsWithDocuments,
+ std::move(pipeline),
+ boost::none /*explain*/),
AssertionException,
ErrorCodes::StaleConfig);
});
@@ -197,6 +214,7 @@ TEST_F(DispatchShardPipelineTest, WrappedDispatchDoesRetryOnStaleConfigError) {
const Document serializedCommand = aggregation_request_helper::serializeToCommandDoc(
AggregateCommandRequest(expCtx()->ns, stages));
const bool hasChangeStream = false;
+ const bool startsWithDocuments = false;
auto future = launchAsync([&] {
// Shouldn't throw.
sharding::router::CollectionRouter router(getServiceContext(), kTestAggregateNss);
@@ -204,7 +222,11 @@ TEST_F(DispatchShardPipelineTest, WrappedDispatchDoesRetryOnStaleConfigError) {
"dispatch shard pipeline"_sd,
[&](OperationContext* opCtx, const ChunkManager& cm) {
return sharded_agg_helpers::dispatchShardPipeline(
- serializedCommand, hasChangeStream, pipeline->clone());
+ serializedCommand,
+ hasChangeStream,
+ startsWithDocuments,
+ pipeline->clone(),
+ boost::none /*explain*/);
});
ASSERT_EQ(results.remoteCursors.size(), 1UL);
ASSERT(!bool(results.splitPipeline));