Update Node_modules

This commit is contained in:
2019-07-02 16:05:15 +02:00
parent 95a0ff3913
commit 9dab5e2dbc
3495 changed files with 39728 additions and 465154 deletions
-749
View File
@@ -1,749 +0,0 @@
'use strict';
const Logger = require('../connection/logger');
const BSON = require('../connection/utils').retrieveBSON();
const MongoError = require('../error').MongoError;
const MongoNetworkError = require('../error').MongoNetworkError;
const mongoErrorContextSymbol = require('../error').mongoErrorContextSymbol;
const Long = BSON.Long;
const deprecate = require('util').deprecate;
const readPreferenceServerSelector = require('./server_selectors').readPreferenceServerSelector;
const ReadPreference = require('../topologies/read_preference');
/**
* Handle callback (including any exceptions thrown)
*/
function handleCallback(callback, err, result) {
try {
callback(err, result);
} catch (err) {
process.nextTick(function() {
throw err;
});
}
}
/**
* This is a cursor results callback
*
* @callback resultCallback
* @param {error} error An error object. Set to null if no error present
* @param {object} document
*/
/**
* An internal class that embodies a cursor on MongoDB, allowing for iteration over the
* results returned from a query.
*
* @property {number} cursorBatchSize The current cursorBatchSize for the cursor
* @property {number} cursorLimit The current cursorLimit for the cursor
* @property {number} cursorSkip The current cursorSkip for the cursor
*/
class Cursor {
/**
* Create a cursor
*
* @param {object} bson An instance of the BSON parser
* @param {string} ns The MongoDB fully qualified namespace (ex: db1.collection1)
* @param {{object}|Long} cmd The selector (can be a command or a cursorId)
* @param {object} [options=null] Optional settings.
* @param {object} [options.batchSize=1000] Batchsize for the operation
* @param {array} [options.documents=[]] Initial documents list for cursor
* @param {object} [options.transforms=null] Transform methods for the cursor results
* @param {function} [options.transforms.query] Transform the value returned from the initial query
* @param {function} [options.transforms.doc] Transform each document returned from Cursor.prototype.next
* @param {object} topology The server topology instance.
* @param {object} topologyOptions The server topology options.
*/
constructor(bson, ns, cmd, options, topology, topologyOptions) {
options = options || {};
// Cursor pool
this.pool = null;
// Cursor server
this.server = null;
// Do we have a not connected handler
this.disconnectHandler = options.disconnectHandler;
// Set local values
this.bson = bson;
this.ns = ns;
this.cmd = cmd;
this.options = options;
this.topology = topology;
// All internal state
this.s = {
cursorId: null,
cmd: cmd,
documents: options.documents || [],
cursorIndex: 0,
dead: false,
killed: false,
init: false,
notified: false,
limit: options.limit || cmd.limit || 0,
skip: options.skip || cmd.skip || 0,
batchSize: options.batchSize || cmd.batchSize || 1000,
currentLimit: 0,
// Result field name if not a cursor (contains the array of results)
transforms: options.transforms
};
if (typeof options.session === 'object') {
this.s.session = options.session;
}
// Add promoteLong to cursor state
if (typeof topologyOptions.promoteLongs === 'boolean') {
this.s.promoteLongs = topologyOptions.promoteLongs;
} else if (typeof options.promoteLongs === 'boolean') {
this.s.promoteLongs = options.promoteLongs;
}
// Add promoteValues to cursor state
if (typeof topologyOptions.promoteValues === 'boolean') {
this.s.promoteValues = topologyOptions.promoteValues;
} else if (typeof options.promoteValues === 'boolean') {
this.s.promoteValues = options.promoteValues;
}
// Add promoteBuffers to cursor state
if (typeof topologyOptions.promoteBuffers === 'boolean') {
this.s.promoteBuffers = topologyOptions.promoteBuffers;
} else if (typeof options.promoteBuffers === 'boolean') {
this.s.promoteBuffers = options.promoteBuffers;
}
if (topologyOptions.reconnect) {
this.s.reconnect = topologyOptions.reconnect;
}
// Logger
this.logger = Logger('Cursor', topologyOptions);
//
// Did we pass in a cursor id
if (typeof cmd === 'number') {
this.s.cursorId = Long.fromNumber(cmd);
this.s.lastCursorId = this.s.cursorId;
} else if (cmd instanceof Long) {
this.s.cursorId = cmd;
this.s.lastCursorId = cmd;
}
}
setCursorBatchSize(value) {
this.s.batchSize = value;
}
cursorBatchSize() {
return this.s.batchSize;
}
setCursorLimit(value) {
this.s.limit = value;
}
cursorLimit() {
return this.s.limit;
}
setCursorSkip(value) {
this.s.skip = value;
}
cursorSkip() {
return this.s.skip;
}
_endSession(options, callback) {
if (typeof options === 'function') {
callback = options;
options = {};
}
options = options || {};
const session = this.s.session;
if (session && (options.force || session.owner === this)) {
this.s.session = undefined;
session.endSession(callback);
return true;
}
if (callback) {
callback();
}
return false;
}
/**
* Clone the cursor
* @method
* @return {Cursor}
*/
clone() {
return this.topology.cursor(this.ns, this.cmd, this.options);
}
/**
* Checks if the cursor is dead
* @method
* @return {boolean} A boolean signifying if the cursor is dead or not
*/
isDead() {
return this.s.dead === true;
}
/**
* Checks if the cursor was killed by the application
* @method
* @return {boolean} A boolean signifying if the cursor was killed by the application
*/
isKilled() {
return this.s.killed === true;
}
/**
* Checks if the cursor notified it's caller about it's death
* @method
* @return {boolean} A boolean signifying if the cursor notified the callback
*/
isNotified() {
return this.s.notified === true;
}
/**
* Returns current buffered documents length
* @method
* @return {number} The number of items in the buffered documents
*/
bufferedCount() {
return this.s.documents.length - this.s.cursorIndex;
}
/**
* Kill the cursor
*
* @param {resultCallback} callback A callback function
*/
kill(callback) {
// Set cursor to dead
this.s.dead = true;
this.s.killed = true;
// Remove documents
this.s.documents = [];
// If no cursor id just return
if (this.s.cursorId == null || this.s.cursorId.isZero() || this.s.init === false) {
if (callback) callback(null, null);
return;
}
// Default pool
const pool = this.s.server.s.pool;
// Execute command
this.s.server.s.wireProtocolHandler.killCursor(this.bson, this.ns, this.s, pool, callback);
}
/**
* Resets the cursor
*/
rewind() {
if (this.s.init) {
if (!this.s.dead) {
this.kill();
}
this.s.currentLimit = 0;
this.s.init = false;
this.s.dead = false;
this.s.killed = false;
this.s.notified = false;
this.s.documents = [];
this.s.cursorId = null;
this.s.cursorIndex = 0;
}
}
/**
* Returns current buffered documents
* @method
* @return {Array} An array of buffered documents
*/
readBufferedDocuments(number) {
const unreadDocumentsLength = this.s.documents.length - this.s.cursorIndex;
const length = number < unreadDocumentsLength ? number : unreadDocumentsLength;
let elements = this.s.documents.slice(this.s.cursorIndex, this.s.cursorIndex + length);
// Transform the doc with passed in transformation method if provided
if (this.s.transforms && typeof this.s.transforms.doc === 'function') {
// Transform all the elements
for (let i = 0; i < elements.length; i++) {
elements[i] = this.s.transforms.doc(elements[i]);
}
}
// Ensure we do not return any more documents than the limit imposed
// Just return the number of elements up to the limit
if (this.s.limit > 0 && this.s.currentLimit + elements.length > this.s.limit) {
elements = elements.slice(0, this.s.limit - this.s.currentLimit);
this.kill();
}
// Adjust current limit
this.s.currentLimit = this.s.currentLimit + elements.length;
this.s.cursorIndex = this.s.cursorIndex + elements.length;
// Return elements
return elements;
}
/**
* Retrieve the next document from the cursor
*
* @param {resultCallback} callback A callback function
*/
next(callback) {
nextFunction(this, callback);
}
}
Cursor.prototype._find = deprecate(
callback => _find(this, callback),
'_find() is deprecated, please stop using it'
);
Cursor.prototype._getmore = deprecate(
callback => _getmore(this, callback),
'_getmore() is deprecated, please stop using it'
);
function _getmore(cursor, callback) {
if (cursor.logger.isDebug()) {
cursor.logger.debug(`schedule getMore call for query [${JSON.stringify(cursor.query)}]`);
}
// Determine if it's a raw query
const raw = cursor.options.raw || cursor.cmd.raw;
// Set the current batchSize
let batchSize = cursor.s.batchSize;
if (cursor.s.limit > 0 && cursor.s.currentLimit + batchSize > cursor.s.limit) {
batchSize = cursor.s.limit - cursor.s.currentLimit;
}
// Default pool
const pool = cursor.s.server.s.pool;
// We have a wire protocol handler
cursor.s.server.s.wireProtocolHandler.getMore(
cursor.bson,
cursor.ns,
cursor.s,
batchSize,
raw,
pool,
cursor.options,
callback
);
}
function _find(cursor, callback) {
if (cursor.logger.isDebug()) {
cursor.logger.debug(
`issue initial query [${JSON.stringify(cursor.cmd)}] with flags [${JSON.stringify(
cursor.query
)}]`
);
}
const queryCallback = (err, r) => {
if (err) return callback(err);
// Get the raw message
const result = r.message;
// Query failure bit set
if (result.queryFailure) {
return callback(new MongoError(result.documents[0]), null);
}
// Check if we have a command cursor
if (
Array.isArray(result.documents) &&
result.documents.length === 1 &&
(!cursor.cmd.find || (cursor.cmd.find && cursor.cmd.virtual === false)) &&
(result.documents[0].cursor !== 'string' ||
result.documents[0]['$err'] ||
result.documents[0]['errmsg'] ||
Array.isArray(result.documents[0].result))
) {
// We have a an error document return the error
if (result.documents[0]['$err'] || result.documents[0]['errmsg']) {
return callback(new MongoError(result.documents[0]), null);
}
// We have a cursor document
if (result.documents[0].cursor != null && typeof result.documents[0].cursor !== 'string') {
const id = result.documents[0].cursor.id;
// If we have a namespace change set the new namespace for getmores
if (result.documents[0].cursor.ns) {
cursor.ns = result.documents[0].cursor.ns;
}
// Promote id to long if needed
cursor.s.cursorId = typeof id === 'number' ? Long.fromNumber(id) : id;
cursor.s.lastCursorId = cursor.s.cursorId;
// If we have a firstBatch set it
if (Array.isArray(result.documents[0].cursor.firstBatch)) {
cursor.s.documents = result.documents[0].cursor.firstBatch;
}
// Return after processing command cursor
return callback(null, result);
}
if (Array.isArray(result.documents[0].result)) {
cursor.s.documents = result.documents[0].result;
cursor.s.cursorId = Long.ZERO;
return callback(null, result);
}
}
// Otherwise fall back to regular find path
cursor.s.cursorId = result.cursorId;
cursor.s.documents = result.documents;
cursor.s.lastCursorId = result.cursorId;
// Transform the results with passed in transformation method if provided
if (cursor.s.transforms && typeof cursor.s.transforms.query === 'function') {
cursor.s.documents = cursor.s.transforms.query(result);
}
// Return callback
callback(null, result);
};
// Options passed to the pool
const queryOptions = {};
// If we have a raw query decorate the function
if (cursor.options.raw || cursor.cmd.raw) {
queryOptions.raw = cursor.options.raw || cursor.cmd.raw;
}
// Do we have documentsReturnedIn set on the query
if (typeof cursor.query.documentsReturnedIn === 'string') {
queryOptions.documentsReturnedIn = cursor.query.documentsReturnedIn;
}
// Add promote Long value if defined
if (typeof cursor.s.promoteLongs === 'boolean') {
queryOptions.promoteLongs = cursor.s.promoteLongs;
}
// Add promote values if defined
if (typeof cursor.s.promoteValues === 'boolean') {
queryOptions.promoteValues = cursor.s.promoteValues;
}
// Add promote values if defined
if (typeof cursor.s.promoteBuffers === 'boolean') {
queryOptions.promoteBuffers = cursor.s.promoteBuffers;
}
if (typeof cursor.s.session === 'object') {
queryOptions.session = cursor.s.session;
}
// Write the initial command out
cursor.s.server.s.pool.write(cursor.query, queryOptions, queryCallback);
}
/**
* Validate if the pool is dead and return error
*/
function isConnectionDead(cursor, callback) {
if (cursor.pool && cursor.pool.isDestroyed()) {
cursor.s.killed = true;
const err = new MongoNetworkError(
`connection to host ${cursor.pool.host}:${cursor.pool.port} was destroyed`
);
_setCursorNotifiedImpl(cursor, () => callback(err));
return true;
}
return false;
}
/**
* Validate if the cursor is dead but was not explicitly killed by user
*/
function isCursorDeadButNotkilled(cursor, callback) {
// Cursor is dead but not marked killed, return null
if (cursor.s.dead && !cursor.s.killed) {
cursor.s.killed = true;
setCursorNotified(cursor, callback);
return true;
}
return false;
}
/**
* Validate if the cursor is dead and was killed by user
*/
function isCursorDeadAndKilled(cursor, callback) {
if (cursor.s.dead && cursor.s.killed) {
handleCallback(callback, new MongoError('cursor is dead'));
return true;
}
return false;
}
/**
* Validate if the cursor was killed by the user
*/
function isCursorKilled(cursor, callback) {
if (cursor.s.killed) {
setCursorNotified(cursor, callback);
return true;
}
return false;
}
/**
* Mark cursor as being dead and notified
*/
function setCursorDeadAndNotified(cursor, callback) {
cursor.s.dead = true;
setCursorNotified(cursor, callback);
}
/**
* Mark cursor as being notified
*/
function setCursorNotified(cursor, callback) {
_setCursorNotifiedImpl(cursor, () => handleCallback(callback, null, null));
}
function _setCursorNotifiedImpl(cursor, callback) {
cursor.s.notified = true;
cursor.s.documents = [];
cursor.s.cursorIndex = 0;
if (cursor._endSession) {
return cursor._endSession(undefined, () => callback());
}
return callback();
}
function initializeCursorAndRetryNext(cursor, callback) {
cursor.topology.selectServer(
readPreferenceServerSelector(cursor.options.readPreference || ReadPreference.primary),
(err, server) => {
if (err) {
callback(err, null);
return;
}
cursor.s.server = server;
cursor.s.init = true;
// check if server supports collation
// NOTE: this should be a part of the selection predicate!
if (cursor.cmd && cursor.cmd.collation && cursor.server.description.maxWireVersion < 5) {
callback(new MongoError(`server ${cursor.server.name} does not support collation`));
return;
}
try {
cursor.query = cursor.s.server.s.wireProtocolHandler.command(
cursor.bson,
cursor.ns,
cursor.cmd,
cursor.s,
cursor.topology,
cursor.options
);
nextFunction(cursor, callback);
} catch (err) {
callback(err);
return;
}
}
);
}
function nextFunction(cursor, callback) {
// We have notified about it
if (cursor.s.notified) {
return callback(new Error('cursor is exhausted'));
}
// Cursor is killed return null
if (isCursorKilled(cursor, callback)) return;
// Cursor is dead but not marked killed, return null
if (isCursorDeadButNotkilled(cursor, callback)) return;
// We have a dead and killed cursor, attempting to call next should error
if (isCursorDeadAndKilled(cursor, callback)) return;
// We have just started the cursor
if (!cursor.s.init) {
return initializeCursorAndRetryNext(cursor, callback);
}
// If we don't have a cursorId execute the first query
if (cursor.s.cursorId == null) {
// Check if pool is dead and return if not possible to
// execute the query against the db
if (isConnectionDead(cursor, callback)) return;
// query, cmd, options, s, callback
return _find(cursor, function(err) {
if (err) return handleCallback(callback, err, null);
if (cursor.s.cursorId && cursor.s.cursorId.isZero() && cursor._endSession) {
cursor._endSession();
}
if (
cursor.s.documents.length === 0 &&
cursor.s.cursorId &&
cursor.s.cursorId.isZero() &&
!cursor.cmd.tailable &&
!cursor.cmd.awaitData
) {
return setCursorNotified(cursor, callback);
}
nextFunction(cursor, callback);
});
}
if (cursor.s.documents.length === cursor.s.cursorIndex && Long.ZERO.equals(cursor.s.cursorId)) {
setCursorDeadAndNotified(cursor, callback);
return;
}
if (cursor.s.limit > 0 && cursor.s.currentLimit >= cursor.s.limit) {
// Ensure we kill the cursor on the server
cursor.kill();
// Set cursor in dead and notified state
setCursorDeadAndNotified(cursor, callback);
return;
}
if (
cursor.s.documents.length === cursor.s.cursorIndex &&
cursor.cmd.tailable &&
Long.ZERO.equals(cursor.s.cursorId)
) {
return handleCallback(
callback,
new MongoError({
message: 'No more documents in tailed cursor',
tailable: cursor.cmd.tailable,
awaitData: cursor.cmd.awaitData
})
);
}
if (cursor.s.cursorIndex === cursor.s.documents.length && !Long.ZERO.equals(cursor.s.cursorId)) {
// Ensure an empty cursor state
cursor.s.documents = [];
cursor.s.cursorIndex = 0;
// Check if connection is dead and return if not possible to
if (isConnectionDead(cursor, callback)) return;
// Execute the next get more
return _getmore(cursor, function(err, doc, connection) {
if (err) {
if (err instanceof MongoError) {
err[mongoErrorContextSymbol].isGetMore = true;
}
return handleCallback(callback, err);
}
if (cursor.s.cursorId && cursor.s.cursorId.isZero() && cursor._endSession) {
cursor._endSession();
}
// Save the returned connection to ensure all getMore's fire over the same connection
cursor.connection = connection;
// Tailable cursor getMore result, notify owner about it
// No attempt is made here to retry, this is left to the user of the
// core module to handle to keep core simple
if (
cursor.s.documents.length === 0 &&
cursor.cmd.tailable &&
Long.ZERO.equals(cursor.s.cursorId)
) {
// No more documents in the tailed cursor
return handleCallback(
callback,
new MongoError({
message: 'No more documents in tailed cursor',
tailable: cursor.cmd.tailable,
awaitData: cursor.cmd.awaitData
})
);
} else if (
cursor.s.documents.length === 0 &&
cursor.cmd.tailable &&
!Long.ZERO.equals(cursor.s.cursorId)
) {
return nextFunction(cursor, callback);
}
if (cursor.s.limit > 0 && cursor.s.currentLimit >= cursor.s.limit) {
return setCursorDeadAndNotified(cursor, callback);
}
nextFunction(cursor, callback);
});
}
if (cursor.s.limit > 0 && cursor.s.currentLimit >= cursor.s.limit) {
// Ensure we kill the cursor on the server
cursor.kill();
// Set cursor in dead and notified state
return setCursorDeadAndNotified(cursor, callback);
}
// Increment the current cursor limit
cursor.s.currentLimit += 1;
// Get the document
let doc = cursor.s.documents[cursor.s.cursorIndex++];
// Doc overflow
if (!doc || doc.$err) {
// Ensure we kill the cursor on the server
cursor.kill();
// Set cursor in dead and notified state
return setCursorDeadAndNotified(cursor, function() {
handleCallback(callback, new MongoError(doc ? doc.$err : undefined));
});
}
// Transform the doc with passed in transformation method if provided
if (cursor.s.transforms && typeof cursor.s.transforms.doc === 'function') {
doc = cursor.s.transforms.doc(doc);
}
// Return the document
handleCallback(callback, null, doc);
}
module.exports = Cursor;
+29 -18
View File
@@ -122,7 +122,15 @@ class ServerHeartbeatFailedEvent {
*
* @param {Server} server The server to monitor
*/
function monitorServer(server) {
function monitorServer(server, options) {
options = options || {};
const heartbeatFrequencyMS = options.heartbeatFrequencyMS || 10000;
if (options.initial === true) {
server.s.monitorId = setTimeout(() => monitorServer(server), heartbeatFrequencyMS);
return;
}
// executes a single check of a server
const checkServer = callback => {
let start = process.hrtime();
@@ -130,6 +138,9 @@ function monitorServer(server) {
// emit a signal indicating we have started the heartbeat
server.emit('serverHeartbeatStarted', new ServerHeartbeatStartedEvent(server.name));
// NOTE: legacy monitoring event
process.nextTick(() => server.emit('monitoring', server));
server.command(
'admin.$cmd',
{ ismaster: true },
@@ -137,7 +148,7 @@ function monitorServer(server) {
monitoring: true,
socketTimeout: server.s.options.connectionTimeout || 2000
},
function(err, result) {
(err, result) => {
let duration = calculateDurationInMs(start);
if (err) {
@@ -167,10 +178,7 @@ function monitorServer(server) {
server.emit('descriptionReceived', new ServerDescription(server.description.address, isMaster));
// schedule the next monitoring process
server.s.monitorId = setTimeout(
() => monitorServer(server),
server.s.options.heartbeatFrequencyMS
);
server.s.monitorId = setTimeout(() => monitorServer(server), heartbeatFrequencyMS);
};
// run the actual monitoring loop
@@ -184,21 +192,24 @@ function monitorServer(server) {
// According to the SDAM specification's "Network error during server check" section, if
// an ismaster call fails we reset the server's pool. If a server was once connected,
// change its type to `Unknown` only after retrying once.
server.s.pool.reset(() => {
// otherwise re-attempt monitoring once
checkServer((error, isMaster) => {
if (error) {
server.s.monitoring = false;
// TODO: we need to reset the pool here
// we revert to an `Unknown` by emitting a default description with no isMaster
server.emit(
'descriptionReceived',
new ServerDescription(server.description.address, null, { error })
);
return checkServer((err, isMaster) => {
if (err) {
server.s.monitoring = false;
// we do not reschedule monitoring in this case
return;
}
// revert to `Unknown` by emitting a default description with no isMaster
server.emit('descriptionReceived', new ServerDescription(server.description.address));
// do not reschedule monitoring in this case
return;
}
successHandler(isMaster);
successHandler(isMaster);
});
});
});
}
+148 -169
View File
@@ -3,16 +3,49 @@ const EventEmitter = require('events');
const MongoError = require('../error').MongoError;
const Pool = require('../connection/pool');
const relayEvents = require('../utils').relayEvents;
const calculateDurationInMs = require('../utils').calculateDurationInMs;
const Query = require('../connection/commands').Query;
const TwoSixWireProtocolSupport = require('../wireprotocol/2_6_support');
const ThreeTwoWireProtocolSupport = require('../wireprotocol/3_2_support');
const wireProtocol = require('../wireprotocol');
const BSON = require('../connection/utils').retrieveBSON();
const createClientInfo = require('../topologies/shared').createClientInfo;
const Logger = require('../connection/logger');
const ServerDescription = require('./server_description').ServerDescription;
const ReadPreference = require('../topologies/read_preference');
const monitorServer = require('./monitoring').monitorServer;
const MongoParseError = require('../error').MongoParseError;
const MongoNetworkError = require('../error').MongoNetworkError;
const collationNotSupported = require('../utils').collationNotSupported;
const debugOptions = require('../connection/utils').debugOptions;
// Used for filtering out fields for logging
const DEBUG_FIELDS = [
'reconnect',
'reconnectTries',
'reconnectInterval',
'emitError',
'cursorFactory',
'host',
'port',
'size',
'keepAlive',
'keepAliveInitialDelay',
'noDelay',
'connectionTimeout',
'checkServerIdentity',
'socketTimeout',
'ssl',
'ca',
'crl',
'cert',
'key',
'rejectUnauthorized',
'promoteLongs',
'promoteValues',
'promoteBuffers',
'servername'
];
const STATE_DISCONNECTED = 0;
const STATE_CONNECTING = 1;
const STATE_CONNECTED = 2;
/**
*
@@ -27,7 +60,7 @@ class Server extends EventEmitter {
* @param {ServerDescription} description
* @param {Object} options
*/
constructor(description, options) {
constructor(description, options, topology) {
super();
this.s = {
@@ -43,8 +76,14 @@ class Server extends EventEmitter {
clientInfo: createClientInfo(options),
// state variable to determine if there is an active server check in progress
monitoring: false,
// the implementation of the monitoring method
monitorFunction: options.monitorFunction || monitorServer,
// the connection pool
pool: null
pool: null,
// the server state
state: STATE_DISCONNECTED,
credentials: options.credentials,
topology
};
}
@@ -58,8 +97,6 @@ class Server extends EventEmitter {
/**
* Initiate server connect
*
* @param {Array} [options.auth] Array of auth options to apply on connect
*/
connect(options) {
options = options || {};
@@ -70,21 +107,35 @@ class Server extends EventEmitter {
}
// create a pool
this.s.pool = new Pool(this, Object.assign(this.s.options, options, { bson: this.s.bson }));
const addressParts = this.description.address.split(':');
const poolOptions = Object.assign(
{ host: addressParts[0], port: parseInt(addressParts[1], 10) },
this.s.options,
options,
{ bson: this.s.bson }
);
// Set up listeners
// NOTE: this should only be the case if we are connecting to a single server
poolOptions.reconnect = true;
this.s.pool = new Pool(this, poolOptions);
// setup listeners
this.s.pool.on('connect', connectEventHandler(this));
this.s.pool.on('close', closeEventHandler(this));
this.s.pool.on('close', errorEventHandler(this));
this.s.pool.on('error', errorEventHandler(this));
this.s.pool.on('parseError', parseErrorEventHandler(this));
// this.s.pool.on('error', errorEventHandler(this));
// it is unclear whether consumers should even know about these events
// this.s.pool.on('timeout', timeoutEventHandler(this));
// this.s.pool.on('parseError', errorEventHandler(this));
// this.s.pool.on('reconnect', reconnectEventHandler(this));
// this.s.pool.on('reconnectFailed', errorEventHandler(this));
// relay all command monitoring events
relayEvents(this.s.pool, this, ['commandStarted', 'commandSucceeded', 'commandFailed']);
this.s.state = STATE_CONNECTING;
// If auth settings have been provided, use them
if (options.auth) {
this.s.pool.connect.apply(this.s.pool, options.auth);
@@ -97,24 +148,44 @@ class Server extends EventEmitter {
/**
* Destroy the server connection
*
* @param {Boolean} [options.emitClose=false] Emit close event on destroy
* @param {Boolean} [options.emitDestroy=false] Emit destroy event on destroy
* @param {Boolean} [options.force=false] Force destroy the pool
*/
destroy(callback) {
if (typeof callback === 'function') {
callback(null, null);
destroy(options, callback) {
if (typeof options === 'function') (callback = options), (options = {});
options = Object.assign({}, { force: false }, options);
if (!this.s.pool) {
this.s.state = STATE_DISCONNECTED;
if (typeof callback === 'function') {
callback(null, null);
}
return;
}
['close', 'error', 'timeout', 'parseError', 'connect'].forEach(event => {
this.s.pool.removeAllListeners(event);
});
if (this.s.monitorId) {
clearTimeout(this.s.monitorId);
}
this.s.pool.destroy(options.force, err => {
this.s.state = STATE_DISCONNECTED;
callback(err);
});
}
/**
* Immediately schedule monitoring of this server. If there already an attempt being made
* this will be a no-op.
*/
monitor() {
if (this.s.monitoring) return;
monitor(options) {
options = options || {};
if (this.s.state !== STATE_CONNECTED || this.s.monitoring) return;
if (this.s.monitorId) clearTimeout(this.s.monitorId);
monitorServer(this);
this.s.monitorFunction(this, options);
}
/**
@@ -146,39 +217,21 @@ class Server extends EventEmitter {
// Debug log
if (this.s.logger.isDebug()) {
this.s.logger.debug(
`executing command [${JSON.stringify({ ns, cmd, options })}] against ${this.name}`
`executing command [${JSON.stringify({
ns,
cmd,
options: debugOptions(DEBUG_FIELDS, options)
})}] against ${this.name}`
);
}
// Check if we have collation support
if (this.description.maxWireVersion < 5 && cmd.collation) {
// error if collation not supported
if (collationNotSupported(this, cmd)) {
callback(new MongoError(`server ${this.name} does not support collation`));
return;
}
// Are we executing against a specific topology
const topology = options.topology || {};
// Create the query object
const query = this.s.wireProtocolHandler.command(this.s.bson, ns, cmd, {}, topology, options);
// Set slave OK of the query
query.slaveOk = options.readPreference ? options.readPreference.slaveOk() : false;
// write options
const writeOptions = {
raw: typeof options.raw === 'boolean' ? options.raw : false,
promoteLongs: typeof options.promoteLongs === 'boolean' ? options.promoteLongs : true,
promoteValues: typeof options.promoteValues === 'boolean' ? options.promoteValues : true,
promoteBuffers: typeof options.promoteBuffers === 'boolean' ? options.promoteBuffers : false,
command: true,
monitoring: typeof options.monitoring === 'boolean' ? options.monitoring : false,
fullResult: typeof options.fullResult === 'boolean' ? options.fullResult : false,
requestId: query.requestId,
socketTimeout: typeof options.socketTimeout === 'number' ? options.socketTimeout : null,
session: options.session || null
};
// write the operation to the pool
this.s.pool.write(query, writeOptions, callback);
wireProtocol.command(this, ns, cmd, options, callback);
}
/**
@@ -230,6 +283,15 @@ class Server extends EventEmitter {
}
}
Object.defineProperty(Server.prototype, 'clusterTime', {
get: function() {
return this.s.topology.clusterTime;
},
set: function(clusterTime) {
this.s.topology.clusterTime = clusterTime;
}
});
function basicWriteValidations(server) {
if (!server.s.pool) {
return new MongoError('server instance is not connected');
@@ -269,145 +331,62 @@ function executeWriteOperation(args, options, callback) {
return;
}
// Check if we have collation support
if (server.description.maxWireVersion < 5 && options.collation) {
if (collationNotSupported(server, options)) {
callback(new MongoError(`server ${this.name} does not support collation`));
return;
}
// Execute write
return server.s.wireProtocolHandler[op](server.s.pool, ns, server.s.bson, ops, options, callback);
}
function saslSupportedMechs(options) {
if (!options) {
return {};
}
const authArray = options.auth || [];
const authMechanism = authArray[0] || options.authMechanism;
const authSource = authArray[1] || options.authSource || options.dbName || 'admin';
const user = authArray[2] || options.user;
if (typeof authMechanism === 'string' && authMechanism.toUpperCase() !== 'DEFAULT') {
return {};
}
if (!user) {
return {};
}
return { saslSupportedMechs: `${authSource}.${user}` };
}
function extractIsMasterError(err, result) {
if (err) return err;
if (result && result.result && result.result.ok === 0) {
return new MongoError(result.result);
}
}
function executeServerHandshake(server, callback) {
// construct an `ismaster` query
const compressors =
server.s.options.compression && server.s.options.compression.compressors
? server.s.options.compression.compressors
: [];
const queryOptions = { numberToSkip: 0, numberToReturn: -1, checkKeys: false, slaveOk: true };
const query = new Query(
server.s.bson,
'admin.$cmd',
Object.assign(
{ ismaster: true, client: server.s.clientInfo, compression: compressors },
saslSupportedMechs(server.s.options)
),
queryOptions
);
// execute the query
server.s.pool.write(
query,
{ socketTimeout: server.s.options.connectionTimeout || 2000 },
callback
);
}
function configureWireProtocolHandler(ismaster) {
// 3.2 wire protocol handler
if (ismaster.maxWireVersion >= 4) {
return new ThreeTwoWireProtocolSupport();
}
// default to 2.6 wire protocol handler
return new TwoSixWireProtocolSupport();
return wireProtocol[op](server, ns, ops, options, callback);
}
function connectEventHandler(server) {
return function() {
// log information of received information if in info mode
// if (server.s.logger.isInfo()) {
// var object = err instanceof MongoError ? JSON.stringify(err) : {};
// server.s.logger.info(`server ${server.name} fired event ${event} out with message ${object}`);
// }
return function(pool, conn) {
const ismaster = conn.ismaster;
server.s.lastIsMasterMS = conn.lastIsMasterMS;
if (conn.agreedCompressor) {
server.s.pool.options.agreedCompressor = conn.agreedCompressor;
}
// begin initial server handshake
const start = process.hrtime();
executeServerHandshake(server, (err, response) => {
// Set initial lastIsMasterMS - is this needed?
server.s.lastIsMasterMS = calculateDurationInMs(start);
if (conn.zlibCompressionLevel) {
server.s.pool.options.zlibCompressionLevel = conn.zlibCompressionLevel;
}
const serverError = extractIsMasterError(err, response);
if (serverError) {
server.emit('error', serverError);
return;
}
if (conn.ismaster.$clusterTime) {
const $clusterTime = conn.ismaster.$clusterTime;
server.s.sclusterTime = $clusterTime;
}
// extract the ismaster from the server response
const isMaster = response.result;
// compression negotation
if (isMaster && isMaster.compression) {
const localCompressionInfo = server.s.options.compression;
const localCompressors = localCompressionInfo.compressors;
for (var i = 0; i < localCompressors.length; i++) {
if (isMaster.compression.indexOf(localCompressors[i]) > -1) {
server.s.pool.options.agreedCompressor = localCompressors[i];
break;
}
}
if (localCompressionInfo.zlibCompressionLevel) {
server.s.pool.options.zlibCompressionLevel = localCompressionInfo.zlibCompressionLevel;
}
}
// configure the wire protocol handler
server.s.wireProtocolHandler = configureWireProtocolHandler(isMaster);
// log the connection event if requested
if (server.s.logger.isInfo()) {
server.s.logger.info(
`server ${server.name} connected with ismaster [${JSON.stringify(isMaster)}]`
);
}
// emit an event indicating that our description has changed
server.emit(
'descriptionReceived',
new ServerDescription(server.description.address, isMaster)
// log the connection event if requested
if (server.s.logger.isInfo()) {
server.s.logger.info(
`server ${server.name} connected with ismaster [${JSON.stringify(ismaster)}]`
);
}
// emit a connect event
server.emit('connect', isMaster);
});
// emit an event indicating that our description has changed
server.emit('descriptionReceived', new ServerDescription(server.description.address, ismaster));
// we are connected and handshaked (guaranteed by the pool)
server.s.state = STATE_CONNECTED;
server.emit('connect', server);
};
}
function closeEventHandler(server) {
return function() {
function errorEventHandler(server) {
return function(err) {
if (err) {
server.emit('error', new MongoNetworkError(err));
}
server.emit('close');
};
}
function parseErrorEventHandler(server) {
return function(err) {
server.s.state = STATE_DISCONNECTED;
server.emit('error', new MongoParseError(err));
};
}
module.exports = Server;
+10 -2
View File
@@ -22,6 +22,10 @@ const WRITABLE_SERVER_TYPES = new Set([
const ISMASTER_FIELDS = [
'minWireVersion',
'maxWireVersion',
'maxBsonObjectSize',
'maxMessageSizeBytes',
'maxWriteBatchSize',
'compression',
'me',
'hosts',
'passives',
@@ -31,7 +35,10 @@ const ISMASTER_FIELDS = [
'setVersion',
'electionId',
'primary',
'logicalSessionTimeoutMinutes'
'logicalSessionTimeoutMinutes',
'saslSupportedMechs',
'__nodejs_mock_server__',
'$clusterTime'
];
/**
@@ -62,7 +69,7 @@ class ServerDescription {
);
this.address = address;
this.error = null;
this.error = options.error || null;
this.roundTripTime = options.roundTripTime || 0;
this.lastUpdateTime = Date.now();
this.lastWriteDate = ismaster.lastWrite ? ismaster.lastWrite.lastWriteDate : null;
@@ -75,6 +82,7 @@ class ServerDescription {
});
// normalize case for hosts
if (this.me) this.me = this.me.toLowerCase();
this.hosts = this.hosts.map(host => host.toLowerCase());
this.passives = this.passives.map(host => host.toLowerCase());
this.arbiters = this.arbiters.map(host => host.toLowerCase());
+40 -2
View File
@@ -8,13 +8,24 @@ const MongoError = require('../error').MongoError;
const IDLE_WRITE_PERIOD = 10000;
const SMALLEST_MAX_STALENESS_SECONDS = 90;
/**
* Returns a server selector that selects for writable servers
*/
function writableServerSelector() {
return function(topologyDescription, servers) {
return latencyWindowReducer(topologyDescription, servers.filter(s => s.isWritable));
};
}
// reducers
/**
* Reduces the passed in array of servers by the rules of the "Max Staleness" specification
* found here: https://github.com/mongodb/specifications/blob/master/source/max-staleness/max-staleness.rst
*
* @param {ReadPreference} readPreference The read preference providing max staleness guidance
* @param {topologyDescription} topologyDescription The topology description
* @param {ServerDescription[]} servers The list of server descriptions to be reduced
* @return {ServerDescription[]} The list of servers that satisfy the requirements of max staleness
*/
function maxStalenessReducer(readPreference, topologyDescription, servers) {
if (readPreference.maxStalenessSeconds == null || readPreference.maxStalenessSeconds < 0) {
return servers;
@@ -24,7 +35,7 @@ function maxStalenessReducer(readPreference, topologyDescription, servers) {
const maxStalenessVariance =
(topologyDescription.heartbeatFrequencyMS + IDLE_WRITE_PERIOD) / 1000;
if (maxStaleness < maxStalenessVariance) {
throw MongoError(`maxStalenessSeconds must be at least ${maxStalenessVariance} seconds`);
throw new MongoError(`maxStalenessSeconds must be at least ${maxStalenessVariance} seconds`);
}
if (maxStaleness < SMALLEST_MAX_STALENESS_SECONDS) {
@@ -61,6 +72,12 @@ function maxStalenessReducer(readPreference, topologyDescription, servers) {
return servers;
}
/**
* Determines whether a server's tags match a given set of tags
*
* @param {String[]} tagSet The requested tag set to match
* @param {String[]} serverTags The server's tags
*/
function tagSetMatch(tagSet, serverTags) {
const keys = Object.keys(tagSet);
const serverTagKeys = Object.keys(serverTags);
@@ -74,6 +91,13 @@ function tagSetMatch(tagSet, serverTags) {
return true;
}
/**
* Reduces a set of server descriptions based on tags requested by the read preference
*
* @param {ReadPreference} readPreference The read preference providing the requested tags
* @param {ServerDescription[]} servers The list of server descriptions to reduce
* @return {ServerDescription[]} The list of servers matching the requested tags
*/
function tagSetReducer(readPreference, servers) {
if (
readPreference.tags == null ||
@@ -97,6 +121,15 @@ function tagSetReducer(readPreference, servers) {
return [];
}
/**
* Reduces a list of servers to ensure they fall within an acceptable latency window. This is
* further specified in the "Server Selection" specification, found here:
* https://github.com/mongodb/specifications/blob/master/source/server-selection/server-selection.rst
*
* @param {topologyDescription} topologyDescription The topology description
* @param {ServerDescription[]} servers The list of servers to reduce
* @returns {ServerDescription[]} The servers which fall within an acceptable latency window
*/
function latencyWindowReducer(topologyDescription, servers) {
const low = servers.reduce(
(min, server) => (min === -1 ? server.roundTripTime : Math.min(server.roundTripTime, min)),
@@ -128,6 +161,11 @@ function knownFilter(server) {
return server.type !== ServerType.Unknown;
}
/**
* Returns a function which selects servers based on a provided read preference
*
* @param {ReadPreference} readPreference The read preference to select with
*/
function readPreferenceServerSelector(readPreference) {
if (!readPreference.isValid()) {
throw new TypeError('Invalid read preference specified');
+523 -135
View File
@@ -1,22 +1,28 @@
'use strict';
const EventEmitter = require('events');
const ServerDescription = require('./server_description').ServerDescription;
const ServerType = require('./server_description').ServerType;
const TopologyDescription = require('./topology_description').TopologyDescription;
const TopologyType = require('./topology_description').TopologyType;
const monitoring = require('./monitoring');
const calculateDurationInMs = require('../utils').calculateDurationInMs;
const MongoTimeoutError = require('../error').MongoTimeoutError;
const MongoNetworkError = require('../error').MongoNetworkError;
const Server = require('./server');
const relayEvents = require('../utils').relayEvents;
const ReadPreference = require('../topologies/read_preference');
const readPreferenceServerSelector = require('./server_selectors').readPreferenceServerSelector;
const writableServerSelector = require('./server_selectors').writableServerSelector;
const isRetryableWritesSupported = require('../topologies/shared').isRetryableWritesSupported;
const Cursor = require('./cursor');
const Cursor = require('../cursor');
const deprecate = require('util').deprecate;
const BSON = require('../connection/utils').retrieveBSON();
const createCompressionInfo = require('../topologies/shared').createCompressionInfo;
const isRetryableError = require('../error').isRetryableError;
const MongoParseError = require('../error').MongoParseError;
const ClientSession = require('../sessions').ClientSession;
const createClientInfo = require('../topologies/shared').createClientInfo;
const MongoError = require('../error').MongoError;
const resolveClusterTime = require('../topologies/shared').resolveClusterTime;
// Global state
let globalTopologyCounter = 0;
@@ -26,9 +32,31 @@ const TOPOLOGY_DEFAULTS = {
localThresholdMS: 15,
serverSelectionTimeoutMS: 10000,
heartbeatFrequencyMS: 30000,
minHeartbeatIntervalMS: 500
minHeartbeatFrequencyMS: 500
};
// events that we relay to the `Topology`
const SERVER_RELAY_EVENTS = [
'serverHeartbeatStarted',
'serverHeartbeatSucceeded',
'serverHeartbeatFailed',
'commandStarted',
'commandSucceeded',
'commandFailed',
// NOTE: Legacy events
'monitoring'
];
// all events we listen to from `Server` instances
const LOCAL_SERVER_EVENTS = SERVER_RELAY_EVENTS.concat([
'error',
'connect',
'descriptionReceived',
'close',
'ended'
]);
/**
* A container of server instances representing a connection to a MongoDB topology.
*
@@ -54,7 +82,7 @@ class Topology extends EventEmitter {
*/
constructor(seedlist, options) {
super();
if (typeof options === 'undefined') {
if (typeof options === 'undefined' && typeof seedlist !== 'string') {
options = seedlist;
seedlist = [];
@@ -65,11 +93,16 @@ class Topology extends EventEmitter {
}
seedlist = seedlist || [];
if (typeof seedlist === 'string') {
seedlist = parseStringSeedlist(seedlist);
}
options = Object.assign({}, TOPOLOGY_DEFAULTS, options);
const topologyType = topologyTypeFromSeedlist(seedlist, options);
const topologyId = globalTopologyCounter++;
const serverDescriptions = seedlist.reduce((result, seed) => {
if (seed.domain_socket) seed.host = seed.domain_socket;
const address = seed.port ? `${seed.host}:${seed.port}` : `${seed.host}:27017`;
result.set(address, new ServerDescription(address));
return result;
@@ -79,7 +112,7 @@ class Topology extends EventEmitter {
// the id of this topology
id: topologyId,
// passed in options
options: Object.assign({}, options),
options,
// initial seedlist of servers to connect to
seedlist: seedlist,
// the topology description
@@ -89,6 +122,7 @@ class Topology extends EventEmitter {
options.replicaSet,
null,
null,
null,
options
),
serverSelectionTimeoutMS: options.serverSelectionTimeoutMS,
@@ -97,28 +131,24 @@ class Topology extends EventEmitter {
// allow users to override the cursor factory
Cursor: options.cursorFactory || Cursor,
// the bson parser
bson:
options.bson ||
new BSON([
BSON.Binary,
BSON.Code,
BSON.DBRef,
BSON.Decimal128,
BSON.Double,
BSON.Int32,
BSON.Long,
BSON.Map,
BSON.MaxKey,
BSON.MinKey,
BSON.ObjectId,
BSON.BSONRegExp,
BSON.Symbol,
BSON.Timestamp
])
bson: options.bson || new BSON(),
// a map of server instances to normalized addresses
servers: new Map(),
// Server Session Pool
sessionPool: null,
// Active client sessions
sessions: [],
// Promise library
promiseLibrary: options.promiseLibrary || Promise,
credentials: options.credentials,
clusterTime: null
};
// amend options for server instance creation
this.s.options.compression = { compressors: createCompressionInfo(options) };
// add client info
this.s.clientInfo = createClientInfo(options);
}
/**
@@ -128,6 +158,10 @@ class Topology extends EventEmitter {
return this.s.description;
}
get parserType() {
return BSON.native ? 'c++' : 'js';
}
/**
* All raw connections
* @method
@@ -144,8 +178,12 @@ class Topology extends EventEmitter {
*
* @param {Object} [options] Optional settings
* @param {Array} [options.auth=null] Array of auth options to apply on connect
* @param {function} [callback] An optional callback called once on the first connected server
*/
connect(/* options */) {
connect(options, callback) {
if (typeof options === 'function') (callback = options), (options = {});
options = options || {};
// emit SDAM monitoring events
this.emit('topologyOpening', new monitoring.TopologyOpeningEvent(this.s.id));
@@ -161,39 +199,126 @@ class Topology extends EventEmitter {
connectServers(this, Array.from(this.s.description.servers.values()));
this.s.connected = true;
// otherwise, wait for a server to properly connect based on user provided read preference,
// or primary.
translateReadPreference(options);
const readPreference = options.readPreference || ReadPreference.primary;
this.selectServer(readPreferenceServerSelector(readPreference), options, (err, server) => {
if (err) {
if (typeof callback === 'function') {
callback(err, null);
} else {
this.emit('error', err);
}
return;
}
const errorHandler = err => {
server.removeListener('connect', connectHandler);
if (typeof callback === 'function') callback(err, null);
};
const connectHandler = (_, err) => {
server.removeListener('error', errorHandler);
this.emit('open', err, this);
this.emit('connect', this);
if (typeof callback === 'function') callback(err, this);
};
const STATE_CONNECTING = 1;
if (server.s.state === STATE_CONNECTING) {
server.once('error', errorHandler);
server.once('connect', connectHandler);
return;
}
connectHandler();
});
}
/**
* Close this topology
*/
close(callback) {
// destroy all child servers
this.s.servers.forEach(server => server.destroy());
close(options, callback) {
if (typeof options === 'function') (callback = options), (options = {});
options = options || {};
// emit an event for close
this.emit('topologyClosed', new monitoring.TopologyClosedEvent(this.s.id));
this.s.connected = false;
if (typeof callback === 'function') {
callback(null, null);
if (this.s.sessionPool) {
this.s.sessions.forEach(session => session.endSession());
this.s.sessionPool.endAllPooledSessions();
}
const servers = this.s.servers;
if (servers.size === 0) {
this.s.connected = false;
if (typeof callback === 'function') {
callback(null, null);
}
return;
}
// destroy all child servers
let destroyed = 0;
servers.forEach(server =>
destroyServer(server, this, () => {
destroyed++;
if (destroyed === servers.size) {
// emit an event for close
this.emit('topologyClosed', new monitoring.TopologyClosedEvent(this.s.id));
this.s.connected = false;
if (typeof callback === 'function') {
callback(null, null);
}
}
})
);
}
/**
* Selects a server according to the selection predicate provided
*
* @param {function} [selector] An optional selector to select servers by, defaults to a random selection within a latency window
* @param {object} [options] Optional settings related to server selection
* @param {number} [options.serverSelectionTimeoutMS] How long to block for server selection before throwing an error
* @param {function} callback The callback used to indicate success or failure
* @return {Server} An instance of a `Server` meeting the criteria of the predicate provided
*/
selectServer(selector, options, callback) {
if (typeof options === 'function') (callback = options), (options = {});
if (typeof options === 'function') {
callback = options;
if (typeof selector !== 'function') {
options = selector;
translateReadPreference(options);
const readPreference = options.readPreference || ReadPreference.primary;
selector = readPreferenceServerSelector(readPreference);
} else {
options = {};
}
}
options = Object.assign(
{},
{ serverSelectionTimeoutMS: this.s.serverSelectionTimeoutMS },
options
);
const isSharded = this.description.type === TopologyType.Sharded;
const session = options.session;
const transaction = session && session.transaction;
if (isSharded && transaction && transaction.server) {
callback(null, transaction.server);
return;
}
selectServers(
this,
selector,
@@ -201,7 +326,56 @@ class Topology extends EventEmitter {
process.hrtime(),
(err, servers) => {
if (err) return callback(err, null);
callback(null, randomSelection(servers));
const selectedServer = randomSelection(servers);
if (isSharded && transaction && transaction.isActive) {
transaction.pinServer(selectedServer);
}
callback(null, selectedServer);
}
);
}
// Sessions related methods
/**
* @return Whether sessions are supported on the current topology
*/
hasSessionSupport() {
return this.description.logicalSessionTimeoutMinutes != null;
}
/**
* Start a logical session
*/
startSession(options, clientOptions) {
const session = new ClientSession(this, this.s.sessionPool, options, clientOptions);
session.once('ended', () => {
this.s.sessions = this.s.sessions.filter(s => !s.equals(session));
});
this.s.sessions.push(session);
return session;
}
/**
* Send endSessions command(s) with the given session ids
*
* @param {Array} sessions The sessions to end
* @param {function} [callback]
*/
endSessions(sessions, callback) {
if (!Array.isArray(sessions)) {
sessions = [sessions];
}
this.command(
'admin.$cmd',
{ endSessions: sessions },
{ readPreference: ReadPreference.primaryPreferred, noResponse: true },
() => {
// intentionally ignored, per spec
if (typeof callback === 'function') callback();
}
);
}
@@ -222,6 +396,10 @@ class Topology extends EventEmitter {
// first update the TopologyDescription
this.s.description = this.s.description.update(serverDescription);
if (this.s.description.compatibilityError) {
this.emit('error', new MongoError(this.s.description.compatibilityError));
return;
}
// emit monitoring events for this change
this.emit(
@@ -237,6 +415,17 @@ class Topology extends EventEmitter {
// update server list from updated descriptions
updateServers(this, serverDescription);
// Driver Sessions Spec: "Whenever a driver receives a cluster time from
// a server it MUST compare it to the current highest seen cluster time
// for the deployment. If the new cluster time is higher than the
// highest seen cluster time it MUST become the new highest seen cluster
// time. Two cluster times are compared using only the BsonTimestamp
// value of the clusterTime embedded field."
const clusterTime = serverDescription.$clusterTime;
if (clusterTime) {
resolveClusterTime(this, clusterTime);
}
this.emit(
'topologyDescriptionChanged',
new monitoring.TopologyDescriptionChangedEvent(
@@ -247,26 +436,13 @@ class Topology extends EventEmitter {
);
}
/**
* Authenticate using a specified mechanism
*
* @param {String} mechanism The auth mechanism used for authentication
* @param {String} db The db we are authenticating against
* @param {Object} options Optional settings for the authenticating mechanism
* @param {authResultCallback} callback A callback function
*/
auth(mechanism, db, options, callback) {
callback(null, null);
auth(credentials, callback) {
if (typeof credentials === 'function') (callback = credentials), (credentials = null);
if (typeof callback === 'function') callback(null, true);
}
/**
* Logout from a database
*
* @param {String} db The db we are logging out from
* @param {authResultCallback} callback A callback function
*/
logout(db, callback) {
callback(null, null);
logout(callback) {
if (typeof callback === 'function') callback(null, true);
}
// Basic operation support. Eventually this should be moved into command construction
@@ -341,14 +517,44 @@ class Topology extends EventEmitter {
(callback = options), (options = {}), (options = options || {});
}
const readPreference = options.readPreference ? options.readPreference : ReadPreference.primary;
this.selectServer(readPreferenceServerSelector(readPreference), (err, server) => {
translateReadPreference(options);
const readPreference = options.readPreference || ReadPreference.primary;
this.selectServer(readPreferenceServerSelector(readPreference), options, (err, server) => {
if (err) {
callback(err, null);
return;
}
server.command(ns, cmd, options, callback);
const willRetryWrite =
!options.retrying &&
!!options.retryWrites &&
options.session &&
isRetryableWritesSupported(this) &&
!options.session.inTransaction() &&
isWriteCommand(cmd);
const cb = (err, result) => {
if (!err) return callback(null, result);
if (!isRetryableError(err)) {
return callback(err);
}
if (willRetryWrite) {
const newOptions = Object.assign({}, options, { retrying: true });
return this.command(ns, cmd, newOptions, callback);
}
return callback(err);
};
// increment and assign txnNumber
if (willRetryWrite) {
options.session.incrementTransactionNumber();
options.willRetryWrite = willRetryWrite;
}
server.command(ns, cmd, options, cb);
});
}
@@ -372,20 +578,106 @@ class Topology extends EventEmitter {
options = options || {};
const topology = options.topology || this;
const CursorClass = options.cursorFactory || this.s.Cursor;
translateReadPreference(options);
return new CursorClass(this.s.bson, ns, cmd, options, topology, this.s.options);
}
get clientInfo() {
return this.s.clientInfo;
}
// Legacy methods for compat with old topology types
isConnected() {
// console.log('not implemented: `isConnected`');
return true;
}
isDestroyed() {
// console.log('not implemented: `isDestroyed`');
return false;
}
unref() {
console.log('not implemented: `unref`');
}
// NOTE: There are many places in code where we explicitly check the last isMaster
// to do feature support detection. This should be done any other way, but for
// now we will just return the first isMaster seen, which should suffice.
lastIsMaster() {
const serverDescriptions = Array.from(this.description.servers.values());
if (serverDescriptions.length === 0) return {};
const sd = serverDescriptions.filter(sd => sd.type !== ServerType.Unknown)[0];
const result = sd || { maxWireVersion: this.description.commonWireVersion };
return result;
}
get logicalSessionTimeoutMinutes() {
return this.description.logicalSessionTimeoutMinutes;
}
get bson() {
return this.s.bson;
}
}
Object.defineProperty(Topology.prototype, 'clusterTime', {
enumerable: true,
get: function() {
return this.s.clusterTime;
},
set: function(clusterTime) {
this.s.clusterTime = clusterTime;
}
});
// legacy aliases
Topology.prototype.destroy = deprecate(
Topology.prototype.close,
'destroy() is deprecated, please use close() instead'
);
const RETRYABLE_WRITE_OPERATIONS = ['findAndModify', 'insert', 'update', 'delete'];
function isWriteCommand(command) {
return RETRYABLE_WRITE_OPERATIONS.some(op => command[op]);
}
/**
* Destroys a server, and removes all event listeners from the instance
*
* @param {Server} server
*/
function destroyServer(server, topology, callback) {
LOCAL_SERVER_EVENTS.forEach(event => server.removeAllListeners(event));
server.destroy(() => {
topology.emit(
'serverClosed',
new monitoring.ServerClosedEvent(topology.s.id, server.description.address)
);
if (typeof callback === 'function') callback(null, null);
});
}
/**
* Parses a basic seedlist in string form
*
* @param {string} seedlist The seedlist to parse
*/
function parseStringSeedlist(seedlist) {
return seedlist.split(',').map(seed => ({
host: seed.split(':')[0],
port: seed.split(':')[1] || 27017
}));
}
function topologyTypeFromSeedlist(seedlist, options) {
if (seedlist.length === 1 && !options.replicaSet) return TopologyType.Single;
if (options.replicaSet) return TopologyType.ReplicaSetNoPrimary;
const replicaSet = options.replicaSet || options.setName || options.rs_name;
if (seedlist.length === 1 && !replicaSet) return TopologyType.Single;
if (replicaSet) return TopologyType.ReplicaSetNoPrimary;
return TopologyType.Unknown;
}
@@ -404,9 +696,43 @@ function randomSelection(array) {
* @param {function} callback The callback used to convey errors or the resultant servers
*/
function selectServers(topology, selector, timeout, start, callback) {
const duration = calculateDurationInMs(start);
if (duration >= timeout) {
return callback(new MongoTimeoutError(`Server selection timed out after ${timeout} ms`));
}
// ensure we are connected
if (!topology.s.connected) {
topology.connect();
// we want to make sure we're still within the requested timeout window
const failToConnectTimer = setTimeout(() => {
topology.removeListener('connect', connectHandler);
callback(new MongoTimeoutError('Server selection timed out waiting to connect'));
}, timeout - duration);
const connectHandler = () => {
clearTimeout(failToConnectTimer);
selectServers(topology, selector, timeout, process.hrtime(), callback);
};
topology.once('connect', connectHandler);
return;
}
// otherwise, attempt server selection
const serverDescriptions = Array.from(topology.description.servers.values());
let descriptions;
// support server selection by options with readPreference
if (typeof selector === 'object') {
const readPreference = selector.readPreference
? selector.readPreference
: ReadPreference.primary;
selector = readPreferenceServerSelector(readPreference);
}
try {
descriptions = selector
? selector(topology.description, serverDescriptions)
@@ -420,48 +746,56 @@ function selectServers(topology, selector, timeout, start, callback) {
return callback(null, servers);
}
const duration = calculateDurationInMs(start);
if (duration >= timeout) {
return callback(new MongoTimeoutError(`Server selection timed out after ${timeout} ms`));
}
const retrySelection = () => {
// ensure all server monitors attempt monitoring immediately
topology.s.servers.forEach(server => server.monitor());
// ensure all server monitors attempt monitoring soon
topology.s.servers.forEach(server => {
setTimeout(
() => server.monitor({ heartbeatFrequencyMS: topology.description.heartbeatFrequencyMS }),
TOPOLOGY_DEFAULTS.minHeartbeatFrequencyMS
);
});
const iterationTimer = setTimeout(() => {
callback(new MongoTimeoutError('Server selection timed out due to monitoring'));
}, topology.s.minHeartbeatIntervalMS);
topology.once('topologyDescriptionChanged', () => {
const descriptionChangedHandler = () => {
// successful iteration, clear the check timer
clearTimeout(iterationTimer);
if (topology.description.error) {
callback(topology.description.error, null);
return;
}
// topology description has changed due to monitoring, reattempt server selection
selectServers(topology, selector, timeout, start, callback);
});
};
};
// ensure we are connected
if (!topology.s.connected) {
topology.connect();
// we want to make sure we're still within the requested timeout window
const failToConnectTimer = setTimeout(() => {
callback(new MongoTimeoutError('Server selection timed out waiting to connect'));
const iterationTimer = setTimeout(() => {
topology.removeListener('topologyDescriptionChanged', descriptionChangedHandler);
callback(new MongoTimeoutError(`Server selection timed out after ${timeout} ms`));
}, timeout - duration);
topology.once('connect', () => {
clearTimeout(failToConnectTimer);
retrySelection();
});
return;
}
topology.once('topologyDescriptionChanged', descriptionChangedHandler);
};
retrySelection();
}
function createAndConnectServer(topology, serverDescription) {
topology.emit(
'serverOpening',
new monitoring.ServerOpeningEvent(topology.s.id, serverDescription.address)
);
const server = new Server(serverDescription, topology.s.options, topology);
relayEvents(server, topology, SERVER_RELAY_EVENTS);
server.once('connect', serverConnectEventHandler(server, topology));
server.on('descriptionReceived', topology.serverUpdateHandler.bind(topology));
server.on('error', serverErrorEventHandler(server, topology));
server.on('close', () => topology.emit('close', server));
server.connect();
return server;
}
/**
* Create `Server` instances for all initially known servers, connect them, and assign
* them to the passed in `Topology`.
@@ -471,53 +805,24 @@ function selectServers(topology, selector, timeout, start, callback) {
*/
function connectServers(topology, serverDescriptions) {
topology.s.servers = serverDescriptions.reduce((servers, serverDescription) => {
// publish an open event for each ServerDescription created
topology.emit(
'serverOpening',
new monitoring.ServerOpeningEvent(topology.s.id, serverDescription.address)
);
const server = new Server(serverDescription, topology.s.options);
relayEvents(server, topology, [
'serverHeartbeatStarted',
'serverHeartbeatSucceeded',
'serverHeartbeatFailed'
]);
server.on('descriptionReceived', topology.serverUpdateHandler.bind(topology));
server.on('connect', serverConnectEventHandler(server, topology));
const server = createAndConnectServer(topology, serverDescription);
servers.set(serverDescription.address, server);
server.connect();
return servers;
}, new Map());
}
function updateServers(topology, currentServerDescription) {
function updateServers(topology, incomingServerDescription) {
// update the internal server's description
if (topology.s.servers.has(currentServerDescription.address)) {
const server = topology.s.servers.get(currentServerDescription.address);
server.s.description = currentServerDescription;
if (topology.s.servers.has(incomingServerDescription.address)) {
const server = topology.s.servers.get(incomingServerDescription.address);
server.s.description = incomingServerDescription;
}
// add new servers for all descriptions we currently don't know about locally
for (const serverDescription of topology.description.servers.values()) {
if (!topology.s.servers.has(serverDescription.address)) {
topology.emit(
'serverOpening',
new monitoring.ServerOpeningEvent(topology.s.id, serverDescription.address)
);
const server = new Server(serverDescription, topology.s.options);
relayEvents(server, topology, [
'serverHeartbeatStarted',
'serverHeartbeatSucceeded',
'serverHeartbeatFailed'
]);
server.on('descriptionReceived', topology.serverUpdateHandler.bind(topology));
server.on('connect', serverConnectEventHandler(server, topology));
const server = createAndConnectServer(topology, serverDescription);
topology.s.servers.set(serverDescription.address, server);
server.connect();
}
}
@@ -531,15 +836,33 @@ function updateServers(topology, currentServerDescription) {
const server = topology.s.servers.get(serverAddress);
topology.s.servers.delete(serverAddress);
server.destroy(() =>
topology.emit('serverClosed', new monitoring.ServerClosedEvent(topology.s.id, serverAddress))
);
// prepare server for garbage collection
destroyServer(server, topology);
}
}
function serverConnectEventHandler(server, topology) {
return function(/* ismaster */) {
topology.emit('connect', topology);
return function(/* isMaster, err */) {
server.monitor({
initial: true,
heartbeatFrequencyMS: topology.description.heartbeatFrequencyMS
});
};
}
function serverErrorEventHandler(server, topology) {
return function(err) {
topology.emit(
'serverClosed',
new monitoring.ServerClosedEvent(topology.s.id, server.description.address)
);
if (err instanceof MongoParseError) {
resetServerState(server, err, { clearPool: true });
return;
}
resetServerState(server, err);
};
}
@@ -555,12 +878,12 @@ function executeWriteOperation(args, options, callback) {
const willRetryWrite =
!args.retrying &&
options.retryWrites &&
!!options.retryWrites &&
options.session &&
isRetryableWritesSupported(topology) &&
!options.session.inTransaction();
topology.selectServer(writableServerSelector(), (err, server) => {
topology.selectServer(writableServerSelector(), options, (err, server) => {
if (err) {
callback(err, null);
return;
@@ -568,7 +891,7 @@ function executeWriteOperation(args, options, callback) {
const handler = (err, result) => {
if (!err) return callback(null, result);
if (!(err instanceof MongoNetworkError) && !err.message.match(/not master/)) {
if (!isRetryableError(err)) {
return callback(err);
}
@@ -592,14 +915,58 @@ function executeWriteOperation(args, options, callback) {
// execute the write operation
server[op](ns, ops, options, handler);
// we need to increment the statement id if we're in a transaction
if (options.session && options.session.inTransaction()) {
options.session.incrementStatementId(ops.length);
}
});
}
/**
* Resets the internal state of this server to `Unknown` by simulating an empty ismaster
*
* @private
* @param {Server} server
* @param {MongoError} error The error that caused the state reset
* @param {object} [options] Optional settings
* @param {boolean} [options.clearPool=false] Pool should be cleared out on state reset
*/
function resetServerState(server, error, options) {
options = Object.assign({}, { clearPool: false }, options);
function resetState() {
server.emit(
'descriptionReceived',
new ServerDescription(server.description.address, null, { error })
);
}
if (options.clearPool && server.pool) {
server.pool.reset(() => resetState());
return;
}
resetState();
}
function translateReadPreference(options) {
if (options.readPreference == null) {
return;
}
let r = options.readPreference;
if (typeof r === 'string') {
options.readPreference = new ReadPreference(r);
} else if (r && !(r instanceof ReadPreference) && typeof r === 'object') {
const mode = r.mode || r.preference;
if (mode && typeof mode === 'string') {
options.readPreference = new ReadPreference(mode, r.tags, {
maxStalenessSeconds: r.maxStalenessSeconds
});
}
} else if (!(r instanceof ReadPreference)) {
throw new TypeError('Invalid read preference: ' + r);
}
return options;
}
/**
* A server opening SDAM monitoring event
*
@@ -663,4 +1030,25 @@ function executeWriteOperation(args, options, callback) {
* @type {ServerHeartbeatSucceededEvent}
*/
/**
* An event emitted indicating a command was started, if command monitoring is enabled
*
* @event Topology#commandStarted
* @type {object}
*/
/**
* An event emitted indicating a command succeeded, if command monitoring is enabled
*
* @event Topology#commandSucceeded
* @type {object}
*/
/**
* An event emitted indicating a command failed, if command monitoring is enabled
*
* @event Topology#commandFailed
* @type {object}
*/
module.exports = Topology;
+37 -20
View File
@@ -2,11 +2,13 @@
const ServerType = require('./server_description').ServerType;
const ServerDescription = require('./server_description').ServerDescription;
const ReadPreference = require('../topologies/read_preference');
const WIRE_CONSTANTS = require('../wireprotocol/constants');
// contstants related to compatability checks
const MIN_SUPPORTED_SERVER_VERSION = '2.6';
const MIN_SUPPORTED_WIRE_VERSION = 2;
const MAX_SUPPORTED_WIRE_VERSION = 5;
const MIN_SUPPORTED_SERVER_VERSION = WIRE_CONSTANTS.MIN_SUPPORTED_SERVER_VERSION;
const MAX_SUPPORTED_SERVER_VERSION = WIRE_CONSTANTS.MAX_SUPPORTED_SERVER_VERSION;
const MIN_SUPPORTED_WIRE_VERSION = WIRE_CONSTANTS.MIN_SUPPORTED_WIRE_VERSION;
const MAX_SUPPORTED_WIRE_VERSION = WIRE_CONSTANTS.MAX_SUPPORTED_WIRE_VERSION;
// An enumeration of topology types we know about
const TopologyType = {
@@ -28,7 +30,16 @@ class TopologyDescription {
* @param {number} maxSetVersion
* @param {ObjectId} maxElectionId
*/
constructor(topologyType, serverDescriptions, setName, maxSetVersion, maxElectionId, options) {
constructor(
topologyType,
serverDescriptions,
setName,
maxSetVersion,
maxElectionId,
commonWireVersion,
options,
error
) {
options = options || {};
// TODO: consider assigning all these values to a temporary value `s` which
@@ -46,14 +57,18 @@ class TopologyDescription {
this.heartbeatFrequencyMS = options.heartbeatFrequencyMS || 0;
this.localThresholdMS = options.localThresholdMS || 0;
this.options = options;
this.error = error;
this.commonWireVersion = commonWireVersion || null;
// determine server compatibility
for (const serverDescription of this.servers.values()) {
if (serverDescription.type === ServerType.Unknown) continue;
if (serverDescription.minWireVersion > MAX_SUPPORTED_WIRE_VERSION) {
this.compatible = false;
this.compatibilityError = `Server at ${serverDescription.address} requires wire version ${
serverDescription.minWireVersion
}, but this version of the driver only supports up to ${MAX_SUPPORTED_WIRE_VERSION}.`;
}, but this version of the driver only supports up to ${MAX_SUPPORTED_WIRE_VERSION} (MongoDB ${MAX_SUPPORTED_SERVER_VERSION})`;
}
if (serverDescription.maxWireVersion < MIN_SUPPORTED_WIRE_VERSION) {
@@ -78,19 +93,6 @@ class TopologyDescription {
}, null);
}
/**
* @returns The minimum reported wire version of all known servers
*/
get commonWireVersion() {
return Array.from(this.servers.values())
.filter(server => server.type !== ServerType.Unknown)
.reduce(
(min, server) =>
min == null ? server.maxWireVersion : Math.min(min, server.maxWireVersion),
null
);
}
/**
* Returns a copy of this description updated with a given ServerDescription
*
@@ -106,10 +108,21 @@ class TopologyDescription {
let setName = this.setName;
let maxSetVersion = this.maxSetVersion;
let maxElectionId = this.maxElectionId;
let commonWireVersion = this.commonWireVersion;
let error = serverDescription.error || null;
const serverType = serverDescription.type;
let serverDescriptions = new Map(this.servers);
// update common wire version
if (serverDescription.maxWireVersion !== 0) {
if (commonWireVersion == null) {
commonWireVersion = serverDescription.maxWireVersion;
} else {
commonWireVersion = Math.min(commonWireVersion, serverDescription.maxWireVersion);
}
}
// update the actual server description
serverDescriptions.set(address, serverDescription);
@@ -121,7 +134,9 @@ class TopologyDescription {
setName,
maxSetVersion,
maxElectionId,
this.options
commonWireVersion,
this.options,
error
);
}
@@ -201,7 +216,9 @@ class TopologyDescription {
setName,
maxSetVersion,
maxElectionId,
this.options
commonWireVersion,
this.options,
error
);
}