Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
http.py304 linesDownload Raw Back to _sync
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 
codekingpro/portable-devtools · Team Ai