pub mod digest;
pub mod locks;
pub mod messaging;
pub mod notes;
pub mod presence;
pub mod quota;
pub mod tasks;
pub mod attachments;
use sqlx::PgPool;
use uuid::Uuid;
use crate::{auth::AuthCtx, error::BusError, error::BusResult, model::WhoAmI};
pub fn normalize_channel(name: &str) -> String {
name.trim().trim_start_matches('#').trim().to_lowercase()
}
pub const MAX_METADATA_BYTES: usize = 16 * 1024;
pub fn normalize_metadata(metadata: Option<serde_json::Value>) -> Option<serde_json::Value> {
match metadata {
Some(serde_json::Value::String(raw)) => match serde_json::from_str(&raw) {
Ok(serde_json::Value::Object(map)) => Some(serde_json::Value::Object(map)),
_ => Some(serde_json::Value::String(raw)),
},
other => other,
}
}
pub fn check_metadata(field: &str, metadata: Option<&serde_json::Value>) -> BusResult<()> {
let Some(value) = metadata else {
return Ok(());
};
let size = serde_json::to_vec(value)
.map(|v| v.len())
.map_err(|_| BusError::invalid(format!("{field} metadata is not serialisable JSON")))?;
if size > MAX_METADATA_BYTES {
return Err(BusError::invalid(format!(
"{field} metadata is {size} bytes; the limit is {MAX_METADATA_BYTES}. \
Keep metadata to identifiers and short labels — put the payload in \
the body, a note, or an attachment."
)));
}
Ok(())
}
pub fn check_text(field: &str, value: &str, max: usize) -> BusResult<String> {
if value.len() > max {
return Err(BusError::invalid(format!(
"{field} is {} bytes; the limit is {max}",
value.len()
)));
}
let trimmed = value.trim();
if trimmed.len() > max {
return Err(BusError::invalid(format!(
"{field} is {} bytes; the limit is {max}",
trimmed.len()
)));
}
Ok(trimmed.to_owned())
}
pub fn session_label(auth: &AuthCtx) -> Option<String> {
(!auth.session.is_empty()).then(|| auth.session.clone())
}
pub async fn agent_id_by_name(pool: &PgPool, team_id: Uuid, name: &str) -> BusResult<Uuid> {
let name = name.trim();
let row: Option<(Uuid,)> =
sqlx::query_as("SELECT id FROM agents WHERE team_id = $1 AND name = $2")
.bind(team_id)
.bind(name)
.fetch_optional(pool)
.await?;
row.map(|r| r.0)
.ok_or_else(|| BusError::not_found(format!("no agent named '{name}' in this team")))
}
pub async fn channel_id_by_name(pool: &PgPool, team_id: Uuid, name: &str) -> BusResult<Uuid> {
let name = normalize_channel(name);
let row: Option<(Uuid,)> =
sqlx::query_as("SELECT id FROM channels WHERE team_id = $1 AND name = $2")
.bind(team_id)
.bind(&name)
.fetch_optional(pool)
.await?;
row.map(|r| r.0).ok_or_else(|| {
BusError::not_found(format!(
"no channel named '{name}'; call list_channels or create it with create_channel"
))
})
}
pub async fn whoami(pool: &PgPool, auth: &AuthCtx) -> BusResult<WhoAmI> {
let (unread,): (i64,) = sqlx::query_as(
r#"
SELECT count(*)
FROM messages m
LEFT JOIN read_cursors c
ON c.agent_id = $1 AND c.scope = $3
WHERE m.recipient_agent_id = $1
AND m.id > COALESCE(c.last_message_id, 0)
AND (m.recipient_session IS NULL OR m.recipient_session = $2)
"#,
)
.bind(auth.agent_id)
.bind(&auth.session)
.bind(messaging::cursor_scope("inbox", &auth.session))
.fetch_one(pool)
.await?;
let (claimed,): (i64,) = sqlx::query_as(
"SELECT count(*) FROM tasks
WHERE claimed_by = $1 AND status = 'claimed'
AND COALESCE(claimed_session, '') = $2",
)
.bind(auth.agent_id)
.bind(&auth.session)
.fetch_one(pool)
.await?;
let default_channel = messaging::default_channel(pool, auth)
.await?
.map(|(_, name)| name);
Ok(WhoAmI {
agent: auth.agent_name.clone(),
agent_id: auth.agent_id.to_string(),
team: auth.team_slug.clone(),
team_id: auth.team_id.to_string(),
session: session_label(auth),
default_channel,
unread_direct_messages: unread,
open_claimed_tasks: claimed,
})
}