diff --git a/.evergreen/config.yml b/.evergreen/config.yml index 842c93ff0e..e1ffe4abf4 100644 --- a/.evergreen/config.yml +++ b/.evergreen/config.yml @@ -281,11 +281,82 @@ functions: working_dir: "src" script: | ${PREPARE_SHELL} - TEST_CMD="bundle exec rake driver_bench" PERFORMANCE_RESULTS_FILE="$PROJECT_DIRECTORY/perf.json" .evergreen/run-tests.sh + # With YJIT, as production deployments run; the suite fails if this + # ruby was built without it. + TEST_CMD="bundle exec rake driver_bench" RUBY_YJIT_ENABLE=1 PERFORMANCE_RESULTS_FILE="$PROJECT_DIRECTORY/perf.json" .evergreen/run-tests.sh - command: perf.send params: file: "${PROJECT_DIRECTORY}/perf.json" + "run benchmark comparison": + # Runs the micro-benchmarks in ${driver_bench_tasks}, set by each task. + # Fewer iterations per run buy more interleaved repetitions in the same + # time, which is what makes the median across repetitions meaningful. + # YJIT is on, as in production; without it the interpreter's per-call + # cost makes tracing look several times more expensive. The run fails if + # this ruby was built without YJIT. + - command: shell.exec + type: test + params: + shell: bash + working_dir: "src" + script: | + ${PREPARE_SHELL} + TEST_CMD="bundle exec rake driver_bench:compare" \ + CONFIGURATIONS="${driver_bench_configurations|off,api-only,sdk-never,sdk-parent-1pct,sdk-always}" \ + DRIVER_BENCH_TASKS="${driver_bench_tasks}" \ + REPS="${driver_bench_reps|10}" \ + DRIVER_BENCH_MAX_ITERATIONS="${driver_bench_max_iterations|20}" \ + RUBY_YJIT_ENABLE="${driver_bench_yjit|1}" \ + DRIVER_BENCH_MIN_TIME="${driver_bench_min_time|20}" \ + DRIVER_BENCH_ENFORCE_TARGETS="${driver_bench_enforce_targets|false}" \ + PERFORMANCE_RESULTS_FILE="$PROJECT_DIRECTORY/${task_name}.json" \ + .evergreen/run-tests.sh + # Sent to the performance store so that overhead_pct, cpu_us_added_per_op + # and allocs_added_per_op are watched for regressions, and kept as an + # artifact for comparison by hand. + - command: perf.send + params: + file: "${PROJECT_DIRECTORY}/${task_name}.json" + - command: s3.put + params: + aws_key: ${aws_key} + aws_secret: ${aws_secret} + local_file: ./src/${task_name}.json + display_name: ${task_name}.json + remote_file: ${UPLOAD_BUCKET}/${version_id}/${build_id}/artifacts/${build_variant}/${task_name}.json + content_type: application/json + permissions: public-read + bucket: mciuploads + + # Measures the per-attribute cost at the SDK boundary. The untraced + # operation is measured on the same host, so the relative figures are + # comparable with the absolute ones. No waterfall: request it in a patch. + # See DRIVERS-3620. + "run attribute cost": + - command: shell.exec + type: test + params: + shell: bash + working_dir: "src" + script: | + ${PREPARE_SHELL} + TEST_CMD="bundle exec rake otel_attributes:evergreen 2>&1 | tee $PROJECT_DIRECTORY/${task_name}.txt" \ + RUBY_YJIT_ENABLE=1 \ + ITERATIONS="${otel_attribute_iterations|50000}" \ + REPS="${otel_attribute_reps|20}" \ + .evergreen/run-tests.sh + - command: s3.put + params: + aws_key: ${aws_key} + aws_secret: ${aws_secret} + local_file: ./src/${task_name}.txt + display_name: ${task_name}.txt + remote_file: ${UPLOAD_BUCKET}/${version_id}/${build_id}/artifacts/${build_variant}/${task_name}.txt + content_type: text/plain + permissions: public-read + bucket: mciuploads + "run tests with orchestration and drivers tools": - command: subprocess.exec type: test @@ -582,6 +653,40 @@ tasks: - name: "driver-bench" commands: - func: "run benchmarks" + # Runs the two small-document micro-benchmarks the OpenTelemetry spec's + # targets are defined on under every driver configuration, and reports + # what each costs relative to the untraced baseline. Both run on one host, + # so their overheads can be compared with each other: hosts differ in + # speed by more than the overheads do. Find many creates too few spans per + # iteration to show tracing cost, and large-document inserts are dominated + # by the server. About two hours at the default ten reps. + - name: "driver-bench-otel" + exec_timeout_secs: 10800 + commands: + - func: "run benchmark comparison" + vars: + driver_bench_tasks: "small doc insertone,find one by id" + # Measures the cost of the command-span attribute shapes end to end, against + # the same two small-document micro-benchmarks. Every configuration records + # every trace; they differ only in how the command span's attributes are + # built. attr-none is the floor (a recorded span with no attributes), + # sdk-always is the driver as shipped, and the rest bracket it. Not on the + # waterfall: request it in a patch. See DRIVERS-3620. + - name: "driver-bench-otel-attributes" + exec_timeout_secs: 10800 + commands: + - func: "run benchmark comparison" + vars: + driver_bench_tasks: "small doc insertone,find one by id" + driver_bench_configurations: "off,sdk-always,attr-none,attr-creation-only,attr-all-at-creation,attr-none-then-all" + driver_bench_reps: "6" + # Measures the per-attribute cost at the SDK boundary, with the untraced + # operation measured on the same host. Not on the waterfall: request it in a + # patch. See DRIVERS-3620. + - name: "otel-attribute-cost" + exec_timeout_secs: 3600 + commands: + - func: "run attribute cost" - name: "test-csot" commands: - func: "run CSOT tests" @@ -1115,6 +1220,44 @@ buildvariants: tasks: - name: "driver-bench" + # Compares the DriverBench micro-benchmarks with OpenTelemetry disabled + # against several enabled configurations. See DRIVERS-3620. + - matrix_name: DriverBenchOTel + matrix_spec: + ruby: "ruby-4.0" + mongodb-version: "8.0" + topology: standalone + os: ubuntu2204 + display_name: DriverBench OTel + tasks: + - name: "driver-bench-otel" + batchtime: 1440 # run at most once a day + + # Measures the cost of the command-span attribute shapes end to end. Kept off + # the waterfall; request it in a patch. See DRIVERS-3620. + - matrix_name: DriverBenchOTelAttributes + matrix_spec: + ruby: "ruby-4.0" + mongodb-version: "8.0" + topology: standalone + os: ubuntu2204 + display_name: DriverBench OTel attributes + tasks: + - name: "driver-bench-otel-attributes" + + # Measures the per-attribute cost at the SDK boundary, with the untraced + # operation on the same host. Kept off the waterfall; request it in a patch. + # See DRIVERS-3620. + - matrix_name: OTelAttributeCost + matrix_spec: + ruby: "ruby-4.0" + mongodb-version: "8.0" + topology: standalone + os: ubuntu2204 + display_name: OTel attribute cost + tasks: + - name: "otel-attribute-cost" + - matrix_name: "auth/ssl" matrix_spec: auth-and-ssl: ["auth-and-ssl", "noauth-and-nossl"] @@ -1222,6 +1365,8 @@ buildvariants: tasks: - name: test-csot + # OTel spec tests run off-PR (no "pr" tag), like ruby-dev and + # DriverBench OTel: they only run on the master waterfall. - matrix_name: OTel matrix_spec: ruby: "ruby-4.0" @@ -1229,7 +1374,6 @@ buildvariants: topology: replica-set-single-node os: ubuntu2204 display_name: "OTel - ${mongodb-version}" - tags: ["pr"] tasks: - name: test-otel diff --git a/.evergreen/config/common.yml.erb b/.evergreen/config/common.yml.erb index 9d440736d8..018f92e374 100644 --- a/.evergreen/config/common.yml.erb +++ b/.evergreen/config/common.yml.erb @@ -278,11 +278,82 @@ functions: working_dir: "src" script: | ${PREPARE_SHELL} - TEST_CMD="bundle exec rake driver_bench" PERFORMANCE_RESULTS_FILE="$PROJECT_DIRECTORY/perf.json" .evergreen/run-tests.sh + # With YJIT, as production deployments run; the suite fails if this + # ruby was built without it. + TEST_CMD="bundle exec rake driver_bench" RUBY_YJIT_ENABLE=1 PERFORMANCE_RESULTS_FILE="$PROJECT_DIRECTORY/perf.json" .evergreen/run-tests.sh - command: perf.send params: file: "${PROJECT_DIRECTORY}/perf.json" + "run benchmark comparison": + # Runs the micro-benchmarks in ${driver_bench_tasks}, set by each task. + # Fewer iterations per run buy more interleaved repetitions in the same + # time, which is what makes the median across repetitions meaningful. + # YJIT is on, as in production; without it the interpreter's per-call + # cost makes tracing look several times more expensive. The run fails if + # this ruby was built without YJIT. + - command: shell.exec + type: test + params: + shell: bash + working_dir: "src" + script: | + ${PREPARE_SHELL} + TEST_CMD="bundle exec rake driver_bench:compare" \ + CONFIGURATIONS="${driver_bench_configurations|off,api-only,sdk-never,sdk-parent-1pct,sdk-always}" \ + DRIVER_BENCH_TASKS="${driver_bench_tasks}" \ + REPS="${driver_bench_reps|10}" \ + DRIVER_BENCH_MAX_ITERATIONS="${driver_bench_max_iterations|20}" \ + RUBY_YJIT_ENABLE="${driver_bench_yjit|1}" \ + DRIVER_BENCH_MIN_TIME="${driver_bench_min_time|20}" \ + DRIVER_BENCH_ENFORCE_TARGETS="${driver_bench_enforce_targets|false}" \ + PERFORMANCE_RESULTS_FILE="$PROJECT_DIRECTORY/${task_name}.json" \ + .evergreen/run-tests.sh + # Sent to the performance store so that overhead_pct, cpu_us_added_per_op + # and allocs_added_per_op are watched for regressions, and kept as an + # artifact for comparison by hand. + - command: perf.send + params: + file: "${PROJECT_DIRECTORY}/${task_name}.json" + - command: s3.put + params: + aws_key: ${aws_key} + aws_secret: ${aws_secret} + local_file: ./src/${task_name}.json + display_name: ${task_name}.json + remote_file: ${UPLOAD_BUCKET}/${version_id}/${build_id}/artifacts/${build_variant}/${task_name}.json + content_type: application/json + permissions: public-read + bucket: mciuploads + + # Measures the per-attribute cost at the SDK boundary. The untraced + # operation is measured on the same host, so the relative figures are + # comparable with the absolute ones. No waterfall: request it in a patch. + # See DRIVERS-3620. + "run attribute cost": + - command: shell.exec + type: test + params: + shell: bash + working_dir: "src" + script: | + ${PREPARE_SHELL} + TEST_CMD="bundle exec rake otel_attributes:evergreen 2>&1 | tee $PROJECT_DIRECTORY/${task_name}.txt" \ + RUBY_YJIT_ENABLE=1 \ + ITERATIONS="${otel_attribute_iterations|50000}" \ + REPS="${otel_attribute_reps|20}" \ + .evergreen/run-tests.sh + - command: s3.put + params: + aws_key: ${aws_key} + aws_secret: ${aws_secret} + local_file: ./src/${task_name}.txt + display_name: ${task_name}.txt + remote_file: ${UPLOAD_BUCKET}/${version_id}/${build_id}/artifacts/${build_variant}/${task_name}.txt + content_type: text/plain + permissions: public-read + bucket: mciuploads + "run tests with orchestration and drivers tools": - command: subprocess.exec type: test @@ -579,6 +650,40 @@ tasks: - name: "driver-bench" commands: - func: "run benchmarks" + # Runs the two small-document micro-benchmarks the OpenTelemetry spec's + # targets are defined on under every driver configuration, and reports + # what each costs relative to the untraced baseline. Both run on one host, + # so their overheads can be compared with each other: hosts differ in + # speed by more than the overheads do. Find many creates too few spans per + # iteration to show tracing cost, and large-document inserts are dominated + # by the server. About two hours at the default ten reps. + - name: "driver-bench-otel" + exec_timeout_secs: 10800 + commands: + - func: "run benchmark comparison" + vars: + driver_bench_tasks: "small doc insertone,find one by id" + # Measures the cost of the command-span attribute shapes end to end, against + # the same two small-document micro-benchmarks. Every configuration records + # every trace; they differ only in how the command span's attributes are + # built. attr-none is the floor (a recorded span with no attributes), + # sdk-always is the driver as shipped, and the rest bracket it. Not on the + # waterfall: request it in a patch. See DRIVERS-3620. + - name: "driver-bench-otel-attributes" + exec_timeout_secs: 10800 + commands: + - func: "run benchmark comparison" + vars: + driver_bench_tasks: "small doc insertone,find one by id" + driver_bench_configurations: "off,sdk-always,attr-none,attr-creation-only,attr-all-at-creation,attr-none-then-all" + driver_bench_reps: "6" + # Measures the per-attribute cost at the SDK boundary, with the untraced + # operation measured on the same host. Not on the waterfall: request it in a + # patch. See DRIVERS-3620. + - name: "otel-attribute-cost" + exec_timeout_secs: 3600 + commands: + - func: "run attribute cost" - name: "test-csot" commands: - func: "run CSOT tests" diff --git a/.evergreen/config/standard.yml.erb b/.evergreen/config/standard.yml.erb index 1f4e48ae20..95ac6ea067 100644 --- a/.evergreen/config/standard.yml.erb +++ b/.evergreen/config/standard.yml.erb @@ -61,6 +61,44 @@ buildvariants: tasks: - name: "driver-bench" + # Compares the DriverBench micro-benchmarks with OpenTelemetry disabled + # against several enabled configurations. See DRIVERS-3620. + - matrix_name: DriverBenchOTel + matrix_spec: + ruby: <%= latest_ruby %> + mongodb-version: <%= latest_stable_mdb %> + topology: standalone + os: ubuntu2204 + display_name: DriverBench OTel + tasks: + - name: "driver-bench-otel" + batchtime: 1440 # run at most once a day + + # Measures the cost of the command-span attribute shapes end to end. Kept off + # the waterfall; request it in a patch. See DRIVERS-3620. + - matrix_name: DriverBenchOTelAttributes + matrix_spec: + ruby: <%= latest_ruby %> + mongodb-version: <%= latest_stable_mdb %> + topology: standalone + os: ubuntu2204 + display_name: DriverBench OTel attributes + tasks: + - name: "driver-bench-otel-attributes" + + # Measures the per-attribute cost at the SDK boundary, with the untraced + # operation on the same host. Kept off the waterfall; request it in a patch. + # See DRIVERS-3620. + - matrix_name: OTelAttributeCost + matrix_spec: + ruby: <%= latest_ruby %> + mongodb-version: <%= latest_stable_mdb %> + topology: standalone + os: ubuntu2204 + display_name: OTel attribute cost + tasks: + - name: "otel-attribute-cost" + - matrix_name: "auth/ssl" matrix_spec: auth-and-ssl: ["auth-and-ssl", "noauth-and-nossl"] @@ -168,6 +206,8 @@ buildvariants: tasks: - name: test-csot + # OTel spec tests run off-PR (no "pr" tag), like ruby-dev and + # DriverBench OTel: they only run on the master waterfall. - matrix_name: OTel matrix_spec: ruby: <%= latest_ruby %> @@ -175,7 +215,6 @@ buildvariants: topology: replica-set-single-node os: ubuntu2204 display_name: "OTel - ${mongodb-version}" - tags: ["pr"] tasks: - name: test-otel diff --git a/.evergreen/functions.sh b/.evergreen/functions.sh index 2f18a4d23f..c604632ecf 100644 --- a/.evergreen/functions.sh +++ b/.evergreen/functions.sh @@ -59,6 +59,22 @@ set_env_vars() { fi } +# Runs every Ruby process with YJIT when the installed ruby has it, as +# production deployments do. A ruby built without YJIT (JRuby, MRI older +# than 3.2) would only print a warning for RUBY_YJIT_ENABLE, so it is set +# only when YJIT is actually available. A value set by the caller (e.g. +# RUBY_YJIT_ENABLE=0 to run under the interpreter) is left alone. +enable_yjit() { + if test -n "$RUBY_YJIT_ENABLE"; then + echo "RUBY_YJIT_ENABLE=$RUBY_YJIT_ENABLE set by the caller" + elif ruby --yjit -e 'exit(RubyVM::YJIT.enabled? ? 0 : 1)' >/dev/null 2>&1; then + export RUBY_YJIT_ENABLE=1 + echo "YJIT enabled" + else + echo "YJIT not available in `ruby -v`" + fi +} + bundle_install() { args=--quiet diff --git a/.evergreen/run-tests-atlas-full.sh b/.evergreen/run-tests-atlas-full.sh index 99a8ed347a..10fb3093a1 100755 --- a/.evergreen/run-tests-atlas-full.sh +++ b/.evergreen/run-tests-atlas-full.sh @@ -18,6 +18,7 @@ export PATH="$HOME/.rbenv/bin:$PATH" eval "$(rbenv init - bash)" export FULL_RUBY_VERSION=$(ls ~/.rbenv/versions | head -n1) rbenv global $FULL_RUBY_VERSION +enable_yjit export JAVA_HOME=/opt/java/jdk21 export JAVACMD=$JAVA_HOME/bin/java diff --git a/.evergreen/run-tests-atlas.sh b/.evergreen/run-tests-atlas.sh index 038adcb56c..6d76a48efd 100755 --- a/.evergreen/run-tests-atlas.sh +++ b/.evergreen/run-tests-atlas.sh @@ -18,6 +18,7 @@ export PATH="$HOME/.rbenv/bin:$PATH" eval "$(rbenv init - bash)" export FULL_RUBY_VERSION=$(ls ~/.rbenv/versions | head -n1) rbenv global $FULL_RUBY_VERSION +enable_yjit export JAVA_HOME=/opt/java/jdk21 export JAVACMD=$JAVA_HOME/bin/java diff --git a/.evergreen/run-tests-azure.sh b/.evergreen/run-tests-azure.sh index 18cd5f3824..5f477c00fb 100755 --- a/.evergreen/run-tests-azure.sh +++ b/.evergreen/run-tests-azure.sh @@ -18,6 +18,7 @@ export PATH="$HOME/.rbenv/bin:$PATH" eval "$(rbenv init - bash)" export FULL_RUBY_VERSION=$(ls ~/.rbenv/versions | head -n1) rbenv global $FULL_RUBY_VERSION +enable_yjit export JAVA_HOME=/opt/java/jdk21 export JAVACMD=$JAVA_HOME/bin/java diff --git a/.evergreen/run-tests-deployed-lambda.sh b/.evergreen/run-tests-deployed-lambda.sh index 8a700e8eab..43343cb5b1 100755 --- a/.evergreen/run-tests-deployed-lambda.sh +++ b/.evergreen/run-tests-deployed-lambda.sh @@ -19,6 +19,7 @@ export PATH="$HOME/.rbenv/bin:$PATH" eval "$(rbenv init - bash)" export FULL_RUBY_VERSION=$(ls ~/.rbenv/versions | head -n1) rbenv global $FULL_RUBY_VERSION +enable_yjit export JAVA_HOME=/opt/java/jdk21 export JAVACMD=$JAVA_HOME/bin/java diff --git a/.evergreen/run-tests-ecs.sh b/.evergreen/run-tests-ecs.sh index adea6939a7..2bb04d42bc 100755 --- a/.evergreen/run-tests-ecs.sh +++ b/.evergreen/run-tests-ecs.sh @@ -28,6 +28,7 @@ export PATH="$HOME/.rbenv/bin:$PATH" eval "$(rbenv init - bash)" export FULL_RUBY_VERSION=$(ls ~/.rbenv/versions | head -n1) rbenv global $FULL_RUBY_VERSION +enable_yjit export JAVA_HOME=/opt/java/jdk21 export JAVACMD=$JAVA_HOME/bin/java diff --git a/.evergreen/run-tests-gcp.sh b/.evergreen/run-tests-gcp.sh index e8990dc375..77110a2462 100755 --- a/.evergreen/run-tests-gcp.sh +++ b/.evergreen/run-tests-gcp.sh @@ -19,6 +19,7 @@ export PATH="$HOME/.rbenv/bin:$PATH" eval "$(rbenv init - bash)" export FULL_RUBY_VERSION=$(ls ~/.rbenv/versions | head -n1) rbenv global $FULL_RUBY_VERSION +enable_yjit export JAVA_HOME=/opt/java/jdk21 export JAVACMD=$JAVA_HOME/bin/java diff --git a/.evergreen/run-tests-new.sh b/.evergreen/run-tests-new.sh index b331b97301..e1a5cfd433 100755 --- a/.evergreen/run-tests-new.sh +++ b/.evergreen/run-tests-new.sh @@ -29,6 +29,7 @@ export PATH="$HOME/.rbenv/bin:$PATH" eval "$(rbenv init - bash)" export FULL_RUBY_VERSION=$(ls ~/.rbenv/versions | head -n1) rbenv global $FULL_RUBY_VERSION +enable_yjit export JAVA_HOME=/opt/java/jdk21 export JAVACMD=$JAVA_HOME/bin/java diff --git a/.evergreen/run-tests-sfp.sh b/.evergreen/run-tests-sfp.sh index 979cd2a8fb..3c859bdcd6 100755 --- a/.evergreen/run-tests-sfp.sh +++ b/.evergreen/run-tests-sfp.sh @@ -18,6 +18,7 @@ export PATH="$HOME/.rbenv/bin:$PATH" eval "$(rbenv init - bash)" export FULL_RUBY_VERSION=$(ls ~/.rbenv/versions | head -n1) rbenv global $FULL_RUBY_VERSION +enable_yjit export JAVA_HOME=/opt/java/jdk21 export JAVACMD=$JAVA_HOME/bin/java diff --git a/.evergreen/run-tests.sh b/.evergreen/run-tests.sh index 0982136f50..fdd8bb6b2d 100755 --- a/.evergreen/run-tests.sh +++ b/.evergreen/run-tests.sh @@ -41,6 +41,7 @@ export PATH="$HOME/.rbenv/bin:$PATH" eval "$(rbenv init - bash)" export FULL_RUBY_VERSION=$(ls ~/.rbenv/versions | head -n1) rbenv global $FULL_RUBY_VERSION +enable_yjit export JAVA_HOME=/opt/java/jdk21 export JAVACMD=$JAVA_HOME/bin/java diff --git a/.mod/drivers-evergreen-tools b/.mod/drivers-evergreen-tools index 890a93bdb8..ec0b1497b3 160000 --- a/.mod/drivers-evergreen-tools +++ b/.mod/drivers-evergreen-tools @@ -1 +1 @@ -Subproject commit 890a93bdb88c595907813a9bfd0a56156d9f9f5b +Subproject commit ec0b1497b3d351acbfc4fd7579bf3de4de0a1046 diff --git a/Rakefile b/Rakefile index f7b43508da..12e4037a50 100644 --- a/Rakefile +++ b/Rakefile @@ -222,3 +222,4 @@ namespace :docs do end load 'profile/driver_bench/rake/tasks.rake' +load 'profile/otel_attributes/rake/tasks.rake' diff --git a/lib/mongo/monitoring/event/secure.rb b/lib/mongo/monitoring/event/secure.rb index 7cc8f38e68..01a92605d6 100644 --- a/lib/mongo/monitoring/event/secure.rb +++ b/lib/mongo/monitoring/event/secure.rb @@ -37,6 +37,11 @@ module Secure copydb ].freeze + # The hello / legacy hello command names, hoisted to a constant: the + # sensitivity check runs on every command event and an array literal + # would allocate on each call. + HELLO_COMMANDS = %w[hello ismaster isMaster].freeze + # Check whether the command is sensitive in terms of command monitoring # spec. A command is detected as sensitive if it is in the # list or if it is a hello/legacy hello command, and @@ -53,7 +58,7 @@ def sensitive?(command_name:, document:) # when speculativeAuthenticate is present, their commands AND replies # MUST be redacted from the events. # See https://github.com/mongodb/specifications/blob/master/source/command-logging-and-monitoring/command-logging-and-monitoring.md#security - %w[hello ismaster isMaster].include?(command_name.to_s) && + HELLO_COMMANDS.include?(command_name.to_s) && !!document['speculativeAuthenticate'] end diff --git a/lib/mongo/tracing/open_telemetry/command_tracer.rb b/lib/mongo/tracing/open_telemetry/command_tracer.rb index 27585839de..dfcaae3621 100644 --- a/lib/mongo/tracing/open_telemetry/command_tracer.rb +++ b/lib/mongo/tracing/open_telemetry/command_tracer.rb @@ -30,6 +30,9 @@ class CommandTracer # would only add noise to traces. HELLO_COMMANDS = %w[hello ismaster isMaster].freeze + # Upper bound of the lsid UUID memo cache (see #lsid). + LSID_CACHE_MAX = 128 + # Initializes a new CommandTracer. # # @param otel_tracer [ OpenTelemetry::Trace::Tracer ] the OpenTelemetry tracer. @@ -41,6 +44,8 @@ def initialize(otel_tracer, parent_tracer, query_text_max_length: 0) @otel_tracer = otel_tracer @parent_tracer = parent_tracer @query_text_max_length = query_text_max_length + @lsid_cache = {} + @lsid_cache_mutex = Mutex.new end # Starts a span for a MongoDB command. @@ -66,15 +71,30 @@ def start_span(message, operation_context, connection); end # @return [ Object ] the result of the command. # rubocop:disable Lint/RescueException def trace_command(message, _operation_context, connection) - return yield if skip_tracing?(message) + # The command document and name are extracted once and threaded + # through: the extraction helpers below allocate on every call. + doc = message.documents.first + name = command_name(doc) + return yield if skip_tracing?(doc, name) # Commands should always be nested under their operation span, not directly under # the transaction span. Don't pass with_parent to use automatic parent resolution # from the currently active span (the operation span). - span = create_command_span(message, connection) + span = create_command_span(name, doc, connection) + # An invalid context has no trace identity: it cannot be propagated, + # continued, or correlated with anything downstream, so every + # operation on it is waste. This is a state check on the span we + # were handed, not detection of whether the SDK is available — a + # custom API-only provider returning real spans sees the full path. + # Must not key on recording?: an unsampled-but-valid context still + # has to be made current for propagation. + return yield unless span.context.valid? + + cursor = cursor_id(name, doc) + apply_deferred_attributes(span, message, name, doc, cursor) if span.recording? ::OpenTelemetry::Trace.with_span(span) do |s, c| yield.tap do |result| - process_command_result(result, cursor_id(message), c, s) + process_command_result(result, cursor, c, s) end end rescue Exception => e @@ -93,26 +113,27 @@ def trace_command(message, _operation_context, connection) # command spans for them. Hello / legacy hello are also skipped to keep # handshake traffic out of traces. # - # @param message [ Mongo::Protocol::Message ] the command message. + # @param doc [ Hash ] the command document. + # @param name [ String ] the command name. # # @return [ Boolean ] true when no command span should be created. - def skip_tracing?(message) - name = command_name(message) + def skip_tracing?(doc, name) return true if HELLO_COMMANDS.include?(name) - sensitive?(command_name: name, document: message.documents.first) + sensitive?(command_name: name, document: doc) end # Creates a span for a command. # - # @param message [ Mongo::Protocol::Message ] the command message. + # @param name [ String ] the command name. + # @param doc [ Hash ] the command document. # @param connection [ Mongo::Server::Connection ] the connection. # # @return [ OpenTelemetry::Trace::Span ] the created span. - def create_command_span(message, connection) + def create_command_span(name, doc, connection) @otel_tracer.start_span( - command_name(message), - attributes: span_attributes(message, connection), + name, + attributes: span_attributes(doc, name, connection), kind: :client ) end @@ -142,61 +163,78 @@ def handle_command_exception(span, exception) span.status = ::OpenTelemetry::Trace::Status.error("Unhandled exception of type: #{exception.class}") end - # Builds span attributes for the command. + # Builds the attributes passed at span creation: the cheap, + # sampler-plausible set. Expensive attributes are deferred to + # apply_deferred_attributes — building them here would defeat the + # sampler's purpose, and on a non-recording span they would be + # discarded anyway. Keys whose value is nil are omitted rather than + # compacted afterwards, so the common case allocates nothing extra. # - # @param message [ Mongo::Protocol::Message ] the command message. + # @param doc [ Hash ] the command document. + # @param name [ String ] the command name. # @param connection [ Mongo::Server::Connection ] the connection. # # @return [ Hash ] OpenTelemetry span attributes following MongoDB semantic conventions. - def span_attributes(message, connection) - base_attributes(message) - .merge(connection_attributes(connection)) - .merge(session_attributes(message)) - .compact - end - - # Returns base database and command attributes. - # - # @param message [ Mongo::Protocol::Message ] the command message. - # - # @return [ Hash ] base span attributes. - def base_attributes(message) - { + def span_attributes(doc, name, connection) + attrs = { 'db.system.name' => 'mongodb', - 'db.namespace' => database(message), - 'db.collection.name' => collection_name(message), - 'db.command.name' => command_name(message), - 'db.query.summary' => query_summary(message), - 'db.query.text' => query_text(message) + 'db.namespace' => database(doc), + 'db.command.name' => name } + if (coll_name = collection_name(name, doc)) + attrs['db.collection.name'] = coll_name + end + attrs.merge(connection_attributes(connection)) end - # Returns connection-related attributes. + # Returns connection-related attributes, computed once per connection + # and frozen. Setting the ivar from here is a benign race: competing + # threads build identical frozen hashes. The value dies with the + # connection, so no cleanup is needed. # # @param connection [ Mongo::Server::Connection ] the connection. # # @return [ Hash ] connection span attributes. def connection_attributes(connection) - { - 'server.port' => connection.address.port, - 'server.address' => connection.address.host, - 'network.transport' => connection.transport.to_s, - 'db.mongodb.server_connection_id' => connection.server.description.server_connection_id, - 'db.mongodb.driver_connection_id' => connection.id - } + attrs = connection.instance_variable_get(:@otel_connection_attributes) + unless attrs + attrs = { + 'server.port' => connection.address.port, + 'server.address' => connection.address.host, + 'network.transport' => connection.transport.to_s, + 'db.mongodb.server_connection_id' => connection.description.server_connection_id, + 'db.mongodb.driver_connection_id' => connection.id + }.freeze + connection.instance_variable_set(:@otel_connection_attributes, attrs) + end + attrs end - # Returns session and transaction attributes. + # Sets the expensive attributes after span creation. Only called for + # recording spans: on non-recording spans set_attribute discards. + # Note for reviewers: these attributes are invisible to the sampler, + # which only sees what is passed to start_span. The built-in samplers + # do not read attributes, and db.query.text is off by default. # + # @param span [ OpenTelemetry::Trace::Span ] the current span. # @param message [ Mongo::Protocol::Message ] the command message. - # - # @return [ Hash ] session span attributes. - def session_attributes(message) - { - 'db.mongodb.cursor_id' => cursor_id(message), - 'db.mongodb.lsid' => lsid(message), - 'db.mongodb.txn_number' => txn_number(message) - } + # @param name [ String ] the command name. + # @param doc [ Hash ] the command document. + # @param cursor [ Integer | nil ] the cursor id, extracted once per command. + def apply_deferred_attributes(span, message, name, doc, cursor) + span.set_attribute('db.query.summary', query_summary(name, doc)) + if (text = query_text(message)) + span.set_attribute('db.query.text', text) + end + if (lsid_value = lsid(doc)) + span.set_attribute('db.mongodb.lsid', lsid_value) + end + unless cursor.nil? + span.set_attribute('db.mongodb.cursor_id', cursor) + end + if (txn = txn_number(doc)) + span.set_attribute('db.mongodb.txn_number', txn) + end end # Processes cursor context from the command result. @@ -242,51 +280,67 @@ def maybe_trace_error(result, span) # Generates a summary string for the query. # - # @param message [ Mongo::Protocol::Message ] the command message. + # @param name [ String ] the command name. + # @param doc [ Hash ] the command document. # # @return [ String ] summary in format "command_name db.collection" or "command_name db". - def query_summary(message) - if (coll_name = collection_name(message)) - "#{command_name(message)} #{database(message)}.#{coll_name}" + def query_summary(name, doc) + if (coll_name = collection_name(name, doc)) + "#{name} #{database(doc)}.#{coll_name}" else - "#{command_name(message)} #{database(message)}" + "#{name} #{database(doc)}" end end - # Extracts the collection name from the command message. + # Extracts the collection name from the command document. # - # @param message [ Mongo::Protocol::Message ] the command message. + # @param name [ String ] the command name. + # @param doc [ Hash ] the command document. # # @return [ String | nil ] the collection name, or nil if not applicable. - def collection_name(message) - case command_name(message) + def collection_name(name, doc) + case name when 'getMore' - message.documents.first['collection'].to_s + doc['collection'].to_s when 'listCollections', 'listDatabases', 'commitTransaction', 'abortTransaction' nil else - value = message.documents.first.values.first + # Iterate instead of using doc.values.first: the block form is + # allocation-free (see #command_name). + value = nil + # rubocop:disable Lint/UnreachableLoop -- intentional: only the first entry is needed + doc.each_value do |v| + value = v + break + end + # rubocop:enable Lint/UnreachableLoop # Return nil if the value is not a string (e.g., for admin commands that have numeric values) value.is_a?(String) ? value : nil end end - # Extracts the command name from the message. + # Extracts the command name from the command document. Iterates + # instead of using doc.keys.first: the block form is allocation-free, + # while keys builds an array of every top-level key per call — and + # this runs on every traced command. # - # @param message [ Mongo::Protocol::Message ] the command message. + # @param doc [ Hash ] the command document. # # @return [ String ] the command name. - def command_name(message) - message.documents.first.keys.first.to_s + def command_name(doc) + # rubocop:disable Lint/UnreachableLoop -- intentional: only the first entry is needed + doc.each_key { |key| return key.to_s } + # rubocop:enable Lint/UnreachableLoop + '' end - # Extracts the database name from the message. + # Extracts the database name from the command document. # - # @param message [ Mongo::Protocol::Message ] the command message. + # @param doc [ Hash ] the command document. # # @return [ String ] the database name. - def database(message) - message.documents.first['$db'].to_s + def database(doc) + doc['$db'].to_s end # Checks if query text capture is enabled. @@ -298,34 +352,48 @@ def query_text? # Extracts the cursor ID from getMore commands. # - # @param message [ Mongo::Protocol::Message ] the command message. + # @param name [ String ] the command name. + # @param doc [ Hash ] the command document. # # @return [ Integer | nil ] the cursor ID, or nil if not a getMore command. - def cursor_id(message) - return unless command_name(message) == 'getMore' + def cursor_id(name, doc) + return unless name == 'getMore' - message.documents.first['getMore'].value + doc['getMore'].value end - # Extracts the logical session ID from the command. + # Extracts the logical session ID from the command. The UUID string + # is formatted once per session id and memoized in a bounded cache: + # the value is invariant for the life of a session, and formatting it + # per command showed up in the DRIVERS-3620 profile. # - # @param message [ Mongo::Protocol::Message ] the command message. + # @param doc [ Hash ] the command document. # - # @return [ BSON::Binary | nil ] the session ID, or nil if not present. - def lsid(message) - lsid_doc = message.documents.first['lsid'] + # @return [ String | nil ] the session ID as a UUID string, or nil if not present. + def lsid(doc) + lsid_doc = doc['lsid'] return unless lsid_doc - lsid_doc['id'].to_uuid + binary = lsid_doc['id'] + key = binary.data + cached = @lsid_cache_mutex.synchronize { @lsid_cache[key] } + return cached if cached + + uuid = binary.to_uuid + @lsid_cache_mutex.synchronize do + @lsid_cache.clear if @lsid_cache.size >= LSID_CACHE_MAX + @lsid_cache[key] = uuid + end + uuid end # Extracts the transaction number from the command. # - # @param message [ Mongo::Protocol::Message ] the command message. + # @param doc [ Hash ] the command document. # # @return [ Integer | nil ] the transaction number, or nil if not present. - def txn_number(message) - txn_num = message.documents.first['txnNumber'] + def txn_number(doc) + txn_num = doc['txnNumber'] return unless txn_num txn_num.value diff --git a/lib/mongo/tracing/open_telemetry/operation_tracer.rb b/lib/mongo/tracing/open_telemetry/operation_tracer.rb index 2800fcd1e6..4a9b91db5b 100644 --- a/lib/mongo/tracing/open_telemetry/operation_tracer.rb +++ b/lib/mongo/tracing/open_telemetry/operation_tracer.rb @@ -37,6 +37,8 @@ class OperationTracer def initialize(otel_tracer, parent_tracer) @otel_tracer = otel_tracer @parent_tracer = parent_tracer + @operation_names = {} + @operation_names_mutex = Mutex.new end # Trace a MongoDB operation. @@ -57,6 +59,20 @@ def initialize(otel_tracer, parent_tracer) # rubocop:disable Lint/RescueException def trace_operation(operation, operation_context, op_name: nil, &block) span = create_operation_span(operation, operation_context, op_name) + # An invalid context has no trace identity: it cannot be propagated, + # continued, or correlated with anything downstream, so every + # operation on it is waste. This is a state check on the span we + # were handed, not detection of whether the SDK is available — a + # custom API-only provider returning real spans sees the full path. + # Must not key on recording?: an unsampled-but-valid context still + # has to be made current for propagation. Skipping here also leaves + # the cursor context map untouched, consistent with tracing being + # effectively off. + return yield unless span.context.valid? + + if span.recording? && !operation.cursor_id.nil? + span.set_attribute('db.mongodb.cursor_id', operation.cursor_id) + end execute_with_span(span, operation, &block) rescue Exception => e handle_span_exception(span, e) @@ -68,7 +84,9 @@ def trace_operation(operation, operation_context, op_name: nil, &block) private - # Creates an OpenTelemetry span for the operation. + # Creates an OpenTelemetry span for the operation. The operation name, + # collection name and span name are each computed once and reused + # between the span name and the attributes. # # @param operation [ Mongo::Operation ] the operation. # @param operation_context [ Mongo::Operation::Context ] the operation context. @@ -76,10 +94,13 @@ def trace_operation(operation, operation_context, op_name: nil, &block) # # @return [ OpenTelemetry::Trace::Span ] the created span. def create_operation_span(operation, operation_context, op_name) - parent_context = parent_context_for(operation_context, operation.cursor_id) + parent_context = parent_context_for(operation_context) + name = operation_name(operation, op_name) + coll_name = collection_name(operation) + span_name = operation_span_name(name, operation.db_name, coll_name) @otel_tracer.start_span( - operation_span_name(operation, op_name), - attributes: span_attributes(operation, op_name), + span_name, + attributes: span_attributes(operation, name, span_name, coll_name), with_parent: parent_context, kind: :client ) @@ -115,30 +136,45 @@ def handle_span_exception(span, exception) end # Returns the operation name from the provided name or operation class. + # The class-derived name is memoized per class: it is invariant and + # deriving it (split + downcase) allocates on every call. The number + # of operation classes is small and fixed, so the cache is unbounded + # by design. Access is synchronized: unsynchronized Hash mutation is + # not safe on all Ruby runtimes (e.g. JRuby). # # @param operation [ Mongo::Operation ] the operation. # @param op_name [ String | nil ] optional operation name. # # @return [ String ] the operation name in lowercase. def operation_name(operation, op_name = nil) - op_name || operation.class.name.split('::').last.downcase + return op_name if op_name + + klass = operation.class + @operation_names_mutex.synchronize do + @operation_names[klass] ||= klass.name.split('::').last.downcase + end end - # Builds span attributes for the operation. + # Builds the attributes passed at span creation: the cheap set. The + # cursor id is deferred to trace_operation (behind recording?) — on a + # non-recording span set_attribute discards, and the sampler has no + # use for it. # # @param operation [ Mongo::Operation ] the operation. - # @param op_name [ String | nil ] optional operation name. + # @param name [ String ] the operation name. + # @param span_name [ String ] the span name, reused as the summary. + # @param coll_name [ String | nil ] the collection name, computed once. # # @return [ Hash ] OpenTelemetry span attributes following MongoDB semantic conventions. - def span_attributes(operation, op_name) - { + def span_attributes(operation, name, span_name, coll_name) + attrs = { 'db.system.name' => 'mongodb', 'db.namespace' => operation.db_name.to_s, - 'db.collection.name' => collection_name(operation), - 'db.operation.name' => operation_name(operation, op_name), - 'db.operation.summary' => operation_span_name(operation, op_name), - 'db.mongodb.cursor_id' => operation.cursor_id, - }.compact + 'db.operation.name' => name, + 'db.operation.summary' => span_name + } + attrs['db.collection.name'] = coll_name unless coll_name.nil? + attrs end # Processes cursor context after operation execution. @@ -209,16 +245,16 @@ def collection_key_for_operation(operation) # Generates the span name for the operation. # - # @param operation [ Mongo::Operation ] the operation. - # @param op_name [ String | nil ] optional operation name. + # @param name [ String ] the operation name. + # @param db_name [ String ] the database name. + # @param coll_name [ String | nil ] the collection name, if any. # # @return [ String ] span name in format "operation_name db.collection" or "operation_name db". - def operation_span_name(operation, op_name = nil) - coll_name = collection_name(operation) + def operation_span_name(name, db_name, coll_name) if coll_name && !coll_name.empty? - "#{operation_name(operation, op_name)} #{operation.db_name}.#{coll_name}" + "#{name} #{db_name}.#{coll_name}" else - "#{operation_name(operation, op_name)} #{operation.db_name}" + "#{name} #{db_name}" end end end diff --git a/lib/mongo/tracing/open_telemetry/tracer.rb b/lib/mongo/tracing/open_telemetry/tracer.rb index 1def549e52..584df8b2ea 100644 --- a/lib/mongo/tracing/open_telemetry/tracer.rb +++ b/lib/mongo/tracing/open_telemetry/tracer.rb @@ -146,34 +146,20 @@ def cursor_context_map @cursor_context_map ||= {} end - # Generates a unique key for cursor tracking in the context map. - # - # @param session [ Mongo::Session ] the session associated with the cursor. - # @param cursor_id [ Integer ] the cursor ID. - # - # @return [ String | nil ] unique key combining session ID and cursor ID, or nil if either is nil. - def cursor_map_key(session, cursor_id) - return if cursor_id.nil? || session.nil? - - "#{session.session_id['id'].to_uuid}-#{cursor_id}" - end - # Determines the parent OpenTelemetry context for an operation. # - # Returns the transaction context if the operation is part of a transaction, - # otherwise returns nil. Cursor-based context nesting is not currently implemented. + # Returns the transaction context if the operation is part of a + # transaction, otherwise returns nil. Cursor operations deliberately + # have no parent here: a caller-driven getMore must not be nested + # under the span of the operation that created the cursor. # # @param operation_context [ Mongo::Operation::Context ] the operation context. - # @param cursor_id [ Integer ] the cursor ID, if applicable. # # @return [ OpenTelemetry::Context | nil ] parent context or nil. - def parent_context_for(operation_context, cursor_id) - if (key = transaction_map_key(operation_context.session)) - transaction_context_map[key] - elsif (_key = cursor_map_key(operation_context.session, cursor_id)) - # We return nil here unless we decide how to nest cursor operations. - nil - end + def parent_context_for(operation_context) + return unless (key = transaction_map_key(operation_context.session)) + + transaction_context_map[key] end # Returns the transaction context map for tracking active transaction contexts. @@ -200,6 +186,12 @@ def transaction_token_map # Generates a unique key for transaction tracking. # # Returns nil for implicit sessions or sessions not in a transaction. + # The key is invariant for the life of a transaction but is needed on + # every operation span created inside it, and formatting it (a UUID + # plus interpolation) showed the same per-call cost the lsid cache + # removed from the command path. It is memoized on the session, keyed + # by the transaction number, so each transaction formats it once. A + # benign race may format it twice; the results are identical. # # @param session [ Mongo::Session ] the session. # @@ -207,7 +199,13 @@ def transaction_token_map def transaction_map_key(session) return if session.nil? || session.implicit? || !session.in_transaction? - "#{session.session_id['id'].to_uuid}-#{session.txn_num}" + txn_num = session.txn_num + cached = session.instance_variable_get(:@otel_transaction_key) + return cached[1] if cached && cached[0] == txn_num + + key = "#{session.session_id['id'].to_uuid}-#{txn_num}" + session.instance_variable_set(:@otel_transaction_key, [ txn_num, key ].freeze) + key end private diff --git a/profile/driver_bench/attribute_profiles.rb b/profile/driver_bench/attribute_profiles.rb new file mode 100644 index 0000000000..1af1e4af0e --- /dev/null +++ b/profile/driver_bench/attribute_profiles.rb @@ -0,0 +1,111 @@ +# frozen_string_literal: true + +module Mongo + module DriverBench + # Reopens the command tracer so the benchmark can measure alternative + # command-span attribute shapes end to end. + # + # Loaded only when OTEL_ATTRIBUTE_PROFILE names a profile, and only by the + # benchmark. The driver never loads this file, and the tracer is untouched + # for every configuration that does not set a profile. + # + # Profiles vary the command span only. The operation span carries a small, + # fixed set of attributes, and the question this answers is about the + # command span's attribute list. + # + # @api private + module CommandAttributeProfiles + # The profiles, and what each does to a command span. + # + # none zero attributes at creation, nothing deferred + # creation-only the driver's creation-time set, nothing deferred + # all-at-creation every attribute passed to start_span + # none-then-all zero attributes at creation, all set afterwards + # + # A nil profile leaves the tracer alone; that is the driver as shipped. + PROFILES = %w[ none creation-only all-at-creation none-then-all ].freeze + + class << self + # @return [ String | nil ] the active profile. + attr_reader :profile + + # Applies the profile named by OTEL_ATTRIBUTE_PROFILE, if any. + def apply! + name = ENV['OTEL_ATTRIBUTE_PROFILE'] + return if name.nil? || name.empty? + + raise ArgumentError, "unknown attribute profile #{name.inspect}" unless PROFILES.include?(name) + + @profile = name + Mongo::Tracing::OpenTelemetry::CommandTracer.prepend(CommandTracerOverride) + end + end + + # Prepended into CommandTracer when a profile is active. + module CommandTracerOverride + # Builds the creation-time attributes for the active profile. + def span_attributes(doc, name, connection) + case CommandAttributeProfiles.profile + when 'none', 'none-then-all' + # Keep the connection so the deferred pass can build the + # connection attributes; do not build anything now. + Thread.current[:otel_bench_connection] = connection + {} + when 'all-at-creation' + super.merge(deferred_attributes(name, doc)) + else + super + end + end + + # Builds the deferred attributes for the active profile. + def apply_deferred_attributes(span, message, name, doc, cursor) + case CommandAttributeProfiles.profile + when 'none', 'creation-only', 'all-at-creation' + nil + when 'none-then-all' + connection = Thread.current[:otel_bench_connection] + creation_attributes(name, doc, connection).each { |key, value| span.set_attribute(key, value) } + deferred_attributes(name, doc, cursor).each { |key, value| span.set_attribute(key, value) } + else + super + end + end + + private + + # A copy of the driver's creation-time set, built independently so that + # the none-then-all profile can defer it. Kept in step with + # CommandTracer#span_attributes by hand; it is only used by the + # benchmark. + def creation_attributes(name, doc, connection) + attrs = { + 'db.system.name' => 'mongodb', + 'db.namespace' => database(doc), + 'db.command.name' => name + } + if (coll_name = collection_name(name, doc)) + attrs['db.collection.name'] = coll_name + end + attrs.merge(connection_attributes(connection)) + end + + # A copy of the driver's deferred set. query_text is left out: it is + # off by default and its value is built from the message, which + # span_attributes does not receive. + def deferred_attributes(name, doc, cursor = nil) + cursor ||= cursor_id(name, doc) + attrs = { 'db.query.summary' => query_summary(name, doc) } + if (lsid_value = lsid(doc)) + attrs['db.mongodb.lsid'] = lsid_value + end + attrs['db.mongodb.cursor_id'] = cursor unless cursor.nil? + if (txn = txn_number(doc)) + attrs['db.mongodb.txn_number'] = txn + end + attrs + end + end + end + end +end diff --git a/profile/driver_bench/base.rb b/profile/driver_bench/base.rb index b65de595ac..e9f30c16c2 100644 --- a/profile/driver_bench/base.rb +++ b/profile/driver_bench/base.rb @@ -3,6 +3,7 @@ require 'benchmark' require 'mongo' +require_relative 'configuration' require_relative 'percentiles' module Mongo @@ -41,11 +42,20 @@ def self.bench_name(benchmark_name = nil) # Instantiate a new micro-benchmark class. def initialize - @max_iterations = debug_mode? ? 10 : 100 - @min_time = debug_mode? ? 1 : 60 + @max_iterations = Integer(ENV['DRIVER_BENCH_MAX_ITERATIONS'] || (debug_mode? ? 10 : 100)) + @min_time = Float(ENV['DRIVER_BENCH_MIN_TIME'] || (debug_mode? ? 1 : 60)) @max_time = 300 # 5 minutes end + # The number of driver operations one iteration performs, or nil when + # the task has no meaningful per-operation unit. Per-operation metrics + # are only reported for tasks that define it. + # + # @return [ Integer | nil ] operations per iteration. + def ops_per_iteration + nil + end + def debug_mode? ENV['PERF_DEBUG'] end @@ -60,8 +70,32 @@ def run score = dataset_size / percentiles[50] / 1_000_000.0 { name: self.class.bench_name, + configuration: Configuration.current.name, score: score, - percentiles: percentiles } + percentiles: percentiles, + per_op: per_op_metrics(percentiles) } + end + + # Runs one iteration with every started span counted, and returns the + # number of spans per operation. Kept apart from #run: counting needs a + # span processor, which the timed runs deliberately do without. + # + # @param counter [ #reset, #count ] the span counter the tracer + # provider reports to. + # + # @return [ Float | nil ] spans per operation, or nil when the task + # defines no per-operation unit. + def count_spans(counter) + return nil unless ops_per_iteration + + setup + before_task + counter.reset + do_task + spans = counter.count + after_task + teardown + spans.to_f / ops_per_iteration end private @@ -82,9 +116,13 @@ def run_benchmark setup + @cpu_times = [] + @gc_times = [] + @allocations = [] + loop do before_task - timing = consider_gc { Benchmark.realtime { debug_mode? ? sleep(0.1) : do_task } } + timing = consider_gc { measure_iteration { debug_mode? ? sleep(0.1) : do_task } } after_task iteration_count += 1 @@ -104,9 +142,64 @@ def run_benchmark end end + # Times one iteration, and records alongside the wall-clock time the + # process CPU time and the number of objects allocated. CPU time leaves + # out the time spent waiting on the server, and allocations are nearly + # deterministic, so both show driver-side cost with far less noise than + # throughput. Process CPU time includes the driver's background threads + # (e.g. server monitoring), which cost the same in every configuration. + # + # GC time is recorded too, where the runtime reports it (Ruby 3.1+). + # When and how long the collector runs varies between processes more + # than the cost of tracing does, so CPU time without GC is the steadier + # measure of the driver's own work. + # + # @return [ Float ] the wall-clock time in seconds. + def measure_iteration(&block) + allocated = GC.stat(:total_allocated_objects) + gc = gc_time + cpu = Process.clock_gettime(Process::CLOCK_PROCESS_CPUTIME_ID) + timing = Benchmark.realtime(&block) + @cpu_times.push(Process.clock_gettime(Process::CLOCK_PROCESS_CPUTIME_ID) - cpu) + @gc_times.push(gc_time - gc) if gc + @allocations.push(GC.stat(:total_allocated_objects) - allocated) + timing + end + + # @return [ Float | nil ] the total time spent in GC so far, in + # seconds, or nil when the runtime does not report it. + def gc_time + ms = GC.stat[:time] + ms && (ms / 1000.0) + end + + # Per-operation medians, for tasks that define ops_per_iteration. + # + # @return [ Hash ] metric name to value. + def per_op_metrics(timings) + return {} unless ops_per_iteration + + ops = ops_per_iteration.to_f + { + 'wall_us_per_op' => timings[50] / ops * 1_000_000, + 'cpu_us_per_op' => Percentiles.new(@cpu_times)[50] / ops * 1_000_000, + 'allocs_per_op' => Percentiles.new(@allocations)[50] / ops + }.merge(gc_metrics(ops)) + end + + def gc_metrics(ops) + return {} if @gc_times.empty? + + cpu_ex_gc = @cpu_times.zip(@gc_times).map { |cpu, gc| cpu - gc } + { + 'gc_us_per_op' => Percentiles.new(@gc_times)[50] / ops * 1_000_000, + 'cpu_ex_gc_us_per_op' => Percentiles.new(cpu_ex_gc)[50] / ops * 1_000_000 + } + end + # Instantiate a new client. def new_client(uri = ENV['MONGODB_URI']) - Mongo::Client.new(uri) + Mongo::Client.new(uri, Configuration.current.client_options) end # Takes care of garbage collection considerations before diff --git a/profile/driver_bench/comparison.rb b/profile/driver_bench/comparison.rb new file mode 100644 index 0000000000..7f09cb5d30 --- /dev/null +++ b/profile/driver_bench/comparison.rb @@ -0,0 +1,341 @@ +# frozen_string_literal: true + +require 'fileutils' +require 'json' +require 'tmpdir' + +require_relative 'configuration' +require_relative 'percentiles' +require_relative 'suite' + +module Mongo + module DriverBench + # Runs the DriverBench micro-benchmarks under several driver + # configurations and reports what each configuration costs relative to the + # baseline. + # + # OpenTelemetry can only be installed into a process once, and the + # "api-only" configuration requires that the SDK was never loaded, so each + # configuration is measured in its own subprocess. + # + # Configurations are run interleaved -- the whole set, then the whole set + # again -- rather than one configuration to completion and then the next. + # The quantity being measured is a difference between two configurations, + # and a machine that slows down halfway through a run would otherwise put + # all of one configuration's samples in the fast half and all of another's + # in the slow half, which shows up as overhead that is not there. + # + # The order is also rotated by one position each repetition, so that no + # configuration always runs in the same slot. With a fixed order, a + # slot-dependent effect showed up as a bias of its own: sdk-parent-1pct + # measured cheaper than sdk-never in every run, though it does strictly + # more work. When the number of repetitions is a multiple of the number + # of configurations, each one runs in every slot equally often. + # + # Parameterised by the environment: + # + # CONFIGURATIONS comma-separated configuration names (default: all) + # REPS repetitions of the whole set (default 1) + # MONGODB_URI server to benchmark + # DRIVER_BENCH_TASKS + # DRIVER_BENCH_EXCLUDE_TASKS + # narrow the micro-benchmarks, as for Suite. When + # neither is set, BSON micro-benchmarks are skipped: + # they never talk to a server, so they create no spans + # and can only add noise to the comparison. + # PERFORMANCE_RESULTS_FILE + # where the combined results are written + # DRIVER_BENCH_ENFORCE_TARGETS + # fail when a configuration exceeds its target + # DRIVER_BENCH_MAX_ITERATIONS, DRIVER_BENCH_MIN_TIME + # passed through to every run, see Base + # RUBY_YJIT_ENABLE + # run with YJIT; fails if this ruby was built without it + # + # @api private + class Comparison + DEFAULT_EXCLUDED_TASKS = 'BSON' + DEFAULT_RESULTS_FILE = 'perf-comparison.json' + + # Metrics carried from every child run into the comparison. Each is + # reduced to its median across repetitions. + CARRIED_METRICS = %w[ + score wall_us_per_op cpu_us_per_op cpu_ex_gc_us_per_op gc_us_per_op allocs_per_op + time-90% time-99% spans_per_op + ].freeze + + TABLE_FORMAT = '%-22s %-16s %9s %8s %-4s ' \ + '%15s %11s %9s %9s %9s %6s' + + def self.run! + new.run + end + + def initialize(env = ENV) + @env = env + @reps = Integer(env['REPS'] || 1) + @configurations = requested_configurations + @samples = Hash.new { |hash, key| hash[key] = Hash.new { |h, k| h[k] = [] } } + end + + # Runs every configuration, @reps times, and reports the comparison. + # + # @return [ String ] the summary table. + def run + raise 'the baseline configuration is required to compare against' unless + @configurations.any?(&:baseline?) + + # Checked here as well as in every child, to fail before the first + # of many long runs rather than inside it. + Configuration.check_jit! + + Dir.mktmpdir('driver-bench') do |dir| + 1.upto(@reps) do |rep| + @configurations.rotate(rep - 1).each { |configuration| measure(configuration, rep, dir) } + end + end + + save_perf_data(compile_perf_data) + summary = summarize + enforce_targets!(summary) + summary + end + + private + + def requested_configurations + names = (@env['CONFIGURATIONS'] || '').split(',').map(&:strip).reject(&:empty?) + return Configuration::ALL if names.empty? + + names.map { |name| Configuration[name] } + end + + def baseline + @baseline ||= @configurations.find(&:baseline?) + end + + # Runs one configuration once, in its own process, and records the + # metrics of every micro-benchmark it reported. + def measure(configuration, rep, dir) + results_file = File.join(dir, "#{configuration.name}-#{rep}.json") + + puts format("\n----- rep %d/%d: %s -----", rep, @reps, configuration.name) + unless system(child_env(configuration, rep, results_file), 'bundle', 'exec', 'rake', 'driver_bench:run') + raise "configuration #{configuration.name} failed in rep #{rep}" + end + + JSON.parse(File.read(results_file)).each do |entry| + samples = @samples[[ configuration.task_name(entry['info']['test_name']), configuration.name ]] + entry['metrics'].each do |metric| + samples[metric['name']] << metric['value'] if CARRIED_METRICS.include?(metric['name']) + end + end + end + + def child_env(configuration, rep, results_file) + { + Configuration::ENV_VAR => configuration.name, + 'PERFORMANCE_RESULTS_FILE' => results_file, + 'DRIVER_BENCH_EXCLUDE_TASKS' => excluded_tasks, + # Span counts do not vary between repetitions; count them once. + 'DRIVER_BENCH_COUNT_SPANS' => (rep == 1).to_s + } + end + + def excluded_tasks + return @env['DRIVER_BENCH_EXCLUDE_TASKS'].to_s if + @env['DRIVER_BENCH_EXCLUDE_TASKS'] || @env['DRIVER_BENCH_TASKS'] + + DEFAULT_EXCLUDED_TASKS + end + + # The median of one metric across repetitions, for the same reason the + # spec takes the median across iterations. Unlike the spec's + # nearest-rank percentile, which suits its iteration counts, this + # averages the middle two of an even number of repetitions: nearest + # rank would always pick the lower one. + def median(task, configuration, metric) + values = @samples[[ task, configuration.name ]][metric].sort + return nil if values.empty? + + mid = values.length / 2 + values.length.odd? ? values[mid] : (values[mid - 1] + values[mid]) / 2.0 + end + + def spread(task, configuration) + values = @samples[[ task, configuration.name ]]['score'] + return nil if values.empty? + + values.minmax + end + + # Scores are throughput, so enabling a feature shows up as a loss. + # + # @return [ Float | nil ] the percentage of throughput given up, + # relative to the baseline. + def loss_for(task, configuration) + reference = median(task, baseline, 'score') + score = median(task, configuration, 'score') + return nil if reference.nil? || score.nil? || reference.zero? + + (reference - score) / reference * 100.0 + end + + # CPU time and allocations grow with cost, so the overhead is the + # difference over the baseline. + def added(task, configuration, metric) + reference = median(task, baseline, metric) + value = median(task, configuration, metric) + return nil if reference.nil? || value.nil? + + value - reference + end + + def cpu_overhead_pct(task, configuration, metric = 'cpu_us_per_op') + reference = median(task, baseline, metric) + delta = added(task, configuration, metric) + return nil if delta.nil? || reference.zero? + + delta / reference * 100.0 + end + + # Spans are only countable under a configuration that records every + # span, but the count is a property of the task, so it is reported + # for the task as a whole. + def spans_per_op(task) + @configurations.each do |configuration| + value = median(task, configuration, 'spans_per_op') + return value unless value.nil? + end + nil + end + + # Targets are checked against the CPU overhead (excluding GC where it + # is reported), not the throughput loss. Throughput includes waiting on + # the server and moved by 2-3 points between runs, enough to flip a + # configuration across its target; it also ranked configurations in an + # order their work rules out (sdk-never above api-only). + def target_met(task, configuration) + overhead = cpu_overhead_pct(task, configuration, cpu_metric(task)) + return nil if overhead.nil? || configuration.target_pct.nil? + + overhead <= configuration.target_pct + end + + def tasks + @tasks ||= @samples.keys.map(&:first).uniq + end + + def compile_perf_data + tasks.flat_map do |task| + @configurations.filter_map do |configuration| + next if median(task, configuration, 'score').nil? + + { + 'info' => { + 'test_name' => configuration.perf_test_name(task), + 'args' => {} + }, + 'metrics' => metrics_for(task, configuration) + } + end + end + end + + def metrics_for(task, configuration) + metrics = CARRIED_METRICS.filter_map do |name| + value = median(task, configuration, name) + { 'name' => name, 'value' => value } unless value.nil? + end + min, max = spread(task, configuration) + metrics << { 'name' => 'score_min', 'value' => min } << { 'name' => 'score_max', 'value' => max } + return metrics if configuration.baseline? + + # The overheads are recorded as metrics of their own, rather than + # left to be derived by comparing two time series later, so that + # they can be watched for regressions directly and so that + # host-to-host variation cancels out of them. + { + 'overhead_pct' => loss_for(task, configuration), + 'cpu_overhead_pct' => cpu_overhead_pct(task, configuration), + 'cpu_us_added_per_op' => added(task, configuration, 'cpu_us_per_op'), + 'cpu_ex_gc_overhead_pct' => cpu_overhead_pct(task, configuration, 'cpu_ex_gc_us_per_op'), + 'cpu_ex_gc_us_added_per_op' => added(task, configuration, 'cpu_ex_gc_us_per_op'), + 'gc_us_added_per_op' => added(task, configuration, 'gc_us_per_op'), + 'allocs_added_per_op' => added(task, configuration, 'allocs_per_op'), + 'target_pct' => configuration.target_pct + }.each do |name, value| + metrics << { 'name' => name, 'value' => value } unless value.nil? + end + metrics + end + + def save_perf_data(data, file_name: @env['PERFORMANCE_RESULTS_FILE'] || DEFAULT_RESULTS_FILE) + File.write(file_name, data.to_json) + end + + def summarize + lines = [ format("\n===== Configuration comparison (%d rep%s, median) =====", + @reps, (@reps == 1) ? '' : 's') ] + lines << 'loss: throughput given up vs off; target: max cpu %; spread: min-max MB/s across reps; ' \ + 'cpu: CPU us added per op, excluding GC; gc: GC us added per op; ' \ + 'allocs: objects added per op' + lines << '' + lines << format(TABLE_FORMAT, + task: 'micro-benchmark', config: 'configuration', loss: 'loss %', + target: 'target', verdict: 'ok?', spread: 'spread MB/s', + cpu: 'cpu us/op', cpu_pct: 'cpu %', gc: 'gc us/op', allocs: 'allocs', spans: 'spans') + tasks.each do |task| + @configurations.reject(&:baseline?).each { |configuration| lines << row(task, configuration) } + end + + lines.join("\n") + end + + def row(task, configuration) + min, max = spread(task, configuration) + + format(TABLE_FORMAT, + task: task, config: configuration.name, + loss: signed(loss_for(task, configuration)), + target: configuration.target_pct ? format('%g', configuration.target_pct) : '-', + verdict: verdict(target_met(task, configuration)), + spread: min ? format('%.4g-%.4g', min, max) : 'n/a', + cpu: signed(added(task, configuration, cpu_metric(task))), + cpu_pct: signed(cpu_overhead_pct(task, configuration, cpu_metric(task))), + gc: signed(added(task, configuration, 'gc_us_per_op')), + allocs: signed(added(task, configuration, 'allocs_per_op'), '%+.0f'), + spans: (spans = spans_per_op(task)) ? format('%.2g', spans) : 'n/a') + end + + # CPU time excluding GC where the runtime reports GC time, total CPU + # time otherwise. + def cpu_metric(task) + median(task, baseline, 'cpu_ex_gc_us_per_op') ? 'cpu_ex_gc_us_per_op' : 'cpu_us_per_op' + end + + def signed(value, pattern = '%+.2f') + value.nil? ? 'n/a' : format(pattern, value) + end + + def verdict(within) + return '-' if within.nil? + + within ? 'yes' : 'NO' + end + + # Fails the run when a configuration misses its target, if asked to + # (DRIVER_BENCH_ENFORCE_TARGETS). The results are saved and printed + # first, so a failing run still shows why. + def enforce_targets!(summary) + return unless %w[1 true yes].include?(@env['DRIVER_BENCH_ENFORCE_TARGETS'].to_s.downcase) + + misses = tasks.product(@configurations).select { |task, configuration| target_met(task, configuration) == false } + return if misses.empty? + + puts summary + raise "over target: #{misses.map { |task, configuration| "#{task} / #{configuration.name}" }.join(', ')}" + end + end + end +end diff --git a/profile/driver_bench/configuration.rb b/profile/driver_bench/configuration.rb new file mode 100644 index 0000000000..6f8143925a --- /dev/null +++ b/profile/driver_bench/configuration.rb @@ -0,0 +1,265 @@ +# frozen_string_literal: true + +module Mongo + module DriverBench + # One driver configuration under which the whole DriverBench task list can + # be run. + # + # A benchmark result is identified by the pair (task, configuration), so + # every configuration becomes its own time series and can be watched for + # regressions independently. Comparing a configuration against the + # baseline is what catches a performance regression in an opt-in feature + # such as OpenTelemetry: no benchmark that only ever runs the default + # configuration will execute that code at all. + # + # OpenTelemetry cannot be reconfigured once its SDK has been installed + # into a process, and the "api-only" configuration requires that the SDK + # was never loaded, so exactly one configuration can be measured per + # process. Comparison runs one subprocess per configuration. + # + # @api private + class Configuration + # The environment variable naming the configuration to run under. + ENV_VAR = 'DRIVER_BENCH_CONFIGURATION' + + # The configuration every other one is compared against. + BASELINE = 'off' + + # @return [ String ] the short name of the configuration, used to key + # results and to select the configuration from the environment. + attr_reader :name + + # @return [ String ] a human-readable description of what the + # configuration measures. + attr_reader :description + + # @return [ Hash ] options to pass to every Mongo::Client the + # benchmarks construct. + attr_reader :client_options + + # @return [ Numeric | nil ] the most throughput, in percent of the + # baseline, this configuration may give up on a small-document task + # (see "Performance Targets" in the OpenTelemetry spec). nil for the + # baseline. + attr_reader :target_pct + + # @return [ String | nil ] the command-span attribute profile to apply + # (see AttributeProfiles), or nil to leave the tracer as shipped. + attr_reader :attribute_profile + + # @param name [ String ] the short name of the configuration. + # @param description [ String ] what the configuration measures. + # @param otel [ Symbol ] how much of OpenTelemetry to load: +:none+ for + # nothing, +:api+ for the API without an SDK (so that spans are + # non-recording), +:sdk+ for a configured SDK. + # @param client_options [ Hash ] options for every client. + # @param sampler_env [ Hash ] OpenTelemetry environment variables that + # select the sampler, applied before the SDK is configured. + # @param target_pct [ Numeric | nil ] the overhead target. + # @param attribute_profile [ String | nil ] the command-span attribute + # profile to apply. + def initialize(name:, description:, otel:, client_options: {}, sampler_env: {}, target_pct: nil, + attribute_profile: nil) + @name = name + @target_pct = target_pct + @description = description + @otel = otel + @client_options = client_options + @sampler_env = sampler_env + @attribute_profile = attribute_profile + end + + # Whether this configuration asks the driver to trace. + def tracing? + @otel != :none + end + + # Whether every span is recorded, which is what makes counting spans + # per operation possible: a span processor only sees sampled spans. + def records_every_span? + @otel == :sdk && @sampler_env['OTEL_TRACES_SAMPLER'] == 'always_on' + end + + # The test name to report a task's results under. perf.send only + # accepts integer args, so the configuration goes into the name. The + # baseline keeps the bare task name, so that its series continues the + # one recorded before configurations existed. + # + # @param task [ String ] the micro-benchmark name. + # + # @return [ String ] the reported test name. + def perf_test_name(task) + baseline? ? task : "#{task} [#{name}]" + end + + # The inverse of #perf_test_name. + # + # @param test_name [ String ] a reported test name. + # + # @return [ String ] the micro-benchmark name. + def task_name(test_name) + test_name.delete_suffix(" [#{name}]") + end + + # Whether this configuration is the baseline that others are compared + # against. + def baseline? + name == BASELINE + end + + # Loads and configures OpenTelemetry for this configuration. + # + # Must run before any Mongo::Client is constructed: a client decides + # whether tracing is active when it builds its tracer, and that + # decision depends on whether ::OpenTelemetry is defined. + def install! + case @otel + when :none then nil + when :api then require 'opentelemetry-api' + when :sdk then install_sdk! + end + apply_attribute_profile! + end + + # All configurations, baseline first. + ALL = [ + new(name: 'off', + description: 'tracing disabled; the baseline', + otel: :none), + + new(name: 'api-only', + description: 'tracing enabled with the OpenTelemetry API but no SDK, ' \ + 'so spans are non-recording; what a user who installs ' \ + 'nothing pays', + otel: :api, + client_options: { tracing: { enabled: true } }, + target_pct: 5), + + new(name: 'sdk-never', + description: 'SDK installed, sampler drops every trace', + otel: :sdk, + client_options: { tracing: { enabled: true } }, + sampler_env: { 'OTEL_TRACES_SAMPLER' => 'always_off' }, + target_pct: 10), + + new(name: 'sdk-parent-1pct', + description: 'SDK installed, parent-based sampler recording 1% of traces', + otel: :sdk, + client_options: { tracing: { enabled: true } }, + sampler_env: { 'OTEL_TRACES_SAMPLER' => 'parentbased_traceidratio', + 'OTEL_TRACES_SAMPLER_ARG' => '0.01' }, + target_pct: 10), + + new(name: 'sdk-always', + description: 'SDK installed, every trace recorded', + otel: :sdk, + client_options: { tracing: { enabled: true } }, + sampler_env: { 'OTEL_TRACES_SAMPLER' => 'always_on' }, + target_pct: 15), + + # Attribute profiles: every trace recorded, but the command span is + # given a different attribute shape. See AttributeProfiles. + new(name: 'attr-none', + description: 'SDK recording, command spans created with no attributes', + otel: :sdk, + client_options: { tracing: { enabled: true } }, + sampler_env: { 'OTEL_TRACES_SAMPLER' => 'always_on' }, + attribute_profile: 'none'), + + new(name: 'attr-creation-only', + description: 'SDK recording, only the creation-time attributes', + otel: :sdk, + client_options: { tracing: { enabled: true } }, + sampler_env: { 'OTEL_TRACES_SAMPLER' => 'always_on' }, + attribute_profile: 'creation-only'), + + new(name: 'attr-all-at-creation', + description: 'SDK recording, every attribute passed to start_span', + otel: :sdk, + client_options: { tracing: { enabled: true } }, + sampler_env: { 'OTEL_TRACES_SAMPLER' => 'always_on' }, + attribute_profile: 'all-at-creation'), + + new(name: 'attr-none-then-all', + description: 'SDK recording, no attributes at creation and all set afterwards', + otel: :sdk, + client_options: { tracing: { enabled: true } }, + sampler_env: { 'OTEL_TRACES_SAMPLER' => 'always_on' }, + attribute_profile: 'none-then-all') + ].freeze + + # Looks up a configuration by name. + # + # @param name [ String ] the configuration name. + # + # @return [ Configuration ] the named configuration. + # + # @raise [ ArgumentError ] if no such configuration exists. + def self.[](name) + ALL.find { |configuration| configuration.name == name } || + raise(ArgumentError, + "unknown configuration #{name.inspect}; " \ + "known configurations are #{ALL.map(&:name).join(', ')}") + end + + # Fails when YJIT was asked for (RUBY_YJIT_ENABLE) but this ruby runs + # without it. A ruby built without YJIT only warns and carries on, and + # the benchmark would then quietly measure the interpreter, where the + # cost of tracing is several times higher. + def self.check_jit! + return unless %w[1 true yes].include?(ENV['RUBY_YJIT_ENABLE'].to_s.downcase) + return if jit? + + raise "RUBY_YJIT_ENABLE is set, but #{RUBY_DESCRIPTION} runs without YJIT" + end + + # @return [ true | false ] whether YJIT is enabled in this process. + def self.jit? + !!(defined?(RubyVM::YJIT) && RubyVM::YJIT.enabled?) + end + + # The configuration the current process is measuring, named by the + # DRIVER_BENCH_CONFIGURATION environment variable. + # + # @return [ Configuration ] the current configuration. + def self.current + @current ||= self[ENV[ENV_VAR] || BASELINE] + end + + private + + def install_sdk! + require 'opentelemetry-sdk' + + # The SDK reads its sampler and exporter from the environment when it + # is configured, so set them here rather than trusting the caller to: + # a configuration must mean the same thing however the suite was + # invoked. + @sampler_env.each { |key, value| ENV[key] = value } + + # No exporter, and so no span processor either. The cost to attribute + # to the driver is the cost of creating and recording spans, which is + # where the regression this work exists to catch actually lives; a + # span processor would additionally charge the SDK's SpanData + # conversion, its background thread and its queue to the driver's + # account, and add variance that hides small driver-side changes. + # Sampled spans are still fully recorded without one, so the + # attribute-building path is exercised. + ENV['OTEL_TRACES_EXPORTER'] = 'none' + + ::OpenTelemetry::SDK.configure + end + + # Applies the command-span attribute profile, if this configuration has + # one. Must run before any client is built: a client decides whether + # tracing is active when it builds its tracer. + def apply_attribute_profile! + return unless @attribute_profile + + require_relative 'attribute_profiles' + ENV['OTEL_ATTRIBUTE_PROFILE'] = @attribute_profile + Mongo::DriverBench::CommandAttributeProfiles.apply! + end + end + end +end diff --git a/profile/driver_bench/rake/tasks.rake b/profile/driver_bench/rake/tasks.rake index 16d038b85b..4fb3f34c7e 100644 --- a/profile/driver_bench/rake/tasks.rake +++ b/profile/driver_bench/rake/tasks.rake @@ -40,6 +40,22 @@ namespace :driver_bench do Mongo::DriverBench::Suite.run! end + desc 'Compares the DriverBench suite across driver configurations' + task compare: :data do + require_relative '../comparison' + + puts Mongo::DriverBench::Comparison.run! + end + + desc 'Lists the driver configurations the suite can run under' + task :configurations do + require_relative '../configuration' + + Mongo::DriverBench::Configuration::ALL.each do |configuration| + puts format('%-16s %s', configuration.name, configuration.description) + end + end + desc 'Runs the crypto benchmark' task :crypto do require_relative '../crypto/decrypt' diff --git a/profile/driver_bench/single_doc/base.rb b/profile/driver_bench/single_doc/base.rb index cf540a9a4e..3f8a009bd9 100644 --- a/profile/driver_bench/single_doc/base.rb +++ b/profile/driver_bench/single_doc/base.rb @@ -43,6 +43,12 @@ def cleanup_client @client.database.drop end + # Every single-doc task performs one driver operation per unit of + # scale. + def ops_per_iteration + scale + end + # Returns the name of the file that contains # the dataset to use. def file_name diff --git a/profile/driver_bench/suite.rb b/profile/driver_bench/suite.rb index 40259872fb..c4895f79a7 100644 --- a/profile/driver_bench/suite.rb +++ b/profile/driver_bench/suite.rb @@ -1,6 +1,7 @@ # frozen_string_literal: true require_relative 'bson' +require_relative 'configuration' require_relative 'multi_doc' require_relative 'parallel' require_relative 'single_doc' @@ -41,28 +42,108 @@ module DriverBench class Suite PERCENTILES = [ 10, 25, 50, 75, 90, 95, 98, 99 ].freeze + # A span processor that only counts started spans. Installed after the + # timed runs, so it costs them nothing. + class SpanCounter + attr_reader :count + + def initialize + @count = 0 + @mutex = Mutex.new + end + + def reset + @mutex.synchronize { @count = 0 } + end + + def on_start(_span, _parent_context) + @mutex.synchronize { @count += 1 } + end + + def on_finish(_span); end + + def force_flush(*) + ::OpenTelemetry::SDK::Trace::Export::SUCCESS + end + + def shutdown(*) + ::OpenTelemetry::SDK::Trace::Export::SUCCESS + end + end + def self.run! new.run end def run + Configuration.check_jit! + configuration.install! + announce_configuration + perf_data = [] benches = Hash.new { |h, k| h[k] = [] } - ALL.each do |klass| + tasks.each do |klass| result = run_benchmark(klass) perf_data << compile_perf_data(result) append_to_benchmarks(klass, result, benches) end - perf_data += compile_benchmarks(benches) + count_spans(perf_data) if count_spans? + + # The composites average a fixed list of micro-benchmarks, so they are + # only meaningful when every micro-benchmark ran. + perf_data += compile_benchmarks(benches) if tasks == ALL save_perf_data(perf_data) summarize_perf_data(perf_data) end + # The micro-benchmarks to run. + # + # By default this is every micro-benchmark, but DRIVER_BENCH_TASKS and + # DRIVER_BENCH_EXCLUDE_TASKS narrow it to those whose names do (or do + # not) contain one of a comma-separated list of substrings. Narrowing + # the list is what makes an optimize-and-remeasure loop practical: the + # full suite spends at least a minute per micro-benchmark per + # configuration. + # + # @return [ Array ] the micro-benchmark classes to run. + def tasks + @tasks ||= begin + included = patterns('DRIVER_BENCH_TASKS') + excluded = patterns('DRIVER_BENCH_EXCLUDE_TASKS') + + selected = ALL.select do |klass| + name = klass.bench_name.downcase + (included.empty? || included.any? { |pattern| name.include?(pattern) }) && + excluded.none? { |pattern| name.include?(pattern) } + end + + raise 'no micro-benchmarks match the requested task filters' if selected.empty? + + (selected == ALL) ? ALL : selected.freeze + end + end + private + def configuration + @configuration ||= Configuration.current + end + + def patterns(variable) + (ENV[variable] || '').split(',').map { |pattern| pattern.strip.downcase }.reject(&:empty?) + end + + def announce_configuration + puts format('===== DriverBench: %s =====', configuration.name) + puts configuration.description + puts RUBY_DESCRIPTION + puts format('%d of %d micro-benchmarks', tasks.length, ALL.length) unless tasks == ALL + puts + end + def run_benchmark(klass) print klass.bench_name, ': ' $stdout.flush @@ -72,21 +153,48 @@ def run_benchmark(klass) end end + # Spans are counted when asked to (DRIVER_BENCH_COUNT_SPANS), and only + # under a configuration that records every span. + def count_spans? + %w[1 true yes].include?(ENV['DRIVER_BENCH_COUNT_SPANS'].to_s.downcase) && + configuration.records_every_span? + end + + # Runs one extra iteration of every task with a counting span + # processor, and adds spans_per_op to the task's metrics. A task that + # creates no spans must not be read as proof of an efficient + # implementation, so the count is reported next to the overhead. + def count_spans(perf_data) + counter = SpanCounter.new + ::OpenTelemetry.tracer_provider.add_span_processor(counter) + + tasks.each do |klass| + spans = klass.new.count_spans(counter) + next if spans.nil? + + test_name = configuration.perf_test_name(klass.bench_name) + entry = perf_data.find { |item| item['info']['test_name'] == test_name } + entry['metrics'] << { 'name' => 'spans_per_op', 'value' => spans } + end + end + def compile_perf_data(result) percentile_data = PERCENTILES.map do |percentile| { 'name' => "time-#{percentile}%", 'value' => result[:percentiles][percentile] } end + per_op_data = result[:per_op].map { |name, value| { 'name' => name, 'value' => value } } { 'info' => { - 'test_name' => result[:name], + 'test_name' => configuration.perf_test_name(result[:name]), 'args' => {}, }, 'metrics' => [ { 'name' => 'score', 'value' => result[:score] }, - *percentile_data + *percentile_data, + *per_op_data ] } end @@ -107,7 +215,7 @@ def compile_benchmarks(benches) benches.map do |bench, score| { 'info' => { - 'test_name' => bench, + 'test_name' => configuration.perf_test_name(bench), 'args' => {} }, 'metrics' => [ @@ -119,7 +227,7 @@ def compile_benchmarks(benches) end def summarize_perf_data(data) - puts '===== Performance Results =====' + puts format('===== Performance Results (%s) =====', configuration.name) data.each do |item| puts format('%s : %4.4g', item['info']['test_name'], item['metrics'][0]['value']) next unless item['metrics'].length > 1 diff --git a/profile/otel_attributes/construction_cost.rb b/profile/otel_attributes/construction_cost.rb new file mode 100644 index 0000000000..b3926b52fe --- /dev/null +++ b/profile/otel_attributes/construction_cost.rb @@ -0,0 +1,179 @@ +# frozen_string_literal: true + +require 'opentelemetry-sdk' +require 'mongo' + +module Mongo + module OtelAttributes + # Measures the driver's own cost of producing each span attribute value: + # the work that happens before the SDK sees an attribute and that the + # driver can gate on the sampling decision. + # + # The span-shape sweep in SpanCost fixes the attribute values as literals + # and measures the SDK's handling of them. This measures the other half: + # extracting the value from the command document, the connection, or the + # session. Together they add up to what an attribute costs end to end. + # + # No server is needed: the helpers are called directly with realistic + # inputs. Caches (the lsid UUID, the per-connection attribute hash) are + # warmed first, so the numbers are the steady-state cost the benchmark + # tasks actually pay. + # + # Parameterised by the environment: + # ITERATIONS calls per sample (default 200000) + # REPS repetitions per measurement (default 5) + # + # @api private + class ConstructionCost + DEFAULT_ITERATIONS = 200_000 + DEFAULT_REPS = 5 + FORMAT = '%-44s %12s %14s' + + UUID = '6f1d0f2e-0f9a-4f3b-9a3e-2f2b1c4d5e6f' + FIND_DOC = { + 'find' => 'corpus', + 'filter' => { '_id' => 1 }, + 'lsid' => { 'id' => BSON::Binary.new([ UUID.delete('-') ].pack('H*'), :uuid) }, + 'txnNumber' => BSON::Int64.new(3), + '$db' => 'perftest' + }.freeze + GET_MORE_DOC = { + 'getMore' => BSON::Int64.new(123_456_789), + 'collection' => 'corpus', + '$db' => 'perftest' + }.freeze + + # A stand-in for Mongo::Server::Connection that answers the methods + # connection_attributes reads. + class FakeConnection + Address = Struct.new(:host, :port) + Description = Struct.new(:server_connection_id) + + def initialize + @address = Address.new('localhost', 27_017) + @description = Description.new(42) + end + + attr_reader :address, :description + + def transport + :tcp + end + + def id + 7 + end + end + + # A stand-in for a Mongo::Operation that answers what the operation + # tracer reads. + class FakeOperation + attr_reader :db_name, :coll_name + + def initialize + @db_name = 'perftest' + @coll_name = 'corpus' + end + end + + def self.run! + new.run + end + + def initialize(env = ENV) + @iterations = Integer(env['ITERATIONS'] || DEFAULT_ITERATIONS) + @reps = Integer(env['REPS'] || DEFAULT_REPS) + @otel_tracer = ::OpenTelemetry::SDK::Trace::TracerProvider + .new(sampler: ::OpenTelemetry::SDK::Trace::Samplers::ALWAYS_ON) + .tracer('otel-attribute-cost', '1.0') + @command_tracer = Mongo::Tracing::OpenTelemetry::CommandTracer.new(@otel_tracer, nil) + @operation_tracer = Mongo::Tracing::OpenTelemetry::OperationTracer.new(@otel_tracer, nil) + @connection = FakeConnection.new + @operation = FakeOperation.new + end + + # Runs every measurement and returns the table. + # + # @return [ String ] the summary table. + def run + prepare + lines = [ header, column_header ] + measurements.each do |label, block| + ns, allocs = sample(&block) + lines << format(FORMAT, attribute: label, ns: format('%.1f', ns), allocs: format('%.2f', allocs)) + end + lines.join("\n") + end + + private + + # Warms every cache the helpers use, so the measurements are steady + # state rather than first-call cost. + def prepare + @command_tracer.send(:connection_attributes, @connection) + @command_tracer.send(:lsid, FIND_DOC) + @span = @otel_tracer.start_span('warmup', kind: :client) + end + + def header + format('===== Driver cost of building one command span attribute (%d calls, median of %d reps) =====', + @iterations, @reps) + end + + def column_header + format(FORMAT, attribute: 'attribute', ns: 'ns/call', allocs: 'allocs/call') + end + + # Each entry is a label and a callable. The callables take the iteration + # index that Integer#times yields and ignore it. + def measurements + [ + [ 'database(doc)', ->(_i) { @command_tracer.send(:database, FIND_DOC) } ], + [ 'command_name(doc)', ->(_i) { @command_tracer.send(:command_name, FIND_DOC) } ], + [ 'collection_name(name, doc)', ->(_i) { @command_tracer.send(:collection_name, 'find', FIND_DOC) } ], + [ 'query_summary(name, doc)', ->(_i) { @command_tracer.send(:query_summary, 'find', FIND_DOC) } ], + [ 'txn_number(doc)', ->(_i) { @command_tracer.send(:txn_number, FIND_DOC) } ], + [ 'cursor_id(name, doc)', ->(_i) { @command_tracer.send(:cursor_id, 'getMore', GET_MORE_DOC) } ], + [ 'lsid(doc) (warm cache)', ->(_i) { @command_tracer.send(:lsid, FIND_DOC) } ], + [ 'connection_attributes(conn) (warm)', + ->(_i) { @command_tracer.send(:connection_attributes, @connection) } ], + [ 'span_attributes(doc, name, conn)', + ->(_i) { @command_tracer.send(:span_attributes, FIND_DOC, 'find', @connection) } ], + [ 'apply_deferred_attributes(recording span)', + ->(_i) { @command_tracer.send(:apply_deferred_attributes, @span, nil, 'find', FIND_DOC, 123_456_789) } ], + [ 'operation span_attributes (creation set)', + ->(_i) { @operation_tracer.send(:span_attributes, @operation, 'find', 'find perftest.corpus', 'corpus') } ], + [ 'operation collection_name(op)', ->(_i) { @operation_tracer.send(:collection_name, @operation) } ] + ] + end + + # Times a block +@reps+ times and returns the median per-call CPU + # nanoseconds excluding GC, and allocations per call. + # + # @return [ Array(Float, Float) ] [ ns/call, allocs/call ]. + def sample(&block) + runs = Array.new(@reps) { time(&block) } + [ median(runs.map(&:first)), median(runs.map(&:last)) ] + end + + def time(&block) + GC.start + allocs_before = GC.stat(:total_allocated_objects) + gc_before = GC.stat(:time) + cpu_before = Process.clock_gettime(Process::CLOCK_PROCESS_CPUTIME_ID) + @iterations.times(&block) + cpu_ns = (Process.clock_gettime(Process::CLOCK_PROCESS_CPUTIME_ID) - cpu_before) * 1e9 / @iterations + gc_ns = (GC.stat(:time) - gc_before) * 1_000_000.0 / @iterations + [ cpu_ns - gc_ns, (GC.stat(:total_allocated_objects) - allocs_before).to_f / @iterations ] + end + + def median(values) + sorted = values.sort + mid = sorted.length / 2 + sorted.length.odd? ? sorted[mid] : (sorted[mid - 1] + sorted[mid]) / 2.0 + end + end + end +end + +puts Mongo::OtelAttributes::ConstructionCost.run! if $PROGRAM_NAME == __FILE__ diff --git a/profile/otel_attributes/profiles.rb b/profile/otel_attributes/profiles.rb new file mode 100644 index 0000000000..8a9dc993c7 --- /dev/null +++ b/profile/otel_attributes/profiles.rb @@ -0,0 +1,81 @@ +# frozen_string_literal: true + +module Mongo + # Harnesses that measure the cost of OpenTelemetry span attributes: what the + # SDK charges for carrying one, and what the driver charges for building one. + module OtelAttributes + # The attributes the driver can put on a span, with one realistic sample + # value each. + # + # The values are literals on purpose. This harness measures the + # OpenTelemetry SDK's cost of carrying an attribute, so building the value + # must not be part of the measurement. The driver's cost of building the + # value is measured separately, by ConstructionCost. + SAMPLE = { + 'db.system.name' => 'mongodb', + 'db.namespace' => 'perftest', + 'db.command.name' => 'find', + 'db.collection.name' => 'corpus', + 'server.port' => 27_017, + 'server.address' => 'localhost', + 'network.transport' => 'tcp', + 'db.mongodb.server_connection_id' => 42, + 'db.mongodb.driver_connection_id' => 7, + 'db.query.summary' => 'find perftest.corpus', + 'db.mongodb.lsid' => '6f1d0f2e-0f9a-4f3b-9a3e-2f2b1c4d5e6f', + 'db.mongodb.cursor_id' => 123_456_789, + 'db.mongodb.txn_number' => 3 + }.freeze + + # The attributes the driver passes to start_span today: the cheap, + # sampler-plausible set. The connection attributes are memoized per + # connection in the driver; at the SDK boundary they are just five more + # entries in the hash. + CREATION_KEYS = %w[ + db.system.name db.namespace db.command.name db.collection.name + server.port server.address network.transport + db.mongodb.server_connection_id db.mongodb.driver_connection_id + ].freeze + + # The attributes the driver sets after span creation, only when the span is + # recording. + DEFERRED_KEYS = %w[ + db.query.summary db.mongodb.lsid db.mongodb.cursor_id db.mongodb.txn_number + ].freeze + + # One span shape: which attributes go to start_span, which are set + # afterwards, and whether the attributes argument is passed at all. + Profile = Struct.new(:name, :creation_keys, :deferred_keys, :bare) do + # @return [ Hash | nil ] the attributes argument, or nil to omit it. + def creation_attributes + return nil if bare + + creation_keys.to_h { |key| [ key, SAMPLE[key] ] } + end + + # @return [ Array ] key/value pairs to set after creation. + def deferred_pairs + deferred_keys.map { |key| [ key, SAMPLE[key] ] } + end + end + + # Every profile the sweep measures. Order is the order they are reported + # in; the sweep rotates it each repetition. + # + # @return [ Array ] the profiles. + def self.profiles + list = [ + Profile.new('none', [], [], true), + Profile.new('empty-hash', [], [], false) + ] + SAMPLE.each_key { |key| list << Profile.new("one:#{key}", [ key ], [], false) } + list + [ + Profile.new('creation-current', CREATION_KEYS, [], false), + Profile.new('deferred-current', [], DEFERRED_KEYS, false), + Profile.new('current-full', CREATION_KEYS, DEFERRED_KEYS, false), + Profile.new('all-at-creation', SAMPLE.keys, [], false), + Profile.new('none-then-all', [], SAMPLE.keys, false) + ] + end + end +end diff --git a/profile/otel_attributes/rake/tasks.rake b/profile/otel_attributes/rake/tasks.rake new file mode 100644 index 0000000000..2c461f7dc1 --- /dev/null +++ b/profile/otel_attributes/rake/tasks.rake @@ -0,0 +1,60 @@ +# frozen_string_literal: true + +$LOAD_PATH.unshift File.expand_path('../../../lib', __dir__) + +namespace :otel_attributes do + desc 'Measures the SDK cost of each OpenTelemetry span attribute shape' + task :span do + require_relative '../span_cost' + + puts Mongo::OtelAttributes::SpanCost.run! + end + + desc 'Measures the driver cost of building each command span attribute' + task :construction do + require_relative '../construction_cost' + + puts Mongo::OtelAttributes::ConstructionCost.run! + end + + desc 'Runs both attribute-cost harnesses' + task all: %i[ span construction ] + + desc 'Measures per-attribute cost relative to the untraced operation on this host' + task :evergreen do + require 'json' + require 'tmpdir' + + baseline_tasks = 'small doc insertone,find one by id' + baseline_file = File.join(Dir.mktmpdir('otel-attribute-baseline'), 'baseline.json') + + # The relative figures are a percentage of the untraced operation, so that + # operation has to be measured on the same host as the spans. + puts '===== Untraced baseline (this host) =====' + ok = system( + { 'DRIVER_BENCH_TASKS' => baseline_tasks, 'PERFORMANCE_RESULTS_FILE' => baseline_file }, + 'bundle', 'exec', 'rake', 'driver_bench' + ) + raise 'the untraced baseline run failed' unless ok + + baseline = JSON.parse(File.read(baseline_file)) + baseline_cpu = lambda do |data, test_name| + entry = data.find { |item| item['info']['test_name'] == test_name } + raise "no untraced baseline for #{test_name.inspect}" unless entry + + metric = entry['metrics'].find { |item| item['name'] == 'cpu_ex_gc_us_per_op' } + raise "no cpu_ex_gc_us_per_op for #{test_name.inspect}" unless metric + + metric['value'] + end + ENV['REFERENCE_FIND_CPU_US'] = baseline_cpu.call(baseline, 'Find one by ID').to_s + ENV['REFERENCE_INSERT_CPU_US'] = baseline_cpu.call(baseline, 'Small doc insertOne').to_s + + require_relative '../span_cost' + require_relative '../construction_cost' + + puts Mongo::OtelAttributes::SpanCost.run! + puts + puts Mongo::OtelAttributes::ConstructionCost.run! + end +end diff --git a/profile/otel_attributes/span_cost.rb b/profile/otel_attributes/span_cost.rb new file mode 100644 index 0000000000..9b67fa68c9 --- /dev/null +++ b/profile/otel_attributes/span_cost.rb @@ -0,0 +1,219 @@ +# frozen_string_literal: true + +require 'fileutils' +require 'json' +require 'opentelemetry-sdk' + +require_relative 'profiles' + +module Mongo + module OtelAttributes + # Measures what each OpenTelemetry span attribute costs at the SDK + # boundary: the cost of carrying an attribute through start_span, and of + # setting one after creation. + # + # Values are fixed literals (see Profiles::SAMPLE), so only the SDK's + # handling of the attribute is measured. The driver's cost of building the + # value is measured by ConstructionCost. + # + # Every profile runs in one process and the order rotates each repetition, + # the control the DriverBench comparison uses: the quantity of interest is + # a difference between profiles, and a machine that slows down partway + # through a run would otherwise charge the drift to whichever profile ran + # in the slow half. + # + # The sweep runs twice, under a recording sampler and under a dropping one, + # because a deferred attribute costs nothing when the span is not + # recording: set_attribute discards. + # + # Parameterised by the environment: + # ITERATIONS spans per sample (default 50000) + # REPS repetitions of the whole profile set (default 7) + # RESULTS_FILE where the JSONL rows are written + # REFERENCE_FIND_CPU_US CPU us of the untraced find operation, for the + # REFERENCE_INSERT_CPU_US relative columns. Both optional: without them + # the percentage columns read n/a. The rake + # task measures them on the same host. + # + # @api private + class SpanCost + SPAN_NAME = 'find perftest.corpus' + DEFAULT_ITERATIONS = 50_000 + DEFAULT_REPS = 7 + DEFAULT_RESULTS_FILE = File.expand_path('../../tmp/otel-attribute-span-cost.jsonl', __dir__) + FORMAT = '%-36s %8s %8s %8s %8s %8s ' \ + '%8s %8s %7s %8s' + + # One timed sample: the cost of one profile over one repetition, per span. + Sample = Struct.new(:cpu_us_per_span, :gc_us_per_span, :wall_us_per_span, :allocs_per_span) + + def self.run! + new.run + end + + def initialize(env = ENV) + @iterations = Integer(env['ITERATIONS'] || DEFAULT_ITERATIONS) + @reps = Integer(env['REPS'] || DEFAULT_REPS) + @results_file = env['RESULTS_FILE'] || DEFAULT_RESULTS_FILE + @reference_find = reference(env['REFERENCE_FIND_CPU_US']) + @reference_insert = reference(env['REFERENCE_INSERT_CPU_US']) + @samples = Hash.new { |hash, key| hash[key] = [] } + end + + # Runs the sweep and returns the summary tables. + # + # @return [ String ] the summary. + def run + FileUtils.mkdir_p(File.dirname(@results_file)) + File.write(@results_file, '') + samplers.each { |label, tracer| sweep(label, tracer) } + summarize + end + + private + + def samplers + { + 'recording' => tracer(::OpenTelemetry::SDK::Trace::Samplers::ALWAYS_ON), + 'non-recording' => tracer(::OpenTelemetry::SDK::Trace::Samplers::ALWAYS_OFF) + } + end + + def tracer(sampler) + ::OpenTelemetry::SDK::Trace::TracerProvider.new(sampler: sampler) + .tracer('otel-attribute-cost', '1.0') + end + + def sweep(label, tracer) + profiles = OtelAttributes.profiles + profiles.each { |profile| measure(tracer, profile, [ @iterations / 10, 1_000 ].max) } + 1.upto(@reps) do |rep| + profiles.rotate(rep - 1).each do |profile| + record(label, profile, rep, measure(tracer, profile)) + end + end + end + + # Times one profile over +iterations+ spans. + # + # @return [ Sample ] the per-span costs. + def measure(tracer, profile, iterations = @iterations) + creation = profile.creation_attributes + deferred = profile.deferred_pairs + GC.start + allocs_before = GC.stat(:total_allocated_objects) + gc_before = GC.stat(:time) + cpu_before = Process.clock_gettime(Process::CLOCK_PROCESS_CPUTIME_ID) + wall_before = Process.clock_gettime(Process::CLOCK_MONOTONIC) + iterations.times { run_once(tracer, profile.bare, creation, deferred) } + cpu = (Process.clock_gettime(Process::CLOCK_PROCESS_CPUTIME_ID) - cpu_before) / iterations * 1e6 + wall = (Process.clock_gettime(Process::CLOCK_MONOTONIC) - wall_before) / iterations * 1e6 + gc = (GC.stat(:time) - gc_before) / 1000.0 / iterations * 1e6 + Sample.new(cpu - gc, gc, wall, (GC.stat(:total_allocated_objects) - allocs_before).to_f / iterations) + end + + def run_once(tracer, bare, creation, deferred) + span = if bare + tracer.start_span(SPAN_NAME, kind: :client) + else + tracer.start_span(SPAN_NAME, attributes: creation, kind: :client) + end + deferred.each { |key, value| span.set_attribute(key, value) } if span.recording? + span.finish + end + + def record(label, profile, rep, sample) + @samples[[ label, profile.name ]] << sample + File.open(@results_file, 'a') do |file| + file.puts(JSON.generate(sampler: label, profile: profile.name, rep: rep, + cpu_us_per_span: sample.cpu_us_per_span, + gc_us_per_span: sample.gc_us_per_span, + wall_us_per_span: sample.wall_us_per_span, + allocs_per_span: sample.allocs_per_span)) + end + end + + def median(label, profile_name, field) + values = @samples[[ label, profile_name ]].map(&field).sort + return nil if values.empty? + + mid = values.length / 2 + values.length.odd? ? values[mid] : (values[mid - 1] + values[mid]) / 2.0 + end + + def summarize + labels = @samples.keys.map(&:first).uniq + labels.flat_map { |label| [ table(label), '' ] }.join("\n") + end + + def table(label) + # The driver always passes an attributes hash, so the driver-shape + # zero-attribute span (empty-hash) is the baseline each attribute is + # measured over, not the bare span. + reference_cpu = median(label, 'empty-hash', :cpu_us_per_span) + reference_allocs = median(label, 'empty-hash', :allocs_per_span) + lines = [ + format('===== Span attribute cost at the SDK boundary: %s (median of %d reps, %d spans/sample) =====', + label, @reps, @iterations), + 'cpu: CPU us per span excluding GC; d-cpu: added over the driver-shape ' \ + 'zero-attribute span (empty-hash); d-cpu%: that delta as a percentage of it; ' \ + '%find/%insert: that delta as a percentage of the untraced operation' + ] + lines << format(FORMAT, profile: 'profile', cpu: 'cpu us', dcpu: 'd-cpu', dpct: 'd-cpu%', + find: '%find', insert: '%insert', allocs: 'allocs', dallocs: 'd-allocs', + gc: 'gc us', wall: 'wall us') + OtelAttributes.profiles.each do |profile| + lines << row(label, profile.name, reference_cpu, reference_allocs) + end + lines.join("\n") + end + + def row(label, profile_name, reference_cpu, reference_allocs) + cpu = median(label, profile_name, :cpu_us_per_span) + allocs = median(label, profile_name, :allocs_per_span) + format(FORMAT, + profile: profile_name, + cpu: number(cpu), dcpu: delta(cpu, reference_cpu), + dpct: percent(delta_value(cpu, reference_cpu), reference_cpu, signed: true), + find: percent(delta_value(cpu, reference_cpu), @reference_find), + insert: percent(delta_value(cpu, reference_cpu), @reference_insert), + allocs: number(allocs), dallocs: delta(allocs, reference_allocs), + gc: number(median(label, profile_name, :gc_us_per_span)), + wall: number(median(label, profile_name, :wall_us_per_span))) + end + + def number(value) + value.nil? ? 'n/a' : format('%.3f', value) + end + + def delta(value, reference) + return 'n/a' if value.nil? || reference.nil? + + format('%+.3f', value - reference) + end + + def delta_value(value, reference) + return nil if value.nil? || reference.nil? + + value - reference + end + + # @param value [ Float | nil ] the cost to express. + # @param reference [ Float | nil ] the baseline it is a percentage of. + # @param signed [ Boolean ] whether to show the sign, as for a delta. + def percent(value, reference, signed: false) + return 'n/a' if value.nil? || reference.nil? || reference.zero? + + format(signed ? '%+.1f%%' : '%.3f%%', value / reference * 100.0) + end + + def reference(value) + return nil if value.nil? || value.empty? + + Float(value) + end + end + end +end + +puts Mongo::OtelAttributes::SpanCost.run! if $PROGRAM_NAME == __FILE__ diff --git a/spec/mongo/tracing/open_telemetry/command_tracer_spec.rb b/spec/mongo/tracing/open_telemetry/command_tracer_spec.rb index 28ac8b36fb..bc1c79a45a 100644 --- a/spec/mongo/tracing/open_telemetry/command_tracer_spec.rb +++ b/spec/mongo/tracing/open_telemetry/command_tracer_spec.rb @@ -18,9 +18,13 @@ id: 123, address: instance_double(Mongo::Address, host: 'localhost', port: 27_017), transport: :tcp, + # The connection's own description carries the id of this + # connection on the server; the server description comes + # from the monitoring connection and must not be used. + description: instance_double(Mongo::Server::Description, server_connection_id: 456), server: instance_double(Mongo::Server, description: instance_double(Mongo::Server::Description, - server_connection_id: 456))) + server_connection_id: 999))) end let(:message) do @@ -65,7 +69,11 @@ end describe '#trace_command' do - let(:span) { instance_double(OpenTelemetry::Trace::Span, finish: nil, set_attribute: nil) } + let(:span_context) { instance_double(OpenTelemetry::Trace::SpanContext, valid?: true) } + let(:span) do + instance_double(OpenTelemetry::Trace::Span, finish: nil, set_attribute: nil, recording?: true, + context: span_context) + end let(:context) { instance_double(Mongo::Operation::Context) } let(:result) { instance_double(Mongo::Operation::Result, cursor_id: 0, successful?: true) } @@ -101,6 +109,94 @@ command_tracer.trace_command(message, operation_context, connection) { result } end + it 'passes only cheap attributes to start_span' do + expect(otel_tracer).to receive(:start_span).with( + 'find', + attributes: hash_excluding('db.query.summary', 'db.query.text', 'db.mongodb.lsid', + 'db.mongodb.cursor_id', 'db.mongodb.txn_number'), + kind: :client + ) + command_tracer.trace_command(message, operation_context, connection) { result } + end + + it 'sets deferred attributes when the span is recording' do + expect(span).to receive(:set_attribute).with('db.query.summary', 'find test_db.users') + expect(span).to receive(:set_attribute).with('db.mongodb.lsid', lsid_value) + command_tracer.trace_command(message, operation_context, connection) { result } + end + + context 'with a non-recording span' do + let(:span) do + instance_double(OpenTelemetry::Trace::Span, finish: nil, set_attribute: nil, recording?: false, + context: span_context) + end + + it 'does not build deferred attributes' do + expect(span).not_to receive(:set_attribute) + command_tracer.trace_command(message, operation_context, connection) { result } + end + end + + it 'memoizes connection attributes per connection' do + command_tracer.trace_command(message, operation_context, connection) { result } + command_tracer.trace_command(message, operation_context, connection) { result } + + attrs = connection.instance_variable_get(:@otel_connection_attributes) + expect(attrs).to be_frozen + end + + it 'extracts the command document once per command' do + count = 0 + allow(message).to receive(:documents) do + count += 1 + [ document ] + end + command_tracer.trace_command(message, operation_context, connection) { result } + expect(count).to eq(1) + end + + it 'formats the lsid UUID once across repeated commands' do + count = 0 + allow(document['lsid']['id']).to receive(:to_uuid) do + count += 1 + lsid_value + end + 3.times do + command_tracer.trace_command(message, operation_context, connection) { result } + end + expect(count).to eq(1) + end + + context 'with an invalid span (no SDK installed)' do + let(:invalid_span) { OpenTelemetry::Trace::Span::INVALID } + + before do + allow(otel_tracer).to receive(:start_span).and_return(invalid_span) + end + + it 'does not attach the span to the context' do + expect(OpenTelemetry::Trace).not_to receive(:with_span) + command_tracer.trace_command(message, operation_context, connection) { result } + end + + it 'returns the block result unchanged' do + return_value = command_tracer.trace_command(message, operation_context, connection) { :done } + expect(return_value).to eq(:done) + end + + it 'finishes the span' do + expect(invalid_span).to receive(:finish) + command_tracer.trace_command(message, operation_context, connection) { result } + end + + # Guards against an opentelemetry-api upgrade changing no-SDK semantics. + # If this fails, re-evaluate the short-circuit: it degrades to a no-op, + # correctness is unaffected. + it 'pins the API behavior: no-SDK spans have invalid contexts' do + expect(invalid_span.context.valid?).to be false + end + end + context 'when result has cursor_id' do let(:result) do instance_double(Mongo::Operation::Result, cursor_id: 789, successful?: true) @@ -288,7 +384,7 @@ end describe '#span_attributes' do - subject { command_tracer.send(:span_attributes, message, connection) } + subject { command_tracer.send(:span_attributes, document, 'find', connection) } it 'includes db.system.name' do expect(subject['db.system.name']).to eq('mongodb') @@ -306,8 +402,11 @@ expect(subject['db.command.name']).to eq('find') end - it 'includes db.query.summary' do - expect(subject['db.query.summary']).to eq('find test_db.users') + it 'excludes deferred attributes' do + %w[db.query.summary db.query.text db.mongodb.lsid db.mongodb.cursor_id + db.mongodb.txn_number].each do |key| + expect(subject).not_to have_key(key) + end end it 'includes server.port' do @@ -330,14 +429,33 @@ expect(subject['db.mongodb.driver_connection_id']).to eq(123) end - it 'includes db.mongodb.lsid' do - expect(subject['db.mongodb.lsid']).to eq(lsid_value) + context 'with an admin command without a collection' do + let(:document) { { 'serverStatus' => 1, '$db' => 'admin' } } + + it 'omits db.collection.name instead of setting nil' do + expect(subject).not_to have_key('db.collection.name') + end + end + end + + describe '#apply_deferred_attributes' do + let(:deferred_span) { instance_double(OpenTelemetry::Trace::Span, set_attribute: nil) } + + it 'sets the query summary' do + expect(deferred_span).to receive(:set_attribute).with('db.query.summary', 'find test_db.users') + command_tracer.send(:apply_deferred_attributes, deferred_span, message, 'find', document, nil) end - it 'does not include nil values' do - expect(subject).not_to have_key('db.mongodb.cursor_id') - expect(subject).not_to have_key('db.mongodb.txn_number') - expect(subject).not_to have_key('db.query.text') + it 'sets the lsid' do + expect(deferred_span).to receive(:set_attribute).with('db.mongodb.lsid', lsid_value) + command_tracer.send(:apply_deferred_attributes, deferred_span, message, 'find', document, nil) + end + + it 'skips attributes that are absent from the command' do + expect(deferred_span).not_to receive(:set_attribute).with('db.mongodb.cursor_id', anything) + expect(deferred_span).not_to receive(:set_attribute).with('db.mongodb.txn_number', anything) + expect(deferred_span).not_to receive(:set_attribute).with('db.query.text', anything) + command_tracer.send(:apply_deferred_attributes, deferred_span, message, 'find', document, nil) end context 'with getMore command' do @@ -349,8 +467,9 @@ } end - it 'includes db.mongodb.cursor_id' do - expect(subject['db.mongodb.cursor_id']).to eq(999) + it 'sets the cursor id' do + expect(deferred_span).to receive(:set_attribute).with('db.mongodb.cursor_id', 999) + command_tracer.send(:apply_deferred_attributes, deferred_span, message, 'getMore', document, 999) end end @@ -363,25 +482,27 @@ } end - it 'includes db.mongodb.txn_number' do - expect(subject['db.mongodb.txn_number']).to eq(42) + it 'sets the transaction number' do + expect(deferred_span).to receive(:set_attribute).with('db.mongodb.txn_number', 42) + command_tracer.send(:apply_deferred_attributes, deferred_span, message, 'find', document, nil) end end context 'with query text enabled' do let(:query_text_max_length) { 1000 } - it 'includes db.query.text' do - expect(subject['db.query.text']).to be_a(String) - expect(subject['db.query.text']).to include('find') + it 'sets the query text' do + expect(deferred_span).to receive(:set_attribute).with('db.query.text', a_string_including('find')) + command_tracer.send(:apply_deferred_attributes, deferred_span, message, 'find', document, nil) end end end describe '#collection_name' do - subject { command_tracer.send(:collection_name, message) } + subject { command_tracer.send(:collection_name, name, document) } context 'with find command' do + let(:name) { 'find' } let(:document) { { 'find' => 'users' } } it 'returns the collection name' do @@ -390,6 +511,7 @@ end context 'with getMore command' do + let(:name) { 'getMore' } let(:document) { { 'getMore' => 123, 'collection' => 'users' } } it 'returns the collection name' do @@ -398,6 +520,7 @@ end context 'with listCollections command' do + let(:name) { 'listCollections' } let(:document) { { 'listCollections' => 1 } } it 'returns nil' do @@ -406,6 +529,7 @@ end context 'with listDatabases command' do + let(:name) { 'listDatabases' } let(:document) { { 'listDatabases' => 1 } } it 'returns nil' do @@ -414,6 +538,7 @@ end context 'with commitTransaction command' do + let(:name) { 'commitTransaction' } let(:document) { { 'commitTransaction' => 1 } } it 'returns nil' do @@ -422,6 +547,7 @@ end context 'with abortTransaction command' do + let(:name) { 'abortTransaction' } let(:document) { { 'abortTransaction' => 1 } } it 'returns nil' do @@ -430,6 +556,7 @@ end context 'with admin command with numeric value' do + let(:name) { 'serverStatus' } let(:document) { { 'serverStatus' => 1 } } it 'returns nil' do @@ -439,7 +566,7 @@ end describe '#command_name' do - subject { command_tracer.send(:command_name, message) } + subject { command_tracer.send(:command_name, document) } let(:document) { { 'find' => 'users' } } @@ -449,7 +576,7 @@ end describe '#database' do - subject { command_tracer.send(:database, message) } + subject { command_tracer.send(:database, document) } let(:document) { { 'find' => 'users', '$db' => 'test_db' } } @@ -459,9 +586,9 @@ end describe '#query_summary' do - subject { command_tracer.send(:query_summary, message) } - context 'with collection name' do + subject { command_tracer.send(:query_summary, 'find', document) } + let(:document) { { 'find' => 'users', '$db' => 'test_db' } } it 'includes collection name' do @@ -470,6 +597,8 @@ end context 'without collection name' do + subject { command_tracer.send(:query_summary, 'listCollections', document) } + let(:document) { { 'listCollections' => 1, '$db' => 'test_db' } } it 'does not include collection name' do @@ -479,9 +608,9 @@ end describe '#cursor_id' do - subject { command_tracer.send(:cursor_id, message) } - context 'with getMore command' do + subject { command_tracer.send(:cursor_id, 'getMore', document) } + let(:document) { { 'getMore' => BSON::Int64.new(999) } } it 'returns the cursor ID' do @@ -490,6 +619,8 @@ end context 'with find command' do + subject { command_tracer.send(:cursor_id, 'find', document) } + let(:document) { { 'find' => 'users' } } it 'returns nil' do @@ -499,7 +630,7 @@ end describe '#lsid' do - subject { command_tracer.send(:lsid, message) } + subject { command_tracer.send(:lsid, document) } context 'with lsid present' do let(:document) { { 'find' => 'users', 'lsid' => { 'id' => BSON::Binary.from_uuid(lsid_value) } } } @@ -519,7 +650,7 @@ end describe '#txn_number' do - subject { command_tracer.send(:txn_number, message) } + subject { command_tracer.send(:txn_number, document) } context 'with txnNumber present' do let(:document) { { 'find' => 'users', 'txnNumber' => BSON::Int64.new(42) } } diff --git a/spec/mongo/tracing/open_telemetry/operation_tracer_spec.rb b/spec/mongo/tracing/open_telemetry/operation_tracer_spec.rb index 0f7f6f2b00..d29c584856 100644 --- a/spec/mongo/tracing/open_telemetry/operation_tracer_spec.rb +++ b/spec/mongo/tracing/open_telemetry/operation_tracer_spec.rb @@ -41,6 +41,8 @@ allow(s).to receive(:set_attribute) allow(s).to receive(:record_exception) allow(s).to receive(:status=) + allow(s).to receive(:context).and_return(double('SpanContext', valid?: true)) + allow(s).to receive(:recording?).and_return(true) end end let(:context) { double('OpenTelemetry::Context') } @@ -91,6 +93,68 @@ operation_tracer.trace_operation(operation, operation_context) { result } end + it 'passes only cheap attributes to start_span' do + expect(otel_tracer).to receive(:start_span).with( + anything, + attributes: hash_excluding('db.mongodb.cursor_id'), + with_parent: nil, + kind: :client + ) + operation_tracer.trace_operation(operation, operation_context) { result } + end + + it 'sets the cursor id after span creation when recording' do + operation = instance_double( + Mongo::Operation::GetMore, + db_name: 'test_db', + coll_name: 'test_collection', + cursor_id: 12_345, + class: class_double(Mongo::Operation::GetMore, name: 'Mongo::Operation::GetMore'), + respond_to?: true + ) + expect(span).to receive(:set_attribute).with('db.mongodb.cursor_id', 12_345) + operation_tracer.trace_operation(operation, operation_context) { result } + end + + context 'with a non-recording span' do + before do + allow(span).to receive(:recording?).and_return(false) + end + + it 'does not set attributes' do + expect(span).not_to receive(:set_attribute) + operation_tracer.trace_operation(operation, operation_context) { result } + end + end + + context 'with an invalid span (no SDK installed)' do + let(:invalid_span) { OpenTelemetry::Trace::Span::INVALID } + + before do + allow(otel_tracer).to receive(:start_span).and_return(invalid_span) + end + + it 'does not attach the span to the context' do + expect(OpenTelemetry::Trace).not_to receive(:with_span) + operation_tracer.trace_operation(operation, operation_context) { result } + end + + it 'returns the block result unchanged' do + return_value = operation_tracer.trace_operation(operation, operation_context) { :done } + expect(return_value).to eq(:done) + end + + it 'leaves the cursor context map untouched' do + operation_tracer.trace_operation(operation, operation_context) { result } + expect(cursor_context_map).to be_empty + end + + it 'finishes the span' do + expect(invalid_span).to receive(:finish) + operation_tracer.trace_operation(operation, operation_context) { result } + end + end + context 'with custom operation name' do it 'uses the provided op_name' do expect(otel_tracer).to receive(:start_span).with( @@ -176,8 +240,11 @@ expect(cursor_context_map).not_to have_key(cursor_id) end - it 'does not set cursor_id attribute' do - expect(span).not_to receive(:set_attribute).with('db.mongodb.cursor_id', anything) + it 'does not set a new cursor_id attribute' do + # The span carries the operation's own cursor id (set after span + # creation, deferred behind recording?); the closed result must not + # add another. + expect(span).to receive(:set_attribute).with('db.mongodb.cursor_id', 999).once operation_tracer.trace_operation(operation, operation_context) { cursor } end end @@ -218,18 +285,19 @@ end describe '#span_attributes' do - subject(:attributes) { operation_tracer.send(:span_attributes, operation, op_name) } + subject(:attributes) do + operation_tracer.send(:span_attributes, operation, name, span_name, coll_name) + end - let(:op_name) { nil } + let(:name) { 'find' } + let(:span_name) { 'find test_db.users' } + let(:coll_name) { 'users' } let(:operation_class) { class_double(Mongo::Operation::Find, name: 'Mongo::Operation::Find') } let(:operation) do instance_double( Mongo::Operation::Find, db_name: 'test_db', - coll_name: 'users', - cursor_id: nil, - class: operation_class, - respond_to?: true + class: operation_class ) end @@ -253,32 +321,22 @@ expect(attributes['db.operation.summary']).to eq('find test_db.users') end - it 'does not include nil values' do + it 'never includes the cursor id' do expect(attributes).not_to have_key('db.mongodb.cursor_id') end - context 'with cursor_id' do - let(:operation_class) { class_double(Mongo::Operation::GetMore, name: 'Mongo::Operation::GetMore') } - let(:operation) do - instance_double( - Mongo::Operation::GetMore, - db_name: 'test_db', - coll_name: 'users', - cursor_id: 12_345, - class: operation_class, - respond_to?: true - ) - end + context 'without a collection name' do + let(:coll_name) { nil } - it 'includes db.mongodb.cursor_id' do - expect(attributes['db.mongodb.cursor_id']).to eq(12_345) + it 'omits db.collection.name instead of setting nil' do + expect(attributes).not_to have_key('db.collection.name') end end - context 'with custom op_name' do - let(:op_name) { 'custom_operation' } + context 'with custom operation name' do + let(:name) { 'custom_operation' } - it 'uses the custom op_name' do + it 'uses the custom operation name' do expect(attributes['db.operation.name']).to eq('custom_operation') end end @@ -297,6 +355,15 @@ expect(op_name_result).to eq('find') end + it 'memoizes the operation name per operation class' do + operation_tracer.send(:operation_name, operation, nil) + operation_tracer.send(:operation_name, operation, nil) + + cache = operation_tracer.instance_variable_get(:@operation_names) + expect(cache.size).to eq(1) + expect(cache.values.first).to eq('find') + end + context 'with custom op_name' do let(:op_name) { 'CustomOperation' } @@ -496,21 +563,8 @@ end describe '#operation_span_name' do - subject(:span_name) { operation_tracer.send(:operation_span_name, operation, op_name) } - - let(:op_name) { nil } - context 'with collection name' do - let(:operation_class) { class_double(Mongo::Operation::Find, name: 'Mongo::Operation::Find') } - let(:operation) do - instance_double( - Mongo::Operation::Find, - db_name: 'test_db', - coll_name: 'users', - class: operation_class, - respond_to?: true - ) - end + subject(:span_name) { operation_tracer.send(:operation_span_name, 'find', 'test_db', 'users') } it 'includes collection name in format' do expect(span_name).to eq('find test_db.users') @@ -518,22 +572,7 @@ end context 'without collection name' do - let(:operation_class) do - class_double(Mongo::Operation::ListCollections, name: 'Mongo::Operation::ListCollections') - end - let(:operation) do - instance_double( - Mongo::Operation::ListCollections, - db_name: 'test_db', - coll_name: nil, - class: operation_class, - respond_to?: true - ) - end - - before do - allow(operation).to receive(:is_a?).and_return(false) - end + subject(:span_name) { operation_tracer.send(:operation_span_name, 'listcollections', 'test_db', nil) } it 'excludes collection name from format' do expect(span_name).to eq('listcollections test_db') @@ -541,16 +580,7 @@ end context 'with empty collection name' do - let(:operation_class) { class_double(Mongo::Operation::Command, name: 'Mongo::Operation::Command') } - let(:operation) do - instance_double( - Mongo::Operation::Command, - db_name: 'test_db', - coll_name: '', - class: operation_class, - respond_to?: true - ) - end + subject(:span_name) { operation_tracer.send(:operation_span_name, 'command', 'test_db', '') } it 'excludes collection name from format' do expect(span_name).to eq('command test_db') diff --git a/spec/mongo/tracing/open_telemetry/tracer_spec.rb b/spec/mongo/tracing/open_telemetry/tracer_spec.rb new file mode 100644 index 0000000000..ddc877ccf6 --- /dev/null +++ b/spec/mongo/tracing/open_telemetry/tracer_spec.rb @@ -0,0 +1,92 @@ +# frozen_string_literal: true + +require 'spec_helper' + +require 'opentelemetry' + +describe Mongo::Tracing::OpenTelemetry::Tracer do + let(:otel_tracer) { double('OpenTelemetry::Trace::Tracer') } + let(:tracer) { described_class.new(enabled: true, otel_tracer: otel_tracer) } + + let(:session_uuid) { 'de124f0e-b9a8-4bfc-8f4e-6d9c8f0a3c1d' } + let(:session_id_binary) { double('BSON::Binary', to_uuid: session_uuid) } + let(:session) do + instance_double( + Mongo::Session, + implicit?: false, + in_transaction?: true, + txn_num: 42, + session_id: { 'id' => session_id_binary } + ) + end + + describe '#transaction_map_key' do + context 'when the session is nil' do + it 'returns nil' do + expect(tracer.transaction_map_key(nil)).to be_nil + end + end + + context 'when the session is implicit' do + let(:session) { instance_double(Mongo::Session, implicit?: true) } + + it 'returns nil' do + expect(tracer.transaction_map_key(session)).to be_nil + end + end + + context 'when the session is not in a transaction' do + let(:session) { instance_double(Mongo::Session, implicit?: false, in_transaction?: false) } + + it 'returns nil' do + expect(tracer.transaction_map_key(session)).to be_nil + end + end + + context 'when the session is in a transaction' do + it 'combines the session UUID and the transaction number' do + expect(tracer.transaction_map_key(session)).to eq("#{session_uuid}-42") + end + + it 'formats the session UUID once across repeated calls' do + expect(session_id_binary).to receive(:to_uuid).once.and_return(session_uuid) + + 3.times { tracer.transaction_map_key(session) } + end + + it 'computes a new key when the transaction number changes' do + expect(tracer.transaction_map_key(session)).to eq("#{session_uuid}-42") + + allow(session).to receive(:txn_num).and_return(43) + expect(tracer.transaction_map_key(session)).to eq("#{session_uuid}-43") + end + end + end + + describe '#parent_context_for' do + let(:operation_context) { instance_double(Mongo::Operation::Context, session: session) } + + context 'when the session is in a transaction' do + it 'returns the transaction context' do + context = double('OpenTelemetry::Context') + tracer.transaction_context_map["#{session_uuid}-42"] = context + + expect(tracer.parent_context_for(operation_context)).to eq(context) + end + + it 'returns nil when the transaction has no stored context' do + expect(tracer.parent_context_for(operation_context)).to be_nil + end + end + + context 'when the session is not in a transaction' do + let(:session) { instance_double(Mongo::Session, implicit?: true) } + + it 'returns nil without formatting a session key' do + expect(session).not_to receive(:session_id) + + expect(tracer.parent_context_for(operation_context)).to be_nil + end + end + end +end diff --git a/spec/shared b/spec/shared index 0a91eb5ed6..6189aac4f0 160000 --- a/spec/shared +++ b/spec/shared @@ -1 +1 @@ -Subproject commit 0a91eb5ed627fc1acc5162d0dba447d5f7d7aff1 +Subproject commit 6189aac4f0092e5692c54f33695e110190898c70 diff --git a/spec/support/tracing.rb b/spec/support/tracing.rb index 6eb96801da..36befc8df9 100644 --- a/spec/support/tracing.rb +++ b/spec/support/tracing.rb @@ -22,6 +22,17 @@ def set_attribute(key, value) @attributes[key] = value end + # Mirrors the OpenTelemetry Span API: these spans always record, so the + # driver's deferred-attribute path runs and the assertions see the + # attributes. + def recording? + true + end + + def context + @context ||= Context.new(self) + end + # rubocop:disable Lint/UnusedMethodArgument def record_exception(exception, attributes: nil) set_attribute('exception.type', exception.class.to_s) @@ -53,6 +64,10 @@ class Context def initialize(span) @span = span end + + def valid? + true + end end class Tracer