opusdev/vector-similarity-api
1
1import type { BSONSerializeOptions, Document } from '../bson';2import type { MongoCredentials } from '../cmap/auth/mongo_credentials';3import type { ConnectionEvents } from '../cmap/connection';4import type { ConnectionPoolEvents } from '../cmap/connection_pool';5import type { ClientMetadata } from '../cmap/handshake/client_metadata';6import { DEFAULT_OPTIONS } from '../connection_string';7import {8 CLOSE,9 CONNECT,10 ERROR,11 LOCAL_SERVER_EVENTS,12 OPEN,13 SERVER_CLOSED,14 SERVER_DESCRIPTION_CHANGED,15 SERVER_OPENING,16 SERVER_RELAY_EVENTS,17 TIMEOUT,18 TOPOLOGY_CLOSED,19 TOPOLOGY_DESCRIPTION_CHANGED,20 TOPOLOGY_OPENING21} from '../constants';22import {23 MongoCompatibilityError,24 type MongoDriverError,25 MongoError,26 MongoErrorLabel,27 MongoOperationTimeoutError,28 MongoRuntimeError,29 MongoServerSelectionError,30 MongoTopologyClosedError31} from '../error';32import type { MongoClient, ServerApi } from '../mongo_client';33import { MongoLoggableComponent, type MongoLogger, SeverityLevel } from '../mongo_logger';34import { type Abortable, TypedEventEmitter } from '../mongo_types';35import { ReadPreference, type ReadPreferenceLike } from '../read_preference';36import type { ClientSession } from '../sessions';37import { Timeout, TimeoutContext, TimeoutError } from '../timeout';38import type { Transaction } from '../transactions';39import {40 addAbortListener,41 type Callback,42 type EventEmitterWithState,43 HostAddress,44 kDispose,45 List,46 makeStateMachine,47 noop,48 now,49 promiseWithResolvers,50 shuffle51} from '../utils';52import {53 _advanceClusterTime,54 type ClusterTime,55 ServerType,56 STATE_CLOSED,57 STATE_CLOSING,58 STATE_CONNECTED,59 STATE_CONNECTING,60 TopologyType61} from './common';62import {63 ServerClosedEvent,64 ServerDescriptionChangedEvent,65 ServerOpeningEvent,66 TopologyClosedEvent,67 TopologyDescriptionChangedEvent,68 TopologyOpeningEvent69} from './events';70import type { ServerMonitoringMode } from './monitor';71import { Server, type ServerEvents, type ServerOptions } from './server';72import { compareTopologyVersion, ServerDescription } from './server_description';73import { readPreferenceServerSelector, type ServerSelector } from './server_selection';74import {75 ServerSelectionFailedEvent,76 ServerSelectionStartedEvent,77 ServerSelectionSucceededEvent,78 WaitingForSuitableServerEvent79} from './server_selection_events';80import { SrvPoller, type SrvPollingEvent } from './srv_polling';81import { TopologyDescription } from './topology_description';82 83// Global state84let globalTopologyCounter = 0;85 86const stateTransition = makeStateMachine({87 [STATE_CLOSED]: [STATE_CLOSED, STATE_CONNECTING],88 [STATE_CONNECTING]: [STATE_CONNECTING, STATE_CLOSING, STATE_CONNECTED, STATE_CLOSED],89 [STATE_CONNECTED]: [STATE_CONNECTED, STATE_CLOSING, STATE_CLOSED],90 [STATE_CLOSING]: [STATE_CLOSING, STATE_CLOSED]91});92 93/** @internal */94export type ServerSelectionCallback = Callback<Server>;95 96/** @internal */97export interface ServerSelectionRequest {98 serverSelector: ServerSelector;99 topologyDescription: TopologyDescription;100 mongoLogger: MongoLogger | undefined;101 transaction?: Transaction;102 startTime: number;103 resolve: (server: Server) => void;104 reject: (error: MongoError) => void;105 cancelled: boolean;106 operationName: string;107 waitingLogged: boolean;108 previousServer?: ServerDescription;109}110 111/** @internal */112export interface TopologyPrivate {113 /** the id of this topology */114 id: number;115 /** passed in options */116 options: TopologyOptions;117 /** initial seedlist of servers to connect to */118 seedlist: HostAddress[];119 /** initial state */120 state: string;121 /** the topology description */122 description: TopologyDescription;123 serverSelectionTimeoutMS: number;124 heartbeatFrequencyMS: number;125 minHeartbeatFrequencyMS: number;126 /** A map of server instances to normalized addresses */127 servers: Map<string, Server>;128 credentials?: MongoCredentials;129 clusterTime?: ClusterTime;130 131 /** related to srv polling */132 srvPoller?: SrvPoller;133 detectShardedTopology: (event: TopologyDescriptionChangedEvent) => void;134 detectSrvRecords: (event: SrvPollingEvent) => void;135}136 137/** @internal */138export interface TopologyOptions extends BSONSerializeOptions, ServerOptions {139 srvMaxHosts: number;140 srvServiceName: string;141 hosts: HostAddress[];142 retryWrites: boolean;143 retryReads: boolean;144 /** How long to block for server selection before throwing an error */145 serverSelectionTimeoutMS: number;146 /** The name of the replica set to connect to */147 replicaSet?: string;148 srvHost?: string;149 srvPoller?: SrvPoller;150 /** Indicates that a client should directly connect to a node without attempting to discover its topology type */151 directConnection: boolean;152 loadBalanced: boolean;153 metadata: ClientMetadata;154 extendedMetadata: Promise<Document>;155 serverMonitoringMode: ServerMonitoringMode;156 /** MongoDB server API version */157 serverApi?: ServerApi;158 __skipPingOnConnect?: boolean;159}160 161/** @public */162export interface ConnectOptions {163 readPreference?: ReadPreference;164}165 166/** @public */167export interface SelectServerOptions {168 readPreference?: ReadPreferenceLike;169 /** How long to block for server selection before throwing an error */170 serverSelectionTimeoutMS?: number;171 session?: ClientSession;172 operationName: string;173 previousServer?: ServerDescription;174 /**175 * @internal176 * TODO(NODE-6496): Make this required by making ChangeStream use LegacyTimeoutContext177 * */178 timeoutContext?: TimeoutContext;179}180 181/** @public */182export type TopologyEvents = {183 /** Top level MongoClient doesn't emit this so it is marked: @internal */184 connect(topology: Topology): void;185 serverOpening(event: ServerOpeningEvent): void;186 serverClosed(event: ServerClosedEvent): void;187 serverDescriptionChanged(event: ServerDescriptionChangedEvent): void;188 topologyClosed(event: TopologyClosedEvent): void;189 topologyOpening(event: TopologyOpeningEvent): void;190 topologyDescriptionChanged(event: TopologyDescriptionChangedEvent): void;191 error(error: Error): void;192 /** @internal */193 open(topology: Topology): void;194 close(): void;195 timeout(): void;196} & Omit<ServerEvents, 'connect'> &197 ConnectionPoolEvents &198 ConnectionEvents &199 EventEmitterWithState;200/**201 * A container of server instances representing a connection to a MongoDB topology.202 * @internal203 */204export class Topology extends TypedEventEmitter<TopologyEvents> {205 /** @internal */206 s: TopologyPrivate;207 /** @internal */208 waitQueue: List<ServerSelectionRequest>;209 /** @internal */210 hello?: Document;211 /** @internal */212 _type?: string;213 214 client!: MongoClient;215 216 /** @internal */217 private connectionLock?: Promise<Topology>;218 219 /** @event */220 static readonly SERVER_OPENING = SERVER_OPENING;221 /** @event */222 static readonly SERVER_CLOSED = SERVER_CLOSED;223 /** @event */224 static readonly SERVER_DESCRIPTION_CHANGED = SERVER_DESCRIPTION_CHANGED;225 /** @event */226 static readonly TOPOLOGY_OPENING = TOPOLOGY_OPENING;227 /** @event */228 static readonly TOPOLOGY_CLOSED = TOPOLOGY_CLOSED;229 /** @event */230 static readonly TOPOLOGY_DESCRIPTION_CHANGED = TOPOLOGY_DESCRIPTION_CHANGED;231 /** @event */232 static readonly ERROR = ERROR;233 /** @event */234 static readonly OPEN = OPEN;235 /** @event */236 static readonly CONNECT = CONNECT;237 /** @event */238 static readonly CLOSE = CLOSE;239 /** @event */240 static readonly TIMEOUT = TIMEOUT;241 242 /**243 * @param seedlist - a list of HostAddress instances to connect to244 */245 constructor(246 client: MongoClient,247 seeds: string | string[] | HostAddress | HostAddress[],248 options: TopologyOptions249 ) {250 super();251 this.on('error', noop);252 253 this.client = client;254 // Options should only be undefined in tests, MongoClient will always have defined options255 options = options ?? {256 hosts: [HostAddress.fromString('localhost:27017')],257 ...Object.fromEntries(DEFAULT_OPTIONS.entries())258 };259 260 if (typeof seeds === 'string') {261 seeds = [HostAddress.fromString(seeds)];262 } else if (!Array.isArray(seeds)) {263 seeds = [seeds];264 }265 266 const seedlist: HostAddress[] = [];267 for (const seed of seeds) {268 if (typeof seed === 'string') {269 seedlist.push(HostAddress.fromString(seed));270 } else if (seed instanceof HostAddress) {271 seedlist.push(seed);272 } else {273 // FIXME(NODE-3483): May need to be a MongoParseError274 throw new MongoRuntimeError(`Topology cannot be constructed from ${JSON.stringify(seed)}`);275 }276 }277 278 const topologyType = topologyTypeFromOptions(options);279 const topologyId = globalTopologyCounter++;280 281 const selectedHosts =282 options.srvMaxHosts == null ||283 options.srvMaxHosts === 0 ||284 options.srvMaxHosts >= seedlist.length285 ? seedlist286 : shuffle(seedlist, options.srvMaxHosts);287 288 const serverDescriptions = new Map();289 for (const hostAddress of selectedHosts) {290 serverDescriptions.set(hostAddress.toString(), new ServerDescription(hostAddress));291 }292 293 this.waitQueue = new List();294 this.s = {295 // the id of this topology296 id: topologyId,297 // passed in options298 options,299 // initial seedlist of servers to connect to300 seedlist,301 // initial state302 state: STATE_CLOSED,303 // the topology description304 description: new TopologyDescription(305 topologyType,306 serverDescriptions,307 options.replicaSet,308 undefined,309 undefined,310 undefined,311 options312 ),313 serverSelectionTimeoutMS: options.serverSelectionTimeoutMS,314 heartbeatFrequencyMS: options.heartbeatFrequencyMS,315 minHeartbeatFrequencyMS: options.minHeartbeatFrequencyMS,316 // a map of server instances to normalized addresses317 servers: new Map(),318 credentials: options?.credentials,319 clusterTime: undefined,320 321 detectShardedTopology: ev => this.detectShardedTopology(ev),322 detectSrvRecords: ev => this.detectSrvRecords(ev)323 };324 325 this.mongoLogger = client.mongoLogger;326 this.component = 'topology';327 328 if (options.srvHost && !options.loadBalanced) {329 this.s.srvPoller =330 options.srvPoller ??331 new SrvPoller({332 heartbeatFrequencyMS: this.s.heartbeatFrequencyMS,333 srvHost: options.srvHost,334 srvMaxHosts: options.srvMaxHosts,335 srvServiceName: options.srvServiceName336 });337 338 this.on(Topology.TOPOLOGY_DESCRIPTION_CHANGED, this.s.detectShardedTopology);339 }340 this.connectionLock = undefined;341 }342 343 private detectShardedTopology(event: TopologyDescriptionChangedEvent) {344 const previousType = event.previousDescription.type;345 const newType = event.newDescription.type;346 347 const transitionToSharded =348 previousType !== TopologyType.Sharded && newType === TopologyType.Sharded;349 const srvListeners = this.s.srvPoller?.listeners(SrvPoller.SRV_RECORD_DISCOVERY);350 const listeningToSrvPolling = !!srvListeners?.includes(this.s.detectSrvRecords);351 352 if (transitionToSharded && !listeningToSrvPolling) {353 this.s.srvPoller?.on(SrvPoller.SRV_RECORD_DISCOVERY, this.s.detectSrvRecords);354 this.s.srvPoller?.start();355 }356 }357 358 private detectSrvRecords(ev: SrvPollingEvent) {359 const previousTopologyDescription = this.s.description;360 this.s.description = this.s.description.updateFromSrvPollingEvent(361 ev,362 this.s.options.srvMaxHosts363 );364 if (this.s.description === previousTopologyDescription) {365 // Nothing changed, so return366 return;367 }368 369 updateServers(this);370 371 this.emitAndLog(372 Topology.TOPOLOGY_DESCRIPTION_CHANGED,373 new TopologyDescriptionChangedEvent(374 this.s.id,375 previousTopologyDescription,376 this.s.description377 )378 );379 }380 381 /**382 * @returns A `TopologyDescription` for this topology383 */384 get description(): TopologyDescription {385 return this.s.description;386 }387 388 get loadBalanced(): boolean {389 return this.s.options.loadBalanced;390 }391 392 get serverApi(): ServerApi | undefined {393 return this.s.options.serverApi;394 }395 396 get capabilities(): ServerCapabilities {397 return new ServerCapabilities(this.lastHello());398 }399 400 /** Initiate server connect */401 async connect(options?: ConnectOptions): Promise<Topology> {402 this.connectionLock ??= this._connect(options);403 try {404 await this.connectionLock;405 return this;406 } finally {407 this.connectionLock = undefined;408 }409 }410 411 private async _connect(options?: ConnectOptions): Promise<Topology> {412 options = options ?? {};413 if (this.s.state === STATE_CONNECTED) {414 return this;415 }416 417 stateTransition(this, STATE_CONNECTING);418 419 // emit SDAM monitoring events420 this.emitAndLog(Topology.TOPOLOGY_OPENING, new TopologyOpeningEvent(this.s.id));421 422 // emit an event for the topology change423 this.emitAndLog(424 Topology.TOPOLOGY_DESCRIPTION_CHANGED,425 new TopologyDescriptionChangedEvent(426 this.s.id,427 new TopologyDescription(TopologyType.Unknown), // initial is always Unknown428 this.s.description429 )430 );431 432 // connect all known servers, then attempt server selection to connect433 const serverDescriptions = Array.from(this.s.description.servers.values());434 this.s.servers = new Map(435 serverDescriptions.map(serverDescription => [436 serverDescription.address,437 createAndConnectServer(this, serverDescription)438 ])439 );440 441 // In load balancer mode we need to fake a server description getting442 // emitted from the monitor, since the monitor doesn't exist.443 if (this.s.options.loadBalanced) {444 for (const description of serverDescriptions) {445 const newDescription = new ServerDescription(description.hostAddress, undefined, {446 loadBalanced: this.s.options.loadBalanced447 });448 this.serverUpdateHandler(newDescription);449 }450 }451 452 const serverSelectionTimeoutMS = this.client.s.options.serverSelectionTimeoutMS;453 const readPreference = options.readPreference ?? ReadPreference.primary;454 const timeoutContext = TimeoutContext.create({455 // TODO(NODE-6448): auto-connect ignores timeoutMS; potential future feature456 timeoutMS: undefined,457 serverSelectionTimeoutMS,458 waitQueueTimeoutMS: this.client.s.options.waitQueueTimeoutMS459 });460 const selectServerOptions = {461 operationName: 'handshake',462 ...options,463 timeoutContext464 };465 466 try {467 const server = await this.selectServer(468 readPreferenceServerSelector(readPreference),469 selectServerOptions470 );471 472 const skipPingOnConnect = this.s.options.__skipPingOnConnect === true;473 if (!skipPingOnConnect && this.s.credentials) {474 const connection = await server.pool.checkOut({ timeoutContext: timeoutContext });475 server.pool.checkIn(connection);476 stateTransition(this, STATE_CONNECTED);477 this.emit(Topology.OPEN, this);478 this.emit(Topology.CONNECT, this);479 480 return this;481 }482 483 stateTransition(this, STATE_CONNECTED);484 this.emit(Topology.OPEN, this);485 this.emit(Topology.CONNECT, this);486 487 return this;488 } catch (error) {489 this.close();490 throw error;491 }492 }493 494 closeCheckedOutConnections() {495 for (const server of this.s.servers.values()) {496 return server.closeCheckedOutConnections();497 }498 }499 500 /** Close this topology */501 close(): void {502 if (this.s.state === STATE_CLOSED || this.s.state === STATE_CLOSING) {503 return;504 }505 506 for (const server of this.s.servers.values()) {507 closeServer(server, this);508 }509 510 this.s.servers.clear();511 512 stateTransition(this, STATE_CLOSING);513 514 drainWaitQueue(this.waitQueue, new MongoTopologyClosedError());515 516 if (this.s.srvPoller) {517 this.s.srvPoller.stop();518 this.s.srvPoller.removeListener(SrvPoller.SRV_RECORD_DISCOVERY, this.s.detectSrvRecords);519 }520 521 this.removeListener(Topology.TOPOLOGY_DESCRIPTION_CHANGED, this.s.detectShardedTopology);522 523 stateTransition(this, STATE_CLOSED);524 525 // emit an event for close526 this.emitAndLog(Topology.TOPOLOGY_CLOSED, new TopologyClosedEvent(this.s.id));527 }528 529 /**530 * Selects a server according to the selection predicate provided531 *532 * @param selector - An optional selector to select servers by, defaults to a random selection within a latency window533 * @param options - Optional settings related to server selection534 * @param callback - The callback used to indicate success or failure535 * @returns An instance of a `Server` meeting the criteria of the predicate provided536 */537 async selectServer(538 selector: string | ReadPreference | ServerSelector,539 options: SelectServerOptions & Abortable540 ): Promise<Server> {541 let serverSelector;542 if (typeof selector !== 'function') {543 if (typeof selector === 'string') {544 serverSelector = readPreferenceServerSelector(ReadPreference.fromString(selector));545 } else {546 let readPreference;547 if (selector instanceof ReadPreference) {548 readPreference = selector;549 } else {550 ReadPreference.translate(options);551 readPreference = options.readPreference || ReadPreference.primary;552 }553 554 serverSelector = readPreferenceServerSelector(readPreference as ReadPreference);555 }556 } else {557 serverSelector = selector;558 }559 560 options = { serverSelectionTimeoutMS: this.s.serverSelectionTimeoutMS, ...options };561 if (562 this.client.mongoLogger?.willLog(MongoLoggableComponent.SERVER_SELECTION, SeverityLevel.DEBUG)563 ) {564 this.client.mongoLogger?.debug(565 MongoLoggableComponent.SERVER_SELECTION,566 new ServerSelectionStartedEvent(selector, this.description, options.operationName)567 );568 }569 let timeout;570 if (options.timeoutContext) timeout = options.timeoutContext.serverSelectionTimeout;571 else {572 timeout = Timeout.expires(options.serverSelectionTimeoutMS ?? 0);573 }574 575 const isSharded = this.description.type === TopologyType.Sharded;576 const session = options.session;577 const transaction = session && session.transaction;578 579 if (isSharded && transaction && transaction.server) {580 if (581 this.client.mongoLogger?.willLog(582 MongoLoggableComponent.SERVER_SELECTION,583 SeverityLevel.DEBUG584 )585 ) {586 this.client.mongoLogger?.debug(587 MongoLoggableComponent.SERVER_SELECTION,588 new ServerSelectionSucceededEvent(589 selector,590 this.description,591 transaction.server.pool.address,592 options.operationName593 )594 );595 }596 if (options.timeoutContext?.clearServerSelectionTimeout) timeout?.clear();597 return transaction.server;598 }599 600 const { promise: serverPromise, resolve, reject } = promiseWithResolvers<Server>();601 602 const waitQueueMember: ServerSelectionRequest = {603 serverSelector,604 topologyDescription: this.description,605 mongoLogger: this.client.mongoLogger,606 transaction,607 resolve,608 reject,609 cancelled: false,610 startTime: now(),611 operationName: options.operationName,612 waitingLogged: false,613 previousServer: options.previousServer614 };615 616 const abortListener = addAbortListener(options.signal, function () {617 waitQueueMember.cancelled = true;618 reject(this.reason);619 });620 621 this.waitQueue.push(waitQueueMember);622 processWaitQueue(this);623 624 try {625 timeout?.throwIfExpired();626 const server = await (timeout ? Promise.race([serverPromise, timeout]) : serverPromise);627 if (options.timeoutContext?.csotEnabled() && server.description.minRoundTripTime !== 0) {628 options.timeoutContext.minRoundTripTime = server.description.minRoundTripTime;629 }630 return server;631 } catch (error) {632 if (TimeoutError.is(error)) {633 // Timeout634 waitQueueMember.cancelled = true;635 const timeoutError = new MongoServerSelectionError(636 `Server selection timed out after ${timeout?.duration} ms`,637 this.description638 );639 if (640 this.client.mongoLogger?.willLog(641 MongoLoggableComponent.SERVER_SELECTION,642 SeverityLevel.DEBUG643 )644 ) {645 this.client.mongoLogger?.debug(646 MongoLoggableComponent.SERVER_SELECTION,647 new ServerSelectionFailedEvent(648 selector,649 this.description,650 timeoutError,651 options.operationName652 )653 );654 }655 656 if (options.timeoutContext?.csotEnabled()) {657 throw new MongoOperationTimeoutError('Timed out during server selection', {658 cause: timeoutError659 });660 }661 throw timeoutError;662 }663 // Other server selection error664 throw error;665 } finally {666 abortListener?.[kDispose]();667 if (options.timeoutContext?.clearServerSelectionTimeout) timeout?.clear();668 }669 }670 /**671 * Update the internal TopologyDescription with a ServerDescription672 *673 * @param serverDescription - The server to update in the internal list of server descriptions674 */675 serverUpdateHandler(serverDescription: ServerDescription): void {676 if (!this.s.description.hasServer(serverDescription.address)) {677 return;678 }679 680 // ignore this server update if its from an outdated topologyVersion681 if (isStaleServerDescription(this.s.description, serverDescription)) {682 return;683 }684 685 // these will be used for monitoring events later686 const previousTopologyDescription = this.s.description;687 const previousServerDescription = this.s.description.servers.get(serverDescription.address);688 if (!previousServerDescription) {689 return;690 }691 692 // Driver Sessions Spec: "Whenever a driver receives a cluster time from693 // a server it MUST compare it to the current highest seen cluster time694 // for the deployment. If the new cluster time is higher than the695 // highest seen cluster time it MUST become the new highest seen cluster696 // time. Two cluster times are compared using only the BsonTimestamp697 // value of the clusterTime embedded field."698 const clusterTime = serverDescription.$clusterTime;699 if (clusterTime) {700 _advanceClusterTime(this, clusterTime);701 }702 703 // If we already know all the information contained in this updated description, then704 // we don't need to emit SDAM events, but still need to update the description, in order705 // to keep client-tracked attributes like last update time and round trip time up to date706 const equalDescriptions =707 previousServerDescription && previousServerDescription.equals(serverDescription);708 709 // first update the TopologyDescription710 this.s.description = this.s.description.update(serverDescription);711 if (this.s.description.compatibilityError) {712 this.emit(Topology.ERROR, new MongoCompatibilityError(this.s.description.compatibilityError));713 return;714 }715 716 // emit monitoring events for this change717 if (!equalDescriptions) {718 const newDescription = this.s.description.servers.get(serverDescription.address);719 if (newDescription) {720 this.emit(721 Topology.SERVER_DESCRIPTION_CHANGED,722 new ServerDescriptionChangedEvent(723 this.s.id,724 serverDescription.address,725 previousServerDescription,726 newDescription727 )728 );729 }730 }731 732 // update server list from updated descriptions733 updateServers(this, serverDescription);734 735 // attempt to resolve any outstanding server selection attempts736 if (this.waitQueue.length > 0) {737 processWaitQueue(this);738 }739 740 if (!equalDescriptions) {741 this.emitAndLog(742 Topology.TOPOLOGY_DESCRIPTION_CHANGED,743 new TopologyDescriptionChangedEvent(744 this.s.id,745 previousTopologyDescription,746 this.s.description747 )748 );749 }750 }751 752 auth(credentials?: MongoCredentials, callback?: Callback): void {753 if (typeof credentials === 'function') ((callback = credentials), (credentials = undefined));754 if (typeof callback === 'function') callback(undefined, true);755 }756 757 get clientMetadata(): ClientMetadata {758 return this.s.options.metadata;759 }760 761 isConnected(): boolean {762 return this.s.state === STATE_CONNECTED;763 }764 765 isDestroyed(): boolean {766 return this.s.state === STATE_CLOSED;767 }768 769 // NOTE: There are many places in code where we explicitly check the last hello770 // to do feature support detection. This should be done any other way, but for771 // now we will just return the first hello seen, which should suffice.772 lastHello(): Document {773 const serverDescriptions = Array.from(this.description.servers.values());774 if (serverDescriptions.length === 0) return {};775 const sd = serverDescriptions.filter(776 (sd: ServerDescription) => sd.type !== ServerType.Unknown777 )[0];778 779 const result = sd || { maxWireVersion: this.description.commonWireVersion };780 return result;781 }782 783 get commonWireVersion(): number | undefined {784 return this.description.commonWireVersion;785 }786 787 get logicalSessionTimeoutMinutes(): number | null {788 return this.description.logicalSessionTimeoutMinutes;789 }790 791 get clusterTime(): ClusterTime | undefined {792 return this.s.clusterTime;793 }794 795 set clusterTime(clusterTime: ClusterTime | undefined) {796 this.s.clusterTime = clusterTime;797 }798}799 800/** Destroys a server, and removes all event listeners from the instance */801function closeServer(server: Server, topology: Topology) {802 for (const event of LOCAL_SERVER_EVENTS) {803 server.removeAllListeners(event);804 }805 806 server.close();807 topology.emitAndLog(808 Topology.SERVER_CLOSED,809 new ServerClosedEvent(topology.s.id, server.description.address)810 );811 812 for (const event of SERVER_RELAY_EVENTS) {813 server.removeAllListeners(event);814 }815}816 817/** Predicts the TopologyType from options */818function topologyTypeFromOptions(options?: TopologyOptions) {819 if (options?.directConnection) {820 return TopologyType.Single;821 }822 823 if (options?.replicaSet) {824 return TopologyType.ReplicaSetNoPrimary;825 }826 827 if (options?.loadBalanced) {828 return TopologyType.LoadBalanced;829 }830 831 return TopologyType.Unknown;832}833 834/**835 * Creates new server instances and attempts to connect them836 *837 * @param topology - The topology that this server belongs to838 * @param serverDescription - The description for the server to initialize and connect to839 */840function createAndConnectServer(topology: Topology, serverDescription: ServerDescription) {841 topology.emitAndLog(842 Topology.SERVER_OPENING,843 new ServerOpeningEvent(topology.s.id, serverDescription.address)844 );845 846 const server = new Server(topology, serverDescription, topology.s.options);847 for (const event of SERVER_RELAY_EVENTS) {848 server.on(event, (e: any) => topology.emit(event, e));849 }850 851 server.on(Server.DESCRIPTION_RECEIVED, description => topology.serverUpdateHandler(description));852 853 server.connect();854 return server;855}856 857/**858 * @param topology - Topology to update.859 * @param incomingServerDescription - New server description.860 */861function updateServers(topology: Topology, incomingServerDescription?: ServerDescription) {862 // update the internal server's description863 if (incomingServerDescription && topology.s.servers.has(incomingServerDescription.address)) {864 const server = topology.s.servers.get(incomingServerDescription.address);865 if (server) {866 server.s.description = incomingServerDescription;867 if (868 incomingServerDescription.error instanceof MongoError &&869 incomingServerDescription.error.hasErrorLabel(MongoErrorLabel.ResetPool)870 ) {871 const interruptInUseConnections = incomingServerDescription.error.hasErrorLabel(872 MongoErrorLabel.InterruptInUseConnections873 );874 875 server.pool.clear({ interruptInUseConnections });876 } else if (incomingServerDescription.error == null) {877 const newTopologyType = topology.s.description.type;878 const shouldMarkPoolReady =879 incomingServerDescription.isDataBearing ||880 (incomingServerDescription.type !== ServerType.Unknown &&881 newTopologyType === TopologyType.Single);882 if (shouldMarkPoolReady) {883 server.pool.ready();884 }885 }886 }887 }888 889 // add new servers for all descriptions we currently don't know about locally890 for (const serverDescription of topology.description.servers.values()) {891 if (!topology.s.servers.has(serverDescription.address)) {892 const server = createAndConnectServer(topology, serverDescription);893 topology.s.servers.set(serverDescription.address, server);894 }895 }896 897 // for all servers no longer known, remove their descriptions and destroy their instances898 for (const entry of topology.s.servers) {899 const serverAddress = entry[0];900 if (topology.description.hasServer(serverAddress)) {901 continue;902 }903 904 if (!topology.s.servers.has(serverAddress)) {905 continue;906 }907 908 const server = topology.s.servers.get(serverAddress);909 topology.s.servers.delete(serverAddress);910 911 // prepare server for garbage collection912 if (server) {913 closeServer(server, topology);914 }915 }916}917 918function drainWaitQueue(queue: List<ServerSelectionRequest>, drainError: MongoDriverError) {919 while (queue.length) {920 const waitQueueMember = queue.shift();921 if (!waitQueueMember) {922 continue;923 }924 925 if (!waitQueueMember.cancelled) {926 if (927 waitQueueMember.mongoLogger?.willLog(928 MongoLoggableComponent.SERVER_SELECTION,929 SeverityLevel.DEBUG930 )931 ) {932 waitQueueMember.mongoLogger?.debug(933 MongoLoggableComponent.SERVER_SELECTION,934 new ServerSelectionFailedEvent(935 waitQueueMember.serverSelector,936 waitQueueMember.topologyDescription,937 drainError,938 waitQueueMember.operationName939 )940 );941 }942 waitQueueMember.reject(drainError);943 }944 }945}946 947function processWaitQueue(topology: Topology) {948 if (topology.s.state === STATE_CLOSED) {949 drainWaitQueue(topology.waitQueue, new MongoTopologyClosedError());950 return;951 }952 953 const isSharded = topology.description.type === TopologyType.Sharded;954 const serverDescriptions = Array.from(topology.description.servers.values());955 const membersToProcess = topology.waitQueue.length;956 for (let i = 0; i < membersToProcess; ++i) {957 const waitQueueMember = topology.waitQueue.shift();958 if (!waitQueueMember) {959 continue;960 }961 962 if (waitQueueMember.cancelled) {963 continue;964 }965 966 let selectedDescriptions;967 try {968 const serverSelector = waitQueueMember.serverSelector;969 const previousServer = waitQueueMember.previousServer;970 selectedDescriptions = serverSelector971 ? serverSelector(972 topology.description,973 serverDescriptions,974 previousServer ? [previousServer] : []975 )976 : serverDescriptions;977 } catch (selectorError) {978 if (979 topology.client.mongoLogger?.willLog(980 MongoLoggableComponent.SERVER_SELECTION,981 SeverityLevel.DEBUG982 )983 ) {984 topology.client.mongoLogger?.debug(985 MongoLoggableComponent.SERVER_SELECTION,986 new ServerSelectionFailedEvent(987 waitQueueMember.serverSelector,988 topology.description,989 selectorError,990 waitQueueMember.operationName991 )992 );993 }994 waitQueueMember.reject(selectorError);995 continue;996 }997 998 let selectedServer: Server | undefined;999 if (selectedDescriptions.length === 0) {1000 if (!waitQueueMember.waitingLogged) {1001 if (1002 topology.client.mongoLogger?.willLog(1003 MongoLoggableComponent.SERVER_SELECTION,1004 SeverityLevel.INFORMATIONAL1005 )1006 ) {1007 topology.client.mongoLogger?.info(1008 MongoLoggableComponent.SERVER_SELECTION,1009 new WaitingForSuitableServerEvent(1010 waitQueueMember.serverSelector,1011 topology.description,1012 topology.s.serverSelectionTimeoutMS !== 01013 ? topology.s.serverSelectionTimeoutMS - (now() - waitQueueMember.startTime)1014 : -1,1015 waitQueueMember.operationName1016 )1017 );1018 }1019 waitQueueMember.waitingLogged = true;1020 }1021 topology.waitQueue.push(waitQueueMember);1022 continue;1023 } else if (selectedDescriptions.length === 1) {1024 selectedServer = topology.s.servers.get(selectedDescriptions[0].address);1025 } else {1026 const descriptions = shuffle(selectedDescriptions, 2);1027 const server1 = topology.s.servers.get(descriptions[0].address);1028 const server2 = topology.s.servers.get(descriptions[1].address);1029 1030 selectedServer =1031 server1 && server2 && server1.s.operationCount < server2.s.operationCount1032 ? server11033 : server2;1034 }1035 1036 if (!selectedServer) {1037 const serverSelectionError = new MongoServerSelectionError(1038 'server selection returned a server description but the server was not found in the topology',1039 topology.description1040 );1041 if (1042 topology.client.mongoLogger?.willLog(1043 MongoLoggableComponent.SERVER_SELECTION,1044 SeverityLevel.DEBUG1045 )1046 ) {1047 topology.client.mongoLogger?.debug(1048 MongoLoggableComponent.SERVER_SELECTION,1049 new ServerSelectionFailedEvent(1050 waitQueueMember.serverSelector,1051 topology.description,1052 serverSelectionError,1053 waitQueueMember.operationName1054 )1055 );1056 }1057 waitQueueMember.reject(serverSelectionError);1058 return;1059 }1060 const transaction = waitQueueMember.transaction;1061 if (isSharded && transaction && transaction.isActive && selectedServer) {1062 transaction.pinServer(selectedServer);1063 }1064 1065 if (1066 topology.client.mongoLogger?.willLog(1067 MongoLoggableComponent.SERVER_SELECTION,1068 SeverityLevel.DEBUG1069 )1070 ) {1071 topology.client.mongoLogger?.debug(1072 MongoLoggableComponent.SERVER_SELECTION,1073 new ServerSelectionSucceededEvent(1074 waitQueueMember.serverSelector,1075 waitQueueMember.topologyDescription,1076 selectedServer.pool.address,1077 waitQueueMember.operationName1078 )1079 );1080 }1081 waitQueueMember.resolve(selectedServer);1082 }1083 1084 if (topology.waitQueue.length > 0) {1085 // ensure all server monitors attempt monitoring soon1086 for (const [, server] of topology.s.servers) {1087 process.nextTick(function scheduleServerCheck() {1088 return server.requestCheck();1089 });1090 }1091 }1092}1093 1094function isStaleServerDescription(1095 topologyDescription: TopologyDescription,1096 incomingServerDescription: ServerDescription1097) {1098 const currentServerDescription = topologyDescription.servers.get(1099 incomingServerDescription.address1100 );1101 const currentTopologyVersion = currentServerDescription?.topologyVersion;1102 return (1103 compareTopologyVersion(currentTopologyVersion, incomingServerDescription.topologyVersion) > 01104 );1105}1106 1107/**1108 * @public1109 * @deprecated This class will be removed as dead code in the next major version.1110 */1111export class ServerCapabilities {1112 maxWireVersion: number;1113 minWireVersion: number;1114 1115 constructor(hello: Document) {1116 this.minWireVersion = hello.minWireVersion || 0;1117 this.maxWireVersion = hello.maxWireVersion || 0;1118 }1119 1120 get hasAggregationCursor(): boolean {1121 return true;1122 }1123 1124 get hasWriteCommands(): boolean {1125 return true;1126 }1127 get hasTextSearch(): boolean {1128 return true;1129 }1130 1131 get hasAuthCommands(): boolean {1132 return true;1133 }1134 1135 get hasListCollectionsCommand(): boolean {1136 return true;1137 }1138 1139 get hasListIndexesCommand(): boolean {1140 return true;1141 }1142 1143 get supportsSnapshotReads(): boolean {1144 return this.maxWireVersion >= 13;1145 }1146 1147 get commandsTakeWriteConcern(): boolean {1148 return true;1149 }1150 1151 get commandsTakeCollation(): boolean {1152 return true;1153 }1154}1155 