Skip to content

Commit 44e6b62

Browse files
author
Reflex
committed
fix(transport): normalize lifecycle incident failures
1 parent e8c3d90 commit 44e6b62

7 files changed

Lines changed: 587 additions & 20 deletions

File tree

src/runloop_api_client/_base_client.py

Lines changed: 44 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -94,6 +94,7 @@
9494
APIResponseValidationError,
9595
)
9696
from ._utils._json import openapi_dumps
97+
from .lib.error_contract import is_safe_transport_retry
9798

9899
log: logging.Logger = logging.getLogger(__name__)
99100

@@ -665,7 +666,10 @@ def _enforce_trailing_slash(self, url: URL) -> URL:
665666
def _make_status_error_from_response(
666667
self,
667668
response: httpx.Response,
669+
*,
670+
attempts: int = 1,
668671
) -> APIStatusError:
672+
body: object | None
669673
if response.is_closed and not response.is_stream_consumed:
670674
# We can't read the response body as it has been closed
671675
# before it was read. This can happen if an event hook
@@ -677,12 +681,19 @@ def _make_status_error_from_response(
677681
body = err_text
678682

679683
try:
680-
body = json.loads(err_text)
681-
err_msg = f"Error code: {response.status_code} - {body}"
684+
body = cast(object, json.loads(err_text))
685+
body_mapping = cast(Mapping[str, object], body) if isinstance(body, dict) else None
686+
body_message = body_mapping.get("message") if body_mapping is not None else None
687+
if isinstance(body_message, str):
688+
err_msg = body_message
689+
else:
690+
err_msg = f"Error code: {response.status_code} - {body}"
682691
except Exception:
683692
err_msg = err_text or f"Error code: {response.status_code}"
684693

685-
return self._make_status_error(err_msg, body=body, response=response)
694+
error = self._make_status_error(err_msg, body=cast(object, body), response=response)
695+
error.attempts = attempts
696+
return error
686697

687698
def _make_status_error(
688699
self,
@@ -1045,6 +1056,26 @@ def _should_retry(self, response: httpx.Response) -> bool:
10451056
log.debug("Not retrying as header `x-should-retry` is set to `false`")
10461057
return False
10471058

1059+
# These failures happen after a request or response may have been
1060+
# partially transferred. Retrying them implicitly can duplicate an
1061+
# execute or replay a multipart stream. Servers may explicitly opt in
1062+
# with X-Should-Retry when an idempotency record makes that safe.
1063+
try:
1064+
raw_payload = response.json()
1065+
except Exception:
1066+
raw_payload = None
1067+
payload = cast(Mapping[str, object], raw_payload) if isinstance(raw_payload, dict) else None
1068+
code = response.headers.get("x-runloop-error-code")
1069+
if code is None and isinstance(payload, dict) and isinstance(payload.get("error"), str):
1070+
code = payload["error"]
1071+
if code in {
1072+
"upload_request_body_idle_timeout",
1073+
"tunnel_backend_idle_timeout",
1074+
"tunnel_backend_connection_reset",
1075+
}:
1076+
log.debug("Not retrying ambiguous transfer failure %s", code)
1077+
return False
1078+
10481079
# Retry on request timeouts.
10491080
if response.status_code == 408:
10501081
log.debug("Retrying due to status code %i", response.status_code)
@@ -1366,7 +1397,7 @@ def request(
13661397
except httpx.TimeoutException as err:
13671398
log.debug("Encountered httpx.TimeoutException", exc_info=True)
13681399

1369-
if remaining_retries > 0:
1400+
if remaining_retries > 0 and is_safe_transport_retry(err):
13701401
self._sleep_for_retry(
13711402
retries_taken=retries_taken,
13721403
max_retries=max_retries,
@@ -1377,11 +1408,11 @@ def request(
13771408
continue
13781409

13791410
log.debug("Raising timeout error")
1380-
raise APITimeoutError(request=request) from err
1411+
raise APITimeoutError(request=request, cause=err, attempts=retries_taken + 1) from err
13811412
except Exception as err:
13821413
log.debug("Encountered Exception", exc_info=True)
13831414

1384-
if remaining_retries > 0:
1415+
if remaining_retries > 0 and is_safe_transport_retry(err):
13851416
self._sleep_for_retry(
13861417
retries_taken=retries_taken,
13871418
max_retries=max_retries,
@@ -1392,7 +1423,7 @@ def request(
13921423
continue
13931424

13941425
log.debug("Raising connection error")
1395-
raise APIConnectionError(request=request) from err
1426+
raise APIConnectionError(request=request, cause=err, attempts=retries_taken + 1) from err
13961427

13971428
log.debug(
13981429
'HTTP Response: %s %s "%i %s" %s',
@@ -1424,7 +1455,7 @@ def request(
14241455
err.response.read()
14251456

14261457
log.debug("Re-raising status error")
1427-
raise self._make_status_error_from_response(err.response) from None
1458+
raise self._make_status_error_from_response(err.response, attempts=retries_taken + 1) from None
14281459

14291460
break
14301461

@@ -2076,7 +2107,7 @@ async def request(
20762107
except httpx.TimeoutException as err:
20772108
log.debug("Encountered httpx.TimeoutException", exc_info=True)
20782109

2079-
if remaining_retries > 0:
2110+
if remaining_retries > 0 and is_safe_transport_retry(err):
20802111
await self._sleep_for_retry(
20812112
retries_taken=retries_taken,
20822113
max_retries=max_retries,
@@ -2087,11 +2118,11 @@ async def request(
20872118
continue
20882119

20892120
log.debug("Raising timeout error")
2090-
raise APITimeoutError(request=request) from err
2121+
raise APITimeoutError(request=request, cause=err, attempts=retries_taken + 1) from err
20912122
except Exception as err:
20922123
log.debug("Encountered Exception", exc_info=True)
20932124

2094-
if remaining_retries > 0:
2125+
if remaining_retries > 0 and is_safe_transport_retry(err):
20952126
await self._sleep_for_retry(
20962127
retries_taken=retries_taken,
20972128
max_retries=max_retries,
@@ -2102,7 +2133,7 @@ async def request(
21022133
continue
21032134

21042135
log.debug("Raising connection error")
2105-
raise APIConnectionError(request=request) from err
2136+
raise APIConnectionError(request=request, cause=err, attempts=retries_taken + 1) from err
21062137

21072138
log.debug(
21082139
'HTTP Response: %s %s "%i %s" %s',
@@ -2134,7 +2165,7 @@ async def request(
21342165
await err.response.aread()
21352166

21362167
log.debug("Re-raising status error")
2137-
raise self._make_status_error_from_response(err.response) from None
2168+
raise self._make_status_error_from_response(err.response, attempts=retries_taken + 1) from None
21382169

21392170
break
21402171

src/runloop_api_client/_exceptions.py

Lines changed: 69 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,8 @@
66

77
import httpx
88

9+
from .lib.error_contract import status_error_details, transport_error_details
10+
911
__all__ = [
1012
"BadRequestError",
1113
"AuthenticationError",
@@ -37,11 +39,39 @@ class APIError(RunloopError):
3739
If there was no response associated with this error then it will be `None`.
3840
"""
3941

40-
def __init__(self, message: str, request: httpx.Request, *, body: object | None) -> None: # noqa: ARG002
42+
code: str
43+
phase: str
44+
retryable: bool
45+
request_id: str | None
46+
retry_after: float | None
47+
attempts: int
48+
cause: BaseException | None
49+
50+
def __init__(
51+
self,
52+
message: str,
53+
request: httpx.Request,
54+
*,
55+
body: object | None,
56+
code: str = "runloop_error",
57+
phase: str = "unknown",
58+
retryable: bool = False,
59+
request_id: str | None = None,
60+
retry_after: float | None = None,
61+
attempts: int = 1,
62+
cause: BaseException | None = None,
63+
) -> None:
4164
super().__init__(message)
4265
self.request = request
4366
self.message = message
4467
self.body = body
68+
self.code = code
69+
self.phase = phase
70+
self.retryable = retryable
71+
self.request_id = request_id
72+
self.retry_after = retry_after
73+
self.attempts = attempts
74+
self.cause = cause
4575

4676

4777
class APIResponseValidationError(APIError):
@@ -60,20 +90,52 @@ class APIStatusError(APIError):
6090
response: httpx.Response
6191
status_code: int
6292

63-
def __init__(self, message: str, *, response: httpx.Response, body: object | None) -> None:
64-
super().__init__(message, response.request, body=body)
93+
def __init__(self, message: str, *, response: httpx.Response, body: object | None, attempts: int = 1) -> None:
94+
details = status_error_details(response, body)
95+
super().__init__(
96+
message,
97+
response.request,
98+
body=body,
99+
code=details.code,
100+
phase=details.phase,
101+
retryable=details.retryable,
102+
request_id=details.request_id,
103+
retry_after=details.retry_after,
104+
attempts=attempts,
105+
)
65106
self.response = response
66107
self.status_code = response.status_code
67108

68109

69110
class APIConnectionError(APIError):
70-
def __init__(self, *, message: str = "Connection error.", request: httpx.Request) -> None:
71-
super().__init__(message, request, body=None)
111+
def __init__(
112+
self,
113+
*,
114+
message: str = "Connection error.",
115+
request: httpx.Request,
116+
cause: BaseException | None = None,
117+
attempts: int = 1,
118+
) -> None:
119+
details = transport_error_details(cause) if cause is not None else transport_error_details(Exception())
120+
super().__init__(
121+
message,
122+
request,
123+
body=None,
124+
code=details.code,
125+
phase=details.phase,
126+
retryable=details.retryable,
127+
attempts=attempts,
128+
cause=cause,
129+
)
72130

73131

74132
class APITimeoutError(APIConnectionError):
75-
def __init__(self, request: httpx.Request) -> None:
76-
super().__init__(message="Request timed out.", request=request)
133+
def __init__(self, request: httpx.Request, *, cause: BaseException | None = None, attempts: int = 1) -> None:
134+
super().__init__(message="Request timed out.", request=request, cause=cause, attempts=attempts)
135+
if cause is None:
136+
self.code = "connection_timeout"
137+
self.phase = "connect"
138+
self.retryable = True
77139

78140

79141
class BadRequestError(APIStatusError):
Lines changed: 97 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,97 @@
1+
"""Stable normalization for Runloop API and HTTPX transport failures.
2+
3+
This module is handwritten and intentionally lives under ``lib`` so generated
4+
client updates only need a small integration point.
5+
"""
6+
7+
from __future__ import annotations
8+
9+
from typing import Mapping, cast
10+
from dataclasses import dataclass
11+
12+
import httpx
13+
14+
15+
@dataclass(frozen=True)
16+
class ErrorDetails:
17+
code: str
18+
phase: str
19+
retryable: bool
20+
request_id: str | None = None
21+
retry_after: float | None = None
22+
23+
24+
def _number(value: object) -> float | None:
25+
try:
26+
parsed = float(value) # type: ignore[arg-type]
27+
except (TypeError, ValueError):
28+
return None
29+
return parsed if parsed >= 0 else None
30+
31+
32+
def parse_retry_after(headers: httpx.Headers, body: object = None) -> float | None:
33+
"""Parse Retry-After while accepting the SDK's millisecond extension."""
34+
milliseconds = _number(headers.get("retry-after-ms"))
35+
if milliseconds is not None:
36+
return milliseconds / 1000
37+
seconds = _number(headers.get("retry-after"))
38+
if seconds is not None:
39+
return seconds
40+
if isinstance(body, Mapping):
41+
payload = cast(Mapping[str, object], body)
42+
details = payload.get("details")
43+
if isinstance(details, Mapping):
44+
return _number(cast(Mapping[str, object], details).get("retry_after"))
45+
return None
46+
47+
48+
def status_error_details(response: httpx.Response, body: object) -> ErrorDetails:
49+
payload: Mapping[str, object] = cast(Mapping[str, object], body) if isinstance(body, Mapping) else {}
50+
header_code = response.headers.get("x-runloop-error-code")
51+
body_code = payload.get("error")
52+
code = header_code or (body_code if isinstance(body_code, str) else None) or f"http_{response.status_code}"
53+
body_phase = payload.get("phase")
54+
phase = body_phase if isinstance(body_phase, str) else "api"
55+
retryable_value = payload.get("retryable")
56+
retryable = retryable_value if isinstance(retryable_value, bool) else response.status_code in {408, 409, 429}
57+
if response.status_code >= 500 and not isinstance(retryable_value, bool):
58+
retryable = True
59+
request_id: str | None = response.headers.get("x-runloop-request-id")
60+
body_request_id = payload.get("request_id")
61+
if request_id is None and isinstance(body_request_id, str):
62+
request_id = body_request_id
63+
return ErrorDetails(
64+
code=code,
65+
phase=phase,
66+
retryable=retryable,
67+
request_id=request_id,
68+
retry_after=parse_retry_after(response.headers, cast(object, body)),
69+
)
70+
71+
72+
def transport_error_details(error: BaseException) -> ErrorDetails:
73+
if isinstance(error, httpx.ConnectTimeout):
74+
return ErrorDetails("connection_timeout", "connect", True)
75+
if isinstance(error, httpx.WriteTimeout):
76+
return ErrorDetails("request_write_timeout", "request_write", False)
77+
if isinstance(error, httpx.WriteError):
78+
return ErrorDetails("request_write_failed", "request_write", False)
79+
if isinstance(error, httpx.ReadTimeout):
80+
return ErrorDetails("response_read_timeout", "response_read", False)
81+
if isinstance(error, httpx.RemoteProtocolError):
82+
if "idle_timeout" in str(error).lower():
83+
return ErrorDetails("http2_idle_timeout", "response_read", False)
84+
return ErrorDetails("http2_protocol_error", "transport", False)
85+
if isinstance(error, httpx.TimeoutException):
86+
return ErrorDetails("connection_timeout", "connect", True)
87+
return ErrorDetails("connection_failed", "connect", isinstance(error, httpx.ConnectError))
88+
89+
90+
def is_safe_transport_retry(error: BaseException) -> bool:
91+
"""Only retry failures that prove the request body was not partially sent."""
92+
# Preserve the generated client's handling of non-HTTPX exceptions (for
93+
# example a pre-send auth hook failure). HTTPX errors carry enough phase
94+
# information for the stricter partial-write audit below.
95+
if not isinstance(error, httpx.HTTPError):
96+
return True
97+
return isinstance(error, (httpx.ConnectTimeout, httpx.ConnectError))

0 commit comments

Comments
 (0)