diff --git a/docs/channels/telegram.md b/docs/channels/telegram.md new file mode 100644 index 000000000..346100630 --- /dev/null +++ b/docs/channels/telegram.md @@ -0,0 +1,19 @@ +# Telegram Channel + +The Telegram channel allows OpenHuman to interact with users via a Telegram bot. + +## Silent Streaming + +While the bot is thinking or streaming a reply, updates are sent silently by default to minimize notification spam on the user's device. This means: +- The initial "thinking..." placeholder is sent without a notification sound. +- Intermediate streaming updates (edits to the message) do not trigger new notifications. +- Standalone messages and final fallback messages (if a message needs to be re-sent instead of edited) will still trigger a notification normally. + +This behavior can be controlled via the `silent_streaming` option in the `[channels.telegram]` section of `config.toml`. It defaults to `true`. + +```toml +[channels.telegram] +bot_token = "YOUR_BOT_TOKEN" +allowed_users = ["your_username"] +silent_streaming = true # Set to false to receive notifications for every update +``` diff --git a/src/openhuman/channels/commands.rs b/src/openhuman/channels/commands.rs index ad21c13ae..bd12e1934 100644 --- a/src/openhuman/channels/commands.rs +++ b/src/openhuman/channels/commands.rs @@ -53,7 +53,11 @@ pub async fn doctor_channels(config: Config) -> Result<()> { tg.allowed_users.clone(), tg.mention_only, ) - .with_streaming(tg.stream_mode, tg.draft_update_interval_ms), + .with_streaming( + tg.stream_mode, + tg.draft_update_interval_ms, + tg.silent_streaming, + ), ), )); } @@ -319,6 +323,7 @@ mod tests { allowed_users: vec!["user1".into()], stream_mode: StreamMode::default(), draft_update_interval_ms: 2000, + silent_streaming: true, mention_only: false, }); let _ = doctor_channels(config).await; @@ -375,6 +380,7 @@ mod tests { allowed_users: vec![], stream_mode: StreamMode::default(), draft_update_interval_ms: 2000, + silent_streaming: true, mention_only: false, }); config.channels_config.discord = Some(DiscordConfig { diff --git a/src/openhuman/channels/controllers/ops.rs b/src/openhuman/channels/controllers/ops.rs index dc5bf1994..a5393d456 100644 --- a/src/openhuman/channels/controllers/ops.rs +++ b/src/openhuman/channels/controllers/ops.rs @@ -234,15 +234,21 @@ pub async fn connect_channel( let allowed_users_count = allowed_users.len(); let mut persisted = config.clone(); - let (stream_mode, draft_update_interval_ms, mention_only) = + let (stream_mode, draft_update_interval_ms, silent_streaming, mention_only) = if let Some(existing) = persisted.channels_config.telegram.as_ref() { ( existing.stream_mode, existing.draft_update_interval_ms, + existing.silent_streaming, existing.mention_only, ) } else { - (crate::openhuman::config::StreamMode::default(), 1000, false) + ( + crate::openhuman::config::StreamMode::default(), + 1000, + true, + false, + ) }; persisted.channels_config.telegram = Some(TelegramConfig { @@ -250,6 +256,7 @@ pub async fn connect_channel( allowed_users, stream_mode, draft_update_interval_ms, + silent_streaming, mention_only, }); diff --git a/src/openhuman/channels/providers/telegram/channel.rs b/src/openhuman/channels/providers/telegram/channel.rs index 58d8965a6..8d7ff79fb 100644 --- a/src/openhuman/channels/providers/telegram/channel.rs +++ b/src/openhuman/channels/providers/telegram/channel.rs @@ -54,6 +54,7 @@ pub struct TelegramChannel { typing_handle: Mutex>, stream_mode: StreamMode, draft_update_interval_ms: u64, + silent_streaming: bool, last_draft_edit: Mutex>, mention_only: bool, bot_username: Mutex>, @@ -81,6 +82,7 @@ impl TelegramChannel { client: reqwest::Client::new(), stream_mode: StreamMode::Off, draft_update_interval_ms: 1000, + silent_streaming: true, last_draft_edit: Mutex::new(std::collections::HashMap::new()), typing_handle: Mutex::new(None), mention_only, @@ -94,9 +96,11 @@ impl TelegramChannel { mut self, stream_mode: StreamMode, draft_update_interval_ms: u64, + silent_streaming: bool, ) -> Self { self.stream_mode = stream_mode; self.draft_update_interval_ms = draft_update_interval_ms; + self.silent_streaming = silent_streaming; self } @@ -913,6 +917,7 @@ Allowlist Telegram username (without '@') or numeric user ID.", chat_id: &str, thread_id: Option<&str>, reply_to_message_id: Option, + disable_notification: bool, ) -> anyhow::Result<()> { let chunks = split_message_for_telegram(message); @@ -932,7 +937,8 @@ Allowlist Telegram username (without '@') or numeric user ID.", let mut markdown_body = serde_json::json!({ "chat_id": chat_id, "text": text, - "parse_mode": "Markdown" + "parse_mode": "Markdown", + "disable_notification": disable_notification, }); // Add message_thread_id for forum topic support @@ -969,6 +975,7 @@ Allowlist Telegram username (without '@') or numeric user ID.", let mut plain_body = serde_json::json!({ "chat_id": chat_id, "text": text, + "disable_notification": disable_notification, }); // Add message_thread_id for forum topic support @@ -1533,6 +1540,7 @@ impl Channel for TelegramChannel { let mut body = serde_json::json!({ "chat_id": chat_id, "text": initial_text, + "disable_notification": self.silent_streaming, }); if let Some(tid) = thread_id { body["message_thread_id"] = serde_json::Value::String(tid.to_string()); @@ -1613,6 +1621,7 @@ impl Channel for TelegramChannel { "chat_id": chat_id, "message_id": message_id_parsed, "text": display_text, + "disable_notification": self.silent_streaming, }); let resp = self @@ -1656,7 +1665,13 @@ impl Channel for TelegramChannel { Err(e) => { tracing::warn!("Invalid Telegram message_id '{message_id}': {e}"); return self - .send_text_chunks(text, &chat_id, thread_id.as_deref(), parent_message_id) + .send_text_chunks( + text, + &chat_id, + thread_id.as_deref(), + parent_message_id, + false, + ) .await; } }; @@ -1674,7 +1689,13 @@ impl Channel for TelegramChannel { // Fall back to chunked send return self - .send_text_chunks(text, &chat_id, thread_id.as_deref(), parent_message_id) + .send_text_chunks( + text, + &chat_id, + thread_id.as_deref(), + parent_message_id, + false, + ) .await; } @@ -1683,7 +1704,13 @@ impl Channel for TelegramChannel { Err(e) => { tracing::warn!("Invalid Telegram message_id '{message_id}': {e}"); return self - .send_text_chunks(text, &chat_id, thread_id.as_deref(), parent_message_id) + .send_text_chunks( + text, + &chat_id, + thread_id.as_deref(), + parent_message_id, + false, + ) .await; } }; @@ -1727,8 +1754,14 @@ impl Channel for TelegramChannel { // Edit failed entirely — fall back to new message tracing::warn!("Telegram finalize_draft edit failed; falling back to sendMessage"); - self.send_text_chunks(text, &chat_id, thread_id.as_deref(), parent_message_id) - .await + self.send_text_chunks( + text, + &chat_id, + thread_id.as_deref(), + parent_message_id, + false, + ) + .await } async fn send(&self, message: &SendMessage) -> anyhow::Result<()> { @@ -1777,8 +1810,14 @@ impl Channel for TelegramChannel { if !attachments.is_empty() { if !text_without_markers.is_empty() { - self.send_text_chunks(&text_without_markers, chat_id, thread_id, parent_message_id) - .await?; + self.send_text_chunks( + &text_without_markers, + chat_id, + thread_id, + parent_message_id, + false, + ) + .await?; } for attachment in &attachments { @@ -1794,8 +1833,14 @@ impl Channel for TelegramChannel { return Ok(()); } - self.send_text_chunks(&reactionless_content, chat_id, thread_id, parent_message_id) - .await + self.send_text_chunks( + &reactionless_content, + chat_id, + thread_id, + parent_message_id, + false, + ) + .await } async fn listen(&self, tx: tokio::sync::mpsc::Sender) -> anyhow::Result<()> { diff --git a/src/openhuman/channels/providers/telegram/channel_tests.rs b/src/openhuman/channels/providers/telegram/channel_tests.rs index 264b8c2a4..76ca42977 100644 --- a/src/openhuman/channels/providers/telegram/channel_tests.rs +++ b/src/openhuman/channels/providers/telegram/channel_tests.rs @@ -74,9 +74,10 @@ fn supports_draft_updates_respects_stream_mode() { assert!(!off.supports_draft_updates()); let partial = TelegramChannel::new("fake-token".into(), vec!["*".into()], false) - .with_streaming(StreamMode::Partial, 750); + .with_streaming(StreamMode::Partial, 750, true); assert!(partial.supports_draft_updates()); assert_eq!(partial.draft_update_interval_ms, 750); + assert!(partial.silent_streaming); } #[tokio::test] @@ -91,8 +92,11 @@ async fn send_draft_returns_none_when_stream_mode_off() { #[tokio::test] async fn update_draft_rate_limit_short_circuits_network() { - let ch = TelegramChannel::new("fake-token".into(), vec!["*".into()], false) - .with_streaming(StreamMode::Partial, 60_000); + let ch = TelegramChannel::new("fake-token".into(), vec!["*".into()], false).with_streaming( + StreamMode::Partial, + 60_000, + true, + ); ch.last_draft_edit .lock() .insert("123".to_string(), std::time::Instant::now()); @@ -103,8 +107,11 @@ async fn update_draft_rate_limit_short_circuits_network() { #[tokio::test] async fn update_draft_utf8_truncation_is_safe_for_multibyte_text() { - let ch = TelegramChannel::new("fake-token".into(), vec!["*".into()], false) - .with_streaming(StreamMode::Partial, 0); + let ch = TelegramChannel::new("fake-token".into(), vec!["*".into()], false).with_streaming( + StreamMode::Partial, + 0, + true, + ); let long_emoji_text = "😀".repeat(TELEGRAM_MAX_MESSAGE_LENGTH + 20); // Invalid message_id returns early after building display_text. @@ -117,8 +124,11 @@ async fn update_draft_utf8_truncation_is_safe_for_multibyte_text() { #[tokio::test] async fn finalize_draft_invalid_message_id_falls_back_to_chunk_send() { - let ch = TelegramChannel::new("fake-token".into(), vec!["*".into()], false) - .with_streaming(StreamMode::Partial, 0); + let ch = TelegramChannel::new("fake-token".into(), vec!["*".into()], false).with_streaming( + StreamMode::Partial, + 0, + true, + ); let long_text = "a".repeat(TELEGRAM_MAX_MESSAGE_LENGTH + 64); // For oversized text + invalid draft message_id, finalize_draft should @@ -1382,6 +1392,23 @@ fn track_update_id_large_volume_beyond_cache_does_not_panic() { ); } +#[test] +fn silent_streaming_is_configurable() { + let silent = TelegramChannel::new("fake-token".into(), vec!["*".into()], false).with_streaming( + StreamMode::Partial, + 1000, + true, + ); + assert!(silent.silent_streaming); + + let noisy = TelegramChannel::new("fake-token".into(), vec!["*".into()], false).with_streaming( + StreamMode::Partial, + 1000, + false, + ); + assert!(!noisy.silent_streaming); +} + // ── Reply-target parsing unit tests ──────────────────────────── #[test] diff --git a/src/openhuman/channels/runtime/startup.rs b/src/openhuman/channels/runtime/startup.rs index eb9654533..310a4af7f 100644 --- a/src/openhuman/channels/runtime/startup.rs +++ b/src/openhuman/channels/runtime/startup.rs @@ -241,7 +241,11 @@ pub async fn start_channels(config: Config) -> Result<()> { tg.allowed_users.clone(), tg.mention_only, ) - .with_streaming(tg.stream_mode, tg.draft_update_interval_ms), + .with_streaming( + tg.stream_mode, + tg.draft_update_interval_ms, + tg.silent_streaming, + ), )); } else { tracing::info!( diff --git a/src/openhuman/config/mod.rs b/src/openhuman/config/mod.rs index 117b3723b..d1f320f9f 100644 --- a/src/openhuman/config/mod.rs +++ b/src/openhuman/config/mod.rs @@ -73,6 +73,7 @@ mod tests { allowed_users: vec!["alice".into()], stream_mode: StreamMode::default(), draft_update_interval_ms: 1000, + silent_streaming: true, mention_only: false, }; diff --git a/src/openhuman/config/schema/channels.rs b/src/openhuman/config/schema/channels.rs index 4c532a506..89339ece8 100644 --- a/src/openhuman/config/schema/channels.rs +++ b/src/openhuman/config/schema/channels.rs @@ -102,6 +102,10 @@ pub(crate) fn default_draft_update_interval_ms() -> u64 { 1000 } +fn default_silent_streaming() -> bool { + true +} + #[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)] pub struct TelegramConfig { pub bot_token: String, @@ -110,6 +114,8 @@ pub struct TelegramConfig { pub stream_mode: StreamMode, #[serde(default = "default_draft_update_interval_ms")] pub draft_update_interval_ms: u64, + #[serde(default = "default_silent_streaming")] + pub silent_streaming: bool, #[serde(default)] pub mention_only: bool, } diff --git a/src/openhuman/config/schema/channels_tests.rs b/src/openhuman/config/schema/channels_tests.rs index 8c82d1bef..8b2d745f3 100644 --- a/src/openhuman/config/schema/channels_tests.rs +++ b/src/openhuman/config/schema/channels_tests.rs @@ -44,6 +44,7 @@ fn has_listening_integrations_detects_telegram() { allowed_users: vec![], stream_mode: StreamMode::Off, draft_update_interval_ms: 1000, + silent_streaming: true, mention_only: false, }); assert!(cfg.has_listening_integrations());