Skip to content

OTLP gRPC exporter never recovers from a wedged HTTP/2 stream: DEADLINE_EXCEEDED is retryable in name only, so the channel is never reinitialized #5683

Description

@sebastien-glon-cko

Describe your environment

OS: macOS 26.7 (arm64) for the reproducer below; also observed on Linux x86_64 in Kubernetes
Python version: 3.11.13
SDK version: 1.40.0
API version: 1.40.0
opentelemetry-exporter-otlp-proto-grpc: 1.40.0
grpcio: 1.76.0

The relevant code in exporter/opentelemetry-exporter-otlp-proto-grpc/src/opentelemetry/exporter/otlp/proto/grpc/exporter.py is unchanged in the current release v1.44.0, which I checked before filing.

What happened?

If an HTTP/2 stream stalls while the underlying connection stays healthy, OTLPSpanExporter fails permanently: it keeps reusing the same broken connection forever and every subsequent export fails with DEADLINE_EXCEEDED, even long after the network fault has cleared.

Two independent guards in _export() combine to produce this:

  1. DEADLINE_EXCEEDED is retryable in principle, but can never actually be retried. Since Make exporter timeout encompass retries/backoffs, add jitter to backoffs, cleanup code a bit #4564, timeout is a total budget:

    deadline_sec = time() + self._timeout
    for retry_num in range(_MAX_RETRYS):
        ...
        self._client.Export(..., timeout=deadline_sec - time())

    The first attempt is given the entire remaining budget, so when it times out, deadline_sec - time() is already <= 0. The loop then exits through

    or backoff_seconds > (deadline_sec - time())

    on the very first iteration. StatusCode.DEADLINE_EXCEEDED is in _RETRYABLE_ERROR_CODES, but in practice a deadline expiry is always a single attempt: retry_num never reaches 1.

  2. Channel reinitialization is therefore unreachable. The recovery path added in Fix: Reinitialize gRPC channel on UNAVAILABLE error #4825 is guarded by:

    if error.code() == StatusCode.UNAVAILABLE and retry_num == 0:
        ...
        self._channel.close()
        self._initialize_channel_and_stub()

    Because the stall is observed as DEADLINE_EXCEEDED and never as UNAVAILABLE, the channel is never closed and never recreated. The exporter holds the wedged connection for the lifetime of the process.

gRPC keepalive does not help, and this is the crux of the problem: keepalive operates at the connection level. PING frames are sent on stream 0 and are ACKed normally by a server that is alive and reachable — the transport is genuinely healthy. Only the stream carrying the Export RPC is stalled, and no connection-level mechanism can observe that. In the reproducer below, 11 PING ACKs are received while every single export fails.

This is not a theoretical concern. It affects long-lived processes with a single persistent exporter — for us, Celery-style background workers, where one wedged channel meant a worker emitted DEADLINE_EXCEEDED continuously for 53 minutes and dropped every span in that window. Short-lived or frequently recycled processes (e.g. WSGI workers cycled by max-requests) mask the bug, because a new process gets a new channel.

Steps to Reproduce

Fully self-contained — no collector and no external service required. It runs an in-process OTLP gRPC server, and puts a TCP proxy in front of it that withholds only the frames carrying the RPC (HEADERS / DATA / CONTINUATION / RST_STREAM on non-zero streams) while forwarding all connection-level frames (SETTINGS / PING / WINDOW_UPDATE on stream 0). The connection therefore stays healthy and PINGs are ACKed by the real server, while no Export RPC ever arrives — so the only status the client can observe is DEADLINE_EXCEEDED.

The wedge is connection-sticky: connections accepted in the first HEAL_AFTER seconds are wedged, later ones are fully transparent. After t=20s the fault is gone and any new connection would succeed immediately — which isolates the bug from the fault itself.

repro.py (stdlib + opentelemetry-exporter-otlp-proto-grpc only)
"""OTLP gRPC exporter never recovers from a wedged HTTP/2 stream.

    exporter --> [frame-withholding TCP proxy] --> in-process OTLP gRPC server

Run:  python repro.py
"""

import socket
import struct
import threading
import time
from concurrent import futures

import grpc
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.proto.collector.trace.v1 import trace_service_pb2 as pb2
from opentelemetry.proto.collector.trace.v1 import trace_service_pb2_grpc as pb2_grpc
from opentelemetry.sdk.trace import SpanProcessor, TracerProvider
from opentelemetry.sdk.trace.export import SpanExportResult

_DATA, _HEADERS, _RST_STREAM, _CONTINUATION, _PING = 0x0, 0x1, 0x3, 0x9, 0x6
_PREFACE_LEN = 24

HEAL_AFTER = 20.0       # seconds after which NEW connections are transparent
EXPORT_TIMEOUT = 5.0    # exporter timeout
RUN_FOR = 60.0          # total duration of the export loop

