#![warn(clippy::pedantic)]
#![cfg_attr(not(test), windows_subsystem = "windows")]
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::agent::message_router;
use mahbot::channels::broadcast_and_persist_incoming_message;
use mahbot::channels::telegram::{ControlInput, control_input, user_command_entries};
use mahbot::config::CONFIG;
use mahbot::gui::{BOOT_LOG_STORE, Dashboard, JETBRAINS_MONO, Message as DashboardMessage};
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 JETBRAINS_MONO_ITALIC_FONT_BYTES: &[u8] = include_bytes!("gui/JetBrainsMono-Italic.ttf");
const TOP_LEVEL_USAGE: &str = "\
mahbot — autonomous agentic engineering system with a GUI dashboard daemon
Usage:
mahbot Launch the GUI dashboard daemon
mahbot paused Launch the daemon with every registered workspace
paused (no ticket is claimed or advanced)
mahbot chrome <args> Browser automation CLI over the shared chrome core
mahbot debug Read-only SQL query tool against the live stores
mahbot bench-openrouter <args>
Standalone OpenRouter provider benchmark
Options:
-h, --help Print this help and exit
-V, --version Print the version and exit";
const LOG_RETENTION_HOURS: i64 = 8;
async fn bootstrap_mahbot_safe(start_paused: bool) -> Result<(), String> {
match AssertUnwindSafe(bootstrap_mahbot(start_paused))
.catch_unwind()
.await
{
Ok(Ok(())) => Ok(()),
Ok(Err(e)) => Err(format!("{e:#}")),
Err(payload) => {
let message = mahbot::util::panic_message(&*payload);
let boot_error = mahbot::boot::record_startup_failure(
"bootstrap",
mahbot::boot::startup_panic_error(&message),
);
Err(format!("{boot_error:#}"))
}
}
}
async fn bootstrap_mahbot(start_paused: bool) -> Result<()> {
mahbot::config::load_or_init()
.await
.map_err(|e| mahbot::boot::record_startup_failure("config::load_or_init", e))?;
let log_store = mahbot::boot::open_stores().await?;
mahbot::pipeline::chronicle::start_subscriber();
mahbot::config::reload_from_db()
.await
.map_err(|e| mahbot::boot::record_startup_failure("config::reload_from_db", e))?;
#[cfg(target_os = "macos")]
mahbot::audio::local_transcriber::spawn_background_init_if_enabled();
mahbot::providers::init_global()
.map_err(|e| mahbot::boot::record_startup_failure("providers::init_global", e))?;
#[cfg(target_os = "macos")]
if mahbot::audio::tts::is_config_enabled() {
let _ = mahbot::audio::tts::ensure_audio_output();
if !mahbot::audio::tts::try_load_cached() {
mahbot::audio::tts::spawn_download();
}
}
BOOT_LOG_STORE
.set(log_store.as_ref().clone())
.map_err(|_| {
mahbot::boot::record_startup_failure(
"BOOT_LOG_STORE::set",
anyhow::anyhow!("BOOT_LOG_STORE already set"),
)
})?;
if start_paused {
let workspaces = mahbot::workspace::store()
.pause_all()
.await
.map_err(|e| mahbot::boot::record_startup_failure("workspace::pause_all", e))?;
info!(
workspaces,
"Started with the paused argument: all registered workspaces are on pause"
);
}
spawn_background_tasks(log_store.clone());
info!(
version = mahbot::self_update::VERSION,
"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::self_update::back_online_message(start_paused),
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,
"shell-env-read",
mahbot::shell_env::run_reader_loop(),
);
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.clone();
move |cutoff| {
let store = store.clone();
async move { store.delete_older_than("INFO", &cutoff).await }
}
}),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"maintainer",
mahbot::agent::maintainer::run_maintainer_loop(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"archive-cancelled",
mahbot::pipeline::board::run_archive_cancelled_loop(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"alarm-sweep",
mahbot::alarms::run_alarm_sweep_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,
"temp-cleanup",
mahbot::temp::run_temp_cleanup_loop(),
);
let ipc_root = mahbot::config::CONFIG.global_storage_root();
spawn_cancellable(&mut tasks, &shutdown_token, "debug-ipc", async move {
mahbot::db::ipc::run_ipc_listener(&ipc_root, log_store.clone()).await;
});
spawn_cancellable(
&mut tasks,
&shutdown_token,
"search-engine-init",
mahbot::search_engine::init_all_engines(),
);
#[cfg(not(target_os = "macos"))]
spawn_cancellable(
&mut tasks,
&shutdown_token,
"computer-capabilities",
mahbot::tools::computer::warm_capture_probe(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"update-availability",
mahbot::self_update::run_update_availability_refresh(),
);
tokio::spawn(mahbot::self_update::run_env_named_update());
spawn_cancellable(
&mut tasks,
&shutdown_token,
"chrome-daemon",
mahbot::tools::chrome_daemon::run_watchdog(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"chrome-run-releases",
mahbot::tools::chrome_release::run_session_release_queue(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"chrome-use-manager",
mahbot::tools::chrome_daemon::run_chrome_use_management(),
);
spawn_cancellable(
&mut tasks,
&shutdown_token,
"bun-manager",
mahbot::tools::bun::run_bun_management(),
);
#[cfg(target_os = "macos")]
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::pipeline::run_management(),
);
tasks.spawn(async move {
let result = AssertUnwindSafe(mahbot::shutdown::wait_for_shutdown_signal())
.catch_unwind()
.await;
match result {
Ok(Ok(cause)) => {
info!(
"{} — force-cancelling drain",
cause.unwrap_or("Second signal received")
);
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::db::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");
#[cfg(target_os = "macos")]
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::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;
}
});
}
#[cfg(target_os = "macos")]
mahbot::channels::register_global();
rx
}
async fn shutdown_after_dashboard() {
info!(
"{} — shutting down",
mahbot::shutdown::exit_trigger().unwrap_or("Dashboard window closed")
);
mahbot::agent::registry::AGENT_REGISTRY.shutdown_all();
let release = mahbot::tools::chrome_release::flush_and_close_all_chrome_sessions();
match mahbot::shutdown::urgent_release_budget() {
Some(budget) => {
if tokio::time::timeout(budget, release).await.is_err() {
warn!("urgent shutdown: browser session release cut at its {budget:?} budget");
}
}
None => release.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::db::checkpoint::checkpoint_all_databases().await;
}
fn main() -> Result<()> {
mahbot::shutdown::install_fatal_signal_handlers();
mahbot::shutdown::install_panic_hook();
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::db::debug::run_debug()) {
Ok(()) => std::process::exit(0),
Err(e) => {
mahbot::util::print_stderr(&format!("Error: {e:#}"));
std::process::exit(1);
}
}
}
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("__env-dump") {
std::process::exit(mahbot::shell_env::dump_environment(
&std::env::args().skip(2).collect::<Vec<_>>(),
));
}
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);
}
if std::env::args().nth(1).as_deref() == Some("chrome") {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
let args = std::env::args().skip(2).collect::<Vec<_>>();
let code = rt.block_on(mahbot::run_chrome_cli(&args));
std::process::exit(code);
}
match std::env::args().nth(1).as_deref() {
Some("-h" | "--help") => {
mahbot::util::print_stdout(TOP_LEVEL_USAGE);
return Ok(());
}
Some("-V" | "--version") => {
mahbot::util::print_stdout(concat!("mahbot ", env!("CARGO_PKG_VERSION")));
return Ok(());
}
_ => {}
}
let start_paused = std::env::args().nth(1).as_deref() == Some("paused");
mahbot::shutdown::install_stop_request_sources();
let storage_root = mahbot::config::default_config_dir()
.map_err(|e| mahbot::boot::record_startup_failure("config::default_config_dir", e))?;
mahbot::temp::init_temp_root().map_err(|e| {
mahbot::boot::record_launch_failure(&storage_root, "temp::init_temp_root", e)
})?;
mahbot::self_update::acquire_lock(&storage_root).map_err(|e| {
mahbot::boot::record_launch_failure(&storage_root, "self_update::acquire_lock", e)
})?;
let window_state = mahbot::gui::read_window_state();
mahbot::gui::init_git_file_change_tx();
mahbot::gui::init_git_commit_tx();
mahbot::gui::init_board_change_tx();
mahbot::gui::init_workspace_tx();
mahbot::gui::init_users_tx();
mahbot::gui::init_user_channels_tx();
mahbot::gui::init_runtime_event_tx();
mahbot::app_icon::install_desktop_integration();
iced::application(
move || {
(
Dashboard::loading(),
iced::Task::perform(bootstrap_mahbot_safe(start_paused), 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)
.font(JETBRAINS_MONO_ITALIC_FONT_BYTES)
.default_font(JETBRAINS_MONO)
.subscription(Dashboard::subscription)
.theme(Dashboard::theme)
.window(dashboard_window_settings(&window_state))
.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(())
}
fn dashboard_window_settings(state: &mahbot::gui::WindowState) -> iced::window::Settings {
iced::window::Settings {
size: iced::Size::new(state.width, state.height),
position: state.position(),
min_size: Some(iced::Size::new(800.0, 500.0)),
icon: mahbot::app_icon::window_icon(),
#[cfg(target_os = "linux")]
platform_specific: iced::window::settings::PlatformSpecific {
application_id: mahbot::app_icon::LINUX_APP_ID.to_owned(),
..Default::default()
},
..iced::window::Settings::default()
}
}
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,
},
};
match control_input(&msg) {
Some(ControlInput::Action(decoded)) => {
spawn(handle_action_callback(msg, decoded));
}
Some(ControlInput::Command(cmd)) => {
handle_bot_command(&msg, cmd).await;
}
None => {
spawn(process_channel_message(msg));
}
}
}
}
async fn handle_bot_command(msg: &ChannelMessage, cmd: BotCommand) {
match cmd {
BotCommand::Start => handle_start_command(msg).await,
BotCommand::Clear => handle_clear_session(msg).await,
BotCommand::ImageModels | BotCommand::VideoModels => {
handle_models_command(msg, cmd == BotCommand::ImageModels).await;
}
BotCommand::Agents => handle_agents_command(msg).await,
BotCommand::Update => mahbot::self_update::handle_update_command(msg).await,
BotCommand::Board
| BotCommand::Archive
| BotCommand::Pause
| BotCommand::Unpause
| BotCommand::Maintenance
| BotCommand::MaintenanceOn
| BotCommand::MaintenanceOff
| BotCommand::Workspace => {
if mahbot::users::is_admin(&msg.user_name).await {
match cmd {
BotCommand::Workspace => handle_workspace_command(msg).await,
BotCommand::Board => handle_board_listing(msg).await,
BotCommand::Archive => handle_archive_command(msg).await,
_ => handle_admin_command(msg, cmd).await,
}
} else {
send_telegram_reply(msg, mahbot::self_update::ADMIN_ONLY_CMD_MSG.to_string()).await;
}
}
}
}
fn with_tail(base: String, tail: Option<&str>) -> String {
match tail {
Some(tail) => format!("{base} {tail}"),
None => base,
}
}
async fn switcher_pointer() -> Option<&'static str> {
mahbot::users::workspace_switcher_available()
.await
.then_some("Use /workspace to choose the active workspace.")
}
async fn send_telegram_reply(msg: &ChannelMessage, content: String) {
mahbot::channels::telegram::send_reply(&msg.reply_target, &content).await;
}
async fn handle_agents_command(msg: &ChannelMessage) {
send_telegram_reply(
msg,
"You have only the Assistant role — there is nothing to switch.".to_string(),
)
.await;
}
async fn handle_workspace_command(msg: &ChannelMessage) {
let workspaces = match mahbot::users::registered_workspaces().await {
Ok(workspaces) => workspaces,
Err(e) => {
send_telegram_reply(msg, format!("Failed to load workspaces: {e}")).await;
return;
}
};
let unavailable = if mahbot::users::switcher_exists(workspaces.len()) {
None
} else if workspaces.is_empty() {
Some("No shared workspace is registered — there is nothing to switch.")
} else {
Some("Only one shared workspace is registered — there is nothing to switch.")
};
if let Some(refusal) = stray_text_refusal(&msg.content) {
let tail = unavailable.unwrap_or("Send it on its own to get the list.");
send_telegram_reply(msg, with_tail(refusal, Some(tail))).await;
return;
}
if let Some(unavailable) = unavailable {
send_telegram_reply(msg, unavailable.to_string()).await;
return;
}
let active = match mahbot::users::get_raw_selected_workspace(&msg.user_name).await {
Ok(name) => name,
Err(e) => {
send_telegram_reply(msg, format!("Failed to read the active workspace: {e}")).await;
return;
}
};
let keyboard =
mahbot::channels::telegram::workspace_picker_keyboard(&workspaces, active.as_deref());
if let Err(e) = mahbot::channels::telegram::send_direct(
&msg.reply_target,
"Select the active workspace:".to_string(),
Some(keyboard),
)
.await
{
warn!(error = %e, "Failed to send the workspace picker");
}
}
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 = match clear_session(&msg.user_name, effective_role.as_str(), &ws.name).await {
Ok(reply) => reply,
Err(e) => {
tracing::warn!(user = %msg.user_name, error = %e, "/clear failed — session kept");
format!("Failed to clear the session: {e:#}")
}
};
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::MessageKind::UserMessage,
role: effective_role,
reply_target: Some(msg.reply_target.clone()),
pending_job_id: None,
originating_workspace: None,
},
&effective_role,
&[],
)
.await;
}
async fn handle_models_command(msg: &ChannelMessage, is_image: bool) {
let reply_markup = build_models_keyboard(is_image, &msg.user_name).await;
let content = if is_image {
"Select an image model:".to_string()
} else {
"Select a video model:".to_string()
};
if let Err(e) =
mahbot::channels::telegram::send_direct(&msg.reply_target, content, Some(reply_markup))
.await
{
warn!(error = %e, "Failed to send the model picker");
}
}
async fn build_models_keyboard(is_image: bool, user_name: &str) -> 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(),
mahbot::users::resolve_image_gen_model(user_name).await,
"__act__set_image_model",
)
} else {
(
CONFIG.video_models(),
mahbot::users::resolve_video_model(user_name).await,
"__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<Workspace>, 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(ws) => Ok(Some(ws)),
None => Err(format!("Active workspace '{name}' no longer exists.")),
}
}
_ => Ok(None),
}
}
async fn handle_admin_command(msg: &ChannelMessage, cmd: mahbot::BotCommand) {
if matches!(cmd, BotCommand::Pause | BotCommand::Unpause)
&& let Some(refusal) = stray_text_refusal(&msg.content)
{
let tail = switcher_pointer().await;
send_telegram_reply(msg, with_tail(refusal, tail)).await;
return;
}
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 = match resolve_admin_workspace(msg).await {
Ok(Some(ws)) => ws,
Ok(None) => {
let text = with_tail(
"No active workspace is selected.".to_string(),
switcher_pointer().await,
);
send_telegram_reply(msg, text).await;
return;
}
Err(e) => {
send_telegram_reply(msg, e).await;
return;
}
};
match (cmd, maintenance_arg) {
(BotCommand::Pause, _) => toggle_workspace_state(msg, &ws, true, false).await,
(BotCommand::Unpause, _) => toggle_workspace_state(msg, &ws, false, false).await,
(BotCommand::Maintenance, Some(enable)) => {
toggle_workspace_state(msg, &ws, enable, true).await;
}
(BotCommand::MaintenanceOn, _) => toggle_workspace_state(msg, &ws, true, true).await,
(BotCommand::MaintenanceOff, _) => toggle_workspace_state(msg, &ws, false, true).await,
_ => unreachable!(),
}
}
fn stray_text_refusal(content: &str) -> Option<String> {
let mut words = content.split_whitespace();
let cmd_word = words.next().unwrap_or_default();
words.next()?;
Some(format!("`{cmd_word}` takes no other text."))
}
async fn toggle_workspace_state(
msg: &ChannelMessage,
ws: &Workspace,
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 '{}': {e}", ws.display_name()),
)
.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.display_name())).await;
}
async fn handle_board_listing(msg: &ChannelMessage) {
let tickets = match mahbot::pipeline::board::store()
.list_all_tickets(None, None)
.await
{
Ok(t) => t,
Err(e) => {
send_telegram_reply(msg, format!("Failed to load the board: {e}")).await;
return;
}
};
let ordered = mahbot::pipeline::board::BoardStore::board_display_order(&tickets);
send_telegram_reply(
msg,
mahbot::channels::telegram::board_listing_text(&ordered),
)
.await;
}
async fn handle_archive_command(msg: &ChannelMessage) {
if let Some(refusal) = stray_text_refusal(&msg.content) {
send_telegram_reply(
msg,
with_tail(
refusal,
Some("It has no target — it archives done and cancelled tickets across all workspaces."),
),
)
.await;
return;
}
match mahbot::pipeline::board::store()
.archive_all_done_and_cancelled()
.await
{
Ok(n) => {
send_telegram_reply(
msg,
format!(
"Archived {n} ticket{} across all workspaces.",
if n == 1 { "" } else { "s" }
),
)
.await;
}
Err(e) => send_telegram_reply(msg, format!("Failed to archive tickets: {e}")).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, "Image generation", true).await;
}
"set_video_model" => {
handle_set_model_action(&msg, &payload, "Video", false).await;
}
"clear_session" => {
answer_telegram_callback(&msg, None).await;
handle_clear_session(&msg).await;
}
"set_workspace" => handle_set_workspace_action(&msg, &payload).await,
_ => {
answer_telegram_callback(
&msg,
Some("This action is no longer available.".to_string()),
)
.await;
tracing::warn!(action = %action, "Unknown __act__ action — ignoring");
}
}
}
static PICKER_WRITE_LOCK: OnceLock<tokio::sync::Mutex<()>> = OnceLock::new();
async fn handle_set_model_action(
msg: &ChannelMessage,
payload: &str,
display_name: &str,
validate_image: bool,
) {
let toast = {
let _guard = PICKER_WRITE_LOCK
.get_or_init(|| tokio::sync::Mutex::new(()))
.lock()
.await;
match set_user_model(&msg.user_name, payload, display_name, validate_image).await {
Ok(toast) => {
refresh_pressed_models_keyboard(msg, validate_image).await;
Some(toast)
}
Err(toast) => Some(toast),
}
};
answer_telegram_callback(msg, toast).await;
}
async fn handle_set_workspace_action(msg: &ChannelMessage, payload: &str) {
{
let _guard = PICKER_WRITE_LOCK
.get_or_init(|| tokio::sync::Mutex::new(()))
.lock()
.await;
let reply = match mahbot::users::set_active_workspace(&msg.user_name, payload).await {
Ok(ws) => {
refresh_pressed_workspaces_keyboard(msg).await;
format!("Active workspace: {}", ws.display_name())
}
Err(e) => e.to_string(),
};
send_telegram_reply(msg, reply).await;
}
answer_telegram_callback(msg, None).await;
}
async fn refresh_pressed_picker_keyboard(
msg: &ChannelMessage,
keyboard: &serde_json::Value,
what: &str,
) {
let (Some(chat_id), Some(message_id)) = (&msg.chat_id, msg.message_id) else {
return;
};
if let Some(channel) = mahbot::channel_registry().get("telegram")
&& let Some(tc) = channel
.as_any()
.downcast_ref::<mahbot::channels::telegram::TelegramChannel>()
&& let Err(e) = tc.edit_reply_markup(chat_id, message_id, keyboard).await
{
tracing::debug!(?e, what, "picker keyboard refresh skipped");
}
}
async fn refresh_pressed_models_keyboard(msg: &ChannelMessage, is_image: bool) {
let keyboard = build_models_keyboard(is_image, &msg.user_name).await;
refresh_pressed_picker_keyboard(msg, &keyboard, "model").await;
}
async fn refresh_pressed_workspaces_keyboard(msg: &ChannelMessage) {
let (Ok(workspaces), Ok(active)) = (
mahbot::users::registered_workspaces().await,
mahbot::users::get_raw_selected_workspace(&msg.user_name).await,
) else {
return;
};
let keyboard =
mahbot::channels::telegram::workspace_picker_keyboard(&workspaces, active.as_deref());
refresh_pressed_picker_keyboard(msg, &keyboard, "workspace").await;
}
async fn set_user_model(
user_name: &str,
payload: &str,
display_name: &str,
validate_image: bool,
) -> Result<String, String> {
if payload.is_empty() {
tracing::warn!(display_name, "{display_name} action with empty payload");
return Err("No model specified.".to_string());
}
if validate_image
&& let Err(e) = mahbot::tools::media_catalog::image::validate_image_model(payload).await
{
return Err(format!("Invalid image model: {e}"));
}
let store = mahbot::users::store();
if !store.user_exists(user_name).await.unwrap_or(false) {
return Err("User is not registered.".to_string());
}
let result = if validate_image {
store.set_image_gen_model(user_name, payload).await
} else {
store.set_video_model(user_name, payload).await
};
match result {
Ok(()) => Ok(format!("{display_name} model set to: {payload}")),
Err(e) => {
tracing::error!(user_name, error = %e, "Failed to save per-user model");
Err(format!("Failed to save model: {e}"))
}
}
}
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(msg: ChannelMessage) {
let user_name = msg.user_name.clone();
let channel = msg.channel.clone();
let reply_target = msg.reply_target.clone();
let mut parts = if msg.parts.is_empty() {
vec![msg]
} else {
msg.parts
};
let ws = mahbot::users::resolve_workspace_for_user_name(&user_name).await;
let ws = mahbot::users::effective_workspace_for_role(mahbot::Role::Assistant, ws, &user_name);
for part in &mut parts {
part.workspace.clone_from(&ws.name);
}
let mut originals = Vec::with_capacity(parts.len());
for part in &parts {
tracing::info!(
"💬 [{}] from {}: {}",
part.channel,
part.user_name,
mahbot::util::truncate(&part.content, 80)
);
originals.push(part.content.clone());
}
futures_util::future::join_all(
parts
.iter_mut()
.map(|part| mahbot::channels::enrich_message(part, Some(ws.as_path()))),
)
.await;
for (part, original_content) in parts.iter_mut().zip(&originals) {
let content_for_history =
mahbot::channels::persist_content(original_content, &part.content);
broadcast_and_persist_incoming_message(part, &part.content, &content_for_history).await;
let enriched = mahbot::channels::enrich_links(&part.content).await;
if let Cow::Owned(s) = enriched {
tracing::info!(
channel = %part.channel,
user_name = %part.user_name,
"Link enricher: prepended URL summaries to message"
);
part.content = s;
}
if let Some(reply) = part.reply_reference.clone() {
part.content = mahbot::channels::apply_reply_marker(&part.content, &reply);
}
}
let content = mahbot::channels::compose_group_content(parts);
message_router::route_user_message(
content,
ws.name,
user_name,
channel,
mahbot::Role::Assistant,
Some(reply_target),
)
.await;
}