Skip to content

Commit 4a34f7e

Browse files
authored
1 parent cf74a0b commit 4a34f7e

12 files changed

Lines changed: 844 additions & 1623 deletions

File tree

sentry_sdk/integrations/redis/_async_common.py

Lines changed: 31 additions & 77 deletions
Original file line numberDiff line numberDiff line change
@@ -14,10 +14,7 @@
1414
_extract_key,
1515
_get_safe_command,
1616
_set_client_data,
17-
_set_pipeline_data,
1817
)
19-
from sentry_sdk.tracing import Span
20-
from sentry_sdk.tracing_utils import has_span_streaming_enabled
2118
from sentry_sdk.utils import capture_internal_exceptions
2219

2320
if TYPE_CHECKING:
@@ -34,7 +31,7 @@ def patch_redis_async_pipeline(
3431
pipeline_cls: "Union[type[Pipeline[Any]], type[ClusterPipeline[Any]]]",
3532
is_cluster: bool,
3633
get_command_args_fn: "Any",
37-
set_db_data_fn: "Callable[[Union[Span, StreamedSpan], Any], None]",
34+
set_db_data_fn: "Callable[[StreamedSpan, Any], None]",
3835
) -> None:
3936
old_execute = pipeline_cls.execute
4037

@@ -55,44 +52,20 @@ async def _sentry_execute(self: "Any", *args: "Any", **kwargs: "Any") -> "Any":
5552
},
5653
)
5754

58-
span_streaming = has_span_streaming_enabled(client.options)
55+
if sentry_sdk.traces.get_current_span() is None:
56+
return await old_execute(self, *args, **kwargs)
5957

60-
span: "Union[Span, StreamedSpan]"
61-
if span_streaming:
62-
if sentry_sdk.traces.get_current_span() is None:
63-
return await old_execute(self, *args, **kwargs)
64-
span = sentry_sdk.traces.start_span(
65-
name="redis.pipeline.execute",
66-
attributes={
67-
"sentry.origin": SPAN_ORIGIN,
68-
"sentry.op": OP.DB_REDIS,
69-
},
70-
)
71-
else:
72-
span = sentry_sdk.start_span(
73-
op=OP.DB_REDIS,
74-
name="redis.pipeline.execute",
75-
origin=SPAN_ORIGIN,
76-
)
58+
span = sentry_sdk.traces.start_span(
59+
name="redis.pipeline.execute",
60+
attributes={
61+
"sentry.origin": SPAN_ORIGIN,
62+
"sentry.op": OP.DB_REDIS,
63+
},
64+
)
7765

7866
with span:
7967
with capture_internal_exceptions():
80-
try:
81-
command_seq = self._execution_strategy._command_queue
82-
except AttributeError:
83-
if is_cluster:
84-
command_seq = self._command_stack
85-
else:
86-
command_seq = self.command_stack
87-
8868
set_db_data_fn(span, self)
89-
_set_pipeline_data(
90-
span,
91-
is_cluster,
92-
get_command_args_fn,
93-
False if is_cluster else self.is_transaction,
94-
command_seq,
95-
)
9669

9770
return await old_execute(self, *args, **kwargs)
9871

@@ -102,7 +75,7 @@ async def _sentry_execute(self: "Any", *args: "Any", **kwargs: "Any") -> "Any":
10275
def patch_redis_async_client(
10376
cls: "Union[type[Redis[Any]], type[RedisCluster[Any]]]",
10477
is_cluster: bool,
105-
set_db_data_fn: "Callable[[Union[Span, StreamedSpan], Any], None]",
78+
set_db_data_fn: "Callable[[StreamedSpan, Any], None]",
10679
) -> None:
10780
old_execute_command = cls.execute_command
10881

@@ -134,9 +107,7 @@ async def _sentry_execute_command(
134107
data=breadcrumb_data,
135108
)
136109

137-
span_streaming = has_span_streaming_enabled(client.options)
138-
139-
if span_streaming and sentry_sdk.traces.get_current_span() is None:
110+
if sentry_sdk.traces.get_current_span() is None:
140111
return await old_execute_command(self, name, *args, **kwargs)
141112

