Skip to content

Commit 9dfcd72

Browse files
authored
Merge branch 'main' into fix-pika-duplicate-consumer-spans
2 parents b24425f + 5d3791e commit 9dfcd72

9 files changed

Lines changed: 323 additions & 48 deletions

File tree

.changelog/4739.added

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
`opentelemetry-instrumentation-pymemcache`: add database semconv stability migration support

.changelog/4747.added

166 Bytes
Binary file not shown.

instrumentation/README.md

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@
1313
| [opentelemetry-instrumentation-aws-lambda](./opentelemetry-instrumentation-aws-lambda) | aws_lambda | No | development
1414
| [opentelemetry-instrumentation-boto3sqs](./opentelemetry-instrumentation-boto3sqs) | boto3 ~= 1.0 | No | development
1515
| [opentelemetry-instrumentation-botocore](./opentelemetry-instrumentation-botocore) | botocore~=1.0,aiobotocore>=2.0,<4.0 | No | development
16-
| [opentelemetry-instrumentation-cassandra](./opentelemetry-instrumentation-cassandra) | cassandra-driver ~= 3.25,scylla-driver ~= 3.25 | No | development
16+
| [opentelemetry-instrumentation-cassandra](./opentelemetry-instrumentation-cassandra) | cassandra-driver ~= 3.25,scylla-driver ~= 3.25 | No | migration
1717
| [opentelemetry-instrumentation-celery](./opentelemetry-instrumentation-celery) | celery >= 4.0, < 6.0 | No | development
1818
| [opentelemetry-instrumentation-click](./opentelemetry-instrumentation-click) | click >= 8.1.3, < 9.0.0 | No | development
1919
| [opentelemetry-instrumentation-confluent-kafka](./opentelemetry-instrumentation-confluent-kafka) | confluent-kafka >= 1.8.2, < 3.0.0 | No | development
@@ -33,7 +33,7 @@
3333
| [opentelemetry-instrumentation-pika](./opentelemetry-instrumentation-pika) | pika >= 0.12.0 | No | development
3434
| [opentelemetry-instrumentation-psycopg](./opentelemetry-instrumentation-psycopg) | psycopg >= 3.1.0 | No | migration
3535
| [opentelemetry-instrumentation-psycopg2](./opentelemetry-instrumentation-psycopg2) | psycopg2 >= 2.7.3.1,psycopg2-binary >= 2.7.3.1 | No | migration
36-
| [opentelemetry-instrumentation-pymemcache](./opentelemetry-instrumentation-pymemcache) | pymemcache >= 1.3.5, < 5 | No | development
36+
| [opentelemetry-instrumentation-pymemcache](./opentelemetry-instrumentation-pymemcache) | pymemcache >= 1.3.5, < 5 | No | migration
3737
| [opentelemetry-instrumentation-pymongo](./opentelemetry-instrumentation-pymongo) | pymongo >= 3.1, < 5.0 | No | development
3838
| [opentelemetry-instrumentation-pymssql](./opentelemetry-instrumentation-pymssql) | pymssql >= 2.1.5, < 3 | No | migration
3939
| [opentelemetry-instrumentation-pymysql](./opentelemetry-instrumentation-pymysql) | PyMySQL < 2 | No | migration

instrumentation/opentelemetry-instrumentation-cassandra/src/opentelemetry/instrumentation/cassandra/__init__.py

Lines changed: 28 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,15 @@
3333
from wrapt import wrap_function_wrapper
3434

