Lumen / tests /test_bot_streaming.py
SilverElixir
Test quality: assert TikTok error shown, behavioral SSE/trim tests, explicit no-raise (A9-1, A9-2, A9-6)
821f63a
Raw History Blame Contribute Delete
44.5 kB
"""
test_bot_streaming.py — Стриминг: куски, точки ожидания, пейсинг, rich-правки.
Выделено из test_bot.py (P2 аудита); общие фейки — в bot_test_helpers.py.
"""
from types import SimpleNamespace
from unittest.mock import MagicMock
import asyncio
import bot
import lumen_router_config
import lumen_streaming
import pytest
import time
from tests.bot_test_helpers import (
_FakeIncomingMessage,
_FakeRichBot,
_FakeRichFailingBot,
_FakeRichMessage,
_FakeSSEResponse,
_FakeSentMessage,
_FakeSessionForSSE,
)
def test_stream_wait_caps_chunk_timeout_by_remaining_budget():
# Регрессия (аудит 26.09.2026): кусок стрима ждал STREAM_CHUNK_TIMEOUT_SEC даже
# после исчерпания бюджета маршрута — маршрут вылезал за ROUTE_TOTAL_BUDGET_SEC.
assert lumen_streaming._stream_wait(bot, None) == bot.STREAM_CHUNK_TIMEOUT_SEC
far = time.monotonic() + 300.0
assert lumen_streaming._stream_wait(bot, far) == bot.STREAM_CHUNK_TIMEOUT_SEC
near = time.monotonic() + 3.0
assert 0 < lumen_streaming._stream_wait(bot, near) <= 3.0
# Истёкший бюджет — быстрый отказ, а не ожидание в 30с.
assert lumen_streaming._stream_wait(bot, time.monotonic() - 10.0) <= 0.5
def _fake_gemini_stream(pieces=None, *, raises=None, hang_sec=0.0):
"""Мок generate_content_stream по НАСТОЯЩЕМУ контракту google-genai 2.24.0:
обычная функция, возвращающая асинхронный генератор (тело SDK — `return
stream_generator()`), а НЕ корутина.
Раньше фейки здесь были `async def ... return gen()`, из-за чего тесты
подпирали фиктивное рукопожатие wait_for вокруг вызова: реальный SDK сети
в корутине не делает, значит тот таймаут не срабатывал никогда
(враждебное ревью 27.09.2026)."""
def _stream(*, model, contents, config=None):
async def gen():
if hang_sec:
await asyncio.sleep(hang_sec)
if raises is not None:
raise raises
for piece in (pieces or []):
yield SimpleNamespace(text=piece)
return gen()
return _stream
def test_gemini_stream_timeout_covers_a_hanging_first_chunk():
# Настоящая защита: сетевое ожидание происходит на __anext__, и его ловит
# _stream_wait под каждый кусок. Раньше проверялось только создание корутины —
# сценарий, невозможный в проде (см. докстринг _fake_gemini_stream).
fake_client = MagicMock()
fake_client.aio.models.generate_content_stream = _fake_gemini_stream(hang_sec=30)
original_client = bot.client
original_cap = bot.STREAM_CHUNK_TIMEOUT_SEC
bot.client = fake_client
bot.STREAM_CHUNK_TIMEOUT_SEC = 0.2
try:
async def _drain():
async for _ in lumen_streaming._gemini_stream_pieces("m", [], None):
pass
with pytest.raises(asyncio.TimeoutError):
asyncio.run(_drain())
finally:
bot.client = original_client
bot.STREAM_CHUNK_TIMEOUT_SEC = original_cap
def test_gemini_stream_call_is_not_awaited():
# Контракт SDK: generate_content_stream — обычная функция. Если бы обёртка снова
# стала await-ить её, тест падал бы с TypeError, а не проходил вхолостую.
calls = []
def sync_only(*, model, contents, config=None):
calls.append(model)
async def gen():
yield SimpleNamespace(text="ok")
return gen()
fake_client = MagicMock()
fake_client.aio.models.generate_content_stream = sync_only
original_client = bot.client
bot.client = fake_client
try:
async def _drain():
return [p async for p in lumen_streaming._gemini_stream_pieces("m1", [], None)]
assert asyncio.run(_drain()) == ["ok"]
assert calls == ["m1"]
finally:
bot.client = original_client
def test_try_gemini_streaming_happy_path_accumulates_and_finalizes():
chat_id = 999101
fake_stream = _fake_gemini_stream(["Привет", ", как ", "дела?"])
fake_client = MagicMock()
fake_client.aio.models.generate_content_stream = fake_stream
incoming = _FakeIncomingMessage(chat_id)
original_client = bot.client
bot.client = fake_client
try:
answer, placeholder = asyncio.run(bot._try_gemini_streaming(chat_id, "Привет!", incoming, bot.DEFAULT_GEMINI_MODEL))
assert answer == "Привет, как дела?"
assert placeholder is None # успех — плейсхолдер уже отредактирован до финального текста
# финальная правка должна прийти с HTML parse_mode (полная markdown-конвертация)
assert incoming.sent[0].edits[-1][1] == bot.ParseMode.HTML
history = bot.chat_state[chat_id]["history"]
assert history[-1] == {"role": "assistant", "content": "Привет, как дела?"}
finally:
bot.client = original_client
bot.chat_state.pop(chat_id, None)
def test_try_gemini_streaming_aborts_on_identity_leak_mid_stream():
# Регрессия на самый чувствительный сценарий: утечка должна обрываться ДО того,
# как накопленный текст попадёт хоть в один edit_text — иначе пользователь успеет
# увидеть утёкший текст на экране ещё до финального завершения потока.
chat_id = 999104
fake_stream = _fake_gemini_stream(["Привет! ", "На самом деле я работаю ", "на базе Gemini от Google."])
fake_client = MagicMock()
fake_client.aio.models.generate_content_stream = fake_stream
incoming = _FakeIncomingMessage(chat_id)
original_client = bot.client
bot.client = fake_client
try:
answer, placeholder = asyncio.run(bot._try_gemini_streaming(chat_id, "Кто ты на самом деле?", incoming, bot.DEFAULT_GEMINI_MODEL))
assert answer == bot._IDENTITY_LEAK_FALLBACK
assert placeholder is None
# Ни в одной показанной пользователю правке НЕ должно быть слова "gemini" —
# проверяем ВСЕ edit_text вызовы единственного отправленного сообщения, а не
# только последний, т.к. именно промежуточные правки могли бы "мигнуть" утечкой.
for shown_text, _parse_mode in incoming.sent[0].edits:
assert "gemini" not in shown_text.lower()
history = bot.chat_state[chat_id]["history"]
assert history[-1] == {"role": "assistant", "content": bot._IDENTITY_LEAK_FALLBACK}
finally:
bot.client = original_client
bot.chat_state.pop(chat_id, None)
def test_try_gemini_streaming_returns_none_on_early_failure():
chat_id = 999102
fake_stream_raises = _fake_gemini_stream(raises=RuntimeError("boom before any content"))
fake_client = MagicMock()
fake_client.aio.models.generate_content_stream = fake_stream_raises
incoming = _FakeIncomingMessage(chat_id)
original_client = bot.client
bot.client = fake_client
try:
answer, placeholder = asyncio.run(bot._try_gemini_streaming(chat_id, "Привет!", incoming, bot.DEFAULT_GEMINI_MODEL))
assert answer is None
# Плейсхолдер теперь НЕ удаляется на этом уровне — он возвращается
# вызывающему коду (_run_route), чтобы тот попробовал доправить в него
# ответ следующей модели по цепочке, а не создавать новое сообщение.
assert placeholder is incoming.sent[0]
assert incoming.sent[0].deleted is False
# история НЕ должна была обновиться — вызывающий код откатится на ask_gemini
assert chat_id not in bot.chat_state or not bot.chat_state[chat_id].get("history")
finally:
bot.client = original_client
bot.chat_state.pop(chat_id, None)
def test_try_gemini_streaming_failed_continuation_does_not_corrupt_first_message():
# Регрессионный тест на баг, найденный код-ревью: при сбое отправки сообщения-
# продолжения (текст длиннее лимита Telegram) первое, уже корректно показанное
# сообщение раньше перезаписывалось чужим (последним) куском текста.
chat_id = 999103
long_piece = "А" * (bot.TG_MAX_LEN + 100) # гарантированно требует второе сообщение
fake_stream = _fake_gemini_stream([long_piece])
fake_client = MagicMock()
fake_client.aio.models.generate_content_stream = fake_stream
class _FailingBot:
async def send_message(self, **kwargs):
return None # имитируем неудачную отправку продолжения
incoming = _FakeIncomingMessage(chat_id)
original_client = bot.client
original_bot = bot.bot
bot.client = fake_client
bot.bot = _FailingBot()
try:
answer, placeholder = asyncio.run(bot._try_gemini_streaming(chat_id, "Напиши длинный текст", incoming, bot.DEFAULT_GEMINI_MODEL))
# Функция должна была вернуть накопленный текст, а не None и не бросить исключение
assert answer is not None
assert placeholder is None # это уже финализированный успех, а не ранний сбой
# Первое сообщение должно содержать ИМЕННО первый кусок (плюс пометка об обрыве),
# а НЕ последний/другой кусок текста — это и была суть бага.
first_msg_final_text = incoming.sent[0].edits[-1][0]
assert first_msg_final_text.startswith("А")
assert "couldn't send the rest of the message" in first_msg_final_text
finally:
bot.client = original_client
bot.bot = original_bot
bot.chat_state.pop(chat_id, None)
def test_openrouter_stream_pieces_parses_sse_chunks():
lines = [
'data: {"choices":[{"delta":{"content":"Привет"}}]}\n'.encode("utf-8"),
'data: {"choices":[{"delta":{"content":", мир"}}]}\n'.encode("utf-8"),
b"data: [DONE]\n",
]
fake_resp = _FakeSSEResponse(lines)
fake_session = _FakeSessionForSSE(fake_resp)
async def fake_get_http_session():
return fake_session
original_get_session = bot._get_http_session
original_key = bot.OPENROUTER_API_KEY
bot._get_http_session = fake_get_http_session
bot.OPENROUTER_API_KEY = "fake-key"
try:
async def collect():
pieces = []
async for piece in bot._openrouter_stream_pieces("meta-llama/llama-3.3-70b-instruct:free", [{"role": "user", "content": "hi"}]):
pieces.append(piece)
return pieces
pieces = asyncio.run(collect())
assert pieces == ["Привет", ", мир"]
finally:
bot._get_http_session = original_get_session
bot.OPENROUTER_API_KEY = original_key
def test_groq_stream_pieces_parses_sse_chunks():
# Groq-подключение 21.09.2026: тот же OpenAI-SSE, что у OpenRouter — куски собираются, [DONE] завершает.
lines = [
'data: {"choices":[{"delta":{"content":"Кан"}}]}\n'.encode("utf-8"),
'data: {"choices":[{"delta":{"content":"берра"}}]}\n'.encode("utf-8"),
b"data: [DONE]\n",
]
fake_resp = _FakeSSEResponse(lines)
fake_session = _FakeSessionForSSE(fake_resp)
async def fake_get_http_session():
return fake_session
original_get_session = bot._get_http_session
original_key = bot.GROQ_API_KEY
bot._get_http_session = fake_get_http_session
bot.GROQ_API_KEY = "fake-key"
try:
async def collect():
pieces = []
async for piece in bot._groq_stream_pieces("qwen/qwen3.8-27b", [{"role": "user", "content": "hi"}]):
pieces.append(piece)
return pieces
pieces = asyncio.run(collect())
assert pieces == ["Кан", "берра"]
finally:
bot._get_http_session = original_get_session
bot.GROQ_API_KEY = original_key
def test_openrouter_stream_pieces_raises_on_http_error_status():
fake_resp = _FakeSSEResponse([], status=500)
fake_session = _FakeSessionForSSE(fake_resp)
async def fake_get_http_session():
return fake_session
original_get_session = bot._get_http_session
original_key = bot.OPENROUTER_API_KEY
bot._get_http_session = fake_get_http_session
bot.OPENROUTER_API_KEY = "fake-key"
try:
async def collect():
async for _ in bot._openrouter_stream_pieces("meta-llama/llama-3.3-70b-instruct:free", []):
pass
with pytest.raises(bot.OpenRouterAPIError):
asyncio.run(collect())
finally:
bot._get_http_session = original_get_session
bot.OPENROUTER_API_KEY = original_key
def test_openrouter_stream_pieces_raises_on_midstream_error_chunk():
# РЕГРЕССИЯ (аудит стриминга): провайдер за OpenRouter может упасть УЖЕ ПОСЛЕ
# старта генерации — HTTP-статус к этому моменту давно 200 (стрим открыт), и
# ошибка приходит не кодом ответа, а прямо внутри SSE-чанка:
# {"error": {...}} вместо {"choices": [...]}. Раньше это тихо пропускалось
# как "пустой чанк" (choices нет -> continue) — пользователь получал молча
# укороченный ответ без единого намёка на причину. Первый кусок ("Начало")
# должен успеть уйти до ошибки — проверяем, что она поднимается уже ПОСЛЕ
# частичного контента, а не глушится.
lines = [
'data: {"choices":[{"delta":{"content":"Начало"}}]}\n'.encode("utf-8"),
'data: {"error":{"message":"Provider returned error","code":502}}\n'.encode("utf-8"),
]
fake_resp = _FakeSSEResponse(lines)
fake_session = _FakeSessionForSSE(fake_resp)
async def fake_get_http_session():
return fake_session
original_get_session = bot._get_http_session
original_key = bot.OPENROUTER_API_KEY
bot._get_http_session = fake_get_http_session
bot.OPENROUTER_API_KEY = "fake-key"
try:
collected = []
async def collect_partial():
agen = bot._openrouter_stream_pieces("meta-llama/llama-3.3-70b-instruct:free", [{"role": "user", "content": "hi"}])
async for piece in agen:
collected.append(piece)
with pytest.raises(bot.OpenRouterAPIError) as exc_info:
asyncio.run(collect_partial())
assert collected == ["Начало"]
assert exc_info.value.status_code == 502
finally:
bot._get_http_session = original_get_session
bot.OPENROUTER_API_KEY = original_key
def test_sse_parsers_share_one_implementation(monkeypatch):
# Было grep-тестом по исходникам (запрещённый класс): теперь проверяем
# поведение — обе обёртки разбирают SSE через общий _sse_pieces.
real_sse = lumen_streaming._sse_pieces
calls = []
async def _counting_sse(*args, **kwargs):
calls.append(1)
async for piece in real_sse(*args, **kwargs):
yield piece
monkeypatch.setattr(lumen_streaming, "_sse_pieces", _counting_sse)
lines = [
'data: {"choices":[{"delta":{"content":"Раз"}}]}\n'.encode("utf-8"),
'data: {"choices":[{"delta":{"content":" два"}}]}\n'.encode("utf-8"),
b"data: [DONE]\n",
]
async def _fake_session():
return _FakeSessionForSSE(_FakeSSEResponse(lines))
async def collect(gen):
return [piece async for piece in gen]
monkeypatch.setattr(bot, "OPENROUTER_API_KEY", "fake-key")
monkeypatch.setattr(bot, "GROQ_API_KEY", "fake-key")
monkeypatch.setattr(bot, "_get_http_session", _fake_session)
or_pieces = asyncio.run(collect(bot._openrouter_stream_pieces("m:free", [])))
groq_pieces = asyncio.run(collect(bot._groq_stream_pieces("m", [])))
assert or_pieces == ["Раз", " два"]
assert groq_pieces == ["Раз", " два"]
assert len(calls) == 2
def test_streaming_history_goes_through_summarizing_trim(monkeypatch):
# Поведенческая пара grep-теста test_streaming_history_uses_summarizing_trim:
# длинная история в стриминге идёт через _trim_history с саммари,
# а не молчаливым срезом.
chat_id = 999981
state = bot.get_state(chat_id)
state["history"] = [{"role": "user", "content": f"q{i}"} for i in range(120)]
async def fake_pieces():
yield "готово"
calls = {}
async def fake_trim(hist):
calls["n"] = calls.get("n", 0) + 1
hist[:] = hist[-10:]
monkeypatch.setattr(bot, "_trim_history", fake_trim)
incoming = _FakeIncomingMessage(chat_id)
try:
answer, _ = asyncio.run(bot._run_streaming_reply(
chat_id, "вопрос", incoming, provider="openrouter", model_id="trim:free",
piece_agen=fake_pieces(),
))
assert answer == "готово"
assert calls.get("n", 0) >= 1
finally:
bot.chat_state.pop(chat_id, None)
@pytest.mark.parametrize("deadline_mode", ["none", "expired"])
def test_sse_pieces_hanging_post_entry_times_out(deadline_mode, monkeypatch):
# Регрессия A4-01: вход в session.post шёл с total=None без общего лимита —
# зависший провайдер держал per-chat lock вечно. Теперь заголовки ждут не
# дольше попытки маршрута, висящий вход быстро даёт TimeoutError.
class _HangingPostCM:
def __init__(self):
self.exited = False
async def __aenter__(self):
await asyncio.sleep(3600)
return self
async def __aexit__(self, *args):
self.exited = True
return False
class _HangingPostSession:
def __init__(self, cm):
self._cm = cm
def post(self, *args, **kwargs):
return self._cm
monkeypatch.setattr(bot, "ROUTE_MODEL_TIMEOUT_SEC", 0.2)
cm = _HangingPostCM()
session = _HangingPostSession(cm)
deadline = None if deadline_mode == "none" else time.monotonic() - 1.0
async def collect():
async for _ in lumen_streaming._sse_pieces(
session, "http://example.test/chat", {}, {},
err_cls=bot.OpenRouterAPIError, provider_label="OpenRouter", deadline=deadline,
):
pass
started = time.monotonic()
with pytest.raises(asyncio.TimeoutError):
asyncio.run(asyncio.wait_for(collect(), timeout=5.0))
elapsed = time.monotonic() - started
# Старый код висел все 5с внешнего wait_for; новый обрывает вход за ~0.2-0.5с.
assert elapsed < 2.5, f"вход в POST ждал мимо лимита попытки: {elapsed:.1f}с"
assert cm.exited
def test_groq_stream_pieces_raises_on_http_error_status():
# У Groq раньше не было теста на HTTP>=400 (ветка-копия без покрытия, аудит 26.09.2026).
fake_resp = _FakeSSEResponse([], status=500)
fake_session = _FakeSessionForSSE(fake_resp)
async def fake_get_http_session():
return fake_session
original_get_session = bot._get_http_session
original_key = bot.GROQ_API_KEY
bot._get_http_session = fake_get_http_session
bot.GROQ_API_KEY = "fake-key"
try:
async def collect():
async for _ in bot._groq_stream_pieces("qwen/qwen3.8-27b", []):
pass
with pytest.raises(bot.GroqAPIError):
asyncio.run(collect())
finally:
bot._get_http_session = original_get_session
bot.GROQ_API_KEY = original_key
def test_groq_stream_pieces_raises_on_midstream_error_chunk():
# Та же midstream-ошибка, что у OpenRouter, но в копии Groq её не ловили.
lines = [
'data: {"choices":[{"delta":{"content":"Начало"}}]}\n'.encode("utf-8"),
'data: {"error":{"message":"Provider returned error","code":503}}\n'.encode("utf-8"),
]
fake_resp = _FakeSSEResponse(lines)
fake_session = _FakeSessionForSSE(fake_resp)
async def fake_get_http_session():
return fake_session
original_get_session = bot._get_http_session
original_key = bot.GROQ_API_KEY
bot._get_http_session = fake_get_http_session
bot.GROQ_API_KEY = "fake-key"
try:
collected = []
async def collect_partial():
agen = bot._groq_stream_pieces("qwen/qwen3.8-27b", [{"role": "user", "content": "hi"}])
async for piece in agen:
collected.append(piece)
with pytest.raises(bot.GroqAPIError) as exc_info:
asyncio.run(collect_partial())
assert collected == ["Начало"]
assert exc_info.value.status_code == 503
finally:
bot._get_http_session = original_get_session
bot.GROQ_API_KEY = original_key
def test_try_openrouter_streaming_happy_path_accumulates_and_finalizes():
chat_id = 999105
async def fake_stream_pieces(model_id, messages, *, deadline=None):
for piece in ["Привет", ", как ", "дела?"]:
yield piece
incoming = _FakeIncomingMessage(chat_id)
original_gen = bot._openrouter_stream_pieces
bot._openrouter_stream_pieces = fake_stream_pieces
try:
answer, placeholder = asyncio.run(bot._try_openrouter_streaming(chat_id, "Привет!", incoming, "meta-llama/llama-3.3-70b-instruct:free"))
assert answer == "Привет, как дела?"
assert placeholder is None
assert incoming.sent[0].edits[-1][1] == bot.ParseMode.HTML
history = bot.chat_state[chat_id]["history"]
assert history[-1] == {"role": "assistant", "content": "Привет, как дела?"}
assert bot.GLOBAL_QUOTA["openrouter"]["meta-llama/llama-3.3-70b-instruct:free"]["used"] >= 1
finally:
bot._openrouter_stream_pieces = original_gen
bot.chat_state.pop(chat_id, None)
def test_streaming_abandons_hung_first_chunk_within_limit(monkeypatch):
# Висящий первый кусок — TimeoutError и плейсхолдер дальше по цепочке (раньше предела не было — лок держался до 30с/навсегда).
#
# Патчим lumen_streaming, а НЕ bot: _run_streaming_reply читает имя из своего
# модуля, и прошлый тест патчил bot._model_first_chunk_limit вхолостую — он
# честно ждал штатные 25с и проходил, потому что граница была 30с
# (враждебное ревью 27.09.2026). Теперь предел жёстче выставленного лимита.
chat_id = 999301
async def hanging_pieces():
await asyncio.sleep(3600)
yield "never arrives"
monkeypatch.setattr(lumen_streaming, "_model_first_chunk_limit", lambda key, floor: 0.05)
incoming = _FakeIncomingMessage(chat_id)
try:
started = time.monotonic()
answer, placeholder = asyncio.run(bot._run_streaming_reply(
chat_id, "Привет!", incoming, provider="openrouter", model_id="x:free",
piece_agen=hanging_pieces(),
))
elapsed = time.monotonic() - started
assert answer is None
assert placeholder is incoming.sent[0]
# 0.05с лимита + небольшой запас на ввод-вывод фейков. Прежние 30с
# пропускали настоящие 25с ожидания, то есть тест ничего не проверял.
assert elapsed < 5, f"тест ждёт настоящего таймаута, а не подменённого: {elapsed:.1f}с"
finally:
bot.chat_state.pop(chat_id, None)
def test_streaming_respects_route_deadline():
# Внешний аудит: капающий по куску стрим жил мимо ROUTE_TOTAL_BUDGET_SEC и держал lock.
chat_id = 999309
async def slow_drip_pieces():
yield "начало"
await asyncio.sleep(3600)
yield "никогда"
incoming = _FakeIncomingMessage(chat_id)
try:
started = time.monotonic()
answer, placeholder = asyncio.run(bot._run_streaming_reply(
chat_id, "Привет!", incoming, provider="openrouter", model_id="y:free",
piece_agen=slow_drip_pieces(), deadline=time.monotonic() - 1.0,
))
elapsed = time.monotonic() - started
# Бюджет уже прошёл: либо плейсхолдер дальше, либо финал с пометкой — но не вечное ожидание.
assert placeholder is not None or answer is not None
# Висящий второй кусок (3600с) обязан оборваться бюджетом, а не зависнуть.
assert elapsed < 30, f"стрим ждал висящий кусок мимо дедлайна: {elapsed:.1f}с"
finally:
bot.chat_state.pop(chat_id, None)
def _streaming_rate_limit_quotas(chat_id, model_id, error_text):
async def pieces_429():
raise bot.OpenRouterAPIError(error_text, status_code=429)
yield ""
incoming = _FakeIncomingMessage(chat_id)
try:
answer, _ = asyncio.run(bot._run_streaming_reply(
chat_id, "Привет!", incoming, provider="openrouter", model_id=model_id,
piece_agen=pieces_429(),
))
assert answer is None
return dict(bot.GLOBAL_QUOTA["openrouter"][model_id])
finally:
bot.chat_state.pop(chat_id, None)
def test_streaming_daily_rate_limit_marks_model_exhausted():
# Суточный лимит аккаунта помечает модель до утра — как и раньше.
try:
entry = _streaming_rate_limit_quotas(999310, "z:free", "free-models-per-day limit reached")
assert entry["exhausted_at"] is not None
assert not entry.get("cooldown_until")
finally:
bot.GLOBAL_QUOTA.get("openrouter", {}).pop("z:free", None)
def test_streaming_burst_rate_limit_only_cools_down():
# Регрессия (враждебное ревью 27.09.2026): минутный всплеск 429 ставил суточную
# метку, и модель выпадала из роута до полуночи — снять метку было нечем, потому
# что заведомо мёртвую модель не зовут. Теперь это короткая остывка.
provider = bot.GLOBAL_QUOTA.setdefault("openrouter", {})
had_model = "y:free" in provider
old_model = dict(provider.get("y:free", {}))
try:
entry = _streaming_rate_limit_quotas(999311, "y:free", "Rate limit reached")
assert entry["exhausted_at"] is None
assert entry["cooldown_until"] > time.time()
assert lumen_router_config._is_quota_exhausted("openrouter", "y:free") is True
# По истечении остывки модель возвращается в роут сама, без смены суток.
provider["y:free"]["cooldown_until"] = time.time() - 1.0
assert lumen_router_config._is_quota_exhausted("openrouter", "y:free") is False
finally:
if had_model:
provider["y:free"] = old_model
else:
provider.pop("y:free", None)
def test_waiting_dots_cycles_frames_then_stops_on_cancel(monkeypatch):
# Юнит-тест самого тикера: первый кадр только после _DOTS_START_AFTER_SEC,
# дальше по кадру каждые _DOTS_TICK_SEC; отмена — штатная остановка.
calls = []
async def fake_sleep(delay):
calls.append(delay)
if len(calls) > 3:
raise asyncio.CancelledError
monkeypatch.setattr(bot, "_dots_sleep", fake_sleep)
placeholder = _FakeSentMessage()
async def run():
with pytest.raises(asyncio.CancelledError):
await bot._tick_waiting_dots(placeholder)
asyncio.run(run())
assert calls[0] == bot._DOTS_START_AFTER_SEC
assert calls[1:] == [bot._DOTS_TICK_SEC] * 3
assert [text for text, _ in placeholder.edits] == list(bot._DOTS_FRAMES[:3])
def test_streaming_fast_path_shows_no_dots_frames():
# Мгновенно ответившая модель: тикер гасится до первого кадра (2.5с тишины
# не наступает) — в правках только контент, никаких "." / "..".
chat_id = 999302
async def instant_pieces():
yield "Привет"
incoming = _FakeIncomingMessage(chat_id)
try:
answer, _ = asyncio.run(bot._run_streaming_reply(
chat_id, "Привет!", incoming, provider="openrouter", model_id="y:free",
piece_agen=instant_pieces(),
))
assert answer == "Привет"
for text, _ in incoming.sent[0].edits:
assert text not in bot._DOTS_FRAMES
finally:
bot.chat_state.pop(chat_id, None)
def test_streaming_whitespace_only_returns_none_not_empty_response():
# Стрим из одних пробелов — тоже "ничего не прислал": плейсхолдер уезжает
# дальше по цепочке, а не превращается в "Empty response" для пользователя.
chat_id = 999305
async def whitespace_pieces():
yield " "
incoming = _FakeIncomingMessage(chat_id)
try:
answer, placeholder = asyncio.run(bot._run_streaming_reply(
chat_id, "Привет!", incoming, provider="openrouter", model_id="z:free",
piece_agen=whitespace_pieces(),
))
assert answer is None
assert placeholder is incoming.sent[0]
finally:
bot.chat_state.pop(chat_id, None)
def test_streaming_reveal_follows_arrival_pace_not_full_dump():
# Куски капают постепенно (10 × 20 симв. с паузами): первая правка обязана показать
# ЧАСТЬ ответа, а не весь текст разом — часы показа идут от первого куска.
chat_id = 999308
full = "x" * 200
async def paced_pieces():
for _ in range(10):
await asyncio.sleep(0.05)
yield "x" * 20
incoming = _FakeIncomingMessage(chat_id)
try:
answer, _ = asyncio.run(bot._run_streaming_reply(
chat_id, "Привет!", incoming, provider="openrouter", model_id="pace:free",
piece_agen=paced_pieces(),
))
assert answer == full
first_edit = incoming.sent[0].edits[0][0]
assert 0 < len(first_edit) < len(full)
finally:
bot.chat_state.pop(chat_id, None)
def test_try_openrouter_streaming_returns_none_on_early_failure():
chat_id = 999106
async def fake_stream_pieces_raises(model_id, messages, *, deadline=None):
raise RuntimeError("boom before any content")
yield "" # делает функцию async-генератором (недостижимо)
incoming = _FakeIncomingMessage(chat_id)
original_gen = bot._openrouter_stream_pieces
bot._openrouter_stream_pieces = fake_stream_pieces_raises
try:
answer, placeholder = asyncio.run(bot._try_openrouter_streaming(chat_id, "Привет!", incoming, "meta-llama/llama-3.3-70b-instruct:free"))
assert answer is None
assert placeholder is incoming.sent[0]
assert incoming.sent[0].deleted is False
assert chat_id not in bot.chat_state or not bot.chat_state[chat_id].get("history")
finally:
bot._openrouter_stream_pieces = original_gen
bot.chat_state.pop(chat_id, None)
def test_run_streaming_reply_paces_reveal_for_burst_instead_of_dumping_full_text():
import lumen_typing_pace
chat_id = 999110
pace_key = lumen_typing_pace.speed_key("gemini", bot.DEFAULT_GEMINI_MODEL)
original_ema = lumen_typing_pace._speed_ema.pop(pace_key, None)
# Реалистичная имитация бэкенда, который не стримит токен-в-токен, а отдаёт
# весь ответ ОДНИМ куском (см. докстринг lumen_typing_pace.py про то, почему
# это обычное дело для бесплатных моделей OpenRouter).
long_text = "Слово " * 80 # ~480 символов одним SSE-куском
fake_stream = _fake_gemini_stream([long_text])
fake_client = MagicMock()
fake_client.aio.models.generate_content_stream = fake_stream
incoming = _FakeIncomingMessage(chat_id)
original_client = bot.client
bot.client = fake_client
try:
answer, placeholder = asyncio.run(bot._try_gemini_streaming(chat_id, "Привет!", incoming, bot.DEFAULT_GEMINI_MODEL))
assert answer == long_text.strip()
assert placeholder is None
plain_edits = [text for text, parse_mode in incoming.sent[0].edits if parse_mode is None]
# Хотя бы одна промежуточная правка должна была показать ЧАСТЬ текста, а
# не весь ответ разом — иначе пейсинг не сработал (регрессия на исходную
# проблему: "скачками по 15-20 слов" вместо плавного набора).
assert any(0 < len(p) < len(long_text) for p in plain_edits)
# Финальная правка — уже HTML с полным текстом, как и раньше.
assert incoming.sent[0].edits[-1][1] == bot.ParseMode.HTML
history = bot.chat_state[chat_id]["history"]
assert history[-1] == {"role": "assistant", "content": long_text.strip()}
finally:
bot.client = original_client
bot.chat_state.pop(chat_id, None)
if original_ema is not None:
lumen_typing_pace._speed_ema[pace_key] = original_ema
else:
lumen_typing_pace._speed_ema.pop(pace_key, None)
def test_run_streaming_reply_records_observed_speed_on_success():
import lumen_typing_pace
chat_id = 999111
pace_key = lumen_typing_pace.speed_key("gemini", bot.DEFAULT_GEMINI_MODEL)
original_ema = lumen_typing_pace._speed_ema.pop(pace_key, None)
fake_stream = _fake_gemini_stream(["Привет", ", мир!"])
fake_client = MagicMock()
fake_client.aio.models.generate_content_stream = fake_stream
incoming = _FakeIncomingMessage(chat_id)
original_client = bot.client
bot.client = fake_client
try:
asyncio.run(bot._try_gemini_streaming(chat_id, "Привет!", incoming, bot.DEFAULT_GEMINI_MODEL))
# После успешного стрима у модели должна появиться собственная запись в
# EMA — самокалибровка происходит без единой ручной правки таблицы.
assert pace_key in lumen_typing_pace._speed_ema
finally:
bot.client = original_client
bot.chat_state.pop(chat_id, None)
if original_ema is not None:
lumen_typing_pace._speed_ema[pace_key] = original_ema
else:
lumen_typing_pace._speed_ema.pop(pace_key, None)
def test_run_streaming_reply_catchup_never_exceeds_max_ticks_even_for_long_slow_burst():
# РЕГРЕССИЯ на ключевое требование: сколько бы ни оценивалась скорость модели
# заниженно и сколько бы символов ни осталось "довыводить" одним куском, число
# искусственных пауз ограничено STREAM_TYPING_MAX_CATCHUP_TICKS — реальная
# скорость ответа не должна страдать ради красивости набора текста.
import lumen_typing_pace
chat_id = 999112
pace_key = lumen_typing_pace.speed_key("gemini", bot.DEFAULT_GEMINI_MODEL)
original_ema = lumen_typing_pace._speed_ema.pop(pace_key, None)
lumen_typing_pace._speed_ema[pace_key] = lumen_typing_pace.MIN_CHARS_PER_SEC # намеренно "медленная" модель
long_text = "Буква " * 500 # ~3000 символов, одним куском, ниже TG_MAX_LEN
fake_stream = _fake_gemini_stream([long_text])
fake_client = MagicMock()
fake_client.aio.models.generate_content_stream = fake_stream
incoming = _FakeIncomingMessage(chat_id)
sleep_calls: list[float] = []
original_sleep = bot._typing_sleep
async def counting_sleep(seconds: float) -> None:
sleep_calls.append(seconds)
bot._typing_sleep = counting_sleep
original_client = bot.client
bot.client = fake_client
try:
answer, _ = asyncio.run(bot._try_gemini_streaming(chat_id, "Привет!", incoming, bot.DEFAULT_GEMINI_MODEL))
assert answer == long_text.strip()
assert len(sleep_calls) <= bot.STREAM_TYPING_MAX_CATCHUP_TICKS
finally:
bot._typing_sleep = original_sleep
bot.client = original_client
bot.chat_state.pop(chat_id, None)
if original_ema is not None:
lumen_typing_pace._speed_ema[pace_key] = original_ema
else:
lumen_typing_pace._speed_ema.pop(pace_key, None)
def test_rich_send_used_for_final_answer_with_table():
chat_id = 999401
incoming = _FakeIncomingMessage(chat_id)
original_bot = bot.bot
bot.bot = _FakeRichBot()
try:
asyncio.run(bot._safe_reply(incoming, "| A | B |\n|---|---|\n| 1 | 2 |"))
assert len(bot.bot.rich_sent) == 1
assert bot.bot.rich_sent[0]["rich_message"].html.startswith("<table bordered>")
assert incoming.sent == []
finally:
bot.bot = original_bot
bot.chat_state.pop(chat_id, None)
def test_rich_failure_falls_back_to_legacy_html():
chat_id = 999402
incoming = _FakeIncomingMessage(chat_id)
original_bot = bot.bot
bot.bot = _FakeRichFailingBot()
try:
asyncio.run(bot._safe_reply(incoming, "**жирный** текст"))
assert len(incoming.sent) == 1
assert incoming.sent[0].edits == []
finally:
bot.bot = original_bot
bot.chat_state.pop(chat_id, None)
def test_rich_disabled_flag_uses_legacy_path(monkeypatch):
chat_id = 999403
incoming = _FakeIncomingMessage(chat_id)
original_bot = bot.bot
bot.bot = _FakeRichBot()
monkeypatch.setattr(bot, "RICH_MESSAGES_ENABLED", False)
try:
asyncio.run(bot._safe_reply(incoming, "**жирный** текст"))
assert bot.bot.rich_sent == []
assert len(incoming.sent) == 1
finally:
bot.bot = original_bot
bot.chat_state.pop(chat_id, None)
def test_rich_edit_used_for_final_message_edit():
msg = _FakeRichMessage()
original_bot = bot.bot
bot.bot = _FakeRichBot()
try:
assert asyncio.run(bot._edit_message_quietly(msg, "## Заголовок")) is True
assert len(bot.bot.rich_edited) == 1
assert bot.bot.rich_edited[0]["rich_message"].html == "<h3>Заголовок</h3>"
assert msg.edits == []
finally:
bot.bot = original_bot
def test_rich_edit_falls_back_to_legacy_on_failure():
msg = _FakeRichMessage()
original_bot = bot.bot
bot.bot = _FakeRichFailingBot()
try:
assert asyncio.run(bot._edit_message_quietly(msg, "**жирный**")) is True
assert msg.edits and msg.edits[0][0] == "<b>жирный</b>"
finally:
bot.bot = original_bot
def test_streaming_limits_chunk_resplits(monkeypatch):
# Аудит A4-11: _split_text_chunks на каждый кусок давал O(n²) в loop —
# теперь пересчёт только при заметном приросте текста.
real_split = lumen_streaming._split_text_chunks
calls = {"n": 0}
def _counting_split(text, limit):
calls["n"] += 1
return real_split(text, limit)
monkeypatch.setattr(lumen_streaming, "_split_text_chunks", _counting_split)
chat_id = 999971
async def many_pieces():
for i in range(30):
yield f"кусок{i:02d} " + "x" * 60
incoming = _FakeIncomingMessage(chat_id)
try:
answer, _ = asyncio.run(bot._run_streaming_reply(
chat_id, "Привет!", incoming, provider="openrouter", model_id="split:free",
piece_agen=many_pieces(),
))
assert answer is not None and "кусок29" in answer
assert calls["n"] < 30, f"разбиение на каждый кусок: {calls['n']}"
finally:
bot.chat_state.pop(chat_id, None)
def test_streaming_falls_back_to_reply_when_edits_die(monkeypatch):
# Аудит A4-12: сообщение снесли посреди стрима — серия неуспешных правок
# ведёт к досылке ответа новым сообщением, а не к истории под невидимый текст.
chat_id = 999972
incoming = _FakeIncomingMessage(chat_id)
async def many_big_pieces():
for i in range(5):
yield f"часть{i} " + "y" * 5000
async def _failing_edits_tg_call(method, *args, **kwargs):
if getattr(method, "__name__", "") in ("reply", "send_message"):
return await incoming.reply("…")
return None
async def send_message(*args, **kwargs):
return await incoming.reply("…")
sent_fallback = {}
async def _recorder_fallback(message, text, **kwargs):
sent_fallback["text"] = text
monkeypatch.setattr(bot, "bot", SimpleNamespace(send_message=send_message))
monkeypatch.setattr(bot, "_tg_call", _failing_edits_tg_call)
monkeypatch.setattr(bot, "_safe_reply", _recorder_fallback)
try:
answer, _ = asyncio.run(bot._run_streaming_reply(
chat_id, "Привет!", incoming, provider="openrouter", model_id="dead:edit",
piece_agen=many_big_pieces(),
))
assert answer is not None and "часть4" in answer
assert "часть4" in sent_fallback.get("text", "")
history = bot.chat_state[chat_id]["history"]
assert history[-1]["content"] == answer
finally:
bot.chat_state.pop(chat_id, None)