summaryrefslogtreecommitdiff
path: root/src/mongo/executor/task_executor_cursor.h
diff options
context:
space:
mode:
Diffstat (limited to 'src/mongo/executor/task_executor_cursor.h')
-rw-r--r--src/mongo/executor/task_executor_cursor.h108
1 files changed, 28 insertions, 80 deletions
diff --git a/src/mongo/executor/task_executor_cursor.h b/src/mongo/executor/task_executor_cursor.h
index 6de3a7664c1..b18ea481c88 100644
--- a/src/mongo/executor/task_executor_cursor.h
+++ b/src/mongo/executor/task_executor_cursor.h
@@ -30,7 +30,6 @@
#pragma once
#include <boost/optional.hpp>
-#include <memory>
#include <vector>
#include "mongo/base/status_with.h"
@@ -41,9 +40,9 @@
#include "mongo/db/query/cursor_response.h"
#include "mongo/executor/remote_command_request.h"
#include "mongo/executor/task_executor.h"
-#include "mongo/executor/task_executor_cursor_parameters_gen.h"
#include "mongo/util/duration.h"
-#include "mongo/util/future.h"
+#include "mongo/util/net/hostandport.h"
+#include "mongo/util/producer_consumer_queue.h"
namespace mongo {
namespace executor {
@@ -54,8 +53,8 @@ namespace executor {
*
* The main differentiator for this type over DBClientCursor is the use of a task executor (which
* provides access to a different connection pool, as well as interruptibility) and the ability to
- * overlap getMores. This starts fetching the next batch as soon as the previous one is received
- * (rather than on a call to 'getNext()').
+ * overlap getMores. This starts fetching the next batch as soon as one is exhausted (rather than
+ * on a call to getNext()).
*/
class TaskExecutorCursor {
public:
@@ -69,44 +68,27 @@ public:
struct Options {
boost::optional<int64_t> batchSize;
- bool pinConnection{gPinTaskExecCursorConns.load()};
- // If true, we will fetch the next batch as soon as the current one is recieved.
- // If false, we will fetch the next batch when the current batch is exhausted and
- // 'getNext()' is invoked.
- bool preFetchNextBatch{true};
-
- // This function, if specified, may modify a getMore request to include additional
- // information.
- std::function<void(BSONObjBuilder& bob)> getMoreAugmentationWriter;
-
- Options() {}
};
/**
- * Construct the cursor with a RemoteCommandRequest wrapping the initial command.
- *
- * Doesn't retry the command if we fail to establish the cursor. To create a TaskExecutorCursor
- * with the option to retry the initial command, see `makeTaskExecutorCursor`below.
+ * Construct the cursor with a RemoteCommandRequest wrapping the initial command
*
* One value is carried over in successive calls to getMore/killCursor:
*
* opCtx - The Logical Session Id from the initial command is carried over in all later stages.
* NOTE - the actual command must not include the lsid
*/
- TaskExecutorCursor(std::shared_ptr<executor::TaskExecutor> executor,
- const RemoteCommandRequest& rcr,
- Options options = {});
+ explicit TaskExecutorCursor(executor::TaskExecutor* executor,
+ const RemoteCommandRequest& rcr,
+ Options&& options = {});
/**
* Construct the cursor from a cursor response from a previously executed RemoteCommandRequest.
* The executor is used for subsequent getMore calls. Uses the original RemoteCommandRequest
* to build subsequent commands. Takes ownership of the CursorResponse and gives it to the new
* cursor.
- * If the cursor should reuse the original transport connection that opened the original
- * cursor, make sure the pinning executor that was used to open that cursor is provided.
*/
- TaskExecutorCursor(std::shared_ptr<executor::TaskExecutor> executor,
- std::shared_ptr<executor::TaskExecutor> underlyingExec,
+ TaskExecutorCursor(executor::TaskExecutor* executor,
CursorResponse&& response,
RemoteCommandRequest& rcr,
Options&& options = {});
@@ -176,6 +158,14 @@ public:
return _additionalCursors.size();
}
+ /**
+ * Return the callback that this cursor is waiting on. Can be used to block on getting a
+ * response to this request. Can be boost::none.
+ */
+ auto getCallbackHandle() {
+ return _cbHandle;
+ }
+
private:
/**
* Runs a remote command and pipes the output back to this object
@@ -183,34 +173,24 @@ private:
void _runRemoteCommand(const RemoteCommandRequest& rcr);
/**
- * Gets the next batch with interruptibility via the opCtx.
+ * Gets the next batch with interruptibility via the opCtx
*/
void _getNextBatch(OperationContext* opCtx);
/**
* Helper for '_getNextBatch' that handles the reading of the 'CursorResponse' object and
- * storing of relevant values.
+ * storing of relevant values. This is also responsible for issuing a getMore request if it
+ * is required to populate the next batch.
*/
void _processResponse(OperationContext* opCtx, CursorResponse&& response);
+
/**
* Create a new request, annotating with lsid and current opCtx
*/
const RemoteCommandRequest& _createRequest(OperationContext* opCtx, const BSONObj& cmd);
- /**
- * Schedules a 'GetMore' request to run asyncronously.
- * This function can only be invoked when:
- * - There is no in-flight request ('_cmdState' is null).
- * - We have an open '_cursorId'.
- */
- void _scheduleGetMore(OperationContext* opCtx);
-
- std::shared_ptr<executor::TaskExecutor> _executor;
- // If we are pinning connections, we need to keep a separate reference to the
- // non-pinning, normal executor, so that we can shut down the pinned executor
- // out-of-line.
- std::shared_ptr<executor::TaskExecutor> _underlyingExecutor;
+ executor::TaskExecutor* _executor;
// Used as a scratch pad for the successive scheduleRemoteCommand calls
RemoteCommandRequest _rcr;
@@ -220,20 +200,8 @@ private:
// If the opCtx is in our initial request, re-use it for all subsequent operations
boost::optional<LogicalSessionId> _lsid;
- struct CommandState {
- TaskExecutor::CallbackHandle cbHandle;
- SharedPromise<BSONObj> promise;
- };
-
- /**
- * Maintains the state for the in progress command (if there is any):
- * - Handle for the task scheduled on `_executor`.
- * - A promise that will be emplaced by the result of running the command.
- *
- * The state may outlive `TaskExecutorCursor` and is shared with the callback that runs on
- * `_executor` upon completion of the remote command.
- */
- std::shared_ptr<CommandState> _cmdState;
+ // Stash the callbackhandle for the current outstanding operation
+ boost::optional<TaskExecutor::CallbackHandle> _cbHandle;
CursorId _cursorId = kUnitializedCursorId;
@@ -254,33 +222,13 @@ private:
decltype(_batch)::iterator _batchIter;
long long _batchNum = 0;
+ // Multi producer because we hold onto the producer side in this object, as well as placing it
+ // into callbacks for the task executor
+ MultiProducerSingleConsumerQueue<StatusWith<BSONObj>>::Pipe _pipe;
+
// Cursors built from the responses returned alongside the results for this cursor.
std::vector<TaskExecutorCursor> _additionalCursors;
};
-// Make a new TaskExecutorCursor using the provided executor, RCR, and options. If we fail to create
-// the cursor, the retryPolicy can inspect the error and make a decision as to whether we should
-// retry. If we do retry, the error is swallowed and another attempt is made. If we don't retry,
-// this function throws the error we failed with.
-inline TaskExecutorCursor makeTaskExecutorCursor(
- OperationContext* opCtx,
- std::shared_ptr<executor::TaskExecutor> executor,
- const RemoteCommandRequest& rcr,
- TaskExecutorCursor::Options options = {},
- std::function<bool(Status)> retryPolicy = nullptr) {
- for (;;) {
- try {
- TaskExecutorCursor tec(executor, rcr, options);
- tec.populateCursor(opCtx);
- return tec;
- } catch (const DBException& ex) {
- bool shouldRetry = retryPolicy && retryPolicy(ex.toStatus());
- if (!shouldRetry) {
- throw;
- }
- }
- }
-}
-
} // namespace executor
} // namespace mongo