# Client keepalive, so the connection is actively probed while the stream is wedged.
CLIENT_KEEPALIVE = (
    ("grpc.keepalive_time_ms", 5000),
    ("grpc.keepalive_timeout_ms", 2000),
    ("grpc.keepalive_permit_without_calls", 1),
    ("grpc.http2.max_pings_without_data", 0),
)

# The server must PERMIT those pings, otherwise it answers GOAWAY(ENHANCE_YOUR_CALM,
# "too_many_pings"), the transport is torn down, and the exporter appears to recover for
# a reason that has nothing to do with the SDK. Do not drop these options.
SERVER_KEEPALIVE = (
    ("grpc.keepalive_permit_without_calls", 1),
    ("grpc.http2.min_ping_interval_without_data_ms", 1000),
    ("grpc.http2.max_ping_strikes", 0),
)

stats = {"conns": 0, "wedged": 0, "ping_acks": 0, "withheld": 0}


class _Collector(pb2_grpc.TraceServiceServicer):
    def Export(self, request, context):
        return pb2.ExportTraceServiceResponse()


def _recv_exact(sock, n):
    buf = b""
    while len(buf) < n:
        try:
            chunk = sock.recv(n - len(buf))
        except OSError:
            return None
        if not chunk:
            return None
        buf += chunk
    return buf


def _client_to_server(cs, ss, wedged):
    preface = _recv_exact(cs, _PREFACE_LEN)
    if preface is None:
        return
    ss.sendall(preface)
    while True:
        hdr = _recv_exact(cs, 9)
        if hdr is None:
            return
        length = int.from_bytes(hdr[0:3], "big")
        ftype = hdr[3]
        sid = struct.unpack("!I", hdr[5:9])[0] & 0x7FFFFFFF
        payload = _recv_exact(cs, length) if length else b""
        if payload is None:
            return
        if wedged and sid != 0 and ftype in (_DATA, _HEADERS, _RST_STREAM, _CONTINUATION):
            stats["withheld"] += 1
            continue
        try:
            ss.sendall(hdr + payload)
        except OSError:
            return


def _server_to_client(ss, cs):
    while True:
        hdr = _recv_exact(ss, 9)
        if hdr is None:
            return
        length = int.from_bytes(hdr[0:3], "big")
        ftype, flags = hdr[3], hdr[4]
        payload = _recv_exact(ss, length) if length else b""
        if payload is None:
            return
        if ftype == _PING and flags & 0x1:
            stats["ping_acks"] += 1
        try:
            cs.sendall(hdr + payload)
        except OSError:
            return


def _proxy(listener, upstream_port, started_at):
    while True:
        try:
            cs, _ = listener.accept()
        except OSError:
            return
        wedged = (time.monotonic() - started_at) < HEAL_AFTER
        stats["conns"] += 1
        stats["wedged"] += 1 if wedged else 0
        ss = socket.create_connection(("127.0.0.1", upstream_port))
        for fn, args in ((_client_to_server, (cs, ss, wedged)), (_server_to_client, (ss, cs))):
            threading.Thread(target=fn, args=args, daemon=True).start()


def _spans():
    captured = []

    class _Capture(SpanProcessor):
        def on_end(self, span):
            captured.append(span)

    provider = TracerProvider()
    provider.add_span_processor(_Capture())
    tracer = provider.get_tracer("repro")
    for i in range(3):
        with tracer.start_as_current_span(f"span-{i}"):
            pass
    return captured


def main():
    server = grpc.server(
        futures.ThreadPoolExecutor(max_workers=4), options=SERVER_KEEPALIVE
    )
    pb2_grpc.add_TraceServiceServicer_to_server(_Collector(), server)
    upstream_port = server.add_insecure_port("127.0.0.1:0")
    server.start()

    listener = socket.socket()
    listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    listener.bind(("127.0.0.1", 0))
    listener.listen(8)
    proxy_port = listener.getsockname()[1]

    started_at = time.monotonic()
    threading.Thread(
        target=_proxy, args=(listener, upstream_port, started_at), daemon=True
    ).start()

    exporter = OTLPSpanExporter(
        endpoint=f"127.0.0.1:{proxy_port}",
        insecure=True,
        timeout=EXPORT_TIMEOUT,
        channel_options=CLIENT_KEEPALIVE,
    )
    spans = _spans()
    print(f"=== heal_after={HEAL_AFTER}s export_timeout={EXPORT_TIMEOUT}s ===", flush=True)

    ok = False
    while time.monotonic() - started_at < RUN_FOR:
        t0 = time.monotonic()
        ok = exporter.export(spans) is SpanExportResult.SUCCESS
        print(
            f"  t={time.monotonic() - started_at:5.1f}s "
            f"{'SUCCESS' if ok else 'FAILURE'} ({time.monotonic() - t0:.1f}s) "
            f"conns={stats['conns']} wedged={stats['wedged']} "
            f"ping_acks={stats['ping_acks']} withheld={stats['withheld']}",
            flush=True,
        )
        if ok:
            break
        time.sleep(0.5)

    print(
        f"=== recovered={'YES' if ok else 'NO'} conns={stats['conns']} "
        f"ping_acks={stats['ping_acks']} withheld={stats['withheld']} ===",
        flush=True,
    )
    server.stop(0)


