use uuid::Uuid;
use crate::application::service::chatter_acl::MessagingIdentity;
use crate::domain::event::constants::{discuss_channel, stage_bus_event};
use crate::infrastructure::persistence::channel_repository::{
ChannelRepository, NewChannelRow,
};
#[derive(Debug, thiserror::Error)]
pub enum ChannelError {
#[error("db: {0}")]
Db(#[from] sqlx::Error),
#[error("invalid input: {0}")]
Invalid(String),
#[error("not found: channel {0}")]
NotFound(Uuid),
}
#[derive(Debug, Clone, Default)]
pub struct CreateChannelCommand {
pub name: Option<String>,
pub channel_type: String,
pub uuid: Option<String>,
pub default_access_mode: Option<String>,
pub email_send: bool,
pub creator: Option<MessagingIdentity>,
}
pub struct ChannelWriteService {
pool: sqlx::PgPool,
}
impl ChannelWriteService {
pub fn new(pool: sqlx::PgPool) -> Self {
Self { pool }
}
pub async fn create(&self, cmd: CreateChannelCommand) -> Result<Uuid, ChannelError> {
let channel_type = match cmd.channel_type.as_str() {
"" | "channel" => "channel",
"chat" => "chat",
"group" => "group",
other => {
return Err(ChannelError::Invalid(format!(
"unknown channel_type {other:?}"
)))
}
};
if channel_type == "chat" {
return Err(ChannelError::Invalid(
"chats are minted via get_or_create_chat, not create (1:1 dedup)".into(),
));
}
let mut tx = self.pool.begin().await?;
let channel_id = Uuid::new_v4();
ChannelRepository::insert_channel(
&mut tx,
&NewChannelRow {
id: channel_id,
name: cmd.name.as_deref(),
channel_type,
uuid: cmd.uuid.as_deref(),
default_access_mode: cmd.default_access_mode.as_deref(),
email_send: cmd.email_send,
},
)
.await?;
if let Some(creator) = &cmd.creator {
ChannelRepository::insert_creator_member(&mut tx, channel_id, creator).await?;
}
stage_bus_event(
&mut tx,
"DiscussChannelCreated",
"DiscussChannel",
channel_id,
discuss_channel(channel_id),
"discuss.channel/insert",
serde_json::json!({
"channel_id": channel_id,
"channel_type": channel_type,
"name": cmd.name,
}),
)
.await?;
tx.commit().await?;
Ok(channel_id)
}
pub async fn get_or_create_chat(
&self,
partner_a: Uuid,
partner_b: Uuid,
) -> Result<(Uuid, bool), ChannelError> {
if partner_a == partner_b {
return Err(ChannelError::Invalid(
"a 1:1 chat needs two distinct partners (Odoo allows self-chat via a different path; not ported)".into(),
));
}
let set = [partner_a, partner_b];
let mut tx = self.pool.begin().await?;
if let Some(id) = ChannelRepository::find_chat_by_member_set(&mut tx, &set).await? {
tx.commit().await?;
return Ok((id, false));
}
tx.commit().await?;
let mut tx = self.pool.begin().await?;
if let Some(id) = ChannelRepository::find_chat_by_member_set(&mut tx, &set).await? {
tx.commit().await?;
return Ok((id, false));
}
let channel_id = Uuid::new_v4();
ChannelRepository::insert_channel(
&mut tx,
&NewChannelRow {
id: channel_id,
name: None, channel_type: "chat",
uuid: None,
default_access_mode: None,
email_send: false,
},
)
.await?;
ChannelRepository::insert_creator_member(
&mut tx,
channel_id,
&MessagingIdentity::User { partner_id: partner_a },
)
.await?;
ChannelRepository::insert_creator_member(
&mut tx,
channel_id,
&MessagingIdentity::User { partner_id: partner_b },
)
.await?;
stage_bus_event(
&mut tx,
"DiscussChannelCreated",
"DiscussChannel",
channel_id,
discuss_channel(channel_id),
"discuss.channel/insert",
serde_json::json!({
"channel_id": channel_id,
"channel_type": "chat",
"members": set,
}),
)
.await?;
tx.commit().await?;
Ok((channel_id, true))
}
pub async fn update_fields(
&self,
channel_id: Uuid,
name: Option<&str>,
description: Option<&str>,
email_send: Option<bool>,
) -> Result<(), ChannelError> {
if name.is_none() && description.is_none() && email_send.is_none() {
return Err(ChannelError::Invalid("nothing to update".into()));
}
let mut tx = self.pool.begin().await?;
Self::must_exist(&mut tx, channel_id).await?;
ChannelRepository::update_channel_fields(&mut tx, channel_id, name, description, email_send)
.await?;
stage_bus_event(
&mut tx,
"DiscussChannelUpdated",
"DiscussChannel",
channel_id,
discuss_channel(channel_id),
"discuss.channel/update",
serde_json::json!({
"channel_id": channel_id,
"name": name,
"description": description,
"email_send": email_send,
}),
)
.await?;
tx.commit().await?;
Ok(())
}
pub async fn set_message_pinned(
&self,
channel_id: Uuid,
message_id: Uuid,
pinned: bool,
) -> Result<(), ChannelError> {
let mut tx = self.pool.begin().await?;
Self::must_exist(&mut tx, channel_id).await?;
let owner: Option<(Option<String>, Option<Uuid>)> = sqlx::query_as(
r#"SELECT model, res_id FROM messaging.mail_messages WHERE id = $1"#,
)
.bind(message_id)
.fetch_optional(&mut *tx)
.await?;
match owner {
Some((model, res_id))
if model.as_deref() == Some("discuss.channel") && res_id == Some(channel_id) => {}
_ => {
return Err(ChannelError::Invalid(format!(
"message {message_id} does not belong to channel {channel_id}"
)))
}
}
ChannelRepository::set_message_pinned(&mut tx, message_id, pinned).await?;
stage_bus_event(
&mut tx,
"MessagePinnedChanged",
"MailMessage",
message_id,
discuss_channel(channel_id),
if pinned { "mail.message/pin" } else { "mail.message/unpin" },
serde_json::json!({
"channel_id": channel_id,
"message_id": message_id,
"pinned": pinned,
}),
)
.await?;
tx.commit().await?;
Ok(())
}
pub async fn archive(&self, channel_id: Uuid) -> Result<(), ChannelError> {
let mut tx = self.pool.begin().await?;
Self::must_exist(&mut tx, channel_id).await?;
ChannelRepository::soft_delete_channel(&mut tx, channel_id).await?;
stage_bus_event(
&mut tx,
"DiscussChannelDeleted",
"DiscussChannel",
channel_id,
discuss_channel(channel_id),
"discuss.channel/delete",
serde_json::json!({ "channel_id": channel_id }),
)
.await?;
tx.commit().await?;
Ok(())
}
async fn must_exist(tx: &mut sqlx::PgConnection, channel_id: Uuid) -> Result<(), ChannelError> {
ChannelRepository::channel_type(tx, channel_id)
.await?
.map(|_| ())
.ok_or(ChannelError::NotFound(channel_id))
}
}