use super::*;
#[tokio::test]
async fn chat_creation_requires_an_existing_bot_after_workspace_selection() {
let root = tempfile::tempdir().expect("root");
let workspace = root.path().join("workspace");
std::fs::create_dir(&workspace).expect("workspace");
let (server, grant) = GatewayServer::bootstrap(
root.path().join("state"),
std::net::SocketAddr::from(([127, 0, 0, 1], 0)),
)
.await
.expect("bootstrap gateway");
let listen = server.listen_addr();
let (shutdown, signal) = tokio::sync::oneshot::channel();
let serving = tokio::spawn(server.serve_until(async move {
let _ = signal.await;
}));
let endpoint = format!("tcp://{listen}")
.parse::<Endpoint>()
.expect("endpoint");
let (connection, _) = GatewayClient::pair(&endpoint, grant.code, "Bot test", ClientKind::Ios)
.await
.expect("connect");
let (sender, mut events) = connection.into_parts();
wait_gateway_ready(&mut events).await;
sender
.send(ClientMessage::CreateSession {
request_id: "create".into(),
workspace,
bot_id: Uuid::new_v4().to_string(),
})
.await
.expect("send create");
loop {
if let ServerMessage::Rejected {
request_id, code, ..
} = next_gateway_message(&mut events).await
&& request_id == "create"
{
assert_eq!(code, "invalid_bot");
break;
}
}
shutdown.send(()).expect("stop gateway");
serving.await.expect("gateway task").expect("gateway stop");
}
#[tokio::test]
async fn deleting_a_bot_clears_its_selected_chat_on_the_requesting_connection() {
let root = tempfile::tempdir().expect("root");
let workspace = root.path().join("workspace");
std::fs::create_dir(&workspace).expect("workspace");
let (server, grant) = configured_test_server(root.path().join("state")).await;
let listen = server.listen_addr();
let bots = Arc::clone(&server.bots);
let (shutdown, signal) = tokio::sync::oneshot::channel();
let serving = tokio::spawn(server.serve_until(async move {
let _ = signal.await;
}));
let endpoint = format!("tcp://{listen}")
.parse::<Endpoint>()
.expect("endpoint");
let (connection, _) =
GatewayClient::pair(&endpoint, grant.code, "Bot deletion test", ClientKind::Ios)
.await
.expect("connect");
let (sender, mut events) = connection.into_parts();
wait_gateway_ready(&mut events).await;
let (session_id, bot_id) = create_bot_chat(&sender, &mut events, &workspace).await;
let bot = bots.bot(&bot_id).expect("created Bot");
sender
.send(ClientMessage::DeleteBot {
request_id: "delete-bot".into(),
id: bot.id,
expected_revision: bot.config.revision,
})
.await
.expect("delete Bot");
loop {
match next_gateway_message(&mut events).await {
ServerMessage::Bots {
request_id: Some(request_id),
..
} if request_id == "delete-bot" => break,
ServerMessage::Rejected {
request_id,
code,
message,
..
} if request_id == "delete-bot" => {
panic!("Bot deletion rejected ({code}): {message}")
}
_ => {}
}
}
sender
.send(ClientMessage::GetSessionHistory {
request_id: "deleted-history".into(),
session_id,
before_sequence: None,
})
.await
.expect("request deleted history");
loop {
if let ServerMessage::Rejected {
request_id, code, ..
} = next_gateway_message(&mut events).await
&& request_id == "deleted-history"
{
assert_eq!(code, "session_required");
break;
}
}
shutdown.send(()).expect("stop gateway");
serving.await.expect("gateway task").expect("gateway stop");
}
#[tokio::test]
async fn bot_catalog_broadcasts_do_not_reintroduce_a_deleted_bot() {
let root = tempfile::tempdir().expect("root");
let (server, grant) = configured_test_server(root.path().join("state")).await;
let bot = server
.host
.create_bot("Original", "Bot catalog ordering test")
.await
.expect("create Bot");
let (client, stream) = tokio::io::duplex(1024 * 1024);
let (reader, mut writer) = tokio::io::split(client);
let mut reader = FrameReader::new(reader);
for message in [
ClientMessage::Pair {
code: grant.code,
client_label: "Bot catalog test".into(),
client_kind: ClientKind::Ios,
},
ClientMessage::UpdateBot {
request_id: "rename".into(),
id: bot.id.clone(),
expected_revision: bot.config.revision,
name: "Renamed".into(),
description: bot.description,
tint: bot.tint,
config: bot.config.config,
},
ClientMessage::DeleteBot {
request_id: "delete".into(),
id: bot.id.clone(),
expected_revision: bot.config.revision + 1,
},
] {
write_frame(&mut writer, &ClientFrame::new(message))
.await
.expect("queue request");
}
let (client_revocations, _) = broadcast::channel(MAX_CONNECTIONS);
let serving = tokio::spawn(serve_connection(
stream,
ConnectionContext {
auth: server.auth,
host: server.host,
bots: server.bots,
client_connections: Arc::new(ClientConnections::default()),
client_revocations,
admission: ConnectionAdmission::new(1, 1).admit().await,
},
Instant::now() + PRE_AUTH_TIMEOUT,
None,
));
let mut deleted = false;
loop {
let frame = tokio::time::timeout(
Duration::from_secs(5),
read_frame::<ServerFrame>(&mut reader),
)
.await
.expect("response timeout")
.expect("read response")
.expect("connection open");
match frame.message {
ServerMessage::Bots { request_id, bots } => {
deleted |= request_id.as_deref() == Some("delete");
if deleted {
assert!(
bots.iter().all(|candidate| candidate.id != bot.id),
"a stale catalog reintroduced the deleted Bot after its deletion response"
);
if request_id.is_none() {
break;
}
}
}
ServerMessage::Rejected { code, message, .. } => {
panic!("Bot mutation rejected ({code}): {message}")
}
_ => {}
}
}
drop(reader);
drop(writer);
serving
.await
.expect("connection task")
.expect("connection stop");
}
#[tokio::test]
async fn swarm_catalog_broadcasts_use_the_latest_renamed_bot_handle() {
let root = tempfile::tempdir().expect("root");
let (server, grant) = configured_test_server(root.path().join("state")).await;
let mut composition = crate::wire::AgentComposition::default();
composition.middleware.set_setting(
"bots",
"collaboration",
Some(mobius::protocol::FrontendSettingValue::String(
"swarm".into(),
)),
);
let leader = server
.bots
.create_bot("Leader", "Swarm leader", composition.clone())
.expect("create leader");
let member = server
.bots
.create_bot("Member", "Swarm member", composition)
.expect("create member");
let swarm = server
.host
.create_swarm(
"Rename team".into(),
leader.id.clone(),
vec![member.id.clone()],
)
.await
.expect("create swarm")
.into_iter()
.next()
.expect("created swarm");
let (client, stream) = tokio::io::duplex(1024 * 1024);
let (reader, mut writer) = tokio::io::split(client);
let mut reader = FrameReader::new(reader);
for message in [
ClientMessage::Pair {
code: grant.code,
client_label: "Swarm rename test".into(),
client_kind: ClientKind::Ios,
},
ClientMessage::UpdateBot {
request_id: "rename-one".into(),
id: leader.id.clone(),
expected_revision: leader.config.revision,
name: "Leader One".into(),
description: leader.description.clone(),
tint: leader.tint,
config: leader.config.config.clone(),
},
ClientMessage::UpdateBot {
request_id: "rename-two".into(),
id: leader.id.clone(),
expected_revision: leader.config.revision + 1,
name: "Leader Two".into(),
description: leader.description.clone(),
tint: leader.tint,
config: leader.config.config.clone(),
},
] {
write_frame(&mut writer, &ClientFrame::new(message))
.await
.expect("queue rename");
}
let (client_revocations, _) = broadcast::channel(MAX_CONNECTIONS);
let serving = tokio::spawn(serve_connection(
stream,
ConnectionContext {
auth: server.auth,
host: server.host,
bots: server.bots,
client_connections: Arc::new(ClientConnections::default()),
client_revocations,
admission: ConnectionAdmission::new(1, 1).admit().await,
},
Instant::now() + PRE_AUTH_TIMEOUT,
None,
));
let mut final_rename_seen = false;
let mut fresh_swarms = 0;
loop {
let frame = tokio::time::timeout(
Duration::from_secs(5),
read_frame::<ServerFrame>(&mut reader),
)
.await
.expect("response timeout")
.expect("read response")
.expect("connection open");
match frame.message {
ServerMessage::Bots {
request_id: Some(request_id),
bots,
} if request_id == "rename-two" => {
let bot = bots
.into_iter()
.find(|candidate| candidate.id == leader.id)
.expect("final renamed Bot");
assert_eq!(bot.name, "Leader Two");
assert_eq!(bot.handle, "leader-two");
final_rename_seen = true;
}
ServerMessage::Swarms {
request_id: None,
swarms,
} if final_rename_seen => {
let swarm = swarms
.iter()
.find(|candidate| candidate.id == swarm.id)
.expect("renamed Swarm");
let member = swarm
.members
.iter()
.find(|candidate| candidate.bot_id == leader.id)
.expect("renamed Swarm member");
assert_eq!(member.handle, "leader-two");
fresh_swarms += 1;
if fresh_swarms == 2 {
break;
}
}
ServerMessage::Rejected { code, message, .. } => {
panic!("Bot mutation rejected ({code}): {message}")
}
_ => {}
}
}
drop(reader);
drop(writer);
serving
.await
.expect("connection task")
.expect("connection stop");
}