use std::path::PathBuf;
use std::process::Command;
use std::sync::Arc;
use std::time::Duration;
use anyhow::{Context, Result};
use clap::Parser;
use rusqlite::{params, Connection};
use teloxide::prelude::*;
use teloxide::types::{BotCommand, ChatId, InlineKeyboardButton, InlineKeyboardMarkup, InputFile};
use tokio::sync::Mutex;
#[derive(Parser, Clone)]
#[command(name = "team-bot", version, about = "Telegram interface for teamctl")]
struct Cli {
#[arg(long, env = "TEAMCTL_MAILBOX")]
mailbox: PathBuf,
#[arg(long, env = "TEAMCTL_TELEGRAM_TOKEN")]
token: String,
#[arg(long, env = "TEAMCTL_TELEGRAM_CHATS")]
authorized_chat_ids: Option<String>,
#[arg(long, env = "TEAMCTL_MANAGER")]
manager: Option<String>,
#[arg(long, env = "TEAMCTL_TMUX_PREFIX", default_value = "t-")]
tmux_prefix: String,
}
struct State {
conn: Mutex<Connection>,
allow: Vec<i64>,
manager: Option<String>,
tmux_prefix: String,
}
impl State {
fn manager_project(&self) -> Option<&str> {
self.manager
.as_deref()
.and_then(|m| m.split_once(':').map(|(p, _)| p))
}
}
impl State {
fn is_authorized(&self, chat: i64) -> bool {
self.allow.is_empty() || self.allow.contains(&chat)
}
}
#[tokio::main]
async fn main() -> Result<()> {
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_env("TEAM_BOT_LOG")
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")),
)
.init();
let cli = Cli::parse();
let bot = Bot::new(&cli.token);
let conn = open_mailbox(&cli.mailbox)?;
let allow: Vec<i64> = cli
.authorized_chat_ids
.as_deref()
.unwrap_or("")
.split(',')
.map(str::trim)
.filter(|s| !s.is_empty())
.filter_map(|s| s.parse().ok())
.collect();
let state = Arc::new(State {
conn: Mutex::new(conn),
allow,
manager: cli.manager,
tmux_prefix: cli.tmux_prefix,
});
let runtime = if let Some(mgr) = state.manager.as_deref() {
let c = state.conn.lock().await;
agent_runtime(&c, mgr)
} else {
None
};
let commands = commands_for_runtime(runtime.as_deref());
if !commands.is_empty() {
if let Err(e) = bot.set_my_commands(commands).await {
tracing::warn!(
"set_my_commands failed (operator gets no autocomplete; \
slash-passthrough still works manually): {e}"
);
}
}
{
let bot = bot.clone();
let state = state.clone();
tokio::spawn(async move { outbound_loop(bot, state).await });
}
let bot_inbound = bot.clone();
let handler = dptree::entry()
.branch(Update::filter_message().endpoint({
let state = state.clone();
move |bot: Bot, msg: Message| {
let state = state.clone();
async move { handle_message(bot, msg, state).await }
}
}))
.branch(Update::filter_callback_query().endpoint({
let state = state.clone();
move |bot: Bot, q: CallbackQuery| {
let state = state.clone();
async move { handle_callback(bot, q, state).await }
}
}));
Dispatcher::builder(bot_inbound, handler)
.enable_ctrlc_handler()
.build()
.dispatch()
.await;
Ok(())
}
fn open_mailbox(path: &std::path::Path) -> Result<Connection> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).ok();
}
let conn = Connection::open(path).context("open mailbox")?;
conn.busy_timeout(Duration::from_secs(5))?;
conn.pragma_update(None, "journal_mode", "WAL")?;
team_core::mailbox::ensure(&conn)?;
Ok(conn)
}
async fn handle_message(bot: Bot, msg: Message, state: Arc<State>) -> ResponseResult<()> {
let chat_id = msg.chat.id.0;
let trimmed = msg.text().map(str::trim).unwrap_or("");
if !state.allow.contains(&chat_id) && trimmed == "/start" {
bot.send_message(
msg.chat.id,
format!(
"This chat isn't authorized yet.\n\n\
Your chat id: {chat_id}\n\n\
Add it to .env next to your team-compose.yaml:\n\
TEAMCTL_TELEGRAM_CHATS={chat_id}\n\n\
Then restart team-bot."
),
)
.await?;
return Ok(());
}
if !state.is_authorized(chat_id) {
return Ok(());
}
if let Some(rest) = trimmed.strip_prefix("/dm ") {
if let Some((target, body)) = rest.split_once(' ') {
if let Some((project, _)) = target.split_once(':') {
let c = state.conn.lock().await;
let _ = c.execute(
"INSERT INTO messages (project_id, sender, recipient, text, sent_at)
VALUES (?1, 'user:telegram', ?2, ?3, strftime('%s','now'))",
params![project, target, body],
);
drop(c);
bot.send_message(msg.chat.id, format!("→ {target}")).await?;
}
}
} else if !trimmed.is_empty() && !trimmed.starts_with('/') && state.manager.is_some() {
let target = state.manager.as_deref().unwrap();
if let Some((project, _)) = target.split_once(':') {
let c = state.conn.lock().await;
let _ = c.execute(
"INSERT INTO messages (project_id, sender, recipient, text, sent_at)
VALUES (?1, 'user:telegram', ?2, ?3, strftime('%s','now'))",
params![project, target, trimmed],
);
drop(c);
bot.send_message(msg.chat.id, format!("→ {target}")).await?;
}
} else if trimmed == "/pending" {
let c = state.conn.lock().await;
let rows: Vec<(i64, String, String, String)> = {
let mut stmt = c
.prepare(
"SELECT id, agent_id, action, summary FROM approvals WHERE status='pending' ORDER BY id",
)
.unwrap();
stmt.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)))
.unwrap()
.flatten()
.collect()
};
drop(c);
if rows.is_empty() {
bot.send_message(msg.chat.id, "No pending approvals.")
.await?;
} else {
let mut out = String::from("Pending approvals:\n");
for (id, agent, action, summary) in rows {
out.push_str(&format!(
"#{id} {agent} · {action}: {}\n",
render_plain(&summary)
));
}
bot.send_message(msg.chat.id, out).await?;
}
} else if trimmed == "/start" || trimmed == "/help" {
let body = match state.manager.as_deref() {
Some(mgr) => format!(
"teamctl bot — connected to {mgr}\n\
Just type a message and it goes straight to {mgr}.\n\
/pending — show pending approvals\n\
/dm <project>:<agent> <text> — send to a different agent (rare)\n\
/<cmd> — slash-passthrough to {mgr}'s tmux session (Claude Code only)"
),
None => "teamctl — Telegram interface\n\
/dm <project>:<agent> <message> — send a DM\n\
/pending — show pending approvals"
.into(),
};
bot.send_message(msg.chat.id, body).await?;
} else if trimmed.starts_with('/') && state.manager.is_some() {
let manager = state.manager.as_deref().unwrap();
let runtime_opt = {
let c = state.conn.lock().await;
agent_runtime(&c, manager)
};
let Some(runtime) = runtime_opt else {
bot.send_message(
msg.chat.id,
format!("unknown manager `{manager}` — slash-passthrough aborted"),
)
.await?;
return Ok(());
};
match slash_outcome(manager, &runtime, &state.tmux_prefix) {
SlashOutcome::Passthrough { session } => match tmux_send_keys(&session, trimmed) {
Ok(()) => {
bot.send_message(msg.chat.id, format!("→ {manager}"))
.await?;
}
Err(err) => {
bot.send_message(msg.chat.id, format!("tmux error: {err}"))
.await?;
}
},
SlashOutcome::Reject { reason } => {
bot.send_message(msg.chat.id, reason).await?;
}
}
}
Ok(())
}
async fn handle_callback(bot: Bot, q: CallbackQuery, state: Arc<State>) -> ResponseResult<()> {
let chat_id = q.message.as_ref().map(|m| m.chat().id.0).unwrap_or(0);
if !state.is_authorized(chat_id) {
return Ok(());
}
let Some(data) = q.data.clone() else {
return Ok(());
};
let Some((verb, id_str)) = data.split_once(':') else {
return Ok(());
};
let Ok(id) = id_str.parse::<i64>() else {
return Ok(());
};
let approved = verb == "approve";
let decided_now = {
let c = state.conn.lock().await;
let n = c
.execute(
"UPDATE approvals SET status=?1, decided_at=strftime('%s','now'), decided_by='user:telegram'
WHERE id=?2 AND status='pending'",
params![if approved { "approved" } else { "denied" }, id],
)
.map(|n| n > 0)
.unwrap_or(false);
if n {
let _ = c.execute(
"UPDATE approvals SET delivered_at=strftime('%s','now')
WHERE id=?1 AND delivered_at IS NULL",
params![id],
);
}
n
};
if !decided_now {
bot.answer_callback_query(q.id)
.text(format!("#{id} already resolved"))
.await?;
return Ok(());
}
if let Some(msg) = q.message.as_ref() {
let chat = msg.chat().id;
let mid = msg.id();
let original = msg.regular_message().and_then(|m| m.text()).unwrap_or("");
let outcome = if approved {
"✅ Approved by Alireza"
} else {
"❌ Rejected by Alireza"
};
let new_text = if original.is_empty() {
outcome.to_string()
} else {
format!("{original}\n\n{outcome}")
};
let _ = bot.edit_message_text(chat, mid, new_text).await;
let _ = bot
.edit_message_reply_markup(chat, mid)
.reply_markup(InlineKeyboardMarkup::new(Vec::<Vec<_>>::new()))
.await;
}
bot.answer_callback_query(q.id)
.text(format!("{} #{id}", if approved { "✅" } else { "❌" }))
.await?;
Ok(())
}
async fn outbound_loop(bot: Bot, state: Arc<State>) {
let Some(&primary) = state.allow.first() else {
tracing::warn!("no authorized_chat_ids — outbound disabled");
return;
};
let chat = ChatId(primary);
let mut last_approval_id: i64 = current_max(&state, "approvals").await;
let mut last_msg_id: i64 = current_max(&state, "messages").await;
loop {
tokio::time::sleep(Duration::from_millis(500)).await;
let approvals: Vec<(i64, String, String, String)> = {
let c = state.conn.lock().await;
let rows: Vec<(i64, String, String, String)> = match state.manager_project() {
Some(project) => {
let mut stmt = c
.prepare(
"SELECT id, agent_id, action, summary FROM approvals
WHERE status='pending' AND id > ?1 AND project_id = ?2
ORDER BY id",
)
.unwrap();
stmt.query_map(params![last_approval_id, project], |r| {
Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?))
})
.unwrap()
.flatten()
.collect()
}
None => {
let mut stmt = c
.prepare(
"SELECT id, agent_id, action, summary FROM approvals
WHERE status='pending' AND id > ?1 ORDER BY id",
)
.unwrap();
stmt.query_map(params![last_approval_id], |r| {
Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?))
})
.unwrap()
.flatten()
.collect()
}
};
rows
};
for (id, agent, action, summary) in approvals {
last_approval_id = last_approval_id.max(id);
let route_ok = {
let c = state.conn.lock().await;
should_route(state.manager.as_deref(), &agent, &c)
};
if !route_ok {
continue;
}
let kb = InlineKeyboardMarkup::new(vec![vec![
InlineKeyboardButton::callback("Approve", format!("approve:{id}")),
InlineKeyboardButton::callback("Deny", format!("deny:{id}")),
]]);
let text = format!(
"🔐 #{id} {agent}\naction: {action}\n{}",
render_plain(&summary)
);
let send_ok = bot.send_message(chat, text).reply_markup(kb).await.is_ok();
if send_ok {
let c = state.conn.lock().await;
let _ = c.execute(
"UPDATE approvals SET delivered_at=strftime('%s','now')
WHERE id=?1 AND delivered_at IS NULL",
params![id],
);
}
}
let forwardable: Vec<MailboxRow> = {
let c = state.conn.lock().await;
let rows: Vec<MailboxRow> = match state.manager_project() {
Some(project) => {
let mut stmt = c
.prepare(
"SELECT m.id, m.sender, m.text, m.kind, m.structured_payload FROM messages m
WHERE m.id > ?1
AND m.recipient = 'user:telegram'
AND m.acked_at IS NULL
AND m.project_id = ?2
ORDER BY m.id",
)
.unwrap();
stmt.query_map(params![last_msg_id, project], MailboxRow::from_row)
.unwrap()
.flatten()
.collect()
}
None => {
let mut stmt = c
.prepare(
"SELECT m.id, m.sender, m.text, m.kind, m.structured_payload FROM messages m
WHERE m.id > ?1
AND m.recipient = 'user:telegram'
AND m.acked_at IS NULL
ORDER BY m.id",
)
.unwrap();
stmt.query_map(params![last_msg_id], MailboxRow::from_row)
.unwrap()
.flatten()
.collect()
}
};
rows
};
for row in forwardable {
last_msg_id = last_msg_id.max(row.id);
let route_ok = {
let c = state.conn.lock().await;
should_route(state.manager.as_deref(), &row.sender, &c)
};
if !route_ok {
continue;
}
forward_row(&bot, chat, &row).await;
let c = state.conn.lock().await;
let _ = c.execute(
"UPDATE messages SET acked_at = strftime('%s','now') WHERE id = ?1",
params![row.id],
);
}
}
}
#[derive(Debug, Clone)]
struct MailboxRow {
id: i64,
sender: String,
text: String,
kind: Option<String>,
payload: Option<String>,
}
impl MailboxRow {
fn from_row(r: &rusqlite::Row<'_>) -> rusqlite::Result<Self> {
Ok(Self {
id: r.get(0)?,
sender: r.get(1)?,
text: r.get(2)?,
kind: r.get(3)?,
payload: r.get(4)?,
})
}
}
struct MediaPayload {
source: String,
value: String,
caption: Option<String>,
}
fn parse_payload(payload: &str) -> Option<MediaPayload> {
let v: serde_json::Value = serde_json::from_str(payload).ok()?;
let source = v.get("source")?.as_str()?.to_string();
let value = v.get("value")?.as_str()?.to_string();
let caption = v
.get("caption")
.and_then(|c| c.as_str())
.map(|s| s.to_string());
Some(MediaPayload {
source,
value,
caption,
})
}
fn input_file_from(payload: &MediaPayload) -> Option<InputFile> {
match payload.source.as_str() {
"path" => Some(InputFile::file(&payload.value)),
"url" => Some(InputFile::url(payload.value.parse().ok()?)),
_ => None,
}
}
#[derive(Debug, PartialEq, Eq)]
enum DispatchKind {
Text,
Image,
File,
UnknownFallback,
}
fn classify_kind(kind: Option<&str>) -> DispatchKind {
match kind {
None | Some("text") | Some("") => DispatchKind::Text,
Some("image") => DispatchKind::Image,
Some("file") => DispatchKind::File,
_ => DispatchKind::UnknownFallback,
}
}
async fn forward_row(bot: &Bot, chat: ChatId, row: &MailboxRow) {
let kind = classify_kind(row.kind.as_deref());
let attribution = format!("\n\n— replied by {}", row.sender);
match kind {
DispatchKind::Text => {
let _ = bot
.send_message(chat, format!("{}{attribution}", render_plain(&row.text)))
.await;
}
DispatchKind::Image | DispatchKind::File => {
let Some(payload) = row.payload.as_deref().and_then(parse_payload) else {
let _ = bot
.send_message(
chat,
format!(
"{} (media payload unparseable){attribution}",
render_plain(&row.text)
),
)
.await;
return;
};
let Some(input) = input_file_from(&payload) else {
let _ = bot
.send_message(
chat,
format!(
"{} (unsupported media source `{}`){attribution}",
render_plain(&row.text),
payload.source
),
)
.await;
return;
};
let caption_text = payload
.caption
.as_deref()
.map(|c| format!("{}{attribution}", render_plain(c)))
.unwrap_or_else(|| attribution.trim_start().to_string());
let result = match kind {
DispatchKind::Image => bot
.send_photo(chat, input)
.caption(caption_text)
.await
.err(),
DispatchKind::File => bot
.send_document(chat, input)
.caption(caption_text)
.await
.err(),
_ => unreachable!(),
};
if let Some(e) = result {
tracing::warn!(
"send_{} failed for mailbox row {}: {e}",
if kind == DispatchKind::Image {
"photo"
} else {
"document"
},
row.id
);
}
}
DispatchKind::UnknownFallback => {
let _ = bot
.send_message(chat, format!("{}{attribution}", render_plain(&row.text)))
.await;
}
}
}
async fn current_max(state: &Arc<State>, table: &str) -> i64 {
let sql = format!("SELECT COALESCE(MAX(id), 0) FROM {table}");
let c = state.conn.lock().await;
c.query_row(&sql, [], |r| r.get(0)).unwrap_or(0)
}
fn manager_of(conn: &Connection, agent_id: &str) -> Option<String> {
let row: Option<(String, i64, Option<String>)> = conn
.query_row(
"SELECT project_id, is_manager, reports_to FROM agents WHERE id = ?1",
params![agent_id],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.ok();
let (project, is_manager, reports_to) = row?;
if is_manager == 1 {
return Some(agent_id.to_string());
}
let role = reports_to?;
Some(format!("{project}:{role}"))
}
fn should_route(scoped: Option<&str>, agent_id: &str, conn: &Connection) -> bool {
let Some(scoped) = scoped else {
return true;
};
let routed = manager_of(conn, agent_id).unwrap_or_else(|| agent_id.to_string());
routed == scoped
}
fn agent_runtime(conn: &Connection, agent_id: &str) -> Option<String> {
conn.query_row(
"SELECT runtime FROM agents WHERE id = ?1",
params![agent_id],
|r| r.get::<_, String>(0),
)
.ok()
}
#[derive(Debug, PartialEq, Eq)]
enum SlashOutcome {
Passthrough { session: String },
Reject { reason: String },
}
fn slash_outcome(manager: &str, runtime: &str, tmux_prefix: &str) -> SlashOutcome {
if runtime != "claude-code" {
return SlashOutcome::Reject {
reason: format!(
"slash-passthrough is only supported on Claude Code agents \
(this manager runs `{runtime}`)."
),
};
}
let (project, role) = match manager.split_once(':') {
Some((p, r)) => (p, r),
None => {
return SlashOutcome::Reject {
reason: format!("malformed manager id `{manager}` (expected `project:role`)."),
};
}
};
SlashOutcome::Passthrough {
session: format!("{tmux_prefix}{project}-{role}"),
}
}
fn tmux_send_keys_argv<'a>(session: &'a str, body: &'a str) -> [&'a str; 5] {
["send-keys", "-t", session, body, "Enter"]
}
fn tmux_send_keys(session: &str, body: &str) -> Result<(), String> {
let argv = tmux_send_keys_argv(session, body);
let output = Command::new("tmux")
.args(argv)
.output()
.map_err(|e| format!("invoke tmux: {e}"))?;
if !output.status.success() {
let stderr = String::from_utf8_lossy(&output.stderr);
let trimmed = stderr.trim();
if trimmed.is_empty() {
return Err(format!("tmux exit {}", output.status));
}
return Err(format!("tmux exit {}: {trimmed}", output.status));
}
Ok(())
}
const CC_SLASH_COMMANDS: &[(&str, &str)] = &[
("clear", "Clear conversation history"),
(
"compact",
"Compact conversation, optionally with focus instructions",
),
("cost", "Show token usage cost"),
("help", "Show available commands and shortcuts"),
("init", "Initialize a new CLAUDE.md file"),
("mcp", "Manage MCP servers"),
("model", "Set the AI model for Claude Code"),
("permissions", "View and edit permissions"),
("resume", "Resume a previous conversation"),
("review", "Review a pull request"),
("status", "Show Claude Code status"),
("vim", "Toggle between vim and default editing modes"),
];
fn commands_for_runtime(runtime: Option<&str>) -> Vec<BotCommand> {
match runtime {
Some("claude-code") => CC_SLASH_COMMANDS
.iter()
.map(|(c, d)| BotCommand::new(*c, *d))
.collect(),
_ => Vec::new(),
}
}
fn render_plain(s: &str) -> String {
let mut out = String::with_capacity(s.len());
for (idx, line) in s.lines().enumerate() {
if idx > 0 {
out.push('\n');
}
let trimmed = line.trim_start();
let leading = &line[..line.len() - trimmed.len()];
let body = if let Some(rest) = trimmed
.strip_prefix("- ")
.or_else(|| trimmed.strip_prefix("* "))
.or_else(|| trimmed.strip_prefix("+ "))
{
format!("• {rest}")
} else {
trimmed.to_string()
};
out.push_str(leading);
out.push_str(&strip_inline_markdown(&body));
}
out
}
fn strip_inline_markdown(s: &str) -> String {
let mut out = String::with_capacity(s.len());
let mut chars = s.chars().peekable();
while let Some(c) = chars.next() {
if (c == '*' || c == '_') && chars.peek() == Some(&c) {
chars.next();
continue;
}
if c == '*' || c == '_' || c == '`' {
continue;
}
out.push(c);
}
out
}
#[cfg(test)]
mod tests {
use super::*;
use rusqlite::Connection;
fn seed(conn: &Connection) {
team_core::mailbox::ensure(conn).unwrap();
conn.execute(
"INSERT OR IGNORE INTO projects (id, name) VALUES ('p','P')",
[],
)
.unwrap();
conn.execute(
"INSERT OR IGNORE INTO agents (id, project_id, role, runtime, is_manager, reports_to)
VALUES ('p:eng_lead','p','eng_lead','claude-code',1,NULL)",
[],
)
.unwrap();
conn.execute(
"INSERT OR IGNORE INTO agents (id, project_id, role, runtime, is_manager, reports_to)
VALUES ('p:dev1','p','dev1','claude-code',0,'eng_lead')",
[],
)
.unwrap();
conn.execute(
"INSERT OR IGNORE INTO agents (id, project_id, role, runtime, is_manager, reports_to)
VALUES ('p:pm','p','pm','claude-code',1,NULL)",
[],
)
.unwrap();
}
#[test]
fn manager_of_returns_self_for_a_manager() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
assert_eq!(
manager_of(&conn, "p:eng_lead").as_deref(),
Some("p:eng_lead")
);
assert_eq!(manager_of(&conn, "p:pm").as_deref(), Some("p:pm"));
}
#[test]
fn manager_of_resolves_reports_to_for_a_worker() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
assert_eq!(manager_of(&conn, "p:dev1").as_deref(), Some("p:eng_lead"));
}
#[test]
fn manager_of_returns_none_for_unknown_agent() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
assert!(manager_of(&conn, "p:ghost").is_none());
}
#[test]
fn classify_kind_treats_null_and_empty_as_text() {
assert_eq!(classify_kind(None), DispatchKind::Text);
assert_eq!(classify_kind(Some("text")), DispatchKind::Text);
assert_eq!(classify_kind(Some("")), DispatchKind::Text);
}
#[test]
fn classify_kind_routes_image_and_file() {
assert_eq!(classify_kind(Some("image")), DispatchKind::Image);
assert_eq!(classify_kind(Some("file")), DispatchKind::File);
}
#[test]
fn classify_kind_falls_back_for_unknown_kinds() {
assert_eq!(
classify_kind(Some("reaction")),
DispatchKind::UnknownFallback
);
assert_eq!(
classify_kind(Some("garbage")),
DispatchKind::UnknownFallback
);
}
#[test]
fn parse_payload_extracts_source_value_and_caption() {
let p = parse_payload(r#"{"source":"path","value":"/tmp/x.png","caption":"hi"}"#)
.expect("payload parses");
assert_eq!(p.source, "path");
assert_eq!(p.value, "/tmp/x.png");
assert_eq!(p.caption.as_deref(), Some("hi"));
}
#[test]
fn parse_payload_handles_missing_caption() {
let p = parse_payload(r#"{"source":"url","value":"https://x.test/a.png"}"#)
.expect("payload parses");
assert_eq!(p.source, "url");
assert!(p.caption.is_none());
}
#[test]
fn parse_payload_returns_none_on_garbage() {
assert!(parse_payload("not json").is_none());
assert!(
parse_payload(r#"{"value":"x"}"#).is_none(),
"missing source"
);
assert!(
parse_payload(r#"{"source":"path"}"#).is_none(),
"missing value"
);
}
#[test]
fn input_file_from_path_and_url_both_construct() {
let p = parse_payload(r#"{"source":"path","value":"/tmp/x.png"}"#).unwrap();
assert!(input_file_from(&p).is_some());
let p = parse_payload(r#"{"source":"url","value":"https://x.test/a.png"}"#).unwrap();
assert!(input_file_from(&p).is_some());
}
#[test]
fn input_file_from_unknown_source_returns_none() {
let p = MediaPayload {
source: "bytes".into(),
value: "abc".into(),
caption: None,
};
assert!(input_file_from(&p).is_none());
}
fn insert_row(
conn: &Connection,
sender: &str,
text: &str,
kind: Option<&str>,
payload: Option<&str>,
) -> i64 {
let project = sender.split_once(':').map(|(p, _)| p).unwrap_or("p");
conn.execute(
"INSERT INTO messages (project_id, sender, recipient, text, sent_at, kind, structured_payload)
VALUES (?1, ?2, 'user:telegram', ?3, strftime('%s','now'), ?4, ?5)",
params![project, sender, text, kind, payload],
)
.unwrap();
conn.last_insert_rowid()
}
#[test]
fn outbound_select_returns_kind_and_payload_for_structured_rows() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
let id = insert_row(
&conn,
"p:eng_lead",
"shot",
Some("image"),
Some(r#"{"source":"path","value":"/tmp/a.png"}"#),
);
let mut stmt = conn
.prepare(
"SELECT m.id, m.sender, m.text, m.kind, m.structured_payload FROM messages m
WHERE m.id > ?1
AND m.recipient = 'user:telegram'
AND m.acked_at IS NULL
ORDER BY m.id",
)
.unwrap();
let rows: Vec<MailboxRow> = stmt
.query_map(params![0i64], MailboxRow::from_row)
.unwrap()
.flatten()
.collect();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].id, id);
assert_eq!(rows[0].kind.as_deref(), Some("image"));
assert!(rows[0].payload.as_deref().unwrap().contains("/tmp/a.png"));
}
#[test]
fn outbound_select_returns_null_kind_for_legacy_text_rows() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
let id = insert_row(&conn, "p:eng_lead", "hello", None, None);
let mut stmt = conn
.prepare(
"SELECT m.id, m.sender, m.text, m.kind, m.structured_payload FROM messages m
WHERE m.id > ?1
AND m.recipient = 'user:telegram'
AND m.acked_at IS NULL
ORDER BY m.id",
)
.unwrap();
let rows: Vec<MailboxRow> = stmt
.query_map(params![0i64], MailboxRow::from_row)
.unwrap()
.flatten()
.collect();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].id, id);
assert!(rows[0].kind.is_none());
assert!(rows[0].payload.is_none());
assert_eq!(classify_kind(rows[0].kind.as_deref()), DispatchKind::Text);
}
#[test]
fn render_plain_strips_paired_emphasis() {
assert_eq!(render_plain("**bold** text"), "bold text");
assert_eq!(render_plain("__also bold__"), "also bold");
assert_eq!(render_plain("plain `code` here"), "plain code here");
}
#[test]
fn render_plain_strips_single_emphasis() {
assert_eq!(render_plain("*italic* text"), "italic text");
assert_eq!(render_plain("_underscored_"), "underscored");
}
#[test]
fn render_plain_translates_list_bullets() {
let input = "- one\n- two\n * nested\n+ three";
let expected = "• one\n• two\n • nested\n• three";
assert_eq!(render_plain(input), expected);
}
#[test]
fn render_plain_preserves_emoji_and_plain_prose() {
let input = "🔐 deploy\nrouting prompt to one channel — the **right** one";
let expected = "🔐 deploy\nrouting prompt to one channel — the right one";
assert_eq!(render_plain(input), expected);
}
fn decide_sql(conn: &Connection, id: i64, approved: bool) -> bool {
let status = if approved { "approved" } else { "denied" };
let n = conn
.execute(
"UPDATE approvals SET status=?1, decided_at=strftime('%s','now'), decided_by='user:telegram'
WHERE id=?2 AND status='pending'",
params![status, id],
)
.map(|n| n > 0)
.unwrap_or(false);
if n {
let _ = conn.execute(
"UPDATE approvals SET delivered_at=strftime('%s','now')
WHERE id=?1 AND delivered_at IS NULL",
params![id],
);
}
n
}
fn insert_approval(conn: &Connection, status: &str, delivered_at: Option<f64>) -> i64 {
conn.execute(
"INSERT INTO approvals (project_id, agent_id, action, summary, status,
requested_at, expires_at, delivered_at)
VALUES ('p', 'eng_lead', 'publish', 's', ?1, 0.0, 999999999.0, ?2)",
params![status, delivered_at],
)
.unwrap();
conn.last_insert_rowid()
}
#[test]
fn stale_tap_on_undeliverable_does_not_flip_delivered_at() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
let id = insert_approval(&conn, "undeliverable", None);
let decided = decide_sql(&conn, id, true);
assert!(!decided, "stale tap should report no live decision");
let (status, delivered_at): (String, Option<f64>) = conn
.query_row(
"SELECT status, delivered_at FROM approvals WHERE id = ?1",
params![id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(status, "undeliverable");
assert!(
delivered_at.is_none(),
"delivered_at must stay NULL on undeliverable row (invariant)"
);
}
#[test]
fn live_tap_on_pending_flips_status_and_delivered_at() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
let id = insert_approval(&conn, "pending", None);
let decided = decide_sql(&conn, id, true);
assert!(decided, "live tap should report decision");
let (status, delivered_at): (String, Option<f64>) = conn
.query_row(
"SELECT status, delivered_at FROM approvals WHERE id = ?1",
params![id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(status, "approved");
assert!(
delivered_at.is_some(),
"live decision implies delivery acknowledgement"
);
}
#[test]
fn unscoped_bot_routes_every_approval() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
assert!(should_route(None, "p:dev1", &conn));
assert!(should_route(None, "p:eng_lead", &conn));
assert!(should_route(None, "p:ghost", &conn));
assert!(should_route(None, "other:agent", &conn));
}
#[test]
fn scoped_bot_routes_only_its_managers_chain() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
assert!(should_route(Some("p:eng_lead"), "p:dev1", &conn));
assert!(should_route(Some("p:eng_lead"), "p:eng_lead", &conn));
assert!(!should_route(Some("p:eng_lead"), "p:pm", &conn));
}
#[test]
fn scoped_bot_with_unknown_agent_falls_back_to_self_routing() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
assert!(!should_route(Some("p:eng_lead"), "p:ghost", &conn));
}
fn insert_reply(conn: &Connection, sender: &str, text: &str) -> i64 {
let project = sender.split_once(':').map(|(p, _)| p).unwrap_or("p");
conn.execute(
"INSERT INTO messages (project_id, sender, recipient, text, sent_at)
VALUES (?1, ?2, 'user:telegram', ?3, strftime('%s','now'))",
params![project, sender, text],
)
.unwrap();
conn.last_insert_rowid()
}
#[test]
fn reply_routes_only_to_its_senders_bot() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
let pm_msg = insert_reply(&conn, "p:pm", "from pm");
let eng_msg = insert_reply(&conn, "p:eng_lead", "from eng");
let mut stmt = conn
.prepare(
"SELECT m.id, m.sender, m.text FROM messages m
WHERE m.id > 0
AND m.recipient = 'user:telegram'
AND m.acked_at IS NULL
AND m.project_id = 'p'
ORDER BY m.id",
)
.unwrap();
let rows: Vec<(i64, String, String)> = stmt
.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))
.unwrap()
.flatten()
.collect();
assert_eq!(rows.len(), 2, "both replies share the project pre-filter");
let pm_routed: Vec<i64> = rows
.iter()
.filter(|(_, sender, _)| should_route(Some("p:pm"), sender, &conn))
.map(|(id, _, _)| *id)
.collect();
assert_eq!(pm_routed, vec![pm_msg]);
let eng_routed: Vec<i64> = rows
.iter()
.filter(|(_, sender, _)| should_route(Some("p:eng_lead"), sender, &conn))
.map(|(id, _, _)| *id)
.collect();
assert_eq!(eng_routed, vec![eng_msg]);
let unscoped: Vec<i64> = rows
.iter()
.filter(|(_, sender, _)| should_route(None, sender, &conn))
.map(|(id, _, _)| *id)
.collect();
assert_eq!(unscoped, vec![pm_msg, eng_msg]);
}
#[test]
fn live_tap_keeps_existing_delivered_at_unchanged() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
let id = insert_approval(&conn, "pending", Some(1234.5));
let decided = decide_sql(&conn, id, false);
assert!(decided);
let delivered_at: f64 = conn
.query_row(
"SELECT delivered_at FROM approvals WHERE id = ?1",
params![id],
|r| r.get(0),
)
.unwrap();
assert!(
(delivered_at - 1234.5).abs() < 1e-6,
"previously-set delivered_at must not be overwritten ({delivered_at})"
);
}
#[test]
fn agent_runtime_returns_runtime_for_known_agent() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
assert_eq!(
agent_runtime(&conn, "p:eng_lead"),
Some("claude-code".into())
);
}
#[test]
fn agent_runtime_returns_runtime_when_runtime_varies() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
conn.execute(
"INSERT OR IGNORE INTO agents (id, project_id, role, runtime, is_manager, reports_to)
VALUES ('p:codex_mgr','p','codex_mgr','codex',1,NULL)",
[],
)
.unwrap();
assert_eq!(agent_runtime(&conn, "p:codex_mgr"), Some("codex".into()));
}
#[test]
fn agent_runtime_returns_none_for_unknown_agent() {
let conn = Connection::open_in_memory().unwrap();
seed(&conn);
assert_eq!(agent_runtime(&conn, "p:ghost"), None);
}
#[test]
fn slash_outcome_passes_through_for_claude_code_runtime() {
let outcome = slash_outcome("writing:manager", "claude-code", "t-");
assert_eq!(
outcome,
SlashOutcome::Passthrough {
session: "t-writing-manager".into(),
}
);
}
#[test]
fn slash_outcome_honours_custom_tmux_prefix() {
let outcome = slash_outcome("news:head_editor", "claude-code", "a-");
assert_eq!(
outcome,
SlashOutcome::Passthrough {
session: "a-news-head_editor".into(),
}
);
}
#[test]
fn slash_outcome_rejects_codex_runtime_with_named_runtime() {
let outcome = slash_outcome("writing:manager", "codex", "t-");
let SlashOutcome::Reject { reason } = outcome else {
panic!("non-CC runtime must reject");
};
assert!(
reason.contains("Claude Code"),
"rejection should reference Claude Code: {reason}"
);
assert!(
reason.contains("codex"),
"rejection should name the actual runtime: {reason}"
);
}
#[test]
fn slash_outcome_rejects_gemini_runtime_with_named_runtime() {
let outcome = slash_outcome("writing:manager", "gemini", "t-");
let SlashOutcome::Reject { reason } = outcome else {
panic!("non-CC runtime must reject");
};
assert!(reason.contains("gemini"), "names the runtime: {reason}");
}
#[test]
fn slash_outcome_rejects_malformed_manager_id() {
let outcome = slash_outcome("not-a-manager-id", "claude-code", "t-");
let SlashOutcome::Reject { reason } = outcome else {
panic!("malformed manager id must reject");
};
assert!(reason.contains("malformed"), "names the failure: {reason}");
}
#[test]
fn tmux_send_keys_argv_pins_send_keys_target_body_enter_shape() {
let argv = tmux_send_keys_argv("t-writing-manager", "/clear");
assert_eq!(
argv,
["send-keys", "-t", "t-writing-manager", "/clear", "Enter"]
);
}
#[test]
fn tmux_send_keys_argv_passes_body_verbatim_no_quote_munging() {
let argv = tmux_send_keys_argv("sess", "/compact focus on the cascade");
assert_eq!(argv[3], "/compact focus on the cascade");
assert_eq!(argv[4], "Enter");
}
#[test]
fn commands_for_runtime_returns_full_cc_list_for_claude_code() {
let cmds = commands_for_runtime(Some("claude-code"));
assert_eq!(
cmds.len(),
CC_SLASH_COMMANDS.len(),
"CC manager registers the full curated list"
);
let names: Vec<&str> = cmds.iter().map(|c| c.command.as_str()).collect();
assert!(names.contains(&"clear"), "must include /clear: {names:?}");
assert!(
names.contains(&"compact"),
"must include /compact: {names:?}"
);
assert!(names.contains(&"help"), "must include /help: {names:?}");
}
#[test]
fn commands_for_runtime_returns_empty_for_codex() {
assert!(commands_for_runtime(Some("codex")).is_empty());
}
#[test]
fn commands_for_runtime_returns_empty_for_gemini() {
assert!(commands_for_runtime(Some("gemini")).is_empty());
}
#[test]
fn commands_for_runtime_returns_empty_for_unknown_runtime() {
assert!(commands_for_runtime(Some("a-future-runtime")).is_empty());
}
#[test]
fn commands_for_runtime_returns_empty_for_unscoped_bot() {
assert!(commands_for_runtime(None).is_empty());
}
#[test]
fn cc_slash_command_names_satisfy_telegram_constraints() {
for (cmd, _desc) in CC_SLASH_COMMANDS {
assert!(
!cmd.is_empty() && cmd.len() <= 32,
"command `{cmd}` violates 1-32 char limit"
);
assert!(
cmd.chars()
.all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '_'),
"command `{cmd}` contains chars Telegram rejects (only [a-z0-9_])"
);
}
}
#[test]
fn cc_slash_command_descriptions_satisfy_telegram_constraints() {
for (cmd, desc) in CC_SLASH_COMMANDS {
assert!(
desc.len() >= 3 && desc.len() <= 256,
"description for `{cmd}` violates 3-256 char limit (got {} chars: {desc:?})",
desc.len()
);
}
}
}