nemo_gym.telemetry.connection_pool

View as Markdown

aiohttp connection-pool capacity diagnostics and queue-wait metrics.

Module Contents

Classes

NameDescription
ConnectionPoolCapacityPer-process connector limits and demand after division across workers.
ConnectionPoolConfigConnector budgets and optional demand estimates for one server or rollout CLI.
QueueTimedTCPConnectorTCPConnector that measures how long connections wait for a free pool slot.

Functions

NameDescription
_connect_count_snapshot-
_connector_queue_constraintClassify the binding connector limit when aiohttp reports a queue wait.
_display_limit-
_effective_per_host_limit-
_ephemeral_port_capacityReturn Linux’s approximate per-destination ephemeral-port budget when available.
build_connection_pool_connectorBuild the timed connector only while this process exports metrics.
connection_pool_capacityCalculate aiohttp limits for one process, preserving explicit unlimited values.
report_connection_pool_capacityReport one server/CLI process group’s pool sizing and unsafe capacity mismatches.
reset_server_nameRestore the caller’s destination label.
set_server_nameSet the bounded destination label while one logical request and its retries run.

Data

_CONNECT_COUNTS

_CONNECT_COUNTS_LOCK

_REPORTED_CAPACITIES

_SERVER_NAME

logger

API

class nemo_gym.telemetry.connection_pool.ConnectionPoolCapacity()

Bases: NamedTuple

Per-process connector limits and demand after division across workers.

workers is the size of one server’s FastAPI process group, or one for the CLI. Positive total and per_host limits are rounded down; zero means explicitly unlimited. per_host is the divided configured value, before a finite total constrains effective per-host capacity. intended and intended_per_host are expected demand rounded up per worker, or None when the corresponding estimate is not configured.

intended
Optional[int]
intended_per_host
Optional[int]
per_host
int
total
int
workers
int
class nemo_gym.telemetry.connection_pool.ConnectionPoolConfig()
Protocol

Connector budgets and optional demand estimates for one server or rollout CLI.

Configured limits apply to one server process group before division across its FastAPI workers, not the whole deployment. An explicit limit of zero is unlimited. Intended concurrency is optional expected outbound demand for that same group; None disables the corresponding sizing check.

global_aiohttp_connector_limit
int
global_aiohttp_connector_limit_per_host
int
global_aiohttp_intended_concurrency
Optional[int]
global_aiohttp_intended_concurrency_per_host
Optional[int]
class nemo_gym.telemetry.connection_pool.QueueTimedTCPConnector(
kwargs: typing.Any = {}
)

Bases: TCPConnector

TCPConnector that measures how long connections wait for a free pool slot.

Every connect() call increments connect_total for the current destination. Each queued acquisition records one queue_duration_ms sample. The sample is labelled with the limit that was binding when the wait started. The wait is timed by overriding aiohttp’s private _wait_for_available_connection. test_queue_wait_override_matches_aiohttp_signature fails if aiohttp changes that method.

nemo_gym.telemetry.connection_pool.QueueTimedTCPConnector._wait_for_available_connection(
key: aiohttp.client_reqrep.ConnectionKey,
traces: list[aiohttp.tracing.Trace]
) -> None
async
nemo_gym.telemetry.connection_pool.QueueTimedTCPConnector.connect(
req: aiohttp.client_reqrep.ClientRequest,
args: typing.Any = (),
kwargs: typing.Any = {}
) -> aiohttp.connector.Connection
async
nemo_gym.telemetry.connection_pool._connect_count_snapshot() -> dict[str, int]
nemo_gym.telemetry.connection_pool._connector_queue_constraint(
connector: typing.Any
) -> str

Classify the binding connector limit when aiohttp reports a queue wait.

nemo_gym.telemetry.connection_pool._display_limit(
limit: int
) -> str
nemo_gym.telemetry.connection_pool._effective_per_host_limit(
total: int,
per_host: int
) -> int
nemo_gym.telemetry.connection_pool._ephemeral_port_capacity() -> typing.Optional[int]

Return Linux’s approximate per-destination ephemeral-port budget when available.

nemo_gym.telemetry.connection_pool.build_connection_pool_connector(
kwargs: typing.Any = {}
) -> aiohttp.TCPConnector

Build the timed connector only while this process exports metrics.

nemo_gym.telemetry.connection_pool.connection_pool_capacity(
workers: int

Calculate aiohttp limits for one process, preserving explicit unlimited values.

nemo_gym.telemetry.connection_pool.report_connection_pool_capacity(
visible: bool = False
) -> None

Report one server/CLI process group’s pool sizing and unsafe capacity mismatches.

nemo_gym.telemetry.connection_pool.reset_server_name(
token: contextvars.Token[str]
) -> None

Restore the caller’s destination label.

nemo_gym.telemetry.connection_pool.set_server_name(
server_name: str
) -> contextvars.Token[str]

Set the bounded destination label while one logical request and its retries run.

nemo_gym.telemetry.connection_pool._CONNECT_COUNTS: dict[str, int] = {}
nemo_gym.telemetry.connection_pool._CONNECT_COUNTS_LOCK = Lock()
nemo_gym.telemetry.connection_pool._REPORTED_CAPACITIES: set[tuple[object, ...]] = set()
nemo_gym.telemetry.connection_pool._SERVER_NAME: ContextVar[str] = ContextVar('nemo_gym_http_server_name', default='external')
nemo_gym.telemetry.connection_pool.logger = logging.getLogger(__name__)