use futures_util::FutureExt;
use serde::{Deserialize, Serialize};
use std::borrow::Cow;
use std::collections::{HashMap, HashSet};
use std::sync::{OnceLock, RwLock};
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info, warn};
use crate::channels::{
broadcast_and_persist_agent_response, spawn_scoped_typing_task, stop_typing,
};
use crate::turso;
use crate::users::UserRecord;
use crate::util::UnwrapPoison;
use crate::{Channel, ChatEvent, Role, SendMessage, Workspace};
const AGENT_FAILURE_EMOJI: &str = "🤖⚠️🔄";
fn telegram_role_emoji(role: Role) -> &'static str {
match role {
Role::Manager => "🤖",
Role::Engineer => "🔧",
Role::Analyst => "🔍",
Role::Coder => "💻",
Role::Qa => "🔨",
Role::Reviewer => "✅",
Role::Discovery => "🔎",
Role::Artist => "🎨",
Role::Maintainer => "⚙️",
Role::Sanitation => "🧼",
Role::Assistant => "💬",
}
}
#[must_use]
fn telegram_delivery_content<'a>(
channel: &str,
role: Role,
recipient_roles: &[String],
response: &'a str,
) -> Cow<'a, str> {
if channel != "telegram" || recipient_roles.len() < 2 {
return Cow::Borrowed(response);
}
Cow::Owned(format!(
"{} {}:\n{}",
telegram_role_emoji(role),
crate::role::role_info(&role).display_label,
response
))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum JobKind {
UserMessage,
TicketNotify,
AnalyzeToolResult,
ResearchResult,
TicketComment,
RecoveryRetry,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentJob {
pub content: String,
pub workspace_name: String,
pub user_name: String,
pub channel: String,
pub kind: JobKind,
pub role: Role,
pub reply_target: Option<String>,
#[serde(default)]
pub pending_job_id: Option<String>,
}
static ROUTER: OnceLock<RwLock<HashMap<String, mpsc::UnboundedSender<AgentJob>>>> = OnceLock::new();
pub fn init_global() -> anyhow::Result<()> {
ROUTER
.set(RwLock::new(HashMap::new()))
.map_err(|_| anyhow::anyhow!("ROUTER already initialised"))?;
Ok(())
}
pub fn route(agent_id: &str, job: AgentJob) {
{
let map = ROUTER
.get()
.expect("ROUTER not initialised — call init_global() first");
let guard = map.read().unwrap_poison();
if let Some(tx) = guard.get(agent_id) {
if tx.send(job).is_err() {
error!(agent_id = %agent_id, "Router: consumer dropped — failed to route job");
}
return;
}
}
let map = ROUTER
.get()
.expect("ROUTER not initialised — call init_global() first");
let mut guard = map.write().unwrap_poison();
if let Some(tx) = guard.get(agent_id) {
if tx.send(job).is_err() {
error!(agent_id = %agent_id, "Router: consumer dropped (double-check) — failed to route job");
}
return;
}
let (tx, rx) = mpsc::unbounded_channel::<AgentJob>();
let agent_id_for_consumer = agent_id.to_string();
let agent_id_for_cleanup = agent_id_for_consumer.clone();
tokio::spawn(async move {
let result = std::panic::AssertUnwindSafe(consumer_loop(agent_id_for_consumer, rx))
.catch_unwind()
.await;
if let Some(map) = ROUTER.get() {
let mut guard = map.write().unwrap_poison();
guard.remove(&agent_id_for_cleanup);
}
if let Err(panic) = result {
error!(
agent_id = %agent_id_for_cleanup,
"Consumer loop panicked — entry removed from router table",
);
error!(agent_id = %agent_id_for_cleanup, panic = %crate::util::panic_message(&*panic), "Consumer loop panic message");
} else {
debug!(
agent_id = %agent_id_for_cleanup,
"Consumer loop exited — removed from router table",
);
}
});
if tx.send(job).is_err() {
error!(agent_id = %agent_id, "Router: brand-new consumer dropped immediately — failed to route job");
}
guard.insert(agent_id.to_string(), tx);
}
pub async fn route_user_message(
content: String,
workspace_name: String,
user_name: String,
channel: String,
role: Role,
reply_target: Option<String>,
) {
let user_name =
crate::session::normalize_user_name(&user_name, "route_user_message").to_string();
let agent_id = crate::session::resolve_agent_id(&user_name, role.as_str(), &workspace_name);
let mut job = AgentJob {
content,
workspace_name,
user_name,
channel,
kind: JobKind::UserMessage,
role,
reply_target,
pending_job_id: None,
};
let mut persisted = false;
if job.role == Role::Manager {
let id = crate::generate_id();
match persist_pending(&job, id.clone()).await {
Ok(()) => {
job.pending_job_id = Some(id);
persisted = true;
}
Err(e) => {
warn!(
error = %e,
"Failed to persist manager message — routing best-effort (at-most-once)",
);
}
}
if crate::shutdown::aborting() && persisted {
return;
}
}
route(&agent_id, job);
}
async fn persist_pending(job: &AgentJob, id: String) -> anyhow::Result<()> {
let now = turso::now();
crate::session::store()
.conn
.execute(
crate::jobs::PENDING_JOB_INSERT_SQL,
crate::jobs::pending_job_params(&id, job, &now)?,
)
.await?;
Ok(())
}
pub fn register_agent(agent_id: &str) -> mpsc::UnboundedReceiver<AgentJob> {
let (tx, rx) = mpsc::unbounded_channel::<AgentJob>();
let map = ROUTER
.get()
.expect("ROUTER not initialised — call init_global() first");
let mut guard = map.write().unwrap_poison();
guard.insert(agent_id.to_string(), tx);
rx
}
pub fn unregister_agent(agent_id: &str) {
let Some(map) = ROUTER.get() else { return };
let mut guard = map.write().unwrap_poison();
guard.remove(agent_id);
}
#[cfg(test)]
pub(crate) fn router_contains(agent_id: &str) -> bool {
let Some(map) = ROUTER.get() else {
return false;
};
map.read().unwrap_poison().contains_key(agent_id)
}
pub fn try_route(agent_id: &str, job: AgentJob) -> bool {
let Some(map) = ROUTER.get() else {
return false;
};
let guard = map.read().unwrap_poison();
if let Some(tx) = guard.get(agent_id) {
tx.send(job).is_ok()
} else {
false
}
}
async fn resolve_workspace(workspace_name: &str) -> anyhow::Result<Option<Workspace>> {
match crate::workspace::get_by_name(workspace_name).await? {
Some(ws) => Ok(Some(ws)),
None if crate::users::is_personal_workspace(workspace_name) => {
let user_name = crate::users::personal_user_name(workspace_name)
.expect("invariant: is_personal_workspace checked prefix");
let path = crate::users::personal_workspace_path(user_name);
Ok(Some(crate::users::personal_workspace_struct(
user_name, &path,
)))
}
None => Ok(None),
}
}
#[expect(clippy::too_many_lines)]
async fn consumer_loop(agent_id: String, mut rx: mpsc::UnboundedReceiver<AgentJob>) {
let shutdown = crate::shutdown::shutdown_token();
loop {
if shutdown.is_cancelled() {
info!(agent_id = %agent_id, "Message router: shutting down — queue drained");
break;
}
if crate::shutdown::aborting() {
info!(agent_id = %agent_id, "Message router: draining — no new jobs pulled");
break;
}
let job = tokio::select! {
job = rx.recv() => {
match job {
Some(job) => job,
None => break,
}
}
() = shutdown.cancelled() => {
info!(agent_id = %agent_id, "Message router: shutting down (global shutdown)");
break;
}
};
debug!(
agent_id = %agent_id,
workspace = %job.workspace_name,
user = %job.user_name,
kind = ?job.kind,
"Message router: processing job",
);
let ws = match resolve_workspace(&job.workspace_name).await {
Ok(Some(ws)) => ws,
Ok(None) => {
error!(
agent_id = %agent_id,
workspace = %job.workspace_name,
"Message router: workspace not found — skipping job",
);
continue;
}
Err(e) => {
error!(
agent_id = %agent_id,
workspace = %job.workspace_name,
error = %e,
"Message router: failed to look up workspace — skipping job",
);
continue;
}
};
let role = job.role;
let users: Vec<UserRecord> = if role == Role::Manager {
match crate::users::USER_STORE.get() {
Some(store) => store
.find_by_workspace(&job.workspace_name)
.await
.unwrap_or_default(),
None => Vec::new(),
}
} else {
match resolve_single_user(&job.user_name).await {
Some(user) => vec![user],
None => Vec::new(),
}
};
let typing_tasks = setup_telegram_typing(&users).await;
broadcast_typing(&users, &job.workspace_name, true);
let message = match (role, job.kind) {
(Role::Manager, JobKind::UserMessage) => {
let drained = crate::ticket_buffer::drain(&job.workspace_name);
if drained.is_empty() {
job.content.clone()
} else {
format!("{drained}\n{content}", content = job.content)
}
}
(_, JobKind::TicketComment) => {
warn!(
agent_id = %agent_id,
"Consumer loop received TicketComment — was try_route() used instead of route()? Discarding",
);
continue;
}
_ => job.content.clone(),
};
let (agent, response) = crate::agent::run_agent(
agent_id.clone(),
role,
&ws,
None,
&message,
job.user_name.clone(),
job.channel.clone(),
None,
false,
None,
None,
None,
)
.await;
for (cancel, handle) in typing_tasks {
cancel.cancel();
stop_typing(handle).await;
}
broadcast_typing(&users, &job.workspace_name, false);
let Some(response) = response else {
if job.kind == JobKind::UserMessage
&& !agent.is_cancelled()
&& !crate::shutdown::aborting()
{
deliver_unregistered_user_response(AGENT_FAILURE_EMOJI, &job, &role).await;
}
confirm_pending_delivery(&job).await;
continue;
};
match role {
Role::Manager => {
deliver_manager_response(&response, &users, &job).await;
}
_ => {
if users.is_empty() {
deliver_unregistered_user_response(&response, &job, &role).await;
} else {
deliver_single_user_response(&response, &users[0], &job, &role).await;
}
}
}
confirm_pending_delivery(&job).await;
}
}
async fn confirm_pending_delivery(job: &AgentJob) {
let Some(id) = job.pending_job_id.as_deref() else {
return;
};
for attempt in 0..2 {
match crate::jobs::delete_pending_job(&crate::session::store().conn, id).await {
Ok(()) => return,
Err(e) if attempt == 0 => {
warn!(
pending_job = %id,
error = %e,
"Failed to confirm pending delivery — retrying once",
);
}
Err(e) => {
warn!(
pending_job = %id,
error = %e,
"Pending delivery confirm failed — duplicate re-delivery at next boot accepted",
);
}
}
}
}
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!("Message router: 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
}
fn broadcast_typing(users: &[UserRecord], workspace: &str, is_typing: bool) {
if let Some(tx) = crate::CHAT_BROADCAST.get() {
for user in users {
let _ = tx.send(ChatEvent::Typing {
user_name: user.name.clone(),
is_typing,
workspace: workspace.to_string(),
});
}
}
}
enum DeliverOutcome {
Sent,
Unresolvable,
Failed(anyhow::Error),
}
async fn deliver_on_channel(
channel: &dyn Channel,
user_name: &str,
reply_target: &str,
response: &str,
) -> DeliverOutcome {
let Some(recipient) = channel.resolve_recipient(user_name, reply_target) else {
return DeliverOutcome::Unresolvable;
};
match channel
.send(&SendMessage {
content: response.to_string(),
recipient,
reply_markup: None,
})
.await
{
Ok(()) => DeliverOutcome::Sent,
Err(e) => DeliverOutcome::Failed(e),
}
}
async fn deliver_manager_response(response: &str, users: &[UserRecord], job: &AgentJob) {
if users.is_empty() {
warn!(
workspace = %job.workspace_name,
"Message router [manager]: no users with workspace — response delivered to nobody",
);
}
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,
response,
agent_role.clone(),
workspace,
)
.await;
}
}
let channels = crate::channel_registry().list();
if channels.is_empty() {
error!("Message router [manager]: no channels registered");
return;
}
for (channel_name, channel) in &channels {
for user in users {
let content =
telegram_delivery_content(channel_name, Role::Manager, &user.roles, response);
for binding in &user.channels {
let reply_target = binding.reply_target.as_deref().unwrap_or(&user.name);
if let DeliverOutcome::Failed(e) =
deliver_on_channel(channel.as_ref(), &user.name, reply_target, &content).await
{
error!(
channel = %channel_name,
user = %user.name,
"Message router [manager]: failed to send response to {}: {e}",
user.name,
);
}
}
}
}
}
async fn deliver_single_user_response(
response: &str,
user: &UserRecord,
job: &AgentJob,
role: &Role,
) {
let channel = job.channel.as_str();
broadcast_and_persist_agent_response(
&user.name,
channel,
response,
Some(role.as_str().to_string()),
&job.workspace_name,
)
.await;
let channels = crate::channel_registry().list();
if channels.is_empty() {
error!(
workspace = %job.workspace_name,
user = %user.name,
"Message router [{role}]: no channels registered",
);
return;
}
for (channel_name, chan) in &channels {
let content = telegram_delivery_content(channel_name, *role, &user.roles, response);
for binding in &user.channels {
if binding.channel != *channel_name {
continue;
}
let target_addr = binding.reply_target.as_deref().unwrap_or(&user.name);
match deliver_on_channel(chan.as_ref(), &user.name, target_addr, &content).await {
DeliverOutcome::Failed(e) => error!(
channel = %channel_name,
user = %user.name,
"Message router [{role}]: failed to send response to {}: {e}",
user.name,
),
DeliverOutcome::Unresolvable => warn!(
channel = %channel_name,
user = %user.name,
"Message router [{role}]: cannot resolve recipient on {} for {} — \
response was persisted but not delivered via transport",
channel_name,
user.name,
),
DeliverOutcome::Sent => {}
}
}
}
}
pub async fn deliver_unregistered_user_response(response: &str, job: &AgentJob, role: &Role) {
let ch = job.channel.as_str();
let user_name =
crate::session::normalize_user_name(&job.user_name, "deliver_unregistered_user_response");
broadcast_and_persist_agent_response(
user_name,
ch,
response,
Some(role.as_str().to_string()),
&job.workspace_name,
)
.await;
let Some(chan) = crate::channel_registry().get(ch) else {
return;
};
let reply_target = job.reply_target.as_deref().unwrap_or(user_name);
match deliver_on_channel(chan.as_ref(), user_name, reply_target, response).await {
DeliverOutcome::Unresolvable => warn!(
workspace = %job.workspace_name,
user = %user_name,
channel = %ch,
"Message router [{role}]: cannot resolve recipient for unregistered user — \
response was persisted but not delivered via transport",
),
DeliverOutcome::Failed(e) => error!(
channel = %ch,
user = %user_name,
"Message router [{role}]: failed to send response to unregistered user: {e}",
),
DeliverOutcome::Sent => {}
}
}
async fn resolve_single_user(user_name: &str) -> Option<UserRecord> {
let store = crate::users::USER_STORE.get()?;
match store.find_by_name(user_name).await {
Ok(Some(user)) => Some(user),
_ => None,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::Role;
use crate::channels::gui::GuiChannel;
use std::sync::Arc;
fn make_job(role: Role, workspace: &str, user: &str, channel: &str) -> AgentJob {
AgentJob {
content: String::new(),
workspace_name: workspace.to_string(),
user_name: user.to_string(),
channel: channel.to_string(),
kind: JobKind::UserMessage,
role,
reply_target: None,
pending_job_id: None,
}
}
#[tokio::test]
async fn test_route_creates_consumer_entry() {
let _ = init_global();
let id = "_test_ua_creates_entry";
route(id, make_job(Role::Assistant, "", "", ""));
let map = ROUTER.get().unwrap();
let guard = map.read().unwrap_poison();
assert!(guard.contains_key(id));
drop(guard);
let map = ROUTER.get().unwrap();
let mut guard = map.write().unwrap_poison();
guard.remove(id);
}
#[tokio::test]
async fn test_route_reuses_existing_consumer() {
let _ = init_global();
let id = "_test_ua_reuses";
route(id, make_job(Role::Assistant, "ws", "alice", "gui"));
route(id, make_job(Role::Assistant, "ws", "bob", "gui"));
let map = ROUTER.get().unwrap();
let guard = map.read().unwrap_poison();
assert!(guard.contains_key(id));
drop(guard);
let map = ROUTER.get().unwrap();
let mut guard = map.write().unwrap_poison();
guard.remove(id);
}
#[tokio::test]
async fn test_route_multiple_agents_get_separate_consumers() {
let _ = init_global();
let id_a = "_test_ua_mult_a";
let id_b = "_test_ua_mult_b";
route(id_a, make_job(Role::Assistant, "ws", "alice", "gui"));
route(id_b, make_job(Role::Engineer, "ws", "bob", "gui"));
let map = ROUTER.get().unwrap();
let guard = map.read().unwrap_poison();
assert!(guard.contains_key(id_a));
assert!(guard.contains_key(id_b));
drop(guard);
let map = ROUTER.get().unwrap();
let mut guard = map.write().unwrap_poison();
guard.remove(id_a);
guard.remove(id_b);
}
#[tokio::test]
async fn test_consumer_loop_exits_gracefully_on_sender_drop() {
let _ = init_global();
let id = "_test_ua_on_close";
route(id, make_job(Role::Assistant, "", "", ""));
assert!(
ROUTER
.get()
.unwrap()
.read()
.unwrap_poison()
.contains_key(id),
"consumer should be registered after first route",
);
let dropped_tx = ROUTER
.get()
.unwrap()
.write()
.unwrap_poison()
.remove(id)
.expect("sender should exist");
drop(dropped_tx);
let (dummy_tx, _) = mpsc::unbounded_channel::<AgentJob>();
ROUTER
.get()
.unwrap()
.write()
.unwrap_poison()
.insert(id.to_string(), dummy_tx);
let deadline = tokio::time::Instant::now() + tokio::time::Duration::from_secs(2);
let mut cleaned_up = false;
while tokio::time::Instant::now() < deadline {
if !ROUTER
.get()
.unwrap()
.read()
.unwrap_poison()
.contains_key(id)
{
cleaned_up = true;
break;
}
tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
}
assert!(
cleaned_up,
"consumer should have exited and its cleanup wrapper should have removed the entry",
);
}
#[tokio::test]
async fn test_resolve_workspace_found() {
crate::util::test::init_management_test_stores().await;
crate::util::test::create_test_workspace("/tmp/test_resolve_ws", "test_resolve_ws").await;
let result = resolve_workspace("test_resolve_ws").await;
let resolved = result.expect("resolve should succeed for DB workspace");
assert!(resolved.is_some(), "DB workspace should be found");
assert_eq!(resolved.unwrap().name, "test_resolve_ws");
}
#[tokio::test]
async fn test_resolve_workspace_personal() {
crate::util::test::init_management_test_stores().await;
let result = resolve_workspace("personal:liliana").await;
let resolved = result.expect("resolve should succeed for personal workspace");
let ws = resolved.expect("personal workspace should be constructed on the fly");
assert_eq!(ws.name, "personal:liliana");
assert_eq!(ws.status, crate::WorkspaceStatus::Ready);
let expected_path = crate::users::personal_workspace_path("liliana");
assert_eq!(ws.path, expected_path);
}
#[tokio::test]
async fn test_resolve_workspace_not_found() {
crate::util::test::init_management_test_stores().await;
let result = resolve_workspace("nonexistent_workspace").await;
let resolved = result.expect("resolve should succeed (no error) for missing workspace");
assert!(
resolved.is_none(),
"nonexistent workspace should not be found",
);
}
const TEST_SPY_CHANNEL: &str = "__test_spy_channel";
struct TestSpyChannel {
sent: Arc<std::sync::Mutex<Vec<SendMessage>>>,
}
#[async_trait::async_trait]
impl crate::Channel for TestSpyChannel {
async fn send(&self, message: &SendMessage) -> anyhow::Result<()> {
self.sent.lock().unwrap_poison().push(message.clone());
Ok(())
}
async fn listen(
&self,
_tx: tokio::sync::mpsc::Sender<crate::ChannelMessage>,
) -> anyhow::Result<()> {
Ok(())
}
fn name(&self) -> &'static str {
TEST_SPY_CHANNEL
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
async fn setup_response_test_infra() {
crate::util::test::init_management_test_stores().await;
let _ = crate::CHANNEL_REGISTRY.set(crate::ChannelRegistry::default());
let (gui_channel, _gui_tx) = GuiChannel::new();
crate::channel_registry().register(Arc::new(gui_channel));
}
#[tokio::test]
async fn test_catch_unwind_cleanup_removes_entry_on_panic() {
let _ = init_global();
let id = "_test_panic_cleanup";
let (tx, rx) = mpsc::unbounded_channel::<AgentJob>();
{
let mut guard = ROUTER.get().unwrap().write().unwrap_poison();
guard.insert(id.to_string(), tx);
}
let agent_id = id.to_string();
let agent_id_for_cleanup = agent_id.clone();
tokio::spawn(async move {
let _result = std::panic::AssertUnwindSafe(async {
drop(rx);
panic!("simulated consumer panic");
})
.catch_unwind()
.await;
if let Some(map) = ROUTER.get() {
let mut guard = map.write().unwrap_poison();
guard.remove(&agent_id_for_cleanup);
}
});
let deadline = tokio::time::Instant::now() + tokio::time::Duration::from_secs(2);
let mut cleaned_up = false;
while tokio::time::Instant::now() < deadline {
let found = {
let guard = ROUTER.get().unwrap().read().unwrap_poison();
guard.contains_key(id)
};
if !found {
cleaned_up = true;
break;
}
tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
}
assert!(
cleaned_up,
"entry should have been removed from router table after consumer panic",
);
}
#[tokio::test]
async fn test_resolve_single_user_found() {
setup_response_test_infra().await;
let user = resolve_single_user("admin").await;
assert!(user.is_some(), "admin user should exist after store init");
assert_eq!(user.as_ref().unwrap().name, "admin");
}
#[tokio::test]
async fn test_resolve_single_user_not_found() {
setup_response_test_infra().await;
let user = resolve_single_user("nonexistent_user").await;
assert!(user.is_none(), "non-existent user should return None");
}
#[tokio::test]
async fn test_deliver_unregistered_user_response() {
setup_response_test_infra().await;
let job = AgentJob {
content: "hello from unregistered user".to_string(),
workspace_name: "default".to_string(),
user_name: "unregistered_alice".to_string(),
channel: "gui".to_string(),
kind: JobKind::UserMessage,
role: Role::Assistant,
reply_target: Some("chat_123".to_string()),
pending_job_id: None,
};
deliver_unregistered_user_response("response to unregistered user", &job, &Role::Assistant)
.await;
}
#[tokio::test]
async fn test_deliver_single_user_response() {
setup_response_test_infra().await;
let store = crate::users::USER_STORE.get().unwrap();
store
.bind_channel("admin", "gui", "admin")
.await
.expect("bind admin to gui channel");
let user = resolve_single_user("admin").await.unwrap();
let job = AgentJob {
content: "hello from registered user".to_string(),
workspace_name: "default".to_string(),
user_name: "admin".to_string(),
channel: "gui".to_string(),
kind: JobKind::UserMessage,
role: Role::Assistant,
reply_target: None,
pending_job_id: None,
};
deliver_single_user_response("response to registered user", &user, &job, &Role::Assistant)
.await;
}
#[tokio::test]
async fn test_deliver_single_user_no_bindings() {
setup_response_test_infra().await;
let user = resolve_single_user("admin").await.unwrap();
let job = AgentJob {
content: "hello".to_string(),
workspace_name: "default".to_string(),
user_name: "admin".to_string(),
channel: "gui".to_string(),
kind: JobKind::UserMessage,
role: Role::Assistant,
reply_target: None,
pending_job_id: None,
};
deliver_single_user_response(
"response to registered user without matching binding",
&user,
&job,
&Role::Assistant,
)
.await;
}
#[tokio::test]
async fn test_deliver_single_user_broadcasts_to_all_bindings() {
setup_response_test_infra().await;
let store = crate::users::USER_STORE.get().unwrap();
store
.bind_channel("admin", "gui", "admin")
.await
.expect("bind admin to gui channel");
let sent = Arc::new(std::sync::Mutex::new(Vec::new()));
crate::channel_registry().register(Arc::new(TestSpyChannel {
sent: Arc::clone(&sent),
}) as Arc<dyn crate::Channel>);
store
.bind_channel("admin", TEST_SPY_CHANNEL, "admin")
.await
.expect("bind admin to spy channel");
let user = resolve_single_user("admin").await.unwrap();
let job = AgentJob {
content: "hello".to_string(),
workspace_name: "default".to_string(),
user_name: "admin".to_string(),
channel: "gui".to_string(),
kind: JobKind::UserMessage,
role: Role::Assistant,
reply_target: None,
pending_job_id: None,
};
deliver_single_user_response("broadcast to all bindings", &user, &job, &Role::Assistant)
.await;
let captured = sent.lock().unwrap_poison();
assert!(
!captured.is_empty(),
"spy channel should have captured at least one broadcast send",
);
}
#[tokio::test]
async fn test_deliver_manager_response_with_users() {
setup_response_test_infra().await;
let store = crate::users::USER_STORE.get().unwrap();
store
.bind_channel("admin", "gui", "admin")
.await
.expect("bind admin to gui channel");
let user = resolve_single_user("admin").await.unwrap();
let job = AgentJob {
content: "manager broadcast".to_string(),
workspace_name: "default".to_string(),
user_name: "admin".to_string(),
channel: "gui".to_string(),
kind: JobKind::TicketNotify,
role: Role::Manager,
reply_target: None,
pending_job_id: None,
};
deliver_manager_response("manager response", &[user], &job).await;
}
#[tokio::test]
async fn test_register_agent_try_route_found() {
let _ = init_global();
let agent_id = "_test_register_agent_found";
let _rx = register_agent(agent_id);
let job = make_job(Role::Assistant, "hello", "user", "gui");
assert!(
try_route(agent_id, job),
"try_route should return true for a registered agent",
);
unregister_agent(agent_id);
}
#[tokio::test]
async fn test_try_route_agent_not_found() {
let _ = init_global();
let agent_id = "_test_try_route_not_found";
let job = make_job(Role::Assistant, "hello", "user", "gui");
assert!(
!try_route(agent_id, job),
"try_route should return false for an unregistered agent",
);
}
#[tokio::test]
async fn test_try_route_receiver_dropped() {
let _ = init_global();
let agent_id = "_test_try_route_dropped";
let rx = register_agent(agent_id);
drop(rx);
let job = make_job(Role::Assistant, "hello", "user", "gui");
assert!(
!try_route(agent_id, job),
"try_route should return false when receiver is dropped",
);
unregister_agent(agent_id);
}
#[tokio::test]
async fn test_unregister_agent_removes_entry() {
let _ = init_global();
let agent_id = "_test_unregister_agent_removes";
let _rx = register_agent(agent_id);
unregister_agent(agent_id);
let job = make_job(Role::Assistant, "hello", "user", "gui");
assert!(
!try_route(agent_id, job),
"try_route should return false after unregister_agent",
);
}
#[tokio::test]
async fn test_register_agent_multiple_agents() {
let _ = init_global();
let id_a = "_test_multi_a";
let id_b = "_test_multi_b";
let _rx_a = register_agent(id_a);
let _rx_b = register_agent(id_b);
assert!(try_route(id_a, make_job(Role::Assistant, "a", "u", "g")));
assert!(try_route(id_b, make_job(Role::Engineer, "b", "u", "g")));
unregister_agent(id_a);
unregister_agent(id_b);
}
#[tokio::test]
async fn test_register_agent_replaces_stale_entry() {
let _ = init_global();
let agent_id = "_test_replace_stale";
let rx = register_agent(agent_id);
drop(rx);
let _rx2 = register_agent(agent_id);
let job = make_job(Role::Assistant, "hello", "user", "gui");
assert!(
try_route(agent_id, job),
"try_route should succeed after replacing stale entry",
);
unregister_agent(agent_id);
}
#[test]
fn test_agent_failure_emoji_constant() {
assert!(!AGENT_FAILURE_EMOJI.is_empty(), "emoji should be non-empty");
assert!(
AGENT_FAILURE_EMOJI.chars().count() >= 3,
"emoji should be at least 3 characters"
);
}
#[test]
fn test_telegram_delivery_content() {
let response = "plain answer";
let content = telegram_delivery_content("telegram", Role::Manager, &[], response);
assert_eq!(content, "plain answer");
assert!(
matches!(content, Cow::Borrowed(_)),
"no-prefix deliveries must not allocate"
);
let content = telegram_delivery_content(
"telegram",
Role::Manager,
&["manager".to_string()],
response,
);
assert_eq!(content, "plain answer");
assert!(
matches!(content, Cow::Borrowed(_)),
"no-prefix deliveries must not allocate"
);
let multi = ["manager".to_string(), "artist".to_string()];
let content = telegram_delivery_content("gui", Role::Manager, &multi, response);
assert_eq!(content, "plain answer");
assert!(
matches!(content, Cow::Borrowed(_)),
"gui/voice deliveries must not allocate"
);
assert_eq!(
telegram_delivery_content("telegram", Role::Manager, &multi, response),
"🤖 Manager:\nplain answer"
);
assert_eq!(
telegram_delivery_content(
"telegram",
Role::Qa,
&["qa".to_string(), "coder".to_string()],
response,
),
"🔨 QA:\nplain answer"
);
}
#[test]
fn test_recovery_retry_kind_invariant() {
assert_ne!(
JobKind::RecoveryRetry,
JobKind::UserMessage,
"RecoveryRetry must be a distinct variant from UserMessage — \
the emoji gate's `== JobKind::UserMessage` check naturally excludes it",
);
}
}