use crate::channels::{
broadcast_and_persist_agent_response, spawn_scoped_typing_task, stop_typing,
};
use crate::session::manager_session_key;
use crate::users::UserRecord;
use crate::{ChatEvent, Role, SendMessage};
use std::collections::HashSet;
use tokio::sync::OnceCell;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info, warn};
#[derive(Debug)]
pub enum JobKind {
UserMessage,
TicketNotify,
AskToolResult,
}
pub struct ManagerJob {
pub content: String,
pub workspace_name: String,
pub kind: JobKind,
}
pub static MANAGER_QUEUE: OnceCell<ManagerQueue> = OnceCell::const_new();
pub fn init_global() -> anyhow::Result<()> {
let (tx, rx) = mpsc::unbounded_channel::<ManagerJob>();
tokio::spawn(consumer_loop(rx));
MANAGER_QUEUE
.set(ManagerQueue { tx })
.map_err(|_| anyhow::anyhow!("MANAGER_QUEUE already initialized"))?;
Ok(())
}
pub fn manager_queue() -> &'static ManagerQueue {
MANAGER_QUEUE
.get()
.expect("MANAGER_QUEUE not initialized — call manager_queue::init_global() in main()")
}
pub struct ManagerQueue {
tx: mpsc::UnboundedSender<ManagerJob>,
}
impl ManagerQueue {
pub fn enqueue(&self, job: ManagerJob) {
if let Err(e) = self.tx.send(job) {
error!("Failed to enqueue Manager job: {e}");
}
}
}
async fn setup_telegram_typing(
users: &[UserRecord],
) -> Vec<(CancellationToken, tokio::task::JoinHandle<()>)> {
let telegram_channel = crate::channel_registry().get("telegram");
let Some(ref tg_channel) = telegram_channel else {
return Vec::new();
};
let mut typing_tasks = Vec::new();
let mut seen_targets = HashSet::new();
for user in users {
let Some(telegram_binding) = user.channels.iter().find(|b| b.channel == "telegram") else {
continue;
};
let Some(reply_target) = &telegram_binding.reply_target else {
continue;
};
if !seen_targets.insert(reply_target.clone()) {
continue;
}
let Some(recipient) = tg_channel.resolve_recipient(&user.name, reply_target) else {
continue;
};
if let Err(e) = tg_channel.start_typing(&recipient).await {
debug!("Manager queue: telegram start_typing failed: {e}");
}
let cancel = CancellationToken::new();
let handle = spawn_scoped_typing_task(recipient, "telegram".to_string(), cancel.clone());
typing_tasks.push((cancel, handle));
}
typing_tasks
}
#[allow(clippy::too_many_lines)]
async fn consumer_loop(mut rx: mpsc::UnboundedReceiver<ManagerJob>) {
let shutdown = crate::shutdown::shutdown_token();
loop {
if shutdown.is_cancelled() {
info!("Manager queue: shutting down — agent queue drained");
break;
}
let job = tokio::select! {
job = rx.recv() => {
match job {
Some(job) => job,
None => break, }
}
() = shutdown.cancelled() => {
info!("Manager queue: shutting down (global shutdown)");
break;
}
};
info!(
workspace = %job.workspace_name,
kind = ?job.kind,
"Manager queue: processing job",
);
let ws = match crate::workspace::get_by_name(&job.workspace_name).await {
Ok(Some(ws)) => ws,
Ok(None) => {
error!(
workspace = %job.workspace_name,
"Manager queue: workspace not found — skipping job",
);
continue;
}
Err(e) => {
error!(
workspace = %job.workspace_name,
error = %e,
"Manager queue: failed to look up workspace — skipping job",
);
continue;
}
};
let session_key = manager_session_key(&ws.name);
let users: Vec<UserRecord> = match crate::users::USER_STORE.get() {
Some(store) => store
.find_by_workspace(&job.workspace_name)
.await
.unwrap_or_default(),
None => Vec::new(),
};
let typing_tasks = setup_telegram_typing(&users).await;
let _ = crate::CHAT_BROADCAST.get().map(|tx| {
for user in &users {
let _ = tx.send(ChatEvent::Typing {
user_name: user.name.clone(),
is_typing: true,
});
}
});
let message = match job.kind {
JobKind::UserMessage => {
let drained = crate::ticket_buffer::drain(&job.workspace_name);
if drained.is_empty() {
job.content
} else {
format!("{drained}\n{content}", content = job.content)
}
}
JobKind::TicketNotify | JobKind::AskToolResult => job.content,
};
let (_agent, response) =
crate::agent::run_agent(session_key, Role::Manager, &ws, None, &message).await;
for (cancel, handle) in typing_tasks {
cancel.cancel();
stop_typing(handle).await;
}
let _ = crate::CHAT_BROADCAST.get().map(|tx| {
for user in &users {
let _ = tx.send(ChatEvent::Typing {
user_name: user.name.clone(),
is_typing: false,
});
}
});
let Some(response) = response else {
continue;
};
let reply_markup: Option<serde_json::Value> = None;
if users.is_empty() {
warn!(
workspace = %job.workspace_name,
"Manager queue: no users with workspace — response delivered to nobody",
);
}
let content = &response;
let agent_role = Some("manager".to_string());
let workspace = &job.workspace_name;
{
let mut seen_names = HashSet::new();
for user in &users {
if !seen_names.insert(&user.name) {
continue;
}
let channel = user.channels.first().map_or("gui", |b| b.channel.as_str());
broadcast_and_persist_agent_response(
&user.name,
channel,
content,
agent_role.clone(),
workspace,
reply_markup.clone(),
)
.await;
}
}
let channels = crate::channel_registry().list();
if channels.is_empty() {
error!("Manager queue: no channels registered");
continue;
}
for (channel_name, channel) in &channels {
for user in &users {
for binding in &user.channels {
let reply_target = binding.reply_target.as_deref().unwrap_or(&user.name);
let Some(recipient) = channel.resolve_recipient(&user.name, reply_target)
else {
continue;
};
if let Err(e) = channel
.send(&SendMessage {
content: content.clone(),
recipient,
reply_markup: reply_markup.clone(),
agent_role: agent_role.clone(),
workspace: workspace.clone(),
})
.await
{
error!(
channel = %channel_name,
user = %user.name,
"Manager queue: failed to send response to {}: {e}",
user.name,
);
}
}
}
}
}
}