SilverElixir commited on
Commit
bd3cb33
·
1 Parent(s): 36d5fbb

Close four security gaps: shared rate limit, proxy auth, subprocess cleanup, log redaction

Browse files
Files changed (7) hide show
  1. bot.py +62 -12
  2. docs/ENVIRONMENT.md +8 -1
  3. lumen_telegram_transport.py +84 -2
  4. lumen_tiktok.py +28 -2
  5. proxy.ts +32 -4
  6. proxy_test.ts +78 -11
  7. test_bot.py +272 -1
bot.py CHANGED
@@ -65,6 +65,29 @@ _LOG_QUEUE: queue.SimpleQueue[logging.LogRecord] = queue.SimpleQueue()
65
  _LOG_QUEUE_HANDLER = logging.handlers.QueueHandler(_LOG_QUEUE)
66
  _LOG_LISTENER: logging.handlers.QueueListener | None = None
67
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
68
  def _setup_logging() -> logging.Logger:
69
  # LOG_LEVEL — раньше был захардкожен INFO везде (root+оба handler'а), из-за
70
  # чего оба существующих log.debug(...) в проекте не печатались никогда, ни в
@@ -87,6 +110,8 @@ def _setup_logging() -> logging.Logger:
87
  console_handler.setFormatter(fmt)
88
  file_handler.setFormatter(fmt)
89
 
 
 
90
  global _LOG_LISTENER
91
  # НАЙДЕНО ПРИ АУДИТЕ (4 сентября 2026, реальный воспроизведённый инцидент — см.
92
  # py-spy дамп стеков зависшего процесса): _setup_logging() вызывается больше
@@ -198,6 +223,7 @@ def _normalize_telegram_base_url(url: str) -> str:
198
  url = "https://" + url
199
  return url
200
 
 
201
  TELEGRAM_API_BASE_URL = _normalize_telegram_base_url(os.getenv("TELEGRAM_API_BASE_URL", "https://api.telegram.org"))
202
  log.info("[setup] Using Telegram API Base URL: %s", TELEGRAM_API_BASE_URL)
203
 
@@ -290,10 +316,7 @@ def _redactable_secrets() -> tuple[str, ...]:
290
  сохранённым историям чатов. honey: имена читаются по значению на момент
291
  вызова — WEBHOOK_SECRET/ADMIN_PANEL_KEY/UPSTASH_REDIS_REST_TOKEN объявлены
292
  ниже по файлу, это безопасно для module-level globals в теле функции."""
293
- return tuple(s for s in (
294
- BOT_TOKEN, GEMINI_API_KEY, OPENROUTER_API_KEY, _ADMIN_SECRET_SEED,
295
- WEBHOOK_SECRET, ADMIN_PANEL_KEY, UPSTASH_REDIS_REST_TOKEN,
296
- ) if s)
297
 
298
  def _sentry_scrub_secrets(event: dict, hint: dict) -> dict | None:
299
  """before_send-хук Sentry — вычищает секреты из события ПЕРЕД отправкой.
@@ -489,6 +512,7 @@ from lumen_telegram_transport import (
489
  _TelegramProxyCircuitBreaker,
490
  _looks_like_proxy_garbage,
491
  IPv4AiohttpSession,
 
492
  get_telegram_session as _lumen_get_telegram_session,
493
  close_telegram_session as _lumen_close_telegram_session,
494
  )
@@ -497,7 +521,10 @@ TG_PROXY_TRIP_THRESHOLD = int(os.getenv("TG_PROXY_TRIP_THRESHOLD", "3"))
497
  _tg_proxy_breaker = _TelegramProxyCircuitBreaker(cooldown_sec=TG_PROXY_COOLDOWN_SEC, trip_threshold=TG_PROXY_TRIP_THRESHOLD)
498
 
499
  async def _get_telegram_session() -> aiohttp.ClientSession:
500
- return await _lumen_get_telegram_session(TELEGRAM_REQUEST_TIMEOUT)
 
 
 
501
 
502
  async def _rotate_telegram_proxy() -> bool:
503
  """Переключается на следующий кандидат из _TELEGRAM_PROXY_CANDIDATES по кругу —
@@ -523,7 +550,10 @@ async def _rotate_telegram_proxy() -> bool:
523
  log.warning('[telegram] Switching to fallback proxy: %s -> %s', old_url, new_url)
524
  if bot is not None:
525
  old_session = bot.session
526
- bot.session = IPv4AiohttpSession(api=TelegramAPIServer.from_base(new_url))
 
 
 
527
  with contextlib.suppress(Exception):
528
  await old_session.close()
529
  return _telegram_proxy_idx != 0
@@ -542,6 +572,10 @@ async def _get_http_session() -> aiohttp.ClientSession:
542
  # используют на порядок меньше одновременных соединений, так что
543
  # повышение лимита их не затрагивает.
544
  connector=aiohttp.TCPConnector(family=socket.AF_INET, limit=40, ttl_dns_cache=300),
 
 
 
 
545
  )
546
  return _http_session
547
 
@@ -2421,6 +2455,7 @@ from lumen_tiktok import (
2421
  _download_url_bin,
2422
  _probe_video_dimensions,
2423
  _generate_video_thumbnail,
 
2424
  _probe_and_thumbnail_from_bytes,
2425
  _resolve_tiktok_short,
2426
  _fetch_tikwm_media_data,
@@ -3974,6 +4009,8 @@ async def cmd_draw(message: Message) -> None:
3974
  if not prompt:
3975
  await _safe_reply(message, "Укажите текст после команды /draw. Пример: /draw космическая станция")
3976
  return
 
 
3977
  await inline_draw(message, prompt)
3978
 
3979
  # ─────────────────── TTS-пайплайн (Fish Audio + Gemini TTS) ───────────────────
@@ -4066,7 +4103,7 @@ async def inline_tts(message: Message, text: str) -> None:
4066
  stdout=asyncio.subprocess.DEVNULL,
4067
  stderr=asyncio.subprocess.DEVNULL,
4068
  )
4069
- await asyncio.wait_for(proc.communicate(), timeout=30)
4070
  if os.path.exists(dst_path) and os.path.getsize(dst_path) > 0:
4071
  with open(dst_path, "rb") as fh:
4072
  final_audio = fh.read()
@@ -4080,7 +4117,7 @@ async def inline_tts(message: Message, text: str) -> None:
4080
  stdout=asyncio.subprocess.PIPE,
4081
  stderr=asyncio.subprocess.DEVNULL,
4082
  )
4083
- probe_out, _ = await asyncio.wait_for(probe.communicate(), timeout=10)
4084
  raw_dur = probe_out.decode().strip()
4085
  voice_duration = max(1, round(float(raw_dur))) if raw_dur else 0
4086
  except Exception as probe_exc:
@@ -4119,6 +4156,8 @@ async def cmd_tts(message: Message) -> None:
4119
  if not text:
4120
  await _safe_reply(message, "Укажите текст после команды /tts. Пример: /tts Добрый день")
4121
  return
 
 
4122
  await inline_tts(message, text)
4123
 
4124
  @dp.message(Command("reset"))
@@ -4537,6 +4576,13 @@ def _check_and_register_rate_limit(user_id: int | None) -> bool:
4537
  return False
4538
 
4539
 
 
 
 
 
 
 
 
4540
  async def _resolve_incoming_media(
4541
  message: Message, state: dict[str, Any], clean_prompt: str, *, is_private: bool,
4542
  ) -> tuple[str | None, str, str, tuple[bytes, str] | None]:
@@ -4623,8 +4669,7 @@ async def _handle_message_core(message: Message, extra_media: list[tuple[bytes,
4623
  return
4624
 
4625
  # Начинаем обработку активного запроса с проверкой rate limit
4626
- if _check_and_register_rate_limit(_rate_limit_key_for_message(message)):
4627
- await _tg_call(message.reply, "Вы отправляете слишком много запросов. Подождите немного.")
4628
  return
4629
 
4630
  # Проверка на ссылки загрузки (TikTok — сразу всегда, даже в группах без упоминания)
@@ -4981,9 +5026,14 @@ async def main() -> None:
4981
  # Мы настраиваем Bot сессию с принудительным IPv4 и таймаутами для hg space
4982
  if TELEGRAM_API_BASE_URL != "https://api.telegram.org":
4983
  api_server = TelegramAPIServer.from_base(TELEGRAM_API_BASE_URL)
4984
- sess = IPv4AiohttpSession(api=api_server)
 
 
 
4985
  else:
4986
- sess = IPv4AiohttpSession()
 
 
4987
  bot = Bot(token=BOT_TOKEN, session=sess)
4988
  client = genai.Client(api_key=GEMINI_API_KEY)
4989
 
 
65
  _LOG_QUEUE_HANDLER = logging.handlers.QueueHandler(_LOG_QUEUE)
66
  _LOG_LISTENER: logging.handlers.QueueListener | None = None
67
 
68
+ _SECRET_NAMES = (
69
+ "BOT_TOKEN", "TELEGRAM_TOKEN", "TELEGRAM_BOT_TOKEN", "GEMINI_API_KEY",
70
+ "OPENROUTER_API_KEY", "OPENROUTER_KEY", "ADMIN_SECRET_SEED",
71
+ "_ADMIN_SECRET_SEED", "WEBHOOK_SECRET", "ADMIN_PANEL_KEY",
72
+ "UPSTASH_REDIS_REST_TOKEN", "LUMEN_PROXY_SECRET",
73
+ )
74
+
75
+
76
+ def _current_log_secrets() -> tuple[str, ...]:
77
+ values = [value for name in _SECRET_NAMES
78
+ for value in (globals().get(name), os.getenv(name, ""))
79
+ if isinstance(value, str) and value]
80
+ return tuple(sorted(set(values), key=len, reverse=True))
81
+
82
+
83
+ class _SecretLogFormatter(logging.Formatter):
84
+ def format(self, record: logging.LogRecord) -> str:
85
+ text = super().format(record)
86
+ for secret in _current_log_secrets():
87
+ text = text.replace(secret, "<REDACTED>")
88
+ return text
89
+
90
+
91
  def _setup_logging() -> logging.Logger:
92
  # LOG_LEVEL — раньше был захардкожен INFO везде (root+оба handler'а), из-за
93
  # чего оба существующих log.debug(...) в проекте не печатались никогда, ни в
 
110
  console_handler.setFormatter(fmt)
111
  file_handler.setFormatter(fmt)
112
 
113
+ _LOG_QUEUE_HANDLER.setFormatter(_SecretLogFormatter())
114
+
115
  global _LOG_LISTENER
116
  # НАЙДЕНО ПРИ АУДИТЕ (4 сентября 2026, реальный воспроизведённый инцидент — см.
117
  # py-spy дамп стеков зависшего процесса): _setup_logging() вызывается больше
 
223
  url = "https://" + url
224
  return url
225
 
226
+ LUMEN_PROXY_SECRET = os.getenv("LUMEN_PROXY_SECRET", "")
227
  TELEGRAM_API_BASE_URL = _normalize_telegram_base_url(os.getenv("TELEGRAM_API_BASE_URL", "https://api.telegram.org"))
228
  log.info("[setup] Using Telegram API Base URL: %s", TELEGRAM_API_BASE_URL)
229
 
 
316
  сохранённым историям чатов. honey: имена читаются по значению на момент
317
  вызова — WEBHOOK_SECRET/ADMIN_PANEL_KEY/UPSTASH_REDIS_REST_TOKEN объявлены
318
  ниже по файлу, это безопасно для module-level globals в теле функции."""
