pub(crate) mod chat_draft;
pub(crate) mod chat_history;
mod enrichment;
pub(crate) mod reply;
pub mod telegram;
mod telegram_group;
use enrichment::has_inbound_temp_marker;
pub use enrichment::{enrich_links, enrich_message};
pub use reply::{ReplyReference, apply_reply_marker};
pub use telegram::mirror_gui_message_to_telegram;
pub use telegram_group::compose_group_content;
use crate::channels::chat_history::ChatHistoryInsert;
use crate::db;
use crate::util::UnwrapPoison;
use crate::{Channel, ChannelMessage, ChatDirection, SendMessage};
use async_trait::async_trait;
#[cfg(target_os = "macos")]
use std::sync::Arc;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
const CHANNEL_TYPING_REFRESH_INTERVAL_SECS: u64 = 4;
#[derive(Debug, Clone)]
struct BroadcastPersistEntry {
user_name: String,
channel: String,
content: String,
direction: ChatDirection,
agent_role: Option<String>,
broadcast_id: Option<String>,
workspace: String,
optimistic_id: Option<String>,
reply_reference: Option<crate::channels::ReplyReference>,
}
impl BroadcastPersistEntry {
async fn broadcast_and_persist(self) {
debug_assert!(
self.direction != ChatDirection::Agent || self.agent_role.is_some(),
"BroadcastPersistEntry: direction=Agent but agent_role is None"
);
let message_id = crate::generate_id();
let timestamp = db::now();
let db_direction = match self.direction {
ChatDirection::Agent => "agent".to_string(),
ChatDirection::User => "user".to_string(),
ChatDirection::Divider => {
unreachable!("Divider markers should not go through broadcast_and_persist")
}
};
broadcast_chat_event(
&message_id,
&self.user_name,
&self.content,
self.direction,
&self.channel,
self.agent_role.clone(),
&self.workspace,
self.optimistic_id.clone(),
self.reply_reference.clone(),
×tamp,
false,
);
let store = crate::channels::chat_history::store();
let _ = store
.insert(&ChatHistoryInsert {
message_id,
user_name: self.user_name,
direction: db_direction,
content: self.content,
agent_role: self.agent_role,
broadcast_id: self.broadcast_id,
workspace: self.workspace,
timestamp: Some(timestamp),
reply_reference: self.reply_reference,
})
.await;
}
}
pub(crate) async fn broadcast_and_persist_agent_response(
user_name: &str,
channel: &str,
content: &str,
agent_role: Option<String>,
workspace: &str,
broadcast_id: Option<String>,
) {
BroadcastPersistEntry {
user_name: user_name.to_string(),
channel: channel.to_string(),
content: content.to_string(),
direction: ChatDirection::Agent,
agent_role, broadcast_id,
workspace: workspace.to_string(),
optimistic_id: None, reply_reference: None, }
.broadcast_and_persist()
.await;
}
#[expect(clippy::too_many_arguments)]
fn broadcast_chat_event(
message_id: &str,
user_name: &str,
content: &str,
direction: ChatDirection,
channel: &str,
agent_role: Option<String>,
workspace: &str,
optimistic_id: Option<String>,
reply_reference: Option<crate::channels::ReplyReference>,
timestamp: &str,
transient: bool,
) {
use crate::ChatEvent;
if let Some(tx) = crate::CHAT_BROADCAST.get() {
let _ = tx.send(ChatEvent::Message {
message_id: message_id.to_string(),
user_name: user_name.to_string(),
content: content.to_string(),
direction,
timestamp: timestamp.to_string(),
channel: channel.to_string(),
agent_role,
workspace: workspace.to_string(),
optimistic_id,
transient,
reply_reference,
});
}
}
#[expect(clippy::too_many_arguments)]
pub(crate) fn broadcast_transient_event(
message_id: &str,
user_name: &str,
content: &str,
direction: ChatDirection,
channel: &str,
agent_role: Option<String>,
workspace: &str,
optimistic_id: Option<String>,
reply_reference: Option<crate::channels::ReplyReference>,
timestamp: &str,
) {
broadcast_chat_event(
message_id,
user_name,
content,
direction,
channel,
agent_role,
workspace,
optimistic_id,
reply_reference,
timestamp,
true,
);
}
#[must_use]
pub fn persist_content(original: &str, enriched: &str) -> String {
if has_inbound_temp_marker(original) {
crate::util::strip_data_uris(enriched)
} else {
original.to_string()
}
}
pub async fn broadcast_and_persist_incoming_message(
msg: &ChannelMessage,
broadcast_content: &str,
content_for_history: &str,
) {
let message_id = crate::generate_id();
let timestamp = db::now();
broadcast_chat_event(
&message_id,
&msg.user_name,
broadcast_content,
ChatDirection::User,
&msg.channel,
None,
&msg.workspace,
msg.optimistic_id.clone(),
msg.reply_reference.clone(),
×tamp,
false,
);
tokio::join!(
async {
let store = crate::channels::chat_history::store();
let _ = store
.insert(&ChatHistoryInsert {
message_id,
user_name: msg.user_name.clone(),
direction: "user".to_string(),
content: content_for_history.to_string(),
agent_role: None,
broadcast_id: None,
workspace: msg.workspace.clone(),
timestamp: Some(timestamp),
reply_reference: msg.reply_reference.clone(),
})
.await;
},
async {
let mut mirror_msg = msg.clone();
mirror_msg.content = content_for_history.to_string();
mirror_gui_message_to_telegram(&mirror_msg).await;
},
);
}
#[must_use]
pub fn spawn_scoped_typing_task(
recipient: String,
channel: String,
cancellation_token: CancellationToken,
) -> tokio::task::JoinHandle<()> {
let refresh_interval = std::time::Duration::from_secs(CHANNEL_TYPING_REFRESH_INTERVAL_SECS);
tokio::spawn(async move {
let Some(ch) = crate::channel_registry().get(&channel) else {
tracing::warn!(
channel = %channel,
"Channel not found in registry — skipping typing indicator"
);
return;
};
let mut interval = tokio::time::interval(refresh_interval);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
() = cancellation_token.cancelled() => break,
_ = interval.tick() => {
if let Err(e) = ch.start_typing(&recipient).await {
tracing::debug!("Failed to start typing on {}: {e}", ch.name());
}
}
}
}
})
}
pub async fn stop_typing(handle: tokio::task::JoinHandle<()>) {
if let Err(error) = handle.await {
tracing::error!("Typing task crashed: {error}");
}
}
pub struct GuiChannel {
gui_rx: std::sync::Mutex<Option<mpsc::UnboundedReceiver<ChannelMessage>>>,
}
impl GuiChannel {
#[must_use]
pub fn new() -> (Self, mpsc::UnboundedSender<ChannelMessage>) {
let (gui_tx, gui_rx) = mpsc::unbounded_channel::<ChannelMessage>();
let channel = Self {
gui_rx: std::sync::Mutex::new(Some(gui_rx)),
};
(channel, gui_tx)
}
}
#[async_trait]
impl Channel for GuiChannel {
async fn send(&self, _message: &SendMessage) -> anyhow::Result<()> {
Ok(())
}
async fn listen(&self, tx: mpsc::Sender<ChannelMessage>) -> anyhow::Result<()> {
let mut gui_rx = self
.gui_rx
.lock()
.unwrap_poison()
.take()
.expect("GuiChannel::listen() called twice");
while let Some(msg) = gui_rx.recv().await {
if tx.send(msg).await.is_err() {
tracing::info!("GuiChannel: pipeline closed — shutting down listener");
break;
}
}
tracing::info!("GuiChannel: listener stopped");
Ok(())
}
fn name(&self) -> &'static str {
"gui"
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
fn resolve_recipient(&self, user_name: &str, _reply_target: &str) -> Option<String> {
Some(user_name.to_string())
}
}
#[cfg(target_os = "macos")]
struct VoiceChannel;
#[cfg(target_os = "macos")]
#[async_trait]
impl Channel for VoiceChannel {
async fn send(&self, _message: &SendMessage) -> anyhow::Result<()> {
Ok(())
}
async fn listen(&self, _tx: mpsc::Sender<ChannelMessage>) -> anyhow::Result<()> {
Ok(())
}
fn name(&self) -> &'static str {
"voice"
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
#[cfg(target_os = "macos")]
pub fn register_global() {
let channel: Arc<dyn Channel> = Arc::new(VoiceChannel);
crate::channel_registry().register(channel);
}
#[cfg(target_os = "macos")]
#[cfg(test)]
mod tests {
use super::*;
use crate::{CHANNEL_REGISTRY, ChannelRegistry, channel_registry};
#[test]
fn test_voice_channel_registration() {
CHANNEL_REGISTRY.get_or_init(ChannelRegistry::default);
register_global();
let found = channel_registry().get("voice");
assert!(
found.is_some(),
"VoiceChannel should be findable by 'voice' name"
);
}
}