use std::{io::Cursor, panic, sync::Arc, time::Duration}; use async_trait::async_trait; use chrono::Utc; use diesel::{pg::upsert::excluded, prelude::*}; use image::{ codecs::png::PngEncoder, ColorType, GenericImageView, ImageEncoder, ImageFormat, ImageReader, }; use pdfium_render::prelude::*; use serde::Deserialize; use serde_json::json; use tokio::task; use tracing::{error, info, warn}; use uuid::Uuid; use crate::{ jobs::JOB_GENERATE_THUMBNAILS, models::{Document, DocumentAsset, DocumentVersion, NewDocumentAsset}, schema::{document_assets, document_versions, documents}, state::AppState, }; use super::{analyze::determine_thumbnail_support, JobExecution, JobHandler}; const THUMBNAIL_WIDTH: u32 = 512; const THUMBNAIL_HEIGHT: u32 = 512; const THUMBNAIL_ASSET_TYPE: &str = "thumbnail"; #[derive(Debug, Deserialize)] struct ThumbnailPayload { document_id: Uuid, document_version_id: Uuid, #[serde(default)] force: bool, } pub struct GenerateThumbnailsJob; impl GenerateThumbnailsJob { pub fn new() -> Self { Self } } #[async_trait] impl JobHandler for GenerateThumbnailsJob { fn job_type(&self) -> &'static str { JOB_GENERATE_THUMBNAILS } async fn handle(&self, state: Arc, job: crate::models::Job) -> JobExecution { let payload: ThumbnailPayload = match serde_json::from_value(job.payload.clone()) { Ok(p) => p, Err(err) => { return JobExecution::Failed { error: format!("invalid thumbnail payload: {err}"), } } }; let state_clone = state.clone(); let initial = match task::spawn_blocking(move || load_thumbnail_context(state_clone, &payload)).await { Ok(Ok(ctx)) => ctx, Ok(Err(err)) => { warn!(job_id = %job.id, error = %err, "thumbnail job will retry"); return JobExecution::Retry { delay: Duration::from_secs(30), error: err, }; } Err(join_err) => { error!(job_id = %job.id, error = %join_err, "thumbnail task panicked"); return JobExecution::Retry { delay: Duration::from_secs(60), error: format!("worker panicked: {join_err}"), }; } }; if initial.skip { info!(job_id = %job.id, "thumbnails already exist; skipping"); return JobExecution::Success; } let bytes = match state.storage.get_object(&initial.version.s3_key).await { Ok(bytes) => bytes, Err(err) => { warn!(job_id = %job.id, error = %err, "thumbnail fetch failed; will retry"); return JobExecution::Retry { delay: Duration::from_secs(30), error: err.to_string(), }; } }; let generation = match generate_thumbnail(&initial.document, &bytes) { Ok(result) => result, Err(err) => { return JobExecution::Failed { error: err }; } }; if let Err(err) = state .storage .put_object( &generation.s3_key, generation.image_bytes.clone(), Some("image/png".into()), ) .await { warn!(job_id = %job.id, error = %err, "failed to upload thumbnail; retrying"); return JobExecution::Retry { delay: Duration::from_secs(30), error: err.to_string(), }; } let state_clone = state.clone(); match task::spawn_blocking(move || { persist_thumbnail_metadata(state_clone, &initial, &generation) }) .await { Ok(Ok(())) => {} Ok(Err(err)) => { warn!(job_id = %job.id, error = %err, "failed to persist thumbnail metadata; retrying"); return JobExecution::Retry { delay: Duration::from_secs(30), error: err, }; } Err(join_err) => { error!(job_id = %job.id, error = %join_err, "thumbnail metadata update panicked"); return JobExecution::Retry { delay: Duration::from_secs(30), error: format!("metadata update panic: {join_err}"), }; } } JobExecution::Success } } struct ThumbnailContext { document: Document, version: DocumentVersion, skip: bool, } struct GeneratedThumbnail { image_bytes: Vec, width: Option, height: Option, s3_key: String, } fn load_thumbnail_context( state: Arc, payload: &ThumbnailPayload, ) -> Result { let mut conn = state.db().map_err(|err| format!("{err:?}"))?; let version: DocumentVersion = document_versions::table .find(payload.document_version_id) .first(&mut conn) .map_err(|err| format!("{err:?}"))?; if version.document_id != payload.document_id { return Err("document/version mismatch".into()); } let document: Document = documents::table .find(payload.document_id) .first(&mut conn) .map_err(|err| format!("{err:?}"))?; let existing: Option = document_assets::table .filter(document_assets::document_version_id.eq(payload.document_version_id)) .filter(document_assets::asset_type.eq(THUMBNAIL_ASSET_TYPE)) .first(&mut conn) .optional() .map_err(|err| format!("{err:?}"))?; let (supported, _) = determine_thumbnail_support(&document); if !supported { return Err("thumbnail generation not supported for this document".into()); } let skip = existing.is_some() && !payload.force; Ok(ThumbnailContext { document, version, skip, }) } fn generate_thumbnail(document: &Document, bytes: &[u8]) -> Result { let is_pdf = document .content_type .as_deref() .map(|mime| mime == "application/pdf") .unwrap_or_else(|| { document .original_name .rsplit('.') .next() .map(|ext| ext.eq_ignore_ascii_case("pdf")) .unwrap_or(false) }); let (png_bytes, width, height) = if is_pdf { generate_pdf_thumbnail(bytes)? } else { generate_image_thumbnail(bytes)? }; let s3_key = format!("thumbnails/{}/{}.png", document.id, Uuid::new_v4()); Ok(GeneratedThumbnail { image_bytes: png_bytes, width, height, s3_key, }) } fn generate_image_thumbnail(bytes: &[u8]) -> Result<(Vec, Option, Option), String> { let reader = ImageReader::new(Cursor::new(bytes)) .with_guessed_format() .map_err(|err| err.to_string())?; let mut image = reader.decode().map_err(|err| err.to_string())?; if image.width() > THUMBNAIL_WIDTH || image.height() > THUMBNAIL_HEIGHT { image = image.thumbnail(THUMBNAIL_WIDTH, THUMBNAIL_HEIGHT); } let (width, height) = image.dimensions(); let mut cursor = Cursor::new(Vec::new()); image .write_to(&mut cursor, ImageFormat::Png) .map_err(|err| err.to_string())?; let buffer = cursor.into_inner(); Ok((buffer, Some(width as i32), Some(height as i32))) } fn generate_pdf_thumbnail(bytes: &[u8]) -> Result<(Vec, Option, Option), String> { let pdfium = panic::catch_unwind(|| Pdfium::default()) .map_err(|_| "failed to initialize PDFium".to_string())?; let document = pdfium .load_pdf_from_byte_slice(bytes, None) .map_err(|err| format!("load pdf: {err}"))?; let page = document .pages() .get(0) .map_err(|err| format!("load first page: {err}"))?; let render_config = PdfRenderConfig::new() .set_target_width(THUMBNAIL_WIDTH as i32) .set_maximum_height(THUMBNAIL_HEIGHT as i32) .render_form_data(true) .rotate_if_landscape(PdfPageRenderRotation::None, true); let bitmap = page .render_with_config(&render_config) .map_err(|err| format!("render pdf page: {err}"))?; let image = bitmap.as_image().to_rgb8(); let (width, height) = image.dimensions(); let mut cursor = Cursor::new(Vec::new()); PngEncoder::new(&mut cursor) .write_image(image.as_raw(), width, height, ColorType::Rgb8.into()) .map_err(|err| format!("encode pdf thumbnail: {err}"))?; Ok((cursor.into_inner(), Some(width as i32), Some(height as i32))) } fn persist_thumbnail_metadata( state: Arc, context: &ThumbnailContext, generated: &GeneratedThumbnail, ) -> Result<(), String> { let mut conn = state.db().map_err(|err| format!("{err:?}"))?; let new_asset = NewDocumentAsset { id: Uuid::new_v4(), document_version_id: context.version.id, asset_type: THUMBNAIL_ASSET_TYPE.to_string(), s3_key: generated.s3_key.clone(), mime_type: "image/png".to_string(), width: generated.width, height: generated.height, metadata: json!({ "generated_at": Utc::now().to_rfc3339(), }), }; diesel::insert_into(document_assets::table) .values(&new_asset) .on_conflict(( document_assets::document_version_id, document_assets::asset_type, )) .do_update() .set(( document_assets::s3_key.eq(excluded(document_assets::s3_key)), document_assets::mime_type.eq(excluded(document_assets::mime_type)), document_assets::width.eq(excluded(document_assets::width)), document_assets::height.eq(excluded(document_assets::height)), document_assets::metadata.eq(excluded(document_assets::metadata)), )) .execute(&mut conn) .map_err(|err| format!("{err:?}"))?; Ok(()) }