319
+ return _current_log_secrets()
 
 
 
320
 
321
  def _sentry_scrub_secrets(event: dict, hint: dict) -> dict | None:
322
  """before_send-хук Sentry — вычищает секреты из события ПЕРЕД отправкой.
 
512
  _TelegramProxyCircuitBreaker,
513
  _looks_like_proxy_garbage,
514
  IPv4AiohttpSession,
515
+ proxy_auth_middlewares,
516
  get_telegram_session as _lumen_get_telegram_session,
517
  close_telegram_session as _lumen_close_telegram_session,
518
  )
 
521
  _tg_proxy_breaker = _TelegramProxyCircuitBreaker(cooldown_sec=TG_PROXY_COOLDOWN_SEC, trip_threshold=TG_PROXY_TRIP_THRESHOLD)
522
 
523
  async def _get_telegram_session() -> aiohttp.ClientSession:
524
+ return await _lumen_get_telegram_session(
525
+ TELEGRAM_REQUEST_TIMEOUT,
526
+ proxy_secret=LUMEN_PROXY_SECRET, proxy_base_urls=_TELEGRAM_PROXY_CANDIDATES,
527
+ )
528
 
529
  async def _rotate_telegram_proxy() -> bool:
530
  """Переключается на следующий кандидат из _TELEGRAM_PROXY_CANDIDATES по кругу —
 
550
  log.warning('[telegram] Switching to fallback proxy: %s -> %s', old_url, new_url)
551
  if bot is not None:
552
  old_session = bot.session
553
+ bot.session = IPv4AiohttpSession(
554
+ api=TelegramAPIServer.from_base(new_url),
555
+ proxy_secret=LUMEN_PROXY_SECRET, proxy_base_urls=_TELEGRAM_PROXY_CANDIDATES,
556
+ )
557
  with contextlib.suppress(Exception):
558
  await old_session.close()
559
  return _telegram_proxy_idx != 0
 
572
  # используют на порядок меньше одновременных соединений, так что
573
  # повышение лимита их не затрагивает.
574
  connector=aiohttp.TCPConnector(family=socket.AF_INET, limit=40, ttl_dns_cache=300),
575
+ middlewares=proxy_auth_middlewares(
576
+ proxy_secret=LUMEN_PROXY_SECRET,
577
+ proxy_base_urls=(*_TELEGRAM_PROXY_CANDIDATES, *_tikwm_proxy_candidates()),
578
+ ),
579
  )
580
  return _http_session
581
 
 
2455
  _download_url_bin,
2456
  _probe_video_dimensions,
2457
  _generate_video_thumbnail,
2458
+ _communicate_process,
2459
  _probe_and_thumbnail_from_bytes,
2460
  _resolve_tiktok_short,
2461
  _fetch_tikwm_media_data,
 
4009
  if not prompt:
4010
  await _safe_reply(message, "Укажите текст после команды /draw. Пример: /draw космическая станция")
4011
  return
4012
+ if await _reject_rate_limited_message(message):
4013
+ return
4014
  await inline_draw(message, prompt)
4015
 
4016
  # ─────────────────── TTS-пайплайн (Fish Audio + Gemini TTS) ───────────────────
 
4103
  stdout=asyncio.subprocess.DEVNULL,
4104
  stderr=asyncio.subprocess.DEVNULL,
4105
  )
4106
+ await _communicate_process(proc, timeout=30)
4107
  if os.path.exists(dst_path) and os.path.getsize(dst_path) > 0:
4108
  with open(dst_path, "rb") as fh:
4109
  final_audio = fh.read()
 
4117
  stdout=asyncio.subprocess.PIPE,
4118
  stderr=asyncio.subprocess.DEVNULL,
4119
  )
4120
+ probe_out, _ = await _communicate_process(probe, timeout=10)
4121
  raw_dur = probe_out.decode().strip()
4122
  voice_duration = max(1, round(float(raw_dur))) if raw_dur else 0
4123
  except Exception as probe_exc:
 
4156
  if not text:
4157
  await _safe_reply(message, "Укажите текст после команды /tts. Пример: /tts Добрый день")
4158
  return
4159
+ if await _reject_rate_limited_message(message):
4160
+ return
4161
  await inline_tts(message, text)
4162
 
4163
  @dp.message(Command("reset"))
 
4576
  return False
4577
 
4578
 
