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()за вас.
sender.stream(typing=False) открывает стрим, который не участвует в показе «печатает…». Нужно для служебных сообщений, которые живут в чате долго: висящий индикатор рядом с кнопкой врёт пользователю, а лимит на typing всё равно общий на чат.
Тонкость, которую легко проглядеть: если из блока вылетело исключение, __aexit__ вызывает finish(raise_on_failure=False). Ваша ошибка не будет замаскирована ошибкой доставки. Приоритет у исходного исключения.
Хотите несколько сообщений? Просто передайте несколько чанков:
stream.update([{"text": "Часть 1"}, {"text": "Часть 2"}])
Стало меньше чанков, чем было? Лишние хвосты помечаются на удаление и удаляются сами. Это не костыль, а нормальный сценарий: LLM переписала ответ короче — библиотека приведёт чат в соответствие.
Chunk — что вообще можно попросить
Чанк — обычный маппинг, библиотека сама превратит его в Chunk. Кроме text и entities он понимает:
reply_markup— инлайн-клавиатура, как её сериализует aiogram:markup.model_dump(mode="json", exclude_none=True). Наeditклавиатура выставляется заново, так что убрать её — это просто перестать её присылать.link_preview— телоLinkPreviewOptions, например{"is_disabled": True}.reply_to—message_id, на который отвечаем. Уходит какreply_parametersсallow_sending_without_reply, поэтому удалённый оригинал не убьёт отправку.parse_mode— для статичных текстов, у которых уже есть разметка ("HTML"). Когда он задан,entitiesне отправляются: у Telegram это две взаимоисключающие формы одного и того же.key— идентичность сообщения, см. ниже.
Всё, кроме key, входит в content_hash: поменяли клавиатуру — поедет edit, не поменяли — не поедет ничего.
key — как подвинуть сообщение вниз
Telegram не умеет перемещать сообщения. А сценарий «служебное сообщение с кнопкой всегда должно быть последним в чате» — обычный: пользователь пишет, а мы каждый раз переспрашиваем внизу.
Сделать это руками — значит звать delete_message и send_message мимо планировщика, то есть мимо всех лимитов. Поэтому перемещение живёт внутри библиотеки и выражается одним полем.
key — это ответ на вопрос «то же самое это сообщение или уже другое». Хеш отвечает за содержимое, key — за личность:
stream.update(
[{"text": "Отправить сообщение?", "reply_markup": markup, "key": str(revision)}]
)
Тот же key — обычный edit на месте. Новый key — сообщение уничтожается и рождается заново, уже внизу чата. Ровно одно такое сообщение существует в любой момент времени.
Внутри это не «удали и отправь одним действием», а два независимых шага: DropIntent (удаление, кинд DELETE) и следом обычный SendIntent. Каждый проходит через свой лимит, каждый ретраится сам по себе, и за один plan по-прежнему выполняется ровно одно действие.
Из этого следует приятное свойство: пока идёт перемещение, update можно звать сколько угодно. Десять сообщений подряд от пользователя не превратятся в десять пар delete/send — состояние схлопнется, и в чат уйдёт одна пара с самым свежим содержимым.
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() получит исключение, а не вечное зависание. То же и при отмене: раньше finish() ждал Event, который отменённый воркер уже некому было поставить, и ждал он до конца жизни процесса. Теперь ожидание — это гонка события и самой задачи воркера, так что смерть воркера тоже считается ответом.
Контекст — кому принадлежит отправка
Воркер один на чат и живёт дольше любого стрима. asyncio копирует контекст вызывающего в каждую задачу, поэтому воркер, созданный внутри обработки одного апдейта, навсегда забирал себе его контекст — и все последующие отправки в этом чате, месяцами, числились за тем самым первым апдейтом.
Поэтому воркер стартует из пустого контекста и не принадлежит никому. А LiveStream при открытии снимает contextvars.copy_context(), и каждый send/edit/delete этого стрима выполняется внутри снимка своего открывателя.
Всё, что живёт в контекстных переменных, доезжает до самого вызова Telegram: контекст трассировки, привязка логов к запросу, что угодно. Библиотеке для этого не нужна ни одна дополнительная зависимость — только stdlib.
Result теперь несёт и сам объект исключения (error), а не только str(error): по строке нельзя записать исключение со стектрейсом ни в трейс, ни куда-либо ещё.
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 считается отдельно и только если есть хоть один не финализированный стрим, открытый с typing=True: «печатает…» показывают, пока текст рождается, а не пока досылаются готовые куски и не пока в чате просто висит кнопка. Врать пользователю незачем.
И главная дисциплина: за один plan — ровно одно действие. Никаких батчей. Именно это делает соблюдение лимитов доказуемым, а не «вроде бы работает».
SenderStream и SenderMessage — diff-движок
SenderStream — список SenderMessage плюс флаги is_final / state.
update() делает выравнивание списков: совпавшие по индексу обновляет, новые добавляет, лишние помечает на удаление. Хвосты, которые ещё не были отправлены, просто выбрасываются — удалять нечего. Те, что уже живут в чате, остаются в списке ждать своего DeleteIntent.
Исключение — сообщение, которое прямо сейчас летит в Telegram. Машина помечает его in_flight, когда отдаёт действие исполнителю, и снимает пометку в apply. Такой хвост не выбрасывается, даже если у него ещё нет message_id: иначе ответ API вернулся бы в пустоту, а в чате осталось бы сообщение, про которое никто уже не помнит. Пока хоть одно сообщение в полёте, стрим не считается завершённым — finish() дождётся и его отправки, и последующего удаления.
SenderMessage.intent() — вся логика в семи строках сравнения:
| Что видим | Что делаем |
|---|---|
нет message_id, есть текст |
SendIntent |
есть message_id, key не совпал |
DropIntent — сообщение переродится ниже |
есть message_id, хеш не совпал |
EditIntent |
есть message_id, текста не хотим |
DeleteIntent |
| хеши совпали | None — всё уже доставлено |
Порядок проверок здесь тоже смысловой: key спрашивают раньше хеша. Если личность сообщения сменилась, редактировать нечего — старое тело всё равно уедет в корзину.
Сравнение по 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», «button_data_invalid» →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.3.0.tar.gz.
File metadata
- Download URL: aiogram_stream_sender-0.3.0.tar.gz
- Upload date:
- Size: 96.1 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via: uv/0.12.5 {"installer":{"name":"uv","version":"0.12.5","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 |
6ee58dc6e68bb5785fa731c043d8328aaf4ba8e8febeba13ea220acff1720c01
|
|
| MD5 |
a76b95f66c096dd4c758c9081088d706
|
|
| BLAKE2b-256 |
95a7b5228b07951880b67825405fd1dc558d4d598ca7dcda9f7e4586c63757ae
|
File details
Details for the file aiogram_stream_sender-0.3.0-py3-none-any.whl.
File metadata
- Download URL: aiogram_stream_sender-0.3.0-py3-none-any.whl
- Upload date:
- Size: 28.4 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via: uv/0.12.5 {"installer":{"name":"uv","version":"0.12.5","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 |
760506f9cdd45368b82ec95d53d238a416f55ed1c72980ee231bba77a3db1a9e
|
|
| MD5 |
7f47945f0855acb7a65ae25aa50ec671
|
|
| BLAKE2b-256 |
a7325bd642fa6200b7cb61ef299f2f3b3849367afc28890f066faf8ada0a3461
|