From 071f663a2fcb3ebe4a2f2b5bfadad6b89fea3fd5 Mon Sep 17 00:00:00 2001 From: Anton Pirker Date: Tue, 16 Apr 2024 15:31:00 +0200 Subject: [PATCH 1/2] Merge baggage instead of concatenating it. --- sentry_sdk/integrations/celery/__init__.py | 37 +++++++++++++--------- 1 file changed, 22 insertions(+), 15 deletions(-) diff --git a/sentry_sdk/integrations/celery/__init__.py b/sentry_sdk/integrations/celery/__init__.py index b3cbfe8acb..e75e54609b 100644 --- a/sentry_sdk/integrations/celery/__init__.py +++ b/sentry_sdk/integrations/celery/__init__.py @@ -16,6 +16,7 @@ from sentry_sdk.tracing import BAGGAGE_HEADER_NAME, TRANSACTION_SOURCE_TASK from sentry_sdk._types import TYPE_CHECKING from sentry_sdk.scope import Scope +from sentry_sdk.tracing_utils import Baggage from sentry_sdk.utils import ( capture_internal_exceptions, ensure_integration_enabled, @@ -177,13 +178,13 @@ def apply_async(*args, **kwargs): ) # type: Union[Span, NoOpMgr] with span_mgr as span: - incoming_headers = kwargs.get("headers") or {} + kwarg_headers = kwargs.get("headers") or {} integration = sentry_sdk.get_client().get_integration(CeleryIntegration) # If Sentry Crons monitoring for Celery Beat tasks is enabled # add start timestamp of task, if integration is not None and integration.monitor_beat_tasks: - incoming_headers.update( + kwarg_headers.update( { "sentry-monitor-start-timestamp-s": "%.9f" % _now_seconds_since_epoch(), @@ -194,7 +195,7 @@ def apply_async(*args, **kwargs): default_propagate_traces = ( integration.propagate_traces if integration is not None else True ) - propagate_traces = incoming_headers.pop( + propagate_traces = kwarg_headers.pop( "sentry-propagate-traces", default_propagate_traces ) @@ -208,21 +209,27 @@ def apply_async(*args, **kwargs): # Set Sentry trace data in the headers of the Celery task if sentry_trace_headers: # Make sure we don't overwrite existing baggage - incoming_baggage = incoming_headers.get(BAGGAGE_HEADER_NAME) + incoming_baggage = kwarg_headers.get(BAGGAGE_HEADER_NAME) sentry_baggage = sentry_trace_headers.get(BAGGAGE_HEADER_NAME) combined_baggage = sentry_baggage or incoming_baggage if sentry_baggage and incoming_baggage: - combined_baggage = "{},{}".format( - incoming_baggage, - sentry_baggage, + # Merge incoming and sentry baggage, where the sentry trace information + # in the incoming baggage takes precedence and the third-party items + # are concatenated. + incoming = Baggage.from_incoming_header(incoming_baggage) + combined = Baggage.from_incoming_header(sentry_baggage) + combined.sentry_items.update(incoming.sentry_items) + combined.third_party_items += ( + "," + incoming.third_party_items ) + combined_baggage = combined.serialize() # Set Sentry trace data to the headers of the Celery task - incoming_headers.update(sentry_trace_headers) + kwarg_headers.update(sentry_trace_headers) if combined_baggage: - incoming_headers[BAGGAGE_HEADER_NAME] = combined_baggage + kwarg_headers[BAGGAGE_HEADER_NAME] = combined_baggage # Set sentry trace data also to the inner headers of the Celery task # https://github.com/celery/celery/issues/4875 @@ -230,11 +237,11 @@ def apply_async(*args, **kwargs): # Need to setdefault the inner headers too since other # tracing tools (dd-trace-py) also employ this exact # workaround and we don't want to break them. - incoming_headers.setdefault("headers", {}).update( + kwarg_headers.setdefault("headers", {}).update( sentry_trace_headers ) if combined_baggage: - incoming_headers["headers"][ + kwarg_headers["headers"][ BAGGAGE_HEADER_NAME ] = combined_baggage @@ -245,13 +252,13 @@ def apply_async(*args, **kwargs): # Need to setdefault the inner headers too since other # tracing tools (dd-trace-py) also employ this exact # workaround and we don't want to break them. - incoming_headers.setdefault("headers", {}) - for key, value in incoming_headers.items(): + kwarg_headers.setdefault("headers", {}) + for key, value in kwarg_headers.items(): if key.startswith("sentry-"): - incoming_headers["headers"][key] = value + kwarg_headers["headers"][key] = value # Run the task (with updated headers in kwargs) - kwargs["headers"] = incoming_headers + kwargs["headers"] = kwarg_headers return f(*args, **kwargs) From a0d1c6f0b9cf5f316856c3dc434ab4031faa77e5 Mon Sep 17 00:00:00 2001 From: Anton Pirker Date: Wed, 17 Apr 2024 09:17:41 +0200 Subject: [PATCH 2/2] Added a new test --- sentry_sdk/integrations/celery/__init__.py | 6 ++-- tests/integrations/celery/test_celery.py | 34 ++++++++++++++++++++++ 2 files changed, 36 insertions(+), 4 deletions(-) diff --git a/sentry_sdk/integrations/celery/__init__.py b/sentry_sdk/integrations/celery/__init__.py index e75e54609b..912c39386f 100644 --- a/sentry_sdk/integrations/celery/__init__.py +++ b/sentry_sdk/integrations/celery/__init__.py @@ -220,10 +220,8 @@ def apply_async(*args, **kwargs): incoming = Baggage.from_incoming_header(incoming_baggage) combined = Baggage.from_incoming_header(sentry_baggage) combined.sentry_items.update(incoming.sentry_items) - combined.third_party_items += ( - "," + incoming.third_party_items - ) - combined_baggage = combined.serialize() + combined.third_party_items = ",".join([x for x in [combined.third_party_items, incoming.third_party_items] if x is not None and x != ""]) + combined_baggage = combined.serialize(include_third_party=True) # Set Sentry trace data to the headers of the Celery task kwarg_headers.update(sentry_trace_headers) diff --git a/tests/integrations/celery/test_celery.py b/tests/integrations/celery/test_celery.py index 7e0b533d4c..a5df6fbe5b 100644 --- a/tests/integrations/celery/test_celery.py +++ b/tests/integrations/celery/test_celery.py @@ -529,6 +529,40 @@ def dummy_task(self, x, y): ) +def test_baggage_propagation_sentry_items(init_celery): + BAGGAGE_VALUE = ( + "other-vendor-value-1=foo;bar;baz, sentry-trace_id=771a43a4192642f0b136d5159a501700, " + "sentry-public_key=49d0f7386ad645858ae85020e393bef3, sentry-sample_rate=0.01337, " + "sentry-user_id=Am%C3%A9lie, other-vendor-value-2=foo;bar;" + ) + + celery = init_celery(traces_sample_rate=1.0, release="abcdef") + + @celery.task(name="dummy_task", bind=True) + def dummy_task(self, x, y): + return _get_headers(self) + + with start_transaction() as transaction: + result = dummy_task.apply_async( + args=(1, 0), + headers={"baggage": BAGGAGE_VALUE}, + ).get() + + assert sorted(result["baggage"].split(",")) == sorted( + [ + "sentry-release=abcdef", + "sentry-trace_id=771a43a4192642f0b136d5159a501700", # NOT transaction.trace_id! + "sentry-environment=production", + "sentry-public_key=49d0f7386ad645858ae85020e393bef3", + "sentry-sample_rate=0.01337", + "sentry-sampled=true", + "other-vendor-value-1=foo;bar;baz", + "other-vendor-value-2=foo;bar;", + "sentry-user_id=Am%C3%A9lie", + ] + ) + + def test_sentry_propagate_traces_override(init_celery): """ Test if the `sentry-propagate-traces` header given to `apply_async`