4579
+ async def _reject_rate_limited_message(message: Message) -> bool:
4580
+ if not _check_and_register_rate_limit(_rate_limit_key_for_message(message)):
4581
+ return False
4582
+ await _tg_call(message.reply, "Вы отправляете слишком много запросов. Подождите немного.")
4583
+ return True
4584
+
4585
+
4586
  async def _resolve_incoming_media(
4587
  message: Message, state: dict[str, Any], clean_prompt: str, *, is_private: bool,
4588
  ) -> tuple[str | None, str, str, tuple[bytes, str] | None]:
 
4669
  return
4670
 
4671
  # Начинаем обработку активного запроса с проверкой rate limit
4672
+ if await _reject_rate_limited_message(message):
 
4673
  return
4674
 
4675
  # Проверка на ссылки загрузки (TikTok — сразу всегда, даже в группах без упоминания)
 
5026
  # Мы настраиваем Bot сессию с принудительным IPv4 и таймаутами для hg space
5027
  if TELEGRAM_API_BASE_URL != "https://api.telegram.org":
5028
  api_server = TelegramAPIServer.from_base(TELEGRAM_API_BASE_URL)
5029
+ sess = IPv4AiohttpSession(
5030
+ api=api_server,
5031
+ proxy_secret=LUMEN_PROXY_SECRET, proxy_base_urls=_TELEGRAM_PROXY_CANDIDATES,
5032
+ )
5033
  else:
5034
+ sess = IPv4AiohttpSession(
5035
+ proxy_secret=LUMEN_PROXY_SECRET, proxy_base_urls=_TELEGRAM_PROXY_CANDIDATES,
5036
+ )
5037
  bot = Bot(token=BOT_TOKEN, session=sess)
5038
  client = genai.Client(api_key=GEMINI_API_KEY)
5039
 
docs/ENVIRONMENT.md CHANGED
@@ -1,6 +1,6 @@
1
  # Environment variables
2
 
3
- Only `BOT_TOKEN` and `GEMINI_API_KEY` are required. Everything else has a working default; most deployments never need to touch the rest of this file.
4
 
5
  ## Required
6
 
@@ -15,6 +15,7 @@ Hugging Face Spaces' outbound IPs are blocked by Telegram's API and rejected (`4
15
 
16
  | Variable | Default | Purpose |
17
  |---|---|---|
 
18
  | `TELEGRAM_API_BASE_URL` | `https://api.telegram.org` | Base URL for the Telegram Bot API. In production this points at the proxy. |
19
  | `TELEGRAM_API_BASE_URL_FALLBACKS` | — | Comma-separated backup proxy addresses. On a circuit-breaker trip the bot rotates through these before pausing. |
20
  | `TG_PROXY_COOLDOWN_SEC` | `20` | Pause (seconds) after the circuit breaker trips, i.e. the proxy is judged unavailable. |
@@ -22,6 +23,12 @@ Hugging Face Spaces' outbound IPs are blocked by Telegram's API and rejected (`4
22
  | `TIKWM_API_BASE_URL` | — (direct requests) | Base URL for a TikWM proxy. Empty means the bot talks to both `tikwm.com` mirrors directly. |
23
  | `TIKWM_API_BASE_URL_FALLBACKS` | — | Comma-separated backup TikWM proxies, tried in order if the primary one fails. |
24
 
 
 
 
 
 
 
25
  ## Admin access & secrets
26
 
27
  | Variable | Default | Purpose |
 
1
  # Environment variables
2
 
3
+ `BOT_TOKEN` and `GEMINI_API_KEY` are required. Deployments using the authenticated proxy also require `LUMEN_PROXY_SECRET` on both the bot and every proxy instance.
4
 
5
  ## Required
6
 
 
15
 
16
  | Variable | Default | Purpose |
17
  |---|---|---|
18
+ | `LUMEN_PROXY_SECRET` | — | Required on every Deno proxy and the bot when using proxies. Use the same independently generated random secret (at least 32 random bytes encoded as hex) for primary and fallback instances. Sent only in `X-Lumen-Proxy-Secret`, never in URLs or `Authorization`. |
19
  | `TELEGRAM_API_BASE_URL` | `https://api.telegram.org` | Base URL for the Telegram Bot API. In production this points at the proxy. |
20
  | `TELEGRAM_API_BASE_URL_FALLBACKS` | — | Comma-separated backup proxy addresses. On a circuit-breaker trip the bot rotates through these before pausing. |
21
  | `TG_PROXY_COOLDOWN_SEC` | `20` | Pause (seconds) after the circuit breaker trips, i.e. the proxy is judged unavailable. |
 
23
  | `TIKWM_API_BASE_URL` | — (direct requests) | Base URL for a TikWM proxy. Empty means the bot talks to both `tikwm.com` mirrors directly. |
24
  | `TIKWM_API_BASE_URL_FALLBACKS` | — | Comma-separated backup TikWM proxies, tried in order if the primary one fails. |
25
 
26
+ The proxy denies requests before forwarding: missing/wrong credentials return `401`; an unset, empty or invalid server secret returns `503`. Secrets must be printable ASCII without spaces. The header is removed before upstream requests; Telegram bot-token paths and upstream `Authorization` are unchanged. Redirects are rejected to keep requests within the host allowlist. Upstream failures return a generic `502` without exception details. Local serving requires permission to read `LUMEN_PROXY_SECRET` (`--allow-env=LUMEN_PROXY_SECRET`) as well as network permission; isolated tests require neither.
27
+
28
+ The bot applies proxy authentication to aiogram, raw Telegram/file requests and TikWM metadata requests, including configured fallbacks. The middleware limits credentials to trusted HTTPS proxy origins and base-path boundaries, never direct Telegram/TikWM or media/CDN hosts. Do not put the secret in shared session headers or the TikTok media headers.
29
+
30
+ Before publishing, configure the same `LUMEN_PROXY_SECRET` on HF and every Deno proxy. Deploy the updated bot client first, then the protected proxy: old clients cannot use the protected proxy. The updated bot refuses configured proxies without a valid secret. Proxy deployment is separate from the HF auto-deploy. The secret is included in log and Sentry redaction.
31
+
32
  ## Admin access & secrets
33
 
34
  | Variable | Default | Purpose |
lumen_telegram_transport.py CHANGED
@@ -26,6 +26,8 @@ from __future__ import annotations
26
  import logging
27
  import socket
28
  import time
 
 
29
 
30
  import aiohttp
31
  from aiogram.client.session.aiohttp import AiohttpSession
@@ -160,32 +162,112 @@ def _build_telegram_connector(limit: int) -> aiohttp.TCPConnector:
160
  )
161
 
162
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
163
  class IPv4AiohttpSession(AiohttpSession):
 
 
 
 
 
 
 
 
164
  async def create_session(self) -> aiohttp.ClientSession:
165
  if self._session is None or self._session.closed:
166
  self._session = aiohttp.ClientSession(
167
  connector=_build_telegram_connector(limit=30),
168
  timeout=aiohttp.ClientTimeout(total=30.0, connect=10.0, sock_read=20.0),
169
  json_serialize=self.json_dumps,
 
170
  )
171
  return self._session
172
 
173
 
174
  _telegram_session: aiohttp.ClientSession | None = None
 
175
 
176
 
177
- async def get_telegram_session(request_timeout: float) -> aiohttp.ClientSession:
 
 
178
  """Кэширующий геттер общей aiohttp-сессии для прямых HTTP-вызовов к Telegram Bot
179
  API (используется telegram_api_call в bot.py). `request_timeout` — значение
180
  TELEGRAM_REQUEST_TIMEOUT из bot.py, передаётся параметром на каждый вызов (а не
181
  импортируется статически), т.к. это часть публичной, потенциально настраиваемой
182
  через env конфигурации bot.py, а не константа этого модуля."""
183
- global _telegram_session
 
 
 
 
 
 
184
  if _telegram_session is None or _telegram_session.closed:
185
  _telegram_session = aiohttp.ClientSession(
186
  connector=_build_telegram_connector(limit=10),
187
  timeout=aiohttp.ClientTimeout(total=request_timeout + 10.0, connect=10.0),
 
188
  )
 
189
  return _telegram_session
190
 
191
 
 
26
  import logging
27
  import socket
28
  import time
29
+ from collections.abc import Sequence
30
+ from urllib.parse import urlsplit
31
 
32
  import aiohttp
33
  from aiogram.client.session.aiohttp import AiohttpSession
 
162
  )
