use std::sync::Arc;
use ai_crew_sync::{
MIGRATOR,
auth::{generate_token, hash_token, token_prefix},
serve::{ServeOptions, build_router},
};
use rmcp::{
ServiceExt,
model::{CallToolRequestParams, ClientInfo},
service::RunningService,
transport::{
StreamableHttpClientTransport, streamable_http_client::StreamableHttpClientTransportConfig,
},
};
use serde_json::{Value, json};
use sqlx::{PgPool, postgres::PgPoolOptions};
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
fn db_url() -> Option<String> {
std::env::var("TEST_DATABASE_URL")
.or_else(|_| std::env::var("DATABASE_URL"))
.ok()
}
struct Harness {
pool: PgPool,
base: String,
ct: CancellationToken,
servers: Vec<tokio::task::JoinHandle<()>>,
schema: String,
}
impl Harness {
async fn add_replica(&mut self) -> String {
let (base, handle) = spawn_server(self.pool.clone(), self.ct.child_token()).await;
self.servers.push(handle);
base
}
async fn shutdown(mut self) {
self.ct.cancel();
let servers = std::mem::take(&mut self.servers);
for mut handle in servers {
match tokio::time::timeout(std::time::Duration::from_secs(5), &mut handle).await {
Ok(Ok(())) => {}
Ok(Err(e)) if e.is_panic() => panic!("a server task panicked: {e}"),
Ok(Err(_)) => {}
Err(_) => {
handle.abort();
panic!("a server task did not stop within 5s of cancellation");
}
}
}
let _ = sqlx::query(&format!("DROP SCHEMA IF EXISTS {} CASCADE", self.schema))
.execute(&self.pool)
.await;
self.pool.close().await;
}
}
impl Drop for Harness {
fn drop(&mut self) {
self.ct.cancel();
}
}
async fn spawn_server(
pool: PgPool,
ct: CancellationToken,
) -> (String, tokio::task::JoinHandle<()>) {
let app = build_router(
pool,
&ServeOptions {
bind: String::new(),
allowed_hosts: vec![],
allowed_origins: vec![],
max_request_bytes: ai_crew_sync::serve::DEFAULT_MAX_REQUEST_BYTES,
rate_limit_per_minute: 0,
dashboard_secret: b"test-dashboard-secret".to_vec(),
},
ct.clone(),
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let handle = tokio::spawn(async move {
let _ = axum::serve(listener, app)
.with_graceful_shutdown(async move { ct.cancelled().await })
.await;
});
(format!("http://{addr}"), handle)
}
async fn setup(schema: &str) -> Option<Harness> {
let url = db_url()?;
let pool = PgPoolOptions::new()
.max_connections(8)
.after_connect({
let schema = schema.to_owned();
move |conn, _| {
let schema = schema.clone();
Box::pin(async move {
sqlx::query(&format!("SET search_path TO {schema}"))
.execute(&mut *conn)
.await?;
Ok(())
})
}
})
.connect(&url)
.await
.expect("connect");
sqlx::query(&format!("DROP SCHEMA IF EXISTS {schema} CASCADE"))
.execute(&pool)
.await
.unwrap();
sqlx::query(&format!("CREATE SCHEMA {schema}"))
.execute(&pool)
.await
.unwrap();
MIGRATOR.run(&pool).await.expect("migrate");
let ct = CancellationToken::new();
let (base, handle) = spawn_server(pool.clone(), ct.child_token()).await;
Some(Harness {
pool,
base,
ct,
servers: vec![handle],
schema: schema.to_owned(),
})
}
async fn seed_agent(pool: &PgPool, team: &str, agent: &str) -> String {
let team_id: (Uuid,) = sqlx::query_as(
"INSERT INTO teams (slug, name) VALUES ($1, $1)
ON CONFLICT (slug) DO UPDATE SET name = EXCLUDED.name RETURNING id",
)
.bind(team)
.fetch_one(pool)
.await
.unwrap();
let agent_id: (Uuid,) =
sqlx::query_as("INSERT INTO agents (team_id, name) VALUES ($1, $2) RETURNING id")
.bind(team_id.0)
.bind(agent)
.fetch_one(pool)
.await
.unwrap();
let raw = generate_token();
sqlx::query("INSERT INTO api_tokens (agent_id, token_hash, prefix) VALUES ($1, $2, $3)")
.bind(agent_id.0)
.bind(hash_token(&raw))
.bind(token_prefix(&raw))
.execute(pool)
.await
.unwrap();
raw
}
type Client = RunningService<rmcp::RoleClient, ClientInfo>;
async fn connect(base: &str, token: &str) -> Client {
let mut config = StreamableHttpClientTransportConfig::with_uri(format!("{base}/mcp"));
config.auth_header = Some(token.to_string());
config.allow_stateless = true;
let transport = StreamableHttpClientTransport::from_config(config);
ClientInfo::default()
.serve(transport)
.await
.expect("mcp handshake")
}
async fn call(client: &Client, name: &str, args: Value) -> Value {
let args: serde_json::Map<String, Value> = serde_json::from_value(args).unwrap();
let result = client
.call_tool(CallToolRequestParams::new(name.to_string()).with_arguments(args))
.await
.unwrap_or_else(|e| panic!("{name} failed: {e}"));
assert_ne!(
result.is_error,
Some(true),
"{name} returned an error: {result:?}"
);
result
.structured_content
.clone()
.unwrap_or_else(|| panic!("{name} returned no structured content: {result:?}"))
}
async fn call_expect_error(client: &Client, name: &str, args: Value) -> String {
let args: serde_json::Map<String, Value> = serde_json::from_value(args).unwrap();
match client
.call_tool(CallToolRequestParams::new(name.to_string()).with_arguments(args))
.await
{
Err(e) => e.to_string(),
Ok(result) => {
assert_eq!(
result.is_error,
Some(true),
"{name} unexpectedly succeeded: {result:?}"
);
format!("{:?}", result.content)
}
}
}
async fn setup_rate_limited(schema: &str, per_minute: u32) -> Option<Harness> {
let url = match db_url() {
Some(url) => url,
None => {
assert!(
!db_required(),
"AI_CREW_SYNC_REQUIRE_DB is set but TEST_DATABASE_URL is not: \
this test would have silently passed without a database"
);
return None;
}
};
let pool = PgPoolOptions::new()
.max_connections(8)
.after_connect({
let schema = schema.to_owned();
move |conn, _| {
let schema = schema.clone();
Box::pin(async move {
sqlx::query(&format!("SET search_path TO {schema}"))
.execute(&mut *conn)
.await?;
Ok(())
})
}
})
.connect(&url)
.await
.expect("connect");
sqlx::query(&format!("DROP SCHEMA IF EXISTS {schema} CASCADE"))
.execute(&pool)
.await
.unwrap();
sqlx::query(&format!("CREATE SCHEMA {schema}"))
.execute(&pool)
.await
.unwrap();
MIGRATOR.run(&pool).await.expect("migrate");
let ct = CancellationToken::new();
let child = ct.child_token();
let app = build_router(
pool.clone(),
&ServeOptions {
bind: String::new(),
allowed_hosts: vec![],
allowed_origins: vec![],
max_request_bytes: 64 * 1024,
rate_limit_per_minute: per_minute,
dashboard_secret: b"test-dashboard-secret".to_vec(),
},
child.clone(),
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let handle = tokio::spawn(async move {
let _ = axum::serve(listener, app)
.with_graceful_shutdown(async move { child.cancelled().await })
.await;
});
Some(Harness {
pool,
base: format!("http://{addr}"),
ct,
servers: vec![handle],
schema: schema.to_owned(),
})
}
fn db_required() -> bool {
std::env::var("AI_CREW_SYNC_REQUIRE_DB").is_ok_and(|v| v != "0")
}
macro_rules! require_db {
($schema:expr) => {
match setup($schema).await {
Some(h) => h,
None => {
assert!(
!db_required(),
"AI_CREW_SYNC_REQUIRE_DB is set but TEST_DATABASE_URL is not: \
the integration suite would have silently passed without \
touching a database"
);
eprintln!("skipping: TEST_DATABASE_URL not set");
return;
}
}
};
}
#[tokio::test]
async fn unauthenticated_requests_are_rejected() {
let h = require_db!("t_auth");
let client = reqwest::Client::new();
let resp = client
.post(format!("{}/mcp", h.base))
.header("Accept", "application/json, text/event-stream")
.json(&json!({"jsonrpc":"2.0","id":1,"method":"tools/list"}))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 401, "no token must be rejected");
let resp = client
.post(format!("{}/mcp", h.base))
.header("Authorization", "Bearer acs_deadbeef")
.header("Accept", "application/json, text/event-stream")
.json(&json!({"jsonrpc":"2.0","id":1,"method":"tools/list"}))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 401, "bogus token must be rejected");
let resp = client
.get(format!("{}/health", h.base))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 200);
}
#[tokio::test]
async fn revoked_token_stops_working() {
let h = require_db!("t_revoke");
let token = seed_agent(&h.pool, "acme", "joaquin").await;
let client = connect(&h.base, &token).await;
call(&client, "whoami", json!({})).await;
let _ = client.cancel().await;
sqlx::query("UPDATE api_tokens SET revoked_at = now()")
.execute(&h.pool)
.await
.unwrap();
let resp = reqwest::Client::new()
.post(format!("{}/mcp", h.base))
.header("Authorization", format!("Bearer {token}"))
.header("Accept", "application/json, text/event-stream")
.json(&json!({"jsonrpc":"2.0","id":1,"method":"tools/list"}))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 401);
}
#[tokio::test]
async fn tools_are_advertised_with_schemas() {
let h = require_db!("t_tools");
let token = seed_agent(&h.pool, "acme", "joaquin").await;
let client = connect(&h.base, &token).await;
let tools = client.list_all_tools().await.unwrap();
let names: Vec<&str> = tools.iter().map(|t| t.name.as_ref()).collect();
for expected in [
"whoami",
"post_message",
"read_messages",
"list_channels",
"create_channel",
"search_messages",
"ask_agent",
"attach_file",
"get_attachment",
"create_task",
"claim_task",
"claim_next_task",
"complete_task",
"release_task",
"renew_task_lease",
"list_tasks",
"get_task",
"heartbeat",
"list_agents",
"set_note",
"get_note",
"list_notes",
"search_notes",
"delete_note",
] {
assert!(
names.contains(&expected),
"missing tool {expected} in {names:?}"
);
}
for tool in &tools {
assert!(
tool.description.as_ref().is_some_and(|d| d.len() > 20),
"tool {} needs a usable description",
tool.name
);
}
let _ = client.cancel().await;
}
#[tokio::test]
async fn direct_messages_and_channels_flow_between_two_agents() {
let h = require_db!("t_msg");
let joaquin_token = seed_agent(&h.pool, "acme", "joaquin").await;
let marta_token = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &joaquin_token).await;
let marta = connect(&h.base, &marta_token).await;
let me = call(&joaquin, "whoami", json!({})).await;
assert_eq!(me["agent"], "joaquin");
assert_eq!(me["team"], "acme");
call(
&joaquin,
"post_message",
json!({"to": "marta", "body": "the auth refactor touches your billing module"}),
)
.await;
let inbox = call(&marta, "read_messages", json!({"scope": "inbox"})).await;
assert_eq!(inbox["messages"].as_array().unwrap().len(), 1);
assert_eq!(inbox["messages"][0]["from"], "joaquin");
assert_eq!(inbox["messages"][0]["to"], "marta");
let again = call(&marta, "read_messages", json!({"scope": "inbox"})).await;
assert_eq!(again["messages"].as_array().unwrap().len(), 0);
let history = call(
&marta,
"read_messages",
json!({"scope": "inbox", "only_new": false}),
)
.await;
assert_eq!(history["messages"].as_array().unwrap().len(), 1);
let joaquin_inbox = call(&joaquin, "read_messages", json!({"scope": "inbox"})).await;
assert_eq!(joaquin_inbox["messages"].as_array().unwrap().len(), 0);
call(
&marta,
"create_channel",
json!({"name": "#Deploys", "topic": "what is going out"}),
)
.await;
let channels = call(&joaquin, "list_channels", json!({})).await;
assert_eq!(
channels["channels"][0]["name"], "deploys",
"name normalised"
);
call(
&marta,
"post_message",
json!({"channel": "deploys", "body": "staging is on 1.4.2"}),
)
.await;
let read = call(&joaquin, "read_messages", json!({"scope": "deploys"})).await;
assert_eq!(read["messages"][0]["body"], "staging is on 1.4.2");
assert_eq!(read["messages"][0]["from"], "marta");
let found = call(&joaquin, "search_messages", json!({"query": "staging"})).await;
assert_eq!(found["messages"].as_array().unwrap().len(), 1);
let err = call_expect_error(
&joaquin,
"post_message",
json!({"to": "nobody", "body": "hi"}),
)
.await;
assert!(
err.contains("nobody"),
"error should name the missing agent: {err}"
);
let err = call_expect_error(&joaquin, "post_message", json!({"body": "hi"})).await;
assert!(err.to_lowercase().contains("channel"), "got: {err}");
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
}
#[tokio::test]
async fn teams_are_isolated_from_each_other() {
let h = require_db!("t_isolation");
let acme = seed_agent(&h.pool, "acme", "joaquin").await;
let other = seed_agent(&h.pool, "globex", "intruder").await;
let acme_client = connect(&h.base, &acme).await;
let other_client = connect(&h.base, &other).await;
call(&acme_client, "create_channel", json!({"name": "secrets"})).await;
call(
&acme_client,
"post_message",
json!({"channel": "secrets", "body": "the api key is in vault"}),
)
.await;
call(
&acme_client,
"set_note",
json!({"scope": "api", "key": "vault-path", "value": "secret/prod/api"}),
)
.await;
call(
&acme_client,
"create_task",
json!({"key": "rotate-keys", "title": "rotate the prod keys"}),
)
.await;
let channels = call(&other_client, "list_channels", json!({})).await;
assert_eq!(channels["channels"].as_array().unwrap().len(), 0);
let msgs = call(&other_client, "read_messages", json!({"scope": "all"})).await;
assert_eq!(msgs["messages"].as_array().unwrap().len(), 0);
let notes = call(&other_client, "list_notes", json!({})).await;
assert_eq!(notes["notes"].as_array().unwrap().len(), 0);
let tasks = call(&other_client, "list_tasks", json!({})).await;
assert_eq!(tasks["tasks"].as_array().unwrap().len(), 0);
let found = call(&other_client, "search_messages", json!({"query": "vault"})).await;
assert_eq!(found["messages"].as_array().unwrap().len(), 0);
let err = call_expect_error(&other_client, "get_task", json!({"key": "rotate-keys"})).await;
assert!(err.contains("not found"), "got: {err}");
let agents = call(&other_client, "list_agents", json!({})).await;
let names: Vec<&str> = agents["agents"]
.as_array()
.unwrap()
.iter()
.map(|a| a["name"].as_str().unwrap())
.collect();
assert_eq!(names, vec!["intruder"]);
let _ = acme_client.cancel().await;
let _ = other_client.cancel().await;
}
#[tokio::test]
async fn a_claimed_task_cannot_be_claimed_by_someone_else() {
let h = require_db!("t_claim");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &a).await;
let marta = connect(&h.base, &b).await;
call(
&joaquin,
"create_task",
json!({"key": "refactor-auth", "title": "rewrite the token refresh flow"}),
)
.await;
let claim = call(
&joaquin,
"claim_task",
json!({"key": "refactor-auth", "lease_seconds": 600}),
)
.await;
assert_eq!(claim["claimed"], true);
assert_eq!(claim["task"]["claimed_by"], "joaquin");
let denied = call(&marta, "claim_task", json!({"key": "refactor-auth"})).await;
assert_eq!(denied["claimed"], false);
assert!(
denied["reason"].as_str().unwrap().contains("joaquin"),
"reason should name the holder: {denied:?}"
);
let again = call(&joaquin, "claim_task", json!({"key": "refactor-auth"})).await;
assert_eq!(again["claimed"], true);
let err = call_expect_error(&marta, "release_task", json!({"key": "refactor-auth"})).await;
assert!(err.contains("do not hold"), "got: {err}");
let err = call_expect_error(&marta, "renew_task_lease", json!({"key": "refactor-auth"})).await;
assert!(err.contains("do not hold"), "got: {err}");
call(&joaquin, "release_task", json!({"key": "refactor-auth"})).await;
let retry = call(&marta, "claim_task", json!({"key": "refactor-auth"})).await;
assert_eq!(retry["claimed"], true);
assert_eq!(retry["task"]["claimed_by"], "marta");
let done = call(
&marta,
"complete_task",
json!({"key": "refactor-auth", "result": "merged in #421"}),
)
.await;
assert_eq!(done["status"], "done");
assert_eq!(done["result"], "merged in #421");
let err = call_expect_error(&joaquin, "complete_task", json!({"key": "refactor-auth"})).await;
assert!(err.contains("already done"), "got: {err}");
let detail = call(&joaquin, "get_task", json!({"key": "refactor-auth"})).await;
let events: Vec<&str> = detail["history"]
.as_array()
.unwrap()
.iter()
.map(|e| e["event"].as_str().unwrap())
.collect();
assert_eq!(
events,
vec![
"created",
"claimed",
"claimed",
"released",
"claimed",
"completed"
]
);
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
}
#[tokio::test]
async fn an_expired_lease_can_be_taken_over() {
let h = require_db!("t_lease");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &a).await;
let marta = connect(&h.base, &b).await;
call(
&joaquin,
"create_task",
json!({"key": "long-job", "title": "reindex everything"}),
)
.await;
call(&joaquin, "claim_task", json!({"key": "long-job"})).await;
sqlx::query("UPDATE tasks SET lease_expires_at = now() - interval '1 minute'")
.execute(&h.pool)
.await
.unwrap();
let listed = call(&marta, "list_tasks", json!({})).await;
assert_eq!(listed["tasks"][0]["lease_expired"], true);
let stolen = call(&marta, "claim_task", json!({"key": "long-job"})).await;
assert_eq!(
stolen["claimed"], true,
"an expired lease must be reclaimable"
);
assert_eq!(stolen["task"]["claimed_by"], "marta");
let renewed = call(
&marta,
"renew_task_lease",
json!({"key": "long-job", "lease_seconds": 3600}),
)
.await;
assert_eq!(renewed["lease_expired"], false);
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
}
#[tokio::test]
async fn concurrent_claim_next_never_hands_out_the_same_task_twice() {
let h = require_db!("t_race");
let mut clients = Vec::new();
for i in 0..4 {
let token = seed_agent(&h.pool, "acme", &format!("agent{i}")).await;
clients.push(Arc::new(connect(&h.base, &token).await));
}
for i in 0..4 {
call(
&clients[0],
"create_task",
json!({"key": format!("job-{i}"), "title": format!("job {i}")}),
)
.await;
}
let mut handles = Vec::new();
for client in &clients {
let client = Arc::clone(client);
handles.push(tokio::spawn(async move {
call(&client, "claim_next_task", json!({})).await
}));
}
let mut keys = Vec::new();
for handle in handles {
let result = handle.await.unwrap();
assert_eq!(result["claimed"], true);
keys.push(result["task"]["key"].as_str().unwrap().to_owned());
}
keys.sort();
keys.dedup();
assert_eq!(
keys.len(),
4,
"each agent must get a distinct task: {keys:?}"
);
let empty = call(&clients[0], "claim_next_task", json!({})).await;
assert_eq!(empty["claimed"], false);
assert!(empty["reason"].as_str().unwrap().contains("no unclaimed"));
}
#[tokio::test]
async fn presence_expires_and_is_visible_to_the_team() {
let h = require_db!("t_presence");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &a).await;
let marta = connect(&h.base, &b).await;
call(
&joaquin,
"heartbeat",
json!({"repo": "acme/api", "branch": "feat/auth", "activity": "rewriting token refresh"}),
)
.await;
let seen = call(&marta, "list_agents", json!({"online_only": true})).await;
assert_eq!(seen["online_count"], 1);
assert_eq!(seen["agents"][0]["name"], "joaquin");
assert_eq!(seen["agents"][0]["activity"], "rewriting token refresh");
assert_eq!(seen["agents"][0]["repo"], "acme/api");
call(&joaquin, "heartbeat", json!({"status": "blocked"})).await;
let seen = call(&marta, "list_agents", json!({"online_only": true})).await;
assert_eq!(seen["agents"][0]["status"], "blocked");
assert_eq!(seen["agents"][0]["repo"], "acme/api", "repo should persist");
sqlx::query("UPDATE agent_presence SET expires_at = now() - interval '1 minute'")
.execute(&h.pool)
.await
.unwrap();
let seen = call(&marta, "list_agents", json!({})).await;
assert_eq!(seen["online_count"], 0);
let joaquin_row = seen["agents"]
.as_array()
.unwrap()
.iter()
.find(|a| a["name"] == "joaquin")
.unwrap();
assert_eq!(joaquin_row["status"], "offline");
let err = call_expect_error(&joaquin, "heartbeat", json!({"status": "vibing"})).await;
assert!(err.contains("active"), "got: {err}");
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
}
#[tokio::test]
async fn notes_are_shared_memory_with_history() {
let h = require_db!("t_notes");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &a).await;
let marta = connect(&h.base, &b).await;
call(
&joaquin,
"set_note",
json!({
"scope": "api",
"key": "why-no-redis",
"value": "we dropped redis in march; the cache lives in postgres now",
"tags": ["Infra", "decision"]
}),
)
.await;
let note = call(
&marta,
"get_note",
json!({"scope": "api", "key": "why-no-redis"}),
)
.await;
assert_eq!(note["found"], true);
assert_eq!(note["note"]["updated_by"], "joaquin");
assert_eq!(note["note"]["tags"][0], "infra");
let missing = call(&marta, "get_note", json!({"key": "nope"})).await;
assert_eq!(missing["found"], false);
assert!(missing["note"].is_null());
call(
&marta,
"set_note",
json!({"scope": "api", "key": "why-no-redis", "value": "correction: valkey, not redis"}),
)
.await;
let note = call(
&joaquin,
"get_note",
json!({"scope": "api", "key": "why-no-redis"}),
)
.await;
assert_eq!(note["note"]["updated_by"], "marta");
let (revisions,): (i64,) = sqlx::query_as("SELECT count(*) FROM note_revisions")
.fetch_one(&h.pool)
.await
.unwrap();
assert_eq!(revisions, 2, "both versions retained");
call(
&joaquin,
"set_note",
json!({"scope": "web", "key": "build", "value": "vite, not webpack", "tags": ["infra"]}),
)
.await;
let api_only = call(&marta, "list_notes", json!({"scope": "api"})).await;
assert_eq!(api_only["notes"].as_array().unwrap().len(), 1);
let all = call(&marta, "list_notes", json!({})).await;
assert_eq!(all["notes"].as_array().unwrap().len(), 2);
let tagged = call(&marta, "list_notes", json!({"tag": "infra"})).await;
assert_eq!(tagged["notes"].as_array().unwrap().len(), 1);
let found = call(&marta, "search_notes", json!({"query": "valkey"})).await;
assert_eq!(found["notes"].as_array().unwrap().len(), 1);
let del = call(
&marta,
"delete_note",
json!({"scope": "api", "key": "why-no-redis"}),
)
.await;
assert_eq!(del["ok"], true);
let del_again = call(
&marta,
"delete_note",
json!({"scope": "api", "key": "why-no-redis"}),
)
.await;
assert_eq!(del_again["ok"], false, "deleting twice is not an error");
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
}
#[tokio::test]
async fn whoami_reports_pending_work() {
let h = require_db!("t_whoami");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &a).await;
let marta = connect(&h.base, &b).await;
call(
&marta,
"post_message",
json!({"to": "joaquin", "body": "ping"}),
)
.await;
call(
&marta,
"post_message",
json!({"to": "joaquin", "body": "ping again"}),
)
.await;
call(
&joaquin,
"create_task",
json!({"key": "t1", "title": "something"}),
)
.await;
call(&joaquin, "claim_task", json!({"key": "t1"})).await;
let me = call(&joaquin, "whoami", json!({})).await;
assert_eq!(me["unread_direct_messages"], 2);
assert_eq!(me["open_claimed_tasks"], 1);
call(&joaquin, "read_messages", json!({"scope": "inbox"})).await;
let me = call(&joaquin, "whoami", json!({})).await;
assert_eq!(me["unread_direct_messages"], 0);
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
}
#[tokio::test]
async fn wait_for_updates_wakes_on_a_teammates_message() {
let h = require_db!("t_wait");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &a).await;
let marta = connect(&h.base, &b).await;
call(&joaquin, "create_channel", json!({"name": "dev"})).await;
let waiter = {
let base = h.base.clone();
let token = a.clone();
tokio::spawn(async move {
let client = connect(&base, &token).await;
let started = std::time::Instant::now();
let result = call(&client, "wait_for_updates", json!({"timeout_seconds": 20})).await;
let _ = client.cancel().await;
(result, started.elapsed())
})
};
tokio::time::sleep(std::time::Duration::from_millis(700)).await;
call(
&marta,
"post_message",
json!({"channel": "dev", "body": "he subido el fix del parser"}),
)
.await;
let (result, elapsed) = waiter.await.unwrap();
assert_eq!(result["woke"], true, "must wake on the message: {result:?}");
assert_eq!(result["timed_out"], false);
assert!(
elapsed < std::time::Duration::from_secs(10),
"woke by event, not by timeout (took {elapsed:?})"
);
let summaries = result["events"].as_array().unwrap();
assert!(
summaries
.iter()
.any(|e| e["summary"].as_str().unwrap().contains("marta")),
"event should name the sender: {summaries:?}"
);
let instant = call(&joaquin, "wait_for_updates", json!({"timeout_seconds": 30})).await;
assert_eq!(instant["woke"], true);
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
}
#[tokio::test]
async fn ask_agent_returns_the_teammates_answer() {
let h = require_db!("t_ask");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &a).await;
let marta = connect(&h.base, &b).await;
let asker = {
let base = h.base.clone();
let token = a.clone();
tokio::spawn(async move {
let client = connect(&base, &token).await;
let started = std::time::Instant::now();
let result = call(
&client,
"ask_agent",
json!({"to": "marta", "question": "does staging run pg16?", "timeout_seconds": 20}),
)
.await;
let _ = client.cancel().await;
(result, started.elapsed())
})
};
tokio::time::sleep(std::time::Duration::from_millis(700)).await;
let inbox = call(&marta, "read_messages", json!({"scope": "inbox"})).await;
let question = inbox["messages"]
.as_array()
.unwrap()
.last()
.unwrap()
.clone();
assert_eq!(
question["metadata"]["question"], true,
"the question DM is marked as such: {question:?}"
);
call(
&marta,
"post_message",
json!({"to": "joaquin", "body": "yes, since yesterday", "reply_to": question["id"]}),
)
.await;
let (result, elapsed) = asker.await.unwrap();
assert_eq!(result["answered"], true, "{result:?}");
assert_eq!(result["answer"]["from"], "marta");
assert_eq!(result["answer"]["body"], "yes, since yesterday");
assert!(
elapsed < std::time::Duration::from_secs(10),
"answered by event, not by timeout (took {elapsed:?})"
);
let timed = call(
&joaquin,
"ask_agent",
json!({"to": "marta", "question": "and prod?", "timeout_seconds": 5}),
)
.await;
assert_eq!(timed["answered"], false, "{timed:?}");
let qid = timed["question_message_id"].as_i64().unwrap();
assert!(
timed["suggestion"]
.as_str()
.unwrap()
.contains(&qid.to_string()),
"timeout suggestion tells how to resume: {timed:?}"
);
call(
&marta,
"post_message",
json!({"to": "joaquin", "body": "prod is still on pg15"}),
)
.await;
let resumed = call(
&joaquin,
"ask_agent",
json!({"to": "marta", "resume_message_id": qid, "timeout_seconds": 5}),
)
.await;
assert_eq!(resumed["answered"], true, "{resumed:?}");
assert_eq!(resumed["answer"]["body"], "prod is still on pg15");
let err = call_expect_error(
&joaquin,
"ask_agent",
json!({"to": "joaquin", "question": "hi"}),
)
.await;
assert!(err.contains("yourself"), "{err}");
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
}
#[tokio::test]
async fn transport_limits_reject_oversized_and_too_frequent_requests() {
let h = match setup_rate_limited("t_limits", 60).await {
Some(h) => h,
None => {
eprintln!("skipping: TEST_DATABASE_URL not set");
return;
}
};
let token = seed_agent(&h.pool, "acme", "joaquin").await;
let http = reqwest::Client::new();
let mcp = format!("{}/mcp", h.base);
let call_body = |body: String| {
let http = http.clone();
let mcp = mcp.clone();
let token = token.clone();
async move {
http.post(&mcp)
.header("Authorization", format!("Bearer {token}"))
.header("Content-Type", "application/json")
.header("Accept", "application/json, text/event-stream")
.body(body)
.send()
.await
.unwrap()
}
};
let huge = "x".repeat(200 * 1024);
let resp = call_body(format!(
r#"{{"jsonrpc":"2.0","id":1,"method":"tools/call","params":{{"name":"whoami","arguments":{{"pad":"{huge}"}}}}}}"#
))
.await;
assert_eq!(resp.status(), 413, "oversized body is rejected");
let text = resp.text().await.unwrap();
assert!(
text.contains("too large") && text.contains("attachments"),
"413 tells the caller what to do: {text}"
);
assert!(
text.contains("65536"),
"413 states the limit this server is configured with: {text}"
);
let small = r#"{"jsonrpc":"2.0","id":1,"method":"tools/call","params":{"name":"whoami","arguments":{}}}"#;
let mut throttled = None;
for _ in 0..40 {
let resp = call_body(small.to_owned()).await;
if resp.status() == 429 {
throttled = Some(resp);
break;
}
}
let resp = throttled.expect("a burst of 40 must exhaust a 60/min bucket");
let retry_after = resp
.headers()
.get("retry-after")
.and_then(|v| v.to_str().ok())
.map(str::to_owned);
let text = resp.text().await.unwrap();
assert!(
retry_after.is_some(),
"429 carries Retry-After: headers missing"
);
assert!(
text.contains("rate limit") && text.contains("wait_for_updates"),
"429 points at the non-polling alternative: {text}"
);
let other = seed_agent(&h.pool, "acme", "marta").await;
let resp = http
.post(&mcp)
.header("Authorization", format!("Bearer {other}"))
.header("Content-Type", "application/json")
.header("Accept", "application/json, text/event-stream")
.body(small)
.send()
.await
.unwrap();
assert_eq!(resp.status(), 200, "another token is unaffected");
}
#[tokio::test]
async fn bounded_fields_reject_oversized_values() {
let h = require_db!("t_bounds");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let joaquin = connect(&h.base, &a).await;
let long = "x".repeat(300);
let err = call_expect_error(
&joaquin,
"create_channel",
json!({"name": "dev", "topic": long.clone()}),
)
.await;
assert!(err.contains("256"), "channel topic bounded: {err}");
let err = call_expect_error(&joaquin, "heartbeat", json!({"activity": long.clone()})).await;
assert!(err.contains("256"), "presence activity bounded: {err}");
let err = call_expect_error(
&joaquin,
"set_note",
json!({"key": "k", "value": "v", "tags": vec!["t"; 20]}),
)
.await;
assert!(err.contains("16"), "tag count bounded: {err}");
call(&joaquin, "create_task", json!({"key": "dep", "title": "t"})).await;
let err = call_expect_error(
&joaquin,
"create_task",
json!({"key": "many-deps", "title": "t", "depends_on": vec!["dep"; 40]}),
)
.await;
assert!(err.contains("32"), "dependency count bounded: {err}");
let _ = joaquin.cancel().await;
}
#[tokio::test]
async fn a_wait_on_one_replica_wakes_on_the_other_replicas_write() {
let mut h = require_db!("t_replicas");
let replica = h.add_replica().await;
assert_ne!(replica, h.base, "a genuinely separate instance");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let marta = connect(&replica, &b).await;
call(&marta, "create_channel", json!({"name": "dev"})).await;
let waiter = {
let base = h.base.clone();
let token = a.clone();
tokio::spawn(async move {
let client = connect(&base, &token).await;
let started = std::time::Instant::now();
let result = call(&client, "wait_for_updates", json!({"timeout_seconds": 20})).await;
let _ = client.cancel().await;
(result, started.elapsed())
})
};
tokio::time::sleep(std::time::Duration::from_millis(700)).await;
call(
&marta,
"post_message",
json!({"channel": "dev", "body": "posted through the other replica"}),
)
.await;
let (result, elapsed) = waiter.await.unwrap();
assert_eq!(
result["woke"], true,
"must wake across replicas: {result:?}"
);
assert!(
elapsed < std::time::Duration::from_secs(10),
"woken by the NOTIFY, not by the timeout (took {elapsed:?})"
);
let joaquin = connect(&h.base, &a).await;
let via_one = call(&joaquin, "read_messages", json!({"scope": "dev"})).await;
assert_eq!(via_one["messages"].as_array().map(Vec::len), Some(1));
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
h.shutdown().await;
}
#[tokio::test]
async fn webhook_delivery_is_exactly_one_row_per_hook_across_replicas() {
let mut h = require_db!("t_outbox");
let _replica = h.add_replica().await;
let _replica_two = h.add_replica().await;
let token = seed_agent(&h.pool, "acme", "joaquin").await;
let client = connect(&h.base, &token).await;
call(&client, "create_channel", json!({"name": "dev"})).await;
let team: (Uuid,) = sqlx::query_as("SELECT id FROM teams WHERE slug = 'acme'")
.fetch_one(&h.pool)
.await
.unwrap();
sqlx::query(
"INSERT INTO webhooks (team_id, url, kind, events)
VALUES ($1, 'http://127.0.0.1:9/hook', 'generic', ARRAY['message','task'])",
)
.bind(team.0)
.execute(&h.pool)
.await
.unwrap();
call(
&client,
"post_message",
json!({"channel": "dev", "body": "one event, three replicas"}),
)
.await;
tokio::time::sleep(std::time::Duration::from_millis(1500)).await;
let (rows,): (i64,) =
sqlx::query_as("SELECT count(*) FROM webhook_deliveries WHERE kind = 'message'")
.fetch_one(&h.pool)
.await
.unwrap();
assert_eq!(
rows, 1,
"one channel message must enqueue exactly one delivery, not one per replica"
);
let (status, attempts, err): (String, i32, Option<String>) = sqlx::query_as(
"SELECT status, attempts, last_error FROM webhook_deliveries WHERE kind = 'message'",
)
.fetch_one(&h.pool)
.await
.unwrap();
assert!(attempts >= 1, "the delivery was attempted");
assert!(err.is_some(), "the failure was recorded: {err:?}");
assert!(
status == "pending" || status == "failed",
"a failed delivery is retried or parked, never lost (got {status})"
);
let _marta = seed_agent(&h.pool, "acme", "marta").await;
call(
&client,
"post_message",
json!({"to": "marta", "body": "private"}),
)
.await;
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let (total,): (i64,) = sqlx::query_as("SELECT count(*) FROM webhook_deliveries")
.fetch_one(&h.pool)
.await
.unwrap();
assert_eq!(total, 1, "a DM must never reach the outbox");
call(&client, "create_task", json!({"key": "t1", "title": "t"})).await;
call(&client, "claim_task", json!({"key": "t1"})).await;
call(
&client,
"renew_task_lease",
json!({"key": "t1", "lease_seconds": 600}),
)
.await;
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let (task_rows,): (i64,) =
sqlx::query_as("SELECT count(*) FROM webhook_deliveries WHERE kind = 'task'")
.fetch_one(&h.pool)
.await
.unwrap();
assert_eq!(
task_rows, 2,
"created + claimed enqueue one each; renewing the lease enqueues none"
);
let _ = client.cancel().await;
h.shutdown().await;
}
#[tokio::test]
async fn the_digest_cost_does_not_grow_with_channel_count() {
let h = require_db!("t_digest_scale");
let token = seed_agent(&h.pool, "acme", "joaquin").await;
let client = connect(&h.base, &token).await;
let seed_channels = |from: usize, to: usize| {
let client = &client;
async move {
for c in from..to {
let name = format!("chan{c}");
call(client, "create_channel", json!({"name": name})).await;
for m in 0..6 {
call(
client,
"post_message",
json!({"channel": name, "body": format!("message {m} in {name}")}),
)
.await;
}
}
}
};
seed_channels(0, 2).await;
call(&client, "team_digest", json!({"hours": 24})).await;
let started = std::time::Instant::now();
let small = call(&client, "team_digest", json!({"hours": 24})).await;
let small_elapsed = started.elapsed();
assert_eq!(small["channels"].as_array().map(Vec::len), Some(2));
seed_channels(2, 20).await;
let started = std::time::Instant::now();
let large = call(&client, "team_digest", json!({"hours": 24})).await;
let large_elapsed = started.elapsed();
let channels = large["channels"].as_array().expect("channels");
assert_eq!(channels.len(), 20, "every channel is reported");
for c in channels {
let tail = c["last_messages"].as_array().expect("tail");
assert!(!tail.is_empty() && tail.len() <= 5, "tail is 1..=5: {c:?}");
assert_eq!(c["message_count"], 6, "counts survive the rewrite: {c:?}");
let first = tail[0]["body"].as_str().unwrap_or("");
assert!(
first.contains("message 1"),
"oldest of the tail first: {tail:?}"
);
}
assert!(
large_elapsed < small_elapsed * 4 + std::time::Duration::from_millis(50),
"digest over 20 channels took {large_elapsed:?} against {small_elapsed:?} over 2: \
the cost is scaling with channel count"
);
let _ = client.cancel().await;
h.shutdown().await;
}
#[tokio::test]
async fn attachment_quotas_are_enforced_per_team() {
use base64::Engine;
let b64 = |data: &[u8]| base64::engine::general_purpose::STANDARD.encode(data);
let h = require_db!("t_quota");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let other = seed_agent(&h.pool, "rival", "spy").await;
let joaquin = connect(&h.base, &a).await;
let spy = connect(&h.base, &other).await;
call(&joaquin, "create_channel", json!({"name": "dev"})).await;
call(&spy, "create_channel", json!({"name": "dev"})).await;
sqlx::query("UPDATE teams SET attachment_bytes_limit = $1 WHERE slug = 'acme'")
.bind(300 * 1024i64)
.execute(&h.pool)
.await
.unwrap();
let file = |n: usize| json!([{"filename": "f.bin", "data_base64": b64(&vec![b'x'; n])}]);
for _ in 0..2 {
call(
&joaquin,
"post_message",
json!({"channel": "dev", "body": "chunk", "attachments": file(128 * 1024)}),
)
.await;
}
let err = call_expect_error(
&joaquin,
"post_message",
json!({"channel": "dev", "body": "chunk", "attachments": file(128 * 1024)}),
)
.await;
assert!(err.contains("quota"), "names the problem: {err}");
assert!(
err.contains("307200") && err.contains("raise the quota"),
"states the limit and the way out: {err}"
);
let msgs = call(
&joaquin,
"read_messages",
json!({"scope": "dev", "only_new": false}),
)
.await;
assert_eq!(
msgs["messages"].as_array().map(Vec::len),
Some(2),
"a quota rejection must not leave the message behind: {msgs:?}"
);
call(
&spy,
"post_message",
json!({"channel": "dev", "body": "unbounded", "attachments": file(200 * 1024)}),
)
.await;
let team: (Uuid,) = sqlx::query_as("SELECT id FROM teams WHERE slug = 'acme'")
.fetch_one(&h.pool)
.await
.unwrap();
let usage = ai_crew_sync::store::quota::usage(&h.pool, team.0)
.await
.expect("usage");
assert_eq!(usage.attachment_count, 2);
assert_eq!(usage.attachment_bytes, 256 * 1024);
assert_eq!(usage.attachment_bytes_limit, Some(300 * 1024));
sqlx::query("UPDATE teams SET attachment_bytes_limit = $1 WHERE slug = 'acme'")
.bind(256 * 1024i64 + 100 * 1024)
.execute(&h.pool)
.await
.unwrap();
let racers: Vec<_> = (0..4)
.map(|_| {
let base = h.base.clone();
let token = a.clone();
tokio::spawn(async move {
let c = connect(&base, &token).await;
let args: serde_json::Map<String, serde_json::Value> =
serde_json::from_value(json!({"channel": "dev", "body": "race",
"attachments": [{"filename": "r.bin",
"data_base64": base64::engine::general_purpose::STANDARD
.encode(vec![b'y'; 90 * 1024])}]}))
.unwrap();
let ok = c
.call_tool(
CallToolRequestParams::new("post_message".to_string()).with_arguments(args),
)
.await
.map(|r| r.is_error != Some(true))
.unwrap_or(false);
let _ = c.cancel().await;
ok
})
})
.collect();
let mut accepted = 0;
for r in racers {
if r.await.unwrap_or(false) {
accepted += 1;
}
}
assert_eq!(
accepted, 1,
"only one of four racing 90 KiB uploads fits in 100 KiB of room"
);
let usage = ai_crew_sync::store::quota::usage(&h.pool, team.0)
.await
.expect("usage");
assert!(
usage.attachment_bytes <= usage.attachment_bytes_limit.unwrap(),
"the quota was never exceeded: {} > {:?}",
usage.attachment_bytes,
usage.attachment_bytes_limit
);
let _ = joaquin.cancel().await;
let _ = spy.cancel().await;
h.shutdown().await;
}
#[tokio::test]
async fn pruning_is_dry_by_default_and_keeps_durable_state() {
let h = require_db!("t_prune");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let joaquin = connect(&h.base, &a).await;
call(&joaquin, "create_channel", json!({"name": "dev"})).await;
call(
&joaquin,
"post_message",
json!({"channel": "dev", "body": "old"}),
)
.await;
call(&joaquin, "set_note", json!({"key": "k", "value": "v1"})).await;
call(&joaquin, "set_note", json!({"key": "k", "value": "v2"})).await;
call(&joaquin, "create_task", json!({"key": "t", "title": "t"})).await;
let team: (Uuid,) = sqlx::query_as("SELECT id FROM teams WHERE slug = 'acme'")
.fetch_one(&h.pool)
.await
.unwrap();
for sql in [
"UPDATE messages SET created_at = now() - interval '200 days'",
"UPDATE note_revisions SET created_at = now() - interval '200 days'",
"UPDATE task_events SET created_at = now() - interval '200 days'",
] {
sqlx::query(sql).execute(&h.pool).await.unwrap();
}
let dry = ai_crew_sync::store::quota::prune(&h.pool, team.0, 90, true)
.await
.expect("dry run");
assert!(dry.dry_run);
assert_eq!(dry.messages, 1, "reports what it would delete");
let (still,): (i64,) = sqlx::query_as("SELECT count(*) FROM messages")
.fetch_one(&h.pool)
.await
.unwrap();
assert_eq!(still, 1, "a dry run deletes nothing");
let applied = ai_crew_sync::store::quota::prune(&h.pool, team.0, 90, false)
.await
.expect("apply");
assert_eq!(
applied.messages, dry.messages,
"the dry run's count was the real one"
);
let note = call(&joaquin, "get_note", json!({"key": "k"})).await;
assert_eq!(
note["note"]["value"], "v2",
"notes are not pruned: {note:?}"
);
let task = call(&joaquin, "get_task", json!({"key": "t"})).await;
assert_eq!(task["task"]["key"], "t", "tasks are not pruned: {task:?}");
let err = ai_crew_sync::store::quota::prune(&h.pool, team.0, 0, true).await;
assert!(err.is_err(), "older_than_days must be at least 1");
let err = ai_crew_sync::store::quota::prune(&h.pool, team.0, 2_147_483_648, true).await;
assert!(err.is_err(), "a day count above i32::MAX must be refused");
let (survived,): (i64,) = sqlx::query_as("SELECT count(*) FROM notes")
.fetch_one(&h.pool)
.await
.unwrap();
assert_eq!(survived, 1, "a refused prune deletes nothing");
let _ = joaquin.cancel().await;
h.shutdown().await;
}
#[tokio::test]
async fn shutdown_stops_the_server_and_its_background_tasks() {
let h = require_db!("t_shutdown");
let base = h.base.clone();
let pool = h.pool.clone();
let token = seed_agent(&h.pool, "acme", "joaquin").await;
let http = reqwest::Client::new();
let resp = http.get(format!("{base}/health")).send().await.unwrap();
assert_eq!(resp.status(), 200);
h.shutdown().await;
let resp = tokio::time::timeout(
std::time::Duration::from_secs(5),
http.get(format!("{base}/health")).send(),
)
.await
.expect("the request must not hang after shutdown");
assert!(
resp.is_err(),
"the server should no longer accept connections"
);
let seeded = sqlx::query("SELECT 1").execute(&pool).await;
assert!(
seeded.is_err(),
"the harness pool must be closed after shutdown"
);
assert!(!token.is_empty(), "the agent was seeded before shutdown");
}
#[tokio::test]
async fn the_database_refuses_cross_team_references() {
let h = require_db!("t_teamfk");
let team_a: (Uuid,) =
sqlx::query_as("INSERT INTO teams (slug, name) VALUES ('a', 'A') RETURNING id")
.fetch_one(&h.pool)
.await
.unwrap();
let team_b: (Uuid,) =
sqlx::query_as("INSERT INTO teams (slug, name) VALUES ('b', 'B') RETURNING id")
.fetch_one(&h.pool)
.await
.unwrap();
let agent_a: (Uuid,) =
sqlx::query_as("INSERT INTO agents (team_id, name) VALUES ($1, 'a') RETURNING id")
.bind(team_a.0)
.fetch_one(&h.pool)
.await
.unwrap();
let agent_b: (Uuid,) =
sqlx::query_as("INSERT INTO agents (team_id, name) VALUES ($1, 'b') RETURNING id")
.bind(team_b.0)
.fetch_one(&h.pool)
.await
.unwrap();
let channel_a: (Uuid,) =
sqlx::query_as("INSERT INTO channels (team_id, name) VALUES ($1, 'dev') RETURNING id")
.bind(team_a.0)
.fetch_one(&h.pool)
.await
.unwrap();
let task_a: (Uuid,) = sqlx::query_as(
"INSERT INTO tasks (team_id, key, title) VALUES ($1, 'ta', 't') RETURNING id",
)
.bind(team_a.0)
.fetch_one(&h.pool)
.await
.unwrap();
let task_b: (Uuid,) = sqlx::query_as(
"INSERT INTO tasks (team_id, key, title) VALUES ($1, 'tb', 't') RETURNING id",
)
.bind(team_b.0)
.fetch_one(&h.pool)
.await
.unwrap();
let err = sqlx::query(
"INSERT INTO messages (team_id, channel_id, sender_agent_id, body) VALUES ($1,$2,$3,'x')",
)
.bind(team_a.0)
.bind(channel_a.0)
.bind(agent_b.0)
.execute(&h.pool)
.await;
assert!(err.is_err(), "a sender from another team must be rejected");
let err = sqlx::query(
"INSERT INTO messages (team_id, sender_agent_id, channel_id, body) VALUES ($1,$2,$3,'x')",
)
.bind(team_b.0)
.bind(agent_b.0)
.bind(channel_a.0)
.execute(&h.pool)
.await;
assert!(err.is_err(), "another team's channel must be rejected");
let err = sqlx::query(
"INSERT INTO messages (team_id, sender_agent_id, recipient_agent_id, body)
VALUES ($1,$2,$3,'x')",
)
.bind(team_a.0)
.bind(agent_a.0)
.bind(agent_b.0)
.execute(&h.pool)
.await;
assert!(err.is_err(), "a cross-team DM must be rejected");
let err = sqlx::query(
"INSERT INTO locks (team_id, name, holder_agent_id, expires_at)
VALUES ($1,'x',$2, now() + interval '1 hour')",
)
.bind(team_a.0)
.bind(agent_b.0)
.execute(&h.pool)
.await;
assert!(err.is_err(), "a holder from another team must be rejected");
let err = sqlx::query("INSERT INTO task_deps (task_id, blocked_by_task_id) VALUES ($1,$2)")
.bind(task_a.0)
.bind(task_b.0)
.execute(&h.pool)
.await;
assert!(err.is_err(), "a cross-team dependency must be rejected");
let msg_a: (i64,) = sqlx::query_as(
"INSERT INTO messages (team_id, channel_id, sender_agent_id, body)
VALUES ($1,$2,$3,'legit') RETURNING id",
)
.bind(team_a.0)
.bind(channel_a.0)
.bind(agent_a.0)
.fetch_one(&h.pool)
.await
.expect("a same-team message is still accepted");
let err = sqlx::query(
"INSERT INTO attachments (team_id, message_id, uploader_agent_id, filename, size_bytes, data)
VALUES ($1,$2,$3,'f',1,'\\x00')",
)
.bind(team_b.0)
.bind(msg_a.0)
.bind(agent_b.0)
.execute(&h.pool)
.await;
assert!(
err.is_err(),
"an attachment on another team's message must be rejected"
);
}
#[tokio::test]
async fn oversized_fields_are_rejected_with_their_limit() {
let h = require_db!("t_caps");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let joaquin = connect(&h.base, &a).await;
call(&joaquin, "create_channel", json!({"name": "dev"})).await;
let fat = "x".repeat(20 * 1024);
let err = call_expect_error(
&joaquin,
"post_message",
json!({"channel": "dev", "body": "hi", "metadata": {"blob": fat}}),
)
.await;
assert!(err.contains("16384"), "names the metadata limit: {err}");
let err = call_expect_error(
&joaquin,
"create_task",
json!({"key": "fat-meta", "title": "t", "metadata": {"blob": fat}}),
)
.await;
assert!(err.contains("16384"), "same limit on tasks: {err}");
let long_title = "t".repeat(600);
let err = call_expect_error(
&joaquin,
"create_task",
json!({"key": "long-title", "title": long_title}),
)
.await;
assert!(err.contains("512"), "names the title limit: {err}");
let long_text = "d".repeat(70 * 1024);
let err = call_expect_error(
&joaquin,
"create_task",
json!({"key": "long-desc", "title": "t", "description": long_text.clone()}),
)
.await;
assert!(err.contains("65536"), "names the description limit: {err}");
call(
&joaquin,
"create_task",
json!({"key": "capped", "title": "fits"}),
)
.await;
call(&joaquin, "claim_task", json!({"key": "capped"})).await;
let err = call_expect_error(
&joaquin,
"complete_task",
json!({"key": "capped", "result": long_text}),
)
.await;
assert!(err.contains("65536"), "names the result limit: {err}");
let task = call(&joaquin, "get_task", json!({"key": "capped"})).await;
assert_eq!(task["task"]["status"], "claimed", "{task:?}");
let one_mib = "b".repeat(1024 * 1024);
call(
&joaquin,
"post_message",
json!({"channel": "dev", "body": one_mib.clone()}),
)
.await;
let err = call_expect_error(
&joaquin,
"post_message",
json!({"channel": "dev", "body": format!("{one_mib}x")}),
)
.await;
assert!(err.contains("1048576"), "names the body limit: {err}");
call(
&joaquin,
"set_note",
json!({"key": "big-note", "value": one_mib.clone()}),
)
.await;
let err = call_expect_error(
&joaquin,
"set_note",
json!({"key": "big-note", "value": format!("{one_mib}x")}),
)
.await;
assert!(err.contains("1048576"), "names the note limit: {err}");
let padded = format!("{}fits", " ".repeat(700));
let err = call_expect_error(
&joaquin,
"create_task",
json!({"key": "padded", "title": padded}),
)
.await;
assert!(
err.contains("512"),
"raw size counts, not just trimmed: {err}"
);
call(
&joaquin,
"post_message",
json!({"channel": "dev", "body": "ok", "metadata": {"k": "v"}}),
)
.await;
call(
&joaquin,
"create_task",
json!({"key": "ok-task", "title": "t".repeat(512), "description": "d".repeat(1000)}),
)
.await;
let _ = joaquin.cancel().await;
}
#[tokio::test]
async fn attachments_travel_with_messages_and_tasks() {
use base64::Engine;
let b64 = |data: &[u8]| base64::engine::general_purpose::STANDARD.encode(data);
let h = require_db!("t_attach");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let c = seed_agent(&h.pool, "acme", "pedro").await;
let joaquin = connect(&h.base, &a).await;
let marta = connect(&h.base, &b).await;
let pedro = connect(&h.base, &c).await;
let diff = "diff --git a/src/lib.rs b/src/lib.rs\n-old\n+new\n";
call(&joaquin, "create_channel", json!({"name": "dev"})).await;
let posted = call(
&joaquin,
"post_message",
json!({
"channel": "dev", "body": "parser fix attached",
"attachments": [{
"filename": "fix.diff", "content_type": "text/plain",
"data_base64": b64(diff.as_bytes())
}]
}),
)
.await;
let att = &posted["message"]["attachments"][0];
assert_eq!(att["filename"], "fix.diff", "{posted:?}");
let att_id = att["id"].as_i64().unwrap();
let read = call(&marta, "read_messages", json!({"scope": "dev"})).await;
let msg = read["messages"].as_array().unwrap().last().unwrap().clone();
assert_eq!(msg["attachments"][0]["id"], att_id, "{msg:?}");
let got = call(&marta, "get_attachment", json!({"id": att_id})).await;
assert_eq!(got["data_base64"].as_str().unwrap(), b64(diff.as_bytes()));
assert_eq!(got["uploaded_by"], "joaquin");
let dm = call(
&joaquin,
"post_message",
json!({
"to": "marta", "body": "the failing log",
"attachments": [{"filename": "secret.log", "data_base64": b64(b"boom")}]
}),
)
.await;
let dm_att = dm["message"]["attachments"][0]["id"].as_i64().unwrap();
call(&marta, "get_attachment", json!({"id": dm_att})).await;
let err = call_expect_error(&pedro, "get_attachment", json!({"id": dm_att})).await;
assert!(
err.contains("not found"),
"third party must not see it: {err}"
);
call(
&joaquin,
"create_task",
json!({"key": "fix-parser", "title": "Fix the parser"}),
)
.await;
call(
&marta,
"attach_file",
json!({"task": "fix-parser", "filename": "repro.log", "data_base64": b64(b"repro")}),
)
.await;
let task = call(&pedro, "get_task", json!({"key": "fix-parser"})).await;
assert_eq!(
task["task"]["attachments"][0]["filename"], "repro.log",
"{task:?}"
);
let big = vec![b'x'; 300 * 1024];
let err = call_expect_error(
&joaquin,
"attach_file",
json!({"task": "fix-parser", "filename": "big.bin", "data_base64": b64(&big)}),
)
.await;
assert!(err.contains("262144"), "must state the limit: {err}");
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
let _ = pedro.cancel().await;
}
#[tokio::test]
async fn blocked_tasks_wait_for_their_dependencies() {
let h = require_db!("t_deps");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &a).await;
let marta = connect(&h.base, &b).await;
call(
&joaquin,
"create_task",
json!({"key": "migrate-schema", "title": "migrate the users schema"}),
)
.await;
let dependent = call(
&joaquin,
"create_task",
json!({
"key": "update-clients",
"title": "update the API clients",
"depends_on": ["migrate-schema"]
}),
)
.await;
assert_eq!(dependent["blocked"], true);
assert_eq!(dependent["depends_on"][0], "migrate-schema");
let err = call_expect_error(
&joaquin,
"create_task",
json!({"key": "x", "title": "x", "depends_on": ["nope"]}),
)
.await;
assert!(err.contains("nope"), "got: {err}");
let denied = call(&marta, "claim_task", json!({"key": "update-clients"})).await;
assert_eq!(denied["claimed"], false);
assert!(
denied["reason"]
.as_str()
.unwrap()
.contains("migrate-schema"),
"reason should name the blocker: {denied:?}"
);
let next = call(&marta, "claim_next_task", json!({})).await;
assert_eq!(next["claimed"], true);
assert_eq!(next["task"]["key"], "migrate-schema");
call(&marta, "complete_task", json!({"key": "migrate-schema"})).await;
let now_free = call(&joaquin, "claim_task", json!({"key": "update-clients"})).await;
assert_eq!(
now_free["claimed"], true,
"unblocked after dep done: {now_free:?}"
);
assert_eq!(now_free["task"]["blocked"], false);
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
}
#[tokio::test]
async fn locks_are_exclusive_expiring_and_visible() {
let h = require_db!("t_locks");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &a).await;
let marta = connect(&h.base, &b).await;
let got = call(
&joaquin,
"acquire_lock",
json!({"name": "Deploy:Staging", "ttl_seconds": 120, "purpose": "rolling out 1.4.2"}),
)
.await;
assert_eq!(got["acquired"], true);
assert_eq!(got["lock"]["name"], "deploy:staging", "name normalised");
let denied = call(&marta, "acquire_lock", json!({"name": "deploy:staging"})).await;
assert_eq!(denied["acquired"], false);
assert!(denied["reason"].as_str().unwrap().contains("joaquin"));
let extended = call(
&joaquin,
"acquire_lock",
json!({"name": "deploy:staging", "ttl_seconds": 600}),
)
.await;
assert_eq!(extended["acquired"], true);
let listed = call(&marta, "list_locks", json!({})).await;
assert_eq!(listed["locks"][0]["holder"], "joaquin");
assert_eq!(listed["locks"][0]["purpose"], "rolling out 1.4.2");
let err = call_expect_error(&marta, "release_lock", json!({"name": "deploy:staging"})).await;
assert!(err.contains("joaquin"), "got: {err}");
call(&joaquin, "release_lock", json!({"name": "deploy:staging"})).await;
let now = call(&marta, "acquire_lock", json!({"name": "deploy:staging"})).await;
assert_eq!(now["acquired"], true);
sqlx::query("UPDATE locks SET expires_at = now() - interval '1 second'")
.execute(&h.pool)
.await
.unwrap();
let stolen = call(&joaquin, "acquire_lock", json!({"name": "deploy:staging"})).await;
assert_eq!(stolen["acquired"], true, "expired lock must be stealable");
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
}
#[tokio::test]
async fn team_digest_summarises_recent_activity() {
let h = require_db!("t_digest");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &a).await;
let marta = connect(&h.base, &b).await;
call(&joaquin, "create_channel", json!({"name": "deploys"})).await;
call(
&joaquin,
"post_message",
json!({"channel": "deploys", "body": "staging lleva la 1.4.2"}),
)
.await;
call(
&marta,
"create_task",
json!({"key": "hotfix", "title": "hotfix the parser"}),
)
.await;
call(&marta, "claim_task", json!({"key": "hotfix"})).await;
call(
&marta,
"complete_task",
json!({"key": "hotfix", "result": "merged in #99"}),
)
.await;
call(
&joaquin,
"set_note",
json!({"scope": "api", "key": "deploy-runbook", "value": "step 1..."}),
)
.await;
call(&marta, "heartbeat", json!({"activity": "reviewing PRs"})).await;
call(
&joaquin,
"post_message",
json!({"to": "marta", "body": "esto es privado"}),
)
.await;
let digest = call(&joaquin, "team_digest", json!({"hours": 24})).await;
assert_eq!(digest["channels"][0]["name"], "deploys");
assert_eq!(digest["channels"][0]["message_count"], 1);
let tasks: Vec<&str> = digest["tasks_moved"]
.as_array()
.unwrap()
.iter()
.map(|t| t["key"].as_str().unwrap())
.collect();
assert!(tasks.contains(&"hotfix"));
assert_eq!(digest["notes_updated"][0]["key"], "deploy-runbook");
assert!(
digest["agents_seen"]
.as_array()
.unwrap()
.iter()
.any(|a| a["name"] == "marta" && a["online"] == true)
);
let serialized = serde_json::to_string(&digest).unwrap();
assert!(
!serialized.contains("privado"),
"digest must never contain direct messages"
);
let _ = joaquin.cancel().await;
let _ = marta.cancel().await;
}
#[tokio::test]
async fn webhooks_forward_channel_messages_but_never_dms() {
let h = require_db!("t_webhooks");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let b = seed_agent(&h.pool, "acme", "marta").await;
let joaquin = connect(&h.base, &a).await;
let received: std::sync::Arc<tokio::sync::Mutex<Vec<Value>>> = Default::default();
let catcher = {
let received = received.clone();
let app = axum::Router::new().route(
"/hook",
axum::routing::post(move |axum::Json(v): axum::Json<Value>| {
let received = received.clone();
async move {
received.lock().await.push(v);
"ok"
}
}),
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
let _ = axum::serve(listener, app).await;
});
addr
};
ai_crew_sync::webhooks::webhook_add(
&h.pool,
"acme",
&format!("http://{catcher}/hook"),
"slack",
"message,task",
None,
)
.await
.unwrap();
call(&joaquin, "create_channel", json!({"name": "deploys"})).await;
call(
&joaquin,
"post_message",
json!({"channel": "deploys", "body": "canary verde"}),
)
.await;
call(
&joaquin,
"post_message",
json!({"to": "marta", "body": "secreto entre nosotros"}),
)
.await;
call(
&joaquin,
"create_task",
json!({"key": "rotate", "title": "rotate keys"}),
)
.await;
let _ = b;
let mut tries = 0;
loop {
tokio::time::sleep(std::time::Duration::from_millis(300)).await;
let got = received.lock().await;
if got.len() >= 2 || tries > 20 {
break;
}
drop(got);
tries += 1;
}
let got = received.lock().await;
let texts: Vec<String> = got
.iter()
.map(|v| v["text"].as_str().unwrap_or_default().to_owned())
.collect();
assert!(
texts
.iter()
.any(|t| t.contains("#deploys") && t.contains("canary verde")),
"channel message must be forwarded in Slack format: {texts:?}"
);
assert!(
texts.iter().any(|t| t.contains("rotate")),
"task event must be forwarded: {texts:?}"
);
assert!(
!texts.iter().any(|t| t.contains("secreto")),
"a DM must NEVER reach a webhook: {texts:?}"
);
let _ = joaquin.cancel().await;
}
#[tokio::test]
async fn dashboard_requires_a_token_and_renders_team_state() {
let h = require_db!("t_dash");
let a = seed_agent(&h.pool, "acme", "joaquin").await;
let joaquin = connect(&h.base, &a).await;
call(&joaquin, "heartbeat", json!({"activity": "smoke testing"})).await;
call(&joaquin, "create_channel", json!({"name": "dev"})).await;
call(
&joaquin,
"post_message",
json!({"channel": "dev", "body": "<script>alert(1)</script> it's here"}),
)
.await;
let http = reqwest::Client::builder()
.redirect(reqwest::redirect::Policy::none())
.build()
.unwrap();
let base = &h.base;
let resp = http.get(format!("{base}/dashboard")).send().await.unwrap();
assert_eq!(resp.status(), 401);
let body = resp.text().await.unwrap();
assert!(body.contains("<form"), "offers a form to sign in: {body}");
assert!(!body.contains("smoke testing"), "leaks no team state");
let resp = http
.get(format!("{base}/dashboard?token={a}"))
.send()
.await
.unwrap();
assert_eq!(
resp.status(),
401,
"query-string tokens must not authenticate"
);
let resp = http
.post(format!("{base}/dashboard/login"))
.header("Content-Type", "application/x-www-form-urlencoded")
.body(format!("token={a}"))
.send()
.await
.unwrap();
assert!(
resp.status().is_redirection(),
"successful login redirects: {}",
resp.status()
);
let cookie = resp
.headers()
.get("set-cookie")
.and_then(|v| v.to_str().ok())
.expect("a session cookie")
.to_owned();
assert!(cookie.contains("HttpOnly"), "cookie is HttpOnly: {cookie}");
assert!(
cookie.contains("SameSite=Strict"),
"cookie is SameSite=Strict: {cookie}"
);
assert!(
!cookie.contains(&a),
"the agent token itself must never be the cookie value"
);
let grant = cookie.split(';').next().expect("cookie pair").to_owned();
let resp = http
.post(format!("{base}/dashboard/login"))
.header("Content-Type", "application/x-www-form-urlencoded")
.body("token=acs_bogus")
.send()
.await
.unwrap();
assert_eq!(resp.status(), 401);
assert!(resp.headers().get("set-cookie").is_none());
let resp = http
.get(format!("{base}/dashboard"))
.header("Cookie", &grant)
.send()
.await
.unwrap();
assert_eq!(resp.status(), 200);
assert_eq!(
resp.headers()
.get("cache-control")
.and_then(|v| v.to_str().ok()),
Some("no-store"),
"team activity is never cached"
);
assert_eq!(
resp.headers()
.get("referrer-policy")
.and_then(|v| v.to_str().ok()),
Some("no-referrer")
);
let body = resp.text().await.unwrap();
assert!(body.contains("joaquin"), "shows the agent");
assert!(body.contains("smoke testing"), "shows the activity");
assert!(
!body.contains("<script>alert(1)</script>"),
"message bodies must be HTML-escaped"
);
assert!(body.contains("<script>"), "escaped form present");
assert!(!body.contains("it's here"), "single quotes escaped too");
assert!(body.contains("it's here"), "escaped quote present");
let grant_value = grant.split_once('=').expect("cookie pair").1.to_owned();
for attempt in [
http.post(format!("{base}/mcp"))
.header("Cookie", &grant)
.header("Accept", "application/json, text/event-stream")
.json(&json!({"jsonrpc":"2.0","id":1,"method":"tools/list"})),
http.post(format!("{base}/mcp"))
.header("Authorization", format!("Bearer {grant_value}"))
.header("Accept", "application/json, text/event-stream")
.json(&json!({"jsonrpc":"2.0","id":1,"method":"tools/list"})),
] {
let resp = attempt.send().await.unwrap();
assert_eq!(
resp.status(),
401,
"a dashboard grant must not authenticate an MCP call"
);
}
let resp = http
.get(format!("{base}/dashboard"))
.header("Authorization", format!("Bearer {a}"))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 200, "curl/script access keeps working");
let _ = joaquin.cancel().await;
}