use super::*;
pub(super) async fn poll_telegram(bot: Bot, state: AppState) -> anyhow::Result<()> {
let dispatcher = update_dispatch::UpdateDispatcher::new(bot.clone(), state.clone());
loop {
let next_update_id = state.transport.polling_offset()?;
let request_offset = i32::try_from(next_update_id)
.context("durable Telegram polling offset exceeds Bot API range")?;
let updates = match bot
.get_updates()
.offset(request_offset)
.timeout(TELEGRAM_POLL_TIMEOUT_SECONDS)
.allowed_updates(vec![
AllowedUpdate::Message,
AllowedUpdate::EditedMessage,
AllowedUpdate::MyChatMember,
AllowedUpdate::ChatMember,
])
.send()
.await
{
Ok(updates) => updates,
Err(error) => {
tracing::debug!(
error_class = telegram_requests::request_error_class(&error),
"Telegram poll retry"
);
tokio::time::sleep(Duration::from_secs(2)).await;
continue;
}
};
let dispatch = dispatcher.clone();
if let Err(error) = process_polled_updates(&state, updates, move |update| {
let dispatch = dispatch.clone();
async move { dispatch.enqueue(update).await }
})
.await
{
tracing::warn!(
error_class = telegram_requests::anyhow_error_class(&error),
"Telegram cursor persistence failed; polling will retry"
);
tokio::time::sleep(Duration::from_secs(2)).await;
}
}
}
async fn process_polled_updates<F, Fut>(
state: &AppState,
mut updates: Vec<Update>,
mut process: F,
) -> anyhow::Result<()>
where
F: FnMut(Update) -> Fut,
Fut: std::future::Future<Output = anyhow::Result<()>>,
{
updates.sort_by_key(|update| update.id.0);
let mut next_update_id = state.transport.polling_offset()?;
for update in updates {
let update_id = i64::from(update.id.0);
if update_id < next_update_id {
continue;
}
if let Err(error) = process(update).await {
tracing::warn!(
update_id,
error_class = telegram_requests::anyhow_error_class(&error),
"Telegram update dispatch failed; advancing the lossy transport cursor"
);
}
next_update_id = state.transport.advance_polling_offset(update_id)?;
}
Ok(())
}
pub(super) async fn process_update(
bot: &Bot,
state: &AppState,
update: Update,
) -> anyhow::Result<()> {
let update_id = i64::from(update.id.0);
match update.kind {
UpdateKind::Message(message) => {
if message.chat.is_private() {
process_private_message(bot, state, update_id, message).await
} else if message.chat.is_group() || message.chat.is_supergroup() {
process_group_message(bot, state, update_id, message, false).await
} else {
Ok(())
}
}
UpdateKind::EditedMessage(message) => {
if message.chat.is_private() {
process_private_message_edit(bot, state, update_id, message).await
} else if message.chat.is_group() || message.chat.is_supergroup() {
process_group_message(bot, state, update_id, message, true).await
} else {
Ok(())
}
}
UpdateKind::ChatMember(change) | UpdateKind::MyChatMember(change) => {
process_group_membership(state, change)
}
_ => Ok(()),
}
}
struct MessageInput {
kind: &'static str,
text: Option<String>,
media_bytes: Option<Vec<u8>>,
mime_type: Option<String>,
file_name: Option<String>,
duration_seconds: Option<i64>,
}
impl From<MessageInput> for MessageContent {
fn from(input: MessageInput) -> Self {
Self {
kind: input.kind.into(),
text: input.text,
media_bytes: input.media_bytes,
mime_type: input.mime_type,
file_name: input.file_name,
duration_seconds: input.duration_seconds,
}
}
}
fn reset_command(message: &Message) -> bool {
message.text().is_some_and(|text| {
text.split_whitespace().next().is_some_and(|command| {
command.eq_ignore_ascii_case("/reset")
|| command.to_ascii_lowercase().starts_with("/reset@")
})
})
}
async fn download_message_file(
bot: &Bot,
chat_id: ChatId,
file_id: teloxide::types::FileId,
expected_size: Option<u32>,
maximum_bytes: usize,
label: &str,
notify_errors: bool,
) -> anyhow::Result<Option<Vec<u8>>> {
if expected_size.is_some_and(|size| u64::from(size) > maximum_bytes as u64) {
if notify_errors {
send_telegram_message(
bot,
chat_id,
format!("That {label} is too large for Kennedy to process."),
)
.await?;
}
return Ok(None);
}
let file =
telegram_requests::retry_request("get_file", || bot.get_file(file_id.clone()).send())
.await?;
let bytes = telegram_requests::retry_download("download_file", || {
let file_path = file.path.clone();
async move {
let mut stream = bot.download_file_stream(&file_path);
let mut bytes =
Vec::with_capacity(expected_size.map(|size| size as usize).unwrap_or(0));
while let Some(chunk) = stream.next().await {
let chunk = chunk?;
if bytes.len().saturating_add(chunk.len()) > maximum_bytes {
return Ok(None);
}
bytes.extend_from_slice(&chunk);
}
Ok(Some(bytes))
}
})
.await?;
if bytes.is_none() && notify_errors {
send_telegram_message(
bot,
chat_id,
format!("That {label} is too large for Kennedy to process."),
)
.await?;
}
Ok(bytes)
}
async fn parse_message_input(
bot: &Bot,
state: &AppState,
message: &Message,
) -> anyhow::Result<Option<MessageInput>> {
parse_message_input_with_feedback(bot, state, message, true).await
}
async fn parse_message_input_with_feedback(
bot: &Bot,
state: &AppState,
message: &Message,
feedback: bool,
) -> anyhow::Result<Option<MessageInput>> {
if let Some(text) = message.text() {
return Ok(Some(if reset_command(message) {
MessageInput {
kind: "reset",
text: None,
media_bytes: None,
mime_type: None,
file_name: None,
duration_seconds: None,
}
} else {
MessageInput {
kind: "text",
text: Some(text.to_owned()),
media_bytes: None,
mime_type: None,
file_name: None,
duration_seconds: None,
}
}));
}
if let Some(media) = native_media::classify_message(message) {
let Some(bytes) = download_message_file(
bot,
message.chat.id,
media.file_id,
media.declared_size,
state.max_voice_bytes,
media.label,
feedback,
)
.await?
else {
return Ok(None);
};
return Ok(Some(MessageInput {
kind: media.kind,
text: media.text,
media_bytes: Some(bytes),
mime_type: media.mime_type,
file_name: media.file_name,
duration_seconds: media.duration_seconds,
}));
}
if feedback {
send_telegram_message(
bot,
message.chat.id,
"Kennedy accepts text, voice notes, native Telegram media, and bounded files here. Use /reset to end this Telegram session.",
)
.await?;
}
Ok(None)
}
fn group_message_text(message: &Message) -> String {
if message.photo().is_some()
|| message.video().is_some()
|| message.animation().is_some()
|| message.audio().is_some()
{
return message.caption().unwrap_or("").to_owned();
}
if let Some(sticker) = message.sticker() {
return sticker.emoji.clone().unwrap_or_default();
}
if message.video_note().is_some() {
return String::new();
}
if message.voice().is_some() {
return "[Voice note]".into();
}
if let Some(document) = message.document() {
let label = format!(
"[File: {}]",
document.file_name.as_deref().unwrap_or("telegram-file")
);
return message
.caption()
.map(|caption| format!("{label} {caption}"))
.unwrap_or(label);
}
if let Some(text) = message.text().or_else(|| message.caption()) {
return text.to_owned();
}
"[Non-text Telegram message]".into()
}
fn report_identity(
sink: &dyn IdentitySink,
telegram_user_id: i64,
username: Option<&str>,
display_name: &str,
) -> anyhow::Result<bool> {
sink.observe_identity(&IdentityObservation {
telegram_user_id,
username: username.map(ToOwned::to_owned),
display_name: display_name.to_owned(),
})?;
Ok(sink.whitelist()?.contains(telegram_user_id))
}
fn identity(user: &teloxide::types::User) -> anyhow::Result<Identity> {
Ok(Identity {
telegram_user_id: i64::try_from(user.id.0)
.context("Telegram user ID exceeds SQLite range")?,
username: user.username.clone(),
display_name: user.full_name(),
})
}
async fn process_private_message(
bot: &Bot,
state: &AppState,
update_id: i64,
message: Message,
) -> anyhow::Result<()> {
let Some(user) = message.from.as_ref() else {
return Ok(());
};
let identity = identity(user)?;
let authorized = report_identity(
state.identity_sink.as_ref(),
identity.telegram_user_id,
identity.username.as_deref(),
&identity.display_name,
)?;
if !authorized {
send_telegram_message(bot, message.chat.id, UNAUTHORIZED_MESSAGE).await?;
return Ok(());
}
if let Some(text) = message.text()
&& text.split_whitespace().next().is_some_and(|command| {
command.eq_ignore_ascii_case("/adduser")
|| command.to_ascii_lowercase().starts_with("/adduser@")
})
{
let Some(handle) = text.split_whitespace().nth(1) else {
send_telegram_message(bot, message.chat.id, "Usage: /adduser @theirHandle").await?;
return Ok(());
};
let status = match state
.identity_sink
.request_add_user(identity.telegram_user_id, handle)?
{
AddUserOutcome::Forbidden => {
"Only the Kennedy administrator can use /adduser.".to_owned()
}
AddUserOutcome::Whitelisted {
handle,
telegram_user_id: Some(id),
} => format!("Whitelisted @{handle} and pinned Telegram user ID {id}."),
AddUserOutcome::Whitelisted {
handle,
telegram_user_id: None,
} => format!(
"Whitelisted @{handle}. Kennedy will pin its numeric Telegram user ID by TOFU the first time that handle is observed."
),
};
send_telegram_message(bot, message.chat.id, status).await?;
return Ok(());
}
state
.transport
.observe_private_chat(identity.telegram_user_id, message.chat.id.0)?;
let Some(input) = parse_message_input(bot, state, &message).await? else {
return Ok(());
};
state.transport.accept_private_message(AcceptedMessage {
update_id,
message_id: i64::from(message.id.0),
chat_id: message.chat.id.0,
identity,
content: input.into(),
created_at: message.date.to_rfc3339(),
edited: false,
})?;
Ok(())
}
async fn process_private_message_edit(
bot: &Bot,
state: &AppState,
update_id: i64,
message: Message,
) -> anyhow::Result<()> {
let Some(user) = message.from.as_ref() else {
return Ok(());
};
let identity = identity(user)?;
if !report_identity(
state.identity_sink.as_ref(),
identity.telegram_user_id,
identity.username.as_deref(),
&identity.display_name,
)? {
return Ok(());
}
let chat_id = message.chat.id.0;
let message_id = i64::from(message.id.0);
if !state.transport.has_private_message(chat_id, message_id)? {
return Ok(());
}
let content = parse_message_input_with_feedback(bot, state, &message, false)
.await?
.map(Into::into);
state.transport.revise_private_message(MessageRevision {
update_id,
message_id,
chat_id,
identity,
content,
})?;
Ok(())
}
fn member_status(kind: &ChatMemberKind) -> (&'static str, bool) {
match kind {
ChatMemberKind::Owner(_) => ("creator", true),
ChatMemberKind::Administrator(_) => ("administrator", true),
ChatMemberKind::Member(_) => ("member", true),
ChatMemberKind::Restricted(member) if member.is_member => ("member", true),
ChatMemberKind::Restricted(_) | ChatMemberKind::Left => ("left", false),
ChatMemberKind::Banned(_) => ("kicked", false),
}
}
fn process_group_membership(
state: &AppState,
change: teloxide::types::ChatMemberUpdated,
) -> anyhow::Result<()> {
if !(change.chat.is_group() || change.chat.is_supergroup()) {
return Ok(());
}
let chat_id = change.chat.id.0;
let group_id = state
.transport
.observe_group(chat_id, change.chat.title().unwrap_or("Telegram group"))?;
let target = &change.new_chat_member.user;
let target_identity = identity(target)?;
let (membership, _) = member_status(&change.new_chat_member.kind);
let is_kennedy = state.bot_user_id == Some(target_identity.telegram_user_id);
if !target.is_anonymous() && !is_kennedy {
let authorized = report_identity(
state.identity_sink.as_ref(),
target_identity.telegram_user_id,
target_identity.username.as_deref(),
&target_identity.display_name,
)?;
state.transport.observe_group_membership(
&group_id,
MembershipObservation {
identity: target_identity,
membership: membership.into(),
},
authorized,
)?;
} else if is_kennedy {
let was_administrator = matches!(
change.old_chat_member.kind,
ChatMemberKind::Owner(_) | ChatMemberKind::Administrator(_)
);
let is_administrator = matches!(
change.new_chat_member.kind,
ChatMemberKind::Owner(_) | ChatMemberKind::Administrator(_)
);
state.transport.observe_bot_administrator_status(
chat_id,
was_administrator,
is_administrator,
)?;
}
state.identity_sink.observe_group(&group_id)?;
Ok(())
}
pub(super) async fn group_security_snapshot(
bot: &Bot,
state: &AppState,
chat_id: i64,
) -> anyhow::Result<GroupSecuritySnapshot> {
let Some(bot_user_id) = state.bot_user_id else {
return Ok(GroupSecuritySnapshot {
bot_is_administrator: false,
telegram_member_count: 0,
administrators: Vec::new(),
authorized_user_ids: HashSet::new(),
});
};
let bot_member = telegram_requests::retry_request("get_chat_member", || {
bot.get_chat_member(
ChatId(chat_id),
teloxide::types::UserId(
u64::try_from(bot_user_id).expect("the Telegram bot ID is nonnegative"),
),
)
.send()
})
.await?;
let bot_is_administrator = matches!(
bot_member.kind,
ChatMemberKind::Owner(_) | ChatMemberKind::Administrator(_)
);
if !bot_is_administrator {
return Ok(GroupSecuritySnapshot {
bot_is_administrator: false,
telegram_member_count: 0,
administrators: Vec::new(),
authorized_user_ids: HashSet::new(),
});
}
let administrators = telegram_requests::retry_request("get_chat_administrators", || {
bot.get_chat_administrators(ChatId(chat_id)).send()
})
.await?;
let telegram_member_count = i64::from(
telegram_requests::retry_request("get_chat_member_count", || {
bot.get_chat_member_count(ChatId(chat_id)).send()
})
.await?,
);
let mut observations = Vec::new();
for administrator in administrators {
let user = administrator.user;
let identity = identity(&user)?;
if identity.telegram_user_id == bot_user_id || user.is_anonymous() {
continue;
}
state.identity_sink.observe_identity(&IdentityObservation {
telegram_user_id: identity.telegram_user_id,
username: identity.username.clone(),
display_name: identity.display_name.clone(),
})?;
observations.push(MembershipObservation {
identity,
membership: member_status(&administrator.kind).0.into(),
});
}
Ok(GroupSecuritySnapshot {
bot_is_administrator,
telegram_member_count,
administrators: observations,
authorized_user_ids: state.identity_sink.whitelist()?.telegram_user_ids,
})
}
#[derive(Debug)]
struct GroupMessageAuthor {
identity: Option<Identity>,
display_name: String,
group_authored: bool,
}
fn is_group_authored_message(message: &Message) -> bool {
message
.sender_chat
.as_ref()
.is_some_and(|sender| sender.id == message.chat.id)
|| message
.from
.as_ref()
.is_some_and(|user| user.is_anonymous())
}
fn group_message_author(message: &Message) -> anyhow::Result<Option<GroupMessageAuthor>> {
if is_group_authored_message(message) {
let display_name = message
.author_signature()
.unwrap_or("Anonymous group administrator")
.to_owned();
return Ok(Some(GroupMessageAuthor {
identity: None,
display_name,
group_authored: true,
}));
}
let Some(user) = message.from.as_ref() else {
return Ok(None);
};
let identity = identity(user)?;
Ok(Some(GroupMessageAuthor {
display_name: identity.display_name.clone(),
identity: Some(identity),
group_authored: false,
}))
}
fn group_invokes_kennedy(message: &Message, bot_user_id: i64, bot_username: Option<&str>) -> bool {
if message
.reply_to_message()
.and_then(|reply| reply.from.as_ref())
.and_then(|user| i64::try_from(user.id.0).ok())
== Some(bot_user_id)
{
return true;
}
let expected = bot_username.map(normalize_username);
message
.parse_entities()
.into_iter()
.flatten()
.chain(message.parse_caption_entities().into_iter().flatten())
.any(|entity| {
matches!(
entity.kind(),
MessageEntityKind::Mention | MessageEntityKind::BotCommand
) && expected.as_deref().is_some_and(|name| {
normalize_username(entity.text().rsplit('@').next().unwrap_or("")) == name
})
})
}
async fn process_group_message(
bot: &Bot,
state: &AppState,
update_id: i64,
message: Message,
edited: bool,
) -> anyhow::Result<()> {
let chat_id = message.chat.id.0;
let title = message.chat.title().unwrap_or("Telegram group");
let migrated = if let Some(new_chat_id) = message.migrate_to_chat_id() {
Some(
state
.transport
.migrate_group(chat_id, new_chat_id.0, title)?,
)
} else if let Some(old_chat_id) = message.migrate_from_chat_id() {
Some(
state
.transport
.migrate_group(old_chat_id.0, chat_id, title)?,
)
} else {
None
};
if let Some(group_id) = migrated {
state.identity_sink.observe_group(&group_id)?;
return Ok(());
}
let group_id = state.transport.observe_group(chat_id, title)?;
state.identity_sink.observe_group(&group_id)?;
let Some(author) = group_message_author(&message)? else {
return Ok(());
};
if let Some(identity) = author.identity.as_ref() {
let authorized = report_identity(
state.identity_sink.as_ref(),
identity.telegram_user_id,
identity.username.as_deref(),
&identity.display_name,
)?;
state.transport.observe_group_membership(
&group_id,
MembershipObservation {
identity: identity.clone(),
membership: "member".into(),
},
authorized,
)?;
}
for member in message.new_chat_members().unwrap_or_default() {
let identity = identity(member)?;
if state.bot_user_id == Some(identity.telegram_user_id) || member.is_anonymous() {
continue;
}
let authorized = report_identity(
state.identity_sink.as_ref(),
identity.telegram_user_id,
identity.username.as_deref(),
&identity.display_name,
)?;
state.transport.observe_group_membership(
&group_id,
MembershipObservation {
identity,
membership: "member".into(),
},
authorized,
)?;
}
if let Some(member) = message.left_chat_member() {
let identity = identity(member)?;
if state.bot_user_id != Some(identity.telegram_user_id) && !member.is_anonymous() {
let authorized = report_identity(
state.identity_sink.as_ref(),
identity.telegram_user_id,
identity.username.as_deref(),
&identity.display_name,
)?;
state.transport.observe_group_membership(
&group_id,
MembershipObservation {
identity,
membership: "left".into(),
},
authorized,
)?;
}
}
if author.group_authored && !matches!(&message.kind, MessageKind::Common(_)) {
return Ok(());
}
let snapshot = group_security_snapshot(bot, state, chat_id).await?;
let admission = match state.transport.apply_group_security(chat_id, snapshot)? {
GroupAdmission::Quarantined => return Ok(()),
GroupAdmission::Admitted(admission) => admission,
};
let Some(bot_user_id) = state.bot_user_id else {
return Ok(());
};
let invoked = !author.group_authored
&& (admission.behaves_as_direct_message()
|| reset_command(&message)
|| group_invokes_kennedy(&message, bot_user_id, state.bot_username.as_deref()));
let input = parse_message_input_with_feedback(bot, state, &message, invoked && !edited).await?;
let archive_text = input
.as_ref()
.and_then(|input| input.text.clone())
.unwrap_or_else(|| group_message_text(&message));
let content = input.map(Into::into);
admission.accept(AcceptedGroupMessage {
update_id,
message_id: i64::from(message.id.0),
chat_id,
author: author.identity,
author_display_name: author.display_name,
content,
created_at: message.date.to_rfc3339(),
reply_to_message_id: message
.reply_to_message()
.map(|reply| i64::from(reply.id.0)),
archive_text,
invoked,
edited,
})?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
fn message(json: &str) -> Message {
serde_json::from_str(json).unwrap()
}
#[test]
fn reset_detection_accepts_addressed_commands_without_trimming_content() {
let reset = message(
r#"{
"message_id":1,"date":1629404938,
"from":{"id":42,"is_bot":false,"first_name":"David"},
"chat":{"id":42,"first_name":"David","type":"private"},
"text":"/reset@KennedyBot because"
}"#,
);
assert!(reset_command(&reset));
assert_eq!(group_message_text(&reset), "/reset@KennedyBot because");
}
#[test]
fn group_invocation_uses_structured_mentions_and_reply_identity() {
let mention = message(
r#"{
"message_id":7,"date":1629404938,
"from":{"id":42,"is_bot":false,"first_name":"David"},
"chat":{"id":-100,"title":"Friends","type":"supergroup"},
"text":"@KennedyBot hello",
"entities":[{"type":"mention","offset":0,"length":11}]
}"#,
);
assert!(group_invokes_kennedy(&mention, 999, Some("KennedyBot")));
assert!(!group_invokes_kennedy(&mention, 999, Some("OtherBot")));
let reply = message(
r#"{
"message_id":8,"date":1629404938,
"from":{"id":42,"is_bot":false,"first_name":"David"},
"chat":{"id":-100,"title":"Friends","type":"supergroup"},
"text":"following up",
"reply_to_message":{
"message_id":6,"date":1629404937,
"from":{"id":999,"is_bot":true,"first_name":"Kennedy"},
"chat":{"id":-100,"title":"Friends","type":"supergroup"},
"text":"answer"
}
}"#,
);
assert!(group_invokes_kennedy(&reply, 999, Some("OtherBot")));
}
#[test]
fn group_authored_messages_never_invent_a_human_identity() {
let authored = message(
r#"{
"message_id":9,"date":1629404938,
"sender_chat":{"id":-100,"title":"Friends","type":"supergroup"},
"chat":{"id":-100,"title":"Friends","type":"supergroup"},
"author_signature":"Moderator","text":"announcement"
}"#,
);
let author = group_message_author(&authored).unwrap().unwrap();
assert!(author.group_authored);
assert!(author.identity.is_none());
assert_eq!(author.display_name, "Moderator");
}
}