use std::sync::Arc;
use uuid::Uuid;
use crate::application::service::chatter_acl::{MessagingIdentity, ThreadAccessResolver};
use crate::domain::event::constants::partner_channel;
use crate::infrastructure::persistence::chatter_repository::ChatterRepository;
use crate::infrastructure::persistence::follower_repository::FollowerRepository;
#[derive(Debug, Clone, serde::Serialize)]
pub struct SuggestedRecipient {
pub partner_id: Uuid,
pub reason: &'static str,
}
#[derive(Debug, thiserror::Error)]
pub enum RecipientQueryError {
#[error("db: {0}")]
Db(#[from] sqlx::Error),
#[error("forbidden: no read access to {0} {1}")]
Forbidden(String, Uuid),
#[error("recipients need a partner identity")]
NeedsPartner,
}
pub struct RecipientQueryService {
pool: sqlx::PgPool,
thread_acl: Arc<dyn ThreadAccessResolver>,
}
impl RecipientQueryService {
pub fn new(pool: sqlx::PgPool, thread_acl: Arc<dyn ThreadAccessResolver>) -> Self {
Self { pool, thread_acl }
}
pub async fn suggested(
&self,
identity: &MessagingIdentity,
model: &str,
res_id: Uuid,
) -> Result<Vec<SuggestedRecipient>, RecipientQueryError> {
if !self.thread_acl.can_read(&self.pool, identity, model, res_id).await {
return Err(RecipientQueryError::Forbidden(model.into(), res_id));
}
let mut conn = self.pool.acquire().await?;
let followers = FollowerRepository::list_followers(&mut conn, model, res_id).await?;
let authors = ChatterRepository::thread_authors(&mut conn, model, res_id).await?;
let mut out: Vec<SuggestedRecipient> = followers
.iter()
.map(|(partner_id, _)| SuggestedRecipient { partner_id: *partner_id, reason: "follower" })
.collect();
for author in authors {
if !followers.iter().any(|(pid, _)| *pid == author) {
out.push(SuggestedRecipient { partner_id: author, reason: "author" });
}
}
Ok(out)
}
pub fn partner_channel_of(partner_id: Uuid) -> String {
partner_channel(partner_id)
}
}