DEV Community

Freddie Thompson
Freddie Thompson

Posted on

Context Propagation Across gRPC and REST: One Broken Hop Splits Every Trace

Your trace shows api-gateway calling orders, then stops. orders definitely called inventory, but inventory appears as a separate, parentless trace.

The context did not propagate. It happens at four specific boundaries, and each needs a different fix.

Boundary 1: protocol translation

W3C Trace Context travels as a traceparent HTTP header. gRPC has metadata, not headers, and the mapping is not automatic unless an instrumented interceptor does it.

from __future__ import annotations

import grpc
from opentelemetry import context as otel_context
from opentelemetry import propagate, trace
from opentelemetry.trace import SpanKind


class _MetadataSetter:
    def set(self, carrier: list, key: str, value: str) -> None:
        # gRPC metadata keys MUST be lowercase. A capitalised key is
        # silently dropped by some implementations and raises in others,
        # which is a fun way to lose context on one language's client only.
        carrier.append((key.lower(), value))


class _MetadataGetter:
    def get(self, carrier, key: str):
        if carrier is None:
            return None
        vals = [v for k, v in carrier if k.lower() == key.lower()]
        return vals or None

    def keys(self, carrier):
        return [k for k, _ in (carrier or [])]


_SETTER, _GETTER = _MetadataSetter(), _MetadataGetter()
_tracer = trace.get_tracer(__name__)


class TracingClientInterceptor(grpc.UnaryUnaryClientInterceptor):
    def intercept_unary_unary(self, continuation, client_call_details, request):
        with _tracer.start_as_current_span(
            client_call_details.method, kind=SpanKind.CLIENT
        ) as span:
            metadata = list(client_call_details.metadata or [])
            propagate.inject(metadata, setter=_SETTER)

            new_details = client_call_details._replace(metadata=metadata)
            response = continuation(new_details, request)

            code = getattr(response, "code", lambda: None)()
            if code not in (None, grpc.StatusCode.OK):
                span.set_attribute("rpc.grpc.status_code", str(code))

            return response


class TracingServerInterceptor(grpc.ServerInterceptor):
    def intercept_service(self, continuation, handler_call_details):
        ctx = propagate.extract(
            list(handler_call_details.invocation_metadata or []),
            getter=_GETTER
        )

        token = otel_context.attach(ctx)

        try:
            return continuation(handler_call_details)
        finally:
            # Detaching is not optional. A leaked context attaches itself to
            # whatever request the worker thread handles next, and you get
            # spans from request B parented under request A.
            otel_context.detach(token)
Enter fullscreen mode Exit fullscreen mode

The lowercase-metadata rule is worth internalising. It fails asymmetrically — a Go client talking to a Python server may work while the reverse does not — which makes it look like a language bug rather than a casing bug.

Boundary 2: message queues

A trace that crosses Kafka or SQS has a producer and a consumer that may run hours apart. Parent-child is the wrong relationship; use a link.

from opentelemetry.trace import Link, NonRecordingSpan, SpanContext


def produce(topic: str, payload: bytes, producer) -> None:
    with _tracer.start_as_current_span(
        f"publish {topic}",
        kind=SpanKind.PRODUCER
    ):
        headers: list = []
        propagate.inject(headers, setter=_SETTER)
        producer.send(topic, value=payload, headers=headers)


def consume(msg) -> None:
    parent_ctx = propagate.extract(
        list(msg.headers or []),
        getter=_GETTER
    )

    parent_span = trace.get_current_span(parent_ctx)

    # A link, not a parent. The consumer may run long after the producer's
    # trace ended; making it a child produces traces with multi-hour gaps
    # that break every latency percentile you compute from span duration.
    links = []

    sc = parent_span.get_span_context()

    if sc.is_valid:
        links.append(Link(sc))

    with _tracer.start_as_current_span(
        f"process {msg.topic}",
        kind=SpanKind.CONSUMER,
        links=links
    ) as span:
        span.set_attribute("messaging.kafka.partition", msg.partition)
        span.set_attribute("messaging.kafka.offset", msg.offset)
        handle(msg)
Enter fullscreen mode Exit fullscreen mode

Boundary 3: thread pools and async

Context is stored in a context-local. Hand work to a thread pool and the new thread has a fresh, empty one.

import contextvars
from concurrent.futures import ThreadPoolExecutor
from functools import partial


class ContextPreservingExecutor(ThreadPoolExecutor):
    """
    Copies the caller's context into the worker thread.

    Without this, every span produced inside a pool task is a new root and
    your trace loses everything that happens in a background worker - which
    is usually the slow part you were trying to see.
    """

    def submit(self, fn, /, *args, **kwargs):
        ctx = contextvars.copy_context()
        return super().submit(
            partial(ctx.run, fn),
            *args,
            **kwargs
        )
Enter fullscreen mode Exit fullscreen mode

asyncio.create_task copies context automatically. loop.run_in_executor and raw Thread do not. Neither does most third-party library code that manages its own pool, which is why a trace can break inside a database driver you did not write.

Boundary 4: sampling decisions

If service A samples a trace out and service B makes its own independent decision, you get a child span with no parent. Always use ParentBased sampling so the root’s decision is authoritative downstream.

Finding the broken hop

Do not hunt for these manually. Query your tracing backend for spans with a parent ID that does not resolve:

def find_orphans(spans: list[dict]) -> dict[str, int]:
    """
    Group orphaned spans by service. The service that shows up is the
    one RECEIVING broken context - so the bug is in its caller, or in its
    own extract path. Run this weekly; it catches a regression the day a
    new service ships without the interceptor.
    """
    known = {s["span_id"] for s in spans}
    orphans: dict[str, int] = {}

    for s in spans:
        parent = s.get("parent_span_id")

        if parent and parent not in known:
            svc = s.get("service_name", "unknown")
            orphans[svc] = orphans.get(svc, 0) + 1

    return dict(
        sorted(
            orphans.items(),
            key=lambda kv: -kv[1]
        )
    )
Enter fullscreen mode Exit fullscreen mode

Caveat: a parent can be legitimately missing because it was sampled out or is still in flight. Run this over a completed window and treat a rate as the signal, not individual cases. A service that jumps from 2% orphans to 60% after a deploy has a propagation bug; one that sits steadily at 3% probably has sampling.

Propagate more than trace context

traceparent gets you the trace. Baggage carries application context — tenant ID, feature flag cohort, request priority — to every downstream service without threading it through every function signature.

Two warnings. Baggage travels in headers to every hop, including third parties, so it must never contain anything sensitive. And it costs bytes on every call; a few small keys are fine, a serialised user object is not.


We build distributed backend systems at SoluLab — more on our custom software development work.

Top comments (0)