use std::collections::HashSet; use std::sync::Arc; use std::time::Duration; use async_trait::async_trait; use diesel::prelude::*; use diesel::result::Error as DieselError; use serde::Deserialize; use tracing::{error, warn}; use uuid::Uuid; use crate::jobs::JOB_PURGE_DOCUMENT; use crate::models::{Document, DocumentVersion}; use crate::schema::{document_asset_objects, document_assets, document_versions}; use crate::state::AppState; use crate::storage::TenantStorage; use super::{JobExecution, JobHandler}; #[derive(Debug, Deserialize)] struct PurgeDocumentPayload { document_id: Uuid, } #[derive(Debug)] struct PurgeContext { document_id: Uuid, version_keys: Vec, asset_keys: Vec, } pub struct PurgeDocumentJob; impl PurgeDocumentJob { pub fn new() -> Self { Self } } #[async_trait] impl JobHandler for PurgeDocumentJob { fn job_type(&self) -> &'static str { JOB_PURGE_DOCUMENT } async fn handle( &self, state: Arc, job: crate::models::Job, storage: TenantStorage, ) -> JobExecution { let payload: PurgeDocumentPayload = match serde_json::from_value(job.payload.clone()) { Ok(payload) => payload, Err(err) => { return JobExecution::Failed { error: format!("invalid purge payload: {err}"), }; } }; let tenant_id = job.tenant_id; let document_id = payload.document_id; let state_for_prepare = state.clone(); let preparation = tokio::task::spawn_blocking(move || { prepare_purge_context(state_for_prepare, tenant_id, document_id) }) .await; let context = match preparation { Ok(Ok(Some(ctx))) => ctx, Ok(Ok(None)) => { // Document already gone or restored; nothing to do. return JobExecution::Success; } Ok(Err(err)) => { warn!(job_id = %job.id, error = %err, "purge preparation failed"); return JobExecution::Retry { delay: Duration::from_secs(30), error: err, }; } Err(join_err) => { error!(job_id = %job.id, error = %join_err, "purge preparation task panicked"); return JobExecution::Retry { delay: Duration::from_secs(60), error: format!("purge preparation panicked: {join_err}"), }; } }; if let Err(err) = delete_storage_objects(&storage, &context).await { warn!(job_id = %job.id, error = %err, "failed to delete storage objects for purge"); return JobExecution::Retry { delay: Duration::from_secs(30), error: err, }; } let PurgeContext { document_id, .. } = context; let state_for_finalize = state.clone(); let finalize = tokio::task::spawn_blocking(move || { finalize_purge(state_for_finalize, tenant_id, document_id) }) .await; match finalize { Ok(Ok(())) => JobExecution::Success, Ok(Err(err)) => { warn!(job_id = %job.id, error = %err, "failed to finalize purge"); JobExecution::Retry { delay: Duration::from_secs(30), error: err, } } Err(join_err) => { error!(job_id = %job.id, error = %join_err, "purge finalize task panicked"); JobExecution::Retry { delay: Duration::from_secs(60), error: format!("purge finalize panicked: {join_err}"), } } } } } fn prepare_purge_context( state: Arc, tenant_id: Uuid, document_id: Uuid, ) -> Result, String> { let mut conn = state .db_for_tenant(tenant_id) .map_err(|err| format!("failed to scope tenant connection: {err:?}"))?; conn.transaction(|conn| { use crate::schema::documents::dsl as doc_dsl; let doc_opt = doc_dsl::documents .filter(doc_dsl::tenant_id.eq(tenant_id)) .find(document_id) .for_update() .first::(conn) .optional()?; let Some(document) = doc_opt else { return Ok(None); }; if document.deleted_at.is_none() { return Ok(None); } let versions: Vec = document_versions::table .filter(document_versions::document_id.eq(document_id)) .filter(document_versions::tenant_id.eq(tenant_id)) .load(conn)?; let version_keys: Vec = versions .iter() .map(|version| version.s3_key.clone()) .collect(); let version_ids: Vec = versions.iter().map(|version| version.id).collect(); let asset_keys = if version_ids.is_empty() { Vec::new() } else { let asset_ids: Vec = document_assets::table .filter(document_assets::document_version_id.eq_any(&version_ids)) .filter(document_assets::tenant_id.eq(tenant_id)) .select(document_assets::id) .load(conn)?; if asset_ids.is_empty() { Vec::new() } else { document_asset_objects::table .filter(document_asset_objects::asset_id.eq_any(&asset_ids)) .filter(document_asset_objects::tenant_id.eq(tenant_id)) .select(document_asset_objects::s3_key) .load(conn)? } }; Ok(Some(PurgeContext { document_id, version_keys, asset_keys, })) }) .map_err(|err: DieselError| format!("failed to prepare purge: {err}")) } async fn delete_storage_objects( storage: &TenantStorage, context: &PurgeContext, ) -> Result<(), String> { let mut keys = HashSet::new(); keys.extend(context.version_keys.iter().cloned()); keys.extend(context.asset_keys.iter().cloned()); for key in keys { if let Err(err) = storage.delete_object(&key).await { return Err(format!("failed to delete object {}: {err:?}", key)); } } Ok(()) } fn finalize_purge(state: Arc, tenant_id: Uuid, document_id: Uuid) -> Result<(), String> { let mut conn = state .db_for_tenant(tenant_id) .map_err(|err| format!("failed to scope tenant connection: {err:?}"))?; conn.transaction(|conn| { use crate::schema::documents::dsl as doc_dsl; let doc_opt = doc_dsl::documents .filter(doc_dsl::tenant_id.eq(tenant_id)) .find(document_id) .for_update() .first::(conn) .optional()?; let Some(document) = doc_opt else { return Ok(()); }; if document.deleted_at.is_none() { return Ok(()); } diesel::delete(doc_dsl::documents.filter(doc_dsl::id.eq(document_id))).execute(conn)?; Ok(()) }) .map_err(|err: DieselError| format!("failed to finalize purge: {err}")) }