tenant gate active

This commit is contained in:
2025-11-11 02:45:05 +01:00
parent ac9e688e7a
commit 4f29ff73c6
12 changed files with 168 additions and 123 deletions
+8 -1
View File
@@ -6,7 +6,8 @@ use serde::Deserialize;
use uuid::Uuid;
use crate::{
jobs::JOB_ANALYZE_DOCUMENT, models::Document, state::AppState, storage::TenantStorage,
auth::ensure_active_tenant, jobs::JOB_ANALYZE_DOCUMENT, models::Document, state::AppState,
storage::TenantStorage,
};
use super::{
@@ -49,6 +50,12 @@ impl JobHandler for AnalyzeDocumentJob {
job: crate::models::Job,
storage: TenantStorage,
) -> JobExecution {
if let Err(err) = ensure_active_tenant(&state, job.tenant_id) {
return JobExecution::Failed {
error: err.to_string(),
};
}
let payload: AnalyzePayload = match serde_json::from_value(job.payload.clone()) {
Ok(payload) => payload,
Err(err) => {
+2 -88
View File
@@ -1,101 +1,15 @@
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use reqwest::Client;
use serde::Deserialize;
use uuid::Uuid;
use crate::{
documents::search::{build_quickwit_ingest_record, quickwit_ingest},
jobs::JOB_INDEX_DOCUMENT_TEXT,
state::AppState,
storage::TenantStorage,
};
use crate::documents::search::{build_quickwit_ingest_record, quickwit_ingest};
use super::{
job_execution_from_task_error,
ocr::OCR_TEXT_ASSET_TYPE,
taskflow::{
document::DocumentVersionTaskContext, BoxedTask, Task, TaskError, TaskExecutor,
TaskPlanner, TaskResult,
},
JobExecution, JobHandler,
taskflow::{document::DocumentVersionTaskContext, Task, TaskError, TaskResult},
};
#[derive(Debug, Deserialize)]
struct IndexPayload {
document_id: Uuid,
document_version_id: Uuid,
}
pub struct IndexDocumentTextJob;
impl IndexDocumentTextJob {
pub fn new() -> Self {
Self
}
}
#[async_trait]
impl JobHandler for IndexDocumentTextJob {
fn job_type(&self) -> &'static str {
JOB_INDEX_DOCUMENT_TEXT
}
async fn handle(
&self,
state: Arc<AppState>,
job: crate::models::Job,
storage: TenantStorage,
) -> JobExecution {
let payload: IndexPayload = match serde_json::from_value(job.payload.clone()) {
Ok(payload) => payload,
Err(err) => {
return JobExecution::Failed {
error: format!("invalid index payload: {err}"),
};
}
};
let mut context = DocumentVersionTaskContext::new(
job.id,
JOB_INDEX_DOCUMENT_TEXT,
job.tenant_id,
payload.document_id,
payload.document_version_id,
false,
state.config.worker_max_document_bytes,
state.clone(),
storage,
);
let planner = IndexPlanner::new();
match TaskExecutor::run(&planner, &mut context).await {
Ok(()) => JobExecution::Success,
Err(err) => job_execution_from_task_error(err),
}
}
}
struct IndexPlanner;
impl IndexPlanner {
fn new() -> Self {
Self
}
}
#[async_trait]
impl TaskPlanner<DocumentVersionTaskContext> for IndexPlanner {
async fn plan(
&self,
_ctx: &mut DocumentVersionTaskContext,
) -> TaskResult<Vec<BoxedTask<DocumentVersionTaskContext>>> {
Ok(vec![Box::new(IndexDocumentTask::new())])
}
}
pub struct IndexDocumentTask;
impl IndexDocumentTask {
-1
View File
@@ -152,7 +152,6 @@ pub fn default_handlers() -> Vec<Arc<dyn JobHandler>> {
vec![
Arc::new(analyze::AnalyzeDocumentJob::new()),
Arc::new(purge::PurgeDocumentJob::new()),
Arc::new(index::IndexDocumentTextJob::new()),
Arc::new(ProvisionTenantJob::new()),
]
}
+7
View File
@@ -8,6 +8,7 @@ use diesel::result::Error as DieselError;
use serde::Deserialize;
use uuid::Uuid;
use crate::auth::ensure_active_tenant;
use crate::jobs::JOB_PURGE_DOCUMENT;
use crate::models::{Document, DocumentVersion};
use crate::schema::{document_asset_objects, document_assets, document_versions};
@@ -52,6 +53,12 @@ impl JobHandler for PurgeDocumentJob {
job: crate::models::Job,
storage: TenantStorage,
) -> JobExecution {
if let Err(err) = ensure_active_tenant(&state, job.tenant_id) {
return JobExecution::Failed {
error: err.to_string(),
};
}
let payload: PurgeDocumentPayload = match serde_json::from_value(job.payload.clone()) {
Ok(payload) => payload,
Err(err) => {