use std::future::Future;
use std::pin::Pin;
use std::sync::{Arc, Mutex};
use serde_json::json;
use tokio::sync::mpsc;
use crate::cli_runtime::drain::FollowUpSource;
pub(crate) const PARENT_MESSAGE_QUEUE_CAPACITY: usize = 8;
#[derive(Clone, Debug)]
pub(crate) struct ClaudeMessageHandle {
gate: Arc<Mutex<Option<mpsc::Sender<String>>>>,
}
impl ClaudeMessageHandle {
pub(crate) async fn send(&self, text: String) -> Result<(), ClaudeMessageSendError> {
let sender = {
let guard = self.gate.lock().expect("claude message gate");
guard.clone().ok_or(ClaudeMessageSendError::Closed)?
};
sender
.send(text)
.await
.map_err(|_| ClaudeMessageSendError::Closed)
}
#[cfg(test)]
fn clone_sender_for_test(&self) -> Option<mpsc::Sender<String>> {
self.gate.lock().expect("claude message gate").clone()
}
}
pub(crate) struct ClaudeMessageInbox {
gate: Arc<Mutex<Option<mpsc::Sender<String>>>>,
receiver: mpsc::Receiver<String>,
}
impl ClaudeMessageInbox {
pub(crate) fn seal(&self) {
let mut guard = self.gate.lock().expect("claude message gate");
*guard = None;
}
pub(crate) fn try_recv(&mut self) -> Result<String, mpsc::error::TryRecvError> {
self.receiver.try_recv()
}
pub(crate) async fn recv(&mut self) -> Option<String> {
self.receiver.recv().await
}
}
#[derive(Clone, Debug, PartialEq, Eq, thiserror::Error)]
pub(crate) enum ClaudeMessageSendError {
#[error("delegated Claude run is no longer accepting parent messages")]
Closed,
}
pub(crate) fn message_channel() -> (ClaudeMessageHandle, ClaudeMessageInbox) {
let (sender, receiver) = mpsc::channel(PARENT_MESSAGE_QUEUE_CAPACITY);
let gate = Arc::new(Mutex::new(Some(sender)));
(
ClaudeMessageHandle {
gate: Arc::clone(&gate),
},
ClaudeMessageInbox { gate, receiver },
)
}
pub(crate) struct ClaudeFollowUpSource {
inbox: ClaudeMessageInbox,
}
impl ClaudeFollowUpSource {
pub(crate) fn new(inbox: ClaudeMessageInbox) -> Self {
Self { inbox }
}
}
impl FollowUpSource for ClaudeFollowUpSource {
fn try_recv(&mut self) -> Result<String, mpsc::error::TryRecvError> {
self.inbox
.try_recv()
.map(|text| encode_user_turn(&frame_parent_message(&text)))
}
fn recv(&mut self) -> Pin<Box<dyn Future<Output = Option<String>> + Send + '_>> {
Box::pin(async {
self.inbox
.recv()
.await
.map(|text| encode_user_turn(&frame_parent_message(&text)))
})
}
fn seal(&self) {
self.inbox.seal();
}
}
pub(crate) fn encode_user_turn(text: &str) -> String {
let mut line = json!({
"type": "user",
"message": {
"role": "user",
"content": text,
},
})
.to_string();
line.push('\n');
line
}
pub(crate) fn frame_parent_message(text: &str) -> String {
format!(
"Message from the parent session (not a new task - incorporate this into your current work):\n\n{text}"
)
}
#[cfg(test)]
#[path = "messaging_tests.rs"]
mod tests;