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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 18 additions & 1 deletion .evergreen/generated_configs/variants.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
48 changes: 36 additions & 12 deletions .evergreen/scripts/generate_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"),
Comment thread
Copilot marked this conversation as resolved.
),
# 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
Expand All @@ -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"),
),
]


Expand Down
28 changes: 28 additions & 0 deletions .github/workflows/test-python.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
72 changes: 50 additions & 22 deletions pymongo/_otel.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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"


Expand All @@ -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:
Expand Down Expand Up @@ -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)

Expand Down Expand Up @@ -591,21 +593,37 @@ 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,
span: Span,
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(
Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand All @@ -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
Expand All @@ -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)


Expand All @@ -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())


Expand All @@ -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)


Expand All @@ -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)


Expand Down
12 changes: 10 additions & 2 deletions pymongo/_telemetry.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,9 @@
_CommandStatusMessage,
_ConnectionStatusMessage,
_debug_log,
_info_log,
_is_debug_enabled,
_is_info_enabled,
_SDAMStatusMessage,
_ServerSelectionStatusMessage,
_verbose_connection_error_reason,
Expand Down Expand Up @@ -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,
Expand All @@ -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):
Comment thread
blink1073 marked this conversation as resolved.
self._emit_log(
_ServerSelectionStatusMessage.WAITING,
self._topology_description,
info=True,
remainingTimeMS=remaining_time_ms,
)

Expand Down
2 changes: 1 addition & 1 deletion pymongo/asynchronous/change_stream.py
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
Loading