163
 
164
 
165
+ PROXY_AUTH_HEADER = "X-Lumen-Proxy-Secret"
166
+
167
+
168
+ def proxy_auth_middlewares(
169
+ *, proxy_secret: str = "", proxy_base_urls: Sequence[str] = (),
170
+ ) -> tuple:
171
+ scopes = []
172
+ for base_url in proxy_base_urls:
173
+ if not base_url:
174
+ continue
175
+ try:
176
+ parsed = urlsplit(base_url)
177
+ port = parsed.port or 443
178
+ except ValueError:
179
+ raise ValueError("Invalid proxy base URL") from None
180
+ if (
181
+ parsed.scheme != "https" or not parsed.hostname
182
+ or parsed.username is not None or parsed.password is not None
183
+ or parsed.query or parsed.fragment
184
+ ):
185
+ raise ValueError("Proxy base URL must use HTTPS without credentials, query or fragment")
186
+ if parsed.hostname in {"api.telegram.org", "www.tikwm.com", "tikwm.com"}:
187
+ continue
188
+ scopes.append((parsed.hostname, port, parsed.path.rstrip("/")))
189
+ if scopes and (not proxy_secret or any(not 33 <= ord(c) <= 126 for c in proxy_secret)):
190
+ raise ValueError("LUMEN_PROXY_SECRET is required for configured proxies and must be printable ASCII without spaces")
191
+
192
+ async def authenticate(request: aiohttp.ClientRequest, handler):
193
+ request.headers.popall(PROXY_AUTH_HEADER, None)
194
+ url = request.url
195
+ authenticated = url.scheme == "https" and any(
196
+ url.host == host and url.port == port
197
+ and (url.path == path or url.path.startswith(path + "/"))
198
+ for host, port, path in scopes
199
+ )
200
+ if authenticated:
201
+ request.headers[PROXY_AUTH_HEADER] = proxy_secret
202
+ try:
203
+ response = await handler(request)
204
+ except Exception:
205
+ if authenticated:
206
+ raise RuntimeError("Authenticated proxy request failed") from None
207
+ raise
208
+ finally:
209
+ request.headers.popall(PROXY_AUTH_HEADER, None)
210
+ if authenticated:
211
+ safe_headers = response.request_info.headers.copy()
212
+ safe_headers.popall(PROXY_AUTH_HEADER, None)
213
+ response._request_info = aiohttp.RequestInfo(
214
+ url=response.request_info.url, method=response.request_info.method,
215
+ headers=safe_headers, real_url=response.request_info.real_url,
216
+ )
217
+ if authenticated and 300 <= response.status < 400:
218
+ response.close()
219
+ raise RuntimeError("Authenticated proxy redirects are disabled")
220
+ return response
221
+
222
+ return (authenticate,)
223
+
224
+
225
  class IPv4AiohttpSession(AiohttpSession):
226
+ def __init__(
227
+ self, *, proxy_secret: str = "", proxy_base_urls: Sequence[str] = (), **kwargs,
228
+ ) -> None:
229
+ self._proxy_middlewares = proxy_auth_middlewares(
230
+ proxy_secret=proxy_secret, proxy_base_urls=proxy_base_urls,
231
+ )
232
+ super().__init__(**kwargs)
233
+
234
  async def create_session(self) -> aiohttp.ClientSession:
235
  if self._session is None or self._session.closed:
236
  self._session = aiohttp.ClientSession(
237
  connector=_build_telegram_connector(limit=30),
238
  timeout=aiohttp.ClientTimeout(total=30.0, connect=10.0, sock_read=20.0),
239
  json_serialize=self.json_dumps,
240
+ middlewares=self._proxy_middlewares,
241
  )
242
  return self._session
243
 
244
 
245
  _telegram_session: aiohttp.ClientSession | None = None
246
+ _telegram_session_auth: tuple[str, tuple[str, ...]] | None = None
247
 
248
 
249
+ async def get_telegram_session(
250
+ request_timeout: float, *, proxy_secret: str = "", proxy_base_urls: Sequence[str] = (),
251
+ ) -> aiohttp.ClientSession:
252
  """Кэширующий геттер общей aiohttp-сессии для прямых HTTP-вызовов к Telegram Bot
253
  API (используется telegram_api_call в bot.py). `request_timeout` — значение
254
  TELEGRAM_REQUEST_TIMEOUT из bot.py, передаётся параметром на каждый вызов (а не
255
  импортируется статически), т.к. это часть публичной, потенциально настраиваемой
256
  через env конфигурации bot.py, а не константа этого модуля."""
257
+ global _telegram_session, _telegram_session_auth
258
+ auth = (proxy_secret, tuple(proxy_base_urls))
259
+ middlewares = proxy_auth_middlewares(
260
+ proxy_secret=proxy_secret, proxy_base_urls=auth[1],
261
+ )
262
+ if _telegram_session is not None and not _telegram_session.closed and _telegram_session_auth != auth:
263
+ await _telegram_session.close()
264
  if _telegram_session is None or _telegram_session.closed:
265
  _telegram_session = aiohttp.ClientSession(
266
  connector=_build_telegram_connector(limit=10),
267
  timeout=aiohttp.ClientTimeout(total=request_timeout + 10.0, connect=10.0),
268
+ middlewares=middlewares,
269
  )
270
+ _telegram_session_auth = auth
271
  return _telegram_session
272
 
273
 
lumen_tiktok.py CHANGED
@@ -522,6 +522,32 @@ async def _download_url_bin(session: aiohttp.ClientSession, url: str, headers: d
522
  return None
523
 
524
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
525
  async def _probe_video_dimensions(path: str) -> tuple[int, int, int]:
526
  """Возвращает (duration_seconds, width, height). Без этих полей Telegram иногда
527
  не может сам распознать видео и показывает его как "сырой файл" с 0:00 вместо
@@ -536,7 +562,7 @@ async def _probe_video_dimensions(path: str) -> tuple[int, int, int]:
536
  stdout=asyncio.subprocess.PIPE,
537
  stderr=asyncio.subprocess.PIPE,
538
  )
539
- stdout, _ = await asyncio.wait_for(proc.communicate(), timeout=15)
540
  width = height = 0
541
  duration = 0
542
  for line in stdout.decode(errors="replace").splitlines():
@@ -567,7 +593,7 @@ async def _generate_video_thumbnail(path: str, duration: int) -> bytes | None:
567
  stdout=asyncio.subprocess.DEVNULL,
568
  stderr=asyncio.subprocess.DEVNULL,
569
  )
570
- await asyncio.wait_for(proc.wait(), timeout=15)
571
  if os.path.exists(thumb_path) and os.path.getsize(thumb_path) > 0:
572
  with open(thumb_path, "rb") as f:
573
  return f.read()
 
522
  return None
523
 
524
 
525
+ async def _communicate_process(proc: asyncio.subprocess.Process, *, timeout: float):
526
+ async def finish():
527
+ try:
528
+ return await proc.communicate()
529
+ finally:
530
+ await proc.wait()
531
+
532
+ completion = asyncio.create_task(finish())
533
+ try:
534
+ return await asyncio.wait_for(asyncio.shield(completion), timeout=timeout)
535
+ except (asyncio.TimeoutError, asyncio.CancelledError):
536
+ if proc.returncode is None:
537
+ with contextlib.suppress(ProcessLookupError):
538
+ proc.kill()
539
+ while not completion.done():
540
+ try:
541
+ await asyncio.shield(completion)
542
+ except asyncio.CancelledError:
543
+ continue
544
+ except Exception:
545
+ break
546
+ with contextlib.suppress(Exception, asyncio.CancelledError):
547
+ completion.result()
548
+ raise
549
+
550
+
551
  async def _probe_video_dimensions(path: str) -> tuple[int, int, int]:
552
  """Возвращает (duration_seconds, width, height). Без этих полей Telegram иногда
553
  не может сам распознать видео и показывает его как "сырой файл" с 0:00 вместо
 
562
  stdout=asyncio.subprocess.PIPE,
563
  stderr=asyncio.subprocess.PIPE,
564
  )
