use crate::cli_runtime::parent_messages::{frame_parent_message, ParentMessageInbox};
use std::future::Future;
use std::pin::Pin;
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, NotificationDelivery};
use crate::run_artifacts::AttachmentEvent;
pub(crate) struct ClaudeFollowUpSource {
inbox: ParentMessageInbox,
}
impl ClaudeFollowUpSource {
pub(crate) fn new(inbox: ParentMessageInbox) -> 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,
NotificationDelivery::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
}
#[cfg(test)]
#[path = "messaging_tests.rs"]
mod tests;