Skip to content
Merged
411 changes: 357 additions & 54 deletions src/runloop_api_client/_base_client.py

Large diffs are not rendered by default.

37 changes: 37 additions & 0 deletions src/runloop_api_client/_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
)
from ._compat import cached_property
from ._version import __version__
from ._constants import DEFAULT_API_POOL_SHARDS, DEFAULT_TRANSFER_POOL_SHARDS, DEFAULT_BACKGROUND_POOL_SHARDS
from ._streaming import Stream as Stream, AsyncStream as AsyncStream
from ._exceptions import RunloopError, APIStatusError
from ._base_client import (
Expand Down Expand Up @@ -96,6 +97,11 @@ def __init__(
# Enables HTTP/2 multiplexing and avoids ConnectTimeout storms under high concurrency.
# Set to False to create a private connection pool (old behavior).
shared_http_pool: bool = True,
# Sharded H2 pools by workload (API / long-polls / transfers). Each shard ≈
# one H2 connection; requests round-robin per client from a random offset.
api_pool_shards: int = DEFAULT_API_POOL_SHARDS,
background_pool_shards: int = DEFAULT_BACKGROUND_POOL_SHARDS,
transfer_pool_shards: int = DEFAULT_TRANSFER_POOL_SHARDS,
# Enable or disable schema validation for data returned by the API.
# When enabled an error APIResponseValidationError is raised
# if the API responds with invalid data for the expected schema.
Expand Down Expand Up @@ -142,6 +148,9 @@ def __init__(
custom_query=default_query,
_strict_response_validation=_strict_response_validation,
shared_http_pool=shared_http_pool,
api_pool_shards=api_pool_shards,
background_pool_shards=background_pool_shards,
transfer_pool_shards=transfer_pool_shards,
)

self._idempotency_header = "x-request-id"
Expand Down Expand Up @@ -284,6 +293,9 @@ def copy(
timeout: float | Timeout | None | NotGiven = not_given,
http_client: httpx.Client | None = None,
shared_http_pool: bool | None = None,
api_pool_shards: int | None = None,
background_pool_shards: int | None = None,
transfer_pool_shards: int | None = None,
max_retries: int | NotGiven = not_given,
default_headers: Mapping[str, str] | None = None,
set_default_headers: Mapping[str, str] | None = None,
Expand Down Expand Up @@ -325,6 +337,13 @@ def copy(
timeout=self.timeout if isinstance(timeout, NotGiven) else timeout,
http_client=http_client,
shared_http_pool=resolved_shared,
api_pool_shards=(api_pool_shards if api_pool_shards is not None else self._api_pool_shards),
background_pool_shards=(
background_pool_shards if background_pool_shards is not None else self._background_pool_shards
),
transfer_pool_shards=(
transfer_pool_shards if transfer_pool_shards is not None else self._transfer_pool_shards
),
max_retries=max_retries if is_given(max_retries) else self.max_retries,
default_headers=headers,
default_query=params,
Expand Down Expand Up @@ -390,6 +409,11 @@ def __init__(
# Enables HTTP/2 multiplexing and avoids ConnectTimeout storms under high concurrency.
# Set to False to create a private connection pool (old behavior).
shared_http_pool: bool = True,
# Sharded H2 pools by workload (API / long-polls / transfers). Each shard ≈
# one H2 connection; requests round-robin per client from a random offset.
api_pool_shards: int = DEFAULT_API_POOL_SHARDS,
background_pool_shards: int = DEFAULT_BACKGROUND_POOL_SHARDS,
transfer_pool_shards: int = DEFAULT_TRANSFER_POOL_SHARDS,
# Enable or disable schema validation for data returned by the API.
# When enabled an error APIResponseValidationError is raised
# if the API responds with invalid data for the expected schema.
Expand Down Expand Up @@ -436,6 +460,9 @@ def __init__(
custom_query=default_query,
_strict_response_validation=_strict_response_validation,
shared_http_pool=shared_http_pool,
api_pool_shards=api_pool_shards,
background_pool_shards=background_pool_shards,
transfer_pool_shards=transfer_pool_shards,
)

self._idempotency_header = "x-request-id"
Expand Down Expand Up @@ -578,6 +605,9 @@ def copy(
timeout: float | Timeout | None | NotGiven = not_given,
http_client: httpx.AsyncClient | None = None,
shared_http_pool: bool | None = None,
api_pool_shards: int | None = None,
background_pool_shards: int | None = None,
transfer_pool_shards: int | None = None,
max_retries: int | NotGiven = not_given,
default_headers: Mapping[str, str] | None = None,
set_default_headers: Mapping[str, str] | None = None,
Expand Down Expand Up @@ -619,6 +649,13 @@ def copy(
timeout=self.timeout if isinstance(timeout, NotGiven) else timeout,
http_client=http_client,
shared_http_pool=resolved_shared,
api_pool_shards=(api_pool_shards if api_pool_shards is not None else self._api_pool_shards),
background_pool_shards=(
background_pool_shards if background_pool_shards is not None else self._background_pool_shards
),
transfer_pool_shards=(
transfer_pool_shards if transfer_pool_shards is not None else self._transfer_pool_shards
),
max_retries=max_retries if is_given(max_retries) else self.max_retries,
default_headers=headers,
default_query=params,
Expand Down
7 changes: 7 additions & 0 deletions src/runloop_api_client/_constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,13 @@
DEFAULT_MAX_RETRIES = 5
DEFAULT_CONNECTION_LIMITS = httpx.Limits(max_connections=20, max_keepalive_connections=10)

# Separate H2 connection pools by workload. Each "shard" is its own httpx
# client / transport (≈ one H2 connection). Defaults sized for Jetty's ~128
# streams/connection: enough API/wait concurrency without oversized transfer.
DEFAULT_API_POOL_SHARDS = 8
DEFAULT_BACKGROUND_POOL_SHARDS = 16
DEFAULT_TRANSFER_POOL_SHARDS = 2

INITIAL_RETRY_DELAY = 1.0
MAX_RETRY_DELAY = 60.0

Expand Down
13 changes: 13 additions & 0 deletions src/runloop_api_client/sdk/async_.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
from .._client import DEFAULT_MAX_RETRIES, AsyncRunloop
from ._helpers import detect_content_type
from .async_axon import AsyncAxon
from .._constants import DEFAULT_API_POOL_SHARDS, DEFAULT_TRANSFER_POOL_SHARDS, DEFAULT_BACKGROUND_POOL_SHARDS
from .async_agent import AsyncAgent
from .async_devbox import AsyncDevbox
from .async_scorer import AsyncScorer
Expand Down Expand Up @@ -1329,6 +1330,9 @@ def __init__(
default_headers: Mapping[str, str] | None = None,
default_query: Mapping[str, object] | None = None,
http_client: httpx.AsyncClient | None = None,
api_pool_shards: int = DEFAULT_API_POOL_SHARDS,
background_pool_shards: int = DEFAULT_BACKGROUND_POOL_SHARDS,
transfer_pool_shards: int = DEFAULT_TRANSFER_POOL_SHARDS,
) -> None:
"""Configure the asynchronous SDK wrapper.

Expand All @@ -1346,6 +1350,12 @@ def __init__(
:type default_query: Mapping[str, object] | None, optional
:param http_client: Custom ``httpx.AsyncClient`` instance to reuse, defaults to None
:type http_client: httpx.AsyncClient | None, optional
:param api_pool_shards: H2 shards for short RPCs (round-robin), defaults to 8
:type api_pool_shards: int, optional
:param background_pool_shards: H2 shards for long-polls (round-robin), defaults to 16
:type background_pool_shards: int, optional
:param transfer_pool_shards: H2 shards for upload/download (round-robin), defaults to 2
:type transfer_pool_shards: int, optional
"""
self.api = AsyncRunloop(
bearer_token=bearer_token,
Expand All @@ -1355,6 +1365,9 @@ def __init__(
default_headers=default_headers,
default_query=default_query,
http_client=http_client,
api_pool_shards=api_pool_shards,
background_pool_shards=background_pool_shards,
transfer_pool_shards=transfer_pool_shards,
)

self.agent = AsyncAgentOps(self.api)
Expand Down
13 changes: 13 additions & 0 deletions src/runloop_api_client/sdk/sync.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@
from .benchmark import Benchmark
from .blueprint import Blueprint
from .mcp_config import McpConfig
from .._constants import DEFAULT_API_POOL_SHARDS, DEFAULT_TRANSFER_POOL_SHARDS, DEFAULT_BACKGROUND_POOL_SHARDS
from .gateway_config import GatewayConfig
from .network_policy import NetworkPolicy
from .storage_object import StorageObject
Expand Down Expand Up @@ -1354,6 +1355,9 @@ def __init__(
default_headers: Mapping[str, str] | None = None,
default_query: Mapping[str, object] | None = None,
http_client: httpx.Client | None = None,
api_pool_shards: int = DEFAULT_API_POOL_SHARDS,
background_pool_shards: int = DEFAULT_BACKGROUND_POOL_SHARDS,
transfer_pool_shards: int = DEFAULT_TRANSFER_POOL_SHARDS,
) -> None:
"""Configure the synchronous SDK wrapper.

Expand All @@ -1371,6 +1375,12 @@ def __init__(
:type default_query: Mapping[str, object] | None, optional
:param http_client: Custom ``httpx.Client`` instance to reuse, defaults to None
:type http_client: httpx.Client | None, optional
:param api_pool_shards: H2 shards for short RPCs (round-robin), defaults to 8
:type api_pool_shards: int, optional
:param background_pool_shards: H2 shards for long-polls (round-robin), defaults to 16
:type background_pool_shards: int, optional
:param transfer_pool_shards: H2 shards for upload/download (round-robin), defaults to 2
:type transfer_pool_shards: int, optional
"""
self.api = Runloop(
bearer_token=bearer_token,
Expand All @@ -1380,6 +1390,9 @@ def __init__(
default_headers=default_headers,
default_query=default_query,
http_client=http_client,
api_pool_shards=api_pool_shards,
background_pool_shards=background_pool_shards,
transfer_pool_shards=transfer_pool_shards,
)

self.agent = AgentOps(self.api)
Expand Down
Loading