142113
cache_properties = _compile_cache_span_properties(
@@ -152,59 +123,42 @@ async def _sentry_execute_command(
152123
_get_safe_command(name, args)
153124
)
154125

155-
cache_span: "Optional[Union[Span, StreamedSpan]]" = None
126+
cache_span: "Optional[StreamedSpan]" = None
156127
if cache_properties["is_cache_key"] and cache_properties["op"] is not None:
157-
if span_streaming:
158-
cache_span = sentry_sdk.traces.start_span(
159-
name=cache_properties["description"],
160-
attributes={
161-
"sentry.op": cache_properties["op"],
162-
"sentry.origin": SPAN_ORIGIN,
163-
**additional_cache_span_attributes,
164-
},
165-
)
166-
else:
167-
cache_span = sentry_sdk.start_span(
168-
op=cache_properties["op"],
169-
name=cache_properties["description"],
170-
origin=SPAN_ORIGIN,
171-
)
172-
cache_span.__enter__()
128+
cache_span = sentry_sdk.traces.start_span(
129+
name=cache_properties["description"],
130+
attributes={
131+
"sentry.op": cache_properties["op"],
132+
"sentry.origin": SPAN_ORIGIN,
133+
**additional_cache_span_attributes,
134+
},
135+
)
173136

174137
additional_db_span_attributes = {}
175138
with capture_internal_exceptions():
176139
additional_db_span_attributes[SPANDATA.DB_QUERY_TEXT] = _get_safe_command(
177140
name, args
178141
)
179142

180-
db_span: "Union[Span, StreamedSpan]"
181-
if span_streaming:
182-
db_span = sentry_sdk.traces.start_span(
183-
name=db_properties["description"],
184-
attributes={
185-
"sentry.op": db_properties["op"],
186-
"sentry.origin": SPAN_ORIGIN,
187-
**additional_db_span_attributes,
188-
},
189-
)
190-
else:
191-
db_span = sentry_sdk.start_span(
192-
op=db_properties["op"],
193-
name=db_properties["description"],
194-
origin=SPAN_ORIGIN,
195-
)
196-
db_span.__enter__()
143+
db_span = sentry_sdk.traces.start_span(
144+
name=db_properties["description"],
145+
attributes={
146+
"sentry.op": db_properties["op"],
147+
"sentry.origin": SPAN_ORIGIN,
148+
**additional_db_span_attributes,
149+
},
150+
)
197151

198152
set_db_data_fn(db_span, self)
199153
_set_client_data(db_span, is_cluster, name, *args)
200154

201155
value = await old_execute_command(self, name, *args, **kwargs)
202156

203-
db_span.__exit__(None, None, None)
157+
db_span.end()
204158

205159
if cache_span:
206160
_set_cache_data(cache_span, self, cache_properties, value)
207-
cache_span.__exit__(None, None, None)
161+
cache_span.end()
208162

209163
return value
210164

sentry_sdk/integrations/redis/_sync_common.py

Lines changed: 33 additions & 76 deletions
Original file line numberDiff line numberDiff line change
@@ -14,15 +14,12 @@
1414
_extract_key,
1515
_get_safe_command,
1616
_set_client_data,
17-
_set_pipeline_data,
1817
)
19-
from sentry_sdk.tracing import Span
20-
from sentry_sdk.tracing_utils import has_span_streaming_enabled
2118
from sentry_sdk.utils import capture_internal_exceptions
2219

2320
if TYPE_CHECKING:
2421
from collections.abc import Callable
25-
from typing import Any, Optional, Union
22+
from typing import Any, Optional
2623

2724
from sentry_sdk.traces import StreamedSpan
2825

@@ -31,7 +28,7 @@ def patch_redis_pipeline(
3128
pipeline_cls: "Any",
3229
is_cluster: bool,
3330
get_command_args_fn: "Any",
34-
set_db_data_fn: "Callable[[Union[Span, StreamedSpan], Any], None]",
31+
set_db_data_fn: "Callable[[StreamedSpan, Any], None]",
3532
) -> None:
3633
old_execute = pipeline_cls.execute
3734

@@ -52,41 +49,20 @@ def sentry_patched_execute(self: "Any", *args: "Any", **kwargs: "Any") -> "Any":
5249
},
5350
)
5451

55-
span_streaming = has_span_streaming_enabled(client.options)
56-
span: "Union[Span, StreamedSpan]"
57-
if span_streaming:
58-
if sentry_sdk.traces.get_current_span() is None:
59-
return old_execute(self, *args, **kwargs)
60-
span = sentry_sdk.traces.start_span(
61-
name="redis.pipeline.execute",
62-
attributes={
63-
"sentry.origin": SPAN_ORIGIN,
64-
"sentry.op": OP.DB_REDIS,
65-
},
66-
)
67-
else:
68-
span = sentry_sdk.start_span(
69-
op=OP.DB_REDIS,
70-
name="redis.pipeline.execute",
71-
origin=SPAN_ORIGIN,
72-
)
52+
if sentry_sdk.traces.get_current_span() is None:
53+
return old_execute(self, *args, **kwargs)
54+
55+
span = sentry_sdk.traces.start_span(
56+
name="redis.pipeline.execute",
57+
attributes={
58+
"sentry.origin": SPAN_ORIGIN,
59+
"sentry.op": OP.DB_REDIS,
60+
},
61+
)
7362

