codekingpro/portable-devtools
114k
1"""Synchronous HTTP client for LangGraph API."""2 3from __future__ import annotations4 5import logging6import sys7import warnings8from collections.abc import Callable, Iterator, Mapping9from typing import Any, cast10 11import httpx12import orjson13 14from langgraph_sdk._shared.utilities import (15 _orjson_default,16 _validate_reconnect_location,17)18from langgraph_sdk.errors import _raise_for_status_typed19from langgraph_sdk.schema import QueryParamTypes, StreamPart20from langgraph_sdk.sse import SSEDecoder, iter_lines_raw21 22logger = logging.getLogger(__name__)23 24 25class SyncHttpClient:26 """Handle synchronous requests to the LangGraph API.27 28 Provides error messaging and content handling enhancements above the29 underlying httpx client, mirroring the interface of [HttpClient](#HttpClient)30 but for sync usage.31 32 Attributes:33 client (httpx.Client): Underlying HTTPX sync client.34 """35 36 def __init__(self, client: httpx.Client) -> None:37 self.client = client38 39 def get(40 self,41 path: str,42 *,43 params: QueryParamTypes | None = None,44 headers: Mapping[str, str] | None = None,45 on_response: Callable[[httpx.Response], None] | None = None,46 ) -> Any:47 """Send a `GET` request."""48 r = self.client.get(path, params=params, headers=headers)49 if on_response:50 on_response(r)51 _raise_for_status_typed(r)52 return _decode_json(r)53 54 def post(55 self,56 path: str,57 *,58 json: dict[str, Any] | list | None,59 params: QueryParamTypes | None = None,60 headers: Mapping[str, str] | None = None,61 on_response: Callable[[httpx.Response], None] | None = None,62 ) -> Any:63 """Send a `POST` request."""64 if json is not None:65 request_headers, content = _encode_json(json)66 else:67 request_headers, content = {}, b""68 if headers:69 request_headers.update(headers)70 r = self.client.post(71 path, headers=request_headers, content=content, params=params72 )73 if on_response:74 on_response(r)75 _raise_for_status_typed(r)76 return _decode_json(r)77 78 def put(79 self,80 path: str,81 *,82 json: dict,83 params: QueryParamTypes | None = None,84 headers: Mapping[str, str] | None = None,85 on_response: Callable[[httpx.Response], None] | None = None,86 ) -> Any:87 """Send a `PUT` request."""88 request_headers, content = _encode_json(json)89 if headers:90 request_headers.update(headers)91 92 r = self.client.put(93 path, headers=request_headers, content=content, params=params94 )95 if on_response:96 on_response(r)97 _raise_for_status_typed(r)98 return _decode_json(r)99 100 def patch(101 self,102 path: str,103 *,104 json: dict,105 params: QueryParamTypes | None = None,106 headers: Mapping[str, str] | None = None,107 on_response: Callable[[httpx.Response], None] | None = None,108 ) -> Any:109 """Send a `PATCH` request."""110 request_headers, content = _encode_json(json)111 if headers:112 request_headers.update(headers)113 r = self.client.patch(114 path, headers=request_headers, content=content, params=params115 )116 if on_response:117 on_response(r)118 _raise_for_status_typed(r)119 return _decode_json(r)120 121 def delete(122 self,123 path: str,124 *,125 json: Any | None = None,126 params: QueryParamTypes | None = None,127 headers: Mapping[str, str] | None = None,128 on_response: Callable[[httpx.Response], None] | None = None,129 ) -> None:130 """Send a `DELETE` request."""131 r = self.client.request(132 "DELETE", path, json=json, params=params, headers=headers133 )134 if on_response:135 on_response(r)136 _raise_for_status_typed(r)137 138 def request_reconnect(139 self,140 path: str,141 method: str,142 *,143 json: dict[str, Any] | None = None,144 params: QueryParamTypes | None = None,145 headers: Mapping[str, str] | None = None,146 on_response: Callable[[httpx.Response], None] | None = None,147 reconnect_limit: int = 5,148 ) -> Any:149 """Send a request that automatically reconnects to Location header."""150 request_headers, content = _encode_json(json)151 if headers:152 request_headers.update(headers)153 with self.client.stream(154 method, path, headers=request_headers, content=content, params=params155 ) as r:156 if on_response:157 on_response(r)158 try:159 r.raise_for_status()160 except httpx.HTTPStatusError as e:161 body = r.read().decode()162 if sys.version_info >= (3, 11):163 e.add_note(body)164 else:165 logger.error(f"Error from langgraph-api: {body}", exc_info=e)166 raise e167 loc = r.headers.get("location")168 if reconnect_limit <= 0 or not loc:169 return _decode_json(r)170 _validate_reconnect_location(self.client.base_url, loc)171 try:172 return _decode_json(r)173 except httpx.HTTPError:174 warnings.warn(175 f"Request failed, attempting reconnect to Location: {loc}",176 stacklevel=2,177 )178 r.close()179 return self.request_reconnect(180 loc,181 "GET",182 headers=request_headers,183 # don't pass on_response so it's only called once184 reconnect_limit=reconnect_limit - 1,185 )186 187 def stream(188 self,189 path: str,190 method: str,191 *,192 json: dict[str, Any] | None = None,193 params: QueryParamTypes | None = None,194 headers: Mapping[str, str] | None = None,195 on_response: Callable[[httpx.Response], None] | None = None,196 ) -> Iterator[StreamPart]:197 """Stream the results of a request using SSE."""198 if json is not None:199 request_headers, content = _encode_json(json)200 else:201 request_headers, content = {}, None202 request_headers["Accept"] = "text/event-stream"203 request_headers["Cache-Control"] = "no-store"204 if headers:205 request_headers.update(headers)206 207 reconnect_headers = {208 key: value209 for key, value in request_headers.items()210 if key.lower() not in {"content-length", "content-type"}211 }212 213 last_event_id: str | None = None214 reconnect_path: str | None = None215 reconnect_attempts = 0216 max_reconnect_attempts = 5217 218 while True:219 current_headers = dict(220 request_headers if reconnect_path is None else reconnect_headers221 )222 if last_event_id is not None:223 current_headers["Last-Event-ID"] = last_event_id224 225 current_method = method if reconnect_path is None else "GET"226 current_content = content if reconnect_path is None else None227 current_params = params if reconnect_path is None else None228 229 retry = False230 with self.client.stream(231 current_method,232 reconnect_path or path,233 headers=current_headers,234 content=current_content,235 params=current_params,236 ) as res:237 if reconnect_path is None and on_response:238 on_response(res)239 # check status240 _raise_for_status_typed(res)241 # check content type242 content_type = res.headers.get("content-type", "").partition(";")[0]243 if "text/event-stream" not in content_type:244 raise httpx.TransportError(245 "Expected response header Content-Type to contain 'text/event-stream', "246 f"got {content_type!r}"247 )248 249 reconnect_location = res.headers.get("location")250 if reconnect_location:251 _validate_reconnect_location(252 self.client.base_url, reconnect_location253 )254 reconnect_path = reconnect_location255 256 decoder = SSEDecoder()257 try:258 for line in iter_lines_raw(res):259 sse = decoder.decode(cast(bytes, line).rstrip(b"\n"))260 if sse is not None:261 if decoder.last_event_id is not None:262 last_event_id = decoder.last_event_id263 if sse.event or sse.data is not None:264 yield sse265 except httpx.HTTPError:266 # httpx.TransportError inherits from HTTPError, so transient267 # disconnects during streaming land here.268 if reconnect_path is None:269 raise270 retry = True271 else:272 if sse := decoder.decode(b""):273 if decoder.last_event_id is not None:274 last_event_id = decoder.last_event_id275 if sse.event or sse.data is not None:276 # See async stream implementation for rationale on277 # skipping empty flush events.278 yield sse279 if retry:280 reconnect_attempts += 1281 if reconnect_attempts > max_reconnect_attempts:282 raise httpx.TransportError(283 "Exceeded maximum SSE reconnection attempts"284 )285 continue286 break287 288 289def _encode_json(json: Any) -> tuple[dict[str, str], bytes]:290 body = orjson.dumps(291 json,292 _orjson_default,293 orjson.OPT_SERIALIZE_NUMPY | orjson.OPT_NON_STR_KEYS,294 )295 content_length = str(len(body))296 content_type = "application/json"297 headers = {"Content-Length": content_length, "Content-Type": content_type}298 return headers, body299 300 301def _decode_json(r: httpx.Response) -> Any:302 body = r.read()303 return orjson.loads(body) if body else None304 