#![warn(clippy::pedantic)]
use anyhow::Result;
use chrono::{Duration as ChronoDuration, Utc};
use futures_util::FutureExt;
use std::panic::AssertUnwindSafe;
use std::sync::{Arc, OnceLock};
use std::time::Duration;
use tokio::task::spawn;
use tracing::{debug, error, info, warn};
use std::borrow::Cow;
use tokio::task::JoinSet;
use tokio_util::sync::CancellationToken;
use mahbot::channels::broadcast_and_persist_incoming_message;
use mahbot::channels::telegram::{decode_action, user_command_entries};
use mahbot::config::{CONFIG, CONFIG_KEY_IMAGE_GEN_MODEL, CONFIG_KEY_VIDEO_MODEL};
use mahbot::gui::{BOOT_LOG_STORE, Dashboard, JETBRAINS_MONO, Message as DashboardMessage};
use mahbot::message_router;
use mahbot::parse_bot_command;
use mahbot::session::clear_session;
use mahbot::util::UnwrapPoison;
use mahbot::{BotCommand, Channel, ChannelMessage, Role, Workspace};
const JETBRAINS_MONO_FONT_BYTES: &[u8] = include_bytes!("gui/JetBrainsMono-Regular.ttf");
const JETBRAINS_MONO_BOLD_FONT_BYTES: &[u8] = include_bytes!("gui/JetBrainsMono-Bold.ttf");
const LOG_RETENTION_HOURS: i64 = 8;
async fn bootstrap_mahbot_safe() -> Result<(), String> {
match AssertUnwindSafe(bootstrap_mahbot()).catch_unwind().await {
Ok(Ok(())) => Ok(()),
Ok(Err(e)) => Err(e.to_string()),
Err(payload) => Err(format_startup_panic(&*payload)),
}
}
fn format_startup_panic(payload: &(dyn std::any::Any + Send)) -> String {
format!("Startup panicked: {}", mahbot::util::panic_message(payload))
}
async fn bootstrap_mahbot() -> Result<()> {
mahbot::config::load_or_init().await?;
let storage_root = mahbot::config::CONFIG.global_storage_root();
mahbot::wal_guard::diagnose_all_stores(std::path::Path::new(&storage_root));
let (log_store, log_broadcast) = mahbot::logs::init_tracing(&storage_root).await?;
let _ = mahbot::gui::LOG_BROADCAST.set(log_broadcast);
mahbot::search_engine::init_global(); mahbot::ticket_buffer::init_global(); mahbot::message_router::init_global()?;
mahbot::audio::voice::init_global()?;
mahbot::audio::tts::init_global()?;
mahbot::turso::init_all_stores().await?;
mahbot::config::reload_from_db().await?;
mahbot::audio::local_transcriber::spawn_background_init_if_enabled();
mahbot::providers::init_global()?;
if mahbot::audio::tts::is_config_enabled() && !mahbot::audio::tts::try_load_cached() {
mahbot::audio::tts::spawn_download();
}
BOOT_LOG_STORE
.set(log_store.as_ref().clone())
.map_err(|_| anyhow::anyhow!("BOOT_LOG_STORE already set"))?;
spawn_background_tasks(log_store.clone());
info!("MahBot initialized β dashboard ready");
let admin_target = mahbot::self_update::resolve_admin_telegram_target().await;
tokio::spawn(async move {
mahbot::self_update::notify_admin("β
MahBot is back online.", admin_target.as_deref())
.await;
});
Ok(())
}
static BACKGROUND_TASKS: std::sync::Mutex<Option<JoinSet<()>>> = std::sync::Mutex::new(None);
fn spawn_cancellable<F>(
tasks: &mut JoinSet<()>,
shutdown_token: &CancellationToken,
name: &'static str,
fut: F,
) where
F: Future<Output = ()> + Send + 'static,
{
let cancel = shutdown_token.clone();
tasks.spawn(async move {
tokio::select! {
result = AssertUnwindSafe(fut).catch_unwind() => {
if let Err(payload) = result {
error!(
"Background task panicked [{name}]: {}",
mahbot::util::panic_message(&*payload),
);
}
}
() = cancel.cancelled() => {},
}
});
}
#[expect(clippy::too_many_lines)]
fn spawn_background_tasks(log_store: Arc<mahbot::logs::LogStore>) {
let mut tasks = JoinSet::<()>::new();
let shutdown_token = mahbot::shutdown::shutdown_token();
spawn_cancellable(
&mut tasks,
&shutdown_token,
"session-cleanup",
run_cleanup_loop(
"Session cleanup",
mahbot::jobs::PURGE_CUTOFF_HOURS,
|cutoff| async move {
let purged = mahbot::jobs::purge_stale_jobs(&cutoff).await?;
let cleaned = mahbot::session::cleanup_old_transient_sessions(&cutoff).await?;
let media = mahbot::research_cleanup::sweep_media().await?;
Ok(purged + cleaned + media)
},
),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"log-cleanup",
run_cleanup_loop("Log cleanup", LOG_RETENTION_HOURS, {
let store = log_store;
move |cutoff| {
let store = store.clone();
async move { store.delete_older_than("INFO", &cutoff).await }
}
}),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"maintainer",
mahbot::maintainer::run_maintainer_loop(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"archive-cancelled",
mahbot::board::run_archive_cancelled_loop(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"dead-session-recovery",
mahbot::session::dead_session::run_dead_session_recovery_loop(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"nightly-check",
mahbot::workspace::run_nightly_check_loop(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"wal-guard",
mahbot::wal_guard::run_wal_guard_loop(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"search-engine-init",
mahbot::search_engine::init_all_engines(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"browser-daemon",
mahbot::tools::browser_daemon::run_watchdog(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"voice-pipeline",
mahbot::audio::voice::run_voice_pipeline(),
);
let rx = init_message_pipeline(&mut tasks, &shutdown_token);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"message-handler",
run_message_dispatch_loop(rx),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"management",
mahbot::management::run_management(),
);
tasks.spawn(async move {
let result = AssertUnwindSafe(mahbot::shutdown::wait_for_shutdown_signal())
.catch_unwind()
.await;
match result {
Ok(Ok(())) => {
info!("Second signal received β force-cancelling drain");
mahbot::shutdown::force_cancel();
}
Ok(Err(e)) => {
error!("Signal handler failed to set up: {e}");
}
Err(payload) => {
error!(
"Signal handler panicked: {}",
mahbot::util::panic_message(&*payload),
);
}
}
});
spawn_cancellable(
&mut tasks,
&shutdown_token,
"drain-watch",
mahbot::jobs::run_drain_watch(),
);
spawn_cancellable(&mut tasks, &shutdown_token, "auto-checkpoint", async {
loop {
if !mahbot::shutdown::sleep_or_shutdown_or_drain(Duration::from_mins(5)).await {
break;
}
if mahbot::self_update::update_is_finalizing() {
break;
}
mahbot::checkpoint::periodic_checkpoint_and_verify().await;
}
});
{
let mut guard = BACKGROUND_TASKS.lock().unwrap_poison();
let _ = guard.insert(tasks);
}
}
fn init_message_pipeline(
tasks: &mut JoinSet<()>,
cancel: &CancellationToken,
) -> tokio::sync::mpsc::Receiver<ChannelMessage> {
let (tx, rx) = tokio::sync::mpsc::channel::<ChannelMessage>(100);
mahbot::MESSAGE_TX
.set(tx.clone())
.expect("MESSAGE_TX already set β should be first init");
let gui_pipeline_tx = tx.clone();
let (chat_tx, _chat_rx) = tokio::sync::broadcast::channel::<mahbot::ChatEvent>(256);
mahbot::CHAT_BROADCAST
.set(chat_tx)
.expect("CHAT_BROADCAST already set β should be first init");
mahbot::audio::tts::init_listener();
let _ = mahbot::CHANNEL_REGISTRY.set(mahbot::ChannelRegistry::default());
if let Some(token) = CONFIG.telegram_bot_token() {
use mahbot::channels::telegram::TelegramChannel;
let channel: std::sync::Arc<TelegramChannel> =
std::sync::Arc::new(TelegramChannel::new(token));
tokio::spawn({
let tc = std::sync::Arc::clone(&channel);
async move {
tc.set_my_commands().await;
}
});
let channel: Arc<dyn Channel> = channel;
mahbot::channel_registry().register(Arc::clone(&channel));
spawn_cancellable(tasks, cancel, "telegram-listener", {
let channel = Arc::clone(&channel);
async move {
let _ = channel.listen(tx).await;
}
});
} else {
info!("No Telegram bot token configured β running in dashboard-only mode");
}
{
use mahbot::channels::gui::GuiChannel;
let (gui_channel, gui_tx) = GuiChannel::new();
mahbot::GUI_MESSAGE_TX
.set(gui_tx)
.expect("GUI_MESSAGE_TX already set β should be first init");
let gui_channel: Arc<dyn Channel> = Arc::new(gui_channel);
mahbot::channel_registry().register(Arc::clone(&gui_channel));
spawn_cancellable(tasks, cancel, "gui-listener", {
let channel = Arc::clone(&gui_channel);
async move {
let _ = channel.listen(gui_pipeline_tx).await;
}
});
}
mahbot::channels::voice::register_global();
rx
}
async fn shutdown_after_dashboard() {
info!("Dashboard window closed β shutting down");
mahbot::registry::AGENT_REGISTRY.shutdown_all();
mahbot::tools::browser::close_all_browser_sessions().await;
let maybe_tasks = {
let mut guard = BACKGROUND_TASKS.lock().unwrap_poison();
guard.take()
};
if let Some(mut tasks) = maybe_tasks {
while let Some(result) = tasks.join_next().await {
match result {
Ok(()) => {}
Err(e) if e.is_cancelled() => {
debug!("background task cancelled during shutdown");
}
Err(e) => {
warn!("background task panicked: {e}");
}
}
}
}
mahbot::checkpoint::checkpoint_all_databases().await;
}
fn main() -> Result<()> {
mahbot::shutdown::install_fatal_signal_handlers();
if std::env::args().nth(1).as_deref() == Some("debug") {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
match rt.block_on(mahbot::debug::run_debug()) {
Ok(()) => std::process::exit(0),
Err(e) if e.downcast_ref::<mahbot::debug::GateRefusal>().is_some() => {
eprintln!("Error: {e:#}");
std::process::exit(2);
}
Err(e) => {
eprintln!("Error: {e:#}");
std::process::exit(1);
}
}
}
#[cfg(unix)]
if std::env::args().nth(1).as_deref() == Some("__grep-engine") {
let code = mahbot::run_grep_engine(&std::env::args().skip(2).collect::<Vec<_>>());
std::process::exit(code);
}
if std::env::args().nth(1).as_deref() == Some("bench-openrouter") {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
let code = rt.block_on(mahbot::bench_openrouter::run_cli());
std::process::exit(code);
}
mahbot::temp_root::init_temp_root()?;
let update_mode = mahbot::self_update::update_mode();
let storage_root = mahbot::config::default_config_dir()?;
mahbot::self_update::acquire_lock(&storage_root)?;
let window_state = mahbot::gui::read_window_state();
iced::application(
move || {
(
Dashboard::loading(update_mode),
iced::Task::perform(bootstrap_mahbot_safe(), DashboardMessage::Boot),
)
},
Dashboard::update,
Dashboard::view,
)
.title(Dashboard::title)
.font(iced_fonts::LUCIDE_FONT_BYTES)
.font(JETBRAINS_MONO_FONT_BYTES)
.font(JETBRAINS_MONO_BOLD_FONT_BYTES)
.default_font(JETBRAINS_MONO)
.subscription(Dashboard::subscription)
.theme(Dashboard::theme)
.window(iced::window::Settings {
size: iced::Size::new(window_state.width, window_state.height),
position: window_state.position(),
min_size: Some(iced::Size::new(800.0, 500.0)),
..iced::window::Settings::default()
})
.exit_on_close_request(false)
.run()
.map_err(|e| anyhow::anyhow!("Iced application error: {e}"))?;
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| anyhow::anyhow!("shutdown runtime: {e}"))?;
rt.block_on(shutdown_after_dashboard());
Ok(())
}
async fn run_cleanup_loop<F, Fut>(label: &'static str, cutoff_hours: i64, cleanup: F)
where
F: Fn(String) -> Fut + Send + 'static,
Fut: Future<Output = Result<u64>> + Send,
{
loop {
if !mahbot::shutdown::sleep_or_shutdown_or_drain(Duration::from_mins(10)).await {
break;
}
let cutoff = (Utc::now() - ChronoDuration::hours(cutoff_hours)).to_rfc3339();
match cleanup(cutoff).await {
Ok(n) if n > 0 => tracing::debug!(deleted = n, "{label}: deleted old entries"),
Ok(_) => tracing::debug!("{label}: nothing to delete"),
Err(e) => warn!(error = %e, "{label} failed"),
}
}
}
async fn run_message_dispatch_loop(mut rx: tokio::sync::mpsc::Receiver<ChannelMessage>) {
let shutdown_token = mahbot::shutdown::shutdown_token();
loop {
let msg = tokio::select! {
() = shutdown_token.cancelled() => break,
msg = rx.recv() => match msg {
Some(msg) => msg,
None => break,
},
};
if let Some(decoded) = decode_action(&msg.content) {
spawn(handle_action_callback(msg, decoded));
continue;
}
if handle_bot_command(&msg).await {
continue;
}
spawn(process_channel_message(msg));
}
}
async fn handle_bot_command(msg: &ChannelMessage) -> bool {
let Some(cmd) = parse_bot_command(&msg.content) else {
return false;
};
if msg.channel != "telegram" {
return false;
}
match cmd {
BotCommand::Start => handle_start_command(msg).await,
BotCommand::Clear => handle_clear_session(msg).await,
BotCommand::ImageModels | BotCommand::VideoModels => {
if mahbot::users::role_pool(&msg.user_name)
.await
.contains(&Role::Artist)
{
handle_models_command(msg, cmd == BotCommand::ImageModels).await;
} else {
send_telegram_reply(
msg,
"This command is only available to Artist users.".to_string(),
)
.await;
}
}
BotCommand::SwitchRole(role) => handle_role_switch(msg, role).await,
BotCommand::Board
| BotCommand::Archive
| BotCommand::Pause
| BotCommand::Unpause
| BotCommand::Maintenance
| BotCommand::MaintenanceOn
| BotCommand::MaintenanceOff => {
if mahbot::users::is_admin(&msg.user_name).await {
handle_admin_command(msg, cmd).await;
} else {
send_telegram_reply(
msg,
"This command is only available to admin users.".to_string(),
)
.await;
}
}
}
true
}
async fn send_telegram_reply(msg: &ChannelMessage, content: String) {
let _ = mahbot::channels::telegram::send_direct(&msg.reply_target, content, None).await;
}
async fn handle_role_switch(msg: &ChannelMessage, role: Role) {
if !mahbot::users::role_pool(&msg.user_name)
.await
.contains(&role)
{
send_telegram_reply(
msg,
format!(
"Role '{}' is not in your allowed roles β ask an admin to add it.",
role.as_str()
),
)
.await;
return;
}
match mahbot::users::switch_active_role(&msg.user_name, role).await {
Ok(()) => {
let text = format!("Active role switched to {}.", role.display_label());
let Some(channel) = mahbot::channel_registry().get("telegram") else {
return;
};
let reply_target = msg.reply_target.clone();
tokio::spawn(async move {
let tc = channel
.as_any()
.downcast_ref::<mahbot::channels::telegram::TelegramChannel>()
.expect("registered telegram channel");
tc.send_role_switch_notification(&reply_target, &text).await;
});
}
Err(e) => send_telegram_reply(msg, format!("Failed to switch role: {e}")).await,
}
}
async fn handle_start_command(msg: &ChannelMessage) {
let mut lines = vec![
"\u{1F916} Welcome to MahBot!\n\nAvailable commands:".to_string(),
"/start β Show this message".to_string(),
];
for (cmd, desc) in user_command_entries(&msg.user_name).await {
lines.push(format!("/{cmd} β {desc}"));
}
send_telegram_reply(msg, lines.join("\n")).await;
}
async fn handle_clear_session(msg: &ChannelMessage) {
let (effective_role, ws) = mahbot::users::resolve_session_target(&msg.user_name).await;
let reply = clear_session(&msg.user_name, effective_role.as_str(), &ws.name).await;
deliver_clear_reply(&reply, msg, &ws, effective_role).await;
}
async fn deliver_clear_reply(
reply: &str,
msg: &ChannelMessage,
ws: &Workspace,
effective_role: Role,
) {
message_router::deliver_unregistered_user_response(
reply,
&message_router::AgentJob {
content: reply.to_string(),
workspace_name: ws.name.clone(),
user_name: msg.user_name.clone(),
channel: msg.channel.clone(),
kind: message_router::JobKind::UserMessage,
role: effective_role,
reply_target: Some(msg.reply_target.clone()),
pending_job_id: None,
},
&effective_role,
)
.await;
}
async fn handle_models_command(msg: &ChannelMessage, is_image: bool) {
let reply_markup = build_models_keyboard(is_image);
let content = if is_image {
"Select an image model:".to_string()
} else {
"Select a video model:".to_string()
};
let _ = mahbot::channels::telegram::send_direct(&msg.reply_target, content, Some(reply_markup))
.await;
}
fn build_models_keyboard(is_image: bool) -> serde_json::Value {
let mut rows: Vec<serde_json::Value> = Vec::new();
let (mut models, active, action_prefix) = if is_image {
(
CONFIG.image_gen_models(),
CONFIG.image_gen_model(),
"__act__set_image_model",
)
} else {
(
CONFIG.video_models(),
CONFIG.video_model(),
"__act__set_video_model",
)
};
if !models.iter().any(|m| m == &active) {
models.push(active.clone());
}
build_model_button_rows(&mut rows, &models, &active, action_prefix);
rows.push(serde_json::json!([{
"text": "Clear session",
"callback_data": "__act__clear_session|",
}]));
serde_json::json!({ "inline_keyboard": rows })
}
fn build_model_button_rows(
rows: &mut Vec<serde_json::Value>,
models: &[String],
active_model: &str,
action_prefix: &str,
) {
for model in models {
let label = if model == active_model {
format!("\u{2713} {model}")
} else {
model.clone()
};
rows.push(serde_json::json!([{
"text": label,
"callback_data": format!("{action_prefix}|{model}"),
}]));
}
}
async fn resolve_admin_workspace(msg: &ChannelMessage) -> Result<Option<String>, String> {
let selected = mahbot::users::get_raw_selected_workspace(&msg.user_name)
.await
.map_err(|e| format!("Failed to read workspace selection: {e}"))?;
match selected {
Some(name) if !name.trim().is_empty() => {
let ws = mahbot::workspace::get_by_name(&name)
.await
.map_err(|e| format!("Failed to look up workspace: {e}"))?;
match ws {
Some(_) => Ok(Some(name)),
None => Err(format!("Active workspace '{name}' no longer exists.")),
}
}
_ => Ok(None),
}
}
async fn handle_admin_command(msg: &ChannelMessage, cmd: mahbot::BotCommand) {
let maintenance_arg = if cmd == BotCommand::Maintenance {
let arg = msg.content.trim().to_ascii_lowercase();
match arg.strip_prefix("/maintenance").map_or("", str::trim) {
"on" => Some(true),
"off" => Some(false),
_ => {
send_telegram_reply(msg, "Usage: /maintenance on|off".to_string()).await;
return;
}
}
} else {
None
};
let ws_name = match resolve_admin_workspace(msg).await {
Ok(Some(name)) => name,
Ok(None) => {
send_telegram_reply(
msg,
"No active workspace β select a shared workspace in Settings β Users.".to_string(),
)
.await;
return;
}
Err(e) => {
send_telegram_reply(msg, e).await;
return;
}
};
match (cmd, maintenance_arg) {
(BotCommand::Board, _) => handle_board_listing(msg, &ws_name).await,
(BotCommand::Archive, _) => {
let count = mahbot::board::store()
.archive_all_done_and_cancelled(Some(&ws_name))
.await;
match count {
Ok(n) => send_telegram_reply(msg, format!("Archived {n} tickets.")).await,
Err(e) => send_telegram_reply(msg, format!("Failed to archive tickets: {e}")).await,
}
}
(BotCommand::Pause, _) => toggle_workspace_state(msg, &ws_name, true, false).await,
(BotCommand::Unpause, _) => toggle_workspace_state(msg, &ws_name, false, false).await,
(BotCommand::Maintenance, Some(enable)) => {
toggle_workspace_state(msg, &ws_name, enable, true).await;
}
(BotCommand::MaintenanceOn, _) => toggle_workspace_state(msg, &ws_name, true, true).await,
(BotCommand::MaintenanceOff, _) => toggle_workspace_state(msg, &ws_name, false, true).await,
_ => unreachable!(),
}
}
async fn toggle_workspace_state(
msg: &ChannelMessage,
ws_name: &str,
enable: bool,
is_maintenance: bool,
) {
let store = mahbot::workspace::store();
let result = if is_maintenance {
store.set_maintenance_enabled(ws_name, enable).await
} else {
store.set_paused(ws_name, enable).await
};
if let Err(e) = result {
send_telegram_reply(msg, format!("Failed to update workspace '{ws_name}': {e}")).await;
return;
}
let verb = match (is_maintenance, enable) {
(true, true) => "Maintenance enabled",
(true, false) => "Maintenance disabled",
(false, true) => "Workspace pipeline paused",
(false, false) => "Workspace pipeline resumed",
};
send_telegram_reply(msg, format!("{verb} for '{ws_name}'.")).await;
}
async fn handle_board_listing(msg: &ChannelMessage, ws_name: &str) {
let tickets = match mahbot::board::store()
.list_all_tickets(Some(ws_name), None)
.await
{
Ok(t) => t,
Err(e) => {
send_telegram_reply(msg, format!("Failed to load board: {e}")).await;
return;
}
};
let ordered = mahbot::board::BoardStore::board_display_order(&tickets);
if ordered.is_empty() {
send_telegram_reply(msg, "No tickets.".to_string()).await;
return;
}
let listing = ordered
.iter()
.map(|t| mahbot::channels::telegram::format_board_line(&t.phase, &t.id, &t.title))
.collect::<Vec<_>>()
.join("\n");
send_telegram_reply(msg, listing).await;
}
async fn handle_action_callback(msg: ChannelMessage, decoded: (String, String)) {
let (action, payload) = decoded;
match action.as_str() {
"set_image_model" => {
handle_set_model_action(
&msg,
&payload,
CONFIG_KEY_IMAGE_GEN_MODEL,
"Image generation",
true,
)
.await;
}
"set_video_model" => {
handle_set_model_action(&msg, &payload, CONFIG_KEY_VIDEO_MODEL, "Video", false).await;
}
"clear_session" => {
answer_telegram_callback(&msg, None).await;
handle_clear_session(&msg).await;
}
_ => {
answer_telegram_callback(&msg, None).await;
tracing::warn!(action = %action, "Unknown __act__ action β ignoring");
}
}
}
static MODEL_WRITE_LOCK: OnceLock<tokio::sync::Mutex<()>> = OnceLock::new();
async fn handle_set_model_action(
msg: &ChannelMessage,
payload: &str,
config_key: &str,
display_name: &str,
validate_image: bool,
) {
let _guard = MODEL_WRITE_LOCK
.get_or_init(|| tokio::sync::Mutex::new(()))
.lock()
.await;
if payload.is_empty() {
tracing::warn!(config_key, "{config_key} action with empty payload");
answer_telegram_callback(msg, Some("No model specified.".to_string())).await;
return;
}
if validate_image
&& let Err(e) = mahbot::tools::image_catalog::validate_image_model(payload).await
{
answer_telegram_callback(msg, Some(format!("Invalid image model: {e}"))).await;
return;
}
let store = mahbot::config_db::store();
if let Err(e) = store.set_kv(config_key, payload).await {
tracing::error!(config_key, error = %e, "Failed to save {config_key}");
answer_telegram_callback(msg, Some(format!("Failed to save model: {e}"))).await;
return;
}
let _ = CONFIG.set_string_field(config_key, payload);
answer_telegram_callback(msg, Some(format!("{display_name} model set to: {payload}"))).await;
}
async fn answer_telegram_callback(msg: &ChannelMessage, toast: Option<String>) {
let Some(cq_id) = &msg.callback_query_id else {
return;
};
if let Some(channel) = mahbot::channel_registry().get("telegram")
&& let Some(tc) = channel
.as_any()
.downcast_ref::<mahbot::channels::telegram::TelegramChannel>()
{
tc.answer_callback_query(cq_id, toast.as_deref()).await;
}
}
async fn process_channel_message(mut msg: ChannelMessage) {
tracing::info!(
"π¬ [{}] from {}: {}",
msg.channel,
msg.user_name,
mahbot::util::truncate(&msg.content, 80)
);
let (ws, (pool, pool_read_failed)) = tokio::join!(
mahbot::users::resolve_workspace_for_user_name(&msg.user_name),
mahbot::users::role_pool_status(&msg.user_name),
);
let role = mahbot::users::resolve_active_role_from_pool(&msg.user_name, &pool).await;
let (effective_role, ws) = match role {
Some(role) => {
let (effective_role, ws) =
mahbot::users::effective_role_and_workspace(role, ws, &msg.user_name, &pool);
(Some(effective_role), ws)
}
None => (None, ws),
};
msg.workspace = ws.name.clone();
let original_content = msg.content.clone();
let is_artist = matches!(effective_role, Some(mahbot::Role::Artist));
let strategy = mahbot::channels::EnrichmentStrategy {
workspace_path: effective_role.is_some().then(|| ws.as_path().to_path_buf()),
compress_images: !is_artist,
};
mahbot::channels::enrich_message(&mut msg, &strategy).await;
let persist_content = if mahbot::channels::has_only_audio_markers(&original_content) {
&msg.content
} else {
&original_content
};
broadcast_and_persist_incoming_message(&msg, &msg.content, persist_content).await;
let enriched = mahbot::channels::enrich_links(&msg.content).await;
if let Cow::Owned(s) = enriched {
tracing::info!(
channel = %msg.channel,
user_name = %msg.user_name,
"Link enricher: prepended URL summaries to message"
);
msg.content = s;
}
let Some(effective_role) = effective_role else {
if msg.channel == "telegram" && !pool_read_failed && pool.is_empty() {
send_telegram_reply(
&msg,
"You have no active role assigned β ask an admin to assign roles \
in Settings β Users."
.to_string(),
)
.await;
}
return;
};
message_router::route_user_message(
msg.content,
ws.name,
msg.user_name,
msg.channel,
effective_role,
Some(msg.reply_target),
)
.await;
}