565
+ stdout, _ = await _communicate_process(proc, timeout=15)
566
  width = height = 0
567
  duration = 0
568
  for line in stdout.decode(errors="replace").splitlines():
 
593
  stdout=asyncio.subprocess.DEVNULL,
594
  stderr=asyncio.subprocess.DEVNULL,
595
  )
596
+ await _communicate_process(proc, timeout=15)
597
  if os.path.exists(thumb_path) and os.path.getsize(thumb_path) > 0:
598
  with open(thumb_path, "rb") as f:
599
  return f.read()
proxy.ts CHANGED
@@ -41,6 +41,8 @@
41
  * изменений в Python-коде для смены адреса прокси не требуется.
42
  */
43
 
 
 
44
  export const ALLOWED_HOSTS = new Set([
45
  "api.telegram.org",
46
  "www.tikwm.com",
@@ -94,6 +96,7 @@ export function resolveTarget(pathname: string, search: string): TargetResolutio
94
  export function buildForwardHeaders(reqHeaders: Headers): Headers {
95
  const headers = new Headers(reqHeaders);
96
  for (const name of HOP_BY_HOP_REQUEST_HEADERS) headers.delete(name);
 
97
  return headers;
98
  }
99
 
@@ -105,7 +108,30 @@ export function buildResponseHeaders(upstreamHeaders: Headers): Headers {
105
 
106
  // fetchImpl — точка подмены для тестов (тот же приём, что bot._get_http_session
107
  // и т.п. в Python-части проекта) — реальная сеть не нужна ни одному юнит-тесту.
108
- export async function handleRequest(req: Request, fetchImpl: typeof fetch = fetch): Promise<Response> {
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
109
  const url = new URL(req.url);
110
  const target = resolveTarget(url.pathname, url.search);
111
  if (!target.ok) {
@@ -118,6 +144,7 @@ export async function handleRequest(req: Request, fetchImpl: typeof fetch = fetc
118
  upstreamResp = await fetchImpl(target.url, {
119
  method: req.method,
120
  headers: forwardHeaders,
 
121
  // GET/HEAD не могут иметь тело запроса (fetch бросит исключение, если
122
  // передать body для них) — для остальных методов пробрасываем тело
123
  // напрямую как поток, не буферизуя целиком в памяти (важно для
@@ -127,8 +154,8 @@ export async function handleRequest(req: Request, fetchImpl: typeof fetch = fetc
127
  // (часть стандарта WHATWG fetch для body типа ReadableStream).
128
  duplex: "half",
129
  });
130
- } catch (e) {
131
- return new Response(`Upstream fetch failed: ${e}`, { status: 502 });
132
  }
133
 
134
  const respHeaders = buildResponseHeaders(upstreamResp.headers);
@@ -141,5 +168,6 @@ export async function handleRequest(req: Request, fetchImpl: typeof fetch = fetc
141
  // Реальный сервер стартует только при прямом запуске файла (deno run/deploy),
142
  // не при импорте из proxy_test.ts — иначе тесты пытались бы забиндить порт.
143
  if (import.meta.main) {
144
- Deno.serve((req) => handleRequest(req));
 
145
  }
 
41
  * изменений в Python-коде для смены адреса прокси не требуется.
42
  */
43
 
44
+ export const PROXY_AUTH_HEADER = "X-Lumen-Proxy-Secret";
45
+
46
  export const ALLOWED_HOSTS = new Set([
47
  "api.telegram.org",
48
  "www.tikwm.com",
 
96
  export function buildForwardHeaders(reqHeaders: Headers): Headers {
97
  const headers = new Headers(reqHeaders);
98
  for (const name of HOP_BY_HOP_REQUEST_HEADERS) headers.delete(name);
99
+ headers.delete(PROXY_AUTH_HEADER);
100
  return headers;
101
  }
102
 
 
108
 
109
  // fetchImpl — точка подмены для тестов (тот же приём, что bot._get_http_session
110
  // и т.п. в Python-части проекта) — реальная сеть не нужна ни одному юнит-тесту.
111
+ export async function handleRequest(
112
+ req: Request,
113
+ fetchImpl: typeof fetch = fetch,
114
+ proxySecret: string | undefined = undefined,
115
+ ): Promise<Response> {
116
+ if (!proxySecret || !/^[\x21-\x7e]+$/.test(proxySecret)) {
117
+ return new Response("Proxy authentication unavailable", { status: 503 });
118
+ }
119
+ const suppliedSecret = req.headers.get(PROXY_AUTH_HEADER) ?? "";
120
+ const encoder = new TextEncoder();
121
+ const [expected, supplied] = await Promise.all([
122
+ crypto.subtle.digest("SHA-256", encoder.encode(proxySecret)),
123
+ crypto.subtle.digest("SHA-256", encoder.encode(suppliedSecret)),
124
+ ]);
125
+ const expectedBytes = new Uint8Array(expected);
126
+ const suppliedBytes = new Uint8Array(supplied);
127
+ let difference = 0;
128
+ for (let i = 0; i < expectedBytes.length; i++) {
129
+ difference |= expectedBytes[i] ^ suppliedBytes[i];
130
+ }
131
+ if (!suppliedSecret || difference !== 0) {
132
+ return new Response("Unauthorized", { status: 401 });
133
+ }
134
+
135
  const url = new URL(req.url);
136
  const target = resolveTarget(url.pathname, url.search);
137
  if (!target.ok) {
 
144
  upstreamResp = await fetchImpl(target.url, {
145
  method: req.method,
146
  headers: forwardHeaders,
147
+ redirect: "error",
148
  // GET/HEAD не могут иметь тело запроса (fetch бросит исключение, если
149
  // передать body для них) — для остальных методов пробрасываем тело
150
  // напрямую как поток, не буферизуя целиком в памяти (важно для
 
154
  // (часть стандарта WHATWG fetch для body типа ReadableStream).
155
  duplex: "half",
156
  });
157
+ } catch {
158
+ return new Response("Upstream fetch failed", { status: 502 });
159
  }
160
 
161
  const respHeaders = buildResponseHeaders(upstreamResp.headers);
 
168
  // Реальный сервер стартует только при прямом запуске файла (deno run/deploy),
169
  // не при импорте из proxy_test.ts — иначе тесты пытались бы забиндить порт.
170
  if (import.meta.main) {
171
+ const proxySecret = Deno.env.get("LUMEN_PROXY_SECRET");
172
+ Deno.serve((req) => handleRequest(req, fetch, proxySecret));
173
  }
proxy_test.ts CHANGED
@@ -7,12 +7,16 @@
7
 
8
  import {
9
  ALLOWED_HOSTS,
 
10
  buildForwardHeaders,
11
  buildResponseHeaders,
12
  handleRequest,
13
  resolveTarget,
14
  } from "./proxy.ts";
15
 
 
 
 
16
  function assert(condition: boolean, message: string): void {
17
  if (!condition) throw new Error(message);
18
  }
@@ -132,8 +136,8 @@ Deno.test("handleRequest пробрасывает GET без тела и воз
132
  headers: { "Content-Type": "application/json" },
133
  }));
134
  };
135
- const req = new Request("https://proxy.example/fetch/www.tikwm.com/api/?url=x", { method: "GET" });
136
- const resp = await handleRequest(req, fakeFetch);
137
  assertEquals(resp.status, 200);
138
  const body = await resp.json();
139
  assertEquals(body, { code: 0 });
@@ -145,16 +149,16 @@ Deno.test("handleRequest возвращает 403 для неразрешённ
145
  fetchCalled = true;
146
  return Promise.resolve(new Response("не должно случиться"));
147
  };
148
- const req = new Request("https://proxy.example/fetch/evil.example.com/x", { method: "GET" });
149
- const resp = await handleRequest(req, fakeFetch);
150
  assertEquals(resp.status, 403);
151
  assertEquals(fetchCalled, false, "fetch не должен был вызываться для неразрешённого хоста");
152
  });
153
 
154
  Deno.test("handleRequest возвращает 404 для пути без /fetch/ префикса", async () => {
155
- const req = new Request("https://proxy.example/something-else", { method: "GET" });
156
  const unusedFetch: typeof fetch = () => Promise.resolve(new Response("unused"));
157
- const resp = await handleRequest(req, unusedFetch);
158
  assertEquals(resp.status, 404);
159
  });
160
 
@@ -169,9 +173,9 @@ Deno.test("handleRequest пробрасывает метод и тело для
169
  const req = new Request("https://proxy.example/fetch/api.telegram.org/bot123/sendMessage", {
170
  method: "POST",
171
  body: JSON.stringify({ chat_id: 1, text: "hi" }),
172
- headers: { "Content-Type": "application/json" },
173
  });
174
- const resp = await handleRequest(req, fakeFetch);
175
  assertEquals(resp.status, 200);
176
  assertEquals(capturedMethod, "POST");
177
  assert(capturedBodyIsStream, "тело POST-запроса должно передаваться апстриму как поток (без буферизации целиком)");
@@ -179,9 +183,72 @@ Deno.test("handleRequest пробрасывает метод и тело для
179
 
180
  Deno.test("handleRequest возвращает 502, если апстрим-fetch упал с исключением", async () => {
181
  const fakeFetch: typeof fetch = () => {
182
- throw new Error("network unreachable");
183
  };
184
- const req = new Request("https://proxy.example/fetch/api.telegram.org/getMe", { method: "GET" });
185
- const resp = await handleRequest(req, fakeFetch);
186
  assertEquals(resp.status, 502);
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
187
  });
 
7
 
8
  import {
9
  ALLOWED_HOSTS,
10
+ PROXY_AUTH_HEADER,
11
  buildForwardHeaders,
12
  buildResponseHeaders,
13
  handleRequest,
14
  resolveTarget,
15
  } from "./proxy.ts";
16
 
17
+ const TEST_SECRET = "isolated-test-secret";
18
+ const AUTH_HEADERS = { [PROXY_AUTH_HEADER]: TEST_SECRET };
19
+
20
  function assert(condition: boolean, message: string): void {
21
  if (!condition) throw new Error(message);
22
  }
 
136
  headers: { "Content-Type": "application/json" },
137
  }));
138
  };
139
+ const req = new Request("https://proxy.example/fetch/www.tikwm.com/api/?url=x", { method: "GET", headers: AUTH_HEADERS });
140
+ const resp = await handleRequest(req, fakeFetch, TEST_SECRET);
141
  assertEquals(resp.status, 200);
142
  const body = await resp.json();
143
  assertEquals(body, { code: 0 });
 
149
  fetchCalled = true;
150
  return Promise.resolve(new Response("не должно случиться"));
151
  };
152
+ const req = new Request("https://proxy.example/fetch/evil.example.com/x", { method: "GET", headers: AUTH_HEADERS });
153
+ const resp = await handleRequest(req, fakeFetch, TEST_SECRET);
154
  assertEquals(resp.status, 403);
155
  assertEquals(fetchCalled, false, "fetch не должен был вызываться для неразрешённого хоста");
156
  });
