use async_trait::async_trait;
use backbone_messaging::{EventError, IntegrationEventEnvelope, IntegrationEventHandler};
use backbone_outbox::inbox;
use chrono::NaiveDate;
use rust_decimal::Decimal;
use sqlx::PgPool;
use uuid::Uuid;
const CONSUMER: &str = "promotion.salary";
pub struct PromotionSalaryHandler {
pool: PgPool,
}
impl PromotionSalaryHandler {
fn rpool(&self) -> sqlx::PgPool {
crate::request_pool::current().unwrap_or_else(|| self.pool.clone())
}
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
}
#[async_trait]
impl IntegrationEventHandler for PromotionSalaryHandler {
async fn handle(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError> {
let event_id = Uuid::parse_str(&envelope.id)
.map_err(|e| handler_err(format!("bad envelope id '{}': {e}", envelope.id)))?;
let p = &envelope.payload;
let employee_id: Uuid = json_field(p, "employee_id")?;
let promotion_id: Option<Uuid> = serde_json::from_value(p["promotion_id"].clone()).ok();
let effective_date: Option<NaiveDate> = serde_json::from_value(p["effective_date"].clone()).ok();
let proposed_salary: Option<Decimal> = serde_json::from_value::<String>(p["proposed_salary"].clone())
.ok()
.and_then(|s| s.parse().ok());
let mut tx = self.rpool().begin().await.map_err(map_db)?;
if let Some(scope) = backbone_orm::org_scope::current_org_scope() {
backbone_orm::org_scope::bind_org_scope_on(&mut tx, &scope)
.await
.map_err(|e| handler_err(format!("org scope bind: {e}")))?;
}
let first_time = inbox::once(&mut *tx, "payroll", CONSUMER, event_id)
.await
.map_err(|e| handler_err(format!("inbox claim: {e}")))?;
if first_time {
if let Some(amount) = proposed_salary {
sqlx::query(
r#"INSERT INTO payroll.compensation_changes (employee_id, change_type, new_amount, effective_date,
reference_id, note,
org_unit_id)
VALUES ($1, 'promotion'::compensation_change_type, $2, $3, $4, $5, $6::uuid)"#,
)
.bind(employee_id)
.bind(amount)
.bind(effective_date)
.bind(promotion_id)
.bind("promotion.effective")
.bind(
backbone_orm::org_scope::current_org_scope()
.map(|s| s.acting_unit_id())
.or_else(|| {
p.get("company_id")
.and_then(|v| v.as_str())
.and_then(|v| v.parse().ok())
}),
)
.execute(&mut *tx)
.await
.map_err(map_db)?;
}
}
tx.commit().await.map_err(map_db)?;
Ok(())
}
fn event_patterns(&self) -> Vec<&'static str> {
vec!["promotion.effective"]
}
fn name(&self) -> &'static str {
"PromotionSalaryHandler"
}
}
fn json_field<T>(p: &serde_json::Value, field: &str) -> Result<T, EventError>
where
T: serde::de::DeserializeOwned,
{
serde_json::from_value(p[field].clone())
.map_err(|e| handler_err(format!("payload.{field}: {e}")))
}
fn map_db(e: sqlx::Error) -> EventError {
handler_err(format!("db: {e}"))
}
fn handler_err(message: String) -> EventError {
EventError::handler(CONSUMER, message)
}