Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions config/spotbugs/exclude.xml
Original file line number Diff line number Diff line change
Expand Up @@ -288,6 +288,16 @@
</Match>

<!-- Void method returning null but @NotNull API -->
<Match>
<Class name="com.mongodb.internal.operation.CreateCollectionOperation"/>
<Method name="execute"/>
<Bug pattern="NP_NONNULL_RETURN_VIOLATION"/>
</Match>
<Match>
<Class name="com.mongodb.internal.operation.DropCollectionOperation"/>
<Method name="execute"/>
<Bug pattern="NP_NONNULL_RETURN_VIOLATION"/>
</Match>
<Match>
<Class name="com.mongodb.internal.operation.DropIndexOperation"/>
<Method name="execute"/>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -40,39 +48,64 @@
*/
abstract class AbstractWriteSearchIndexOperation implements WriteOperation<Void> {
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<SpecRetryPolicy> retryControl = createSpecRetryControl(
overloadForWrite(retryWrites, maxAdaptiveRetriesSetting),
operationContext);
Supplier<Void> 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<Void> 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<SpecRetryPolicy> retryControl = createSpecRetryControl(
overloadForWrite(retryWrites, maxAdaptiveRetriesSetting),
operationContext);
AsyncCallbackSupplier<Void> 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);
}

/**
Expand Down Expand Up @@ -101,4 +134,5 @@ <E extends Throwable> void swallowOrThrow(@Nullable final E mongoExecutionExcept
public MongoNamespace getNamespace() {
return namespace;
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -64,6 +65,9 @@ public class AggregateToCollectionOperation implements ReadOperationSimple<Void>
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;
Expand All @@ -79,11 +83,19 @@ public AggregateToCollectionOperation(final MongoNamespace namespace, final List

public AggregateToCollectionOperation(final MongoNamespace namespace, final List<BsonDocument> pipeline,
@Nullable final ReadConcern readConcern, @Nullable final WriteConcern writeConcern, final AggregationLevel aggregationLevel) {
this(namespace, pipeline, readConcern, writeConcern, aggregationLevel, false, null);
Comment thread
vbabanin marked this conversation as resolved.
}

public AggregateToCollectionOperation(final MongoNamespace namespace, final List<BsonDocument> 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());
}
Expand Down Expand Up @@ -178,8 +190,7 @@ public Void execute(final ReadBinding binding, final OperationContext operationC
getCommandCreator(),
new BsonDocumentCodec(),
transformer(),
false,
null);
overloadForWrite(retryWrites, maxAdaptiveRetriesSetting));
}

@Override
Expand All @@ -194,8 +205,7 @@ public void executeAsync(final AsyncReadBinding binding, final OperationContext
getCommandCreator(),
new BsonDocumentCodec(),
asyncTransformer(),
false,
null,
overloadForWrite(retryWrites, maxAdaptiveRetriesSetting),
callback);
}

Expand Down Expand Up @@ -251,4 +261,5 @@ private static CommandReadTransformerAsync<BsonDocument, Void> asyncTransformer(
return null;
};
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -191,9 +192,8 @@ private void getMoreLoop(final ServerCursor localServerCursor,
}

private void getMore(final ServerCursor cursor, final OperationContext operationContext, final SingleResultCallback<List<T>> callback) {
SpecRetryPolicy.IndividualPolicies policies = new SpecRetryPolicy.IndividualPolicies(retryReads)
.includeOverload(maxAdaptiveRetriesSetting, SpecRetryPolicy.ErrorPropagation.AS_READ_POLICY);
RetryControl<SpecRetryPolicy> retryControl = createSpecRetryControl(policies, operationContext);
RetryControl<SpecRetryPolicy> retryControl = createSpecRetryControl(
overloadForRead(retryReads, maxAdaptiveRetriesSetting), operationContext);
AsyncCallbackSupplier<List<T>> retryingCommandExecutor = decorateWithRetriesAsync(retryControl, operationContext, attemptCallback ->
resourceManager.executeWithConnection(operationContext, (connection, wrappedCallback) ->
executeGetMoreCommand(assertNotNull(connection), cursor, operationContext, retryControl, wrappedCallback),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<SpecRetryPolicy> retryControl = createSpecRetryControl(policies, operationContext);
RetryControl<SpecRetryPolicy> retryControl = createSpecRetryControl(
overloadForRead(retryReads, maxAdaptiveRetriesSetting), operationContext);
Supplier<Void> retryingCommandExecutor = decorateWithRetries(retryControl, operationContext, () -> {
resourceManager.executeWithConnection(connection -> {
ServerCursor nextServerCursor;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

/**
Expand Down Expand Up @@ -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);
}
}
Loading