mirror of
https://github.com/NousResearch/hermes-agent.git
synced 2026-07-22 16:25:58 +00:00
* feat(attribution): conflict-free contributor mappings via contributors/emails/ directory
The AUTHOR_MAP dict in scripts/release.py was a merge-conflict magnet:
every concurrent salvage PR appended entries to the same lines of the
same file, so parallel PRs re-conflicted on every merge to main.
New system: one file per email under contributors/emails/ — filename is
the commit-author email, first non-comment line is the GitHub login.
File additions never conflict, so any number of PRs can add mappings
concurrently.
- scripts/release.py: AUTHOR_MAP is now LEGACY_AUTHOR_MAP (frozen)
merged with the directory at import time (directory wins). All
existing consumers (resolve_author, contributor_audit.py) unchanged.
- scripts/add_contributor.py: idempotent CLI to add a mapping; refuses
conflicting reassignments (incl. against the legacy map), validates
email/login shapes.
- contributor-check.yml: attribution gate now accepts a mapping file OR
a legacy entry; failure message prints the exact add_contributor
command. Also auto-resolves bare <login>@users.noreply.github.com
emails is intentionally NOT added (kept id+login form only, matching
previous behavior).
- contributor_audit.py: guidance now points at add_contributor.py.
- tests/scripts/test_contributor_map.py: 12 tests covering loader,
merge precedence, CLI idempotency/conflict/validation, subprocess E2E.
* feat(ci): one-shot per-file flake retry in the parallel test runner
A failing test FILE is re-run once in a fresh subprocess. Pass-on-retry
counts as green but is loudly reported in a '⚠ FLAKY' summary section
(with both attempts' output preserved) so the flake gets fixed instead
of eating a full-run rerun. Deterministic failures fail both attempts —
regressions cannot be laundered green.
- --file-retries N / HERMES_TEST_FILE_RETRIES (default 1, 0 disables)
- E2E verified: simulated first-run-fail flake goes green with banner;
deterministic failure still exits 1; retries=0 restores old behavior.
This converts the dominant CI failure mode (one timing-sensitive test
flaking a 4600-test shard, requiring a manual 10-minute rerun and an
agent triage loop) into a self-healing retry that costs one file's
runtime.
* test(approval): loosen wall-clock perf bounds 0.15s -> 2.0s
These guard against catastrophic regex backtracking (seconds-to-minutes
class), but 0.15s is within scheduler-stall noise on loaded shared CI
runners — test_max_accepted_separator_free_input_is_fast failed a CI
shard this week on runner load alone. 2.0s still catches the regression
class with zero flake surface.
* fix(ci): job timeouts everywhere + retries on all network installs
Reliability pass over every workflow:
- timeout-minutes on all 21 jobs that lacked one (a hung job previously
burned the 6-hour default runner budget)
- ./.github/actions/retry wrapped around every network-fetching install
that lacked it: pip installs (deploy-site, skills-index), npm ci
(deploy-site website, upload_to_pypi web + ui-tui), uv sync (docker
test deps). Deterministic build steps (npm run build) deliberately
NOT retried — split into separate steps so a real build failure fails
fast instead of retrying 3x.
* docs(agents): document the file-retry flake policy
* fix(ci): curl retries on deploy hook + skills-index probe
* fix(ci): kill the remaining transient-failure classes in workflows + Dockerfile
From the workflow reliability audit:
- tests.yml: duration-cache restore had NO restore-keys while saves use
run_id-suffixed keys — the cache never matched once, so LPT slicing
always ran blind and unbalanced slices pushed heavy files toward the
per-file timeout. One-line restore-keys fixes slice balancing.
- Label gates (lint ci-reviewed, supply-chain mcp-catalog-reviewed):
'gh pr view || true' turned an API blip into 'label absent' → false
BLOCKING failure. Now 3x retry, and API failure is reported as an API
failure instead of a missing label.
- detect-changes action: compare API retried before failing open (was
silently running all lanes on any blip).
- uv-lockfile-check: 'uv lock --check' resolves against PyPI — retried
so registry blips don't read as 'lockfile stale'.
- docker.yml merge job: imagetools create retried (Docker Hub eventual
consistency on just-pushed digests).
- Dockerfile: apt-get Acquire::Retries=3; s6-overlay ADDs converted to
curl --retry 3 (ADD cannot retry; checksums still enforced); npm
--fetch-retries=5; playwright chromium fetch retried 3x.
- Advisory artifact uploads (per-slice durations, ci-timings report)
get continue-on-error so an artifact-service blip can't fail a green
test slice.
* fix(tests): kill the two root-cause flakes — leaking pre-warm timer + env-dependent provider list
- test_tui_gateway_server.py: session.create / non-eager session.resume
arm a 50ms threading.Timer (_schedule_agent_build) that outlives its
test and fires into the NEXT test's _make_agent mock, racily
corrupting captured state (the recurring session_resume shard
failures). Replaced the per-test whack-a-mole stub with a module-wide
autouse fixture; the 3 worker-lifecycle tests that genuinely need the
deferred build opt back in via @pytest.mark.real_agent_prewarm (new
marker in pyproject).
- test_api_key_providers.py: PROVIDER_ENV_VARS is now derived from the
live PROVIDER_REGISTRY instead of a hand-list that had drifted
(missing HF_TOKEN / DEEPINFRA_API_KEY) — resolve_provider('auto')
tests failed on any machine with HF_TOKEN exported. E2E-verified with
HF_TOKEN/DEEPINFRA_API_KEY set: 42/42 pass.
* test: de-flake 30 timing-sensitive test files for loaded CI runners
Root-cause fixes from the flake audit (session-DB mining + repo sweep):
Event-based sync instead of sleep-sync:
- title_generator: mock sets threading.Event, wait(10) replaces
sleep(0.3) hoping the daemon thread got scheduled
- docker zombie_reaping / profile_gateway: poll-for-state helpers
replace fixed 1-3s sleeps (s6 transitions + SIGCHLD reaping are async)
- process_registry tree test: select()-bounded readline replaces an
unbounded blocking read (parent wedge now fails THIS test with a clear
message instead of an opaque rc=124 file kill); SIGTERM grace 1s->2s
(the 1s partition window mid-interpreter-startup is how a child PID
escaped the live-system guard in CI)
Timeout raises (loaded 8-way-sliced runners see ~5s scheduling floors;
all of these complete in ms-to-1s when healthy so the raises cost
nothing on green runs):
- subprocess/thread waits <= 2s raised to 10-15s across mcp_tool,
mcp_circuit_breaker, mcp_reconnect_retry_reset, mcp_parked_self_probe,
mcp_cancelled_error_propagation, registry, clarify_gateway, interrupt,
voice_cli_integration, docker_environment, session_store_lock_io,
planned_stop_watcher, cli_interrupt_subagent, thread_scoped_output
(joins now also assert not is_alive() so stragglers fail loudly)
- wall-clock discrimination ceilings loosened where the guarded hang is
10x larger: local_background_child_hang 4s->10s, interrupt_cleanup
setup 5s->20s + pgid-exit 30s->60s, mcp_stability grandchild spinup
5s->15s, protocol/gil-starvation fast-handler 0.5s->2s,
iso_certify_seam 1.5s->5s, wait_for_mcp_discovery 0.1s->1s
- narrow assertion windows widened: honcho first-turn wait 0.4..0.65 ->
0.25..2.0 (property is bounded-not-hung, not an exact wall-clock);
compression fork-lock TTL 1s->3s (12 refresh chances per lease);
compression-lock expiry margins symmetric (ttl 0.05->0.5, sleep 1.0)
- telegram hung-DNS bound 1.0->1.4 (fake hang is 1.5s — must stay under)
* fix(tests): repair indentation from de-flake batch edit
* fix(tests): harden env isolation and replace remaining sleep-sync races
The full 42k-test run and complete npm check surfaced three more classes:
- Environment isolation: local ~/.honcho defaultHost and SSH_* variables
leaked into Python/TUI tests. Pin the default Honcho host in the
hermetic fixture, isolate the one fallback test from ~/.honcho, and
blank SSH_* around terminalSetup tests. This flipped 20 false failures
back to deterministic behavior on developer machines.
- Background-thread sleep-sync: Honcho async writer tests patched
time.sleep globally, then busy-polled with that same mocked sleep. Under
full-suite load the poller could starve the writer. Each test now waits
on an Event emitted by the exact flush/retry transition; 30/30 passed
under 15-way contention.
- Desktop streaming: the test slept 80ms and assumed a 500ms timer could
not fire before its assertion. A loaded runner descheduled the test for
>500ms and both chunks arrived. Producer controls now gate second-chunk
and completion transitions explicitly.
Also make file-retry observability complete: a self-healed flaky file now
prints BOTH attempts' full output in the FLAKY summary. Two behavioral
runner tests prove pass-on-retry is green+loud+traceback-preserving, while
a deterministic failure remains red.
* refactor(ci): use gh bot pat, better retries
refactor(ci): use retry action for PR label fetch
the retry action now captures stdout as a step output, so it can serve
double duty: retry + output capture for commands like 'gh pr view' whose
result must be consumed by later steps.
Retry action gains:
- 'stdout' output (heredoc-delimited to preserve newlines)
- tee to temp file so stdout still streams to the job log
- step id 'retry' for output reference
Both lint.yml and supply-chain-audit.yml now use the retry action
directly with 'command: gh pr view ...' and read
steps.<id>.outputs.stdout.
ci: use AUTOFIX_BOT_PAT for all gh CLI / GitHub API auth
Replace secrets.GITHUB_TOKEN and github.token with
secrets.AUTOFIX_BOT_PAT across all workflows and composite actions
that use the gh CLI or GitHub API. The PAT has consistent permissions
across fork PRs (where GITHUB_TOKEN is read-only), avoids API rate
limit sharing with the default token, and is already used by
js-autofix.yml for the same reasons.
19 sites swapped across 9 files:
- lint.yml (3): label fetch, comment post/edit, comment update
- supply-chain-audit.yml (5): scan, critical comment, unbounded dep
comment, label fetch, mcp-catalog comment
- lockfile-diff.yml (1): PR comment post/update
- skills-index-freshness.yml (1): issue creation on degraded probe
- skills-index.yml (2): index build, trigger deploy workflow
- upload_to_pypi.yml (2): release view poll, release upload
- ci.yml (1): timings report
- deploy-site.yml (2): skills index crawl
- detect-changes/action.yml (1): compare API call
---------
Co-authored-by: ethernet <arilotter@gmail.com>
758 lines
33 KiB
Python
758 lines
33 KiB
Python
"""Tests for plugins.platforms.telegram.telegram_network – fallback transport layer.
|
||
|
||
Background
|
||
----------
|
||
api.telegram.org resolves to an IP (e.g. 149.154.166.110) that is unreachable
|
||
from some networks. The workaround: route TCP through a different IP in the
|
||
same Telegram-owned 149.154.160.0/20 block (e.g. 149.154.167.220) while
|
||
keeping TLS SNI and the Host header as api.telegram.org so Telegram's edge
|
||
servers still accept the request. This is the programmatic equivalent of:
|
||
|
||
curl --resolve api.telegram.org:443:149.154.167.220 https://api.telegram.org/bot<token>/getMe
|
||
|
||
The TelegramFallbackTransport implements this: try the primary (DNS-resolved)
|
||
path first, and on ConnectTimeout / ConnectError fall through to configured
|
||
fallback IPs in order, then "stick" to whichever IP works.
|
||
"""
|
||
|
||
import httpx
|
||
import pytest
|
||
|
||
import plugins.platforms.telegram.telegram_network as tnet
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# Helpers
|
||
# ---------------------------------------------------------------------------
|
||
|
||
class FakeTransport(httpx.AsyncBaseTransport):
|
||
"""Records calls and raises / returns based on a host→action mapping."""
|
||
|
||
def __init__(self, calls, behavior):
|
||
self.calls = calls
|
||
self.behavior = behavior
|
||
self.closed = False
|
||
|
||
async def handle_async_request(self, request: httpx.Request) -> httpx.Response:
|
||
self.calls.append(
|
||
{
|
||
"url_host": request.url.host,
|
||
"host_header": request.headers.get("host"),
|
||
"sni_hostname": request.extensions.get("sni_hostname"),
|
||
"path": request.url.path,
|
||
}
|
||
)
|
||
action = self.behavior.get(request.url.host, "ok")
|
||
if action == "timeout":
|
||
raise httpx.ConnectTimeout("timed out")
|
||
if action == "connect_error":
|
||
raise httpx.ConnectError("connect error")
|
||
if isinstance(action, Exception):
|
||
raise action
|
||
return httpx.Response(200, request=request, text="ok")
|
||
|
||
async def aclose(self) -> None:
|
||
self.closed = True
|
||
|
||
|
||
def _fake_transport_factory(calls, behavior):
|
||
"""Returns a factory that creates FakeTransport instances."""
|
||
instances = []
|
||
|
||
def factory(**kwargs):
|
||
t = FakeTransport(calls, behavior)
|
||
instances.append(t)
|
||
return t
|
||
|
||
factory.instances = instances
|
||
return factory
|
||
|
||
|
||
def _telegram_request(path="/botTOKEN/getMe"):
|
||
return httpx.Request("GET", f"https://api.telegram.org{path}")
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
# IP parsing & validation
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
|
||
class TestParseFallbackIpEnv:
|
||
def test_filters_invalid_and_ipv6(self, caplog):
|
||
ips = tnet.parse_fallback_ip_env("149.154.167.220, bad, 2001:67c:4e8:f004::9,149.154.167.220")
|
||
assert ips == ["149.154.167.220", "149.154.167.220"]
|
||
assert "Ignoring invalid Telegram fallback IP" in caplog.text
|
||
assert "Ignoring non-IPv4 Telegram fallback IP" in caplog.text
|
||
|
||
def test_none_returns_empty(self):
|
||
assert tnet.parse_fallback_ip_env(None) == []
|
||
|
||
def test_empty_string_returns_empty(self):
|
||
assert tnet.parse_fallback_ip_env("") == []
|
||
|
||
def test_whitespace_only_returns_empty(self):
|
||
assert tnet.parse_fallback_ip_env(" , , ") == []
|
||
|
||
def test_single_valid_ip(self):
|
||
assert tnet.parse_fallback_ip_env("149.154.167.220") == ["149.154.167.220"]
|
||
|
||
def test_multiple_valid_ips(self):
|
||
ips = tnet.parse_fallback_ip_env("149.154.167.220, 149.154.167.221")
|
||
assert ips == ["149.154.167.220", "149.154.167.221"]
|
||
|
||
def test_rejects_leading_zeros(self, caplog):
|
||
"""Leading zeros are ambiguous (octal?) so ipaddress rejects them."""
|
||
ips = tnet.parse_fallback_ip_env("149.154.167.010")
|
||
assert ips == []
|
||
assert "Ignoring invalid" in caplog.text
|
||
|
||
|
||
class TestNormalizeFallbackIps:
|
||
def test_deduplication_happens_at_transport_level(self):
|
||
"""_normalize does not dedup; TelegramFallbackTransport.__init__ does."""
|
||
raw = ["149.154.167.220", "149.154.167.220"]
|
||
assert tnet._normalize_fallback_ips(raw) == ["149.154.167.220", "149.154.167.220"]
|
||
|
||
def test_empty_strings_skipped(self):
|
||
assert tnet._normalize_fallback_ips(["", " ", "149.154.167.220"]) == ["149.154.167.220"]
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
# Request rewriting
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
|
||
class TestRewriteRequestForIp:
|
||
def test_preserves_host_and_sni(self):
|
||
request = _telegram_request()
|
||
rewritten = tnet._rewrite_request_for_ip(request, "149.154.167.220")
|
||
|
||
assert rewritten.url.host == "149.154.167.220"
|
||
assert rewritten.headers["host"] == "api.telegram.org"
|
||
assert rewritten.extensions["sni_hostname"] == "api.telegram.org"
|
||
assert rewritten.url.path == "/botTOKEN/getMe"
|
||
|
||
def test_preserves_method_and_path(self):
|
||
request = httpx.Request("POST", "https://api.telegram.org/botTOKEN/sendMessage")
|
||
rewritten = tnet._rewrite_request_for_ip(request, "149.154.167.220")
|
||
|
||
assert rewritten.method == "POST"
|
||
assert rewritten.url.path == "/botTOKEN/sendMessage"
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
# Fallback transport – core behavior
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
|
||
class TestFallbackTransport:
|
||
"""Primary path fails → try fallback IPs → stick to whichever works."""
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_falls_back_on_connect_timeout_and_becomes_sticky(self, monkeypatch):
|
||
calls = []
|
||
behavior = {"api.telegram.org": "timeout", "149.154.167.220": "ok"}
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", _fake_transport_factory(calls, behavior))
|
||
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220"])
|
||
resp = await transport.handle_async_request(_telegram_request())
|
||
|
||
assert resp.status_code == 200
|
||
assert transport._sticky_ip == "149.154.167.220"
|
||
# First attempt was primary (api.telegram.org), second was fallback
|
||
assert calls[0]["url_host"] == "api.telegram.org"
|
||
assert calls[1]["url_host"] == "149.154.167.220"
|
||
assert calls[1]["host_header"] == "api.telegram.org"
|
||
assert calls[1]["sni_hostname"] == "api.telegram.org"
|
||
|
||
# Second request goes straight to sticky IP
|
||
calls.clear()
|
||
resp2 = await transport.handle_async_request(_telegram_request())
|
||
assert resp2.status_code == 200
|
||
assert calls[0]["url_host"] == "149.154.167.220"
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_falls_back_on_connect_error(self, monkeypatch):
|
||
calls = []
|
||
behavior = {"api.telegram.org": "connect_error", "149.154.167.220": "ok"}
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", _fake_transport_factory(calls, behavior))
|
||
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220"])
|
||
resp = await transport.handle_async_request(_telegram_request())
|
||
|
||
assert resp.status_code == 200
|
||
assert transport._sticky_ip == "149.154.167.220"
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_does_not_fallback_on_non_connect_error(self, monkeypatch):
|
||
"""Errors like ReadTimeout are not connection issues — don't retry."""
|
||
calls = []
|
||
behavior = {"api.telegram.org": httpx.ReadTimeout("read timeout"), "149.154.167.220": "ok"}
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", _fake_transport_factory(calls, behavior))
|
||
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220"])
|
||
|
||
with pytest.raises(httpx.ReadTimeout):
|
||
await transport.handle_async_request(_telegram_request())
|
||
|
||
assert [c["url_host"] for c in calls] == ["api.telegram.org"]
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_all_ips_fail_raises_last_error(self, monkeypatch):
|
||
calls = []
|
||
behavior = {"api.telegram.org": "timeout", "149.154.167.220": "timeout"}
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", _fake_transport_factory(calls, behavior))
|
||
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220"])
|
||
|
||
with pytest.raises(httpx.ConnectTimeout):
|
||
await transport.handle_async_request(_telegram_request())
|
||
|
||
assert [c["url_host"] for c in calls] == ["api.telegram.org", "149.154.167.220"]
|
||
assert transport._sticky_ip is None
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_multiple_fallback_ips_tried_in_order(self, monkeypatch):
|
||
calls = []
|
||
behavior = {
|
||
"api.telegram.org": "timeout",
|
||
"149.154.167.220": "timeout",
|
||
"149.154.167.221": "ok",
|
||
}
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", _fake_transport_factory(calls, behavior))
|
||
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220", "149.154.167.221"])
|
||
resp = await transport.handle_async_request(_telegram_request())
|
||
|
||
assert resp.status_code == 200
|
||
assert transport._sticky_ip == "149.154.167.221"
|
||
assert [c["url_host"] for c in calls] == [
|
||
"api.telegram.org",
|
||
"149.154.167.220",
|
||
"149.154.167.221",
|
||
]
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_sticky_ip_tried_first_but_falls_through_if_stale(self, monkeypatch):
|
||
"""If the sticky IP stops working, the transport retries others."""
|
||
calls = []
|
||
behavior = {
|
||
"api.telegram.org": "timeout",
|
||
"149.154.167.220": "ok",
|
||
"149.154.167.221": "ok",
|
||
}
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", _fake_transport_factory(calls, behavior))
|
||
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220", "149.154.167.221"])
|
||
|
||
# First request: primary fails → .220 works → becomes sticky
|
||
await transport.handle_async_request(_telegram_request())
|
||
assert transport._sticky_ip == "149.154.167.220"
|
||
|
||
# Now .220 goes bad too
|
||
calls.clear()
|
||
behavior["149.154.167.220"] = "timeout"
|
||
|
||
resp = await transport.handle_async_request(_telegram_request())
|
||
assert resp.status_code == 200
|
||
# After #24511: when sticky fails the transport also resets and
|
||
# re-tries the primary DNS path before falling through to other IPs.
|
||
# Path: sticky (.220) → primary (api.telegram.org) → .221
|
||
assert [c["url_host"] for c in calls] == ["149.154.167.220", "api.telegram.org", "149.154.167.221"]
|
||
assert transport._sticky_ip == "149.154.167.221"
|
||
|
||
|
||
class TestFallbackTransportPassthrough:
|
||
"""Requests that don't need fallback behavior."""
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_non_telegram_host_bypasses_fallback(self, monkeypatch):
|
||
calls = []
|
||
behavior = {}
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", _fake_transport_factory(calls, behavior))
|
||
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220"])
|
||
request = httpx.Request("GET", "https://example.com/path")
|
||
resp = await transport.handle_async_request(request)
|
||
|
||
assert resp.status_code == 200
|
||
assert calls[0]["url_host"] == "example.com"
|
||
assert transport._sticky_ip is None
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_empty_fallback_list_uses_primary_only(self, monkeypatch):
|
||
calls = []
|
||
behavior = {}
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", _fake_transport_factory(calls, behavior))
|
||
|
||
transport = tnet.TelegramFallbackTransport([])
|
||
resp = await transport.handle_async_request(_telegram_request())
|
||
|
||
assert resp.status_code == 200
|
||
assert calls[0]["url_host"] == "api.telegram.org"
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_primary_succeeds_no_fallback_needed(self, monkeypatch):
|
||
calls = []
|
||
behavior = {"api.telegram.org": "ok"}
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", _fake_transport_factory(calls, behavior))
|
||
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220"])
|
||
resp = await transport.handle_async_request(_telegram_request())
|
||
|
||
assert resp.status_code == 200
|
||
assert transport._sticky_ip is None
|
||
assert len(calls) == 1
|
||
|
||
|
||
class TestFallbackTransportInit:
|
||
def test_deduplicates_fallback_ips(self, monkeypatch):
|
||
monkeypatch.setattr(
|
||
tnet.httpx, "AsyncHTTPTransport", lambda **kw: FakeTransport([], {})
|
||
)
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220", "149.154.167.220"])
|
||
assert transport._fallback_ips == ["149.154.167.220"]
|
||
|
||
def test_filters_invalid_ips_at_init(self, monkeypatch):
|
||
monkeypatch.setattr(
|
||
tnet.httpx, "AsyncHTTPTransport", lambda **kw: FakeTransport([], {})
|
||
)
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220", "not-an-ip"])
|
||
assert transport._fallback_ips == ["149.154.167.220"]
|
||
|
||
def test_uses_proxy_env_for_primary_and_fallback_transports(self, monkeypatch):
|
||
seen_kwargs = []
|
||
|
||
def factory(**kwargs):
|
||
seen_kwargs.append(kwargs.copy())
|
||
return FakeTransport([], {})
|
||
|
||
for key in ("HTTPS_PROXY", "HTTP_PROXY", "ALL_PROXY", "https_proxy", "http_proxy", "all_proxy", "TELEGRAM_PROXY", "NO_PROXY", "no_proxy"):
|
||
monkeypatch.delenv(key, raising=False)
|
||
monkeypatch.setenv("HTTPS_PROXY", "http://proxy.example:8080")
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", factory)
|
||
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220"])
|
||
|
||
assert transport._fallback_ips == ["149.154.167.220"]
|
||
assert len(seen_kwargs) == 2
|
||
assert all(kwargs["proxy"] == "http://proxy.example:8080" for kwargs in seen_kwargs)
|
||
|
||
def test_no_proxy_bypasses_fallback_ip_cidr(self, monkeypatch):
|
||
seen_kwargs = []
|
||
|
||
def factory(**kwargs):
|
||
seen_kwargs.append(kwargs.copy())
|
||
return FakeTransport([], {})
|
||
|
||
for key in ("HTTPS_PROXY", "HTTP_PROXY", "ALL_PROXY", "https_proxy", "http_proxy", "all_proxy", "TELEGRAM_PROXY", "NO_PROXY", "no_proxy"):
|
||
monkeypatch.delenv(key, raising=False)
|
||
monkeypatch.setenv("HTTPS_PROXY", "http://proxy.example:8080")
|
||
monkeypatch.setenv("NO_PROXY", "149.154.160.0/20")
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", factory)
|
||
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220"])
|
||
|
||
assert transport._fallback_ips == ["149.154.167.220"]
|
||
assert len(seen_kwargs) == 2
|
||
assert all("proxy" not in kwargs for kwargs in seen_kwargs)
|
||
|
||
def test_forwards_limits_to_inner_transports(self, monkeypatch):
|
||
"""Verify that caller-supplied limits reach the inner
|
||
AsyncHTTPTransport instances (#58790). httpx ignores the
|
||
client-level limits kwarg when a custom transport is
|
||
supplied, so the limits must be forwarded via transport_kwargs.
|
||
"""
|
||
seen_kwargs = []
|
||
|
||
def factory(**kwargs):
|
||
seen_kwargs.append(kwargs.copy())
|
||
return FakeTransport([], {})
|
||
|
||
for key in ("HTTPS_PROXY", "HTTP_PROXY", "ALL_PROXY", "https_proxy", "http_proxy", "all_proxy", "TELEGRAM_PROXY", "NO_PROXY", "no_proxy"):
|
||
monkeypatch.delenv(key, raising=False)
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", factory)
|
||
|
||
custom_limits = httpx.Limits(
|
||
max_connections=42,
|
||
max_keepalive_connections=10,
|
||
keepalive_expiry=30.0,
|
||
)
|
||
transport = tnet.TelegramFallbackTransport(
|
||
["149.154.167.220"], limits=custom_limits
|
||
)
|
||
|
||
# 1 primary + 1 fallback = 2 AsyncHTTPTransport instances
|
||
assert len(seen_kwargs) == 2
|
||
for kw in seen_kwargs:
|
||
assert "limits" in kw
|
||
assert kw["limits"] is custom_limits
|
||
|
||
|
||
class TestFallbackTransportClose:
|
||
@pytest.mark.asyncio
|
||
async def test_aclose_closes_all_transports(self, monkeypatch):
|
||
factory = _fake_transport_factory([], {})
|
||
monkeypatch.setattr(tnet.httpx, "AsyncHTTPTransport", factory)
|
||
|
||
transport = tnet.TelegramFallbackTransport(["149.154.167.220", "149.154.167.221"])
|
||
await transport.aclose()
|
||
|
||
# 1 primary + 2 fallback transports
|
||
assert len(factory.instances) == 3
|
||
assert all(t.closed for t in factory.instances)
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
# Config layer – TELEGRAM_FALLBACK_IPS env → config.extra
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
|
||
class TestConfigFallbackIps:
|
||
def test_env_var_populates_config_extra(self, monkeypatch):
|
||
from gateway.config import GatewayConfig, Platform, PlatformConfig, _apply_env_overrides
|
||
|
||
monkeypatch.setenv("TELEGRAM_FALLBACK_IPS", "149.154.167.220,149.154.167.221")
|
||
config = GatewayConfig(platforms={Platform.TELEGRAM: PlatformConfig(enabled=True, token="tok")})
|
||
_apply_env_overrides(config)
|
||
|
||
assert config.platforms[Platform.TELEGRAM].extra["fallback_ips"] == [
|
||
"149.154.167.220", "149.154.167.221",
|
||
]
|
||
|
||
def test_env_var_creates_platform_if_missing(self, monkeypatch):
|
||
from gateway.config import GatewayConfig, Platform, _apply_env_overrides
|
||
|
||
monkeypatch.setenv("TELEGRAM_FALLBACK_IPS", "149.154.167.220")
|
||
config = GatewayConfig(platforms={})
|
||
_apply_env_overrides(config)
|
||
|
||
assert Platform.TELEGRAM in config.platforms
|
||
assert config.platforms[Platform.TELEGRAM].extra["fallback_ips"] == ["149.154.167.220"]
|
||
|
||
def test_env_var_strips_whitespace(self, monkeypatch):
|
||
from gateway.config import GatewayConfig, Platform, PlatformConfig, _apply_env_overrides
|
||
|
||
monkeypatch.setenv("TELEGRAM_FALLBACK_IPS", " 149.154.167.220 , 149.154.167.221 ")
|
||
config = GatewayConfig(platforms={Platform.TELEGRAM: PlatformConfig(enabled=True, token="tok")})
|
||
_apply_env_overrides(config)
|
||
|
||
assert config.platforms[Platform.TELEGRAM].extra["fallback_ips"] == [
|
||
"149.154.167.220", "149.154.167.221",
|
||
]
|
||
|
||
def test_empty_env_var_does_not_populate(self, monkeypatch):
|
||
from gateway.config import GatewayConfig, Platform, PlatformConfig, _apply_env_overrides
|
||
|
||
monkeypatch.setenv("TELEGRAM_FALLBACK_IPS", "")
|
||
config = GatewayConfig(platforms={Platform.TELEGRAM: PlatformConfig(enabled=True, token="tok")})
|
||
_apply_env_overrides(config)
|
||
|
||
assert "fallback_ips" not in config.platforms[Platform.TELEGRAM].extra
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
# Adapter layer – _fallback_ips() reads config correctly
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
|
||
class TestAdapterFallbackIps:
|
||
def _make_adapter(self, extra=None):
|
||
import sys
|
||
from unittest.mock import MagicMock
|
||
|
||
# Ensure telegram mock is in place
|
||
if "telegram" not in sys.modules or not hasattr(sys.modules["telegram"], "__file__"):
|
||
mod = MagicMock()
|
||
mod.ext.ContextTypes.DEFAULT_TYPE = type(None)
|
||
mod.constants.ParseMode.MARKDOWN_V2 = "MarkdownV2"
|
||
mod.constants.ChatType.GROUP = "group"
|
||
mod.constants.ChatType.SUPERGROUP = "supergroup"
|
||
mod.constants.ChatType.CHANNEL = "channel"
|
||
mod.constants.ChatType.PRIVATE = "private"
|
||
for name in ("telegram", "telegram.ext", "telegram.constants", "telegram.request"):
|
||
sys.modules.setdefault(name, mod)
|
||
|
||
from gateway.config import PlatformConfig
|
||
from plugins.platforms.telegram.adapter import TelegramAdapter
|
||
|
||
config = PlatformConfig(enabled=True, token="test-token")
|
||
if extra:
|
||
config.extra.update(extra)
|
||
return TelegramAdapter(config)
|
||
|
||
def test_list_in_extra(self):
|
||
adapter = self._make_adapter(extra={"fallback_ips": ["149.154.167.220"]})
|
||
assert adapter._fallback_ips() == ["149.154.167.220"]
|
||
|
||
def test_csv_string_in_extra(self):
|
||
adapter = self._make_adapter(extra={"fallback_ips": "149.154.167.220,149.154.167.221"})
|
||
assert adapter._fallback_ips() == ["149.154.167.220", "149.154.167.221"]
|
||
|
||
def test_empty_extra(self):
|
||
adapter = self._make_adapter()
|
||
assert adapter._fallback_ips() == []
|
||
|
||
def test_no_extra_attr(self):
|
||
adapter = self._make_adapter()
|
||
adapter.config.extra = None
|
||
assert adapter._fallback_ips() == []
|
||
|
||
def test_invalid_ips_filtered(self):
|
||
adapter = self._make_adapter(extra={"fallback_ips": ["149.154.167.220", "not-valid"]})
|
||
assert adapter._fallback_ips() == ["149.154.167.220"]
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
# DoH auto-discovery
|
||
# ═══════════════════════════════════════════════════════════════════════════
|
||
|
||
def _doh_answer(*ips: str) -> dict:
|
||
"""Build a minimal DoH JSON response with A records."""
|
||
return {"Answer": [{"type": 1, "data": ip} for ip in ips]}
|
||
|
||
|
||
class FakeDoHClient:
|
||
"""Mock httpx.AsyncClient for DoH queries."""
|
||
|
||
def __init__(self, responses: dict):
|
||
# responses: URL prefix → (status, json_body) | Exception
|
||
self._responses = responses
|
||
self.requests_made: list[dict] = []
|
||
|
||
@staticmethod
|
||
def _make_response(status, body, url):
|
||
"""Build an httpx.Response with a request attached (needed for raise_for_status)."""
|
||
request = httpx.Request("GET", url)
|
||
return httpx.Response(status, json=body, request=request)
|
||
|
||
async def get(self, url, *, params=None, headers=None, **kwargs):
|
||
self.requests_made.append({"url": url, "params": params, "headers": headers})
|
||
for prefix, action in self._responses.items():
|
||
if url.startswith(prefix):
|
||
if isinstance(action, Exception):
|
||
raise action
|
||
status, body = action
|
||
return self._make_response(status, body, url)
|
||
return self._make_response(200, {}, url)
|
||
|
||
async def __aenter__(self):
|
||
return self
|
||
|
||
async def __aexit__(self, *args):
|
||
pass
|
||
|
||
|
||
class TestDiscoverFallbackIps:
|
||
"""Tests for discover_fallback_ips() — DoH-based auto-discovery."""
|
||
|
||
def _patch_doh(self, monkeypatch, responses, system_dns_ips=None):
|
||
"""Wire up fake DoH client and system DNS."""
|
||
client = FakeDoHClient(responses)
|
||
monkeypatch.setattr(tnet.httpx, "AsyncClient", lambda **kw: client)
|
||
|
||
if system_dns_ips is not None:
|
||
addrs = [(None, None, None, None, (ip, 443)) for ip in system_dns_ips]
|
||
monkeypatch.setattr(tnet.socket, "getaddrinfo", lambda *a, **kw: addrs)
|
||
else:
|
||
def _fail(*a, **kw):
|
||
raise OSError("dns failed")
|
||
monkeypatch.setattr(tnet.socket, "getaddrinfo", _fail)
|
||
return client
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_google_and_cloudflare_ips_collected(self, monkeypatch):
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": (200, _doh_answer("149.154.167.220")),
|
||
"https://cloudflare-dns.com": (200, _doh_answer("149.154.167.221")),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
|
||
ips = await tnet.discover_fallback_ips()
|
||
assert "149.154.167.220" in ips
|
||
assert "149.154.167.221" in ips
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_system_dns_ip_kept_when_doh_confirms(self, monkeypatch):
|
||
"""DoH-confirmed IPs are kept even when they match system DNS (#14520).
|
||
|
||
The system-DNS IP is often the most reliable path; including it as a
|
||
fallback lets the IP-rewrite retry recover from transient primary-path
|
||
failures instead of jumping straight to the hardcoded seed list.
|
||
"""
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": (200, _doh_answer("149.154.166.110", "149.154.167.220")),
|
||
"https://cloudflare-dns.com": (200, _doh_answer("149.154.166.110")),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
|
||
ips = await tnet.discover_fallback_ips()
|
||
assert ips == ["149.154.166.110", "149.154.167.220"]
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_doh_results_deduplicated(self, monkeypatch):
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": (200, _doh_answer("149.154.167.220")),
|
||
"https://cloudflare-dns.com": (200, _doh_answer("149.154.167.220")),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
|
||
ips = await tnet.discover_fallback_ips()
|
||
assert ips == ["149.154.167.220"]
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_doh_timeout_falls_back_to_seed(self, monkeypatch):
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": httpx.TimeoutException("timeout"),
|
||
"https://cloudflare-dns.com": httpx.TimeoutException("timeout"),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
|
||
ips = await tnet.discover_fallback_ips()
|
||
assert ips == tnet._SEED_FALLBACK_IPS
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_doh_connect_error_falls_back_to_seed(self, monkeypatch):
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": httpx.ConnectError("refused"),
|
||
"https://cloudflare-dns.com": httpx.ConnectError("refused"),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
|
||
ips = await tnet.discover_fallback_ips()
|
||
assert ips == tnet._SEED_FALLBACK_IPS
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_doh_malformed_json_falls_back_to_seed(self, monkeypatch):
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": (200, {"Status": 0}), # no Answer key
|
||
"https://cloudflare-dns.com": (200, {"garbage": True}),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
|
||
ips = await tnet.discover_fallback_ips()
|
||
assert ips == tnet._SEED_FALLBACK_IPS
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_one_provider_fails_other_succeeds(self, monkeypatch):
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": httpx.TimeoutException("timeout"),
|
||
"https://cloudflare-dns.com": (200, _doh_answer("149.154.167.220")),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
|
||
ips = await tnet.discover_fallback_ips()
|
||
assert ips == ["149.154.167.220"]
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_system_dns_failure_keeps_all_doh_ips(self, monkeypatch):
|
||
"""If system DNS fails, nothing gets excluded — all DoH IPs kept."""
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": (200, _doh_answer("149.154.166.110", "149.154.167.220")),
|
||
"https://cloudflare-dns.com": (200, _doh_answer()),
|
||
}, system_dns_ips=None) # triggers OSError
|
||
|
||
ips = await tnet.discover_fallback_ips()
|
||
assert "149.154.166.110" in ips
|
||
assert "149.154.167.220" in ips
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_all_doh_ips_same_as_system_dns_kept(self, monkeypatch):
|
||
"""DoH agrees with system DNS — keep that IP instead of seed list (#14520).
|
||
|
||
Previous behavior fell through to ``_SEED_FALLBACK_IPS`` here, but the
|
||
seed addresses are not routable on every network. When DoH confirms
|
||
the system IP, that IP is the best candidate we have and should be
|
||
used as the fallback target.
|
||
"""
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": (200, _doh_answer("149.154.166.110")),
|
||
"https://cloudflare-dns.com": (200, _doh_answer("149.154.166.110")),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
|
||
ips = await tnet.discover_fallback_ips()
|
||
assert ips == ["149.154.166.110"]
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_cloudflare_gets_accept_header(self, monkeypatch):
|
||
client = self._patch_doh(monkeypatch, {
|
||
"https://dns.google": (200, _doh_answer("149.154.167.220")),
|
||
"https://cloudflare-dns.com": (200, _doh_answer("149.154.167.221")),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
|
||
await tnet.discover_fallback_ips()
|
||
|
||
cf_reqs = [r for r in client.requests_made if "cloudflare" in r["url"]]
|
||
assert cf_reqs
|
||
assert cf_reqs[0]["headers"]["Accept"] == "application/dns-json"
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_non_a_records_ignored(self, monkeypatch):
|
||
"""AAAA records (type 28) and CNAME (type 5) should be skipped."""
|
||
answer = {
|
||
"Answer": [
|
||
{"type": 5, "data": "telegram.org"}, # CNAME
|
||
{"type": 28, "data": "2001:67c:4e8:f004::9"}, # AAAA
|
||
{"type": 1, "data": "149.154.167.220"}, # A ✓
|
||
]
|
||
}
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": (200, answer),
|
||
"https://cloudflare-dns.com": (200, _doh_answer()),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
|
||
ips = await tnet.discover_fallback_ips()
|
||
assert ips == ["149.154.167.220"]
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_invalid_ip_in_doh_response_skipped(self, monkeypatch):
|
||
answer = {"Answer": [
|
||
{"type": 1, "data": "not-an-ip"},
|
||
{"type": 1, "data": "149.154.167.220"},
|
||
]}
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": (200, answer),
|
||
"https://cloudflare-dns.com": (200, _doh_answer()),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
|
||
ips = await tnet.discover_fallback_ips()
|
||
assert ips == ["149.154.167.220"]
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_hung_system_dns_does_not_gate_doh_results(self, monkeypatch):
|
||
"""#63309: socket.getaddrinfo has no timeout of its own — a wedged OS
|
||
resolver must not stall discovery. DoH answers must come back promptly
|
||
even while the system-DNS worker thread is still hanging."""
|
||
import time as _time
|
||
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": (200, _doh_answer("149.154.167.220")),
|
||
"https://cloudflare-dns.com": (200, _doh_answer()),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
monkeypatch.setattr(tnet, "_DOH_TIMEOUT", 0.2)
|
||
|
||
def _hung_getaddrinfo(*a, **kw):
|
||
_time.sleep(1.5) # far beyond the discovery bound
|
||
raise OSError("resolver wedged")
|
||
|
||
monkeypatch.setattr(tnet.socket, "getaddrinfo", _hung_getaddrinfo)
|
||
|
||
start = _time.monotonic()
|
||
ips = await tnet.discover_fallback_ips()
|
||
elapsed = _time.monotonic() - start
|
||
|
||
assert ips == ["149.154.167.220"]
|
||
assert elapsed < 1.4, f"discovery gated on hung system DNS ({elapsed:.2f}s)"
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_hung_system_dns_with_no_doh_answers_bounded_seed_fallback(self, monkeypatch):
|
||
"""Worst case — resolver wedged AND no DoH answers — must still return
|
||
the seed list within the bound instead of hanging connect()."""
|
||
import time as _time
|
||
|
||
self._patch_doh(monkeypatch, {
|
||
"https://dns.google": (200, {"Status": 0}),
|
||
"https://cloudflare-dns.com": (200, {"garbage": True}),
|
||
}, system_dns_ips=["149.154.166.110"])
|
||
monkeypatch.setattr(tnet, "_DOH_TIMEOUT", 0.2)
|
||
|
||
def _hung_getaddrinfo(*a, **kw):
|
||
_time.sleep(1.5)
|
||
raise OSError("resolver wedged")
|
||
|
||
monkeypatch.setattr(tnet.socket, "getaddrinfo", _hung_getaddrinfo)
|
||
|
||
start = _time.monotonic()
|
||
ips = await tnet.discover_fallback_ips()
|
||
elapsed = _time.monotonic() - start
|
||
|
||
assert ips == tnet._SEED_FALLBACK_IPS
|
||
assert elapsed < 1.4, f"seed fallback gated on hung system DNS ({elapsed:.2f}s)"
|