From 7f53bc900aef3a025fbf1512bcdb53f00685ac01 Mon Sep 17 00:00:00 2001 From: Nils Schneider Date: Thu, 6 Nov 2025 22:52:25 +0100 Subject: [PATCH] refactor --- backend/src/documents/asset.rs | 13 +- backend/src/routes/documents.rs | 7 +- backend/src/services/documents.rs | 372 ++++++++++++++++-------------- 3 files changed, 211 insertions(+), 181 deletions(-) diff --git a/backend/src/documents/asset.rs b/backend/src/documents/asset.rs index f75657d..546cbdb 100644 --- a/backend/src/documents/asset.rs +++ b/backend/src/documents/asset.rs @@ -10,7 +10,7 @@ use uuid::Uuid; use crate::error::{AppError, AppResult}; use crate::models::{Document, DocumentAsset, DocumentAssetObject, DocumentVersion}; use crate::schema::{document_asset_objects, document_assets, document_versions}; -use crate::state::AppState; +use crate::state::{AppState, PgPooledConnection}; use crate::utils::time::to_iso; #[derive(Serialize, Clone, ToSchema)] @@ -158,6 +158,14 @@ pub async fn load_asset_responses( version_id: Uuid, ) -> AppResult> { let mut conn = state.db_for_tenant(tenant_id)?; + load_asset_responses_with_conn(&mut conn, tenant_id, version_id) +} + +pub fn load_asset_responses_with_conn( + conn: &mut PgPooledConnection, + tenant_id: Uuid, + version_id: Uuid, +) -> AppResult> { let assets: Vec<(DocumentAsset, Option)> = document_assets::table .left_outer_join( document_asset_objects::table.on(document_asset_objects::asset_id @@ -171,8 +179,7 @@ pub async fn load_asset_responses( document_assets::all_columns, document_asset_objects::all_columns.nullable(), )) - .load(&mut conn)?; - drop(conn); + .load(conn)?; Ok(assets .into_iter() diff --git a/backend/src/routes/documents.rs b/backend/src/routes/documents.rs index 51ba3fd..16ad92c 100644 --- a/backend/src/routes/documents.rs +++ b/backend/src/routes/documents.rs @@ -186,11 +186,12 @@ pub async fn get_document( )] pub async fn upload_document( State(state): State, - TenantScopedConn { - tenant_id, user_id, .. - }: TenantScopedConn, + scoped_conn: TenantScopedConn, mut multipart: Multipart, ) -> AppResult> { + let tenant_id = scoped_conn.tenant_id; + let user_id = scoped_conn.user_id; + drop(scoped_conn); let mut file_bytes: Option> = None; let mut original_name: Option = None; let mut content_type: Option = None; diff --git a/backend/src/services/documents.rs b/backend/src/services/documents.rs index 7bf9bd5..ef1997c 100644 --- a/backend/src/services/documents.rs +++ b/backend/src/services/documents.rs @@ -20,9 +20,10 @@ use uuid::Uuid; use crate::documents::{ asset::{ build_download_path, derive_document_title, filename_with_retained_extension, - load_asset_responses, load_primary_assets, to_asset_detail_response, - to_asset_object_response, to_version_response, DocumentAssetDetailResponse, - DocumentAssetResponse, DocumentVersionDetailResponse, DocumentVersionResponse, + load_asset_responses, load_asset_responses_with_conn, load_primary_assets, + to_asset_detail_response, to_asset_object_response, to_version_response, + DocumentAssetDetailResponse, DocumentAssetResponse, DocumentVersionDetailResponse, + DocumentVersionResponse, }, correspondents::{ insert_document_correspondents, normalize_correspondent_ids, DocumentCorrespondentResponse, @@ -660,9 +661,10 @@ impl<'a> DocumentsService<'a> { skip_if_existing, } = request; + let mut tenant_conn = self.state.db_for_tenant(tenant_id)?; + if let Some(folder) = folder_id { - let mut conn = self.state.db_for_tenant(tenant_id)?; - ensure_folder_exists_on_conn(&mut conn, tenant_id, folder)?; + ensure_folder_exists_on_conn(&mut tenant_conn, tenant_id, folder)?; } let doc_id = Uuid::new_v4(); @@ -681,106 +683,20 @@ impl<'a> DocumentsService<'a> { let size_bytes = bytes.len() as i64; let s3_key = document_version_object_key(doc_id, version_number, version_id); + if let Some(reused) = self + .try_reuse_existing_document( + &mut tenant_conn, + tenant_id, + user_id, + &checksum_hex, + issued_at_override, + skip_if_existing, + &tag_ids, + &correspondents, + ) + .await? { - let mut conn = self.state.db_for_tenant(tenant_id)?; - - let existing = documents::table - .inner_join( - document_versions::table - .on(document_versions::id.eq(documents::current_version_id)), - ) - .filter(documents::tenant_id.eq(tenant_id)) - .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 skip_if_existing { - info!( - document_id = %document.id, - checksum = %checksum_hex, - "upload rejected because document already exists", - ); - return Err(AppError::conflict( - "a document with the same contents already exists", - ) - .with_code("duplicate_document") - .with_details(json!({ - "conflict_document_id": document.id, - }))); - } - - if let Some(issued_at) = issued_at_override { - if document.issued_at != Some(issued_at) { - diesel::update( - documents::table - .find(document.id) - .filter(documents::tenant_id.eq(tenant_id)), - ) - .set(( - documents::issued_at.eq(Some(issued_at)), - documents::updated_at.eq(Utc::now().naive_utc()), - )) - .execute(&mut conn)?; - document.issued_at = Some(issued_at); - } - } - - assign_tags_to_document(&mut conn, tenant_id, &document, &tag_ids, Some(user_id))?; - - if !correspondents.is_empty() { - let raw_ids: Vec = correspondents - .iter() - .map(|assignment| assignment.correspondent_id) - .collect(); - let correspondent_ids = normalize_correspondent_ids(&raw_ids)?; - insert_document_correspondents( - &mut conn, - tenant_id, - document.id, - user_id, - &correspondent_ids, - )?; - } - - 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_and_correspondents = - load_tags_and_correspondents(&mut conn, &[document.id])?; - let (tags, correspondents) = tags_and_correspondents - .get(&document.id) - .cloned() - .unwrap_or_else(|| (Vec::new(), Vec::new())); - let assets = load_asset_responses(self.state, tenant_id, version.id).await?; - let version_response = to_version_response(version.clone()); - - info!( - document_id = %document.id, - checksum = %checksum_hex, - "upload deduplicated existing document" - ); - - return Ok(DocumentUploadOutcome::Reused(DocumentDetailResponse { - document: self.to_document_response( - user_id, - document, - tags, - correspondents, - Some((version_response, assets)), - )?, - })); - } + return Ok(DocumentUploadOutcome::Reused(reused)); } let content_disposition = inline_content_disposition(&stored_filename); @@ -803,66 +719,61 @@ impl<'a> DocumentsService<'a> { metadata }; - let (document, version) = { - let mut conn = self.state.db_for_tenant(tenant_id)?; - let transaction_result = 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_id: version_id, - metadata: metadata_value.clone(), - issued_at: issued_at_override, - title: derived_title.clone(), - tenant_id, - }; - diesel::insert_into(documents::table) - .values(&new_document) - .execute(conn)?; + let (document, version) = match tenant_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_id: version_id, + metadata: metadata_value.clone(), + issued_at: issued_at_override, + title: derived_title.clone(), + tenant_id, + }; + 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(), - metadata: Value::Object(Default::default()), - tenant_id, - }; + let new_version = NewDocumentVersion { + id: version_id, + document_id: doc_id, + version_number, + s3_key: s3_key.clone(), + size_bytes, + checksum: checksum_hex.clone(), + metadata: Value::Object(Default::default()), + tenant_id, + }; - diesel::insert_into(document_versions::table) - .values(&new_version) - .execute(conn)?; + 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)?; + 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)) - }); - - match transaction_result { - Ok(result) => result, - Err(diesel::result::Error::DatabaseError( - DatabaseErrorKind::UniqueViolation, - _, - )) => { - return Err(AppError::conflict( - "another document in this folder already uses that filename", - ) - .with_code("duplicate_filename")) - } - Err(err) => return Err(AppError::from(err)), + Ok::<_, diesel::result::Error>((document, version)) + }) { + Ok(result) => result, + Err(diesel::result::Error::DatabaseError(DatabaseErrorKind::UniqueViolation, _)) => { + return Err(AppError::conflict( + "another document in this folder already uses that filename", + ) + .with_code("duplicate_filename")) } + Err(err) => return Err(AppError::from(err)), }; let detail = { - let mut conn = self.state.db_for_tenant(tenant_id)?; - - assign_tags_to_document(&mut conn, tenant_id, &document, &tag_ids, Some(user_id))?; + assign_tags_to_document( + &mut tenant_conn, + tenant_id, + &document, + &tag_ids, + Some(user_id), + )?; if !correspondents.is_empty() { let raw_ids: Vec = correspondents @@ -871,7 +782,7 @@ impl<'a> DocumentsService<'a> { .collect(); let correspondent_ids = normalize_correspondent_ids(&raw_ids)?; insert_document_correspondents( - &mut conn, + &mut tenant_conn, tenant_id, document.id, user_id, @@ -879,7 +790,8 @@ impl<'a> DocumentsService<'a> { )?; } - let tags_and_correspondents = load_tags_and_correspondents(&mut conn, &[doc_id])?; + let tags_and_correspondents = + load_tags_and_correspondents(&mut tenant_conn, &[doc_id])?; let (tags, correspondents) = tags_and_correspondents .get(&doc_id) .cloned() @@ -896,22 +808,18 @@ impl<'a> DocumentsService<'a> { } }; - if let Ok(mut conn) = self.state.db_for_tenant(tenant_id) { - if let Err(err) = enqueue_job( - &mut conn, - tenant_id, - 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"); + if let Err(err) = enqueue_job( + &mut tenant_conn, + tenant_id, + 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"); } Ok(DocumentUploadOutcome::Created(detail)) @@ -1106,7 +1014,7 @@ impl<'a> DocumentsService<'a> { .filter(document_versions::tenant_id.eq(tenant_id)) .first(conn)?; - let assets = load_asset_responses(self.state, tenant_id, version.id).await?; + let assets = load_asset_responses_with_conn(conn, tenant_id, version.id)?; let download_path = build_download_path(self.state, &document, user_id)?; let version_core = to_version_response(version); @@ -1450,6 +1358,120 @@ impl<'a> DocumentsService<'a> { current_version, }) } + + async fn try_reuse_existing_document( + &self, + conn: &mut PgPooledConnection, + tenant_id: Uuid, + user_id: Uuid, + checksum_hex: &str, + issued_at_override: Option, + skip_if_existing: bool, + tag_ids: &[Uuid], + correspondents: &[CorrespondentAssignmentInput], + ) -> AppResult> { + let existing = documents::table + .inner_join( + document_versions::table + .on(document_versions::id.eq(documents::current_version_id)), + ) + .filter(documents::tenant_id.eq(tenant_id)) + .filter(document_versions::checksum.eq(checksum_hex)) + .select((documents::all_columns, document_versions::all_columns)) + .first::<(Document, DocumentVersion)>(conn) + .optional()?; + + let Some((mut document, version)) = existing else { + return Ok(None); + }; + + if skip_if_existing { + info!( + document_id = %document.id, + checksum = %checksum_hex, + "upload rejected because document already exists", + ); + return Err( + AppError::conflict("a document with the same contents already exists") + .with_code("duplicate_document") + .with_details(json!({ + "conflict_document_id": document.id, + })), + ); + } + + if let Some(issued_at) = issued_at_override { + if document.issued_at != Some(issued_at) { + diesel::update( + documents::table + .find(document.id) + .filter(documents::tenant_id.eq(tenant_id)), + ) + .set(( + documents::issued_at.eq(Some(issued_at)), + documents::updated_at.eq(Utc::now().naive_utc()), + )) + .execute(conn)?; + document.issued_at = Some(issued_at); + } + } + + assign_tags_to_document(conn, tenant_id, &document, tag_ids, Some(user_id))?; + + if !correspondents.is_empty() { + let raw_ids: Vec = correspondents + .iter() + .map(|assignment| assignment.correspondent_id) + .collect(); + let correspondent_ids = normalize_correspondent_ids(&raw_ids)?; + insert_document_correspondents( + conn, + tenant_id, + document.id, + user_id, + &correspondent_ids, + )?; + } + + 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(conn)?; + document.deleted_at = None; + document.updated_at = now; + } + + let relations = load_tags_and_correspondents(conn, &[document.id])?; + let (tags, correspondents_list) = relations + .get(&document.id) + .cloned() + .unwrap_or_else(|| (Vec::new(), Vec::new())); + + let assets = load_asset_responses_with_conn(conn, tenant_id, version.id)?; + let version_response = to_version_response(version.clone()); + + info!( + document_id = %document.id, + checksum = %checksum_hex, + "upload deduplicated existing document" + ); + + let detail = DocumentDetailResponse { + document: self.to_document_response( + user_id, + document, + tags, + correspondents_list, + Some((version_response, assets)), + )?, + }; + + Ok(Some(detail)) + } } fn presign_disposition_for_asset(