use crate::channels::ReplyReference;
use crate::channels::reply::normalize_reply_text;
use crate::pipeline::board::{Ticket, TicketPhase};
use crate::util::html::{decode_html_entities, escape_html, push_escaped};
use crate::util::media_target::{self, MediaTarget};
use crate::util::{
FILE_MAX_BYTES, MediaMarkerKind, TELEGRAM_MEDIA_MARKER_RE, UnwrapPoison, file_name_or_path,
is_http_url, parse_media_marker,
};
use crate::{Channel, ChannelMessage, SendMessage, Workspace};
use anyhow::Context;
use async_trait::async_trait;
use reqwest::multipart::{Form, Part};
use std::collections::HashMap;
use std::fmt::Write as _;
use std::path::{Path, PathBuf};
use std::time::Duration;
use strum::IntoEnumIterator;
const TELEGRAM_MAX_MESSAGE_LENGTH: usize = 4096;
const TELEGRAM_CONTINUATION_OVERHEAD: usize = 30;
const MAX_ATTACHMENT_FILENAME_BYTES: usize = 180;
const MEDIA_TRANSFER_TIMEOUT: Duration = Duration::from_mins(5);
const CLEAR_COMMAND_DESC: &str = "Reset your session";
const IMAGE_MODELS_COMMAND_DESC: &str = "Select image generation model";
const VIDEO_MODELS_COMMAND_DESC: &str = "Select video model";
const BOARD_COMMAND_DESC: &str = "List tickets from all workspaces";
const ARCHIVE_COMMAND_DESC: &str = "Archive done & cancelled tickets (all workspaces)";
const PAUSE_COMMAND_DESC: &str = "Pause the workspace pipeline";
const UNPAUSE_COMMAND_DESC: &str = "Resume the workspace pipeline";
const MAINTENANCE_ON_COMMAND_DESC: &str = "Enable workspace maintenance";
const MAINTENANCE_OFF_COMMAND_DESC: &str = "Disable workspace maintenance";
const UPDATE_COMMAND_DESC: &str = "Update MahBot to the latest version";
const WORKSPACE_COMMAND_DESC: &str = "Select the active workspace";
const ACTION_PREFIX: &str = "__act__";
#[must_use]
pub(crate) fn decode_action(content: &str) -> Option<(String, String)> {
let rest = content.strip_prefix(ACTION_PREFIX)?;
match rest.split_once('|') {
Some((action, payload)) => Some((action.to_string(), payload.to_string())),
None => Some((rest.to_string(), String::new())),
}
}
#[derive(Debug)]
pub enum ControlInput {
Action((String, String)),
Command(crate::BotCommand),
}
#[must_use]
pub fn control_input(msg: &ChannelMessage) -> Option<ControlInput> {
if let Some(action) = decode_action(&msg.content) {
return Some(ControlInput::Action(action));
}
if msg.channel == "telegram"
&& let Some(cmd) = crate::parse_bot_command(&msg.content)
{
return Some(ControlInput::Command(cmd));
}
None
}
#[must_use]
pub(crate) fn is_control_message(msg: &ChannelMessage) -> bool {
msg.callback_query_id.is_some() || control_input(msg).is_some()
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct IncomingAttachment {
file_id: String,
file_name: Option<String>,
file_size: Option<u64>,
caption: Option<String>,
kind: IncomingAttachmentKind,
mime_type: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum IncomingAttachmentKind {
Document,
Photo,
Video,
Audio,
Voice,
}
fn split_message_for_telegram(message: &str) -> Vec<String> {
if message.chars().count() <= TELEGRAM_MAX_MESSAGE_LENGTH {
return vec![message.to_string()];
}
let mut chunks = Vec::new();
let mut remaining = message;
let chunk_limit = TELEGRAM_MAX_MESSAGE_LENGTH - TELEGRAM_CONTINUATION_OVERHEAD;
while !remaining.is_empty() {
let hard_split = remaining
.char_indices()
.nth(chunk_limit)
.map_or(remaining.len(), |(idx, _)| idx);
let mut chunk_end = if hard_split == remaining.len() {
hard_split
} else {
find_split_boundary(remaining, hard_split)
};
if let Some(adjusted) = extend_past_open_tag(remaining, chunk_end) {
chunk_end = adjusted.min(hard_split);
}
chunks.push(remaining[..chunk_end].to_string());
remaining = &remaining[chunk_end..];
}
chunks
}
fn wrap_chunk(chunk: &str, index: usize, total: usize) -> String {
if total > 1 {
if index == 0 {
format!("{chunk}\n\n(continues...)")
} else if index == total - 1 {
format!("(continued)\n\n{chunk}")
} else {
format!("(continued)\n\n{chunk}\n\n(continues...)")
}
} else {
chunk.to_string()
}
}
fn find_split_boundary(text: &str, hard_split: usize) -> usize {
let search_area = &text[..hard_split];
search_area
.rfind('\n')
.max(search_area.rfind(' '))
.map_or(hard_split, |p| p + 1)
}
fn extend_past_open_tag(text: &str, pos: usize) -> Option<usize> {
let prefix = &text[..pos];
let last_open = prefix.rfind('<')?;
let mut in_quote = false;
let mut quote_char = '"';
for (i, c) in text[last_open..].char_indices() {
match c {
'"' | '\'' if !in_quote => {
in_quote = true;
quote_char = c;
}
'"' | '\'' if in_quote && c == quote_char => {
in_quote = false;
}
'>' if !in_quote => {
let gt_absolute = last_open + i;
if gt_absolute < pos {
return None; }
return Some(gt_absolute + 1); }
_ => {}
}
}
None
}
fn extract_sender_user_name(message: &serde_json::Value) -> String {
message
.get("from")
.and_then(|from| from.get("username"))
.and_then(serde_json::Value::as_str)
.unwrap_or(crate::users::TELEGRAM_UNKNOWN_SENTINEL)
.to_string()
}
struct TelegramSender {
numeric_id: Option<String>,
nickname: Option<String>,
}
fn telegram_sender_person(source: &serde_json::Value) -> Option<TelegramSender> {
if source.get("sender_chat").is_some() {
return None;
}
let from = source.get("from")?;
if from.get("is_bot").and_then(serde_json::Value::as_bool) == Some(true) {
return None;
}
let numeric_id = from
.get("id")
.and_then(serde_json::Value::as_i64)
.map(|id| id.to_string());
if let Some(id) = &numeric_id
&& crate::users::is_telegram_service_id(id)
{
return None;
}
let nickname = from
.get("username")
.and_then(serde_json::Value::as_str)
.filter(|nickname| *nickname != crate::users::TELEGRAM_UNKNOWN_SENTINEL)
.map(String::from);
Some(TelegramSender {
numeric_id,
nickname,
})
}
struct MessageContext {
user_name: String,
chat_id: String,
message_id: i64,
reply_target: String,
reply_reference: Option<ReplyReference>,
}
impl MessageContext {
fn into_channel_message(self, content: String, cq_id: Option<String>) -> ChannelMessage {
ChannelMessage {
user_name: self.user_name,
reply_target: self.reply_target,
content,
channel: "telegram".to_string(),
workspace: String::new(),
optimistic_id: None,
callback_query_id: cq_id,
reply_reference: self.reply_reference,
chat_id: (!self.chat_id.is_empty()).then_some(self.chat_id),
message_id: (self.message_id != 0).then_some(self.message_id),
attachment_dirs: Vec::new(),
parts: Vec::new(),
}
}
fn into_attachment_message(self, content: String, staging_dir: String) -> ChannelMessage {
ChannelMessage {
attachment_dirs: vec![staging_dir],
..self.into_channel_message(content, None)
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum TelegramAttachmentKind {
Image,
Document,
Video,
Audio,
}
impl From<MediaMarkerKind> for TelegramAttachmentKind {
fn from(kind: MediaMarkerKind) -> Self {
match kind {
MediaMarkerKind::Image => Self::Image,
MediaMarkerKind::Audio => Self::Audio,
MediaMarkerKind::Video => Self::Video,
MediaMarkerKind::File => Self::Document,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct TelegramAttachment {
kind: TelegramAttachmentKind,
target: String,
}
#[derive(Debug, Clone, Copy)]
struct AttachmentMeta {
api_method: &'static str,
form_field: &'static str,
default_filename: &'static str,
label: &'static str,
disable_content_type_detection: bool,
}
impl TelegramAttachmentKind {
const fn meta(self) -> AttachmentMeta {
match self {
Self::Image => AttachmentMeta {
api_method: "sendPhoto",
form_field: "photo",
default_filename: "photo.jpg",
label: "Image",
disable_content_type_detection: false,
},
Self::Document => AttachmentMeta {
api_method: "sendDocument",
form_field: "document",
default_filename: "file",
label: "Document",
disable_content_type_detection: true,
},
Self::Video => AttachmentMeta {
api_method: "sendVideo",
form_field: "video",
default_filename: "video.mp4",
label: "Video",
disable_content_type_detection: false,
},
Self::Audio => AttachmentMeta {
api_method: "sendAudio",
form_field: "audio",
default_filename: "audio.mp3",
label: "Audio",
disable_content_type_detection: false,
},
}
}
const fn file_meta(self) -> AttachmentMeta {
match self {
Self::Image => AttachmentMeta {
api_method: "sendDocument",
form_field: "document",
default_filename: "image",
label: "Image",
disable_content_type_detection: true,
},
_ => self.meta(),
}
}
}
const IMAGE_EXTENSIONS: &[&str] = &["png", "jpg", "jpeg", "webp"];
const IMAGE_MIME_PREFIXES: &[&str] = &["image/png", "image/jpeg", "image/jpg", "image/webp"];
fn sender_label_with_fallback(from: &serde_json::Value, fallback: &str) -> String {
from.get("username")
.and_then(serde_json::Value::as_str)
.map_or_else(
|| {
from.get("first_name")
.and_then(serde_json::Value::as_str)
.unwrap_or(fallback)
.to_string()
},
|u| format!("@{u}"),
)
}
#[must_use]
fn format_sender_label(from: &serde_json::Value) -> String {
sender_label_with_fallback(from, crate::users::TELEGRAM_UNKNOWN_SENTINEL)
}
#[must_use]
fn replied_to_sender_label(message: &serde_json::Value) -> String {
let label = if let Some(sender_chat) = message.get("sender_chat") {
sender_chat
.get("title")
.and_then(serde_json::Value::as_str)
.unwrap_or("channel")
.to_string()
} else if let Some(from) = message.get("from") {
if from
.get("is_bot")
.and_then(serde_json::Value::as_bool)
.unwrap_or(false)
{
from.get("username")
.and_then(serde_json::Value::as_str)
.map_or_else(|| "bot".to_string(), |u| format!("@{u}"))
} else {
sender_label_with_fallback(from, "user")
}
} else {
"user".to_string()
};
label.replace(['<', '>'], "")
}
#[must_use]
fn media_kind_placeholder(reply_to: &serde_json::Value) -> String {
if reply_to.get("voice").is_some() || reply_to.get("audio").is_some() {
"[Voice message]".to_string()
} else if reply_to.get("photo").is_some() {
"[Photo]".to_string()
} else if reply_to.get("document").is_some() {
"[Document]".to_string()
} else if reply_to.get("video").is_some() {
"[Video]".to_string()
} else if reply_to.get("sticker").is_some() {
"[Sticker]".to_string()
} else {
"[Message]".to_string()
}
}
#[must_use]
fn reply_to_snippet(reply_to: &serde_json::Value) -> String {
if let Some(text) = reply_to.get("text").and_then(serde_json::Value::as_str)
&& !text.is_empty()
{
normalize_reply_text(text)
} else if let Some(caption) = reply_to.get("caption").and_then(serde_json::Value::as_str)
&& !caption.is_empty()
{
normalize_reply_text(caption)
} else {
media_kind_placeholder(reply_to)
}
}
#[must_use]
fn build_reply_reference(message: &serde_json::Value) -> Option<ReplyReference> {
let reply_to = message.get("reply_to_message")?;
if !reply_to.is_object() {
return None;
}
Some(ReplyReference {
author: replied_to_sender_label(reply_to),
snippet: reply_to_snippet(reply_to),
})
}
fn format_attachment_content(
kind: IncomingAttachmentKind,
local_path: &Path,
mime_type: Option<&str>,
) -> String {
let is_image = crate::util::has_extension(local_path, IMAGE_EXTENSIONS)
|| mime_type.is_some_and(|m| IMAGE_MIME_PREFIXES.iter().any(|p| m.starts_with(p)));
let is_video = crate::util::is_video_extension(local_path)
|| mime_type.is_some_and(|m| m.starts_with("video/"));
match kind {
IncomingAttachmentKind::Photo | IncomingAttachmentKind::Document if is_image => {
format!("[IMAGE:{}]", local_path.display())
}
IncomingAttachmentKind::Video | IncomingAttachmentKind::Document if is_video => {
format!("[VIDEO:{}]", local_path.display())
}
IncomingAttachmentKind::Audio | IncomingAttachmentKind::Voice => {
let kind = if cfg!(target_os = "macos") {
"AUDIO"
} else {
"FILE"
};
format!("[{kind}:{}]", local_path.display())
}
_ => {
let target = local_path.to_string_lossy().to_string();
if crate::util::with_block_in_place(|| {
crate::util::media_target::classify_media_image_target(&target)
}) == MediaTarget::LocalImage
{
format!("[IMAGE:{}]", local_path.display())
} else {
format!("[FILE:{}]", local_path.display())
}
}
}
}
fn normalize_video_filename(
kind: IncomingAttachmentKind,
filename: &str,
mime_type: Option<&str>,
) -> String {
let is_video =
kind == IncomingAttachmentKind::Video || mime_type.is_some_and(|m| m.starts_with("video/"));
if is_video && !crate::util::is_video_extension(std::path::Path::new(filename)) {
let Some(stem) = std::path::Path::new(filename)
.file_stem()
.and_then(|s| s.to_str())
else {
return filename.to_string();
};
let ext = match mime_type {
Some("video/webm") => "webm",
Some("video/quicktime") => "mov",
_ => "mp4",
};
format!("{stem}.{ext}")
} else {
filename.to_string()
}
}
fn sanitize_attachment_filename(name: &str, fallback: &str) -> String {
let mut cleaned = [name, fallback]
.into_iter()
.find_map(|name| {
let component = Path::new(name).file_name().and_then(|n| n.to_str())?;
Some(crate::util::neutralized_name(component))
})
.unwrap_or_else(|| "file".to_string());
if cleaned.starts_with('.') {
cleaned.insert(0, '_');
}
if cleaned.len() > MAX_ATTACHMENT_FILENAME_BYTES {
let stem = crate::util::name_stem(&cleaned);
let ext = &cleaned[stem.len()..];
let truncated = crate::util::truncate_bytes(
stem,
MAX_ATTACHMENT_FILENAME_BYTES.saturating_sub(ext.len()),
);
cleaned = if truncated.is_empty() {
crate::util::truncate_bytes(cleaned.as_str(), MAX_ATTACHMENT_FILENAME_BYTES).to_string()
} else {
format!("{truncated}{ext}")
};
}
cleaned
}
fn sanitize_remote_extension(remote_ext: &str) -> Option<String> {
let cleaned: String = remote_ext
.chars()
.filter(char::is_ascii_alphanumeric)
.collect();
(!cleaned.is_empty()).then_some(cleaned)
}
fn local_attachment_name(
attachment: &IncomingAttachment,
chat_id: &str,
message_id: i64,
remote_ext: Option<&str>,
) -> String {
let remote_ext = remote_ext.and_then(sanitize_remote_extension);
let ext = remote_ext.as_deref().unwrap_or(match attachment.kind {
IncomingAttachmentKind::Video => "mp4",
IncomingAttachmentKind::Document => "bin",
_ => "jpg",
});
let prefix = match attachment.kind {
IncomingAttachmentKind::Photo => "photo",
IncomingAttachmentKind::Video => "video",
IncomingAttachmentKind::Audio | IncomingAttachmentKind::Voice => "audio",
IncomingAttachmentKind::Document => "file",
};
let fallback = format!("{prefix}_{chat_id}_{message_id}.{ext}");
let chosen = normalize_video_filename(
attachment.kind,
attachment.file_name.as_deref().unwrap_or(&fallback),
attachment.mime_type.as_deref(),
);
sanitize_attachment_filename(&chosen, &fallback)
}
#[must_use]
fn attachment_rejection_content(
display_name: Option<&str>,
outcome: &str,
reason: &str,
caption: Option<&str>,
) -> String {
let note = match display_name {
Some(name) => format!("[File {name}: {outcome} — {reason}]"),
None => format!("[File: {outcome} — {reason}]"),
};
match caption.filter(|c| !c.is_empty()) {
Some(caption) => format!("{caption}\n\n{note}"),
None => note,
}
}
fn is_attachable_target(kind: TelegramAttachmentKind, target: &str) -> bool {
match kind {
TelegramAttachmentKind::Image => matches!(
media_target::classify_media_image_target(target),
MediaTarget::LocalImage | MediaTarget::RemoteUrl
),
TelegramAttachmentKind::Video | TelegramAttachmentKind::Audio => {
is_http_url(target) || Path::new(target).is_file()
}
TelegramAttachmentKind::Document => false,
}
}
fn file_refusal(name: &str, reason: &str) -> String {
format!("Could not send \"{name}\": {reason}")
}
fn resolve_file_target(target: &str, file_roots: &[PathBuf]) -> Result<PathBuf, String> {
let name = crate::util::truncate(file_name_or_path(target), 60);
if is_http_url(target) {
return Err(file_refusal(
&name,
"FILE targets are local files, not URLs.",
));
}
if file_roots.is_empty() {
return Err(file_refusal(
&name,
"no workspace is available for file delivery.",
));
}
let roots: Vec<PathBuf> = file_roots
.iter()
.map(|root| std::fs::canonicalize(root).unwrap_or_else(|_| root.clone()))
.collect();
let raw = Path::new(target);
let joined = if raw.is_absolute() {
raw.to_path_buf()
} else {
roots
.iter()
.map(|root| root.join(raw))
.find(|candidate| candidate.exists())
.unwrap_or_else(|| roots[0].join(raw))
};
let Ok(canonical) = std::fs::canonicalize(&joined) else {
return Err(file_refusal(&name, "file not found."));
};
if !crate::tools::path::is_path_under_roots(&canonical, &roots) {
return Err(file_refusal(&name, "it is outside your workspace."));
}
let Ok(meta) = std::fs::metadata(&canonical) else {
return Err(file_refusal(&name, "file not found."));
};
if meta.is_dir() {
return Err(file_refusal(&name, "it is a directory."));
}
if meta.len() > FILE_MAX_BYTES {
return Err(file_refusal(
&name,
&format!(
"it is {} MB, larger than the {} MB limit.",
meta.len().div_ceil(1024 * 1024),
mb(FILE_MAX_BYTES)
),
));
}
Ok(crate::util::strip_verbatim_prefix(&canonical))
}
fn parse_attachment_markers(
message: &str,
file_roots: &[PathBuf],
) -> (String, Vec<TelegramAttachment>, Vec<String>) {
let mut attachments: Vec<TelegramAttachment> = Vec::new();
let mut refusals: Vec<String> = Vec::new();
let cleaned = TELEGRAM_MEDIA_MARKER_RE
.replace_all(message, |caps: ®ex::Captures| {
let (marker_kind, path) = parse_media_marker(caps);
let path = path.trim();
if marker_kind == MediaMarkerKind::File {
match resolve_file_target(path, file_roots) {
Ok(canonical) => attachments.push(TelegramAttachment {
kind: TelegramAttachmentKind::Document,
target: canonical.to_string_lossy().into_owned(),
}),
Err(notice) => refusals.push(notice),
}
return String::new();
}
let kind = TelegramAttachmentKind::from(marker_kind);
if !is_attachable_target(kind, path) {
return caps.get_match().as_str().to_string();
}
attachments.push(TelegramAttachment {
kind,
target: path.to_string(),
});
String::new()
})
.to_string();
(cleaned.trim().to_string(), attachments, refusals)
}
const API_BASE: &str = "https://api.telegram.org";
fn bot_api_url(token: &str, method: &str) -> String {
format!("{API_BASE}/bot{token}/{method}")
}
const TELEGRAM_MAX_FILE_DOWNLOAD_BYTES: u64 = 20 * 1024 * 1024;
fn mb(bytes: u64) -> u64 {
bytes / (1024 * 1024)
}
fn file_too_large_reason() -> String {
format!("the file exceeds the {} MB limit", mb(FILE_MAX_BYTES))
}
fn telegram_download_limit_reason() -> String {
format!(
"Telegram only lets the bot download files up to {} MB",
mb(TELEGRAM_MAX_FILE_DOWNLOAD_BYTES)
)
}
const TELEGRAM_TRANSFER_REFUSED_REASON: &str =
"Telegram refused the transfer (it may exceed what the bot can download)";
const ATTACHMENT_DOWNLOAD_TIMEOUT_REASON: &str = "the download took too long and was stopped";
const ATTACHMENT_DOWNLOAD_FAILED_REASON: &str = "the connection to Telegram failed";
const TELEGRAM_LOOKUP_REFUSED_REASON: &str = "Telegram could not provide the file";
const NOT_RECEIVED: &str = "not received";
const NOT_STORED: &str = "could not be stored";
const ATTACHMENT_DIR_REFUSED_REASON: &str =
"the local temp directory for the file could not be created";
const ATTACHMENT_SAVE_REFUSED_REASON: &str = "the file could not be saved locally";
const ATTACHMENT_UPLOAD_FAILED_REASON: &str = "the upload failed";
#[cfg(not(target_os = "macos"))]
const VOICE_UNSUPPORTED_REASON: &str = "voice messages are not supported on this platform";
#[must_use]
fn declared_size_refusal(size: u64) -> Option<String> {
if size > FILE_MAX_BYTES {
Some(file_too_large_reason())
} else if size > TELEGRAM_MAX_FILE_DOWNLOAD_BYTES {
Some(telegram_download_limit_reason())
} else {
None
}
}
fn download_failure_reason(error: &anyhow::Error) -> &'static str {
match error
.chain()
.find_map(|cause| cause.downcast_ref::<reqwest::Error>())
{
Some(e) if e.is_timeout() => ATTACHMENT_DOWNLOAD_TIMEOUT_REASON,
Some(_) => ATTACHMENT_DOWNLOAD_FAILED_REASON,
None => TELEGRAM_TRANSFER_REFUSED_REASON,
}
}
enum ChatMenuState {
Registered(Option<String>),
InFlight,
}
pub struct TelegramChannel {
bot_token: String,
http_client: reqwest::Client,
cancel: std::sync::Arc<tokio_util::sync::CancellationToken>,
offset: std::sync::Arc<std::sync::atomic::AtomicI64>,
menu_cache: std::sync::Arc<std::sync::Mutex<std::collections::HashMap<String, ChatMenuState>>>,
}
fn extract_chat_context(message: &serde_json::Value) -> Option<(String, String)> {
let chat_id = message.get("chat")?.get("id")?.as_i64()?.to_string();
let thread_id = message
.get("message_thread_id")
.and_then(serde_json::Value::as_i64)
.map(|id| id.to_string());
let reply_target = match &thread_id {
Some(tid) => format!("{chat_id}:{tid}"),
None => chat_id.clone(),
};
Some((chat_id, reply_target))
}
fn set_thread_id_on_json(body: &mut serde_json::Value, thread_id: Option<&str>) {
if let Some(tid) = thread_id {
body["message_thread_id"] = serde_json::Value::String(tid.to_string());
}
}
fn parse_recipient(recipient: &str) -> (&str, Option<&str>) {
match recipient.split_once(':') {
Some((chat, thread)) => (chat, Some(thread)),
None => (recipient, None),
}
}
async fn resolve_authorized_sender(
sender_source: &serde_json::Value,
chat_source: &serde_json::Value,
) -> Option<(String, String, String)> {
let sender = telegram_sender_person(sender_source)?;
let matched = match crate::users::store()
.resolve_telegram_user(sender.numeric_id.as_deref(), sender.nickname.as_deref())
.await
{
Ok(matched) => matched?,
Err(e) => {
tracing::warn!(error = %e, "Failed to resolve Telegram sender to an account");
return None;
}
};
let (chat_id, reply_target) = extract_chat_context(chat_source)?;
let _ =
crate::users::update_channel_contact("telegram", &matched.identifier, &reply_target).await;
Some((matched.user_name, chat_id, reply_target))
}
fn parse_markdown_link(text: &str, i: usize) -> Option<(&str, &str, usize)> {
let after_open = text.get(i..)?.strip_prefix('[')?;
let bracket_end = after_open.find(']')?;
let label = &after_open[..bracket_end];
let after_bracket = i + bracket_end + 2;
let after_url_open = text.get(after_bracket..)?.strip_prefix('(')?;
let url = &after_url_open[..after_url_open.find(')')?];
if !is_http_url(url) || is_marker_label(label) {
return None;
}
Some((label, url, after_bracket + url.len() + 2))
}
fn is_marker_label(label: &str) -> bool {
let Some((kind, _)) = label.split_once(':') else {
return false;
};
MediaMarkerKind::iter().any(|marker| marker.token().eq_ignore_ascii_case(kind))
}
const ANCHOR_PREFIX: &str = "<a href=\"";
fn render_if_link(text: &str) -> Option<String> {
let rendered = render_inline(text);
rendered.contains(ANCHOR_PREFIX).then_some(rendered)
}
fn try_format_inline(
text: &str,
i: &mut usize,
out: &mut String,
delim: &str,
tag: &str,
linkify: bool,
) -> bool {
if text[*i..].starts_with(delim) {
let content_start = *i + delim.len();
if let Some(end) = text[content_start..].find(delim)
&& end > 0
{
let inner = &text[content_start..content_start + end];
if linkify && let Some(rendered) = render_if_link(inner) {
out.push_str(&rendered);
} else {
let _ = write!(out, "<{tag}>{}</{tag}>", escape_html(inner));
}
*i += delim.len() * 2 + end;
return true;
}
}
false
}
fn render_inline(text: &str) -> String {
let mut out = String::with_capacity(text.len());
let bytes = text.as_bytes();
let len = bytes.len();
let mut i = 0;
while i < len {
if try_format_inline(text, &mut i, &mut out, "**", "b", true) {
continue;
}
if try_format_inline(text, &mut i, &mut out, "__", "b", true) {
continue;
}
if (i == 0 || bytes[i - 1] != b'*')
&& try_format_inline(text, &mut i, &mut out, "*", "i", true)
{
continue;
}
if (i == 0 || bytes[i - 1] != b'`')
&& try_format_inline(text, &mut i, &mut out, "`", "code", false)
{
continue;
}
if bytes[i] == b'['
&& let Some((label, url, next)) = parse_markdown_link(text, i)
{
let _ = write!(
out,
"{ANCHOR_PREFIX}{}\">{}</a>",
escape_html(url),
escape_html(label)
);
i = next;
continue;
}
if try_format_inline(text, &mut i, &mut out, "~~", "s", true) {
continue;
}
let ch = text[i..].chars().next().unwrap();
push_escaped(ch, &mut out);
i += ch.len_utf8();
}
out
}
fn markdown_to_telegram_html(text: &str) -> String {
let mut out = String::with_capacity(text.len());
let mut in_code_block = false;
let mut code_buf = String::new();
for line in text.split('\n') {
let trimmed = line.trim_start();
if trimmed.starts_with("```") {
if in_code_block {
in_code_block = false;
let escaped = escape_html(code_buf.trim_end_matches('\n'));
let _ = writeln!(out, "<pre><code>{escaped}</code></pre>");
} else {
in_code_block = true;
}
code_buf.clear();
continue;
}
if in_code_block {
code_buf.push_str(line);
code_buf.push('\n');
continue;
}
if trimmed == "<blockquote>" || trimmed == "</blockquote>" {
out.push_str(trimmed);
out.push('\n');
continue;
}
let stripped = line.trim_start_matches('#');
let header_level = line.len() - stripped.len();
if header_level > 0 && stripped.starts_with(' ') {
let title = stripped.trim();
match render_if_link(title) {
Some(rendered) => {
let _ = writeln!(out, "{rendered}");
}
None => {
let _ = writeln!(out, "<b>{}</b>", escape_html(title));
}
}
continue;
}
out.push_str(&render_inline(line));
out.push('\n');
}
if in_code_block && !code_buf.is_empty() {
let _ = writeln!(
out,
"<pre><code>{}</code></pre>",
escape_html(code_buf.trim_end())
);
}
out.trim_end_matches('\n').to_string()
}
fn strip_html_tags(s: &str) -> String {
let mut out = String::with_capacity(s.len());
let mut in_tag = false;
let mut in_quote = false;
let mut quote_char = '"';
for c in s.chars() {
match c {
'<' if !in_tag => in_tag = true,
'>' if in_tag && !in_quote => in_tag = false,
'"' | '\'' if in_tag => {
if in_quote && c == quote_char {
in_quote = false;
} else if !in_quote {
in_quote = true;
quote_char = c;
}
}
_ if !in_tag => out.push(c),
_ => {}
}
}
out
}
#[derive(Debug)]
pub enum EditMessageFailure {
NotFound,
CannotEdit,
NotModified,
Other,
}
fn classify_edit_failure(description: &str) -> EditMessageFailure {
let lower = description.to_lowercase();
if lower.contains("not found") {
EditMessageFailure::NotFound
} else if lower.contains("not modified") {
EditMessageFailure::NotModified
} else if lower.contains("can't be edited") || lower.contains("cant be edited") {
EditMessageFailure::CannotEdit
} else {
EditMessageFailure::Other
}
}
fn to_telegram_html(text: &str) -> String {
markdown_to_telegram_html(&decode_html_entities(text))
}
pub(super) const FORWARD_ATTRIBUTION_PREFIX: &str = "[Forwarded from ";
impl TelegramChannel {
#[must_use]
fn new_with(bot_token: String, offset: std::sync::Arc<std::sync::atomic::AtomicI64>) -> Self {
Self {
bot_token,
http_client: crate::util::http::build_http_client(Duration::from_mins(1)),
cancel: std::sync::Arc::new(tokio_util::sync::CancellationToken::new()),
offset,
menu_cache: std::sync::Arc::new(
std::sync::Mutex::new(std::collections::HashMap::new()),
),
}
}
#[must_use]
pub fn new(bot_token: String) -> Self {
Self::new_with(
bot_token,
std::sync::Arc::new(std::sync::atomic::AtomicI64::new(0)),
)
}
#[must_use]
fn with_offset(
bot_token: String,
inherited_offset: std::sync::Arc<std::sync::atomic::AtomicI64>,
) -> Self {
Self::new_with(bot_token, inherited_offset)
}
pub async fn answer_callback_query(&self, callback_query_id: &str, text: Option<&str>) {
let mut body = serde_json::json!({
"callback_query_id": callback_query_id,
});
if let Some(txt) = text {
body["text"] = serde_json::Value::String(txt.to_string());
}
if let Err((status, error)) = self
.post_telegram_json("answerCallbackQuery", body, "answerCallbackQuery error")
.await
{
tracing::warn!(
callback_query_id = %callback_query_id,
status = %status,
error = %error,
"answerCallbackQuery failed"
);
}
}
async fn parse_callback_query(&self, cq: &serde_json::Value) -> Option<ChannelMessage> {
let data = cq.get("data").and_then(serde_json::Value::as_str)?;
let msg = cq.get("message")?;
let callback_query_id = cq
.get("id")
.and_then(serde_json::Value::as_str)
.map(String::from);
let Some((user_name, chat_id, reply_target)) = resolve_authorized_sender(cq, msg).await
else {
tracing::debug!(
"Telegram: ignoring callback query from unknown user '{}'",
extract_sender_user_name(cq)
);
return None;
};
let message_id = msg
.get("message_id")
.and_then(serde_json::Value::as_i64)
.unwrap_or(0);
let ctx = MessageContext {
user_name,
chat_id,
message_id,
reply_target,
reply_reference: None,
};
Some(ctx.into_channel_message(data.to_string(), callback_query_id))
}
fn extract_update_message_target(update: &serde_json::Value) -> Option<(String, i64)> {
let message = update.get("message")?;
let chat_id = extract_chat_context(message)?.0;
let message_id = message
.get("message_id")
.and_then(serde_json::Value::as_i64)?;
Some((chat_id, message_id))
}
async fn extract_message_context(&self, message: &serde_json::Value) -> Option<MessageContext> {
let Some((user_name, chat_id, reply_target)) =
resolve_authorized_sender(message, message).await
else {
tracing::debug!(
"Telegram: ignoring message from unknown user '{}'",
extract_sender_user_name(message)
);
return None;
};
let message_id = message
.get("message_id")
.and_then(serde_json::Value::as_i64)
.unwrap_or(0);
Some(MessageContext {
user_name,
chat_id,
message_id,
reply_target,
reply_reference: build_reply_reference(message),
})
}
fn prepend_forward_attribution(content: String, message: &serde_json::Value) -> String {
if let Some(attr) = Self::format_forward_attribution(message) {
format!("{attr}{content}")
} else {
content
}
}
fn try_add_ack_reaction_nonblocking(&self, chat_id: String, message_id: i64) {
let client = self.http_client().clone();
let url = self.api_url("setMessageReaction");
let body = serde_json::json!({
"chat_id": &chat_id,
"message_id": message_id,
"reaction": [{"type": "emoji", "emoji": "👀"}]
});
tokio::spawn(async move {
let response = match client.post(&url).json(&body).send().await {
Ok(resp) => resp,
Err(err) => {
tracing::warn!(
"Telegram: failed to add ACK reaction to chat_id={chat_id}, message_id={message_id}: {err}"
);
return;
}
};
if !response.status().is_success() {
let status = response.status();
let err_body =
crate::util::http::read_error_body(response, "ACK reaction error").await;
tracing::warn!(
"Telegram: add ACK reaction failed for chat_id={chat_id}, message_id={message_id}: status={status}, body={err_body}"
);
}
});
}
const fn http_client(&self) -> &reqwest::Client {
&self.http_client
}
fn api_url(&self, method: &str) -> String {
bot_api_url(&self.bot_token, method)
}
fn cancel_own(&self) {
self.cancel.cancel();
}
pub async fn set_my_commands(&self) {
let body = serde_json::json!({
"commands": [
{"command": "clear", "description": CLEAR_COMMAND_DESC},
]
});
post_set_my_commands(self.http_client(), &self.api_url("setMyCommands"), &body).await;
}
fn spawn_menu_refresh(&self, chat_id: &str) {
{
let cache = self.menu_cache.lock().unwrap_poison();
if matches!(cache.get(chat_id), Some(ChatMenuState::InFlight)) {
return;
}
}
let bot_token = self.bot_token.clone();
let http_client = self.http_client.clone();
let cache = std::sync::Arc::clone(&self.menu_cache);
let chat_id = chat_id.to_string();
tokio::spawn(async move {
let Some(user_name) =
crate::users::resolve_user_by_reply_target("telegram", &chat_id).await
else {
return;
};
let entries = user_command_entries(&user_name).await;
let payload = serde_json::json!({
"commands": entries
.iter()
.map(|(cmd, desc)| serde_json::json!({ "command": cmd, "description": desc }))
.collect::<Vec<_>>(),
});
let payload_str = payload.to_string();
let should_send = {
let mut cache = cache.lock().unwrap_poison();
match cache.get(&chat_id) {
Some(ChatMenuState::Registered(last))
if last.as_deref() == Some(&payload_str) =>
{
false
}
Some(ChatMenuState::InFlight) => false,
_ => {
cache.insert(chat_id.clone(), ChatMenuState::InFlight);
true
}
}
};
if !should_send {
return;
}
let url = bot_api_url(&bot_token, "setMyCommands");
let body = serde_json::json!({
"scope": {
"type": "chat",
"chat_id": chat_id.parse::<i64>().map_or_else(
|_| serde_json::Value::String(chat_id.clone()),
serde_json::Value::from,
),
},
"commands": payload["commands"].clone(),
});
let ok = post_set_my_commands(&http_client, &url, &body).await;
let mut cache = cache.lock().unwrap_poison();
cache.insert(
chat_id,
if ok {
ChatMenuState::Registered(Some(payload_str))
} else {
ChatMenuState::Registered(None)
},
);
});
}
pub async fn validate_token(token: &str) -> anyhow::Result<()> {
if token.trim().is_empty() {
anyhow::bail!("Telegram bot token is empty");
}
let url = bot_api_url(token, "getMe");
let client = crate::util::http::build_http_client(std::time::Duration::from_secs(10));
let resp = client
.get(&url)
.send()
.await
.context("Failed to reach Telegram API")?;
let status = resp.status();
let body: serde_json::Value = match resp.json().await {
Ok(v) => v,
Err(e) => anyhow::bail!("Failed to parse Telegram API response: {e}"),
};
if !status.is_success() || body.get("ok").and_then(serde_json::Value::as_bool) != Some(true)
{
let desc = body
.get("description")
.and_then(serde_json::Value::as_str)
.unwrap_or("unknown error");
anyhow::bail!("Invalid Telegram bot token: {desc}");
}
Ok(())
}
fn handle_non_parseable_message(update: &serde_json::Value) {
let Some(message) = update.get("message") else {
return;
};
let text = message
.get("text")
.and_then(serde_json::Value::as_str)
.unwrap_or("<non-text content>");
tracing::debug!("Telegram: message not parseable (unsupported type), skipping: {text}");
}
async fn get_file_path(&self, file_id: &str) -> anyhow::Result<String> {
let url = self.api_url("getFile");
let resp = self
.http_client()
.get(&url)
.query(&[("file_id", file_id)])
.send()
.await
.context("Failed to call Telegram getFile")?;
let data: serde_json::Value = resp.json().await?;
data.get("result")
.and_then(|r| r.get("file_path"))
.and_then(serde_json::Value::as_str)
.map(String::from)
.context("Telegram getFile: missing file_path in response")
}
async fn download_file(&self, file_path: &str) -> anyhow::Result<Vec<u8>> {
let url = format!("{API_BASE}/file/bot{}/{file_path}", self.bot_token);
let resp = self
.http_client()
.get(&url)
.timeout(MEDIA_TRANSFER_TIMEOUT)
.send()
.await
.context("Failed to download Telegram file")?;
if !resp.status().is_success() {
anyhow::bail!("Telegram file download failed: {}", resp.status());
}
anyhow::ensure!(
resp.content_length()
.is_none_or(|len| len <= FILE_MAX_BYTES),
"Telegram file is larger than the {} MB product cap",
mb(FILE_MAX_BYTES)
);
let bytes = resp.bytes().await?;
anyhow::ensure!(
bytes.len() as u64 <= FILE_MAX_BYTES,
"Telegram file is larger than the {} MB product cap",
mb(FILE_MAX_BYTES)
);
Ok(bytes.to_vec())
}
fn parse_attachment_metadata(message: &serde_json::Value) -> Option<IncomingAttachment> {
for (key, kind) in [
("document", IncomingAttachmentKind::Document),
("video", IncomingAttachmentKind::Video),
("video_note", IncomingAttachmentKind::Video),
("animation", IncomingAttachmentKind::Video),
("audio", IncomingAttachmentKind::Audio),
("voice", IncomingAttachmentKind::Voice),
] {
if let Some(v) = message.get(key) {
return Self::build_attachment(v, message, kind);
}
}
if let Some(photos) = message.get("photo").and_then(serde_json::Value::as_array) {
let best = photos.last()?;
return Self::build_attachment(best, message, IncomingAttachmentKind::Photo);
}
None
}
fn build_attachment(
sub_obj: &serde_json::Value,
message: &serde_json::Value,
kind: IncomingAttachmentKind,
) -> Option<IncomingAttachment> {
let file_id = sub_obj.get("file_id")?.as_str()?.to_string();
let file_name = sub_obj
.get("file_name")
.and_then(serde_json::Value::as_str)
.map(String::from);
let file_size = sub_obj.get("file_size").and_then(serde_json::Value::as_u64);
let caption = message
.get("caption")
.and_then(serde_json::Value::as_str)
.map(String::from);
let mime_type = sub_obj
.get("mime_type")
.and_then(serde_json::Value::as_str)
.map(String::from);
Some(IncomingAttachment {
file_id,
file_name,
file_size,
caption,
kind,
mime_type,
})
}
async fn try_parse_attachment_message(
&self,
update: &serde_json::Value,
) -> Option<ChannelMessage> {
let message = update.get("message")?;
let ctx = self.extract_message_context(message).await?;
let attachment = Self::parse_attachment_metadata(message)?;
#[cfg(not(target_os = "macos"))]
if attachment.kind == IncomingAttachmentKind::Voice {
return Some(
self.reject_voice_message(ctx, attachment.caption.as_deref())
.await,
);
}
if let Some(size) = attachment.file_size
&& let Some(reason) = declared_size_refusal(size)
{
tracing::info!("Rejecting attachment: declared size {size} bytes ({reason})");
return Some(
self.reject_attachment(ctx, &attachment, None, NOT_RECEIVED, &reason)
.await,
);
}
let tg_file_path = match self.get_file_path(&attachment.file_id).await {
Ok(p) => p,
Err(e) => {
tracing::warn!("Failed to get attachment file path: {e}");
return Some(
self.reject_attachment(
ctx,
&attachment,
None,
NOT_RECEIVED,
TELEGRAM_LOOKUP_REFUSED_REASON,
)
.await,
);
}
};
let remote_ext = Path::new(&tg_file_path)
.extension()
.and_then(|e| e.to_str());
let file_data = match self.download_file(&tg_file_path).await {
Ok(d) => d,
Err(e) => {
tracing::warn!("Failed to download attachment: {e}");
return Some(
self.reject_attachment(
ctx,
&attachment,
remote_ext,
NOT_RECEIVED,
download_failure_reason(&e),
)
.await,
);
}
};
let staging_dir = crate::util::telegram_staging_dir_name(&ctx.chat_id, ctx.message_id);
let save_dir = crate::util::telegram_files_root().join(&staging_dir);
if let Err(e) = tokio::fs::create_dir_all(&save_dir).await {
tracing::warn!("Failed to create telegram attachment directory: {e}");
return Some(
self.reject_attachment(
ctx,
&attachment,
remote_ext,
NOT_STORED,
ATTACHMENT_DIR_REFUSED_REASON,
)
.await,
);
}
let local_filename =
local_attachment_name(&attachment, &ctx.chat_id, ctx.message_id, remote_ext);
let local_path = save_dir.join(&local_filename);
if let Err(e) = tokio::fs::write(&local_path, &file_data).await {
tracing::warn!("Failed to save attachment to {}: {e}", local_path.display());
let _ = tokio::fs::remove_dir_all(&save_dir).await;
return Some(
self.reject_attachment(
ctx,
&attachment,
remote_ext,
NOT_STORED,
ATTACHMENT_SAVE_REFUSED_REASON,
)
.await,
);
}
let mut content = format_attachment_content(
attachment.kind,
&local_path,
attachment.mime_type.as_deref(),
);
if let Some(caption) = &attachment.caption
&& !caption.is_empty()
{
let _ = write!(content, "\n\n{caption}");
}
let content = Self::prepend_forward_attribution(content, message);
Some(ctx.into_attachment_message(content, staging_dir))
}
async fn reject_attachment(
&self,
ctx: MessageContext,
attachment: &IncomingAttachment,
remote_ext: Option<&str>,
outcome: &str,
reason: &str,
) -> ChannelMessage {
let display_name =
local_attachment_name(attachment, &ctx.chat_id, ctx.message_id, remote_ext);
let notice = format!("Could not process the file \"{display_name}\": {reason}.");
let content = attachment_rejection_content(
Some(&display_name),
outcome,
reason,
attachment.caption.as_deref(),
);
self.refuse_with_notice(ctx, ¬ice, content).await
}
async fn refuse_with_notice(
&self,
ctx: MessageContext,
notice: &str,
content: String,
) -> ChannelMessage {
let (_, thread_id) = parse_recipient(&ctx.reply_target);
if let Err(e) = self
.send_text_chunks(notice, &ctx.chat_id, thread_id, None)
.await
{
tracing::warn!(%notice, "Failed to send refusal notice: {e}");
}
ctx.into_channel_message(content, None)
}
#[cfg(not(target_os = "macos"))]
async fn reject_voice_message(
&self,
ctx: MessageContext,
caption: Option<&str>,
) -> ChannelMessage {
let notice = format!("Could not process the voice message: {VOICE_UNSUPPORTED_REASON}.");
let content =
attachment_rejection_content(None, NOT_RECEIVED, VOICE_UNSUPPORTED_REASON, caption);
self.refuse_with_notice(ctx, ¬ice, content).await
}
fn format_forward_attribution(message: &serde_json::Value) -> Option<String> {
if let Some(from_chat) = message.get("forward_from_chat") {
let title = from_chat
.get("title")
.and_then(serde_json::Value::as_str)
.unwrap_or("unknown channel");
Some(format!("{FORWARD_ATTRIBUTION_PREFIX}channel: {title}] "))
} else if let Some(from_user) = message.get("forward_from") {
let label = format_sender_label(from_user);
Some(format!("{FORWARD_ATTRIBUTION_PREFIX}{label}] "))
} else {
message
.get("forward_sender_name")
.and_then(serde_json::Value::as_str)
.map(|name| format!("{FORWARD_ATTRIBUTION_PREFIX}{name}] "))
}
}
async fn parse_update_message(&self, update: &serde_json::Value) -> Option<ChannelMessage> {
let message = update.get("message")?;
let text = message.get("text").and_then(serde_json::Value::as_str)?;
let ctx = self.extract_message_context(message).await?;
let text = if text.starts_with('/') {
text.split('@').next().unwrap_or(text)
} else {
text
};
let content = text.to_string();
let content = Self::prepend_forward_attribution(content, message);
Some(ctx.into_channel_message(content, None))
}
async fn post_telegram_json(
&self,
method: &str,
body: serde_json::Value,
err_label: &str,
) -> Result<reqwest::Response, (reqwest::StatusCode, String)> {
let resp = self
.http_client()
.post(self.api_url(method))
.json(&body)
.send()
.await
.map_err(|e| (reqwest::StatusCode::BAD_GATEWAY, e.to_string()))?;
let status = resp.status();
if status.is_success() {
Ok(resp)
} else {
let err_body = crate::util::http::read_error_body(resp, err_label).await;
Err((status, err_body))
}
}
async fn send_message_get_id(
&self,
chat_id: &str,
thread_id: Option<&str>,
text: &str,
parse_mode: Option<&str>,
reply_markup: Option<serde_json::Value>,
) -> Result<Option<i64>, (reqwest::StatusCode, String)> {
let mut body = serde_json::json!({
"chat_id": chat_id,
"text": text,
});
if let Some(mode) = parse_mode {
body["parse_mode"] = serde_json::Value::String(mode.to_string());
}
set_thread_id_on_json(&mut body, thread_id);
if let Some(markup) = reply_markup {
body["reply_markup"] = markup;
}
body["link_preview_options"] = serde_json::json!({ "is_disabled": true });
let resp = self
.post_telegram_json("sendMessage", body, "sendMessage error")
.await?;
let status = resp.status();
let id = match resp.json::<serde_json::Value>().await {
Ok(body) => body["result"]["message_id"].as_i64(),
Err(e) => {
tracing::warn!(
status = ?status,
error = %e,
"sendMessage: unparseable response body — treating as delivered"
);
None
}
};
Ok(id)
}
async fn send_single_message(
&self,
chat_id: &str,
thread_id: Option<&str>,
text: &str,
parse_mode: Option<&str>,
reply_markup: Option<serde_json::Value>,
) -> Result<(), (reqwest::StatusCode, String)> {
self.send_message_get_id(chat_id, thread_id, text, parse_mode, reply_markup)
.await
.map(|_| ())
}
async fn send_text_chunks(
&self,
message: &str,
chat_id: &str,
thread_id: Option<&str>,
reply_markup: Option<serde_json::Value>,
) -> anyhow::Result<()> {
let html = to_telegram_html(message);
let chunks = split_message_for_telegram(&html);
for (index, chunk) in chunks.iter().enumerate() {
let text = wrap_chunk(chunk, index, chunks.len());
let chunk_reply_markup = if index == chunks.len() - 1 {
reply_markup.clone()
} else {
None
};
if let Err((html_status, html_err)) = self
.send_single_message(
chat_id,
thread_id,
&text,
Some("HTML"),
chunk_reply_markup.clone(),
)
.await
{
tracing::info!(
status = ?html_status,
"Telegram sendMessage with HTML parse_mode failed; retrying without parse_mode"
);
let clean_text = strip_html_tags(&text);
self.send_single_message(chat_id, thread_id, &clean_text, None, chunk_reply_markup)
.await
.map_err(|(plain_status, plain_err)| {
anyhow::anyhow!(
"Telegram sendMessage failed (html {html_status}: {html_err}; plain {plain_status}: {plain_err})"
)
})?;
}
if index < chunks.len() - 1 {
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
Ok(())
}
pub async fn edit_reply_markup(
&self,
chat_id: &str,
message_id: i64,
reply_markup: &serde_json::Value,
) -> Result<(), EditMessageFailure> {
self.post_telegram_json(
"editMessageReplyMarkup",
serde_json::json!({
"chat_id": chat_id,
"message_id": message_id,
"reply_markup": reply_markup,
}),
"editMessageReplyMarkup error",
)
.await
.map(|_| ())
.map_err(|(_, desc)| classify_edit_failure(&desc))
}
async fn send_attachment(
&self,
chat_id: &str,
thread_id: Option<&str>,
attachment: &TelegramAttachment,
) -> anyhow::Result<()> {
let target = attachment.target.trim();
if is_http_url(target) {
let result = self
.send_media_by_url(chat_id, thread_id, attachment.kind, target)
.await;
if let Err(e) = result {
tracing::warn!(
url = target,
error = %e,
"Telegram send media by URL failed; falling back to text link"
);
let fallback_text = format!("{}: {target}", attachment.kind.meta().label);
self.send_text_chunks(&fallback_text, chat_id, thread_id, None)
.await?;
}
return Ok(());
}
let path = Path::new(&target);
if !path.is_file() {
anyhow::bail!("Telegram attachment target is not an existing regular file: {target}");
}
self.send_media_file(chat_id, thread_id, attachment.kind, path)
.await
}
async fn send_media(
&self,
chat_id: &str,
api_method: &'static str,
request: reqwest::RequestBuilder,
label: &str,
) -> anyhow::Result<()> {
let resp = request.send().await?;
if !resp.status().is_success() {
let err = resp.text().await?;
anyhow::bail!("Telegram {api_method} failed: {err}");
}
tracing::debug!("Telegram {api_method} sent to {chat_id}: {label}");
Ok(())
}
async fn send_media_file(
&self,
chat_id: &str,
thread_id: Option<&str>,
kind: TelegramAttachmentKind,
file_path: &Path,
) -> anyhow::Result<()> {
let meta = kind.file_meta();
let file_name = file_path
.file_name()
.and_then(|n| n.to_str())
.unwrap_or(meta.default_filename);
let file_bytes = tokio::fs::read(file_path).await?;
let part = Part::bytes(file_bytes).file_name(file_name.to_string());
let mut form = Form::new()
.text("chat_id", chat_id.to_string())
.part(meta.form_field, part);
if meta.disable_content_type_detection {
form = form.text("disable_content_type_detection", "true");
}
if let Some(tid) = thread_id {
form = form.text("message_thread_id", tid.to_string());
}
let request = self
.http_client()
.post(self.api_url(meta.api_method))
.multipart(form)
.timeout(MEDIA_TRANSFER_TIMEOUT);
self.send_media(chat_id, meta.api_method, request, file_name)
.await
}
async fn send_media_by_url(
&self,
chat_id: &str,
thread_id: Option<&str>,
kind: TelegramAttachmentKind,
url: &str,
) -> anyhow::Result<()> {
let meta = kind.meta();
let mut body = serde_json::json!({ "chat_id": chat_id });
body[meta.form_field] = serde_json::Value::String(url.to_string());
set_thread_id_on_json(&mut body, thread_id);
let request = self
.http_client()
.post(self.api_url(meta.api_method))
.json(&body);
self.send_media(chat_id, meta.api_method, request, url)
.await
}
}
enum PollOutcome {
Updates(Vec<serde_json::Value>),
Conflict,
Error(String),
Transport,
}
impl TelegramChannel {
async fn poll_get_updates(
&self,
offset: &mut i64,
timeout: u64,
ok_default: bool,
) -> PollOutcome {
let url = self.api_url("getUpdates");
let body = serde_json::json!({
"offset": *offset,
"timeout": timeout,
"allowed_updates": ["message", "callback_query"]
});
let resp = match self.http_client().post(&url).json(&body).send().await {
Ok(r) => r,
Err(e) => {
tracing::info!("Telegram poll error: {e}");
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
return PollOutcome::Transport;
}
};
let data: serde_json::Value = match resp.json().await {
Ok(d) => d,
Err(e) => {
tracing::warn!("Telegram parse error: {e}");
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
return PollOutcome::Transport;
}
};
let ok = data
.get("ok")
.and_then(serde_json::Value::as_bool)
.unwrap_or(ok_default);
if ok {
if let Some(results) = data.get("result").and_then(serde_json::Value::as_array) {
for update in results {
if let Some(uid) = update.get("update_id").and_then(serde_json::Value::as_i64) {
*offset = (*offset).max(uid + 1);
}
}
return PollOutcome::Updates(results.clone());
}
return PollOutcome::Updates(Vec::new());
}
let error_code = data
.get("error_code")
.and_then(serde_json::Value::as_i64)
.unwrap_or_default();
if error_code == 409 {
PollOutcome::Conflict
} else {
let desc = data
.get("description")
.and_then(serde_json::Value::as_str)
.unwrap_or("unknown Telegram API error");
PollOutcome::Error(desc.to_string())
}
}
async fn probe_startup_slot(&self, offset: &mut i64) -> bool {
loop {
if self.cancel.is_cancelled() {
tracing::info!("Telegram channel cancelled during startup probe");
return false;
}
match self.poll_get_updates(offset, 0, false).await {
PollOutcome::Updates(_) => return true,
PollOutcome::Conflict => {
tracing::debug!("Startup probe: slot busy (409), retrying in 5s");
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
}
PollOutcome::Error(desc) => {
tracing::warn!("Startup probe: API error: {desc}; retrying in 5s");
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
}
PollOutcome::Transport => {} }
}
}
async fn parse_polled_update(&self, update: &serde_json::Value) -> Option<ChannelMessage> {
let msg = if let Some(msg) = self.parse_update_message(update).await {
msg
} else if let Some(msg) = self.try_parse_attachment_message(update).await {
msg
} else {
Self::handle_non_parseable_message(update);
return None;
};
if let Some((reaction_chat_id, reaction_message_id)) =
Self::extract_update_message_target(update)
{
self.try_add_ack_reaction_nonblocking(reaction_chat_id, reaction_message_id);
}
Some(msg)
}
async fn process_updates(
&self,
tx: &tokio::sync::mpsc::Sender<ChannelMessage>,
updates: Vec<serde_json::Value>,
) -> bool {
let plan = poll_plan(&updates);
for (update, placement) in updates.iter().zip(plan) {
match placement {
PollPlacement::Callback => {
if !self.handle_callback_press(tx, update).await {
return false;
}
}
PollPlacement::Skip => {}
PollPlacement::Single => {
if let Some(msg) = self.parse_polled_update(update).await
&& tx.send(msg).await.is_err()
{
return false;
}
}
PollPlacement::Album(members) => {
let mut album = None;
for member in members {
if let Some(msg) = self.parse_polled_update(&updates[member]).await {
album = Some(merge_album_member(album, msg));
}
}
if let Some(album) = album
&& tx.send(album).await.is_err()
{
return false;
}
}
}
}
true
}
async fn handle_callback_press(
&self,
tx: &tokio::sync::mpsc::Sender<ChannelMessage>,
update: &serde_json::Value,
) -> bool {
let cq = &update["callback_query"];
let cq_id = cq["id"].as_str().map(ToString::to_string);
let cq_data = cq["data"].as_str().unwrap_or("");
if !cq_data.starts_with(ACTION_PREFIX)
&& let Some(ref id) = cq_id
{
self.answer_callback_query(id, None).await;
}
let Some(msg) = self.parse_callback_query(cq).await else {
return true;
};
tx.send(msg).await.is_ok()
}
}
fn merge_album_member(album: Option<ChannelMessage>, member: ChannelMessage) -> ChannelMessage {
match album {
Some(mut album) => {
album.content.push('\n');
album.content.push_str(&member.content);
album.attachment_dirs.extend(member.attachment_dirs);
album
}
None => member,
}
}
#[derive(Debug, PartialEq, Eq)]
enum PollPlacement {
Single,
Album(Vec<usize>),
Skip,
Callback,
}
fn poll_plan(updates: &[serde_json::Value]) -> Vec<PollPlacement> {
let mut plan: Vec<PollPlacement> = updates
.iter()
.map(|update| match update.get("callback_query") {
Some(_) => PollPlacement::Callback,
None => PollPlacement::Single,
})
.collect();
let mut albums: HashMap<&str, Vec<usize>> = HashMap::new();
for (at, update) in updates.iter().enumerate() {
if plan[at] == PollPlacement::Callback {
continue;
}
if let Some(group) = update
.get("message")
.and_then(|message| message.get("media_group_id"))
.and_then(serde_json::Value::as_str)
{
albums.entry(group).or_default().push(at);
}
}
for members in albums.into_values() {
let (first, rest) = members.split_first().expect("an album has a member");
let first = *first;
for &later in rest {
plan[later] = PollPlacement::Skip;
}
plan[first] = PollPlacement::Album(members);
}
plan
}
#[async_trait]
impl Channel for TelegramChannel {
fn name(&self) -> &'static str {
"telegram"
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
async fn send(&self, message: &SendMessage) -> anyhow::Result<()> {
let content = message.content.trim();
if content.is_empty() {
tracing::warn!("TelegramChannel: attempted to send empty message – skipping");
return Ok(()); }
let (chat_id, thread_id) = parse_recipient(&message.recipient);
self.spawn_menu_refresh(chat_id);
let (text_without_markers, attachments, refusals) =
crate::util::with_block_in_place(|| {
parse_attachment_markers(content, &message.file_roots)
});
if attachments.is_empty() && refusals.is_empty() {
return self
.send_text_chunks(content, chat_id, thread_id, message.reply_markup.clone())
.await;
}
let body = if refusals.is_empty() {
text_without_markers
} else {
let notices = refusals.join("\n\n");
if text_without_markers.is_empty() {
notices
} else {
format!("{text_without_markers}\n\n{notices}")
}
};
if !body.is_empty() {
self.send_text_chunks(&body, chat_id, thread_id, message.reply_markup.clone())
.await?;
}
let mut first_error: Option<anyhow::Error> = None;
for attachment in &attachments {
if let Err(e) = self.send_attachment(chat_id, thread_id, attachment).await {
tracing::warn!(
error = %e,
file = %attachment.target,
"Telegram: attachment delivery failed"
);
let notice = file_refusal(
file_name_or_path(&attachment.target),
ATTACHMENT_UPLOAD_FAILED_REASON,
);
if let Err(notice_err) = self
.send_text_chunks(¬ice, chat_id, thread_id, None)
.await
{
tracing::warn!(
error = %notice_err,
"Telegram: failed to report attachment failure to the user"
);
}
first_error.get_or_insert(e);
}
}
match first_error {
Some(e) => Err(e),
None => Ok(()),
}
}
async fn listen(&self, tx: tokio::sync::mpsc::Sender<ChannelMessage>) -> anyhow::Result<()> {
use std::sync::atomic::Ordering;
let mut offset = self.offset.load(Ordering::Acquire);
if offset > 0 {
tracing::info!(offset, "Telegram channel resuming from previous offset");
}
tracing::info!("Telegram channel listening for messages...");
if !self.probe_startup_slot(&mut offset).await {
return Ok(());
}
tracing::debug!("Startup probe succeeded; entering main long-poll loop.");
let (parts_tx, parts_rx) =
tokio::sync::mpsc::channel(crate::channels::telegram_group::QUEUE_CAPACITY);
crate::channels::telegram_group::spawn_collector(parts_rx, tx);
let shutdown_token = crate::shutdown::shutdown_token();
let per_channel_cancel = self.cancel.clone();
loop {
tokio::select! {
() = shutdown_token.cancelled() => {
tracing::info!("Telegram channel shutting down (global shutdown)");
self.offset.store(offset, Ordering::Release);
return Ok(());
}
() = per_channel_cancel.cancelled() => {
tracing::info!("Telegram channel shutting down (token hot-reload)");
self.offset.store(offset, Ordering::Release);
return Ok(());
}
poll_result = self.poll_get_updates(&mut offset, 30, true) => {
self.offset.store(offset, Ordering::Release);
let updates = match poll_result {
PollOutcome::Updates(updates) => updates,
PollOutcome::Conflict => {
tracing::warn!(
"Telegram polling conflict (409). \
Ensure only one `mahbot` process is using this bot token."
);
tokio::time::sleep(std::time::Duration::from_secs(35)).await;
continue;
}
PollOutcome::Error(desc) => {
tracing::info!("Telegram getUpdates API error: {desc}");
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
continue;
}
PollOutcome::Transport => continue,
};
if !self.process_updates(&parts_tx, updates).await {
tracing::warn!(
"Telegram listener stopping: the message pipeline is closed"
);
return Ok(());
}
}
}
}
}
async fn start_typing(&self, recipient: &str) -> anyhow::Result<()> {
let url = self.api_url("sendChatAction");
let (chat_id, thread_id) = parse_recipient(recipient);
let mut body = serde_json::json!({
"chat_id": chat_id,
"action": "typing"
});
set_thread_id_on_json(&mut body, thread_id);
self.http_client().post(&url).json(&body).send().await?;
Ok(())
}
}
fn cancel_old_listener(old: Option<&std::sync::Arc<dyn Channel>>) {
if let Some(old) = old
&& let Some(tc) = old.as_any().downcast_ref::<TelegramChannel>()
{
tc.cancel_own();
}
}
pub async fn restart_telegram_listener(new_token: Option<&str>) -> anyhow::Result<()> {
let registry = crate::channel_registry();
let old_channel = registry.get("telegram");
if let Some(token) = new_token.filter(|t| !t.trim().is_empty()) {
TelegramChannel::validate_token(token).await?;
let offset = old_channel
.as_ref()
.and_then(|c| c.as_any().downcast_ref::<TelegramChannel>())
.map(|tc| std::sync::Arc::clone(&tc.offset));
let telegram_channel = if let Some(inherited_offset) = offset {
TelegramChannel::with_offset(token.to_string(), inherited_offset)
} else {
TelegramChannel::new(token.to_string())
};
let telegram_arc = std::sync::Arc::new(telegram_channel);
tokio::spawn({
let tc = std::sync::Arc::clone(&telegram_arc);
async move {
tc.set_my_commands().await;
}
});
let new_channel: std::sync::Arc<dyn Channel> = telegram_arc;
registry.replace(std::sync::Arc::clone(&new_channel));
cancel_old_listener(old_channel.as_ref());
if let Some(tx) = crate::MESSAGE_TX.get() {
let tx = tx.clone();
tokio::spawn(async move {
if let Err(e) = new_channel.listen(tx).await {
tracing::error!(error = %e, "Telegram listener error after hot-reload");
}
});
} else {
tracing::error!("MESSAGE_TX not set — cannot spawn Telegram listener");
}
tracing::info!("Telegram bot listener restarted with new token");
} else {
cancel_old_listener(old_channel.as_ref());
registry.unregister("telegram");
tracing::info!("Telegram bot token cleared — listener stopped");
}
Ok(())
}
pub async fn send_direct(
recipient: &str,
content: String,
reply_markup: Option<serde_json::Value>,
) -> anyhow::Result<()> {
let Some(channel) = crate::channel_registry().get("telegram") else {
anyhow::bail!("Telegram channel not found in registry");
};
let reply = SendMessage {
content,
recipient: recipient.to_string(),
reply_markup,
file_roots: Vec::new(),
};
channel.send(&reply).await
}
pub async fn send_reply(recipient: &str, content: &str) {
let _ = send_direct(recipient, content.to_string(), None).await;
}
pub async fn mirror_gui_message_to_telegram(msg: &ChannelMessage) {
if msg.channel != "gui" && msg.channel != "voice" {
return;
}
let trimmed = msg.content.trim();
if trimmed.is_empty() {
return;
}
let Some(channel) = crate::channel_registry().get("telegram") else {
return;
};
let bindings = match crate::users::store()
.get_user_channels(&msg.user_name)
.await
{
Ok(b) => b,
Err(e) => {
tracing::warn!(
user = %msg.user_name,
error = %e,
"Failed to look up user channels for message mirror"
);
return;
}
};
let telegram_bindings: Vec<_> = bindings
.into_iter()
.filter(|b| b.channel == "telegram")
.collect();
if telegram_bindings.is_empty() {
return; }
let content = TELEGRAM_MEDIA_MARKER_RE
.replace_all(trimmed, "")
.to_string();
let content = content.trim().to_string();
let quoted = if let Some(reply) = &msg.reply_reference {
let header = format!("↩ {}: {}", reply.author, reply.snippet);
if content.is_empty() {
format!("<blockquote>\n{header}\n</blockquote>")
} else {
format!("<blockquote>\n{header}\n{content}\n</blockquote>")
}
} else {
if content.is_empty() {
return; }
format!("<blockquote>\n{content}\n</blockquote>")
};
for binding in &telegram_bindings {
let Some(reply_target) = &binding.reply_target else {
continue; };
let reply = SendMessage {
content: quoted.clone(),
recipient: reply_target.clone(),
reply_markup: None,
file_roots: Vec::new(),
};
if let Err(e) = channel.send(&reply).await {
tracing::error!(
user = %msg.user_name,
recipient = %reply_target,
error = %e,
"Failed to mirror local message to Telegram"
);
}
}
}
async fn post_set_my_commands(
client: &reqwest::Client,
url: &str,
body: &serde_json::Value,
) -> bool {
match client.post(url).json(body).send().await {
Ok(resp) if resp.status().is_success() => {
tracing::debug!("Telegram bot commands registered successfully");
true
}
Ok(resp) => {
let status = resp.status();
let body = resp.text().await.unwrap_or_default();
tracing::warn!(
status = %status,
body = %body,
"Telegram setMyCommands returned unsuccessful status"
);
false
}
Err(e) => {
tracing::warn!(error = %e, "Failed to call Telegram setMyCommands");
false
}
}
}
fn phase_emoji(phase: TicketPhase) -> &'static str {
match phase {
TicketPhase::Backlog => "📥",
TicketPhase::Analysis => "🔍",
TicketPhase::Planning => "🧭",
TicketPhase::Queued => "⏳",
TicketPhase::InDevelopment => "🔨",
TicketPhase::InSanitation => "🧹",
TicketPhase::Verification => "🎯",
TicketPhase::Done => "✅",
TicketPhase::Cancelled => "🚫",
TicketPhase::Failed => "❌",
}
}
#[must_use]
pub fn format_board_line(phase: &TicketPhase, id: &str, title: &str) -> String {
format!("{} `{}` {}", phase_emoji(*phase), id, title)
}
#[must_use]
pub fn board_listing_text(tickets: &[&Ticket]) -> String {
if tickets.is_empty() {
return "All workspaces — no tickets".to_string();
}
let listing = tickets
.iter()
.map(|t| format_board_line(&t.phase, &t.id, &t.title))
.collect::<Vec<_>>()
.join("\n");
let noun = if tickets.len() == 1 {
"ticket"
} else {
"tickets"
};
format!("All workspaces — {} {noun}\n{listing}", tickets.len())
}
#[must_use]
fn workspace_picker_label(ws: &Workspace, active: bool) -> String {
let name = if active {
format!("\u{2713} {}", ws.display_name())
} else {
ws.display_name()
};
let mut label = format!("{name} — {}", ws.status);
if ws.paused {
label.push_str(", paused");
}
label
}
#[must_use]
pub fn workspace_picker_keyboard(
workspaces: &[Workspace],
active: Option<&str>,
) -> serde_json::Value {
let rows = workspaces
.iter()
.map(|ws| {
serde_json::json!([{
"text": workspace_picker_label(ws, active == Some(ws.name.as_str())),
"callback_data": format!("{ACTION_PREFIX}set_workspace|{}", ws.name),
}])
})
.collect::<Vec<_>>();
serde_json::json!({ "inline_keyboard": rows })
}
#[must_use]
pub async fn user_command_entries(user_name: &str) -> Vec<(String, String)> {
let mut entries: Vec<(String, String)> = Vec::new();
if crate::users::is_admin(user_name).await {
if crate::users::workspace_switcher_available().await {
entries.push(("workspace".to_string(), WORKSPACE_COMMAND_DESC.to_string()));
}
entries.push(("board".to_string(), BOARD_COMMAND_DESC.to_string()));
entries.push(("archive".to_string(), ARCHIVE_COMMAND_DESC.to_string()));
if crate::self_update::should_show_update(crate::self_update::update_availability()) {
entries.push(("update".to_string(), UPDATE_COMMAND_DESC.to_string()));
}
if let Ok(Some(ws_name)) = crate::users::get_raw_selected_workspace(user_name).await
&& let Ok(Some(ws)) = crate::workspace::get_by_name(&ws_name).await
{
if ws.paused {
entries.push(("unpause".to_string(), UNPAUSE_COMMAND_DESC.to_string()));
} else {
entries.push(("pause".to_string(), PAUSE_COMMAND_DESC.to_string()));
}
if ws.maintenance_enabled {
entries.push((
"maintenance_off".to_string(),
MAINTENANCE_OFF_COMMAND_DESC.to_string(),
));
} else {
entries.push((
"maintenance_on".to_string(),
MAINTENANCE_ON_COMMAND_DESC.to_string(),
));
}
}
}
entries.push((
"image_models".to_string(),
IMAGE_MODELS_COMMAND_DESC.to_string(),
));
entries.push((
"video_models".to_string(),
VIDEO_MODELS_COMMAND_DESC.to_string(),
));
entries.push(("clear".to_string(), CLEAR_COMMAND_DESC.to_string()));
entries
}
#[cfg(test)]
#[path = "telegram_tests.rs"]
mod tests;