From d1612c5f01261ff335a2a07065776ec05b00f991 Mon Sep 17 00:00:00 2001 From: pnagre05 Date: Mon, 24 Aug 2026 13:46:00 -0500 Subject: [PATCH] chore(huey): drop transaction based tracing --- sentry_sdk/integrations/huey.py | 94 +--- tests/integrations/huey/test_huey.py | 742 ++++++++++----------------- 2 files changed, 313 insertions(+), 523 deletions(-) diff --git a/sentry_sdk/integrations/huey.py b/sentry_sdk/integrations/huey.py index ea1aeed323..7eb97033f3 100644 --- a/sentry_sdk/integrations/huey.py +++ b/sentry_sdk/integrations/huey.py @@ -3,17 +3,15 @@ from typing import TYPE_CHECKING import sentry_sdk -from sentry_sdk.api import continue_trace, get_baggage, get_traceparent -from sentry_sdk.consts import OP, SPANDATA, SPANSTATUS +from sentry_sdk.api import get_baggage, get_traceparent +from sentry_sdk.consts import OP, SPANDATA from sentry_sdk.integrations import DidNotEnable, Integration, _check_minimum_version from sentry_sdk.scope import should_send_default_pii from sentry_sdk.traces import SegmentNameSource, SpanStatus, StreamedSpan from sentry_sdk.tracing import ( BAGGAGE_HEADER_NAME, SENTRY_TRACE_HEADER_NAME, - TransactionSource, ) -from sentry_sdk.tracing_utils import has_span_streaming_enabled from sentry_sdk.utils import ( SENSITIVE_DATA_SUBSTITUTE, _register_control_flow_exception, @@ -81,15 +79,11 @@ def _sentry_enqueue( else: span_name = item.name - is_span_streaming_enabled = has_span_streaming_enabled( - sentry_sdk.get_client().options - ) - no_headers_types = (PeriodicTask,) + tuple( t for t in [HueyGroup, HueyChord] if t is not None ) - if is_span_streaming_enabled and sentry_sdk.traces.get_current_span() is None: + if sentry_sdk.traces.get_current_span() is None: if not isinstance(item, no_headers_types): item.kwargs["sentry_headers"] = { BAGGAGE_HEADER_NAME: get_baggage(), @@ -97,22 +91,14 @@ def _sentry_enqueue( } return old_enqueue(self, item) - if is_span_streaming_enabled: - span_ctx = sentry_sdk.traces.start_span( - name=span_name, - attributes={ - "sentry.op": OP.QUEUE_SUBMIT_HUEY, - "sentry.origin": HueyIntegration.origin, - SPANDATA.MESSAGING_DESTINATION_NAME: self.name, - }, - ) - else: - span_ctx = sentry_sdk.start_span( - op=OP.QUEUE_SUBMIT_HUEY, - name=span_name, - origin=HueyIntegration.origin, - ) - span_ctx.set_data(SPANDATA.MESSAGING_DESTINATION_NAME, self.name) + span_ctx = sentry_sdk.traces.start_span( + name=span_name, + attributes={ + "sentry.op": OP.QUEUE_SUBMIT_HUEY, + "sentry.origin": HueyIntegration.origin, + SPANDATA.MESSAGING_DESTINATION_NAME: self.name, + }, + ) with span_ctx: if not isinstance(item, no_headers_types): @@ -166,20 +152,12 @@ def event_processor(event: "Event", hint: "Hint") -> "Optional[Event]": def _capture_exception(exc_info: "ExcInfo") -> None: scope = sentry_sdk.get_current_scope() - is_span_streaming_enabled = has_span_streaming_enabled( - sentry_sdk.get_client().options - ) - if exc_info[0] in HUEY_CONTROL_FLOW_EXCEPTIONS: - if not is_span_streaming_enabled: - scope.transaction.set_status(SPANSTATUS.ABORTED) - elif type(scope._span) is StreamedSpan: + if type(scope._span) is StreamedSpan: scope._span._segment.status = SpanStatus.OK return - if not is_span_streaming_enabled: - scope.transaction.set_status(SPANSTATUS.INTERNAL_ERROR) - elif type(scope._span) is StreamedSpan: + if type(scope._span) is StreamedSpan: scope._span._segment.status = SpanStatus.ERROR event, hint = event_from_exception( @@ -219,39 +197,23 @@ def _sentry_execute( scope.add_event_processor(_make_event_processor(task)) sentry_headers = task.kwargs.pop("sentry_headers", None) - is_span_streaming_enabled = has_span_streaming_enabled( - sentry_sdk.get_client().options + headers = sentry_headers or {} + sentry_sdk.traces.continue_trace(headers) + span_ctx = sentry_sdk.traces.start_span( + name=task.name, + attributes={ + "sentry.op": OP.QUEUE_TASK_HUEY, + "sentry.origin": HueyIntegration.origin, + "sentry.segment.name.source": SegmentNameSource.TASK, + SPANDATA.MESSAGING_DESTINATION_NAME: self.name, + "messaging.message.id": task.id, + "messaging.message.system": "huey", + "messaging.message.retry.count": (task.default_retries or 0) + - task.retries, + }, + parent_span=None, ) - if is_span_streaming_enabled: - headers = sentry_headers or {} - sentry_sdk.traces.continue_trace(headers) - span_ctx = sentry_sdk.traces.start_span( - name=task.name, - attributes={ - "sentry.op": OP.QUEUE_TASK_HUEY, - "sentry.origin": HueyIntegration.origin, - "sentry.segment.name.source": SegmentNameSource.TASK, - SPANDATA.MESSAGING_DESTINATION_NAME: self.name, - "messaging.message.id": task.id, - "messaging.message.system": "huey", - "messaging.message.retry.count": (task.default_retries or 0) - - task.retries, - }, - parent_span=None, - ) - else: - transaction = continue_trace( - sentry_headers or {}, - name=task.name, - op=OP.QUEUE_TASK_HUEY, - source=TransactionSource.TASK, - origin=HueyIntegration.origin, - ) - transaction.set_status(SPANSTATUS.OK) - span_ctx = sentry_sdk.start_transaction(transaction) - span_ctx.set_data(SPANDATA.MESSAGING_DESTINATION_NAME, self.name) - if not getattr(task, "_sentry_is_patched", False): task.execute = _wrap_task_execute(task.execute) task._sentry_is_patched = True diff --git a/tests/integrations/huey/test_huey.py b/tests/integrations/huey/test_huey.py index a81a348d87..46d6fe5bf3 100644 --- a/tests/integrations/huey/test_huey.py +++ b/tests/integrations/huey/test_huey.py @@ -6,7 +6,6 @@ from huey.exceptions import CancelExecution, RetryTask import sentry_sdk -from sentry_sdk import start_transaction from sentry_sdk.consts import OP, SPANDATA from sentry_sdk.integrations.huey import HueyIntegration from sentry_sdk.traces import SegmentNameSource, SpanStatus @@ -24,12 +23,12 @@ @pytest.fixture def init_huey(sentry_init): - def inner(has_span_streaming=None, init_kwargs=None): + def inner(init_kwargs=None): sentry_init_kwargs = { "integrations": [HueyIntegration()], "traces_sample_rate": 1.0, "send_default_pii": True, - "trace_lifecycle": "stream" if has_span_streaming else "static", + "trace_lifecycle": "stream", } sentry_init_kwargs.update(init_kwargs or {}) sentry_init(**sentry_init_kwargs) @@ -76,75 +75,35 @@ def increase(num): @pytest.mark.parametrize("task_fails", [True, False], ids=["error", "success"]) -@pytest.mark.parametrize( - "has_span_streaming", [True, False], ids=["streaming", "no_streaming"] -) -def test_task_transaction_or_segment( - capture_events, capture_items, init_huey, task_fails, has_span_streaming -): - huey = init_huey(has_span_streaming=has_span_streaming) +def test_task_transaction_or_segment(capture_items, init_huey, task_fails): + huey = init_huey() @huey.task() def division(a, b): return a / b - if has_span_streaming: - items = capture_items("span") - execute_huey_task( - huey, division, 1, int(not task_fails), exceptions=(DivisionByZero,) - ) - sentry_sdk.get_client().flush() - - payloads = [i.payload for i in items] - # The task is enqueued without a wrapping span, so in streaming mode no - # producer (queue.submit.huey) span is created (see the early return in - # patch_enqueue). Only the consumer segment is emitted. - assert len(payloads) == 1 - (execute_span,) = payloads - - assert execute_span["is_segment"] - assert execute_span["attributes"]["sentry.op"] == OP.QUEUE_TASK_HUEY - assert ( - execute_span["attributes"][SPANDATA.MESSAGING_DESTINATION_NAME] == huey.name - ) - assert execute_span["name"] == "division" - assert execute_span["status"] == ( - SpanStatus.ERROR if task_fails else SpanStatus.OK - ) - else: - events = capture_events() - execute_huey_task( - huey, division, 1, int(not task_fails), exceptions=(DivisionByZero,) - ) - - if task_fails: - error_event = events.pop(0) - assert error_event["exception"]["values"][0]["type"] == "ZeroDivisionError" - assert error_event["exception"]["values"][0]["mechanism"]["type"] == "huey" - - (event,) = events - assert event["type"] == "transaction" - assert event["transaction"] == "division" - assert event["transaction_info"] == {"source": "task"} - assert ( - event["contexts"]["trace"]["data"][SPANDATA.MESSAGING_DESTINATION_NAME] - == huey.name - ) - - if task_fails: - assert event["contexts"]["trace"]["status"] == "internal_error" - else: - assert event["contexts"]["trace"]["status"] == "ok" - - assert "huey_task_id" in event["tags"] - assert "huey_task_retry" in event["tags"] + items = capture_items("span") + execute_huey_task( + huey, division, 1, int(not task_fails), exceptions=(DivisionByZero,) + ) + sentry_sdk.get_client().flush() + payloads = [i.payload for i in items] + # The task is enqueued without a wrapping span, so in streaming mode no + # producer (queue.submit.huey) span is created (see the early return in + # patch_enqueue). Only the consumer segment is emitted. + assert len(payloads) == 1 + (execute_span,) = payloads -@pytest.mark.parametrize( - "has_span_streaming", [True, False], ids=["streaming", "no_streaming"] -) -def test_task_retry(capture_events, capture_items, init_huey, has_span_streaming): - huey = init_huey(has_span_streaming=has_span_streaming) + assert execute_span["is_segment"] + assert execute_span["attributes"]["sentry.op"] == OP.QUEUE_TASK_HUEY + assert execute_span["attributes"][SPANDATA.MESSAGING_DESTINATION_NAME] == huey.name + assert execute_span["name"] == "division" + assert execute_span["status"] == (SpanStatus.ERROR if task_fails else SpanStatus.OK) + + +def test_task_retry(capture_items, init_huey): + huey = init_huey() context = {"retry": True} @huey.task() @@ -153,108 +112,70 @@ def retry_task(context): context["retry"] = False raise RetryTask() - if has_span_streaming: - items = capture_items("span") - execute_huey_task(huey, retry_task, context) - sentry_sdk.get_client().flush() - - payloads = [i.payload for i in items] - # The initial enqueue happens without a wrapping span, so no producer - # span is created for it. The re-enqueue triggered by RetryTask happens - # inside the running consumer segment, so it does get a child span. - assert len(payloads) == 2 - - re_enqueue_span, execute_span = payloads + items = capture_items("span") + execute_huey_task(huey, retry_task, context) + sentry_sdk.get_client().flush() - assert re_enqueue_span["attributes"]["sentry.op"] == OP.QUEUE_SUBMIT_HUEY - assert not re_enqueue_span["is_segment"] + payloads = [i.payload for i in items] + # The initial enqueue happens without a wrapping span, so no producer + # span is created for it. The re-enqueue triggered by RetryTask happens + # inside the running consumer segment, so it does get a child span. + assert len(payloads) == 2 - assert execute_span["attributes"]["sentry.op"] == OP.QUEUE_TASK_HUEY - assert execute_span["is_segment"] - assert execute_span["name"] == "retry_task" - assert execute_span["status"] == SpanStatus.OK + re_enqueue_span, execute_span = payloads - assert len(huey) == 1 + assert re_enqueue_span["attributes"]["sentry.op"] == OP.QUEUE_SUBMIT_HUEY + assert not re_enqueue_span["is_segment"] - task = huey.dequeue() - huey.execute(task) + assert execute_span["attributes"]["sentry.op"] == OP.QUEUE_TASK_HUEY + assert execute_span["is_segment"] + assert execute_span["name"] == "retry_task" + assert execute_span["status"] == SpanStatus.OK - sentry_sdk.get_client().flush() + assert len(huey) == 1 - all_payloads = [i.payload for i in items] + task = huey.dequeue() + huey.execute(task) - assert len(all_payloads) == 3 - retry_span = all_payloads[2] + sentry_sdk.get_client().flush() - assert retry_span["is_segment"] - assert retry_span["name"] == "retry_task" - assert retry_span["status"] == SpanStatus.OK - assert len(huey) == 0 - else: - events = capture_events() - result = execute_huey_task(huey, retry_task, context) - (event,) = events + all_payloads = [i.payload for i in items] - assert event["transaction"] == "retry_task" - assert event["tags"]["huey_task_id"] == result.task.id - assert len(huey) == 1 + assert len(all_payloads) == 3 + retry_span = all_payloads[2] - task = huey.dequeue() - huey.execute(task) - (event, _) = events - - assert event["transaction"] == "retry_task" - assert event["tags"]["huey_task_id"] == result.task.id - assert len(huey) == 0 + assert retry_span["is_segment"] + assert retry_span["name"] == "retry_task" + assert retry_span["status"] == SpanStatus.OK + assert len(huey) == 0 -@pytest.mark.parametrize( - "has_span_streaming", [True, False], ids=["streaming", "no_streaming"] -) -def test_task_cancel_does_not_override_status( - capture_events, capture_items, init_huey, has_span_streaming -): - huey = init_huey(has_span_streaming=has_span_streaming) +def test_task_cancel_does_not_override_status(capture_items, init_huey): + huey = init_huey() @huey.task() def cancel_task(): raise CancelExecution() - if has_span_streaming: - items = capture_items("span") - execute_huey_task(huey, cancel_task) - sentry_sdk.get_client().flush() + items = capture_items("span") + execute_huey_task(huey, cancel_task) + sentry_sdk.get_client().flush() - payloads = [i.payload for i in items] - # Enqueued without a wrapping span -> no producer span in streaming mode. - assert len(payloads) == 1 - (execute_span,) = payloads - - assert execute_span["attributes"]["sentry.op"] == OP.QUEUE_TASK_HUEY - assert execute_span["is_segment"] - assert execute_span["name"] == "cancel_task" - assert execute_span["status"] == SpanStatus.OK - else: - events = capture_events() - execute_huey_task(huey, cancel_task) + payloads = [i.payload for i in items] + # Enqueued without a wrapping span -> no producer span in streaming mode. + assert len(payloads) == 1 + (execute_span,) = payloads - if HUEY_VERSION < (3, 0, 1): - (event, _) = events - else: - (event,) = events - assert event["transaction"] == "cancel_task" - assert event["contexts"]["trace"]["status"] == "aborted" + assert execute_span["attributes"]["sentry.op"] == OP.QUEUE_TASK_HUEY + assert execute_span["is_segment"] + assert execute_span["name"] == "cancel_task" + assert execute_span["status"] == SpanStatus.OK @pytest.mark.parametrize("lock_name", ["lock.a", "lock.b"], ids=["locked", "unlocked"]) -@pytest.mark.parametrize( - "has_span_streaming", [True, False], ids=["streaming", "no_streaming"] -) @pytest.mark.skipif(HUEY_VERSION < (2, 5), reason="is_locked was added in 2.5") -def test_task_lock( - capture_events, capture_items, init_huey, lock_name, has_span_streaming -): - huey = init_huey(has_span_streaming=has_span_streaming) +def test_task_lock(capture_items, init_huey, lock_name): + huey = init_huey() task_lock_name = "lock.a" should_be_locked = task_lock_name == lock_name @@ -264,39 +185,22 @@ def test_task_lock( def maybe_locked_task(): pass - if has_span_streaming: - items = capture_items("span") - with huey.lock_task(lock_name): - assert huey.is_locked(task_lock_name) == should_be_locked - execute_huey_task(huey, maybe_locked_task) - sentry_sdk.get_client().flush() - - payloads = [i.payload for i in items] - # Enqueued without a wrapping span -> no producer span in streaming mode. - assert len(payloads) == 1 - (execute_span,) = payloads - - assert execute_span["attributes"]["sentry.op"] == OP.QUEUE_TASK_HUEY - - assert execute_span["is_segment"] - assert execute_span["name"] == "maybe_locked_task" - assert execute_span["status"] == SpanStatus.OK - else: - events = capture_events() + items = capture_items("span") + with huey.lock_task(lock_name): + assert huey.is_locked(task_lock_name) == should_be_locked + execute_huey_task(huey, maybe_locked_task) + sentry_sdk.get_client().flush() - with huey.lock_task(lock_name): - assert huey.is_locked(task_lock_name) == should_be_locked - result = execute_huey_task(huey, maybe_locked_task) + payloads = [i.payload for i in items] + # Enqueued without a wrapping span -> no producer span in streaming mode. + assert len(payloads) == 1 + (execute_span,) = payloads - (event,) = events + assert execute_span["attributes"]["sentry.op"] == OP.QUEUE_TASK_HUEY - assert event["transaction"] == "maybe_locked_task" - assert event["tags"]["huey_task_id"] == result.task.id - assert ( - event["contexts"]["trace"]["status"] == "aborted" - if should_be_locked - else event["contexts"]["trace"]["status"] == "ok" - ) + assert execute_span["is_segment"] + assert execute_span["name"] == "maybe_locked_task" + assert execute_span["status"] == SpanStatus.OK assert len(huey) == 0 @@ -304,33 +208,23 @@ def maybe_locked_task(): "init_kwargs,expected_args,expected_kwargs", DATA_COLLECTION_QUEUES_CASES, ) -@pytest.mark.parametrize( - "has_span_streaming", [True, False], ids=["streaming", "no_streaming"] -) def test_task_args_kwargs_data_collection( - capture_events, capture_items, init_huey, - has_span_streaming, init_kwargs, expected_args, expected_kwargs, ): - huey = init_huey(has_span_streaming=has_span_streaming, init_kwargs=init_kwargs) + huey = init_huey(init_kwargs=init_kwargs) @huey.task() def division(a, b): return a / b - if has_span_streaming: - items = capture_items("event") - execute_huey_task(huey, division, 1, b=0, exceptions=(DivisionByZero,)) - sentry_sdk.get_client().flush() - events = [item.payload for item in items] - else: - events = capture_events() - execute_huey_task(huey, division, 1, b=0, exceptions=(DivisionByZero,)) - + items = capture_items("event") + execute_huey_task(huey, division, 1, b=0, exceptions=(DivisionByZero,)) + sentry_sdk.get_client().flush() + events = [item.payload for item in items] (event,) = [event for event in events if "exception" in event] huey_job = event["extra"]["huey-job"] @@ -343,69 +237,92 @@ def division(a, b): assert huey_job["kwargs"] == expected_kwargs -def test_huey_enqueue(init_huey, capture_events): +def test_huey_enqueue(init_huey, capture_items): huey = init_huey() @huey.task(name="different_task_name") def dummy_task(): pass - events = capture_events() + items = capture_items("span") - with start_transaction() as transaction: + with sentry_sdk.traces.start_span(name="test"): dummy_task() - (event,) = events + sentry_sdk.get_client().flush() + + payloads = [i.payload for i in items] - assert event["contexts"]["trace"]["trace_id"] == transaction.trace_id - assert event["contexts"]["trace"]["span_id"] == transaction.span_id + enqueue_span = next( + p + for p in payloads + if p.get("attributes", {}).get("sentry.op") == OP.QUEUE_SUBMIT_HUEY + ) + segment_span = next(p for p in payloads if p.get("is_segment")) - assert len(event["spans"]) - assert event["spans"][0]["op"] == "queue.submit.huey" - assert event["spans"][0]["description"] == "different_task_name" - assert event["spans"][0]["data"][SPANDATA.MESSAGING_DESTINATION_NAME] == huey.name + assert enqueue_span["trace_id"] == segment_span["trace_id"] + assert enqueue_span["parent_span_id"] == segment_span["span_id"] + assert enqueue_span["name"] == "different_task_name" + assert enqueue_span["attributes"]["sentry.op"] == OP.QUEUE_SUBMIT_HUEY + assert enqueue_span["attributes"][SPANDATA.MESSAGING_DESTINATION_NAME] == huey.name -def test_huey_propagate_trace(init_huey, capture_events): +def test_huey_propagate_trace(init_huey, capture_items): huey = init_huey() - events = capture_events() + items = capture_items("span") @huey.task() def propagated_trace_task(): pass - with start_transaction() as outer_transaction: + with sentry_sdk.traces.start_span(name="producer"): execute_huey_task(huey, propagated_trace_task) - assert ( - events[0]["transaction"] == "propagated_trace_task" - ) # the "inner" transaction - assert events[0]["contexts"]["trace"]["trace_id"] == outer_transaction.trace_id + sentry_sdk.get_client().flush() + + payloads = [i.payload for i in items] + + producer_span = next(p for p in payloads if p.get("name") == "producer") + execute_span = next( + p + for p in payloads + if p.get("attributes", {}).get("sentry.op") == OP.QUEUE_TASK_HUEY + ) + + assert execute_span["name"] == "propagated_trace_task" + assert execute_span["trace_id"] == producer_span["trace_id"] -def test_span_origin_producer(init_huey, capture_events): +def test_span_origin_producer(init_huey, capture_items): huey = init_huey() @huey.task(name="different_task_name") def dummy_task(): pass - events = capture_events() + items = capture_items("span") - with start_transaction(): + with sentry_sdk.traces.start_span(name="test"): dummy_task() - (event,) = events + sentry_sdk.get_client().flush() - assert event["contexts"]["trace"]["origin"] == "manual" - assert event["spans"][0]["origin"] == "auto.queue.huey" + payloads = [i.payload for i in items] + enqueue_span = next( + p + for p in payloads + if p.get("attributes", {}).get("sentry.op") == OP.QUEUE_SUBMIT_HUEY + ) -def test_span_origin_consumer(init_huey, capture_events): + assert enqueue_span["attributes"]["sentry.origin"] == "auto.queue.huey" + + +def test_span_origin_consumer(init_huey, capture_items): huey = init_huey() - events = capture_events() + items = capture_items("span") @huey.task() def propagated_trace_task(): @@ -413,17 +330,22 @@ def propagated_trace_task(): execute_huey_task(huey, propagated_trace_task) - (event,) = events + sentry_sdk.get_client().flush() + + payloads = [i.payload for i in items] + + execute_span = next( + p + for p in payloads + if p.get("attributes", {}).get("sentry.op") == OP.QUEUE_TASK_HUEY + ) - assert event["contexts"]["trace"]["origin"] == "auto.queue.huey" + assert execute_span["attributes"]["sentry.origin"] == "auto.queue.huey" -@pytest.mark.parametrize("has_span_streaming", [True, False]) @pytest.mark.skipif(HUEY_VERSION < (3, 0), reason="group was added in 3.0") -def test_huey_enqueue_group( - init_huey, capture_events, capture_items, has_span_streaming -): - huey = init_huey(has_span_streaming=has_span_streaming) +def test_huey_enqueue_group(init_huey, capture_items): + huey = init_huey() @huey.task() def task1(): @@ -433,130 +355,85 @@ def task1(): def task2(): pass - if has_span_streaming: - items = capture_items("span") + items = capture_items("span") - with sentry_sdk.traces.start_span(name="submission"): - huey.enqueue(group([task1.s(), task2.s()])) - - for _ in range(2): - task = huey.dequeue() - huey.execute(task) + with sentry_sdk.traces.start_span(name="submission"): + huey.enqueue(group([task1.s(), task2.s()])) - sentry_sdk.get_client().flush() - assert len(items) == 6 - - ( - task1_enqueue_span, - task2_enqueue_span, - group_span, - submission_span, - task1_execute_span, - task2_execute_span, - ) = [i.payload for i in items] - - # The enqueue happens inside a wrapping span, so the group producer - # tree is created and parented under that segment. - assert submission_span["is_segment"] - assert submission_span["name"] == "submission" - assert not group_span["is_segment"] - assert not task1_enqueue_span["is_segment"] - assert not task2_enqueue_span["is_segment"] - assert task1_execute_span["is_segment"] - assert task2_execute_span["is_segment"] - - assert group_span["parent_span_id"] == submission_span["span_id"] - assert group_span["name"] == "Huey Task Group" - assert group_span["status"] == "ok" - assert group_span["attributes"]["sentry.op"] == "queue.submit.huey" - assert group_span["attributes"]["sentry.origin"] == "auto.queue.huey" - - assert task1_enqueue_span["name"] == "task1" - assert task1_enqueue_span["status"] == "ok" - assert task1_enqueue_span["parent_span_id"] == group_span["span_id"] - assert task1_enqueue_span["attributes"]["sentry.segment.name"] == "submission" - assert task1_enqueue_span["attributes"]["sentry.op"] == "queue.submit.huey" - assert task1_enqueue_span["attributes"]["sentry.origin"] == "auto.queue.huey" - - assert task2_enqueue_span["name"] == "task2" - assert task2_enqueue_span["status"] == "ok" - assert task2_enqueue_span["parent_span_id"] == group_span["span_id"] - assert task2_enqueue_span["attributes"]["sentry.segment.name"] == "submission" - assert task2_enqueue_span["attributes"]["sentry.op"] == "queue.submit.huey" - assert task2_enqueue_span["attributes"]["sentry.origin"] == "auto.queue.huey" - - assert task1_execute_span["name"] == "task1" - assert task1_execute_span["status"] == "ok" - assert task1_execute_span["attributes"]["messaging.message.system"] == "huey" - assert task1_execute_span["parent_span_id"] == task1_enqueue_span["span_id"] - assert task1_execute_span["attributes"]["sentry.op"] == "queue.task.huey" - assert task1_execute_span["attributes"]["sentry.origin"] == "auto.queue.huey" - assert ( - task1_execute_span["attributes"]["sentry.segment.name.source"] - == SegmentNameSource.TASK - ) - assert task1_execute_span["attributes"]["messaging.message.id"] is not None - assert task1_execute_span["attributes"]["messaging.message.retry.count"] == 0 - - assert task2_execute_span["name"] == "task2" - assert task2_execute_span["status"] == "ok" - assert task2_execute_span["parent_span_id"] == task2_enqueue_span["span_id"] - assert task2_execute_span["attributes"]["messaging.message.system"] == "huey" - assert task2_execute_span["attributes"]["sentry.op"] == "queue.task.huey" - assert task2_execute_span["attributes"]["sentry.origin"] == "auto.queue.huey" - assert ( - task2_execute_span["attributes"]["sentry.segment.name.source"] - == SegmentNameSource.TASK - ) + for _ in range(2): + task = huey.dequeue() + huey.execute(task) - else: - events = capture_events() - with start_transaction() as transaction: - huey.enqueue(group([task1.s(), task2.s()])) + sentry_sdk.get_client().flush() + assert len(items) == 6 + + ( + task1_enqueue_span, + task2_enqueue_span, + group_span, + submission_span, + task1_execute_span, + task2_execute_span, + ) = [i.payload for i in items] + + # The enqueue happens inside a wrapping span, so the group producer + # tree is created and parented under that segment. + assert submission_span["is_segment"] + assert submission_span["name"] == "submission" + assert not group_span["is_segment"] + assert not task1_enqueue_span["is_segment"] + assert not task2_enqueue_span["is_segment"] + assert task1_execute_span["is_segment"] + assert task2_execute_span["is_segment"] + + assert group_span["parent_span_id"] == submission_span["span_id"] + assert group_span["name"] == "Huey Task Group" + assert group_span["status"] == "ok" + assert group_span["attributes"]["sentry.op"] == "queue.submit.huey" + assert group_span["attributes"]["sentry.origin"] == "auto.queue.huey" + + assert task1_enqueue_span["name"] == "task1" + assert task1_enqueue_span["status"] == "ok" + assert task1_enqueue_span["parent_span_id"] == group_span["span_id"] + assert task1_enqueue_span["attributes"]["sentry.segment.name"] == "submission" + assert task1_enqueue_span["attributes"]["sentry.op"] == "queue.submit.huey" + assert task1_enqueue_span["attributes"]["sentry.origin"] == "auto.queue.huey" + + assert task2_enqueue_span["name"] == "task2" + assert task2_enqueue_span["status"] == "ok" + assert task2_enqueue_span["parent_span_id"] == group_span["span_id"] + assert task2_enqueue_span["attributes"]["sentry.segment.name"] == "submission" + assert task2_enqueue_span["attributes"]["sentry.op"] == "queue.submit.huey" + assert task2_enqueue_span["attributes"]["sentry.origin"] == "auto.queue.huey" + + assert task1_execute_span["name"] == "task1" + assert task1_execute_span["status"] == "ok" + assert task1_execute_span["attributes"]["messaging.message.system"] == "huey" + assert task1_execute_span["parent_span_id"] == task1_enqueue_span["span_id"] + assert task1_execute_span["attributes"]["sentry.op"] == "queue.task.huey" + assert task1_execute_span["attributes"]["sentry.origin"] == "auto.queue.huey" + assert ( + task1_execute_span["attributes"]["sentry.segment.name.source"] + == SegmentNameSource.TASK + ) + assert task1_execute_span["attributes"]["messaging.message.id"] is not None + assert task1_execute_span["attributes"]["messaging.message.retry.count"] == 0 + + assert task2_execute_span["name"] == "task2" + assert task2_execute_span["status"] == "ok" + assert task2_execute_span["parent_span_id"] == task2_enqueue_span["span_id"] + assert task2_execute_span["attributes"]["messaging.message.system"] == "huey" + assert task2_execute_span["attributes"]["sentry.op"] == "queue.task.huey" + assert task2_execute_span["attributes"]["sentry.origin"] == "auto.queue.huey" + assert ( + task2_execute_span["attributes"]["sentry.segment.name.source"] + == SegmentNameSource.TASK + ) - for _ in range(2): - task = huey.dequeue() - huey.execute(task) - assert len(events) == 3 - - # Assert enqueue spans were successfully recorded - producer_event = events[0] - assert producer_event["type"] == "transaction" - assert producer_event["contexts"]["trace"]["trace_id"] == transaction.trace_id - assert producer_event["contexts"]["trace"]["origin"] == "manual" - - spans = producer_event["spans"] - assert len(spans) == 3 - assert spans[0]["op"] == "queue.submit.huey" - assert spans[0]["description"] == "Huey Task Group" - assert spans[1]["op"] == "queue.submit.huey" - assert spans[1]["description"] == "task1" - assert spans[2]["op"] == "queue.submit.huey" - assert spans[2]["description"] == "task2" - - # Consumer transaction assertions (one per task) - consumer_events = events[1:] - for consumer_event, expected_name in zip(consumer_events, ["task1", "task2"]): - assert consumer_event["type"] == "transaction" - assert consumer_event["transaction"] == expected_name - assert consumer_event["transaction_info"] == {"source": "task"} - assert consumer_event["contexts"]["trace"]["op"] == "queue.task.huey" - assert consumer_event["contexts"]["trace"]["origin"] == "auto.queue.huey" - assert consumer_event["contexts"]["trace"]["status"] == "ok" - assert ( - consumer_event["contexts"]["trace"]["trace_id"] == transaction.trace_id - ) - assert "huey_task_id" in consumer_event["tags"] - assert consumer_event["tags"]["huey_task_retry"] is False - - -@pytest.mark.parametrize("has_span_streaming", [True, False]) @pytest.mark.skipif(HUEY_VERSION < (3, 0), reason="chord was added in 3.0") -def test_huey_enqueue_chord( - init_huey, capture_events, capture_items, has_span_streaming -): - huey = init_huey(has_span_streaming=has_span_streaming) +def test_huey_enqueue_chord(init_huey, capture_items): + huey = init_huey() @huey.task() def task1(): @@ -566,123 +443,74 @@ def task1(): def task2(results): pass - if has_span_streaming: - items = capture_items("span") - with sentry_sdk.traces.start_span(name="submission"): - huey.enqueue(chord([task1.s()], task2.s())) + items = capture_items("span") + with sentry_sdk.traces.start_span(name="submission"): + huey.enqueue(chord([task1.s()], task2.s())) - for _ in range(2): - task = huey.dequeue() - huey.execute(task) - - sentry_sdk.get_client().flush() - assert len(items) == 6 - - ( - task1_enqueue_span, - chord_span, - submission_span, - task2_enqueue_span, - task1_execute_span, - task2_execute_span, - ) = [i.payload for i in items] - - # The enqueue happens inside a wrapping span, so the chord producer - # tree is created and parented under that segment. - assert submission_span["is_segment"] - assert submission_span["name"] == "submission" - assert not chord_span["is_segment"] - assert not task1_enqueue_span["is_segment"] - assert not task2_enqueue_span["is_segment"] - assert task1_execute_span["is_segment"] - assert task2_execute_span["is_segment"] - - assert chord_span["parent_span_id"] == submission_span["span_id"] - assert chord_span["name"] == "Huey Chord" - assert chord_span["status"] == "ok" - assert chord_span["attributes"]["sentry.op"] == "queue.submit.huey" - assert chord_span["attributes"]["sentry.origin"] == "auto.queue.huey" - - assert task1_enqueue_span["name"] == "task1" - assert task1_enqueue_span["status"] == "ok" - assert task1_enqueue_span["parent_span_id"] == chord_span["span_id"] - assert task1_enqueue_span["attributes"]["sentry.segment.name"] == "submission" - assert task1_enqueue_span["attributes"]["sentry.op"] == "queue.submit.huey" - assert task1_enqueue_span["attributes"]["sentry.origin"] == "auto.queue.huey" - - assert task1_execute_span["name"] == "task1" - assert task1_execute_span["status"] == "ok" - assert task1_execute_span["attributes"]["messaging.message.system"] == "huey" - assert task1_execute_span["parent_span_id"] == task1_enqueue_span["span_id"] - assert task1_execute_span["attributes"]["sentry.op"] == "queue.task.huey" - assert task1_execute_span["attributes"]["sentry.origin"] == "auto.queue.huey" - assert ( - task1_execute_span["attributes"]["sentry.segment.name.source"] - == SegmentNameSource.TASK - ) - # chord callback (task2) is enqueued during task1's execution - assert task2_enqueue_span["name"] == "task2" - assert task2_enqueue_span["status"] == "ok" - assert task2_enqueue_span["parent_span_id"] == task1_execute_span["span_id"] - assert task2_enqueue_span["attributes"]["sentry.segment.name"] == "task1" - assert task2_enqueue_span["attributes"]["sentry.op"] == "queue.submit.huey" - assert task2_enqueue_span["attributes"]["sentry.origin"] == "auto.queue.huey" - - assert task2_execute_span["name"] == "task2" - assert task2_execute_span["status"] == "ok" - assert task2_execute_span["parent_span_id"] == task2_enqueue_span["span_id"] - assert task2_execute_span["attributes"]["messaging.message.system"] == "huey" - assert task2_execute_span["attributes"]["sentry.op"] == "queue.task.huey" - assert task2_execute_span["attributes"]["sentry.origin"] == "auto.queue.huey" - assert ( - task2_execute_span["attributes"]["sentry.segment.name.source"] - == SegmentNameSource.TASK - ) - else: - events = capture_events() - with start_transaction() as transaction: - huey.enqueue(chord([task1.s()], task2.s())) - - for _ in range(2): - task = huey.dequeue() - huey.execute(task) + for _ in range(2): + task = huey.dequeue() + huey.execute(task) - assert len(events) == 3 - - # Enqueue spans - producer_event = events[0] - assert producer_event["contexts"]["trace"]["trace_id"] == transaction.trace_id - assert producer_event["contexts"]["trace"]["origin"] == "manual" - - spans = producer_event["spans"] - assert len(spans) == 2 - assert spans[0]["op"] == "queue.submit.huey" - assert spans[0]["description"] == "Huey Chord" - assert spans[1]["op"] == "queue.submit.huey" - assert spans[1]["description"] == "task1" - - task1_event = events[1] - # Confirm the first task enqueued the chord callback - assert len(task1_event["spans"]) == 1 - assert task1_event["spans"][0]["op"] == "queue.submit.huey" - assert task1_event["spans"][0]["description"] == "task2" - assert task1_event["type"] == "transaction" - assert task1_event["transaction"] == "task1" - assert task1_event["transaction_info"] == {"source": "task"} - assert task1_event["contexts"]["trace"]["op"] == "queue.task.huey" - assert task1_event["contexts"]["trace"]["origin"] == "auto.queue.huey" - assert task1_event["contexts"]["trace"]["status"] == "ok" - assert task1_event["contexts"]["trace"]["trace_id"] == transaction.trace_id - assert "huey_task_id" in task1_event["tags"] - assert task1_event["tags"]["huey_task_retry"] is False - - task2_event = events[2] - assert task2_event["type"] == "transaction" - assert task2_event["transaction"] == "task2" - assert task2_event["transaction_info"] == {"source": "task"} - assert task2_event["contexts"]["trace"]["op"] == "queue.task.huey" - assert task2_event["contexts"]["trace"]["origin"] == "auto.queue.huey" - assert task2_event["contexts"]["trace"]["status"] == "ok" - assert task2_event["contexts"]["trace"]["trace_id"] == transaction.trace_id - assert "huey_task_id" in task2_event["tags"] - assert task2_event["tags"]["huey_task_retry"] is False + sentry_sdk.get_client().flush() + assert len(items) == 6 + + ( + task1_enqueue_span, + chord_span, + submission_span, + task2_enqueue_span, + task1_execute_span, + task2_execute_span, + ) = [i.payload for i in items] + + # The enqueue happens inside a wrapping span, so the chord producer + # tree is created and parented under that segment. + assert submission_span["is_segment"] + assert submission_span["name"] == "submission" + assert not chord_span["is_segment"] + assert not task1_enqueue_span["is_segment"] + assert not task2_enqueue_span["is_segment"] + assert task1_execute_span["is_segment"] + assert task2_execute_span["is_segment"] + + assert chord_span["parent_span_id"] == submission_span["span_id"] + assert chord_span["name"] == "Huey Chord" + assert chord_span["status"] == "ok" + assert chord_span["attributes"]["sentry.op"] == "queue.submit.huey" + assert chord_span["attributes"]["sentry.origin"] == "auto.queue.huey" + + assert task1_enqueue_span["name"] == "task1" + assert task1_enqueue_span["status"] == "ok" + assert task1_enqueue_span["parent_span_id"] == chord_span["span_id"] + assert task1_enqueue_span["attributes"]["sentry.segment.name"] == "submission" + assert task1_enqueue_span["attributes"]["sentry.op"] == "queue.submit.huey" + assert task1_enqueue_span["attributes"]["sentry.origin"] == "auto.queue.huey" + + assert task1_execute_span["name"] == "task1" + assert task1_execute_span["status"] == "ok" + assert task1_execute_span["attributes"]["messaging.message.system"] == "huey" + assert task1_execute_span["parent_span_id"] == task1_enqueue_span["span_id"] + assert task1_execute_span["attributes"]["sentry.op"] == "queue.task.huey" + assert task1_execute_span["attributes"]["sentry.origin"] == "auto.queue.huey" + assert ( + task1_execute_span["attributes"]["sentry.segment.name.source"] + == SegmentNameSource.TASK + ) + # chord callback (task2) is enqueued during task1's execution + assert task2_enqueue_span["name"] == "task2" + assert task2_enqueue_span["status"] == "ok" + assert task2_enqueue_span["parent_span_id"] == task1_execute_span["span_id"] + assert task2_enqueue_span["attributes"]["sentry.segment.name"] == "task1" + assert task2_enqueue_span["attributes"]["sentry.op"] == "queue.submit.huey" + assert task2_enqueue_span["attributes"]["sentry.origin"] == "auto.queue.huey" + + assert task2_execute_span["name"] == "task2" + assert task2_execute_span["status"] == "ok" + assert task2_execute_span["parent_span_id"] == task2_enqueue_span["span_id"] + assert task2_execute_span["attributes"]["messaging.message.system"] == "huey" + assert task2_execute_span["attributes"]["sentry.op"] == "queue.task.huey" + assert task2_execute_span["attributes"]["sentry.origin"] == "auto.queue.huey" + assert ( + task2_execute_span["attributes"]["sentry.segment.name.source"] + == SegmentNameSource.TASK + )