7463
with span:
7564
with capture_internal_exceptions():
76-
command_seq = None
77-
try:
78-
command_seq = self._execution_strategy.command_queue
79-
except AttributeError:
80-
command_seq = self.command_stack
81-
8265
set_db_data_fn(span, self)
83-
_set_pipeline_data(
84-
span,
85-
is_cluster,
86-
get_command_args_fn,
87-
False if is_cluster else self.transaction,
88-
command_seq,
89-
)
9066

9167
return old_execute(self, *args, **kwargs)
9268

@@ -96,7 +72,7 @@ def sentry_patched_execute(self: "Any", *args: "Any", **kwargs: "Any") -> "Any":
9672
def patch_redis_client(
9773
cls: "Any",
9874
is_cluster: bool,
99-
set_db_data_fn: "Callable[[Union[Span, StreamedSpan], Any], None]",
75+
set_db_data_fn: "Callable[[StreamedSpan, Any], None]",
10076
) -> None:
10177
"""
10278
This function can be used to instrument custom redis client classes or
@@ -132,9 +108,7 @@ def sentry_patched_execute_command(
132108
data=breadcrumb_data,
133109
)
134110

135-
span_streaming = has_span_streaming_enabled(client.options)
136-
137-
if span_streaming and sentry_sdk.traces.get_current_span() is None:
111+
if sentry_sdk.traces.get_current_span() is None:
138112
return old_execute_command(self, name, *args, **kwargs)
139113

140114
cache_properties = _compile_cache_span_properties(
@@ -150,59 +124,42 @@ def sentry_patched_execute_command(
150124
_get_safe_command(name, args)
151125
)
152126

153-
cache_span: "Optional[Union[Span, StreamedSpan]]" = None
127+
cache_span: "Optional[StreamedSpan]" = None
154128
if cache_properties["is_cache_key"] and cache_properties["op"] is not None:
155-
if span_streaming:
156-
cache_span = sentry_sdk.traces.start_span(
157-
name=cache_properties["description"],
158-
attributes={
159-
"sentry.op": cache_properties["op"],
160-
"sentry.origin": SPAN_ORIGIN,
161-
**additional_cache_span_attributes,
162-
},
163-
)
164-
else:
165-
cache_span = sentry_sdk.start_span(
166-
op=cache_properties["op"],
167-
name=cache_properties["description"],
168-
origin=SPAN_ORIGIN,
169-
)
170-
cache_span.__enter__()
129+
cache_span = sentry_sdk.traces.start_span(
130+
name=cache_properties["description"],
131+
attributes={
132+
"sentry.op": cache_properties["op"],
133+
"sentry.origin": SPAN_ORIGIN,
134+
**additional_cache_span_attributes,
135+
},
136+
)
171137

172138
additional_db_span_attributes = {}
173139
with capture_internal_exceptions():
174140
additional_db_span_attributes[SPANDATA.DB_QUERY_TEXT] = _get_safe_command(
175141
name, args
176142
)
177143

178-
db_span: "Union[Span, StreamedSpan]"
179-
if span_streaming:
180-
db_span = sentry_sdk.traces.start_span(
181-
name=db_properties["description"],
182-
attributes={
183-
"sentry.op": db_properties["op"],
184-
"sentry.origin": SPAN_ORIGIN,
185-
**additional_db_span_attributes,
186-
},
187-
)
188-
else:
189-
db_span = sentry_sdk.start_span(
190-
op=db_properties["op"],
191-
name=db_properties["description"],
192-
origin=SPAN_ORIGIN,
193-
)
194-
db_span.__enter__()
144+
db_span = sentry_sdk.traces.start_span(
145+
name=db_properties["description"],
146+
attributes={
147+
"sentry.op": db_properties["op"],
148+
"sentry.origin": SPAN_ORIGIN,
149+
**additional_db_span_attributes,
150+
},
151+
)
195152

196153
set_db_data_fn(db_span, self)
197154
_set_client_data(db_span, is_cluster, name, *args)
198155

199156
value = old_execute_command(self, name, *args, **kwargs)
200157

201-
db_span.__exit__(None, None, None)
158+
db_span.end()
202159

203160
if cache_span:
204161
_set_cache_data(cache_span, self, cache_properties, value)
205-
cache_span.__exit__(None, None, None)
162+
cache_span.end()
206163

207164
return value
208165

sentry_sdk/integrations/redis/consts.py

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,4 +15,3 @@
1515
"auth",
1616
]
1717
_MAX_NUM_ARGS = 10 # Trim argument lists to this many values
18-
_MAX_NUM_COMMANDS = 10 # Trim command lists to this many values

0 commit comments

Comments
 (0)