refactor(adapters/telegram): 11957->8819; unify send/edit retry classification, media caching, auth gates, model-picker branches, polling recovery; compact docs

This commit is contained in:
Teknium
2026-09-02 14:06:32 -07:00
parent 581d97e545
commit cc11c6bfdd
5 changed files with 2647 additions and 5786 deletions

File diff suppressed because it is too large Load Diff

View File

@@ -1,26 +1,14 @@
#!/usr/bin/env python3
"""Telegram inline command picker — searchable access to EVERY command/skill.
Telegram's BotCommand menu is capped (100 per scope, ~4KB payload; Hermes
defaults to 60 slots), so most skill commands can never appear in the ``/``
menu. Inline mode has no such cap: typing ``@yourbot <query>`` in any chat
asks the bot for results live, per keystroke, paginated 50 at a time — the
same trick Discord's ``/skill`` autocomplete uses (options fetched
dynamically, nothing pre-registered).
The BotCommand menu is capped (100/scope, Hermes uses 60), so inline mode
(``@yourbot <query>``, live per keystroke, 50 per page) exposes the rest. Tapping a
result sends ``/cmd args`` as the user; it starts with ``/`` so it arrives even under
privacy mode and dispatches through the normal command path.
Tapping a result sends the command text (e.g. ``/plan migrate the auth``)
into the chat as the user. Because the sent message starts with ``/``, the
bot receives it even under Telegram's default privacy mode ("messages with
commands meant for the bot" are always delivered), and it dispatches through
the existing command path — zero new dispatch code.
This module is PTB-object-free on purpose: it returns plain dicts so the
catalog/filter/pagination logic is unit-testable without python-telegram-bot
installed. The adapter converts dicts to ``InlineQueryResultArticle``.
Setup note (docs): inline mode must be enabled once per bot via BotFather's
``/setinline``. Until then Telegram never delivers ``inline_query`` updates,
so the registered handler is inert — safe to ship enabled by default.
PTB-object-free on purpose: plain dicts keep catalog/filter/pagination unit-testable
without python-telegram-bot; the adapter converts to ``InlineQueryResultArticle``.
Inline mode must be enabled via BotFather ``/setinline``; until then the handler is inert.
"""
from __future__ import annotations
@@ -33,24 +21,15 @@ logger = logging.getLogger(__name__)
# Telegram hard limit: max 50 results per answerInlineQuery call.
PAGE_SIZE = 50
# Results depend on the caller's auth and the install's skill set — never
# share cached results across users, and keep the cache short so freshly
# installed skills appear quickly.
# Results depend on caller auth + installed skills: never share across users, keep short.
CACHE_TIME_SECONDS = 10
def collect_inline_catalog() -> List[Dict[str, str]]:
"""Return every dispatchable command as ``{name, description}`` dicts.
"""Every dispatchable command as ``{name, description}``, first occurrence wins.
Sources, deduped in priority order (first occurrence wins):
1. Core gateway-visible ``CommandDef`` commands (Telegram-sanitized
names, same gating as the BotCommand menu).
2. Plugin slash commands + built-in skill commands via the shared
collector — with ``max_slots=None`` so NOTHING is trimmed. This is
the whole point: the inline picker has no cap.
Skill entries honor the same filtering as the menu (hub excluded,
per-platform disabled excluded, external-dir allowlist).
Core gateway commands first (menu gating), then plugin + skill commands via the
shared collector with ``max_slots=None`` — the inline picker has no cap.
"""
catalog: List[Dict[str, str]] = []
seen: set[str] = set()
@@ -95,10 +74,7 @@ def collect_inline_catalog() -> List[Dict[str, str]]:
def filter_catalog(catalog: List[Dict[str, str]], term: str) -> List[Dict[str, str]]:
"""Rank *catalog* against *term*: prefix > name-substring > description.
Empty term returns the full catalog in its collection order (core first,
then plugins, then skills alphabetically) — the "browse" view.
"""
Empty term returns the full catalog in collection order (the "browse" view)."""
term = (term or "").strip().lower().lstrip("/")
if not term:
return list(catalog)
@@ -120,20 +96,13 @@ def filter_catalog(catalog: List[Dict[str, str]], term: str) -> List[Dict[str, s
def build_inline_results(
query: str,
offset: str = "",
page_size: int = PAGE_SIZE,
query: str, offset: str = "", page_size: int = PAGE_SIZE,
) -> Tuple[List[Dict[str, Any]], str]:
"""Build one page of inline results for *query*.
"""One page of inline results for *query*.
The first whitespace-separated token of *query* filters the catalog; any
remainder is carried into the sent command as its argument. Example:
``@bot plan migrate auth to OIDC`` → filter ``plan``, and tapping the
``/plan`` result sends ``/plan migrate auth to OIDC``.
Returns ``(results, next_offset)`` where each result is
``{"id", "title", "description", "message_text"}`` and *next_offset* is
``""`` when this is the last page (Telegram's stop signal).
First token filters the catalog; the remainder becomes the command argument
(``@bot plan migrate auth`` → tapping ``/plan`` sends ``/plan migrate auth``).
Returns ``(results, next_offset)``; ``next_offset == ""`` means last page.
"""
query = (query or "").strip()
parts = query.split(None, 1)
@@ -154,13 +123,10 @@ def build_inline_results(
message_text = f"/{item['name']}"
if args:
message_text += f" {args}"
results.append(
{
# Offset-scoped ids stay unique across pages of one query.
"id": f"{start}:{item['name']}"[:64],
"title": f"/{item['name']}",
"description": (item.get("description") or "")[:100],
"message_text": message_text[:4096],
}
)
results.append({
"id": f"{start}:{item['name']}"[:64], # offset-scoped: unique across pages
"title": f"/{item['name']}",
"description": (item.get("description") or "")[:100],
"message_text": message_text[:4096],
})
return results, next_offset

View File

@@ -35,11 +35,6 @@ def normalize_telegram_chat_id(chat_id: Any) -> Union[int, str]:
return chat_id_str
def telegram_chat_id_key(chat_id: Any) -> str:
"""Stable string key for a chat_id (for dict keys / persisted state)."""
return str(normalize_telegram_chat_id(chat_id))
def looks_like_telegram_username(chat_id: Any) -> bool:
"""True when the value is an ``@username``-format Telegram chat identifier."""
return bool(_TELEGRAM_USERNAME_RE.fullmatch(str(chat_id).strip()))

View File

@@ -1,11 +1,5 @@
"""Telegram-specific network helpers.
Provides a hostname-preserving fallback transport for networks where
api.telegram.org resolves to an endpoint that is unreachable from the current
host. The transport keeps the logical request host and TLS SNI as
api.telegram.org while retrying the TCP connection against one or more fallback
IPv4 addresses.
"""
"""Telegram network helpers: a hostname-preserving fallback transport (Host + SNI stay
api.telegram.org while TCP retries known IPv4 literals) plus DoH-based IP discovery."""
from __future__ import annotations
@@ -21,26 +15,18 @@ logger = logging.getLogger(__name__)
_TELEGRAM_API_HOST = "api.telegram.org"
# TCP keepalive so a half-open or CLOSE-WAIT long-poll errors out instead of
# blocking getUpdates indefinitely. Windows does not enable SO_KEEPALIVE on
# new sockets by default, so a dead api.telegram.org peer can hang forever
# (#87057). Idle/interval knobs are best-effort — not every Python/OS combo
# exposes TCP_KEEPIDLE / TCP_KEEPALIVE.
# TCP keepalive so a half-open/CLOSE-WAIT long-poll errors out instead of blocking
# getUpdates forever (Windows leaves SO_KEEPALIVE off by default). Idle/interval knobs
# are best-effort — not every Python/OS combo exposes TCP_KEEPIDLE / TCP_KEEPALIVE.
_TCP_KEEPALIVE_IDLE_S = 30
_TCP_KEEPALIVE_INTERVAL_S = 10
_TCP_KEEPALIVE_COUNT = 3
def tcp_keepalive_socket_options() -> list[tuple[int, int, int]]:
"""Return ``setsockopt`` tuples that enable TCP keepalive on new sockets.
Pure data for httpx/httpcore ``socket_options``. Safe on every host: the
list always includes ``SO_KEEPALIVE`` and adds idle/interval/count only
when the running interpreter exposes those option names.
"""
options: list[tuple[int, int, int]] = [
(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1),
]
"""``setsockopt`` tuples for httpx ``socket_options``: always SO_KEEPALIVE, plus
idle/interval/count when the interpreter exposes those option names."""
options: list[tuple[int, int, int]] = [(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)]
idle = getattr(socket, "TCP_KEEPIDLE", None) or getattr(socket, "TCP_KEEPALIVE", None)
if idle is not None:
options.append((socket.IPPROTO_TCP, idle, _TCP_KEEPALIVE_IDLE_S))
@@ -52,50 +38,38 @@ def tcp_keepalive_socket_options() -> list[tuple[int, int, int]]:
options.append((socket.IPPROTO_TCP, count, _TCP_KEEPALIVE_COUNT))
return options
# DNS-over-HTTPS providers used to discover Telegram API IPs that may differ
# from the (potentially unreachable) IP returned by the local system resolver.
_DOH_TIMEOUT = 4.0 # seconds — bounded so connect() isn't noticeably delayed
# DNS-over-HTTPS providers: discover Telegram API IPs that may differ from the
# (possibly unreachable) one the local resolver returns. Bounded so connect() isn't delayed.
_DOH_TIMEOUT = 4.0
_DOH_PROVIDERS: list[dict] = [
{
"url": "https://dns.google/resolve",
"params": {"name": _TELEGRAM_API_HOST, "type": "A"},
"headers": {},
},
{"url": "https://dns.google/resolve", "params": {"name": _TELEGRAM_API_HOST, "type": "A"}, "headers": {}},
{
"url": "https://cloudflare-dns.com/dns-query",
"params": {"name": _TELEGRAM_API_HOST, "type": "A"},
"headers": {"Accept": "application/dns-json"},
},
]
# Last-resort IPv4 Telegram Bot API endpoints in 149.154.160.0/20
# (same seed used by OpenClaw). Used when DoH is blocked AND as the
# first-try connect targets so a blackholed IPv6 AAAA for the hostname
# cannot pin initialize() (#87015).
# Last-resort IPv4 Bot API endpoints (149.154.160.0/20). Used when DoH is blocked AND as
# first-try connect targets so a blackholed IPv6 AAAA for the hostname can't pin initialize().
SEED_FALLBACK_IPS: list[str] = ["149.154.166.110", "149.154.167.220"]
_UNSET = object()
def _resolve_proxy_url(target_hosts=None) -> str | None:
# Delegate to shared implementation (env vars + macOS system proxy detection)
from gateway.platforms.base import resolve_proxy_url
from gateway.platforms.base import resolve_proxy_url # env vars + macOS system proxy
return resolve_proxy_url("TELEGRAM_PROXY", target_hosts=target_hosts)
class TelegramFallbackTransport(httpx.AsyncBaseTransport):
"""Reach Telegram Bot API via known IPv4 literals first, hostname last.
"""Reach the Bot API via known IPv4 literals first, dual-stack hostname last.
Requests still target https://api.telegram.org/... logically (Host + SNI
stay on the hostname). TCP connects to a known A-record IP first so a
blackholed IPv6 AAAA cannot pin initialize(). Equivalent to
``curl --resolve api.telegram.org:443:<ip>``. The dual-stack hostname
is last resort for IPv6-only networks.
Logically requests still target https://api.telegram.org (Host + SNI stay on the
hostname) — like ``curl --resolve api.telegram.org:443:<ip>`` — so a blackholed
IPv6 AAAA can't pin initialize(); the hostname remains for IPv6-only networks.
"""
# Bound every pool. httpx defaults to 100 connections per pool, so a wedged
# endpoint plus the seed IPs can outgrow the process file-descriptor limit
# on its own (#63311).
# Bound every pool: httpx's 100-connection default × (wedged endpoint + seed IPs)
# can outgrow the process fd limit on its own.
_POOL_LIMITS = httpx.Limits(max_connections=8, max_keepalive_connections=4)
def __init__(self, fallback_ips: Iterable[str], **transport_kwargs):
@@ -112,9 +86,7 @@ class TelegramFallbackTransport(httpx.AsyncBaseTransport):
# Built on demand and discarded on failure — see _reset_fallback.
self._fallbacks: dict[str, httpx.AsyncHTTPTransport] = {}
self._fallback_lock = asyncio.Lock()
# ``_UNSET`` vs ``None`` vs ``str``: unset / sticky hostname / sticky IPv4.
# ``None`` cannot mean both "no sticky yet" and "sticky dual-stack
# hostname" (#87015).
# ``_UNSET`` / ``None`` / ``str`` = no sticky yet / sticky hostname / sticky IPv4.
self._sticky_ip: object = _UNSET
self._sticky_lock = asyncio.Lock()
@@ -127,8 +99,8 @@ class TelegramFallbackTransport(httpx.AsyncBaseTransport):
return transport
async def _reset_primary(self, transport: httpx.AsyncHTTPTransport) -> None:
# Retryable primary failures can leave half-closed sockets in the pool;
# replace and close the failed generation before trying fallback.
# Retryable primary failures leave half-closed sockets in the pool; replace the
# generation and close the old one before trying fallback.
async with self._primary_lock:
if self._primary_closed or transport is not self._primary:
return
@@ -139,13 +111,8 @@ class TelegramFallbackTransport(httpx.AsyncBaseTransport):
logger.debug("[Telegram] Error closing primary transport: %s", exc)
async def _reset_fallback(self, ip: str) -> None:
"""Discard a failed fallback pool so its dead sockets are released.
A connect that reaches ESTABLISHED and is then closed by the peer leaves
its socket in CLOSE_WAIT inside the pool. Retaining the poisoned pool
leaks one descriptor per retry until the process hits its file limit and
can no longer accept connections or resolve DNS (#63311).
"""
"""Discard a failed fallback pool: a peer-closed connect leaves a CLOSE_WAIT socket
in it, and keeping the poisoned pool leaks one fd per retry until the process limit."""
async with self._fallback_lock:
transport = self._fallbacks.pop(ip, None)
if transport is None:
@@ -156,13 +123,10 @@ class TelegramFallbackTransport(httpx.AsyncBaseTransport):
logger.debug("[Telegram] Error closing fallback transport %s: %s", ip, exc)
def _attempt_order(self) -> list[Optional[str]]:
"""IPv4 literals first; dual-stack hostname last.
"""Sticky path first, then IPv4 literals, dual-stack hostname last.
A blackholed IPv6 path to ``api.telegram.org`` never errors — Happy
Eyeballs waits on AAAA until the OS TCP timeout, which can pin the
event loop so ``_await_with_thread_deadline`` never fires (#87015).
Known A-record IPs connect over IPv4 immediately. The hostname is
kept as a last resort for IPv6-only networks.
A blackholed IPv6 path never errors — Happy Eyeballs waits on AAAA until the OS
TCP timeout and can pin the loop so the thread deadline never fires.
"""
order: list[Optional[str]] = []
if self._sticky_ip is not _UNSET:
@@ -195,8 +159,7 @@ class TelegramFallbackTransport(httpx.AsyncBaseTransport):
log = logger.warning if last_error is not None else logger.info
log(
"[Telegram] Using sticky IPv4 Telegram API path %s "
"(dual-stack hostname tried last — #87015)",
ip,
"(dual-stack hostname tried last — #87015)", ip,
)
return response
except Exception as exc:
@@ -214,10 +177,7 @@ class TelegramFallbackTransport(httpx.AsyncBaseTransport):
)
if ip is None:
await self._reset_primary(transport)
logger.warning(
"[Telegram] Dual-stack api.telegram.org path failed (%s)",
exc,
)
logger.warning("[Telegram] Dual-stack api.telegram.org path failed (%s)", exc)
continue
logger.warning("[Telegram] IPv4 Telegram API IP %s failed: %s", ip, exc)
await self._reset_fallback(ip)
@@ -276,14 +236,10 @@ def _resolve_system_dns() -> set[str]:
return set()
async def _query_doh_provider(
client: httpx.AsyncClient, provider: dict
) -> list[str]:
async def _query_doh_provider(client: httpx.AsyncClient, provider: dict) -> list[str]:
"""Query one DoH provider and return A-record IPs."""
try:
resp = await client.get(
provider["url"], params=provider["params"], headers=provider["headers"]
)
resp = await client.get(provider["url"], params=provider["params"], headers=provider["headers"])
resp.raise_for_status()
data = resp.json()
ips: list[str] = []
@@ -303,26 +259,19 @@ async def _query_doh_provider(
async def discover_fallback_ips() -> list[str]:
"""Auto-discover Telegram API IPs via DNS-over-HTTPS.
"""Resolve api.telegram.org via Google + Cloudflare DoH; unique A records, in order.
Resolves api.telegram.org through Google and Cloudflare DoH and returns all
unique A records. IPs that match the local system resolver are kept rather
than excluded: in many networks the system-DNS IP is the most reliable path
to api.telegram.org and a transient primary-path failure should be retried
against the same address via the IP-rewrite path before the seed list is
consulted (#14520). Falls back to a hardcoded seed list only when DoH
yields no usable answers.
IPs matching the system resolver are deliberately KEPT (often the most reliable
path; a transient primary failure should retry it via IP-rewrite before the seed
list). Falls back to ``SEED_FALLBACK_IPS`` only when DoH yields nothing usable.
"""
async with httpx.AsyncClient(timeout=httpx.Timeout(_DOH_TIMEOUT)) as client:
doh_tasks = [_query_doh_provider(client, p) for p in _DOH_PROVIDERS]
system_dns_task = asyncio.ensure_future(asyncio.to_thread(_resolve_system_dns))
results = await asyncio.gather(*doh_tasks, return_exceptions=True)
# The system-resolver leg runs socket.getaddrinfo in a worker thread with
# no timeout of its own — a wedged OS resolver (broken VPN/DNS) can sit for
# minutes. Its result only feeds the no-usable-answers log line below, so
# it must never gate discovery: bound it and move on (#63309). The DoH legs
# are already bounded by the client timeout above.
# The getaddrinfo leg has no timeout of its own (a wedged resolver can sit for
# minutes) and only feeds the log line below — bound it, never gate discovery on it.
system_ips: set[str] = set()
try:
system_result = await asyncio.wait_for(system_dns_task, timeout=_DOH_TIMEOUT)
@@ -335,18 +284,7 @@ async def discover_fallback_ips() -> list[str]:
for r in results:
if isinstance(r, list):
doh_ips.extend(r)
# Deduplicate preserving order
seen: set[str] = set()
candidates: list[str] = []
for ip in doh_ips:
if ip not in seen:
seen.add(ip)
candidates.append(ip)
# Validate through existing normalization
validated = _normalize_fallback_ips(candidates)
validated = _normalize_fallback_ips(list(dict.fromkeys(doh_ips))) # dedupe, keep order
if validated:
logger.debug("Discovered Telegram fallback IPs via DoH: %s", ", ".join(validated))
return validated
@@ -367,11 +305,7 @@ def _rewrite_request_for_ip(request: httpx.Request, ip: str) -> httpx.Request:
extensions = dict(request.extensions)
extensions["sni_hostname"] = original_host
return httpx.Request(
method=request.method,
url=url,
headers=headers,
stream=request.stream,
extensions=extensions,
method=request.method, url=url, headers=headers, stream=request.stream, extensions=extensions,
)

View File

@@ -18,7 +18,6 @@ from plugins.platforms.telegram.telegram_ids import (
looks_like_telegram_username,
normalize_telegram_chat_id,
parse_telegram_username_target,
telegram_chat_id_key,
)