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::{FollowUp, FollowUpSource};
use crate::cli_runtime::stream_effect::StreamEffect;
use crate::presentation::{parent_message_card, MessageDelivery};
use crate::run_artifacts::AttachmentEvent;
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 }
}
fn encode(text: String) -> FollowUp {
let line = encode_user_turn(&frame_parent_message(&text));
FollowUp {
line,
written: Some(StreamEffect::Attachment(AttachmentEvent::Message(
Box::new(parent_message_card(
text,
MessageDelivery::Queued,
"written to Claude stdin; awaiting its next turn".into(),
)),
))),
}
}
}
impl FollowUpSource for ClaudeFollowUpSource {
fn try_recv(&mut self) -> Result<FollowUp, mpsc::error::TryRecvError> {
let text = self.inbox.try_recv()?;
Ok(Self::encode(text))
}
fn recv(&mut self) -> Pin<Box<dyn Future<Output = Option<FollowUp>> + Send + '_>> {
Box::pin(async {
let text = self.inbox.recv().await?;
Some(Self::encode(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;