use std::{collections::HashMap, path::Path as FsPath, time::Duration}; use axum::extract::{Json, Multipart, Path, Query, State}; use axum::http::StatusCode; use axum::response::IntoResponse; use chrono::{DateTime, NaiveDateTime, Utc}; use diesel::dsl::exists; use diesel::{prelude::*, select, PgConnection}; use serde::{Deserialize, Serialize}; use serde_json::{json, Value}; use sha2::{Digest, Sha256}; use tracing::{error, info, warn}; use uuid::Uuid; use crate::auth::AuthenticatedUser; use crate::error::{AppError, AppResult}; use crate::jobs::{enqueue_job, JOB_ANALYZE_DOCUMENT}; use crate::models::{ Document, DocumentAsset, DocumentVersion, NewDocument, NewDocumentTag, NewDocumentVersion, Tag, }; use crate::schema::{ document_assets, document_tags, document_versions, documents, folders, refresh_tokens::dsl as refresh_dsl, tags, }; use crate::state::AppState; const PRESIGNED_URL_EXPIRY_SECONDS: u64 = 300; fn inline_content_disposition(filename: &str) -> Option { if filename.is_empty() { return None; } let sanitized: String = filename .chars() .map(|ch| match ch { '"' | '\\' => '_', _ => ch, }) .collect(); let encoded = percent_encoding::utf8_percent_encode(&sanitized, percent_encoding::NON_ALPHANUMERIC); Some(format!( "inline; filename=\"{}\"; filename*=UTF-8''{}", sanitized, encoded )) } #[derive(Deserialize)] pub struct DocumentListQuery { pub folder_id: Option, #[serde(default)] pub include_deleted: bool, } #[derive(Deserialize)] pub struct AssetRequestQuery { #[serde(default)] pub force: bool, } #[derive(Serialize)] pub struct TagResponse { pub id: Uuid, pub label: String, pub color: Option, } impl From for TagResponse { fn from(tag: Tag) -> Self { Self { id: tag.id, label: tag.label, color: tag.color, } } } #[derive(Serialize)] pub struct DocumentResponse { pub id: Uuid, pub filename: String, pub title: String, pub original_name: String, pub content_type: Option, pub folder_id: Option, pub current_version: i32, pub uploaded_at: String, pub updated_at: String, pub deleted_at: Option, pub issued_at: Option, pub metadata: Value, pub tags: Vec, pub thumbnail: Option, pub download_path: String, } #[derive(Serialize, Clone)] pub struct DocumentVersionResponse { pub id: Uuid, pub version_number: i32, pub s3_key: String, pub size_bytes: i64, pub checksum: String, pub created_at: String, pub operations_summary: Value, } #[derive(Serialize, Clone)] pub struct DocumentAssetResponse { pub id: Uuid, pub asset_type: String, pub mime_type: String, pub width: Option, pub height: Option, pub url: String, pub created_at: String, } #[derive(Serialize)] pub struct DocumentDetailResponse { pub document: DocumentResponse, pub current_version: DocumentVersionResponse, pub assets: Vec, } #[derive(Serialize)] pub struct DocumentDownloadResponse { pub url: String, pub expires_in: u64, pub filename: String, pub content_type: Option, pub size_bytes: i64, } #[derive(Serialize)] pub struct BulkReanalyzeResponse { pub queued: usize, } #[derive(Deserialize)] pub struct BulkMoveRequest { pub document_ids: Vec, pub folder_id: Option, } #[derive(Serialize)] pub struct BulkMoveResponse { pub updated: usize, } #[derive(Deserialize)] #[serde(rename_all = "snake_case")] pub enum BulkTagAction { Add, Remove, } #[derive(Deserialize)] pub struct BulkTagRequest { pub document_ids: Vec, pub tag_ids: Vec, pub action: BulkTagAction, } #[derive(Serialize)] pub struct BulkTagResponse { pub added: usize, pub removed: usize, } #[derive(Deserialize)] pub struct BulkReanalyzeSelectionRequest { pub document_ids: Vec, #[serde(default = "default_true")] pub force: bool, } fn default_true() -> bool { true } struct UploadRequest { bytes: Vec, original_name: String, content_type: Option, folder_id: Option, metadata: Value, } struct UploadOutcome { detail: DocumentDetailResponse, created: bool, } #[derive(Deserialize)] pub struct MoveDocumentRequest { pub folder_id: Option, } #[derive(Deserialize)] pub struct AssignTagsRequest { pub tag_ids: Vec, } pub async fn list_documents( State(state): State, Query(query): Query, user: AuthenticatedUser, ) -> AppResult>> { let mut conn = state.db()?; let mut base_query = documents::table.into_boxed(); if !query.include_deleted { base_query = base_query.filter(documents::deleted_at.is_null()); } match query.folder_id { Some(folder_id) => { base_query = base_query.filter(documents::folder_id.eq(Some(folder_id))); } None => { base_query = base_query.filter(documents::folder_id.is_null()); } } let docs: Vec = base_query .order(documents::uploaded_at.desc()) .load(&mut conn)?; let doc_ids: Vec = docs.iter().map(|doc| doc.id).collect(); let tags_map = load_tags_for_documents(&mut conn, &doc_ids)?; drop(conn); let thumbnails = load_primary_thumbnails(&state, &docs).await?; let mut response = Vec::with_capacity(doc_ids.len()); for doc in docs { let tags = tags_map.get(&doc.id).cloned(); let thumbnail = thumbnails.get(&doc.id).cloned(); response.push(to_document_response( &state, user.user_id, doc, tags, thumbnail, )?); } Ok(Json(response)) } pub async fn get_document( State(state): State, Path(document_id): Path, user: AuthenticatedUser, ) -> AppResult> { let mut conn = state.db()?; let doc: Document = documents::table.find(document_id).first(&mut conn)?; if doc.deleted_at.is_some() { return Err(AppError::not_found()); } let current_version: DocumentVersion = document_versions::table .filter(document_versions::document_id.eq(document_id)) .filter(document_versions::version_number.eq(doc.current_version)) .first(&mut conn)?; let tags_map = load_tags_for_documents(&mut conn, &[document_id])?; let version_id = current_version.id; drop(conn); let assets = load_asset_responses(&state, version_id).await?; let thumbnail = assets .iter() .find(|asset| asset.asset_type == "thumbnail") .cloned(); Ok(Json(DocumentDetailResponse { document: to_document_response( &state, user.user_id, doc, tags_map.get(&document_id).cloned(), thumbnail, )?, current_version: to_version_response(current_version), assets, })) } pub async fn upload_document( State(state): State, user: AuthenticatedUser, mut multipart: Multipart, ) -> AppResult<(StatusCode, Json)> { let mut file_bytes: Option> = None; let mut original_name: Option = None; let mut content_type: Option = None; let mut folder_id: Option = None; let mut metadata: Value = Value::Object(Default::default()); while let Some(field) = multipart.next_field().await.map_err(|err| { let msg = format!("invalid multipart data: {err}"); error!(error = %err, "invalid multipart data"); AppError::bad_request(msg) })? { let name = field.name().map(|n| n.to_string()); match name.as_deref() { Some("file") => { let file_name = field.file_name().map(|n| n.to_string()); original_name = file_name.clone(); content_type = field.content_type().map(|mime| mime.to_string()); let data = field.bytes().await.map_err(|err| { let msg = format!("failed to read file bytes: {err}"); error!(error = %err, "failed to read file bytes"); AppError::bad_request(msg) })?; file_bytes = Some(data.to_vec()); } Some("folder_id") => { let value = field.text().await.map_err(|err| { let msg = format!("invalid folder id: {err}"); error!(error = %err, "invalid folder id"); AppError::bad_request(msg) })?; if !value.trim().is_empty() { let parsed = Uuid::parse_str(value.trim()) .map_err(|_| AppError::bad_request("folder_id must be a valid UUID"))?; folder_id = Some(parsed); } } Some("metadata") => { let value = field.text().await.map_err(|err| { let msg = format!("invalid metadata: {err}"); error!(error = %err, "invalid metadata payload"); AppError::bad_request(msg) })?; metadata = serde_json::from_str(&value).map_err(|err| { let msg = format!("metadata must be valid JSON: {err}"); error!(error = %err, "metadata parse failure"); AppError::bad_request(msg) })?; } _ => {} } } let file_bytes = file_bytes.ok_or_else(|| { error!("upload rejected: missing file field"); AppError::bad_request("file field is required") })?; if file_bytes.is_empty() { error!("upload rejected: empty file payload"); return Err(AppError::bad_request("file field must not be empty")); } let original_name = original_name.ok_or_else(|| { error!("upload rejected: missing original filename"); AppError::bad_request("filename is required") })?; let original_name_for_log = original_name.clone(); let request = UploadRequest { bytes: file_bytes, original_name, content_type, folder_id, metadata, }; let outcome = match process_upload(&state, request, user.user_id).await { Ok(outcome) => { info!( document_id = %outcome.detail.document.id, original_name = %outcome.detail.document.original_name, created = outcome.created, reused_existing = !outcome.created, "document upload succeeded" ); outcome } Err(err) => { error!(error = ?err, original_name = %original_name_for_log, "document upload failed"); return Err(err); } }; let status = if outcome.created { StatusCode::CREATED } else { StatusCode::OK }; Ok((status, Json(outcome.detail))) } pub async fn request_document_assets( State(state): State, Path(document_id): Path, Query(query): Query, ) -> AppResult { let mut conn = state.db()?; let document: Document = documents::table.find(document_id).first(&mut conn)?; if document.deleted_at.is_some() { return Err(AppError::not_found()); } let version: DocumentVersion = document_versions::table .filter(document_versions::document_id.eq(document_id)) .filter(document_versions::version_number.eq(document.current_version)) .first(&mut conn)?; enqueue_job( &mut conn, JOB_ANALYZE_DOCUMENT, json!({ "document_id": document_id, "document_version_id": version.id, "force": query.force, }), None, ) .map_err(|err| AppError::internal(format!("failed to enqueue analyze job: {err}")))?; Ok(StatusCode::ACCEPTED) } pub async fn reanalyze_all_documents( State(state): State, _user: AuthenticatedUser, ) -> AppResult<(StatusCode, Json)> { let mut conn = state.db()?; let targets: Vec<(Uuid, Uuid)> = document_versions::table .inner_join(documents::table.on(document_versions::document_id.eq(documents::id))) .filter(documents::deleted_at.is_null()) .filter(document_versions::version_number.eq(documents::current_version)) .select((documents::id, document_versions::id)) .load(&mut conn)?; let mut queued = 0usize; for (document_id, version_id) in targets { enqueue_job( &mut conn, JOB_ANALYZE_DOCUMENT, json!({ "document_id": document_id, "document_version_id": version_id, "force": true, }), None, ) .map_err(|err| AppError::internal(format!("failed to enqueue analyze job: {err}")))?; queued += 1; } Ok((StatusCode::ACCEPTED, Json(BulkReanalyzeResponse { queued }))) } pub async fn reanalyze_selected_documents( State(state): State, Json(payload): Json, ) -> AppResult<(StatusCode, Json)> { let BulkReanalyzeSelectionRequest { mut document_ids, force, } = payload; if document_ids.is_empty() { return Err(AppError::bad_request("document_ids must not be empty")); } document_ids.sort(); document_ids.dedup(); let mut conn = state.db()?; let targets: Vec<(Uuid, Uuid)> = document_versions::table .inner_join(documents::table.on(document_versions::document_id.eq(documents::id))) .filter(documents::id.eq_any(&document_ids)) .filter(documents::deleted_at.is_null()) .filter(document_versions::version_number.eq(documents::current_version)) .select((documents::id, document_versions::id)) .load(&mut conn)?; if targets.len() != document_ids.len() { return Err(AppError::bad_request( "one or more documents do not exist or are inaccessible", )); } let mut queued = 0usize; for (document_id, version_id) in targets { enqueue_job( &mut conn, JOB_ANALYZE_DOCUMENT, json!({ "document_id": document_id, "document_version_id": version_id, "force": force, }), None, ) .map_err(|err| AppError::internal(format!("failed to enqueue analyze job: {err}")))?; queued += 1; } Ok((StatusCode::ACCEPTED, Json(BulkReanalyzeResponse { queued }))) } pub async fn list_document_assets( State(state): State, Path(document_id): Path, ) -> AppResult>> { let mut conn = state.db()?; let document: Document = documents::table.find(document_id).first(&mut conn)?; if document.deleted_at.is_some() { return Err(AppError::not_found()); } let version: DocumentVersion = document_versions::table .filter(document_versions::document_id.eq(document_id)) .filter(document_versions::version_number.eq(document.current_version)) .first(&mut conn)?; drop(conn); let assets = load_asset_responses(&state, version.id).await?; Ok(Json(assets)) } pub async fn download_document( State(state): State, Path(document_id): Path, ) -> AppResult> { let mut conn = state.db()?; let doc: Document = documents::table.find(document_id).first(&mut conn)?; if doc.deleted_at.is_some() { return Err(AppError::not_found()); } let version: DocumentVersion = document_versions::table .filter(document_versions::document_id.eq(document_id)) .filter(document_versions::version_number.eq(doc.current_version)) .first(&mut conn)?; let presigned_url = state .storage .presign_get_object( &version.s3_key, Duration::from_secs(PRESIGNED_URL_EXPIRY_SECONDS), ) .await .map_err(|err| AppError::internal(format!("failed to generate download URL: {err}")))?; Ok(Json(DocumentDownloadResponse { url: presigned_url, expires_in: PRESIGNED_URL_EXPIRY_SECONDS, filename: doc.original_name.clone(), content_type: doc.content_type.clone(), size_bytes: version.size_bytes, })) } pub async fn download_with_token( State(state): State, Path(token): Path, ) -> AppResult { let claims = state .jwt .verify_download_token(&token) .map_err(|_| AppError::unauthorized())?; let mut conn = state.db()?; let doc: Document = documents::table.find(claims.doc_id).first(&mut conn)?; if doc.deleted_at.is_some() { return Err(AppError::not_found()); } let version: DocumentVersion = document_versions::table .filter(document_versions::document_id.eq(claims.doc_id)) .filter(document_versions::version_number.eq(doc.current_version)) .first(&mut conn)?; let now = Utc::now().naive_utc(); let has_active_refresh: bool = select(exists( refresh_dsl::refresh_tokens .filter(refresh_dsl::user_id.eq(claims.user_id)) .filter(refresh_dsl::revoked_at.is_null()) .filter(refresh_dsl::expires_at.gt(now)), )) .get_result(&mut conn)?; if !has_active_refresh { return Err(AppError::unauthorized()); } drop(conn); let presigned_url = state .storage .presign_get_object( &version.s3_key, Duration::from_secs(PRESIGNED_URL_EXPIRY_SECONDS), ) .await .map_err(|err| AppError::internal(format!("failed to generate download URL: {err}")))?; Ok(axum::response::Redirect::temporary(&presigned_url)) } pub async fn delete_document( State(state): State, Path(document_id): Path, ) -> AppResult { let mut conn = state.db()?; let now = Utc::now().naive_utc(); diesel::update(documents::table.find(document_id)) .set(( documents::deleted_at.eq(Some(now)), documents::updated_at.eq(now), )) .execute(&mut conn)?; Ok(StatusCode::NO_CONTENT) } pub async fn move_document( State(state): State, Path(document_id): Path, Json(payload): Json, ) -> AppResult { if let Some(folder_id) = payload.folder_id { ensure_folder_exists(&state, folder_id)?; } let mut conn = state.db()?; let now = Utc::now().naive_utc(); diesel::update(documents::table.find(document_id)) .set(( documents::folder_id.eq(payload.folder_id), documents::updated_at.eq(now), )) .execute(&mut conn)?; Ok(StatusCode::NO_CONTENT) } pub async fn bulk_move_documents( State(state): State, Json(payload): Json, ) -> AppResult<(StatusCode, Json)> { let BulkMoveRequest { mut document_ids, folder_id, } = payload; if document_ids.is_empty() { return Err(AppError::bad_request("document_ids must not be empty")); } document_ids.sort(); document_ids.dedup(); if let Some(target_folder) = folder_id { ensure_folder_exists(&state, target_folder)?; } let mut conn = state.db()?; let existing: Vec<(Uuid, Option)> = documents::table .filter(documents::id.eq_any(&document_ids)) .select((documents::id, documents::deleted_at)) .load(&mut conn)?; if existing.len() != document_ids.len() { return Err(AppError::bad_request( "one or more documents do not exist or are inaccessible", )); } if existing.iter().any(|(_, deleted)| deleted.is_some()) { return Err(AppError::bad_request("cannot move deleted documents")); } let now = Utc::now().naive_utc(); let updated = diesel::update(documents::table.filter(documents::id.eq_any(&document_ids))) .set(( documents::folder_id.eq(folder_id), documents::updated_at.eq(now), )) .execute(&mut conn)?; Ok((StatusCode::OK, Json(BulkMoveResponse { updated }))) } pub async fn assign_tags( State(state): State, Path(document_id): Path, user: AuthenticatedUser, Json(payload): Json, ) -> AppResult { if payload.tag_ids.is_empty() { return Err(AppError::bad_request("tag_ids must not be empty")); } let mut conn = state.db()?; // Ensure document exists documents::table .find(document_id) .first::(&mut conn)?; // Ensure tags exist let existing_tags: Vec = tags::table .filter(tags::id.eq_any(&payload.tag_ids)) .load(&mut conn)?; if existing_tags.len() != payload.tag_ids.len() { return Err(AppError::bad_request("one or more tags do not exist")); } let new_tags: Vec = payload .tag_ids .iter() .map(|tag_id| NewDocumentTag { document_id, tag_id: *tag_id, assigned_by: Some(user.user_id), }) .collect(); diesel::insert_into(document_tags::table) .values(&new_tags) .on_conflict_do_nothing() .execute(&mut conn)?; Ok(StatusCode::NO_CONTENT) } pub async fn bulk_update_tags( State(state): State, user: AuthenticatedUser, Json(payload): Json, ) -> AppResult<(StatusCode, Json)> { let BulkTagRequest { mut document_ids, mut tag_ids, action, } = payload; if document_ids.is_empty() { return Err(AppError::bad_request("document_ids must not be empty")); } if tag_ids.is_empty() { return Err(AppError::bad_request("tag_ids must not be empty")); } document_ids.sort(); document_ids.dedup(); tag_ids.sort(); tag_ids.dedup(); let mut conn = state.db()?; let docs: Vec<(Uuid, Option)> = documents::table .filter(documents::id.eq_any(&document_ids)) .select((documents::id, documents::deleted_at)) .load(&mut conn)?; if docs.len() != document_ids.len() { return Err(AppError::bad_request( "one or more documents do not exist or are inaccessible", )); } if docs.iter().any(|(_, deleted)| deleted.is_some()) { return Err(AppError::bad_request( "cannot assign or remove tags from deleted documents", )); } let existing_tags: Vec = tags::table .filter(tags::id.eq_any(&tag_ids)) .load(&mut conn)?; if existing_tags.len() != tag_ids.len() { return Err(AppError::bad_request("one or more tags do not exist")); } let response = match action { BulkTagAction::Add => { let mut inserts = Vec::with_capacity(document_ids.len() * tag_ids.len()); for doc_id in &document_ids { for tag_id in &tag_ids { inserts.push(NewDocumentTag { document_id: *doc_id, tag_id: *tag_id, assigned_by: Some(user.user_id), }); } } let added = if inserts.is_empty() { 0 } else { diesel::insert_into(document_tags::table) .values(&inserts) .on_conflict_do_nothing() .execute(&mut conn)? }; BulkTagResponse { added, removed: 0 } } BulkTagAction::Remove => { let removed = diesel::delete( document_tags::table .filter(document_tags::document_id.eq_any(&document_ids)) .filter(document_tags::tag_id.eq_any(&tag_ids)), ) .execute(&mut conn)?; BulkTagResponse { added: 0, removed } } }; Ok((StatusCode::OK, Json(response))) } pub async fn remove_tag( State(state): State, Path((document_id, tag_id)): Path<(Uuid, Uuid)>, ) -> AppResult { let mut conn = state.db()?; diesel::delete( document_tags::table .filter(document_tags::document_id.eq(document_id)) .filter(document_tags::tag_id.eq(tag_id)), ) .execute(&mut conn)?; Ok(StatusCode::NO_CONTENT) } async fn process_upload( state: &AppState, request: UploadRequest, user_id: Uuid, ) -> AppResult { let UploadRequest { bytes, original_name, content_type, folder_id, metadata, } = request; if let Some(folder) = folder_id { ensure_folder_exists(state, folder)?; } let doc_id = Uuid::new_v4(); let version_id = Uuid::new_v4(); let version_number = 1; let stored_filename = original_name.clone(); let checksum = Sha256::digest(&bytes); let checksum_hex = hex::encode(checksum); let size_bytes = bytes.len() as i64; let s3_key = format!("documents/{doc_id}/v{version_number}/{version_id}"); { let mut conn = state.db()?; let existing = documents::table .inner_join( document_versions::table.on(document_versions::document_id .eq(documents::id) .and(document_versions::version_number.eq(documents::current_version))), ) .filter(document_versions::checksum.eq(&checksum_hex)) .select((documents::all_columns, document_versions::all_columns)) .first::<(Document, DocumentVersion)>(&mut conn) .optional()?; if let Some((mut document, version)) = existing { if document.deleted_at.is_some() { let now = Utc::now().naive_utc(); diesel::update(documents::table.find(document.id)) .set(( documents::deleted_at.eq(None::), documents::updated_at.eq(now), )) .execute(&mut conn)?; document.deleted_at = None; document.updated_at = now; } let tags_map = load_tags_for_documents(&mut conn, &[document.id])?; let tags = tags_map.get(&document.id).cloned(); drop(conn); let assets = load_asset_responses(state, version.id).await?; let thumbnail = assets .iter() .find(|asset| asset.asset_type == "thumbnail") .cloned(); info!( document_id = %document.id, checksum = %checksum_hex, "upload deduplicated existing document" ); return Ok(UploadOutcome { detail: DocumentDetailResponse { document: to_document_response(state, user_id, document, tags, thumbnail)?, current_version: to_version_response(version), assets, }, created: false, }); } } let content_disposition = inline_content_disposition(&original_name); state .storage .put_object( &s3_key, bytes.clone(), content_type.clone(), content_disposition.clone(), ) .await .map_err(|err| { error!(error = %err, key = %s3_key, "failed to store document"); AppError::internal(format!("failed to store document: {err}")) })?; let metadata_value = if metadata.is_null() { Value::Object(Default::default()) } else { metadata }; let (document, version) = { let mut conn = state.db()?; conn.transaction(|conn| { let new_document = NewDocument { id: doc_id, filename: stored_filename.clone(), original_name: original_name.clone(), content_type: content_type.clone(), folder_id, current_version: version_number, issued_at: None, title: derive_document_title(&original_name), metadata: metadata_value.clone(), }; diesel::insert_into(documents::table) .values(&new_document) .execute(conn)?; let new_version = NewDocumentVersion { id: version_id, document_id: doc_id, version_number, s3_key: s3_key.clone(), size_bytes, checksum: checksum_hex.clone(), operations_summary: Value::Object(Default::default()), }; diesel::insert_into(document_versions::table) .values(&new_version) .execute(conn)?; let document: Document = documents::table.find(doc_id).first(conn)?; let version: DocumentVersion = document_versions::table.find(version_id).first(conn)?; Ok::<_, diesel::result::Error>((document, version)) })? }; let detail = DocumentDetailResponse { document: to_document_response(state, user_id, document, None, None)?, current_version: to_version_response(version.clone()), assets: Vec::new(), }; if let Ok(mut conn) = state.db() { if let Err(err) = enqueue_job( &mut conn, JOB_ANALYZE_DOCUMENT, json!({ "document_id": doc_id, "document_version_id": version.id, "force": false, }), None, ) { warn!(document_id = %doc_id, error = %err, "failed to enqueue analyze job"); } } else { warn!(document_id = %doc_id, "failed to enqueue analyze job due to pool error"); } Ok(UploadOutcome { detail, created: true, }) } fn ensure_folder_exists(state: &AppState, folder_id: Uuid) -> AppResult<()> { let mut conn = state.db()?; let exists: bool = diesel::select(exists(folders::table.filter(folders::id.eq(folder_id)))) .get_result(&mut conn)?; if !exists { return Err(AppError::bad_request("folder does not exist")); } Ok(()) } pub(crate) fn load_tags_for_documents( conn: &mut PgConnection, document_ids: &[Uuid], ) -> AppResult>> { if document_ids.is_empty() { return Ok(HashMap::new()); } let rows: Vec<(Uuid, Tag)> = document_tags::table .inner_join(tags::table) .filter(document_tags::document_id.eq_any(document_ids)) .select((document_tags::document_id, tags::all_columns)) .load(conn)?; let mut map: HashMap> = HashMap::new(); for (doc_id, tag) in rows { map.entry(doc_id).or_default().push(tag); } Ok(map) } pub(crate) async fn load_primary_thumbnails( state: &AppState, documents: &[Document], ) -> AppResult> { if documents.is_empty() { return Ok(HashMap::new()); } let mut current_versions: HashMap = HashMap::new(); let mut doc_ids = Vec::with_capacity(documents.len()); for doc in documents { current_versions.insert(doc.id, doc.current_version); doc_ids.push(doc.id); } let mut conn = state.db()?; let versions: Vec = document_versions::table .filter(document_versions::document_id.eq_any(&doc_ids)) .load(&mut conn)?; let mut version_id_by_doc: HashMap = HashMap::new(); let mut doc_id_by_version: HashMap = HashMap::new(); for version in versions { if let Some(current_number) = current_versions.get(&version.document_id) { if *current_number == version.version_number { version_id_by_doc.insert(version.document_id, version.id); doc_id_by_version.insert(version.id, version.document_id); } } } if version_id_by_doc.is_empty() { return Ok(HashMap::new()); } let version_ids: Vec = version_id_by_doc.values().copied().collect(); let assets: Vec = document_assets::table .filter(document_assets::document_version_id.eq_any(&version_ids)) .filter(document_assets::asset_type.eq("thumbnail")) .order(( document_assets::document_version_id.asc(), document_assets::created_at.asc(), )) .load(&mut conn)?; drop(conn); let mut first_assets: HashMap = HashMap::new(); for asset in assets { if let Some(doc_id) = doc_id_by_version.get(&asset.document_version_id) { first_assets.entry(*doc_id).or_insert(asset); } } let mut responses = HashMap::with_capacity(first_assets.len()); for (doc_id, asset) in first_assets { let url = state .storage .presign_get_object( &asset.s3_key, Duration::from_secs(PRESIGNED_URL_EXPIRY_SECONDS), ) .await .map_err(|err| AppError::internal(format!("failed to sign asset URL: {err}")))?; responses.insert(doc_id, to_asset_response(asset, url)); } Ok(responses) } pub(crate) fn to_document_response( state: &AppState, user_id: Uuid, doc: Document, tags: Option>, thumbnail: Option, ) -> AppResult { let download_path = build_download_path(state, doc.id, user_id)?; Ok(DocumentResponse { id: doc.id, filename: doc.filename, title: doc.title, original_name: doc.original_name, content_type: doc.content_type, folder_id: doc.folder_id, current_version: doc.current_version, uploaded_at: to_iso(doc.uploaded_at), updated_at: to_iso(doc.updated_at), deleted_at: doc.deleted_at.map(to_iso), issued_at: doc.issued_at.map(to_iso), metadata: doc.metadata, tags: tags .unwrap_or_default() .into_iter() .map(TagResponse::from) .collect(), thumbnail, download_path, }) } fn build_download_path(state: &AppState, document_id: Uuid, user_id: Uuid) -> AppResult { state .jwt .generate_download_token(document_id, user_id) .map(|token| format!("/download/{token}")) .map_err(|err| AppError::internal(format!("failed to generate download token: {err}"))) } fn to_version_response(version: DocumentVersion) -> DocumentVersionResponse { DocumentVersionResponse { id: version.id, version_number: version.version_number, s3_key: version.s3_key, size_bytes: version.size_bytes, checksum: version.checksum, created_at: to_iso(version.created_at), operations_summary: version.operations_summary, } } fn to_asset_response(asset: DocumentAsset, url: String) -> DocumentAssetResponse { DocumentAssetResponse { id: asset.id, asset_type: asset.asset_type, mime_type: asset.mime_type, width: asset.width, height: asset.height, url, created_at: to_iso(asset.created_at), } } fn derive_document_title(original: &str) -> String { let trimmed = original.trim(); if trimmed.is_empty() { return "Document".to_string(); } let stem = FsPath::new(trimmed) .file_stem() .and_then(|s| s.to_str()) .map(|s| s.trim()) .filter(|s| !s.is_empty()) .map(|s| s.to_string()); stem.unwrap_or_else(|| trimmed.to_string()) } async fn load_asset_responses( state: &AppState, version_id: Uuid, ) -> AppResult> { let mut conn = state.db()?; let assets: Vec = document_assets::table .filter(document_assets::document_version_id.eq(version_id)) .order(document_assets::created_at.asc()) .load(&mut conn)?; drop(conn); let mut responses = Vec::with_capacity(assets.len()); for asset in assets { let url = state .storage .presign_get_object( &asset.s3_key, Duration::from_secs(PRESIGNED_URL_EXPIRY_SECONDS), ) .await .map_err(|err| AppError::internal(format!("failed to sign asset URL: {err}")))?; responses.push(to_asset_response(asset, url)); } Ok(responses) } pub(crate) fn to_iso(dt: NaiveDateTime) -> String { DateTime::::from_naive_utc_and_offset(dt, Utc).to_rfc3339() }