1use std::{
2 collections::HashSet,
3 path::PathBuf,
4 sync::Arc,
5 time::{Duration, Instant},
6};
7
8use anyhow::Context;
9use futures::StreamExt;
10use kcode_telegram_text_delivery::send_telegram_message;
11use kcode_telegram_transport_state::{
12 AcceptedGroupMessage, AcceptedMessage, AdmittedGroup, Event, GroupAdmission,
13 GroupSecuritySnapshot, Identity, Media as StoredMedia, MembershipObservation, MessageContent,
14 MessagePreparation, MessageRevision, ReplyDelivery, ResetDelivery, SentMessage, StateError,
15 TransportState,
16};
17use serde::Serialize;
18use serde_json::{Value, json};
19use teloxide::{
20 net::Download,
21 prelude::*,
22 requests::Request,
23 types::{
24 AllowedUpdate, ChatMemberKind, Message, MessageEntityKind, MessageKind, Update, UpdateKind,
25 },
26};
27use zeroize::Zeroize;
28
29mod inbound;
30mod native_media;
31mod telegram_requests;
32mod transport_extensions;
33mod update_dispatch;
34
35const UNAUTHORIZED_MESSAGE: &str =
36 "Sorry, this Kennedy bot is private and your Telegram handle is not whitelisted.";
37const TELEGRAM_POLL_TIMEOUT_SECONDS: u32 = 90;
38const TELEGRAM_HTTP_TIMEOUT_SECONDS: u64 = 120;
39
40pub struct BotToken(String);
41
42impl BotToken {
43 pub fn new(value: String) -> anyhow::Result<Self> {
44 anyhow::ensure!(
45 !value.trim().is_empty(),
46 "Telegram bot token must not be empty"
47 );
48 Ok(Self(value))
49 }
50
51 fn expose(&self) -> &str {
52 &self.0
53 }
54}
55
56impl std::fmt::Debug for BotToken {
57 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
58 formatter.write_str("BotToken([REDACTED])")
59 }
60}
61
62impl Drop for BotToken {
63 fn drop(&mut self) {
64 self.0.zeroize();
65 }
66}
67
68#[derive(Clone, Debug, PartialEq, Eq)]
69pub struct IdentityObservation {
70 pub telegram_user_id: i64,
71 pub username: Option<String>,
72 pub display_name: String,
73}
74
75#[derive(Clone, Debug, Default)]
76pub struct WhitelistSnapshot {
77 pub telegram_user_ids: HashSet<i64>,
78}
79
80impl WhitelistSnapshot {
81 pub fn contains(&self, telegram_user_id: i64) -> bool {
82 self.telegram_user_ids.contains(&telegram_user_id)
83 }
84}
85
86#[derive(Clone, Debug, PartialEq, Eq)]
87pub enum AddUserOutcome {
88 Forbidden,
89 Whitelisted {
90 handle: String,
91 telegram_user_id: Option<i64>,
92 },
93}
94
95pub trait IdentitySink: Send + Sync {
96 fn observe_identity(&self, observation: &IdentityObservation) -> anyhow::Result<()>;
97 fn whitelist(&self) -> anyhow::Result<WhitelistSnapshot>;
98 fn request_add_user(
99 &self,
100 requested_by_telegram_user_id: i64,
101 handle: &str,
102 ) -> anyhow::Result<AddUserOutcome>;
103 fn observe_group(&self, group_id: &str) -> anyhow::Result<()>;
104}
105
106pub struct Config {
107 pub database: PathBuf,
108 pub bot_token: Option<BotToken>,
109 pub identity_sink: Arc<dyn IdentitySink>,
110 pub max_voice_bytes: usize,
111}
112
113#[derive(Clone)]
114pub struct Service {
115 state: AppState,
116}
117
118pub struct Runtime {
119 service: Service,
120}
121
122#[derive(Clone)]
123struct AppState {
124 transport: TransportState,
125 identity_sink: Arc<dyn IdentitySink>,
126 bot: Option<Bot>,
127 max_voice_bytes: usize,
128 bot_user_id: Option<i64>,
129 bot_username: Option<String>,
130}
131
132#[derive(Debug)]
133pub struct Error {
134 code: &'static str,
135 message: String,
136}
137
138type ApiError = Error;
139
140impl Error {
141 fn new(code: &'static str, message: impl Into<String>) -> Self {
142 Self {
143 code,
144 message: message.into(),
145 }
146 }
147
148 fn bad(message: impl Into<String>) -> Self {
149 Self::new("invalid_request", message)
150 }
151
152 fn unavailable() -> Self {
153 Self::new(
154 "telegram_unavailable",
155 "The Telegram bot token is not configured.",
156 )
157 }
158
159 fn internal(error: impl std::fmt::Display) -> Self {
160 tracing::warn!(error=%error, "Telegram relay request failed");
161 Self::new(
162 "internal_error",
163 "An unexpected Telegram relay error occurred.",
164 )
165 }
166
167 fn state(error: StateError) -> Self {
168 if error.code() == "internal_error" {
169 Self::internal(error)
170 } else {
171 Self::new(error.code(), error.message())
172 }
173 }
174
175 pub fn code(&self) -> &'static str {
176 self.code
177 }
178
179 pub fn message(&self) -> &str {
180 &self.message
181 }
182}
183
184impl std::fmt::Display for Error {
185 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
186 formatter.write_str(&self.message)
187 }
188}
189
190impl std::error::Error for Error {}
191
192#[derive(Clone, Debug, Serialize)]
193#[serde(rename_all = "camelCase")]
194pub struct PrivateSession {
195 pub telegram_user_id: i64,
196 pub current_conversation_id: Option<String>,
197}
198
199#[derive(Clone, Debug, Serialize)]
200#[serde(rename_all = "camelCase")]
201pub struct Status {
202 pub service: &'static str,
203 pub status: &'static str,
204 pub telegram: &'static str,
205 pub inbound_media_kinds: &'static [&'static str],
206 pub outbound_media_kinds: &'static [&'static str],
207 pub max_media_bytes: usize,
208}
209
210#[derive(Clone, Debug)]
211pub struct Media {
212 pub bytes: Vec<u8>,
213 pub media_type: String,
214}
215
216#[derive(Clone, Debug)]
217pub struct MediaMetadata {
218 pub size_bytes: u64,
219 pub media_type: String,
220}
221
222#[derive(Clone, Debug)]
223pub struct Attachment {
224 pub bytes: Vec<u8>,
225 pub file_name: Option<String>,
226 pub media_type: Option<String>,
227 pub kind: Option<String>,
228 pub caption: Option<String>,
229}
230
231pub async fn open(config: Config) -> anyhow::Result<Runtime> {
232 if config.max_voice_bytes == 0 {
233 anyhow::bail!("telegram max_voice_bytes must be greater than zero");
234 }
235 let transport =
236 TransportState::open(&config.database).context("opening Telegram transport state")?;
237 let bot = match config.bot_token.as_ref() {
238 Some(token) => {
239 let client = teloxide::net::default_reqwest_settings()
240 .timeout(Duration::from_secs(TELEGRAM_HTTP_TIMEOUT_SECONDS))
241 .build()
242 .context("building Telegram HTTP client")?;
243 Some(Bot::with_client(token.expose(), client))
244 }
245 None => None,
246 };
247 let (bot_user_id, bot_username) = if let Some(bot) = bot.as_ref() {
248 let me = telegram_requests::retry_request("get_me", || bot.get_me().send())
249 .await
250 .map_err(|error| {
251 anyhow::anyhow!(
252 "validating Telegram bot token failed ({})",
253 telegram_requests::request_error_class(&error)
254 )
255 })?;
256 (
257 Some(i64::try_from(me.id.0).context("Telegram bot ID exceeds SQLite range")?),
258 me.username.clone(),
259 )
260 } else {
261 (None, None)
262 };
263 let service = Service {
264 state: AppState {
265 transport,
266 identity_sink: config.identity_sink,
267 bot: bot.clone(),
268 max_voice_bytes: config.max_voice_bytes,
269 bot_user_id,
270 bot_username,
271 },
272 };
273 tracing::info!(enabled = bot.is_some(), "Telegram transport ready");
274 Ok(Runtime { service })
275}
276
277impl Runtime {
278 pub fn service(&self) -> Service {
279 self.service.clone()
280 }
281
282 pub async fn run(self) -> anyhow::Result<()> {
283 let Some(bot) = self.service.state.bot.clone() else {
284 return std::future::pending::<anyhow::Result<()>>().await;
285 };
286 inbound::poll_telegram(bot, self.service.state)
287 .await
288 .context("polling Telegram")
289 }
290}
291
292impl Service {
293 pub fn status(&self) -> Status {
294 Status {
295 service: "kcode-tg-kennedy-bot",
296 status: "ok",
297 telegram: if self.state.bot.is_some() {
298 "ready"
299 } else {
300 "disabled"
301 },
302 inbound_media_kinds: &native_media::INBOUND_MEDIA_KINDS,
303 outbound_media_kinds: &native_media::OUTBOUND_MEDIA_KINDS,
304 max_media_bytes: self.state.max_voice_bytes,
305 }
306 }
307
308 pub async fn list_private_sessions(&self) -> Result<Vec<PrivateSession>, Error> {
309 self.state
310 .transport
311 .list_private_sessions()
312 .map_err(ApiError::state)
313 .map(|sessions| {
314 sessions
315 .into_iter()
316 .map(|session| PrivateSession {
317 telegram_user_id: session.telegram_user_id,
318 current_conversation_id: session.current_conversation_id,
319 })
320 .collect()
321 })
322 }
323
324 pub async fn send_private_message(
325 &self,
326 telegram_user_id: i64,
327 conversation_id: String,
328 expected_conversation_id: Option<String>,
329 text: String,
330 ) -> Result<Value, Error> {
331 send_private_message(
332 self.state.clone(),
333 telegram_user_id,
334 conversation_id,
335 expected_conversation_id,
336 text,
337 )
338 .await
339 }
340
341 pub async fn send_cold_private_message(
342 &self,
343 telegram_user_id: i64,
344 text: String,
345 ) -> Result<Value, Error> {
346 send_cold_private_message(self.state.clone(), telegram_user_id, text).await
347 }
348
349 pub async fn send_group_message(&self, group_id: String, text: String) -> Result<Value, Error> {
350 send_group_message(self.state.clone(), group_id, text).await
351 }
352
353 pub async fn list_group_ingress(&self) -> Result<Value, Error> {
354 let batches = self
355 .state
356 .transport
357 .list_group_ingress()
358 .map_err(ApiError::state)?;
359 Ok(json!({"batches":batches}))
360 }
361
362 pub async fn complete_group_ingress(&self, batch_id: String) -> Result<Value, Error> {
363 self.state
364 .transport
365 .complete_group_ingress(&batch_id)
366 .map_err(ApiError::state)?;
367 Ok(json!({"id":batch_id,"status":"complete"}))
368 }
369
370 pub async fn list_group_session_updates(&self) -> Result<Value, Error> {
371 let updates = self
372 .state
373 .transport
374 .list_group_session_updates()
375 .map_err(ApiError::state)?;
376 Ok(json!({"updates":updates}))
377 }
378
379 pub async fn detach_group_session(
380 &self,
381 conversation_id: String,
382 group_id: String,
383 telegram_user_id: i64,
384 ) -> Result<Value, Error> {
385 self.state
386 .transport
387 .detach_group_session(&conversation_id, &group_id, telegram_user_id)
388 .map_err(ApiError::state)?;
389 Ok(json!({
390 "conversationId":conversation_id,
391 "groupId":group_id,
392 "telegramUserId":telegram_user_id,
393 "status":"detached",
394 }))
395 }
396
397 pub async fn acknowledge_group_session_context(
398 &self,
399 conversation_id: String,
400 through_message_id: i64,
401 ) -> Result<Value, Error> {
402 self.state
403 .transport
404 .acknowledge_group_session_context(&conversation_id, through_message_id)
405 .map_err(ApiError::state)?;
406 Ok(json!({
407 "conversationId":conversation_id,
408 "throughMessageId":through_message_id,
409 }))
410 }
411
412 pub async fn complete_silent_group_reset(
413 &self,
414 conversation_id: String,
415 ) -> Result<Value, Error> {
416 self.state
417 .transport
418 .complete_silent_group_reset(&conversation_id)
419 .map_err(ApiError::state)?;
420 Ok(json!({"conversationId":conversation_id,"status":"complete"}))
421 }
422
423 pub async fn save_group_message_preparation(
424 &self,
425 chat_id: i64,
426 message_id: i64,
427 text: String,
428 model: Option<String>,
429 format: Option<String>,
430 truncated: bool,
431 ) -> Result<Value, Error> {
432 self.state
433 .transport
434 .save_group_message_preparation(
435 chat_id,
436 message_id,
437 MessagePreparation {
438 text: text.clone(),
439 model,
440 format,
441 truncated,
442 },
443 )
444 .map_err(ApiError::state)?;
445 Ok(json!({"chatId":chat_id,"messageId":message_id,"text":text}))
446 }
447
448 pub async fn list_events(&self) -> Result<Value, Error> {
449 let events = self
450 .state
451 .transport
452 .pending_events()
453 .map_err(ApiError::state)?;
454 Ok(json!({"events":events}))
455 }
456
457 pub async fn bind_event(
458 &self,
459 event_id: String,
460 conversation_id: String,
461 expected_conversation_id: Option<String>,
462 ) -> Result<Value, Error> {
463 let event = self
464 .state
465 .transport
466 .bind_event(
467 &event_id,
468 &conversation_id,
469 expected_conversation_id.as_deref(),
470 )
471 .map_err(ApiError::state)?;
472 serde_json::to_value(event).map_err(ApiError::internal)
473 }
474
475 pub async fn save_transcription(
476 &self,
477 event_id: String,
478 text: String,
479 transcription_model: String,
480 ) -> Result<Value, Error> {
481 let event = self
482 .state
483 .transport
484 .save_transcription(&event_id, &text, &transcription_model)
485 .map_err(ApiError::state)?;
486 serde_json::to_value(event).map_err(ApiError::internal)
487 }
488
489 pub async fn reply_event(
490 &self,
491 event_id: String,
492 conversation_id: String,
493 text: String,
494 context_warning: Option<String>,
495 ) -> Result<Value, Error> {
496 let event = reply_event(
497 self.state.clone(),
498 event_id,
499 conversation_id,
500 text,
501 context_warning,
502 )
503 .await?;
504 serde_json::to_value(event).map_err(ApiError::internal)
505 }
506
507 pub async fn abort_event(
508 &self,
509 event_id: String,
510 conversation_id: Option<String>,
511 message: String,
512 ) -> Result<Value, Error> {
513 let event = abort_event(self.state.clone(), event_id, conversation_id, message).await?;
514 serde_json::to_value(event).map_err(ApiError::internal)
515 }
516
517 pub async fn interrupt_event(
518 &self,
519 event_id: String,
520 conversation_id: String,
521 ) -> Result<Value, Error> {
522 let event = self
523 .state
524 .transport
525 .interrupt_event(&event_id, &conversation_id)
526 .map_err(ApiError::state)?;
527 serde_json::to_value(event).map_err(ApiError::internal)
528 }
529
530 pub async fn complete_reset(
531 &self,
532 event_id: String,
533 message: Option<String>,
534 ) -> Result<Value, Error> {
535 let event = complete_reset(self.state.clone(), event_id, message).await?;
536 serde_json::to_value(event).map_err(ApiError::internal)
537 }
538
539 pub fn event_media(&self, event_id: &str) -> Result<Media, Error> {
540 self.state
541 .transport
542 .event_media(event_id)
543 .map(convert_media)
544 .map_err(ApiError::state)
545 }
546
547 pub fn event_media_metadata(&self, event_id: &str) -> Result<MediaMetadata, Error> {
548 self.state
549 .transport
550 .event_media_metadata(event_id)
551 .map(|media| MediaMetadata {
552 size_bytes: media.size_bytes,
553 media_type: media.media_type,
554 })
555 .map_err(ApiError::state)
556 }
557
558 pub fn group_message_media(&self, chat_id: i64, message_id: i64) -> Result<Media, Error> {
559 self.state
560 .transport
561 .group_message_media(chat_id, message_id)
562 .map(convert_media)
563 .map_err(ApiError::state)
564 }
565
566 pub fn group_message_media_metadata(
567 &self,
568 chat_id: i64,
569 message_id: i64,
570 ) -> Result<MediaMetadata, Error> {
571 self.state
572 .transport
573 .group_message_media_metadata(chat_id, message_id)
574 .map(|media| MediaMetadata {
575 size_bytes: media.size_bytes,
576 media_type: media.media_type,
577 })
578 .map_err(ApiError::state)
579 }
580}
581
582pub fn migrate_storage(database: &std::path::Path) -> anyhow::Result<()> {
583 kcode_telegram_transport_state::migrate_storage(database)
584}
585
586fn convert_media(media: StoredMedia) -> Media {
587 Media {
588 bytes: media.bytes,
589 media_type: media.media_type,
590 }
591}
592
593fn normalize_username(value: &str) -> String {
594 value.trim().trim_start_matches('@').to_ascii_lowercase()
595}
596
597fn nonempty_verbatim(value: &str) -> Option<&str> {
598 (!value.trim().is_empty()).then_some(value)
599}
600
601fn validate_opaque_group_id(group_id: &str) -> Result<&str, ApiError> {
602 let group_id = group_id.trim();
603 if group_id.is_empty() || group_id.len() > 200 || group_id.chars().any(char::is_control) {
604 return Err(ApiError::bad("groupId is not a valid opaque group ID."));
605 }
606 Ok(group_id)
607}
608
609async fn send_telegram_text(
610 bot: &Bot,
611 chat_id: i64,
612 text: &str,
613 reply_to_message_id: Option<i64>,
614) -> Result<Vec<Message>, ApiError> {
615 kcode_telegram_text_delivery::send_telegram_text(bot, chat_id, text, reply_to_message_id)
616 .await
617 .map_err(|error| {
618 tracing::warn!(
619 %chat_id,
620 error_class = telegram_requests::request_error_class(&error),
621 "Telegram reply failed"
622 );
623 ApiError::new("telegram_send_failed", "Telegram did not accept the reply.")
624 })
625}
626
627fn sent_messages(messages: Vec<Message>) -> Vec<SentMessage> {
628 messages
629 .into_iter()
630 .map(|message| SentMessage {
631 message_id: i64::from(message.id.0),
632 text: message.text().unwrap_or("").to_owned(),
633 sent_at: message.date.to_rfc3339(),
634 })
635 .collect()
636}
637
638async fn send_cold_private_message(
639 state: AppState,
640 telegram_user_id: i64,
641 text: String,
642) -> Result<Value, ApiError> {
643 let started = Instant::now();
644 let text = nonempty_verbatim(&text).ok_or_else(|| ApiError::bad("text must not be empty."))?;
645 let delivery = state
646 .transport
647 .cold_private_delivery(telegram_user_id)
648 .map_err(ApiError::state)?;
649 let bot = state.bot.as_ref().ok_or_else(ApiError::unavailable)?;
650 let sent = send_telegram_text(bot, delivery.chat_id(), text, None).await?;
651 let message_ids = sent
652 .iter()
653 .map(|message| i64::from(message.id.0))
654 .collect::<Vec<_>>();
655 tracing::info!(
656 %telegram_user_id,
657 duration_ms=started.elapsed().as_millis(),
658 "Telegram cold direct message"
659 );
660 Ok(json!({"telegramUserId":telegram_user_id,"messageIds":message_ids}))
661}
662
663async fn send_private_message(
664 state: AppState,
665 telegram_user_id: i64,
666 conversation_id: String,
667 expected_conversation_id: Option<String>,
668 text: String,
669) -> Result<Value, ApiError> {
670 let started = Instant::now();
671 let text = nonempty_verbatim(&text).ok_or_else(|| ApiError::bad("text must not be empty."))?;
672 let delivery = state
673 .transport
674 .private_delivery(
675 telegram_user_id,
676 conversation_id.clone(),
677 expected_conversation_id,
678 )
679 .map_err(ApiError::state)?;
680 let bot = state.bot.as_ref().ok_or_else(ApiError::unavailable)?;
681 let sent = send_telegram_text(bot, delivery.chat_id(), text, None).await?;
682 let message_ids = sent
683 .iter()
684 .map(|message| i64::from(message.id.0))
685 .collect::<Vec<_>>();
686 delivery.record_accepted().map_err(ApiError::state)?;
687 tracing::info!(
688 %telegram_user_id,
689 %conversation_id,
690 duration_ms=started.elapsed().as_millis(),
691 "Telegram cold direct message"
692 );
693 Ok(json!({
694 "telegramUserId":telegram_user_id,
695 "conversationId":conversation_id,
696 "messageIds":message_ids,
697 }))
698}
699
700async fn validated_group_delivery(
701 state: &AppState,
702 group_id: &str,
703) -> Result<AdmittedGroup, ApiError> {
704 let group_id = validate_opaque_group_id(group_id)?;
705 let candidate = state
706 .transport
707 .group_delivery_candidate(group_id)
708 .map_err(ApiError::state)?;
709 let bot = state.bot.as_ref().ok_or_else(ApiError::unavailable)?;
710 let snapshot = inbound::group_security_snapshot(bot, state, candidate.chat_id())
711 .await
712 .map_err(|error| {
713 tracing::warn!(
714 %group_id,
715 error_class = telegram_requests::anyhow_error_class(&error),
716 "Telegram group delivery authorization refresh failed"
717 );
718 ApiError::new(
719 "group_validation_failed",
720 "Telegram group membership could not be revalidated.",
721 )
722 })?;
723 candidate
724 .authorize(snapshot)
725 .map_err(ApiError::state)?
726 .ok_or_else(|| {
727 ApiError::new(
728 "group_not_allowed",
729 "Kennedy may send only when she is an administrator and every historical group member is whitelisted.",
730 )
731 })
732}
733
734async fn send_group_message(
735 state: AppState,
736 group_id: String,
737 text: String,
738) -> Result<Value, ApiError> {
739 let started = Instant::now();
740 let group_id = validate_opaque_group_id(&group_id)?.to_owned();
741 let text = nonempty_verbatim(&text).ok_or_else(|| ApiError::bad("text must not be empty."))?;
742 let delivery = validated_group_delivery(&state, &group_id).await?;
743 let bot = state.bot.as_ref().ok_or_else(ApiError::unavailable)?;
744 let sent = send_telegram_text(bot, delivery.chat_id(), text, None).await?;
745 let message_ids = sent
746 .iter()
747 .map(|message| i64::from(message.id.0))
748 .collect::<Vec<_>>();
749 tracing::info!(
750 %group_id,
751 duration_ms=started.elapsed().as_millis(),
752 "Telegram cold group message"
753 );
754 Ok(json!({"groupId":group_id,"messageIds":message_ids}))
755}
756
757async fn reply_event(
758 state: AppState,
759 id: String,
760 conversation_id: String,
761 text: String,
762 context_warning: Option<String>,
763) -> Result<Event, ApiError> {
764 let started = Instant::now();
765 let text = nonempty_verbatim(&text).ok_or_else(|| ApiError::bad("text must not be empty."))?;
766 let delivery = match state
767 .transport
768 .reply_delivery(&id, &conversation_id)
769 .map_err(ApiError::state)?
770 {
771 ReplyDelivery::Complete(event) => return Ok(event),
772 ReplyDelivery::Pending(delivery) => delivery,
773 };
774 let event = delivery.event();
775 let bot = state.bot.as_ref().ok_or_else(ApiError::unavailable)?;
776 let group_reply = (event.session_kind == "group").then_some(event.message_id);
777 let mut sent = send_telegram_text(bot, event.chat_id, text, group_reply).await?;
778 if let Some(warning) = context_warning.as_deref().and_then(nonempty_verbatim) {
779 sent.extend(send_telegram_text(bot, event.chat_id, warning, None).await?);
780 }
781 let event = delivery
782 .record_text(sent_messages(sent), true)
783 .map_err(ApiError::state)?;
784 tracing::info!(event_id=%id, duration_ms=started.elapsed().as_millis(), "Telegram reply");
785 Ok(event)
786}
787
788async fn abort_event(
789 state: AppState,
790 id: String,
791 conversation_id: Option<String>,
792 message: String,
793) -> Result<Event, ApiError> {
794 let message =
795 nonempty_verbatim(&message).ok_or_else(|| ApiError::bad("message must not be empty."))?;
796 let outcome = state
797 .transport
798 .abort_event(&id, conversation_id.as_deref())
799 .map_err(ApiError::state)?;
800 let event = outcome.event().clone();
801 if !outcome.newly_aborted() {
802 return Ok(event);
803 }
804 let sent = if let Some(bot) = state.bot.as_ref() {
805 let group_reply = (event.session_kind == "group").then_some(event.message_id);
806 match send_telegram_text(bot, event.chat_id, message, group_reply).await {
807 Ok(sent) => sent_messages(sent),
808 Err(error) => {
809 tracing::warn!(event_id=%id, error=%error.message, "Telegram timeout notice could not be delivered");
810 Vec::new()
811 }
812 }
813 } else {
814 Vec::new()
815 };
816 outcome.record_notice(sent).map_err(ApiError::state)?;
817 tracing::warn!(event_id=%id, "Telegram response aborted at its hard timeout");
818 Ok(event)
819}
820
821async fn complete_reset(
822 state: AppState,
823 id: String,
824 message: Option<String>,
825) -> Result<Event, ApiError> {
826 let delivery = match state
827 .transport
828 .reset_delivery(&id)
829 .map_err(ApiError::state)?
830 {
831 ResetDelivery::Complete(event) => return Ok(event),
832 ResetDelivery::Pending(delivery) => delivery,
833 };
834 let event = delivery.event();
835 let bot = state.bot.as_ref().ok_or_else(ApiError::unavailable)?;
836 let message = message.as_deref().and_then(nonempty_verbatim).unwrap_or(
837 "Conversation reset. Your previous Telegram session has been queued for memory ingress.",
838 );
839 let group_reply = (event.session_kind == "group").then_some(event.message_id);
840 let sent = send_telegram_text(bot, event.chat_id, message, group_reply).await?;
841 delivery
842 .record_reset(sent_messages(sent))
843 .map_err(ApiError::state)
844}
845
846#[cfg(test)]
847mod tests {
848 use super::*;
849
850 #[derive(Default)]
851 struct TestIdentities;
852
853 impl IdentitySink for TestIdentities {
854 fn observe_identity(&self, _observation: &IdentityObservation) -> anyhow::Result<()> {
855 Ok(())
856 }
857
858 fn whitelist(&self) -> anyhow::Result<WhitelistSnapshot> {
859 Ok(WhitelistSnapshot::default())
860 }
861
862 fn request_add_user(
863 &self,
864 _requested_by_telegram_user_id: i64,
865 _handle: &str,
866 ) -> anyhow::Result<AddUserOutcome> {
867 Ok(AddUserOutcome::Forbidden)
868 }
869
870 fn observe_group(&self, _group_id: &str) -> anyhow::Result<()> {
871 Ok(())
872 }
873 }
874
875 #[test]
876 fn token_debug_is_redacted() {
877 let token = BotToken::new("123:secret".into()).unwrap();
878 assert_eq!(format!("{token:?}"), "BotToken([REDACTED])");
879 assert!(BotToken::new(" ".into()).is_err());
880 }
881
882 #[test]
883 fn disabled_status_preserves_the_public_capability_projection() {
884 let service = Service {
885 state: AppState {
886 transport: TransportState::open_in_memory().unwrap(),
887 identity_sink: Arc::new(TestIdentities),
888 bot: None,
889 max_voice_bytes: 4096,
890 bot_user_id: None,
891 bot_username: None,
892 },
893 };
894 let status = service.status();
895 assert_eq!(status.service, "kcode-tg-kennedy-bot");
896 assert_eq!(status.telegram, "disabled");
897 assert_eq!(status.max_media_bytes, 4096);
898 assert!(status.inbound_media_kinds.contains(&"document"));
899 assert!(status.outbound_media_kinds.contains(&"photo"));
900 }
901
902 #[test]
903 fn verbatim_text_validation_does_not_trim_delivery_content() {
904 assert_eq!(nonempty_verbatim(" hello\n"), Some(" hello\n"));
905 assert_eq!(nonempty_verbatim(" \n\t"), None);
906 }
907}