157
 
158
  Deno.test("handleRequest возвращает 404 для пути без /fetch/ префикса", async () => {
159
+ const req = new Request("https://proxy.example/something-else", { method: "GET", headers: AUTH_HEADERS });
160
  const unusedFetch: typeof fetch = () => Promise.resolve(new Response("unused"));
161
+ const resp = await handleRequest(req, unusedFetch, TEST_SECRET);
162
  assertEquals(resp.status, 404);
163
  });
164
 
 
173
  const req = new Request("https://proxy.example/fetch/api.telegram.org/bot123/sendMessage", {
174
  method: "POST",
175
  body: JSON.stringify({ chat_id: 1, text: "hi" }),
176
+ headers: { "Content-Type": "application/json", ...AUTH_HEADERS },
177
  });
178
+ const resp = await handleRequest(req, fakeFetch, TEST_SECRET);
179
  assertEquals(resp.status, 200);
180
  assertEquals(capturedMethod, "POST");
181
  assert(capturedBodyIsStream, "тело POST-запроса должно передаваться апстриму как поток (без буферизации целиком)");
 
183
 
184
  Deno.test("handleRequest возвращает 502, если апстрим-fetch упал с исключением", async () => {
185
  const fakeFetch: typeof fetch = () => {
186
+ throw new Error(`network unreachable ${TEST_SECRET}`);
187
  };
188
+ const req = new Request("https://proxy.example/fetch/api.telegram.org/getMe", { method: "GET", headers: AUTH_HEADERS });
189
+ const resp = await handleRequest(req, fakeFetch, TEST_SECRET);
190
  assertEquals(resp.status, 502);
191
+ assertEquals(await resp.text(), "Upstream fetch failed");
192
+ });
193
+
194
+ for (const suppliedSecret of [undefined, "", "incorrect-test-secret"]) {
195
+ Deno.test(`handleRequest rejects missing or wrong secret: ${String(suppliedSecret)}`, async () => {
196
+ let fetchCalled = false;
197
+ const fakeFetch: typeof fetch = () => {
198
+ fetchCalled = true;
199
+ return Promise.resolve(new Response("unexpected"));
200
+ };
201
+ const headers = new Headers();
202
+ if (suppliedSecret !== undefined) headers.set(PROXY_AUTH_HEADER, suppliedSecret);
203
+ const req = new Request("https://proxy.example/fetch/api.telegram.org/bot123/getMe", { headers });
204
+ const resp = await handleRequest(req, fakeFetch, TEST_SECRET);
205
+ assertEquals(resp.status, 401);
206
+ assertEquals(await resp.text(), "Unauthorized");
207
+ assertEquals(fetchCalled, false);
208
+ });
209
+ }
210
+
211
+ for (const configuredSecret of [undefined, "", " ", "invalid secret"]) {
212
+ Deno.test(`handleRequest fails closed for invalid configuration: ${String(configuredSecret)}`, async () => {
213
+ let fetchCalled = false;
214
+ const fakeFetch: typeof fetch = () => {
215
+ fetchCalled = true;
216
+ return Promise.resolve(new Response("unexpected"));
217
+ };
218
+ const req = new Request("https://proxy.example/fetch/www.tikwm.com/api/", { headers: AUTH_HEADERS });
219
+ const resp = await handleRequest(req, fakeFetch, configuredSecret);
220
+ assertEquals(resp.status, 503);
221
+ assertEquals(await resp.text(), "Proxy authentication unavailable");
222
+ assertEquals(fetchCalled, false);
223
+ });
224
+ }
225
+
226
+ Deno.test("correct secret is stripped without changing Telegram authentication", async () => {
227
+ let fetchCalled = false;
228
+ const fakeFetch: typeof fetch = (url, init) => {
229
+ fetchCalled = true;
230
+ assertEquals(String(url), "https://api.telegram.org/bot123/getMe");
231
+ const headers = new Headers(init?.headers);
232
+ assertEquals(headers.has(PROXY_AUTH_HEADER), false);
233
+ assertEquals(headers.get("authorization"), "Bearer upstream-test-token");
234
+ assertEquals(init?.redirect, "error");
235
+ return Promise.resolve(new Response('{"ok":true}', { status: 200 }));
236
+ };
237
+ const req = new Request("https://proxy.example/fetch/api.telegram.org/bot123/getMe", {
238
+ headers: {
239
+ "x-lumen-proxy-secret": TEST_SECRET,
240
+ "Authorization": "Bearer upstream-test-token",
241
+ },
242
+ });
243
+ const resp = await handleRequest(req, fakeFetch, TEST_SECRET);
244
+ assertEquals(resp.status, 200);
245
+ assertEquals(fetchCalled, true);
246
+ assertEquals(await resp.json(), { ok: true });
247
+ assertEquals(req.headers.get(PROXY_AUTH_HEADER), TEST_SECRET);
248
+ });
249
+
250
+ Deno.test("buildForwardHeaders removes mixed-case proxy secret without mutating input", () => {
251
+ const original = new Headers({ "X-LuMeN-PrOxY-SeCrEt": TEST_SECRET });
252
+ assertEquals(buildForwardHeaders(original).has(PROXY_AUTH_HEADER), false);
253
+ assertEquals(original.get(PROXY_AUTH_HEADER), TEST_SECRET);
254
  });
