use std::sync::{Arc, Mutex};
use serde_json::json;
use tokio::sync::mpsc;
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) 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;