1use std::{
2 collections::{HashMap, HashSet},
3 ops::Deref,
4 sync::{
5 Arc,
6 atomic::{AtomicBool, Ordering},
7 },
8 time::{Duration, Instant},
9};
10
11use anyhow::Context as _;
12use chrono::{DateTime, Duration as ChronoDuration, Utc};
13use kcode_kennedy_orchestration::{
14 AgentMode, ApiError, Orchestrator, Session, TurnCompletion, data_url, persist_record,
15 telegram_caption_for,
16};
17use kcode_kennedy_roots::DirectoryRoots;
18use kcode_kennedy_sessions::{ResolvedObject, SessionOptions, validate_delivery_file_name};
19use kcode_server_object_envelopes::sanitize_file_name;
20use kcode_session_history::SessionRecord;
21use serde_json::{Value, json};
22use tokio::sync::{Mutex, RwLock};
23use uuid::Uuid;
24
25const POLL_INTERVAL: Duration = Duration::from_secs(1);
26const TELEGRAM_TIMEOUT: Duration = Duration::from_secs(90 * 60);
27const TELEGRAM_SESSION_MAX_AGE: ChronoDuration = ChronoDuration::hours(6);
28const TELEGRAM_TIMEOUT_NOTICE: &str = "Kennedy could not complete a response within 90 minutes, so this request was stopped. Please send it again if you want to retry it.";
29
30#[derive(Clone, Debug)]
31pub struct Config {
32 pub telegram_max_media_bytes: usize,
33 pub telegram_web_user_handle: String,
34}
35
36#[derive(Debug, Eq, PartialEq)]
37enum MissingGroupSessionRecovery {
38 CompleteSilentReset,
39 DetachCurrent {
40 group_id: String,
41 telegram_user_id: i64,
42 },
43}
44
45struct TelegramEventRetry {
46 failures: u32,
47 not_before: Instant,
48 last_error: String,
49}
50
51enum TelegramDelivery {
52 Object {
53 object_id: String,
54 file_name: Option<String>,
55 },
56 Text {
57 text: String,
58 response_warning: Value,
59 captionable: bool,
60 },
61}
62
63pub struct Runtime {
64 config: Config,
65 control: Arc<Orchestrator>,
66 roots: DirectoryRoots,
67 writer_job_active: AtomicBool,
68 events_in_flight: Mutex<HashSet<String>>,
69 event_retries: Mutex<HashMap<String, TelegramEventRetry>>,
70 group_updates_in_flight: Mutex<HashSet<String>>,
71 group_ingress_in_flight: Mutex<HashSet<String>>,
72 last_poll_error: RwLock<Option<String>>,
73}
74
75impl Runtime {
76 pub fn new(config: Config, control: Arc<Orchestrator>, roots: DirectoryRoots) -> Self {
77 Self {
78 config,
79 control,
80 roots,
81 writer_job_active: AtomicBool::new(false),
82 events_in_flight: Mutex::new(HashSet::new()),
83 event_retries: Mutex::new(HashMap::new()),
84 group_updates_in_flight: Mutex::new(HashSet::new()),
85 group_ingress_in_flight: Mutex::new(HashSet::new()),
86 last_poll_error: RwLock::new(None),
87 }
88 }
89
90 pub async fn run(self: Arc<Self>) -> anyhow::Result<()> {
91 self.initialize_until_ready().await;
92 self.control.api().telegram_health();
93 self.queue_detached_private_telegram_sessions().await?;
94 let wakeups = self.clone();
95 tokio::spawn(async move { wakeups.run_wakeup_scheduler().await });
96 loop {
97 let result = async {
98 self.roots.reconcile_pending().await?;
99 let histories = self.list_history().await?;
100 self.queue_expired_telegram_sessions(&histories, Utc::now())
101 .await?;
102 self.sync_group_updates().await?;
103 self.sync_group_ingress().await?;
104 self.sync_telegram_events().await?;
105 self.schedule_wakeup_job(&histories).await;
106 anyhow::Ok(())
107 }
108 .await;
109 match result {
110 Ok(()) => *self.last_poll_error.write().await = None,
111 Err(error) => {
112 let message = error.to_string();
113 let mut previous = self.last_poll_error.write().await;
114 if previous.as_deref() != Some(message.as_str()) {
115 tracing::warn!(error=%error, "Telegram session runtime poll will retry");
116 *previous = Some(message);
117 }
118 }
119 }
120 tokio::time::sleep(POLL_INTERVAL).await;
121 }
122 }
123
124 async fn schedule_wakeup_job(self: &Arc<Self>, histories: &[SessionRecord]) {
125 if self.writer_job_active.load(Ordering::Acquire) {
126 return;
127 }
128 let Some(record) = histories
129 .iter()
130 .find(|record| record.phase == "active" && session_type(record) == "wakeup")
131 .cloned()
132 else {
133 return;
134 };
135 if self
136 .writer_job_active
137 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
138 .is_err()
139 {
140 return;
141 }
142 let worker = self.clone();
143 tokio::spawn(async move {
144 let writer = worker.control.writer().clone();
145 let _writer_guard = writer.lock().await;
146 let id = record.id;
147 let result = async {
148 let Some(record) = worker.get_listed_conversation(&id).await? else {
149 return Ok(());
150 };
151 worker.process_wakeup(record).await
152 }
153 .await;
154 if let Err(error) = result {
155 tracing::warn!(error=%bounded_error(&error), "Scheduled wakeup will retry");
156 }
157 worker.writer_job_active.store(false, Ordering::Release);
158 });
159 }
160
161 async fn run_wakeup_scheduler(self: Arc<Self>) {
162 loop {
163 let marker = next_wakeup_marker(Utc::now());
164 let delay = (marker - Utc::now()).to_std().unwrap_or(Duration::ZERO);
165 tokio::time::sleep(delay).await;
166 if let Err(error) = self.create_wakeup_sessions(marker).await {
167 tracing::warn!(
168 marker=%marker.to_rfc3339(),
169 error=%error,
170 "Scheduled wakeup session creation failed; this marker will not be retried"
171 );
172 }
173 }
174 }
175
176 async fn create_wakeup_sessions(&self, marker: DateTime<Utc>) -> anyhow::Result<()> {
177 let private_sessions = self.control.api().telegram_private_sessions().await?;
178 for private_session in private_sessions {
179 let telegram_user_id = private_session.telegram_user_id;
180 if let Err(error) = self.create_wakeup_session(telegram_user_id, marker).await {
181 tracing::warn!(
182 %telegram_user_id,
183 marker=%marker.to_rfc3339(),
184 error=%error,
185 "Could not create this user's scheduled wakeup session"
186 );
187 }
188 }
189 Ok(())
190 }
191
192 async fn create_wakeup_session(
193 &self,
194 telegram_user_id: i64,
195 marker: DateTime<Utc>,
196 ) -> anyhow::Result<()> {
197 let runtime = self.runtime()?.clone();
198 let user = self.roots.ensure_user(telegram_user_id).await?;
199 let user_root = user
200 .root_node_id
201 .context("Telegram user root is not ready for a wakeup session")?;
202 let mut options = SessionOptions::conversation(
203 "wakeup",
204 vec![user_root, runtime.kennedy_root_node_id.clone()],
205 );
206 options.mode = AgentMode::Wakeup;
207 options.channel = json!({
208 "kind":"wakeup",
209 "telegramUserId":telegram_user_id,
210 "username":user.current_username.or(Some(user.handle)),
211 "displayName":user.display_name,
212 "wakeupMarker":marker.to_rfc3339(),
213 });
214 options.orchestration = json!({"owner":"backend","status":"scheduled"});
215 let mut session = self.open_session(runtime, options, None).await?;
216 session.stage_wakeup_opening()?;
217 let state = session.snapshot()?;
218 self.control
219 .api()
220 .history_register(kcode_session_history::RegisterSession {
221 id: required_string(&state, "sessionId")?,
222 started_at: session.started_at.clone(),
223 state,
224 })
225 .await?;
226 Ok(())
227 }
228
229 async fn process_wakeup(&self, record: SessionRecord) -> anyhow::Result<()> {
230 let id = record.id.clone();
231 let mut session = self.session_for_record(&record).await?;
232 session.stage_wakeup_opening()?;
233 let record = Arc::new(Mutex::new(record));
234 persist_record(self.control.api(), &record, session.snapshot()?, false).await?;
235 let api = self.control.api().clone();
236 let saved = record.clone();
237 let completion = self
238 .run_session_turn(&id, &mut session, Uuid::new_v4(), move |state| {
239 let api = api.clone();
240 let record = saved.clone();
241 async move {
242 persist_record(&api, &record, state, false).await?;
243 Ok(())
244 }
245 })
246 .await?;
247 if matches!(completion, TurnCompletion::Stopped) {
248 session.interrupt_current_turn()?;
249 }
250 session.commit_current_write_session()?;
251 persist_record(self.control.api(), &record, session.snapshot()?, false).await?;
252 session.release_managed_sources().await;
253 let mut locked = record.lock().await;
254 let completed = self
255 .control
256 .api()
257 .history_complete(
258 &id,
259 kcode_session_history::Checkpoint {
260 expected_version: locked.version,
261 state: locked.state.clone(),
262 user_activity: false,
263 },
264 )
265 .await?;
266 *locked = completed;
267 Ok(())
268 }
269
270 async fn directory_user(&self, event: &Value) -> anyhow::Result<kcode_telegram_identity::User> {
271 let id = event
272 .get("telegramUserId")
273 .map(value_string)
274 .context("Telegram event omitted user ID")?
275 .parse::<i64>()
276 .context("Telegram event has an invalid user ID")?;
277 self.roots.ensure_user(id).await
278 }
279
280 async fn directory_group(
281 &self,
282 group_id: &str,
283 ) -> anyhow::Result<kcode_telegram_identity::Group> {
284 self.roots.ensure_group(group_id).await
285 }
286
287 async fn decorate_group_context(
288 &self,
289 mut context: Value,
290 group_id: &str,
291 ) -> anyhow::Result<Value> {
292 let group = self.directory_group(group_id).await?;
293 context["groupId"] = json!(group_id);
294 context["groupRootNodeId"] = json!(group.root_node_id);
295 context["groupRootReady"] = json!(group.root_ready);
296 let mut participants = Vec::new();
297 for participant in context
298 .get("participants")
299 .and_then(Value::as_array)
300 .cloned()
301 .unwrap_or_default()
302 {
303 let user = self.directory_user(&participant).await?;
304 let mut participant = participant;
305 participant["rootNodeId"] = json!(user.root_node_id);
306 participant["rootReady"] = json!(user.root_ready);
307 participants.push(participant);
308 }
309 context["participants"] = json!(participants);
310 Ok(context)
311 }
312
313 async fn prepare_group_context(
314 &self,
315 mut context: Value,
316 excluded_message_id: Option<&str>,
317 group_id: &str,
318 ) -> anyhow::Result<Value> {
319 let chat_id = context
320 .get("chatId")
321 .and_then(Value::as_i64)
322 .context("Telegram group context omitted its numeric chat ID")?;
323 let mut messages = Vec::new();
324 for mut message in context
325 .get("messages")
326 .and_then(Value::as_array)
327 .cloned()
328 .unwrap_or_default()
329 {
330 let message_id = message
331 .get("messageId")
332 .map(value_string)
333 .unwrap_or_default();
334 let numeric_message_id = message
335 .get("messageId")
336 .and_then(Value::as_i64)
337 .context("Telegram group context message omitted its numeric message ID")?;
338 let excluded = excluded_message_id == Some(message_id.as_str());
339 let kind = message
340 .get("kind")
341 .and_then(Value::as_str)
342 .unwrap_or("text")
343 .to_owned();
344 if !excluded
345 && message.get("sentByKennedy").and_then(Value::as_bool) != Some(true)
346 && kind == "document"
347 && message
348 .get("preparedText")
349 .and_then(Value::as_str)
350 .is_none()
351 && message.get("hasMedia").and_then(Value::as_bool) == Some(true)
352 {
353 let prepared = async {
354 let (bytes, mime) = self
355 .control
356 .api()
357 .telegram_group_message_media(chat_id, numeric_message_id)?;
358 let result = self
359 .control
360 .api()
361 .extract_document(
362 bytes,
363 message
364 .get("fileName")
365 .and_then(Value::as_str)
366 .unwrap_or("telegram-document")
367 .to_owned(),
368 &mime,
369 )
370 .await?;
371 Ok::<_, anyhow::Error>((
372 result.text,
373 None::<String>,
374 Some(result.format),
375 result.truncated,
376 ))
377 }
378 .await;
379 let (text, model, format, truncated) = match prepared {
380 Ok(value) => value,
381 Err(error) => (
382 format!("Document extraction failed: {error}"),
383 Some("preparation-error".into()),
384 None,
385 false,
386 ),
387 };
388 message["preparedText"] = json!(text);
389 message["preparationModel"] = json!(model);
390 message["documentFormat"] = json!(format);
391 message["preparationTruncated"] = json!(truncated);
392 let _ = self
393 .control
394 .api()
395 .telegram_save_group_message_preparation(
396 chat_id,
397 numeric_message_id,
398 &text,
399 model.as_deref(),
400 format.as_deref(),
401 truncated,
402 )
403 .await;
404 }
405 if !excluded
406 && matches!(
407 kind.as_str(),
408 "voice"
409 | "document"
410 | "photo"
411 | "video"
412 | "animation"
413 | "audio"
414 | "video_note"
415 | "sticker"
416 )
417 {
418 let has_media = message.get("hasMedia").and_then(Value::as_bool) == Some(true);
419 if has_media {
420 let (size_bytes, downloaded_mime_type) = self
421 .control
422 .api()
423 .telegram_group_message_media_metadata(chat_id, numeric_message_id)?;
424 let mime_type = normalized_file_mime_type(
425 message
426 .get("mimeType")
427 .and_then(Value::as_str)
428 .unwrap_or(&downloaded_mime_type),
429 );
430 let supplied_file_name = message
431 .get("fileName")
432 .and_then(Value::as_str)
433 .filter(|value| !value.trim().is_empty());
434 let file_name_supplied = supplied_file_name.is_some();
435 let file_name = telegram_group_context_file_name(
436 supplied_file_name,
437 &kind,
438 &mime_type,
439 &message_id,
440 );
441 message["fileName"] = json!(file_name);
442 message["fileNameSource"] = json!(if file_name_supplied {
443 "transport"
444 } else {
445 "synthesized"
446 });
447 message["mimeType"] = json!(mime_type);
448 message["sizeBytes"] = json!(size_bytes);
449 message["mediaRef"] = json!({"kind":kind,"source":"telegram-group","chatId":context.get("chatId").cloned().unwrap_or(Value::Null),"messageId":message.get("messageId").cloned().unwrap_or(Value::Null),"fileName":file_name,"fileNameSource":message.get("fileNameSource").cloned().unwrap_or(Value::Null),"mimeType":mime_type,"sizeBytes":size_bytes,"durationSeconds":message.get("durationSeconds").cloned().unwrap_or(Value::Null)});
450 }
451 let base = message.get("text").and_then(Value::as_str).unwrap_or("");
452 let prepared = message
453 .get("preparedText")
454 .and_then(Value::as_str)
455 .unwrap_or("Document text extraction unavailable.");
456 message["text"] = json!(if has_media {
457 let file_metadata = telegram_group_context_file_metadata(&message);
458 if kind == "voice" {
459 format!(
460 "{base}\n\n{file_metadata}\nThe voice note was not automatically transcribed."
461 )
462 } else if kind == "document" {
463 format!("{base}\n\n{file_metadata}\n\n{prepared}")
464 } else {
465 format!("{base}\n\n{file_metadata}")
466 }
467 } else {
468 format!("{base}\n\n[The Telegram {kind} file is unavailable.]")
469 });
470 }
471 messages.push(message);
472 }
473 context["messages"] = json!(messages);
474 self.decorate_group_context(context, group_id).await
475 }
476
477 async fn sync_telegram_events(self: &Arc<Self>) -> anyhow::Result<()> {
478 let events = self
479 .control
480 .api()
481 .telegram_events()
482 .await?
483 .get("events")
484 .and_then(Value::as_array)
485 .cloned()
486 .unwrap_or_default();
487 let listed_ids = events
488 .iter()
489 .filter_map(|event| event.get("id").and_then(Value::as_str))
490 .map(str::to_owned)
491 .collect::<HashSet<_>>();
492 self.event_retries
493 .lock()
494 .await
495 .retain(|id, _| listed_ids.contains(id));
496 for event in events {
497 let id = required_string(&event, "id")?;
498 if self
499 .event_retries
500 .lock()
501 .await
502 .get(&id)
503 .is_some_and(|retry| Instant::now() < retry.not_before)
504 {
505 continue;
506 }
507 let mut set = self.events_in_flight.lock().await;
508 if !set.insert(id.clone()) {
509 continue;
510 }
511 drop(set);
512 let worker = self.clone();
513 tokio::spawn(async move {
514 worker.run_telegram_event(event).await;
515 worker.events_in_flight.lock().await.remove(&id);
516 });
517 }
518 Ok(())
519 }
520
521 async fn run_telegram_event(&self, event: Value) {
522 let id = event
523 .get("id")
524 .and_then(Value::as_str)
525 .unwrap_or("unknown")
526 .to_owned();
527 let operation_id = Uuid::new_v4();
528 let initial_conversation_id = event
529 .get("conversationId")
530 .and_then(Value::as_str)
531 .map(str::to_owned);
532 if let Some(conversation_id) = &initial_conversation_id {
533 self.register_operation(conversation_id, operation_id).await;
534 }
535 let conversation_id = Arc::new(Mutex::new(initial_conversation_id.clone()));
536 let result = tokio::time::timeout(
537 telegram_timeout(&event),
538 self.process_telegram_event(&event, operation_id, conversation_id.clone()),
539 )
540 .await;
541 if let Some(conversation_id) = conversation_id.lock().await.clone() {
542 self.remove_operation(&conversation_id, operation_id).await;
543 }
544 if let Some(initial_conversation_id) = initial_conversation_id {
545 self.remove_operation(&initial_conversation_id, operation_id)
546 .await;
547 }
548 match result {
549 Ok(Ok(())) => {
550 self.event_retries.lock().await.remove(&id);
551 }
552 Ok(Err(error)) => {
553 let message = bounded_error(&error);
554 let (attempt, delay, should_warn) =
555 self.record_telegram_event_retry(&id, &message).await;
556 if should_warn {
557 tracing::warn!(
558 event_id=%id,
559 attempt,
560 retry_in_seconds=delay.as_secs(),
561 error=%message,
562 "Telegram event will retry"
563 );
564 } else {
565 tracing::debug!(
566 event_id=%id,
567 attempt,
568 retry_in_seconds=delay.as_secs(),
569 error=%message,
570 "Telegram event retry remains unsuccessful"
571 );
572 }
573 }
574 Err(_) => {
575 let _ = self.control.api().cancel_intelligence(operation_id);
576 let conversation = conversation_id.lock().await.clone();
577 if let Some(conversation_id) = &conversation
578 && let Err(error) = self
579 .transition_timed_out_telegram_to_ingress(conversation_id)
580 .await
581 {
582 let message = bounded_error(&error);
583 let (attempt, delay, should_warn) =
584 self.record_telegram_event_retry(&id, &message).await;
585 if should_warn {
586 tracing::warn!(
587 event_id=%id,
588 attempt,
589 retry_in_seconds=delay.as_secs(),
590 error=%message,
591 "Timed-out Telegram event retained until its conversation can be queued for ingress"
592 );
593 } else {
594 tracing::debug!(
595 event_id=%id,
596 attempt,
597 retry_in_seconds=delay.as_secs(),
598 error=%message,
599 "Timed-out Telegram ingress handoff remains unsuccessful"
600 );
601 }
602 return;
603 }
604 self.event_retries.lock().await.remove(&id);
605 let _ = self
606 .control
607 .api()
608 .telegram_abort_event(&id, conversation.as_deref(), TELEGRAM_TIMEOUT_NOTICE)
609 .await;
610 tracing::error!(event_id=%id,"Telegram event reached its 90-minute deadline and was aborted");
611 }
612 }
613 }
614
615 async fn record_telegram_event_retry(&self, id: &str, error: &str) -> (u32, Duration, bool) {
616 let mut retries = self.event_retries.lock().await;
617 let failures = retries
618 .get(id)
619 .map_or(1, |retry| retry.failures.saturating_add(1));
620 let delay = telegram_event_retry_delay(failures);
621 let should_warn = telegram_event_retry_should_warn(
622 retries.get(id).map(|retry| retry.last_error.as_str()),
623 error,
624 failures,
625 );
626 retries.insert(
627 id.to_owned(),
628 TelegramEventRetry {
629 failures,
630 not_before: Instant::now() + delay,
631 last_error: error.to_owned(),
632 },
633 );
634 (failures, delay, should_warn)
635 }
636
637 async fn process_telegram_event(
638 &self,
639 event: &Value,
640 operation_id: Uuid,
641 bound_conversation_id: Arc<Mutex<Option<String>>>,
642 ) -> anyhow::Result<()> {
643 let id = required_string(event, "id")?;
644 let _private_user_guard =
645 if event.get("sessionKind").and_then(Value::as_str) != Some("group") {
646 let telegram_user_id = event
647 .get("telegramUserId")
648 .and_then(Value::as_i64)
649 .context("private Telegram event is missing its numeric user identity")?;
650 Some(
651 self.control
652 .api()
653 .telegram_user_lock(telegram_user_id)
654 .await
655 .lock_owned()
656 .await,
657 )
658 } else {
659 None
660 };
661 self.directory_user(event).await?;
662 if event.get("kind").and_then(Value::as_str) == Some("reset") {
663 return self.process_telegram_reset(event).await;
664 }
665 let (record_arc, _) = self.telegram_session(event).await?;
666 let conversation_id = {
667 let locked = record_arc.lock().await;
668 locked.id.clone()
669 };
670 *bound_conversation_id.lock().await = Some(conversation_id.clone());
671 self.register_operation(&conversation_id, operation_id)
672 .await;
673 let lock = self.conversation_lock(&conversation_id).await;
674 let _guard = lock.lock().await;
675 let mut session = {
676 let record = record_arc.lock().await;
677 self.session_for_record(&record).await?
678 };
679 if session.answer_for_external_event(&id).is_none() {
680 if session.pending_turn && session.pending_external_event_id.as_deref() != Some(&id) {
681 anyhow::bail!("This Telegram session has an earlier saved query to finish.");
682 }
683 if !session.pending_turn {
684 let input = self.telegram_input(event).await;
685 let (text, metadata) = match input {
686 Ok(input) => input,
687 Err(error) if event.get("kind").and_then(Value::as_str) == Some("document") => {
688 let filename = event
689 .get("fileName")
690 .and_then(Value::as_str)
691 .unwrap_or("that document");
692 self.control.api()
693 .telegram_reply_event(
694 &id,
695 &conversation_id,
696 &format!(
697 "I couldn't read {filename}: {error} Please try sending it again."
698 ),
699 None,
700 )
701 .await?;
702 return Ok(());
703 }
704 Err(error) => return Err(error),
705 };
706 session.begin_user_turn(&text, &metadata);
707 persist_record(self.control.api(), &record_arc, session.snapshot()?, true).await?;
708 }
709 let api = self.control.api().clone();
710 let saved = record_arc.clone();
711 let completion = self
712 .run_session_turn(&conversation_id, &mut session, operation_id, move |state| {
713 let api = api.clone();
714 let record = saved.clone();
715 async move {
716 persist_record(&api, &record, state, false).await?;
717 Ok(())
718 }
719 })
720 .await?;
721 if matches!(completion, TurnCompletion::Stopped) {
722 session.interrupt_current_turn()?;
723 persist_record(self.control.api(), &record_arc, session.snapshot()?, false).await?;
724 self.control
725 .api()
726 .telegram_interrupt_event(&id, &conversation_id)
727 .await?;
728 self.complete_pending_stop(
729 &conversation_id,
730 json!({"status":"stopped","scope":"turn"}),
731 )
732 .await?;
733 return Ok(());
734 }
735 persist_record(self.control.api(), &record_arc, session.snapshot()?, false).await?;
736 if session.requires_history_ingress() {
737 session.orchestration =
738 json!({"owner":"backend","status":"ending","reason":"context-limit"});
739 persist_record(self.control.api(), &record_arc, session.snapshot()?, false).await?;
740 self.request_conversation_ingress(&record_arc, None).await?;
741 self.deliver_telegram_responses(&mut session, &id, &conversation_id)
742 .await?;
743 self.complete_pending_stop(
744 &conversation_id,
745 json!({"status":"already-completed","scope":"turn"}),
746 )
747 .await?;
748 return Ok(());
749 }
750 }
751 self.deliver_telegram_responses(&mut session, &id, &conversation_id)
752 .await?;
753 self.complete_pending_stop(
754 &conversation_id,
755 json!({"status":"already-completed","scope":"turn"}),
756 )
757 .await?;
758 Ok(())
759 }
760
761 async fn deliver_telegram_responses(
762 &self,
763 session: &mut Session,
764 event_id: &str,
765 conversation_id: &str,
766 ) -> anyhow::Result<()> {
767 let mut deliveries = Vec::new();
768 for response in session.responses_for_external_event(event_id) {
769 for (object_id, file_name) in telegram_response_object_deliveries(response) {
770 deliveries.push(TelegramDelivery::Object {
771 object_id,
772 file_name,
773 });
774 }
775 if let Some(text) = response
776 .get("content")
777 .and_then(Value::as_str)
778 .filter(|text| !text.is_empty())
779 {
780 deliveries.push(TelegramDelivery::Text {
781 text: text.to_owned(),
782 response_warning: response
783 .get("contextWarning")
784 .cloned()
785 .unwrap_or(Value::Null),
786 captionable: response.get("role").and_then(Value::as_str) == Some("kennedy"),
787 });
788 }
789 }
790 anyhow::ensure!(
791 !deliveries.is_empty(),
792 "Kennedy completed the turn without a recoverable Telegram response"
793 );
794 let delivery_count = deliveries.len();
795 let mut index = 0;
796 while index < delivery_count {
797 match &deliveries[index] {
798 TelegramDelivery::Object {
799 object_id,
800 file_name,
801 } => {
802 let mut file = session.resolve_object(object_id)?;
803 if let Some(file_name) = file_name {
804 validate_delivery_file_name(file_name)?;
805 file.file_name = file_name.clone();
806 }
807 anyhow::ensure!(
808 file.bytes.len() <= self.config.telegram_max_media_bytes,
809 "object {object_id} is {} bytes, over the configured {}-byte Telegram media limit",
810 file.bytes.len(),
811 self.config.telegram_max_media_bytes
812 );
813 let caption = telegram_reply_caption(&deliveries, index, &file);
814 let complete = caption.is_some() || index + 1 == delivery_count;
815 self.control
816 .api()
817 .telegram_send_object(event_id, conversation_id, &file, caption, complete)
818 .await?;
819 index += if caption.is_some() { 2 } else { 1 };
820 }
821 TelegramDelivery::Text {
822 text,
823 response_warning,
824 ..
825 } => {
826 self.control
827 .api()
828 .telegram_reply_event(
829 event_id,
830 conversation_id,
831 text,
832 response_warning.as_str(),
833 )
834 .await?;
835 index += 1;
836 }
837 }
838 }
839 Ok(())
840 }
841
842 async fn transition_timed_out_telegram_to_ingress(
843 &self,
844 conversation_id: &str,
845 ) -> anyhow::Result<()> {
846 let record = match self.get_conversation(conversation_id).await {
847 Ok(record) => record,
848 Err(error)
849 if error
850 .downcast_ref::<ApiError>()
851 .is_some_and(|error| error.code == "not_found") =>
852 {
853 return Ok(());
854 }
855 Err(error) => return Err(error),
856 };
857 if record.phase != "active" {
858 return Ok(());
859 }
860 let mut state = record.state.clone();
861 state["orchestration"] =
862 json!({"owner":"backend","status":"stopped","reason":"telegram-timeout"});
863 self.control
864 .api()
865 .history_request_ingress(
866 conversation_id,
867 kcode_session_history::Checkpoint {
868 expected_version: record.version,
869 state: state.clone(),
870 user_activity: false,
871 },
872 )
873 .await?;
874 if let Some(session_id) = state.get("rustLibSessionId").and_then(Value::as_str) {
875 self.control.api().release_managed_sources(session_id).await;
876 }
877 Ok(())
878 }
879
880 async fn queue_detached_private_telegram_sessions(&self) -> anyhow::Result<()> {
881 let bound = self
882 .control
883 .api()
884 .telegram_private_sessions()
885 .await?
886 .into_iter()
887 .filter_map(|session| session.current_conversation_id)
888 .collect::<HashSet<_>>();
889 let histories = self.list_history().await?;
890 for record in histories.iter().filter(|record| {
891 record.phase == "active"
892 && session_type(record) == "telegram"
893 && !bound.contains(&record.id)
894 }) {
895 self.queue_telegram_session_for_ingress(record, "telegram-detached")
896 .await?;
897 }
898 Ok(())
899 }
900
901 async fn queue_expired_telegram_sessions(
902 &self,
903 histories: &[SessionRecord],
904 now: DateTime<Utc>,
905 ) -> anyhow::Result<()> {
906 for record in histories
907 .iter()
908 .filter(|record| telegram_session_is_expired(record, now))
909 {
910 self.queue_telegram_session_for_ingress(record, "telegram-session-timeout")
911 .await?;
912 }
913 Ok(())
914 }
915
916 async fn queue_telegram_session_for_ingress(
917 &self,
918 summary: &SessionRecord,
919 reason: &str,
920 ) -> anyhow::Result<bool> {
921 let id = summary.id.clone();
922 let lock = self.conversation_lock(&id).await;
923 let _guard = lock.lock().await;
924 let record = self.get_conversation(&id).await?;
925 if record.phase != "active"
926 || !matches!(
927 session_type(&record).as_str(),
928 "telegram" | "telegram-group"
929 )
930 {
931 return Ok(false);
932 }
933 if reason == "telegram-session-timeout" && !telegram_session_is_expired(&record, Utc::now())
934 {
935 return Ok(false);
936 }
937 let mut state = record.state.clone();
938 state["orchestration"] = json!({
939 "owner":"backend",
940 "status":"stopped",
941 "reason":reason,
942 });
943 self.control
944 .api()
945 .history_request_ingress(
946 &id,
947 kcode_session_history::Checkpoint {
948 expected_version: record.version,
949 state: state.clone(),
950 user_activity: false,
951 },
952 )
953 .await?;
954 if let Some(session_id) = state.get("rustLibSessionId").and_then(Value::as_str) {
955 self.control.api().release_managed_sources(session_id).await;
956 }
957 tracing::info!(session_id=%id, %reason, "Queued Telegram session for history ingress");
958 Ok(true)
959 }
960
961 async fn telegram_session(
962 &self,
963 event: &Value,
964 ) -> anyhow::Result<(Arc<Mutex<SessionRecord>>, Session)> {
965 let histories = self.list_history().await?;
966 let group = event.get("sessionKind").and_then(Value::as_str) == Some("group");
967 let user_id = event
968 .get("telegramUserId")
969 .map(value_string)
970 .unwrap_or_default();
971 let group_id = event.get("groupId").and_then(Value::as_str);
972 let mut record = event
973 .get("conversationId")
974 .and_then(Value::as_str)
975 .and_then(|id| histories.iter().find(|record| record.id == id).cloned());
976 if record.is_none() {
977 record = histories.into_iter().find(|record| {
978 record.phase == "active"
979 && if group {
980 session_type(record) == "telegram-group"
981 && record_group_id(record) == group_id
982 && record_user_id(record) == user_id
983 } else {
984 session_type(record) == "telegram" && record_user_id(record) == user_id
985 }
986 });
987 }
988 let created = record
989 .as_ref()
990 .is_none_or(|record| record.phase != "active");
991 let (record, mut session) =
992 if let Some(record) = record.filter(|record| record.phase == "active") {
993 let record = self.get_conversation(&record.id).await?;
994 let session = self.session_for_record(&record).await?;
995 (record, session)
996 } else {
997 self.create_telegram_session(event).await?
998 };
999 let record = Arc::new(Mutex::new(record));
1000 let id = {
1001 let locked = record.lock().await;
1002 locked.id.clone()
1003 };
1004 if event.get("conversationId").and_then(Value::as_str) != Some(&id)
1005 || event.get("processingStartedAt").is_none()
1006 {
1007 self.control
1008 .api()
1009 .telegram_bind_event(
1010 &required_string(event, "id")?,
1011 &id,
1012 event.get("conversationId").and_then(Value::as_str),
1013 )
1014 .await?;
1015 }
1016 if group && !created {
1017 let group_id = required_string(event, "groupId")?;
1018 if let Some(context) = event.get("groupContext") {
1019 let context = self
1020 .prepare_group_context(
1021 context.clone(),
1022 event.get("messageId").map(value_string).as_deref(),
1023 &group_id,
1024 )
1025 .await?;
1026 session.refresh_telegram_group_context(
1027 &context,
1028 event.get("messageId").map(value_string).as_deref(),
1029 )?;
1030 persist_record(self.control.api(), &record, session.snapshot()?, false).await?;
1031 }
1032 }
1033 Ok((record, session))
1034 }
1035
1036 async fn create_telegram_session(
1037 &self,
1038 event: &Value,
1039 ) -> anyhow::Result<(SessionRecord, Session)> {
1040 let runtime = self.runtime()?.clone();
1041 let user = self.directory_user(event).await?;
1042 let group = event.get("sessionKind").and_then(Value::as_str) == Some("group");
1043 let mut roots = vec![
1044 user.root_node_id
1045 .context("Telegram user root is not ready")?,
1046 ];
1047 let mut channel = json!({"kind":if group{"telegram-group"}else{"telegram"},"telegramUserId":event.get("telegramUserId").cloned().unwrap_or(Value::Null),"chatId":event.get("chatId").cloned().unwrap_or(Value::Null),"groupId":event.get("groupId").cloned().unwrap_or(Value::Null),"username":event.get("username").cloned().unwrap_or(Value::Null),"displayName":event.get("displayName").cloned().unwrap_or(Value::Null),"maxObjectBytes":self.config.telegram_max_media_bytes});
1048 let mut references = Vec::new();
1049 if group {
1050 let group_id = required_string(event, "groupId")?;
1051 let group_record = self.directory_group(&group_id).await?;
1052 let group_root = group_record
1053 .root_node_id
1054 .context("Telegram group root is not ready")?;
1055 roots.push(group_root.clone());
1056 if let Some(context) = event.get("groupContext") {
1057 let context = self
1058 .prepare_group_context(
1059 context.clone(),
1060 event.get("messageId").map(value_string).as_deref(),
1061 &group_id,
1062 )
1063 .await?;
1064 channel["groupContext"] = context.clone();
1065 channel["groupRootNodeId"] = json!(group_root);
1066 references = participant_references(&context, &roots);
1067 }
1068 }
1069 roots.push(runtime.kennedy_root_node_id.clone());
1070 references.retain(|id| !roots.contains(id));
1071 let mut options =
1072 SessionOptions::conversation(if group { "telegram-group" } else { "telegram" }, roots);
1073 options.channel = channel;
1074 options.reference_root_node_ids = references;
1075 let session = self.open_session(runtime, options, None).await?;
1076 let state = session.snapshot()?;
1077 let record = self
1078 .control
1079 .api()
1080 .history_register(kcode_session_history::RegisterSession {
1081 id: required_string(&state, "sessionId")?,
1082 started_at: session.started_at.clone(),
1083 state,
1084 })
1085 .await?;
1086 Ok((record, session))
1087 }
1088
1089 async fn telegram_input(&self, event: &Value) -> anyhow::Result<(String, Value)> {
1090 let Some(batch) = event
1091 .get("batchedEvents")
1092 .and_then(Value::as_array)
1093 .filter(|batch| batch.len() > 1)
1094 else {
1095 return self.telegram_event_input(event).await;
1096 };
1097 let mut inputs = Vec::with_capacity(batch.len());
1098 for batched_event in batch {
1099 inputs.push(self.telegram_event_input(batched_event).await?);
1100 }
1101 merge_telegram_batch_inputs(event, inputs)
1102 }
1103
1104 async fn telegram_event_input(&self, event: &Value) -> anyhow::Result<(String, Value)> {
1105 let id = required_string(event, "id")?;
1106 match event.get("kind").and_then(Value::as_str).unwrap_or("text") {
1107 "voice" => {
1108 let (bytes, mime) = self.control.api().telegram_event_media(&id)?;
1109 let filename = event
1110 .get("fileName")
1111 .and_then(Value::as_str)
1112 .filter(|value| !value.trim().is_empty())
1113 .unwrap_or("telegram-voice.ogg");
1114 Ok((
1115 event
1116 .get("text")
1117 .and_then(Value::as_str)
1118 .unwrap_or("")
1119 .into(),
1120 json!({"externalEventId":id,"inputKind":"voice","media":{"id":format!("telegram:{id}"),"kind":"voice","source":"telegram","mimeType":mime,"fileName":filename,"dataUrl":data_url(&mime,&bytes),"sizeBytes":bytes.len(),"durationSeconds":event.get("durationSeconds").cloned().unwrap_or(Value::Null)}}),
1121 ))
1122 }
1123 "document" => {
1124 let (bytes, mime) = self.control.api().telegram_event_media(&id)?;
1125 let filename = event
1126 .get("fileName")
1127 .and_then(Value::as_str)
1128 .unwrap_or("telegram-document")
1129 .to_owned();
1130 let extraction = self
1131 .control
1132 .api()
1133 .extract_document(bytes.clone(), filename.clone(), &mime)
1134 .await;
1135 let mut attachment = json!({
1136 "id":format!("telegram:{id}"),
1137 "kind":"document",
1138 "source":"telegram",
1139 "fileName":filename,
1140 "mimeType":mime,
1141 "sizeBytes":bytes.len(),
1142 "dataUrl":data_url(&mime,&bytes),
1143 });
1144 match extraction {
1145 Ok(result) => {
1146 attachment["format"] = json!(result.format);
1147 attachment["text"] = json!(result.text);
1148 attachment["characters"] = json!(result.characters);
1149 attachment["truncated"] = json!(result.truncated);
1150 }
1151 Err(error) => {
1152 attachment["extractionError"] = json!(error.to_string());
1153 }
1154 }
1155 Ok((
1156 event
1157 .get("text")
1158 .and_then(Value::as_str)
1159 .unwrap_or("")
1160 .into(),
1161 json!({"externalEventId":id,"inputKind":"document","attachments":[attachment]}),
1162 ))
1163 }
1164 kind @ ("photo" | "video" | "animation" | "audio" | "video_note" | "sticker") => {
1165 let (bytes, downloaded_mime) = self.control.api().telegram_event_media(&id)?;
1166 let mime = event
1167 .get("mimeType")
1168 .and_then(Value::as_str)
1169 .filter(|value| !value.trim().is_empty())
1170 .unwrap_or(&downloaded_mime)
1171 .to_owned();
1172 let extension = match kind {
1173 "photo" => "jpg",
1174 "video" | "video_note" => "mp4",
1175 "animation" => "gif",
1176 "audio" => "mp3",
1177 "sticker" => "webp",
1178 _ => "bin",
1179 };
1180 let filename = event
1181 .get("fileName")
1182 .and_then(Value::as_str)
1183 .filter(|value| !value.trim().is_empty())
1184 .map(str::to_owned)
1185 .unwrap_or_else(|| format!("telegram-{kind}.{extension}"));
1186 let mut attachment = json!({
1187 "id":format!("telegram:{id}"),
1188 "kind":kind,
1189 "source":"telegram",
1190 "fileName":filename,
1191 "mimeType":mime,
1192 "sizeBytes":bytes.len(),
1193 "dataUrl":data_url(&mime,&bytes),
1194 });
1195 if let Some(value) = event.get("durationSeconds") {
1196 attachment["durationSeconds"] = value.clone();
1197 }
1198 Ok((
1199 event
1200 .get("text")
1201 .and_then(Value::as_str)
1202 .unwrap_or("")
1203 .into(),
1204 json!({
1205 "externalEventId":id,
1206 "inputKind":kind,
1207 "attachments":[attachment],
1208 }),
1209 ))
1210 }
1211 _ => Ok((
1212 event
1213 .get("text")
1214 .and_then(Value::as_str)
1215 .unwrap_or("")
1216 .into(),
1217 json!({"externalEventId":id,"inputKind":"text"}),
1218 )),
1219 }
1220 }
1221
1222 async fn process_telegram_reset(&self, event: &Value) -> anyhow::Result<()> {
1223 let id = required_string(event, "id")?;
1224 let Some(conversation_id) = event.get("conversationId").and_then(Value::as_str) else {
1225 self.control.api()
1226 .telegram_complete_reset(
1227 &id,
1228 Some(
1229 "There is no active Telegram session to reset. Your next message will begin one.",
1230 ),
1231 )
1232 .await?;
1233 return Ok(());
1234 };
1235 let record = match self.get_conversation(conversation_id).await {
1236 Ok(record) => record,
1237 Err(error)
1238 if error
1239 .downcast_ref::<ApiError>()
1240 .is_some_and(|error| error.code == "not_found") =>
1241 {
1242 self.control.api()
1243 .telegram_complete_reset(
1244 &id,
1245 Some(
1246 "There is no active Telegram session to reset. Your next message will begin one.",
1247 ),
1248 )
1249 .await?;
1250 return Ok(());
1251 }
1252 Err(error) => return Err(error),
1253 };
1254 if record.phase != "active" {
1255 self.control.api()
1256 .telegram_complete_reset(
1257 &id,
1258 Some(
1259 "There is no active Telegram session to reset. Your next message will begin one.",
1260 ),
1261 )
1262 .await?;
1263 return Ok(());
1264 }
1265 let session = self.session_for_record(&record).await?;
1266 session.release_managed_sources().await;
1267 self.control
1268 .api()
1269 .history_request_ingress(
1270 conversation_id,
1271 kcode_session_history::Checkpoint {
1272 expected_version: record.version,
1273 state: record.state,
1274 user_activity: false,
1275 },
1276 )
1277 .await?;
1278 self.control.api()
1279 .telegram_complete_reset(
1280 &id,
1281 Some(
1282 "Conversation reset. The Telegram session has been queued for memory ingress; your next message will begin a new session.",
1283 ),
1284 )
1285 .await?;
1286 Ok(())
1287 }
1288
1289 async fn sync_group_updates(self: &Arc<Self>) -> anyhow::Result<()> {
1290 let updates = self
1291 .control
1292 .api()
1293 .telegram_group_session_updates()
1294 .await?
1295 .get("updates")
1296 .and_then(Value::as_array)
1297 .cloned()
1298 .unwrap_or_default();
1299 for update in updates {
1300 let id = required_string(&update, "conversationId")?;
1301 let mut set = self.group_updates_in_flight.lock().await;
1302 if !set.insert(id.clone()) {
1303 continue;
1304 }
1305 drop(set);
1306 let worker = self.clone();
1307 tokio::spawn(async move {
1308 if let Err(error) = worker.process_group_update(update).await {
1309 tracing::warn!(conversation_id=%id,error=%error,"Telegram group context update will retry");
1310 }
1311 worker.group_updates_in_flight.lock().await.remove(&id);
1312 });
1313 }
1314 Ok(())
1315 }
1316
1317 async fn process_group_update(&self, update: Value) -> anyhow::Result<()> {
1318 let id = required_string(&update, "conversationId")?;
1319 let lock = self.conversation_lock(&id).await;
1320 let _guard = lock.lock().await;
1321 let record = match self.get_conversation(&id).await {
1322 Ok(record) => record,
1323 Err(error)
1324 if error
1325 .downcast_ref::<ApiError>()
1326 .is_some_and(|error| error.code == "not_found") =>
1327 {
1328 self.reconcile_missing_group_session(&update, &id).await?;
1329 return Ok(());
1330 }
1331 Err(error) => return Err(error),
1332 };
1333 if record.phase != "active" {
1334 if update.get("resetRequired").and_then(Value::as_bool) == Some(true) {
1335 self.control
1336 .api()
1337 .telegram_complete_silent_group_reset(&id)
1338 .await?;
1339 }
1340 return Ok(());
1341 }
1342 let mut session = self.session_for_record(&record).await?;
1343 if session.pending_turn {
1344 return Ok(());
1345 }
1346 let group_id = required_string(&update, "groupId")?;
1347 let context = self
1348 .prepare_group_context(
1349 update.get("groupContext").cloned().unwrap_or(Value::Null),
1350 None,
1351 &group_id,
1352 )
1353 .await?;
1354 session.refresh_telegram_group_context(&context, None)?;
1355 let record = Arc::new(Mutex::new(record));
1356 persist_record(self.control.api(), &record, session.snapshot()?, false).await?;
1357 if update.get("resetRequired").and_then(Value::as_bool) == Some(true) {
1358 self.close_conversation(&record, &session).await?;
1359 self.control
1360 .api()
1361 .telegram_complete_silent_group_reset(&id)
1362 .await?;
1363 } else {
1364 self.control
1365 .api()
1366 .telegram_acknowledge_group_context(
1367 &id,
1368 update
1369 .get("throughMessageId")
1370 .and_then(Value::as_i64)
1371 .unwrap_or(0),
1372 )
1373 .await?;
1374 }
1375 Ok(())
1376 }
1377
1378 async fn reconcile_missing_group_session(
1379 &self,
1380 update: &Value,
1381 conversation_id: &str,
1382 ) -> anyhow::Result<()> {
1383 match missing_group_session_recovery(update)? {
1384 MissingGroupSessionRecovery::CompleteSilentReset => {
1385 self.control
1386 .api()
1387 .telegram_complete_silent_group_reset(conversation_id)
1388 .await?;
1389 tracing::info!(
1390 %conversation_id,
1391 "Completed orphaned Telegram group reset"
1392 );
1393 }
1394 MissingGroupSessionRecovery::DetachCurrent {
1395 group_id,
1396 telegram_user_id,
1397 } => {
1398 let result = self
1399 .control
1400 .api()
1401 .telegram_detach_group_session(conversation_id, &group_id, telegram_user_id)
1402 .await;
1403 match result {
1404 Ok(_) => {
1405 tracing::info!(
1406 %conversation_id,
1407 %group_id,
1408 telegram_user_id,
1409 "Detached orphaned Telegram group session"
1410 );
1411 }
1412 Err(error) if error.code == "state_conflict" => {
1413 tracing::info!(
1414 %conversation_id,
1415 %group_id,
1416 telegram_user_id,
1417 "Telegram group session was already detached or rebound"
1418 );
1419 }
1420 Err(error) => return Err(error.into()),
1421 }
1422 }
1423 }
1424 Ok(())
1425 }
1426
1427 async fn sync_group_ingress(self: &Arc<Self>) -> anyhow::Result<()> {
1428 let batches = self
1429 .control
1430 .api()
1431 .telegram_group_ingress()
1432 .await?
1433 .get("batches")
1434 .and_then(Value::as_array)
1435 .cloned()
1436 .unwrap_or_default();
1437 for batch in batches {
1438 let id = required_string(&batch, "id")?;
1439 let mut set = self.group_ingress_in_flight.lock().await;
1440 if !set.insert(id.clone()) {
1441 continue;
1442 }
1443 drop(set);
1444 let worker = self.clone();
1445 tokio::spawn(async move {
1446 if let Err(error) = worker.process_group_ingress(batch).await {
1447 tracing::warn!(batch_id=%id,error=%error,"Telegram group ingress preparation will retry");
1448 }
1449 worker.group_ingress_in_flight.lock().await.remove(&id);
1450 });
1451 }
1452 Ok(())
1453 }
1454
1455 async fn process_group_ingress(&self, batch: Value) -> anyhow::Result<()> {
1456 let runtime = self.runtime()?.clone();
1457 let id = required_string(&batch, "id")?;
1458 if let Some(existing) = self.list_history().await?.into_iter().find(|record| {
1459 record
1460 .state
1461 .get("channel")
1462 .and_then(|channel| channel.get("groupIngressBatchId"))
1463 .and_then(Value::as_str)
1464 == Some(&id)
1465 }) {
1466 match existing.phase.as_str() {
1467 "complete" => {
1468 self.control
1469 .api()
1470 .telegram_complete_group_ingress(&id)
1471 .await?;
1472 }
1473 "active" => {
1474 let existing = self.get_conversation(&existing.id).await?;
1475 self.control
1476 .api()
1477 .history_request_ingress(
1478 &existing.id,
1479 kcode_session_history::Checkpoint {
1480 expected_version: existing.version,
1481 state: existing.state,
1482 user_activity: false,
1483 },
1484 )
1485 .await?;
1486 }
1487 _ => {}
1488 }
1489 return Ok(());
1490 }
1491 let group_id = required_string(&batch, "groupId")?;
1492 let group = self.directory_group(&group_id).await?;
1493 let raw_context = json!({"groupTitle":batch.get("groupTitle").cloned().unwrap_or(json!("Telegram group")),"chatId":batch.get("chatId").cloned().unwrap_or(Value::Null),"participants":batch.get("participants").cloned().unwrap_or(json!([])),"messages":batch.get("messages").cloned().unwrap_or(json!([]))});
1494 let context = self
1495 .prepare_group_context(raw_context, None, &group_id)
1496 .await?;
1497 let group_root = group
1498 .root_node_id
1499 .context("Telegram group root is not ready")?;
1500 let roots = vec![group_root.clone(), runtime.kennedy_root_node_id.clone()];
1501 let channel = json!({"kind":"telegram-group","chatId":batch.get("chatId").cloned().unwrap_or(Value::Null),"groupId":group_id,"groupRootNodeId":group_root,"groupIngressBatchId":id,"backgroundIngress":true,"groupContext":context});
1502 let mut options = SessionOptions::conversation("telegram-group", roots.clone());
1503 options.channel = channel;
1504 options.reference_root_node_ids = participant_references(&context, &roots);
1505 options.source_session_type = Some("telegram-group".into());
1506 let mut session = self.open_session(runtime, options, None).await?;
1507 for message in context
1508 .get("messages")
1509 .and_then(Value::as_array)
1510 .into_iter()
1511 .flatten()
1512 {
1513 session.stage_source_message(
1514 message.get("sentByKennedy").and_then(Value::as_bool) == Some(true),
1515 message
1516 .get("text")
1517 .and_then(Value::as_str)
1518 .unwrap_or_default(),
1519 message.clone(),
1520 )?;
1521 }
1522 let state = session.snapshot()?;
1523 let record = self
1524 .control
1525 .api()
1526 .history_register(kcode_session_history::RegisterSession {
1527 id: required_string(&state, "sessionId")?,
1528 started_at: batch
1529 .get("createdAt")
1530 .and_then(Value::as_str)
1531 .map(str::to_owned)
1532 .unwrap_or_else(|| Utc::now().to_rfc3339()),
1533 state,
1534 })
1535 .await?;
1536 self.control
1537 .api()
1538 .history_request_ingress(
1539 &record.id,
1540 kcode_session_history::Checkpoint {
1541 expected_version: record.version,
1542 state: record.state,
1543 user_activity: false,
1544 },
1545 )
1546 .await?;
1547 Ok(())
1548 }
1549}
1550
1551impl Deref for Runtime {
1552 type Target = Orchestrator;
1553
1554 fn deref(&self) -> &Self::Target {
1555 &self.control
1556 }
1557}
1558fn session_type(record: &SessionRecord) -> String {
1559 record
1560 .state
1561 .get("sessionType")
1562 .and_then(Value::as_str)
1563 .unwrap_or("conversation")
1564 .into()
1565}
1566
1567fn next_wakeup_marker(now: DateTime<Utc>) -> DateTime<Utc> {
1568 let day_start = now
1569 .date_naive()
1570 .and_hms_opt(0, 0, 0)
1571 .expect("midnight is always a valid UTC time")
1572 .and_utc();
1573 for hour in [0_i64, 4, 8, 12, 16, 20] {
1574 let candidate = day_start + ChronoDuration::hours(hour);
1575 if candidate > now {
1576 return candidate;
1577 }
1578 }
1579 day_start + ChronoDuration::days(1)
1580}
1581
1582fn telegram_session_is_expired(record: &SessionRecord, now: DateTime<Utc>) -> bool {
1583 if record.phase != "active"
1584 || !matches!(session_type(record).as_str(), "telegram" | "telegram-group")
1585 || record.state.get("pendingTurn").and_then(Value::as_bool) == Some(true)
1586 {
1587 return false;
1588 }
1589 DateTime::parse_from_rfc3339(&record.started_at)
1590 .ok()
1591 .is_some_and(|started| now >= started.with_timezone(&Utc) + TELEGRAM_SESSION_MAX_AGE)
1592}
1593fn record_channel(record: &SessionRecord) -> Option<&Value> {
1594 record.state.get("channel")
1595}
1596fn record_group_id(record: &SessionRecord) -> Option<&str> {
1597 record_channel(record)
1598 .and_then(|channel| {
1599 channel.get("groupId").or_else(|| {
1600 channel
1601 .get("groupContext")
1602 .and_then(|context| context.get("groupId"))
1603 })
1604 })
1605 .and_then(Value::as_str)
1606}
1607fn record_user_id(record: &SessionRecord) -> String {
1608 record_channel(record)
1609 .and_then(|channel| channel.get("telegramUserId"))
1610 .map(value_string)
1611 .unwrap_or_default()
1612}
1613fn participant_references(context: &Value, roots: &[String]) -> Vec<String> {
1614 let mut values = context
1615 .get("participants")
1616 .and_then(Value::as_array)
1617 .into_iter()
1618 .flatten()
1619 .filter_map(|participant| participant.get("rootNodeId").and_then(Value::as_str))
1620 .filter(|id| !roots.iter().any(|root| root == id))
1621 .map(str::to_owned)
1622 .collect::<Vec<_>>();
1623 values.sort();
1624 values.dedup();
1625 values
1626}
1627fn normalized_file_mime_type(value: &str) -> String {
1628 let value = value
1629 .split(';')
1630 .next()
1631 .unwrap_or(value)
1632 .trim()
1633 .to_ascii_lowercase();
1634 if value.contains('/')
1635 && !value.is_empty()
1636 && !value.chars().any(char::is_whitespace)
1637 && !value.chars().any(char::is_control)
1638 {
1639 value
1640 } else {
1641 "application/octet-stream".into()
1642 }
1643}
1644fn file_extension_for_mime_type(mime_type: &str) -> &'static str {
1645 match normalized_file_mime_type(mime_type).as_str() {
1646 "image/jpeg" => "jpg",
1647 "image/png" => "png",
1648 "image/webp" => "webp",
1649 "image/gif" => "gif",
1650 "audio/ogg" | "audio/opus" | "application/ogg" => "ogg",
1651 "audio/mpeg" | "audio/mp3" => "mp3",
1652 "audio/mp4" | "video/mp4" => "mp4",
1653 "audio/webm" | "video/webm" => "webm",
1654 "audio/wav" | "audio/x-wav" => "wav",
1655 "application/pdf" => "pdf",
1656 _ => "bin",
1657 }
1658}
1659fn telegram_group_context_file_name(
1660 supplied: Option<&str>,
1661 kind: &str,
1662 mime_type: &str,
1663 message_id: &str,
1664) -> String {
1665 let fallback = format!(
1666 "telegram-group-{kind}-{message_id}.{}",
1667 file_extension_for_mime_type(mime_type)
1668 );
1669 sanitize_file_name(supplied.unwrap_or_default(), &fallback)
1670}
1671fn file_name_extension(file_name: &str) -> String {
1672 file_name
1673 .rsplit_once('.')
1674 .and_then(|(stem, extension)| {
1675 (!stem.is_empty() && !extension.is_empty()).then_some(extension)
1676 })
1677 .map(|extension| format!(".{extension}"))
1678 .unwrap_or_else(|| "(none)".into())
1679}
1680fn telegram_group_context_file_metadata(message: &Value) -> String {
1681 let file_name = message
1682 .get("fileName")
1683 .and_then(Value::as_str)
1684 .unwrap_or("telegram-file");
1685 let source_note =
1686 if message.get("fileNameSource").and_then(Value::as_str) == Some("synthesized") {
1687 " (synthesized because Telegram supplied no filename)"
1688 } else {
1689 ""
1690 };
1691 let mime_type = normalized_file_mime_type(
1692 message
1693 .get("mimeType")
1694 .and_then(Value::as_str)
1695 .unwrap_or("application/octet-stream"),
1696 );
1697 let size_bytes = message
1698 .get("sizeBytes")
1699 .and_then(Value::as_u64)
1700 .unwrap_or_default();
1701 format!(
1702 "User-provided file\nOriginal filename: {file_name}{source_note}\nExtension: {}\nMIME type: {mime_type}\nSize: {size_bytes} bytes",
1703 file_name_extension(file_name),
1704 )
1705}
1706
1707fn telegram_response_object_deliveries(response: &Value) -> Vec<(String, Option<String>)> {
1708 let objects = response
1709 .get("objects")
1710 .and_then(Value::as_array)
1711 .map(Vec::as_slice)
1712 .unwrap_or_default();
1713 let attachments = response
1714 .get("attachments")
1715 .and_then(Value::as_array)
1716 .map(Vec::as_slice)
1717 .unwrap_or_default();
1718 objects
1719 .iter()
1720 .filter_map(Value::as_str)
1721 .enumerate()
1722 .map(|(index, object_id)| {
1723 let descriptor = attachments
1724 .iter()
1725 .find(|candidate| {
1726 ["objectId", "pendingId", "id"]
1727 .iter()
1728 .any(|key| candidate.get(key).and_then(Value::as_str) == Some(object_id))
1729 })
1730 .or_else(|| attachments.get(index));
1731 let file_name = descriptor
1732 .and_then(|descriptor| descriptor.get("fileName"))
1733 .and_then(Value::as_str)
1734 .map(str::to_owned);
1735 (object_id.to_owned(), file_name)
1736 })
1737 .collect()
1738}
1739
1740fn telegram_reply_caption<'a>(
1741 deliveries: &'a [TelegramDelivery],
1742 object_index: usize,
1743 file: &ResolvedObject,
1744) -> Option<&'a str> {
1745 if object_index + 2 != deliveries.len() {
1746 return None;
1747 }
1748 match &deliveries[object_index + 1] {
1749 TelegramDelivery::Text {
1750 text,
1751 response_warning,
1752 captionable: true,
1753 } if response_warning.is_null() || response_warning.as_str() == Some("") => {
1754 telegram_caption_for(file, text)
1755 }
1756 _ => None,
1757 }
1758}
1759
1760fn required_string(value: &Value, key: &str) -> anyhow::Result<String> {
1761 value
1762 .get(key)
1763 .and_then(Value::as_str)
1764 .filter(|value| !value.is_empty())
1765 .map(str::to_owned)
1766 .with_context(|| format!("backend response omitted {key}"))
1767}
1768
1769fn merge_telegram_batch_inputs(
1770 event: &Value,
1771 inputs: Vec<(String, Value)>,
1772) -> anyhow::Result<(String, Value)> {
1773 let mut texts = Vec::new();
1774 let mut attachments = Vec::new();
1775 let mut event_ids = Vec::with_capacity(inputs.len());
1776 for (text, metadata) in inputs {
1777 if !text.is_empty() {
1778 texts.push(text);
1779 }
1780 if let Some(media) = metadata.get("media").filter(|value| value.is_object()) {
1781 attachments.push(media.clone());
1782 }
1783 if let Some(items) = metadata.get("attachments").and_then(Value::as_array) {
1784 attachments.extend(items.iter().cloned());
1785 }
1786 if let Some(id) = metadata.get("externalEventId").and_then(Value::as_str) {
1787 event_ids.push(id.to_owned());
1788 }
1789 }
1790 Ok((
1791 texts.join("\n\n"),
1792 json!({
1793 "externalEventId":required_string(event, "id")?,
1794 "externalEventIds":event_ids,
1795 "inputKind":"batch",
1796 "attachments":attachments,
1797 }),
1798 ))
1799}
1800
1801fn missing_group_session_recovery(update: &Value) -> anyhow::Result<MissingGroupSessionRecovery> {
1802 if update.get("resetRequired").and_then(Value::as_bool) == Some(true) {
1803 return Ok(MissingGroupSessionRecovery::CompleteSilentReset);
1804 }
1805 Ok(MissingGroupSessionRecovery::DetachCurrent {
1806 group_id: required_string(update, "groupId")?,
1807 telegram_user_id: update
1808 .get("telegramUserId")
1809 .and_then(Value::as_i64)
1810 .context("backend response omitted telegramUserId")?,
1811 })
1812}
1813fn value_string(value: &Value) -> String {
1814 value
1815 .as_str()
1816 .map(str::to_owned)
1817 .unwrap_or_else(|| value.to_string())
1818}
1819fn bounded_error(error: &anyhow::Error) -> String {
1820 format!("{error:#}").chars().take(1_000).collect()
1821}
1822fn telegram_event_retry_delay(failures: u32) -> Duration {
1823 let exponent = failures.saturating_sub(1).min(5);
1824 Duration::from_secs((2_u64 << exponent).min(60))
1825}
1826fn telegram_event_retry_should_warn(
1827 previous_error: Option<&str>,
1828 error: &str,
1829 failures: u32,
1830) -> bool {
1831 previous_error != Some(error) || failures.is_multiple_of(10)
1832}
1833fn telegram_timeout(event: &Value) -> Duration {
1834 let elapsed = event
1835 .get("processingStartedAt")
1836 .and_then(Value::as_str)
1837 .and_then(|value| DateTime::parse_from_rfc3339(value).ok())
1838 .map(|value| {
1839 (Utc::now() - value.with_timezone(&Utc))
1840 .to_std()
1841 .unwrap_or_default()
1842 })
1843 .unwrap_or_default();
1844 TELEGRAM_TIMEOUT.saturating_sub(elapsed)
1845}
1846
1847#[cfg(test)]
1848mod tests {
1849 use super::*;
1850
1851 fn session_record(overrides: Value) -> SessionRecord {
1852 let mut record = json!({
1853 "id":"session",
1854 "phase":"active",
1855 "started_at":"2026-07-30T00:00:00Z",
1856 "updated_at":"2026-07-30T00:00:00Z",
1857 "state":{},
1858 "provenance_id":null,
1859 "version":1,
1860 "last_user_message_at":null,
1861 "ended_at":null,
1862 "ingress_failure_count":0,
1863 "ingress_failures":[],
1864 "ingress_next_attempt_at":null
1865 });
1866 record
1867 .as_object_mut()
1868 .unwrap()
1869 .extend(overrides.as_object().unwrap().clone());
1870 serde_json::from_value(record).unwrap()
1871 }
1872
1873 #[test]
1874 fn batches_keep_order_text_and_every_attachment() {
1875 let (text, metadata) = merge_telegram_batch_inputs(
1876 &json!({"id":"batch"}),
1877 vec![
1878 ("first".into(), json!({"externalEventId":"one"})),
1879 (
1880 String::new(),
1881 json!({"externalEventId":"two","media":{"id":"voice"}}),
1882 ),
1883 (
1884 "third".into(),
1885 json!({"externalEventId":"three","attachments":[{"id":"document"}]}),
1886 ),
1887 ],
1888 )
1889 .unwrap();
1890 assert_eq!(text, "first\n\nthird");
1891 assert_eq!(metadata["externalEventIds"], json!(["one", "two", "three"]));
1892 assert_eq!(
1893 metadata["attachments"],
1894 json!([{"id":"voice"}, {"id":"document"}])
1895 );
1896 }
1897
1898 #[test]
1899 fn wakeups_use_strictly_future_four_hour_utc_boundaries() {
1900 let before = DateTime::parse_from_rfc3339("2026-07-28T03:59:59Z")
1901 .unwrap()
1902 .with_timezone(&Utc);
1903 assert_eq!(
1904 next_wakeup_marker(before).to_rfc3339(),
1905 "2026-07-28T04:00:00+00:00"
1906 );
1907 let exactly = DateTime::parse_from_rfc3339("2026-07-28T20:00:00Z")
1908 .unwrap()
1909 .with_timezone(&Utc);
1910 assert_eq!(
1911 next_wakeup_marker(exactly).to_rfc3339(),
1912 "2026-07-29T00:00:00+00:00"
1913 );
1914 }
1915
1916 #[test]
1917 fn idle_telegram_sessions_roll_over_after_six_hours() {
1918 let now = DateTime::parse_from_rfc3339("2026-07-30T12:00:00Z")
1919 .unwrap()
1920 .with_timezone(&Utc);
1921 assert!(telegram_session_is_expired(
1922 &session_record(json!({
1923 "started_at":"2026-07-30T06:00:00Z",
1924 "state":{"sessionType":"telegram","pendingTurn":false}
1925 })),
1926 now
1927 ));
1928 assert!(!telegram_session_is_expired(
1929 &session_record(json!({
1930 "started_at":"2026-07-30T05:00:00Z",
1931 "state":{"sessionType":"telegram","pendingTurn":true}
1932 })),
1933 now
1934 ));
1935 }
1936
1937 #[test]
1938 fn missing_group_sessions_choose_reset_or_detach() {
1939 assert_eq!(
1940 missing_group_session_recovery(&json!({"resetRequired":true})).unwrap(),
1941 MissingGroupSessionRecovery::CompleteSilentReset
1942 );
1943 assert_eq!(
1944 missing_group_session_recovery(&json!({
1945 "groupId":"group-1","telegramUserId":42,"resetRequired":false
1946 }))
1947 .unwrap(),
1948 MissingGroupSessionRecovery::DetachCurrent {
1949 group_id: "group-1".into(),
1950 telegram_user_id: 42,
1951 }
1952 );
1953 }
1954
1955 #[test]
1956 fn response_objects_keep_their_filename_overrides() {
1957 assert_eq!(
1958 telegram_response_object_deliveries(&json!({
1959 "objects":["pending:2","AAECAwQF"],
1960 "attachments":[
1961 {"objectId":"AAECAwQF","fileName":"canonical.pdf"},
1962 {"objectId":"pending:2","fileName":"draft.pdf"}
1963 ]
1964 })),
1965 vec![
1966 ("pending:2".into(), Some("draft.pdf".into())),
1967 ("AAECAwQF".into(), Some("canonical.pdf".into())),
1968 ]
1969 );
1970 }
1971
1972 #[test]
1973 fn exact_final_text_is_used_only_as_a_supported_caption() {
1974 let file = ResolvedObject {
1975 object_id: "object".into(),
1976 bytes: vec![1],
1977 file_name: "photo.jpg".into(),
1978 media_type: "image/jpeg".into(),
1979 transport_kind: Some("photo".into()),
1980 };
1981 let deliveries = vec![
1982 TelegramDelivery::Object {
1983 object_id: "object".into(),
1984 file_name: None,
1985 },
1986 TelegramDelivery::Text {
1987 text: " exact caption\n".into(),
1988 response_warning: Value::Null,
1989 captionable: true,
1990 },
1991 ];
1992 assert_eq!(
1993 telegram_reply_caption(&deliveries, 0, &file),
1994 Some(" exact caption\n")
1995 );
1996 }
1997
1998 #[test]
1999 fn event_retries_back_off_and_repeat_warnings_periodically() {
2000 assert_eq!(telegram_event_retry_delay(1), Duration::from_secs(2));
2001 assert_eq!(telegram_event_retry_delay(6), Duration::from_secs(60));
2002 assert_eq!(
2003 telegram_event_retry_delay(u32::MAX),
2004 Duration::from_secs(60)
2005 );
2006 assert!(!telegram_event_retry_should_warn(
2007 Some("failure"),
2008 "failure",
2009 2
2010 ));
2011 assert!(telegram_event_retry_should_warn(
2012 Some("failure"),
2013 "failure",
2014 10
2015 ));
2016 }
2017}