use async_trait::async_trait;
use backbone_messaging::{EventError, IntegrationEventEnvelope, IntegrationEventHandler};
use std::sync::Arc;
use uuid::Uuid;
use crate::application::service::digest_write_service::DigestWriteService;
use crate::application::service::engagement_port::{RecipientContextPort, RecipientContextSlot};
const HANDLER: &str = "digest.user_created";
pub struct UserCreatedHandler {
write: Arc<DigestWriteService>,
port: RecipientContextSlot,
default_digest: Option<Uuid>,
}
impl UserCreatedHandler {
pub fn new(
write: Arc<DigestWriteService>,
port: RecipientContextSlot,
default_digest: Option<Uuid>,
) -> Self {
Self { write, port, default_digest }
}
pub const EVENT_TYPE: &'static str = "sapiens.user.created";
}
#[async_trait]
impl IntegrationEventHandler for UserCreatedHandler {
async fn handle(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError> {
let user_id: Uuid = serde_json::from_value(envelope.payload["user_id"].clone())
.map_err(|e| {
EventError::handler(
HANDLER,
format!("payload.user_id: {e} (envelope {})", envelope.id),
)
})?;
let Some(default_digest) = self.default_digest else {
tracing::warn!(
target: "digest::audit",
event = "auto_subscribe_refused",
reason = "default_digest_not_configured",
user_id = %user_id,
envelope_id = %envelope.id
);
return Ok(());
};
let ctx = match self.port.resolve(user_id).await {
Ok(ctx) => ctx,
Err(e) if e.is_not_wired() => {
return Err(EventError::handler(
HANDLER,
format!("recipient port not wired; cannot resolve {user_id}: {e}"),
));
}
Err(e) => {
return Err(EventError::handler(
HANDLER,
format!("recipient {user_id} unresolvable: {e}"),
));
}
};
if !ctx.is_internal {
tracing::warn!(
target: "digest::audit",
event = "auto_subscribe_refused",
reason = "non_internal",
user_id = %user_id,
envelope_id = %envelope.id
);
return Ok(());
}
let first_time = self
.write
.subscribe_user(
default_digest,
user_id,
Some(serde_json::json!({
"auto_subscribed": true,
"auto_subscribed_at": chrono::Utc::now().to_rfc3339(),
"auto_subscribed_envelope": envelope.id,
})),
)
.await
.map_err(|e| EventError::handler(HANDLER, format!("subscribe verb: {e}")))?;
if first_time {
tracing::info!(
target: "digest::audit",
event = "digest_user_auto_subscribed",
digest_id = %default_digest,
user_id = %user_id,
envelope_id = %envelope.id
);
}
Ok(())
}
fn event_patterns(&self) -> Vec<&'static str> {
vec![Self::EVENT_TYPE]
}
fn name(&self) -> &'static str {
"DigestUserCreatedHandler"
}
}