diff --git a/.evergreen/generated_configs/variants.yml b/.evergreen/generated_configs/variants.yml index e47af42b83..8e83f00d65 100644 --- a/.evergreen/generated_configs/variants.yml +++ b/.evergreen/generated_configs/variants.yml @@ -470,11 +470,28 @@ buildvariants: # Otel tests - name: otel-rhel8 + tasks: + - name: .test-non-standard .replica_set-noauth-ssl .server-4.4 .python-3.11 + - name: .test-non-standard .standalone-noauth-nossl .server-5.0 .python-3.13 + - name: .test-non-standard .sharded_cluster-auth-ssl .server-6.0 .python-pypy3.11 + - name: .test-non-standard .replica_set-noauth-ssl .server-7.0 .python-3.12 + - name: .test-non-standard .replica_set-noauth-ssl .server-8.0 .python-3.14 + - name: .test-non-standard .standalone-noauth-nossl .server-9.0 .python-3.15 + - name: .test-non-standard .standalone-noauth-nossl .server-rapid .python-3.12 + - name: .test-non-standard .sharded_cluster-auth-ssl .server-latest .python-3.15 + display_name: OTel RHEL8 + run_on: + - rhel8.10-small + expansions: + TEST_NAME: otel + COVERAGE: "1" + tags: [pr] + - name: otel-server-trace-rhel8 tasks: - name: .test-non-standard .replica_set-noauth-ssl .server-latest - name: .test-non-standard .sharded_cluster-auth-ssl .server-latest .python-pypy3.11 - name: .test-non-standard .standalone-noauth-nossl .server-latest .python-3.10 - display_name: OTel RHEL8 + display_name: OTel server trace RHEL8 run_on: - rhel8.10-small expansions: diff --git a/.evergreen/scripts/generate_config.py b/.evergreen/scripts/generate_config.py index e3a144a741..d6282238a0 100644 --- a/.evergreen/scripts/generate_config.py +++ b/.evergreen/scripts/generate_config.py @@ -516,19 +516,43 @@ def create_doctests_variants(): def create_otel_variants(): host = DEFAULT_HOST # Merge otel's coverage into the combined report; see setup_tests.py's COVERAGE handling. - # OTEL=1 makes drivers-evergreen-tools enable the server's OpenTelemetry file exporter - # and export OTEL_TRACE_DIR, which TestServerTraceContext requires. - expansions = dict(TEST_NAME="otel", COVERAGE="1", OTEL="1") + # One task per supported server version: rotate the main three topology/auth/ssl combos + # and the CPython versions across them. OTEL=1 (the server's OpenTelemetry file exporter) + # is not set here: it requires MongoDB 9.0+ and this variant must also run on older + # versions. The client-side tracing tests still run without it; the server + # trace-context tests run in the dedicated variant below. + version_matrix = { + "4.4": (".replica_set-noauth-ssl", ".python-3.11"), + "5.0": (".standalone-noauth-nossl", ".python-3.13"), + "6.0": (".sharded_cluster-auth-ssl", ".python-pypy3.11"), + "7.0": (".replica_set-noauth-ssl", ".python-3.12"), + "8.0": (".replica_set-noauth-ssl", ".python-3.14"), + "9.0": (".standalone-noauth-nossl", ".python-3.15"), + "rapid": (".standalone-noauth-nossl", ".python-3.12"), + "latest": (".sharded_cluster-auth-ssl", ".python-3.15"), + } + + def task_for(version): + combo, python = version_matrix[version] + return f".test-non-standard {combo} .server-{version} {python}" + return [ + create_variant( + [task_for(version) for version in version_matrix], + get_variant_name("OTel", host), + host=host, + tags=["pr"], + expansions=dict(TEST_NAME="otel", COVERAGE="1"), + ), + # OTEL=1 makes drivers-evergreen-tools enable the server's OpenTelemetry file + # exporter and export OTEL_TRACE_DIR, which TestServerTraceContext requires. The + # exporter requires MongoDB 9.0+ and a binary that accepts every OTel setParameter; + # only the latest nightly qualifies (the v9.0 nightly rejects + # openTelemetryTracingFileFlushCount), so only latest tasks are selected. Keep the + # same three tasks the single-variant config ran before the version sweep was added. create_variant( [ # All three topologies, one task each to keep the task count small. - # - # OTEL=1 enables the server's OpenTelemetry file exporter, which - # requires MongoDB 9.0+ and a binary that accepts every OTel - # setParameter; only the latest nightly qualifies (the v9.0 - # nightly rejects openTelemetryTracingFileFlushCount), so only - # latest tasks are selected. ".test-non-standard .replica_set-noauth-ssl .server-latest", # Sharded adds mongos, which rewrites commands and reports a # different server.address, plus auth and ssl, which exercise @@ -539,11 +563,11 @@ def create_otel_variants(): # opentelemetry-api down to the floor in requirements/. ".test-non-standard .standalone-noauth-nossl .server-latest .python-3.10", ], - get_variant_name("OTel", host), + get_variant_name("OTel server trace", host), host=host, tags=["pr"], - expansions=expansions, - ) + expansions=dict(TEST_NAME="otel", COVERAGE="1", OTEL="1"), + ), ] diff --git a/.github/workflows/test-python.yml b/.github/workflows/test-python.yml index e1a4cdbf02..32d8291081 100644 --- a/.github/workflows/test-python.yml +++ b/.github/workflows/test-python.yml @@ -318,3 +318,31 @@ jobs: run: TEST_NAME=otel just setup-tests - name: Run tests run: just run-tests + + # The otel tests above run without authentication; this variant runs the same + # suite against an authenticated replica set, exercising sensitive-command + # redaction and authentication-related spans and log messages. + otel-auth: + runs-on: ubuntu-latest + name: OTel Tests (auth replica set) + steps: + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + with: + persist-credentials: false + - name: Install uv + uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0 + with: + enable-cache: true + python-version: "3.11" + - id: setup-mongodb + uses: mongodb-labs/drivers-evergreen-tools@c70d29e2102cd048ad800932209776499553870c # v1.0.1 + with: + version: "8.0" + topology: "replica_set" + auth: "auth" + - name: Install just + run: uv tool install rust-just + - name: Setup tests + run: TEST_NAME=otel AUTH=auth just setup-tests + - name: Run tests + run: just run-tests diff --git a/pymongo/_otel.py b/pymongo/_otel.py index 744769091b..38c0bb29d9 100644 --- a/pymongo/_otel.py +++ b/pymongo/_otel.py @@ -64,6 +64,14 @@ "_CURRENT_OPERATION_NAME", default=None ) +# Whether the active operation span already carries its namespace attributes, +# set eagerly by start_operation_span. Sibling to _CURRENT_OPERATION_NAME: set +# and reset at the same sites, so start_command_span knows the namespace is +# already initialized without inspecting the span. The OpenTelemetry Span API +# does not expose the span's attributes, so an implementation-specific property +# cannot be relied on here. +_NAMESPACE_INITIALIZED: ContextVar[bool] = ContextVar("_NAMESPACE_INITIALIZED", default=False) + # True while the driver is draining a cursor of its own to build one public API # call's return value. See internal_cursor_iteration. _INTERNAL_CURSOR_ITERATION: ContextVar[bool] = ContextVar( @@ -289,17 +297,8 @@ def _build_query_summary(command_name: str, dbname: str, collection: Optional[st return f"{command_name} {dbname}" -# LUT of db.operation.name where the spec's name differs from our `_Op` value; -# the nested command span still reports the wire name in db.command.name. -# DRIVERS-3625 and PYTHON-6054 cover the table's gaps and inconsistencies. -_OPERATION_NAME_OVERRIDES = { - "drop": "dropCollection", - "create": "createCollection", - "dropSearchIndexes": "dropSearchIndex", -} - -# The spec names anything sent through the generic `Database.command()` API "runCommand", -# not after the command it carries. +# The spec names anything sent through the generic `Database.command()` API +# (and its cursor-returning sibling) "runCommand", not after the command it carries. _RUN_COMMAND_OPERATION_NAME = "runCommand" @@ -320,8 +319,7 @@ def _build_operation_name(operation: Any, is_run_command: bool = False) -> str: """Return the ``db.operation.name`` the spec wants for this operation.""" if is_run_command: return _RUN_COMMAND_OPERATION_NAME - name = _normalize_operation_name(operation) - return _OPERATION_NAME_OVERRIDES.get(name, name) + return _normalize_operation_name(operation) def _is_sensitive_command(command_name: str, speculative_hello: bool) -> bool: @@ -376,12 +374,16 @@ def start_command_span( if current_operation is not None: current_span = trace.get_current_span() if current_span.is_recording(): - summary = _build_query_summary(current_operation, dbname, collection) - current_span.update_name(summary) - current_span.set_attribute("db.namespace", dbname) - current_span.set_attribute("db.operation.summary", summary) - if collection: - current_span.set_attribute("db.collection.name", collection) + # Call sites that know the operation's namespace up front already set it + # (e.g. a rename targets a collection while its command runs against + # admin), so only backfill what the span has not learned yet. + if not _NAMESPACE_INITIALIZED.get(): + summary = _build_query_summary(current_operation, dbname, collection) + current_span.update_name(summary) + current_span.set_attribute("db.namespace", dbname) + current_span.set_attribute("db.operation.summary", summary) + if collection: + current_span.set_attribute("db.collection.name", collection) if sent_cursor_id: current_span.set_attribute("db.mongodb.cursor_id", sent_cursor_id) @@ -591,9 +593,21 @@ class _OperationSpanHandle: ``_cm`` is the ``start_as_current_span`` context manager, or ``None`` in detached mode, where ``use_operation_span`` makes the span current per use. + ``namespace_initialized`` records whether the span was created with eager + namespace attributes, so ``start_command_span`` skips its backfill without + inspecting the span (the OpenTelemetry Span API does not expose attributes). + ``_ns_token`` is the ``_NAMESPACE_INITIALIZED`` token to reset when the + span stops being current, ``None`` in detached mode. """ - __slots__ = ("_cm", "_name_token", "operation_name", "span") + __slots__ = ( + "_cm", + "_name_token", + "_ns_token", + "namespace_initialized", + "operation_name", + "span", + ) def __init__( self, @@ -601,11 +615,15 @@ def __init__( cm: Any, name_token: Any, operation_name: str, + namespace_initialized: bool = False, + ns_token: Any = None, ) -> None: self.span = span self._cm = cm self._name_token = name_token self.operation_name = operation_name + self.namespace_initialized = namespace_initialized + self._ns_token = ns_token def start_operation_span( @@ -656,7 +674,9 @@ def start_operation_span( span = _TRACER.start_span( name, kind=SpanKind.CLIENT, context=context, attributes=attributes ) - return _OperationSpanHandle(span, None, None, operation) + return _OperationSpanHandle( + span, None, None, operation, namespace_initialized=dbname is not None + ) cm = _TRACER.start_as_current_span( name, kind=SpanKind.CLIENT, @@ -665,7 +685,10 @@ def start_operation_span( ) span = cm.__enter__() name_token = _CURRENT_OPERATION_NAME.set(operation) - return _OperationSpanHandle(span, cm, name_token, operation) + ns_token = _NAMESPACE_INITIALIZED.set(dbname is not None) + return _OperationSpanHandle( + span, cm, name_token, operation, namespace_initialized=dbname is not None, ns_token=ns_token + ) @contextlib.contextmanager @@ -679,6 +702,7 @@ def use_operation_span(handle: Optional[_OperationSpanHandle]) -> Iterator[None] yield return token = _CURRENT_OPERATION_NAME.set(handle.operation_name) + ns_token = _NAMESPACE_INITIALIZED.set(handle.namespace_initialized) try: # Left on, these would auto-record any exception leaving the block and # set ERROR status, duplicating the caller's end_operation_span_failure @@ -691,6 +715,7 @@ def use_operation_span(handle: Optional[_OperationSpanHandle]) -> Iterator[None] ): yield finally: + _NAMESPACE_INITIALIZED.reset(ns_token) _CURRENT_OPERATION_NAME.reset(token) @@ -706,6 +731,7 @@ def reset_context() -> None: if not _HAS_OPENTELEMETRY: return _CURRENT_OPERATION_NAME.set(None) + _NAMESPACE_INITIALIZED.set(False) context.attach(context.Context()) @@ -717,6 +743,7 @@ def end_operation_span_success(handle: Optional[_OperationSpanHandle]) -> None: handle.span.end() return _CURRENT_OPERATION_NAME.reset(handle._name_token) + _NAMESPACE_INITIALIZED.reset(handle._ns_token) handle._cm.__exit__(None, None, None) @@ -736,6 +763,7 @@ def end_operation_span_failure(handle: Optional[_OperationSpanHandle], exc: Base handle.span.end() else: _CURRENT_OPERATION_NAME.reset(handle._name_token) + _NAMESPACE_INITIALIZED.reset(handle._ns_token) handle._cm.__exit__(None, None, None) diff --git a/pymongo/_telemetry.py b/pymongo/_telemetry.py index 208f107955..67d58f4e8b 100644 --- a/pymongo/_telemetry.py +++ b/pymongo/_telemetry.py @@ -32,7 +32,9 @@ _CommandStatusMessage, _ConnectionStatusMessage, _debug_log, + _info_log, _is_debug_enabled, + _is_info_enabled, _SDAMStatusMessage, _ServerSelectionStatusMessage, _verbose_connection_error_reason, @@ -808,9 +810,11 @@ def _emit_log( self, message: _ServerSelectionStatusMessage, topology_description: TopologyDescription, + info: bool = False, **extra: Any, ) -> None: - _debug_log( + log_fn = _info_log if info else _debug_log + log_fn( _SERVER_SELECTION_LOGGER, message=message, clientId=self._topology_id, @@ -828,10 +832,14 @@ def started(self) -> None: def waiting(self, remaining_time_ms: int) -> None: """Emit the server selection WAITING log entry.""" - if self._should_log: + # Unlike the other server selection messages, which are debug level, the + # WAITING message MUST be logged at info level per the server selection + # logging specification. + if _is_info_enabled(_SERVER_SELECTION_LOGGER): self._emit_log( _ServerSelectionStatusMessage.WAITING, self._topology_description, + info=True, remainingTimeMS=remaining_time_ms, ) diff --git a/pymongo/asynchronous/change_stream.py b/pymongo/asynchronous/change_stream.py index 5f156a99ba..cf46862e32 100644 --- a/pymongo/asynchronous/change_stream.py +++ b/pymongo/asynchronous/change_stream.py @@ -223,7 +223,7 @@ async def _run_aggregation_cmd( cmd.get_cursor, self._target._read_preference_for(session), session, - operation=_Op.AGGREGATE, + operation=_Op.WATCH, ) async def _create_cursor(self) -> AsyncCommandCursor: # type: ignore[type-arg] diff --git a/pymongo/asynchronous/collection.py b/pymongo/asynchronous/collection.py index ded61264ee..d4f92ec566 100644 --- a/pymongo/asynchronous/collection.py +++ b/pymongo/asynchronous/collection.py @@ -766,14 +766,29 @@ async def bulk_write( common.validate_list("requests", requests) blk = _AsyncBulk(self, ordered, bypass_document_validation, comment=comment, let=let) + # Derive the operation name from the write models: the write type's name + # when every request is of the same type, ``bulkWrite`` when mixed. + operation: Optional[_Op] = None for request in requests: try: request._add_to_bulk(blk) except AttributeError: raise TypeError(f"{request!r} is not a valid request") from None + if isinstance(request, InsertOne): + current = _Op.INSERT + elif isinstance(request, (UpdateOne, UpdateMany, ReplaceOne)): + current = _Op.UPDATE + elif isinstance(request, (DeleteOne, DeleteMany)): + current = _Op.DELETE + else: + current = _Op.BULK_WRITE + if operation is None: + operation = current + elif operation != current: + operation = _Op.BULK_WRITE write_concern = self._write_concern_for(session) - bulk_api_result = await blk.execute(write_concern, session, _Op.INSERT) + bulk_api_result = await blk.execute(write_concern, session, operation or _Op.BULK_WRITE) if bulk_api_result is not None: return BulkWriteResult(bulk_api_result, True) return BulkWriteResult({}, False) @@ -2148,7 +2163,7 @@ async def _cmd( return 0 return result["n"] - return await self._retryable_non_cursor_read(_cmd, session, _Op.COUNT) + return await self._retryable_non_cursor_read(_cmd, session, _Op.COUNT_DOCUMENTS) async def _retryable_non_cursor_read( self, @@ -3157,7 +3172,14 @@ async def inner( client=client, ) - return await client._retryable_write(False, inner, session, _Op.RENAME) + return await client._retryable_write( + False, + inner, + session, + _Op.RENAME, + dbname=self._database.name, + collection=self._name, + ) async def distinct( self, diff --git a/pymongo/asynchronous/database.py b/pymongo/asynchronous/database.py index 330ba68bb6..785d796d59 100644 --- a/pymongo/asynchronous/database.py +++ b/pymongo/asynchronous/database.py @@ -33,7 +33,7 @@ from bson.dbref import DBRef from bson.timestamp import Timestamp from pymongo import _csot, common -from pymongo._otel import internal_cursor_iteration +from pymongo._otel import _extract_collection_name, internal_cursor_iteration from pymongo.asynchronous.aggregation import _DatabaseAggregationCommand from pymongo.asynchronous.change_stream import AsyncDatabaseChangeStream from pymongo.asynchronous.collection import AsyncCollection @@ -1054,6 +1054,10 @@ async def inner( else: raise InvalidOperation("Command does not return a cursor.") + # A cursor-returning generic command usually targets a collection; the + # operation span reports it when it does (DRIVERS-3625). + cmd_doc = {command_name: value} if isinstance(command, str) else command + collection_name = _extract_collection_name(command_name, self.name, cmd_doc) return await self.client._retryable_read( inner, read_preference, @@ -1062,6 +1066,7 @@ async def inner( None, False, dbname=self.name, + collection=collection_name, is_run_command=True, attach_operation_telemetry=True, ) diff --git a/pymongo/asynchronous/mongo_client.py b/pymongo/asynchronous/mongo_client.py index 071cdd6114..af6f895167 100644 --- a/pymongo/asynchronous/mongo_client.py +++ b/pymongo/asynchronous/mongo_client.py @@ -56,7 +56,7 @@ from bson.codec_options import DEFAULT_CODEC_OPTIONS, CodecOptions, TypeRegistry from bson.timestamp import Timestamp from pymongo import _csot, _op_id, common, helpers_shared, periodic_executor -from pymongo._otel import is_internal_cursor_iteration +from pymongo._otel import _build_operation_name, is_internal_cursor_iteration from pymongo._telemetry import ( _generate_op_id_or_none, _operation_telemetry_or_none, @@ -2013,6 +2013,8 @@ async def _retry_with_session( bulk: Optional[Union[_AsyncBulk, _AsyncClientBulk]], operation: str, operation_id: Optional[int] = None, + dbname: Optional[str] = None, + collection: Optional[str] = None, ) -> T: """Execute an operation with at most one consecutive retries @@ -2033,6 +2035,8 @@ async def _retry_with_session( operation=operation, retryable=retryable, operation_id=operation_id, + dbname=dbname, + collection=collection, ) @_csot.apply @@ -2208,6 +2212,8 @@ async def _retryable_write( operation: str, bulk: Optional[Union[_AsyncBulk, _AsyncClientBulk]] = None, operation_id: Optional[int] = None, + dbname: Optional[str] = None, + collection: Optional[str] = None, ) -> T: """Execute an operation with consecutive retries if possible @@ -2224,7 +2230,9 @@ async def _retryable_write( :param operation_id: Stable operation id shared across retries, defaults to None """ async with self._tmp_session(session) as s: - return await self._retry_with_session(retryable, func, s, bulk, operation, operation_id) + return await self._retry_with_session( + retryable, func, s, bulk, operation, operation_id, dbname, collection + ) def _cleanup_cursor_no_lock( self, @@ -3007,7 +3015,9 @@ def __init__( self._address = address self._server: Server = None # type: ignore self._deprioritized_servers: Optional[list[Server]] = None - self._operation = operation + # Normalize once so server selection logging and the operation span + # report the same operation name (DRIVERS-3625). + self._operation = _build_operation_name(operation, is_run_command) # Only generate an operation id when APM/logging is enabled. if operation_id is None: operation_id = _generate_op_id_or_none(self._client._event_listeners) diff --git a/pymongo/asynchronous/topology.py b/pymongo/asynchronous/topology.py index c4f1a3fa96..e4507e9c05 100644 --- a/pymongo/asynchronous/topology.py +++ b/pymongo/asynchronous/topology.py @@ -55,7 +55,11 @@ _async_create_condition, _async_create_lock, ) -from pymongo.logger import _SERVER_SELECTION_LOGGER, _is_debug_enabled +from pymongo.logger import ( + _SERVER_SELECTION_LOGGER, + _is_debug_enabled, + _is_info_enabled, +) from pymongo.pool_options import PoolOptions from pymongo.server_description import ServerDescription from pymongo.server_selectors import ( @@ -272,9 +276,14 @@ async def _select_servers_loop( now = time.monotonic() end_time = now + timeout logged_waiting = False - # Server selection does not have APM events, gate only on logging + # Server selection does not have APM events, gate only on logging. The + # WAITING message is info level while the others are debug level, so the + # telemetry object is created for users of either level; each method + # checks the level it needs. ss: Optional[_ServerSelectionTelemetry] = None - if _is_debug_enabled(_SERVER_SELECTION_LOGGER): + if _is_debug_enabled(_SERVER_SELECTION_LOGGER) or _is_info_enabled( + _SERVER_SELECTION_LOGGER + ): ss = _ServerSelectionTelemetry( self._topology_id, selector, operation, operation_id, self.description ) diff --git a/pymongo/logger.py b/pymongo/logger.py index 43256d18a5..8c52ea11f5 100644 --- a/pymongo/logger.py +++ b/pymongo/logger.py @@ -100,6 +100,10 @@ def _is_debug_enabled(logger: logging.Logger) -> bool: return logger.isEnabledFor(logging.DEBUG) +def _is_info_enabled(logger: logging.Logger) -> bool: + return logger.isEnabledFor(logging.INFO) + + def _log_client_error() -> None: # This is called from a daemon thread so check for None to account for interpreter shutdown. logger = _CLIENT_LOGGER diff --git a/pymongo/operations.py b/pymongo/operations.py index 99b6a00c52..7ffef31384 100644 --- a/pymongo/operations.py +++ b/pymongo/operations.py @@ -55,15 +55,19 @@ class _Op(str, enum.Enum): BULK_WRITE = "bulkWrite" COMMIT = "commitTransaction" COUNT = "count" - CREATE = "create" + COUNT_DOCUMENTS = "countDocuments" + # Values MUST match the covered operations table of the OpenTelemetry + # specification (DRIVERS-3625), which is the definitive list of operation + # names for both tracing and server selection logging. + CREATE = "createCollection" CREATE_INDEXES = "createIndexes" CREATE_SEARCH_INDEXES = "createSearchIndexes" DELETE = "delete" DISTINCT = "distinct" - DROP = "drop" + DROP = "dropCollection" DROP_DATABASE = "dropDatabase" DROP_INDEXES = "dropIndexes" - DROP_SEARCH_INDEXES = "dropSearchIndexes" + DROP_SEARCH_INDEXES = "dropSearchIndex" END_SESSIONS = "endSessions" FIND_AND_MODIFY = "findAndModify" FIND = "find" @@ -75,24 +79,27 @@ class _Op(str, enum.Enum): UPDATE = "update" UPDATE_INDEX = "updateIndex" UPDATE_SEARCH_INDEX = "updateSearchIndex" - RENAME = "rename" + RENAME = "renameCollection" + WATCH = "watch" GETMORE = "getMore" KILL_CURSORS = "killCursors" TEST = "testOperation" +# Literal wire command names, not _Op values: the consumer matches against +# next(iter(command)), so _Op renames must never change this set (PYTHON-5809). _WRITES_WITH_CLUSTER_TIME = frozenset( { - _Op.INSERT.value, - _Op.UPDATE.value, - _Op.FIND_AND_MODIFY.value, - _Op.DELETE.value, - _Op.BULK_WRITE.value, - _Op.CREATE.value, - _Op.CREATE_INDEXES.value, - _Op.DROP.value, - _Op.DROP_DATABASE.value, - _Op.DROP_INDEXES.value, + "insert", + "update", + "findAndModify", + "delete", + "bulkWrite", + "create", + "createIndexes", + "drop", + "dropDatabase", + "dropIndexes", } ) diff --git a/pymongo/synchronous/change_stream.py b/pymongo/synchronous/change_stream.py index 1be0e0308c..f575dc44fe 100644 --- a/pymongo/synchronous/change_stream.py +++ b/pymongo/synchronous/change_stream.py @@ -221,7 +221,7 @@ def _run_aggregation_cmd(self, session: Optional[ClientSession]) -> CommandCurso cmd.get_cursor, self._target._read_preference_for(session), session, - operation=_Op.AGGREGATE, + operation=_Op.WATCH, ) def _create_cursor(self) -> CommandCursor: # type: ignore[type-arg] diff --git a/pymongo/synchronous/collection.py b/pymongo/synchronous/collection.py index 81dea58b0b..9b89279874 100644 --- a/pymongo/synchronous/collection.py +++ b/pymongo/synchronous/collection.py @@ -766,14 +766,29 @@ def bulk_write( common.validate_list("requests", requests) blk = _Bulk(self, ordered, bypass_document_validation, comment=comment, let=let) + # Derive the operation name from the write models: the write type's name + # when every request is of the same type, ``bulkWrite`` when mixed. + operation: Optional[_Op] = None for request in requests: try: request._add_to_bulk(blk) except AttributeError: raise TypeError(f"{request!r} is not a valid request") from None + if isinstance(request, InsertOne): + current = _Op.INSERT + elif isinstance(request, (UpdateOne, UpdateMany, ReplaceOne)): + current = _Op.UPDATE + elif isinstance(request, (DeleteOne, DeleteMany)): + current = _Op.DELETE + else: + current = _Op.BULK_WRITE + if operation is None: + operation = current + elif operation != current: + operation = _Op.BULK_WRITE write_concern = self._write_concern_for(session) - bulk_api_result = blk.execute(write_concern, session, _Op.INSERT) + bulk_api_result = blk.execute(write_concern, session, operation or _Op.BULK_WRITE) if bulk_api_result is not None: return BulkWriteResult(bulk_api_result, True) return BulkWriteResult({}, False) @@ -2146,7 +2161,7 @@ def _cmd( return 0 return result["n"] - return self._retryable_non_cursor_read(_cmd, session, _Op.COUNT) + return self._retryable_non_cursor_read(_cmd, session, _Op.COUNT_DOCUMENTS) def _retryable_non_cursor_read( self, @@ -3153,7 +3168,14 @@ def inner( client=client, ) - return client._retryable_write(False, inner, session, _Op.RENAME) + return client._retryable_write( + False, + inner, + session, + _Op.RENAME, + dbname=self._database.name, + collection=self._name, + ) def distinct( self, diff --git a/pymongo/synchronous/database.py b/pymongo/synchronous/database.py index ad89f4c389..f44c0e40ec 100644 --- a/pymongo/synchronous/database.py +++ b/pymongo/synchronous/database.py @@ -33,7 +33,7 @@ from bson.dbref import DBRef from bson.timestamp import Timestamp from pymongo import _csot, common -from pymongo._otel import internal_cursor_iteration +from pymongo._otel import _extract_collection_name, internal_cursor_iteration from pymongo.common import _ecoc_coll_name, _esc_coll_name from pymongo.database_shared import _check_name, _CodecDocumentType from pymongo.errors import CollectionInvalid, InvalidOperation @@ -1054,6 +1054,10 @@ def inner( else: raise InvalidOperation("Command does not return a cursor.") + # A cursor-returning generic command usually targets a collection; the + # operation span reports it when it does (DRIVERS-3625). + cmd_doc = {command_name: value} if isinstance(command, str) else command + collection_name = _extract_collection_name(command_name, self.name, cmd_doc) return self.client._retryable_read( inner, read_preference, @@ -1062,6 +1066,7 @@ def inner( None, False, dbname=self.name, + collection=collection_name, is_run_command=True, attach_operation_telemetry=True, ) diff --git a/pymongo/synchronous/mongo_client.py b/pymongo/synchronous/mongo_client.py index a3f8585435..c9cd0dcfa5 100644 --- a/pymongo/synchronous/mongo_client.py +++ b/pymongo/synchronous/mongo_client.py @@ -56,7 +56,7 @@ from bson.codec_options import DEFAULT_CODEC_OPTIONS, CodecOptions, TypeRegistry from bson.timestamp import Timestamp from pymongo import _csot, _op_id, common, helpers_shared, periodic_executor -from pymongo._otel import is_internal_cursor_iteration +from pymongo._otel import _build_operation_name, is_internal_cursor_iteration from pymongo._telemetry import ( _generate_op_id_or_none, _operation_telemetry_or_none, @@ -2008,6 +2008,8 @@ def _retry_with_session( bulk: Optional[Union[_Bulk, _ClientBulk]], operation: str, operation_id: Optional[int] = None, + dbname: Optional[str] = None, + collection: Optional[str] = None, ) -> T: """Execute an operation with at most one consecutive retries @@ -2028,6 +2030,8 @@ def _retry_with_session( operation=operation, retryable=retryable, operation_id=operation_id, + dbname=dbname, + collection=collection, ) @_csot.apply @@ -2203,6 +2207,8 @@ def _retryable_write( operation: str, bulk: Optional[Union[_Bulk, _ClientBulk]] = None, operation_id: Optional[int] = None, + dbname: Optional[str] = None, + collection: Optional[str] = None, ) -> T: """Execute an operation with consecutive retries if possible @@ -2219,7 +2225,9 @@ def _retryable_write( :param operation_id: Stable operation id shared across retries, defaults to None """ with self._tmp_session(session) as s: - return self._retry_with_session(retryable, func, s, bulk, operation, operation_id) + return self._retry_with_session( + retryable, func, s, bulk, operation, operation_id, dbname, collection + ) def _cleanup_cursor_no_lock( self, @@ -2996,7 +3004,9 @@ def __init__( self._address = address self._server: Server = None # type: ignore self._deprioritized_servers: Optional[list[Server]] = None - self._operation = operation + # Normalize once so server selection logging and the operation span + # report the same operation name (DRIVERS-3625). + self._operation = _build_operation_name(operation, is_run_command) # Only generate an operation id when APM/logging is enabled. if operation_id is None: operation_id = _generate_op_id_or_none(self._client._event_listeners) diff --git a/pymongo/synchronous/topology.py b/pymongo/synchronous/topology.py index c6468c8912..9c9f92b0c9 100644 --- a/pymongo/synchronous/topology.py +++ b/pymongo/synchronous/topology.py @@ -52,7 +52,11 @@ _create_condition, _create_lock, ) -from pymongo.logger import _SERVER_SELECTION_LOGGER, _is_debug_enabled +from pymongo.logger import ( + _SERVER_SELECTION_LOGGER, + _is_debug_enabled, + _is_info_enabled, +) from pymongo.pool_options import PoolOptions from pymongo.server_description import ServerDescription from pymongo.server_selectors import ( @@ -272,9 +276,14 @@ def _select_servers_loop( now = time.monotonic() end_time = now + timeout logged_waiting = False - # Server selection does not have APM events, gate only on logging + # Server selection does not have APM events, gate only on logging. The + # WAITING message is info level while the others are debug level, so the + # telemetry object is created for users of either level; each method + # checks the level it needs. ss: Optional[_ServerSelectionTelemetry] = None - if _is_debug_enabled(_SERVER_SELECTION_LOGGER): + if _is_debug_enabled(_SERVER_SELECTION_LOGGER) or _is_info_enabled( + _SERVER_SELECTION_LOGGER + ): ss = _ServerSelectionTelemetry( self._topology_id, selector, operation, operation_id, self.description ) diff --git a/test/asynchronous/test_logger.py b/test/asynchronous/test_logger.py index c9a80b9eb6..11edf9529b 100644 --- a/test/asynchronous/test_logger.py +++ b/test/asynchronous/test_logger.py @@ -17,8 +17,12 @@ from unittest.mock import patch from bson import json_util -from pymongo.errors import OperationFailure -from pymongo.logger import _DEFAULT_DOCUMENT_LENGTH, _CommandStatusMessage +from pymongo.errors import OperationFailure, ServerSelectionTimeoutError +from pymongo.logger import ( + _DEFAULT_DOCUMENT_LENGTH, + _CommandStatusMessage, + _ServerSelectionStatusMessage, +) from test import unittest from test.asynchronous import AsyncIntegrationTest, async_client_context @@ -124,6 +128,25 @@ async def test_logging_without_listeners(self): await c.db.coll.insert_one({"x": "1"}) self.assertGreater(len(cm.records), 0) + async def test_server_selection_waiting_message_info_level(self): + # The "Waiting for suitable server to become available" message MUST be logged at + # info level, unlike the other server selection messages, which are debug level. + # A user observing server selection logs at INFO (only) must still receive the + # waiting message, which is only possible if the telemetry object is created for + # info-enabled users too. + client = await self.async_single_client(p=27999, serverSelectionTimeoutMS=500) + try: + with self.assertLogs("pymongo.serverSelection", level="INFO") as cm: + with self.assertRaises(ServerSelectionTimeoutError): + await client.pymongo_test.command("ping") + finally: + await client.close() + self.assertEqual(len(cm.records), 1) + for record in cm.records: + self.assertEqual(record.levelname, "INFO") + log = json_util.loads(record.getMessage()) + self.assertEqual(log["message"], _ServerSelectionStatusMessage.WAITING) + @async_client_context.require_failCommand_fail_point async def test_logging_retry_read_attempts(self): await self.db.coll.insert_one({"x": "1"}) diff --git a/test/asynchronous/test_otel.py b/test/asynchronous/test_otel.py index cf9436aac1..f89ac26a2b 100644 --- a/test/asynchronous/test_otel.py +++ b/test/asynchronous/test_otel.py @@ -22,7 +22,8 @@ import subprocess import sys import time -from typing import Callable, Optional +from collections.abc import Mapping +from typing import Any, Callable, Optional from unittest.mock import MagicMock, patch sys.path[0:0] = [""] @@ -30,6 +31,7 @@ import pytest import pymongo._otel as _otel +from bson.objectid import ObjectId from pymongo import _telemetry, common from pymongo._telemetry import _OperationTelemetry from pymongo.cursor_shared import CursorType @@ -54,10 +56,12 @@ _HAS_OTEL_TEST_DEPS = False if _otel._HAS_OPENTELEMETRY: try: + from opentelemetry import context as otel_context from opentelemetry import trace from opentelemetry.sdk.trace.export import SimpleSpanProcessor from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter - from opentelemetry.trace import StatusCode + from opentelemetry.trace import Span, SpanContext, StatusCode + from opentelemetry.util import types as otel_types _HAS_OTEL_TEST_DEPS = True except ImportError: @@ -86,6 +90,64 @@ def _qualified_name(exc_type: type) -> str: return f"{exc_type.__module__}.{exc_type.__qualname__}" +if _HAS_OTEL_TEST_DEPS: + + class _ApiOnlySpan(Span): + """A recording span implementing the OpenTelemetry Span interface without exposing an + ``attributes`` property, like a non-SDK span implementation may.""" + + def __init__(self, attributes: Optional[otel_types.Attributes] = None) -> None: + self._attributes: dict[str, otel_types.AttributeValue] = dict(attributes or {}) + self._name = "" + + def get_span_context(self) -> SpanContext: + return SpanContext(trace_id=1, span_id=1, is_remote=False) + + def is_recording(self) -> bool: + return True + + def end(self, end_time: Optional[int] = None) -> None: + pass + + def set_attributes(self, attributes: Mapping[str, otel_types.AttributeValue]) -> None: + self._attributes.update(attributes) + + def set_attribute(self, key: str, value: otel_types.AttributeValue) -> None: + self._attributes[key] = value + + def add_event( + self, + name: str, + attributes: otel_types.Attributes = None, + timestamp: Optional[int] = None, + ) -> None: + pass + + def record_exception( + self, + exception: BaseException, + attributes: otel_types.Attributes = None, + timestamp: Optional[int] = None, + escaped: bool = False, + ) -> None: + pass + + def update_name(self, name: str) -> None: + self._name = name + + def set_status(self, status: Any, description: Optional[str] = None) -> None: + pass + + class _FakeConnInfo: + """Minimal stand-in for ``_ConnectionTelemetryInfo`` as read by ``start_command_span``.""" + + id: int = 1 + address: tuple[str, Optional[int]] = ("localhost", 27017) + server_connection_id: Optional[int] = 1 + service_id: Optional[ObjectId] = None + max_wire_version: int = 0 + + @unittest.skipUnless(_HAS_OTEL_TEST_DEPS, "opentelemetry-sdk is not installed") class TestOTelOperationSpanPrimitives(unittest.TestCase): """Unit tests for the pymongo._otel operation-span primitives.""" @@ -128,6 +190,61 @@ def test_start_operation_span_failure_records_exception(self): self.assertEqual(len(span.events), 1) self.assertEqual(span.events[0].name, "exception") + def test_start_command_span_backfill_skips_initialized_namespace(self): + """A rename targets a collection while its command runs against admin. The + operation span's namespace is set eagerly, and the command backfill must not + overwrite it, even when the span implementation does not expose its attributes + (the OpenTelemetry Span API does not require an ``attributes`` getter).""" + handle = _otel.start_operation_span( + _tracing_opts(), "renameCollection", None, dbname="db", collection="coll" + ) + fake = _ApiOnlySpan( + { + "db.system.name": "mongodb", + "db.operation.name": "renameCollection", + "db.namespace": "db", + "db.collection.name": "coll", + "db.operation.summary": "renameCollection db.coll", + } + ) + token = otel_context.attach(trace.set_span_in_context(fake)) + try: + _otel.start_command_span( + _tracing_opts(), + _FakeConnInfo(), + {"renameCollection": "coll"}, + "admin", + "renameCollection", + False, + ) + finally: + otel_context.detach(token) + _otel.end_operation_span_success(handle) + self.assertEqual(fake._attributes["db.namespace"], "db") + self.assertEqual(fake._attributes["db.collection.name"], "coll") + self.assertEqual(fake._name, "") + + def test_start_command_span_backfills_unknown_namespace(self): + """An operation that started without a namespace learns it from the first command.""" + handle = _otel.start_operation_span(_tracing_opts(), "runCommand", None) + fake = _ApiOnlySpan( + { + "db.system.name": "mongodb", + "db.operation.name": "runCommand", + "db.operation.summary": "runCommand", + } + ) + token = otel_context.attach(trace.set_span_in_context(fake)) + try: + _otel.start_command_span( + _tracing_opts(), _FakeConnInfo(), {"ping": 1}, "admin", "ping", False + ) + finally: + otel_context.detach(token) + _otel.end_operation_span_success(handle) + self.assertEqual(fake._attributes["db.namespace"], "admin") + self.assertNotIn("db.collection.name", fake._attributes) + def test_start_operation_span_with_parent(self): parent_handle = _otel.start_operation_span(_tracing_opts(), "transaction", None) handle = _otel.start_operation_span(_tracing_opts(), "insert", parent_handle.span) @@ -848,6 +965,7 @@ def test_no_client_options_is_never_traced(self): with patch.dict(os.environ, {"OTEL_PYTHON_INSTRUMENTATION_MONGODB_ENABLED": "true"}): self.assertFalse(_otel._is_tracing_enabled(None)) + @async_client_context.require_version_min(8, 0, 0, -24) async def test_unacknowledged_bulk_write_query_text_includes_documents(self): # db.query.text is built from the document published in # CommandStartedEvent, which carries the write documents that the @@ -899,6 +1017,7 @@ async def task(): (span,) = self.spans("getMore") self.assertEqual(span.status.status_code, trace.StatusCode.ERROR) + @async_client_context.require_version_min(8, 0, 0, -24) async def test_bulk_write_unacknowledged_gets_operation_span(self): client = await self.async_rs_or_single_client(tracing={"enabled": True}, w=0) self.exporter.clear() @@ -1001,7 +1120,7 @@ async def test_operation_span_falls_back_to_bare_name_when_no_command_is_sent(se self.assertEqual(span.status.status_code, StatusCode.ERROR) async def test_operation_span_name_can_differ_from_command_name(self): - # count_documents' operation span is named "count" but sends an + # count_documents' operation span is named "countDocuments" but sends an # aggregate, so an operation span name is not the command beneath it. # count.json covers estimated_document_count, where the two coincide. client = await self.async_rs_or_single_client(tracing={"enabled": True}) @@ -1010,8 +1129,8 @@ async def test_operation_span_name_can_differ_from_command_name(self): self.exporter.clear() await db.mycoll.count_documents({}) - (op_span,) = self.spans("count pymongo_test.mycoll") - self.assertEqual(op_span.attributes["db.operation.name"], "count") + (op_span,) = self.spans("countDocuments pymongo_test.mycoll") + self.assertEqual(op_span.attributes["db.operation.name"], "countDocuments") self.assertEqual(op_span.attributes["db.namespace"], "pymongo_test") (cmd_span,) = self.spans("aggregate") self.assertEqual(cmd_span.attributes["db.command.name"], "aggregate") diff --git a/test/asynchronous/test_otel_getmore.py b/test/asynchronous/test_otel_getmore.py index a4a3ec63df..115fac010c 100644 --- a/test/asynchronous/test_otel_getmore.py +++ b/test/asynchronous/test_otel_getmore.py @@ -121,11 +121,11 @@ def ping_spans(self): or s.attributes.get("db.operation.name") == "runCommand" ] - def _aggregate_operation_span(self): + def _watch_operation_span(self): matching = [ s for s in self.exporter.get_finished_spans() - if s.attributes.get("db.operation.name") == "aggregate" + if s.attributes.get("db.operation.name") == "watch" ] self.assertEqual(len(matching), 1) return matching[0] @@ -474,7 +474,7 @@ async def test_change_stream_collection_level_operation_span_has_full_namespace( self.exporter.clear() async with await coll.watch(): pass - span = self._aggregate_operation_span() + span = self._watch_operation_span() self.assertEqual(span.attributes["db.namespace"], "pymongo_test") self.assertEqual(span.attributes["db.collection.name"], "test_otel_change_stream_coll") @@ -486,7 +486,7 @@ async def test_change_stream_database_level_operation_span_omits_collection_name self.exporter.clear() async with await db.watch(): pass - span = self._aggregate_operation_span() + span = self._watch_operation_span() self.assertEqual(span.attributes["db.namespace"], "pymongo_test") self.assertNotIn("db.collection.name", span.attributes) @@ -497,6 +497,6 @@ async def test_change_stream_cluster_level_operation_span_targets_admin(self): self.exporter.clear() async with await client.watch(): pass - span = self._aggregate_operation_span() + span = self._watch_operation_span() self.assertEqual(span.attributes["db.namespace"], "admin") self.assertNotIn("db.collection.name", span.attributes) diff --git a/test/asynchronous/unified_format.py b/test/asynchronous/unified_format.py index 7720856c9b..fab55cbddc 100644 --- a/test/asynchronous/unified_format.py +++ b/test/asynchronous/unified_format.py @@ -71,6 +71,7 @@ CommandStartedEvent, ) from pymongo.operations import ( + IndexModel, SearchIndexModel, ) from pymongo.read_concern import ReadConcern @@ -964,6 +965,24 @@ async def _collectionOperation_createFindCursor(self, target, *args, **kwargs): def _collectionOperation_count(self, target, *args, **kwargs): self.skipTest("PyMongo does not support collection.count()") + async def _collectionOperation_renameCollection(self, target, *args, **kwargs): + # PyMongo exposes the renameCollection command as Collection.rename(). + kwargs["new_name"] = kwargs.pop("to") + return await target.rename(*args, **kwargs) + + async def _collectionOperation_createIndexes(self, target, *args, **kwargs): + models = [IndexModel(**i) for i in kwargs.pop("indexes")] + return await target.create_indexes(models, *args, **kwargs) + + async def _collectionOperation_dropIndexes(self, target, *args, **kwargs): + index = kwargs.pop("index", None) or kwargs.pop("indexes", None) + if index is None: + # No argument drops all indexes. + return await target.drop_indexes(*args, **kwargs) + if isinstance(index, str): + index = [index] + return await target.drop_index(index[0], *args, **kwargs) + async def _collectionOperation_listIndexes(self, target, *args, **kwargs): if "batch_size" in kwargs: self.skipTest("PyMongo does not support batch_size for list_indexes") diff --git a/test/open_telemetry/operation/collection_bulk_write.json b/test/open_telemetry/operation/collection_bulk_write.json new file mode 100644 index 0000000000..9fd2a8355a --- /dev/null +++ b/test/open_telemetry/operation/collection_bulk_write.json @@ -0,0 +1,537 @@ +{ + "description": "operation collection_bulk_write", + "schemaVersion": "1.27", + "createEntities": [ + { + "client": { + "id": "client0", + "useMultipleMongoses": false, + "observeTracingMessages": { + "enableCommandPayload": true + } + } + }, + { + "database": { + "id": "database0", + "client": "client0", + "databaseName": "operation-collection-bulk-write" + } + }, + { + "collection": { + "id": "collection0", + "database": "database0", + "collectionName": "test" + } + } + ], + "initialData": [ + { + "collectionName": "test", + "databaseName": "operation-collection-bulk-write", + "documents": [] + } + ], + "tests": [ + { + "description": "mixed write models report bulkWrite", + "operations": [ + { + "name": "bulkWrite", + "object": "collection0", + "arguments": { + "requests": [ + { + "insertOne": { + "document": { + "_id": 1, + "x": 1 + } + } + }, + { + "updateOne": { + "filter": { + "_id": 2 + }, + "update": { + "$set": { + "x": 2 + } + } + } + }, + { + "deleteOne": { + "filter": { + "_id": 3 + } + } + } + ] + } + } + ], + "expectTracingMessages": [ + { + "client": "client0", + "ignoreExtraSpans": false, + "spans": [ + { + "name": "bulkWrite operation-collection-bulk-write.test", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-collection-bulk-write", + "db.collection.name": "test", + "db.operation.name": "bulkWrite", + "db.operation.summary": "bulkWrite operation-collection-bulk-write.test" + }, + "nested": [ + { + "name": "insert", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-collection-bulk-write", + "db.collection.name": "test", + "db.command.name": "insert", + "network.transport": "tcp", + "db.mongodb.cursor_id": { + "$$exists": false + }, + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "insert operation-collection-bulk-write.test", + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + }, + { + "name": "update", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-collection-bulk-write", + "db.collection.name": "test", + "db.command.name": "update", + "network.transport": "tcp", + "db.mongodb.cursor_id": { + "$$exists": false + }, + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "update operation-collection-bulk-write.test", + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + }, + { + "name": "delete", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-collection-bulk-write", + "db.collection.name": "test", + "db.command.name": "delete", + "network.transport": "tcp", + "db.mongodb.cursor_id": { + "$$exists": false + }, + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "delete operation-collection-bulk-write.test", + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] + }, + { + "description": "uniform insert write models report insert", + "operations": [ + { + "name": "bulkWrite", + "object": "collection0", + "arguments": { + "requests": [ + { + "insertOne": { + "document": { + "_id": 1, + "x": 1 + } + } + }, + { + "insertOne": { + "document": { + "_id": 2, + "x": 2 + } + } + } + ] + } + } + ], + "expectTracingMessages": [ + { + "client": "client0", + "ignoreExtraSpans": false, + "spans": [ + { + "name": "insert operation-collection-bulk-write.test", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-collection-bulk-write", + "db.collection.name": "test", + "db.operation.name": "insert", + "db.operation.summary": "insert operation-collection-bulk-write.test" + }, + "nested": [ + { + "name": "insert", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-collection-bulk-write", + "db.collection.name": "test", + "db.command.name": "insert", + "network.transport": "tcp", + "db.mongodb.cursor_id": { + "$$exists": false + }, + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "insert operation-collection-bulk-write.test", + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] + }, + { + "description": "uniform update write models report update", + "operations": [ + { + "name": "bulkWrite", + "object": "collection0", + "arguments": { + "requests": [ + { + "updateOne": { + "filter": { + "_id": 1 + }, + "update": { + "$set": { + "x": 1 + } + } + } + }, + { + "updateMany": { + "filter": { + "_id": { + "$gt": 1 + } + }, + "update": { + "$inc": { + "x": 1 + } + } + } + } + ] + } + } + ], + "expectTracingMessages": [ + { + "client": "client0", + "ignoreExtraSpans": false, + "spans": [ + { + "name": "update operation-collection-bulk-write.test", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-collection-bulk-write", + "db.collection.name": "test", + "db.operation.name": "update", + "db.operation.summary": "update operation-collection-bulk-write.test" + }, + "nested": [ + { + "name": "update", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-collection-bulk-write", + "db.collection.name": "test", + "db.command.name": "update", + "network.transport": "tcp", + "db.mongodb.cursor_id": { + "$$exists": false + }, + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "update operation-collection-bulk-write.test", + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] + }, + { + "description": "uniform delete write models report delete", + "operations": [ + { + "name": "bulkWrite", + "object": "collection0", + "arguments": { + "requests": [ + { + "deleteOne": { + "filter": { + "_id": 1 + } + } + }, + { + "deleteMany": { + "filter": { + "_id": { + "$gt": 1 + } + } + } + } + ] + } + } + ], + "expectTracingMessages": [ + { + "client": "client0", + "ignoreExtraSpans": false, + "spans": [ + { + "name": "delete operation-collection-bulk-write.test", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-collection-bulk-write", + "db.collection.name": "test", + "db.operation.name": "delete", + "db.operation.summary": "delete operation-collection-bulk-write.test" + }, + "nested": [ + { + "name": "delete", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-collection-bulk-write", + "db.collection.name": "test", + "db.command.name": "delete", + "network.transport": "tcp", + "db.mongodb.cursor_id": { + "$$exists": false + }, + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "delete operation-collection-bulk-write.test", + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] + } + ] +} diff --git a/test/open_telemetry/operation/count_documents.json b/test/open_telemetry/operation/count_documents.json new file mode 100644 index 0000000000..acf65ed076 --- /dev/null +++ b/test/open_telemetry/operation/count_documents.json @@ -0,0 +1,132 @@ +{ + "description": "operation count_documents", + "schemaVersion": "1.27", + "createEntities": [ + { + "client": { + "id": "client0", + "useMultipleMongoses": false, + "observeTracingMessages": { + "enableCommandPayload": true + } + } + }, + { + "database": { + "id": "database0", + "client": "client0", + "databaseName": "operation-count-documents" + } + }, + { + "collection": { + "id": "collection0", + "database": "database0", + "collectionName": "test" + } + } + ], + "initialData": [ + { + "collectionName": "test", + "databaseName": "operation-count-documents", + "documents": [ + { + "_id": 1, + "x": 1 + } + ] + } + ], + "tests": [ + { + "description": "countDocuments", + "operations": [ + { + "name": "countDocuments", + "object": "collection0", + "arguments": { + "filter": { + "_id": 1 + } + }, + "expectResult": 1 + } + ], + "expectTracingMessages": [ + { + "client": "client0", + "ignoreExtraSpans": false, + "spans": [ + { + "name": "countDocuments operation-count-documents.test", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-count-documents", + "db.collection.name": "test", + "db.operation.name": "countDocuments", + "db.operation.summary": "countDocuments operation-count-documents.test" + }, + "nested": [ + { + "name": "aggregate", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-count-documents", + "db.collection.name": "test", + "db.command.name": "aggregate", + "network.transport": "tcp", + "db.mongodb.cursor_id": { + "$$exists": false + }, + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "aggregate operation-count-documents.test", + "db.query.text": { + "$$matchAsDocument": { + "$$matchAsRoot": { + "aggregate": "test" + } + } + }, + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] + } + ] +} diff --git a/test/open_telemetry/operation/error_type.json b/test/open_telemetry/operation/error_type.json index 50faa921df..7dd5a0b9aa 100644 --- a/test/open_telemetry/operation/error_type.json +++ b/test/open_telemetry/operation/error_type.json @@ -111,7 +111,9 @@ "db.command.name": "find", "network.transport": "tcp", "db.response.status_code": "8", - "error.type": "8", + "error.type": { + "$$type": "string" + }, "exception.message": { "$$type": "string" }, @@ -447,7 +449,9 @@ "db.command.name": "find", "network.transport": "tcp", "db.response.status_code": "89", - "error.type": "89", + "error.type": { + "$$type": "string" + }, "exception.message": { "$$type": "string" }, diff --git a/test/open_telemetry/operation/rename_collection.json b/test/open_telemetry/operation/rename_collection.json new file mode 100644 index 0000000000..0e5f241d29 --- /dev/null +++ b/test/open_telemetry/operation/rename_collection.json @@ -0,0 +1,156 @@ +{ + "description": "operation rename_collection", + "schemaVersion": "1.27", + "createEntities": [ + { + "client": { + "id": "client0", + "useMultipleMongoses": false, + "observeTracingMessages": { + "enableCommandPayload": true + } + } + }, + { + "database": { + "id": "database0", + "client": "client0", + "databaseName": "operation-rename-collection" + } + }, + { + "collection": { + "id": "collection0", + "database": "database0", + "collectionName": "test" + } + } + ], + "initialData": [ + { + "collectionName": "test", + "databaseName": "operation-rename-collection", + "documents": [] + } + ], + "tests": [ + { + "description": "renameCollection", + "operations": [ + { + "name": "dropCollection", + "object": "database0", + "arguments": { + "collection": "test1" + } + }, + { + "name": "renameCollection", + "object": "collection0", + "arguments": { + "to": "test1" + } + } + ], + "expectTracingMessages": [ + { + "client": "client0", + "ignoreExtraSpans": false, + "spans": [ + { + "name": "dropCollection operation-rename-collection.test1", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-rename-collection", + "db.collection.name": "test1", + "db.operation.name": "dropCollection", + "db.operation.summary": "dropCollection operation-rename-collection.test1" + }, + "nested": [ + { + "name": "drop", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-rename-collection", + "db.collection.name": "test1", + "db.command.name": "drop", + "db.query.summary": "drop operation-rename-collection.test1" + } + } + ] + }, + { + "name": "renameCollection operation-rename-collection.test", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-rename-collection", + "db.collection.name": "test", + "db.operation.name": "renameCollection", + "db.operation.summary": "renameCollection operation-rename-collection.test" + }, + "nested": [ + { + "name": "renameCollection", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "admin", + "db.collection.name": { + "$$exists": false + }, + "db.command.name": "renameCollection", + "network.transport": "tcp", + "db.mongodb.cursor_id": { + "$$exists": false + }, + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "renameCollection admin", + "db.query.text": { + "$$matchAsDocument": { + "$$matchAsRoot": { + "renameCollection": "operation-rename-collection.test", + "to": "operation-rename-collection.test1" + } + } + }, + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] + } + ] +} diff --git a/test/open_telemetry/operation/run_command.json b/test/open_telemetry/operation/run_command.json new file mode 100644 index 0000000000..5b0a703c64 --- /dev/null +++ b/test/open_telemetry/operation/run_command.json @@ -0,0 +1,117 @@ +{ + "description": "operation run_command", + "schemaVersion": "1.27", + "createEntities": [ + { + "client": { + "id": "client0", + "useMultipleMongoses": false, + "observeTracingMessages": { + "enableCommandPayload": true + } + } + }, + { + "database": { + "id": "database0", + "client": "client0", + "databaseName": "operation-run-command" + } + } + ], + "tests": [ + { + "description": "runCommand reports runCommand regardless of the command sent", + "operations": [ + { + "name": "runCommand", + "object": "database0", + "arguments": { + "command": { + "ping": 1 + }, + "commandName": "ping" + } + } + ], + "expectTracingMessages": [ + { + "client": "client0", + "ignoreExtraSpans": false, + "spans": [ + { + "name": "runCommand operation-run-command", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-run-command", + "db.collection.name": { + "$$exists": false + }, + "db.operation.name": "runCommand", + "db.operation.summary": "runCommand operation-run-command" + }, + "nested": [ + { + "name": "ping", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-run-command", + "db.collection.name": { + "$$exists": false + }, + "db.command.name": "ping", + "network.transport": "tcp", + "db.mongodb.cursor_id": { + "$$exists": false + }, + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "ping operation-run-command", + "db.query.text": { + "$$matchAsDocument": { + "$$matchAsRoot": { + "ping": 1 + } + } + }, + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] + } + ] +} diff --git a/test/open_telemetry/operation/run_cursor_command.json b/test/open_telemetry/operation/run_cursor_command.json new file mode 100644 index 0000000000..a19a6047fd --- /dev/null +++ b/test/open_telemetry/operation/run_cursor_command.json @@ -0,0 +1,236 @@ +{ + "description": "operation run_cursor_command", + "schemaVersion": "1.27", + "createEntities": [ + { + "client": { + "id": "client0", + "useMultipleMongoses": false, + "observeTracingMessages": { + "enableCommandPayload": true + } + } + }, + { + "database": { + "id": "database0", + "client": "client0", + "databaseName": "operation-run-cursor-command" + } + }, + { + "collection": { + "id": "collection0", + "database": "database0", + "collectionName": "test" + } + } + ], + "initialData": [ + { + "collectionName": "test", + "databaseName": "operation-run-cursor-command", + "documents": [ + { + "_id": 1 + } + ] + } + ], + "tests": [ + { + "description": "cursor-returning runCommand on a database-level command", + "operations": [ + { + "name": "runCursorCommand", + "object": "database0", + "arguments": { + "command": { + "listCollections": 1, + "nameOnly": true + }, + "commandName": "listCollections" + }, + "expectResult": [ + { + "name": "test" + } + ] + } + ], + "expectTracingMessages": [ + { + "client": "client0", + "ignoreExtraSpans": false, + "spans": [ + { + "name": "runCommand operation-run-cursor-command", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-run-cursor-command", + "db.collection.name": { + "$$exists": false + }, + "db.operation.name": "runCommand", + "db.operation.summary": "runCommand operation-run-cursor-command" + }, + "nested": [ + { + "name": "listCollections", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-run-cursor-command", + "db.collection.name": { + "$$exists": false + }, + "db.command.name": "listCollections", + "network.transport": "tcp", + "db.mongodb.cursor_id": { + "$$exists": false + }, + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "listCollections operation-run-cursor-command", + "db.query.text": { + "$$matchAsDocument": { + "$$matchAsRoot": { + "listCollections": 1 + } + } + }, + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] + }, + { + "description": "cursor-returning runCommand on a collection-targeting command", + "operations": [ + { + "name": "runCursorCommand", + "object": "database0", + "arguments": { + "command": { + "find": "test", + "filter": {} + }, + "commandName": "find" + }, + "expectResult": [ + { + "_id": 1 + } + ] + } + ], + "expectTracingMessages": [ + { + "client": "client0", + "ignoreExtraSpans": false, + "spans": [ + { + "name": "runCommand operation-run-cursor-command.test", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-run-cursor-command", + "db.collection.name": "test", + "db.operation.name": "runCommand", + "db.operation.summary": "runCommand operation-run-cursor-command.test" + }, + "nested": [ + { + "name": "find", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-run-cursor-command", + "db.collection.name": "test", + "db.command.name": "find", + "network.transport": "tcp", + "db.mongodb.cursor_id": { + "$$exists": false + }, + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "find operation-run-cursor-command.test", + "db.query.text": { + "$$matchAsDocument": { + "$$matchAsRoot": { + "find": "test" + } + } + }, + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] + } + ] +} diff --git a/test/open_telemetry/operation/watch.json b/test/open_telemetry/operation/watch.json new file mode 100644 index 0000000000..b565654b67 --- /dev/null +++ b/test/open_telemetry/operation/watch.json @@ -0,0 +1,132 @@ +{ + "description": "operation watch", + "schemaVersion": "1.27", + "runOnRequirements": [ + { + "minServerVersion": "4.0.0", + "topologies": [ + "replicaset", + "sharded", + "load-balanced" + ], + "serverless": "forbid" + } + ], + "createEntities": [ + { + "client": { + "id": "client0", + "useMultipleMongoses": false, + "observeTracingMessages": { + "enableCommandPayload": true + } + } + }, + { + "database": { + "id": "database0", + "client": "client0", + "databaseName": "operation-watch" + } + }, + { + "collection": { + "id": "collection0", + "database": "database0", + "collectionName": "test" + } + } + ], + "initialData": [ + { + "collectionName": "test", + "databaseName": "operation-watch", + "documents": [] + } + ], + "tests": [ + { + "description": "watch", + "operations": [ + { + "name": "createChangeStream", + "object": "collection0", + "arguments": { + "pipeline": [] + } + } + ], + "expectTracingMessages": [ + { + "client": "client0", + "ignoreExtraSpans": false, + "spans": [ + { + "name": "watch operation-watch.test", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-watch", + "db.collection.name": "test", + "db.operation.name": "watch", + "db.operation.summary": "watch operation-watch.test" + }, + "nested": [ + { + "name": "aggregate", + "attributes": { + "db.system.name": "mongodb", + "db.namespace": "operation-watch", + "db.collection.name": "test", + "db.command.name": "aggregate", + "network.transport": "tcp", + "db.response.status_code": { + "$$exists": false + }, + "exception.message": { + "$$exists": false + }, + "exception.type": { + "$$exists": false + }, + "exception.stacktrace": { + "$$exists": false + }, + "server.address": { + "$$type": "string" + }, + "server.port": { + "$$type": [ + "int", + "long" + ] + }, + "db.query.summary": "aggregate operation-watch.test", + "db.query.text": { + "$$matchAsDocument": { + "$$matchAsRoot": { + "aggregate": "test" + } + } + }, + "db.mongodb.server_connection_id": { + "$$type": [ + "int", + "long" + ] + }, + "db.mongodb.driver_connection_id": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] + } + ] +} diff --git a/test/server_selection_logging/operation-id.json b/test/server_selection_logging/operation-id.json index ccc2623166..72ebff60d8 100644 --- a/test/server_selection_logging/operation-id.json +++ b/test/server_selection_logging/operation-id.json @@ -197,7 +197,7 @@ } }, { - "level": "debug", + "level": "info", "component": "serverSelection", "data": { "message": "Waiting for suitable server to become available", @@ -383,7 +383,7 @@ } }, { - "level": "debug", + "level": "info", "component": "serverSelection", "data": { "message": "Waiting for suitable server to become available", diff --git a/test/server_selection_logging/operation-names-transactions.json b/test/server_selection_logging/operation-names-transactions.json new file mode 100644 index 0000000000..10c4f2b7d5 --- /dev/null +++ b/test/server_selection_logging/operation-names-transactions.json @@ -0,0 +1,565 @@ +{ + "description": "operation-names-transactions", + "schemaVersion": "1.14", + "runOnRequirements": [ + { + "minServerVersion": "4.4", + "topologies": [ + "replicaset" + ], + "serverless": "forbid" + } + ], + "createEntities": [ + { + "client": { + "id": "client", + "uriOptions": { + "retryWrites": false, + "heartbeatFrequencyMS": 500, + "appName": "loggingClient", + "serverSelectionTimeoutMS": 2000 + }, + "observeLogMessages": { + "serverSelection": "debug" + }, + "observeEvents": [ + "serverDescriptionChangedEvent", + "topologyDescriptionChangedEvent" + ] + } + }, + { + "database": { + "id": "database", + "client": "client", + "databaseName": "logging-tests" + } + }, + { + "collection": { + "id": "collection", + "database": "database", + "collectionName": "server-selection" + } + }, + { + "session": { + "id": "session", + "client": "client" + } + } + ], + "initialData": [ + { + "collectionName": "server-selection", + "databaseName": "logging-tests", + "documents": [ + { + "_id": 1, + "x": 1 + } + ] + } + ], + "tests": [ + { + "description": "withTransaction logs the callback's operations and commitTransaction", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 4 + } + }, + { + "name": "withTransaction", + "object": "session", + "arguments": { + "callback": [ + { + "name": "insertOne", + "object": "collection", + "arguments": { + "document": { + "_id": 2 + }, + "session": "session" + } + } + ] + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "insert", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "insert", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "commitTransaction", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "commitTransaction", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "core transaction API logs the callback's operations and commitTransaction", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 4 + } + }, + { + "name": "startTransaction", + "object": "session" + }, + { + "name": "insertOne", + "object": "collection", + "arguments": { + "document": { + "_id": 3 + }, + "session": "session" + } + }, + { + "name": "commitTransaction", + "object": "session" + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "insert", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "insert", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "commitTransaction", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "commitTransaction", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "abortTransaction", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 4 + } + }, + { + "name": "startTransaction", + "object": "session" + }, + { + "name": "insertOne", + "object": "collection", + "arguments": { + "document": { + "_id": 4 + }, + "session": "session" + } + }, + { + "name": "abortTransaction", + "object": "session" + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "insert", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "insert", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "abortTransaction", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "abortTransaction", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "watch reports watch", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 4 + } + }, + { + "name": "createChangeStream", + "object": "collection", + "arguments": { + "pipeline": [] + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "watch", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "watch", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "client bulkWrite reports bulkWrite", + "runOnRequirements": [ + { + "minServerVersion": "8.0", + "topologies": [ + "replicaset" + ], + "serverless": "forbid" + } + ], + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 4 + } + }, + { + "name": "clientBulkWrite", + "object": "client", + "arguments": { + "models": [ + { + "insertOne": { + "namespace": "logging-tests.server-selection", + "document": { + "_id": 5, + "x": 5 + } + } + } + ] + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "bulkWrite", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "bulkWrite", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] +} diff --git a/test/server_selection_logging/operation-names.json b/test/server_selection_logging/operation-names.json new file mode 100644 index 0000000000..88d086cb43 --- /dev/null +++ b/test/server_selection_logging/operation-names.json @@ -0,0 +1,1740 @@ +{ + "description": "operation-names-standalone", + "schemaVersion": "1.14", + "runOnRequirements": [ + { + "topologies": [ + "single" + ] + } + ], + "createEntities": [ + { + "client": { + "id": "client", + "uriOptions": { + "retryWrites": false, + "heartbeatFrequencyMS": 500, + "appName": "loggingClient", + "serverSelectionTimeoutMS": 2000 + }, + "observeLogMessages": { + "serverSelection": "debug" + }, + "observeEvents": [ + "serverDescriptionChangedEvent", + "topologyDescriptionChangedEvent" + ] + } + }, + { + "database": { + "id": "database", + "client": "client", + "databaseName": "logging-tests" + } + }, + { + "collection": { + "id": "collection", + "database": "database", + "collectionName": "server-selection" + } + } + ], + "initialData": [ + { + "collectionName": "server-selection", + "databaseName": "logging-tests", + "documents": [ + { + "_id": 1, + "x": 1 + } + ] + } + ], + "tests": [ + { + "description": "insert", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "insertOne", + "object": "collection", + "arguments": { + "document": { + "_id": 2, + "x": 2 + } + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "insert", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "insert", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "update", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "updateOne", + "object": "collection", + "arguments": { + "filter": { + "_id": 1 + }, + "update": { + "$set": { + "x": 2 + } + } + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "update", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "update", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "delete", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "deleteOne", + "object": "collection", + "arguments": { + "filter": { + "_id": 2 + } + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "delete", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "delete", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "find", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "find", + "object": "collection", + "arguments": { + "filter": {} + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "find", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "find", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "aggregate", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "aggregate", + "object": "collection", + "arguments": { + "pipeline": [ + { + "$match": { + "_id": 1 + } + } + ] + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "aggregate", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "aggregate", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "distinct", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "distinct", + "object": "collection", + "arguments": { + "fieldName": "x" + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "distinct", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "distinct", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "estimatedDocumentCount reports count", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "estimatedDocumentCount", + "object": "collection" + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "count", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "count", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "countDocuments", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "countDocuments", + "object": "collection", + "arguments": { + "filter": { + "_id": 1 + } + }, + "expectResult": 1 + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "countDocuments", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "countDocuments", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "createCollection", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "createCollection", + "object": "database", + "arguments": { + "collection": "created-collection" + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "createCollection", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "createCollection", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "dropCollection", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "dropCollection", + "object": "database", + "arguments": { + "collection": "created-collection" + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "dropCollection", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "dropCollection", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "renameCollection", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "dropCollection", + "object": "database", + "arguments": { + "collection": "renamed-collection" + } + }, + { + "name": "renameCollection", + "object": "collection", + "arguments": { + "to": "renamed-collection" + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "dropCollection", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "dropCollection", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "renameCollection", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "renameCollection", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "createIndexes and dropIndexes", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "createIndexes", + "object": "collection", + "arguments": { + "indexes": [ + { + "keys": { + "x": 1 + }, + "name": "x_1" + } + ] + } + }, + { + "name": "dropIndexes", + "object": "collection", + "arguments": { + "index": "x_1" + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "createIndexes", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "createIndexes", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "dropIndexes", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "dropIndexes", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "listCollections", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "listCollections", + "object": "database" + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "listCollections", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "listCollections", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "listDatabases", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "listDatabases", + "object": "client" + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "listDatabases", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "listDatabases", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "listIndexes", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "listIndexes", + "object": "collection" + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "listIndexes", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "listIndexes", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "findOneAndUpdate reports findAndModify", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "findOneAndUpdate", + "object": "collection", + "arguments": { + "filter": { + "_id": 1 + }, + "update": { + "$set": { + "x": 3 + } + } + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "findAndModify", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "findAndModify", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "collection bulkWrite with mixed write models reports bulkWrite", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "bulkWrite", + "object": "collection", + "arguments": { + "requests": [ + { + "insertOne": { + "document": { + "_id": 3, + "x": 3 + } + } + }, + { + "updateOne": { + "filter": { + "_id": 1 + }, + "update": { + "$set": { + "x": 4 + } + } + } + }, + { + "deleteOne": { + "filter": { + "_id": 3 + } + } + } + ] + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "bulkWrite", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "bulkWrite", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "collection bulkWrite with uniform insert write models reports insert", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "bulkWrite", + "object": "collection", + "arguments": { + "requests": [ + { + "insertOne": { + "document": { + "_id": 4, + "x": 4 + } + } + }, + { + "insertOne": { + "document": { + "_id": 5, + "x": 5 + } + } + } + ] + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "insert", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "insert", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "collection bulkWrite with uniform update write models reports update", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "bulkWrite", + "object": "collection", + "arguments": { + "requests": [ + { + "updateOne": { + "filter": { + "_id": 4 + }, + "update": { + "$set": { + "x": 6 + } + } + } + }, + { + "updateMany": { + "filter": { + "_id": { + "$gt": 1 + } + }, + "update": { + "$inc": { + "x": 1 + } + } + } + } + ] + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "update", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "update", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "collection bulkWrite with uniform delete write models reports delete", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "bulkWrite", + "object": "collection", + "arguments": { + "requests": [ + { + "deleteOne": { + "filter": { + "_id": 5 + } + } + }, + { + "deleteMany": { + "filter": { + "_id": { + "$gt": 100 + } + } + } + } + ] + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "delete", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "delete", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "runCommand", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "runCommand", + "object": "database", + "arguments": { + "command": { + "ping": 1 + }, + "commandName": "ping" + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "runCommand", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "runCommand", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + }, + { + "description": "cursor-returning runCommand", + "operations": [ + { + "name": "waitForEvent", + "object": "testRunner", + "arguments": { + "client": "client", + "event": { + "topologyDescriptionChangedEvent": {} + }, + "count": 2 + } + }, + { + "name": "runCursorCommand", + "object": "database", + "arguments": { + "command": { + "listCollections": 1, + "nameOnly": true + }, + "commandName": "listCollections" + } + } + ], + "expectLogMessages": [ + { + "client": "client", + "messages": [ + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection started", + "selector": { + "$$exists": true + }, + "operation": "runCommand", + "topologyDescription": { + "$$exists": true + } + } + }, + { + "level": "debug", + "component": "serverSelection", + "data": { + "message": "Server selection succeeded", + "selector": { + "$$exists": true + }, + "operation": "runCommand", + "topologyDescription": { + "$$exists": true + }, + "serverHost": { + "$$type": "string" + }, + "serverPort": { + "$$type": [ + "int", + "long" + ] + } + } + } + ] + } + ] + } + ] +} diff --git a/test/server_selection_logging/replica-set.json b/test/server_selection_logging/replica-set.json index 830b1ea51a..5eba784bf2 100644 --- a/test/server_selection_logging/replica-set.json +++ b/test/server_selection_logging/replica-set.json @@ -184,7 +184,7 @@ } }, { - "level": "debug", + "level": "info", "component": "serverSelection", "data": { "message": "Waiting for suitable server to become available", diff --git a/test/server_selection_logging/sharded.json b/test/server_selection_logging/sharded.json index 346c050f9e..d42fba9100 100644 --- a/test/server_selection_logging/sharded.json +++ b/test/server_selection_logging/sharded.json @@ -193,7 +193,7 @@ } }, { - "level": "debug", + "level": "info", "component": "serverSelection", "data": { "message": "Waiting for suitable server to become available", diff --git a/test/server_selection_logging/standalone.json b/test/server_selection_logging/standalone.json index fa01ad9911..3b3eddd841 100644 --- a/test/server_selection_logging/standalone.json +++ b/test/server_selection_logging/standalone.json @@ -191,7 +191,7 @@ } }, { - "level": "debug", + "level": "info", "component": "serverSelection", "data": { "message": "Waiting for suitable server to become available", diff --git a/test/test_logger.py b/test/test_logger.py index 781b7b5f11..b55d1e7b35 100644 --- a/test/test_logger.py +++ b/test/test_logger.py @@ -17,8 +17,12 @@ from unittest.mock import patch from bson import json_util -from pymongo.errors import OperationFailure -from pymongo.logger import _DEFAULT_DOCUMENT_LENGTH, _CommandStatusMessage +from pymongo.errors import OperationFailure, ServerSelectionTimeoutError +from pymongo.logger import ( + _DEFAULT_DOCUMENT_LENGTH, + _CommandStatusMessage, + _ServerSelectionStatusMessage, +) from test import IntegrationTest, client_context, unittest _IS_SYNC = True @@ -123,6 +127,25 @@ def test_logging_without_listeners(self): c.db.coll.insert_one({"x": "1"}) self.assertGreater(len(cm.records), 0) + def test_server_selection_waiting_message_info_level(self): + # The "Waiting for suitable server to become available" message MUST be logged at + # info level, unlike the other server selection messages, which are debug level. + # A user observing server selection logs at INFO (only) must still receive the + # waiting message, which is only possible if the telemetry object is created for + # info-enabled users too. + client = self.single_client(p=27999, serverSelectionTimeoutMS=500) + try: + with self.assertLogs("pymongo.serverSelection", level="INFO") as cm: + with self.assertRaises(ServerSelectionTimeoutError): + client.pymongo_test.command("ping") + finally: + client.close() + self.assertEqual(len(cm.records), 1) + for record in cm.records: + self.assertEqual(record.levelname, "INFO") + log = json_util.loads(record.getMessage()) + self.assertEqual(log["message"], _ServerSelectionStatusMessage.WAITING) + @client_context.require_failCommand_fail_point def test_logging_retry_read_attempts(self): self.db.coll.insert_one({"x": "1"}) diff --git a/test/test_otel.py b/test/test_otel.py index 3dfa6f1628..c1af931798 100644 --- a/test/test_otel.py +++ b/test/test_otel.py @@ -22,7 +22,8 @@ import subprocess import sys import time -from typing import Callable, Optional +from collections.abc import Mapping +from typing import Any, Callable, Optional from unittest.mock import MagicMock, patch sys.path[0:0] = [""] @@ -30,6 +31,7 @@ import pytest import pymongo._otel as _otel +from bson.objectid import ObjectId from pymongo import _telemetry, common from pymongo._telemetry import _OperationTelemetry from pymongo.cursor_shared import CursorType @@ -54,10 +56,12 @@ _HAS_OTEL_TEST_DEPS = False if _otel._HAS_OPENTELEMETRY: try: + from opentelemetry import context as otel_context from opentelemetry import trace from opentelemetry.sdk.trace.export import SimpleSpanProcessor from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter - from opentelemetry.trace import StatusCode + from opentelemetry.trace import Span, SpanContext, StatusCode + from opentelemetry.util import types as otel_types _HAS_OTEL_TEST_DEPS = True except ImportError: @@ -86,6 +90,64 @@ def _qualified_name(exc_type: type) -> str: return f"{exc_type.__module__}.{exc_type.__qualname__}" +if _HAS_OTEL_TEST_DEPS: + + class _ApiOnlySpan(Span): + """A recording span implementing the OpenTelemetry Span interface without exposing an + ``attributes`` property, like a non-SDK span implementation may.""" + + def __init__(self, attributes: Optional[otel_types.Attributes] = None) -> None: + self._attributes: dict[str, otel_types.AttributeValue] = dict(attributes or {}) + self._name = "" + + def get_span_context(self) -> SpanContext: + return SpanContext(trace_id=1, span_id=1, is_remote=False) + + def is_recording(self) -> bool: + return True + + def end(self, end_time: Optional[int] = None) -> None: + pass + + def set_attributes(self, attributes: Mapping[str, otel_types.AttributeValue]) -> None: + self._attributes.update(attributes) + + def set_attribute(self, key: str, value: otel_types.AttributeValue) -> None: + self._attributes[key] = value + + def add_event( + self, + name: str, + attributes: otel_types.Attributes = None, + timestamp: Optional[int] = None, + ) -> None: + pass + + def record_exception( + self, + exception: BaseException, + attributes: otel_types.Attributes = None, + timestamp: Optional[int] = None, + escaped: bool = False, + ) -> None: + pass + + def update_name(self, name: str) -> None: + self._name = name + + def set_status(self, status: Any, description: Optional[str] = None) -> None: + pass + + class _FakeConnInfo: + """Minimal stand-in for ``_ConnectionTelemetryInfo`` as read by ``start_command_span``.""" + + id: int = 1 + address: tuple[str, Optional[int]] = ("localhost", 27017) + server_connection_id: Optional[int] = 1 + service_id: Optional[ObjectId] = None + max_wire_version: int = 0 + + @unittest.skipUnless(_HAS_OTEL_TEST_DEPS, "opentelemetry-sdk is not installed") class TestOTelOperationSpanPrimitives(unittest.TestCase): """Unit tests for the pymongo._otel operation-span primitives.""" @@ -128,6 +190,61 @@ def test_start_operation_span_failure_records_exception(self): self.assertEqual(len(span.events), 1) self.assertEqual(span.events[0].name, "exception") + def test_start_command_span_backfill_skips_initialized_namespace(self): + """A rename targets a collection while its command runs against admin. The + operation span's namespace is set eagerly, and the command backfill must not + overwrite it, even when the span implementation does not expose its attributes + (the OpenTelemetry Span API does not require an ``attributes`` getter).""" + handle = _otel.start_operation_span( + _tracing_opts(), "renameCollection", None, dbname="db", collection="coll" + ) + fake = _ApiOnlySpan( + { + "db.system.name": "mongodb", + "db.operation.name": "renameCollection", + "db.namespace": "db", + "db.collection.name": "coll", + "db.operation.summary": "renameCollection db.coll", + } + ) + token = otel_context.attach(trace.set_span_in_context(fake)) + try: + _otel.start_command_span( + _tracing_opts(), + _FakeConnInfo(), + {"renameCollection": "coll"}, + "admin", + "renameCollection", + False, + ) + finally: + otel_context.detach(token) + _otel.end_operation_span_success(handle) + self.assertEqual(fake._attributes["db.namespace"], "db") + self.assertEqual(fake._attributes["db.collection.name"], "coll") + self.assertEqual(fake._name, "") + + def test_start_command_span_backfills_unknown_namespace(self): + """An operation that started without a namespace learns it from the first command.""" + handle = _otel.start_operation_span(_tracing_opts(), "runCommand", None) + fake = _ApiOnlySpan( + { + "db.system.name": "mongodb", + "db.operation.name": "runCommand", + "db.operation.summary": "runCommand", + } + ) + token = otel_context.attach(trace.set_span_in_context(fake)) + try: + _otel.start_command_span( + _tracing_opts(), _FakeConnInfo(), {"ping": 1}, "admin", "ping", False + ) + finally: + otel_context.detach(token) + _otel.end_operation_span_success(handle) + self.assertEqual(fake._attributes["db.namespace"], "admin") + self.assertNotIn("db.collection.name", fake._attributes) + def test_start_operation_span_with_parent(self): parent_handle = _otel.start_operation_span(_tracing_opts(), "transaction", None) handle = _otel.start_operation_span(_tracing_opts(), "insert", parent_handle.span) @@ -844,6 +961,7 @@ def test_no_client_options_is_never_traced(self): with patch.dict(os.environ, {"OTEL_PYTHON_INSTRUMENTATION_MONGODB_ENABLED": "true"}): self.assertFalse(_otel._is_tracing_enabled(None)) + @client_context.require_version_min(8, 0, 0, -24) def test_unacknowledged_bulk_write_query_text_includes_documents(self): # db.query.text is built from the document published in # CommandStartedEvent, which carries the write documents that the @@ -895,6 +1013,7 @@ def task(): (span,) = self.spans("getMore") self.assertEqual(span.status.status_code, trace.StatusCode.ERROR) + @client_context.require_version_min(8, 0, 0, -24) def test_bulk_write_unacknowledged_gets_operation_span(self): client = self.rs_or_single_client(tracing={"enabled": True}, w=0) self.exporter.clear() @@ -995,7 +1114,7 @@ def test_operation_span_falls_back_to_bare_name_when_no_command_is_sent(self): self.assertEqual(span.status.status_code, StatusCode.ERROR) def test_operation_span_name_can_differ_from_command_name(self): - # count_documents' operation span is named "count" but sends an + # count_documents' operation span is named "countDocuments" but sends an # aggregate, so an operation span name is not the command beneath it. # count.json covers estimated_document_count, where the two coincide. client = self.rs_or_single_client(tracing={"enabled": True}) @@ -1004,8 +1123,8 @@ def test_operation_span_name_can_differ_from_command_name(self): self.exporter.clear() db.mycoll.count_documents({}) - (op_span,) = self.spans("count pymongo_test.mycoll") - self.assertEqual(op_span.attributes["db.operation.name"], "count") + (op_span,) = self.spans("countDocuments pymongo_test.mycoll") + self.assertEqual(op_span.attributes["db.operation.name"], "countDocuments") self.assertEqual(op_span.attributes["db.namespace"], "pymongo_test") (cmd_span,) = self.spans("aggregate") self.assertEqual(cmd_span.attributes["db.command.name"], "aggregate") diff --git a/test/test_otel_getmore.py b/test/test_otel_getmore.py index 12a48e1374..0145b3e604 100644 --- a/test/test_otel_getmore.py +++ b/test/test_otel_getmore.py @@ -121,11 +121,11 @@ def ping_spans(self): or s.attributes.get("db.operation.name") == "runCommand" ] - def _aggregate_operation_span(self): + def _watch_operation_span(self): matching = [ s for s in self.exporter.get_finished_spans() - if s.attributes.get("db.operation.name") == "aggregate" + if s.attributes.get("db.operation.name") == "watch" ] self.assertEqual(len(matching), 1) return matching[0] @@ -474,7 +474,7 @@ def test_change_stream_collection_level_operation_span_has_full_namespace(self): self.exporter.clear() with coll.watch(): pass - span = self._aggregate_operation_span() + span = self._watch_operation_span() self.assertEqual(span.attributes["db.namespace"], "pymongo_test") self.assertEqual(span.attributes["db.collection.name"], "test_otel_change_stream_coll") @@ -486,7 +486,7 @@ def test_change_stream_database_level_operation_span_omits_collection_name(self) self.exporter.clear() with db.watch(): pass - span = self._aggregate_operation_span() + span = self._watch_operation_span() self.assertEqual(span.attributes["db.namespace"], "pymongo_test") self.assertNotIn("db.collection.name", span.attributes) @@ -497,6 +497,6 @@ def test_change_stream_cluster_level_operation_span_targets_admin(self): self.exporter.clear() with client.watch(): pass - span = self._aggregate_operation_span() + span = self._watch_operation_span() self.assertEqual(span.attributes["db.namespace"], "admin") self.assertNotIn("db.collection.name", span.attributes) diff --git a/test/unified_format.py b/test/unified_format.py index e532b8042e..ae26478306 100644 --- a/test/unified_format.py +++ b/test/unified_format.py @@ -64,6 +64,7 @@ CommandStartedEvent, ) from pymongo.operations import ( + IndexModel, SearchIndexModel, ) from pymongo.read_concern import ReadConcern @@ -961,6 +962,24 @@ def _collectionOperation_createFindCursor(self, target, *args, **kwargs): def _collectionOperation_count(self, target, *args, **kwargs): self.skipTest("PyMongo does not support collection.count()") + def _collectionOperation_renameCollection(self, target, *args, **kwargs): + # PyMongo exposes the renameCollection command as Collection.rename(). + kwargs["new_name"] = kwargs.pop("to") + return target.rename(*args, **kwargs) + + def _collectionOperation_createIndexes(self, target, *args, **kwargs): + models = [IndexModel(**i) for i in kwargs.pop("indexes")] + return target.create_indexes(models, *args, **kwargs) + + def _collectionOperation_dropIndexes(self, target, *args, **kwargs): + index = kwargs.pop("index", None) or kwargs.pop("indexes", None) + if index is None: + # No argument drops all indexes. + return target.drop_indexes(*args, **kwargs) + if isinstance(index, str): + index = [index] + return target.drop_index(index[0], *args, **kwargs) + def _collectionOperation_listIndexes(self, target, *args, **kwargs): if "batch_size" in kwargs: self.skipTest("PyMongo does not support batch_size for list_indexes")