Skip to content

Commit e7a481c

Browse files
authored
fix(celery): Don't set span status to error on control flow exception (#7912)
1 parent 9666b9c commit e7a481c

2 files changed

Lines changed: 40 additions & 11 deletions

File tree

‎sentry_sdk/integrations/celery/__init__.py‎

Lines changed: 6 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
from sentry_sdk.traces import BAGGAGE_HEADER_NAME, SegmentNameSource, Span
2020
from sentry_sdk.tracing_utils import Baggage
2121
from sentry_sdk.utils import (
22+
_register_control_flow_exception,
2223
capture_internal_exceptions,
2324
event_from_exception,
2425
parse_version,
@@ -78,6 +79,7 @@ def setup_once() -> None:
7879
_patch_celery_send_task()
7980
_patch_worker_exit()
8081
_patch_producer_publish()
82+
_register_control_flow_exception(list(CELERY_CONTROL_FLOW_EXCEPTIONS))
8183

8284
# This logger logs every status of every task that ran on the worker.
8385
# Meaning that every task's breadcrumbs are full of stuff like "Task
@@ -90,25 +92,18 @@ def setup_once() -> None:
9092
ignore_logger_for_events("celery.redirected")
9193

9294

93-
def _set_status(status: str) -> None:
94-
with capture_internal_exceptions():
95-
span = sentry_sdk.get_current_span()
96-
97-
if span is not None:
98-
span.status = "ok" if status == "ok" else "error"
99-
100-
10195
def _capture_exception(task: "Any", exc_info: "ExcInfo") -> None:
10296
client = sentry_sdk.get_client()
10397
if client.get_integration(CeleryIntegration) is None:
10498
return
10599

106100
if isinstance(exc_info[1], CELERY_CONTROL_FLOW_EXCEPTIONS):
107-
# ??? Doesn't map to anything
108-
_set_status("aborted")
101+
# Expected control flow exits, not errors
109102
return
110103

111-
_set_status("internal_error")
104+
span = sentry_sdk.get_current_span()
105+
if span is not None:
106+
span.status = "error"
112107

113108
if hasattr(task, "throws") and isinstance(exc_info[1], task.throws):
114109
return

‎tests/integrations/celery/test_celery.py‎

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66
import pytest
77
from celery import VERSION, Celery
88
from celery.bin import worker
9+
from celery.exceptions import Ignore, Reject, Retry
910

1011
import sentry_sdk
1112
from sentry_sdk.integrations.celery import (
@@ -409,6 +410,39 @@ def dummy_task(self):
409410
assert e["type"] == "ZeroDivisionError"
410411

411412

413+
@pytest.mark.parametrize(
414+
"exception",
415+
[Retry, Ignore, Reject],
416+
ids=["retry", "ignore", "reject"],
417+
)
418+
def test_control_flow_exceptions_not_errors(init_celery, capture_items, exception):
419+
celery = init_celery(traces_sample_rate=1.0)
420+
should_raise = [False]
421+
422+
@celery.task(name="dummy_task")
423+
def dummy_task():
424+
if should_raise[0]:
425+
raise exception()
426+
427+
# XXX: For some reason the first call does not get instrumented properly.
428+
dummy_task.delay()
429+
sentry_sdk.flush()
430+
431+
items = capture_items("event", "span")
432+
should_raise[0] = True
433+
434+
dummy_task.delay()
435+
sentry_sdk.flush()
436+
437+
assert not [item for item in items if item.type == "event"]
438+
439+
process_span, execution_span = [item.payload for item in items]
440+
assert process_span["attributes"]["sentry.op"] == "queue.process"
441+
assert process_span["status"] == "ok"
442+
assert execution_span["is_segment"] is True
443+
assert execution_span["status"] == "ok"
444+
445+
412446
@pytest.mark.forked
413447
@pytest.mark.parametrize("newrelic_order", ["sentry_first", "sentry_last"])
414448
def test_newrelic_interference(init_celery, newrelic_order, celery_invocation):

0 commit comments

Comments
 (0)