3535
from opentelemetry import trace
36+
from opentelemetry.instrumentation._semconv import (
37+
_get_schema_url_for_signal_types,
38+
_OpenTelemetrySemanticConventionStability,
39+
_OpenTelemetryStabilitySignalType,
40+
_set_db_name,
41+
_set_db_statement,
42+
_set_db_system,
43+
_set_http_net_peer_name_client,
44+
)
3645
from opentelemetry.instrumentation.cassandra.package import (
3746
_instruments_any,
3847
_instruments_cassandra_driver,
@@ -41,17 +50,11 @@
4150
from opentelemetry.instrumentation.cassandra.version import __version__
4251
from opentelemetry.instrumentation.instrumentor import BaseInstrumentor
4352
from opentelemetry.instrumentation.utils import unwrap
44-
from opentelemetry.semconv._incubating.attributes.db_attributes import (
45-
DB_NAME,
46-
DB_STATEMENT,
47-
DB_SYSTEM,
48-
)
49-
from opentelemetry.semconv._incubating.attributes.net_attributes import (
50-
NET_PEER_NAME,
51-
)
5253

5354

54-
def _instrument(tracer_provider, include_db_statement=False):
55+
def _instrument(
56+
tracer_provider, include_db_statement=False, sem_conv_opt_in_mode=None
57+
):
5558
"""Instruments the cassandra-driver/scylla-driver module
5659
5760
Wraps cassandra.cluster.Session.execute_async().
@@ -60,7 +63,9 @@ def _instrument(tracer_provider, include_db_statement=False):
6063
__name__,
6164
__version__,
6265
tracer_provider,
63-
schema_url="https://opentelemetry.io/schemas/1.11.0",
66+
schema_url=_get_schema_url_for_signal_types(
67+
[_OpenTelemetryStabilitySignalType.DATABASE]
68+
),
6469
)
6570
name = "Cassandra"
6671

@@ -69,17 +74,18 @@ def _traced_execute_async(func, instance, args, kwargs):
6974
name, kind=trace.SpanKind.CLIENT
7075
) as span:
7176
if span.is_recording():
72-
span.set_attribute(DB_NAME, instance.keyspace)
73-
span.set_attribute(DB_SYSTEM, "cassandra")
74-
span.set_attribute(
75-
NET_PEER_NAME,
77+
attrs = {}
78+
_set_db_system(attrs, "cassandra", sem_conv_opt_in_mode)
79+
_set_db_name(attrs, instance.keyspace, sem_conv_opt_in_mode)
80+
_set_http_net_peer_name_client(
81+
attrs,
7682
instance.cluster.contact_points,
83+
sem_conv_opt_in_mode,
7784
)
78-
7985
if include_db_statement:
8086
query = args[0]
81-
span.set_attribute(DB_STATEMENT, str(query))
82-
87+
_set_db_statement(attrs, str(query), sem_conv_opt_in_mode)
88+
span.set_attributes(attrs)
8389
response = func(*args, **kwargs)
8490
return response
8591

@@ -108,9 +114,14 @@ def instrumentation_dependencies(self) -> Collection[str]:
108114
return _instruments_any
109115

110116
def _instrument(self, **kwargs):
117+
_OpenTelemetrySemanticConventionStability._initialize()
118+
sem_conv_opt_in_mode = _OpenTelemetrySemanticConventionStability._get_opentelemetry_stability_opt_in_mode(
119+
_OpenTelemetryStabilitySignalType.DATABASE,
120+
)
111121
_instrument(
112122
tracer_provider=kwargs.get("tracer_provider"),
113123
include_db_statement=kwargs.get("include_db_statement"),
124+
sem_conv_opt_in_mode=sem_conv_opt_in_mode,
114125
)
115126

116127
def _uninstrument(self, **kwargs):

instrumentation/opentelemetry-instrumentation-cassandra/src/opentelemetry/instrumentation/cassandra/package.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,3 +7,4 @@
77

88
_instruments = ()
99
_instruments_any = (_instruments_cassandra_driver, _instruments_scylla_driver)
10+
_semconv_status = "migration"

instrumentation/opentelemetry-instrumentation-cassandra/tests/test_cassandra_integration.py

Lines changed: 139 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,12 +10,31 @@
1010

1111
import opentelemetry.instrumentation.cassandra
1212
from opentelemetry import trace as trace_api
13+
from opentelemetry.instrumentation._semconv import (
14+
_OpenTelemetrySemanticConventionStability,
15+
)
1316
from opentelemetry.instrumentation.cassandra import CassandraInstrumentor
1417
from opentelemetry.instrumentation.cassandra.package import (
1518
_instruments_cassandra_driver,
1619
_instruments_scylla_driver,
1720
)
1821
from opentelemetry.sdk import resources
22+
from opentelemetry.semconv._incubating.attributes.db_attributes import (
23+
DB_NAME,
24+
DB_STATEMENT,
25+
DB_SYSTEM,
26+
)
27+
from opentelemetry.semconv._incubating.attributes.net_attributes import (
28+
NET_PEER_NAME,
29+
)
30+
from opentelemetry.semconv.attributes.db_attributes import (
31+
DB_NAMESPACE,
32+
DB_QUERY_TEXT,
33+
DB_SYSTEM_NAME,
34+
)
35+
from opentelemetry.semconv.attributes.server_attributes import (
36+
SERVER_ADDRESS,
37+
)
1938
from opentelemetry.test.test_base import TestBase
2039
from opentelemetry.trace import SpanKind
2140

@@ -241,3 +260,123 @@ def _distribution(name):
241260
package_to_instrument,
242261
(_instruments_cassandra_driver, _instruments_scylla_driver),
243262
)
263+
264+
265+
class TestCassandraSemconvStability(TestBase):
266+
def tearDown(self):
267+
super().tearDown()
268+
with self.disable_logging():
269+
CassandraInstrumentor().uninstrument()
270+
_OpenTelemetrySemanticConventionStability._initialized = False
271+
272+
@property
273+
def _mocked_session(self):
274+
return cassandra.cluster.Session(cluster=mock.Mock(), hosts=[])
275+
276+
@mock.patch("cassandra.cluster.Cluster.connect")
277+
@mock.patch("cassandra.cluster.Session.__init__")
278+
@mock.patch("cassandra.cluster.Session._create_response_future")
279+
def test_default_semconv(
280+
self, mock_create_response_future, mock_session_init, mock_connect
281+
):
282+
mock_create_response_future.return_value = mock.Mock()
283+
mock_session_init.return_value = None
284+
mock_connect.return_value = self._mocked_session
285+
286+
CassandraInstrumentor().instrument(include_db_statement=True)
287+
connect_and_execute_query()
288+
289+
spans = self.memory_exporter.get_finished_spans()
290+
self.assertEqual(len(spans), 1)
291+
span = spans[0]
292+
293+
self.assertEqual(span.attributes[DB_NAME], "test")
294+
self.assertEqual(span.attributes[DB_SYSTEM], "cassandra")
295+
self.assertEqual(span.attributes[DB_STATEMENT], "SELECT * FROM test")
296+
self.assertIn(NET_PEER_NAME, span.attributes)
297+
self.assertNotIn(DB_NAMESPACE, span.attributes)
298+
self.assertNotIn(DB_SYSTEM_NAME, span.attributes)
299+
self.assertNotIn(DB_QUERY_TEXT, span.attributes)
300+
self.assertNotIn(SERVER_ADDRESS, span.attributes)
301+
self.assertEqual(
302+
span.instrumentation_scope.schema_url,
303+
"https://opentelemetry.io/schemas/1.11.0",
304+
)
305+
306+
@mock.patch("cassandra.cluster.Cluster.connect")
307+
@mock.patch("cassandra.cluster.Session.__init__")
308+
@mock.patch("cassandra.cluster.Session._create_response_future")
309+
def test_new_semconv(
310+
self, mock_create_response_future, mock_session_init, mock_connect
311+
):
312+
mock_create_response_future.return_value = mock.Mock()
313+
mock_session_init.return_value = None
314+
mock_connect.return_value = self._mocked_session
315+
316+
with mock.patch.dict(
317+
"os.environ",
318+
{"OTEL_SEMCONV_STABILITY_OPT_IN": "database"},
319+
):
320+
_OpenTelemetrySemanticConventionStability._initialized = False
321+
CassandraInstrumentor().instrument(include_db_statement=True)
322+
connect_and_execute_query()
323+
324+
spans = self.memory_exporter.get_finished_spans()
325+
self.assertEqual(len(spans), 1)
326+
span = spans[0]
327+
328+
self.assertEqual(span.attributes[DB_NAMESPACE], "test")
329+
self.assertEqual(span.attributes[DB_SYSTEM_NAME], "cassandra")
330+
self.assertEqual(
331+
span.attributes[DB_QUERY_TEXT], "SELECT * FROM test"
332+
)
333+
self.assertIn(SERVER_ADDRESS, span.attributes)
334+
self.assertNotIn(DB_NAME, span.attributes)
335+
self.assertNotIn(DB_SYSTEM, span.attributes)
336+
self.assertNotIn(DB_STATEMENT, span.attributes)
337+
self.assertNotIn(NET_PEER_NAME, span.attributes)
338+
self.assertEqual(
339+
span.instrumentation_scope.schema_url,
340+
"https://opentelemetry.io/schemas/1.25.0",
341+
)
342+
_OpenTelemetrySemanticConventionStability._initialized = False
343+
344+
@mock.patch("cassandra.cluster.Cluster.connect")
345+
@mock.patch("cassandra.cluster.Session.__init__")
346+
@mock.patch("cassandra.cluster.Session._create_response_future")
347+
def test_dup_semconv(
348+
self, mock_create_response_future, mock_session_init, mock_connect
349+
):
350+
mock_create_response_future.return_value = mock.Mock()
351+
mock_session_init.return_value = None
352+
mock_connect.return_value = self._mocked_session
353+
354+
with mock.patch.dict(
355+
"os.environ",
356+
{"OTEL_SEMCONV_STABILITY_OPT_IN": "database/dup"},
357+
):
358+
_OpenTelemetrySemanticConventionStability._initialized = False
359+
CassandraInstrumentor().instrument(include_db_statement=True)
360+
connect_and_execute_query()
361+
362+
spans = self.memory_exporter.get_finished_spans()
363+
self.assertEqual(len(spans), 1)
364+
span = spans[0]
365+
366+
self.assertEqual(span.attributes[DB_NAME], "test")
367+
self.assertEqual(span.attributes[DB_SYSTEM], "cassandra")
368+
self.assertEqual(
369+
span.attributes[DB_STATEMENT], "SELECT * FROM test"
370+
)
371+
self.assertIn(NET_PEER_NAME, span.attributes)
372+
self.assertEqual(span.attributes[DB_NAMESPACE], "test")
373+
self.assertEqual(span.attributes[DB_SYSTEM_NAME], "cassandra")
374+
self.assertEqual(
375+
span.attributes[DB_QUERY_TEXT], "SELECT * FROM test"
376+
)
377+
self.assertIn(SERVER_ADDRESS, span.attributes)
378+
self.assertEqual(
379+
span.instrumentation_scope.schema_url,
380+
"https://opentelemetry.io/schemas/1.25.0",
381+
)
382+
_OpenTelemetrySemanticConventionStability._initialized = False

0 commit comments

Comments
 (0)