diff --git a/sentry_sdk/integrations/asyncpg.py b/sentry_sdk/integrations/asyncpg.py index c746c3a8f6..e99560eecf 100644 --- a/sentry_sdk/integrations/asyncpg.py +++ b/sentry_sdk/integrations/asyncpg.py @@ -2,16 +2,14 @@ import contextlib import re -from typing import Any, Awaitable, Callable, Iterator, TypeVar, Union +from typing import Any, Awaitable, Callable, Iterator, TypeVar import sentry_sdk from sentry_sdk.consts import OP, SPANDATA from sentry_sdk.integrations import DidNotEnable, Integration, _check_minimum_version from sentry_sdk.traces import StreamedSpan -from sentry_sdk.tracing import Span from sentry_sdk.tracing_utils import ( add_query_source, - has_span_streaming_enabled, record_sql_queries, ) from sentry_sdk.utils import ( @@ -95,11 +93,6 @@ async def _inner(*args: "Any", **kwargs: "Any") -> "T": span_origin=AsyncPGIntegration.origin, ) as span: res = await f(*args, **kwargs) - if isinstance(span, StreamedSpan): - with capture_internal_exceptions(): - add_query_source(span) - - if not isinstance(span, StreamedSpan): with capture_internal_exceptions(): add_query_source(span) @@ -118,7 +111,7 @@ def _record( params_list: "tuple[Any, ...] | None", *, executemany: bool = False, -) -> "Iterator[Union[Span, StreamedSpan]]": +) -> "Iterator[StreamedSpan]": client = sentry_sdk.get_client() integration = client.get_integration(AsyncPGIntegration) if integration is not None and not integration._record_params: @@ -152,11 +145,6 @@ async def _inner(*args: "Any", **kwargs: "Any") -> "T": res = await f(*args, **kwargs) - if isinstance(span, StreamedSpan): - with capture_internal_exceptions(): - add_query_source(span) - - if not isinstance(span, StreamedSpan): with capture_internal_exceptions(): add_query_source(span) @@ -194,11 +182,6 @@ async def _inner(*args: "Any", **kwargs: "Any") -> "T": _set_db_data(span, cursor._connection) res = await f(*args, **kwargs) - if isinstance(span, StreamedSpan): - with capture_internal_exceptions(): - add_query_source(span) - - if not isinstance(span, StreamedSpan): with capture_internal_exceptions(): add_query_source(span) @@ -219,95 +202,51 @@ async def _inner(*args: "Any", **kwargs: "Any") -> "T": database = kwargs["params"].database addr = kwargs.get("addr") - if has_span_streaming_enabled(client.options): - span_attributes = { - "sentry.op": OP.DB, - "sentry.origin": AsyncPGIntegration.origin, - SPANDATA.DB_SYSTEM_NAME: "postgresql", - SPANDATA.DB_USER: user, - SPANDATA.DB_NAMESPACE: database, - SPANDATA.DB_DRIVER_NAME: "asyncpg", - } - if addr: - try: - span_attributes[SPANDATA.SERVER_ADDRESS] = addr[0] - span_attributes[SPANDATA.SERVER_PORT] = addr[1] - except IndexError: - pass + span_attributes = { + "sentry.op": OP.DB, + "sentry.origin": AsyncPGIntegration.origin, + SPANDATA.DB_SYSTEM_NAME: "postgresql", + SPANDATA.DB_USER: user, + SPANDATA.DB_NAMESPACE: database, + SPANDATA.DB_DRIVER_NAME: "asyncpg", + } + if addr: + try: + span_attributes[SPANDATA.SERVER_ADDRESS] = addr[0] + span_attributes[SPANDATA.SERVER_PORT] = addr[1] + except IndexError: + pass + + with capture_internal_exceptions(): + sentry_sdk.add_breadcrumb( + message="connect", category="query", data=span_attributes + ) - with capture_internal_exceptions(): - sentry_sdk.add_breadcrumb( - message="connect", category="query", data=span_attributes - ) - - if sentry_sdk.traces.get_current_span() is None: - return await f(*args, **kwargs) - - with sentry_sdk.traces.start_span( - name="connect", attributes=span_attributes - ): - return await f(*args, **kwargs) - - with sentry_sdk.start_span( - op=OP.DB, - name="connect", - origin=AsyncPGIntegration.origin, - ) as span: - span.set_data(SPANDATA.DB_SYSTEM, "postgresql") - if addr: - try: - span.set_data(SPANDATA.SERVER_ADDRESS, addr[0]) - span.set_data(SPANDATA.SERVER_PORT, addr[1]) - except IndexError: - pass - span.set_data(SPANDATA.DB_NAME, database) - span.set_data(SPANDATA.DB_USER, user) - span.set_data(SPANDATA.DB_DRIVER_NAME, "asyncpg") + if sentry_sdk.traces.get_current_span() is None: + return await f(*args, **kwargs) - with capture_internal_exceptions(): - sentry_sdk.add_breadcrumb( - message="connect", category="query", data=span._data - ) + with sentry_sdk.traces.start_span(name="connect", attributes=span_attributes): return await f(*args, **kwargs) return _inner -def _set_db_data(span: "Union[Span, StreamedSpan]", conn: "Any") -> None: +def _set_db_data(span: "StreamedSpan", conn: "Any") -> None: addr = conn._addr database = conn._params.database user = conn._params.user - if isinstance(span, StreamedSpan): - span.set_attribute(SPANDATA.DB_SYSTEM_NAME, "postgresql") - span.set_attribute(SPANDATA.DB_DRIVER_NAME, "asyncpg") - if addr: - try: - span.set_attribute(SPANDATA.SERVER_ADDRESS, addr[0]) - span.set_attribute(SPANDATA.SERVER_PORT, addr[1]) - except IndexError: - pass - - if database: - span.set_attribute(SPANDATA.DB_NAMESPACE, database) - - if user: - span.set_attribute(SPANDATA.DB_USER, user) - else: - # Remove this else block once we've completely migrated to streamed spans - # The use of deprecated attributes here is to ensure backwards compatibility - span.set_data(SPANDATA.DB_SYSTEM, "postgresql") - span.set_data(SPANDATA.DB_DRIVER_NAME, "asyncpg") - - if addr: - try: - span.set_data(SPANDATA.SERVER_ADDRESS, addr[0]) - span.set_data(SPANDATA.SERVER_PORT, addr[1]) - except IndexError: - pass + span.set_attribute(SPANDATA.DB_SYSTEM_NAME, "postgresql") + span.set_attribute(SPANDATA.DB_DRIVER_NAME, "asyncpg") + if addr: + try: + span.set_attribute(SPANDATA.SERVER_ADDRESS, addr[0]) + span.set_attribute(SPANDATA.SERVER_PORT, addr[1]) + except IndexError: + pass - if database: - span.set_data(SPANDATA.DB_NAME, database) + if database: + span.set_attribute(SPANDATA.DB_NAMESPACE, database) - if user: - span.set_data(SPANDATA.DB_USER, user) + if user: + span.set_attribute(SPANDATA.DB_USER, user) diff --git a/tests/integrations/asyncpg/test_asyncpg.py b/tests/integrations/asyncpg/test_asyncpg.py index 1991b9af05..e847c1c717 100644 --- a/tests/integrations/asyncpg/test_asyncpg.py +++ b/tests/integrations/asyncpg/test_asyncpg.py @@ -44,21 +44,6 @@ def _get_db_name(): PG_USER, PG_PASSWORD, PG_HOST, PG_NAME ) CRUMBS_CONNECT = { - "category": "query", - "data": ApproxDict( - { - "db.name": PG_NAME, - "db.system": "postgresql", - "db.user": PG_USER, - "db.driver.name": "asyncpg", - "server.address": PG_HOST, - "server.port": PG_PORT, - } - ), - "message": "connect", - "type": "default", -} -CRUMBS_CONNECT_STREAMING = { "category": "query", "data": ApproxDict( { @@ -110,22 +95,19 @@ async def _clean_pg(): @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_connect( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ) -> None: sentry_init( integrations=[AsyncPGIntegration()], - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", _experiments={ "record_sql_params": True, }, ) - if span_streaming: - items = capture_items("event") - else: - events = capture_events() + items = capture_items("event") conn: Connection = await connect(PG_CONNECTION_URI) @@ -133,37 +115,29 @@ async def test_connect( capture_message("hi") - if span_streaming: - event = items[0].payload - else: - (event,) = events + event = items[0].payload for crumb in event["breadcrumbs"]["values"]: del crumb["timestamp"] - expected_crumbs_connect = ( - CRUMBS_CONNECT_STREAMING if span_streaming else CRUMBS_CONNECT - ) + expected_crumbs_connect = CRUMBS_CONNECT assert event["breadcrumbs"]["values"] == [expected_crumbs_connect] @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_execute( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ) -> None: sentry_init( integrations=[AsyncPGIntegration()], - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", _experiments={ "record_sql_params": True, }, ) - if span_streaming: - items = capture_items("event") - else: - events = capture_events() + items = capture_items("event") conn: Connection = await connect(PG_CONNECTION_URI) @@ -188,17 +162,12 @@ async def test_execute( capture_message("hi") - if span_streaming: - event = items[0].payload - else: - (event,) = events + event = items[0].payload for crumb in event["breadcrumbs"]["values"]: del crumb["timestamp"] - expected_crumbs_connect = ( - CRUMBS_CONNECT_STREAMING if span_streaming else CRUMBS_CONNECT - ) + expected_crumbs_connect = CRUMBS_CONNECT assert event["breadcrumbs"]["values"] == [ expected_crumbs_connect, { @@ -229,22 +198,19 @@ async def test_execute( @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_execute_many( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ) -> None: sentry_init( integrations=[AsyncPGIntegration()], - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", _experiments={ "record_sql_params": True, }, ) - if span_streaming: - items = capture_items("event") - else: - events = capture_events() + items = capture_items("event") conn: Connection = await connect(PG_CONNECTION_URI) @@ -260,19 +226,13 @@ async def test_execute_many( capture_message("hi") - if span_streaming: - event = items[0].payload - else: - (event,) = events + event = items[0].payload for crumb in event["breadcrumbs"]["values"]: del crumb["timestamp"] - expected_crumbs_connect = ( - CRUMBS_CONNECT_STREAMING if span_streaming else CRUMBS_CONNECT - ) assert event["breadcrumbs"]["values"] == [ - expected_crumbs_connect, + CRUMBS_CONNECT, { "category": "query", "data": {"db.executemany": True}, @@ -287,6 +247,7 @@ async def test_record_params(sentry_init, capture_events) -> None: sentry_init( integrations=[AsyncPGIntegration(record_params=True)], _experiments={"record_sql_params": True}, + trace_lifecycle="stream", ) events = capture_events() @@ -327,6 +288,7 @@ async def test_cursor(sentry_init, capture_events) -> None: sentry_init( integrations=[AsyncPGIntegration()], _experiments={"record_sql_params": True}, + trace_lifecycle="stream", ) events = capture_events() @@ -346,7 +308,7 @@ async def test_cursor(sentry_init, capture_events) -> None: async for record in conn.cursor( "SELECT * FROM users WHERE dob > $1", datetime.date(1970, 1, 1) ): - print(record) + pass await conn.close() @@ -381,6 +343,7 @@ async def test_cursor_manual(sentry_init, capture_events) -> None: sentry_init( integrations=[AsyncPGIntegration()], _experiments={"record_sql_params": True}, + trace_lifecycle="stream", ) events = capture_events() @@ -400,11 +363,9 @@ async def test_cursor_manual(sentry_init, capture_events) -> None: cur = await conn.cursor( "SELECT * FROM users WHERE dob > $1", datetime.date(1970, 1, 1) ) - record = await cur.fetchrow() - print(record) + await cur.fetchrow() while await cur.forward(1): - record = await cur.fetchrow() - print(record) + await cur.fetchrow() await conn.close() @@ -445,6 +406,7 @@ async def test_prepared_stmt(sentry_init, capture_events) -> None: sentry_init( integrations=[AsyncPGIntegration()], _experiments={"record_sql_params": True}, + trace_lifecycle="stream", ) events = capture_events() @@ -460,8 +422,8 @@ async def test_prepared_stmt(sentry_init, capture_events) -> None: stmt = await conn.prepare("SELECT * FROM users WHERE name = $1") - print(await stmt.fetchval("Bob")) - print(await stmt.fetchval("Alice")) + await stmt.fetchval("Bob") + await stmt.fetchval("Alice") await conn.close() @@ -494,6 +456,7 @@ async def test_connection_pool(sentry_init, capture_events) -> None: sentry_init( integrations=[AsyncPGIntegration()], _experiments={"record_sql_params": True}, + trace_lifecycle="stream", ) events = capture_events() @@ -555,60 +518,41 @@ async def test_connection_pool(sentry_init, capture_events) -> None: @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_query_source_disabled( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ): - sentry_options = { - "integrations": [AsyncPGIntegration()], - "traces_sample_rate": 1.0, - "enable_db_query_source": False, - "db_query_source_threshold_ms": 0, - "trace_lifecycle": "stream" if span_streaming else "static", - } - - sentry_init(**sentry_options) - - if span_streaming: - items = capture_items("span") - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) - - await conn.execute( - "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", - ) - - await conn.close() - sentry_sdk.flush() - - spans = [item.payload for item in items] + sentry_init( + integrations=[AsyncPGIntegration()], + traces_sample_rate=1.0, + enable_db_query_source=False, + db_query_source_threshold_ms=0, + trace_lifecycle="stream", + ) - assert len(spans) == 3 + items = capture_items("span") + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) - connect_span = spans[0] - insert_span = spans[1] - segment = spans[2] + await conn.execute( + "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", + ) - assert segment["name"] == "test_segment" - assert insert_span["name"].startswith("INSERT INTO") - assert connect_span["name"] == "connect" - data = insert_span.get("attributes", {}) - else: - events = capture_events() + await conn.close() + sentry_sdk.flush() - with start_transaction(name="test_transaction", sampled=True): - conn: Connection = await connect(PG_CONNECTION_URI) + spans = [item.payload for item in items] - await conn.execute( - "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", - ) + assert len(spans) == 3 - await conn.close() + connect_span = spans[0] + insert_span = spans[1] + segment = spans[2] - (event,) = events - span = event["spans"][-1] - assert span["description"].startswith("INSERT INTO") - data = span.get("data", {}) + assert segment["name"] == "test_segment" + assert insert_span["name"].startswith("INSERT INTO") + assert connect_span["name"] == "connect" + data = insert_span.get("attributes", {}) assert SPANDATA.CODE_LINENO not in data assert SPANDATA.CODE_NAMESPACE not in data @@ -618,148 +562,110 @@ async def test_query_source_disabled( @pytest.mark.asyncio @pytest.mark.parametrize("enable_db_query_source", [None, True]) -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_query_source_enabled( - sentry_init, capture_events, capture_items, enable_db_query_source, span_streaming + sentry_init, + capture_items, + enable_db_query_source, ): sentry_options = { "integrations": [AsyncPGIntegration()], "traces_sample_rate": 1.0, "db_query_source_threshold_ms": 0, - "trace_lifecycle": "stream" if span_streaming else "static", + "trace_lifecycle": "stream", } if enable_db_query_source is not None: sentry_options["enable_db_query_source"] = enable_db_query_source sentry_init(**sentry_options) - if span_streaming: - items = capture_items("span") - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) - - await conn.execute( - "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", - ) - - await conn.close() - sentry_sdk.flush() - - spans = [item.payload for item in items] - - assert len(spans) == 3 + items = capture_items("span") + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) - connect_span = spans[0] - insert_span = spans[1] - segment = spans[2] + await conn.execute( + "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", + ) - assert segment["name"] == "test_segment" - assert insert_span["name"].startswith("INSERT INTO") - assert connect_span["name"] == "connect" - data = insert_span.get("attributes", {}) - else: - events = capture_events() + await conn.close() + sentry_sdk.flush() - with start_transaction(name="test_transaction", sampled=True): - conn: Connection = await connect(PG_CONNECTION_URI) + spans = [item.payload for item in items] - await conn.execute( - "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", - ) + assert len(spans) == 3 - await conn.close() + connect_span = spans[0] + insert_span = spans[1] + segment = spans[2] - (event,) = events - span = event["spans"][-1] - assert span["description"].startswith("INSERT INTO") - data = span.get("data", {}) + assert segment["name"] == "test_segment" + assert insert_span["name"].startswith("INSERT INTO") + assert connect_span["name"] == "connect" + data = insert_span.get("attributes", {}) - lineno_key = "code.line.number" if span_streaming else SPANDATA.CODE_LINENO - filepath_key = "code.file.path" if span_streaming else SPANDATA.CODE_FILEPATH - - assert lineno_key in data - assert filepath_key in data + assert "code.line.number" in data + assert "code.file.path" in data assert SPANDATA.CODE_NAMESPACE in data assert SPANDATA.CODE_FUNCTION in data @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) -async def test_query_source(sentry_init, capture_events, capture_items, span_streaming): +async def test_query_source(sentry_init, capture_items): sentry_init( integrations=[AsyncPGIntegration()], traces_sample_rate=1.0, enable_db_query_source=True, db_query_source_threshold_ms=0, - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", ) - if span_streaming: - items = capture_items("span") - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) - - await conn.execute( - "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", - ) - - await conn.close() - sentry_sdk.flush() - - spans = [item.payload for item in items] - - assert len(spans) == 3 - - connect_span = spans[0] - insert_span = spans[1] - segment = spans[2] + items = capture_items("span") + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) - assert segment["name"] == "test_segment" - assert insert_span["name"].startswith("INSERT INTO") - assert connect_span["name"] == "connect" - data = insert_span.get("attributes", {}) - else: - events = capture_events() + await conn.execute( + "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", + ) - with start_transaction(name="test_transaction", sampled=True): - conn: Connection = await connect(PG_CONNECTION_URI) + await conn.close() + sentry_sdk.flush() - await conn.execute( - "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", - ) + spans = [item.payload for item in items] - await conn.close() + assert len(spans) == 3 - (event,) = events - span = event["spans"][-1] - assert span["description"].startswith("INSERT INTO") - data = span.get("data", {}) + connect_span = spans[0] + insert_span = spans[1] + segment = spans[2] - lineno_key = "code.line.number" if span_streaming else SPANDATA.CODE_LINENO - filepath_key = "code.file.path" if span_streaming else SPANDATA.CODE_FILEPATH + assert segment["name"] == "test_segment" + assert insert_span["name"].startswith("INSERT INTO") + assert connect_span["name"] == "connect" + data = insert_span.get("attributes", {}) - assert lineno_key in data - assert filepath_key in data + assert "code.line.number" in data + assert "code.file.path" in data assert SPANDATA.CODE_NAMESPACE in data assert SPANDATA.CODE_FUNCTION in data - assert type(data.get(lineno_key)) == int - assert data.get(lineno_key) > 0 + assert type(data.get("code.line.number")) == int + assert data.get("code.line.number") > 0 assert ( data.get(SPANDATA.CODE_NAMESPACE) == "tests.integrations.asyncpg.test_asyncpg" ) - assert data.get(filepath_key).endswith("tests/integrations/asyncpg/test_asyncpg.py") + assert data.get("code.file.path").endswith( + "tests/integrations/asyncpg/test_asyncpg.py" + ) - is_relative_path = data.get(filepath_key)[0] != os.sep + is_relative_path = data.get("code.file.path")[0] != os.sep assert is_relative_path assert data.get(SPANDATA.CODE_FUNCTION) == "test_query_source" @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_query_source_with_module_in_search_path( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ): """ Test that query source is relative to the path of the module it ran in @@ -769,156 +675,101 @@ async def test_query_source_with_module_in_search_path( traces_sample_rate=1.0, enable_db_query_source=True, db_query_source_threshold_ms=0, - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", ) from asyncpg_helpers.helpers import execute_query_in_connection - if span_streaming: - items = capture_items("span") - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) - - await execute_query_in_connection( - "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", - conn, - ) - - await conn.close() - sentry_sdk.flush() - - spans = [item.payload for item in items] - - assert len(spans) == 3 - - connect_span = spans[0] - insert_span = spans[1] - segment = spans[2] - - assert segment["name"] == "test_segment" - assert insert_span["name"].startswith("INSERT INTO") - assert connect_span["name"] == "connect" - data = insert_span.get("attributes", {}) - else: - events = capture_events() + items = capture_items("span") + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) - with start_transaction(name="test_transaction", sampled=True): - conn: Connection = await connect(PG_CONNECTION_URI) + await execute_query_in_connection( + "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", + conn, + ) - await execute_query_in_connection( - "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", - conn, - ) + await conn.close() + sentry_sdk.flush() - await conn.close() + spans = [item.payload for item in items] - (event,) = events + assert len(spans) == 3 - span = event["spans"][-1] - assert span["description"].startswith("INSERT INTO") - data = span.get("data", {}) + connect_span = spans[0] + insert_span = spans[1] + segment = spans[2] - lineno_key = "code.line.number" if span_streaming else SPANDATA.CODE_LINENO - filepath_key = "code.file.path" if span_streaming else SPANDATA.CODE_FILEPATH + assert segment["name"] == "test_segment" + assert insert_span["name"].startswith("INSERT INTO") + assert connect_span["name"] == "connect" + data = insert_span.get("attributes", {}) - assert lineno_key in data - assert filepath_key in data + assert "code.line.number" in data + assert "code.file.path" in data assert SPANDATA.CODE_NAMESPACE in data assert SPANDATA.CODE_FUNCTION in data - assert type(data.get(lineno_key)) == int - assert data.get(lineno_key) > 0 - assert data.get(filepath_key) == "asyncpg_helpers/helpers.py" + assert type(data.get("code.line.number")) == int + assert data.get("code.line.number") > 0 + assert data.get("code.file.path") == "asyncpg_helpers/helpers.py" assert data.get(SPANDATA.CODE_NAMESPACE) == "asyncpg_helpers.helpers" - is_relative_path = data.get(filepath_key)[0] != os.sep + is_relative_path = data.get("code.file.path")[0] != os.sep assert is_relative_path assert data.get(SPANDATA.CODE_FUNCTION) == "execute_query_in_connection" @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_no_query_source_if_duration_too_short( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ): sentry_init( integrations=[AsyncPGIntegration()], traces_sample_rate=1.0, enable_db_query_source=True, db_query_source_threshold_ms=100, - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", ) - if span_streaming: - items = capture_items("span") - - @contextmanager - def fake_record_sql_queries_streaming(*args, **kwargs): - with record_sql_queries(*args, **kwargs) as span: - pass - span._start_timestamp = datetime.datetime(2024, 1, 1, microsecond=0) - if span_streaming: - span._end_timestamp = datetime.datetime(2024, 1, 1, microsecond=99999) - else: - span._timestamp = datetime.datetime(2024, 1, 1, microsecond=99999) - yield span - - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) - - with mock.patch( - "sentry_sdk.integrations.asyncpg.record_sql_queries", - fake_record_sql_queries_streaming, - ): - await conn.execute( - "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", - ) - - await conn.close() - sentry_sdk.flush() + items = capture_items("span") - spans = [item.payload for item in items] + @contextmanager + def fake_record_sql_queries(*args, **kwargs): + with record_sql_queries(*args, **kwargs) as span: + pass + span._start_timestamp = datetime.datetime(2024, 1, 1, microsecond=0) + span._end_timestamp = datetime.datetime(2024, 1, 1, microsecond=99999) + yield span - assert len(spans) == 3 - - connect_span = spans[0] - insert_span = spans[1] - segment = spans[2] - - assert segment["name"] == "test_segment" - assert insert_span["name"].startswith("INSERT INTO") - assert connect_span["name"] == "connect" - data = insert_span.get("attributes", {}) - else: - events = capture_events() + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) - with start_transaction(name="test_transaction", sampled=True): - conn: Connection = await connect(PG_CONNECTION_URI) + with mock.patch( + "sentry_sdk.integrations.asyncpg.record_sql_queries", + fake_record_sql_queries, + ): + await conn.execute( + "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", + ) - @contextmanager - def fake_record_sql_queries(*args, **kwargs): - with record_sql_queries(*args, **kwargs) as span: - pass - span.start_timestamp = datetime.datetime(2024, 1, 1, microsecond=0) - span.timestamp = datetime.datetime(2024, 1, 1, microsecond=99999) - yield span + await conn.close() + sentry_sdk.flush() - with mock.patch( - "sentry_sdk.integrations.asyncpg.record_sql_queries", - fake_record_sql_queries, - ): - await conn.execute( - "INSERT INTO users(name, password, dob) VALUES ('Alice', 'secret', '1990-12-25')", - ) + spans = [item.payload for item in items] - await conn.close() + assert len(spans) == 3 - (event,) = events + connect_span = spans[0] + insert_span = spans[1] + segment = spans[2] - span = event["spans"][-1] - assert span["description"].startswith("INSERT INTO") - data = span.get("data", {}) + assert segment["name"] == "test_segment" + assert insert_span["name"].startswith("INSERT INTO") + assert connect_span["name"] == "connect" + data = insert_span.get("attributes", {}) assert SPANDATA.CODE_LINENO not in data assert SPANDATA.CODE_NAMESPACE not in data @@ -989,129 +840,81 @@ def fake_record_sql_queries(*args, **kwargs): @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) -async def test_span_origin(sentry_init, capture_events, capture_items, span_streaming): +async def test_span_origin(sentry_init, capture_items): sentry_init( integrations=[AsyncPGIntegration()], traces_sample_rate=1.0, - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", ) - if span_streaming: - items = capture_items("span") - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) - - await conn.execute("SELECT 1") - await conn.fetchrow("SELECT 2") - await conn.close() - sentry_sdk.flush() - - spans = [item.payload for item in items] - - assert len(spans) == 4 - - connect_span = spans[0] - select1_span = spans[1] - select2_span = spans[2] - segment = spans[3] - - assert segment["name"] == "test_segment" - assert connect_span["name"] == "connect" - assert select1_span["name"] == "SELECT 1" - assert select2_span["name"] == "SELECT 2" + items = capture_items("span") + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) - assert segment["attributes"]["sentry.origin"] == "manual" - assert connect_span["attributes"]["sentry.origin"] == "auto.db.asyncpg" - assert select1_span["attributes"]["sentry.origin"] == "auto.db.asyncpg" - assert select2_span["attributes"]["sentry.origin"] == "auto.db.asyncpg" - else: - events = capture_events() + await conn.execute("SELECT 1") + await conn.fetchrow("SELECT 2") + await conn.close() + sentry_sdk.flush() - with start_transaction(name="test_transaction"): - conn: Connection = await connect(PG_CONNECTION_URI) + spans = [item.payload for item in items] - await conn.execute("SELECT 1") - await conn.fetchrow("SELECT 2") - await conn.close() + assert len(spans) == 4 - (event,) = events + connect_span = spans[0] + select1_span = spans[1] + select2_span = spans[2] + segment = spans[3] - assert event["contexts"]["trace"]["origin"] == "manual" + assert segment["name"] == "test_segment" + assert connect_span["name"] == "connect" + assert select1_span["name"] == "SELECT 1" + assert select2_span["name"] == "SELECT 2" - for span in event["spans"]: - assert span["origin"] == "auto.db.asyncpg" + assert segment["attributes"]["sentry.origin"] == "manual" + assert connect_span["attributes"]["sentry.origin"] == "auto.db.asyncpg" + assert select1_span["attributes"]["sentry.origin"] == "auto.db.asyncpg" + assert select2_span["attributes"]["sentry.origin"] == "auto.db.asyncpg" @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_multiline_query_description_normalized( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ): sentry_init( integrations=[AsyncPGIntegration()], traces_sample_rate=1.0, - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", ) - if span_streaming: - items = capture_items("span") - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) - await conn.execute( - """ - SELECT - id, - name - FROM - users - WHERE - name = 'Alice' - """ - ) - await conn.close() - sentry_sdk.flush() - - spans = [item.payload for item in items] + items = capture_items("span") + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) + await conn.execute( + """ + SELECT + id, + name + FROM + users + WHERE + name = 'Alice' + """ + ) + await conn.close() + sentry_sdk.flush() - assert len(spans) == 3 + spans = [item.payload for item in items] - connect_span = spans[0] - select_span = spans[1] - segment = spans[2] + assert len(spans) == 3 - assert segment["name"] == "test_segment" - assert connect_span["name"] == "connect" - assert select_span["name"] == "SELECT id, name FROM users WHERE name = 'Alice'" - else: - events = capture_events() + connect_span = spans[0] + select_span = spans[1] + segment = spans[2] - with start_transaction(name="test_transaction"): - conn: Connection = await connect(PG_CONNECTION_URI) - await conn.execute( - """ - SELECT - id, - name - FROM - users - WHERE - name = 'Alice' - """ - ) - await conn.close() - - (event,) = events - - spans = [ - s - for s in event["spans"] - if s["op"] == "db" and "SELECT" in s.get("description", "") - ] - assert len(spans) == 1 - assert ( - spans[0]["description"] == "SELECT id, name FROM users WHERE name = 'Alice'" - ) + assert segment["name"] == "test_segment" + assert connect_span["name"] == "connect" + assert select_span["name"] == "SELECT id, name FROM users WHERE name = 'Alice'" @pytest.mark.asyncio @@ -1156,198 +959,143 @@ def before_send_transaction(event, hint): assert spans[0]["description"] == "filtered" -def _assert_query_source(span, span_streaming, expected_function): - if span_streaming: - data = span.get("attributes", {}) - lineno_key = "code.line.number" - filepath_key = "code.file.path" - else: - data = span.get("data", {}) - lineno_key = SPANDATA.CODE_LINENO - filepath_key = SPANDATA.CODE_FILEPATH +def _assert_query_source(span, expected_function): + data = span.get("attributes", {}) - assert lineno_key in data - assert filepath_key in data + assert "code.line.number" in data + assert "code.file.path" in data assert SPANDATA.CODE_NAMESPACE in data assert SPANDATA.CODE_FUNCTION in data - assert type(data.get(lineno_key)) == int - assert data.get(lineno_key) > 0 + assert type(data.get("code.line.number")) == int + assert data.get("code.line.number") > 0 assert data[SPANDATA.CODE_NAMESPACE] == "tests.integrations.asyncpg.test_asyncpg" - assert data.get(filepath_key).endswith("tests/integrations/asyncpg/test_asyncpg.py") - assert data.get(filepath_key)[0] != os.sep + assert data.get("code.file.path").endswith( + "tests/integrations/asyncpg/test_asyncpg.py" + ) + assert data.get("code.file.path")[0] != os.sep assert data[SPANDATA.CODE_FUNCTION] == expected_function @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_query_source_execute( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ): sentry_init( integrations=[AsyncPGIntegration()], traces_sample_rate=1.0, enable_db_query_source=True, db_query_source_threshold_ms=0, - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", ) - if span_streaming: - items = capture_items("span") - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) - await conn.execute( - "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", - "Alice", - "pw", - datetime.date(1990, 12, 25), - ) - await conn.close() - sentry_sdk.flush() - - spans = [item.payload for item in items] - assert len(spans) == 3 - - connect_span = spans[0] - query_span = spans[1] - segment = spans[2] - - assert connect_span["name"] == "connect" - assert query_span["name"].startswith("INSERT INTO") - assert segment["name"] == "test_segment" - assert segment["is_segment"] is True - else: - events = capture_events() - with start_transaction(name="test_transaction", sampled=True): - conn: Connection = await connect(PG_CONNECTION_URI) - await conn.execute( - "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", - "Alice", - "pw", - datetime.date(1990, 12, 25), - ) - await conn.close() + items = capture_items("span") + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) + await conn.execute( + "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", + "Alice", + "pw", + datetime.date(1990, 12, 25), + ) + await conn.close() + sentry_sdk.flush() - (event,) = events - spans = event["spans"] - assert len(spans) == 2 - assert spans[0]["description"] == "connect" - assert spans[1]["description"].startswith("INSERT INTO") - query_span = spans[1] + spans = [item.payload for item in items] + assert len(spans) == 3 - _assert_query_source(query_span, span_streaming, "test_query_source_execute") + connect_span = spans[0] + query_span = spans[1] + segment = spans[2] + + assert connect_span["name"] == "connect" + assert query_span["name"].startswith("INSERT INTO") + assert segment["name"] == "test_segment" + assert segment["is_segment"] is True + + _assert_query_source(query_span, "test_query_source_execute") @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_query_source_executemany( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ): sentry_init( integrations=[AsyncPGIntegration()], traces_sample_rate=1.0, enable_db_query_source=True, db_query_source_threshold_ms=0, - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", ) - if span_streaming: - items = capture_items("span") - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) - await conn.executemany( - "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", - [("Bob", "secret_pw", datetime.date(1984, 3, 1))], - ) - await conn.close() - sentry_sdk.flush() - - spans = [item.payload for item in items] - assert len(spans) == 3 - - connect_span = spans[0] - query_span = spans[1] - segment = spans[2] - - assert connect_span["name"] == "connect" - assert query_span["name"].startswith("INSERT INTO") - assert segment["name"] == "test_segment" - else: - events = capture_events() - with start_transaction(name="test_transaction", sampled=True): - conn: Connection = await connect(PG_CONNECTION_URI) - await conn.executemany( - "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", - [("Bob", "secret_pw", datetime.date(1984, 3, 1))], - ) - await conn.close() + items = capture_items("span") + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) + await conn.executemany( + "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", + [("Bob", "secret_pw", datetime.date(1984, 3, 1))], + ) + await conn.close() + sentry_sdk.flush() + + spans = [item.payload for item in items] + assert len(spans) == 3 + + connect_span = spans[0] + query_span = spans[1] + segment = spans[2] - (event,) = events - spans = event["spans"] - assert len(spans) == 2 - assert spans[0]["description"] == "connect" - assert spans[1]["description"].startswith("INSERT INTO") - query_span = spans[1] + assert connect_span["name"] == "connect" + assert query_span["name"].startswith("INSERT INTO") + assert segment["name"] == "test_segment" - _assert_query_source(query_span, span_streaming, "test_query_source_executemany") + _assert_query_source(query_span, "test_query_source_executemany") @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_query_source_prepare( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ): sentry_init( integrations=[AsyncPGIntegration()], traces_sample_rate=1.0, enable_db_query_source=True, db_query_source_threshold_ms=0, - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", ) - if span_streaming: - items = capture_items("span") - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) - await conn.prepare("SELECT * FROM users WHERE name = $1") - await conn.close() - sentry_sdk.flush() - - spans = [item.payload for item in items] - assert len(spans) == 3 - connect_span = spans[0] - query_span = spans[1] - segment = spans[2] - - assert connect_span["name"] == "connect" - assert query_span["name"] == "SELECT * FROM users WHERE name = $1" - assert segment["name"] == "test_segment" - - assert ( - query_span["attributes"][SPANDATA.DB_QUERY_TEXT] - == "SELECT * FROM users WHERE name = $1" - ) - else: - events = capture_events() - with start_transaction(name="test_transaction", sampled=True): - conn: Connection = await connect(PG_CONNECTION_URI) - await conn.prepare("SELECT * FROM users WHERE name = $1") - await conn.close() + items = capture_items("span") + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) + await conn.prepare("SELECT * FROM users WHERE name = $1") + await conn.close() + sentry_sdk.flush() - (event,) = events - spans = event["spans"] - assert len(spans) == 2 - assert spans[0]["description"] == "connect" - assert spans[1]["description"] == "SELECT * FROM users WHERE name = $1" - query_span = spans[1] + spans = [item.payload for item in items] + assert len(spans) == 3 + connect_span = spans[0] + query_span = spans[1] + segment = spans[2] - _assert_query_source(query_span, span_streaming, "test_query_source_prepare") + assert connect_span["name"] == "connect" + assert query_span["name"] == "SELECT * FROM users WHERE name = $1" + assert segment["name"] == "test_segment" + + assert ( + query_span["attributes"][SPANDATA.DB_QUERY_TEXT] + == "SELECT * FROM users WHERE name = $1" + ) + + _assert_query_source(query_span, "test_query_source_prepare") @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_cursor_iteration_creates_db_cursor_iter_spans( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ) -> None: """ Regression test for https://github.com/getsentry/sentry-python/issues/6576 @@ -1364,196 +1112,115 @@ async def test_cursor_iteration_creates_db_cursor_iter_spans( sentry_init( integrations=[AsyncPGIntegration()], traces_sample_rate=1.0, - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", ) - if span_streaming: - items = capture_items("span") - - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) + items = capture_items("span") - await conn.executemany( - "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", - [(f"user-{i}", "pw", datetime.date(1990, 1, 1)) for i in range(20)], - ) - - async with conn.transaction(): - async for _record in conn.cursor("SELECT * FROM users", prefetch=5): - pass - - await conn.close() - - sentry_sdk.flush() - - cursor_iter_spans = [ - item.payload - for item in items - if item.payload.get("name") == "SELECT * FROM users" - ] - - assert len(cursor_iter_spans) == 5 - for span in cursor_iter_spans: - assert span["attributes"]["sentry.op"] == OP.DB_CURSOR_ITERATOR - assert span["attributes"][SPANDATA.DB_QUERY_TEXT] == "SELECT * FROM users" - else: - events = capture_events() - - with start_transaction(name="test_transaction", sampled=True): - conn: Connection = await connect(PG_CONNECTION_URI) + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) - await conn.executemany( - "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", - [(f"user-{i}", "pw", datetime.date(1990, 1, 1)) for i in range(20)], - ) + await conn.executemany( + "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", + [(f"user-{i}", "pw", datetime.date(1990, 1, 1)) for i in range(20)], + ) - async with conn.transaction(): - async for _record in conn.cursor("SELECT * FROM users", prefetch=5): - pass + async with conn.transaction(): + async for _record in conn.cursor("SELECT * FROM users", prefetch=5): + pass - await conn.close() + await conn.close() - (event,) = events + sentry_sdk.flush() - cursor_iter_spans = [ - s for s in event["spans"] if s.get("description") == "SELECT * FROM users" - ] + cursor_iter_spans = [ + item.payload + for item in items + if item.payload.get("name") == "SELECT * FROM users" + ] - assert len(cursor_iter_spans) == 5 - for span in cursor_iter_spans: - assert span["op"] == OP.DB_CURSOR_ITERATOR + assert len(cursor_iter_spans) == 5 + for span in cursor_iter_spans: + assert span["attributes"]["sentry.op"] == OP.DB_CURSOR_ITERATOR + assert span["attributes"][SPANDATA.DB_QUERY_TEXT] == "SELECT * FROM users" @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_cursor_fetch_methods_create_spans( - sentry_init, capture_events, capture_items, span_streaming + sentry_init, + capture_items, ) -> None: sentry_init( integrations=[AsyncPGIntegration()], traces_sample_rate=1.0, enable_db_query_source=True, db_query_source_threshold_ms=0, - trace_lifecycle="stream" if span_streaming else "static", + trace_lifecycle="stream", ) - if span_streaming: - items = capture_items() - - with sentry_sdk.traces.start_span(name="test_segment"): - conn: Connection = await connect(PG_CONNECTION_URI) - - await conn.executemany( - "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", - [ - ("Bob", "secret_pw", datetime.date(1984, 3, 1)), - ("Alice", "pw", datetime.date(1990, 12, 25)), - ], - ) - - async with conn.transaction(): - cur = await conn.cursor( - "SELECT * FROM users WHERE dob > $1", datetime.date(1970, 1, 1) - ) - # These exercise the `_exec` patch - await cur.fetchrow() - await cur.fetchrow() - - await conn.close() + items = capture_items() - sentry_sdk.flush() - - spans = [item.payload for item in items] - - assert len(spans) == 7 - - connect_span = spans[0] - executemany_span = spans[1] - begin_span = spans[2] - fetchrow_span_1 = spans[3] - fetchrow_span_2 = spans[4] - commit_span = spans[5] - _segment_span = spans[6] + with sentry_sdk.traces.start_span(name="test_segment"): + conn: Connection = await connect(PG_CONNECTION_URI) - assert connect_span["name"] == "connect" - assert ( - executemany_span["name"] - == "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)" - ) - assert begin_span["name"] == "BEGIN;" - assert fetchrow_span_1["name"] == "SELECT * FROM users WHERE dob > $1" - assert fetchrow_span_2["name"] == "SELECT * FROM users WHERE dob > $1" - assert commit_span["name"] == "COMMIT;" - - assert ( - fetchrow_span_1["attributes"][SPANDATA.DB_QUERY_TEXT] - == "SELECT * FROM users WHERE dob > $1" - ) - assert ( - fetchrow_span_2["attributes"][SPANDATA.DB_QUERY_TEXT] - == "SELECT * FROM users WHERE dob > $1" + await conn.executemany( + "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", + [ + ("Bob", "secret_pw", datetime.date(1984, 3, 1)), + ("Alice", "pw", datetime.date(1990, 12, 25)), + ], ) - for span in (fetchrow_span_1, fetchrow_span_2): - assert span["attributes"][SPANDATA.DB_SYSTEM_NAME] == "postgresql" - assert span["attributes"][SPANDATA.DB_DRIVER_NAME] == "asyncpg" - assert span["attributes"]["sentry.op"] == OP.DB_CURSOR_FETCH - assert span["attributes"]["sentry.origin"] == "auto.db.asyncpg" - - else: - events = capture_events() + async with conn.transaction(): + cur = await conn.cursor( + "SELECT * FROM users WHERE dob > $1", datetime.date(1970, 1, 1) + ) + # These exercise the `_exec` patch + await cur.fetchrow() + await cur.fetchrow() - with start_transaction(name="test_transaction"): - conn: Connection = await connect(PG_CONNECTION_URI) + await conn.close() - await conn.executemany( - "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)", - [ - ("Bob", "secret_pw", datetime.date(1984, 3, 1)), - ("Alice", "pw", datetime.date(1990, 12, 25)), - ], - ) + sentry_sdk.flush() - async with conn.transaction(): - cur = await conn.cursor( - "SELECT * FROM users WHERE dob > $1", datetime.date(1970, 1, 1) - ) - # These exercise the `_exec` patch - await cur.fetchrow() - await cur.fetchrow() + spans = [item.payload for item in items] - await conn.close() + assert len(spans) == 7 - (event,) = events + connect_span = spans[0] + executemany_span = spans[1] + begin_span = spans[2] + fetchrow_span_1 = spans[3] + fetchrow_span_2 = spans[4] + commit_span = spans[5] + _segment_span = spans[6] - assert len(event["spans"]) == 6 + assert connect_span["name"] == "connect" + assert ( + executemany_span["name"] + == "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)" + ) + assert begin_span["name"] == "BEGIN;" + assert fetchrow_span_1["name"] == "SELECT * FROM users WHERE dob > $1" + assert fetchrow_span_2["name"] == "SELECT * FROM users WHERE dob > $1" + assert commit_span["name"] == "COMMIT;" - connect_span = event["spans"][0] - executemany_span = event["spans"][1] - begin_span = event["spans"][2] - fetchrow_span_1 = event["spans"][3] - fetchrow_span_2 = event["spans"][4] - commit_span = event["spans"][5] + assert ( + fetchrow_span_1["attributes"][SPANDATA.DB_QUERY_TEXT] + == "SELECT * FROM users WHERE dob > $1" + ) + assert ( + fetchrow_span_2["attributes"][SPANDATA.DB_QUERY_TEXT] + == "SELECT * FROM users WHERE dob > $1" + ) - assert connect_span["description"] == "connect" - assert ( - executemany_span["description"] - == "INSERT INTO users(name, password, dob) VALUES($1, $2, $3)" - ) - assert begin_span["description"] == "BEGIN;" - assert fetchrow_span_1["description"] == "SELECT * FROM users WHERE dob > $1" - assert fetchrow_span_2["description"] == "SELECT * FROM users WHERE dob > $1" - assert commit_span["description"] == "COMMIT;" - - for span in (fetchrow_span_1, fetchrow_span_2): - assert span["data"]["db.cursor"] is not None - assert span["data"]["db.system"] == "postgresql" - assert span["data"]["db.driver.name"] == "asyncpg" - assert span["op"] == OP.DB_CURSOR_FETCH - assert span["origin"] == "auto.db.asyncpg" + for span in (fetchrow_span_1, fetchrow_span_2): + assert span["attributes"][SPANDATA.DB_SYSTEM_NAME] == "postgresql" + assert span["attributes"][SPANDATA.DB_DRIVER_NAME] == "asyncpg" + assert span["attributes"]["sentry.op"] == OP.DB_CURSOR_FETCH + assert span["attributes"]["sentry.origin"] == "auto.db.asyncpg" _assert_query_source( span, - span_streaming, "test_cursor_fetch_methods_create_spans", )