Skip to main content

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(), где результат исполнения раскладывается по пяти веткам. Порядок веток тут не стилистика, а семантика:

  1. retry_after пришёл — это 429. Ставим hold_until на весь чат, эмитим ChatHold, выходим немедленно. Ключевое: last_at не обновляется и попытка не засчитывается. Мы ведь ничего не сделали, нас отшили. Если бы 429 считался попыткой, стримы дохли бы от max_attempts во время обычного троттлинга — худший из возможных исходов.
  2. typing — обновили last_at, и всё. У него нет ни успеха, ни провала.
  3. успех — снять бэкофф, записать доставленное состояние.
  4. STREAM_DEAD — бот заблокирован, чат удалён, топик закрыт. Стрим убит целиком, летит StreamFailed.
  5. 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

aiogram_stream_sender-0.1.2.tar.gz (90.4 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

aiogram_stream_sender-0.1.2-py3-none-any.whl (24.4 kB view details)

Uploaded Python 3

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

Hashes for aiogram_stream_sender-0.1.2.tar.gz
Algorithm Hash digest
SHA256 e52fc1a45e9494d0358e2ff239f6bcb5144594c841fdc430c9dc5c050a264036
MD5 5a24d0a51428a092bcbb5782649804a8
BLAKE2b-256 405784131b7ca96e4ee648b2f610714a31636bb0d0648aa429e2be36b48d5f7b

See more details on using hashes here.

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

Hashes for aiogram_stream_sender-0.1.2-py3-none-any.whl
Algorithm Hash digest
SHA256 87285300aaad7bf71ac2e1edc7814a7bcab247cfe5fd124a0d54649bee278bbe
MD5 57ad6bf88458730e6a7031ce97572fbf
BLAKE2b-256 6452a36fea7db1abe21fa8b870bfe95eaee2994f49f46f8c9b07c772cfe60a1e

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page