use std::{
collections::{HashMap, HashSet},
ops::Deref,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
time::{Duration, Instant},
};
use anyhow::Context as _;
use chrono::{DateTime, Duration as ChronoDuration, Utc};
use kcode_kennedy_orchestration::{
AgentMode, ApiError, Orchestrator, Session, TurnCompletion, data_url, persist_record,
telegram_caption_for,
};
use kcode_kennedy_roots::DirectoryRoots;
use kcode_kennedy_sessions::{ResolvedObject, SessionOptions, validate_delivery_file_name};
use kcode_server_object_envelopes::sanitize_file_name;
use kcode_session_history::SessionRecord;
use serde_json::{Value, json};
use tokio::sync::{Mutex, RwLock};
use uuid::Uuid;
const POLL_INTERVAL: Duration = Duration::from_secs(1);
const TELEGRAM_TIMEOUT: Duration = Duration::from_secs(90 * 60);
const TELEGRAM_SESSION_MAX_AGE: ChronoDuration = ChronoDuration::hours(6);
const 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.";
#[derive(Clone, Debug)]
pub struct Config {
pub telegram_max_media_bytes: usize,
pub telegram_web_user_handle: String,
}
#[derive(Debug, Eq, PartialEq)]
enum MissingGroupSessionRecovery {
CompleteSilentReset,
DetachCurrent {
group_id: String,
telegram_user_id: i64,
},
}
struct TelegramEventRetry {
failures: u32,
not_before: Instant,
last_error: String,
}
enum TelegramDelivery {
Object {
object_id: String,
file_name: Option<String>,
},
Text {
text: String,
response_warning: Value,
captionable: bool,
},
}
pub struct Runtime {
config: Config,
control: Arc<Orchestrator>,
roots: DirectoryRoots,
writer_job_active: AtomicBool,
events_in_flight: Mutex<HashSet<String>>,
event_retries: Mutex<HashMap<String, TelegramEventRetry>>,
group_updates_in_flight: Mutex<HashSet<String>>,
group_ingress_in_flight: Mutex<HashSet<String>>,
last_poll_error: RwLock<Option<String>>,
}
impl Runtime {
pub fn new(config: Config, control: Arc<Orchestrator>, roots: DirectoryRoots) -> Self {
Self {
config,
control,
roots,
writer_job_active: AtomicBool::new(false),
events_in_flight: Mutex::new(HashSet::new()),
event_retries: Mutex::new(HashMap::new()),
group_updates_in_flight: Mutex::new(HashSet::new()),
group_ingress_in_flight: Mutex::new(HashSet::new()),
last_poll_error: RwLock::new(None),
}
}
pub async fn run(self: Arc<Self>) -> anyhow::Result<()> {
self.initialize_until_ready().await;
self.control.api().telegram_health();
self.queue_detached_private_telegram_sessions().await?;
let wakeups = self.clone();
tokio::spawn(async move { wakeups.run_wakeup_scheduler().await });
loop {
let result = async {
self.roots.reconcile_pending().await?;
let histories = self.list_history().await?;
self.queue_expired_telegram_sessions(&histories, Utc::now())
.await?;
self.sync_group_updates().await?;
self.sync_group_ingress().await?;
self.sync_telegram_events().await?;
self.schedule_wakeup_job(&histories).await;
anyhow::Ok(())
}
.await;
match result {
Ok(()) => *self.last_poll_error.write().await = None,
Err(error) => {
let message = error.to_string();
let mut previous = self.last_poll_error.write().await;
if previous.as_deref() != Some(message.as_str()) {
tracing::warn!(error=%error, "Telegram session runtime poll will retry");
*previous = Some(message);
}
}
}
tokio::time::sleep(POLL_INTERVAL).await;
}
}
async fn schedule_wakeup_job(self: &Arc<Self>, histories: &[SessionRecord]) {
if self.writer_job_active.load(Ordering::Acquire) {
return;
}
let Some(record) = histories
.iter()
.find(|record| record.phase == "active" && session_type(record) == "wakeup")
.cloned()
else {
return;
};
if self
.writer_job_active
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return;
}
let worker = self.clone();
tokio::spawn(async move {
let writer = worker.control.writer().clone();
let _writer_guard = writer.lock().await;
let id = record.id;
let result = async {
let Some(record) = worker.get_listed_conversation(&id).await? else {
return Ok(());
};
worker.process_wakeup(record).await
}
.await;
if let Err(error) = result {
tracing::warn!(error=%bounded_error(&error), "Scheduled wakeup will retry");
}
worker.writer_job_active.store(false, Ordering::Release);
});
}
async fn run_wakeup_scheduler(self: Arc<Self>) {
loop {
let marker = next_wakeup_marker(Utc::now());
let delay = (marker - Utc::now()).to_std().unwrap_or(Duration::ZERO);
tokio::time::sleep(delay).await;
if let Err(error) = self.create_wakeup_sessions(marker).await {
tracing::warn!(
marker=%marker.to_rfc3339(),
error=%error,
"Scheduled wakeup session creation failed; this marker will not be retried"
);
}
}
}
async fn create_wakeup_sessions(&self, marker: DateTime<Utc>) -> anyhow::Result<()> {
let private_sessions = self.control.api().telegram_private_sessions().await?;
for private_session in private_sessions {
let telegram_user_id = private_session.telegram_user_id;
if let Err(error) = self.create_wakeup_session(telegram_user_id, marker).await {
tracing::warn!(
%telegram_user_id,
marker=%marker.to_rfc3339(),
error=%error,
"Could not create this user's scheduled wakeup session"
);
}
}
Ok(())
}
async fn create_wakeup_session(
&self,
telegram_user_id: i64,
marker: DateTime<Utc>,
) -> anyhow::Result<()> {
let runtime = self.runtime()?.clone();
let user = self.roots.ensure_user(telegram_user_id).await?;
let user_root = user
.root_node_id
.context("Telegram user root is not ready for a wakeup session")?;
let mut options = SessionOptions::conversation(
"wakeup",
vec![user_root, runtime.kennedy_root_node_id.clone()],
);
options.mode = AgentMode::Wakeup;
options.channel = json!({
"kind":"wakeup",
"telegramUserId":telegram_user_id,
"username":user.current_username.or(Some(user.handle)),
"displayName":user.display_name,
"wakeupMarker":marker.to_rfc3339(),
});
options.orchestration = json!({"owner":"backend","status":"scheduled"});
let mut session = self.open_session(runtime, options, None).await?;
session.stage_wakeup_opening()?;
let state = session.snapshot()?;
self.control
.api()
.history_register(kcode_session_history::RegisterSession {
id: required_string(&state, "sessionId")?,
started_at: session.started_at.clone(),
state,
})
.await?;
Ok(())
}
async fn process_wakeup(&self, record: SessionRecord) -> anyhow::Result<()> {
let id = record.id.clone();
let mut session = self.session_for_record(&record).await?;
session.stage_wakeup_opening()?;
let record = Arc::new(Mutex::new(record));
persist_record(self.control.api(), &record, session.snapshot()?, false).await?;
let api = self.control.api().clone();
let saved = record.clone();
let completion = self
.run_session_turn(&id, &mut session, Uuid::new_v4(), move |state| {
let api = api.clone();
let record = saved.clone();
async move {
persist_record(&api, &record, state, false).await?;
Ok(())
}
})
.await?;
if matches!(completion, TurnCompletion::Stopped) {
session.interrupt_current_turn()?;
}
session.commit_current_write_session()?;
persist_record(self.control.api(), &record, session.snapshot()?, false).await?;
session.release_managed_sources().await;
let mut locked = record.lock().await;
let completed = self
.control
.api()
.history_complete(
&id,
kcode_session_history::Checkpoint {
expected_version: locked.version,
state: locked.state.clone(),
user_activity: false,
},
)
.await?;
*locked = completed;
Ok(())
}
async fn directory_user(&self, event: &Value) -> anyhow::Result<kcode_telegram_identity::User> {
let id = event
.get("telegramUserId")
.map(value_string)
.context("Telegram event omitted user ID")?
.parse::<i64>()
.context("Telegram event has an invalid user ID")?;
self.roots.ensure_user(id).await
}
async fn directory_group(
&self,
group_id: &str,
) -> anyhow::Result<kcode_telegram_identity::Group> {
self.roots.ensure_group(group_id).await
}
async fn decorate_group_context(
&self,
mut context: Value,
group_id: &str,
) -> anyhow::Result<Value> {
let group = self.directory_group(group_id).await?;
context["groupId"] = json!(group_id);
context["groupRootNodeId"] = json!(group.root_node_id);
context["groupRootReady"] = json!(group.root_ready);
let mut participants = Vec::new();
for participant in context
.get("participants")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default()
{
let user = self.directory_user(&participant).await?;
let mut participant = participant;
participant["rootNodeId"] = json!(user.root_node_id);
participant["rootReady"] = json!(user.root_ready);
participants.push(participant);
}
context["participants"] = json!(participants);
Ok(context)
}
async fn prepare_group_context(
&self,
mut context: Value,
excluded_message_id: Option<&str>,
group_id: &str,
) -> anyhow::Result<Value> {
let chat_id = context
.get("chatId")
.and_then(Value::as_i64)
.context("Telegram group context omitted its numeric chat ID")?;
let mut messages = Vec::new();
for mut message in context
.get("messages")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default()
{
let message_id = message
.get("messageId")
.map(value_string)
.unwrap_or_default();
let numeric_message_id = message
.get("messageId")
.and_then(Value::as_i64)
.context("Telegram group context message omitted its numeric message ID")?;
let excluded = excluded_message_id == Some(message_id.as_str());
let kind = message
.get("kind")
.and_then(Value::as_str)
.unwrap_or("text")
.to_owned();
if !excluded
&& message.get("sentByKennedy").and_then(Value::as_bool) != Some(true)
&& kind == "document"
&& message
.get("preparedText")
.and_then(Value::as_str)
.is_none()
&& message.get("hasMedia").and_then(Value::as_bool) == Some(true)
{
let prepared = async {
let (bytes, mime) = self
.control
.api()
.telegram_group_message_media(chat_id, numeric_message_id)?;
let result = self
.control
.api()
.extract_document(
bytes,
message
.get("fileName")
.and_then(Value::as_str)
.unwrap_or("telegram-document")
.to_owned(),
&mime,
)
.await?;
Ok::<_, anyhow::Error>((
result.text,
None::<String>,
Some(result.format),
result.truncated,
))
}
.await;
let (text, model, format, truncated) = match prepared {
Ok(value) => value,
Err(error) => (
format!("Document extraction failed: {error}"),
Some("preparation-error".into()),
None,
false,
),
};
message["preparedText"] = json!(text);
message["preparationModel"] = json!(model);
message["documentFormat"] = json!(format);
message["preparationTruncated"] = json!(truncated);
let _ = self
.control
.api()
.telegram_save_group_message_preparation(
chat_id,
numeric_message_id,
&text,
model.as_deref(),
format.as_deref(),
truncated,
)
.await;
}
if !excluded
&& matches!(
kind.as_str(),
"voice"
| "document"
| "photo"
| "video"
| "animation"
| "audio"
| "video_note"
| "sticker"
)
{
let has_media = message.get("hasMedia").and_then(Value::as_bool) == Some(true);
if has_media {
let (size_bytes, downloaded_mime_type) = self
.control
.api()
.telegram_group_message_media_metadata(chat_id, numeric_message_id)?;
let mime_type = normalized_file_mime_type(
message
.get("mimeType")
.and_then(Value::as_str)
.unwrap_or(&downloaded_mime_type),
);
let supplied_file_name = message
.get("fileName")
.and_then(Value::as_str)
.filter(|value| !value.trim().is_empty());
let file_name_supplied = supplied_file_name.is_some();
let file_name = telegram_group_context_file_name(
supplied_file_name,
&kind,
&mime_type,
&message_id,
);
message["fileName"] = json!(file_name);
message["fileNameSource"] = json!(if file_name_supplied {
"transport"
} else {
"synthesized"
});
message["mimeType"] = json!(mime_type);
message["sizeBytes"] = json!(size_bytes);
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)});
}
let base = message.get("text").and_then(Value::as_str).unwrap_or("");
let prepared = message
.get("preparedText")
.and_then(Value::as_str)
.unwrap_or("Document text extraction unavailable.");
message["text"] = json!(if has_media {
let file_metadata = telegram_group_context_file_metadata(&message);
if kind == "voice" {
format!(
"{base}\n\n{file_metadata}\nThe voice note was not automatically transcribed."
)
} else if kind == "document" {
format!("{base}\n\n{file_metadata}\n\n{prepared}")
} else {
format!("{base}\n\n{file_metadata}")
}
} else {
format!("{base}\n\n[The Telegram {kind} file is unavailable.]")
});
}
messages.push(message);
}
context["messages"] = json!(messages);
self.decorate_group_context(context, group_id).await
}
async fn sync_telegram_events(self: &Arc<Self>) -> anyhow::Result<()> {
let events = self
.control
.api()
.telegram_events()
.await?
.get("events")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default();
let listed_ids = events
.iter()
.filter_map(|event| event.get("id").and_then(Value::as_str))
.map(str::to_owned)
.collect::<HashSet<_>>();
self.event_retries
.lock()
.await
.retain(|id, _| listed_ids.contains(id));
for event in events {
let id = required_string(&event, "id")?;
if self
.event_retries
.lock()
.await
.get(&id)
.is_some_and(|retry| Instant::now() < retry.not_before)
{
continue;
}
let mut set = self.events_in_flight.lock().await;
if !set.insert(id.clone()) {
continue;
}
drop(set);
let worker = self.clone();
tokio::spawn(async move {
worker.run_telegram_event(event).await;
worker.events_in_flight.lock().await.remove(&id);
});
}
Ok(())
}
async fn run_telegram_event(&self, event: Value) {
let id = event
.get("id")
.and_then(Value::as_str)
.unwrap_or("unknown")
.to_owned();
let operation_id = Uuid::new_v4();
let initial_conversation_id = event
.get("conversationId")
.and_then(Value::as_str)
.map(str::to_owned);
if let Some(conversation_id) = &initial_conversation_id {
self.register_operation(conversation_id, operation_id).await;
}
let conversation_id = Arc::new(Mutex::new(initial_conversation_id.clone()));
let result = tokio::time::timeout(
telegram_timeout(&event),
self.process_telegram_event(&event, operation_id, conversation_id.clone()),
)
.await;
if let Some(conversation_id) = conversation_id.lock().await.clone() {
self.remove_operation(&conversation_id, operation_id).await;
}
if let Some(initial_conversation_id) = initial_conversation_id {
self.remove_operation(&initial_conversation_id, operation_id)
.await;
}
match result {
Ok(Ok(())) => {
self.event_retries.lock().await.remove(&id);
}
Ok(Err(error)) => {
let message = bounded_error(&error);
let (attempt, delay, should_warn) =
self.record_telegram_event_retry(&id, &message).await;
if should_warn {
tracing::warn!(
event_id=%id,
attempt,
retry_in_seconds=delay.as_secs(),
error=%message,
"Telegram event will retry"
);
} else {
tracing::debug!(
event_id=%id,
attempt,
retry_in_seconds=delay.as_secs(),
error=%message,
"Telegram event retry remains unsuccessful"
);
}
}
Err(_) => {
let _ = self.control.api().cancel_intelligence(operation_id);
let conversation = conversation_id.lock().await.clone();
if let Some(conversation_id) = &conversation
&& let Err(error) = self
.transition_timed_out_telegram_to_ingress(conversation_id)
.await
{
let message = bounded_error(&error);
let (attempt, delay, should_warn) =
self.record_telegram_event_retry(&id, &message).await;
if should_warn {
tracing::warn!(
event_id=%id,
attempt,
retry_in_seconds=delay.as_secs(),
error=%message,
"Timed-out Telegram event retained until its conversation can be queued for ingress"
);
} else {
tracing::debug!(
event_id=%id,
attempt,
retry_in_seconds=delay.as_secs(),
error=%message,
"Timed-out Telegram ingress handoff remains unsuccessful"
);
}
return;
}
self.event_retries.lock().await.remove(&id);
let _ = self
.control
.api()
.telegram_abort_event(&id, conversation.as_deref(), TELEGRAM_TIMEOUT_NOTICE)
.await;
tracing::error!(event_id=%id,"Telegram event reached its 90-minute deadline and was aborted");
}
}
}
async fn record_telegram_event_retry(&self, id: &str, error: &str) -> (u32, Duration, bool) {
let mut retries = self.event_retries.lock().await;
let failures = retries
.get(id)
.map_or(1, |retry| retry.failures.saturating_add(1));
let delay = telegram_event_retry_delay(failures);
let should_warn = telegram_event_retry_should_warn(
retries.get(id).map(|retry| retry.last_error.as_str()),
error,
failures,
);
retries.insert(
id.to_owned(),
TelegramEventRetry {
failures,
not_before: Instant::now() + delay,
last_error: error.to_owned(),
},
);
(failures, delay, should_warn)
}
async fn process_telegram_event(
&self,
event: &Value,
operation_id: Uuid,
bound_conversation_id: Arc<Mutex<Option<String>>>,
) -> anyhow::Result<()> {
let id = required_string(event, "id")?;
let _private_user_guard =
if event.get("sessionKind").and_then(Value::as_str) != Some("group") {
let telegram_user_id = event
.get("telegramUserId")
.and_then(Value::as_i64)
.context("private Telegram event is missing its numeric user identity")?;
Some(
self.control
.api()
.telegram_user_lock(telegram_user_id)
.await
.lock_owned()
.await,
)
} else {
None
};
self.directory_user(event).await?;
if event.get("kind").and_then(Value::as_str) == Some("reset") {
return self.process_telegram_reset(event).await;
}
let (record_arc, _) = self.telegram_session(event).await?;
let conversation_id = {
let locked = record_arc.lock().await;
locked.id.clone()
};
*bound_conversation_id.lock().await = Some(conversation_id.clone());
self.register_operation(&conversation_id, operation_id)
.await;
let lock = self.conversation_lock(&conversation_id).await;
let _guard = lock.lock().await;
let mut session = {
let record = record_arc.lock().await;
self.session_for_record(&record).await?
};
if session.answer_for_external_event(&id).is_none() {
if session.pending_turn && session.pending_external_event_id.as_deref() != Some(&id) {
anyhow::bail!("This Telegram session has an earlier saved query to finish.");
}
if !session.pending_turn {
let input = self.telegram_input(event).await;
let (text, metadata) = match input {
Ok(input) => input,
Err(error) if event.get("kind").and_then(Value::as_str) == Some("document") => {
let filename = event
.get("fileName")
.and_then(Value::as_str)
.unwrap_or("that document");
self.control.api()
.telegram_reply_event(
&id,
&conversation_id,
&format!(
"I couldn't read {filename}: {error} Please try sending it again."
),
None,
)
.await?;
return Ok(());
}
Err(error) => return Err(error),
};
session.begin_user_turn(&text, &metadata);
persist_record(self.control.api(), &record_arc, session.snapshot()?, true).await?;
}
let api = self.control.api().clone();
let saved = record_arc.clone();
let completion = self
.run_session_turn(&conversation_id, &mut session, operation_id, move |state| {
let api = api.clone();
let record = saved.clone();
async move {
persist_record(&api, &record, state, false).await?;
Ok(())
}
})
.await?;
if matches!(completion, TurnCompletion::Stopped) {
session.interrupt_current_turn()?;
persist_record(self.control.api(), &record_arc, session.snapshot()?, false).await?;
self.control
.api()
.telegram_interrupt_event(&id, &conversation_id)
.await?;
self.complete_pending_stop(
&conversation_id,
json!({"status":"stopped","scope":"turn"}),
)
.await?;
return Ok(());
}
persist_record(self.control.api(), &record_arc, session.snapshot()?, false).await?;
if session.requires_history_ingress() {
session.orchestration =
json!({"owner":"backend","status":"ending","reason":"context-limit"});
persist_record(self.control.api(), &record_arc, session.snapshot()?, false).await?;
self.request_conversation_ingress(&record_arc, None).await?;
self.deliver_telegram_responses(&mut session, &id, &conversation_id)
.await?;
self.complete_pending_stop(
&conversation_id,
json!({"status":"already-completed","scope":"turn"}),
)
.await?;
return Ok(());
}
}
self.deliver_telegram_responses(&mut session, &id, &conversation_id)
.await?;
self.complete_pending_stop(
&conversation_id,
json!({"status":"already-completed","scope":"turn"}),
)
.await?;
Ok(())
}
async fn deliver_telegram_responses(
&self,
session: &mut Session,
event_id: &str,
conversation_id: &str,
) -> anyhow::Result<()> {
let mut deliveries = Vec::new();
for response in session.responses_for_external_event(event_id) {
for (object_id, file_name) in telegram_response_object_deliveries(response) {
deliveries.push(TelegramDelivery::Object {
object_id,
file_name,
});
}
if let Some(text) = response
.get("content")
.and_then(Value::as_str)
.filter(|text| !text.is_empty())
{
deliveries.push(TelegramDelivery::Text {
text: text.to_owned(),
response_warning: response
.get("contextWarning")
.cloned()
.unwrap_or(Value::Null),
captionable: response.get("role").and_then(Value::as_str) == Some("kennedy"),
});
}
}
anyhow::ensure!(
!deliveries.is_empty(),
"Kennedy completed the turn without a recoverable Telegram response"
);
let delivery_count = deliveries.len();
let mut index = 0;
while index < delivery_count {
match &deliveries[index] {
TelegramDelivery::Object {
object_id,
file_name,
} => {
let mut file = session.resolve_object(object_id)?;
if let Some(file_name) = file_name {
validate_delivery_file_name(file_name)?;
file.file_name = file_name.clone();
}
anyhow::ensure!(
file.bytes.len() <= self.config.telegram_max_media_bytes,
"object {object_id} is {} bytes, over the configured {}-byte Telegram media limit",
file.bytes.len(),
self.config.telegram_max_media_bytes
);
let caption = telegram_reply_caption(&deliveries, index, &file);
let complete = caption.is_some() || index + 1 == delivery_count;
self.control
.api()
.telegram_send_object(event_id, conversation_id, &file, caption, complete)
.await?;
index += if caption.is_some() { 2 } else { 1 };
}
TelegramDelivery::Text {
text,
response_warning,
..
} => {
self.control
.api()
.telegram_reply_event(
event_id,
conversation_id,
text,
response_warning.as_str(),
)
.await?;
index += 1;
}
}
}
Ok(())
}
async fn transition_timed_out_telegram_to_ingress(
&self,
conversation_id: &str,
) -> anyhow::Result<()> {
let record = match self.get_conversation(conversation_id).await {
Ok(record) => record,
Err(error)
if error
.downcast_ref::<ApiError>()
.is_some_and(|error| error.code == "not_found") =>
{
return Ok(());
}
Err(error) => return Err(error),
};
if record.phase != "active" {
return Ok(());
}
let mut state = record.state.clone();
state["orchestration"] =
json!({"owner":"backend","status":"stopped","reason":"telegram-timeout"});
self.control
.api()
.history_request_ingress(
conversation_id,
kcode_session_history::Checkpoint {
expected_version: record.version,
state: state.clone(),
user_activity: false,
},
)
.await?;
if let Some(session_id) = state.get("rustLibSessionId").and_then(Value::as_str) {
self.control.api().release_managed_sources(session_id).await;
}
Ok(())
}
async fn queue_detached_private_telegram_sessions(&self) -> anyhow::Result<()> {
let bound = self
.control
.api()
.telegram_private_sessions()
.await?
.into_iter()
.filter_map(|session| session.current_conversation_id)
.collect::<HashSet<_>>();
let histories = self.list_history().await?;
for record in histories.iter().filter(|record| {
record.phase == "active"
&& session_type(record) == "telegram"
&& !bound.contains(&record.id)
}) {
self.queue_telegram_session_for_ingress(record, "telegram-detached")
.await?;
}
Ok(())
}
async fn queue_expired_telegram_sessions(
&self,
histories: &[SessionRecord],
now: DateTime<Utc>,
) -> anyhow::Result<()> {
for record in histories
.iter()
.filter(|record| telegram_session_is_expired(record, now))
{
self.queue_telegram_session_for_ingress(record, "telegram-session-timeout")
.await?;
}
Ok(())
}
async fn queue_telegram_session_for_ingress(
&self,
summary: &SessionRecord,
reason: &str,
) -> anyhow::Result<bool> {
let id = summary.id.clone();
let lock = self.conversation_lock(&id).await;
let _guard = lock.lock().await;
let record = self.get_conversation(&id).await?;
if record.phase != "active"
|| !matches!(
session_type(&record).as_str(),
"telegram" | "telegram-group"
)
{
return Ok(false);
}
if reason == "telegram-session-timeout" && !telegram_session_is_expired(&record, Utc::now())
{
return Ok(false);
}
let mut state = record.state.clone();
state["orchestration"] = json!({
"owner":"backend",
"status":"stopped",
"reason":reason,
});
self.control
.api()
.history_request_ingress(
&id,
kcode_session_history::Checkpoint {
expected_version: record.version,
state: state.clone(),
user_activity: false,
},
)
.await?;
if let Some(session_id) = state.get("rustLibSessionId").and_then(Value::as_str) {
self.control.api().release_managed_sources(session_id).await;
}
tracing::info!(session_id=%id, %reason, "Queued Telegram session for history ingress");
Ok(true)
}
async fn telegram_session(
&self,
event: &Value,
) -> anyhow::Result<(Arc<Mutex<SessionRecord>>, Session)> {
let histories = self.list_history().await?;
let group = event.get("sessionKind").and_then(Value::as_str) == Some("group");
let user_id = event
.get("telegramUserId")
.map(value_string)
.unwrap_or_default();
let group_id = event.get("groupId").and_then(Value::as_str);
let mut record = event
.get("conversationId")
.and_then(Value::as_str)
.and_then(|id| histories.iter().find(|record| record.id == id).cloned());
if record.is_none() {
record = histories.into_iter().find(|record| {
record.phase == "active"
&& if group {
session_type(record) == "telegram-group"
&& record_group_id(record) == group_id
&& record_user_id(record) == user_id
} else {
session_type(record) == "telegram" && record_user_id(record) == user_id
}
});
}
let created = record
.as_ref()
.is_none_or(|record| record.phase != "active");
let (record, mut session) =
if let Some(record) = record.filter(|record| record.phase == "active") {
let record = self.get_conversation(&record.id).await?;
let session = self.session_for_record(&record).await?;
(record, session)
} else {
self.create_telegram_session(event).await?
};
let record = Arc::new(Mutex::new(record));
let id = {
let locked = record.lock().await;
locked.id.clone()
};
if event.get("conversationId").and_then(Value::as_str) != Some(&id)
|| event.get("processingStartedAt").is_none()
{
self.control
.api()
.telegram_bind_event(
&required_string(event, "id")?,
&id,
event.get("conversationId").and_then(Value::as_str),
)
.await?;
}
if group && !created {
let group_id = required_string(event, "groupId")?;
if let Some(context) = event.get("groupContext") {
let context = self
.prepare_group_context(
context.clone(),
event.get("messageId").map(value_string).as_deref(),
&group_id,
)
.await?;
session.refresh_telegram_group_context(
&context,
event.get("messageId").map(value_string).as_deref(),
)?;
persist_record(self.control.api(), &record, session.snapshot()?, false).await?;
}
}
Ok((record, session))
}
async fn create_telegram_session(
&self,
event: &Value,
) -> anyhow::Result<(SessionRecord, Session)> {
let runtime = self.runtime()?.clone();
let user = self.directory_user(event).await?;
let group = event.get("sessionKind").and_then(Value::as_str) == Some("group");
let mut roots = vec![
user.root_node_id
.context("Telegram user root is not ready")?,
];
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});
let mut references = Vec::new();
if group {
let group_id = required_string(event, "groupId")?;
let group_record = self.directory_group(&group_id).await?;
let group_root = group_record
.root_node_id
.context("Telegram group root is not ready")?;
roots.push(group_root.clone());
if let Some(context) = event.get("groupContext") {
let context = self
.prepare_group_context(
context.clone(),
event.get("messageId").map(value_string).as_deref(),
&group_id,
)
.await?;
channel["groupContext"] = context.clone();
channel["groupRootNodeId"] = json!(group_root);
references = participant_references(&context, &roots);
}
}
roots.push(runtime.kennedy_root_node_id.clone());
references.retain(|id| !roots.contains(id));
let mut options =
SessionOptions::conversation(if group { "telegram-group" } else { "telegram" }, roots);
options.channel = channel;
options.reference_root_node_ids = references;
let session = self.open_session(runtime, options, None).await?;
let state = session.snapshot()?;
let record = self
.control
.api()
.history_register(kcode_session_history::RegisterSession {
id: required_string(&state, "sessionId")?,
started_at: session.started_at.clone(),
state,
})
.await?;
Ok((record, session))
}
async fn telegram_input(&self, event: &Value) -> anyhow::Result<(String, Value)> {
let Some(batch) = event
.get("batchedEvents")
.and_then(Value::as_array)
.filter(|batch| batch.len() > 1)
else {
return self.telegram_event_input(event).await;
};
let mut inputs = Vec::with_capacity(batch.len());
for batched_event in batch {
inputs.push(self.telegram_event_input(batched_event).await?);
}
merge_telegram_batch_inputs(event, inputs)
}
async fn telegram_event_input(&self, event: &Value) -> anyhow::Result<(String, Value)> {
let id = required_string(event, "id")?;
match event.get("kind").and_then(Value::as_str).unwrap_or("text") {
"voice" => {
let (bytes, mime) = self.control.api().telegram_event_media(&id)?;
let filename = event
.get("fileName")
.and_then(Value::as_str)
.filter(|value| !value.trim().is_empty())
.unwrap_or("telegram-voice.ogg");
Ok((
event
.get("text")
.and_then(Value::as_str)
.unwrap_or("")
.into(),
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)}}),
))
}
"document" => {
let (bytes, mime) = self.control.api().telegram_event_media(&id)?;
let filename = event
.get("fileName")
.and_then(Value::as_str)
.unwrap_or("telegram-document")
.to_owned();
let extraction = self
.control
.api()
.extract_document(bytes.clone(), filename.clone(), &mime)
.await;
let mut attachment = json!({
"id":format!("telegram:{id}"),
"kind":"document",
"source":"telegram",
"fileName":filename,
"mimeType":mime,
"sizeBytes":bytes.len(),
"dataUrl":data_url(&mime,&bytes),
});
match extraction {
Ok(result) => {
attachment["format"] = json!(result.format);
attachment["text"] = json!(result.text);
attachment["characters"] = json!(result.characters);
attachment["truncated"] = json!(result.truncated);
}
Err(error) => {
attachment["extractionError"] = json!(error.to_string());
}
}
Ok((
event
.get("text")
.and_then(Value::as_str)
.unwrap_or("")
.into(),
json!({"externalEventId":id,"inputKind":"document","attachments":[attachment]}),
))
}
kind @ ("photo" | "video" | "animation" | "audio" | "video_note" | "sticker") => {
let (bytes, downloaded_mime) = self.control.api().telegram_event_media(&id)?;
let mime = event
.get("mimeType")
.and_then(Value::as_str)
.filter(|value| !value.trim().is_empty())
.unwrap_or(&downloaded_mime)
.to_owned();
let extension = match kind {
"photo" => "jpg",
"video" | "video_note" => "mp4",
"animation" => "gif",
"audio" => "mp3",
"sticker" => "webp",
_ => "bin",
};
let filename = event
.get("fileName")
.and_then(Value::as_str)
.filter(|value| !value.trim().is_empty())
.map(str::to_owned)
.unwrap_or_else(|| format!("telegram-{kind}.{extension}"));
let mut attachment = json!({
"id":format!("telegram:{id}"),
"kind":kind,
"source":"telegram",
"fileName":filename,
"mimeType":mime,
"sizeBytes":bytes.len(),
"dataUrl":data_url(&mime,&bytes),
});
if let Some(value) = event.get("durationSeconds") {
attachment["durationSeconds"] = value.clone();
}
Ok((
event
.get("text")
.and_then(Value::as_str)
.unwrap_or("")
.into(),
json!({
"externalEventId":id,
"inputKind":kind,
"attachments":[attachment],
}),
))
}
_ => Ok((
event
.get("text")
.and_then(Value::as_str)
.unwrap_or("")
.into(),
json!({"externalEventId":id,"inputKind":"text"}),
)),
}
}
async fn process_telegram_reset(&self, event: &Value) -> anyhow::Result<()> {
let id = required_string(event, "id")?;
let Some(conversation_id) = event.get("conversationId").and_then(Value::as_str) else {
self.control.api()
.telegram_complete_reset(
&id,
Some(
"There is no active Telegram session to reset. Your next message will begin one.",
),
)
.await?;
return Ok(());
};
let record = match self.get_conversation(conversation_id).await {
Ok(record) => record,
Err(error)
if error
.downcast_ref::<ApiError>()
.is_some_and(|error| error.code == "not_found") =>
{
self.control.api()
.telegram_complete_reset(
&id,
Some(
"There is no active Telegram session to reset. Your next message will begin one.",
),
)
.await?;
return Ok(());
}
Err(error) => return Err(error),
};
if record.phase != "active" {
self.control.api()
.telegram_complete_reset(
&id,
Some(
"There is no active Telegram session to reset. Your next message will begin one.",
),
)
.await?;
return Ok(());
}
let session = self.session_for_record(&record).await?;
session.release_managed_sources().await;
self.control
.api()
.history_request_ingress(
conversation_id,
kcode_session_history::Checkpoint {
expected_version: record.version,
state: record.state,
user_activity: false,
},
)
.await?;
self.control.api()
.telegram_complete_reset(
&id,
Some(
"Conversation reset. The Telegram session has been queued for memory ingress; your next message will begin a new session.",
),
)
.await?;
Ok(())
}
async fn sync_group_updates(self: &Arc<Self>) -> anyhow::Result<()> {
let updates = self
.control
.api()
.telegram_group_session_updates()
.await?
.get("updates")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default();
for update in updates {
let id = required_string(&update, "conversationId")?;
let mut set = self.group_updates_in_flight.lock().await;
if !set.insert(id.clone()) {
continue;
}
drop(set);
let worker = self.clone();
tokio::spawn(async move {
if let Err(error) = worker.process_group_update(update).await {
tracing::warn!(conversation_id=%id,error=%error,"Telegram group context update will retry");
}
worker.group_updates_in_flight.lock().await.remove(&id);
});
}
Ok(())
}
async fn process_group_update(&self, update: Value) -> anyhow::Result<()> {
let id = required_string(&update, "conversationId")?;
let lock = self.conversation_lock(&id).await;
let _guard = lock.lock().await;
let record = match self.get_conversation(&id).await {
Ok(record) => record,
Err(error)
if error
.downcast_ref::<ApiError>()
.is_some_and(|error| error.code == "not_found") =>
{
self.reconcile_missing_group_session(&update, &id).await?;
return Ok(());
}
Err(error) => return Err(error),
};
if record.phase != "active" {
if update.get("resetRequired").and_then(Value::as_bool) == Some(true) {
self.control
.api()
.telegram_complete_silent_group_reset(&id)
.await?;
}
return Ok(());
}
let mut session = self.session_for_record(&record).await?;
if session.pending_turn {
return Ok(());
}
let group_id = required_string(&update, "groupId")?;
let context = self
.prepare_group_context(
update.get("groupContext").cloned().unwrap_or(Value::Null),
None,
&group_id,
)
.await?;
session.refresh_telegram_group_context(&context, None)?;
let record = Arc::new(Mutex::new(record));
persist_record(self.control.api(), &record, session.snapshot()?, false).await?;
if update.get("resetRequired").and_then(Value::as_bool) == Some(true) {
self.close_conversation(&record, &session).await?;
self.control
.api()
.telegram_complete_silent_group_reset(&id)
.await?;
} else {
self.control
.api()
.telegram_acknowledge_group_context(
&id,
update
.get("throughMessageId")
.and_then(Value::as_i64)
.unwrap_or(0),
)
.await?;
}
Ok(())
}
async fn reconcile_missing_group_session(
&self,
update: &Value,
conversation_id: &str,
) -> anyhow::Result<()> {
match missing_group_session_recovery(update)? {
MissingGroupSessionRecovery::CompleteSilentReset => {
self.control
.api()
.telegram_complete_silent_group_reset(conversation_id)
.await?;
tracing::info!(
%conversation_id,
"Completed orphaned Telegram group reset"
);
}
MissingGroupSessionRecovery::DetachCurrent {
group_id,
telegram_user_id,
} => {
let result = self
.control
.api()
.telegram_detach_group_session(conversation_id, &group_id, telegram_user_id)
.await;
match result {
Ok(_) => {
tracing::info!(
%conversation_id,
%group_id,
telegram_user_id,
"Detached orphaned Telegram group session"
);
}
Err(error) if error.code == "state_conflict" => {
tracing::info!(
%conversation_id,
%group_id,
telegram_user_id,
"Telegram group session was already detached or rebound"
);
}
Err(error) => return Err(error.into()),
}
}
}
Ok(())
}
async fn sync_group_ingress(self: &Arc<Self>) -> anyhow::Result<()> {
let batches = self
.control
.api()
.telegram_group_ingress()
.await?
.get("batches")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default();
for batch in batches {
let id = required_string(&batch, "id")?;
let mut set = self.group_ingress_in_flight.lock().await;
if !set.insert(id.clone()) {
continue;
}
drop(set);
let worker = self.clone();
tokio::spawn(async move {
if let Err(error) = worker.process_group_ingress(batch).await {
tracing::warn!(batch_id=%id,error=%error,"Telegram group ingress preparation will retry");
}
worker.group_ingress_in_flight.lock().await.remove(&id);
});
}
Ok(())
}
async fn process_group_ingress(&self, batch: Value) -> anyhow::Result<()> {
let runtime = self.runtime()?.clone();
let id = required_string(&batch, "id")?;
if let Some(existing) = self.list_history().await?.into_iter().find(|record| {
record
.state
.get("channel")
.and_then(|channel| channel.get("groupIngressBatchId"))
.and_then(Value::as_str)
== Some(&id)
}) {
match existing.phase.as_str() {
"complete" => {
self.control
.api()
.telegram_complete_group_ingress(&id)
.await?;
}
"active" => {
let existing = self.get_conversation(&existing.id).await?;
self.control
.api()
.history_request_ingress(
&existing.id,
kcode_session_history::Checkpoint {
expected_version: existing.version,
state: existing.state,
user_activity: false,
},
)
.await?;
}
_ => {}
}
return Ok(());
}
let group_id = required_string(&batch, "groupId")?;
let group = self.directory_group(&group_id).await?;
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!([]))});
let context = self
.prepare_group_context(raw_context, None, &group_id)
.await?;
let group_root = group
.root_node_id
.context("Telegram group root is not ready")?;
let roots = vec![group_root.clone(), runtime.kennedy_root_node_id.clone()];
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});
let mut options = SessionOptions::conversation("telegram-group", roots.clone());
options.channel = channel;
options.reference_root_node_ids = participant_references(&context, &roots);
options.source_session_type = Some("telegram-group".into());
let mut session = self.open_session(runtime, options, None).await?;
for message in context
.get("messages")
.and_then(Value::as_array)
.into_iter()
.flatten()
{
session.stage_source_message(
message.get("sentByKennedy").and_then(Value::as_bool) == Some(true),
message
.get("text")
.and_then(Value::as_str)
.unwrap_or_default(),
message.clone(),
)?;
}
let state = session.snapshot()?;
let record = self
.control
.api()
.history_register(kcode_session_history::RegisterSession {
id: required_string(&state, "sessionId")?,
started_at: batch
.get("createdAt")
.and_then(Value::as_str)
.map(str::to_owned)
.unwrap_or_else(|| Utc::now().to_rfc3339()),
state,
})
.await?;
self.control
.api()
.history_request_ingress(
&record.id,
kcode_session_history::Checkpoint {
expected_version: record.version,
state: record.state,
user_activity: false,
},
)
.await?;
Ok(())
}
}
impl Deref for Runtime {
type Target = Orchestrator;
fn deref(&self) -> &Self::Target {
&self.control
}
}
fn session_type(record: &SessionRecord) -> String {
record
.state
.get("sessionType")
.and_then(Value::as_str)
.unwrap_or("conversation")
.into()
}
fn next_wakeup_marker(now: DateTime<Utc>) -> DateTime<Utc> {
let day_start = now
.date_naive()
.and_hms_opt(0, 0, 0)
.expect("midnight is always a valid UTC time")
.and_utc();
for hour in [0_i64, 4, 8, 12, 16, 20] {
let candidate = day_start + ChronoDuration::hours(hour);
if candidate > now {
return candidate;
}
}
day_start + ChronoDuration::days(1)
}
fn telegram_session_is_expired(record: &SessionRecord, now: DateTime<Utc>) -> bool {
if record.phase != "active"
|| !matches!(session_type(record).as_str(), "telegram" | "telegram-group")
|| record.state.get("pendingTurn").and_then(Value::as_bool) == Some(true)
{
return false;
}
DateTime::parse_from_rfc3339(&record.started_at)
.ok()
.is_some_and(|started| now >= started.with_timezone(&Utc) + TELEGRAM_SESSION_MAX_AGE)
}
fn record_channel(record: &SessionRecord) -> Option<&Value> {
record.state.get("channel")
}
fn record_group_id(record: &SessionRecord) -> Option<&str> {
record_channel(record)
.and_then(|channel| {
channel.get("groupId").or_else(|| {
channel
.get("groupContext")
.and_then(|context| context.get("groupId"))
})
})
.and_then(Value::as_str)
}
fn record_user_id(record: &SessionRecord) -> String {
record_channel(record)
.and_then(|channel| channel.get("telegramUserId"))
.map(value_string)
.unwrap_or_default()
}
fn participant_references(context: &Value, roots: &[String]) -> Vec<String> {
let mut values = context
.get("participants")
.and_then(Value::as_array)
.into_iter()
.flatten()
.filter_map(|participant| participant.get("rootNodeId").and_then(Value::as_str))
.filter(|id| !roots.iter().any(|root| root == id))
.map(str::to_owned)
.collect::<Vec<_>>();
values.sort();
values.dedup();
values
}
fn normalized_file_mime_type(value: &str) -> String {
let value = value
.split(';')
.next()
.unwrap_or(value)
.trim()
.to_ascii_lowercase();
if value.contains('/')
&& !value.is_empty()
&& !value.chars().any(char::is_whitespace)
&& !value.chars().any(char::is_control)
{
value
} else {
"application/octet-stream".into()
}
}
fn file_extension_for_mime_type(mime_type: &str) -> &'static str {
match normalized_file_mime_type(mime_type).as_str() {
"image/jpeg" => "jpg",
"image/png" => "png",
"image/webp" => "webp",
"image/gif" => "gif",
"audio/ogg" | "audio/opus" | "application/ogg" => "ogg",
"audio/mpeg" | "audio/mp3" => "mp3",
"audio/mp4" | "video/mp4" => "mp4",
"audio/webm" | "video/webm" => "webm",
"audio/wav" | "audio/x-wav" => "wav",
"application/pdf" => "pdf",
_ => "bin",
}
}
fn telegram_group_context_file_name(
supplied: Option<&str>,
kind: &str,
mime_type: &str,
message_id: &str,
) -> String {
let fallback = format!(
"telegram-group-{kind}-{message_id}.{}",
file_extension_for_mime_type(mime_type)
);
sanitize_file_name(supplied.unwrap_or_default(), &fallback)
}
fn file_name_extension(file_name: &str) -> String {
file_name
.rsplit_once('.')
.and_then(|(stem, extension)| {
(!stem.is_empty() && !extension.is_empty()).then_some(extension)
})
.map(|extension| format!(".{extension}"))
.unwrap_or_else(|| "(none)".into())
}
fn telegram_group_context_file_metadata(message: &Value) -> String {
let file_name = message
.get("fileName")
.and_then(Value::as_str)
.unwrap_or("telegram-file");
let source_note =
if message.get("fileNameSource").and_then(Value::as_str) == Some("synthesized") {
" (synthesized because Telegram supplied no filename)"
} else {
""
};
let mime_type = normalized_file_mime_type(
message
.get("mimeType")
.and_then(Value::as_str)
.unwrap_or("application/octet-stream"),
);
let size_bytes = message
.get("sizeBytes")
.and_then(Value::as_u64)
.unwrap_or_default();
format!(
"User-provided file\nOriginal filename: {file_name}{source_note}\nExtension: {}\nMIME type: {mime_type}\nSize: {size_bytes} bytes",
file_name_extension(file_name),
)
}
fn telegram_response_object_deliveries(response: &Value) -> Vec<(String, Option<String>)> {
let objects = response
.get("objects")
.and_then(Value::as_array)
.map(Vec::as_slice)
.unwrap_or_default();
let attachments = response
.get("attachments")
.and_then(Value::as_array)
.map(Vec::as_slice)
.unwrap_or_default();
objects
.iter()
.filter_map(Value::as_str)
.enumerate()
.map(|(index, object_id)| {
let descriptor = attachments
.iter()
.find(|candidate| {
["objectId", "pendingId", "id"]
.iter()
.any(|key| candidate.get(key).and_then(Value::as_str) == Some(object_id))
})
.or_else(|| attachments.get(index));
let file_name = descriptor
.and_then(|descriptor| descriptor.get("fileName"))
.and_then(Value::as_str)
.map(str::to_owned);
(object_id.to_owned(), file_name)
})
.collect()
}
fn telegram_reply_caption<'a>(
deliveries: &'a [TelegramDelivery],
object_index: usize,
file: &ResolvedObject,
) -> Option<&'a str> {
if object_index + 2 != deliveries.len() {
return None;
}
match &deliveries[object_index + 1] {
TelegramDelivery::Text {
text,
response_warning,
captionable: true,
} if response_warning.is_null() || response_warning.as_str() == Some("") => {
telegram_caption_for(file, text)
}
_ => None,
}
}
fn required_string(value: &Value, key: &str) -> anyhow::Result<String> {
value
.get(key)
.and_then(Value::as_str)
.filter(|value| !value.is_empty())
.map(str::to_owned)
.with_context(|| format!("backend response omitted {key}"))
}
fn merge_telegram_batch_inputs(
event: &Value,
inputs: Vec<(String, Value)>,
) -> anyhow::Result<(String, Value)> {
let mut texts = Vec::new();
let mut attachments = Vec::new();
let mut event_ids = Vec::with_capacity(inputs.len());
for (text, metadata) in inputs {
if !text.is_empty() {
texts.push(text);
}
if let Some(media) = metadata.get("media").filter(|value| value.is_object()) {
attachments.push(media.clone());
}
if let Some(items) = metadata.get("attachments").and_then(Value::as_array) {
attachments.extend(items.iter().cloned());
}
if let Some(id) = metadata.get("externalEventId").and_then(Value::as_str) {
event_ids.push(id.to_owned());
}
}
Ok((
texts.join("\n\n"),
json!({
"externalEventId":required_string(event, "id")?,
"externalEventIds":event_ids,
"inputKind":"batch",
"attachments":attachments,
}),
))
}
fn missing_group_session_recovery(update: &Value) -> anyhow::Result<MissingGroupSessionRecovery> {
if update.get("resetRequired").and_then(Value::as_bool) == Some(true) {
return Ok(MissingGroupSessionRecovery::CompleteSilentReset);
}
Ok(MissingGroupSessionRecovery::DetachCurrent {
group_id: required_string(update, "groupId")?,
telegram_user_id: update
.get("telegramUserId")
.and_then(Value::as_i64)
.context("backend response omitted telegramUserId")?,
})
}
fn value_string(value: &Value) -> String {
value
.as_str()
.map(str::to_owned)
.unwrap_or_else(|| value.to_string())
}
fn bounded_error(error: &anyhow::Error) -> String {
format!("{error:#}").chars().take(1_000).collect()
}
fn telegram_event_retry_delay(failures: u32) -> Duration {
let exponent = failures.saturating_sub(1).min(5);
Duration::from_secs((2_u64 << exponent).min(60))
}
fn telegram_event_retry_should_warn(
previous_error: Option<&str>,
error: &str,
failures: u32,
) -> bool {
previous_error != Some(error) || failures.is_multiple_of(10)
}
fn telegram_timeout(event: &Value) -> Duration {
let elapsed = event
.get("processingStartedAt")
.and_then(Value::as_str)
.and_then(|value| DateTime::parse_from_rfc3339(value).ok())
.map(|value| {
(Utc::now() - value.with_timezone(&Utc))
.to_std()
.unwrap_or_default()
})
.unwrap_or_default();
TELEGRAM_TIMEOUT.saturating_sub(elapsed)
}
#[cfg(test)]
mod tests {
use super::*;
fn session_record(overrides: Value) -> SessionRecord {
let mut record = json!({
"id":"session",
"phase":"active",
"started_at":"2026-07-30T00:00:00Z",
"updated_at":"2026-07-30T00:00:00Z",
"state":{},
"provenance_id":null,
"version":1,
"last_user_message_at":null,
"ended_at":null,
"ingress_failure_count":0,
"ingress_failures":[],
"ingress_next_attempt_at":null
});
record
.as_object_mut()
.unwrap()
.extend(overrides.as_object().unwrap().clone());
serde_json::from_value(record).unwrap()
}
#[test]
fn batches_keep_order_text_and_every_attachment() {
let (text, metadata) = merge_telegram_batch_inputs(
&json!({"id":"batch"}),
vec![
("first".into(), json!({"externalEventId":"one"})),
(
String::new(),
json!({"externalEventId":"two","media":{"id":"voice"}}),
),
(
"third".into(),
json!({"externalEventId":"three","attachments":[{"id":"document"}]}),
),
],
)
.unwrap();
assert_eq!(text, "first\n\nthird");
assert_eq!(metadata["externalEventIds"], json!(["one", "two", "three"]));
assert_eq!(
metadata["attachments"],
json!([{"id":"voice"}, {"id":"document"}])
);
}
#[test]
fn wakeups_use_strictly_future_four_hour_utc_boundaries() {
let before = DateTime::parse_from_rfc3339("2026-07-28T03:59:59Z")
.unwrap()
.with_timezone(&Utc);
assert_eq!(
next_wakeup_marker(before).to_rfc3339(),
"2026-07-28T04:00:00+00:00"
);
let exactly = DateTime::parse_from_rfc3339("2026-07-28T20:00:00Z")
.unwrap()
.with_timezone(&Utc);
assert_eq!(
next_wakeup_marker(exactly).to_rfc3339(),
"2026-07-29T00:00:00+00:00"
);
}
#[test]
fn idle_telegram_sessions_roll_over_after_six_hours() {
let now = DateTime::parse_from_rfc3339("2026-07-30T12:00:00Z")
.unwrap()
.with_timezone(&Utc);
assert!(telegram_session_is_expired(
&session_record(json!({
"started_at":"2026-07-30T06:00:00Z",
"state":{"sessionType":"telegram","pendingTurn":false}
})),
now
));
assert!(!telegram_session_is_expired(
&session_record(json!({
"started_at":"2026-07-30T05:00:00Z",
"state":{"sessionType":"telegram","pendingTurn":true}
})),
now
));
}
#[test]
fn missing_group_sessions_choose_reset_or_detach() {
assert_eq!(
missing_group_session_recovery(&json!({"resetRequired":true})).unwrap(),
MissingGroupSessionRecovery::CompleteSilentReset
);
assert_eq!(
missing_group_session_recovery(&json!({
"groupId":"group-1","telegramUserId":42,"resetRequired":false
}))
.unwrap(),
MissingGroupSessionRecovery::DetachCurrent {
group_id: "group-1".into(),
telegram_user_id: 42,
}
);
}
#[test]
fn response_objects_keep_their_filename_overrides() {
assert_eq!(
telegram_response_object_deliveries(&json!({
"objects":["pending:2","AAECAwQF"],
"attachments":[
{"objectId":"AAECAwQF","fileName":"canonical.pdf"},
{"objectId":"pending:2","fileName":"draft.pdf"}
]
})),
vec![
("pending:2".into(), Some("draft.pdf".into())),
("AAECAwQF".into(), Some("canonical.pdf".into())),
]
);
}
#[test]
fn exact_final_text_is_used_only_as_a_supported_caption() {
let file = ResolvedObject {
object_id: "object".into(),
bytes: vec![1],
file_name: "photo.jpg".into(),
media_type: "image/jpeg".into(),
transport_kind: Some("photo".into()),
};
let deliveries = vec![
TelegramDelivery::Object {
object_id: "object".into(),
file_name: None,
},
TelegramDelivery::Text {
text: " exact caption\n".into(),
response_warning: Value::Null,
captionable: true,
},
];
assert_eq!(
telegram_reply_caption(&deliveries, 0, &file),
Some(" exact caption\n")
);
}
#[test]
fn event_retries_back_off_and_repeat_warnings_periodically() {
assert_eq!(telegram_event_retry_delay(1), Duration::from_secs(2));
assert_eq!(telegram_event_retry_delay(6), Duration::from_secs(60));
assert_eq!(
telegram_event_retry_delay(u32::MAX),
Duration::from_secs(60)
);
assert!(!telegram_event_retry_should_warn(
Some("failure"),
"failure",
2
));
assert!(telegram_event_retry_should_warn(
Some("failure"),
"failure",
10
));
}
}