test_bot.py CHANGED
@@ -26,15 +26,18 @@ BOT_LOG_PATH) ДО импорта bot.py, так что реальные сек
26
  import asyncio
27
  import base64
28
  import json
 
29
  import os
 
30
  import time
31
  from types import SimpleNamespace
32
- from unittest.mock import MagicMock, patch
33
 
34
  import pytest
35
  import sentry_sdk
36
 
37
  import bot
 
38
  import lumen_tiktok
39
 
40
 
@@ -4501,8 +4504,276 @@ def test_download_url_bin_ignores_malformed_content_length_header():
4501
  assert result == b"ok"
4502
 
4503
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
4504
  def test_download_url_bin_returns_none_on_non_200_status():
4505
  resp = _FakeDownloadResponse([b"error page"], status=404)
4506
  session = _FakeDownloadSession(resp)
4507
  result = asyncio.run(lumen_tiktok._download_url_bin(session, "https://tikwm.com/missing.jpg"))
4508
  assert result is None
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
26
  import asyncio
27
  import base64
28
  import json
29
+ import logging
30
  import os
31
+ import sys
32
  import time
33
  from types import SimpleNamespace
34
+ from unittest.mock import AsyncMock, MagicMock, patch
35
 
36
  import pytest
37
  import sentry_sdk
38
 
39
  import bot
40
+ import lumen_telegram_transport
41
  import lumen_tiktok
42
 
43
 
 
4504
  assert result == b"ok"
4505
 
4506
 
4507
+ @pytest.fixture
4508
+ def rate_guard_setup(monkeypatch):
4509
+ def _setup():
4510
+ monkeypatch.setattr(bot, "user_rate_limits", {})
4511
+ monkeypatch.setattr(bot, "RATE_LIMIT_MAX_REQUESTS", 2)
4512
+ monkeypatch.setattr(bot, "get_state", lambda chat_id: {})
4513
+ monkeypatch.setattr(bot, "is_guest_message", lambda message: False)
4514
+ monkeypatch.setattr(bot, "message_mentions_bot", lambda message: False)
4515
+ monkeypatch.setattr(bot, "_safe_reply", AsyncMock())
4516
+ monkeypatch.setattr(bot, "_tg_call", AsyncMock())
4517
+ monkeypatch.setattr(bot, "inline_draw", AsyncMock())
4518
+ monkeypatch.setattr(bot, "inline_tts", AsyncMock())
4519
+ return SimpleNamespace(
4520
+ text="", caption=None, chat=SimpleNamespace(id=123, type=bot.ChatType.PRIVATE),
4521
+ from_user=SimpleNamespace(id=456), reply=AsyncMock(), reply_to_message=None,
4522
+ )
4523
+ return _setup
4524
+
4525
+
4526
+ def test_draw_command_rejects_exhausted_quota(rate_guard_setup, monkeypatch):
4527
+ message = rate_guard_setup()
4528
+ message.text = "/draw кот"
4529
+ monkeypatch.setattr(bot, "_check_and_register_rate_limit", lambda user_id: True)
4530
+
4531
+ asyncio.run(bot.cmd_draw(message))
4532
+
4533
+ bot.inline_draw.assert_not_awaited()
4534
+ bot._tg_call.assert_awaited_once()
4535
+
4536
+
4537
+ def test_mixed_commands_share_rate_quota_and_reject_before_work(rate_guard_setup, monkeypatch):
4538
+ message = rate_guard_setup()
4539
+ media = AsyncMock()
4540
+ monkeypatch.setattr(bot, "_resolve_incoming_media", media)
4541
+
4542
+ async def run():
4543
+ message.text = "/draw кот"
4544
+ await bot.cmd_draw(message)
4545
+ message.text = "/tts привет"
4546
+ await bot.cmd_tts(message)
4547
+ await bot.cmd_tts(message)
4548
+ message.text = "/draw собака"
4549
+ await bot.cmd_draw(message)
4550
+ message.text = "обычный вопрос"
4551
+ await bot._handle_message_core(message)
4552
+
4553
+ asyncio.run(run())
4554
+ bot.inline_draw.assert_awaited_once_with(message, "кот")
4555
+ bot.inline_tts.assert_awaited_once_with(message, "привет")
4556
+ media.assert_not_awaited()
4557
+ assert bot._tg_call.await_count == 3
4558
+ assert len(bot.user_rate_limits[456]) == 2
4559
+
4560
+
4561
+ @pytest.mark.parametrize("command", ["draw", "tts"])
4562
+ def test_blank_commands_do_not_consume_quota(rate_guard_setup, command):
4563
+ message = rate_guard_setup()
4564
+ message.text = f"/{command} "
4565
+ asyncio.run(getattr(bot, f"cmd_{command}")(message))
4566
+ assert bot.user_rate_limits == {}
4567
+ bot._safe_reply.assert_awaited_once()
4568
+ bot.inline_draw.assert_not_awaited()
4569
+ bot.inline_tts.assert_not_awaited()
4570
+
4571
+
4572
+ @pytest.mark.parametrize("text, handler, content", [
4573
+ ("нарисуй кота", "inline_draw", "кота"),
4574
+ ("озвучь привет", "inline_tts", "привет"),
4575
+ ])
4576
+ def test_natural_language_trigger_consumes_one_slot(rate_guard_setup, text, handler, content):
4577
+ message = rate_guard_setup()
4578
+ message.text = text
4579
+ asyncio.run(bot._handle_message_core(message))
4580
+ getattr(bot, handler).assert_awaited_once_with(message, content)
4581
+ assert len(bot.user_rate_limits[456]) == 1
4582
+
4583
+
4584
+ def test_passive_group_message_does_not_consume_quota(rate_guard_setup, monkeypatch):
4585
+ message = rate_guard_setup()
4586
+ message.chat.type = bot.ChatType.GROUP
4587
+ message.text = "обычный разговор"
4588
+ record = MagicMock()
4589
+ monkeypatch.setattr(bot, "_record_passive_group_context", record)
4590
+ asyncio.run(bot._handle_message_core(message))
4591
+ record.assert_called_once()
4592
+ assert bot.user_rate_limits == {}
4593
+ bot._tg_call.assert_not_awaited()
4594
+
4595
+
4596
  def test_download_url_bin_returns_none_on_non_200_status():
4597
  resp = _FakeDownloadResponse([b"error page"], status=404)
4598
  session = _FakeDownloadSession(resp)
4599
  result = asyncio.run(lumen_tiktok._download_url_bin(session, "https://tikwm.com/missing.jpg"))
4600
  assert result is None
