mod config;
mod lifecycle;
mod turn;
mod workspaces;
pub use config::{ServeConfig, SessionSource};
use agent_client_protocol::{
Agent, Client, ConnectionTo, Handled,
schema::v1::{
AgentCapabilities, AvailableCommand, AvailableCommandsUpdate, CancelNotification,
CloseSessionRequest, Error, InitializeRequest, InitializeResponse, ListSessionsRequest,
LoadSessionRequest, LoadSessionResponse, McpCapabilities, NewSessionRequest,
PromptCapabilities, PromptRequest, ResumeSessionRequest, ResumeSessionResponse,
SessionCapabilities, SessionCloseCapabilities, SessionId, SessionListCapabilities,
SessionNotification, SessionResumeCapabilities, SessionUpdate, SetSessionModeRequest,
},
};
use crate::session::SessionRegistry;
use lifecycle::Replay;
pub async fn serve<T>(config: ServeConfig, transport: T) -> Result<(), Error>
where
T: agent_client_protocol::ConnectTo<Agent> + 'static,
{
let sessions = SessionRegistry::new();
Agent
.builder()
.name("basis")
.on_receive_request(
{
let config = config.clone();
async move |request: InitializeRequest, responder, _connection| {
responder.respond(initialize(&request, &config))
}
},
agent_client_protocol::on_receive_request!(),
)
.on_receive_request(
{
let config = config.clone();
async move |request: ListSessionsRequest, responder, _connection| {
responder.respond_with_result(lifecycle::list_sessions(&config, request).await)
}
},
agent_client_protocol::on_receive_request!(),
)
.on_receive_request(
{
let sessions = sessions.clone();
let config = config.clone();
async move |request: NewSessionRequest,
responder,
connection: ConnectionTo<Client>| {
let opened = lifecycle::new_session(&config, &sessions, request).await;
let announcement = opened.as_ref().ok().map(|opened| {
(opened.response.session_id.clone(), opened.commands.clone())
});
responder.respond_with_result(opened.map(|opened| opened.response))?;
if let Some((session_id, commands)) = announcement {
let _ = announce_commands(&connection, &session_id, commands);
}
Ok(Handled::Yes)
}
},
agent_client_protocol::on_receive_request!(),
)
.on_receive_request(
{
let sessions = sessions.clone();
let config = config.clone();
async move |request: LoadSessionRequest,
responder,
connection: ConnectionTo<Client>| {
let sessions = sessions.clone();
let config = config.clone();
let spawned = connection.clone();
connection.spawn(async move {
let modes = lifecycle::open_persisted(
&config,
&sessions,
&spawned,
request.session_id,
request.cwd,
request.mcp_servers,
Replay::Yes,
)
.await;
responder.respond_with_result(
modes.map(|modes| LoadSessionResponse::new().modes(modes)),
)
})?;
Ok(Handled::Yes)
}
},
agent_client_protocol::on_receive_request!(),
)
.on_receive_request(
{
let sessions = sessions.clone();
let config = config.clone();
async move |request: ResumeSessionRequest,
responder,
connection: ConnectionTo<Client>| {
let sessions = sessions.clone();
let config = config.clone();
let spawned = connection.clone();
connection.spawn(async move {
let modes = lifecycle::open_persisted(
&config,
&sessions,
&spawned,
request.session_id,
request.cwd,
request.mcp_servers,
Replay::No,
)
.await;
responder.respond_with_result(
modes.map(|modes| ResumeSessionResponse::new().modes(modes)),
)
})?;
Ok(Handled::Yes)
}
},
agent_client_protocol::on_receive_request!(),
)
.on_receive_request(
{
let sessions = sessions.clone();
async move |request: SetSessionModeRequest,
responder,
connection: ConnectionTo<Client>| {
responder.respond_with_result(lifecycle::set_mode(
&sessions,
&connection,
&request,
))
}
},
agent_client_protocol::on_receive_request!(),
)
.on_receive_request(
{
let sessions = sessions.clone();
async move |request: CloseSessionRequest, responder, _connection| {
responder.respond_with_result(lifecycle::close_session(&sessions, &request))
}
},
agent_client_protocol::on_receive_request!(),
)
.on_receive_request(
{
let sessions = sessions.clone();
async move |request: PromptRequest, responder, connection: ConnectionTo<Client>| {
let sessions = sessions.clone();
let spawned = connection.clone();
connection.spawn(async move {
let outcome = turn::prompt(&sessions, &spawned, request).await;
responder.respond_with_result(outcome)
})?;
Ok(Handled::Yes)
}
},
agent_client_protocol::on_receive_request!(),
)
.on_receive_notification(
{
let sessions = sessions.clone();
async move |notification: CancelNotification, _connection| {
if let Some(session) = sessions.get(¬ification.session_id) {
session.cancel();
}
Ok(())
}
},
agent_client_protocol::on_receive_notification!(),
)
.connect_to(transport)
.await
}
fn initialize(request: &InitializeRequest, config: &ServeConfig) -> InitializeResponse {
InitializeResponse::new(request.protocol_version)
.agent_capabilities(
AgentCapabilities::new()
.load_session(true)
.mcp_capabilities(McpCapabilities::new().sse(true))
.prompt_capabilities(PromptCapabilities::new().embedded_context(true))
.session_capabilities(
SessionCapabilities::new()
.resume(SessionResumeCapabilities::new())
.close(SessionCloseCapabilities::new())
.list(
config
.source
.lists_sessions()
.then(SessionListCapabilities::new),
),
),
)
.auth_methods(Vec::new())
.agent_info(agent_client_protocol::schema::v1::Implementation::new(
"basis",
env!("CARGO_PKG_VERSION"),
))
}
fn announce_commands(
connection: &ConnectionTo<Client>,
session_id: &SessionId,
commands: Vec<AvailableCommand>,
) -> Result<(), Error> {
if commands.is_empty() {
return Ok(());
}
notify(
connection,
session_id,
SessionUpdate::AvailableCommandsUpdate(AvailableCommandsUpdate::new(commands)),
)
}
fn notify(
connection: &ConnectionTo<Client>,
session_id: &SessionId,
update: SessionUpdate,
) -> Result<(), Error> {
connection.send_notification(SessionNotification::new(session_id.clone(), update))
}
#[cfg(test)]
mod tests;