diff --git a/.changelog/27.fixed b/.changelog/27.fixed new file mode 100644 index 00000000000..16e2d987e75 --- /dev/null +++ b/.changelog/27.fixed @@ -0,0 +1,5 @@ +OTLP HTTP exporters now retry on HTTP `429 Too Many Requests`, honor the +`Retry-After` response header (both delay-seconds and HTTP-date forms) when +choosing the delay before the next retry, and clamp the exponential backoff to +a maximum interval. The gRPC exporter backoff is now clamped to the same +maximum. diff --git a/exporter/opentelemetry-exporter-otlp-proto-grpc/src/opentelemetry/exporter/otlp/proto/grpc/exporter.py b/exporter/opentelemetry-exporter-otlp-proto-grpc/src/opentelemetry/exporter/otlp/proto/grpc/exporter.py index f7fa0b8697d..50670e1405a 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-grpc/src/opentelemetry/exporter/otlp/proto/grpc/exporter.py +++ b/exporter/opentelemetry-exporter-otlp-proto-grpc/src/opentelemetry/exporter/otlp/proto/grpc/exporter.py @@ -119,6 +119,10 @@ ] ) _MAX_RETRYS = 6 +# Upper bound (in seconds) applied to the exponential backoff between retries so +# it cannot grow without limit. This mirrors the behavior of the Go and Java +# OTLP exporters, which both cap the backoff interval. +_MAX_BACKOFF = 32 logger = getLogger(__name__) # This prevents logs generated when a log fails to be written to generate another log which fails to be written etc. etc. logger.addFilter(DuplicateFilter()) @@ -473,7 +477,10 @@ def _export( "google.rpc.retryinfo-bin" # type: ignore [reportArgumentType] ) # multiplying by a random number between .8 and 1.2 introduces a +/20% jitter to each backoff. - backoff_seconds = 2**retry_num * random.uniform(0.8, 1.2) + # The backoff is clamped to _MAX_BACKOFF so it cannot grow without bound. + backoff_seconds = min( + 2**retry_num * random.uniform(0.8, 1.2), _MAX_BACKOFF + ) if retry_info_bin is not None: retry_info = RetryInfo() retry_info.ParseFromString(retry_info_bin) diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_common/__init__.py b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_common/__init__.py index 57bd7ca065a..b2626e0621d 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_common/__init__.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_common/__init__.py @@ -1,6 +1,8 @@ # Copyright The OpenTelemetry Authors # SPDX-License-Identifier: Apache-2.0 +from datetime import datetime, timezone +from email.utils import parsedate_to_datetime from os import environ from typing import Literal @@ -11,15 +13,56 @@ ) from opentelemetry.util._importlib_metadata import entry_points +# Upper bound (in seconds) applied to the exponential backoff between retries so +# it cannot grow without limit. This mirrors the behavior of the Go and Java +# OTLP exporters, which both cap the backoff interval. +_MAX_BACKOFF = 32 + def _is_retryable(resp: requests.Response) -> bool: if resp.status_code == 408: return True + if resp.status_code == 429: + return True if resp.status_code >= 500 and resp.status_code <= 599: return True return False +def _parse_retry_after_header(resp: requests.Response) -> float | None: + """Return the ``Retry-After`` delay in seconds, or ``None`` if absent/invalid. + + The ``Retry-After`` header can be either an integer number of seconds + (delay-seconds form) or an HTTP-date. Both forms are defined by RFC 9110 + and the OpenTelemetry OTLP specification requires that the exporter honor + the value when present on a retryable response (e.g. 429 or 503). + """ + retry_after = resp.headers.get("Retry-After") + if retry_after is None: + return None + retry_after = retry_after.strip() + if not retry_after: + return None + try: + # delay-seconds form, e.g. "Retry-After: 120" + return float(int(retry_after)) + except ValueError: + pass + try: + # HTTP-date form, e.g. "Retry-After: Wed, 21 Oct 2015 07:28:00 GMT" + retry_date = parsedate_to_datetime(retry_after) + except (TypeError, ValueError): + return None + if retry_date is None: + return None + if retry_date.tzinfo is None: + retry_date = retry_date.replace(tzinfo=timezone.utc) + delay = (retry_date - datetime.now(timezone.utc)).total_seconds() + if delay < 0: + return 0.0 + return delay + + def _load_session_from_envvar( cred_envvar: Literal[ "OTEL_PYTHON_EXPORTER_OTLP_HTTP_LOGS_CREDENTIAL_PROVIDER", diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_log_exporter/__init__.py b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_log_exporter/__init__.py index a133874fcfa..ad39c062de3 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_log_exporter/__init__.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_log_exporter/__init__.py @@ -25,8 +25,10 @@ Compression, ) from opentelemetry.exporter.otlp.proto.http._common import ( + _MAX_BACKOFF, _is_retryable, _load_session_from_envvar, + _parse_retry_after_header, ) from opentelemetry.metrics import MeterProvider from opentelemetry.sdk._logs import ReadableLogRecord @@ -203,7 +205,10 @@ def export( deadline_sec = time() + self._timeout for retry_num in range(_MAX_RETRYS): # multiplying by a random number between .8 and 1.2 introduces a +/20% jitter to each backoff. - backoff_seconds = 2**retry_num * random.uniform(0.8, 1.2) + # The backoff is clamped to _MAX_BACKOFF so it cannot grow without bound. + backoff_seconds = min( + 2**retry_num * random.uniform(0.8, 1.2), _MAX_BACKOFF + ) export_error: Exception | None = None try: resp = self._export(serialized_data, deadline_sec - time()) @@ -218,6 +223,11 @@ def export( reason = resp.reason retryable = _is_retryable(resp) status_code = resp.status_code + # Honor a Retry-After header when present, overriding the + # computed backoff with the server-requested delay. + retry_after_seconds = _parse_retry_after_header(resp) + if retry_after_seconds is not None: + backoff_seconds = retry_after_seconds if not retryable: _logger.error( diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/metric_exporter/__init__.py b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/metric_exporter/__init__.py index 6cf13deae17..3a1c088b5af 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/metric_exporter/__init__.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/metric_exporter/__init__.py @@ -39,8 +39,10 @@ Compression, ) from opentelemetry.exporter.otlp.proto.http._common import ( + _MAX_BACKOFF, _is_retryable, _load_session_from_envvar, + _parse_retry_after_header, ) from opentelemetry.metrics import MeterProvider from opentelemetry.proto.collector.metrics.v1.metrics_service_pb2 import ( # noqa: F401 @@ -274,7 +276,10 @@ def _export_with_retries( deadline_sec = time() + self._timeout for retry_num in range(_MAX_RETRYS): # multiplying by a random number between .8 and 1.2 introduces a +/20% jitter to each backoff. - backoff_seconds = 2**retry_num * random.uniform(0.8, 1.2) + # The backoff is clamped to _MAX_BACKOFF so it cannot grow without bound. + backoff_seconds = min( + 2**retry_num * random.uniform(0.8, 1.2), _MAX_BACKOFF + ) export_error: Exception | None = None try: resp = self._export(serialized_data, deadline_sec - time()) @@ -289,6 +294,11 @@ def _export_with_retries( reason = resp.reason retryable = _is_retryable(resp) status_code = resp.status_code + # Honor a Retry-After header when present, overriding the + # computed backoff with the server-requested delay. + retry_after_seconds = _parse_retry_after_header(resp) + if retry_after_seconds is not None: + backoff_seconds = retry_after_seconds if not retryable: _logger.error( diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/trace_exporter/__init__.py b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/trace_exporter/__init__.py index 2be240103c0..06c55c0c86f 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/trace_exporter/__init__.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/trace_exporter/__init__.py @@ -27,8 +27,10 @@ Compression, ) from opentelemetry.exporter.otlp.proto.http._common import ( + _MAX_BACKOFF, _is_retryable, _load_session_from_envvar, + _parse_retry_after_header, ) from opentelemetry.metrics import MeterProvider from opentelemetry.sdk.environment_variables import ( @@ -196,7 +198,10 @@ def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: deadline_sec = time() + self._timeout for retry_num in range(_MAX_RETRYS): # multiplying by a random number between .8 and 1.2 introduces a +/20% jitter to each backoff. - backoff_seconds = 2**retry_num * random.uniform(0.8, 1.2) + # The backoff is clamped to _MAX_BACKOFF so it cannot grow without bound. + backoff_seconds = min( + 2**retry_num * random.uniform(0.8, 1.2), _MAX_BACKOFF + ) export_error: Exception | None = None try: resp = self._export(serialized_data, deadline_sec - time()) @@ -211,6 +216,11 @@ def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: reason = resp.reason retryable = _is_retryable(resp) status_code = resp.status_code + # Honor a Retry-After header when present, overriding the + # computed backoff with the server-requested delay. + retry_after_seconds = _parse_retry_after_header(resp) + if retry_after_seconds is not None: + backoff_seconds = retry_after_seconds if not retryable: _logger.error( diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/tests/metrics/test_otlp_metrics_exporter.py b/exporter/opentelemetry-exporter-otlp-proto-http/tests/metrics/test_otlp_metrics_exporter.py index eb9839d07c9..ed90cecd3b2 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/tests/metrics/test_otlp_metrics_exporter.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/tests/metrics/test_otlp_metrics_exporter.py @@ -1351,6 +1351,42 @@ def test_preferred_aggregation_override(self): exporter._preferred_aggregation[Histogram], histogram_aggregation ) + @patch.object(Session, "post") + def test_429_is_retryable(self, mock_post): + exporter = OTLPMetricExporter(timeout=1.5) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + self.assertEqual( + exporter.export(self.metrics["sum_int"]), + MetricExportResult.FAILURE, + ) + self.assertGreater(mock_post.call_count, 1) + self.assertIn( + "Transient error TOO_MANY_REQUESTS encountered while " + "exporting metrics batch, retrying in", + warning.records[0].message, + ) + + @patch.object(Session, "post") + def test_retry_after_header_seconds_is_honored(self, mock_post): + exporter = OTLPMetricExporter(timeout=10) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + resp.headers["Retry-After"] = "2" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + exporter.export(self.metrics["sum_int"]) + self.assertIn( + "retrying in 2.00s", + warning.records[0].message, + ) + @patch.dict( "os.environ", {OTEL_PYTHON_SDK_INTERNAL_METRICS_ENABLED: "true"} ) diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_log_exporter.py b/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_log_exporter.py index d7f3592e288..971b1c7a47e 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_log_exporter.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_log_exporter.py @@ -529,6 +529,42 @@ def test_retry_timeout(self, mock_post): 503, ) + @patch.object(Session, "post") + def test_429_is_retryable(self, mock_post): + exporter = OTLPLogExporter(timeout=1.5) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + self.assertEqual( + exporter.export(self._get_sdk_log_data()), + LogRecordExportResult.FAILURE, + ) + self.assertGreater(mock_post.call_count, 1) + self.assertIn( + "Transient error TOO_MANY_REQUESTS encountered while " + "exporting logs batch, retrying in", + warning.records[0].message, + ) + + @patch.object(Session, "post") + def test_retry_after_header_seconds_is_honored(self, mock_post): + exporter = OTLPLogExporter(timeout=10) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + resp.headers["Retry-After"] = "2" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + exporter.export(self._get_sdk_log_data()) + self.assertIn( + "retrying in 2.00s", + warning.records[0].message, + ) + @patch.object(Session, "post") def test_export_no_collector_available_retryable(self, mock_post): exporter = OTLPLogExporter(timeout=1.5) diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_span_exporter.py b/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_span_exporter.py index 1580e5a1802..1d67bde91a1 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_span_exporter.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_span_exporter.py @@ -4,6 +4,7 @@ import threading import time import unittest +from datetime import datetime, timedelta, timezone from logging import WARNING from unittest.mock import MagicMock, Mock, patch @@ -13,6 +14,11 @@ from requests.models import Response from opentelemetry.exporter.otlp.proto.http import Compression +from opentelemetry.exporter.otlp.proto.http._common import ( + _MAX_BACKOFF, + _is_retryable, + _parse_retry_after_header, +) from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( DEFAULT_COMPRESSION, DEFAULT_ENDPOINT, @@ -486,6 +492,94 @@ def test_shutdown_interrupts_retry_backoff(self, mock_post): assert after - before < 0.2 + @patch.object(Session, "post") + def test_429_is_retryable(self, mock_post): + exporter = OTLPSpanExporter(timeout=1.5) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + self.assertEqual( + exporter.export([BASIC_SPAN]), + SpanExportResult.FAILURE, + ) + # A 429 must be retried, so more than a single POST is expected. + self.assertGreater(mock_post.call_count, 1) + self.assertIn( + "Transient error TOO_MANY_REQUESTS encountered while " + "exporting span batch, retrying in", + warning.records[0].message, + ) + + @patch.object(Session, "post") + def test_retry_after_header_seconds_is_honored(self, mock_post): + exporter = OTLPSpanExporter(timeout=10) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + resp.headers["Retry-After"] = "2" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + exporter.export([BASIC_SPAN]) + self.assertIn( + "retrying in 2.00s", + warning.records[0].message, + ) + + @patch.object(Session, "post") + def test_retry_after_header_http_date_is_honored(self, mock_post): + exporter = OTLPSpanExporter(timeout=100) + + resp = Response() + resp.status_code = 503 + resp.reason = "UNAVAILABLE" + retry_at = datetime.now(timezone.utc) + timedelta(seconds=3) + resp.headers["Retry-After"] = retry_at.strftime( + "%a, %d %b %Y %H:%M:%S GMT" + ) + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + exporter.export([BASIC_SPAN]) + message = warning.records[0].message + self.assertIn("retrying in", message) + reported = float( + message.split("retrying in ")[1].rstrip("s.").rstrip("s") + ) + # The HTTP-date is ~3 seconds in the future. + self.assertTrue(1.5 < reported <= 3.0) + + @patch("opentelemetry.exporter.otlp.proto.http.trace_exporter.random") + @patch.object(Session, "post") + def test_backoff_is_clamped_to_max(self, mock_post, mock_random): + # No jitter, so the raw backoff is a clean power of two. + mock_random.uniform.return_value = 1.0 + + resp = Response() + resp.status_code = 503 + resp.reason = "UNAVAILABLE" + mock_post.return_value = resp + exporter = OTLPSpanExporter(timeout=1000) + with self.assertLogs(level=WARNING) as warning: + exporter.export([BASIC_SPAN]) + reported_backoffs = [] + for record in warning.records: + if "retrying in" in record.message: + reported_backoffs.append( + float( + record.message.split("retrying in ")[1] + .rstrip("s.") + .rstrip("s") + ) + ) + # Uncapped, retry #5 would ask for 2**5 == 32 and #6 would be + # larger; every reported backoff must be clamped at _MAX_BACKOFF. + self.assertTrue(reported_backoffs) + for backoff in reported_backoffs: + self.assertLessEqual(backoff, _MAX_BACKOFF) + def assert_standard_metric_attrs(self, attributes): self.assertEqual( attributes["otel.component.type"], "otlp_http_span_exporter" @@ -497,3 +591,55 @@ def assert_standard_metric_attrs(self, attributes): ) self.assertEqual(attributes["server.address"], "localhost") self.assertEqual(attributes["server.port"], 4318) + + +class TestCommonRetryHelpers(unittest.TestCase): + def _response(self, status_code, headers=None): + resp = Response() + resp.status_code = status_code + if headers: + resp.headers.update(headers) + return resp + + def test_is_retryable_status_codes(self): + for code in (408, 429, 500, 502, 503, 504): + self.assertTrue(_is_retryable(self._response(code)), code) + for code in (200, 400, 401, 403, 404): + self.assertFalse(_is_retryable(self._response(code)), code) + + def test_parse_retry_after_missing(self): + self.assertIsNone(_parse_retry_after_header(self._response(429))) + + def test_parse_retry_after_seconds(self): + self.assertEqual( + _parse_retry_after_header( + self._response(429, {"Retry-After": "120"}) + ), + 120.0, + ) + + def test_parse_retry_after_http_date(self): + retry_at = datetime.now(timezone.utc) + timedelta(seconds=30) + header = retry_at.strftime("%a, %d %b %Y %H:%M:%S GMT") + delay = _parse_retry_after_header( + self._response(503, {"Retry-After": header}) + ) + self.assertIsNotNone(delay) + self.assertTrue(25 < delay <= 30) + + def test_parse_retry_after_http_date_in_past(self): + retry_at = datetime.now(timezone.utc) - timedelta(seconds=30) + header = retry_at.strftime("%a, %d %b %Y %H:%M:%S GMT") + self.assertEqual( + _parse_retry_after_header( + self._response(503, {"Retry-After": header}) + ), + 0.0, + ) + + def test_parse_retry_after_invalid(self): + self.assertIsNone( + _parse_retry_after_header( + self._response(429, {"Retry-After": "not-a-date"}) + ) + )