diff --git a/config/spotbugs/exclude.xml b/config/spotbugs/exclude.xml index e2f1925fd7..151845d59b 100644 --- a/config/spotbugs/exclude.xml +++ b/config/spotbugs/exclude.xml @@ -288,6 +288,16 @@ + + + + + + + + + + diff --git a/driver-core/src/main/com/mongodb/internal/operation/AbstractWriteSearchIndexOperation.java b/driver-core/src/main/com/mongodb/internal/operation/AbstractWriteSearchIndexOperation.java index a409da75f6..375280c90c 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/AbstractWriteSearchIndexOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/AbstractWriteSearchIndexOperation.java @@ -20,15 +20,23 @@ import com.mongodb.MongoCommandException; import com.mongodb.MongoNamespace; import com.mongodb.internal.async.SingleResultCallback; +import com.mongodb.internal.async.function.AsyncCallbackSupplier; +import com.mongodb.internal.async.function.RetryControl; import com.mongodb.internal.binding.AsyncWriteBinding; import com.mongodb.internal.binding.WriteBinding; import com.mongodb.internal.connection.OperationContext; import com.mongodb.lang.Nullable; import org.bson.BsonDocument; +import java.util.function.Supplier; + +import static com.mongodb.internal.operation.AsyncOperationHelper.decorateWithRetriesAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.executeCommandAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.withAsyncSourceAndConnection; import static com.mongodb.internal.operation.AsyncOperationHelper.writeConcernErrorTransformerAsync; +import static com.mongodb.internal.operation.CommandOperationHelper.createSpecRetryControl; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForWrite; +import static com.mongodb.internal.operation.SyncOperationHelper.decorateWithRetries; import static com.mongodb.internal.operation.SyncOperationHelper.executeCommand; import static com.mongodb.internal.operation.SyncOperationHelper.withConnection; import static com.mongodb.internal.operation.SyncOperationHelper.writeConcernErrorTransformer; @@ -40,39 +48,64 @@ */ abstract class AbstractWriteSearchIndexOperation implements WriteOperation { private final MongoNamespace namespace; + private final boolean retryWrites; + @Nullable + private final Integer maxAdaptiveRetriesSetting; AbstractWriteSearchIndexOperation(final MongoNamespace namespace) { + this(namespace, false, null); + } + + AbstractWriteSearchIndexOperation(final MongoNamespace namespace, final boolean retryWrites, + @Nullable final Integer maxAdaptiveRetriesSetting) { this.namespace = namespace; + this.retryWrites = retryWrites; + this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting; } @Override public Void execute(final WriteBinding binding, final OperationContext operationContext) { - return withConnection(binding, operationContext, (connection, operationContextWithMinRtt) -> { - try { - executeCommand(binding, operationContextWithMinRtt, namespace.getDatabaseName(), buildCommand(), - connection, - writeConcernErrorTransformer(operationContextWithMinRtt.getTimeoutContext())); - } catch (MongoCommandException mongoCommandException) { - swallowOrThrow(mongoCommandException); - } - return null; + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + Supplier retryingCommandExecutor = decorateWithRetries(retryControl, operationContext, () -> { + retryControl.getPolicy().onCommand(this::getCommandName); + return withConnection(binding, operationContext, (connection, operationContextWithMinRtt) -> { + try { + executeCommand(binding, operationContextWithMinRtt, namespace.getDatabaseName(), buildCommand(), + connection, + writeConcernErrorTransformer(operationContextWithMinRtt.getTimeoutContext())); + } catch (MongoCommandException mongoCommandException) { + swallowOrThrow(mongoCommandException); + } + return null; + }); }); + return retryingCommandExecutor.get(); } @Override public void executeAsync(final AsyncWriteBinding binding, final OperationContext operationContext, final SingleResultCallback callback) { - withAsyncSourceAndConnection(binding::getWriteConnectionSource, false, operationContext, callback, - (connectionSource, connection, operationContextWithMinRtt, cb) -> - executeCommandAsync(binding, operationContextWithMinRtt, namespace.getDatabaseName(), buildCommand(), connection, - writeConcernErrorTransformerAsync(operationContextWithMinRtt.getTimeoutContext()), (result, commandExecutionError) -> { - try { - swallowOrThrow(commandExecutionError); - cb.onResult(result, null); - } catch (Throwable mongoCommandException) { - cb.onResult(null, mongoCommandException); + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + AsyncCallbackSupplier retryingCommandExecutor = decorateWithRetriesAsync(retryControl, operationContext, supplierCallback -> { + retryControl.getPolicy().onCommand(this::getCommandName); + withAsyncSourceAndConnection(binding::getWriteConnectionSource, false, operationContext, supplierCallback, + (connectionSource, connection, operationContextWithMinRtt, cb) -> + executeCommandAsync(binding, operationContextWithMinRtt, namespace.getDatabaseName(), buildCommand(), + connection, writeConcernErrorTransformerAsync(operationContextWithMinRtt.getTimeoutContext()), + (result, commandExecutionError) -> { + try { + swallowOrThrow(commandExecutionError); + cb.onResult(result, null); + } catch (Throwable mongoCommandException) { + cb.onResult(null, mongoCommandException); + } } - } - )); + )); + }); + retryingCommandExecutor.get(callback); } /** @@ -101,4 +134,5 @@ void swallowOrThrow(@Nullable final E mongoExecutionExcept public MongoNamespace getNamespace() { return namespace; } + } diff --git a/driver-core/src/main/com/mongodb/internal/operation/AggregateToCollectionOperation.java b/driver-core/src/main/com/mongodb/internal/operation/AggregateToCollectionOperation.java index 69332c409b..763c9d1a78 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/AggregateToCollectionOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/AggregateToCollectionOperation.java @@ -43,6 +43,7 @@ import static com.mongodb.internal.operation.AsyncOperationHelper.CommandReadTransformerAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.executeRetryableReadAsync; import static com.mongodb.internal.operation.ServerVersionHelper.FIVE_DOT_ZERO_WIRE_VERSION; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForWrite; import static com.mongodb.internal.operation.SyncOperationHelper.CommandReadTransformer; import static com.mongodb.internal.operation.SyncOperationHelper.executeRetryableRead; import static com.mongodb.internal.operation.WriteConcernHelper.appendWriteConcernToCommand; @@ -64,6 +65,9 @@ public class AggregateToCollectionOperation implements ReadOperationSimple private final WriteConcern writeConcern; private final ReadConcern readConcern; private final AggregationLevel aggregationLevel; + private final boolean retryWrites; + @Nullable + private final Integer maxAdaptiveRetriesSetting; private Boolean allowDiskUse; private Boolean bypassDocumentValidation; @@ -79,11 +83,19 @@ public AggregateToCollectionOperation(final MongoNamespace namespace, final List public AggregateToCollectionOperation(final MongoNamespace namespace, final List pipeline, @Nullable final ReadConcern readConcern, @Nullable final WriteConcern writeConcern, final AggregationLevel aggregationLevel) { + this(namespace, pipeline, readConcern, writeConcern, aggregationLevel, false, null); + } + + public AggregateToCollectionOperation(final MongoNamespace namespace, final List pipeline, + @Nullable final ReadConcern readConcern, @Nullable final WriteConcern writeConcern, final AggregationLevel aggregationLevel, + final boolean retryWrites, @Nullable final Integer maxAdaptiveRetriesSetting) { this.namespace = notNull("namespace", namespace); this.pipeline = notNull("pipeline", pipeline); this.writeConcern = writeConcern; this.readConcern = readConcern; this.aggregationLevel = notNull("aggregationLevel", aggregationLevel); + this.retryWrites = retryWrites; + this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting; isTrueArgument("pipeline is not empty", !pipeline.isEmpty()); } @@ -178,8 +190,7 @@ public Void execute(final ReadBinding binding, final OperationContext operationC getCommandCreator(), new BsonDocumentCodec(), transformer(), - false, - null); + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting)); } @Override @@ -194,8 +205,7 @@ public void executeAsync(final AsyncReadBinding binding, final OperationContext getCommandCreator(), new BsonDocumentCodec(), asyncTransformer(), - false, - null, + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), callback); } @@ -251,4 +261,5 @@ private static CommandReadTransformerAsync asyncTransformer( return null; }; } + } diff --git a/driver-core/src/main/com/mongodb/internal/operation/AsyncCommandCursor.java b/driver-core/src/main/com/mongodb/internal/operation/AsyncCommandCursor.java index 0e33e79351..6a4ef41a43 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/AsyncCommandCursor.java +++ b/driver-core/src/main/com/mongodb/internal/operation/AsyncCommandCursor.java @@ -54,6 +54,7 @@ import static com.mongodb.internal.async.SingleResultCallback.THEN_DO_NOTHING; import static com.mongodb.internal.operation.AsyncOperationHelper.decorateWithRetriesAsync; import static com.mongodb.internal.operation.CommandOperationHelper.createSpecRetryControl; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForRead; import static com.mongodb.internal.operation.CommandBatchCursorHelper.FIRST_BATCH; import static com.mongodb.internal.operation.CommandBatchCursorHelper.MESSAGE_IF_CLOSED_AS_CURSOR; import static com.mongodb.internal.operation.CommandBatchCursorHelper.NEXT_BATCH; @@ -191,9 +192,8 @@ private void getMoreLoop(final ServerCursor localServerCursor, } private void getMore(final ServerCursor cursor, final OperationContext operationContext, final SingleResultCallback> callback) { - SpecRetryPolicy.IndividualPolicies policies = new SpecRetryPolicy.IndividualPolicies(retryReads) - .includeOverload(maxAdaptiveRetriesSetting, SpecRetryPolicy.ErrorPropagation.AS_READ_POLICY); - RetryControl retryControl = createSpecRetryControl(policies, operationContext); + RetryControl retryControl = createSpecRetryControl( + overloadForRead(retryReads, maxAdaptiveRetriesSetting), operationContext); AsyncCallbackSupplier> retryingCommandExecutor = decorateWithRetriesAsync(retryControl, operationContext, attemptCallback -> resourceManager.executeWithConnection(operationContext, (connection, wrappedCallback) -> executeGetMoreCommand(assertNotNull(connection), cursor, operationContext, retryControl, wrappedCallback), diff --git a/driver-core/src/main/com/mongodb/internal/operation/CommandCursor.java b/driver-core/src/main/com/mongodb/internal/operation/CommandCursor.java index 2f8a98d619..75f4f3bbf8 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/CommandCursor.java +++ b/driver-core/src/main/com/mongodb/internal/operation/CommandCursor.java @@ -49,6 +49,7 @@ import static com.mongodb.assertions.Assertions.assertTrue; import static com.mongodb.internal.VisibleForTesting.AccessModifier.PRIVATE; import static com.mongodb.internal.operation.CommandOperationHelper.createSpecRetryControl; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForRead; import static com.mongodb.internal.operation.SyncOperationHelper.decorateWithRetries; import static com.mongodb.internal.operation.CommandBatchCursorHelper.FIRST_BATCH; import static com.mongodb.internal.operation.CommandBatchCursorHelper.MESSAGE_IF_CLOSED_AS_CURSOR; @@ -231,9 +232,8 @@ public int getMaxWireVersion() { private void getMore(final OperationContext operationContext) { ServerCursor serverCursor = assertNotNull(resourceManager.getServerCursor()); - SpecRetryPolicy.IndividualPolicies policies = new SpecRetryPolicy.IndividualPolicies(retryReads) - .includeOverload(maxAdaptiveRetriesSetting, SpecRetryPolicy.ErrorPropagation.AS_READ_POLICY); - RetryControl retryControl = createSpecRetryControl(policies, operationContext); + RetryControl retryControl = createSpecRetryControl( + overloadForRead(retryReads, maxAdaptiveRetriesSetting), operationContext); Supplier retryingCommandExecutor = decorateWithRetries(retryControl, operationContext, () -> { resourceManager.executeWithConnection(connection -> { ServerCursor nextServerCursor; diff --git a/driver-core/src/main/com/mongodb/internal/operation/CommandReadOperation.java b/driver-core/src/main/com/mongodb/internal/operation/CommandReadOperation.java index 8e623d5dc7..d4f4c56708 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/CommandReadOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/CommandReadOperation.java @@ -25,6 +25,7 @@ import org.bson.codecs.Decoder; import static com.mongodb.internal.operation.AsyncOperationHelper.executeRetryableReadAsync; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForWrite; import static com.mongodb.internal.operation.SyncOperationHelper.executeRetryableRead; /** @@ -76,8 +77,6 @@ public void executeAsync(final AsyncReadBinding binding, final OperationContext } private SpecRetryPolicy.IndividualPolicies createRetryPolicy() { - boolean retryPolicyEnabled = retryReads && retryWrites; - return new SpecRetryPolicy.IndividualPolicies(retryPolicyEnabled) - .includeOverload(maxAdaptiveRetriesSetting, SpecRetryPolicy.ErrorPropagation.AS_WRITE_POLICY); + return overloadForWrite(retryReads && retryWrites, maxAdaptiveRetriesSetting); } } diff --git a/driver-core/src/main/com/mongodb/internal/operation/CreateCollectionOperation.java b/driver-core/src/main/com/mongodb/internal/operation/CreateCollectionOperation.java index 9f82effafc..a2758733d8 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/CreateCollectionOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/CreateCollectionOperation.java @@ -28,9 +28,10 @@ import com.mongodb.client.model.ValidationLevel; import com.mongodb.connection.ConnectionDescription; import com.mongodb.internal.async.SingleResultCallback; +import com.mongodb.internal.async.function.AsyncCallbackSupplier; +import com.mongodb.internal.async.function.RetryControl; import com.mongodb.internal.binding.AsyncWriteBinding; import com.mongodb.internal.binding.WriteBinding; -import com.mongodb.internal.connection.AsyncConnection; import com.mongodb.internal.connection.OperationContext; import com.mongodb.lang.Nullable; import org.bson.BsonArray; @@ -46,16 +47,18 @@ import java.util.function.Supplier; import static com.mongodb.assertions.Assertions.notNull; -import static com.mongodb.internal.async.ErrorHandlingResultCallback.errorHandlingCallback; +import static com.mongodb.internal.operation.AsyncOperationHelper.decorateWithRetriesAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.executeCommandAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.releasingCallback; import static com.mongodb.internal.operation.AsyncOperationHelper.withAsyncConnection; import static com.mongodb.internal.operation.AsyncOperationHelper.writeConcernErrorTransformerAsync; +import static com.mongodb.internal.operation.CommandOperationHelper.createSpecRetryControl; import static com.mongodb.internal.operation.DocumentHelper.putIfFalse; import static com.mongodb.internal.operation.DocumentHelper.putIfNotNull; import static com.mongodb.internal.operation.DocumentHelper.putIfNotZero; -import static com.mongodb.internal.operation.OperationHelper.LOGGER; import static com.mongodb.internal.operation.ServerVersionHelper.serverIsLessThanVersionSevenDotZero; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForWrite; +import static com.mongodb.internal.operation.SyncOperationHelper.decorateWithRetries; import static com.mongodb.internal.operation.SyncOperationHelper.executeCommand; import static com.mongodb.internal.operation.SyncOperationHelper.withConnection; import static com.mongodb.internal.operation.SyncOperationHelper.writeConcernErrorTransformer; @@ -76,6 +79,9 @@ public class CreateCollectionOperation implements WriteOperation { private final String databaseName; private final String collectionName; private final WriteConcern writeConcern; + private final boolean retryWrites; + @Nullable + private final Integer maxAdaptiveRetriesSetting; private boolean capped = false; private long sizeInBytes = 0; private boolean autoIndex = true; @@ -95,9 +101,16 @@ public class CreateCollectionOperation implements WriteOperation { private BsonDocument encryptedFields; public CreateCollectionOperation(final String databaseName, final String collectionName, @Nullable final WriteConcern writeConcern) { + this(databaseName, collectionName, writeConcern, false, null); + } + + public CreateCollectionOperation(final String databaseName, final String collectionName, @Nullable final WriteConcern writeConcern, + final boolean retryWrites, @Nullable final Integer maxAdaptiveRetriesSetting) { this.databaseName = notNull("databaseName", databaseName); this.collectionName = notNull("collectionName", collectionName); this.writeConcern = writeConcern; + this.retryWrites = retryWrites; + this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting; } public String getCollectionName() { @@ -245,31 +258,26 @@ public MongoNamespace getNamespace() { @Override public Void execute(final WriteBinding binding, final OperationContext operationContext) { - return withConnection(binding, operationContext, (connection, operationContextWithMinRtt)-> { - checkEncryptedFieldsSupported(connection.getDescription()); - getCommandFunctions().forEach(commandCreator -> - executeCommand(binding, operationContextWithMinRtt, databaseName, commandCreator.get(), connection, - writeConcernErrorTransformer(operationContextWithMinRtt.getTimeoutContext())) - ); - return null; + getCommandFunctions().forEach(commandCreator -> { + RetryControl retryControl = createSpecRetryControl(overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + Supplier retryingCommandExecutor = decorateWithRetries(retryControl, operationContext, () -> { + retryControl.getPolicy().onCommand(this::getCommandName); + return withConnection(binding, operationContext, (connection, connectionScopedOperationContext) -> { + checkEncryptedFieldsSupported(connection.getDescription()); + executeCommand(binding, connectionScopedOperationContext, databaseName, commandCreator.get(), connection, + writeConcernErrorTransformer(connectionScopedOperationContext.getTimeoutContext())); + return null; + }); + }); + retryingCommandExecutor.get(); }); + return null; } @Override public void executeAsync(final AsyncWriteBinding binding, final OperationContext operationContext, final SingleResultCallback callback) { - withAsyncConnection(binding, operationContext, (connection, operationContextWithMinRtt, t) -> { - SingleResultCallback errHandlingCallback = errorHandlingCallback(callback, LOGGER); - if (t != null) { - errHandlingCallback.onResult(null, t); - } else { - SingleResultCallback releasingCallback = releasingCallback(errHandlingCallback, connection); - if (!checkEncryptedFieldsSupported(connection.getDescription(), releasingCallback)) { - return; - } - new ProcessCommandsCallback(binding, operationContextWithMinRtt, connection, releasingCallback) - .onResult(null, null); - } - }); + new ProcessCommandsCallback(binding, operationContext, callback).onResult(null, null); } private String getGranularityAsString(final TimeSeriesGranularity granularity) { @@ -411,15 +419,13 @@ private boolean checkEncryptedFieldsSupported(final ConnectionDescription connec class ProcessCommandsCallback implements SingleResultCallback { private final AsyncWriteBinding binding; private final OperationContext operationContext; - private final AsyncConnection connection; private final SingleResultCallback finalCallback; private final Deque> commands; ProcessCommandsCallback( - final AsyncWriteBinding binding, final OperationContext operationContext, final AsyncConnection connection, final SingleResultCallback finalCallback) { + final AsyncWriteBinding binding, final OperationContext operationContext, final SingleResultCallback finalCallback) { this.binding = binding; this.operationContext = operationContext; - this.connection = connection; this.finalCallback = finalCallback; this.commands = new ArrayDeque<>(getCommandFunctions()); } @@ -434,8 +440,27 @@ public void onResult(@Nullable final Void result, @Nullable final Throwable t) { if (nextCommandFunction == null) { finalCallback.onResult(null, null); } else { - executeCommandAsync(binding, operationContext, databaseName, nextCommandFunction.get(), - connection, writeConcernErrorTransformerAsync(operationContext.getTimeoutContext()), this); + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), operationContext); + AsyncCallbackSupplier retryingCommandExecutor = decorateWithRetriesAsync(retryControl, operationContext, + supplierCallback -> { + retryControl.getPolicy().onCommand(CreateCollectionOperation.this::getCommandName); + withAsyncConnection(binding, operationContext, (connection, connectionScopedOperationContext, t1) -> { + if (t1 != null) { + supplierCallback.onResult(null, t1); + } else { + SingleResultCallback connectionReleasingCallback = releasingCallback(supplierCallback, connection); + if (!checkEncryptedFieldsSupported(connection.getDescription(), connectionReleasingCallback)) { + return; + } + executeCommandAsync(binding, connectionScopedOperationContext, databaseName, + nextCommandFunction.get(), connection, + writeConcernErrorTransformerAsync(connectionScopedOperationContext.getTimeoutContext()), + connectionReleasingCallback); + } + }); + }); + retryingCommandExecutor.get(this); } } } diff --git a/driver-core/src/main/com/mongodb/internal/operation/CreateIndexesOperation.java b/driver-core/src/main/com/mongodb/internal/operation/CreateIndexesOperation.java index 8ad1280369..d096be3abd 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/CreateIndexesOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/CreateIndexesOperation.java @@ -26,6 +26,8 @@ import com.mongodb.WriteConcern; import com.mongodb.WriteConcernResult; import com.mongodb.internal.async.SingleResultCallback; +import com.mongodb.internal.async.function.AsyncCallbackSupplier; +import com.mongodb.internal.async.function.RetryControl; import com.mongodb.internal.binding.AsyncWriteBinding; import com.mongodb.internal.binding.WriteBinding; import com.mongodb.internal.bulk.IndexRequest; @@ -42,13 +44,18 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.TimeUnit; +import java.util.function.Supplier; import static com.mongodb.assertions.Assertions.assertNotNull; import static com.mongodb.assertions.Assertions.notNull; +import static com.mongodb.internal.operation.AsyncOperationHelper.decorateWithRetriesAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.executeCommandAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.writeConcernErrorTransformerAsync; +import static com.mongodb.internal.operation.CommandOperationHelper.createSpecRetryControl; import static com.mongodb.internal.operation.IndexHelper.generateIndexName; import static com.mongodb.internal.operation.ServerVersionHelper.serverIsAtLeastVersionFourDotFour; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForWrite; +import static com.mongodb.internal.operation.SyncOperationHelper.decorateWithRetries; import static com.mongodb.internal.operation.SyncOperationHelper.executeCommand; import static com.mongodb.internal.operation.SyncOperationHelper.writeConcernErrorTransformer; import static com.mongodb.internal.operation.WriteConcernHelper.appendWriteConcernToCommand; @@ -63,13 +70,24 @@ public class CreateIndexesOperation implements WriteOperation { private final MongoNamespace namespace; private final List requests; private final WriteConcern writeConcern; + private final boolean retryWrites; + @Nullable + private final Integer maxAdaptiveRetriesSetting; private CreateIndexCommitQuorum commitQuorum; public CreateIndexesOperation(final MongoNamespace namespace, final List requests, @Nullable final WriteConcern writeConcern) { + this(namespace, requests, writeConcern, false, null); + } + + public CreateIndexesOperation(final MongoNamespace namespace, final List requests, + @Nullable final WriteConcern writeConcern, final boolean retryWrites, + @Nullable final Integer maxAdaptiveRetriesSetting) { this.namespace = notNull("namespace", namespace); this.requests = notNull("indexRequests", requests); this.writeConcern = writeConcern; + this.retryWrites = retryWrites; + this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting; } public WriteConcern getWriteConcern() { @@ -113,9 +131,16 @@ public MongoNamespace getNamespace() { @Override public Void execute(final WriteBinding binding, final OperationContext operationContext) { + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + Supplier retryingCommandExecutor = decorateWithRetries(retryControl, operationContext, () -> { + retryControl.getPolicy().onCommand(this::getCommandName); + return executeCommand(binding, operationContext, namespace.getDatabaseName(), getCommandCreator(), + writeConcernErrorTransformer(operationContext.getTimeoutContext())); + }); try { - return executeCommand(binding, operationContext, namespace.getDatabaseName(), getCommandCreator(), writeConcernErrorTransformer( - operationContext.getTimeoutContext())); + return retryingCommandExecutor.get(); } catch (MongoCommandException e) { throw checkForDuplicateKeyError(e); } @@ -123,14 +148,21 @@ public Void execute(final WriteBinding binding, final OperationContext operation @Override public void executeAsync(final AsyncWriteBinding binding, final OperationContext operationContext, final SingleResultCallback callback) { - executeCommandAsync(binding, operationContext, namespace.getDatabaseName(), getCommandCreator(), writeConcernErrorTransformerAsync(operationContext.getTimeoutContext()), - ((result, t) -> { - if (t != null) { - callback.onResult(null, translateException(t)); - } else { - callback.onResult(result, null); - } - })); + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + AsyncCallbackSupplier retryingCommandExecutor = decorateWithRetriesAsync(retryControl, operationContext, supplierCallback -> { + retryControl.getPolicy().onCommand(this::getCommandName); + executeCommandAsync(binding, operationContext, namespace.getDatabaseName(), getCommandCreator(), + writeConcernErrorTransformerAsync(operationContext.getTimeoutContext()), supplierCallback); + }); + retryingCommandExecutor.get((result, t) -> { + if (t != null) { + callback.onResult(null, translateException(t)); + } else { + callback.onResult(result, null); + } + }); } @SuppressWarnings("deprecation") @@ -232,4 +264,5 @@ private MongoException checkForDuplicateKeyError(final MongoCommandException e) return e; } } + } diff --git a/driver-core/src/main/com/mongodb/internal/operation/CreateSearchIndexesOperation.java b/driver-core/src/main/com/mongodb/internal/operation/CreateSearchIndexesOperation.java index bf75ee88b0..3847afaae5 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/CreateSearchIndexesOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/CreateSearchIndexesOperation.java @@ -18,6 +18,7 @@ import com.mongodb.MongoNamespace; import com.mongodb.client.model.SearchIndexType; +import com.mongodb.lang.Nullable; import org.bson.BsonArray; import org.bson.BsonDocument; import org.bson.BsonString; @@ -37,7 +38,12 @@ public final class CreateSearchIndexesOperation extends AbstractWriteSearchIndex private final List indexRequests; public CreateSearchIndexesOperation(final MongoNamespace namespace, final List indexRequests) { - super(namespace); + this(namespace, indexRequests, false, null); + } + + public CreateSearchIndexesOperation(final MongoNamespace namespace, final List indexRequests, + final boolean retryWrites, @Nullable final Integer maxAdaptiveRetriesSetting) { + super(namespace, retryWrites, maxAdaptiveRetriesSetting); this.indexRequests = assertNotNull(indexRequests); } diff --git a/driver-core/src/main/com/mongodb/internal/operation/CreateViewOperation.java b/driver-core/src/main/com/mongodb/internal/operation/CreateViewOperation.java index b80129093b..b2e95daf7e 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/CreateViewOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/CreateViewOperation.java @@ -20,6 +20,8 @@ import com.mongodb.WriteConcern; import com.mongodb.client.model.Collation; import com.mongodb.internal.async.SingleResultCallback; +import com.mongodb.internal.async.function.AsyncCallbackSupplier; +import com.mongodb.internal.async.function.RetryControl; import com.mongodb.internal.binding.AsyncWriteBinding; import com.mongodb.internal.binding.WriteBinding; import com.mongodb.internal.connection.OperationContext; @@ -30,14 +32,19 @@ import org.bson.codecs.BsonDocumentCodec; import java.util.List; +import java.util.function.Supplier; import static com.mongodb.assertions.Assertions.notNull; import static com.mongodb.internal.async.ErrorHandlingResultCallback.errorHandlingCallback; +import static com.mongodb.internal.operation.AsyncOperationHelper.decorateWithRetriesAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.executeCommandAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.releasingCallback; import static com.mongodb.internal.operation.AsyncOperationHelper.withAsyncConnection; import static com.mongodb.internal.operation.AsyncOperationHelper.writeConcernErrorTransformerAsync; +import static com.mongodb.internal.operation.CommandOperationHelper.createSpecRetryControl; import static com.mongodb.internal.operation.OperationHelper.LOGGER; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForWrite; +import static com.mongodb.internal.operation.SyncOperationHelper.decorateWithRetries; import static com.mongodb.internal.operation.SyncOperationHelper.executeCommand; import static com.mongodb.internal.operation.SyncOperationHelper.withConnection; import static com.mongodb.internal.operation.SyncOperationHelper.writeConcernErrorTransformer; @@ -54,15 +61,25 @@ public class CreateViewOperation implements WriteOperation { private final String viewOn; private final List pipeline; private final WriteConcern writeConcern; + private final boolean retryWrites; + @Nullable + private final Integer maxAdaptiveRetriesSetting; private Collation collation; public CreateViewOperation(final String databaseName, final String viewName, final String viewOn, final List pipeline, final WriteConcern writeConcern) { + this(databaseName, viewName, viewOn, pipeline, writeConcern, false, null); + } + + public CreateViewOperation(final String databaseName, final String viewName, final String viewOn, final List pipeline, + final WriteConcern writeConcern, final boolean retryWrites, @Nullable final Integer maxAdaptiveRetriesSetting) { this.databaseName = notNull("databaseName", databaseName); this.viewName = notNull("viewName", viewName); this.viewOn = notNull("viewOn", viewOn); this.pipeline = notNull("pipeline", pipeline); this.writeConcern = notNull("writeConcern", writeConcern); + this.retryWrites = retryWrites; + this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting; } public String getDatabaseName() { @@ -137,26 +154,38 @@ public MongoNamespace getNamespace() { @Override public Void execute(final WriteBinding binding, final OperationContext operationContext) { - return withConnection(binding, operationContext, (connection, operationContextWithMinRtt) -> { - executeCommand(binding, operationContextWithMinRtt, databaseName, getCommand(), new BsonDocumentCodec(), - writeConcernErrorTransformer(operationContextWithMinRtt.getTimeoutContext())); - return null; + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + Supplier retryingCommandExecutor = decorateWithRetries(retryControl, operationContext, () -> { + retryControl.getPolicy().onCommand(this::getCommandName); + return withConnection(binding, operationContext, (connection, operationContextWithMinRtt) -> { + executeCommand(binding, operationContextWithMinRtt, databaseName, getCommand(), new BsonDocumentCodec(), + writeConcernErrorTransformer(operationContextWithMinRtt.getTimeoutContext())); + return null; + }); }); + return retryingCommandExecutor.get(); } @Override public void executeAsync(final AsyncWriteBinding binding, final OperationContext operationContext, final SingleResultCallback callback) { - withAsyncConnection(binding, operationContext, (connection, operationContextWithMinRtt, t) -> { - SingleResultCallback errHandlingCallback = errorHandlingCallback(callback, LOGGER); - if (t != null) { - errHandlingCallback.onResult(null, t); - } else { - SingleResultCallback wrappedCallback = releasingCallback(errHandlingCallback, connection); - executeCommandAsync(binding, operationContextWithMinRtt, databaseName, getCommand(), connection, - writeConcernErrorTransformerAsync(operationContextWithMinRtt.getTimeoutContext()), - wrappedCallback); - } - }); + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + AsyncCallbackSupplier retryingCommandExecutor = decorateWithRetriesAsync(retryControl, operationContext, supplierCallback -> + withAsyncConnection(binding, operationContext, (connection, operationContextWithMinRtt, t) -> { + SingleResultCallback errHandlingCallback = errorHandlingCallback(supplierCallback, LOGGER); + if (t != null) { + errHandlingCallback.onResult(null, t); + } else { + SingleResultCallback wrappedCallback = releasingCallback(errHandlingCallback, connection); + executeCommandAsync(binding, operationContextWithMinRtt, databaseName, getCommand(), connection, + writeConcernErrorTransformerAsync(operationContextWithMinRtt.getTimeoutContext()), + wrappedCallback); + } + })); + retryingCommandExecutor.get(callback); } private BsonDocument getCommand() { @@ -170,4 +199,5 @@ private BsonDocument getCommand() { appendWriteConcernToCommand(writeConcern, commandDocument); return commandDocument; } + } diff --git a/driver-core/src/main/com/mongodb/internal/operation/DropCollectionOperation.java b/driver-core/src/main/com/mongodb/internal/operation/DropCollectionOperation.java index af6b73e6da..fcb13a02ca 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/DropCollectionOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/DropCollectionOperation.java @@ -18,14 +18,14 @@ import com.mongodb.MongoCommandException; import com.mongodb.MongoNamespace; -import com.mongodb.MongoOperationTimeoutException; import com.mongodb.WriteConcern; import com.mongodb.internal.async.SingleResultCallback; +import com.mongodb.internal.async.function.AsyncCallbackSupplier; +import com.mongodb.internal.async.function.RetryControl; import com.mongodb.internal.binding.AsyncReadWriteBinding; import com.mongodb.internal.binding.AsyncWriteBinding; import com.mongodb.internal.binding.ReadWriteBinding; import com.mongodb.internal.binding.WriteBinding; -import com.mongodb.internal.connection.AsyncConnection; import com.mongodb.internal.connection.OperationContext; import com.mongodb.lang.Nullable; import org.bson.BsonDocument; @@ -40,13 +40,17 @@ import static com.mongodb.assertions.Assertions.notNull; import static com.mongodb.internal.async.ErrorHandlingResultCallback.errorHandlingCallback; +import static com.mongodb.internal.operation.AsyncOperationHelper.decorateWithRetriesAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.executeCommandAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.releasingCallback; import static com.mongodb.internal.operation.AsyncOperationHelper.withAsyncConnection; import static com.mongodb.internal.operation.AsyncOperationHelper.writeConcernErrorTransformerAsync; +import static com.mongodb.internal.operation.CommandOperationHelper.createSpecRetryControl; import static com.mongodb.internal.operation.CommandOperationHelper.isNamespaceError; import static com.mongodb.internal.operation.CommandOperationHelper.rethrowIfNotNamespaceError; import static com.mongodb.internal.operation.OperationHelper.LOGGER; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForWrite; +import static com.mongodb.internal.operation.SyncOperationHelper.decorateWithRetries; import static com.mongodb.internal.operation.SyncOperationHelper.executeCommand; import static com.mongodb.internal.operation.SyncOperationHelper.withConnection; import static com.mongodb.internal.operation.SyncOperationHelper.writeConcernErrorTransformer; @@ -65,12 +69,22 @@ public class DropCollectionOperation implements WriteOperation { private static final BsonValueCodec BSON_VALUE_CODEC = new BsonValueCodec(); private final MongoNamespace namespace; private final WriteConcern writeConcern; + private final boolean retryWrites; + @Nullable + private final Integer maxAdaptiveRetriesSetting; private BsonDocument encryptedFields; private boolean autoEncryptedFields; public DropCollectionOperation(final MongoNamespace namespace, @Nullable final WriteConcern writeConcern) { + this(namespace, writeConcern, false, null); + } + + public DropCollectionOperation(final MongoNamespace namespace, @Nullable final WriteConcern writeConcern, + final boolean retryWrites, @Nullable final Integer maxAdaptiveRetriesSetting) { this.namespace = notNull("namespace", namespace); this.writeConcern = writeConcern; + this.retryWrites = retryWrites; + this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting; } public WriteConcern getWriteConcern() { @@ -100,38 +114,35 @@ public MongoNamespace getNamespace() { @Override public Void execute(final WriteBinding binding, final OperationContext operationContext) { BsonDocument localEncryptedFields = getEncryptedFields((ReadWriteBinding) binding, operationContext); - return withConnection(binding, operationContext, (connection, operationContextWithMinRtt) -> { - getCommands(localEncryptedFields).forEach(command -> { - try { - executeCommand(binding, operationContextWithMinRtt, namespace.getDatabaseName(), command.get(), - connection, writeConcernErrorTransformer(operationContextWithMinRtt.getTimeoutContext())); - } catch (MongoCommandException e) { - rethrowIfNotNamespaceError(e); - } + getCommands(localEncryptedFields).forEach(commandCreator -> { + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), operationContext); + Supplier retryingCommandExecutor = decorateWithRetries(retryControl, operationContext, () -> { + retryControl.getPolicy().onCommand(this::getCommandName); + return withConnection(binding, operationContext, (connection, connectionScopedOperationContext) -> { + try { + executeCommand(binding, connectionScopedOperationContext, namespace.getDatabaseName(), commandCreator.get(), + connection, writeConcernErrorTransformer(connectionScopedOperationContext.getTimeoutContext())); + } catch (MongoCommandException e) { + rethrowIfNotNamespaceError(e); + } + return null; + }); }); - return null; + retryingCommandExecutor.get(); }); + return null; } @Override public void executeAsync(final AsyncWriteBinding binding, final OperationContext operationContext, final SingleResultCallback callback) { - SingleResultCallback errHandlingCallback = errorHandlingCallback(callback, LOGGER); - getEncryptedFields((AsyncReadWriteBinding) binding, operationContext, (result, t) -> { + getEncryptedFields((AsyncReadWriteBinding) binding, operationContext, (localEncryptedFields, t) -> { if (t != null) { - errHandlingCallback.onResult(null, t); - } else { - withAsyncConnection(binding, operationContext, (connection, operationContextWithMinRtt, t1) -> { - if (t1 != null) { - errHandlingCallback.onResult(null, t1); - } else { - new ProcessCommandsCallback(binding, operationContextWithMinRtt, connection, getCommands(result), - releasingCallback(errHandlingCallback, - connection)) - .onResult(null, null); - } - }); + errorHandlingCallback(callback, LOGGER).onResult(null, t); + return; } + new ProcessCommandsCallback(binding, operationContext, getCommands(localEncryptedFields), callback).onResult(null, null); }); } @@ -239,19 +250,16 @@ private ListCollectionsOperation listCollectionOperation() { class ProcessCommandsCallback implements SingleResultCallback { private final AsyncWriteBinding binding; private final OperationContext operationContext; - private final AsyncConnection connection; private final SingleResultCallback finalCallback; private final Deque> commands; ProcessCommandsCallback( final AsyncWriteBinding binding, final OperationContext operationContext, - final AsyncConnection connection, final List> commands, final SingleResultCallback finalCallback) { this.binding = binding; this.operationContext = operationContext; - this.connection = connection; this.finalCallback = finalCallback; this.commands = new ArrayDeque<>(commands); } @@ -266,14 +274,27 @@ public void onResult(@Nullable final Void result, @Nullable final Throwable t) { if (nextCommandFunction == null) { finalCallback.onResult(null, null); } else { - try { - executeCommandAsync(binding, operationContext, namespace.getDatabaseName(), nextCommandFunction.get(), - connection, writeConcernErrorTransformerAsync(operationContext.getTimeoutContext()), this); - } catch (MongoOperationTimeoutException operationTimeoutException) { - finalCallback.onResult(null, operationTimeoutException); - } + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), operationContext); + AsyncCallbackSupplier retryingCommandExecutor = decorateWithRetriesAsync(retryControl, operationContext, + supplierCallback -> { + retryControl.getPolicy().onCommand(DropCollectionOperation.this::getCommandName); + withAsyncConnection(binding, operationContext, (connection, connectionScopedOperationContext, t1) -> { + if (t1 != null) { + supplierCallback.onResult(null, t1); + } else { + SingleResultCallback connectionReleasingCallback = releasingCallback(supplierCallback, connection); + executeCommandAsync(binding, connectionScopedOperationContext, namespace.getDatabaseName(), + nextCommandFunction.get(), connection, + writeConcernErrorTransformerAsync(connectionScopedOperationContext.getTimeoutContext()), + connectionReleasingCallback); + } + }); + }); + retryingCommandExecutor.get(this); } } } + } diff --git a/driver-core/src/main/com/mongodb/internal/operation/DropDatabaseOperation.java b/driver-core/src/main/com/mongodb/internal/operation/DropDatabaseOperation.java index ba76be13e5..d44549d339 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/DropDatabaseOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/DropDatabaseOperation.java @@ -20,6 +20,8 @@ import com.mongodb.WriteConcern; import com.mongodb.internal.MongoNamespaceHelper; import com.mongodb.internal.async.SingleResultCallback; +import com.mongodb.internal.async.function.AsyncCallbackSupplier; +import com.mongodb.internal.async.function.RetryControl; import com.mongodb.internal.binding.AsyncWriteBinding; import com.mongodb.internal.binding.WriteBinding; import com.mongodb.internal.connection.OperationContext; @@ -27,13 +29,19 @@ import org.bson.BsonDocument; import org.bson.BsonInt32; +import java.util.function.Supplier; + import static com.mongodb.assertions.Assertions.notNull; import static com.mongodb.internal.async.ErrorHandlingResultCallback.errorHandlingCallback; +import static com.mongodb.internal.operation.AsyncOperationHelper.decorateWithRetriesAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.executeCommandAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.releasingCallback; import static com.mongodb.internal.operation.AsyncOperationHelper.withAsyncConnection; import static com.mongodb.internal.operation.AsyncOperationHelper.writeConcernErrorTransformerAsync; +import static com.mongodb.internal.operation.CommandOperationHelper.createSpecRetryControl; import static com.mongodb.internal.operation.OperationHelper.LOGGER; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForWrite; +import static com.mongodb.internal.operation.SyncOperationHelper.decorateWithRetries; import static com.mongodb.internal.operation.SyncOperationHelper.executeCommand; import static com.mongodb.internal.operation.SyncOperationHelper.withConnection; import static com.mongodb.internal.operation.SyncOperationHelper.writeConcernErrorTransformer; @@ -48,10 +56,20 @@ public class DropDatabaseOperation implements WriteOperation { private final String databaseName; private final WriteConcern writeConcern; + private final boolean retryWrites; + @Nullable + private final Integer maxAdaptiveRetriesSetting; public DropDatabaseOperation(final String databaseName, @Nullable final WriteConcern writeConcern) { + this(databaseName, writeConcern, false, null); + } + + public DropDatabaseOperation(final String databaseName, @Nullable final WriteConcern writeConcern, + final boolean retryWrites, @Nullable final Integer maxAdaptiveRetriesSetting) { this.databaseName = notNull("databaseName", databaseName); this.writeConcern = writeConcern; + this.retryWrites = retryWrites; + this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting; } public WriteConcern getWriteConcern() { @@ -70,26 +88,37 @@ public MongoNamespace getNamespace() { @Override public Void execute(final WriteBinding binding, final OperationContext operationContext) { - return withConnection(binding, operationContext, (connection, operationContextWithMinRtt) -> { - executeCommand(binding, operationContextWithMinRtt, databaseName, getCommand(), connection, writeConcernErrorTransformer(operationContextWithMinRtt - .getTimeoutContext())); - return null; + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + Supplier retryingCommandExecutor = decorateWithRetries(retryControl, operationContext, () -> { + retryControl.getPolicy().onCommand(this::getCommandName); + return withConnection(binding, operationContext, (connection, operationContextWithMinRtt) -> { + executeCommand(binding, operationContextWithMinRtt, databaseName, getCommand(), connection, + writeConcernErrorTransformer(operationContextWithMinRtt.getTimeoutContext())); + return null; + }); }); + return retryingCommandExecutor.get(); } @Override public void executeAsync(final AsyncWriteBinding binding, final OperationContext operationContext, final SingleResultCallback callback) { - withAsyncConnection(binding, operationContext, (connection, operationContextWithMinRtt, t) -> { - SingleResultCallback errHandlingCallback = errorHandlingCallback(callback, LOGGER); - if (t != null) { - errHandlingCallback.onResult(null, t); - } else { - executeCommandAsync(binding, operationContextWithMinRtt, databaseName, getCommand(), connection, - writeConcernErrorTransformerAsync(operationContextWithMinRtt.getTimeoutContext()), - releasingCallback(errHandlingCallback, connection)); - - } - }); + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + AsyncCallbackSupplier retryingCommandExecutor = decorateWithRetriesAsync(retryControl, operationContext, supplierCallback -> + withAsyncConnection(binding, operationContext, (connection, operationContextWithMinRtt, t) -> { + SingleResultCallback errHandlingCallback = errorHandlingCallback(supplierCallback, LOGGER); + if (t != null) { + errHandlingCallback.onResult(null, t); + } else { + executeCommandAsync(binding, operationContextWithMinRtt, databaseName, getCommand(), connection, + writeConcernErrorTransformerAsync(operationContextWithMinRtt.getTimeoutContext()), + releasingCallback(errHandlingCallback, connection)); + } + })); + retryingCommandExecutor.get(callback); } private BsonDocument getCommand() { @@ -97,4 +126,5 @@ private BsonDocument getCommand() { appendWriteConcernToCommand(writeConcern, commandDocument); return commandDocument; } + } diff --git a/driver-core/src/main/com/mongodb/internal/operation/DropIndexOperation.java b/driver-core/src/main/com/mongodb/internal/operation/DropIndexOperation.java index d77074c738..83efcdc692 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/DropIndexOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/DropIndexOperation.java @@ -20,6 +20,8 @@ import com.mongodb.MongoNamespace; import com.mongodb.WriteConcern; import com.mongodb.internal.async.SingleResultCallback; +import com.mongodb.internal.async.function.AsyncCallbackSupplier; +import com.mongodb.internal.async.function.RetryControl; import com.mongodb.internal.binding.AsyncWriteBinding; import com.mongodb.internal.binding.WriteBinding; import com.mongodb.internal.connection.OperationContext; @@ -27,11 +29,17 @@ import org.bson.BsonDocument; import org.bson.BsonString; +import java.util.function.Supplier; + import static com.mongodb.assertions.Assertions.notNull; +import static com.mongodb.internal.operation.AsyncOperationHelper.decorateWithRetriesAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.executeCommandAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.writeConcernErrorTransformerAsync; +import static com.mongodb.internal.operation.CommandOperationHelper.createSpecRetryControl; import static com.mongodb.internal.operation.CommandOperationHelper.isNamespaceError; import static com.mongodb.internal.operation.CommandOperationHelper.rethrowIfNotNamespaceError; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForWrite; +import static com.mongodb.internal.operation.SyncOperationHelper.decorateWithRetries; import static com.mongodb.internal.operation.SyncOperationHelper.executeCommand; import static com.mongodb.internal.operation.SyncOperationHelper.writeConcernErrorTransformer; import static com.mongodb.internal.operation.WriteConcernHelper.appendWriteConcernToCommand; @@ -47,19 +55,36 @@ public class DropIndexOperation implements WriteOperation { private final String indexName; private final BsonDocument indexKeys; private final WriteConcern writeConcern; + private final boolean retryWrites; + @Nullable + private final Integer maxAdaptiveRetriesSetting; public DropIndexOperation(final MongoNamespace namespace, final String indexName, @Nullable final WriteConcern writeConcern) { + this(namespace, indexName, writeConcern, false, null); + } + + public DropIndexOperation(final MongoNamespace namespace, final BsonDocument indexKeys, @Nullable final WriteConcern writeConcern) { + this(namespace, indexKeys, writeConcern, false, null); + } + + public DropIndexOperation(final MongoNamespace namespace, final String indexName, @Nullable final WriteConcern writeConcern, + final boolean retryWrites, @Nullable final Integer maxAdaptiveRetriesSetting) { this.namespace = notNull("namespace", namespace); this.indexName = notNull("indexName", indexName); this.indexKeys = null; this.writeConcern = writeConcern; + this.retryWrites = retryWrites; + this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting; } - public DropIndexOperation(final MongoNamespace namespace, final BsonDocument indexKeys, @Nullable final WriteConcern writeConcern) { + public DropIndexOperation(final MongoNamespace namespace, final BsonDocument indexKeys, @Nullable final WriteConcern writeConcern, + final boolean retryWrites, @Nullable final Integer maxAdaptiveRetriesSetting) { this.namespace = notNull("namespace", namespace); this.indexKeys = notNull("indexKeys", indexKeys); this.indexName = null; this.writeConcern = writeConcern; + this.retryWrites = retryWrites; + this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting; } public WriteConcern getWriteConcern() { @@ -78,9 +103,17 @@ public MongoNamespace getNamespace() { @Override public Void execute(final WriteBinding binding, final OperationContext operationContext) { - try { + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + Supplier retryingCommandExecutor = decorateWithRetries(retryControl, operationContext, () -> { + retryControl.getPolicy().onCommand(this::getCommandName); executeCommand(binding, operationContext, namespace.getDatabaseName(), getCommandCreator(), writeConcernErrorTransformer( operationContext.getTimeoutContext())); + return null; + }); + try { + retryingCommandExecutor.get(); } catch (MongoCommandException e) { rethrowIfNotNamespaceError(e); } @@ -90,8 +123,15 @@ public Void execute(final WriteBinding binding, final OperationContext operation @Override public void executeAsync(final AsyncWriteBinding binding, final OperationContext operationContext, final SingleResultCallback callback) { - executeCommandAsync(binding, operationContext, namespace.getDatabaseName(), getCommandCreator(), - writeConcernErrorTransformerAsync(operationContext.getTimeoutContext()), (result, t) -> { + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + AsyncCallbackSupplier retryingCommandExecutor = decorateWithRetriesAsync(retryControl, operationContext, supplierCallback -> { + retryControl.getPolicy().onCommand(this::getCommandName); + executeCommandAsync(binding, operationContext, namespace.getDatabaseName(), getCommandCreator(), + writeConcernErrorTransformerAsync(operationContext.getTimeoutContext()), supplierCallback); + }); + retryingCommandExecutor.get((result, t) -> { if (t != null && !isNamespaceError(t)) { callback.onResult(null, t); } else { @@ -112,4 +152,5 @@ private CommandOperationHelper.CommandCreator getCommandCreator() { return command; }; } + } diff --git a/driver-core/src/main/com/mongodb/internal/operation/DropSearchIndexOperation.java b/driver-core/src/main/com/mongodb/internal/operation/DropSearchIndexOperation.java index a440dbd0e7..98f597bdc4 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/DropSearchIndexOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/DropSearchIndexOperation.java @@ -33,7 +33,12 @@ final class DropSearchIndexOperation extends AbstractWriteSearchIndexOperation { private final String indexName; DropSearchIndexOperation(final MongoNamespace namespace, final String indexName) { - super(namespace); + this(namespace, indexName, false, null); + } + + DropSearchIndexOperation(final MongoNamespace namespace, final String indexName, + final boolean retryWrites, @Nullable final Integer maxAdaptiveRetriesSetting) { + super(namespace, retryWrites, maxAdaptiveRetriesSetting); this.indexName = indexName; } diff --git a/driver-core/src/main/com/mongodb/internal/operation/Operations.java b/driver-core/src/main/com/mongodb/internal/operation/Operations.java index 1b014119a6..8e1e4ab61c 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/Operations.java +++ b/driver-core/src/main/com/mongodb/internal/operation/Operations.java @@ -337,7 +337,8 @@ public ReadOperationSimple aggregateToCollection(final List ReadOperationSimple commandRead(final Bson command, final Class public WriteOperation dropDatabase() { return new DropDatabaseOperation(assertNotNull(namespace).getDatabaseName(), - getWriteConcern()); + getWriteConcern(), isRetryWrites(), maxAdaptiveRetriesSetting); } public WriteOperation createCollection(final String collectionName, final CreateCollectionOptions createCollectionOptions, @Nullable final AutoEncryptionSettings autoEncryptionSettings) { CreateCollectionOperation operation = new CreateCollectionOperation( - assertNotNull(namespace).getDatabaseName(), collectionName, writeConcern) + assertNotNull(namespace).getDatabaseName(), collectionName, writeConcern, isRetryWrites(), maxAdaptiveRetriesSetting) .collation(createCollectionOptions.getCollation()) .capped(createCollectionOptions.isCapped()) .sizeInBytes(createCollectionOptions.getSizeInBytes()) @@ -662,7 +663,7 @@ public WriteOperation dropCollection( final DropCollectionOptions dropCollectionOptions, @Nullable final AutoEncryptionSettings autoEncryptionSettings) { DropCollectionOperation operation = new DropCollectionOperation( - assertNotNull(namespace), writeConcern); + assertNotNull(namespace), writeConcern, isRetryWrites(), maxAdaptiveRetriesSetting); Bson encryptedFields = dropCollectionOptions.getEncryptedFields(); if (encryptedFields != null) { operation.encryptedFields(assertNotNull(toBsonDocument(encryptedFields))); @@ -680,7 +681,8 @@ public WriteOperation dropCollection( public WriteOperation renameCollection(final MongoNamespace newCollectionNamespace, final RenameCollectionOptions renameCollectionOptions) { return new RenameCollectionOperation(assertNotNull(namespace), - newCollectionNamespace, writeConcern).dropTarget(renameCollectionOptions.isDropTarget()); + newCollectionNamespace, writeConcern, isRetryWrites(), maxAdaptiveRetriesSetting) + .dropTarget(renameCollectionOptions.isDropTarget()); } public WriteOperation createView(final String viewName, final String viewOn, final List pipeline, @@ -688,7 +690,8 @@ public WriteOperation createView(final String viewName, final String viewO notNull("options", createViewOptions); notNull("pipeline", pipeline); return new CreateViewOperation(assertNotNull(namespace).getDatabaseName(), viewName, - viewOn, assertNotNull(toBsonDocumentList(pipeline)), writeConcern).collation(createViewOptions.getCollation()); + viewOn, assertNotNull(toBsonDocumentList(pipeline)), writeConcern, isRetryWrites(), maxAdaptiveRetriesSetting) + .collation(createViewOptions.getCollation()); } public WriteOperation createIndexes(final List indexes, final CreateIndexOptions createIndexOptions) { @@ -722,7 +725,7 @@ public WriteOperation createIndexes(final List indexes, final ); } return new CreateIndexesOperation( - assertNotNull(namespace), indexRequests, writeConcern) + assertNotNull(namespace), indexRequests, writeConcern, isRetryWrites(), maxAdaptiveRetriesSetting) .commitQuorum(createIndexOptions.getCommitQuorum()); } @@ -730,18 +733,18 @@ public WriteOperation createSearchIndexes(final List ind List indexRequests = indexes.stream() .map(this::createSearchIndexRequest) .collect(Collectors.toList()); - return new CreateSearchIndexesOperation(assertNotNull(namespace), indexRequests); + return new CreateSearchIndexesOperation(assertNotNull(namespace), indexRequests, isRetryWrites(), maxAdaptiveRetriesSetting); } public WriteOperation updateSearchIndex(final String indexName, final Bson definition) { BsonDocument definitionDocument = assertNotNull(toBsonDocument(definition)); SearchIndexRequest searchIndexRequest = new SearchIndexRequest(definitionDocument, indexName); - return new UpdateSearchIndexesOperation(assertNotNull(namespace), searchIndexRequest); + return new UpdateSearchIndexesOperation(assertNotNull(namespace), searchIndexRequest, isRetryWrites(), maxAdaptiveRetriesSetting); } public WriteOperation dropSearchIndex(final String indexName) { - return new DropSearchIndexOperation(assertNotNull(namespace), indexName); + return new DropSearchIndexOperation(assertNotNull(namespace), indexName, isRetryWrites(), maxAdaptiveRetriesSetting); } @@ -753,11 +756,12 @@ public ReadOperationExplainable listSearchIndexes(final Class resultCl } public WriteOperation dropIndex(final String indexName, final DropIndexOptions ignoredOptions) { - return new DropIndexOperation(assertNotNull(namespace), indexName, writeConcern); + return new DropIndexOperation(assertNotNull(namespace), indexName, writeConcern, isRetryWrites(), maxAdaptiveRetriesSetting); } public WriteOperation dropIndex(final Bson keys, final DropIndexOptions ignoredOptions) { - return new DropIndexOperation(assertNotNull(namespace), keys.toBsonDocument(BsonDocument.class, codecRegistry), writeConcern); + return new DropIndexOperation(assertNotNull(namespace), keys.toBsonDocument(BsonDocument.class, codecRegistry), writeConcern, + isRetryWrites(), maxAdaptiveRetriesSetting); } public ReadOperationCursor listCollections(final String databaseName, final Class resultClass, diff --git a/driver-core/src/main/com/mongodb/internal/operation/RenameCollectionOperation.java b/driver-core/src/main/com/mongodb/internal/operation/RenameCollectionOperation.java index dfc21c3b7e..81a5fa0803 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/RenameCollectionOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/RenameCollectionOperation.java @@ -19,6 +19,8 @@ import com.mongodb.MongoNamespace; import com.mongodb.WriteConcern; import com.mongodb.internal.async.SingleResultCallback; +import com.mongodb.internal.async.function.AsyncCallbackSupplier; +import com.mongodb.internal.async.function.RetryControl; import com.mongodb.internal.binding.AsyncWriteBinding; import com.mongodb.internal.binding.WriteBinding; import com.mongodb.internal.connection.OperationContext; @@ -27,14 +29,20 @@ import org.bson.BsonDocument; import org.bson.BsonString; +import java.util.function.Supplier; + import static com.mongodb.assertions.Assertions.assertNotNull; import static com.mongodb.assertions.Assertions.notNull; import static com.mongodb.internal.async.ErrorHandlingResultCallback.errorHandlingCallback; +import static com.mongodb.internal.operation.AsyncOperationHelper.decorateWithRetriesAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.executeCommandAsync; import static com.mongodb.internal.operation.AsyncOperationHelper.releasingCallback; import static com.mongodb.internal.operation.AsyncOperationHelper.withAsyncConnection; import static com.mongodb.internal.operation.AsyncOperationHelper.writeConcernErrorTransformerAsync; +import static com.mongodb.internal.operation.CommandOperationHelper.createSpecRetryControl; import static com.mongodb.internal.operation.OperationHelper.LOGGER; +import static com.mongodb.internal.operation.SpecRetryPolicy.IndividualPolicies.overloadForWrite; +import static com.mongodb.internal.operation.SyncOperationHelper.decorateWithRetries; import static com.mongodb.internal.operation.SyncOperationHelper.executeCommand; import static com.mongodb.internal.operation.SyncOperationHelper.withConnection; import static com.mongodb.internal.operation.SyncOperationHelper.writeConcernErrorTransformer; @@ -53,13 +61,24 @@ public class RenameCollectionOperation implements WriteOperation { private final MongoNamespace originalNamespace; private final MongoNamespace newNamespace; private final WriteConcern writeConcern; + private final boolean retryWrites; + @Nullable + private final Integer maxAdaptiveRetriesSetting; private boolean dropTarget; public RenameCollectionOperation(final MongoNamespace originalNamespace, final MongoNamespace newNamespace, @Nullable final WriteConcern writeConcern) { + this(originalNamespace, newNamespace, writeConcern, false, null); + } + + public RenameCollectionOperation(final MongoNamespace originalNamespace, final MongoNamespace newNamespace, + @Nullable final WriteConcern writeConcern, final boolean retryWrites, + @Nullable final Integer maxAdaptiveRetriesSetting) { this.originalNamespace = notNull("originalNamespace", originalNamespace); this.newNamespace = notNull("newNamespace", newNamespace); this.writeConcern = writeConcern; + this.retryWrites = retryWrites; + this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting; } public WriteConcern getWriteConcern() { @@ -87,24 +106,36 @@ public MongoNamespace getNamespace() { @Override public Void execute(final WriteBinding binding, final OperationContext operationContext) { - return withConnection(binding, operationContext, (connection, operationContextWithMinRtt) -> - executeCommand(binding, - operationContextWithMinRtt, "admin", getCommand(), connection, - writeConcernErrorTransformer(operationContextWithMinRtt.getTimeoutContext()))); + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + Supplier retryingCommandExecutor = decorateWithRetries(retryControl, operationContext, () -> { + retryControl.getPolicy().onCommand(this::getCommandName); + return withConnection(binding, operationContext, (connection, operationContextWithMinRtt) -> + executeCommand(binding, + operationContextWithMinRtt, "admin", getCommand(), connection, + writeConcernErrorTransformer(operationContextWithMinRtt.getTimeoutContext()))); + }); + return retryingCommandExecutor.get(); } @Override public void executeAsync(final AsyncWriteBinding binding, final OperationContext operationContext, final SingleResultCallback callback) { - withAsyncConnection(binding, operationContext, (connection, operationContextWithMinRtt, t) -> { - SingleResultCallback errHandlingCallback = errorHandlingCallback(callback, LOGGER); - if (t != null) { - errHandlingCallback.onResult(null, t); - } else { - executeCommandAsync(binding, operationContextWithMinRtt, "admin", getCommand(), assertNotNull(connection), - writeConcernErrorTransformerAsync(operationContextWithMinRtt.getTimeoutContext()), - releasingCallback(errHandlingCallback, connection)); - } - }); + RetryControl retryControl = createSpecRetryControl( + overloadForWrite(retryWrites, maxAdaptiveRetriesSetting), + operationContext); + AsyncCallbackSupplier retryingCommandExecutor = decorateWithRetriesAsync(retryControl, operationContext, supplierCallback -> + withAsyncConnection(binding, operationContext, (connection, operationContextWithMinRtt, t) -> { + SingleResultCallback errHandlingCallback = errorHandlingCallback(supplierCallback, LOGGER); + if (t != null) { + errHandlingCallback.onResult(null, t); + } else { + executeCommandAsync(binding, operationContextWithMinRtt, "admin", getCommand(), assertNotNull(connection), + writeConcernErrorTransformerAsync(operationContextWithMinRtt.getTimeoutContext()), + releasingCallback(errHandlingCallback, connection)); + } + })); + retryingCommandExecutor.get(callback); } private BsonDocument getCommand() { @@ -114,4 +145,5 @@ private BsonDocument getCommand() { appendWriteConcernToCommand(writeConcern, commandDocument); return commandDocument; } + } diff --git a/driver-core/src/main/com/mongodb/internal/operation/SpecRetryPolicy.java b/driver-core/src/main/com/mongodb/internal/operation/SpecRetryPolicy.java index 35adcf13e2..650423c7d5 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/SpecRetryPolicy.java +++ b/driver-core/src/main/com/mongodb/internal/operation/SpecRetryPolicy.java @@ -331,6 +331,16 @@ static final class IndividualPolicies { this.policies = new EnumMap<>(Descriptor.class); } + static IndividualPolicies overloadForWrite(final boolean retryWrites, @Nullable final Integer maxAdaptiveRetriesSetting) { + return new IndividualPolicies(retryWrites) + .includeOverload(maxAdaptiveRetriesSetting, ErrorPropagation.AS_WRITE_POLICY); + } + + static IndividualPolicies overloadForRead(final boolean retryReads, @Nullable final Integer maxAdaptiveRetriesSetting) { + return new IndividualPolicies(retryReads) + .includeOverload(maxAdaptiveRetriesSetting, ErrorPropagation.AS_READ_POLICY); + } + private IndividualPolicies assertValid() { assertFalse(policies.isEmpty()); assertNoConflicts(policies.keySet()); @@ -692,7 +702,7 @@ private int maxAttempts(final IndividualPolicies policies) { * Selects the error propagation shape for overload-only policy * compositions (see {@link IndividualPolicies#includeOverload(Integer, ErrorPropagation)}). */ - enum ErrorPropagation { + private enum ErrorPropagation { AS_READ_POLICY, AS_WRITE_POLICY } diff --git a/driver-core/src/main/com/mongodb/internal/operation/UpdateSearchIndexesOperation.java b/driver-core/src/main/com/mongodb/internal/operation/UpdateSearchIndexesOperation.java index ca23fd8e50..429850e752 100644 --- a/driver-core/src/main/com/mongodb/internal/operation/UpdateSearchIndexesOperation.java +++ b/driver-core/src/main/com/mongodb/internal/operation/UpdateSearchIndexesOperation.java @@ -17,6 +17,7 @@ package com.mongodb.internal.operation; import com.mongodb.MongoNamespace; +import com.mongodb.lang.Nullable; import org.bson.BsonDocument; import org.bson.BsonString; @@ -30,7 +31,12 @@ final class UpdateSearchIndexesOperation extends AbstractWriteSearchIndexOperati private final SearchIndexRequest request; UpdateSearchIndexesOperation(final MongoNamespace namespace, final SearchIndexRequest request) { - super(namespace); + this(namespace, request, false, null); + } + + UpdateSearchIndexesOperation(final MongoNamespace namespace, final SearchIndexRequest request, + final boolean retryWrites, @Nullable final Integer maxAdaptiveRetriesSetting) { + super(namespace, retryWrites, maxAdaptiveRetriesSetting); this.request = request; } diff --git a/driver-legacy/src/main/com/mongodb/DB.java b/driver-legacy/src/main/com/mongodb/DB.java index 9af169194b..494970e261 100644 --- a/driver-legacy/src/main/com/mongodb/DB.java +++ b/driver-legacy/src/main/com/mongodb/DB.java @@ -195,7 +195,8 @@ public DBCollection getCollection(final String name) { */ public void dropDatabase() { try { - getExecutor().execute(new DropDatabaseOperation(getName(), getWriteConcern()), getReadConcern()); + getExecutor().execute(new DropDatabaseOperation(getName(), getWriteConcern(), + mongo.getMongoClientOptions().getRetryWrites(), mongo.getMongoClientOptions().getMaxAdaptiveRetries()), getReadConcern()); } catch (MongoWriteConcernException e) { throw createWriteConcernException(e); } @@ -310,7 +311,8 @@ public DBCollection createView(final String viewName, final String viewOn, final notNull("options", options); DBCollection view = getCollection(viewName); executor.execute(new CreateViewOperation(name, viewName, viewOn, - view.preparePipeline(pipeline), writeConcern) + view.preparePipeline(pipeline), writeConcern, + mongo.getMongoClientOptions().getRetryWrites(), mongo.getMongoClientOptions().getMaxAdaptiveRetries()) .collation(options.getCollation()), getReadConcern()); return view; } catch (MongoWriteConcernException e) { @@ -387,7 +389,7 @@ private CreateCollectionOperation getCreateCollectionOperation(final String coll } Collation collation = DBObjectCollationHelper.createCollationFromOptions(options); return new CreateCollectionOperation(getName(), collectionName, - getWriteConcern()) + getWriteConcern(), mongo.getMongoClientOptions().getRetryWrites(), mongo.getMongoClientOptions().getMaxAdaptiveRetries()) .capped(capped) .collation(collation) .sizeInBytes(sizeInBytes) diff --git a/driver-legacy/src/main/com/mongodb/DBCollection.java b/driver-legacy/src/main/com/mongodb/DBCollection.java index 4084bba17d..fdfa0767b8 100644 --- a/driver-legacy/src/main/com/mongodb/DBCollection.java +++ b/driver-legacy/src/main/com/mongodb/DBCollection.java @@ -32,6 +32,7 @@ import com.mongodb.internal.bulk.InsertRequest; import com.mongodb.internal.bulk.UpdateRequest; import com.mongodb.internal.bulk.WriteRequest.Type; +import com.mongodb.internal.client.model.AggregationLevel; import com.mongodb.internal.connection.PowerOfTwoBufferPool; import com.mongodb.internal.operation.AggregateOperation; import com.mongodb.internal.operation.AggregateToCollectionOperation; @@ -979,7 +980,8 @@ public DBCollection rename(final String newName) { public DBCollection rename(final String newName, final boolean dropTarget) { try { executor.execute(new RenameCollectionOperation(getNamespace(), - new MongoNamespace(getNamespace().getDatabaseName(), newName), getWriteConcern()) + new MongoNamespace(getNamespace().getDatabaseName(), newName), getWriteConcern(), + retryWrites, maxAdaptiveRetriesSetting) .dropTarget(dropTarget), getReadConcern()); return getDB().getCollection(newName); } catch (MongoWriteConcernException e) { @@ -1249,7 +1251,8 @@ public Cursor aggregate(final List pipeline, final Aggregati if (outCollection != null) { AggregateToCollectionOperation operation = new AggregateToCollectionOperation( - getNamespace(), stages, getReadConcern(), getWriteConcern()) + getNamespace(), stages, getReadConcern(), getWriteConcern(), AggregationLevel.COLLECTION, + retryWrites, maxAdaptiveRetriesSetting) .allowDiskUse(options.getAllowDiskUse()) .bypassDocumentValidation(options.getBypassDocumentValidation()) .collation(options.getCollation()); @@ -1818,7 +1821,7 @@ public ReadConcern getReadConcern() { public void drop() { try { executor.execute(new DropCollectionOperation(getNamespace(), - getWriteConcern()), getReadConcern()); + getWriteConcern(), retryWrites, maxAdaptiveRetriesSetting), getReadConcern()); } catch (MongoWriteConcernException e) { throw createWriteConcernException(e); } @@ -1913,7 +1916,7 @@ public OperationExecutor getExecutor() { public void dropIndex(final DBObject index) { try { executor.execute(new DropIndexOperation(getNamespace(), wrap(index), - getWriteConcern()), getReadConcern()); + getWriteConcern(), retryWrites, maxAdaptiveRetriesSetting), getReadConcern()); } catch (MongoWriteConcernException e) { throw createWriteConcernException(e); } @@ -1929,7 +1932,7 @@ public void dropIndex(final DBObject index) { public void dropIndex(final String indexName) { try { executor.execute(new DropIndexOperation(getNamespace(), indexName, - getWriteConcern()), getReadConcern()); + getWriteConcern(), retryWrites, maxAdaptiveRetriesSetting), getReadConcern()); } catch (MongoWriteConcernException e) { throw createWriteConcernException(e); } @@ -2156,7 +2159,7 @@ private CreateIndexesOperation createIndexOperation(final DBObject key, final DB if (options.containsField("collation")) { request.collation(DBObjectCollationHelper.createCollationFromOptions(options)); } - return new CreateIndexesOperation(getNamespace(), singletonList(request), writeConcern); + return new CreateIndexesOperation(getNamespace(), singletonList(request), writeConcern, retryWrites, maxAdaptiveRetriesSetting); } Codec getObjectCodec() { diff --git a/driver-legacy/src/test/functional/com/mongodb/DBCollectionSpecification.groovy b/driver-legacy/src/test/functional/com/mongodb/DBCollectionSpecification.groovy index 785044308e..c546017f5b 100644 --- a/driver-legacy/src/test/functional/com/mongodb/DBCollectionSpecification.groovy +++ b/driver-legacy/src/test/functional/com/mongodb/DBCollectionSpecification.groovy @@ -32,6 +32,7 @@ import com.mongodb.internal.bulk.DeleteRequest import com.mongodb.internal.bulk.IndexRequest import com.mongodb.internal.bulk.InsertRequest import com.mongodb.internal.bulk.UpdateRequest +import com.mongodb.internal.client.model.AggregationLevel import com.mongodb.internal.operation.AggregateOperation import com.mongodb.internal.operation.AggregateToCollectionOperation import com.mongodb.internal.operation.BatchCursor @@ -661,21 +662,22 @@ class DBCollectionSpecification extends Specification { then: expect executor.getReadOperation(), isTheSameAs(new AggregateToCollectionOperation(collection.getNamespace(), - bsonPipeline, collection.getReadConcern(), collection.getWriteConcern())) + bsonPipeline, collection.getReadConcern(), collection.getWriteConcern(), AggregationLevel.COLLECTION, true, null)) when: // Inherits from DB collection.aggregate(pipeline, AggregationOptions.builder().build()) then: expect executor.getReadOperation(), isTheSameAs(new AggregateToCollectionOperation(collection.getNamespace(), - bsonPipeline, collection.getReadConcern(), collection.getWriteConcern())) + bsonPipeline, collection.getReadConcern(), collection.getWriteConcern(), AggregationLevel.COLLECTION, true, null)) when: collection.aggregate(pipeline, AggregationOptions.builder().collation(collation).build()) then: expect executor.getReadOperation(), isTheSameAs(new AggregateToCollectionOperation(collection.getNamespace(), - bsonPipeline, collection.getReadConcern(), collection.getWriteConcern()).collation(collation)) + bsonPipeline, collection.getReadConcern(), collection.getWriteConcern(), + AggregationLevel.COLLECTION, true, null).collation(collation)) } def 'explainAggregate should create the correct AggregateOperation'() { diff --git a/driver-legacy/src/test/unit/com/mongodb/DBSpecification.groovy b/driver-legacy/src/test/unit/com/mongodb/DBSpecification.groovy index c06dd67a3a..832d75f520 100644 --- a/driver-legacy/src/test/unit/com/mongodb/DBSpecification.groovy +++ b/driver-legacy/src/test/unit/com/mongodb/DBSpecification.groovy @@ -88,7 +88,7 @@ class DBSpecification extends Specification { then: def operation = executor.getWriteOperation() as CreateCollectionOperation - expect operation, isTheSameAs(new CreateCollectionOperation('test', 'ctest', db.getWriteConcern())) + expect operation, isTheSameAs(new CreateCollectionOperation('test', 'ctest', db.getWriteConcern(), true, null)) executor.getReadConcern() == ReadConcern.MAJORITY when: @@ -108,7 +108,7 @@ class DBSpecification extends Specification { operation = executor.getWriteOperation() as CreateCollectionOperation then: - expect operation, isTheSameAs(new CreateCollectionOperation('test', 'ctest', db.getWriteConcern()) + expect operation, isTheSameAs(new CreateCollectionOperation('test', 'ctest', db.getWriteConcern(), true, null) .sizeInBytes(100000) .maxDocuments(2000) .capped(true) @@ -136,7 +136,7 @@ class DBSpecification extends Specification { operation = executor.getWriteOperation() as CreateCollectionOperation then: - expect operation, isTheSameAs(new CreateCollectionOperation('test', 'ctest', db.getWriteConcern()) + expect operation, isTheSameAs(new CreateCollectionOperation('test', 'ctest', db.getWriteConcern(), true, null) .collation(collation)) executor.getReadConcern() == ReadConcern.MAJORITY } @@ -167,7 +167,7 @@ class DBSpecification extends Specification { then: def operation = executor.getWriteOperation() as CreateViewOperation expect operation, isTheSameAs(new CreateViewOperation(databaseName, viewName, viewOn, - [new BsonDocument('$match', new BsonDocument('x', BsonBoolean.TRUE))], writeConcern)) + [new BsonDocument('$match', new BsonDocument('x', BsonBoolean.TRUE))], writeConcern, true, null)) executor.getReadConcern() == ReadConcern.MAJORITY when: @@ -176,7 +176,7 @@ class DBSpecification extends Specification { then: expect operation, isTheSameAs(new CreateViewOperation(databaseName, viewName, viewOn, - [new BsonDocument('$match', new BsonDocument('x', BsonBoolean.TRUE))], writeConcern).collation(collation)) + [new BsonDocument('$match', new BsonDocument('x', BsonBoolean.TRUE))], writeConcern, true, null).collation(collation)) executor.getReadConcern() == ReadConcern.MAJORITY } diff --git a/driver-sync/src/main/com/mongodb/client/internal/AggregateIterableImpl.java b/driver-sync/src/main/com/mongodb/client/internal/AggregateIterableImpl.java index e69ca0d78c..841d253231 100644 --- a/driver-sync/src/main/com/mongodb/client/internal/AggregateIterableImpl.java +++ b/driver-sync/src/main/com/mongodb/client/internal/AggregateIterableImpl.java @@ -68,10 +68,11 @@ class AggregateIterableImpl extends MongoIterableImpl resultClass, final CodecRegistry codecRegistry, final ReadPreference readPreference, final ReadConcern readConcern, final WriteConcern writeConcern, final OperationExecutor executor, final List pipeline, final AggregationLevel aggregationLevel, - final boolean retryReads, @Nullable final Integer maxAdaptiveRetriesSetting, + final boolean retryWrites, final boolean retryReads, @Nullable final Integer maxAdaptiveRetriesSetting, final TimeoutSettings timeoutSettings) { this(clientSession, new MongoNamespace(databaseName, "_ignored"), documentClass, resultClass, codecRegistry, readPreference, - readConcern, writeConcern, executor, pipeline, aggregationLevel, retryReads, maxAdaptiveRetriesSetting, timeoutSettings); + readConcern, writeConcern, executor, pipeline, aggregationLevel, retryWrites, retryReads, maxAdaptiveRetriesSetting, + timeoutSettings); } @SuppressWarnings("checkstyle:ParameterNumber") @@ -79,11 +80,11 @@ class AggregateIterableImpl extends MongoIterableImpl resultClass, final CodecRegistry codecRegistry, final ReadPreference readPreference, final ReadConcern readConcern, final WriteConcern writeConcern, final OperationExecutor executor, final List pipeline, final AggregationLevel aggregationLevel, - final boolean retryReads, @Nullable final Integer maxAdaptiveRetriesSetting, + final boolean retryWrites, final boolean retryReads, @Nullable final Integer maxAdaptiveRetriesSetting, final TimeoutSettings timeoutSettings) { super(clientSession, executor, readConcern, readPreference, retryReads, timeoutSettings); this.operations = new Operations<>(namespace, documentClass, readPreference, codecRegistry, readConcern, writeConcern, - true, retryReads, maxAdaptiveRetriesSetting, timeoutSettings); + retryWrites, retryReads, maxAdaptiveRetriesSetting, timeoutSettings); this.namespace = notNull("namespace", namespace); this.documentClass = notNull("documentClass", documentClass); this.resultClass = notNull("resultClass", resultClass); diff --git a/driver-sync/src/main/com/mongodb/client/internal/MongoCollectionImpl.java b/driver-sync/src/main/com/mongodb/client/internal/MongoCollectionImpl.java index 4619289bc2..2641616a2d 100755 --- a/driver-sync/src/main/com/mongodb/client/internal/MongoCollectionImpl.java +++ b/driver-sync/src/main/com/mongodb/client/internal/MongoCollectionImpl.java @@ -364,7 +364,7 @@ private AggregateIterable createAggregateIterable(@Nullable f final Class resultClass) { return new AggregateIterableImpl<>(clientSession, namespace, documentClass, resultClass, codecRegistry, readPreference, readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, - retryReads, maxAdaptiveRetriesSetting, timeoutSettings); + retryWrites, retryReads, maxAdaptiveRetriesSetting, timeoutSettings); } @Override diff --git a/driver-sync/src/main/com/mongodb/client/internal/MongoDatabaseImpl.java b/driver-sync/src/main/com/mongodb/client/internal/MongoDatabaseImpl.java index 141564bc6e..1f8a2cdedb 100644 --- a/driver-sync/src/main/com/mongodb/client/internal/MongoDatabaseImpl.java +++ b/driver-sync/src/main/com/mongodb/client/internal/MongoDatabaseImpl.java @@ -399,7 +399,7 @@ private AggregateIterable createAggregateIterable(@Nullable f final Class resultClass) { return new AggregateIterableImpl<>(clientSession, name, Document.class, resultClass, codecRegistry, readPreference, readConcern, writeConcern, executor, pipeline, AggregationLevel.DATABASE, - retryReads, maxAdaptiveRetriesSetting, timeoutSettings); + retryWrites, retryReads, maxAdaptiveRetriesSetting, timeoutSettings); } private ChangeStreamIterable createChangeStreamIterable(@Nullable final ClientSession clientSession, diff --git a/driver-sync/src/test/functional/com/mongodb/client/BackpressureProseTest.java b/driver-sync/src/test/functional/com/mongodb/client/BackpressureProseTest.java index 2f12788621..bb566ed2fe 100644 --- a/driver-sync/src/test/functional/com/mongodb/client/BackpressureProseTest.java +++ b/driver-sync/src/test/functional/com/mongodb/client/BackpressureProseTest.java @@ -20,31 +20,47 @@ import com.mongodb.MongoCommandException; import com.mongodb.MongoNamespace; import com.mongodb.MongoServerException; +import com.mongodb.client.model.CreateCollectionOptions; +import com.mongodb.client.model.DropCollectionOptions; import com.mongodb.client.model.Filters; +import com.mongodb.client.model.SearchIndexModel; import com.mongodb.client.model.Updates; import com.mongodb.client.model.bulk.ClientBulkWriteResult; import com.mongodb.client.model.bulk.ClientNamespacedWriteModel; import com.mongodb.event.CommandFailedEvent; +import com.mongodb.event.CommandStartedEvent; import com.mongodb.internal.connection.TestCommandListener; import com.mongodb.internal.event.ConfigureFailPointCommandListener; import com.mongodb.internal.time.ExponentialBackoff; import com.mongodb.internal.time.StartTime; import com.mongodb.lang.Nullable; import org.bson.BsonDocument; +import org.bson.BsonString; import org.bson.Document; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; import java.time.Duration; +import java.util.ArrayList; import java.util.List; import java.util.concurrent.ExecutionException; +import java.util.function.Consumer; +import java.util.stream.Collectors; +import java.util.stream.IntStream; +import java.util.stream.Stream; +import static com.mongodb.client.model.Aggregates.match; import static com.mongodb.client.model.bulk.ClientBulkWriteOptions.clientBulkWriteOptions; import static com.mongodb.client.model.bulk.ClientUpdateOneOptions.clientUpdateOneOptions; import static java.lang.String.join; import static java.util.Arrays.asList; import static java.util.Collections.nCopies; +import static java.util.Collections.singletonList; +import static com.mongodb.ClusterFixture.isStandalone; import static com.mongodb.ClusterFixture.serverVersionAtLeast; import static com.mongodb.MongoException.RETRYABLE_ERROR_LABEL; import static com.mongodb.MongoException.SYSTEM_OVERLOADED_ERROR_LABEL; @@ -58,6 +74,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assumptions.assumeFalse; import static org.junit.jupiter.api.Assumptions.assumeTrue; /** @@ -65,6 +82,10 @@ * Prose Tests. */ public class BackpressureProseTest { + private static final String ENCRYPTED_STATE_COLLECTION_PREFIX = "enxcol_."; + private static final int SYSTEM_OVERLOAD_ERROR_CODE = 462; + private static final int RETRYABLE_ERROR_CODE = 11602; + private static final MongoNamespace NAMESPACE = new MongoNamespace(getDefaultDatabaseName(), BackpressureProseTest.class.getSimpleName()); protected MongoClient createClient(final MongoClientSettings mongoClientSettings) { return MongoClients.create(mongoClientSettings); } @@ -259,29 +280,8 @@ void runCommandPropagatesRetryableWriteErrorAfterOverloadRetry() throws Interrup */ @Test void runCommandDoesNotRetryOnRetryableWriteError() throws InterruptedException { - assumeTrue(serverVersionAtLeast(4, 4)); - BsonDocument retryableWriteErrorFailPoint = BsonDocument.parse( - "{\n" - + " configureFailPoint: 'failCommand',\n" - + " mode: {times: 1},\n" - + " data: {\n" - + " failCommands: ['ping'],\n" - + " errorCode: 11602,\n" - + " errorLabels: ['" + RETRYABLE_WRITE_ERROR_LABEL + "']\n" - + " }\n" - + "}\n"); - TestCommandListener commandListener = new TestCommandListener(); - try (MongoClient client = createClient(MongoClientSettings.builder(getMongoClientSettings()) - .addCommandListener(commandListener) - .build()); - FailPoint ignored = FailPoint.enable(retryableWriteErrorFailPoint, getPrimary())) { - MongoServerException exception = assertThrows(MongoServerException.class, - () -> client.getDatabase("admin").runCommand(BsonDocument.parse("{ping: 1}"))); - assertTrue(exception.hasErrorLabel(RETRYABLE_WRITE_ERROR_LABEL), - "Expected RetryableWriteError, got: " + exception); - assertEquals(1, commandListener.getCommandStartedEvents("ping").size(), - "Expected exactly one ping attempt (runCommand overload-only policy does not retry RetryableWriteError)"); - } + assertCommandNotRetriedOnRetryableWriteError("ping", + client -> client.getDatabase("admin").runCommand(BsonDocument.parse("{ping: 1}"))); } /** @@ -333,28 +333,8 @@ void runCommandPropagatesRetryableReadErrorAfterOverloadRetry() throws Interrupt */ @Test void runCommandDoesNotRetryOnRetryableReadError() throws InterruptedException { - assumeTrue(serverVersionAtLeast(4, 4)); - BsonDocument retryableReadErrorFailPoint = BsonDocument.parse( - "{\n" - + " configureFailPoint: 'failCommand',\n" - + " mode: {times: 1},\n" - + " data: {\n" - + " failCommands: ['ping'],\n" - + " errorCode: 11602\n" - + " }\n" - + "}\n"); - TestCommandListener commandListener = new TestCommandListener(); - try (MongoClient client = createClient(MongoClientSettings.builder(getMongoClientSettings()) - .addCommandListener(commandListener) - .build()); - FailPoint ignored = FailPoint.enable(retryableReadErrorFailPoint, getPrimary())) { - MongoServerException exception = assertThrows(MongoServerException.class, - () -> client.getDatabase("admin").runCommand(BsonDocument.parse("{ping: 1}"))); - assertEquals(11602, ((MongoCommandException) exception).getErrorCode(), - "Expected retryable-read-style error, got: " + exception); - assertEquals(1, commandListener.getCommandStartedEvents("ping").size(), - "Expected exactly one ping attempt (runCommand overload-only policy does not retry retryable-read codes)"); - } + assertCommandNotRetriedOnRetryableReadError("ping", + client -> client.getDatabase("admin").runCommand(BsonDocument.parse("{ping: 1}"))); } /** @@ -457,31 +437,7 @@ void clientBulkWriteGetMoreDoesNotRetryNonOverloadError() throws InterruptedExce @Test void clientBulkWriteGetMoreDoesNotRetryOverloadWhenRetryReadsDisabled() throws InterruptedException { assumeTrue(serverVersionAtLeast(8, 0)); - BsonDocument overloadOnGetMoreOnce = BsonDocument.parse( - "{\n" - + " configureFailPoint: 'failCommand',\n" - + " mode: {times: 1},\n" - + " data: {\n" - + " failCommands: ['getMore'],\n" - + " errorCode: 462,\n" - + " errorLabels: ['" + SYSTEM_OVERLOADED_ERROR_LABEL + "', '" + RETRYABLE_ERROR_LABEL + "']\n" - + " }\n" - + "}\n"); - TestCommandListener commandListener = new TestCommandListener(); - try (MongoClient client = createClient(MongoClientSettings.builder(getMongoClientSettings()) - .retryWrites(false) - .retryReads(false) - .addCommandListener(commandListener) - .build())) { - try (FailPoint ignored = FailPoint.enable(overloadOnGetMoreOnce, getPrimary())) { - MongoServerException exception = assertThrows(MongoServerException.class, - () -> executeClientBulkWrite(client)); - assertTrue(exception.hasErrorLabel(SYSTEM_OVERLOADED_ERROR_LABEL), - "Expected propagated overload error, got: " + exception); - } - assertEquals(1, commandListener.getCommandStartedEvents("getMore").size(), - "Expected exactly one getMore attempt (retryReads=false disables overload retry for getMore)"); - } + assertCommandNotRetriedWhenRetryReadsDisabled("getMore", BackpressureProseTest::executeClientBulkWrite); } private static ClientBulkWriteResult executeClientBulkWrite(final MongoClient client) { @@ -504,9 +460,427 @@ private static ClientBulkWriteResult executeClientBulkWrite(final MongoClient cl return client.bulkWrite(models, clientBulkWriteOptions().verboseResults(true)); } + @Test + void createViewExhaustsOverloadRetriesAndThrows() throws InterruptedException { + assertCommandExhaustsOverloadRetriesAndThrows("create", + client -> client.getDatabase(NAMESPACE.getDatabaseName()) + .createView(NAMESPACE.getCollectionName() + "View", NAMESPACE.getCollectionName(), + singletonList(match(Filters.empty())))); + } + + @Test + void dropCollectionExhaustsOverloadRetriesAndThrows() throws InterruptedException { + assertCommandExhaustsOverloadRetriesAndThrows("drop", client -> getCollection(client).drop()); + } + + @Test + void dropDatabaseExhaustsOverloadRetriesAndThrows() throws InterruptedException { + assertCommandExhaustsOverloadRetriesAndThrows("dropDatabase", + client -> client.getDatabase(NAMESPACE.getDatabaseName()).drop()); + } + + @Test + void renameCollectionExhaustsOverloadRetriesAndThrows() throws InterruptedException { + assertCommandExhaustsOverloadRetriesAndThrows("renameCollection", + client -> getCollection(client).renameCollection( + new MongoNamespace(NAMESPACE.getDatabaseName(), NAMESPACE.getCollectionName() + "Renamed"))); + } + + @Test + void createSearchIndexesExhaustsOverloadRetriesAndThrows() throws InterruptedException { + assumeTrue(serverVersionAtLeast(6, 0)); + assertCommandExhaustsOverloadRetriesAndThrows("createSearchIndexes", + client -> getCollection(client).createSearchIndexes( + singletonList(new SearchIndexModel(new Document("mappings", new Document("dynamic", true)))))); + } + + @Test + void updateSearchIndexExhaustsOverloadRetriesAndThrows() throws InterruptedException { + assumeTrue(serverVersionAtLeast(6, 0)); + assertCommandExhaustsOverloadRetriesAndThrows("updateSearchIndex", + client -> getCollection(client).updateSearchIndex("default", new Document("mappings", new Document("dynamic", true)))); + } + + @Test + void dropSearchIndexExhaustsOverloadRetriesAndThrows() throws InterruptedException { + assumeTrue(serverVersionAtLeast(6, 0)); + assertCommandExhaustsOverloadRetriesAndThrows("dropSearchIndex", + client -> getCollection(client).dropSearchIndex("default")); + } + + @Test + void createCollectionExhaustsOverloadRetriesAndThrows() throws InterruptedException { + assertCommandExhaustsOverloadRetriesAndThrows("create", + client -> client.getDatabase(NAMESPACE.getDatabaseName()).createCollection(NAMESPACE.getCollectionName())); + } + + + @Test + void createViewDoesNotRetryOverloadWhenRetryWritesDisabled() throws InterruptedException { + assertCommandNotRetriedWhenRetryWritesDisabled("create", + client -> client.getDatabase(NAMESPACE.getDatabaseName()) + .createView(NAMESPACE.getCollectionName() + "View", NAMESPACE.getCollectionName(), + singletonList(match(Filters.empty())))); + } + + @Test + void dropCollectionDoesNotRetryOverloadWhenRetryWritesDisabled() throws InterruptedException { + assertCommandNotRetriedWhenRetryWritesDisabled("drop", client -> getCollection(client).drop()); + } + + @Test + void dropDatabaseDoesNotRetryOverloadWhenRetryWritesDisabled() throws InterruptedException { + assertCommandNotRetriedWhenRetryWritesDisabled("dropDatabase", + client -> client.getDatabase(NAMESPACE.getDatabaseName()).drop()); + } + + @Test + void renameCollectionDoesNotRetryOverloadWhenRetryWritesDisabled() throws InterruptedException { + assertCommandNotRetriedWhenRetryWritesDisabled("renameCollection", + client -> getCollection(client).renameCollection( + new MongoNamespace(NAMESPACE.getDatabaseName(), NAMESPACE.getCollectionName() + "Renamed"))); + } + + @Test + void createCollectionDoesNotRetryOverloadWhenRetryWritesDisabled() throws InterruptedException { + assertCommandNotRetriedWhenRetryWritesDisabled("create", + client -> client.getDatabase(NAMESPACE.getDatabaseName()).createCollection(NAMESPACE.getCollectionName())); + } + + @Test + void createSearchIndexesDoesNotRetryOverloadWhenRetryWritesDisabled() throws InterruptedException { + assumeTrue(serverVersionAtLeast(6, 0)); + + assertCommandNotRetriedWhenRetryWritesDisabled("createSearchIndexes", + client -> getCollection(client).createSearchIndexes( + singletonList(new SearchIndexModel(new Document("mappings", new Document("dynamic", true)))))); + } + + @Test + void updateSearchIndexDoesNotRetryOverloadWhenRetryWritesDisabled() throws InterruptedException { + assumeTrue(serverVersionAtLeast(6, 0)); + + assertCommandNotRetriedWhenRetryWritesDisabled("updateSearchIndex", + client -> getCollection(client).updateSearchIndex("default", new Document("mappings", new Document("dynamic", true)))); + } + + @Test + void dropSearchIndexDoesNotRetryOverloadWhenRetryWritesDisabled() throws InterruptedException { + assumeTrue(serverVersionAtLeast(6, 0)); + + assertCommandNotRetriedWhenRetryWritesDisabled("dropSearchIndex", + client -> getCollection(client).dropSearchIndex("default")); + } + + @Test + void createViewDoesNotRetryOnRetryableWriteError() throws InterruptedException { + assertCommandNotRetriedOnRetryableWriteError("create", + client -> client.getDatabase(NAMESPACE.getDatabaseName()) + .createView(NAMESPACE.getCollectionName() + "View", NAMESPACE.getCollectionName(), + singletonList(match(Filters.empty())))); + } + + @Test + void dropCollectionDoesNotRetryOnRetryableWriteError() throws InterruptedException { + assertCommandNotRetriedOnRetryableWriteError("drop", client -> getCollection(client).drop()); + } + + @Test + void dropDatabaseDoesNotRetryOnRetryableWriteError() throws InterruptedException { + assertCommandNotRetriedOnRetryableWriteError("dropDatabase", + client -> client.getDatabase(NAMESPACE.getDatabaseName()).drop()); + } + + @Test + void renameCollectionDoesNotRetryOnRetryableWriteError() throws InterruptedException { + assertCommandNotRetriedOnRetryableWriteError("renameCollection", + client -> getCollection(client).renameCollection( + new MongoNamespace(NAMESPACE.getDatabaseName(), NAMESPACE.getCollectionName() + "Renamed"))); + } + + @Test + void createCollectionDoesNotRetryOnRetryableWriteError() throws InterruptedException { + assertCommandNotRetriedOnRetryableWriteError("create", + client -> client.getDatabase(NAMESPACE.getDatabaseName()).createCollection(NAMESPACE.getCollectionName())); + } + + @Test + void createSearchIndexesDoesNotRetryOnRetryableWriteError() throws InterruptedException { + assumeTrue(serverVersionAtLeast(6, 0)); + + assertCommandNotRetriedOnRetryableWriteError("createSearchIndexes", + client -> getCollection(client).createSearchIndexes( + singletonList(new SearchIndexModel(new Document("mappings", new Document("dynamic", true)))))); + } + + @Test + void updateSearchIndexDoesNotRetryOnRetryableWriteError() throws InterruptedException { + assumeTrue(serverVersionAtLeast(6, 0)); + + assertCommandNotRetriedOnRetryableWriteError("updateSearchIndex", + client -> getCollection(client).updateSearchIndex("default", new Document("mappings", new Document("dynamic", true)))); + } + + @Test + void dropSearchIndexDoesNotRetryOnRetryableWriteError() throws InterruptedException { + assumeTrue(serverVersionAtLeast(6, 0)); + // assumeTrue(hasAtlasSearchIndexHelperEnabled(), "Atlas Search Index tests are disabled"); + assertCommandNotRetriedOnRetryableWriteError("dropSearchIndex", + client -> getCollection(client).dropSearchIndex("default")); + } + + private static Stream createEncryptedCollectionRetriesEachCommandIndependently() { + String collectionName = NAMESPACE.getCollectionName(); + List commandSequence = asList( + new BsonDocument("create", new BsonString(ENCRYPTED_STATE_COLLECTION_PREFIX + collectionName + ".esc")), + new BsonDocument("create", new BsonString(ENCRYPTED_STATE_COLLECTION_PREFIX + collectionName + ".ecoc")), + new BsonDocument("create", new BsonString(collectionName)), + new BsonDocument("createIndexes", new BsonString(collectionName))); + // QE createCollection issues the command sequence above; we generate one variant per command in round-robin, + // where that command is the one expected to fail and exhaust its overload retries. + return IntStream.range(0, commandSequence.size()).mapToObj(failingCommandIndex -> { + BsonDocument failingCommand = commandSequence.get(failingCommandIndex); + List expectedCommands = new ArrayList<>(commandSequence.subList(0, failingCommandIndex)); + expectedCommands.addAll(nCopies(DEFAULT_MAX_ADAPTIVE_RETRIES + 1, failingCommand)); + return Arguments.of(failingCommand, failingCommandIndex, expectedCommands); + }); + } + + @ParameterizedTest(name = "createEncryptedCollectionRetriesEachCommandIndependently. failingCommand={0}, failPointSkip=={1}") + @MethodSource + void createEncryptedCollectionRetriesEachCommandIndependently( + final BsonDocument failingCommand, + final int failPointSkip, + final List expectedCommands) throws InterruptedException { + assumeTrue(serverVersionAtLeast(7, 0)); + assumeFalse(isStandalone(), "Encrypted collections are not supported on standalone"); + TestCommandListener commandListener = new TestCommandListener(); + // The failPoint fails every command of the sequence, so `skip` is the number of commands preceding the + // failing one. It lets them pass through and then fails every subsequent one, so that all the retries of a + // single command in the sequence are exhausted. + BsonDocument configureFailPoint = BsonDocument.parse( + "{\n" + + " configureFailPoint: 'failCommand',\n" + + " mode: {skip: " + failPointSkip + "},\n" + + " data: {\n" + + " failCommands: ['create', 'createIndexes'],\n" + + " errorCode: " + SYSTEM_OVERLOAD_ERROR_CODE + ",\n" + + " errorLabels: ['" + SYSTEM_OVERLOADED_ERROR_LABEL + "', '" + RETRYABLE_ERROR_LABEL + "']\n" + + " }\n" + + "}\n"); + try (MongoClient client = createClient(MongoClientSettings.builder(getMongoClientSettings()) + .addCommandListener(commandListener) + .build())) { + MongoDatabase database = client.getDatabase(NAMESPACE.getDatabaseName()); + try (FailPoint ignored = FailPoint.enable(configureFailPoint, getPrimary())) { + commandListener.reset(); + MongoServerException e = assertThrows(MongoServerException.class, () -> database.createCollection( + NAMESPACE.getCollectionName(), encryptedCollectionOptions())); + assertEquals(SYSTEM_OVERLOAD_ERROR_CODE, e.getCode()); + assertCommandsStarted(expectedCommands, commandListener); + } + } + } + + private static Stream dropEncryptedCollectionRetriesEachCommandIndependently() { + String collectionName = NAMESPACE.getCollectionName(); + List commandSequence = asList( + new BsonDocument("drop", new BsonString(ENCRYPTED_STATE_COLLECTION_PREFIX + collectionName + ".esc")), + new BsonDocument("drop", new BsonString(ENCRYPTED_STATE_COLLECTION_PREFIX + collectionName + ".ecoc")), + new BsonDocument("drop", new BsonString(collectionName))); + return IntStream.range(0, commandSequence.size()).mapToObj(failingCommandIndex -> { + BsonDocument failingCommand = commandSequence.get(failingCommandIndex); + List expectedCommands = new ArrayList<>(commandSequence.subList(0, failingCommandIndex)); + expectedCommands.addAll(nCopies(DEFAULT_MAX_ADAPTIVE_RETRIES + 1, failingCommand)); + return Arguments.of(failingCommand, failingCommandIndex, expectedCommands); + }); + } + + @ParameterizedTest(name = "dropEncryptedCollectionRetriesEachCommandIndependently. failingCommand={0}, failPointSkip=={1}") + @MethodSource + void dropEncryptedCollectionRetriesEachCommandIndependently( + final BsonDocument failingCommand, + final int failPointSkip, + final List expectedCommands) throws InterruptedException { + assumeTrue(serverVersionAtLeast(7, 0)); + assumeFalse(isStandalone(), "Encrypted collections are not supported on standalone"); + TestCommandListener commandListener = new TestCommandListener(); + // The failPoint fails every command of the sequence, so `skip` is the number of the commands preceding the + // failing one. It lets them pass through and then fails every subsequent one, so that all the retries of a + // single command in the sequence are exhausted. + BsonDocument configureFailPoint = BsonDocument.parse( + "{\n" + + " configureFailPoint: 'failCommand',\n" + + " mode: {skip: " + failPointSkip + "},\n" + + " data: {\n" + + " failCommands: ['drop'],\n" + + " errorCode: " + SYSTEM_OVERLOAD_ERROR_CODE + ",\n" + + " errorLabels: ['" + SYSTEM_OVERLOADED_ERROR_LABEL + "', '" + RETRYABLE_ERROR_LABEL + "']\n" + + " }\n" + + "}\n"); + try (MongoClient client = createClient(MongoClientSettings.builder(getMongoClientSettings()) + .addCommandListener(commandListener) + .build())) { + try (FailPoint ignored = FailPoint.enable(configureFailPoint, getPrimary())) { + commandListener.reset(); + MongoServerException e = assertThrows(MongoServerException.class, () -> getCollection(client).drop( + new DropCollectionOptions().encryptedFields(encryptedCollectionOptions().getEncryptedFields()))); + assertEquals(SYSTEM_OVERLOAD_ERROR_CODE, e.getCode()); + assertCommandsStarted(expectedCommands, commandListener); + } + } + } + + private void assertCommandExhaustsOverloadRetriesAndThrows(final String failingCommandName, final Consumer operation) + throws InterruptedException { + assumeTrue(serverVersionAtLeast(4, 4)); + TestCommandListener commandListener = new TestCommandListener(); + BsonDocument configureFailPoint = BsonDocument.parse( + "{\n" + + " configureFailPoint: 'failCommand',\n" + + " mode: 'alwaysOn',\n" + + " data: {\n" + + " failCommands: ['" + failingCommandName + "'],\n" + + " errorCode: " + SYSTEM_OVERLOAD_ERROR_CODE + ",\n" + + " errorLabels: ['" + SYSTEM_OVERLOADED_ERROR_LABEL + "', '" + RETRYABLE_ERROR_LABEL + "']\n" + + " }\n" + + "}\n"); + try (MongoClient client = createClient(MongoClientSettings.builder(getMongoClientSettings()) + .addCommandListener(commandListener) + .build())) { + try (FailPoint ignored = FailPoint.enable(configureFailPoint, getPrimary())) { + commandListener.reset(); + MongoServerException exception = assertThrows(MongoServerException.class, () -> operation.accept(client)); + assertEquals(SYSTEM_OVERLOAD_ERROR_CODE, exception.getCode()); + assertTrue(exception.hasErrorLabel(SYSTEM_OVERLOADED_ERROR_LABEL)); + assertTrue(exception.hasErrorLabel(RETRYABLE_ERROR_LABEL)); + assertEquals(DEFAULT_MAX_ADAPTIVE_RETRIES + 1, + commandListener.getCommandStartedEvents(failingCommandName).size(), + "Expected initial attempt plus " + DEFAULT_MAX_ADAPTIVE_RETRIES + " overload retries"); + } + } + } + + private void assertCommandNotRetriedOnRetryableWriteError(final String failingCommandName, final Consumer operation) + throws InterruptedException { + assertCommandNotRetriedOnNonOverloadError(failingCommandName, operation, RETRYABLE_WRITE_ERROR_LABEL); + } + + private void assertCommandNotRetriedOnRetryableReadError(final String failingCommandName, final Consumer operation) + throws InterruptedException { + assertCommandNotRetriedOnNonOverloadError(failingCommandName, operation, null); + } + + private void assertCommandNotRetriedOnNonOverloadError(final String failingCommandName, final Consumer operation, + @Nullable final String errorLabel) + throws InterruptedException { + assumeTrue(serverVersionAtLeast(4, 4)); + TestCommandListener commandListener = new TestCommandListener(); + BsonDocument configureFailPoint = BsonDocument.parse( + "{\n" + + " configureFailPoint: 'failCommand',\n" + + " mode: {times: 1},\n" + + " data: {\n" + + " failCommands: ['" + failingCommandName + "'],\n" + + " errorCode: " + RETRYABLE_ERROR_CODE + ",\n" + + " errorLabels: [" + (errorLabel == null ? "" : "'" + errorLabel + "'") + "]\n" + + " }\n" + + "}\n"); + try (MongoClient client = createClient(MongoClientSettings.builder(getMongoClientSettings()) + .addCommandListener(commandListener) + .build())) { + try (FailPoint ignored = FailPoint.enable(configureFailPoint, getPrimary())) { + commandListener.reset(); + MongoServerException exception = assertThrows(MongoServerException.class, () -> operation.accept(client)); + assertEquals(RETRYABLE_ERROR_CODE, ((MongoCommandException) exception).getErrorCode(), + format("Expected the propagated non-overload error, got: %s", exception)); + if (errorLabel != null) { + assertTrue(exception.hasErrorLabel(errorLabel), + format("Expected the propagated error to have the %s label, got: %s", errorLabel, exception)); + } + assertEquals(1, commandListener.getCommandStartedEvents(failingCommandName).size(), + format("Expected exactly one attempt of %s, as the overload-only policy does not retry" + + " non-overload errors", failingCommandName)); + } + } + } + + private void assertCommandNotRetriedWhenRetryWritesDisabled(final String failingCommandName, final Consumer operation) + throws InterruptedException { + assertCommandNotOverloadRetried(failingCommandName, operation, true, false); + } + + private void assertCommandNotRetriedWhenRetryReadsDisabled(final String failingCommandName, final Consumer operation) + throws InterruptedException { + assertCommandNotOverloadRetried(failingCommandName, operation, false, true); + } + + private void assertCommandNotOverloadRetried(final String failingCommandName, final Consumer operation, + final boolean retryReads, final boolean retryWrites) + throws InterruptedException { + assumeTrue(serverVersionAtLeast(4, 4)); + TestCommandListener commandListener = new TestCommandListener(); + BsonDocument configureFailPoint = BsonDocument.parse( + "{\n" + + " configureFailPoint: 'failCommand',\n" + + " mode: 'alwaysOn',\n" + + " data: {\n" + + " failCommands: ['" + failingCommandName + "'],\n" + + " errorCode: " + SYSTEM_OVERLOAD_ERROR_CODE + ",\n" + + " errorLabels: ['" + SYSTEM_OVERLOADED_ERROR_LABEL + "', '" + RETRYABLE_ERROR_LABEL + "']\n" + + " }\n" + + "}\n"); + try (MongoClient client = createClient(MongoClientSettings.builder(getMongoClientSettings()) + .retryReads(retryReads) + .retryWrites(retryWrites) + .addCommandListener(commandListener) + .build())) { + try (FailPoint ignored = FailPoint.enable(configureFailPoint, getPrimary())) { + commandListener.reset(); + MongoServerException exception = assertThrows(MongoServerException.class, () -> operation.accept(client)); + assertEquals(SYSTEM_OVERLOAD_ERROR_CODE, exception.getCode()); + assertTrue(exception.hasErrorLabel(SYSTEM_OVERLOADED_ERROR_LABEL), + "Expected propagated overload error, got: " + exception); + assertTrue(exception.hasErrorLabel(RETRYABLE_ERROR_LABEL)); + assertEquals(1, commandListener.getCommandStartedEvents(failingCommandName).size(), + format("Expected exactly one attempt of %s, as retryReads=%b, retryWrites=%b disable the" + + " overload retry", failingCommandName, retryReads, retryWrites)); + } + } + } + private static MongoCollection dropAndGetCollection(final String name, final MongoClient client) { MongoCollection result = client.getDatabase(getDefaultDatabaseName()).getCollection(name); result.drop(); return result; } + + /** + * Asserts that the commands started by the {@code commandListener} are exactly the {@code expectedCommands}, in + * order. Each expected command is required to be a subset of the actual one, so that it has to specify only the + * entries identifying the command. + */ + private static void assertCommandsStarted(final List expectedCommands, + final TestCommandListener commandListener) { + List actualCommands = commandListener.getCommandStartedEvents().stream() + .map(CommandStartedEvent::getCommand) + .collect(Collectors.toList()); + assertEquals(expectedCommands.size(), actualCommands.size(), + format("Expected %s but observed %s", expectedCommands, actualCommands)); + for (int i = 0; i < expectedCommands.size(); i++) { + BsonDocument expected = expectedCommands.get(i); + BsonDocument actual = actualCommands.get(i); + assertTrue(actual.entrySet().containsAll(expected.entrySet()), + format("Expected the command at index %d to contain %s but it was %s", i, expected, actual)); + } + } + + private static CreateCollectionOptions encryptedCollectionOptions() { + return new CreateCollectionOptions().encryptedFields(BsonDocument.parse( + "{fields: [{path: 'ssn', bsonType: 'string'," + + " keyId: {$binary: {base64: 'AAAAAAAAAAAAAAAAAAAAAA==', subType: '04'}}}]}")); + } + private static MongoCollection getCollection(final MongoClient client) { + return client.getDatabase(NAMESPACE.getDatabaseName()).getCollection(NAMESPACE.getCollectionName()); + } } diff --git a/driver-sync/src/test/functional/com/mongodb/client/unified/UnifiedTestModifications.java b/driver-sync/src/test/functional/com/mongodb/client/unified/UnifiedTestModifications.java index 32d0a6f811..6a267ed3c2 100644 --- a/driver-sync/src/test/functional/com/mongodb/client/unified/UnifiedTestModifications.java +++ b/driver-sync/src/test/functional/com/mongodb/client/unified/UnifiedTestModifications.java @@ -576,7 +576,7 @@ public static void applyCustomizations(final TestDef def) { // backpressure - def.modify(WAIT_FOR_BATCH_CURSOR_CREATION) + def.modify(WAIT_FOR_BATCH_CURSOR_CREATION, IGNORE_EXTRA_EVENTS) .test("client-backpressure", "tests that operations retry at most maxAttempts=2 times", "client.createChangeStream retries at most maxAttempts=2 times") .test("client-backpressure", "tests that operations retry at most maxAttempts=2 times", @@ -596,24 +596,6 @@ public static void applyCustomizations(final TestDef def) { .test("client-backpressure", "tests that operations respect overload backoff retry loop", "collection.createChangeStream (read) does not retry if retryReads=false"); - // TODO-BACKPRESSURE enable the below tests when JAVA-5956 is done - def.skipJira("https://jira.mongodb.org/browse/JAVA-5956 TODO-JAVA-5956") - .test("client-backpressure", "tests that operations respect overload backoff retry loop", "collection.createIndex retries using operation loop"); - def.skipJira("https://jira.mongodb.org/browse/JAVA-5956 TODO-JAVA-5956") - .test("client-backpressure", "tests that operations respect overload backoff retry loop", "collection.dropIndex retries using operation loop"); - def.skipJira("https://jira.mongodb.org/browse/JAVA-5956 TODO-JAVA-5956") - .test("client-backpressure", "tests that operations respect overload backoff retry loop", "collection.dropIndexes retries using operation loop"); - def.skipJira("https://jira.mongodb.org/browse/JAVA-5956 TODO-JAVA-5956") - .test("client-backpressure", "tests that operations respect overload backoff retry loop", "collection.aggregate write retries using operation loop"); - def.skipJira("https://jira.mongodb.org/browse/JAVA-5956 TODO-JAVA-5956") - .test("client-backpressure", "tests that operations retry at most maxAttempts=2 times", "collection.createIndex retries at most maxAttempts=2 times"); - def.skipJira("https://jira.mongodb.org/browse/JAVA-5956 TODO-JAVA-5956") - .test("client-backpressure", "tests that operations retry at most maxAttempts=2 times", "collection.dropIndex retries at most maxAttempts=2 times"); - def.skipJira("https://jira.mongodb.org/browse/JAVA-5956 TODO-JAVA-5956") - .test("client-backpressure", "tests that operations retry at most maxAttempts=2 times", "collection.dropIndexes retries at most maxAttempts=2 times"); - def.skipJira("https://jira.mongodb.org/browse/JAVA-5956 TODO-JAVA-5956") - .test("client-backpressure", "tests that operations retry at most maxAttempts=2 times", "collection.aggregate write retries at most maxAttempts=2 times"); - // BatchCursorFlux fires closeCursor() then sink.error(e) without awaiting the killCursors reply, // so under reactive the test framework snapshots command events before killCursors succeeded lands. // Equivalent coverage is provided by the reactive BackpressureProseTest. diff --git a/driver-sync/src/test/unit/com/mongodb/client/internal/AggregateIterableSpecification.groovy b/driver-sync/src/test/unit/com/mongodb/client/internal/AggregateIterableSpecification.groovy index 01c94a882f..42ac081943 100644 --- a/driver-sync/src/test/unit/com/mongodb/client/internal/AggregateIterableSpecification.groovy +++ b/driver-sync/src/test/unit/com/mongodb/client/internal/AggregateIterableSpecification.groovy @@ -63,7 +63,7 @@ class AggregateIterableSpecification extends Specification { def pipeline = [new Document('$match', 1)] def aggregationIterable = new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, - true, null, TIMEOUT_SETTINGS) + false, true, null, TIMEOUT_SETTINGS) when: 'default input should be as expected' aggregationIterable.iterator() @@ -98,7 +98,7 @@ class AggregateIterableSpecification extends Specification { when: 'both hint and hint string are set' aggregationIterable = new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, - readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) aggregationIterable .hint(new Document('a', 1)) @@ -122,7 +122,7 @@ class AggregateIterableSpecification extends Specification { when: 'aggregation includes $out' new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) .batchSize(99) .allowDiskUse(true) .collation(collation) @@ -152,7 +152,7 @@ class AggregateIterableSpecification extends Specification { when: 'aggregation includes $out and is at the database level' new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - pipeline, AggregationLevel.DATABASE, false, null, TIMEOUT_SETTINGS) + pipeline, AggregationLevel.DATABASE, false, false, null, TIMEOUT_SETTINGS) .batchSize(99) .maxTime(100, MILLISECONDS) .allowDiskUse(true) @@ -185,7 +185,7 @@ class AggregateIterableSpecification extends Specification { when: 'toCollection should work as expected' new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) .allowDiskUse(true) .collation(collation) .hint(new Document('a', 1)) @@ -212,7 +212,7 @@ class AggregateIterableSpecification extends Specification { when: 'aggregation includes $out and hint string' new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) .hintString('x_1').iterator() def operation = executor.getReadOperation() as AggregateToCollectionOperation @@ -226,7 +226,7 @@ class AggregateIterableSpecification extends Specification { when: 'aggregation includes $out and hint and hint string' executor = new TestOperationExecutor([null, null, null, null, null]) new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) .hint(new BsonDocument('x', new BsonInt32(1))) .hintString('x_1').iterator() @@ -250,7 +250,7 @@ class AggregateIterableSpecification extends Specification { when: 'aggregation includes $merge' new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) .batchSize(99) .allowDiskUse(true) .collation(collation) @@ -281,7 +281,7 @@ class AggregateIterableSpecification extends Specification { when: 'aggregation includes $merge into a different database' new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - pipelineWithIntoDocument, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + pipelineWithIntoDocument, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) .batchSize(99) .maxTime(100, MILLISECONDS) .allowDiskUse(true) @@ -314,7 +314,7 @@ class AggregateIterableSpecification extends Specification { when: 'aggregation includes $merge and is at the database level' new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - pipeline, AggregationLevel.DATABASE, false, null, TIMEOUT_SETTINGS) + pipeline, AggregationLevel.DATABASE, false, false, null, TIMEOUT_SETTINGS) .batchSize(99) .maxTime(100, MILLISECONDS) .allowDiskUse(true) @@ -346,7 +346,7 @@ class AggregateIterableSpecification extends Specification { when: 'toCollection should work as expected' new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) .allowDiskUse(true) .collation(collation) .hint(new Document('a', 1)) @@ -375,7 +375,7 @@ class AggregateIterableSpecification extends Specification { when: new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) .iterator() def operation = executor.getReadOperation() as AggregateToCollectionOperation @@ -419,7 +419,7 @@ class AggregateIterableSpecification extends Specification { when: 'aggregation includes $out' def aggregateIterable = new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, - readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) aggregateIterable.toCollection() def operation = executor.getReadOperation() as AggregateToCollectionOperation @@ -438,7 +438,7 @@ class AggregateIterableSpecification extends Specification { when: 'aggregation includes $out and is at the database level' aggregateIterable = new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, - readConcern, writeConcern, executor, pipeline, AggregationLevel.DATABASE, false, null, TIMEOUT_SETTINGS) + readConcern, writeConcern, executor, pipeline, AggregationLevel.DATABASE, false, false, null, TIMEOUT_SETTINGS) aggregateIterable.toCollection() operation = executor.getReadOperation() as AggregateToCollectionOperation @@ -457,7 +457,7 @@ class AggregateIterableSpecification extends Specification { when: 'toCollection should work as expected' aggregateIterable = new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, - readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) aggregateIterable.toCollection() operation = executor.getReadOperation() as AggregateToCollectionOperation @@ -475,7 +475,7 @@ class AggregateIterableSpecification extends Specification { when: 'aggregation includes $out with namespace' aggregateIterable = new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, - readConcern, writeConcern, executor, outWithDBpipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + readConcern, writeConcern, executor, outWithDBpipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) aggregateIterable.toCollection() operation = executor.getReadOperation() as AggregateToCollectionOperation @@ -502,7 +502,7 @@ class AggregateIterableSpecification extends Specification { def executor = new TestOperationExecutor([batchCursor, batchCursor]) def pipeline = [new Document('$match', 1)] def aggregationIterable = new AggregateIterableImpl(clientSession, namespace, Document, Document, codecRegistry, readPreference, - readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) when: aggregationIterable.first() @@ -528,7 +528,7 @@ class AggregateIterableSpecification extends Specification { def executor = new TestOperationExecutor([null, batchCursor, null, batchCursor, null]) def pipeline = [new Document('$match', 1), new Document('$out', 'collName')] def aggregationIterable = new AggregateIterableImpl(clientSession, namespace, Document, Document, codecRegistry, readPreference, - readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) when: aggregationIterable.first() @@ -559,7 +559,7 @@ class AggregateIterableSpecification extends Specification { def executor = new TestOperationExecutor([new MongoException('failure')]) def pipeline = [new BsonDocument('$match', new BsonInt32(1))] def aggregationIterable = new AggregateIterableImpl(null, namespace, BsonDocument, BsonDocument, codecRegistry, readPreference, - readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS) + readConcern, writeConcern, executor, pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) when: 'The operation fails with an exception' aggregationIterable.iterator() @@ -575,14 +575,14 @@ class AggregateIterableSpecification extends Specification { when: 'a codec is missing' new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - pipeline, AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS).iterator() + pipeline, AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS).iterator() then: thrown(CodecConfigurationException) when: 'pipeline contains null' new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, - [null], AggregationLevel.COLLECTION, false, null, TIMEOUT_SETTINGS).iterator() + [null], AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS).iterator() then: thrown(IllegalArgumentException) @@ -611,7 +611,7 @@ class AggregateIterableSpecification extends Specification { } def executor = new TestOperationExecutor([cursor(), cursor(), cursor(), cursor()]) def mongoIterable = new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, - readConcern, writeConcern, executor, [new Document('$match', 1)], AggregationLevel.COLLECTION, false, null, + readConcern, writeConcern, executor, [new Document('$match', 1)], AggregationLevel.COLLECTION, false, false, null, TIMEOUT_SETTINGS) when: @@ -657,7 +657,7 @@ class AggregateIterableSpecification extends Specification { def batchSize = 5 def mongoIterable = new AggregateIterableImpl(null, namespace, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, Stub(OperationExecutor), [new Document('$match', 1)], AggregationLevel.COLLECTION, - false, null, TIMEOUT_SETTINGS) + false, false, null, TIMEOUT_SETTINGS) then: mongoIterable.getBatchSize() == null diff --git a/driver-sync/src/test/unit/com/mongodb/client/internal/MongoCollectionSpecification.groovy b/driver-sync/src/test/unit/com/mongodb/client/internal/MongoCollectionSpecification.groovy index fa7c868509..69aedcd479 100644 --- a/driver-sync/src/test/unit/com/mongodb/client/internal/MongoCollectionSpecification.groovy +++ b/driver-sync/src/test/unit/com/mongodb/client/internal/MongoCollectionSpecification.groovy @@ -376,7 +376,7 @@ class MongoCollectionSpecification extends Specification { then: expect aggregateIterable, isTheSameAs(new AggregateIterableImpl<>(session, namespace, Document, Document, codecRegistry, readPreference, readConcern, ACKNOWLEDGED, executor, [new Document('$match', 1)], - AggregationLevel.COLLECTION, true, null, TIMEOUT_SETTINGS)) + AggregationLevel.COLLECTION, true, true, null, TIMEOUT_SETTINGS)) when: aggregateIterable = execute(aggregateMethod, session, [new Document('$match', 1)], BsonDocument) @@ -384,7 +384,7 @@ class MongoCollectionSpecification extends Specification { then: expect aggregateIterable, isTheSameAs(new AggregateIterableImpl<>(session, namespace, Document, BsonDocument, codecRegistry, readPreference, readConcern, ACKNOWLEDGED, executor, [new Document('$match', 1)], - AggregationLevel.COLLECTION, true, null, TIMEOUT_SETTINGS)) + AggregationLevel.COLLECTION, true, true, null, TIMEOUT_SETTINGS)) where: session << [null, Stub(ClientSession)] @@ -1149,7 +1149,7 @@ class MongoCollectionSpecification extends Specification { def executor = new TestOperationExecutor([null]) def collection = new MongoCollectionImpl(namespace, Document, codecRegistry, readPreference, ACKNOWLEDGED, true, true, null, readConcern, JAVA_LEGACY, null, TIMEOUT_SETTINGS, executor) - def expectedOperation = new DropCollectionOperation(namespace, ACKNOWLEDGED) + def expectedOperation = new DropCollectionOperation(namespace, ACKNOWLEDGED, true, null) def dropMethod = collection.&drop when: @@ -1174,7 +1174,7 @@ class MongoCollectionSpecification extends Specification { when: def expectedOperation = new CreateIndexesOperation(namespace, - [new IndexRequest(new BsonDocument('key', new BsonInt32(1)))], ACKNOWLEDGED) + [new IndexRequest(new BsonDocument('key', new BsonInt32(1)))], ACKNOWLEDGED, true, null) def indexName = execute(createIndexMethod, session, new Document('key', 1)) def operation = executor.getWriteOperation() as CreateIndexesOperation @@ -1185,7 +1185,7 @@ class MongoCollectionSpecification extends Specification { when: expectedOperation = new CreateIndexesOperation(namespace, [new IndexRequest(new BsonDocument('key', new BsonInt32(1))), - new IndexRequest(new BsonDocument('key1', new BsonInt32(1)))], ACKNOWLEDGED) + new IndexRequest(new BsonDocument('key1', new BsonInt32(1)))], ACKNOWLEDGED, true, null) def indexNames = execute(createIndexesMethod, session, [new IndexModel(new Document('key', 1)), new IndexModel(new Document('key1', 1))]) operation = executor.getWriteOperation() as CreateIndexesOperation @@ -1198,7 +1198,7 @@ class MongoCollectionSpecification extends Specification { when: expectedOperation = new CreateIndexesOperation(namespace, [new IndexRequest(new BsonDocument('key', new BsonInt32(1))), - new IndexRequest(new BsonDocument('key1', new BsonInt32(1)))], ACKNOWLEDGED) + new IndexRequest(new BsonDocument('key1', new BsonInt32(1)))], ACKNOWLEDGED, true, null) indexNames = execute(createIndexesMethod, session, [new IndexModel(new Document('key', 1)), new IndexModel(new Document('key1', 1))], new CreateIndexOptions().maxTime(100, MILLISECONDS)) @@ -1212,7 +1212,7 @@ class MongoCollectionSpecification extends Specification { when: expectedOperation = new CreateIndexesOperation(namespace, [new IndexRequest(new BsonDocument('key', new BsonInt32(1))), - new IndexRequest(new BsonDocument('key1', new BsonInt32(1)))], ACKNOWLEDGED) + new IndexRequest(new BsonDocument('key1', new BsonInt32(1)))], ACKNOWLEDGED, true, null) .commitQuorum(CreateIndexCommitQuorum.VOTING_MEMBERS) indexNames = execute(createIndexesMethod, session, [new IndexModel(new Document('key', 1)), new IndexModel(new Document('key1', 1))], @@ -1246,7 +1246,7 @@ class MongoCollectionSpecification extends Specification { .collation(collation) .wildcardProjection(new BsonDocument('a', new BsonInt32(1))) .hidden(true) - ], ACKNOWLEDGED) + ], ACKNOWLEDGED, true, null) indexName = execute(createIndexMethod, session, new Document('key', 1), new IndexOptions() .background(true) .unique(true) @@ -1342,7 +1342,7 @@ class MongoCollectionSpecification extends Specification { def dropIndexMethod = collection.&dropIndex when: - def expectedOperation = new DropIndexOperation(namespace, 'indexName', ACKNOWLEDGED) + def expectedOperation = new DropIndexOperation(namespace, 'indexName', ACKNOWLEDGED, true, null) execute(dropIndexMethod, session, 'indexName') def operation = executor.getWriteOperation() as DropIndexOperation @@ -1352,7 +1352,7 @@ class MongoCollectionSpecification extends Specification { when: def keys = new BsonDocument('x', new BsonInt32(1)) - expectedOperation = new DropIndexOperation(namespace, keys, ACKNOWLEDGED) + expectedOperation = new DropIndexOperation(namespace, keys, ACKNOWLEDGED, true, null) execute(dropIndexMethod, session, keys) operation = executor.getWriteOperation() as DropIndexOperation @@ -1361,7 +1361,7 @@ class MongoCollectionSpecification extends Specification { executor.getClientSession() == session when: - expectedOperation = new DropIndexOperation(namespace, keys, ACKNOWLEDGED) + expectedOperation = new DropIndexOperation(namespace, keys, ACKNOWLEDGED, true, null) execute(dropIndexMethod, session, keys, new DropIndexOptions().maxTime(100, MILLISECONDS)) operation = executor.getWriteOperation() as DropIndexOperation @@ -1378,7 +1378,7 @@ class MongoCollectionSpecification extends Specification { def executor = new TestOperationExecutor([null, null]) def collection = new MongoCollectionImpl(namespace, Document, codecRegistry, readPreference, ACKNOWLEDGED, true, true, null, readConcern, JAVA_LEGACY, null, TIMEOUT_SETTINGS, executor) - def expectedOperation = new DropIndexOperation(namespace, '*', ACKNOWLEDGED) + def expectedOperation = new DropIndexOperation(namespace, '*', ACKNOWLEDGED, true, null) def dropIndexesMethod = collection.&dropIndexes when: @@ -1390,7 +1390,7 @@ class MongoCollectionSpecification extends Specification { executor.getClientSession() == session when: - expectedOperation = new DropIndexOperation(namespace, '*', ACKNOWLEDGED) + expectedOperation = new DropIndexOperation(namespace, '*', ACKNOWLEDGED, true, null) execute(dropIndexesMethod, session, new DropIndexOptions().maxTime(100, MILLISECONDS)) operation = executor.getWriteOperation() as DropIndexOperation @@ -1409,7 +1409,7 @@ class MongoCollectionSpecification extends Specification { true, true, null, readConcern, JAVA_LEGACY, null, TIMEOUT_SETTINGS, executor) def newNamespace = new MongoNamespace(namespace.getDatabaseName(), 'newName') def renameCollectionOptions = new RenameCollectionOptions().dropTarget(dropTarget) - def expectedOperation = new RenameCollectionOperation(namespace, newNamespace, ACKNOWLEDGED) + def expectedOperation = new RenameCollectionOperation(namespace, newNamespace, ACKNOWLEDGED, true, null) def renameCollection = collection.&renameCollection when: diff --git a/driver-sync/src/test/unit/com/mongodb/client/internal/MongoDatabaseSpecification.groovy b/driver-sync/src/test/unit/com/mongodb/client/internal/MongoDatabaseSpecification.groovy index d606cdfd04..4e8ac62d8c 100644 --- a/driver-sync/src/test/unit/com/mongodb/client/internal/MongoDatabaseSpecification.groovy +++ b/driver-sync/src/test/unit/com/mongodb/client/internal/MongoDatabaseSpecification.groovy @@ -448,7 +448,7 @@ class MongoDatabaseSpecification extends Specification { then: expect aggregateIterable, isTheSameAs(new AggregateIterableImpl<>(session, name, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, [], AggregationLevel.DATABASE, - false, null, TIMEOUT_SETTINGS), ['codec']) + false, false, null, TIMEOUT_SETTINGS), ['codec']) when: aggregateIterable = execute(aggregateMethod, session, [new Document('$match', 1)]) @@ -456,7 +456,7 @@ class MongoDatabaseSpecification extends Specification { then: expect aggregateIterable, isTheSameAs(new AggregateIterableImpl<>(session, name, Document, Document, codecRegistry, readPreference, readConcern, writeConcern, executor, [new Document('$match', 1)], - AggregationLevel.DATABASE, false, null, TIMEOUT_SETTINGS), ['codec']) + AggregationLevel.DATABASE, false, false, null, TIMEOUT_SETTINGS), ['codec']) when: aggregateIterable = execute(aggregateMethod, session, [new Document('$match', 1)], BsonDocument) @@ -464,7 +464,7 @@ class MongoDatabaseSpecification extends Specification { then: expect aggregateIterable, isTheSameAs(new AggregateIterableImpl<>(session, name, Document, BsonDocument, codecRegistry, readPreference, readConcern, writeConcern, executor, [new Document('$match', 1)], - AggregationLevel.DATABASE, false, null, TIMEOUT_SETTINGS), ['codec']) + AggregationLevel.DATABASE, false, false, null, TIMEOUT_SETTINGS), ['codec']) where: session << [null, Stub(ClientSession)]