if __name__ == "__main__":
    main()

Expected Result

Once the fault clears at t=20s, exports should start succeeding again. More generally: a persistent exporter should be able to recover from a broken connection without the process being restarted. Either

  • a DEADLINE_EXCEEDED should be able to trigger an actual retry (each attempt getting a share of the budget rather than all of it), or
  • repeated export failures should force the channel to be closed and recreated, regardless of the specific status code.

Actual Result

Every export fails for the whole run. conns stays at 1 — a new connection is never opened, so the healed path is never taken:

=== heal_after=20.0s export_timeout=5.0s ===
Failed to export traces to 127.0.0.1:49984, error code: StatusCode.DEADLINE_EXCEEDED
  t=  5.0s FAILURE (5.0s) conns=1 wedged=1 ping_acks=0 withheld=2
  t= 10.5s FAILURE (5.0s) conns=1 wedged=1 ping_acks=1 withheld=5
  t= 16.0s FAILURE (5.0s) conns=1 wedged=1 ping_acks=2 withheld=8
Failed to export traces to 127.0.0.1:49984, error code: StatusCode.DEADLINE_EXCEEDED
  t= 21.5s FAILURE (5.0s) conns=1 wedged=1 ping_acks=3 withheld=11
  t= 27.1s FAILURE (5.0s) conns=1 wedged=1 ping_acks=4 withheld=14
  t= 32.6s FAILURE (5.0s) conns=1 wedged=1 ping_acks=5 withheld=17
  t= 38.1s FAILURE (5.0s) conns=1 wedged=1 ping_acks=6 withheld=20
Failed to export traces to 127.0.0.1:49984, error code: StatusCode.DEADLINE_EXCEEDED
  t= 43.6s FAILURE (5.0s) conns=1 wedged=1 ping_acks=8 withheld=24
  t= 49.1s FAILURE (5.0s) conns=1 wedged=1 ping_acks=8 withheld=26
  t= 54.6s FAILURE (5.0s) conns=1 wedged=1 ping_acks=9 withheld=29
Failed to export traces to 127.0.0.1:49984, error code: StatusCode.DEADLINE_EXCEEDED
  t= 60.1s FAILURE (5.0s) conns=1 wedged=1 ping_acks=11 withheld=32
=== recovered=NO conns=1 ping_acks=11 withheld=33 ===

Three details worth highlighting:

  • conns=1 throughout. No reconnection, no channel reinit, across 12 consecutive failures — even though the fault cleared 40 seconds before the end of the run.
  • ping_acks=11. Connection-level keepalive is working perfectly and proves the transport is alive. It cannot see a stream-level stall.
  • Every export takes exactly 5.0s, the full timeout, in a single attempt. The "Transient error %s encountered ... retrying in %.2fs." warning is never emitted, confirming that DEADLINE_EXCEEDED is retryable in name only.

Additional context

Contrast with the UNAVAILABLE path, which works exactly as designed: in a variant of this setup where the proxy eventually drops the connection, the escalation to UNAVAILABLE immediately produces Transient error StatusCode.UNAVAILABLE ... retrying in 1.07s followed by a successful export on a fresh connection. So #4825's recovery mechanism is sound — it is simply unreachable for stream-level stalls.

As a stopgap we wrapped the exporter in a subclass that counts consecutive failures and closes + recreates the channel after two of them, independent of the status code. That restores recovery in roughly 35 s instead of never. Sketch:

class ChannelRecyclingSpanExporter(OTLPSpanExporter):
    def export(self, spans):
        result = super().export(spans)
        if result is SpanExportResult.SUCCESS:
            self._consecutive_failures = 0
            return result
        self._consecutive_failures += 1
        if self._consecutive_failures >= _MAX_CONSECUTIVE_EXPORT_FAILURES:
            self._consecutive_failures = 0
            try:
                if self._channel:
                    self._channel.close()
            except Exception:
                logger.debug("Error closing the OTLP gRPC channel", exc_info=True)
            self._initialize_channel_and_stub()
        return result

It relies on private attributes, which is why it belongs upstream rather than in user code. I am happy to open a PR if maintainers can indicate which direction is preferred — per-attempt deadlines, or a status-code-agnostic channel recycle after N consecutive failures.

One practical note for anyone trying to reproduce this: the client keepalive interval must be permitted by the server (SERVER_KEEPALIVE above). With the gRPC default min_ping_interval_without_data_ms of 5 minutes, a 5 s client ping earns a GOAWAY(ENHANCE_YOUR_CALM, "too_many_pings"), the transport is torn down, and the exporter appears to recover — for reasons that have nothing to do with the SDK. That artefact cost me a full test run before I noticed it.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions