use std::sync::Arc;
use agent_client_protocol as acp;
use crate::agent::ZephAcpAgentState;
pub(crate) async fn handle_prompt(
req: acp::schema::v1::PromptRequest,
responder: acp::Responder<acp::schema::v1::PromptResponse>,
#[cfg_attr(not(feature = "unstable-cancel-request"), allow(unused_variables))]
cx: acp::ConnectionTo<acp::Client>,
state: Arc<ZephAcpAgentState>,
) -> acp::Result<()> {
#[cfg(feature = "unstable-cancel-request")]
let cancel_request_bridge = spawn_cancel_request_bridge(&req, &responder, &cx, &state);
let resp = state.do_prompt(req).await?;
#[cfg(feature = "unstable-cancel-request")]
drop(cancel_request_bridge);
responder.respond(resp)
}
#[cfg(feature = "unstable-cancel-request")]
fn spawn_cancel_request_bridge(
req: &acp::schema::v1::PromptRequest,
responder: &acp::Responder<acp::schema::v1::PromptResponse>,
cx: &acp::ConnectionTo<acp::Client>,
state: &Arc<ZephAcpAgentState>,
) -> tokio::sync::oneshot::Sender<()> {
let (done_tx, done_rx) = tokio::sync::oneshot::channel::<()>();
if let Some(cancel_signal) = state.session_cancel_signal(&req.session_id) {
let cancellation = responder.cancellation();
let spawn_result = cx.spawn(async move {
tokio::select! {
biased;
_ = done_rx => {}
() = cancellation.cancelled() => cancel_signal.notify_one(),
}
Ok(())
});
if let Err(e) = spawn_result {
tracing::debug!(error = %e, "failed to spawn $/cancel_request watcher for prompt");
}
}
done_tx
}