summaryrefslogtreecommitdiff
path: root/src/mongo/s/commands/strategy.cpp
diff options
context:
space:
mode:
authorLucas de Castro Borges <lucas@gnuabordo.com.br>2025-02-11 15:07:35 -0300
committerLucas de Castro Borges <lucas@gnuabordo.com.br>2025-02-11 15:07:35 -0300
commit4cb8841196d0625dfa3825aa326f071cd27c7b8b (patch)
tree1682a647d4463397c119183369ae6f750d5fdcff /src/mongo/s/commands/strategy.cpp
parentaa03c6362cbaa767638e6eed9b031d86dd2643d1 (diff)
parent8f0827553e09872941945a093b647a4211a9db7f (diff)
Update upstream source from tag 'upstream/6.0.0'master
Update to upstream version '6.0.0' with Debian dir 5604a80ec1c96ca76f25f40d78e6ef855abec322
Diffstat (limited to 'src/mongo/s/commands/strategy.cpp')
-rw-r--r--src/mongo/s/commands/strategy.cpp102
1 files changed, 36 insertions, 66 deletions
diff --git a/src/mongo/s/commands/strategy.cpp b/src/mongo/s/commands/strategy.cpp
index e12b3385a13..13d061e9e6d 100644
--- a/src/mongo/s/commands/strategy.cpp
+++ b/src/mongo/s/commands/strategy.cpp
@@ -426,6 +426,14 @@ public:
explicit RunInvocation(ParseAndRunCommand* parc) : _parc(parc) {}
+ ~RunInvocation() {
+ if (!_shouldAffectCommandCounter)
+ return;
+ auto opCtx = _parc->_rec->getOpCtx();
+ Grid::get(opCtx)->catalogCache()->checkAndRecordOperationBlockedByRefresh(
+ opCtx, mongo::LogicalOp::opCommand);
+ }
+
Future<void> run();
private:
@@ -434,6 +442,7 @@ private:
ParseAndRunCommand* const _parc;
boost::optional<RouterOperationContextSession> _routerSession;
+ bool _shouldAffectCommandCounter = false;
};
/*
@@ -698,53 +707,31 @@ Status ParseAndRunCommand::RunInvocation::_setup() {
(opCtx->getClient()->session() &&
(opCtx->getClient()->session()->getTags() & transport::Session::kInternalClient));
- bool canApplyDefaultWC = supportsWriteConcern &&
+ if (supportsWriteConcern && !clientSuppliedWriteConcern &&
(!TransactionRouter::get(opCtx) || isTransactionCommand(_parc->_commandName)) &&
- !opCtx->getClient()->isInDirectClient();
-
- if (canApplyDefaultWC) {
- auto getDefaultWC = ([&]() {
- auto rwcDefaults =
+ !opCtx->getClient()->isInDirectClient()) {
+ if (isInternalClient) {
+ uassert(
+ 5569900,
+ "received command without explicit writeConcern on an internalClient connection {}"_format(
+ redact(request.body.toString())),
+ request.body.hasField(WriteConcernOptions::kWriteConcernField));
+ } else {
+ // This command is not from a DBDirectClient or internal client, and supports WC, but
+ // wasn't given one - so apply the default, if there is one.
+ const auto rwcDefaults =
ReadWriteConcernDefaults::get(opCtx->getServiceContext()).getDefault(opCtx);
- auto wcDefault = rwcDefaults.getDefaultWriteConcern();
- const auto defaultWriteConcernSource = rwcDefaults.getDefaultWriteConcernSource();
- customDefaultWriteConcernWasApplied = defaultWriteConcernSource &&
- defaultWriteConcernSource == DefaultWriteConcernSourceEnum::kGlobal;
- return wcDefault;
- });
-
- if (!clientSuppliedWriteConcern) {
- if (isInternalClient) {
- uassert(
- 5569900,
- "received command without explicit writeConcern on an internalClient connection {}"_format(
- redact(request.body.toString())),
- request.body.hasField(WriteConcernOptions::kWriteConcernField));
- } else {
- // This command is not from a DBDirectClient or internal client, and supports WC,
- // but wasn't given one - so apply the default, if there is one.
- const auto wcDefault = getDefaultWC();
- // Default WC can be 'boost::none' if the implicit default is used and set to 'w:1'.
- if (wcDefault) {
- _parc->_wc = *wcDefault;
- LOGV2_DEBUG(22766,
- 2,
- "Applying default writeConcern on command",
- "command"_attr = request.getCommandName(),
- "writeConcern"_attr = *wcDefault);
- }
- }
- }
- // Client supplied a write concern object without 'w' field.
- else if (_parc->_wc->isExplicitWithoutWField()) {
- const auto wcDefault = getDefaultWC();
- // Default WC can be 'boost::none' if the implicit default is used and set to 'w:1'.
- if (wcDefault) {
- clientSuppliedWriteConcern = false;
- _parc->_wc->w = wcDefault->w;
- if (_parc->_wc->syncMode == WriteConcernOptions::SyncMode::UNSET) {
- _parc->_wc->syncMode = wcDefault->syncMode;
- }
+ if (const auto wcDefault = rwcDefaults.getDefaultWriteConcern()) {
+ _parc->_wc = *wcDefault;
+ const auto defaultWriteConcernSource = rwcDefaults.getDefaultWriteConcernSource();
+ customDefaultWriteConcernWasApplied = defaultWriteConcernSource &&
+ defaultWriteConcernSource == DefaultWriteConcernSourceEnum::kGlobal;
+ LOGV2_DEBUG(22766,
+ 2,
+ "Applying default writeConcern on {command} of {writeConcern}",
+ "Applying default writeConcern on command",
+ "command"_attr = request.getCommandName(),
+ "writeConcern"_attr = *wcDefault);
}
}
}
@@ -911,6 +898,7 @@ Status ParseAndRunCommand::RunInvocation::_setup() {
if (command->shouldAffectCommandCounter()) {
globalOpCounters.gotCommand();
+ _shouldAffectCommandCounter = true;
}
return Status::OK();
@@ -1048,21 +1036,10 @@ void ParseAndRunCommand::RunAndRetry::_onNeedRetargetting(Status& status) {
auto opCtx = _parc->_rec->getOpCtx();
const auto staleNs = staleInfo->getNss();
- const auto& originalNs = _parc->_invocation->ns();
auto catalogCache = Grid::get(opCtx)->catalogCache();
catalogCache->invalidateShardOrEntireCollectionEntryForShardedCollection(
staleNs, staleInfo->getVersionWanted(), staleInfo->getShardId());
- if ((staleNs.isTimeseriesBucketsCollection() || originalNs.isTimeseriesBucketsCollection()) &&
- staleNs != originalNs) {
- // A timeseries might've been created, so we need to invalidate the original namespace
- // version.
- Grid::get(opCtx)
- ->catalogCache()
- ->invalidateShardOrEntireCollectionEntryForShardedCollection(
- originalNs, boost::none, staleInfo->getShardId());
- }
-
catalogCache->setOperationShouldBlockBehindCatalogCacheRefresh(opCtx, true);
_checkRetryForTransaction(status);
@@ -1193,13 +1170,6 @@ public:
Future<DbResponse> run();
private:
- std::string _getDatabaseStringForLogging() const try {
- // `getDatabase` throws if the request doesn't have a '$db' field.
- return _rec->getRequest().getDatabase().toString();
- } catch (const DBException& ex) {
- return ex.toString();
- }
-
void _parseMessage();
Future<void> _execute();
@@ -1242,7 +1212,7 @@ Future<void> ClientCommand::_execute() {
3,
"Command begin db: {db} msg id: {headerId}",
"Command begin",
- "db"_attr = _getDatabaseStringForLogging(),
+ "db"_attr = _rec->getRequest().getDatabase().toString(),
"headerId"_attr = _rec->getMessage().header().getId());
return future_util::makeState<ParseAndRunCommand>(_rec, _errorBuilder)
@@ -1252,7 +1222,7 @@ Future<void> ClientCommand::_execute() {
3,
"Command end db: {db} msg id: {headerId}",
"Command end",
- "db"_attr = _getDatabaseStringForLogging(),
+ "db"_attr = _rec->getRequest().getDatabase().toString(),
"headerId"_attr = _rec->getMessage().header().getId());
})
.tapError([this](Status status) {
@@ -1261,7 +1231,7 @@ Future<void> ClientCommand::_execute() {
1,
"Exception thrown while processing command on {db} msg id: {headerId} {error}",
"Exception thrown while processing command",
- "db"_attr = _getDatabaseStringForLogging(),
+ "db"_attr = _rec->getRequest().getDatabase().toString(),
"headerId"_attr = _rec->getMessage().header().getId(),
"error"_attr = redact(status));