From 514b3e1754bb74f4480098951468afcfa9265fcf Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Mon, 24 Aug 2026 08:52:11 +0200 Subject: [PATCH 1/8] chore(rq): Remove transaction-based tracing --- tests/integrations/arq/test_arq.py | 391 ++++++++--------------------- 1 file changed, 108 insertions(+), 283 deletions(-) diff --git a/tests/integrations/arq/test_arq.py b/tests/integrations/arq/test_arq.py index 5601c6dea3..59c0b28c83 100644 --- a/tests/integrations/arq/test_arq.py +++ b/tests/integrations/arq/test_arq.py @@ -44,7 +44,6 @@ def info(self, section): @pytest.fixture def init_arq(sentry_init): def inner( - span_streaming, cls_functions=None, cls_cron_jobs=None, kw_functions=None, @@ -65,7 +64,7 @@ def inner( "integrations": [ArqIntegration()], "traces_sample_rate": 1.0, "send_default_pii": True, - "trace_lifecycle": "stream" if span_streaming else "static", + "trace_lifecycle": "stream", } sentry_init_kwargs.update(init_kwargs or {}) sentry_init(**sentry_init_kwargs) @@ -94,7 +93,6 @@ class WorkerSettings: @pytest.fixture def init_arq_with_dict_settings(sentry_init): def inner( - span_streaming, cls_functions=None, cls_cron_jobs=None, kw_functions=None, @@ -115,7 +113,7 @@ def inner( "integrations": [ArqIntegration()], "traces_sample_rate": 1.0, "send_default_pii": True, - "trace_lifecycle": "stream" if span_streaming else "static", + "trace_lifecycle": "stream", } sentry_init_kwargs.update(init_kwargs or {}) sentry_init(**sentry_init_kwargs) @@ -147,7 +145,6 @@ def init_arq_with_kwarg_settings(sentry_init): """Test fixture that passes settings_cls as keyword argument only.""" def inner( - span_streaming, cls_functions=None, cls_cron_jobs=None, kw_functions=None, @@ -168,7 +165,7 @@ def inner( "integrations": [ArqIntegration()], "traces_sample_rate": 1.0, "send_default_pii": True, - "trace_lifecycle": "stream" if span_streaming else "static", + "trace_lifecycle": "stream", } sentry_init_kwargs.update(init_kwargs or {}) sentry_init(**sentry_init_kwargs) @@ -200,8 +197,7 @@ class WorkerSettings: "init_arq_settings", ["init_arq", "init_arq_with_dict_settings", "init_arq_with_kwarg_settings"], ) -@pytest.mark.parametrize("span_streaming", [True, False]) -async def test_job_result(init_arq_settings, request, span_streaming): +async def test_job_result(init_arq_settings, request, ): async def increase(ctx, num): return num + 1 @@ -227,14 +223,12 @@ async def increase(ctx, num): @pytest.mark.parametrize( "init_arq_settings", ["init_arq", "init_arq_with_dict_settings"] ) -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_job_retry( capture_events, capture_items, init_arq_settings, request, - span_streaming, -): + ): async def retry_job(ctx): if ctx["job_try"] < 2: raise arq.worker.Retry @@ -246,49 +240,29 @@ async def retry_job(ctx): pool, worker = init_fixture_method(span_streaming, [retry_job]) job = await pool.enqueue_job("retry_job") + items = capture_items("span") - if span_streaming: - items = capture_items("span") - - await worker.run_job(job.job_id, timestamp_ms()) - - sentry_sdk.flush() - spans = [item.payload for item in items] - - # The retry re-enqueue happens without an active span, so no producer - # (queue.submit.arq) span is created for it; only the consumer segment - # is emitted. The consumer segment is preceded by the redis spans for - # the re-enqueue, so it lands at index 2. - assert spans[2]["attributes"]["sentry.op"] == "queue.task.arq" - assert spans[2]["status"] == "ok" - assert spans[2]["name"] == "retry_job" - - await worker.run_job(job.job_id, timestamp_ms()) - - sentry_sdk.flush() - spans = [item.payload for item in items] + await worker.run_job(job.job_id, timestamp_ms()) - assert spans[5]["attributes"]["sentry.op"] == "queue.task.arq" - assert spans[5]["status"] == "ok" - assert spans[5]["name"] == "retry_job" - else: - events = capture_events() + sentry_sdk.flush() + spans = [item.payload for item in items] - await worker.run_job(job.job_id, timestamp_ms()) + # The retry re-enqueue happens without an active span, so no producer + # (queue.submit.arq) span is created for it; only the consumer segment + # is emitted. The consumer segment is preceded by the redis spans for + # the re-enqueue, so it lands at index 2. + assert spans[2]["attributes"]["sentry.op"] == "queue.task.arq" + assert spans[2]["status"] == "ok" + assert spans[2]["name"] == "retry_job" - event = events.pop(0) - assert event["contexts"]["trace"]["status"] == "aborted" - assert event["transaction"] == "retry_job" - assert event["tags"]["arq_task_id"] == job.job_id - assert event["extra"]["arq-job"]["retry"] == 1 + await worker.run_job(job.job_id, timestamp_ms()) - await worker.run_job(job.job_id, timestamp_ms()) + sentry_sdk.flush() + spans = [item.payload for item in items] - event = events.pop(0) - assert event["contexts"]["trace"]["status"] == "ok" - assert event["transaction"] == "retry_job" - assert event["tags"]["arq_task_id"] == job.job_id - assert event["extra"]["arq-job"]["retry"] == 2 + assert spans[5]["attributes"]["sentry.op"] == "queue.task.arq" + assert spans[5]["status"] == "ok" + assert spans[5]["name"] == "retry_job" @pytest.mark.parametrize( @@ -299,7 +273,6 @@ async def retry_job(ctx): "init_arq_settings", ["init_arq", "init_arq_with_dict_settings"] ) @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_job_transaction( capture_events, capture_items, @@ -307,8 +280,7 @@ async def test_job_transaction( source, job_fails, request, - span_streaming, -): + ): async def division(_, a, b=0): return a / b @@ -327,134 +299,62 @@ async def division(_, a, b=0): ) job = await pool.enqueue_job("division", 1, b=int(not job_fails)) + items = capture_items("event", "span") - if span_streaming: - items = capture_items("event", "span") - - await worker.run_job(job.job_id, timestamp_ms()) - - loop = asyncio.get_event_loop() - task = loop.create_task(worker.async_run()) - await asyncio.sleep(1) - - task.cancel() - - await worker.close() - - events = [item.payload for item in items if item.type == "event"] - if job_fails: - error_func_event = events.pop(0) - error_cron_event = events.pop(0) - - assert ( - error_func_event["exception"]["values"][0]["type"] - == "ZeroDivisionError" - ) - assert ( - error_func_event["exception"]["values"][0]["mechanism"]["type"] == "arq" - ) - - func_extra = error_func_event["extra"]["arq-job"] - assert func_extra["task"] == "division" + await worker.run_job(job.job_id, timestamp_ms()) - assert ( - error_cron_event["exception"]["values"][0]["type"] - == "ZeroDivisionError" - ) - assert ( - error_cron_event["exception"]["values"][0]["mechanism"]["type"] == "arq" - ) + loop = asyncio.get_event_loop() + task = loop.create_task(worker.async_run()) + await asyncio.sleep(1) - cron_extra = error_cron_event["extra"]["arq-job"] - assert cron_extra["task"] == "cron:division" + task.cancel() - sentry_sdk.flush() - spans = [item.payload for item in items if item.type == "span"] + await worker.close() - task_spans = [ - span - for span in spans - if span["attributes"].get("sentry.op") == "queue.task.arq" - ] + events = [item.payload for item in items if item.type == "event"] + if job_fails: + error_func_event = events.pop(0) + error_cron_event = events.pop(0) - division_span = next(span for span in task_spans if span["name"] == "division") - assert division_span["attributes"]["sentry.segment.name.source"] == "task" assert ( - division_span["attributes"][SPANDATA.MESSAGING_DESTINATION_NAME] - == worker.queue_name + error_func_event["exception"]["values"][0]["type"] + == "ZeroDivisionError" ) - - assert any(span["name"] == "cron:division" for span in task_spans) - else: - events = capture_events() - - await worker.run_job(job.job_id, timestamp_ms()) - - loop = asyncio.get_event_loop() - task = loop.create_task(worker.async_run()) - await asyncio.sleep(1) - - task.cancel() - - await worker.close() - - if job_fails: - error_func_event = events.pop(0) - error_cron_event = events.pop(1) - - assert ( - error_func_event["exception"]["values"][0]["type"] - == "ZeroDivisionError" - ) - assert ( - error_func_event["exception"]["values"][0]["mechanism"]["type"] == "arq" - ) - - func_extra = error_func_event["extra"]["arq-job"] - assert func_extra["task"] == "division" - - assert ( - error_cron_event["exception"]["values"][0]["type"] - == "ZeroDivisionError" - ) - assert ( - error_cron_event["exception"]["values"][0]["mechanism"]["type"] == "arq" - ) - - cron_extra = error_cron_event["extra"]["arq-job"] - assert cron_extra["task"] == "cron:division" - - [func_event, cron_event] = events - - assert func_event["type"] == "transaction" - assert func_event["transaction"] == "division" - assert func_event["transaction_info"] == {"source": "task"} assert ( - func_event["contexts"]["trace"]["data"][SPANDATA.MESSAGING_DESTINATION_NAME] - == worker.queue_name + error_func_event["exception"]["values"][0]["mechanism"]["type"] == "arq" ) - assert "arq_task_id" in func_event["tags"] - assert "arq_task_retry" in func_event["tags"] + func_extra = error_func_event["extra"]["arq-job"] + assert func_extra["task"] == "division" - func_extra = func_event["extra"]["arq-job"] + assert ( + error_cron_event["exception"]["values"][0]["type"] + == "ZeroDivisionError" + ) + assert ( + error_cron_event["exception"]["values"][0]["mechanism"]["type"] == "arq" + ) - assert func_extra["task"] == "division" - assert func_extra["kwargs"] == {"b": int(not job_fails)} - assert func_extra["retry"] == 1 + cron_extra = error_cron_event["extra"]["arq-job"] + assert cron_extra["task"] == "cron:division" - assert cron_event["type"] == "transaction" - assert cron_event["transaction"] == "cron:division" - assert cron_event["transaction_info"] == {"source": "task"} + sentry_sdk.flush() + spans = [item.payload for item in items if item.type == "span"] - assert "arq_task_id" in cron_event["tags"] - assert "arq_task_retry" in cron_event["tags"] + task_spans = [ + span + for span in spans + if span["attributes"].get("sentry.op") == "queue.task.arq" + ] - cron_extra = cron_event["extra"]["arq-job"] + division_span = next(span for span in task_spans if span["name"] == "division") + assert division_span["attributes"]["sentry.segment.name.source"] == "task" + assert ( + division_span["attributes"][SPANDATA.MESSAGING_DESTINATION_NAME] + == worker.queue_name + ) - assert cron_extra["task"] == "cron:division" - assert cron_extra["kwargs"] == {} - assert cron_extra["retry"] == 1 + assert any(span["name"] == "cron:division" for span in task_spans) @pytest.mark.parametrize( @@ -462,12 +362,10 @@ async def division(_, a, b=0): DATA_COLLECTION_QUEUES_CASES, ) @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_job_args_kwargs_data_collection( capture_events, capture_items, init_arq, - span_streaming, init_kwargs, expected_args, expected_kwargs, @@ -484,16 +382,10 @@ async def division(_, a, b=1): ) job = await pool.enqueue_job("division", 1, b=0) - - if span_streaming: - items = capture_items("event") - else: - events = capture_events() + items = capture_items("event") await worker.run_job(job.job_id, timestamp_ms()) - - if span_streaming: - events = [item.payload for item in items] + events = [item.payload for item in items] (event,) = [event for event in events if "exception" in event] @@ -514,61 +406,41 @@ async def division(_, a, b=1): "init_arq_settings", ["init_arq", "init_arq_with_dict_settings"] ) @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_enqueue_job( capture_events, capture_items, init_arq_settings, source, request, - span_streaming, -): + ): async def dummy_job(_): pass init_fixture_method = request.getfixturevalue(init_arq_settings) pool, _ = init_fixture_method(span_streaming, **{source: [dummy_job]}) + items = capture_items("span") - if span_streaming: - items = capture_items("span") + with sentry_sdk.traces.start_span(name="custom parent") as span: + await pool.enqueue_job("dummy_job") - with sentry_sdk.traces.start_span(name="custom parent") as span: - await pool.enqueue_job("dummy_job") + sentry_sdk.flush() + spans = [item.payload for item in items] - sentry_sdk.flush() - spans = [item.payload for item in items] + assert spans[2]["is_segment"] is True + assert spans[2]["trace_id"] == span.trace_id + assert spans[2]["span_id"] == span.span_id - assert spans[2]["is_segment"] is True - assert spans[2]["trace_id"] == span.trace_id - assert spans[2]["span_id"] == span.span_id - - assert spans[1]["attributes"]["sentry.op"] == "queue.submit.arq" - assert spans[1]["name"] == "dummy_job" - else: - events = capture_events() - - with start_transaction() as transaction: - await pool.enqueue_job("dummy_job") - - (event,) = events - - assert event["contexts"]["trace"]["trace_id"] == transaction.trace_id - assert event["contexts"]["trace"]["span_id"] == transaction.span_id - - assert len(event["spans"]) - assert event["spans"][0]["op"] == "queue.submit.arq" - assert event["spans"][0]["description"] == "dummy_job" + assert spans[1]["attributes"]["sentry.op"] == "queue.submit.arq" + assert spans[1]["name"] == "dummy_job" @pytest.mark.asyncio @pytest.mark.parametrize( "init_arq_settings", ["init_arq", "init_arq_with_dict_settings"] ) -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_execute_job_without_integration( - init_arq_settings, request, span_streaming -): + init_arq_settings, request, ): async def dummy_job(_ctx): pass @@ -592,55 +464,40 @@ async def dummy_job(_ctx): "init_arq_settings", ["init_arq", "init_arq_with_dict_settings"] ) @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_span_origin_producer( capture_events, capture_items, init_arq_settings, source, request, - span_streaming, -): + ): async def dummy_job(_): pass init_fixture_method = request.getfixturevalue(init_arq_settings) pool, _ = init_fixture_method(span_streaming, **{source: [dummy_job]}) + items = capture_items("span") - if span_streaming: - items = capture_items("span") - - with sentry_sdk.traces.start_span(name="custom parent"): - await pool.enqueue_job("dummy_job") + with sentry_sdk.traces.start_span(name="custom parent"): + await pool.enqueue_job("dummy_job") - sentry_sdk.flush() - spans = [item.payload for item in items] - assert spans[2]["attributes"]["sentry.origin"] == "manual" - assert spans[1]["attributes"]["sentry.origin"] == "auto.queue.arq" - else: - events = capture_events() - - with start_transaction(): - await pool.enqueue_job("dummy_job") - - (event,) = events - assert event["contexts"]["trace"]["origin"] == "manual" - assert event["spans"][0]["origin"] == "auto.queue.arq" + sentry_sdk.flush() + spans = [item.payload for item in items] + assert spans[2]["attributes"]["sentry.origin"] == "manual" + assert spans[1]["attributes"]["sentry.origin"] == "auto.queue.arq" @pytest.mark.asyncio @pytest.mark.parametrize( "init_arq_settings", ["init_arq", "init_arq_with_dict_settings"] ) -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_span_origin_consumer( capture_events, capture_items, init_arq_settings, request, - span_streaming, -): + ): async def job(ctx): pass @@ -649,46 +506,30 @@ async def job(ctx): job.__qualname__ = job.__name__ pool, worker = init_fixture_method(span_streaming, [job]) + job = await pool.enqueue_job("job") - if span_streaming: - job = await pool.enqueue_job("job") - - items = capture_items("span") - - await worker.run_job(job.job_id, timestamp_ms()) - - sentry_sdk.flush() - spans = [item.payload for item in items] - - # No producer (queue.submit.arq) span is created for the re-enqueue - # triggered by the retry, since it happens without an active span, so - # the consumer segment lands at index 2. - assert spans[2]["attributes"]["sentry.op"] == "queue.task.arq" - assert spans[2]["attributes"]["sentry.origin"] == "auto.queue.arq" - assert spans[1]["attributes"]["sentry.origin"] == "auto.db.redis" - assert spans[0]["attributes"]["sentry.origin"] == "auto.db.redis" - else: - job = await pool.enqueue_job("job") - - events = capture_events() + items = capture_items("span") - await worker.run_job(job.job_id, timestamp_ms()) + await worker.run_job(job.job_id, timestamp_ms()) - (event,) = events + sentry_sdk.flush() + spans = [item.payload for item in items] - assert event["contexts"]["trace"]["origin"] == "auto.queue.arq" - assert event["spans"][0]["origin"] == "auto.db.redis" - assert event["spans"][1]["origin"] == "auto.db.redis" + # No producer (queue.submit.arq) span is created for the re-enqueue + # triggered by the retry, since it happens without an active span, so + # the consumer segment lands at index 2. + assert spans[2]["attributes"]["sentry.op"] == "queue.task.arq" + assert spans[2]["attributes"]["sentry.origin"] == "auto.queue.arq" + assert spans[1]["attributes"]["sentry.origin"] == "auto.db.redis" + assert spans[0]["attributes"]["sentry.origin"] == "auto.db.redis" @pytest.mark.asyncio -@pytest.mark.parametrize("span_streaming", [True, False]) async def test_job_concurrency( capture_events, capture_items, init_arq, - span_streaming, -): + ): """ 10 - division starts 70 - sleepy starts @@ -715,35 +556,19 @@ async def division(_): await pool.enqueue_job( "sleepy", _job_id="456", _defer_by=timedelta(milliseconds=70) ) + items = capture_items("event") - if span_streaming: - items = capture_items("event") - - loop = asyncio.get_event_loop() - task = loop.create_task(worker.async_run()) - await asyncio.sleep(1) - - task.cancel() - - await worker.close() - - events = [item.payload for item in items] - exception_event = events[0] - assert exception_event["exception"]["values"][0]["type"] == "ZeroDivisionError" - assert exception_event["transaction"] == "division" - else: - events = capture_events() - - loop = asyncio.get_event_loop() - task = loop.create_task(worker.async_run()) - await asyncio.sleep(1) + loop = asyncio.get_event_loop() + task = loop.create_task(worker.async_run()) + await asyncio.sleep(1) - task.cancel() + task.cancel() - await worker.close() + await worker.close() - (exception_event,) = (event for event in events if "exception" in event) - assert exception_event["exception"]["values"][0]["type"] == "ZeroDivisionError" - assert exception_event["transaction"] == "division" + events = [item.payload for item in items] + exception_event = events[0] + assert exception_event["exception"]["values"][0]["type"] == "ZeroDivisionError" + assert exception_event["transaction"] == "division" assert exception_event["extra"]["arq-job"]["task"] == "division" From 714b6619bf9c9311f4485511dee2ee85b8a09abc Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Mon, 24 Aug 2026 08:55:26 +0200 Subject: [PATCH 2/8] . --- sentry_sdk/integrations/arq.py | 60 ++++++++++-------------------- tests/integrations/arq/test_arq.py | 2 +- 2 files changed, 21 insertions(+), 41 deletions(-) diff --git a/sentry_sdk/integrations/arq.py b/sentry_sdk/integrations/arq.py index d318152f88..e7d2eeae92 100644 --- a/sentry_sdk/integrations/arq.py +++ b/sentry_sdk/integrations/arq.py @@ -79,21 +79,15 @@ async def _sentry_enqueue_job( if client.get_integration(ArqIntegration) is None: return await old_enqueue_job(self, function, *args, **kwargs) - if has_span_streaming_enabled(client.options): - if sentry_sdk.traces.get_current_span() is None: - return await old_enqueue_job(self, function, *args, **kwargs) - - with sentry_sdk.traces.start_span( - name=function, - attributes={ - "sentry.op": OP.QUEUE_SUBMIT_ARQ, - "sentry.origin": ArqIntegration.origin, - }, - ): - return await old_enqueue_job(self, function, *args, **kwargs) + if sentry_sdk.traces.get_current_span() is None: + return await old_enqueue_job(self, function, *args, **kwargs) - with sentry_sdk.start_span( - op=OP.QUEUE_SUBMIT_ARQ, name=function, origin=ArqIntegration.origin + with sentry_sdk.traces.start_span( + name=function, + attributes={ + "sentry.op": OP.QUEUE_SUBMIT_ARQ, + "sentry.origin": ArqIntegration.origin, + }, ): return await old_enqueue_job(self, function, *args, **kwargs) @@ -113,34 +107,20 @@ async def _sentry_run_job(self: "Worker", job_id: str, score: int) -> None: scope._name = "arq" scope.clear_breadcrumbs() - if has_span_streaming_enabled(client.options): - with sentry_sdk.traces.start_span( - name="unknown arq task", - attributes={ - "sentry.op": OP.QUEUE_TASK_ARQ, - "sentry.origin": ArqIntegration.origin, - "sentry.segment.name.source": SegmentNameSource.TASK, - SPANDATA.MESSAGING_MESSAGE_ID: job_id, - }, - parent_span=None, - ) as span: - if self.queue_name is not None: - span.set_attribute( - SPANDATA.MESSAGING_DESTINATION_NAME, self.queue_name - ) - return await old_run_job(self, job_id, score) - - transaction = Transaction( + with sentry_sdk.traces.start_span( name="unknown arq task", - status="ok", - op=OP.QUEUE_TASK_ARQ, - source=TransactionSource.TASK, - origin=ArqIntegration.origin, - ) - - with sentry_sdk.start_transaction(transaction) as span: + attributes={ + "sentry.op": OP.QUEUE_TASK_ARQ, + "sentry.origin": ArqIntegration.origin, + "sentry.segment.name.source": SegmentNameSource.TASK, + SPANDATA.MESSAGING_MESSAGE_ID: job_id, + }, + parent_span=None, + ) as span: if self.queue_name is not None: - span.set_data(SPANDATA.MESSAGING_DESTINATION_NAME, self.queue_name) + span.set_attribute( + SPANDATA.MESSAGING_DESTINATION_NAME, self.queue_name + ) return await old_run_job(self, job_id, score) Worker.run_job = _sentry_run_job diff --git a/tests/integrations/arq/test_arq.py b/tests/integrations/arq/test_arq.py index 59c0b28c83..d0a3e9aa1f 100644 --- a/tests/integrations/arq/test_arq.py +++ b/tests/integrations/arq/test_arq.py @@ -548,7 +548,7 @@ async def division(_): sleepy.__qualname__ = sleepy.__name__ division.__qualname__ = division.__name__ - pool, worker = init_arq(span_streaming, [sleepy, division]) + pool, worker = init_arq([sleepy, division]) await pool.enqueue_job( "division", _job_id="123", _defer_by=timedelta(milliseconds=10) From c5f73873fd038ad139ed6acaf1c592cf63f6777a Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Mon, 24 Aug 2026 08:56:34 +0200 Subject: [PATCH 3/8] . --- tests/integrations/arq/test_arq.py | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/tests/integrations/arq/test_arq.py b/tests/integrations/arq/test_arq.py index d0a3e9aa1f..a6a9cd3edf 100644 --- a/tests/integrations/arq/test_arq.py +++ b/tests/integrations/arq/test_arq.py @@ -205,7 +205,7 @@ async def increase(ctx, num): increase.__qualname__ = increase.__name__ - pool, worker = init_fixture_method(span_streaming, [increase]) + pool, worker = init_fixture_method([increase]) job = await pool.enqueue_job("increase", 3) @@ -237,7 +237,7 @@ async def retry_job(ctx): retry_job.__qualname__ = retry_job.__name__ - pool, worker = init_fixture_method(span_streaming, [retry_job]) + pool, worker = init_fixture_method([retry_job]) job = await pool.enqueue_job("retry_job") items = capture_items("span") @@ -376,7 +376,6 @@ async def division(_, a, b=1): division.__qualname__ = division.__name__ pool, worker = init_arq( - span_streaming, cls_functions=[division], init_kwargs=init_kwargs, ) @@ -418,7 +417,7 @@ async def dummy_job(_): init_fixture_method = request.getfixturevalue(init_arq_settings) - pool, _ = init_fixture_method(span_streaming, **{source: [dummy_job]}) + pool, _ = init_fixture_method(**{source: [dummy_job]}) items = capture_items("span") with sentry_sdk.traces.start_span(name="custom parent") as span: @@ -448,7 +447,7 @@ async def dummy_job(_ctx): dummy_job.__qualname__ = dummy_job.__name__ - pool, worker = init_fixture_method(span_streaming, [dummy_job]) + pool, worker = init_fixture_method([dummy_job]) # remove the integration to trigger the edge case get_client().integrations.pop("arq") @@ -476,7 +475,7 @@ async def dummy_job(_): init_fixture_method = request.getfixturevalue(init_arq_settings) - pool, _ = init_fixture_method(span_streaming, **{source: [dummy_job]}) + pool, _ = init_fixture_method(**{source: [dummy_job]}) items = capture_items("span") with sentry_sdk.traces.start_span(name="custom parent"): @@ -505,7 +504,7 @@ async def job(ctx): job.__qualname__ = job.__name__ - pool, worker = init_fixture_method(span_streaming, [job]) + pool, worker = init_fixture_method([job]) job = await pool.enqueue_job("job") items = capture_items("span") From 140ee3e5909d7b75c0da3bb2e7cbd9ae3dc3a2e1 Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Mon, 24 Aug 2026 08:57:42 +0200 Subject: [PATCH 4/8] . --- tests/integrations/arq/test_arq.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/integrations/arq/test_arq.py b/tests/integrations/arq/test_arq.py index a6a9cd3edf..213150f93a 100644 --- a/tests/integrations/arq/test_arq.py +++ b/tests/integrations/arq/test_arq.py @@ -295,7 +295,7 @@ async def division(_, a, b=0): functions_key, cron_jobs_key = source pool, worker = init_fixture_method( - span_streaming, **{functions_key: [division], cron_jobs_key: [cron_job]} + **{functions_key: [division], cron_jobs_key: [cron_job]} ) job = await pool.enqueue_job("division", 1, b=int(not job_fails)) From d10c6fa3b19a3a2c60c1506243306a1640753d80 Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Mon, 24 Aug 2026 08:58:49 +0200 Subject: [PATCH 5/8] remove unused imports --- sentry_sdk/integrations/arq.py | 1 - tests/integrations/arq/test_arq.py | 2 +- 2 files changed, 1 insertion(+), 2 deletions(-) diff --git a/sentry_sdk/integrations/arq.py b/sentry_sdk/integrations/arq.py index e7d2eeae92..c82a01da31 100644 --- a/sentry_sdk/integrations/arq.py +++ b/sentry_sdk/integrations/arq.py @@ -6,7 +6,6 @@ from sentry_sdk.integrations.logging import ignore_logger from sentry_sdk.scope import should_send_default_pii from sentry_sdk.traces import SegmentNameSource -from sentry_sdk.tracing import Transaction, TransactionSource from sentry_sdk.tracing_utils import has_span_streaming_enabled from sentry_sdk.utils import ( SENSITIVE_DATA_SUBSTITUTE, diff --git a/tests/integrations/arq/test_arq.py b/tests/integrations/arq/test_arq.py index 213150f93a..e5fb1a7843 100644 --- a/tests/integrations/arq/test_arq.py +++ b/tests/integrations/arq/test_arq.py @@ -10,7 +10,7 @@ from fakeredis.aioredis import FakeRedis import sentry_sdk -from sentry_sdk import get_client, start_transaction +from sentry_sdk import get_client from sentry_sdk.consts import SPANDATA from sentry_sdk.integrations.arq import ArqIntegration from tests.integrations.utils import DATA_COLLECTION_QUEUES_CASES From 643806b61187b0a701ba4059be56f449d62cf91e Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Mon, 24 Aug 2026 09:00:22 +0200 Subject: [PATCH 6/8] reformat --- tests/integrations/arq/test_arq.py | 39 +++++++++++++----------------- 1 file changed, 17 insertions(+), 22 deletions(-) diff --git a/tests/integrations/arq/test_arq.py b/tests/integrations/arq/test_arq.py index e5fb1a7843..bb0b00ca3a 100644 --- a/tests/integrations/arq/test_arq.py +++ b/tests/integrations/arq/test_arq.py @@ -197,7 +197,10 @@ class WorkerSettings: "init_arq_settings", ["init_arq", "init_arq_with_dict_settings", "init_arq_with_kwarg_settings"], ) -async def test_job_result(init_arq_settings, request, ): +async def test_job_result( + init_arq_settings, + request, +): async def increase(ctx, num): return num + 1 @@ -228,7 +231,7 @@ async def test_job_retry( capture_items, init_arq_settings, request, - ): +): async def retry_job(ctx): if ctx["job_try"] < 2: raise arq.worker.Retry @@ -280,7 +283,7 @@ async def test_job_transaction( source, job_fails, request, - ): +): async def division(_, a, b=0): return a / b @@ -316,24 +319,14 @@ async def division(_, a, b=0): error_func_event = events.pop(0) error_cron_event = events.pop(0) - assert ( - error_func_event["exception"]["values"][0]["type"] - == "ZeroDivisionError" - ) - assert ( - error_func_event["exception"]["values"][0]["mechanism"]["type"] == "arq" - ) + assert error_func_event["exception"]["values"][0]["type"] == "ZeroDivisionError" + assert error_func_event["exception"]["values"][0]["mechanism"]["type"] == "arq" func_extra = error_func_event["extra"]["arq-job"] assert func_extra["task"] == "division" - assert ( - error_cron_event["exception"]["values"][0]["type"] - == "ZeroDivisionError" - ) - assert ( - error_cron_event["exception"]["values"][0]["mechanism"]["type"] == "arq" - ) + assert error_cron_event["exception"]["values"][0]["type"] == "ZeroDivisionError" + assert error_cron_event["exception"]["values"][0]["mechanism"]["type"] == "arq" cron_extra = error_cron_event["extra"]["arq-job"] assert cron_extra["task"] == "cron:division" @@ -411,7 +404,7 @@ async def test_enqueue_job( init_arq_settings, source, request, - ): +): async def dummy_job(_): pass @@ -439,7 +432,9 @@ async def dummy_job(_): "init_arq_settings", ["init_arq", "init_arq_with_dict_settings"] ) async def test_execute_job_without_integration( - init_arq_settings, request, ): + init_arq_settings, + request, +): async def dummy_job(_ctx): pass @@ -469,7 +464,7 @@ async def test_span_origin_producer( init_arq_settings, source, request, - ): +): async def dummy_job(_): pass @@ -496,7 +491,7 @@ async def test_span_origin_consumer( capture_items, init_arq_settings, request, - ): +): async def job(ctx): pass @@ -528,7 +523,7 @@ async def test_job_concurrency( capture_events, capture_items, init_arq, - ): +): """ 10 - division starts 70 - sleepy starts From 47a53457ad7a5063ebc6ec9f9c6baaee876e9da9 Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Mon, 24 Aug 2026 09:20:12 +0200 Subject: [PATCH 7/8] remove unused parameters --- tests/integrations/arq/test_arq.py | 7 ------- 1 file changed, 7 deletions(-) diff --git a/tests/integrations/arq/test_arq.py b/tests/integrations/arq/test_arq.py index bb0b00ca3a..9486b92757 100644 --- a/tests/integrations/arq/test_arq.py +++ b/tests/integrations/arq/test_arq.py @@ -227,7 +227,6 @@ async def increase(ctx, num): "init_arq_settings", ["init_arq", "init_arq_with_dict_settings"] ) async def test_job_retry( - capture_events, capture_items, init_arq_settings, request, @@ -277,7 +276,6 @@ async def retry_job(ctx): ) @pytest.mark.asyncio async def test_job_transaction( - capture_events, capture_items, init_arq_settings, source, @@ -356,7 +354,6 @@ async def division(_, a, b=0): ) @pytest.mark.asyncio async def test_job_args_kwargs_data_collection( - capture_events, capture_items, init_arq, init_kwargs, @@ -399,7 +396,6 @@ async def division(_, a, b=1): ) @pytest.mark.asyncio async def test_enqueue_job( - capture_events, capture_items, init_arq_settings, source, @@ -459,7 +455,6 @@ async def dummy_job(_ctx): ) @pytest.mark.asyncio async def test_span_origin_producer( - capture_events, capture_items, init_arq_settings, source, @@ -487,7 +482,6 @@ async def dummy_job(_): "init_arq_settings", ["init_arq", "init_arq_with_dict_settings"] ) async def test_span_origin_consumer( - capture_events, capture_items, init_arq_settings, request, @@ -520,7 +514,6 @@ async def job(ctx): @pytest.mark.asyncio async def test_job_concurrency( - capture_events, capture_items, init_arq, ): From 47fee4ab56a42c3d68ce147ccbbaf710006a9b9f Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Mon, 24 Aug 2026 09:24:11 +0200 Subject: [PATCH 8/8] rename test to remove transaction reference --- tests/integrations/arq/test_arq.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/integrations/arq/test_arq.py b/tests/integrations/arq/test_arq.py index 9486b92757..1c1ea43a5c 100644 --- a/tests/integrations/arq/test_arq.py +++ b/tests/integrations/arq/test_arq.py @@ -275,7 +275,7 @@ async def retry_job(ctx): "init_arq_settings", ["init_arq", "init_arq_with_dict_settings"] ) @pytest.mark.asyncio -async def test_job_transaction( +async def test_worker_jobs( capture_items, init_arq_settings, source,