From 6951c5ce6a9dceb1e9924ae69f3510696c07b99d Mon Sep 17 00:00:00 2001 From: axisrow Date: Thu, 24 Sep 2026 15:35:00 +0800 Subject: [PATCH 1/3] feat(telegram): calibrate the gate's history category from production logs (#1418) The released telethon-floodgate 0.1.0 default (600/min) never bound: the production app.log shows 209 FLOOD_WAITs on messages.getHistory with collector peaks of 115 channel fetches per minute, far below that guard. Pass the empirically measured boundary (30 req/~30s, the 31st returned FLOOD_WAIT_3, minus a 20% margin) as a category_limits override when the pool constructs its gate. Categories with no flood signal in the logs (send, admin_action, channel_lifecycle, default) keep the package defaults and their needs-calibration note. Regression tests: the pool ships the calibrated spec, and a peak burst defers the 25th history call inside the 30s window (fake clock). Part of #1331 Closes #1418 Co-Authored-By: Claude Code --- src/telegram/client_pool.py | 33 +++++-- tcf-1418-calib-report.md | 158 ++++++++++++++++++++++++++++++++++ tests/test_rate_limit_gate.py | 40 +++++++++ 3 files changed, 226 insertions(+), 5 deletions(-) create mode 100644 tcf-1418-calib-report.md diff --git a/src/telegram/client_pool.py b/src/telegram/client_pool.py index a1f70523..8bcc5ad8 100644 --- a/src/telegram/client_pool.py +++ b/src/telegram/client_pool.py @@ -57,7 +57,12 @@ # #1046 split. The live call sites now live in the mixin modules, which import # these names into their own namespaces. from telethon.tl.types import ChannelForbidden # noqa: F401 -from telethon_floodgate import FloodCircuitBreaker, ResolveRateLimiter, TelegramRateLimitGate +from telethon_floodgate import ( + FloodCircuitBreaker, + RateLimitSpec, + ResolveRateLimiter, + TelegramRateLimitGate, +) from src.config import TelegramRuntimeConfig from src.database import Database @@ -91,6 +96,20 @@ # duplicated here — nothing reads ``client_pool.`` any more (#1046 cleanup). +# Phase 2 calibration of the gate's ``history`` category (#1418, epic #1331). +# The released telethon-floodgate 0.1.0 default (600/min) never bound: the +# production app.log shows 209 FLOOD_WAITs on messages.getHistory with +# collector peaks of 115 channel fetches per minute, far below that guard. +# 24/30s is the empirically measured Telegram boundary (30 requests in ~30s, +# the 31st returned FLOOD_WAIT_3) with a 20% margin. Peak collector bursts +# (p95 74-101 fetches/min) stretch instead of flooding; incremental min_id +# collection catches deferred channels on the next pass. Other categories +# had no flood signal in the logs and keep the package defaults pending a +# larger sample. Keep in sync with the package default; drop the override +# once a released telethon-floodgate ships this value (#1418). +HISTORY_CALIBRATED_SPEC = RateLimitSpec(max_calls=24, window_sec=30.0) + + @dataclass(frozen=True) class StatsClientAvailability: state: str # "available" | "all_flooded" | "no_connected_active" @@ -172,10 +191,14 @@ def __init__( self._dialog_refresh_tasks: dict[tuple[str, str], asyncio.Task[list[dict]]] = {} self._premium_flood_wait_until: dict[str, datetime] = {} self._resolve_rate_limiter = ResolveRateLimiter() - # Central proactive gate. Category limits are conservative operating - # defaults calibrated from the available production signals; keep the - # registry injectable for future recalibration (#1331). - self._rate_limit_gate = TelegramRateLimitGate() + # Central proactive gate. ``history`` is calibrated from the + # production log sample (#1418); the remaining categories are the + # package's conservative operating defaults pending production + # signals — keep the registry injectable for future recalibration + # (#1331). + self._rate_limit_gate = TelegramRateLimitGate( + category_limits={"history": HISTORY_CALIBRATED_SPEC}, + ) # Reactive counterpart to the gate (#1330/#1368): the gate paces calls # against guessed limits, the breaker stops an (operation, phone) pair # that Telegram is already flood-waiting instead of hammering on. diff --git a/tcf-1418-calib-report.md b/tcf-1418-calib-report.md new file mode 100644 index 00000000..f51eca71 --- /dev/null +++ b/tcf-1418-calib-report.md @@ -0,0 +1,158 @@ +# #1418: калибровка rate_limit_gate по прод-логам — отчёт (Phase 2 эпика #1331) + +Дата: 2026-09-24. Источник данных: `data/app.log` (основной чекаут, +10.4 MB, 2026-06-12 … 2026-09-01). Живой Telegram не использовался, код +пакета не менялся. Метод: read-only разбор лога + калибровка донора конфигом +пакета (сумма работ #1432/#1434 вынесла стек в PyPI `telethon-floodgate`). + +## 1. Методика + +Скрипт `/tmp/tcf1418_calib.py` + `/tmp/tcf1418_deep.py` (одноразовые, текст +алгоритма ниже): + +1. Строки `Flood wait for : N seconds` (reactive, pool_flood) — 2734 шт. +2. Строки `: transient Flood wait Ns until ... UTC for ` + (proactive, с тегом операции). +3. Матчинг пар по (минута, телефон) ±1 минута: **2666 из 2734 спарены (97.5 %)**. +4. Нагрузка коллектора: строки `Collecting channel … account=` → + частота на (аккаунт, минута); контекст флуда — число Collecting-строк + за 60 с до него. + +## 2. Что показал лог + +### 2.1 Распределение по операциям (2666 спаренных) + +| Операция | Флудов | Доля | Категория gate | +|---|---|---|---| +| `telegram_warm_dialog_cache` | 2446 | 91.7 % | `dialogs` | +| `telegram_stream_messages` | 209 | 7.8 % | `history` | +| `telegram_stream_dialogs` | 8 | 0.3 % | `dialog_sweep` | +| `telegram_invoke_request` | 3 | 0.1 % | `default` | +| `send` / `admin_action` / `channel_lifecycle` | **0** | — | — | + +### 2.2 Длительности + +- Спаренные (n=2666): медиана 29 с, p90 29 с, p99 30 с, max 30 с — это + транс-пейсер Telethon при пагинации (27–30 с), «мягкое» замедление. +- `stream_messages` (n=209): медиана **14 с**, p90 25 с, max 27 с. +- **Все тяжёлые флуды (≥120 с) — в 68 неспаренных**: медиана 120 с, p90 600 с, + max 53123 с (бан 14.8 ч, +66...2247, 2026-08-26 21:00). По месяцам: + 2026-06 — 4, 2026-08 — 30, 2026-09 — 33. Значительная часть на + тест-телефоне `+70...0001` (600 с, серийные) — тестовый шум в прод-логе; + реальные тяжёлые: 28.06 (+86...9509, 782/691 с), 26.08 (бан 53123 с). + +### 2.3 Динамика (Phase 1 работает) + +| Операция | июнь | август | +|---|---|---| +| `warm_dialog_cache` | 2391 | **55** (падение в 43×) | +| `stream_messages` | 95 | **114** — но 40+74 из них 26–27.08, т.е. день бана и день после | + +Аккаунты: `+66...2247` — 2482 (93 %), далее 75/64/45. `stream_messages` +кластеризуется слабо: 156 окон × 1 флуд, 25 × 2, 1 × 3. + +### 2.4 Нагрузка коллектора (Collecting-строк на аккаунт-минуту) + +| Аккаунт | минут | медиана | p95 | max | минут >48/мин | +|---|---|---|---|---|---| +| +66...2629 | 144 | 44 | 97 | 104 | 66 | +| +66...2531 | 143 | 50 | 101 | 115 | 73 | +| +86...9509 | 105 | 40 | 74 | 92 | 45 | +| +66...2247 | 60 | 26 | 101 | 103 | 32 | + +Перед `stream_messages`-флудом коллектор делал 55–119 сборов/мин (мода 59). +**Вывод: пиковая реальная нагрузка 92–115/мин; старый guard `history` 600/мин +никогда не срабатывал** — медианные активные минуты 26–50/мин, пики в 5–8 раз +ниже лимита. 209 флудов — плата именно за это: гейта фактически не было. + +### 2.5 Спецвопрос issue №4 (history 600/мин) — ответ + +Лимит не «слишком высок», он **не срабатывал вовсе**: максимум наблюдаемой +нагрузки (115/мин) на 5× ниже 600/мин. Действующая граница Telegram +измерена прямо (калибровочный прогон пакета, см. §4): 30 запросов/~30 с, +31-й → FLOOD_WAIT_3. Лог-данные (медиана флуда 14 с при 95–119 сборах/мин) +с этой границей согласуются. + +### 2.6 Спецвопрос issue №5 (warm_dialog_cache под гейтом?) — ответ + +Все вызовы `warm_dialog_cache` в `src/` идут через **один** transport-метод +`TelegramTransportSession.warm_dialog_cache()` (`backends.py:402`), который +сам резервирует слот гейта с тегом `telegram_warm_dialog_cache`. Call-сайтов +после декомпозиции #1046 — 10 (`pool_dialogs` ×5, `collector_mixins` ×3, +`telegram_search`, CLI ×2), каждый передаёт составной тег +(`collect_channel_warm_dialog_cache`, …), который матчится суффиксным +правилом `endswith("_warm_dialog_cache")` — регресс класса #1336 закрыт +в пакете тестом. Т.е. **непокрытых путей нет**; падение 2391→55 — эффект +связки `dialogs` 1/60 + персистентный `dialog_cache`-skip, а не утечки мимо +гейта. 85 % флудов июня — доисторическая эпоха до Phase 1 (#1330). + +## 3. Калибровочная таблица + +Прод-пакет = PyPI `telethon-floodgate` 0.1.0 (pin `>=0.1.0,<0.2`). + +| Категория | Было (прод 0.1.0) | Стало | Обоснование | +|---|---|---|---| +| `history` | 600/60 с — догадка, никогда не связывал | **24/30 с** | Эмпирическая граница Telegram (30 req/30 с, 31-й → FLOOD_WAIT_3) минус 20 % маржа; лог: 209 флудов при пиковой нагрузке 115/мин << 600/мин | +| `dialogs` | 1/60 с | без изменений | #1330; лог подтверждает: 2419→55 флудов warm после введения | +| `dialog_sweep` | 12/60 с | без изменений | #1359; 8 флудов за 3 мес — режим штатный | +| `send` | 30/60 с (+per-peer 1/1.0 с, 20/60) | без изменений, **требует калибровки** | 0 флудов — данных нет; per-peer-правки (1/1.1, 16/60) уже в локальном пакете по live-замеру | +| `admin_action` | 10/60 с | без изменений, **требует калибровки** | 0 флудов — данных нет | +| `channel_lifecycle` | 3/300 с | без изменений, **требует калибровки** | 0 флудов — данных нет | +| `default` | 1000/60 с | без изменений, **требует калибровки** | 3 единичных флуда `invoke_request` за 3 мес — против catch-all-страховки данных нет | + +Дифф (минимальный): `client_pool.py` передаёт гейту +`category_limits={"history": RateLimitSpec(24, 30.0)}` — константа +`HISTORY_CALIBRATED_SPEC` + регресс-тесты (`tests/test_rate_limit_gate.py`): +пул применяет спеку; 25-й вызов в 30-секундном окне деферится до окна +(фейковые часы `_Clock`). Остальные категории — дефолты пакета, пометка +«требует калибровки» живёт в комментарии у конструирования гейта и в этой +таблице. + +## 4. Координация с пакетом (важно) + +В локальной копии `~/Projects/telethon-floodgate` (main, uncommitted) уже +лежит ровно та же правка `HISTORY_SPEC = 24/30` (комментарий «empirical +boundary») — на PyPI она **не выпущена**, поэтому прод с 0.1.0 живёт с +600/60. Донорский оверрайд закрывает прод уже сейчас; после релиза пакета +с этим дефолтом оверрайд можно снять (значения совпадут) — это записано в +комментарии константы. + +## 5. Пропускная способность сбора (критерий AC) + +Живой замер до/после невозможен (запрещены живые вызовы) — оценка по логу: + +- Медианные активные минуты (26–50 сборов/мин) — **ниже 48/мин** (24/30 с), + дефера нет вообще, поведение не меняется. +- Пики (p95 74–101, max 115) — дефер; канал откладывается, а не теряется: + сбор инкрементальный (`min_id`), следующий проход добирает. Дефер warm- + prefetch (`collection.py:491`) вообще не останавливает сбор канала. +- Вместо пиков с флудами (каждый флуд = пауза 14–29 с + риск эскалации до + 120 с/бана) пики растягиваются гейтом до ~48 вызовов/30 с. Полных потерь + пропускной способности нет; платформа-цена — задержка добора пиковых + каналов до следующего прохода. + +## 6. Связка с #1419 (рестарты) + +Замер #1419 дополнительно показывает: 44 % стартов ловят флуд в первые 60 с +(~19× к базовой), рестарты 0.38/день кластеризуются, breaker в проде ни разу +не открывался (деплой после бан-инцидента). Данный PR калибровку лимитов +делает строже, что снижает и послемстартовый залп (gate-бюджеты после +рестарта пустые), но саму персистентность breaker/gate-состояний не решает — +это отдельное согласование (§8 отчёта #1419). + +## 7. Проверки + +- `ruff check src/telegram/client_pool.py tests/test_rate_limit_gate.py` — чисто. +- `tests/test_rate_limit_gate.py` — 13 passed (11 старых + 2 новых). +- Полный сюит: параллельная (`-m "not aiosqlite_serial" -n auto`) и + serial (`-m aiosqlite_serial -n auto --dist=loadfile`) части — см. PR. + +## 8. Риски / follow-ups + +- «Пропускная способность не деградирует» подтверждена оценкой по логу, не + живым замером — если понадобится живая проверка, делать по протоколу + калибровщика пакета (`scripts/calibrate_send_limits.py`), burnable-аккаунт. +- Релиз `telethon-floodgate` 0.1.1 с локальными правками владельца → снять + донорский оверрайд. +- Персистентность breaker-состояний (#1419 §8) — открытый вектор бан-масштаба + при рестарт-циклах. diff --git a/tests/test_rate_limit_gate.py b/tests/test_rate_limit_gate.py index 3272550e..a63524da 100644 --- a/tests/test_rate_limit_gate.py +++ b/tests/test_rate_limit_gate.py @@ -1,6 +1,7 @@ from __future__ import annotations from types import SimpleNamespace +from unittest.mock import MagicMock import pytest from telethon import TelegramClient @@ -16,6 +17,7 @@ ) from src.telegram.backends import TelegramTransportSession +from src.telegram.client_pool import HISTORY_CALIBRATED_SPEC class _Clock: @@ -429,3 +431,41 @@ class Pool: assert await session.send_message(12345, "b") == "ok" # An entity peer_key cannot read at all yields None and also proceeds. assert await session.send_message(object(), "c") == "ok" + + +def test_pool_applies_production_history_calibration() -> None: + """The pool's gate must ship the #1418 calibrated history spec. + + The released floodgate default (600/min) never bound: production peaks ran + at 115 fetches/min. The override lives in the pool, so the pool is what + this pins; the remaining categories must stay untouched package defaults. + """ + from src.telegram.client_pool import ClientPool + + pool = ClientPool(MagicMock(api_id=1, api_hash="h"), MagicMock()) + history = pool._rate_limit_gate._limiters["history"] + assert history._max_calls == 24 + assert history._window_sec == 30.0 + # Untouched categories keep their package defaults. + send = pool._rate_limit_gate._limiters["send"] + assert (send._max_calls, send._window_sec) == (30, 60.0) + + +def test_history_calibration_stops_a_collector_burst_before_telegram() -> None: + """A peak burst (p95 74-101 fetches/min in app.log) defers, not floods. + + 24 calls fit the calibrated 30s window; the 25th is refused before any + Telegram call, and the budget returns once the window slides. + """ + clock = _Clock() + gate = TelegramRateLimitGate( + category_limits={"history": HISTORY_CALIBRATED_SPEC}, + time_func=clock, + ) + for _ in range(HISTORY_CALIBRATED_SPEC.max_calls): + assert gate.try_acquire("+7001", "history") == 0.0 + deferred = gate.try_acquire("+7001", "history") + assert deferred > 0 + # Sliding window: after the window passes, the bucket accepts again. + clock.now += HISTORY_CALIBRATED_SPEC.window_sec + assert gate.try_acquire("+7001", "history") == 0.0 From b298b82c11bb62db69e5f038428984db8a46304a Mon Sep 17 00:00:00 2001 From: axisrow Date: Thu, 24 Sep 2026 15:43:56 +0800 Subject: [PATCH 2/3] chore: move the calibration report out of the repo (#1441 followup) Co-Authored-By: Claude Code --- tcf-1418-calib-report.md | 158 --------------------------------------- 1 file changed, 158 deletions(-) delete mode 100644 tcf-1418-calib-report.md diff --git a/tcf-1418-calib-report.md b/tcf-1418-calib-report.md deleted file mode 100644 index f51eca71..00000000 --- a/tcf-1418-calib-report.md +++ /dev/null @@ -1,158 +0,0 @@ -# #1418: калибровка rate_limit_gate по прод-логам — отчёт (Phase 2 эпика #1331) - -Дата: 2026-09-24. Источник данных: `data/app.log` (основной чекаут, -10.4 MB, 2026-06-12 … 2026-09-01). Живой Telegram не использовался, код -пакета не менялся. Метод: read-only разбор лога + калибровка донора конфигом -пакета (сумма работ #1432/#1434 вынесла стек в PyPI `telethon-floodgate`). - -## 1. Методика - -Скрипт `/tmp/tcf1418_calib.py` + `/tmp/tcf1418_deep.py` (одноразовые, текст -алгоритма ниже): - -1. Строки `Flood wait for : N seconds` (reactive, pool_flood) — 2734 шт. -2. Строки `: transient Flood wait Ns until ... UTC for ` - (proactive, с тегом операции). -3. Матчинг пар по (минута, телефон) ±1 минута: **2666 из 2734 спарены (97.5 %)**. -4. Нагрузка коллектора: строки `Collecting channel … account=` → - частота на (аккаунт, минута); контекст флуда — число Collecting-строк - за 60 с до него. - -## 2. Что показал лог - -### 2.1 Распределение по операциям (2666 спаренных) - -| Операция | Флудов | Доля | Категория gate | -|---|---|---|---| -| `telegram_warm_dialog_cache` | 2446 | 91.7 % | `dialogs` | -| `telegram_stream_messages` | 209 | 7.8 % | `history` | -| `telegram_stream_dialogs` | 8 | 0.3 % | `dialog_sweep` | -| `telegram_invoke_request` | 3 | 0.1 % | `default` | -| `send` / `admin_action` / `channel_lifecycle` | **0** | — | — | - -### 2.2 Длительности - -- Спаренные (n=2666): медиана 29 с, p90 29 с, p99 30 с, max 30 с — это - транс-пейсер Telethon при пагинации (27–30 с), «мягкое» замедление. -- `stream_messages` (n=209): медиана **14 с**, p90 25 с, max 27 с. -- **Все тяжёлые флуды (≥120 с) — в 68 неспаренных**: медиана 120 с, p90 600 с, - max 53123 с (бан 14.8 ч, +66...2247, 2026-08-26 21:00). По месяцам: - 2026-06 — 4, 2026-08 — 30, 2026-09 — 33. Значительная часть на - тест-телефоне `+70...0001` (600 с, серийные) — тестовый шум в прод-логе; - реальные тяжёлые: 28.06 (+86...9509, 782/691 с), 26.08 (бан 53123 с). - -### 2.3 Динамика (Phase 1 работает) - -| Операция | июнь | август | -|---|---|---| -| `warm_dialog_cache` | 2391 | **55** (падение в 43×) | -| `stream_messages` | 95 | **114** — но 40+74 из них 26–27.08, т.е. день бана и день после | - -Аккаунты: `+66...2247` — 2482 (93 %), далее 75/64/45. `stream_messages` -кластеризуется слабо: 156 окон × 1 флуд, 25 × 2, 1 × 3. - -### 2.4 Нагрузка коллектора (Collecting-строк на аккаунт-минуту) - -| Аккаунт | минут | медиана | p95 | max | минут >48/мин | -|---|---|---|---|---|---| -| +66...2629 | 144 | 44 | 97 | 104 | 66 | -| +66...2531 | 143 | 50 | 101 | 115 | 73 | -| +86...9509 | 105 | 40 | 74 | 92 | 45 | -| +66...2247 | 60 | 26 | 101 | 103 | 32 | - -Перед `stream_messages`-флудом коллектор делал 55–119 сборов/мин (мода 59). -**Вывод: пиковая реальная нагрузка 92–115/мин; старый guard `history` 600/мин -никогда не срабатывал** — медианные активные минуты 26–50/мин, пики в 5–8 раз -ниже лимита. 209 флудов — плата именно за это: гейта фактически не было. - -### 2.5 Спецвопрос issue №4 (history 600/мин) — ответ - -Лимит не «слишком высок», он **не срабатывал вовсе**: максимум наблюдаемой -нагрузки (115/мин) на 5× ниже 600/мин. Действующая граница Telegram -измерена прямо (калибровочный прогон пакета, см. §4): 30 запросов/~30 с, -31-й → FLOOD_WAIT_3. Лог-данные (медиана флуда 14 с при 95–119 сборах/мин) -с этой границей согласуются. - -### 2.6 Спецвопрос issue №5 (warm_dialog_cache под гейтом?) — ответ - -Все вызовы `warm_dialog_cache` в `src/` идут через **один** transport-метод -`TelegramTransportSession.warm_dialog_cache()` (`backends.py:402`), который -сам резервирует слот гейта с тегом `telegram_warm_dialog_cache`. Call-сайтов -после декомпозиции #1046 — 10 (`pool_dialogs` ×5, `collector_mixins` ×3, -`telegram_search`, CLI ×2), каждый передаёт составной тег -(`collect_channel_warm_dialog_cache`, …), который матчится суффиксным -правилом `endswith("_warm_dialog_cache")` — регресс класса #1336 закрыт -в пакете тестом. Т.е. **непокрытых путей нет**; падение 2391→55 — эффект -связки `dialogs` 1/60 + персистентный `dialog_cache`-skip, а не утечки мимо -гейта. 85 % флудов июня — доисторическая эпоха до Phase 1 (#1330). - -## 3. Калибровочная таблица - -Прод-пакет = PyPI `telethon-floodgate` 0.1.0 (pin `>=0.1.0,<0.2`). - -| Категория | Было (прод 0.1.0) | Стало | Обоснование | -|---|---|---|---| -| `history` | 600/60 с — догадка, никогда не связывал | **24/30 с** | Эмпирическая граница Telegram (30 req/30 с, 31-й → FLOOD_WAIT_3) минус 20 % маржа; лог: 209 флудов при пиковой нагрузке 115/мин << 600/мин | -| `dialogs` | 1/60 с | без изменений | #1330; лог подтверждает: 2419→55 флудов warm после введения | -| `dialog_sweep` | 12/60 с | без изменений | #1359; 8 флудов за 3 мес — режим штатный | -| `send` | 30/60 с (+per-peer 1/1.0 с, 20/60) | без изменений, **требует калибровки** | 0 флудов — данных нет; per-peer-правки (1/1.1, 16/60) уже в локальном пакете по live-замеру | -| `admin_action` | 10/60 с | без изменений, **требует калибровки** | 0 флудов — данных нет | -| `channel_lifecycle` | 3/300 с | без изменений, **требует калибровки** | 0 флудов — данных нет | -| `default` | 1000/60 с | без изменений, **требует калибровки** | 3 единичных флуда `invoke_request` за 3 мес — против catch-all-страховки данных нет | - -Дифф (минимальный): `client_pool.py` передаёт гейту -`category_limits={"history": RateLimitSpec(24, 30.0)}` — константа -`HISTORY_CALIBRATED_SPEC` + регресс-тесты (`tests/test_rate_limit_gate.py`): -пул применяет спеку; 25-й вызов в 30-секундном окне деферится до окна -(фейковые часы `_Clock`). Остальные категории — дефолты пакета, пометка -«требует калибровки» живёт в комментарии у конструирования гейта и в этой -таблице. - -## 4. Координация с пакетом (важно) - -В локальной копии `~/Projects/telethon-floodgate` (main, uncommitted) уже -лежит ровно та же правка `HISTORY_SPEC = 24/30` (комментарий «empirical -boundary») — на PyPI она **не выпущена**, поэтому прод с 0.1.0 живёт с -600/60. Донорский оверрайд закрывает прод уже сейчас; после релиза пакета -с этим дефолтом оверрайд можно снять (значения совпадут) — это записано в -комментарии константы. - -## 5. Пропускная способность сбора (критерий AC) - -Живой замер до/после невозможен (запрещены живые вызовы) — оценка по логу: - -- Медианные активные минуты (26–50 сборов/мин) — **ниже 48/мин** (24/30 с), - дефера нет вообще, поведение не меняется. -- Пики (p95 74–101, max 115) — дефер; канал откладывается, а не теряется: - сбор инкрементальный (`min_id`), следующий проход добирает. Дефер warm- - prefetch (`collection.py:491`) вообще не останавливает сбор канала. -- Вместо пиков с флудами (каждый флуд = пауза 14–29 с + риск эскалации до - 120 с/бана) пики растягиваются гейтом до ~48 вызовов/30 с. Полных потерь - пропускной способности нет; платформа-цена — задержка добора пиковых - каналов до следующего прохода. - -## 6. Связка с #1419 (рестарты) - -Замер #1419 дополнительно показывает: 44 % стартов ловят флуд в первые 60 с -(~19× к базовой), рестарты 0.38/день кластеризуются, breaker в проде ни разу -не открывался (деплой после бан-инцидента). Данный PR калибровку лимитов -делает строже, что снижает и послемстартовый залп (gate-бюджеты после -рестарта пустые), но саму персистентность breaker/gate-состояний не решает — -это отдельное согласование (§8 отчёта #1419). - -## 7. Проверки - -- `ruff check src/telegram/client_pool.py tests/test_rate_limit_gate.py` — чисто. -- `tests/test_rate_limit_gate.py` — 13 passed (11 старых + 2 новых). -- Полный сюит: параллельная (`-m "not aiosqlite_serial" -n auto`) и - serial (`-m aiosqlite_serial -n auto --dist=loadfile`) части — см. PR. - -## 8. Риски / follow-ups - -- «Пропускная способность не деградирует» подтверждена оценкой по логу, не - живым замером — если понадобится живая проверка, делать по протоколу - калибровщика пакета (`scripts/calibrate_send_limits.py`), burnable-аккаунт. -- Релиз `telethon-floodgate` 0.1.1 с локальными правками владельца → снять - донорский оверрайд. -- Персистентность breaker-состояний (#1419 §8) — открытый вектор бан-масштаба - при рестарт-циклах. From 666c23bd1ee3894b8058cfd4f145096c3816806f Mon Sep 17 00:00:00 2001 From: axisrow Date: Thu, 24 Sep 2026 16:06:41 +0800 Subject: [PATCH 3/3] fix(queue): reschedule instead of FAIL when the calibrated gate defers (#1441) With history calibrated to 24/30s (#1418) the gate legitimately binds on peak collector minutes (median 50 fetches/min vs the 48/min cap), so TelegramRateLimitedError now reaches the queue's main collection path. Handle it like the neighbouring UsernameResolveRateLimitedError branch: reschedule with run_after = now + retry_after + buffer and a pending note, instead of FAILED + logger.exception on every peak. Co-Authored-By: Claude Code --- src/collection_queue.py | 33 +++++++++++++++ tests/test_collection_queue_db_pull.py | 58 ++++++++++++++++++++++++++ 2 files changed, 91 insertions(+) diff --git a/src/collection_queue.py b/src/collection_queue.py index 01cb9fdb..f48ca09d 100644 --- a/src/collection_queue.py +++ b/src/collection_queue.py @@ -8,6 +8,8 @@ from datetime import datetime, timedelta, timezone from typing import Any, cast +from telethon_floodgate import TelegramRateLimitedError + from src.database import Database, DatabaseBusyError from src.database.bundles import ChannelBundle from src.live_runtime_pause import LiveRuntimePauseGate @@ -30,6 +32,10 @@ # must be matched by BOTH type and message. Messages mirror facade._SQLITE_BUSY_MESSAGES. _SQLITE_BUSY_MESSAGES = ("database is locked", "database table is locked", "database is busy") +# Slack added on top of the gate's exact retry_after so the rescheduled run +# does not land a millisecond before the sliding window actually reopens. +GATE_RATE_LIMIT_RETRY_BUFFER_SEC = 5.0 + def _is_transient_busy_error(exc: BaseException) -> bool: """True for a transient SQLite lock from either DB path (#1249). @@ -602,6 +608,33 @@ async def _handle_collection_exception( exc.phone, ) return True, False + if isinstance(exc, TelegramRateLimitedError): + # The calibrated gate (#1418) legitimately binds on peak collector + # minutes (media 50/min vs the 48/min history cap): deferring the + # task is the designed outcome, not a failure — mirror the + # resolve-rate-limited branch above. + run_after = datetime.now(timezone.utc) + timedelta( + seconds=exc.retry_after_sec + GATE_RATE_LIMIT_RETRY_BUFFER_SEC + ) + note = ( + "Отложено: gate " + f"{exc.category} rate-limited до {run_after.astimezone(timezone.utc).isoformat()}" + ) + self._retried_tasks.discard(task_id) + await self._channels.reschedule_collection_task(task_id, run_after=run_after, note=note) + self._schedule_requeue_after_delay( + task_id=task_id, channel=channel, force=force, full=full, run_after=run_after + ) + logger.warning( + "Rescheduled collection task %d for channel %d until %s: " + "gate %s rate-limited on %s", + task_id, + channel.channel_id, + run_after.isoformat(), + exc.category, + exc.phone, + ) + return True, False if isinstance(exc, NoActiveCollectionClientsError): run_after = datetime.now(timezone.utc) + timedelta( seconds=self.NO_CLIENTS_RETRY_DELAY_SEC diff --git a/tests/test_collection_queue_db_pull.py b/tests/test_collection_queue_db_pull.py index 458ee2fe..3960947d 100644 --- a/tests/test_collection_queue_db_pull.py +++ b/tests/test_collection_queue_db_pull.py @@ -12,6 +12,7 @@ from datetime import datetime, timedelta, timezone import pytest +from telethon_floodgate import TelegramRateLimitedError from src.collection_queue import CollectionQueue from src.database import Database @@ -77,6 +78,15 @@ async def collect_single_channel( raise UsernameResolveRateLimitedError("+7001", 28.2) +class _GateRateLimitedCollector(_FakeCollector): + async def collect_single_channel( + self, channel, *, full=False, progress_callback=None, force=False, cancel_event=None + ): + self.calls.append(channel.channel_id) + self.full_calls.append(full) + raise TelegramRateLimitedError("+7001", "history", 17.0) + + class _BlockingCollector: def __init__(self): self.calls: list[int] = [] @@ -331,6 +341,54 @@ async def test_username_resolve_rate_limit_keeps_task_pending(tmp_path): await db.close() +@pytest.mark.anyio +async def test_gate_rate_limit_keeps_task_pending(tmp_path): + """A calibrated-gate deferral reschedules instead of failing (#1418). + + Once the history category binds on peak collector minutes, + TelegramRateLimitedError reaches the queue's main path and must take the + same reschedule route as the neighbouring resolve-rate-limited branch — + not FAILED + logger.exception noise on every peak. + """ + db = Database(str(tmp_path / "queue.db")) + await db.initialize() + try: + await _seed_channel(db) + + collector = _GateRateLimitedCollector() + queue = CollectionQueue(collector, db) + channel = (await db.get_channels(active_only=True))[0] + before = datetime.now(timezone.utc) + task_id = await queue.enqueue(channel) + + deadline = asyncio.get_event_loop().time() + 2.0 + while asyncio.get_event_loop().time() < deadline: + task = await db.get_collection_task(task_id) + if task.status == "pending" and task.run_after is not None: + break + await asyncio.sleep(0.05) + + task = await db.get_collection_task(task_id) + assert task.status == "pending" + assert task.error is None + assert task.run_after is not None + # retry_after (17s) + the 5s reschedule buffer. + assert task.run_after >= before + timedelta(seconds=22) + assert "gate history rate-limited" in (task.note or "") + + queue.start_db_pull(interval=0.02) + try: + await asyncio.sleep(0.12) + finally: + await queue.stop_db_pull() + + assert task_id in queue._known_task_ids + assert len(queue._delayed_requeues) == 1 + finally: + await queue.shutdown() + await db.close() + + @pytest.mark.anyio async def test_db_pull_does_not_double_ingest(tmp_path): db = Database(str(tmp_path / "queue.db"))