use async_trait::async_trait;
use backbone_messaging::{EventError, IntegrationEventEnvelope, IntegrationEventHandler};
use backbone_outbox::inbox;
use chrono::Utc;
use rust_decimal::Decimal;
use sqlx::PgPool;
use uuid::Uuid;
const CONSUMER: &str = "onboarding.enroll";
#[async_trait]
pub trait OnboardingEnrollInputs: Send + Sync {
async fn starting_salary(&self, employee_id: Uuid) -> Result<Option<Decimal>, sqlx::Error>;
}
pub struct PoolOnboardingEnrollInputs {
pool: PgPool,
}
impl PoolOnboardingEnrollInputs {
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 OnboardingEnrollInputs for PoolOnboardingEnrollInputs {
async fn starting_salary(&self, employee_id: Uuid) -> Result<Option<Decimal>, sqlx::Error> {
let row: Option<(Option<Decimal>,)> = backbone_orm::company_scope::fetch_optional_scoped(
&self.pool,
sqlx::query_as(
r#"SELECT base_salary
FROM employee.employees
WHERE id = $1
AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(employee_id),
)
.await?;
Ok(row
.and_then(|(b,)| b)
.filter(|d| *d != Decimal::ZERO))
}
}
pub struct OnboardingEnrolledHandler {
pool: PgPool,
inputs: Box<dyn OnboardingEnrollInputs>,
}
impl OnboardingEnrolledHandler {
fn rpool(&self) -> sqlx::PgPool {
crate::request_pool::current().unwrap_or_else(|| self.pool.clone())
}
pub fn new(pool: PgPool) -> Self {
Self {
inputs: Box::new(PoolOnboardingEnrollInputs::new(pool.clone())),
pool,
}
}
}
#[async_trait]
impl IntegrationEventHandler for OnboardingEnrolledHandler {
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 onboarding_id: Option<Uuid> = serde_json::from_value(p["onboarding_id"].clone()).ok();
let base_salary = self
.inputs
.starting_salary(employee_id)
.await
.map_err(map_db)?;
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) = base_salary {
let period = Utc::now().format("%Y").to_string();
let note = format!("onboarding enrollment: initial compensation (period {period})");
sqlx::query(
r#"INSERT INTO payroll.compensation_changes (employee_id, change_type, new_amount, effective_date,
reference_id, note,
org_unit_id)
VALUES ($1, 'hire'::compensation_change_type, $2, $3, $4, $5, $6::uuid)"#,
)
.bind(employee_id)
.bind(amount)
.bind(Utc::now().date_naive())
.bind(onboarding_id)
.bind(¬e)
.bind(backbone_orm::org_scope::current_org_scope().map(|s| s.acting_unit_id()))
.execute(&mut *tx)
.await
.map_err(map_db)?;
}
else {
tracing::warn!(
target: "payroll.onboarding_enrolled",
employee_id = ?employee_id,
onboarding_id = ?onboarding_id,
"onboarding enrollment SKIPPED: no starting salary on the employee master — \
record base_salary and add the initial compensation change by hand"
);
}
}
tx.commit().await.map_err(map_db)?;
Ok(())
}
fn event_patterns(&self) -> Vec<&'static str> {
vec!["onboarding.completed"]
}
fn name(&self) -> &'static str {
"OnboardingEnrolledHandler"
}
}
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)
}