use uuid::Uuid;
use crate::domain::event::{record_channel, stage_bus_event};
use crate::infrastructure::persistence::message_pipeline_repository::{
MessagePipelineRepository, MintedNotification, NewMailMessageRow, NewMailNotificationRow,
NewMailQueueRow, NewSmsRow,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum NotificationChannel {
Inbox,
Email,
Sms,
}
impl NotificationChannel {
pub fn as_str(self) -> &'static str {
match self {
NotificationChannel::Inbox => "inbox",
NotificationChannel::Email => "email",
NotificationChannel::Sms => "sms",
}
}
}
#[derive(Debug, Clone)]
pub struct PostRecipient {
pub res_partner_id: Option<Uuid>,
pub channel: NotificationChannel,
pub email: Option<String>,
pub number: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct MessagePostCommand {
pub body: String,
pub subject: Option<String>,
pub message_type: String,
pub subtype_id: Option<Uuid>,
pub subtype_name: Option<String>,
pub is_internal: bool,
pub author_id: Option<Uuid>,
pub author_guest_id: Option<Uuid>,
pub email_from: Option<String>,
pub reply_to: Option<String>,
pub model: Option<String>,
pub res_id: Option<Uuid>,
pub record_name: Option<String>,
pub recipients: Vec<PostRecipient>,
}
#[derive(Debug, Clone, Default)]
pub struct PostedMessage {
pub message_id: Uuid,
pub subtype_id: Option<Uuid>,
pub notifications: Vec<MintedNotification>,
pub mail_id: Option<Uuid>,
pub sms: Vec<(Uuid, String)>,
}
#[derive(Debug, thiserror::Error)]
pub enum MailError {
#[error("db: {0}")]
Db(#[from] sqlx::Error),
#[error("invalid input: {0}")]
Invalid(String),
}
pub struct MessageWriteService {
pool: sqlx::PgPool,
}
impl MessageWriteService {
pub fn new(pool: sqlx::PgPool) -> Self {
Self { pool }
}
pub async fn message_post(&self, cmd: MessagePostCommand) -> Result<PostedMessage, MailError> {
if cmd.body.trim().is_empty() {
return Err(MailError::Invalid("message needs a body".into()));
}
if cmd.model.is_some() != cmd.res_id.is_some() {
return Err(MailError::Invalid(
"model and res_id travel together (the chatter edge is a pair)".into(),
));
}
let mut tx = self.pool.begin().await?;
let subtype_id = MessagePipelineRepository::resolve_subtype(
&mut tx, cmd.subtype_id, cmd.subtype_name.as_deref(),
)
.await?;
let message_id = Uuid::new_v4();
MessagePipelineRepository::insert_mail_message(
&mut tx,
&NewMailMessageRow {
id: message_id,
subject: cmd.subject.as_deref(),
body: &cmd.body,
message_type: if cmd.message_type.is_empty() { "comment" } else { &cmd.message_type },
subtype_id,
is_internal: cmd.is_internal,
author_id: cmd.author_id,
author_guest_id: cmd.author_guest_id,
email_from: cmd.email_from.as_deref(),
message_id: None,
reply_to: cmd.reply_to.as_deref(),
model: cmd.model.as_deref(),
res_id: cmd.res_id,
record_name: cmd.record_name.as_deref(),
},
)
.await?;
let ctx = PostContext {
message_id,
body: cmd.body.clone(),
model: cmd.model.clone(),
res_id: cmd.res_id,
reply_to: cmd.reply_to.clone(),
};
let mut posted = PostedMessage { message_id, subtype_id, ..Default::default() };
for (channel, recipients) in group_by_channel(cmd.recipients) {
let minted = dispatch_channel(&mut tx, &ctx, channel, &recipients).await?;
posted.notifications.extend(minted.notifications);
if minted.mail_id.is_some() {
posted.mail_id = minted.mail_id;
}
posted.sms.extend(minted.sms);
}
let channel_key = match (&ctx.model, ctx.res_id) {
(Some(m), Some(r)) => record_channel(m, r),
_ => format!("mail.message_{message_id}"),
};
stage_bus_event(
&mut tx,
"MessagePosted",
"MailMessage",
message_id,
channel_key.clone(),
"MessagePosted",
serde_json::json!({
"message_id": message_id,
"message_type": cmd.message_type,
"model": cmd.model,
"res_id": cmd.res_id,
"notification_count": posted.notifications.len(),
}),
)
.await?;
for (sms_id, sms_uuid) in &posted.sms {
stage_bus_event(
&mut tx,
"SmsCreated",
"Sms",
sms_id,
channel_key.clone(),
"SmsCreated",
serde_json::json!({ "sms_id": sms_id, "uuid": sms_uuid }),
)
.await?;
}
tx.commit().await?;
Ok(posted)
}
}
struct PostContext {
message_id: Uuid,
body: String,
model: Option<String>,
res_id: Option<Uuid>,
reply_to: Option<String>,
}
struct ChannelMint {
notifications: Vec<MintedNotification>,
mail_id: Option<Uuid>,
sms: Vec<(Uuid, String)>,
}
fn group_by_channel(recipients: Vec<PostRecipient>) -> Vec<(NotificationChannel, Vec<PostRecipient>)> {
let mut groups: Vec<(NotificationChannel, Vec<PostRecipient>)> = Vec::new();
for r in recipients {
match groups.iter_mut().find(|(c, _)| *c == r.channel) {
Some((_, list)) => list.push(r),
None => groups.push((r.channel, vec![r])),
}
}
groups
}
async fn dispatch_channel(
tx: &mut sqlx::PgConnection,
ctx: &PostContext,
channel: NotificationChannel,
recipients: &[PostRecipient],
) -> Result<ChannelMint, MailError> {
match channel {
NotificationChannel::Inbox => inbox_strategy(tx, ctx, recipients).await,
NotificationChannel::Email => email_strategy(tx, ctx, recipients).await,
NotificationChannel::Sms => sms_strategy(tx, ctx, recipients).await,
}
}
async fn inbox_strategy(
tx: &mut sqlx::PgConnection,
ctx: &PostContext,
recipients: &[PostRecipient],
) -> Result<ChannelMint, MailError> {
let mut notifications = Vec::new();
for r in recipients {
let id = Uuid::new_v4();
MessagePipelineRepository::insert_mail_notification(
tx,
&NewMailNotificationRow {
id,
mail_message_id: ctx.message_id,
res_partner_id: r.res_partner_id,
notification_type: "inbox",
notification_status: "sent",
mail_mail_id_int: None,
},
)
.await?;
notifications.push(MintedNotification {
id,
res_partner_id: r.res_partner_id,
notification_type: NotificationChannel::Inbox,
notification_status: "sent".into(),
});
}
Ok(ChannelMint { notifications, mail_id: None, sms: Vec::new() })
}
async fn email_strategy(
tx: &mut sqlx::PgConnection,
ctx: &PostContext,
recipients: &[PostRecipient],
) -> Result<ChannelMint, MailError> {
let mut addresses = Vec::new();
for r in recipients {
let Some(email) = r.email.as_deref() else {
return Err(MailError::Invalid(format!(
"email-channel recipient {:?} has no email address",
r.res_partner_id
)));
};
addresses.push(email.to_string());
}
if addresses.is_empty() {
return Ok(ChannelMint { notifications: Vec::new(), mail_id: None, sms: Vec::new() });
}
let mail_id = Uuid::new_v4();
MessagePipelineRepository::insert_mail(
tx,
&NewMailQueueRow {
id: mail_id,
mail_message_id: ctx.message_id,
email_to: &addresses.join(", "),
email_cc: None,
reply_to: ctx.reply_to.as_deref(),
scheduled_date: None,
},
)
.await?;
let mut notifications = Vec::new();
for r in recipients {
let id = Uuid::new_v4();
MessagePipelineRepository::insert_mail_notification(
tx,
&NewMailNotificationRow {
id,
mail_message_id: ctx.message_id,
res_partner_id: r.res_partner_id,
notification_type: "email",
notification_status: "ready",
mail_mail_id_int: Some(mail_id),
},
)
.await?;
notifications.push(MintedNotification {
id,
res_partner_id: r.res_partner_id,
notification_type: NotificationChannel::Email,
notification_status: "ready".into(),
});
}
Ok(ChannelMint { notifications, mail_id: Some(mail_id), sms: Vec::new() })
}
async fn sms_strategy(
tx: &mut sqlx::PgConnection,
ctx: &PostContext,
recipients: &[PostRecipient],
) -> Result<ChannelMint, MailError> {
let mut notifications = Vec::new();
let mut sms = Vec::new();
for r in recipients {
let Some(number) = r.number.as_deref() else {
return Err(MailError::Invalid(format!(
"sms-channel recipient {:?} has no number",
r.res_partner_id
)));
};
let notification_id = Uuid::new_v4();
MessagePipelineRepository::insert_mail_notification(
tx,
&NewMailNotificationRow {
id: notification_id,
mail_message_id: ctx.message_id,
res_partner_id: r.res_partner_id,
notification_type: "sms",
notification_status: "ready",
mail_mail_id_int: None,
},
)
.await?;
notifications.push(MintedNotification {
id: notification_id,
res_partner_id: r.res_partner_id,
notification_type: NotificationChannel::Sms,
notification_status: "ready".into(),
});
let sms_id = Uuid::new_v4();
let sms_uuid = Uuid::new_v4().simple().to_string();
MessagePipelineRepository::insert_sms_with_tracker(
tx,
&NewSmsRow {
id: sms_id,
uuid: sms_uuid.clone(),
number,
body: &ctx.body,
mail_message_id: Some(ctx.message_id),
notification_id: Some(notification_id),
},
)
.await?;
sms.push((sms_id, sms_uuid));
}
Ok(ChannelMint { notifications, mail_id: None, sms })
}