4601
+
4602
+
4603
+ class _FakeProc:
4604
+ def __init__(self, *, hang: bool = False, fail_communicate: bool = False):
4605
+ self._hang = hang
4606
+ self._fail_communicate = fail_communicate
4607
+ self.killed = False
4608
+ self.kill_count = 0
4609
+ self.wait_count = 0
4610
+ self.returncode = None
4611
+
4612
+ async def communicate(self):
4613
+ if self._fail_communicate:
4614
+ raise RuntimeError("pipe broken")
4615
+ while self._hang and not self.killed:
4616
+ await asyncio.sleep(0.01)
4617
+ return b"out", b"err"
4618
+
4619
+ def kill(self):
4620
+ self.kill_count += 1
4621
+ self.killed = True
4622
+ self.returncode = -9
4623
+ self._hang = False
4624
+
4625
+ async def wait(self):
4626
+ self.wait_count += 1
4627
+ return self.returncode
4628
+
4629
+
4630
+ def test_communicate_process_returns_output_on_success():
4631
+ proc = _FakeProc()
4632
+ out, err = asyncio.run(lumen_tiktok._communicate_process(proc, timeout=5))
4633
+ assert (out, err) == (b"out", b"err")
4634
+ assert not proc.killed
4635
+
4636
+
4637
+ def test_communicate_process_kills_hung_process_on_timeout():
4638
+ async def run():
4639
+ proc = _FakeProc(hang=True)
4640
+ started = time.monotonic()
4641
+ with pytest.raises(asyncio.TimeoutError):
4642
+ await lumen_tiktok._communicate_process(proc, timeout=0.05)
4643
+ elapsed = time.monotonic() - started
4644
+ assert proc.killed
4645
+ assert proc.wait_count >= 1
4646
+ return elapsed
4647
+
4648
+ elapsed = asyncio.run(run())
4649
+ assert elapsed < 30
4650
+
4651
+
4652
+ def test_communicate_process_cleans_up_when_outer_task_cancelled():
4653
+ async def run():
4654
+ proc = _FakeProc(hang=True)
4655
+ task = asyncio.create_task(lumen_tiktok._communicate_process(proc, timeout=5))
4656
+ await asyncio.sleep(0.01)
4657
+ task.cancel()
4658
+ with pytest.raises(asyncio.CancelledError):
4659
+ await task
4660
+ for _ in range(5):
4661
+ if proc.killed and task.done():
4662
+ break
4663
+ await asyncio.sleep(0.01)
4664
+ assert proc.killed
4665
+ assert task.cancelled()
4666
+
4667
+ asyncio.run(run())
4668
+
4669
+
4670
+ def test_communicate_process_propagates_communicate_failure():
4671
+ async def run():
4672
+ proc = _FakeProc(fail_communicate=True)
4673
+ with pytest.raises(RuntimeError):
4674
+ await lumen_tiktok._communicate_process(proc, timeout=5)
4675
+ assert proc.wait_count >= 1
4676
+
4677
+ asyncio.run(run())
4678
+
4679
+
4680
+ def test_log_queue_handler_masks_secrets_in_args_and_traceback(monkeypatch):
4681
+ monkeypatch.setattr(bot, "LUMEN_PROXY_SECRET", "fake-proxy-secret-123", raising=False)
4682
+ try:
4683
+ raise RuntimeError("leak fake-proxy-secret-123")
4684
+ except RuntimeError:
4685
+ exc_info = sys.exc_info()
4686
+ record = logging.LogRecord(
4687
+ "bot", logging.ERROR, "path", 1,
4688
+ "failed with token %s", ("fake-proxy-secret-123",), exc_info=exc_info,
4689
+ )
4690
+ prepared = bot._LOG_QUEUE_HANDLER.prepare(record)
4691
+ rendered = bot._LOG_QUEUE_HANDLER.format(prepared)
4692
+ assert "fake-proxy-secret-123" not in rendered
4693
+ assert "<REDACTED>" in rendered
4694
+
4695
+
4696
+ def test_secret_log_formatter_masks_message_before_queue(monkeypatch):
4697
+ formatter = bot._SecretLogFormatter("%(message)s")
4698
+ monkeypatch.setattr(bot, "LUMEN_PROXY_SECRET", "fake-formatter-secret-456", raising=False)
4699
+ record = logging.LogRecord(
4700
+ "bot", logging.INFO, "path", 1,
4701
+ "key is %s", ("fake-formatter-secret-456",), None,
4702
+ )
4703
+ assert "fake-formatter-secret-456" not in formatter.format(record)
4704
+ assert "<REDACTED>" in formatter.format(record)
4705
+
4706
+
4707
+ def test_current_log_secrets_reads_env_before_module_globals(monkeypatch):
4708
+ monkeypatch.delenv("LUMEN_PROXY_SECRET", raising=False)
4709
+ monkeypatch.setattr(bot, "LUMEN_PROXY_SECRET", "", raising=False)
4710
+ monkeypatch.setenv("LUMEN_PROXY_SECRET", "fake-env-only-secret-789")
4711
+ secrets = bot._current_log_secrets()
4712
+ assert "fake-env-only-secret-789" in secrets
4713
+ monkeypatch.delenv("LUMEN_PROXY_SECRET", raising=False)
4714
+
4715
+
4716
+ def _run_proxy_middleware(url, *, secret="proxy-secret-abc", bases=("https://proxy.example/fetch/api.telegram.org",), headers=None):
4717
+ from multidict import CIMultiDict
4718
+ from yarl import URL
4719
+
4720
+ (authenticate,) = lumen_telegram_transport.proxy_auth_middlewares(
4721
+ proxy_secret=secret, proxy_base_urls=bases,
4722
+ )
4723
+ seen = {}
4724
+
4725
+ class _FakeRequest:
4726
+ def __init__(self):
4727
+ self.url = URL(url)
4728
+ self.headers = CIMultiDict(headers or {})
4729
+
4730
+ async def _handler(request):
4731
+ seen["sent"] = dict(request.headers)
4732
+ info_headers = CIMultiDict(request.headers)
4733
+
4734
+ class _FakeResponse:
4735
+ status = 200
4736
+
4737
+ def __init__(self):
4738
+ self.request_info = SimpleNamespace(
4739
+ url=request.url, method="GET", headers=info_headers, real_url=request.url,
4740
+ )
4741
+ self.closed = False
4742
+
4743
+ def close(self):
4744
+ self.closed = True
4745
+
4746
+ return _FakeResponse()
4747
+
4748
+ request = _FakeRequest()
4749
+ response = asyncio.run(authenticate(request, _handler))
4750
+ return request, response, seen
4751
+
4752
+
4753
+ def test_proxy_middleware_sends_secret_only_to_configured_proxy():
4754
+ _, _, seen = _run_proxy_middleware("https://proxy.example/fetch/api.telegram.org/bot123/sendMessage")
4755
+ assert seen["sent"].get("X-Lumen-Proxy-Secret") == "proxy-secret-abc"
4756
+
4757
+
4758
+ def test_proxy_middleware_never_sends_secret_to_direct_or_unrelated_hosts():
4759
+ for url in (
4760
+ "https://api.telegram.org/bot123/sendMessage",
4761
+ "https://www.tikwm.com/api/?url=x",
4762
+ "https://cdn.example.com/file.jpg",
4763
+ "https://proxy.example/other-path",
4764
+ "http://proxy.example/fetch/api.telegram.org/bot123/sendMessage",
4765
+ ):
4766
+ _, _, seen = _run_proxy_middleware(url)
4767
+ assert "X-Lumen-Proxy-Secret" not in seen["sent"], url
4768
+
4769
+
4770
+ def test_proxy_middleware_strips_stale_secret_and_requires_secret():
4771
+ _, _, seen = _run_proxy_middleware(
4772
+ "https://api.telegram.org/bot123/sendMessage",
4773
+ headers={"X-Lumen-Proxy-Secret": "stale"},
4774
+ )
4775
+ assert "X-Lumen-Proxy-Secret" not in seen["sent"]
4776
+ with pytest.raises(ValueError):
4777
+ lumen_telegram_transport.proxy_auth_middlewares(
4778
+ proxy_secret="", proxy_base_urls=("https://proxy.example/fetch/api.telegram.org",),
4779
+ )