Files
papercrate/backend/src/routes/documents.rs
T
nils c17a9484ca
ci / docker (backend, backend/Dockerfile, backend) (push) Failing after 9m32s
ci / docker (frontend, frontend/Dockerfile, frontend) (push) Successful in 9m48s
a lot of things
2025-10-14 00:56:25 +02:00

1296 lines
38 KiB
Rust

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<String> {
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<Uuid>,
#[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<String>,
}
impl From<Tag> for TagResponse {
fn from(tag: Tag) -> Self {
Self {
id: tag.id,
label: tag.label,
color: tag.color,
}
}
}
#[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 metadata: Value,
#[serde(skip_serializing_if = "Option::is_none")]
pub url: Option<String>,
pub created_at: String,
}
#[derive(Serialize, Clone)]
pub struct DocumentCurrentVersionResponse {
#[serde(flatten)]
pub version: DocumentVersionResponse,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub assets: Vec<DocumentAssetResponse>,
pub download_path: String,
}
#[derive(Serialize)]
pub struct DocumentResponse {
pub id: Uuid,
pub filename: String,
pub title: String,
pub original_name: String,
pub content_type: Option<String>,
pub folder_id: Option<Uuid>,
pub uploaded_at: String,
pub updated_at: String,
pub deleted_at: Option<String>,
pub issued_at: Option<String>,
pub metadata: Value,
pub tags: Vec<TagResponse>,
#[serde(skip_serializing_if = "Option::is_none")]
pub current_version: Option<DocumentCurrentVersionResponse>,
}
#[derive(Serialize)]
pub struct DocumentDetailResponse {
pub document: DocumentResponse,
}
#[derive(Serialize)]
pub struct DocumentDownloadResponse {
pub url: String,
pub expires_in: u64,
pub filename: String,
pub content_type: Option<String>,
pub size_bytes: i64,
}
#[derive(Serialize)]
pub struct BulkReanalyzeResponse {
pub queued: usize,
}
#[derive(Deserialize)]
pub struct BulkMoveRequest {
pub document_ids: Vec<Uuid>,
pub folder_id: Option<Uuid>,
}
#[derive(Deserialize)]
pub struct UpdateDocumentRequest {
pub title: Option<String>,
}
#[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<Uuid>,
pub tag_ids: Vec<Uuid>,
pub action: BulkTagAction,
}
#[derive(Serialize)]
pub struct BulkTagResponse {
pub added: usize,
pub removed: usize,
}
#[derive(Deserialize)]
pub struct BulkReanalyzeSelectionRequest {
pub document_ids: Vec<Uuid>,
#[serde(default = "default_true")]
pub force: bool,
}
fn default_true() -> bool {
true
}
struct UploadRequest {
bytes: Vec<u8>,
original_name: String,
content_type: Option<String>,
folder_id: Option<Uuid>,
metadata: Value,
}
struct UploadOutcome {
detail: DocumentDetailResponse,
created: bool,
}
#[derive(Deserialize)]
pub struct MoveDocumentRequest {
pub folder_id: Option<Uuid>,
}
#[derive(Deserialize)]
pub struct AssignTagsRequest {
pub tag_ids: Vec<Uuid>,
}
pub async fn list_documents(
State(state): State<AppState>,
Query(query): Query<DocumentListQuery>,
user: AuthenticatedUser,
) -> AppResult<Json<Vec<DocumentResponse>>> {
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<Document> = base_query
.order(documents::uploaded_at.desc())
.load(&mut conn)?;
let doc_ids: Vec<Uuid> = docs.iter().map(|doc| doc.id).collect();
let tags_map = load_tags_for_documents(&mut conn, &doc_ids)?;
drop(conn);
let primary_versions = load_primary_assets(&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 current_version = primary_versions.get(&doc.id).cloned();
response.push(to_document_response(
&state,
user.user_id,
doc,
tags,
current_version,
)?);
}
Ok(Json(response))
}
pub async fn get_document(
State(state): State<AppState>,
Path(document_id): Path<Uuid>,
user: AuthenticatedUser,
) -> AppResult<Json<DocumentDetailResponse>> {
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
.find(doc.current_version_id)
.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 version_response = to_version_response(current_version);
Ok(Json(DocumentDetailResponse {
document: to_document_response(
&state,
user.user_id,
doc,
tags_map.get(&document_id).cloned(),
Some((version_response, assets)),
)?,
}))
}
pub async fn upload_document(
State(state): State<AppState>,
user: AuthenticatedUser,
mut multipart: Multipart,
) -> AppResult<(StatusCode, Json<DocumentDetailResponse>)> {
let mut file_bytes: Option<Vec<u8>> = None;
let mut original_name: Option<String> = None;
let mut content_type: Option<String> = None;
let mut folder_id: Option<Uuid> = 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<AppState>,
Path(document_id): Path<Uuid>,
Query(query): Query<AssetRequestQuery>,
) -> AppResult<StatusCode> {
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());
}
enqueue_job(
&mut conn,
JOB_ANALYZE_DOCUMENT,
json!({
"document_id": document_id,
"document_version_id": document.current_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<AppState>,
_user: AuthenticatedUser,
) -> AppResult<(StatusCode, Json<BulkReanalyzeResponse>)> {
let mut conn = state.db()?;
let targets: Vec<(Uuid, Uuid)> = documents::table
.filter(documents::deleted_at.is_null())
.select((documents::id, documents::current_version_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<AppState>,
Json(payload): Json<BulkReanalyzeSelectionRequest>,
) -> AppResult<(StatusCode, Json<BulkReanalyzeResponse>)> {
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)> = documents::table
.filter(documents::id.eq_any(&document_ids))
.filter(documents::deleted_at.is_null())
.select((documents::id, documents::current_version_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<AppState>,
Path(document_id): Path<Uuid>,
) -> AppResult<Json<Vec<DocumentAssetResponse>>> {
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_id = document.current_version_id;
drop(conn);
let assets = load_asset_responses(&state, version_id).await?;
Ok(Json(assets))
}
pub async fn get_document_asset(
State(state): State<AppState>,
Path((document_id, asset_id)): Path<(Uuid, Uuid)>,
) -> AppResult<Json<DocumentAssetResponse>> {
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 asset: DocumentAsset = document_assets::table.find(asset_id).first(&mut conn)?;
let version: DocumentVersion = document_versions::table
.find(asset.document_version_id)
.first(&mut conn)?;
if version.document_id != document_id {
return Err(AppError::not_found());
}
let s3_key = asset.s3_key.clone();
drop(conn);
let presigned_url = state
.storage
.presign_get_object(&s3_key, Duration::from_secs(PRESIGNED_URL_EXPIRY_SECONDS))
.await
.map_err(|err| AppError::internal(format!("failed to generate asset URL: {err}")))?;
Ok(Json(to_asset_response(asset, Some(presigned_url))))
}
pub async fn download_document(
State(state): State<AppState>,
Path(document_id): Path<Uuid>,
) -> AppResult<Json<DocumentDownloadResponse>> {
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
.find(doc.current_version_id)
.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<AppState>,
Path(token): Path<String>,
) -> AppResult<impl IntoResponse> {
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
.find(doc.current_version_id)
.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<AppState>,
Path(document_id): Path<Uuid>,
) -> AppResult<impl IntoResponse> {
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 update_document(
State(state): State<AppState>,
Path(document_id): Path<Uuid>,
user: AuthenticatedUser,
Json(payload): Json<UpdateDocumentRequest>,
) -> AppResult<Json<DocumentDetailResponse>> {
let mut conn = state.db()?;
let mut document: Document = documents::table.find(document_id).first(&mut conn)?;
if document.deleted_at.is_some() {
return Err(AppError::not_found());
}
let new_title = match payload.title {
Some(ref title) => {
let trimmed = title.trim();
if trimmed.is_empty() {
return Err(AppError::bad_request("title must not be empty"));
}
Some(trimmed.to_string())
}
None => None,
};
if new_title.is_none() {
return Err(AppError::bad_request("no changes provided"));
}
if let Some(title) = new_title {
let now = Utc::now().naive_utc();
diesel::update(documents::table.find(document_id))
.set((documents::title.eq(title), documents::updated_at.eq(now)))
.execute(&mut conn)?;
document = documents::table.find(document_id).first(&mut conn)?;
}
let current_version: DocumentVersion = document_versions::table
.find(document.current_version_id)
.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 version_response = to_version_response(current_version);
Ok(Json(DocumentDetailResponse {
document: to_document_response(
&state,
user.user_id,
document,
tags_map.get(&document_id).cloned(),
Some((version_response, assets)),
)?,
}))
}
pub async fn move_document(
State(state): State<AppState>,
Path(document_id): Path<Uuid>,
Json(payload): Json<MoveDocumentRequest>,
) -> AppResult<impl IntoResponse> {
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<AppState>,
Json(payload): Json<BulkMoveRequest>,
) -> AppResult<(StatusCode, Json<BulkMoveResponse>)> {
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<NaiveDateTime>)> = 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<AppState>,
Path(document_id): Path<Uuid>,
user: AuthenticatedUser,
Json(payload): Json<AssignTagsRequest>,
) -> AppResult<impl IntoResponse> {
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::<Document>(&mut conn)?;
// Ensure tags exist
let existing_tags: Vec<Tag> = 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<NewDocumentTag> = 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<AppState>,
user: AuthenticatedUser,
Json(payload): Json<BulkTagRequest>,
) -> AppResult<(StatusCode, Json<BulkTagResponse>)> {
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<NaiveDateTime>)> = 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<Tag> = 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<AppState>,
Path((document_id, tag_id)): Path<(Uuid, Uuid)>,
) -> AppResult<impl IntoResponse> {
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<UploadOutcome> {
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::id.eq(documents::current_version_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 document.deleted_at.is_some() {
let now = Utc::now().naive_utc();
diesel::update(documents::table.find(document.id))
.set((
documents::deleted_at.eq(None::<NaiveDateTime>),
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 version_response = to_version_response(version.clone());
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,
Some((version_response, 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_id: version_id,
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,
Some((to_version_response(version.clone()), 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<HashMap<Uuid, Vec<Tag>>> {
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<Uuid, Vec<Tag>> = HashMap::new();
for (doc_id, tag) in rows {
map.entry(doc_id).or_default().push(tag);
}
Ok(map)
}
pub(crate) async fn load_primary_assets(
state: &AppState,
documents: &[Document],
) -> AppResult<HashMap<Uuid, (DocumentVersionResponse, Vec<DocumentAssetResponse>)>> {
if documents.is_empty() {
return Ok(HashMap::new());
}
let mut doc_to_version: HashMap<Uuid, Uuid> = HashMap::with_capacity(documents.len());
let mut version_ids: Vec<Uuid> = Vec::with_capacity(documents.len());
for doc in documents {
doc_to_version.insert(doc.id, doc.current_version_id);
version_ids.push(doc.current_version_id);
}
version_ids.sort();
version_ids.dedup();
let mut conn = state.db()?;
let versions: Vec<DocumentVersion> = document_versions::table
.filter(document_versions::id.eq_any(&version_ids))
.load(&mut conn)?;
let mut version_map: HashMap<Uuid, DocumentVersion> = HashMap::new();
for version in versions {
version_map.insert(version.id, version);
}
let assets: Vec<DocumentAsset> = document_assets::table
.filter(document_assets::document_version_id.eq_any(&version_ids))
.order((
document_assets::document_version_id.asc(),
document_assets::created_at.asc(),
))
.load(&mut conn)?;
let mut assets_by_version: HashMap<Uuid, Vec<DocumentAssetResponse>> = HashMap::new();
for asset in assets {
let version_id = asset.document_version_id;
let response = to_asset_response(asset, None);
assets_by_version
.entry(version_id)
.or_default()
.push(response);
}
drop(conn);
let mut result: HashMap<Uuid, (DocumentVersionResponse, Vec<DocumentAssetResponse>)> =
HashMap::with_capacity(doc_to_version.len());
for (doc_id, version_id) in doc_to_version {
if let Some(version) = version_map.remove(&version_id) {
let assets = assets_by_version.remove(&version_id).unwrap_or_default();
result.insert(doc_id, (to_version_response(version), assets));
}
}
Ok(result)
}
pub(crate) fn to_document_response(
state: &AppState,
user_id: Uuid,
doc: Document,
tags: Option<Vec<Tag>>,
current_version: Option<(DocumentVersionResponse, Vec<DocumentAssetResponse>)>,
) -> AppResult<DocumentResponse> {
let current_version = if let Some((version, assets)) = current_version {
let download_path = build_download_path(state, doc.id, user_id)?;
Some(DocumentCurrentVersionResponse {
version,
assets,
download_path,
})
} else {
None
};
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,
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(),
current_version,
})
}
fn build_download_path(state: &AppState, document_id: Uuid, user_id: Uuid) -> AppResult<String> {
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: Option<String>) -> DocumentAssetResponse {
let metadata = asset.metadata.clone();
DocumentAssetResponse {
id: asset.id,
asset_type: asset.asset_type,
mime_type: asset.mime_type,
metadata,
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<Vec<DocumentAssetResponse>> {
let mut conn = state.db()?;
let assets: Vec<DocumentAsset> = document_assets::table
.filter(document_assets::document_version_id.eq(version_id))
.order(document_assets::created_at.asc())
.load(&mut conn)?;
drop(conn);
Ok(assets
.into_iter()
.map(|asset| to_asset_response(asset, None))
.collect())
}
pub(crate) fn to_iso(dt: NaiveDateTime) -> String {
DateTime::<Utc>::from_naive_utc_and_offset(dt, Utc).to_rfc3339()
}