Skip to main content

kcode_kennedy_telegram_runtime/
lib.rs

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}