backbone_payroll/presentation/http/
compensation_change_handler.rs1use std::collections::HashMap;
8use std::sync::Arc;
9
10use axum::Router;
11use serde::{Deserialize, Serialize};
12use uuid::Uuid;
13use chrono::{NaiveDate};
14use rust_decimal::Decimal;
15
16use backbone_core::http::BackboneCrudHandler;
18
19#[cfg(feature = "auth")]
21use backbone_auth::middleware::AuthContext;
22#[cfg(feature = "auth")]
23use backbone_auth::AuthMiddleware;
24
25use crate::domain::entity::*;
27use crate::application::service::{CompensationChangeService, ServiceError};
28
29use crate::presentation::dto::{CreateCompensationChangeDto, UpdateCompensationChangeDto, PatchCompensationChangeDto, CompensationChangeResponseDto};
31
32
33#[derive(Debug, thiserror::Error)]
35pub enum CompensationChangeError {
36 #[error("Not found: {0}")]
37 NotFound(String),
38 #[error("Validation error: {0}")]
39 Validation(String),
40 #[error("Database error: {0}")]
41 Database(String),
42 #[error("Internal error: {0}")]
43 Internal(String),
44}
45
46impl From<ServiceError> for CompensationChangeError {
47 fn from(err: ServiceError) -> Self {
48 match err {
49 ServiceError::NotFound => Self::NotFound(err.to_string()),
50 ServiceError::Validation(ref msg) => Self::Validation(msg.clone()),
51 ServiceError::AlreadyExists(ref msg) => Self::Validation(msg.clone()),
52 ServiceError::Repository(ref e) => Self::Database(e.to_string()),
53 ServiceError::Internal(ref msg) => Self::Internal(msg.clone()),
54 }
55 }
56}
57
58impl axum::response::IntoResponse for CompensationChangeError {
59 fn into_response(self) -> axum::response::Response {
60 use axum::http::StatusCode;
61 use axum::Json;
62
63 let (status, code) = match &self {
64 Self::NotFound(_) => (StatusCode::NOT_FOUND, "COMPENSATIONCHANGE_NOT_FOUND"),
65 Self::Validation(_) => (StatusCode::BAD_REQUEST, "COMPENSATIONCHANGE_VALIDATION_ERROR"),
66 Self::Database(_) => (StatusCode::INTERNAL_SERVER_ERROR, "COMPENSATIONCHANGE_DATABASE_ERROR"),
67 Self::Internal(_) => (StatusCode::INTERNAL_SERVER_ERROR, "COMPENSATIONCHANGE_INTERNAL_ERROR"),
68 };
69
70 let body = serde_json::json!({
71 "success": false,
72 "error": code,
73 "message": self.to_string(),
74 });
75
76 (status, Json(body)).into_response()
77 }
78}
79
80pub fn create_compensation_change_routes(service: Arc<CompensationChangeService>) -> Router {
113 BackboneCrudHandler::<CompensationChangeService, CompensationChange, CreateCompensationChangeDto, UpdateCompensationChangeDto, CompensationChangeResponseDto>::routes(
114 service,
115 "/compensation_changes",
116 )
117}
118
119pub fn create_compensation_change_read_routes(service: Arc<CompensationChangeService>) -> Router {
125 BackboneCrudHandler::<CompensationChangeService, CompensationChange, CreateCompensationChangeDto, UpdateCompensationChangeDto, CompensationChangeResponseDto>::read_routes(
126 service,
127 "/compensation_changes",
128 )
129}
130
131pub fn create_compensation_change_write_routes(service: Arc<CompensationChangeService>) -> Router {
143 BackboneCrudHandler::<CompensationChangeService, CompensationChange, CreateCompensationChangeDto, UpdateCompensationChangeDto, CompensationChangeResponseDto>::write_routes(
144 service,
145 "/compensation_changes",
146 )
147}
148
149#[cfg(feature = "auth")]
155pub fn create_protected_compensation_change_routes<A: AuthMiddleware + Send + Sync + 'static>(
156 service: Arc<CompensationChangeService>,
157 auth: Arc<A>,
158) -> Router {
159 use axum::middleware;
160 use axum::response::IntoResponse;
161
162 let auth_layer = auth.clone();
163 create_compensation_change_routes(service)
164 .layer(middleware::from_fn(move |mut req: axum::extract::Request, next: axum::middleware::Next| {
165 let auth = auth_layer.clone();
166 async move {
167 let token = req.headers()
168 .get(axum::http::header::AUTHORIZATION)
169 .and_then(|h| h.to_str().ok())
170 .and_then(|raw| raw.strip_prefix("Bearer ").or_else(|| raw.strip_prefix("bearer ")))
171 .unwrap_or("");
172 match auth.authenticate(token).await {
173 Ok(ctx) => {
174 req.extensions_mut().insert(ctx);
175 next.run(req).await
176 }
177 Err(_) => {
178 (axum::http::StatusCode::UNAUTHORIZED,
179 axum::Json(serde_json::json!({
180 "success": false,
181 "error": "unauthorized",
182 "message": "Authentication required"
183 }))
184 ).into_response()
185 }
186 }
187 }
188 }))
189}
190
191#[derive(serde::Deserialize)]
193pub struct CompensationChangeHistoryQuery {
194 pub as_of: Option<chrono::DateTime<chrono::Utc>>,
196}
197
198pub async fn compensation_change_history(
205 axum::extract::State(pool): axum::extract::State<sqlx::PgPool>,
206 axum::extract::Path(id): axum::extract::Path<uuid::Uuid>,
207 axum::extract::Query(q): axum::extract::Query<CompensationChangeHistoryQuery>,
208) -> axum::response::Response {
209 use axum::http::StatusCode;
210 use axum::response::IntoResponse;
211 use serde_json::json;
212 use sqlx::Row;
213 let subject_id = id.to_string();
214 let mut conn = match pool.acquire().await {
215 Ok(c) => c,
216 Err(e) => {
217 return (
218 StatusCode::SERVICE_UNAVAILABLE,
219 axum::Json(json!({
220 "error": "history_pool_unavailable",
221 "message": e.to_string(),
222 })),
223 )
224 .into_response()
225 }
226 };
227 if let Some(scope) = backbone_orm::org_scope::current_org_scope() {
230 if let Err(e) = backbone_orm::org_scope::bind_org_scope_on(&mut conn, &scope).await {
231 return (
232 StatusCode::INTERNAL_SERVER_ERROR,
233 axum::Json(json!({
234 "error": "history_scope_bind",
235 "message": e.to_string(),
236 })),
237 )
238 .into_response()
239 }
240 }
241 const SUBJECT_TYPE: &str = "payroll.compensation_changes";
242 match q.as_of {
243 None => {
244 let rows = match sqlx::query("SELECT action, actor, changed, reason, occurred_at, txid, ROW_NUMBER() OVER (ORDER BY occurred_at, txid) AS position, COUNT(*) OVER () AS total FROM auditlog.audit_trails WHERE subject_type = $1 AND subject_id = $2 ORDER BY occurred_at DESC, txid DESC")
245 .bind(SUBJECT_TYPE)
246 .bind(&subject_id)
247 .fetch_all(&mut *conn)
248 .await
249 {
250 Ok(r) => r,
251 Err(e) => {
252 return (
253 StatusCode::INTERNAL_SERVER_ERROR,
254 axum::Json(json!({
255 "error": "history_read",
256 "message": e.to_string(),
257 })),
258 )
259 .into_response()
260 }
261 };
262 let total: i64 = rows.first().map(|r| r.get("total")).unwrap_or(0);
263 let entries: Vec<serde_json::Value> = rows
264 .iter()
265 .map(|r| {
266 json!({
267 "position": r.get::<i64, _>("position"),
268 "action": r.get::<String, _>("action"),
269 "actor": r.get::<String, _>("actor"),
270 "changed": r.get::<serde_json::Value, _>("changed"),
271 "reason": r.get::<Option<String>, _>("reason"),
272 "occurred_at": r.get::<chrono::DateTime<chrono::Utc>, _>("occurred_at"),
273 "txid": r.get::<String, _>("txid"),
274 })
275 })
276 .collect();
277 (
278 StatusCode::OK,
279 axum::Json(json!({
280 "subject_type": SUBJECT_TYPE,
281 "subject_id": subject_id,
282 "total": total,
283 "entries": entries,
284 })),
285 )
286 .into_response()
287 }
288 Some(as_of) => {
289 let rows = match sqlx::query("SELECT action, changed, occurred_at FROM auditlog.audit_trails WHERE subject_type = $1 AND subject_id = $2 AND occurred_at <= $3 ORDER BY occurred_at ASC, txid ASC")
290 .bind(SUBJECT_TYPE)
291 .bind(&subject_id)
292 .bind(as_of)
293 .fetch_all(&mut *conn)
294 .await
295 {
296 Ok(r) => r,
297 Err(e) => {
298 return (
299 StatusCode::INTERNAL_SERVER_ERROR,
300 axum::Json(json!({
301 "error": "history_read",
302 "message": e.to_string(),
303 })),
304 )
305 .into_response()
306 }
307 };
308 if rows.is_empty() {
309 return (
310 StatusCode::NOT_FOUND,
311 axum::Json(json!({
312 "error": "history_before_subject",
313 "message": "as_of precedes the subject's first captured event",
314 })),
315 )
316 .into_response();
317 }
318 let mut image = serde_json::Map::new();
319 let mut deleted_at: Option<chrono::DateTime<chrono::Utc>> = None;
320 for r in &rows {
321 let action: String = r.get("action");
322 let changed: serde_json::Value = r.get("changed");
323 if action == "delete" {
324 image.clear();
325 deleted_at = Some(r.get("occurred_at"));
326 continue;
327 }
328 deleted_at = None;
329 if let Some(fields) = changed.as_object() {
330 for (field, diff) in fields {
331 if let Some(to) = diff.get("to") {
332 image.insert(field.clone(), to.clone());
333 }
334 }
335 }
336 }
337 if let Some(at) = deleted_at {
338 return (
339 StatusCode::NOT_FOUND,
340 axum::Json(json!({
341 "error": "subject_deleted_before_as_of",
342 "deleted_at": at,
343 })),
344 )
345 .into_response();
346 }
347 (
348 StatusCode::OK,
349 axum::Json(json!({
350 "as_of": as_of,
351 "image": serde_json::Value::Object(image),
352 })),
353 )
354 .into_response()
355 }
356 }
357}
358
359pub fn create_compensation_change_history_route(pool: sqlx::PgPool) -> axum::Router {
362 axum::Router::new()
363 .route("/compensation_changes/:id/history", axum::routing::get(compensation_change_history))
364 .with_state(pool)
365}