From e948301d5961658af1f7c9050ca970665219ebe2 Mon Sep 17 00:00:00 2001 From: andrexibiza <84248988+andrexibiza@users.noreply.github.com> Date: Mon, 3 Aug 2026 02:51:29 -0500 Subject: [PATCH 1/3] refactor(gateway): extract compression-status helpers from run.py (slice 4 of #54962) Fourth slice of the gateway god-file unpacking: extract _status_template_to_regex and the _COMPRESSION_PROGRESS_STATUS_RE matcher into gateway/status_helpers.py. - Byte-identical extraction (AST-verified): the regex is built from the SAME template constants the emit sites format (agent/conversation_compression.py), so wording drift cannot silently diverge - gateway/run.py: -47 lines; the 8 now-unused template imports removed; helpers imported at the extraction point - 78 tests pass (compression-progress notices + replay + resume suites) Follow-up to #77433 (slice 1) and #77438 (slices 2-3). Progress on #54962. Signed-off-by: andrexibiza <84248988+andrexibiza@users.noreply.github.com> --- gateway/run.py | 51 +++---------------------------- gateway/status_helpers.py | 63 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 67 insertions(+), 47 deletions(-) create mode 100644 gateway/status_helpers.py diff --git a/gateway/run.py b/gateway/run.py index 8c162f7dc9a46..14d762461d817 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -47,16 +47,6 @@ from typing import Awaitable, Callable, Dict, Optional, Any, List, Tuple, Union, cast from agent.async_utils import consume_detached_task_result, safe_schedule_threadsafe -from agent.conversation_compression import ( - COMPACTION_STATUS, - COMPRESSION_RETRY_CONTEXT_REDUCED_STATUS_TEMPLATE, - COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE, - COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE, - COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE, - IDLE_COMPACTION_STATUS_TEMPLATE, - PRE_API_COMPRESSION_STATUS_TEMPLATE, - PREFLIGHT_COMPRESSION_STATUS_TEMPLATE, -) from agent.conversation_loop import INTERRUPT_WAITING_FOR_MODEL_PREFIX from agent.i18n import t from agent.interrupt_compat import request_hard_interrupt @@ -142,43 +132,10 @@ def _record_hygiene_cooldown(gateway, session_id: str, cooldown_seconds: float) logger.debug("session hygiene cooldown persist failed: %s", exc) -def _status_template_to_regex(template: str) -> str: - """Compile a compression status template constant into a regex source. - - Literal text is escaped verbatim (so wording drift in - agent/conversation_compression.py cannot silently diverge from this - matcher — the constants ARE the wording) and each ``{field}`` format - placeholder is replaced with a numeric-ish pattern covering every value - the emit sites format in (ints, ``{:,}`` thousands separators). - """ - parts = re.split(r"\{[^{}]*\}", template) - return r"[\d,]+".join(re.escape(part) for part in parts) - - -# ROUTINE compression progress statuses, derived from the SAME template -# constants the emit sites format (agent/conversation_compression.py, #69550) -# — never re-inlined wording. Used ONLY by the opt-in -# ``compression.progress_notices`` gate below (#52995) to decide which of the -# noisy statuses matched by _TELEGRAM_NOISY_STATUS_RE are compression -# progress (deliverable when the user opted in) versus unrelated aux/retry -# chatter (always suppressed on chat surfaces). Failure notices and manual -# /compress feedback never match _TELEGRAM_NOISY_STATUS_RE in the first -# place, so they are unaffected by this gate. -_COMPRESSION_PROGRESS_STATUS_RE = re.compile( - "|".join( - _status_template_to_regex(_template) - for _template in ( - COMPACTION_STATUS, - PRE_API_COMPRESSION_STATUS_TEMPLATE, - PREFLIGHT_COMPRESSION_STATUS_TEMPLATE, - IDLE_COMPACTION_STATUS_TEMPLATE, - COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE, - COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE, - COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE, - COMPRESSION_RETRY_CONTEXT_REDUCED_STATUS_TEMPLATE, - ) - ), - re.IGNORECASE, +# --- Compression-status helpers (extracted to gateway.status_helpers) --- +from gateway.status_helpers import ( # noqa: E402 + _COMPRESSION_PROGRESS_STATUS_RE, + _status_template_to_regex, ) diff --git a/gateway/status_helpers.py b/gateway/status_helpers.py new file mode 100644 index 0000000000000..7f7af8336a111 --- /dev/null +++ b/gateway/status_helpers.py @@ -0,0 +1,63 @@ +"""Compression-status helpers extracted from gateway/run.py (#54962). + +Fourth slice of the gateway god-file unpacking: the status-template-to-regex +compiler and the opt-in compression-progress matcher. Pure derivation — the +regex is built from the SAME template constants the emit sites format, so +wording drift between agent/conversation_compression.py and the gateway +matcher cannot silently diverge. +""" + +from __future__ import annotations + +import re + +from agent.conversation_compression import ( + COMPACTION_STATUS, + COMPRESSION_RETRY_CONTEXT_REDUCED_STATUS_TEMPLATE, + COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE, + COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE, + COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE, + IDLE_COMPACTION_STATUS_TEMPLATE, + PREFLIGHT_COMPRESSION_STATUS_TEMPLATE, + PRE_API_COMPRESSION_STATUS_TEMPLATE, +) + + +def _status_template_to_regex(template: str) -> str: + """Compile a compression status template constant into a regex source. + + Literal text is escaped verbatim (so wording drift in + agent/conversation_compression.py cannot silently diverge from this + matcher — the constants ARE the wording) and each ``{field}`` format + placeholder is replaced with a numeric-ish pattern covering every value + the emit sites format in (ints, ``{:,}`` thousands separators). + """ + parts = re.split(r"\{[^{}]*\}", template) + return r"[\d,]+".join(re.escape(part) for part in parts) + + +# ROUTINE compression progress statuses, derived from the SAME template +# constants the emit sites format (agent/conversation_compression.py, #69550) +# — never re-inlined wording. Used ONLY by the opt-in +# ``compression.progress_notices`` gate below (#52995) to decide which of the +# noisy statuses matched by _TELEGRAM_NOISY_STATUS_RE are compression +# progress (deliverable when the user opted in) versus unrelated aux/retry +# chatter (always suppressed on chat surfaces). Failure notices and manual +# /compress feedback never match _TELEGRAM_NOISY_STATUS_RE in the first +# place, so they are unaffected by this gate. +_COMPRESSION_PROGRESS_STATUS_RE = re.compile( + "|".join( + _status_template_to_regex(_template) + for _template in ( + COMPACTION_STATUS, + PRE_API_COMPRESSION_STATUS_TEMPLATE, + PREFLIGHT_COMPRESSION_STATUS_TEMPLATE, + IDLE_COMPACTION_STATUS_TEMPLATE, + COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE, + COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE, + COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE, + COMPRESSION_RETRY_CONTEXT_REDUCED_STATUS_TEMPLATE, + ) + ), + re.IGNORECASE, +) From 605227ec8774e1740b355c13c459dd1ae3c0414e Mon Sep 17 00:00:00 2001 From: andrexibiza <84248988+andrexibiza@users.noreply.github.com> Date: Mon, 3 Aug 2026 02:56:36 -0500 Subject: [PATCH 2/3] refactor(gateway): extract transient-error helpers from run.py (slice 5 of #54962) Fifth slice of the gateway god-file unpacking: extract _is_transient_network_error and _gateway_loop_exception_handler into gateway/error_helpers.py. - Byte-identical extraction (AST-verified against origin/main): the classifier walks the exception cause chain (#31066/#31110); the loop handler wires it into the event-loop safety net - gateway/run.py: -78 lines; helpers imported at the extraction point - 46 tests pass (loop exception handler + compression notices) Follow-up to #77433 (slice 1), #77438 (slices 2-3), #77450 (slice 4). Progress on #54962. Signed-off-by: andrexibiza <84248988+andrexibiza@users.noreply.github.com> --- gateway/error_helpers.py | 94 ++++++++++++++++++++++++++++++++++++++++ gateway/run.py | 81 +++------------------------------- 2 files changed, 99 insertions(+), 76 deletions(-) create mode 100644 gateway/error_helpers.py diff --git a/gateway/error_helpers.py b/gateway/error_helpers.py new file mode 100644 index 0000000000000..3370ae975e45a --- /dev/null +++ b/gateway/error_helpers.py @@ -0,0 +1,94 @@ +"""Transient-error classification + loop handler extracted from gateway/run.py (#54962). + +Fifth slice of the gateway god-file unpacking: the transient network +error classifier and the asyncio loop-level exception handler that uses +it. The classifier is a pure function (walk the cause chain, match known +transient exception class names); the handler wires it into the event +loop's safety net so a transient peer error can never kill the gateway +process (#31066 / #31110). +""" + +from __future__ import annotations + +import logging +from typing import Any, Dict, Optional + +logger = logging.getLogger(__name__) + + +def _is_transient_network_error(exc: BaseException) -> bool: + """Return True for transient network errors safe to log + swallow. + + The crash class targeted by #31066 / #31110: an unhandled Telegram + ``TimedOut`` (or peer ``NetworkError`` / ``httpx`` connection error) + propagating to the event loop and killing the entire gateway + process. These are by definition transient — the next poll cycle or + user action recovers — so they must never crash the process. + + Walk the exception cause chain so wrapped errors (e.g. PTB's + ``NetworkError`` wrapping ``httpx.ConnectError``) are still + classified. The chain is bounded to avoid pathological cycles. + """ + seen: set[int] = set() + cur: Optional[BaseException] = exc + depth = 0 + transient_class_names = { + "TimedOut", + "NetworkError", + "ReadError", + "WriteError", + "ConnectError", + "ConnectTimeout", + "ReadTimeout", + "WriteTimeout", + "PoolTimeout", + "RemoteProtocolError", + "ServerDisconnectedError", + "ClientConnectorError", + "ClientOSError", + } + while cur is not None and depth < 12: + ident = id(cur) + if ident in seen: + break + seen.add(ident) + depth += 1 + name = type(cur).__name__ + if name in transient_class_names: + return True + cur = cur.__cause__ or cur.__context__ + return False + + +def _gateway_loop_exception_handler( + loop: "asyncio.AbstractEventLoop", context: Dict[str, Any] +) -> None: + """Loop-level safety net for transient network errors. + + Installed once during :func:`start_gateway`. Catches the + ``telegram.error.TimedOut`` crash class (issues #31066 / #31110) + and any peer transient network error before it can kill the + gateway process. Logs at WARNING with full traceback so the + originating call site stays diagnosable; non-transient errors + are forwarded to the default loop handler so real bugs still + surface. + """ + exc = context.get("exception") + if exc is not None and _is_transient_network_error(exc): + task = context.get("future") or context.get("task") + task_name = "" + if task is not None: + try: + task_name = task.get_name() if hasattr(task, "get_name") else repr(task) + except Exception: + task_name = repr(task) + logger.warning( + "Gateway swallowed transient network error from %s: %s: %s", + task_name or "", + type(exc).__name__, + exc, + exc_info=(type(exc), exc, exc.__traceback__), + ) + return + # Fall back to the default handler for anything we don't recognise. + loop.default_exception_handler(context) diff --git a/gateway/run.py b/gateway/run.py index 14d762461d817..492472781427e 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -324,82 +324,11 @@ def _seed_hygiene_system_prompt( return bool(stored_prompt) -def _is_transient_network_error(exc: BaseException) -> bool: - """Return True for transient network errors safe to log + swallow. - - The crash class targeted by #31066 / #31110: an unhandled Telegram - ``TimedOut`` (or peer ``NetworkError`` / ``httpx`` connection error) - propagating to the event loop and killing the entire gateway - process. These are by definition transient — the next poll cycle or - user action recovers — so they must never crash the process. - - Walk the exception cause chain so wrapped errors (e.g. PTB's - ``NetworkError`` wrapping ``httpx.ConnectError``) are still - classified. The chain is bounded to avoid pathological cycles. - """ - seen: set[int] = set() - cur: Optional[BaseException] = exc - depth = 0 - transient_class_names = { - "TimedOut", - "NetworkError", - "ReadError", - "WriteError", - "ConnectError", - "ConnectTimeout", - "ReadTimeout", - "WriteTimeout", - "PoolTimeout", - "RemoteProtocolError", - "ServerDisconnectedError", - "ClientConnectorError", - "ClientOSError", - } - while cur is not None and depth < 12: - ident = id(cur) - if ident in seen: - break - seen.add(ident) - depth += 1 - name = type(cur).__name__ - if name in transient_class_names: - return True - cur = cur.__cause__ or cur.__context__ - return False - - -def _gateway_loop_exception_handler( - loop: "asyncio.AbstractEventLoop", context: Dict[str, Any] -) -> None: - """Loop-level safety net for transient network errors. - - Installed once during :func:`start_gateway`. Catches the - ``telegram.error.TimedOut`` crash class (issues #31066 / #31110) - and any peer transient network error before it can kill the - gateway process. Logs at WARNING with full traceback so the - originating call site stays diagnosable; non-transient errors - are forwarded to the default loop handler so real bugs still - surface. - """ - exc = context.get("exception") - if exc is not None and _is_transient_network_error(exc): - task = context.get("future") or context.get("task") - task_name = "" - if task is not None: - try: - task_name = task.get_name() if hasattr(task, "get_name") else repr(task) - except Exception: - task_name = repr(task) - logger.warning( - "Gateway swallowed transient network error from %s: %s: %s", - task_name or "", - type(exc).__name__, - exc, - exc_info=(type(exc), exc, exc.__traceback__), - ) - return - # Fall back to the default handler for anything we don't recognise. - loop.default_exception_handler(context) +# --- Transient-error helpers (extracted to gateway.error_helpers) --- +from gateway.error_helpers import ( # noqa: E402 + _gateway_loop_exception_handler, + _is_transient_network_error, +) def _redact_gateway_user_facing_secrets(text: str) -> str: From fcbbefa30803f3d07f7ee6f11676c6f59141e7d4 Mon Sep 17 00:00:00 2001 From: andrexibiza <84248988+andrexibiza@users.noreply.github.com> Date: Mon, 3 Aug 2026 03:03:50 -0500 Subject: [PATCH 3/3] Revert "refactor(gateway): extract transient-error helpers from run.py (slice 5 of #54962)" This reverts commit 605227ec8774e1740b355c13c459dd1ae3c0414e. --- gateway/error_helpers.py | 94 ---------------------------------------- gateway/run.py | 81 +++++++++++++++++++++++++++++++--- 2 files changed, 76 insertions(+), 99 deletions(-) delete mode 100644 gateway/error_helpers.py diff --git a/gateway/error_helpers.py b/gateway/error_helpers.py deleted file mode 100644 index 3370ae975e45a..0000000000000 --- a/gateway/error_helpers.py +++ /dev/null @@ -1,94 +0,0 @@ -"""Transient-error classification + loop handler extracted from gateway/run.py (#54962). - -Fifth slice of the gateway god-file unpacking: the transient network -error classifier and the asyncio loop-level exception handler that uses -it. The classifier is a pure function (walk the cause chain, match known -transient exception class names); the handler wires it into the event -loop's safety net so a transient peer error can never kill the gateway -process (#31066 / #31110). -""" - -from __future__ import annotations - -import logging -from typing import Any, Dict, Optional - -logger = logging.getLogger(__name__) - - -def _is_transient_network_error(exc: BaseException) -> bool: - """Return True for transient network errors safe to log + swallow. - - The crash class targeted by #31066 / #31110: an unhandled Telegram - ``TimedOut`` (or peer ``NetworkError`` / ``httpx`` connection error) - propagating to the event loop and killing the entire gateway - process. These are by definition transient — the next poll cycle or - user action recovers — so they must never crash the process. - - Walk the exception cause chain so wrapped errors (e.g. PTB's - ``NetworkError`` wrapping ``httpx.ConnectError``) are still - classified. The chain is bounded to avoid pathological cycles. - """ - seen: set[int] = set() - cur: Optional[BaseException] = exc - depth = 0 - transient_class_names = { - "TimedOut", - "NetworkError", - "ReadError", - "WriteError", - "ConnectError", - "ConnectTimeout", - "ReadTimeout", - "WriteTimeout", - "PoolTimeout", - "RemoteProtocolError", - "ServerDisconnectedError", - "ClientConnectorError", - "ClientOSError", - } - while cur is not None and depth < 12: - ident = id(cur) - if ident in seen: - break - seen.add(ident) - depth += 1 - name = type(cur).__name__ - if name in transient_class_names: - return True - cur = cur.__cause__ or cur.__context__ - return False - - -def _gateway_loop_exception_handler( - loop: "asyncio.AbstractEventLoop", context: Dict[str, Any] -) -> None: - """Loop-level safety net for transient network errors. - - Installed once during :func:`start_gateway`. Catches the - ``telegram.error.TimedOut`` crash class (issues #31066 / #31110) - and any peer transient network error before it can kill the - gateway process. Logs at WARNING with full traceback so the - originating call site stays diagnosable; non-transient errors - are forwarded to the default loop handler so real bugs still - surface. - """ - exc = context.get("exception") - if exc is not None and _is_transient_network_error(exc): - task = context.get("future") or context.get("task") - task_name = "" - if task is not None: - try: - task_name = task.get_name() if hasattr(task, "get_name") else repr(task) - except Exception: - task_name = repr(task) - logger.warning( - "Gateway swallowed transient network error from %s: %s: %s", - task_name or "", - type(exc).__name__, - exc, - exc_info=(type(exc), exc, exc.__traceback__), - ) - return - # Fall back to the default handler for anything we don't recognise. - loop.default_exception_handler(context) diff --git a/gateway/run.py b/gateway/run.py index 492472781427e..14d762461d817 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -324,11 +324,82 @@ def _seed_hygiene_system_prompt( return bool(stored_prompt) -# --- Transient-error helpers (extracted to gateway.error_helpers) --- -from gateway.error_helpers import ( # noqa: E402 - _gateway_loop_exception_handler, - _is_transient_network_error, -) +def _is_transient_network_error(exc: BaseException) -> bool: + """Return True for transient network errors safe to log + swallow. + + The crash class targeted by #31066 / #31110: an unhandled Telegram + ``TimedOut`` (or peer ``NetworkError`` / ``httpx`` connection error) + propagating to the event loop and killing the entire gateway + process. These are by definition transient — the next poll cycle or + user action recovers — so they must never crash the process. + + Walk the exception cause chain so wrapped errors (e.g. PTB's + ``NetworkError`` wrapping ``httpx.ConnectError``) are still + classified. The chain is bounded to avoid pathological cycles. + """ + seen: set[int] = set() + cur: Optional[BaseException] = exc + depth = 0 + transient_class_names = { + "TimedOut", + "NetworkError", + "ReadError", + "WriteError", + "ConnectError", + "ConnectTimeout", + "ReadTimeout", + "WriteTimeout", + "PoolTimeout", + "RemoteProtocolError", + "ServerDisconnectedError", + "ClientConnectorError", + "ClientOSError", + } + while cur is not None and depth < 12: + ident = id(cur) + if ident in seen: + break + seen.add(ident) + depth += 1 + name = type(cur).__name__ + if name in transient_class_names: + return True + cur = cur.__cause__ or cur.__context__ + return False + + +def _gateway_loop_exception_handler( + loop: "asyncio.AbstractEventLoop", context: Dict[str, Any] +) -> None: + """Loop-level safety net for transient network errors. + + Installed once during :func:`start_gateway`. Catches the + ``telegram.error.TimedOut`` crash class (issues #31066 / #31110) + and any peer transient network error before it can kill the + gateway process. Logs at WARNING with full traceback so the + originating call site stays diagnosable; non-transient errors + are forwarded to the default loop handler so real bugs still + surface. + """ + exc = context.get("exception") + if exc is not None and _is_transient_network_error(exc): + task = context.get("future") or context.get("task") + task_name = "" + if task is not None: + try: + task_name = task.get_name() if hasattr(task, "get_name") else repr(task) + except Exception: + task_name = repr(task) + logger.warning( + "Gateway swallowed transient network error from %s: %s: %s", + task_name or "", + type(exc).__name__, + exc, + exc_info=(type(exc), exc, exc.__traceback__), + ) + return + # Fall back to the default handler for anything we don't recognise. + loop.default_exception_handler(context) def _redact_gateway_user_facing_secrets(text: str) -> str: