aiogram-stream-sender
(ридми написан ии)
Ваша LLM выдаёт текст по токену. Telegram разрешает править сообщение примерно раз в секунду и очень не любит, когда его не слушают. Между этими двумя фактами и живёт эта библиотека.
Она берёт на себя всё скучное: троттлинг, 429, разбиение длинного ответа на несколько сообщений, ретраи с бэкоффом, «печатает…», удаление лишних хвостов и корректное завершение при остановке бота. Вы просто говорите «сейчас текст выглядит вот так» — столько раз, сколько хотите.
async with sender.stream() as stream:
async for delta in llm:
buffer += delta
stream.update([{"text": buffer}])
update — синхронный и дешёвый. Он не отправляет ничего, он только меняет желаемое состояние. Что и когда полетит в Telegram, решает планировщик. Хоть тысячу раз в секунду зовите — в чат уйдёт ровно столько запросов, сколько разрешено.
Идея в одном абзаце
Библиотека построена вокруг разделения намерения и исполнения. Есть «как должно быть» (список чанков, который вы обновляете) и «как есть» (что реально доставлено в Telegram). Между ними — конечный автомат, который постоянно считает разницу и выдаёт по одному действию за раз: отправить, отредактировать, удалить, показать typing. Автомат ничего не знает про сеть и не умеет спать — он чистая функция от состояния и времени. Поэтому его можно протестировать целиком за миллисекунды, подсовывая фейковое now.
Дальше — сверху вниз, от того, что вы настраиваете, до того, что реально уходит в aiogram.
Options — все ручки сразу
Единственный конфиг на весь рантайм. Frozen dataclass, дефолты подобраны так, чтобы «просто работало».
Ритм. send_interval, edit_interval, delete_interval (по 1.0 с) и action_interval (5.0 с). Это минимальные паузы между действиями одного вида в одном чате. Виды независимы: send не мешает edit. Telegram считает лимиты примерно так же.
Реакция на сбои. backoff_base 1.0, backoff_max 30.0, backoff_jitter 0.2, max_attempts 5. Экспоненциальный бэкофф с потолком и случайным разбросом ±20% — чтобы сотня стримов, споткнувшихся об один сетевой блип, не ломанулась в API синхронно.
Время жизни. stream_ttl 5.0 — сколько завершённый стрим ещё лежит в памяти, чтобы можно было спросить его результат. machine_ttl 30.0 — сколько пустая машина чата ждёт новых стримов, прежде чем самоликвидироваться.
Прочее. typing_enabled — показывать ли «печатает…». raise_on_failure — кидать StreamFailedError при провале или молча вернуть список ID. shutdown_timeout — сколько ждать на aclose(). timings_capacity 1024 — размер LRU-кэша таймингов по чатам.
ScopedSender и LiveStream — то, что вы держите в руках
SenderMiddleware кладёт в хендлер ScopedSender, уже привязанный к боту, чату и топику. Это фабрика: sender.stream() открывает LiveStream.
LiveStream — async context manager с тремя приёмами:
update(chunks)— задать желаемое состояние. Список маппингов{"text": ..., "entities": [...]}. Один элемент — одно сообщение в чате.finish()— сказать «я всё» и дождаться, пока последний байт реально доедет. Возвращает списокmessage_id.- выход из
async with— вызываетfinish()за вас.
Тонкость, которую легко проглядеть: если из блока вылетело исключение, __aexit__ вызывает finish(raise_on_failure=False). Ваша ошибка не будет замаскирована ошибкой доставки. Приоритет у исходного исключения.
Хотите несколько сообщений? Просто передайте несколько чанков:
stream.update([{"text": "Часть 1"}, {"text": "Часть 2"}])
Стало меньше чанков, чем было? Лишние хвосты помечаются на удаление и удаляются сами. Это не костыль, а нормальный сценарий: LLM переписала ответ короче — библиотека приведёт чат в соответствие.
SenderRuntime — фабрика и владелец фоновых задач
Один на приложение. Хранит словарь (bot_id, chat_id) → MachineWorker и раздаёт стримам сквозные ID.
При open_stream рантайм смотрит, есть ли живой воркер для этого чата. Нет или он уже «stopping» — создаёт новую машину, новый TelegramExecutor и запускает воркер как asyncio.Task. Есть — переиспользует. Так все стримы одного чата автоматически делят общий бюджет rate limit, даже если открыты из разных хендлеров.
Отдельно живёт TimingsCache — LRU по чатам. И вот зачем: машину можно выбросить из памяти как мусор, а знание «в этот чат мы редактировали 0.9 секунды назад» переживёт её и достанется следующей машине. Без этого после каждой эвикции мы бы радостно ловили 429 на первом же сообщении.
aclose() при остановке бота финализирует все стримы, ждёт воркеров до shutdown_timeout, остальных отменяет. Недописанные сообщения успевают дописаться.
MachineWorker — цикл, который крутит всё
Собственно, весь асинхрон сосредоточен здесь, в двадцати строчках:
план = machine.plan(now)
есть действие → выполнить → machine.apply(...) → повторить
нет действия → sweep, проверить эвикцию, спать до дедлайна или до пробуждения
Спит воркер через гонку двух задач: sleep_until(deadline) и wakeup.wait(). Кто первый — тот и разбудил. Поэтому update() из вашего хендлера не ждёт следующего тика: он ставит Event, и воркер просыпается немедленно. Никакого поллинга, никакого «проверяем раз в 100 мс».
Воркер же играет роль почтальона для finish(). Каждому стриму при регистрации выдаётся asyncio.Event; после каждого шага _settle() проверяет, какие стримы уже устаканились, снимает копию их результата и ставит событие. Копия важна: пока finish() просыпается, машина может успеть подмести стрим по TTL, и outcome() вернул бы пустоту. Снапшот в _outcomes закрывает эту гонку.
Если воркер всё-таки упал с необработанным исключением — он не уносит стримы с собой в тишину: kill_all("worker crashed"), всех разбудить, и только потом умереть. Ждущий finish() получит исключение, а не вечное зависание.
SenderMachine — диспетчер одного чата
Ей задают ровно один вопрос, снова и снова: «сейчас now, что делать?» Она держит четыре кучки состояния, каждая на своём горизонте времени: _streams (сейчас), _retry_at (секунды), _done_at (TTL стримов), _idle_since (TTL самой машины).
Главное — apply(), где результат исполнения раскладывается по пяти веткам. Порядок веток тут не стилистика, а семантика:
retry_afterпришёл — это 429. Ставимhold_untilна весь чат, эмитимChatHold, выходим немедленно. Ключевое:last_atне обновляется и попытка не засчитывается. Мы ведь ничего не сделали, нас отшили. Если бы 429 считался попыткой, стримы дохли бы отmax_attemptsво время обычного троттлинга — худший из возможных исходов.- typing — обновили
last_at, и всё. У него нет ни успеха, ни провала. - успех — снять бэкофф, записать доставленное состояние.
STREAM_DEAD— бот заблокирован, чат удалён, топик закрыт. Стрим убит целиком, летитStreamFailed.MESSAGE_DEAD— умерло одно сообщение (например, «message to edit not found»). Стрим живёт дальше, летитMessageFailed. Именно отсюда потом берётся статусpartial.
Всё остальное — временная ошибка: засчитать попытку, назначить retry_at, а на max_attempts — убить стрим.
Уборка идёт в два уровня. Завершённый стрим ещё stream_ttl секунд лежит в памяти (вдруг кто спросит результат), потом выметается. Опустевшая машина ждёт machine_ttl и становится is_evictable. Чтобы это не зависло, _ttl_deadline() подмешивается в дедлайн сна: даже когда делать абсолютно нечего, воркер просит разбудить себя к истечению ближайшего TTL.
scheduler.plan — кто ходит следующим
Сердце библиотеки, страница кода. Читается как «кто первым будет готов».
Сначала общий стоп-кран: если now < hold_until, весь чат молчит, без исключений для кого бы то ни было. Потом обход стримов, у каждого спрашиваем pending(). Для найденного намерения считаем момент готовности:
ready = max(last_at[вид] + interval[вид], # лимит чата по виду действия
retry_at[(стрим, индекс)]) # персональный бэкофф сообщения
Два независимых ограничения, берётся более позднее. Ничего этого вида ещё не делали — -inf, можно прямо сейчас.
Выбор победителя — минимум по кортежу (ready, 0 если финализирован иначе 1, stream_id). Приоритеты, записанные лексикографически:
- раньше готов — раньше идёт; никакой приоритет не отменяет лимиты;
- при равной готовности выигрывает финализированный стрим — его LLM уже договорила, ему осталось только добить последний edit и освободить слот; несправедливо держать его в очереди за тем, что ещё генерируется;
stream_idкак тай-брейк — чтобы порядок был детерминированным, а не зависел от порядка в хеш-таблице. Тесты скажут спасибо.
Typing считается отдельно и только если есть хоть один не финализированный стрим: «печатает…» показывают, пока текст рождается, а не пока досылаются готовые куски. Врать пользователю незачем.
И главная дисциплина: за один plan — ровно одно действие. Никаких батчей. Именно это делает соблюдение лимитов доказуемым, а не «вроде бы работает».
SenderStream и SenderMessage — diff-движок
SenderStream — список SenderMessage плюс флаги is_final / state.
update() делает выравнивание списков: совпавшие по индексу обновляет, новые добавляет, лишние помечает на удаление. Хвосты, которые ещё не были отправлены, просто выбрасываются — удалять нечего. Те, что уже живут в чате, остаются в списке ждать своего DeleteIntent.
SenderMessage.intent() — вся логика в семи строках сравнения:
| Что видим | Что делаем |
|---|---|
нет message_id, есть текст |
SendIntent |
есть message_id, хеш не совпал |
EditIntent |
есть message_id, текста не хотим |
DeleteIntent |
| хеши совпали | None — всё уже доставлено |
Сравнение по content_hash (sha256 от текста + entities, считается в __post_init__ чанка) означает, что повторный update с тем же содержимым не породит ни одного запроса. Можно звать хоть в цикле.
Порядок в pending() заслуживает внимания: сначала все send/edit слева направо, и только потом удаления — с конца. Пользователь видит, как текст растёт естественно, а хвосты подчищаются последними и с конца, чтобы не ломать индексы.
TelegramExecutor и classify — граница с aiogram
Наконец, единственное место во всей библиотеке, которое знает про Bot. Исполнитель тупой и это by design: получил ScopedAction — сделал ровно один вызов API — вернул Result. Никакой логики, никаких решений.
Вся хитрость — в переводе исключений aiogram на язык машины. classify() смотрит на тип и текст ошибки:
TelegramRetryAfter→ transient +retry_after(машина превратит это вhold_untilна весь чат);TelegramForbiddenErrorи маркеры вроде «bot was blocked», «chat not found», «topic_closed» →STREAM_DEAD;TelegramBadRequestс «message to edit not found», «message is too long», «can't parse entities» →MESSAGE_DEAD;- всё остальное →
TRANSIENT, попробуем ещё раз.
Отдельным случаем идёт «message is not modified». Формально это ошибка, по смыслу — успех: то, что мы хотели видеть, уже там. is_not_modified() ловит её до классификации и возвращает Result(ok=True). Иначе безобидная гонка съедала бы попытки и в итоге убивала стрим.
Порядок проверок в classify тоже не случаен: _STREAM_DEAD проверяется до TelegramBadRequest, потому что «not enough rights» приезжает именно как BadRequest, но это приговор всему чату, а не одному сообщению.
События
Опциональный sink: Callable[[Event], None] в конструктор рантайма. Три типа: MessageFailed, StreamFailed, ChatHold. Синхронный, вызывается из воркера — так что метрики инкрементить можно, а вот ходить в базу не стоит. Исключения из sink ловятся и логируются: ваша сломанная телеметрия не уронит доставку.
Итог: три статуса
finish() возвращает список ID, а внутри всё сводится к трём исходам:
ok— доставлено;partial— часть сообщений умерла, но стрим дожил до конца;failed— стрим убит, с причиной; приraise_on_failureлетитStreamFailedErrorс уже доставленными ID внутри.
Разделение уровней отказа — временная ошибка, смерть сообщения, смерть стрима, пауза всего чата — и есть главное, что тут спроектировано. Наивные реализации сваливают это в один except и потом годами ловят странности.
Установка
uv add aiogram-stream-sender
Python 3.14+, aiogram 3.22+. Рабочий бот целиком — в example/bot.py.
MIT.
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file aiogram_stream_sender-0.1.2.tar.gz.
File metadata
- Download URL: aiogram_stream_sender-0.1.2.tar.gz
- Upload date:
- Size: 90.4 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via: uv/0.12.2 {"installer":{"name":"uv","version":"0.12.2","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
e52fc1a45e9494d0358e2ff239f6bcb5144594c841fdc430c9dc5c050a264036
|
|
| MD5 |
5a24d0a51428a092bcbb5782649804a8
|
|
| BLAKE2b-256 |
405784131b7ca96e4ee648b2f610714a31636bb0d0648aa429e2be36b48d5f7b
|
File details
Details for the file aiogram_stream_sender-0.1.2-py3-none-any.whl.
File metadata
- Download URL: aiogram_stream_sender-0.1.2-py3-none-any.whl
- Upload date:
- Size: 24.4 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via: uv/0.12.2 {"installer":{"name":"uv","version":"0.12.2","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
87285300aaad7bf71ac2e1edc7814a7bcab247cfe5fd124a0d54649bee278bbe
|
|
| MD5 |
57ad6bf88458730e6a7031ce97572fbf
|
|
| BLAKE2b-256 |
6452a36fea7db1abe21fa8b870bfe95eaee2994f49f46f8c9b07c772cfe60a1e
|