← Назад к новостям

Я делал «nginx для телеграм‑ботов», а получился конструктор шлюзов на Rust

Привет, Хабр. Это история о том, как я писал балансировщик нагрузки для телеграм‑ботов, на полпути понял, что телеграм тут вообще ни при чём, выкинул половину проекта — и теперь у меня Rust‑конструктор, из которого одинаково собирается реверс‑прокси, самодельный ngrok, мост HTTP → RabbitMQ → Kafka и тот самый балансировщик ботов, с которого всё началось. Проект называется tgin. Раньше это был Telegram Gateway Interface, теперь — Traffic Gateway Interface. С чего всё началось Началось с банальной боли: когда телеграм‑бот перестаёт влезать в один процесс, начинается цирк. Обновления надо раскидывать по нескольким инстансам, следить, чтобы сообщения одного чата не обрабатывались вразнобой, деплоиться без потери апдейтов. Я посмотрел на это и подумал: так это же nginx, только для Bot API. И сел писать «nginx для телеграм‑ботов» — лонгполл и вебхук на входе, round‑robin по инстансам на выходе, конфиг‑файл, вот это вот всё. Работало. Даже неплохо. А потом я сел переделывать это по‑нормальному и заметил вещь, которая сломала мне всё позиционирование. Прозрение Ни одно сложное место моего «телеграмного» балансировщика не было про телеграм. Смотрите, в чем заключалась суть проекта: принять поток сообщений, буферизовать, раскидать по потребителям, понять доставлено или нет, повторить, если нет, не захлебнуться, когда потребители тупят, и умереть красиво по ctrl+c, ничего не потеряв. Теперь замените «обновление от телеграма» на «HTTP‑запрос» или «сообщение из RabbitMQ» — список не изменился ни на пункт. Телеграм оказался просто первой формой трафика, которая протекала через код. И тут же вторая мысль, обиднее первой: этот же клей я уже писал руками на других проектах, и не раз. Триста строк tokio‑лапши, которые перекладывают запросы из HTTP в очередь. Пишутся за вечер, отлаживаются до пенсии: backpressure забыли, ack соврали, graceful shutdown осуществляется через kill -9. Подозреваю, такой файлик есть в каждом втором проекте, просто все делают вид, что его нет. Так родился план: один раз честно решить скучные проблемы в маленьком ядре, а всё остальное вынести в батарейки. Как это выглядит Шлюз на tgin — это ваш обычный Rust‑бинарь строк на двадцать. Вот полноценный реверс‑прокси: let server = HttpServer::new("0.0.0.0:8080"); let http = HttpClient::new(); Tgin::new() .pipeline( HttpIngress::catch_all(&server), Route::new() .when( |r: &RequestData| r.uri.path().starts_with("/api"), RoundRobin::new() .to(HttpEgress::new(&http, "http://10.0.0.1:3000")) .to(HttpEgress::new(&http, "http://10.0.0.2:3000")), ) .otherwise(HttpEgress::new(&http, "http://10.0.0.3:3000")), ) .run() .await; /api балансится между двумя бэкендами, всё остальное едет на третий. Убьёте один бэкенд — трафик молча перетечёт на живой. Убьёте оба — клиент получит 502, а не вечно висящий коннект. Конфиг‑файла нет, конфиг это код. Здесь несовместимая пара «вход‑выход» не сработает, тк все типы должны быть согласованы по типу который принимают и отдают. Три сущности и один конверт Внутри всё стоит на трёх китах. Ingress — откуда трафик приходит. Egress — куда уходит. Между ними ездит Envelope : данные + метаданные + опциональный канал для ответа. Пайплайн = ingress + egress. Больше в ядре ничего нет. А декораторы — это egress`ы, оборачивающие другие egress`ы. RoundRobin , Route , Retry , фан‑аут All — ядро вообще не знает об их существовании, для него это просто ещё один выход. Reply — это ack и backward Самое важное место архитектуры — поле reply в Envelope. Разберем на примере HTTP на входе, RabbitMQ на выходе Приходит POST /input. HTTP‑хендлер кладёт запрос в конверт, вкладывает туда oneshot‑канал и курит бамбук на нём. Конверт доезжает до Rabbit‑egress’а, тот публикует сообщение (persistent, в durable‑очередь) и, ключевой момент — ждёт publisher confirm от брокера. Исход уезжает обратно в oneshot, хендлер просыпается и переводит его на язык HTTP: брокер подтвердил → клиент получает 202 Accepted . Не 200 — мы не сделали работу, мы приняли её на хранение. Двести второй ровно для этого и придуман, просто им никто не пользуется; брокер лежит → 502 . Клиент знает, что доставки не было, и может повторить; внутренняя очередь переполнена → 429 , backpressure вместо пожирания памяти. А теперь разворачиваем стрелку: Rabbit на входе, HTTP‑воркер на выходе. Тот же самый reply теперь работает как ack: воркер ответил — ack, ретраи кончились — nack с requeue, ошибка фатальная — nack без requeue. Одна механика на два мира: для HTTP reply становится статус‑кодом, для очереди — сообщением. Побочный эффект: система различает «апстрим ответил 502» и «я не смог достучаться до апстрима». Звучит банально? Сходите гляньте в свой прод, я подожду. Самодельный ngrok Когда ядро устаканилось, захотелось проверить композицию чем‑то наглым. Я взял и запилил туннель как у ngrok: батарейка на ~370 строк, ноль правок ядра. // на VPS: всё входящее заворачиваем в туннель Tgin::new().pipeline( HttpIngress::catch_all(&server), TunnelEgress::new(&server, "/tunnel", &token), ) // дома: вылезаем наружу и раздаём локалхост Tgin::new().pipeline( TunnelIngress::new("ws://мой-vps:8080/tunnel", &token), HttpEgress::new(&http, "http://127.0.0.1:3000"), ) Локалхост наружу отдаётся, kill клиента — сервер честно отвечает 502, поднял обратно — переподключился сам за секунду. Отдельно смешно, что туннель сам живёт на WebSocket, но проксировать чужие WebSocket’ы не умеет. Об ограничениях ниже. Телеграм никуда не делся Старое позиционирование не умерло, а превратилось в набор батареек, и по дороге прокачалось. Вход — лонгполл или вебхук. Лонгполл‑ingress сдвигает offset у телеграма только после того, как обновление реально доставлено дальше по цепочке: упал ваш бот — телеграм передоставит. Каждому апдейту проставляется ключ = chat_id, и ядро бесплатно даёт порядок внутри чата при параллельности между чатами. Выход интереснее. Понятно, что можно пушить обновления инстансам по HTTP, как делает сам телеграм — это вебхук‑egress. Но есть и наглый вариант: tgin поднимает у себя эндпоинт getUpdates , и немодифицированные лонгполл‑боты поллят его, свято веря, что разговаривают с api.telegram.org . Свои update_id, длинный полл с таймаутом, подтверждение оффсетом, передоставка при падении — всё как у большого. Самозванец as a service. Tgin::new().pipeline( TelegramBotPollingIngress::new(&http, &token), RoundRobin::new() .to(TelegramBotPollingEgress::new(&server, "/bot1")) .to(TelegramBotWebhookEgress::new(&http, "http://10.0.0.2:3000/webhook")), ) Один бот поллит, второму пушится, между ними балансировка — это, кстати, дословный перевод конфига из самой первой версии tgin. Кафка С RabbitMQ всё довольно прямолинейно: durable‑очереди, persistent‑сообщения, publisher confirms, ack по исходу reply. А вот с кафкой пришлось подумать, потому что у кафки нет ack’ов. Поэтому Kafka‑ingress коммитит только непрерывный префикс подтверждённого: на каждой партиции сидит acker, который ждёт reply в порядке offset`ов и стопорит коммит на первом провале, а потом делает seek назад и передоставляет. Медленнее бездумного авто‑коммита? Да. Теряет сообщения? Нет. Я за второе. На выходе симметрично: meta.key конверта становится partition key. Цепочка телеграм‑чат → tgin → кафка сохраняет порядок сообщений чата в партиции, и для этого не написано ни строчки специального кода — ключ просто едет в конверте через весь пайплайн. Чего оно не умеет Стриминг тел. Тело запроса буферизуется целиком. Гонять двухгигабайтные файлы через tgin — пока плохая идея. WebSocket passthrough. Да, туннель сам ездит по WebSocket. Да, проксировать чужие WebSocket’ы нельзя. Hot reload топологии. Конфигурация статическая и компилируется. Новая топология — новый бинарь. Потрогать руками В репозитории лежат демки на docker compose: reverse‑proxy, туннель, HTTP→Rabbit→воркер и HTTP→Kafka→воркер. Каждая — один docker compose up --build плюс README с прогоном. Зачем это всё Я не строю убийцу nginx. По голому перфу glue‑архитектура всегда проиграет специализированному прокси ‑за композицию заплачено микросекундами на внутреннем хопе, и я этот размен сделал сознательно. Я строю инструмент для стыков — мест, где один узкий инструмент заканчивается, другой начинается, а между ними по традиции живёт та самая самописная лапша из трёхсот строк. Теперь вместо неё можно взять ядро, у которого backpressure, ack’и и graceful shutdown сделаны один раз и правильно, досочинить свою батарейку — и получить туннель, прокси, очередь и свою логику одним бинарём. Репозиторий: https://github.com/chesnokpeter/tgin Пишите комментарии, пообщаемся, мне интересно что за штуку я придумал
📊 Источник: Habr | Оригинал