Руководство · боты · SERVICE(10) API
Как создавать ботов Aster
Бот Aster — это headless-клиент: та же криптография, что у обычного собеседника (Noise-канал, PQXDH + Double Ratchet, sealed sender), плюс типизированный формат сообщений SERVICE(10) и цикл опроса. Релей остаётся слепым — он передаёт только запечатанный конверт. От генерации ключа до инлайн-кнопок и совместимости с Telegram Bot API.
Headless-клиент, не вебхуки — бот опрашивает релей Хранение только в RAM E2E: релей не видит ни текста, ни типа кадра
На этой странице
Что такое бот Aster
Бот — это обычный собеседник с точки зрения сети: у него есть адрес ASTER1…, его добавляют в контакты как человека, и вся переписка с ним сквозная. Отличие лишь в том, что за адресом стоит не человек, а программа, которая автоматически отвечает.
- Та же криптография. Бот переиспользует крипто и транспорт клиента: Noise_NK на линке, PQXDH-рукопожатие, Double Ratchet на каждое сообщение, sealed sender. Ничего нового на проводе — релей по-прежнему видит только непрозрачный запечатанный шифртекст.
- Никаких вебхуков. Бот не открывает портов наружу и не ждёт входящих. Он сам опрашивает релей (polling), забирает адресованные ему конверты, расшифровывает, отвечает. Это работает за NAT, без белого IP и без домена.
- Типизированные кадры. Поверх открытого текста живёт кодек
SERVICE(10)— типизированная схема (текст, кнопки, callback, правка/удаление), внутри ратчет-плейнтекста. Релей не отличает кадр с кнопками от простого «привет». - RAM-first. Продакшн-бот ничего не пишет на диск, кроме (по желанию) 64-байтного ключа идентичности. Сессии, контакты, prekey живут только в оперативной памяти. Физический захват диска бота не даёт данных пользователей — их там нет.
aster-services дерева исходников Aster. Он не меняет ни aster-core, ни aster-wire, ни aster-net, ни релей: бот — чисто клиентская надстройка.
Как это устроено
Один оборот цикла бота:
- Публикация бандла. При старте бот выкладывает на релей подписанный prekey-бандл (X25519 SPK + ML-KEM prekey), чтобы новые собеседники могли инициировать PQXDH.
- Опрос. Бот периодически спрашивает релей о конвертах в своих ротируемых почтовых ящиках (адрес ящика уникален на каждый счётчик — релей не связывает их между собой). Опрос адаптивный: чаще во время диалога, реже в тишине.
- Расшифровка. Каждый конверт вскрывается (sealed sender → PQXDH/Ratchet), плейнтекст декодируется в
Update. КадрSERVICE→ типизированное событие; сырой UTF-8 от немодифицированного клиента → просто текст. - Обработчик.
Updateуходит в вашHandler, который возвращаетVec<Outgoing>— ноль или больше ответных элементов.Outgoing— это либо кадрServiceFrame(текст, кнопки, правка, удаление), либо файлFileSpec(фото/видео/документ). ИServiceFrame, иFileSpecпревращаются вOutgoingвызовом.into(). - Ответ. Каждый ответ кодируется, шифруется ратчетом, запечатывается и кладётся в ротируемый ящик собеседника. Файл едет вне сессии (см. «Файлы, фото и видео»), а по ратчету — только его локатор. Релей снова видит только шифртекст.
Требования
| Компонент | Значение |
|---|---|
| Rust | ≥ 1.74 (edition 2021) — как для сборки всего Aster |
| Исходники | дерево репозитория Aster (крейт aster-services) |
| Асинхронность | tokio (runtime уже в зависимостях крейта) |
| Релей | адрес host:порт и закреплённый публичный ключ релея (64 hex). Тот же, что пинят обычные клиенты |
| Идентичность | 64 байта (128 hex): 32 — seed подписи ‖ 32 — секрет DH |
| Порты | боту не нужны входящие порты — он исходящий клиент |
1Идентичность бота
Бот определяется 64-байтным ключом: первые 32 байта — seed для подписи (Ed25519), вторые 32 — секрет обмена ключами (X25519). Публичный адрес ASTER1… выводится из ключа подписи детерминированно. Сгенерировать можно любой источник криптослучайности:
# 64 байта = 128 hex-символов
openssl rand -hex 64
Полученную строку бот читает со стандартного ввода (stdin) при запуске — так секрет никогда не попадает в список аргументов процесса и не пишется на диск. Держите её как пароль: кто владеет этими 64 байтами — тот и есть бот.
2Первый бот — эхо
Самый простой способ — добавить новый бинарь прямо в крейт aster-services (там же живут примеры demobot/statusbot). Создайте файл crates/aster-services/src/bin/mybot.rs:
use std::io::Read;
use aster_net::from_hex;
use aster_services::{Bot, Incoming, ServiceFrame, Update};
#[tokio::main]
async fn main() {
// RAM-защита: без core-dump, non-dumpable, mlock (best-effort).
aster_services::apply_ram_protection();
// Идентичность: 128 hex (64 байта = 32 sign ‖ 32 dh) со stdin.
let mut s = String::new();
std::io::stdin().read_to_string(&mut s).ok();
let b = from_hex(s.trim()).expect("нужно 128 hex");
let (mut sign, mut dh) = ([0u8; 32], [0u8; 32]);
sign.copy_from_slice(&b[..32]);
dh.copy_from_slice(&b[32..]);
let mut bot = Bot::from_identity_secret(sign, dh);
println!("эхо-бот: {}", bot.address()); // печатает ASTER1…-адрес
println!("полная карта: {}", bot.queue_card(0, true)); // ASTER1…#ASTERQ1… — вставляйте ЭТО
// Обработчик: на каждое событие вернуть элементы-ответы (Vec<Outgoing>).
// ServiceFrame → Outgoing через .into(); Vec::new() = ничего не отвечать.
let handler: aster_services::Handler = Box::new(|u: &Update| match &u.event {
Incoming::Message { text } => vec![ServiceFrame::text(format!("Ты написал: {text}")).into()],
Incoming::Command { name, .. } => vec![ServiceFrame::text(format!("Команда /{name}")).into()],
_ => Vec::new(),
});
// Аргументы: <адрес релея> <публичный ключ релея, hex>.
let args: Vec<String> = std::env::args().skip(1).collect();
bot.serve(&args[0], &args[1], handler).await.unwrap();
}
Ключевые типы: Bot::from_identity_secret(sign, dh) строит RAM-first бота (ничего не пишет на диск), bot.address() отдаёт его адрес, Handler — это Box<dyn FnMut(&Update) -> Vec<Outgoing> + Send>, а ServiceFrame::text(...) — быстрый способ вернуть простой текст (не забудьте .into(), чтобы получить Outgoing).
3Подключение к релею
Всё соединение делает один вызов:
bot.serve(addr, server_pubhex, handler).await
| Параметр | Что это |
|---|---|
addr | адрес нативного слушателя релея, host:порт (по умолчанию порт 9977, транспорт — Noise поверх TCP) |
server_pubhex | публичный Noise-ключ релея, 64 hex. Бот закрепляет его (TOFU) — он обязан совпасть с тем, что раздаёт оператор релея |
handler | ваш обработчик Handler |
По умолчанию используется чистый Noise. Если релей маскирует канал под TLS, выставьте переменную окружения ASTER_TLS=1 перед запуском. serve сам переподключается при обрыве связи и возвращает Err только на фатальной ошибке (например, кривой hex ключа).
45.89.60.139:9977 и ключ, начинающийся на 937ee3c4… (полный отпечаток — на странице загрузки). Всегда сверяйте отпечаток вне сети.
4Запуск и добавление в клиент
Соберите и запустите бинарь, передав ключ на stdin, а адрес и ключ релея — аргументами:
# собрать (из корня дерева Aster)
cargo build --release -p aster-services --bin mybot
# запустить: ключ идентичности → stdin, релей+его pubkey → аргументы
printf '<128-hex-ключ>' | ./target/release/mybot 45.89.60.139:9977 <pubkey-hex>
В консоли бот напечатает полную карту — строку вида ASTER1…#ASTERQ1…. Это не голый адрес, а комбинация: адрес + очередь доставки. Клиенту нужна именно полная карта — без неё добавление контакта невозможно.
ASTER1… не содержит очередь доставки и не подходит для добавления контакта. Вставляйте полную карту целиком: ASTER1…#ASTERQ1… (всё до # — адрес, всё после # — очередь). Полная карта печатается ботом при старте и при каждом переподключении.
- Откройте добавление контакта и вставьте полную карту бота (или отсканируйте QR). Отметьте контакт как бота — тогда клиент включает кнопки, кликабельные команды и восстановление после перезапуска.
- В пустом чате появится кнопка ▶ /start — нажмите её (или просто отправьте
/start). Это инициирует PQXDH и «будит» бота. - Команды в ответах бота (текст вида
/help,/status) становятся кликабельными — тап отправляет команду, вводить вручную не нужно.
5Автозапуск (systemd)
Чтобы бот жил постоянно и переживал перезагрузку, оформите его systemd-юнитом. Ключ идентичности держите в отдельном файле (режим 600), секреты — в EnvironmentFile. Пример по образцу боевого aster-statusbot:
# /etc/systemd/system/mybot.service
[Unit]
Description=Мой бот Aster
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
User=botuser
WorkingDirectory=/home/botuser/aster
ExecStart=/bin/bash -c 'exec /home/botuser/aster/target/release/mybot 45.89.60.139:9977 <pubkey-hex> < /home/botuser/.mybot-id'
Restart=always
RestartSec=5
MemoryMax=200M
[Install]
WantedBy=multi-user.target
# ключ идентичности в файл с правами 600
printf '<128-hex-ключ>' > ~/.mybot-id && chmod 600 ~/.mybot-id
# включить автозапуск и стартовать
sudo systemctl daemon-reload
sudo systemctl enable --now mybot
systemctl/df), под низким лимитом MEMLOCK упадёт на mlockall() при первом fork. Добавьте в юнит Environment=ASTER_NO_MLOCK=1 (защиту от свопа тогда обеспечивает сам деплой, например LimitMEMLOCK=infinity).
Сообщения и команды
Всё входящее приходит как Update { peer_sign, peer_dh, event }, где event: Incoming. Два основных вида — текст и команда:
Incoming::Message { text }— обычное текстовое сообщение.Incoming::Command { name, args }— сообщение, начинающееся с/. Например,/start привет→name = "start",args = "привет". Хвост@имяботана слове команды отрезается автоматически.
Ответ — вектор элементов Outgoing. Простой текст (ServiceFrame → Outgoing через .into()):
vec![ServiceFrame::text("Привет 👋").into()]
u.peer_address() вернёт ASTER1…-адрес собеседника (для логов/отображения), а peer_sign — стабильный ключ, по которому удобно вести своё состояние диалога (в RAM).
Ответ на нажатие
Чтобы подтвердить нажатие всплывающей подсказкой (тост или модальное окно) — не обязательно, но приятно, — верните ServiceFrame::CbAnswer:
ServiceFrame::CbAnswer {
cb_id: rand::random(),
text: Some("Готово ✅".into()),
alert: false, // true → модальное окно вместо тоста
url: None,
}
Можно вернуть сразу оба кадра: и CbAnswer (подтверждение), и новый/изменённый Msg.
Файлы, фото и видео
Бот отправляет вложение, добавив в ответ Outgoing::File(FileSpec) — тем же decoupled-путём, что и обычные собеседники. Байты уходят вне сессии: бот кладёт их в свежий одноразовый feed-блоб на релее, а по ратчету едет только маленький кадр-локатор. Получатель сразу видит карточку СКАЧАТЬ и подтягивает байты по тапу. Клиент при этом не меняется — это тот же формат, что у передачи файлов между людьми.
Файл описывается структурой FileSpec:
use aster_services::{FileSpec, Incoming, ServiceFrame};
// Картинка вшита в бинарь — никакого чтения диска в рантайме.
const LOGO: &[u8] = include_bytes!("logo.png");
// …в обработчике: подпись отдельным сообщением + сама картинка.
Incoming::Command { name, .. } if name == "photo" => vec![
ServiceFrame::text("Вот картинка 📷").into(),
FileSpec::new(LOGO.to_vec(), "logo.png", "image/png")
.with_dims(240, 240) // ширина/высота для резервирования места (опц.)
.into(),
],
Тип вложения задаёт только MIME: клиент рисует изображения прямо в ленте, а прочее — карточкой-вложением. Видео и документы — тот же FileSpec:
FileSpec::new(clip, "clip.mp4", "video/mp4").with_dims(1280, 720).into()
FileSpec::new(doc, "отчёт.pdf", "application/pdf").into()
| Метод / поле | Значение |
|---|---|
FileSpec::new(bytes, name, mime) | bytes: Vec<u8> — содержимое файла; name — имя (до 200 симв.); mime — тип (до 100 симв.) |
.with_dims(w, h) | пиксельные размеры для фото/видео — клиент заранее резервирует место под них (необязательно) |
| Предел размера | до 200 МиБ на файл; байты режутся на чанки по 48 КиБ автоматически |
Правка, удаление, профиль, набор
- Редактирование —
ServiceFrame::Edit { msg_ref, msg }: заменяет ранее отправленное сообщение по егоclient_msg_id. Удобно для «живых» экранов (обновляемый статус, счётчики). - Удаление —
ServiceFrame::Delete { client_msg_id }: кооперативно удаляет сообщение у собеседника (семантика WIPE). - Профиль —
ServiceFrame::BotProfile { name, about, commands, menu_button }: анонс имени, описания и списка команд (пары(команда, описание)— аналогsetMyCommands). Клиент показывает их как подсказки. - Набор —
ServiceFrame::Typing: эфемерное «печатает…», без подтверждения.
Справочник ServiceFrame
Полный набор кадров кодека SERVICE(10). Индексы кадров, как и в проводном протоколе, только дописываются в конец — порядок заморожен.
| # | Кадр | Направление | Назначение |
|---|---|---|---|
| 0 | Msg { text, buttons, reply_to, client_msg_id } | оба | сообщение; команда — это Msg, чей текст начинается с / |
| 1 | Callback { msg_ref, data } | польз.→бот | нажата инлайн-кнопка |
| 2 | CbAnswer { cb_id, text, alert, url } | бот→польз. | ответ на нажатие (тост/модалка) |
| 3 | Edit { msg_ref, msg } | бот→польз. | заменить прежнее сообщение по id |
| 4 | Delete { client_msg_id } | оба | кооперативное удаление (WIPE) |
| 5 | BotProfile { name, about, commands, menu_button } | бот→польз. | имя/описание/команды (setMyCommands) |
| 6 | Typing | бот→польз. | эфемерное «печатает…» |
Для ручной работы с кадрами крейт экспортирует encode_frame, decode_frame, is_frame, а также константы SERVICE_TAG (=10) и SERVICE_VERSION (=1). Обычному боту они не нужны — рантайм делает это сам.
Что приходит боту (Incoming)
| Вариант | Когда |
|---|---|
Message { text } | обычный текст |
Command { name, args } | сообщение с / в начале, разобранное на имя и аргументы |
Callback { msg_ref, data } | нажата инлайн-кнопка |
Frame(ServiceFrame) | любой другой декодированный кадр (правка/удаление/typing…) — передаётся как есть |
SERVICE. Рантайм это распознаёт и всё равно отдаёт Incoming::Message. А ответы бот таким собеседникам шлёт простым текстом (без кнопок) — интероперабельность из коробки.
Совместимость с Telegram Bot API
Если у вас уже есть код бота под Telegram — его можно почти не переписывать. Вместо инлайн-Handler вызовите serve_shim: он поднимает на 127.0.0.1:<порт> локальный Telegram-совместимый HTTP-API. Достаточно поменять базовый URL вашего клиента на этот адрес.
bot.serve_shim(
addr, // адрес релея
server_pubhex, // его публичный ключ
8081, // http-порт шима (только loopback)
token, // Bearer-токен: Vec<u8>
"Мой бот".into(), // имя (getMe)
"mybot".into(), // username (getMe)
).await
Поддерживаемые методы: getUpdates (long-poll), sendMessage, sendPhoto, sendDocument, sendVideo, sendAudio, answerCallbackQuery, setMyCommands, getMe. Запросы принимаются только с петлевого интерфейса и требуют Authorization: Bearer <token> (сравнение в постоянное время).
multipart/form-data, поэтому файл передаётся base64-строкой в JSON-теле, в поле по имени метода (photo/document/video/audio). Дополнительно можно указать filename, mime, width/height и caption (подпись приходит отдельным сообщением). Шлите тело как Content-Type: application/json — при form-кодировании символ + из base64 портится, и запрос честно отклоняется. Предел через шим — 12 МиБ на файл; для крупных файлов используйте инлайн-обработчик (Outgoing::File).
# отправить фото (base64) через шим
curl -s http://127.0.0.1:8081/sendPhoto \
-H 'Authorization: Bearer <token>' -H 'Content-Type: application/json' \
-d "{\"chat_id\":$CHAT,\"photo\":\"$(base64 -w0 logo.png)\",\"filename\":\"logo.png\",\"caption\":\"привет\"}"
chat_id выводится из ключа собеседника.
Переменные окружения
| Переменная | Действие |
|---|---|
ASTER_TLS=1 | подключаться к релею через TLS-маскированный канал (по умолчанию — чистый Noise) |
ASTER_NO_MLOCK=1 | пропустить mlockall() — нужно ботам, которые порождают дочерние процессы под низким лимитом MEMLOCK |
Свою конфигурацию бот читает как обычная программа. Например, боевой statusbot берёт из окружения список сервисов, секретное слово владельца и флаг разрешения перезапуска (STATUSBOT_SERVICES, STATUSBOT_SECRET, STATUSBOT_ALLOW_RESTART) — это его собственные переменные, а не часть API.
Модель безопасности
- RAM-first.
from_identity_secretничего не пишет на диск: путь пуст →save()— no-op. Сессии, контакты, prekey живут в памяти и теряются при рестарте — намеренно. - Закалка процесса.
apply_ram_protection()запрещает core-dump (секреты не окажутся в дампе), делает процесс non-dumpable (чужойptrace//proc/self/memзакрыт) и вызываетmlockall(), чтобы ключи не ушли в своп. - Честный рестарт. Каждый старт публикует монотонно растущую генерацию prekey; клиент это видит и переустанавливает сессию сам. Никаких «тихих» дыр — только для контактов, помеченных как бот.
- Один процесс на идентичность. Из-за delete-on-fetch две копии бота с одним ключом теряют сообщения. Держите ровно один живой процесс.
+Исходный код примеров
Полный исходный код ботов, работающих на публичном релеи Aster. Каждый пример — самодостаточный main.rs; скопируйте, подставьте свой ключ и запустите.
echo — минимальный эхо-бот — Самый простой пример: отвечает тем же текстом. Два режима запуска: RAM-first (из stdin) и persist (с файлом состояния). Доказывает полный round-trip (sealed sender + Double Ratchet + SERVICE codec).
crates/aster-services/src/bin/echo.rs
//! Echo bot — replies to every message/command with the same text.
//!
//! Proves the full end-to-end round-trip (sealed sender + Double Ratchet + the
//! SERVICE(10) codec) against any Aster client, through the unmodified blind relay.
//! It interops with an *unmodified* `aster-cli` too: raw-text peers get raw text back.
//!
//! Usage:
//! echo <state_path> <server_addr> <server_pubhex> (dev: keypair generated + persisted to <state_path>)
//! echo - <server_addr> <server_pubhex> (prod: client-issued identity from stdin; NOTHING persisted)
//! e.g.
//! echo /tmp/echo.bot 45.89.60.139:9099 937ee3c4… (set ASTER_TLS=1 for the wss relay)
//! printf '<128-hex sign||dh>' | echo - 45.89.60.139:9099 937ee3c4…
//!
//! In the `-` (RAM-first) mode the bot owner supplies the long-term identity keypair
//! their client issued; this process writes no state to disk (sessions/contacts stay
//! in RAM, lost on restart) — the mission's disk-seizure-minimisation property.
use std::io::Read;
use aster_net::from_hex;
use aster_services::{Bot, Incoming, ServiceFrame, Update};
/// Read a client-issued identity (128 hex chars = 64 bytes: sign‖dh) from stdin.
fn read_identity_from_stdin() -> Result<([u8; 32], [u8; 32]), Box<dyn std::error::Error>> {
let mut s = String::new();
std::io::stdin().read_to_string(&mut s)?;
let bytes = from_hex(s.trim()).ok_or("identity must be hex")?;
if bytes.len() != 64 {
return Err("identity must be 64 bytes (128 hex chars): sign(32) || dh(32)".into());
}
let mut sign = [0u8; 32];
let mut dh = [0u8; 32];
sign.copy_from_slice(&bytes[..32]);
dh.copy_from_slice(&bytes[32..]);
Ok((sign, dh))
}
#[tokio::main]
async fn main() {
let args: Vec<String> = std::env::args().skip(1).collect();
if args.len() != 3 {
eprintln!("usage: echo <state_path|-> <server_addr> <server_pubhex>");
eprintln!(" '-' reads a client-issued identity (128 hex) from stdin; persists nothing");
std::process::exit(2);
}
let (state_path, addr, pubhex) = (&args[0], &args[1], &args[2]);
// Pin RAM, forbid core dumps + ptrace before any user data is decrypted (mission:
// minimise USER-data leakage on physical seizure — see aster_services::harden).
aster_services::apply_ram_protection();
let mut bot = if state_path == "-" {
match read_identity_from_stdin() {
Ok((sign, dh)) => Bot::from_identity_secret(sign, dh),
Err(e) => {
eprintln!("error reading identity from stdin: {e}");
std::process::exit(1);
}
}
} else {
match Bot::load_or_generate(state_path) {
Ok(b) => b,
Err(e) => {
eprintln!("error: {e}");
std::process::exit(1);
}
}
};
println!("echo bot address: {}", bot.address());
println!("share it with a client, then message the bot.");
let handler: aster_services::Handler = Box::new(|u: &Update| -> Vec<aster_services::Outgoing> {
match &u.event {
Incoming::Message { text } => vec![ServiceFrame::text(text.clone()).into()],
Incoming::Command { name, args } => {
let reply = match name.as_str() {
"start" => "echo bot ready — send me anything and I'll echo it back".to_string(),
_ if args.is_empty() => format!("/{name}"),
_ => format!("/{name} {args}"),
};
vec![ServiceFrame::text(reply).into()]
}
Incoming::Callback { data, .. } => {
vec![ServiceFrame::text(format!("callback: {}", String::from_utf8_lossy(data))).into()]
}
// edit/delete/typing/etc. — nothing to echo.
Incoming::Frame(_) => Vec::new(),
}
});
if let Err(e) = bot.serve(addr, pubhex, handler).await {
eprintln!("fatal: {e}");
std::process::exit(1);
}
}
echo-shim — Telegram Bot API шим — Бот без инлайн-обработчика: поднимает локальный HTTP-API (`127.0.0.1:порт`) совместимый с Telegram Bot API. Существующий Telegram-код работает после смены базового URL.
crates/aster-services/src/bin/echo-shim.rs
//! echo-shim — a bot that exposes a localhost Telegram-Bot-API (getUpdates/sendMessage/
//! answerCallbackQuery/setMyCommands/getMe), so existing Telegram bot code drives it by
//! swapping the base URL to http://127.0.0.1:<port>. There is NO inline handler — the HTTP
//! client is the bot logic. RAM-only: the identity is provided at launch on stdin and
//! nothing (queue, sessions, contacts) touches disk.
//!
//! Usage: printf '<128-hex sign||dh>' | echo-shim <server_addr> <server_pubhex> <http_port> <token>
//! (set ASTER_TLS=1 for the wss relay)
use std::io::Read;
use aster_net::from_hex;
use aster_services::Bot;
#[tokio::main]
async fn main() {
let a: Vec<String> = std::env::args().skip(1).collect();
if a.len() != 4 {
eprintln!("usage: printf '<128hex>' | echo-shim <server_addr> <server_pubhex> <http_port> <token>");
std::process::exit(2);
}
let (addr, pubhex, port, token) = (&a[0], &a[1], &a[2], &a[3]);
// Pin RAM / forbid core dumps + ptrace before any user data is decrypted.
aster_services::apply_ram_protection();
// RAM-first ONLY: the identity is issued through the client and provided on stdin; this
// process writes nothing to disk (a file-backed bot would persist sessions/contacts,
// defeating the seizure-minimisation the shim relies on).
let mut s = String::new();
if std::io::stdin().read_to_string(&mut s).is_err() {
eprintln!("error: could not read the identity from stdin");
std::process::exit(1);
}
let bytes = match from_hex(s.trim()) {
Some(b) if b.len() == 64 => b,
_ => {
eprintln!("error: identity must be 128 hex chars (64 bytes: sign || dh)");
std::process::exit(1);
}
};
let (mut sign, mut dh) = ([0u8; 32], [0u8; 32]);
sign.copy_from_slice(&bytes[..32]);
dh.copy_from_slice(&bytes[32..]);
let mut bot = Bot::from_identity_secret(sign, dh);
let http_port: u16 = match port.parse() {
Ok(p) => p,
Err(_) => {
eprintln!("error: bad http_port");
std::process::exit(1);
}
};
println!("bot address: {}", bot.address());
println!("bot-api on http://127.0.0.1:{port} (Authorization: Bearer {token})");
if let Err(e) = bot
.serve_shim(addr, pubhex, http_port, token.as_bytes().to_vec(), "Echo".into(), "echobot".into())
.await
{
eprintln!("fatal: {e}");
std::process::exit(1);
}
}
demobot — кнопки и медиа — Демонстрационный бот: инлайн-кнопки (🔴/🟢), команды `/photo`, `/video`, `/file`, `/media`, `/big` для тестирования отправки файлов. Встраивает тестовые медиа-файлы в бинарь (`include_bytes!`).
crates/aster-services/src/bin/demobot.rs
//! demobot — a tiny showcase/test bot. Any message → a reply with two inline buttons; a tap →
//! a reply naming the button. Media test commands send REAL sample files so the browser client
//! can be verified end to end: `/photo` (JPEG), `/video` (MP4), `/file` (PDF), `/media` (all three),
//! `/big` (~150 KB multi-chunk stress). Proves the browser bot-UI + decoupled media path against the
//! live client. RAM-first: identity from stdin.
//!
//! Usage: printf '<128-hex sign||dh>' | demobot <server_addr> <server_pubhex>
//! (set ASTER_TLS=1 for a TLS-camouflaged relay)
use std::io::Read;
use aster_net::from_hex;
use aster_services::{Bot, Button, FileSpec, Incoming, Msg, Outgoing, ServiceFrame, Update};
// Real sample media, embedded at compile time (crates/aster-services/assets/) — no runtime disk
// reads, no external fixtures. Each is a genuine, valid file so the client renders/opens it.
const TEST_PHOTO: &[u8] = include_bytes!("../../assets/test-photo.jpg"); // 480×360 JPEG test pattern
const TEST_VIDEO: &[u8] = include_bytes!("../../assets/test-video.mp4"); // 320×240 2s MP4 (H.264/yuv420p)
const TEST_DOC: &[u8] = include_bytes!("../../assets/test-doc.pdf"); // minimal 1-page PDF
fn photo() -> Outgoing {
FileSpec::new(TEST_PHOTO.to_vec(), "test-photo.jpg", "image/jpeg").with_dims(480, 360).into()
}
fn video() -> Outgoing {
FileSpec::new(TEST_VIDEO.to_vec(), "test-video.mp4", "video/mp4").with_dims(320, 240).into()
}
fn document() -> Outgoing {
FileSpec::new(TEST_DOC.to_vec(), "test-doc.pdf", "application/pdf").into()
}
/// Command hint shown on plain messages / unknown commands so a tester knows what to try.
const HINT: &str = "Команды: /photo /video /file /media /big — или жми кнопку 👇";
#[tokio::main]
async fn main() {
let a: Vec<String> = std::env::args().skip(1).collect();
if a.len() != 2 {
eprintln!("usage: printf '<128hex>' | demobot <server_addr> <server_pubhex>");
std::process::exit(2);
}
aster_services::apply_ram_protection();
let mut s = String::new();
std::io::stdin().read_to_string(&mut s).ok();
let bytes = match from_hex(s.trim()) {
Some(b) if b.len() == 64 => b,
_ => {
eprintln!("error: identity must be 128 hex chars (64 bytes)");
std::process::exit(1);
}
};
let (mut sign, mut dh) = ([0u8; 32], [0u8; 32]);
sign.copy_from_slice(&bytes[..32]);
dh.copy_from_slice(&bytes[32..]);
let mut bot = Bot::from_identity_secret(sign, dh);
println!("demobot address: {}", bot.address());
let handler: aster_services::Handler = Box::new(|u: &Update| -> Vec<Outgoing> {
match &u.event {
Incoming::Message { text } => vec![menu(&format!("Ты написал: «{text}».\n{HINT}")).into()],
// Media test commands — each sends a REAL decoupled file the client renders/downloads.
Incoming::Command { name, .. } if name == "photo" => {
vec![ServiceFrame::text("📷 Тестовое фото (JPEG)").into(), photo()]
}
Incoming::Command { name, .. } if name == "video" => {
vec![ServiceFrame::text("🎬 Тестовое видео (MP4)").into(), video()]
}
Incoming::Command { name, .. } if name == "file" => {
vec![ServiceFrame::text("📄 Тестовый документ (PDF)").into(), document()]
}
Incoming::Command { name, .. } if name == "media" => {
vec![ServiceFrame::text("Все тестовые вложения 👇").into(), photo(), video(), document()]
}
// /big exercises the MULTI-CHUNK streaming path (≈150 KB → 4 feed-blob chunks).
Incoming::Command { name, .. } if name == "big" => {
let data: Vec<u8> = (0..150_000).map(|i: usize| (i.wrapping_mul(131).wrapping_add(7)) as u8).collect();
vec![
ServiceFrame::text("📦 Большой файл (~150 КБ, многочанковый)").into(),
FileSpec::new(data, "big.bin", "application/octet-stream").into(),
]
}
Incoming::Command { name, .. } if name == "start" || name == "help" => {
vec![menu(&format!("demobot на связи. {HINT}")).into()]
}
Incoming::Command { name, .. } => vec![menu(&format!("Команда /{name}. {HINT}")).into()],
Incoming::Callback { data, .. } => {
let d = String::from_utf8_lossy(data);
let label = match d.as_ref() {
"red" => "🔴 красную",
"green" => "🟢 зелёную",
other => other,
};
vec![ServiceFrame::text(format!("Ты нажал {label} кнопку ✅")).into()]
}
Incoming::Frame(_) => Vec::new(),
}
});
if let Err(e) = bot.serve(&a[0], &a[1], handler).await {
eprintln!("fatal: {e}");
std::process::exit(1);
}
}
/// A text message carrying two inline buttons (with a fresh client_msg_id so taps reference it).
fn menu(text: &str) -> ServiceFrame {
ServiceFrame::Msg(Msg {
text: Some(text.to_string()),
buttons: vec![vec![
Button { text: "🔴 Красная".into(), data: b"red".to_vec() },
Button { text: "🟢 Зелёная".into(), data: b"green".to_vec() },
]],
reply_to: None,
client_msg_id: rand::random(),
})
}
taskbot — таск-трекер — Коллаборативный трекер задач в стиле Trello: пространства с инвайт-кодами, задачи со статусами (todo→doing→review→done), назначение исполнителя, комментарии, пуши участникам. Состояние в JSON.
crates/aster-services/src/bin/taskbot.rs
//! taskbot — коллаборативный таск-трекер для Aster (Trello-стиль).
//! Пространства + инвайт-коды, задачи со статусами/исполнителем/комментариями,
//! пуши участникам. Всё состояние — в постоянном хранилище:
//! .taskbot-botstate.bin (E2E-сессии, persist-runtime) + .taskbot-state.json (данные).
//! Usage: printf '<128hex>' | taskbot <relay_addr> <relay_pubhex>
use std::collections::HashMap;
use std::io::{Read, Write};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use aster_net::from_hex;
use aster_services::{Bot, Button, Handler, Incoming, Msg, Outgoing, ServiceFrame, Update};
const STATE_PATH: &str = "/home/ubu/aster/.taskbot-state.json";
const ROSTER_PATH: &str = "/home/ubu/aster/.taskbot-roster.json";
const SESS_PATH: &str = "/home/ubu/aster/.taskbot-botstate.bin";
const CODE_ALPHABET: &[u8] = b"23456789ABCDEFGHJKMNPQRSTUVWXYZ";
const MAX_WS_PER_USER: usize = 8;
const MAX_TASKS_PER_WS: usize = 100;
// ── модель ──────────────────────────────────────────────────────────────────
#[derive(Clone, serde::Serialize, serde::Deserialize)]
struct Member {
sign: String, // hex peer_sign (= идентификатор юзера)
dh: String, // hex peer_dh (для пушей)
role: String, // owner | member | removed
}
#[derive(Clone, serde::Serialize, serde::Deserialize)]
struct Activity {
ts: u64,
author: String, // короткое имя автора
kind: String, // created|status|assigned|comment|deleted
text: String,
}
#[derive(Clone, serde::Serialize, serde::Deserialize)]
struct Task {
id: u64,
title: String,
desc: String,
status: String, // todo|doing|done|cancel
assignee: Option<String>,
created_by: String,
activity: Vec<Activity>,
}
#[derive(Clone, serde::Serialize, serde::Deserialize)]
struct Workspace {
id: String,
name: String,
code: String,
owner: String,
members: Vec<Member>,
tasks: Vec<Task>,
next_task: u64,
}
#[derive(serde::Serialize, serde::Deserialize, Default)]
struct Db {
workspaces: Vec<Workspace>,
}
#[derive(Clone)]
enum Wizard {
CreateWs,
Join,
NewTask(String, String), // wid, title (накоплено)
Comment(String, u64), // wid, tid
}
type Roster = HashMap<String, ([u8; 32], [u8; 32])>; // sign_hex -> keys (для пушей)
struct App {
db: Mutex<Db>,
roster: Mutex<Roster>,
wizards: Mutex<HashMap<String, Wizard>>,
outbox: Mutex<Vec<aster_services::Push>>, // отложенные пуши другим пирам
}
fn now() -> u64 {
std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).map(|d| d.as_secs()).unwrap_or(0)
}
fn hexs(b: &[u8; 32]) -> String { b.iter().map(|x| format!("{x:02x}")).collect::<String>() }
fn unhex(s: &str) -> [u8; 32] {
let mut o = [0u8; 32];
let s = &s[..s.len().min(64)];
for i in 0..s.len() / 2 { o[i] = u8::from_str_radix(&s[i * 2..i * 2 + 2], 16).unwrap_or(0); }
o
}
fn uname(sign_hex: &str) -> String { format!("@{}", &sign_hex[56..60]) } // короткое имя из хвоста ключа
fn gen_code() -> String {
(0..6).map(|_| CODE_ALPHABET[rand::random::<usize>() % CODE_ALPHABET.len()] as char).collect()
}
impl App {
fn save(&self) {
let Ok(db) = self.db.lock() else { return };
let tmp = format!("{STATE_PATH}.tmp");
if std::fs::write(&tmp, serde_json::to_string(&*db).unwrap_or_default()).is_ok() {
let _ = std::fs::rename(&tmp, STATE_PATH);
}
}
// ── поиск ──
fn ws_by_code(&self, code: &str) -> Option<usize> {
self.db.lock().ok()?.workspaces.iter().position(|w| w.code.eq_ignore_ascii_case(code))
}
fn ws_idx(&self, id: &str) -> Option<usize> {
self.db.lock().ok()?.workspaces.iter().position(|w| w.id == id)
}
fn member(&self, wi: usize, sign: &str) -> Option<Member> {
self.db.lock().ok()?.workspaces.get(wi)?.members.iter().find(|m| m.sign == sign && m.role != "removed").cloned()
}
}
const STATUSES: &[(&str, &str)] = &[("todo", "🆕 К выполнению"), ("doing", "🔨 В работе"), ("review", "🔍 На проверке"), ("done", "✅ Готово"), ("cancel", "⛔ Отменена")];
/// Разрешённые переходы (воронка как в Jira): из любой не-финальной можно в отмену.
fn next_statuses(cur: &str) -> Vec<&'static str> {
match cur {
"todo" => vec!["doing"],
"doing" => vec!["review", "todo"],
"review" => vec!["done", "doing"], // принять / на доработку
"done" => vec!["todo"], // переоткрытие
"cancel" => vec!["todo"],
_ => vec![],
}
}
fn st_emoji(s: &str) -> &'static str { STATUSES.iter().find(|(k, _)| *k == s).map(|(_, v)| *v).unwrap_or("?") }
// ── рендеры ─────────────────────────────────────────────────────────────────
fn msgf(text: String, buttons: Vec<Vec<Button>>) -> Msg {
Msg { text: Some(text), buttons, reply_to: None, client_msg_id: rand::random() }
}
fn frame(text: String, buttons: Vec<Vec<Button>>) -> Outgoing {
ServiceFrame::Msg(msgf(text, buttons)).into()
}
fn outm(m: Msg) -> Outgoing { ServiceFrame::Msg(m).into() }
fn btn(text: &str, data: String) -> Button { Button { text: text.into(), data: data.into_bytes() } }
fn menu_main() -> Msg {
msgf("📋 TaskBot — совместные задачи.\n\n➕ Создайте пространство и поделитесь кодом с командой,\nили войдите по чужому коду.".into(), vec![
vec![btn("📋 Мои задачи", "my".into()), btn("🗂 Мои пространства", "spaces".into())],
vec![btn("➕ Создать пространство", "new".into()), btn("🔑 Войти по коду", "join".into())],
vec![btn("❓ Помощь", "help".into())],
])
}
/// Список АКТИВНЫХ задач по всем моим пространствам: (1) назначенные на меня,
/// (2) созданные мной (даже без исполнителя — иначе свежесозданное «пропадает»).
/// Доступ к карточке повторно проверяет членство (`on_cb` → ["tk"]) — безопасно.
fn render_my_tasks(app: &App, me: &str) -> Msg {
let db = app.db.lock().ok();
let mut rows: Vec<Vec<Button>> = Vec::new();
let mut assigned = String::new();
let mut created = String::new();
let mut n_assigned = 0usize;
let mut n_created = 0usize;
if let Some(db) = db {
for w in &db.workspaces {
for x in &w.tasks {
if x.status == "done" || x.status == "cancel" { continue; }
let is_asg = x.assignee.as_deref() == Some(me);
let is_own = x.created_by == me;
if !is_asg && !is_own { continue; }
let label = format!("#{} {} · {}", x.id, st_emoji(&x.status), short(&w.name, 16));
let row = vec![btn(label.as_str(), format!("tk:{}:{}", w.id, x.id))];
if is_asg {
assigned.push_str(&format!("#{} {} {} (📦 {})\n", x.id, st_emoji(&x.status), x.title, w.name));
rows.push(row); n_assigned += 1;
} else {
created.push_str(&format!("#{} {} {}{}\n", x.id, st_emoji(&x.status), x.title,
if x.assignee.is_some() { format!(" · исп: {}", uname(x.assignee.as_deref().unwrap_or(""))) } else { String::new() }));
rows.push(row); n_created += 1;
}
}
}
}
let mut text = String::from("📋 Мои активные задачи:\n");
if n_assigned > 0 {
text.push_str(&format!("\n🙋 Назначены на вас ({n_assigned}):\n{assigned}"));
} else {
text.push_str("\n🙋 Назначенные на вас: (нет)");
}
if n_created > 0 {
text.push_str(&format!("\n✍️ Созданные вами ({n_created}):\n{created}"));
} else if n_assigned == 0 {
text.push_str("\n✍️ Созданные вами: (нет)");
}
if n_assigned == 0 && n_created == 0 {
text.push_str("\n\nСоздайте первую: откройте пространство → 🆕 Задача.");
}
rows.push(vec![btn("🏠 Меню", "home".into())]);
msgf(text, rows)
}
/// Мои пространства: все, где я участник (owner/member), с кнопкой входа в каждое.
fn render_my_spaces(app: &App, me: &str) -> Msg {
let db = app.db.lock().ok();
let mut rows: Vec<Vec<Button>> = Vec::new();
let mut text = String::from("🗂 Мои пространства:\n");
let mut total = 0usize;
if let Some(db) = db {
for w in &db.workspaces {
let Some(m) = w.members.iter().find(|m| m.sign == me) else { continue };
if m.role == "removed" { continue; }
let active = w.tasks.iter().filter(|x| x.status != "done" && x.status != "cancel").count();
let role_tag = if m.role == "owner" { "👑" } else { "" };
text.push_str(&format!("\n{}📦 {}\n задач: {} (активных: {})\n", role_tag, w.name, w.tasks.len(), active));
rows.push(vec![btn(format!("Открыть «{}»", short(&w.name, 20)).as_str(), format!("ws:{}", w.id))]);
total += 1;
}
}
if total == 0 { text.push_str("\n(вы пока не состоите ни в одном пространстве)"); }
rows.push(vec![btn("➕ Создать пространство", "new".into()), btn("🔑 Войти по коду", "join".into())]);
rows.push(vec![btn("🏠 Меню", "home".into())]);
msgf(text, rows)
}
fn help_msg() -> Msg {
msgf("🤖 TaskBot: совместная работа над задачами.\n\n• Создайте пространство — получите код приглашения\n• Участники входят по коду\n• Задачи: статусы, исполнитель, комментарии\n• Все участники получают уведомления об изменениях\n\nКоманды: /start — меню".into(), vec![
vec![btn("🏠 Меню", "home".into())],
])
}
fn render_ws(app: &App, wi: usize, me: &str) -> Option<Msg> {
let db = app.db.lock().ok()?;
let w = db.workspaces.get(wi)?;
let mut rows = vec![
vec![btn("🆕 Задача", format!("nt:{}", w.id)), btn("📊 Сводка", format!("sum:{}", w.id))],
vec![btn("👥 Участники", format!("mem:{}", w.id)),
btn(format!("🗂 Архив ({})", w.tasks.iter().filter(|x| x.status == "done" || x.status == "cancel").count()).as_str(), format!("arc:{}", w.id))],
];
if w.owner == me {
rows.push(vec![btn("🔗 Новый код", format!("regen:{}", w.id)), btn("🗑 Расформировать", format!("drop:{}", w.id))]);
} else {
rows.push(vec![btn("🚪 Покинуть", format!("leave:{}", w.id))]);
}
let mut t = format!("📦 {}\nКод приглашения: {}\nЗадач: {}\n", w.name, w.code, w.tasks.len());
let mut list: Vec<&Task> = w.tasks.iter().filter(|x| x.status != "done" && x.status != "cancel").collect();
list.sort_by_key(|x| x.id);
for x in list.iter().rev().take(15).rev() {
let asg = x.assignee.as_ref().map(|a| format!(" · {}", uname(a))).unwrap_or_default();
t.push_str(&format!("\n#{} {} {}{}", x.id, st_emoji(&x.status), x.title, asg));
rows.push(vec![btn(format!("#{} {}", x.id, short(&x.title, 18)).as_str(), format!("tk:{}:{}", w.id, x.id))]);
}
if w.tasks.iter().all(|x| x.status == "done" || x.status == "cancel") && !w.tasks.is_empty() {
t.push_str("\n(активных задач нет)");
}
rows.push(vec![btn("🏠 Меню", "home".into())]);
Some(msgf(t, rows))
}
fn short(s: &str, n: usize) -> String {
if s.chars().count() <= n { s.to_string() } else { format!("{}…", s.chars().take(n - 1).collect::<String>()) }
}
fn render_card(app: &App, wi: usize, tid: u64, me: &str) -> Option<Msg> {
let db = app.db.lock().ok()?;
let w = db.workspaces.get(wi)?;
let tk = w.tasks.iter().find(|t| t.id == tid)?;
let is_owner = w.owner == me;
let mut rows: Vec<Vec<Button>> = Vec::new();
let mut strow = Vec::new();
for k in next_statuses(&tk.status) {
let (_, label) = STATUSES.iter().find(|(kk, _)| *kk == k).unwrap();
strow.push(btn(label, format!("st:{}:{}:{}", w.id, tid, k)));
}
if !strow.is_empty() { rows.push(strow); }
rows.push(vec![
btn(if tk.assignee.as_deref() == Some(me) { "🙋 Снять себя" } else { "🙋 Взять себе" }, format!("as:{}:{}", w.id, tid)),
btn("💬 Комментарий", format!("cm:{}:{}", w.id, tid)),
]);
if is_owner || tk.created_by == me {
rows.push(vec![btn("🗑 Удалить задачу", format!("del:{}:{}", w.id, tid))]);
}
rows.push(vec![btn("◀ К пространству", format!("ws:{}", w.id))]);
let acts: Vec<String> = tk.activity.iter().rev().take(8).rev()
.map(|a| format!("· {} {}: {}", fmt_ts(a.ts), a.author,
if a.kind == "comment" { format!("«{}»", short(&a.text, 60)) } else { a.text.clone() }))
.collect();
Some(msgf(format!(
"#{} {}\nСтатус: {}{}\nСоздал: {}, {}\n\n📜 Журнал:\n{}",
tk.id, tk.title, st_emoji(&tk.status),
tk.assignee.as_ref().map(|a| format!(" · Исполнитель: {}", uname(a))).unwrap_or_default(),
uname(&tk.created_by), fmt_ts(tk.activity.first().map(|a| a.ts).unwrap_or(0)),
if acts.is_empty() { "(пусто)".into() } else { acts.join("\n") }), rows))
}
fn fmt_ts(ts: u64) -> String {
let d = ts % 86400; format!("{:02}:{:02} {}", d / 3600, (d % 3600) / 60, ts / 86400 % 100)
}
// ── мутации + уведомления ───────────────────────────────────────────────────
/// Рассылка текста всем активным членам пространства кроме автора.
/// Возвращает исходящие кадры (handler-путь) — тикер не нужен: всё событийно.
fn broadcast(app: &App, w: &Workspace, except: &str, text: &str) {
let Ok(roster) = app.roster.lock() else { return };
let mut ob = app.outbox.lock().unwrap();
for m in &w.members {
if m.role == "removed" || m.sign == except { continue }
if let Some((ps, pd)) = roster.get(&m.sign).copied() {
ob.push(aster_services::Push {
peer_sign: ps, peer_dh: pd,
out: ServiceFrame::Msg(Msg {
text: Some(format!("📦 {}\n{text}", short(&w.name, 24))),
buttons: vec![vec![btn("Открыть пространство", format!("ws:{}", w.id))]],
reply_to: None, client_msg_id: rand::random(),
}).into(),
});
}
}
}
// ── main ────────────────────────────────────────────────────────────────────
#[tokio::main(flavor = "multi_thread")]
async fn main() {
let a: Vec<String> = std::env::args().skip(1).collect();
let mut app_roster: Roster = HashMap::new();
if a.len() != 2 { eprintln!("usage: printf '<128hex>' | taskbot <relay_addr> <relay_pubhex>"); std::process::exit(2); }
aster_services::apply_ram_protection();
let mut inp = String::new(); std::io::stdin().read_to_string(&mut inp).ok();
let bytes = from_hex(inp.trim()).expect("need 128 hex identity");
let (mut sign, mut dh) = ([0u8; 32], [0u8; 32]);
sign.copy_from_slice(&bytes[..32]); dh.copy_from_slice(&bytes[32..]);
let mut bot = Bot::from_identity_secret_persist(sign, dh, SESS_PATH);
// загрузка БД + ростера
if let Ok(txt) = std::fs::read_to_string(ROSTER_PATH) {
if let Ok(v) = serde_json::from_str::<serde_json::Value>(&txt) {
if let Some(a) = v["pairs"].as_array() {
for p2 in a {
if let (Some(k), Some(dh)) = (p2[0].as_str(), p2[1].as_str()) {
app_roster.insert(k.to_string(), (unhex(k), unhex(dh)));
}
}
}
}
}
let db = std::fs::read_to_string(STATE_PATH)
.ok().and_then(|t| serde_json::from_str::<Db>(&t).ok())
.unwrap_or_default();
println!("taskbot ready: {} space(s)", db.workspaces.len());
let app = Arc::new(App {
db: Mutex::new(db),
roster: Mutex::new(app_roster),
wizards: Mutex::new(HashMap::new()),
outbox: Mutex::new(Vec::new()),
});
let app_h = app.clone();
let _ = app_roster;
let handler: Handler = Box::new(move |u: &Update| -> Vec<Outgoing> {
let me = hexs(&u.peer_sign);
// ростер: любой входящий обновляет ключи доставки
if let Ok(mut r) = app_h.roster.lock() {
if r.insert(me.clone(), (u.peer_sign, u.peer_dh)).is_none() {
let pairs: Vec<[String; 2]> = r.iter().map(|(k, (_, d))| [k.clone(), hexs(d)]).collect();
let _ = std::fs::write(ROSTER_PATH, serde_json::json!({ "pairs": pairs }).to_string());
}
}
route(&app_h, u, &me)
});
// тикер: доставка накопленных пушей другим пирам
let app_t = app.clone();
let ticker: aster_services::Ticker = Box::new(move || {
app_t.outbox.lock().ok().map(|mut o| std::mem::take(&mut *o)).unwrap_or_default()
});
if let Err(e) = bot.serve_ticked(&a[0], &a[1], handler, std::time::Duration::from_secs(3), ticker).await {
eprintln!("fatal: {e}"); std::process::exit(1);
}
}
fn route(app: &App, u: &Update, me: &str) -> Vec<Outgoing> {
match &u.event {
Incoming::Command { name, .. } => match name.as_str() {
"start" => vec![outm(menu_main())],
"help" => vec![outm(help_msg())],
_ => vec![frame(format!("Неизвестная команда /{name}. /start — меню."), vec![vec![btn("🏠 Меню", "home".into())]])],
},
Incoming::Message { text } => on_text(app, me, text.trim()),
Incoming::Callback { msg_ref, data } => on_cb(app, me, &String::from_utf8_lossy(data), *msg_ref),
Incoming::Frame(_) => Vec::new(),
}
}
fn on_text(app: &App, me: &str, text: &str) -> Vec<Outgoing> {
let wiz = app.wizards.lock().ok().and_then(|w| w.get(me).cloned());
if let Some(wiz) = wiz {
match wiz {
Wizard::CreateWs => {
app.wizards.lock().ok().as_mut().map(|w| w.remove(me));
if text.len() > 60 || text.is_empty() { return vec![frame("Название 1–60 символов. Попробуйте ещё раз:".into(), vec![vec![btn("Отмена", "cancel".into())]])]; }
let id = format!("w{:06x}", rand::random::<u32>());
let code = gen_code();
let mut db = app.db.lock().unwrap();
if db.workspaces.iter().filter(|w| w.members.iter().any(|m| m.sign == me && m.role != "removed")).count() >= MAX_WS_PER_USER {
return vec![frame("Лимит пространств достигнут (8).".into(), vec![vec![btn("🏠 Меню", "home".into())]])];
}
let wi = db.workspaces.len();
db.workspaces.push(Workspace {
id: id.clone(), name: text.to_string(), code: code.clone(), owner: me.to_string(),
members: vec![Member { sign: me.into(), dh: String::new(), role: "owner".into() }],
tasks: Vec::new(), next_task: 1,
});
drop(db); app.save();
return vec![frame(format!("✅ Пространство «{text}» создано.\nКод приглашения: {code}\nПоделитесь им — участники войдут через 🔑 Войти по коду."),
vec![vec![btn("📦 Открыть", format!("ws:{id}"))], vec![btn("🏠 Меню", "home".into())]])];
}
Wizard::Join => {
app.wizards.lock().ok().as_mut().map(|w| w.remove(me));
if text.len() != 6 { return vec![frame("Код состоит из 6 символов. Ещё раз:".into(), vec![vec![btn("Отмена", "cancel".into())]])]; }
let Some(wi) = app.ws_by_code(text) else {
return vec![frame("❌ Пространство с таким кодом не найдено.".into(), vec![vec![btn("🏠 Меню", "home".into())]])];
};
let mut db = app.db.lock().unwrap();
let w = &mut db.workspaces[wi];
if let Some(m) = w.members.iter_mut().find(|m| m.sign == me) {
if m.role == "removed" { return vec![frame("🚫 Вы были удалены из этого пространства.".into(), vec![vec![btn("🏠 Меню", "home".into())]])]; }
return vec![frame("Вы уже участник.".into(), vec![vec![btn("📦 Открыть", format!("ws:{}", w.id))]])];
}
let id = w.id.clone(); let name = w.name.clone(); let owner = w.owner.clone();
w.members.push(Member { sign: me.into(), dh: String::new(), role: "member".into() });
let count = w.members.iter().filter(|m| m.role != "removed").count();
drop(db); app.save();
// Уведомление создателю пространства о новом участнике
if owner != me {
if let Some((ps, pd)) = app.roster.lock().ok().and_then(|r| r.get(&owner).copied()) {
app.outbox.lock().unwrap().push(aster_services::Push {
peer_sign: ps, peer_dh: pd,
out: ServiceFrame::Msg(msgf(
format!("📦 «{}»: 👤 {} присоединился по коду (участников: {})", name, uname(me), count),
vec![vec![btn("Открыть пространство", format!("ws:{id}"))]]),
).into()});
}
}
return vec![frame(format!("✅ Вы присоединились к «{name}»."), vec![vec![btn("📦 Открыть", format!("ws:{id}"))], vec![btn("🏠 Меню", "home".into())]])];
}
Wizard::NewTask(wid, title) => {
app.wizards.lock().ok().as_mut().map(|w| w.remove(me));
return create_task(app, me, &wid, &title, text);
}
Wizard::Comment(wid, tid) => {
app.wizards.lock().ok().as_mut().map(|w| w.remove(me));
return comment(app, me, &wid, tid, text);
}
}
}
if text.starts_with('/') { return vec![outm(menu_main())]; }
// произвольный код без кнопки? попробуем как инвайт
if text.len() == 6 { return on_cb(app, me, "join", rand::random()); }
vec![outm(menu_main())]
}
fn create_task(app: &App, me: &str, wid: &str, title: &str, desc: &str) -> Vec<Outgoing> {
let Some(wi) = app.ws_idx(wid) else { return vec![frame("Пространство не найдено.".into(), vec![])] };
let mut db = app.db.lock().unwrap();
let w = &mut db.workspaces[wi];
if !w.members.iter().any(|m| m.sign == me && m.role != "removed") { return vec![frame("Нет доступа.".into(), vec![])]; }
if w.tasks.len() >= MAX_TASKS_PER_WS { return vec![frame("Лимит задач пространства (100).".into(), vec![])]; }
let id = w.next_task; w.next_task += 1;
let task = Task {
id, title: title.to_string(), desc: desc.to_string(), status: "todo".into(),
assignee: None, created_by: me.into(),
activity: vec![Activity { ts: now(), author: uname(me), kind: "created".into(), text: "задача создана".into() }],
};
w.tasks.push(task);
let snapshot_name = w.name.clone(); let members = w.members.clone();
drop(db); app.save();
broadcast(app, &Workspace { id: wid.into(), name: snapshot_name, code: String::new(), owner: String::new(), members, tasks: Vec::new(), next_task: 0 }, me,
&format!("🆕 #{id} «{}» — новая задача ({})", title, uname(me)));
vec![frame(format!("✅ Задача #{id} создана."), vec![vec![btn("📦 К пространству", format!("ws:{wid}"))]])]
}
fn comment(app: &App, me: &str, wid: &str, tid: u64, text: &str) -> Vec<Outgoing> {
let Some(wi) = app.ws_idx(wid) else { return vec![frame("Пространство не найдено.".into(), vec![])]; };
let (members, wname) = {
let mut db = app.db.lock().unwrap();
let w = &mut db.workspaces[wi];
let Some(tk) = w.tasks.iter_mut().find(|t| t.id == tid) else { return vec![frame("Задача не найдена.".into(), vec![])]; };
tk.activity.push(Activity { ts: now(), author: uname(me), kind: "comment".into(), text: text.chars().take(4000).collect() });
(w.members.clone(), w.name.clone())
};
app.save();
broadcast(app, &Workspace { id: wid.into(), name: wname, code: String::new(), owner: String::new(), members, tasks: Vec::new(), next_task: 0 }, me,
&format!("💬 #{} — комментарий от {}", tid, uname(me)));
vec![frame("✅ Комментарий добавлен.".into(), vec![vec![btn("↩ К карточке", format!("tk:{wid}:{tid}"))]])]
}
fn on_cb(app: &App, me: &str, data: &str, mref: [u8; 16]) -> Vec<Outgoing> {
let parts: Vec<&str> = data.split(':').collect();
// навигационные экраны редактируют сообщение на месте
let ed = |m: Msg| vec![ServiceFrame::Edit { msg_ref: mref, msg: m }.into()];
match parts.as_slice() {
["home"] | ["cancel"] => { app.wizards.lock().ok().as_mut().map(|w| w.remove(me)); ed(menu_main()) }
["my"] => ed(render_my_tasks(app, me)),
["spaces"] => ed(render_my_spaces(app, me)),
["help"] => ed(help_msg()),
["new"] => {
app.wizards.lock().ok().as_mut().map(|w| w.insert(me.into(), Wizard::CreateWs));
vec![frame("Название нового пространства:".into(), vec![vec![btn("Отмена", "cancel".into())]])]
}
["join"] => {
app.wizards.lock().ok().as_mut().map(|w| w.insert(me.into(), Wizard::Join));
vec![frame("Пришлите код пространства (6 символов):".into(), vec![vec![btn("Отмена", "cancel".into())]])]
}
["ws", wid] => match app.ws_idx(wid) {
Some(wi) if app.member(wi, me).is_some() => match render_ws(app, wi, me) {
Some(m) => ed(m),
None => vec![frame("Ошибка рендера.".into(), vec![])],
},
_ => vec![frame("Пространство не найдено или нет доступа.".into(), vec![vec![btn("🏠 Меню", "home".into())]])],
},
["nt", wid] => {
app.wizards.lock().ok().as_mut().map(|w| w.insert(me.into(), Wizard::NewTask((*wid).into(), String::new())));
vec![frame("Заголовок задачи:".into(), vec![vec![btn("Отмена", "cancel".into())]])]
}
["tk", wid, tid] => match (app.ws_idx(wid), tid.parse::<u64>()) {
(Some(wi), Ok(id)) if app.member(wi, me).is_some() => match render_card(app, wi, id, me) {
Some(m) => ed(m),
None => vec![frame("Карточка недоступна.".into(), vec![])],
},
_ => vec![frame("Карточка недоступна.".into(), vec![])],
},
["st", wid, tid, st] => set_status(app, me, wid, tid.parse().unwrap_or(0), st, mref),
["as", wid, tid] => assign_self(app, me, wid, tid.parse().unwrap_or(0), mref),
["cm", wid, tid] => {
app.wizards.lock().ok().as_mut().map(|w| w.insert(me.into(), Wizard::Comment((*wid).into(), tid.parse().unwrap_or(0))));
vec![frame("Текст комментария:".into(), vec![vec![btn("Отмена", "cancel".into())]])]
}
["del", wid, tid] => del_task(app, me, wid, tid.parse().unwrap_or(0)),
["mem", wid] => members_screen(app, me, wid),
["regen", wid] => regen_code(app, me, wid),
["leave", wid] => leave(app, me, wid),
["drop", wid] => drop_ws(app, me, wid),
["arc", wid] => archive_screen(app, me, wid),
["kick", wid, tgt] => kick(app, me, wid, tgt),
["sum", wid] => summary(app, me, wid),
_ => vec![outm(menu_main())],
}
}
fn with_task<F>(app: &App, me: &str, wid: &str, tid: u64, mref: [u8; 16], mut f: F) -> Vec<Outgoing>
where F: FnMut(&mut Task) -> String {
let Some(wi) = app.ws_idx(wid) else { return vec![frame("Не найдено.".into(), vec![])]; };
let mut db = app.db.lock().unwrap();
let w = &mut db.workspaces[wi];
if !w.members.iter().any(|m| m.sign == me && m.role != "removed") { return vec![frame("Нет доступа.".into(), vec![])]; }
let note = {
let Some(tk) = w.tasks.iter_mut().find(|t| t.id == tid) else { return vec![frame("Задача не найдена.".into(), vec![])] };
f(tk)
};
let (members, wname) = (w.members.clone(), w.name.clone());
let card_open = format!("tk:{wid}:{tid}");
drop(db); app.save();
broadcast(app, &Workspace { id: wid.into(), name: wname, code: String::new(), owner: String::new(), members, tasks: Vec::new(), next_task: 0 }, me, ¬e);
match app.ws_idx(wid) {
Some(wi) => match render_card(app, wi, tid, me) {
Some(m) => vec![ServiceFrame::Edit { msg_ref: mref, msg: m }.into()],
None => vec![frame(note, vec![vec![btn("↩ К карточке", card_open)]])],
},
None => vec![frame(note, vec![])],
}
}
fn set_status(app: &App, me: &str, wid: &str, tid: u64, st: &str, mref: [u8; 16]) -> Vec<Outgoing> {
if !STATUSES.iter().any(|(k, _)| *k == st) { return vec![frame("Неизвестный статус.".into(), vec![])]; }
with_task(app, me, wid, tid, mref, |tk| {
let old = tk.status.clone();
tk.status = st.into();
tk.activity.push(Activity { ts: now(), author: uname(me), kind: "status".into(),
text: format!("{} → {}", st_emoji(&old), st_emoji(st)) });
format!("#{} {}: {} → {}", tid, tk.title, st_emoji(&old), st_emoji(st))
})
}
fn assign_self(app: &App, me: &str, wid: &str, tid: u64, mref: [u8; 16]) -> Vec<Outgoing> {
with_task(app, me, wid, tid, mref, |tk| {
if tk.assignee.as_deref() == Some(me) {
tk.assignee = None;
tk.activity.push(Activity { ts: now(), author: uname(me), kind: "assigned".into(), text: "снял себя".into() });
format!("#{} {}: исполнитель снят", tid, tk.title)
} else {
tk.assignee = Some(me.into());
tk.activity.push(Activity { ts: now(), author: uname(me), kind: "assigned".into(), text: format!("взята в работу: {}", uname(me)) });
if tk.status == "todo" { tk.status = "doing".into(); }
format!("#{} {}: исполнителем стал {}", tid, tk.title, uname(me))
}
})
}
fn del_task(app: &App, me: &str, wid: &str, tid: u64) -> Vec<Outgoing> {
let Some(wi) = app.ws_idx(wid) else { return vec![frame("Не найдено.".into(), vec![])]; };
let ok = {
let mut db = app.db.lock().unwrap();
let w = &mut db.workspaces[wi];
if w.owner != me { return vec![frame("Только владелец может удалять задачи.".into(), vec![])]; }
let before = w.tasks.len();
w.tasks.retain(|t| t.id != tid);
before != w.tasks.len()
};
if ok { app.save(); vec![frame(format!("🗑 Задача #{tid} удалена."), vec![vec![btn("📦 К пространству", format!("ws:{wid}"))]])] }
else { vec![frame("Задача не найдена.".into(), vec![])] }
}
fn members_screen(app: &App, me: &str, wid: &str) -> Vec<Outgoing> {
let Some(wi) = app.ws_idx(wid) else { return vec![] };
let Some(db) = app.db.lock().ok() else { return vec![] };
let w = &db.workspaces[wi];
let is_owner = w.owner == me;
let mut t = format!("👥 {} · код: {}\n\n", w.name, w.code);
let mut rows: Vec<Vec<Button>> = Vec::new();
for m in &w.members {
if m.role == "removed" { continue }
t.push_str(&format!("\n{} {} ({})", if m.role == "owner" { "👑" } else { "👤" }, uname(&m.sign), m.role));
if is_owner && m.role != "owner" {
rows.push(vec![btn(format!("✖ выгнать {}", uname(&m.sign)).as_str(), format!("kick:{}:{}", wid, m.sign))]);
}
}
rows.push(vec![btn("🔗 Сгенерировать новый код", format!("regen:{}", wid))]);
rows.push(vec![btn("◀ К пространству", format!("ws:{wid}"))]);
drop(db);
vec![frame(t, rows)]
}
fn regen_code(app: &App, me: &str, wid: &str) -> Vec<Outgoing> {
let Some(wi) = app.ws_idx(wid) else { return vec![] };
let mut db = app.db.lock().unwrap();
let w = &mut db.workspaces[wi];
if w.owner != me { return vec![frame("Только владелец.".into(), vec![])]; }
w.code = gen_code();
let code = w.code.clone();
drop(db); app.save();
vec![frame(format!("🔗 Новый код: {code}"), vec![vec![btn("◀ Назад", format!("mem:{wid}"))]])]
}
fn kick(app: &App, me: &str, wid: &str, target: &str) -> Vec<Outgoing> {
let Some(wi) = app.ws_idx(wid) else { return vec![] };
let mut db = app.db.lock().unwrap();
let w = &mut db.workspaces[wi];
if w.owner != me || target == w.owner { return vec![frame("Нельзя.".into(), vec![])]; }
if let Some(m) = w.members.iter_mut().find(|m| m.sign == target) { m.role = "removed".into(); }
drop(db); app.save();
members_screen(app, me, wid)
}
fn leave(app: &App, me: &str, wid: &str) -> Vec<Outgoing> {
let Some(wi) = app.ws_idx(wid) else { return vec![] };
let name;
{
let mut db = app.db.lock().unwrap();
let w = &mut db.workspaces[wi];
if w.owner == me { return vec![frame("Владелец не может выйти — расформируйте пространство.".into(), vec![])]; }
name = w.name.clone();
if let Some(m) = w.members.iter_mut().find(|m| m.sign == me) { m.role = "removed".into(); }
}
app.save();
vec![frame(format!("Вы покинули «{name}»."), vec![vec![btn("🏠 Меню", "home".into())]])]
}
fn drop_ws(app: &App, me: &str, wid: &str) -> Vec<Outgoing> {
let mut db = app.db.lock().unwrap();
if let Some(w) = db.workspaces.iter().find(|w| w.id == wid) {
if w.owner != me { return vec![frame("Только владелец.".into(), vec![])]; }
}
db.workspaces.retain(|w| !(w.id == wid && w.owner == me));
drop(db); app.save();
vec![frame("Пространство расформировано.".into(), vec![vec![btn("🏠 Меню", "home".into())]])]
}
fn archive_screen(app: &App, me: &str, wid: &str) -> Vec<Outgoing> {
let Some(wi) = app.ws_idx(wid) else { return vec![] };
let Some(db) = app.db.lock().ok() else { return vec![] };
let w = &db.workspaces[wi];
if !w.members.iter().any(|m| m.sign == me && m.role != "removed") { return vec![] }
let mut rows: Vec<Vec<Button>> = Vec::new();
let mut t = String::from("🗂 Завершённые и отменённые задачи:\n");
let mut done_list: Vec<&Task> = w.tasks.iter().filter(|x| x.status == "done" || x.status == "cancel").collect();
done_list.sort_by_key(|x| x.id);
for x in done_list.iter().rev().take(20).rev() {
t.push_str(&format!("\n#{} {} {} · {}", x.id, st_emoji(&x.status), short(&x.title, 24),
x.activity.last().map(|a| fmt_ts(a.ts)).unwrap_or_default()));
rows.push(vec![btn(format!("#{} {}", x.id, short(&x.title, 18)).as_str(), format!("tk:{}:{}", w.id, x.id))]);
}
if done_list.is_empty() { t.push_str("\n(пока пусто)"); }
rows.push(vec![btn("◀ К пространству", format!("ws:{}", wid))]);
vec![outm(msgf(t, rows))]
}
fn summary(app: &App, me: &str, wid: &str) -> Vec<Outgoing> {
let Some(wi) = app.ws_idx(wid) else { return vec![] };
let Some(db) = app.db.lock().ok() else { return vec![] };
let w = &db.workspaces[wi];
if !w.members.iter().any(|m| m.sign == me && m.role != "removed") { return vec![] }
let c = |s: &str| w.tasks.iter().filter(|t| t.status == s).count();
let mine = w.tasks.iter().filter(|t| t.assignee.as_deref() == Some(me) && t.status != "done" && t.status != "cancel").count();
vec![frame(format!(
"📊 «{}»: всего {}\n🆕 {} · 🔨 {} · ✅ {} · ⛔ {}\n👤 Мои активные: {mine}",
short(&w.name, 30), w.tasks.len(), c("todo"), c("doing"), c("done"), c("cancel")),
vec![vec![btn("◀ К пространству", format!("ws:{wid}"))]])]
}
newsbot — новостной издатель — Публикует новости из JFF в Telegram-канал и владельцам-подписчикам Aster. Периодический поллинг (`/news`, `/mirrors`), дифф-состояние, автоматическая подписка при `/start`. HTTPS в Telegram через `tokio-rustls`.
crates/aster-services/src/bin/newsbot.rs
//! newsbot — публикует новости и новые зеркала из JFF в Telegram-канал и
//! владельцам-подписчикам Aster. Источники — только ПУБЛИЧНЫЕ эндпоинты
//! (GET /news, GET /mirrors по localhost): бот не хранит админ-креденшелов.
//!
//! Usage: printf '<128hex sign||dh>' | newsbot <relay_addr> <relay_pubhex>
//! Env: JFF_API_BASE (default http://127.0.0.1:8000), POLL_SEC=60,
//! TG_BOT_TOKEN, TG_CHAT_ID, NEWSBOT_SECRET (слово владельца),
//! NEWSBOT_SEED_SILENT=1 (первый цикл не постит существующее)
use std::collections::{HashMap, HashSet};
use std::io::{Read, Write};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use aster_net::from_hex;
use aster_services::{
Bot, Button, Handler, Incoming, Msg, Outgoing, Push, ServiceFrame, Ticker, Update,
};
const HTTP_TIMEOUT: Duration = Duration::from_secs(5);
const TG_MAX: usize = 3900;
fn env_s(key: &str, def: &str) -> String {
std::env::var(key).ok().filter(|v| !v.trim().is_empty()).unwrap_or_else(|| def.into())
}
// ── JFF poll (blocking, thread::scope как probe_all у statusbot) ────────────
fn http_get_json(base: &str, path: &str) -> Option<serde_json::Value> {
let full = format!("{}{}", base.trim_end_matches('/'), path);
let (host, port, is_tls) = parse_base(&full)?;
if is_tls { return None } // JFF ходим только по localhost http
let mut stream = match std::net::TcpStream::connect((host.as_str(), port)) {
Ok(x) => x, Err(e) => { eprintln!("[newsbot] fetch connect fail: {e}"); return None }
};
stream.set_read_timeout(Some(HTTP_TIMEOUT)).ok();
let req = format!("GET {path} HTTP/1.0\r\nHost: {host}\r\nConnection: close\r\n\r\n");
stream.write_all(req.as_bytes()).ok()?;
let mut buf = String::new();
if let Err(e) = stream.read_to_string(&mut buf) { eprintln!("[newsbot] fetch read fail: {e}"); return None }
let body = match buf.split("\r\n\r\n").nth(1) { Some(b) => b, None => { eprintln!("[newsbot] fetch: no body in {} bytes", buf.len()); return None } };
match serde_json::from_str(body) { Ok(v) => Some(v), Err(e) => { eprintln!("[newsbot] fetch parse fail: {e} | head: {}", &body[..body.len().min(120)]); None } }
}
/// "http://127.0.0.1:8000" → ("127.0.0.1", 8000, false)
fn parse_base(base: &str) -> Option<(String, u16, bool)> {
let tls = base.starts_with("https://");
let rest = if tls { base.strip_prefix("https://")? } else { base.strip_prefix("http://")? };
let hostport = rest.split('/').next()?; // отрезать путь: host:port/path → host:port
let (h, p) = hostport.split_once(':')?;
Some((h.to_string(), p.parse().ok()?, tls))
}
// ── Telegram (минимальный HTTPS/1.1 на tokio-rustls) ───────────────────────
pub fn tg_escape(s: &str) -> String {
s.replace('&', "&").replace('<', "<").replace('>', ">")
}
pub fn chunk_text(s: &str, max: usize) -> Vec<String> {
let mut out = Vec::new(); let mut cur = String::new();
for line in s.lines() {
if cur.len() + line.len() + 1 > max { out.push(std::mem::take(&mut cur)); }
if !cur.is_empty() { cur.push('\n'); }
cur.push_str(line);
}
if !cur.is_empty() { out.push(cur); }
if out.is_empty() { out.push(String::new()); }
out
}
pub fn tg_send_blocking(token: &str, chat_id: &str, text: &str) -> Result<(), String> {
tokio::task::block_in_place(|| {
tokio::runtime::Handle::current().block_on(tg_send(token, chat_id, text))
})
}
async fn tg_send(token: &str, chat_id: &str, text: &str) -> Result<(), String> {
use tokio_rustls::rustls::{ClientConfig, RootCertStore};
let roots = RootCertStore { roots: webpki_roots::TLS_SERVER_ROOTS.to_vec() };
let cfg = Arc::new(ClientConfig::builder().with_root_certificates(roots).with_no_client_auth());
let name = rustls_pki_names("api.telegram.org");
let server = format!("api.telegram.org:443");
let sock = tokio::net::TcpStream::connect(&server).await.map_err(|e| e.to_string())?;
let connector = tokio_rustls::TlsConnector::from(cfg);
let domain = tokio_rustls::rustls::pki_types::ServerName::try_from(name.clone()).map_err(|e| e.to_string())?;
let mut tls = connector.connect(domain, sock).await.map_err(|e| e.to_string())?;
let payload = serde_json::json!({
"chat_id": chat_id, "text": text, "parse_mode": "HTML",
"disable_web_page_preview": true
});
let body = payload.to_string();
let host = "api.telegram.org";
let req = format!(
"POST /bot{token}/sendMessage HTTP/1.1\r\nHost: {host}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}", body.len());
use tokio::io::{AsyncReadExt, AsyncWriteExt};
tls.write_all(req.as_bytes()).await.map_err(|e| e.to_string())?;
let mut resp = String::new(); tls.read_to_string(&mut resp).await.map_err(|e| e.to_string())?;
if resp.contains("\"ok\":true") { Ok(()) } else { Err(format!("tg: {}", &resp[..resp.len().min(160)])) }
}
fn rustls_pki_names(_h: &str) -> String { "api.telegram.org".into() }
// ── diff state ───────────────────────────────────────────────────────────────
struct State {
mirrors: Mutex<HashSet<String>>,
last_news_id: AtomicU64,
}
const STATE_PATH: &str = "/home/ubu/aster/.newsbot-state.json";
fn hexs(b: &[u8; 32]) -> String { b.iter().map(|x| format!("{x:02x}")).collect() }
fn unhex(s: &str) -> [u8; 32] {
let mut o = [0u8; 32];
for i in 0..32 { o[i] = u8::from_str_radix(&s[i*2..i*2+2], 16).unwrap_or(0); }
o
}
type Roster = HashMap<[u8; 32], [u8; 32]>;
fn save_state(r: &Roster, st: &State) {
let v = serde_json::json!({
"roster": r.iter().map(|(k, v)| [hexs(k), hexs(v)]).collect::<Vec<_>>(),
"last_news_id": st.last_news_id.load(Ordering::Relaxed),
"mirrors": st.mirrors.lock().map(|g| g.iter().cloned().collect::<Vec<_>>()).unwrap_or_default(),
});
let tmp = format!("{STATE_PATH}.tmp");
if std::fs::write(&tmp, v.to_string()).is_ok() {
let _ = std::fs::rename(&tmp, STATE_PATH);
}
}
fn load_state(r: &mut Roster, st: &State) {
let Ok(txt) = std::fs::read_to_string(STATE_PATH) else { return };
let Ok(v) = serde_json::from_str::<serde_json::Value>(&txt) else { return };
if let Some(a) = v["roster"].as_array() {
for pair in a {
if let (Some(k), Some(val)) = (pair[0].as_str(), pair[1].as_str()) {
r.insert(unhex(k), unhex(val));
}
}
}
st.last_news_id.store(v["last_news_id"].as_u64().unwrap_or(0), Ordering::Relaxed);
if let Some(a) = v["mirrors"].as_array() {
*st.mirrors.lock().unwrap() = a.iter().filter_map(|x| x.as_str().map(String::from)).collect();
}
eprintln!("[newsbot] state loaded: {} owners, last_id={}", r.len(), st.last_news_id.load(Ordering::Relaxed));
}
/// Возвращает тексты для публикации (пусто на сид-цикле).
fn poll_and_diff(cfg: &Cfg, st: &State) -> Vec<String> {
let mut texts = Vec::new();
eprintln!("[newsbot] poll_and_diff enter, base={}", cfg.jff);
// зеркала
let mv = http_get_json(&cfg.jff, "/mirrors");
eprintln!("[newsbot] mirrors json some={} details_len={}", mv.is_some(), mv.as_ref().and_then(|v| v.get("details")).and_then(|d| d.as_array()).map(|a| a.len()).unwrap_or(usize::MAX));
if let Some(v) = mv {
if let Some(arr) = v["details"].as_array() {
let now: HashSet<String> = arr.iter().filter_map(|d| d["domain"].as_str().map(String::from)).collect();
let mut seen = match st.mirrors.lock() { Ok(g) => g, Err(po) => po.into_inner() };
eprintln!("[newsbot] diff: fetched mirrors block"); if seen.is_empty() && cfg.seed_silent { *seen = now; }
else {
for d in &now { if seen.insert(d.clone()) {
let label = arr.iter().find(|x| x["domain"].as_str() == Some(d))
.and_then(|x| x["label"].as_str()).unwrap_or("");
texts.push(format!("🪞 Новое зеркало: https://{d}{}",
if label.is_empty() { String::new() } else { format!(" ({label})") }));
}}
for d in seen.difference(&now).cloned().collect::<Vec<_>>() { seen.remove(&d); }
}
}
}
// новости
let after = st.last_news_id.load(Ordering::Relaxed);
if let Some(v) = http_get_json(&cfg.jff, &format!("/news?after={after}&limit=20")) {
let last = v["last_id"].as_u64().unwrap_or(after);
st.last_news_id.store(last, Ordering::Relaxed);
if let Some(items) = v["items"].as_array() {
for it in items.iter().rev() { // старые вперёд
let t = it["title"].as_str().unwrap_or(""); if t.is_empty() { continue }
let b = it["body"].as_str().unwrap_or("");
let date = it["published_at"].as_str().and_then(|s| s.get(..16)).unwrap_or("");
let _ = date;
texts.push(format!("📰 <b>{}</b>\n\n{}", tg_escape(t), tg_escape(b)));
}
}
}
texts
}
struct Cfg { jff: String, seed_silent: bool, tg_token: String, tg_chat: String }
fn publish_all(cfg: &Cfg, texts: &[String]) {
for raw in texts {
let plain_for_tg = raw.replace("<b>", "").replace("</b>", "");
for chunk in chunk_text(&plain_for_tg, TG_MAX) {
let _ = tg_send_blocking(&cfg.tg_token, &cfg.tg_chat, &chunk);
}
}
}
#[tokio::main(flavor = "multi_thread")]
async fn main() {
let a: Vec<String> = std::env::args().skip(1).collect();
if a.len() != 2 { eprintln!("usage: printf '<128hex>' | newsbot <relay_addr> <relay_pubhex>"); std::process::exit(2); }
aster_services::apply_ram_protection();
let mut s = String::new(); std::io::stdin().read_to_string(&mut s).ok();
let bytes = from_hex(s.trim()).expect("need 128 hex identity");
let (mut sign, mut dh) = ([0u8; 32], [0u8; 32]);
sign.copy_from_slice(&bytes[..32]); dh.copy_from_slice(&bytes[32..]);
let mut bot = Bot::from_identity_secret_persist(sign, dh, "/home/ubu/aster/.newsbot-botstate.bin");
let cfg = Arc::new(Cfg {
jff: env_s("JFF_API_BASE", "http://127.0.0.1:8000"),
seed_silent: std::env::var("NEWSBOT_SEED_SILENT").map(|v| v == "1").unwrap_or(true),
tg_token: env_s("TG_BOT_TOKEN", ""),
tg_chat: env_s("TG_CHAT_ID", ""),
});
let poll = Duration::from_secs(std::env::var("POLL_SEC").ok().and_then(|v| v.parse().ok()).unwrap_or(60).max(10));
let roster: Arc<Mutex<HashMap<[u8; 32], [u8; 32]>>> = Arc::new(Mutex::new(HashMap::new()));
println!("newsbot address: {}", bot.address());
// reactive handler
let h_roster = roster.clone(); let h_cfg = cfg.clone();
let handler: Handler = Box::new(move |u: &Update| -> Vec<Outgoing> {
let msg_frame = |t: String| ServiceFrame::Msg(Msg {
text: Some(t), reply_to: None,
buttons: vec![vec![Button { text: "🪞 Зеркала".into(), data: b"mirrors:list".to_vec() },
Button { text: "ℹ️ Помощь".into(), data: b"help".to_vec() }]],
client_msg_id: rand::random(),
});
eprintln!("[newsbot] event: {:?}", match &u.event {
Incoming::Message { .. } => "message", Incoming::Command { name, .. } => name.as_str(),
Incoming::Callback { data, .. } => { let d = String::from_utf8_lossy(data); Box::leak(d.into_owned().into_boxed_str()) }
Incoming::Frame(_) => "frame",
});
// Автоподписка: любой, кто написал боту (обычно /start) получает пуши новостей.
if matches!(&u.event, Incoming::Message { .. } | Incoming::Command { .. }) {
let mut r = h_roster.lock().unwrap();
let fresh = r.insert(u.peer_sign, u.peer_dh).is_none();
drop(r);
if fresh {
return vec![ServiceFrame::Msg(Msg { text: Some("✅ Подписка оформлена: новости и новые зеркала будут приходить сюда автоматически.".into()), reply_to: None, client_msg_id: rand::random(),
buttons: vec![vec![Button { text: "🪞 Зеркала".into(), data: b"mirrors:list".to_vec() }, Button { text: "ℹ️ Помощь".into(), data: b"help".to_vec() }]] }).into()];
}
}
let t = text_of(u);
match &u.event {
Incoming::Command { name, .. } => match name.as_str() {
"start" | "help" => vec![msg_frame("📢 Новостной бот JFF.\nНовости и новые зеркала приходят автоматически.\nКнопка «Зеркала» — актуальный список адресов.".into()).into()],
"mirrors" => vec![mirrors_text(&h_cfg).into()],
_ => vec![msg_frame("Неизвестная команда. Доступно: /start, /help, /mirrors".into()).into()],
},
Incoming::Message { .. } => vec![msg_frame("📢 Бот новостей JFF. Новости и зеркала приходят автоматически — вы уже подписаны.".into()).into()],
Incoming::Callback { msg_ref, data } => {
if String::from_utf8_lossy(data) == "mirrors:list" { eprintln!("[newsbot] sending mirrors edit"); vec![ServiceFrame::Edit { msg_ref: *msg_ref, msg: mirrors_msg(&h_cfg) }.into()] }
else if String::from_utf8_lossy(data) == "help" { vec![msg_frame("📢 Новости и зеркала приходят автоматически всем, кто добавил бота.".into()).into()] }
else { Vec::new() }
}
Incoming::Frame(_) => Vec::new(),
}
});
// ticker: diff → publish to TG + Aster roster
let t_cfg = cfg.clone(); let t_roster = roster.clone();
let st = Arc::new(State { mirrors: Mutex::new(HashSet::new()), last_news_id: AtomicU64::new(0) });
load_state(&mut roster.lock().unwrap(), &st);
let t_state = st.clone();
let ticker: Ticker = Box::new(move || -> Vec<Push> {
eprintln!("[newsbot] tick");
let texts = poll_and_diff(&t_cfg, &t_state);
if !texts.is_empty() {
eprintln!("[newsbot] publishing {} item(s) to {} owner(s)", texts.len(), t_roster.lock().unwrap().len());
save_state(&t_roster.lock().unwrap(), &t_state);
}
if texts.is_empty() { return Vec::new() }
publish_all(&t_cfg, &texts); // Telegram sink
let roster = t_roster.lock().unwrap();
let mut pushes = Vec::new();
for raw in &texts {
let plain = strip_html(raw);
for (sign, dh) in roster.iter() {
pushes.push(Push { peer_sign: *sign, peer_dh: *dh,
out: ServiceFrame::Msg(Msg { text: Some(plain.clone()), buttons: vec![],
reply_to: None, client_msg_id: rand::random() }).into() });
}
}
pushes
});
if let Err(e) = bot.serve_ticked(&a[0], &a[1], handler, poll, ticker).await {
eprintln!("fatal: {e}"); std::process::exit(1);
}
fn text_of<'u>(u: &'u Update) -> String {
match &u.event {
Incoming::Message { text } => text.clone(),
Incoming::Command { name, args } => format!("/{name} {args}"),
_ => String::new(),
}
}
fn strip_html(s: &str) -> String {
s.replace("&", "&").replace("<", "<").replace(">", ">")
.replace("<b>", "").replace("</b>", "")
}
fn mirrors_msg(cfg: &Cfg) -> Msg {
let list = http_get_json(&cfg.jff, "/mirrors")
.map(|v| v["details"].as_array().cloned().unwrap_or_default())
.unwrap_or_default();
eprintln!("[newsbot] mirrors_msg: {} item(s)", list.len());
let mut t = String::from("🪞 Актуальные зеркала JFF:\n");
for d in &list {
let dom = d["domain"].as_str().unwrap_or(""); if dom.is_empty() { continue }
t.push_str(&format!("\n• https://{dom}{}", d["label"].as_str().map(|l| format!(" — {l}")).unwrap_or_default()));
}
if list.is_empty() { t.push_str("\n(список пуст — основной адрес см. на сайте)"); }
Msg { text: Some(t), buttons: vec![vec![
Button { text: "🔄 Обновить".into(), data: b"mirrors:list".to_vec() }]],
reply_to: None, client_msg_id: rand::random() }
}
fn mirrors_text(cfg: &Cfg) -> ServiceFrame { ServiceFrame::Msg(mirrors_msg(cfg)) }
}
statusbot — мониторинг сервера — Полноценный бот состояния сервера: CPU/RAM/диск, `systemctl is-active`, HTTP/TCP health-checks (в т.ч. JFF-пресет), кнопки перезапуска с подтверждением, алерты при падении/восстановлении, ежедневные сводки. Секретное слово для входа владельца.
crates/aster-services/src/bin/statusbot.rs
//! statusbot — an Aster bot that reports the host's live server state (like the JFF Telegram
//! monitor bot): CPU load, RAM, disk, uptime, and per-service `systemctl is-active`, with
//! inline "🔄 Перезапустить" buttons (confirm → `sudo -n systemctl restart`, rate-capped).
//! Only whitelisted ASTER1 admins may see status or restart; everyone else is refused.
//!
//! Usage: printf '<128-hex sign||dh>' | statusbot <server_addr> <server_pubhex>
//! Env:
//! STATUSBOT_SERVICES "unit:Label,unit2:Label2" (default: a small demo set)
//! STATUSBOT_ADMINS "ASTER1…,ASTER1…" (empty = everyone allowed — demo only!)
//! ASTER_TLS=1 for a TLS-camouflaged relay
use std::collections::{HashMap, HashSet};
use std::io::Read;
use std::process::Command;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use aster_net::from_hex;
use aster_services::{Bot, Button, Handler, Incoming, Msg, Outgoing, Push, ServiceFrame, Ticker, Update};
const RESTART_CAP_PER_HOUR: usize = 3;
const ROSTER_FILE: &str = ".statusbot-roster";
/// Save the owner roster to disk (only public keys, no secrets).
fn save_roster(roster: &HashMap<[u8; 32], [u8; 32]>) {
let v: Vec<([u8; 32], [u8; 32])> = roster.iter().map(|(&k, &v)| (k, v)).collect();
if let Ok(bytes) = postcard::to_allocvec(&v) {
let _ = std::fs::write(ROSTER_FILE, &bytes);
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let _ = std::fs::set_permissions(ROSTER_FILE, std::fs::Permissions::from_mode(0o600));
}
}
}
/// Load the owner roster from disk (returns empty map if file is missing/corrupt).
fn load_roster() -> HashMap<[u8; 32], [u8; 32]> {
std::fs::read(ROSTER_FILE)
.ok()
.and_then(|bytes| postcard::from_bytes::<Vec<([u8; 32], [u8; 32])>>(&bytes).ok())
.map(|v| v.into_iter().collect())
.unwrap_or_default()
}
struct Cfg {
checks: Vec<Check>, // what to monitor (unit / HTTP / TCP health probes)
secret: String, // owner login word; empty = deny everyone (fail-closed)
allow_restart: bool, // restart buttons/action only when STATUSBOT_ALLOW_RESTART=1
}
/// How a target's health is probed. HTTP/TCP mirror the JFF Telegram monitor's real endpoint checks
/// (a unit can be `active` in systemd while its HTTP is 500/hung — the probe catches that).
enum Probe {
Unit, // `systemctl is-active <unit>` == "active"
Http { url: String, ok: Vec<u16> }, // GET localhost URL; status must be in `ok`
Tcp { host: String, port: u16 }, // a TCP connect must succeed
}
/// One monitored target. `unit` (if any) is the systemd unit to offer a restart button for.
struct Check {
key: String, // stable alert key (e.g. "api", "svc:nginx")
label: String,
unit: Option<String>,
probe: Probe,
}
#[tokio::main]
async fn main() {
let a: Vec<String> = std::env::args().skip(1).collect();
if a.len() != 2 {
eprintln!("usage: printf '<128hex>' | statusbot <server_addr> <server_pubhex>");
std::process::exit(2);
}
aster_services::apply_ram_protection();
let mut s = String::new();
std::io::stdin().read_to_string(&mut s).ok();
let bytes = match from_hex(s.trim()) {
Some(b) if b.len() == 64 => b,
_ => {
eprintln!("error: identity must be 128 hex chars (64 bytes)");
std::process::exit(1);
}
};
let (mut sign, mut dh) = ([0u8; 32], [0u8; 32]);
sign.copy_from_slice(&bytes[..32]);
dh.copy_from_slice(&bytes[32..]);
let mut bot = Bot::from_identity_secret_persist(sign, dh, "/home/ubu/aster/.statusbot-botstate.bin");
let cfg = Arc::new(Cfg {
checks: build_checks(),
secret: std::env::var("STATUSBOT_SECRET").unwrap_or_default().trim().to_string(),
allow_restart: std::env::var("STATUSBOT_ALLOW_RESTART").map(|v| v == "1").unwrap_or(false),
});
// Proactive push config (all additive; sensible defaults so the digest is on out of the box).
let tick = Duration::from_secs(env_u64("STATUSBOT_TICK_SEC", 30).max(1));
let realert = Duration::from_secs(env_u64("STATUSBOT_REALERT_MIN", 30).max(1) * 60);
let digest_hours: Vec<u8> = match std::env::var("STATUSBOT_DIGEST_HOUR") {
Ok(v) if v.trim().is_empty() || v.trim().eq_ignore_ascii_case("off") => Vec::new(),
Ok(v) => v.trim().split(',').filter_map(|s| s.trim().parse::<u8>().ok().filter(|h| *h < 24)).collect(),
Err(_) => vec![10], // default: a daily digest at 10:00 (server local time)
};
let load1_max = std::env::var("STATUSBOT_LOAD1").ok().and_then(|v| v.trim().parse::<f64>().ok());
let mem_max = std::env::var("STATUSBOT_MEM_PCT").ok().and_then(|v| v.trim().parse::<u64>().ok());
let disk_max = std::env::var("STATUSBOT_DISK_PCT").ok().and_then(|v| v.trim().parse::<u64>().ok());
let cpu_max = std::env::var("STATUSBOT_CPU_PCT").ok().and_then(|v| v.trim().parse::<u64>().ok());
println!("statusbot address: {}", bot.address());
println!(" checks: {}", cfg.checks.iter().map(|c| c.key.as_str()).collect::<Vec<_>>().join(", "));
println!(
" auth: {}",
if cfg.secret.is_empty() { "(НЕТ секрета — бот откажет всем; задай STATUSBOT_SECRET)" } else { "секретное слово задано — владелец логинится, прислав его" }
);
println!(" restart: {}", if cfg.allow_restart { "РАЗРЕШЁН (STATUSBOT_ALLOW_RESTART=1)" } else { "выключен (read-only)" });
println!(
" push: тик {}с · алерты вкл · напоминание {} мин · сводка {} · пороги load={:?} mem={:?} disk={:?} cpu={:?}",
tick.as_secs(), realert.as_secs() / 60,
if digest_hours.is_empty() { "выкл".into() } else { format!("ежедневно в {}:00", digest_hours.iter().map(|h| h.to_string()).collect::<Vec<_>>().join(", ")) },
load1_max, mem_max, disk_max, cpu_max
);
// Owners: peer_sign -> peer_dh, shared between the reactive handler (writer, on login) and the
// periodic ticker (reader, for the push fan-out). Keyed by peer_sign — exactly the tuple the push
// path needs. Persisted to `.statusbot-roster` so owners survive restarts.
let roster_map = load_roster();
let roster_count = roster_map.len();
let roster: Arc<Mutex<HashMap<[u8; 32], [u8; 32]>>> = Arc::new(Mutex::new(roster_map));
let roster_dirty: Arc<std::sync::atomic::AtomicBool> = Arc::new(std::sync::atomic::AtomicBool::new(false));
if roster_count > 0 {
println!(" roster: загружен из {ROSTER_FILE} ({roster_count} влад.)");
}
// ---- reactive handler ----
let mut restarts: HashMap<String, Vec<u64>> = HashMap::new();
let h_cfg = cfg.clone();
let h_roster = roster.clone();
let h_dirty = roster_dirty.clone();
let handler: Handler = Box::new(move |u: &Update| -> Vec<Outgoing> {
// Secret-word login: sending EXACTLY the secret enrolls this peer as an owner (captures BOTH
// its keys so the ticker can push to it). The word travels E2E; FAIL-CLOSED — no secret ⇒ nobody.
let text = match &u.event {
Incoming::Message { text } => Some(text.trim().to_string()),
Incoming::Command { name, args } => {
Some(if args.is_empty() { format!("/{name}") } else { format!("/{name} {args}") })
}
_ => None,
};
if !h_cfg.secret.is_empty() {
if let Some(t) = &text {
if t == &h_cfg.secret {
let mut r = h_roster.lock().unwrap();
r.insert(u.peer_sign, u.peer_dh);
save_roster(&r);
drop(r);
h_dirty.store(true, std::sync::atomic::Ordering::Relaxed);
return vec![
ServiceFrame::text("✅ Владелец распознан. Теперь вы будете получать сводку и алерты о сбоях.").into(),
status_frame(&h_cfg).into(),
];
}
}
}
if !h_roster.lock().unwrap().contains_key(&u.peer_sign) {
return vec![ServiceFrame::text("🔒 Доступ только для владельца. Пришли секретное слово.").into()];
}
match &u.event {
Incoming::Command { name, .. } => match name.as_str() {
"start" | "help" => vec![help_frame().into()],
"status" => vec![status_frame(&h_cfg).into()],
"services" => vec![services_frame(&h_cfg).into()],
"logs" => vec![logs_frame(&h_cfg).into()],
other => vec![menu(format!("Неизвестная команда /{other}. Доступно: /status, /services, /logs, /help")).into()],
},
Incoming::Message { .. } => vec![status_frame(&h_cfg).into()], // any text → show status
Incoming::Callback { msg_ref, data } => {
let d = String::from_utf8_lossy(data);
handle_callback(&h_cfg, &mut restarts, &d, *msg_ref)
}
Incoming::Frame(_) => Vec::new(),
}
});
// ---- periodic ticker: probe services/metrics, emit alerts on transitions + a daily digest,
// fan out to every logged-in owner (fresh frame per owner×message → unique client_msg_id). ----
let t_cfg = cfg.clone();
let t_roster = roster.clone();
let host = hostname();
let mut mon = Monitor::new(realert, digest_hours.clone(), load1_max, mem_max, disk_max, cpu_max);
let ticker: Ticker = Box::new(move || -> Vec<Push> {
// No owners logged in → nobody to alert. Skip probing AND advancing any alert/digest state, so
// the FIRST owner to log in during an ongoing outage still gets a fresh 🔴 (and the day's сводка
// is not silently marked sent). Cheap fast-path; monitoring is active whenever ≥1 owner exists.
if t_roster.lock().unwrap().is_empty() {
return Vec::new();
}
let now = Instant::now();
let (l1, l5, l15) = loadavg();
let ram_full = mem();
let (_, ram) = ((), ram_full.1);
let disk_full = disk_root();
let disk = disk_full.0;
let (hour, day, disp) = local_time();
let mut texts: Vec<String> = Vec::new();
// health-check up/down transitions (the core alert, like the Telegram monitor: HTTP/TCP probes).
// Probes run CONCURRENTLY so a tick can't serialize N timeouts and freeze the poll loop.
for (c, res) in t_cfg.checks.iter().zip(probe_all(&t_cfg.checks)) {
let (failing, detail) = match res {
Ok(()) => (false, String::new()),
Err(reason) => (true, reason),
};
if let Some(m) = mon.transition(&c.key, &c.label, failing, &detail, now) {
texts.push(m);
}
}
// optional host-resource thresholds (off unless configured)
if let Some(max) = mon.load1_max {
if let Some(m) = mon.transition("res:load", "Нагрузка (load1)", l1 > max, &format!("load1 {l1:.2} > {max:.2} (5м: {l5:.2}, 15м: {l15:.2})"), now) { texts.push(m); }
}
if let Some(max) = mon.mem_max {
let (ram_used_s, ram_pct) = mem();
if let Some(m) = mon.transition("res:mem", "Память", ram_pct > max, &format!("{ram_pct}% ({ram_used_s}) > {max}%"), now) { texts.push(m); }
}
if let Some(max) = mon.disk_max {
if let Some(m) = mon.transition("res:disk", "Диск /", disk > max, &format!("{disk}% занято, свободно {} > {max}%", disk_full.1), now) { texts.push(m); }
}
if let Some(max) = mon.cpu_max {
let smp = cpu_sample();
if let Some((pt, pi)) = mon.prev_cpu {
let dt = smp.0.saturating_sub(pt);
if dt > 0 {
let di = smp.1.saturating_sub(pi);
let busy_pct = 100u64.saturating_sub(100 * di / dt);
if let Some(m) = mon.transition("res:cpu", "Процессор", busy_pct > max, &format!("{busy_pct}% > {max}%"), now) { texts.push(m); }
}
}
mon.prev_cpu = Some(smp);
}
// once-a-day digest (сводка) at the configured local hours
if hour != u8::MAX {
if let Some(m) = mon.maybe_digest(hour, &day, l1, ram, disk) { texts.push(m); }
}
if texts.is_empty() {
return Vec::new();
}
let owners = t_roster.lock().unwrap();
eprintln!("[statusbot] ALERT out x{} to {} owner(s): {}", texts.len(), owners.len(), texts[0].replace('\n', " | "));
let mut v = Vec::with_capacity(owners.len() * texts.len());
for (sign, dh) in owners.iter() {
let snap = {
let (l1, l5, l15) = loadavg();
format!("\n⚙️ Load {:.2} / {:.2} / {:.2} · 🧠 RAM {}% ({}) · 💾 Диск {}%", l1, l5, l15, ram_full.1, ram_full.0, disk)
};
for t in &texts {
// fresh ServiceFrame per (owner, message) → unique random client_msg_id (never clone).
let body = format!("{t}{snap}\n{host} · {disp}");
v.push(Push { peer_sign: *sign, peer_dh: *dh, out: ServiceFrame::text(body).into() });
}
}
v
});
if let Err(e) = bot.serve_ticked(&a[0], &a[1], handler, tick, ticker).await {
eprintln!("fatal: {e}");
std::process::exit(1);
}
}
fn env_u64(key: &str, default: u64) -> u64 {
std::env::var(key).ok().and_then(|v| v.trim().parse().ok()).unwrap_or(default)
}
/// Parse `STATUSBOT_SERVICES` = "unit:Label,unit2:Label2" into (unit, label) pairs.
fn parse_services(spec: &str) -> Vec<(String, String)> {
spec.split(',')
.filter_map(|e| {
let e = e.trim();
if e.is_empty() {
return None;
}
let (u, l) = e.split_once(':').unwrap_or((e, e));
Some((u.trim().to_string(), l.trim().to_string()))
})
.collect()
}
/// Build the monitored-check list. `STATUSBOT_PRESET=jff` loads the EXACT JFF health checks the
/// Telegram monitor uses (5 HTTP + 2 TCP, each tied to its systemd unit for a restart button) — so
/// "unit active in systemd but HTTP 500/hung" is still caught. Plain `STATUSBOT_SERVICES` unit:label
/// pairs are added too, deduped against any unit a preset check already covers (never alert twice).
fn build_checks() -> Vec<Check> {
let mut checks: Vec<Check> = Vec::new();
if std::env::var("STATUSBOT_PRESET").map(|v| v.eq_ignore_ascii_case("jff")).unwrap_or(false) {
checks.extend(jff_preset());
}
// Custom HTTP checks: "key|Label|http://host:port/path|200,301|unit ; …"
checks.extend(parse_custom(&std::env::var("STATUSBOT_HTTP").unwrap_or_default(), true));
// Custom TCP checks: "key|Label|host:port|unit ; …"
checks.extend(parse_custom(&std::env::var("STATUSBOT_TCP").unwrap_or_default(), false));
for (unit, label) in parse_services(&std::env::var("STATUSBOT_SERVICES").unwrap_or_default()) {
if !checks.iter().any(|c| c.unit.as_deref() == Some(unit.as_str())) {
checks.push(Check { key: format!("svc:{unit}"), label, unit: Some(unit), probe: Probe::Unit });
}
}
// Dedup by alert key — two checks sharing a key would corrupt each other's AlertState.
let mut seen = std::collections::HashSet::new();
checks.retain(|c| {
if seen.insert(c.key.clone()) {
true
} else {
eprintln!("[statusbot] дубликат ключа проверки '{}' пропущен", c.key);
false
}
});
if checks.is_empty() {
// Demo default (no preset + nothing configured): watch whatever is commonly present.
checks.push(Check { key: "svc:ssh".into(), label: "SSH".into(), unit: Some("ssh".into()), probe: Probe::Unit });
checks.push(Check { key: "svc:cron".into(), label: "Cron".into(), unit: Some("cron".into()), probe: Probe::Unit });
}
checks
}
/// Parse a `;`-separated list of `|`-delimited custom checks. HTTP: `key|Label|url|okCsv|unit`
/// (okCsv defaults to `200`); TCP: `key|Label|host:port|unit`. A trailing `unit` is optional.
fn parse_custom(spec: &str, http: bool) -> Vec<Check> {
spec.split(';')
.filter_map(|e| {
let f: Vec<&str> = e.split('|').map(|x| x.trim()).collect();
if f.len() < 3 || f[0].is_empty() {
return None;
}
let (key, label) = (f[0].to_string(), f[1].to_string());
if http {
parse_http_url(f[2])?; // validate the URL shape up front
let ok: Vec<u16> = f.get(3).map(|s| s.split(',').filter_map(|c| c.trim().parse().ok()).collect()).unwrap_or_default();
let ok = if ok.is_empty() { vec![200] } else { ok };
let unit = f.get(4).filter(|u| !u.is_empty()).map(|u| u.to_string());
Some(Check { key, label, unit, probe: Probe::Http { url: f[2].to_string(), ok } })
} else {
// host:port, incl. bracketed IPv6 ([::1]:5432 → host "::1").
let (host, port) = if let Some(rest6) = f[2].strip_prefix('[') {
let (h, after) = rest6.split_once(']')?;
(h.to_string(), after.strip_prefix(':')?.parse::<u16>().ok()?)
} else {
let (h, p) = f[2].rsplit_once(':')?;
(h.to_string(), p.parse::<u16>().ok()?)
};
let unit = f.get(3).filter(|u| !u.is_empty()).map(|u| u.to_string());
Some(Check { key, label, unit, probe: Probe::Tcp { host, port } })
}
})
.collect()
}
/// The JFF Telegram-monitor's exact target set (jff-monitor/monitor.js http[]+tcp[]): real endpoint
/// health, each carrying its unit so the restart button + `status` line still work.
fn jff_preset() -> Vec<Check> {
let redir = vec![200u16, 301, 302, 307, 308];
let http = |key: &str, label: &str, url: &str, ok: Vec<u16>, unit: &str| Check {
key: key.into(), label: label.into(), unit: Some(unit.into()), probe: Probe::Http { url: url.into(), ok },
};
let tcp = |key: &str, label: &str, port: u16, unit: &str| Check {
key: key.into(), label: label.into(), unit: Some(unit.into()), probe: Probe::Tcp { host: "127.0.0.1".into(), port },
};
vec![
http("site", "Сайт (nginx :80)", "http://127.0.0.1/", redir.clone(), "nginx"),
http("api", "API /health :8000", "http://127.0.0.1:8000/health", vec![200], "jff-backend"),
http("admin_api", "Админ-API /health :8001", "http://127.0.0.1:8001/health", vec![200], "jff-admin-backend"),
http("client", "Клиент SSR :3000", "http://127.0.0.1:3000/", redir.clone(), "jff-client"),
http("admin_front", "Админ-фронт :3100", "http://127.0.0.1:3100/", redir, "jff-admin"),
tcp("postgres", "PostgreSQL :5432", 5432, "postgresql@16-main"),
tcp("redis", "Redis :6379", 6379, "redis-server"),
]
}
// ---- health probes (blocking, short timeouts; localhost targets answer in ms) ----
const HTTP_TIMEOUT: Duration = Duration::from_secs(3);
const TCP_TIMEOUT: Duration = Duration::from_secs(2);
/// Probe EVERY check concurrently (one scoped thread each) → total wall-time ≈ the slowest single
/// timeout, not their sum. Keeps the synchronous ticker (and the reactive /status) from stalling the
/// poll loop for ~19 s when many targets hang. Results are in `checks` order. NOTE: for IP-literal
/// targets `connect_timeout` bounds each thread; a custom check pointed at a HOSTNAME can still block
/// on DNS resolution (no timeout) — configure custom checks with IPs.
fn probe_all(checks: &[Check]) -> Vec<Result<(), String>> {
std::thread::scope(|s| {
let handles: Vec<_> = checks.iter().map(|c| s.spawn(|| probe_check(c))).collect();
handles.into_iter().map(|h| h.join().unwrap_or_else(|_| Err("probe panic".into()))).collect()
})
}
/// Probe one target → Ok if up, Err(short RU/EN reason) if down.
fn probe_check(c: &Check) -> Result<(), String> {
match &c.probe {
Probe::Unit => {
let st = is_active(c.unit.as_deref().unwrap_or(""));
if st == "active" { Ok(()) } else { Err(st) }
}
Probe::Http { url, ok } => http_get(url, ok),
Probe::Tcp { host, port } => tcp_ok(host, *port),
}
}
fn tcp_ok(host: &str, port: u16) -> Result<(), String> {
use std::net::ToSocketAddrs;
let mut addrs = (host, port).to_socket_addrs().map_err(|_| "bad addr".to_string())?;
match addrs.next() {
Some(a) if std::net::TcpStream::connect_timeout(&a, TCP_TIMEOUT).is_ok() => Ok(()),
_ => Err("нет соединения".into()),
}
}
/// Split `http://host[:port]/path` into (host, port, path). `None` if not a plain http URL.
fn parse_http_url(url: &str) -> Option<(String, u16, String)> {
let rest = url.strip_prefix("http://")?;
let (authority, path) = match rest.find('/') {
Some(i) => (&rest[..i], rest[i..].to_string()),
None => (rest, "/".to_string()),
};
if authority.is_empty() {
return None;
}
// Handle bracketed IPv6 literals: [::1]:8000 → host "::1", port 8000.
let (host, port) = if let Some(rest6) = authority.strip_prefix('[') {
let (h, after) = rest6.split_once(']')?;
let port = after.strip_prefix(':').and_then(|p| p.parse::<u16>().ok()).unwrap_or(80);
(h.to_string(), port)
} else {
match authority.rsplit_once(':') {
Some((h, p)) => (h.to_string(), p.parse::<u16>().unwrap_or(80)),
None => (authority.to_string(), 80),
}
};
Some((host, port, path))
}
/// Minimal blocking HTTP/1.0 GET for a localhost health URL (no TLS, no deps). Ok if the response
/// status is in `ok`, else a short reason. Reads only the status line.
fn http_get(url: &str, ok: &[u16]) -> Result<(), String> {
use std::io::{Read, Write};
use std::net::{TcpStream, ToSocketAddrs};
let (host, port, path) = parse_http_url(url).ok_or_else(|| "не http".to_string())?;
let addr = (host.as_str(), port).to_socket_addrs().map_err(|_| "bad addr".to_string())?.next().ok_or_else(|| "bad addr".to_string())?;
let mut s = TcpStream::connect_timeout(&addr, HTTP_TIMEOUT).map_err(|_| "нет соединения".to_string())?;
s.set_read_timeout(Some(HTTP_TIMEOUT)).ok();
s.set_write_timeout(Some(HTTP_TIMEOUT)).ok();
let req = format!("GET {path} HTTP/1.0\r\nHost: {host}\r\nUser-Agent: aster-statusbot\r\nConnection: close\r\n\r\n");
s.write_all(req.as_bytes()).map_err(|_| "write fail".to_string())?;
let mut buf = [0u8; 128];
let n = s.read(&mut buf).map_err(|_| "нет ответа".to_string())?;
if n == 0 {
return Err("пустой ответ".into());
}
let line = String::from_utf8_lossy(&buf[..n]);
let code = line.split_whitespace().nth(1).and_then(|c| c.parse::<u16>().ok()).ok_or_else(|| "bad response".to_string())?;
if ok.contains(&code) { Ok(()) } else { Err(format!("HTTP {code}")) }
}
// ---- frames ----
fn help_msg() -> Msg {
Msg {
text: Some("🖥 Статус-бот сервера\n\nКоманды:\n/status — нагрузка, RAM, диск, аптайм\n/services — сервисы + кнопки перезапуска\n/logs — последние 50 строк journalctl\n/help — эта справка\n\n🔔 Я сам пришлю алерт при падении сервиса (и «восстановлено» при подъёме) и суточные сводки — как только вы вошли владельцем.\n\nИли пользуйтесь кнопками внизу 👇".into()),
buttons: main_menu("help"),
reply_to: None,
client_msg_id: rand::random(),
}
}
fn help_frame() -> ServiceFrame {
ServiceFrame::Msg(help_msg())
}
fn msg(text: String, buttons: Vec<Vec<Button>>) -> ServiceFrame {
ServiceFrame::Msg(Msg { text: Some(text), buttons, reply_to: None, client_msg_id: rand::random() })
}
/// Persistent bottom keyboard appended to EVERY reply — the main functions + a Home button.
/// So the owner can always navigate without retyping a command (Telegram reply-keyboard feel).
/// `refresh` is the callback the "🔄 Обновить" button carries: it re-renders the CURRENT screen
/// (status→status, services→svc:list), so it's a real refresh and not a duplicate of Главная.
fn main_menu(refresh: &str) -> Vec<Vec<Button>> {
vec![
vec![
Button { text: "🏠 Главная".into(), data: b"home".to_vec() },
Button { text: "🔧 Сервисы".into(), data: b"svc:list".to_vec() },
],
vec![
Button { text: "📋 Логи".into(), data: b"logs:list".to_vec() },
Button { text: "🔄 Обновить".into(), data: refresh.as_bytes().to_vec() },
],
vec![
Button { text: "❓ Помощь".into(), data: b"help".to_vec() },
],
]
}
/// A text reply that ALWAYS carries the persistent bottom keyboard. A one-off result has no
/// "current screen", so its refresh lands on the live status dashboard.
fn menu(text: impl Into<String>) -> ServiceFrame {
msg(text.into(), main_menu("status"))
}
/// Truncate a label for a button so a long service name doesn't wrap the tap target.
fn short(s: &str, max: usize) -> String {
if s.chars().count() <= max {
s.to_string()
} else {
let t: String = s.chars().take(max.saturating_sub(1)).collect();
format!("{t}…")
}
}
fn status_msg(cfg: &Cfg) -> Msg {
let (l1, l5, l15) = loadavg();
let (ram_used, ram_pct) = mem();
let (disk_pct, disk_avail) = disk_root();
let up = uptime_h();
let results = probe_all(&cfg.checks);
let ok = results.iter().filter(|r| r.is_ok()).count();
let dot = if ok == cfg.checks.len() { "🟢" } else { "🔴" };
let mut text = format!(
"🖥 Статус сервера\n\n⚙️ Load: {l1:.2} · {l5:.2} · {l15:.2}\n🧠 RAM: {ram_pct}% ({ram_used})\n💾 Диск /: {disk_pct}% занято (свободно {disk_avail})\n⏱ Аптайм: {up}\n{dot} Сервисы: {ok}/{}",
cfg.checks.len()
);
for (c, r) in cfg.checks.iter().zip(&results) {
if let Err(reason) = r {
text.push_str(&format!("\n 🔴 {} — {reason}", c.label));
}
}
Msg { text: Some(text), buttons: main_menu("status"), reply_to: None, client_msg_id: rand::random() }
}
fn status_frame(cfg: &Cfg) -> ServiceFrame {
ServiceFrame::Msg(status_msg(cfg))
}
fn services_msg(cfg: &Cfg) -> Msg {
services_msg_note(cfg, None)
}
/// Services screen with an optional one-line note (restart result / read-only notice) on top —
/// persistent in-place feedback without spawning a new message.
fn services_msg_note(cfg: &Cfg, note: Option<&str>) -> Msg {
let mut lines = match note {
Some(n) => format!("{n}\n\n🔧 Сервисы:\n"),
None => String::from("🔧 Сервисы:\n"),
};
let mut rows: Vec<Vec<Button>> = Vec::new();
for (c, res) in cfg.checks.iter().zip(probe_all(&cfg.checks)) {
let (dot, st) = match res {
Ok(()) => ("🟢", "ок".to_string()),
Err(reason) => ("🔴", reason),
};
lines.push_str(&format!("\n{dot} {} — {st}", c.label));
if cfg.allow_restart {
if let Some(unit) = &c.unit {
rows.push(vec![Button { text: format!("🔄 Перезапустить {}", short(&c.label, 16)), data: format!("svc:ask:{unit}").into_bytes() }]);
}
}
}
if !cfg.allow_restart {
lines.push_str("\n\nℹ️ Режим только для чтения (перезапуск отключён).");
}
rows.extend(main_menu("svc:list"));
Msg { text: Some(lines), buttons: rows, reply_to: None, client_msg_id: rand::random() }
}
fn services_frame(cfg: &Cfg) -> ServiceFrame {
ServiceFrame::Msg(services_msg(cfg))
}
const LOG_DEFAULT_LINES: u32 = 50;
const LOG_MAX_LINES: u32 = 200;
const LOG_MAX_CHARS: usize = 4000;
/// Last N lines of journalctl for a systemd unit. Returns the text (may be truncated).
fn get_logs(unit: &str, lines: u32) -> String {
let lines = lines.min(LOG_MAX_LINES);
let out = Command::new("journalctl")
.args(["-u", unit, "-n", &lines.to_string(), "--no-pager", "-o", "short-iso"])
.output();
match out {
Ok(o) if o.status.success() => {
let text = String::from_utf8_lossy(&o.stdout).to_string();
if text.len() > LOG_MAX_CHARS {
let truncated: String = text.chars().take(LOG_MAX_CHARS).collect();
format!("{truncated}\n\n… (показано последние {lines} строк, обрезано до {LOG_MAX_CHARS} символов)")
} else if text.trim().is_empty() {
format!("📋 Логи {unit}: пусто (нет записей)")
} else {
text
}
}
Ok(o) => {
let stderr = String::from_utf8_lossy(&o.stderr);
format!("❌ Ошибка чтения логов {unit}: {stderr}")
}
Err(e) => format!("❌ Не удалось запустить journalctl: {e}"),
}
}
/// Show the list of services with log-view buttons.
fn logs_msg(cfg: &Cfg) -> Msg {
let mut rows: Vec<Vec<Button>> = Vec::new();
for c in &cfg.checks {
if c.unit.is_some() {
rows.push(vec![Button {
text: format!("📋 {}", short(&c.label, 20)),
data: format!("logs:show:{}", c.unit.as_deref().unwrap()).into_bytes(),
}]);
}
}
if rows.is_empty() {
return Msg {
text: Some("📋 Нет сервисов с systemd-юнитами для просмотра логов.".into()),
buttons: main_menu("logs:list"),
reply_to: None,
client_msg_id: rand::random(),
};
}
rows.extend(main_menu("logs:list"));
Msg { text: Some("📋 Выберите сервис для просмотра логов (последние 50 строк):".into()), buttons: rows, reply_to: None, client_msg_id: rand::random() }
}
fn logs_frame(cfg: &Cfg) -> ServiceFrame {
ServiceFrame::Msg(logs_msg(cfg))
}
/// Show last N lines of journalctl for a specific unit (as an editable Msg body).
/// `None` if the unit is not one we monitor.
fn show_logs_msg(cfg: &Cfg, unit: &str) -> Option<Msg> {
// Verify the unit is one we monitor (never probe an arbitrary attacker-supplied unit).
if !cfg.checks.iter().any(|c| c.unit.as_deref() == Some(unit)) {
return None;
}
let label = label_of(cfg, unit);
let log_text = get_logs(unit, LOG_DEFAULT_LINES);
let header = format!("📋 Логи «{label}» (последние {LOG_DEFAULT_LINES} строк):\n\n");
let full = format!("{header}{log_text}");
// Truncate if still too long after the header
let text = if full.len() > LOG_MAX_CHARS {
let truncated: String = full.chars().take(LOG_MAX_CHARS).collect();
format!("{truncated}\n\n… (обрезано)")
} else {
full
};
let mut rows = main_menu("logs:list");
// Add a "🔄 Обновить" button specific to this unit (its callback edits this message again)
rows.insert(0, vec![Button {
text: "🔄 Обновить логи".into(),
data: format!("logs:show:{unit}").into_bytes(),
}]);
Some(Msg { text: Some(text), buttons: rows, reply_to: None, client_msg_id: rand::random() })
}
/// Restart confirmation screen (replaces the services message in-place; its own buttons then
/// edit THIS message — svc:do / svc:cancel carry this Msg's client_msg_id back).
fn confirm_restart_msg(unit: &str, label: &str) -> Msg {
let mut buttons = vec![vec![
Button { text: "✅ Да, перезапустить".into(), data: format!("svc:do:{unit}").into_bytes() },
Button { text: "❌ Отмена".into(), data: b"svc:cancel".to_vec() },
]];
buttons.extend(main_menu("svc:list"));
Msg {
text: Some(format!("Перезапустить «{label}»?")),
buttons,
reply_to: None,
client_msg_id: rand::random(),
}
}
/// Every button answer EDITS the tapped message in-place: the whole bot UI lives in ONE
/// dynamically re-rendered bubble and nothing new is appended to the chat.
fn handle_callback(cfg: &Cfg, restarts: &mut HashMap<String, Vec<u64>>, data: &str, msg_ref: [u8; 16]) -> Vec<Outgoing> {
let edit = |m: Msg| -> Vec<Outgoing> { vec![ServiceFrame::Edit { msg_ref, msg: m }.into()] };
if data == "home" || data == "status" {
return edit(status_msg(cfg));
}
if data == "svc:list" {
return edit(services_msg(cfg));
}
if data == "logs:list" {
return edit(logs_msg(cfg));
}
if data == "help" {
return edit(help_msg());
}
if let Some(unit) = data.strip_prefix("logs:show:") {
// Unknown unit (can't normally happen — buttons are bot-generated): fall back to services.
return match show_logs_msg(cfg, unit) {
Some(m) => edit(m),
None => edit(services_msg_note(cfg, Some("❌ Неизвестный сервис."))),
};
}
// Restart is a privileged action — only when explicitly enabled at launch.
if (data.starts_with("svc:ask:") || data.starts_with("svc:do:")) && !cfg.allow_restart {
return edit(services_msg_note(cfg, Some("ℹ️ Перезапуск отключён (read-only). Запусти бота с STATUSBOT_ALLOW_RESTART=1.")));
}
if let Some(unit) = data.strip_prefix("svc:ask:") {
return edit(confirm_restart_msg(unit, &label_of(cfg, unit)));
}
if data == "svc:cancel" {
return edit(services_msg(cfg)); // back to the fresh list, no extra chatter
}
if let Some(unit) = data.strip_prefix("svc:do:") {
// unit must be one we manage (never restart an arbitrary attacker-supplied unit).
if !cfg.checks.iter().any(|c| c.unit.as_deref() == Some(unit)) {
return edit(services_msg_note(cfg, Some("❌ Неизвестный сервис.")));
}
let label = label_of(cfg, unit);
let now = now_secs();
let hist = restarts.entry(unit.to_string()).or_default();
hist.retain(|t| now - t < 3600);
if hist.len() >= RESTART_CAP_PER_HOUR {
return edit(services_msg_note(cfg, Some(&format!(
"⚠️ {label}: достигнут лимит {RESTART_CAP_PER_HOUR} перезапусков за час. Проверьте вручную: journalctl -u {unit} -n 100"
))));
}
hist.push(now);
let out = Command::new("sudo").args(["-n", "systemctl", "restart", unit]).output();
let after = is_active(unit);
let ok = matches!(&out, Ok(o) if o.status.success()) && after == "active";
let note = format!(
"{} {label}: systemctl restart → сейчас {after}",
if ok { "✅" } else { "❌" }
);
return edit(services_msg_note(cfg, Some(¬e)));
}
// Unknown callback — recover by rendering the live status dashboard into the same bubble.
edit(status_msg(cfg))
}
fn label_of(cfg: &Cfg, unit: &str) -> String {
cfg.checks.iter().find(|c| c.unit.as_deref() == Some(unit)).map(|c| c.label.clone()).unwrap_or_else(|| unit.to_string())
}
// ---- metrics (Linux /proc + coreutils) ----
fn now_secs() -> u64 {
SystemTime::now().duration_since(UNIX_EPOCH).map(|d| d.as_secs()).unwrap_or(0)
}
fn loadavg() -> (f64, f64, f64) {
let s = std::fs::read_to_string("/proc/loadavg").unwrap_or_default();
let mut it = s.split_whitespace();
let p = |x: Option<&str>| x.and_then(|v| v.parse().ok()).unwrap_or(0.0);
(p(it.next()), p(it.next()), p(it.next()))
}
fn mem() -> (String, u64) {
let s = std::fs::read_to_string("/proc/meminfo").unwrap_or_default();
let get = |k: &str| -> u64 {
s.lines()
.find(|l| l.starts_with(k))
.and_then(|l| l.split_whitespace().nth(1))
.and_then(|v| v.parse::<u64>().ok())
.unwrap_or(0)
};
let total = get("MemTotal:"); // kB
let avail = get("MemAvailable:");
let used = total.saturating_sub(avail);
let pct = if total > 0 { used * 100 / total } else { 0 };
(format!("{} / {}", human_kb(used), human_kb(total)), pct)
}
fn human_kb(kb: u64) -> String {
let mb = kb / 1024;
if mb >= 1024 {
format!("{:.1} ГБ", mb as f64 / 1024.0)
} else {
format!("{mb} МБ")
}
}
fn disk_root() -> (u64, String) {
let out = Command::new("df").args(["-P", "/"]).output();
if let Ok(o) = out {
let t = String::from_utf8_lossy(&o.stdout);
if let Some(l) = t.lines().nth(1) {
let f: Vec<&str> = l.split_whitespace().collect();
if f.len() >= 5 {
let pct = f[4].trim_end_matches('%').parse::<u64>().unwrap_or(0);
let avail_kb = f[3].parse::<u64>().unwrap_or(0);
return (pct, human_kb(avail_kb));
}
}
}
(0, "?".into())
}
/// (total, idle) тики из /proc/stat — для расчёта % между тиками.
fn cpu_sample() -> (u64, u64) {
let Some(line) = std::fs::read_to_string("/proc/stat").ok().and_then(|s| s.lines().next().map(str::to_string)) else { return (0, 0) };
let f: Vec<u64> = line.split_whitespace().skip(1).filter_map(|v| v.parse().ok()).collect();
if f.len() < 4 { return (0, 0) }
let idle = f[3] + f.get(4).copied().unwrap_or(0); // idle + iowait
(f.iter().sum(), idle)
}
fn uptime_h() -> String {
let s = std::fs::read_to_string("/proc/uptime").unwrap_or_default();
let secs: u64 = s.split_whitespace().next().and_then(|v| v.parse::<f64>().ok()).unwrap_or(0.0) as u64;
let (d, h, m) = (secs / 86400, (secs % 86400) / 3600, (secs % 3600) / 60);
if d > 0 {
format!("{d}д {h}ч {m}м")
} else {
format!("{h}ч {m}м")
}
}
fn is_active(unit: &str) -> String {
Command::new("systemctl")
.args(["is-active", unit])
.output()
.map(|o| String::from_utf8_lossy(&o.stdout).trim().to_string())
.unwrap_or_else(|_| "unknown".into())
}
// ---- proactive monitor: alert edge-detection + daily digest (mirrors the JFF Telegram monitor) ----
fn hostname() -> String {
std::fs::read_to_string("/etc/hostname")
.ok()
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.or_else(|| {
Command::new("uname").arg("-n").output().ok().map(|o| String::from_utf8_lossy(&o.stdout).trim().to_string())
})
.filter(|s| !s.is_empty())
.unwrap_or_else(|| "host".into())
}
/// (local hour 0..23, "YYYY-MM-DD" day-key, "DD.MM HH:MM" display). std has no timezone, so we shell
/// out to `date` (like the existing systemctl/df probes). Hour is `u8::MAX` if `date` is unavailable —
/// then the digest simply never fires (it's gated on `hour != u8::MAX`).
fn local_time() -> (u8, String, String) {
if let Ok(o) = Command::new("date").arg("+%H|%F|%d.%m %H:%M").output() {
let s = String::from_utf8_lossy(&o.stdout);
let parts: Vec<&str> = s.trim().split('|').collect();
if parts.len() == 3 {
if let Ok(h) = parts[0].parse::<u8>() {
if h < 24 {
return (h, parts[1].to_string(), parts[2].to_string());
}
}
}
}
(u8::MAX, String::new(), "?".to_string())
}
struct AlertState {
active: bool,
since: Instant,
last_notified: Instant,
}
/// Owns the alert state + digest scheduling. Lives solely inside the ticker (single-threaded on the
/// poll task) → no locking. All alert timing uses monotonic `Instant`, so it is immune to wall-clock
/// steps; only the once-a-day digest touches wall-clock, bounded to one fire per local day-key.
struct Monitor {
alerts: HashMap<String, AlertState>,
realert: Duration,
digest_hours: Vec<u8>,
last_digest_fired: HashSet<(u8, String)>,
load1_max: Option<f64>,
mem_max: Option<u64>,
disk_max: Option<u64>,
cpu_max: Option<u64>,
prev_cpu: Option<(u64, u64)>,
}
impl Monitor {
fn new(realert: Duration, digest_hours: Vec<u8>, load1_max: Option<f64>, mem_max: Option<u64>, disk_max: Option<u64>, cpu_max: Option<u64>) -> Self {
// Startup dedup-seed: if we boot AT/AFTER a digest hour today, mark it as already sent so a
// restart during/after the digest window doesn't re-emit today's сводка for that hour.
let mut last_digest_fired = HashSet::new();
let (hour, day, _) = local_time();
if hour != u8::MAX && !day.is_empty() {
for &h in &digest_hours {
if hour >= h {
last_digest_fired.insert((h, day.clone()));
}
}
}
Monitor { alerts: HashMap::new(), realert, digest_hours, last_digest_fired, load1_max, mem_max, disk_max, cpu_max, prev_cpu: None }
}
/// Alert edge-detector (mirrors monitor.js raiseAlert/clearAlert). Returns a message exactly on:
/// OK→fail (fresh 🔴), still-failing after `realert` (reminder), fail→OK (recovery ✅); otherwise
/// None. `detail` describes the current failing condition. Pure `Instant` math → unit-testable.
fn transition(&mut self, key: &str, label: &str, failing: bool, detail: &str, now: Instant) -> Option<String> {
match self.alerts.get_mut(key) {
Some(a) if a.active => {
if failing {
if now.duration_since(a.last_notified) >= self.realert {
a.last_notified = now;
let mins = now.duration_since(a.since).as_secs() / 60;
Some(format!("🔴 {label} — всё ещё ({detail}, ~{mins} мин)"))
} else {
None
}
} else {
a.active = false;
let mins = now.duration_since(a.since).as_secs() / 60;
Some(format!("✅ {label} — восстановлено (было ~{mins} мин)"))
}
}
_ => {
if failing {
self.alerts.insert(key.to_string(), AlertState { active: true, since: now, last_notified: now });
Some(format!("🔴 {label} — {detail}"))
} else {
None
}
}
}
}
/// Once per (hour, day-key) at the configured digest hours: a "всё в норме / N алертов" summary
/// (сводка). Deduped by `last_digest_fired` so each hour fires once even though the ticker runs
/// every few seconds.
fn maybe_digest(&mut self, hour: u8, day: &str, l1: f64, ram: u64, disk: u64) -> Option<String> {
if !self.digest_hours.contains(&hour) || day.is_empty() || self.last_digest_fired.contains(&(hour, day.to_string())) {
return None;
}
self.last_digest_fired.insert((hour, day.to_string()));
let n = self.alerts.values().filter(|a| a.active).count();
let head = if n == 0 { "✅ Все системы в норме".to_string() } else { format!("⚠️ Активных алертов: {n}") };
Some(format!("📊 Сводка ({hour}:00)\n{head}\nload {l1:.2} · RAM {ram}% · диск {disk}%"))
}
}
#[cfg(test)]
mod tests {
use super::*;
fn mon(realert_secs: u64, digest_hours: Vec<u8>) -> Monitor {
Monitor {
alerts: HashMap::new(),
realert: Duration::from_secs(realert_secs),
digest_hours,
last_digest_fired: HashSet::new(),
load1_max: None,
mem_max: None,
disk_max: None,
cpu_max: None,
prev_cpu: None,
}
}
#[test]
fn alert_state_machine() {
let mut m = mon(60, vec![]);
let t0 = Instant::now();
// OK→fail: exactly one 🔴 with the detail.
let a = m.transition("svc:x", "Сервис X", true, "недоступен", t0);
assert!(a.as_deref().is_some_and(|s| s.contains("🔴") && s.contains("недоступен")), "{a:?}");
// still failing before realert → nothing.
assert!(m.transition("svc:x", "Сервис X", true, "недоступен", t0 + Duration::from_secs(30)).is_none());
// at realert → one reminder.
let r = m.transition("svc:x", "Сервис X", true, "недоступен", t0 + Duration::from_secs(60));
assert!(r.as_deref().is_some_and(|s| s.contains("всё ещё")), "{r:?}");
// fail→OK → one recovery ✅ with duration.
let ok = m.transition("svc:x", "Сервис X", false, "", t0 + Duration::from_secs(90));
assert!(ok.as_deref().is_some_and(|s| s.contains("✅") && s.contains("восстановлено")), "{ok:?}");
// recovery on an already-cleared key → nothing.
assert!(m.transition("svc:x", "Сервис X", false, "", t0 + Duration::from_secs(120)).is_none());
// clear on a never-seen key → nothing.
assert!(m.transition("svc:y", "Y", false, "", t0).is_none());
}
#[test]
fn digest_scheduling_and_dedup() {
let mut m = mon(60, vec![14]);
// fires once at the hour...
let d = m.maybe_digest(14, "2026-08-18", 0.42, 20, 30);
assert!(d.as_deref().is_some_and(|s| s.contains("Сводка (14:00)") && s.contains("Все системы в норме") && s.contains("load 0.42")), "{d:?}");
// ...then dedups the same (hour, day)...
assert!(m.maybe_digest(14, "2026-08-18", 0.42, 20, 30).is_none());
// ...never off-hour...
assert!(m.maybe_digest(15, "2026-08-18", 0.42, 20, 30).is_none());
// ...but fires again the next day.
assert!(m.maybe_digest(14, "2026-08-19", 0.42, 20, 30).is_some());
// active alerts change the head line.
m.alerts.insert("svc:x".into(), AlertState { active: true, since: Instant::now(), last_notified: Instant::now() });
let d2 = m.maybe_digest(14, "2026-08-20", 0.42, 20, 30);
assert!(d2.as_deref().is_some_and(|s| s.contains("Активных алертов: 1")), "{d2:?}");
// digest disabled → never fires.
let mut off = mon(60, vec![]);
assert!(off.maybe_digest(14, "2026-08-18", 0.42, 20, 30).is_none());
}
#[test]
fn multi_hour_digest_fires_at_each_hour() {
let mut m = mon(60, vec![10, 14, 18, 21]);
// Each hour fires once.
assert!(m.maybe_digest(10, "2026-08-21", 0.5, 30, 40).is_some());
assert!(m.maybe_digest(14, "2026-08-21", 0.5, 30, 40).is_some());
assert!(m.maybe_digest(18, "2026-08-21", 0.5, 30, 40).is_some());
assert!(m.maybe_digest(21, "2026-08-21", 0.5, 30, 40).is_some());
// Second call for same (hour, day) → deduped.
assert!(m.maybe_digest(10, "2026-08-21", 0.5, 30, 40).is_none());
assert!(m.maybe_digest(14, "2026-08-21", 0.5, 30, 40).is_none());
// Different day → fires again.
assert!(m.maybe_digest(10, "2026-08-22", 0.5, 30, 40).is_some());
assert!(m.maybe_digest(14, "2026-08-22", 0.5, 30, 40).is_some());
// Not in the configured set → never fires.
assert!(m.maybe_digest(2, "2026-08-21", 0.5, 30, 40).is_none());
assert!(m.maybe_digest(23, "2026-08-21", 0.5, 30, 40).is_none());
}
#[test]
fn jff_preset_mirrors_the_telegram_targets() {
let p = jff_preset();
assert_eq!(p.len(), 7, "5 HTTP + 2 TCP");
let keys: Vec<&str> = p.iter().map(|c| c.key.as_str()).collect();
for k in ["site", "api", "admin_api", "client", "admin_front", "postgres", "redis"] {
assert!(keys.contains(&k), "missing {k}");
}
// every check ties to a systemd unit (for the restart button).
assert!(p.iter().all(|c| c.unit.is_some()));
// api → HTTP :8000/health, ok=[200], unit jff-backend.
let api = p.iter().find(|c| c.key == "api").unwrap();
assert_eq!(api.unit.as_deref(), Some("jff-backend"));
assert!(matches!(&api.probe, Probe::Http { url, ok } if url.contains(":8000/health") && ok == &[200]));
// postgres → TCP 5432.
let pg = p.iter().find(|c| c.key == "postgres").unwrap();
assert!(matches!(&pg.probe, Probe::Tcp { port, .. } if *port == 5432));
}
#[test]
fn custom_check_parsing() {
let http = parse_custom("web|Веб|http://127.0.0.1:8080/health|200,204|myunit", true);
assert_eq!(http.len(), 1);
assert_eq!(http[0].key, "web");
assert_eq!(http[0].unit.as_deref(), Some("myunit"));
assert!(matches!(&http[0].probe, Probe::Http { url, ok } if url.ends_with("/health") && ok == &[200, 204]));
// ok defaults to [200]; unit optional
let h2 = parse_custom("x|X|http://127.0.0.1/", true);
assert!(matches!(&h2[0].probe, Probe::Http { ok, .. } if ok == &[200]) && h2[0].unit.is_none());
let tcp = parse_custom("db|БД|127.0.0.1:5432|pg ; cache|Кэш|127.0.0.1:6379", false);
assert_eq!(tcp.len(), 2);
assert!(matches!(&tcp[0].probe, Probe::Tcp { host, port } if host == "127.0.0.1" && *port == 5432));
assert_eq!(tcp[0].unit.as_deref(), Some("pg"));
assert!(tcp[1].unit.is_none());
// bracketed IPv6 host
let v6 = parse_custom("db6|БД6|[::1]:5432|pg", false);
assert!(matches!(&v6[0].probe, Probe::Tcp { host, port } if host == "::1" && *port == 5432));
// malformed entries are dropped
assert!(parse_custom("bad", false).is_empty());
assert!(parse_custom("k|L|nohost", false).is_empty()); // no :port
assert!(parse_custom("k|L|https://x/|200|u", true).is_empty()); // non-http url rejected
}
#[test]
fn http_url_parsing() {
assert_eq!(parse_http_url("http://127.0.0.1/"), Some(("127.0.0.1".into(), 80, "/".into())));
assert_eq!(parse_http_url("http://127.0.0.1:8000/health"), Some(("127.0.0.1".into(), 8000, "/health".into())));
assert_eq!(parse_http_url("http://127.0.0.1:3100/"), Some(("127.0.0.1".into(), 3100, "/".into())));
assert_eq!(parse_http_url("http://host"), Some(("host".into(), 80, "/".into())));
// bracketed IPv6 literal
assert_eq!(parse_http_url("http://[::1]:8000/health"), Some(("::1".into(), 8000, "/health".into())));
assert_eq!(parse_http_url("http://[::1]/"), Some(("::1".into(), 80, "/".into())));
assert!(parse_http_url("https://x/").is_none()); // TLS not supported by the minimal probe
assert!(parse_http_url("ftp://x").is_none());
}
}
catalogbot — каталог ботов — Каталог ботов сообщества с button-driven UX (без команд). Пошаговый wizard добавления, голосование (👍/👎), пагинация, «мои боты», удаление. Вся навигация — инлайн-кнопки, экраны редактируются на месте (`ServiceFrame::Edit`).
crates/aster-services/src/bin/catalogbot.rs
//! catalogbot — a bot DIRECTORY with a BUTTON-DRIVEN, command-free UX. The whole directory lives in
//! ONE evolving "app card" per user that morphs in place via `ServiceFrame::Edit`: Home ⇄ Catalog
//! (paginated) ⇄ Bot-detail (vote / remove) ⇄ Add-wizard ⇄ My-bots. The user NEVER types a slash —
//! navigation is entirely inline buttons; the only free text is a guided add-a-bot wizard.
//!
//! HOW IT STAYS ROBUST (two verified client facts, web/app.js):
//! • Edit finds the target by `msg_ref == client_msg_id` and KEEPS the original cmid, so the card's
//! identity never drifts across Edits.
//! • every callback hands the bot `msg_ref` = the exact message the finger is on.
//! So the rule is "adopt-on-tap": on any button tap, Edit the tapped message to the target screen.
//! Callback data is fully SELF-DESCRIBING (encodes the target state + a back-origin), so a tap on any
//! old/scrolled/post-restart card is still correct with ~zero navigation state. The only deviation is
//! right after the user TYPES a wizard field (their bubble is now at the bottom and an Edit to the
//! scrolled-up prompt could silently no-op) — a typed step RE-ANCHORS: collapse the old prompt to a
//! breadcrumb and send the next step as a fresh bottom message.
//!
//! STORAGE: a third-party bot MAY keep persistent memory (RAM-first is for the official relay/infra,
//! not bots). With `CATALOGBOT_STATE` set, the FULL catalog + vote map + submitters persist to JSON;
//! per-user UI state (the current card / wizard) is RAM-only (fine to drop on restart).
//!
//! Usage: printf '<128-hex sign||dh>' | catalogbot <server_addr> <server_pubhex>
//! Env: CATALOGBOT_STATE=/path/catalog.json · CATALOGBOT_ADMIN=ASTER1… (may remove any entry) · ASTER_TLS=1
use std::collections::HashMap;
use std::io::Read;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use aster_core::identity::{account_address, parse_account_address};
use aster_net::{from_hex, to_hex};
use aster_services::{Bot, Button, Handler, Incoming, Msg, Outgoing, ServiceFrame, Update};
use serde::{Deserialize, Serialize};
const MAX_ENTRIES: usize = 500;
const MAX_PER_USER: usize = 10;
const NAME_MAX: usize = 60;
const DESC_MAX: usize = 200;
const PAGE_SIZE: usize = 6;
const WIZARD_TTL: Duration = Duration::from_secs(15 * 60);
// ───────────────────────── persisted catalog (unchanged) ─────────────────────────
#[derive(Serialize, Deserialize, Clone)]
struct Entry {
id: u32,
addr: String,
name: String,
desc: String,
up: u32,
down: u32,
}
#[derive(Serialize, Deserialize, Default)]
struct Catalog {
next_id: u32,
entries: Vec<Entry>,
}
/// On-disk snapshot (JSON): the full state. `[u8;32]` keys are hex strings so JSON can hold them.
#[derive(Serialize, Deserialize, Default)]
struct Snapshot {
next_id: u32,
entries: Vec<Entry>,
votes: Vec<(u32, String, i8)>,
submitters: Vec<(u32, String)>,
}
// ───────────────────────── per-user UI state (RAM-only) ─────────────────────────
/// Which field the add-wizard is currently asking for. `Confirm` awaits buttons only.
#[derive(Clone, Default, PartialEq, Debug)]
enum Field {
#[default]
Addr,
Name,
Desc,
Confirm,
}
/// The slash-free add-a-bot wizard: collected fields + what we're asking now. `rtc` (return-to-confirm)
/// is set when a single field is being re-edited from the Confirm screen, so the next value jumps
/// straight back to Confirm.
#[derive(Clone, Default)]
struct Wizard {
addr: Option<String>,
name: Option<String>,
desc: Option<String>,
asking: Field,
rtc: bool,
}
impl Wizard {
/// Fill defaults so a Confirm screen always has all three fields.
fn fill(&mut self) {
self.name.get_or_insert_with(|| "Без названия".to_string());
self.desc.get_or_insert_with(String::new);
}
fn awaits_text(&self) -> bool {
matches!(self.asking, Field::Addr | Field::Name | Field::Desc)
}
}
struct UserUi {
app_cmid: Option<[u8; 16]>,
wizard: Option<Wizard>,
touched: Instant,
}
struct State {
catalog: Catalog,
path: Option<String>,
admin: Option<String>,
votes: HashMap<u32, HashMap<[u8; 32], i8>>,
submitters: HashMap<u32, [u8; 32]>,
uis: HashMap<[u8; 32], UserUi>, // RAM-only per-user UI (never persisted)
}
fn sign32(hexstr: &str) -> Option<[u8; 32]> {
from_hex(hexstr).and_then(|v| <[u8; 32]>::try_from(v.as_slice()).ok())
}
impl State {
fn load(path: Option<String>, admin: Option<String>) -> Self {
let snap: Snapshot = path
.as_ref()
.and_then(|p| std::fs::read(p).ok())
.and_then(|b| serde_json::from_slice::<Snapshot>(&b).ok())
.unwrap_or_default();
let mut votes: HashMap<u32, HashMap<[u8; 32], i8>> = HashMap::new();
for (id, voter, dir) in snap.votes {
if let Some(v) = sign32(&voter) {
votes.entry(id).or_default().insert(v, dir);
}
}
let submitters = snap.submitters.into_iter().filter_map(|(id, s)| sign32(&s).map(|k| (id, k))).collect();
State {
catalog: Catalog { next_id: snap.next_id, entries: snap.entries },
path,
admin,
votes,
submitters,
uis: HashMap::new(),
}
}
fn save(&self) {
let Some(p) = &self.path else { return };
let snap = Snapshot {
next_id: self.catalog.next_id,
entries: self.catalog.entries.clone(),
votes: self.votes.iter().flat_map(|(id, m)| m.iter().map(move |(v, d)| (*id, to_hex(v), *d))).collect(),
submitters: self.submitters.iter().map(|(id, s)| (*id, to_hex(s))).collect(),
};
if let Ok(bytes) = serde_json::to_vec_pretty(&snap) {
let _ = std::fs::write(p, bytes);
}
}
fn entry(&self, id: u32) -> Option<&Entry> {
self.catalog.entries.iter().find(|e| e.id == id)
}
fn entry_by_addr(&self, addr: &str) -> Option<&Entry> {
self.catalog.entries.iter().find(|e| e.addr == addr)
}
fn ui(&mut self, k: [u8; 32]) -> &mut UserUi {
self.uis.entry(k).or_insert_with(|| UserUi { app_cmid: None, wizard: None, touched: Instant::now() })
}
/// Validate + insert a submission from ALREADY-SEPARATE fields (the wizard path). No delimiter is
/// involved, so a name/description may safely contain any character (incl. '|') and is stored
/// exactly as the user confirmed. Returns the new id or a user-facing error.
fn add_fields(&mut self, by: [u8; 32], addr_in: &str, name_in: &str, desc_in: &str) -> Result<u32, String> {
let sign = parse_account_address(addr_in.trim()).ok_or("Неверный адрес (должен начинаться с ASTER1).")?;
let addr = account_address(&sign);
let name = clean(name_in.trim(), NAME_MAX);
let desc = clean(desc_in.trim(), DESC_MAX);
if self.catalog.entries.iter().any(|e| e.addr == addr) {
return Err("Этот бот уже в каталоге.".into());
}
if self.catalog.entries.len() >= MAX_ENTRIES {
return Err("Каталог заполнен.".into());
}
if self.submitters.values().filter(|s| **s == by).count() >= MAX_PER_USER {
return Err(format!("Лимит {MAX_PER_USER} ботов на пользователя."));
}
let name = if name.is_empty() { "Без названия".to_string() } else { name };
let id = self.catalog.next_id;
self.catalog.next_id += 1;
self.catalog.entries.push(Entry { id, addr, name, desc, up: 0, down: 0 });
self.submitters.insert(id, by);
self.save();
Ok(id)
}
/// Thin "addr | name | desc" parser over [`add_fields`] (kept for tests / a plain-text caller).
/// NOTE: the wizard NEVER uses this — it calls `add_fields` directly so a '|' in the name can't
/// bleed into the description.
fn add(&mut self, by: [u8; 32], args: &str) -> Result<u32, String> {
let mut parts = args.splitn(3, '|');
let addr = parts.next().unwrap_or("");
let name = parts.next().unwrap_or("");
let desc = parts.next().unwrap_or("");
self.add_fields(by, addr, name, desc)
}
/// Apply a vote (+1/-1). Tapping the same direction again clears the vote. Returns the voter's
/// resulting vote or None if the entry is gone.
fn vote(&mut self, id: u32, voter: [u8; 32], dir: i8) -> Option<i8> {
let e = self.catalog.entries.iter_mut().find(|e| e.id == id)?;
let vm = self.votes.entry(id).or_default();
let prev = vm.get(&voter).copied().unwrap_or(0);
if prev == 1 {
e.up = e.up.saturating_sub(1);
} else if prev == -1 {
e.down = e.down.saturating_sub(1);
}
let newv = if prev == dir { 0 } else { dir };
match newv {
1 => {
e.up += 1;
vm.insert(voter, 1);
}
-1 => {
e.down += 1;
vm.insert(voter, -1);
}
_ => {
vm.remove(&voter);
}
}
self.save();
Some(newv)
}
fn my_vote(&self, id: u32, voter: &[u8; 32]) -> i8 {
self.votes.get(&id).and_then(|m| m.get(voter)).copied().unwrap_or(0)
}
fn is_owner(&self, id: u32, by: &[u8; 32], by_addr: &str) -> bool {
self.submitters.get(&id).map(|s| s == by).unwrap_or(false) || self.admin.as_deref() == Some(by_addr)
}
fn remove(&mut self, id: u32, by: [u8; 32], by_addr: &str) -> Result<String, String> {
if !self.is_owner(id, &by, by_addr) {
return Err("Удалять может только автор бота или админ каталога.".into());
}
let pos = self.catalog.entries.iter().position(|e| e.id == id).ok_or("Нет такого бота.")?;
let name = self.catalog.entries.remove(pos).name;
self.votes.remove(&id);
self.submitters.remove(&id);
self.save();
Ok(name)
}
}
/// Strip control chars + clamp length so a submission can't inject newlines/huge blobs.
fn clean(s: &str, max: usize) -> String {
s.chars().filter(|c| !c.is_control()).take(max).collect()
}
// ───────────────────────── callback grammar (self-describing) ─────────────────────────
#[derive(Clone, PartialEq, Debug)]
enum Origin {
Catalog(u32),
My,
}
impl Origin {
fn data(&self) -> String {
match self {
Origin::My => "@m".into(),
Origin::Catalog(p) => format!("@c{p}"),
}
}
fn back(&self) -> String {
match self {
Origin::My => "m".into(),
Origin::Catalog(p) => format!("c:{p}"),
}
}
}
/// Parse callback data into (token, optional id, back-origin). Fully self-describing — no server state.
fn parse_cb(data: &str) -> (String, Option<u32>, Origin) {
let (main, orig) = match data.split_once('@') {
Some((m, o)) => {
let orig = if o == "m" {
Origin::My
} else if let Some(p) = o.strip_prefix('c') {
Origin::Catalog(p.parse().unwrap_or(1).max(1))
} else {
Origin::Catalog(1)
};
(m, orig)
}
None => (data, Origin::Catalog(1)),
};
match main.split_once(':') {
Some((t, a)) => (t.to_string(), a.parse::<u32>().ok(), orig),
None => (main.to_string(), None, orig),
}
}
// ───────────────────────── screen builders (pure; each carries a fresh cmid) ─────────────────────────
fn btn(text: impl Into<String>, data: impl Into<String>) -> Button {
Button { text: text.into(), data: data.into().into_bytes() }
}
fn msg(text: String, buttons: Vec<Vec<Button>>) -> Msg {
Msg { text: Some(text), buttons, reply_to: None, client_msg_id: rand::random() }
}
fn net_badge(up: u32, down: u32) -> String {
let net = up as i64 - down as i64;
if net >= 0 {
format!("▲{net}")
} else {
format!("▼{}", -net)
}
}
fn count_mine(st: &State, viewer: &[u8; 32]) -> usize {
st.submitters.values().filter(|s| *s == viewer).count()
}
fn home_msg(st: &State, viewer: &[u8; 32]) -> Msg {
let n = st.catalog.entries.len();
if n == 0 {
return msg(
"📚 Каталог ботов Aster\nЗдесь пока пусто. Станьте первым, кто добавит бота — это займёт меньше минуты.".into(),
vec![vec![btn("➕ Добавить первого бота", "a")]],
);
}
let k = count_mine(st, viewer);
msg(
format!("📚 Каталог ботов Aster\nВитрина ботов сообщества — открывайте, оценивайте, добавляйте свои. Всё делается кнопками.\n\n🤖 Ботов в каталоге: {n}\n📁 Вы добавили: {k}"),
vec![vec![btn("📖 Открыть каталог", "c:1")], vec![btn("➕ Добавить бота", "a"), btn("📁 Мои боты", "m")]],
)
}
/// Entries sorted best-first: net rating desc, id asc (stable).
fn ranked(st: &State) -> Vec<&Entry> {
let mut v: Vec<&Entry> = st.catalog.entries.iter().collect();
v.sort_by_key(|e| (-(e.up as i64 - e.down as i64), e.id));
v
}
fn catalog_msg(st: &State, _viewer: &[u8; 32], page: u32) -> Msg {
let n = st.catalog.entries.len();
if n == 0 {
return msg(
"📭 Каталог пока пуст.\nДобавьте первого бота — это займёт меньше минуты.".into(),
vec![vec![btn("➕ Добавить бота", "a"), btn("🏠 Домой", "h")]],
);
}
let pages = n.div_ceil(PAGE_SIZE).max(1) as u32;
let p = page.clamp(1, pages);
let ordered = ranked(st);
let start = (p as usize - 1) * PAGE_SIZE;
let slice = &ordered[start..(start + PAGE_SIZE).min(n)];
let mut text = format!("📖 Каталог · {n} ботов · стр. {p}/{pages}\nСортировка: по рейтингу\n");
let mut rows: Vec<Vec<Button>> = Vec::new();
for e in slice {
let d = if e.desc.is_empty() { "—" } else { e.desc.as_str() };
text.push_str(&format!("\n🤖 {} · {}\n {}", e.name, net_badge(e.up, e.down), d));
rows.push(vec![btn(format!("{} ›", e.name), format!("b:{}@c{p}", e.id))]);
}
let prev = p.saturating_sub(1).max(1);
let next = (p + 1).min(pages);
rows.push(vec![btn("◀", format!("c:{prev}")), btn(format!("{p}/{pages} ⟳"), format!("c:{p}")), btn("▶", format!("c:{next}"))]);
rows.push(vec![btn("➕ Добавить", "a"), btn("📁 Мои", "m"), btn("🏠", "h")]);
msg(text, rows)
}
fn detail_msg(st: &State, viewer: &[u8; 32], viewer_addr: &str, id: u32, orig: &Origin) -> Msg {
let Some(e) = st.entry(id) else { return vanished_msg() };
let owner = st.is_owner(id, viewer, viewer_addr);
let my = st.my_vote(id, viewer);
let last = if owner {
"Добавили: вы".to_string()
} else {
format!("Ваш голос: {}", vote_label(my))
};
let text = format!(
"🤖 {}\n{}\n\nАдрес (нажмите и удерживайте, чтобы скопировать):\n{}\n\nРейтинг: 👍 {} · 👎 {}\n{}",
e.name,
if e.desc.is_empty() { "—" } else { &e.desc },
e.addr,
e.up,
e.down,
last
);
let os = orig.data();
let up = format!("👍 {}{}", e.up, if my == 1 { " ✓" } else { "" });
let down = format!("👎 {}{}", e.down, if my == -1 { " ✓" } else { "" });
let mut rows = vec![vec![btn(up, format!("u:{id}{os}")), btn(down, format!("d:{id}{os}"))]];
if owner {
rows.push(vec![btn("🗑 Убрать из каталога", format!("rq:{id}{os}"))]);
}
rows.push(vec![btn("◀ К списку", orig.back()), btn("🏠 Домой", "h")]);
msg(text, rows)
}
fn vote_label(v: i8) -> &'static str {
match v {
1 => "👍",
-1 => "👎",
_ => "—",
}
}
fn mybots_msg(st: &State, viewer: &[u8; 32]) -> Msg {
let mine: Vec<&Entry> = ranked(st).into_iter().filter(|e| st.submitters.get(&e.id) == Some(viewer)).collect();
if mine.is_empty() {
return msg(
"📁 Ваши боты\nВы ещё не добавляли ботов.\n(Если бот недавно перезапускался без сохранения состояния, привязка авторства могла сброситься.)".into(),
vec![vec![btn("➕ Добавить бота", "a"), btn("🏠 Домой", "h")]],
);
}
let mut text = format!("📁 Ваши боты · {}\n", mine.len());
let mut rows: Vec<Vec<Button>> = Vec::new();
for e in &mine {
text.push_str(&format!("\n🤖 {} — 👍{} 👎{}", e.name, e.up, e.down));
rows.push(vec![btn(format!("{} ›", e.name), format!("b:{}@m", e.id))]);
}
rows.push(vec![btn("➕ Добавить ещё", "a"), btn("🏠 Домой", "h")]);
msg(text, rows)
}
fn confirm_remove_msg(st: &State, id: u32, orig: &Origin) -> Msg {
let Some(e) = st.entry(id) else { return vanished_msg() };
let os = orig.data();
msg(
format!("🗑 Убрать «{}» из каталога?\nКарточка и её голоса будут удалены. Это не удаляет самого бота — только запись в каталоге.", e.name),
vec![vec![btn("✅ Да, убрать", format!("rm:{id}{os}"))], vec![btn("↩ Отмена", format!("b:{id}{os}"))]],
)
}
fn removed_msg(name: &str) -> Msg {
msg(
format!("🗑 «{name}» убран из каталога."),
vec![vec![btn("📁 Мои боты", "m"), btn("📖 Каталог", "c:1"), btn("🏠", "h")]],
)
}
fn removed_err_msg(reason: &str) -> Msg {
msg(
format!("⚠️ {reason}"),
vec![vec![btn("📁 Мои боты", "m"), btn("📖 Каталог", "c:1"), btn("🏠", "h")]],
)
}
fn vanished_msg() -> Msg {
msg(
"⚠️ Этого бота больше нет в каталоге\n(его убрал автор или админ).".into(),
vec![vec![btn("📖 Каталог", "c:1"), btn("🏠 Домой", "h")]],
)
}
fn breadcrumb_msg() -> Msg {
msg("➕ …".into(), Vec::new())
}
// ── wizard screens ──
fn cancel_row() -> Vec<Button> {
vec![btn("✖ Отмена", "x")]
}
/// The prompt for the wizard's current `asking` field (or the Confirm hub).
fn wizard_ask_msg(w: &Wizard) -> Msg {
match w.asking {
Field::Confirm => wizard_confirm_msg(w),
Field::Addr if w.rtc => msg(
"✏️ Пришлите новый адрес бота — строку вида ASTER1… 👇".into(),
vec![vec![btn("↩ Назад к проверке", "rev")], cancel_row()],
),
Field::Addr => msg(
"➕ Добавление бота · шаг 1 из 3\n\nПришлите адрес бота — строку вида ASTER1… (его выдаёт клиент владельца бота).\nПросто отправьте адрес сообщением 👇".into(),
vec![cancel_row()],
),
Field::Name if w.rtc => msg(
"✏️ Пришлите новое имя бота (до 60 символов) 👇".into(),
vec![vec![btn("↩ Назад к проверке", "rev")], cancel_row()],
),
Field::Name => msg(
"➕ Добавление бота · шаг 2 из 3\n✓ Адрес принят\n\nКак назвать бота? Пришлите короткое имя (до 60 символов) 👇".into(),
vec![vec![btn("⏭ Без названия", "sk")], cancel_row()],
),
Field::Desc if w.rtc => msg(
"✏️ Пришлите новое описание (до 200 символов) 👇".into(),
vec![vec![btn("↩ Назад к проверке", "rev")], cancel_row()],
),
Field::Desc => msg(
format!(
"➕ Добавление бота · шаг 3 из 3\n✓ Адрес принят · Имя: {}\n\nПришлите короткое описание (до 200 символов) — что умеет бот. Можно пропустить 👇",
w.name.as_deref().unwrap_or("Без названия")
),
vec![vec![btn("⏭ Пропустить", "sk")], cancel_row()],
),
}
}
fn wizard_confirm_msg(w: &Wizard) -> Msg {
let name = w.name.as_deref().unwrap_or("Без названия");
let desc = w.desc.as_deref().unwrap_or("");
let addr = w.addr.as_deref().unwrap_or("");
msg(
format!(
"➕ Проверьте перед публикацией\n\n🤖 {name}\n{}\nАдрес: {addr}\n\nБот появится в каталоге под вашим именем и станет виден всем. Всё верно?",
if desc.is_empty() { "—" } else { desc }
),
vec![
vec![btn("✅ Опубликовать", "ok")],
vec![btn("✏️ Имя", "en"), btn("✏️ Описание", "ed"), btn("✏️ Адрес", "ea")],
cancel_row(),
],
)
}
fn wizard_invalid_addr_msg() -> Msg {
msg(
"⚠️ Это не похоже на адрес Aster.\nАдрес начинается с «ASTER1». Скопируйте его из карточки бота и пришлите ещё раз 👇".into(),
vec![cancel_row()],
)
}
fn wizard_duplicate_msg(existing_id: u32) -> Msg {
msg(
"ℹ️ Этот бот уже есть в каталоге.\nМожно открыть его карточку или прислать другой адрес 👇".into(),
vec![vec![btn("Открыть карточку ›", format!("b:{existing_id}@c1"))], cancel_row()],
)
}
fn wizard_limit_msg(reason: &str) -> Msg {
msg(format!("⚠️ {reason}"), vec![vec![btn("🏠 Домой", "h")], cancel_row()])
}
fn wizard_success_msg(st: &State, id: u32) -> Msg {
let Some(e) = st.entry(id) else { return home_msg(st, &[0u8; 32]) };
msg(
format!(
"✅ Готово! «{}» добавлен в каталог.\nТеперь его видят все и могут оценить.\n\n🤖 {}\n{}\nАдрес: {}\nРейтинг: 👍 {} · 👎 {}",
e.name,
e.name,
if e.desc.is_empty() { "—" } else { &e.desc },
e.addr,
e.up,
e.down
),
vec![vec![btn("📖 Открыть в каталоге ›", format!("b:{id}@c1"))], vec![btn("➕ Добавить ещё", "a"), btn("🏠 Домой", "h")]],
)
}
fn wizard_publish_err_msg(reason: &str) -> Msg {
msg(
format!("⚠️ Не удалось опубликовать: {reason}\nВернитесь к проверке или отмените."),
vec![vec![btn("↩ К проверке", "rev")], cancel_row()],
)
}
// ───────────────────────── send helpers ─────────────────────────
/// A callback tap → Edit the tapped (on-screen) card in place; remember it as the app card.
fn edit(st: &mut State, viewer: [u8; 32], msg_ref: [u8; 16], screen: Msg) -> Vec<Outgoing> {
let ui = st.ui(viewer);
ui.app_cmid = Some(msg_ref);
ui.touched = Instant::now();
vec![ServiceFrame::Edit { msg_ref, msg: screen }.into()]
}
/// Open-app / bootstrap → a FRESH bottom card (no breadcrumb); remember its cmid.
fn fresh(st: &mut State, viewer: [u8; 32], m: Msg) -> Vec<Outgoing> {
let cmid = m.client_msg_id;
let ui = st.ui(viewer);
ui.app_cmid = Some(cmid);
ui.touched = Instant::now();
vec![ServiceFrame::Msg(m).into()]
}
/// A typed wizard step → collapse the previous prompt to a breadcrumb (best-effort) and send the next
/// step as a fresh bottom message (feedback must land at the bottom; toasts are ignored by the client).
fn reanchor(st: &mut State, viewer: [u8; 32], next: Msg) -> Vec<Outgoing> {
let cmid = next.client_msg_id;
let old = {
let ui = st.ui(viewer);
let old = ui.app_cmid;
ui.app_cmid = Some(cmid);
ui.touched = Instant::now();
old
};
let mut out = Vec::new();
if let Some(prev) = old {
out.push(ServiceFrame::Edit { msg_ref: prev, msg: breadcrumb_msg() }.into());
}
out.push(ServiceFrame::Msg(next).into());
out
}
// ───────────────────────── main ─────────────────────────
#[tokio::main]
async fn main() {
let a: Vec<String> = std::env::args().skip(1).collect();
if a.len() != 2 {
eprintln!("usage: printf '<128hex>' | catalogbot <server_addr> <server_pubhex>");
std::process::exit(2);
}
aster_services::apply_ram_protection();
let mut s = String::new();
std::io::stdin().read_to_string(&mut s).ok();
let bytes = match from_hex(s.trim()) {
Some(b) if b.len() == 64 => b,
_ => {
eprintln!("error: identity must be 128 hex chars (64 bytes)");
std::process::exit(1);
}
};
let (mut sign, mut dh) = ([0u8; 32], [0u8; 32]);
sign.copy_from_slice(&bytes[..32]);
dh.copy_from_slice(&bytes[32..]);
// Persist sessions when CATALOGBOT_BOTSTATE is set, so a redeploy doesn't drop every conversation.
let botstate = std::env::var("CATALOGBOT_BOTSTATE").ok().filter(|p| !p.is_empty());
let mut bot = match &botstate {
Some(p) => Bot::from_identity_secret_persist(sign, dh, p),
None => Bot::from_identity_secret(sign, dh),
};
let path = std::env::var("CATALOGBOT_STATE").ok().filter(|p| !p.is_empty());
let admin = std::env::var("CATALOGBOT_ADMIN").ok().map(|s| s.trim().to_string()).filter(|s| !s.is_empty());
let state = Arc::new(Mutex::new(State::load(path.clone(), admin.clone())));
println!("catalogbot address: {}", bot.address());
println!(" state: {}", path.as_deref().unwrap_or("(RAM-only, не сохраняется)"));
println!(" sessions: {}", botstate.as_deref().unwrap_or("(RAM-only, теряются при рестарте)"));
println!(" admin: {}", admin.as_deref().unwrap_or("(нет)"));
println!(" bots: {}", state.lock().unwrap().catalog.entries.len());
println!(" ui: кнопочная (без команд); добавление — пошаговый визард");
let handler: Handler = Box::new(move |u: &Update| -> Vec<Outgoing> {
let mut st = state.lock().unwrap();
let viewer = u.peer_sign;
match &u.event {
// Bootstrap ▶/start (or any slashy text) → a fresh Home card. No slash commands exist.
Incoming::Command { .. } => {
st.ui(viewer).wizard = None;
let m = home_msg(&st, &viewer);
fresh(&mut st, viewer, m)
}
Incoming::Message { text } => on_text(&mut st, viewer, text.trim()),
Incoming::Callback { msg_ref, data } => {
let ds = String::from_utf8_lossy(data).into_owned();
on_callback(&mut st, viewer, &u.peer_address(), *msg_ref, &ds)
}
Incoming::Frame(_) => Vec::new(),
}
});
if let Err(e) = bot.serve(&a[0], &a[1], handler).await {
eprintln!("fatal: {e}");
std::process::exit(1);
}
}
/// A plain text message: a guided wizard field, a paste-to-add shortcut, or "open the app".
fn on_text(st: &mut State, viewer: [u8; 32], text: &str) -> Vec<Outgoing> {
// Expire an abandoned wizard.
{
let ui = st.ui(viewer);
if ui.wizard.is_some() && Instant::now().duration_since(ui.touched) > WIZARD_TTL {
ui.wizard = None;
}
}
let taken = st.ui(viewer).wizard.take();
if let Some(mut w) = taken {
if w.awaits_text() {
let next = apply_wizard_text(st, viewer, &mut w, text);
st.ui(viewer).wizard = Some(w);
return reanchor(st, viewer, next);
}
// Confirm awaits buttons; a stray text just re-anchors Confirm unchanged.
let next = wizard_confirm_msg(&w);
st.ui(viewer).wizard = Some(w);
return reanchor(st, viewer, next);
}
// No wizard: paste-to-add if it's an address, else open Home.
if let Some(sign) = parse_account_address(text) {
let addr = account_address(&sign);
if let Some(e) = st.entry_by_addr(&addr) {
let m = wizard_duplicate_msg(e.id);
return reanchor(st, viewer, m);
}
if let Some(reason) = cap_reason(st, &viewer) {
let m = wizard_limit_msg(&reason);
return reanchor(st, viewer, m);
}
let w = Wizard { addr: Some(addr), asking: Field::Name, ..Default::default() };
let next = wizard_ask_msg(&w);
st.ui(viewer).wizard = Some(w);
return reanchor(st, viewer, next);
}
let m = home_msg(st, &viewer);
fresh(st, viewer, m)
}
/// Why a submission would be rejected on caps (checked early for a nicer wizard), or None if ok.
fn cap_reason(st: &State, by: &[u8; 32]) -> Option<String> {
if st.catalog.entries.len() >= MAX_ENTRIES {
Some("Каталог заполнен.".into())
} else if st.submitters.values().filter(|s| *s == by).count() >= MAX_PER_USER {
Some(format!("Лимит {MAX_PER_USER} ботов на пользователя."))
} else {
None
}
}
/// Advance the wizard given a typed field value; returns the next screen (Msg) to re-anchor.
fn apply_wizard_text(st: &State, _viewer: [u8; 32], w: &mut Wizard, text: &str) -> Msg {
let text = text.trim();
match w.asking {
Field::Addr => match parse_account_address(text) {
None => wizard_invalid_addr_msg(),
Some(sign) => {
let addr = account_address(&sign);
if let Some(e) = st.entry_by_addr(&addr) {
wizard_duplicate_msg(e.id)
} else {
w.addr = Some(addr);
if w.rtc {
w.fill();
w.asking = Field::Confirm;
} else {
w.asking = Field::Name;
}
wizard_ask_msg(w)
}
}
},
Field::Name => {
let n = clean(text, NAME_MAX);
w.name = Some(if n.is_empty() { "Без названия".into() } else { n });
if w.rtc {
w.fill();
w.asking = Field::Confirm;
} else {
w.asking = Field::Desc;
}
wizard_ask_msg(w)
}
Field::Desc => {
w.desc = Some(clean(text, DESC_MAX));
w.fill();
w.asking = Field::Confirm;
wizard_ask_msg(w)
}
Field::Confirm => wizard_confirm_msg(w),
}
}
/// A button tap: Edit the tapped card in place to the target screen.
fn on_callback(st: &mut State, viewer: [u8; 32], viewer_addr: &str, msg_ref: [u8; 16], data: &str) -> Vec<Outgoing> {
let (token, id, orig) = parse_cb(data);
// Any navigation away from a wizard implicitly cancels it.
if matches!(token.as_str(), "h" | "c" | "b" | "m" | "u" | "d" | "rq") {
st.ui(viewer).wizard = None;
}
let admin = st.admin.clone();
let _ = admin;
let screen: Msg = match token.as_str() {
"h" => home_msg(st, &viewer),
"c" => catalog_msg(st, &viewer, id.unwrap_or(1)),
"m" => mybots_msg(st, &viewer),
"b" => match id {
Some(i) => detail_msg(st, &viewer, viewer_addr, i, &orig),
None => home_msg(st, &viewer),
},
"u" | "d" => match id {
Some(i) => {
let dir = if token == "u" { 1 } else { -1 };
match st.vote(i, viewer, dir) {
Some(_) => detail_msg(st, &viewer, viewer_addr, i, &orig),
None => vanished_msg(),
}
}
None => home_msg(st, &viewer),
},
"rq" => match id {
Some(i) if st.entry(i).is_some() => confirm_remove_msg(st, i, &orig),
_ => vanished_msg(),
},
"rm" => match id {
Some(i) => match st.remove(i, viewer, viewer_addr) {
Ok(name) => removed_msg(&name),
Err(e) => removed_err_msg(&e),
},
None => home_msg(st, &viewer),
},
"a" => wiz_start(st, viewer),
"sk" => wiz_skip(st, viewer),
"ok" => wiz_publish(st, viewer, viewer_addr),
"en" | "ed" | "ea" => wiz_edit_field(st, viewer, &token),
"rev" => wiz_rev(st, viewer),
"x" => {
st.ui(viewer).wizard = None;
home_msg(st, &viewer)
}
_ => home_msg(st, &viewer),
};
edit(st, viewer, msg_ref, screen)
}
fn wiz_start(st: &mut State, viewer: [u8; 32]) -> Msg {
let w = Wizard::default();
let m = wizard_ask_msg(&w);
st.ui(viewer).wizard = Some(w);
m
}
fn wiz_skip(st: &mut State, viewer: [u8; 32]) -> Msg {
let mut w = match st.ui(viewer).wizard.take() {
Some(w) => w,
None => return home_msg(st, &viewer),
};
match w.asking {
Field::Name => {
w.name = Some("Без названия".into());
if w.rtc {
w.fill();
w.asking = Field::Confirm;
} else {
w.asking = Field::Desc;
}
}
Field::Desc => {
w.desc = Some(String::new());
w.fill();
w.asking = Field::Confirm;
}
_ => {}
}
let m = wizard_ask_msg(&w);
st.ui(viewer).wizard = Some(w);
m
}
fn wiz_edit_field(st: &mut State, viewer: [u8; 32], which: &str) -> Msg {
let mut w = match st.ui(viewer).wizard.take() {
Some(w) => w,
None => return home_msg(st, &viewer),
};
w.rtc = true;
w.asking = match which {
"en" => Field::Name,
"ed" => Field::Desc,
_ => Field::Addr,
};
let m = wizard_ask_msg(&w);
st.ui(viewer).wizard = Some(w);
m
}
fn wiz_rev(st: &mut State, viewer: [u8; 32]) -> Msg {
let mut w = match st.ui(viewer).wizard.take() {
Some(w) => w,
None => return home_msg(st, &viewer),
};
w.fill();
w.asking = Field::Confirm;
let m = wizard_confirm_msg(&w);
st.ui(viewer).wizard = Some(w);
m
}
fn wiz_publish(st: &mut State, viewer: [u8; 32], _viewer_addr: &str) -> Msg {
let w = match st.ui(viewer).wizard.take() {
Some(w) => w,
None => return home_msg(st, &viewer),
};
let addr = w.addr.clone().unwrap_or_default();
let name = w.name.clone().unwrap_or_else(|| "Без названия".into());
let desc = w.desc.clone().unwrap_or_default();
// Insert the fields STRUCTURALLY (no pipe round-trip) so a '|' in the name never bleeds into desc.
match st.add_fields(viewer, &addr, &name, &desc) {
Ok(id) => {
st.ui(viewer).wizard = None;
wizard_success_msg(st, id)
}
Err(reason) => {
// keep the wizard so the user can fix + retry
st.ui(viewer).wizard = Some(w);
wizard_publish_err_msg(&reason)
}
}
}
// ───────────────────────── tests ─────────────────────────
#[cfg(test)]
mod tests {
use super::*;
fn st() -> State {
State { catalog: Catalog::default(), path: None, admin: None, votes: HashMap::new(), submitters: HashMap::new(), uis: HashMap::new() }
}
fn addr_of(seed: u8) -> String {
account_address(&[seed; 32])
}
#[test]
fn add_validates_and_dedups() {
let mut s = st();
let good = addr_of(7);
let id = s.add([1u8; 32], &format!("{good} | Погодный | Прогноз погоды")).unwrap();
assert_eq!(s.entry(id).unwrap().name, "Погодный");
assert!(s.add([1u8; 32], "ASTER1NOTVALID | X | Y").is_err());
assert!(s.add([2u8; 32], &format!("{good} | Дубль |")).is_err());
let id2 = s.add([1u8; 32], &format!("{} | Имя\nсо\tсимволами | {}", addr_of(9), "д".repeat(500))).unwrap();
let e = s.entry(id2).unwrap();
assert!(!e.name.contains('\n') && !e.name.contains('\t'));
assert!(e.desc.chars().count() <= DESC_MAX);
}
#[test]
fn add_fields_keeps_pipes_in_name() {
// A name/description may contain '|'; the structured insert must store them verbatim
// (no bleed of the name tail into the description) so stored == confirmed preview.
let mut s = st();
let id = s.add_fields([1u8; 32], &addr_of(6), "Погода | Москва", "прогноз").unwrap();
let e = s.entry(id).unwrap();
assert_eq!(e.name, "Погода | Москва");
assert_eq!(e.desc, "прогноз");
}
#[test]
fn voting_is_one_per_user_and_toggles() {
let mut s = st();
let id = s.add([1u8; 32], &format!("{} | Бот |", addr_of(3))).unwrap();
let (a, b) = ([10u8; 32], [11u8; 32]);
assert_eq!(s.vote(id, a, 1), Some(1));
assert_eq!((s.entry(id).unwrap().up, s.entry(id).unwrap().down), (1, 0));
assert_eq!(s.vote(id, a, 1), Some(0));
assert_eq!(s.entry(id).unwrap().up, 0);
assert_eq!(s.vote(id, a, -1), Some(-1));
assert_eq!(s.entry(id).unwrap().down, 1);
assert_eq!(s.vote(id, a, 1), Some(1));
assert_eq!((s.entry(id).unwrap().up, s.entry(id).unwrap().down), (1, 0));
assert_eq!(s.vote(id, b, 1), Some(1));
assert_eq!(s.entry(id).unwrap().up, 2);
for _ in 0..5 {
s.vote(id, a, 1);
}
assert!(s.entry(id).unwrap().up <= 2);
}
#[test]
fn remove_authorization() {
let mut s = st();
let owner = [1u8; 32];
let id = s.add(owner, &format!("{} | Бот |", addr_of(4))).unwrap();
assert!(s.remove(id, [2u8; 32], "ASTER1STRANGER").is_err());
assert!(s.remove(id, owner, "ASTER1WHATEVER").is_ok());
assert!(s.entry(id).is_none());
}
#[test]
fn persistence_round_trip() {
let dir = std::env::temp_dir();
let p = dir.join("catalogbot-ux-test-state.json").to_string_lossy().to_string();
let _ = std::fs::remove_file(&p);
let voter = [9u8; 32];
let submitter = [1u8; 32];
let id;
{
let mut s = State::load(Some(p.clone()), None);
id = s.add(submitter, &format!("{} | Сохр |", addr_of(5))).unwrap();
s.vote(id, voter, 1);
}
let s2 = State::load(Some(p.clone()), None);
assert_eq!(s2.catalog.entries.len(), 1);
assert_eq!(s2.catalog.entries[0].up, 1);
assert_eq!(s2.my_vote(id, &voter), 1);
assert_eq!(s2.submitters.get(&id), Some(&submitter));
let _ = std::fs::remove_file(&p);
}
#[test]
fn callback_grammar_parses() {
let (t, id, o) = parse_cb("h");
assert_eq!((t.as_str(), id), ("h", None));
assert_eq!(o, Origin::Catalog(1)); // default origin
let (t, id, _) = parse_cb("c:3");
assert_eq!((t.as_str(), id), ("c", Some(3)));
let (t, id, o) = parse_cb("b:5@c2");
assert_eq!((t.as_str(), id, o), ("b", Some(5), Origin::Catalog(2)));
let (t, id, o) = parse_cb("u:12@m");
assert_eq!((t.as_str(), id, o), ("u", Some(12), Origin::My));
let (t, id, o) = parse_cb("rq:7");
assert_eq!((t.as_str(), id, o), ("rq", Some(7), Origin::Catalog(1)));
// back/data round-trip
assert_eq!(Origin::Catalog(2).back(), "c:2");
assert_eq!(Origin::My.data(), "@m");
// garbage → token with no id (dispatch maps unknown → home)
let (t, id, _) = parse_cb("garbage:x");
assert_eq!((t.as_str(), id), ("garbage", None));
}
#[test]
fn wizard_state_machine_happy_path_and_edits() {
let s = st();
// normal flow: Addr → Name → Desc → Confirm
let mut w = Wizard::default();
assert_eq!(w.asking, Field::Addr);
apply_wizard_text(&s, [1u8; 32], &mut w, &addr_of(8));
assert_eq!(w.asking, Field::Name);
assert!(w.addr.is_some());
apply_wizard_text(&s, [1u8; 32], &mut w, "Мой бот");
assert_eq!(w.asking, Field::Desc);
assert_eq!(w.name.as_deref(), Some("Мой бот"));
apply_wizard_text(&s, [1u8; 32], &mut w, "делает добро");
assert_eq!(w.asking, Field::Confirm);
assert_eq!(w.desc.as_deref(), Some("делает добро"));
// invalid address keeps us on Addr
let mut w2 = Wizard::default();
apply_wizard_text(&s, [1u8; 32], &mut w2, "не адрес");
assert_eq!(w2.asking, Field::Addr);
assert!(w2.addr.is_none());
// edit-a-single-field-from-Confirm returns straight to Confirm (rtc), keeping other fields
let mut w3 = w.clone();
w3.rtc = true;
w3.asking = Field::Name; // simulate [✏️ Имя]
apply_wizard_text(&s, [1u8; 32], &mut w3, "Новое имя");
assert_eq!(w3.asking, Field::Confirm);
assert_eq!(w3.name.as_deref(), Some("Новое имя"));
assert_eq!(w3.desc.as_deref(), Some("делает добро")); // preserved
// empty name → default
let mut w4 = Wizard { addr: Some(addr_of(2)), asking: Field::Name, ..Default::default() };
apply_wizard_text(&s, [1u8; 32], &mut w4, " ");
assert_eq!(w4.name.as_deref(), Some("Без названия"));
}
}
Диагностика
| Симптом | Куда смотреть |
|---|---|
| Бот молчит на сообщения | убедитесь, что запущен ровно один процесс с этим ключом (ps -C mybot); проверьте journalctl -u mybot; релей доступен и pubkey совпадает |
Падает при /-команде, которая зовёт систему | mlockall() под низким MEMLOCK — добавьте ASTER_NO_MLOCK=1 |
| Клиент показал «⟳ бот перезапустился» | норма после рестарта: RAM-сессии сброшены, сессия переустановится сама; нажмите ▶ /start |
| Ошибка «bad server pubkey hex» | второй аргумент — не 64 корректных hex-символа ключа релея |
| Медленные ответы | опрос адаптивный (чаще во время диалога); проверьте задержку до релея и что процесс один |
Документ соответствует крейту aster-services: типы Bot, Handler, Update, Incoming, Outgoing, ServiceFrame, FileSpec, Button, Msg. Готовые примеры — src/bin/demobot.rs (кнопки + команды /photo и /big для медиа) и src/bin/statusbot.rs (полноценный бот состояния сервера).