caps and delete

This commit is contained in:
2025-11-05 12:32:10 +01:00
parent 4ec19dbd70
commit 480bc20ae7
30 changed files with 2702 additions and 262 deletions
+2
View File
@@ -16,6 +16,7 @@ pub mod analyze;
pub mod common;
pub mod index;
pub mod ocr;
pub mod purge;
pub mod tenants;
pub mod thumbnails;
@@ -149,6 +150,7 @@ pub fn default_handlers() -> Vec<Arc<dyn JobHandler>> {
Arc::new(analyze::AnalyzeDocumentJob::new()),
Arc::new(thumbnails::GenerateThumbnailsJob::new()),
Arc::new(ocr::GenerateOcrTextJob::new()),
Arc::new(purge::PurgeDocumentJob::new()),
Arc::new(index::IndexDocumentTextJob::new()),
Arc::new(ProvisionTenantJob::new()),
]
+239
View File
@@ -0,0 +1,239 @@
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<String>,
asset_keys: Vec<String>,
}
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<AppState>,
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<AppState>,
tenant_id: Uuid,
document_id: Uuid,
) -> Result<Option<PurgeContext>, 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::<Document>(conn)
.optional()?;
let Some(document) = doc_opt else {
return Ok(None);
};
if document.deleted_at.is_none() {
return Ok(None);
}
let versions: Vec<DocumentVersion> = 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<String> = versions
.iter()
.map(|version| version.s3_key.clone())
.collect();
let version_ids: Vec<Uuid> = versions.iter().map(|version| version.id).collect();
let asset_keys = if version_ids.is_empty() {
Vec::new()
} else {
let asset_ids: Vec<Uuid> = 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<AppState>, 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::<Document>(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}"))
}
+47
View File
@@ -8,6 +8,9 @@ use serde::Deserialize;
use tracing::warn;
use uuid::Uuid;
use crate::auth::capability_sets::{
ensure_capability_set, owner_capabilities, user_capabilities, webdav_capabilities,
};
use crate::documents::search::ensure_quickwit_index;
use crate::jobs::JOB_PROVISION_TENANT;
use crate::models::{NewUserMembership, TenantStatus};
@@ -128,12 +131,56 @@ impl JobHandler for ProvisionTenantJob {
};
}
let owner_capability_set_id =
match ensure_capability_set(&mut conn, tenant.id, owner_capabilities()) {
Ok(set) => set.id,
Err(err) => {
warn!(
job_id = %job.id,
tenant_id = %tenant.id,
error = ?err,
"failed to ensure owner capability set during provisioning"
);
return JobExecution::Retry {
delay: std::time::Duration::from_secs(30),
error: "owner capability set unavailable".into(),
};
}
};
if let Err(err) = ensure_capability_set(&mut conn, tenant.id, user_capabilities()) {
warn!(
job_id = %job.id,
tenant_id = %tenant.id,
error = ?err,
"failed to ensure user capability set during provisioning"
);
return JobExecution::Retry {
delay: std::time::Duration::from_secs(30),
error: "user capability set unavailable".into(),
};
}
if let Err(err) = ensure_capability_set(&mut conn, tenant.id, webdav_capabilities()) {
warn!(
job_id = %job.id,
tenant_id = %tenant.id,
error = ?err,
"failed to ensure webdav capability set during provisioning"
);
return JobExecution::Retry {
delay: std::time::Duration::from_secs(30),
error: "webdav capability set unavailable".into(),
};
}
if let Some(members) = ProvisionPayload::from_job(&job) {
for member in members {
let new_membership = NewUserMembership {
id: Uuid::new_v4(),
user_id: member,
tenant_id: tenant.id,
capability_set_id: Some(owner_capability_set_id),
};
if let Err(err) = diesel::insert_into(user_memberships::table)