more multi-tenancy
This commit is contained in:
@@ -14,6 +14,7 @@ use crate::{
|
||||
models::{Document, DocumentAsset, DocumentVersion},
|
||||
schema::{document_assets, document_versions, documents},
|
||||
state::AppState,
|
||||
storage::TenantStorage,
|
||||
};
|
||||
|
||||
use super::{JobExecution, JobHandler};
|
||||
@@ -40,7 +41,12 @@ impl JobHandler for AnalyzeDocumentJob {
|
||||
JOB_ANALYZE_DOCUMENT
|
||||
}
|
||||
|
||||
async fn handle(&self, state: Arc<AppState>, job: crate::models::Job) -> JobExecution {
|
||||
async fn handle(
|
||||
&self,
|
||||
state: Arc<AppState>,
|
||||
job: crate::models::Job,
|
||||
_storage: TenantStorage,
|
||||
) -> JobExecution {
|
||||
let payload: AnalyzePayload = match serde_json::from_value(job.payload.clone()) {
|
||||
Ok(payload) => payload,
|
||||
Err(err) => {
|
||||
|
||||
@@ -15,6 +15,7 @@ use crate::{
|
||||
models::{Document, DocumentVersion},
|
||||
schema::{document_asset_objects, document_assets, document_versions, documents},
|
||||
state::AppState,
|
||||
storage::TenantStorage,
|
||||
};
|
||||
|
||||
use super::{ocr::OCR_TEXT_ASSET_TYPE, JobExecution, JobHandler};
|
||||
@@ -39,7 +40,12 @@ impl JobHandler for IndexDocumentTextJob {
|
||||
JOB_INDEX_DOCUMENT_TEXT
|
||||
}
|
||||
|
||||
async fn handle(&self, state: Arc<AppState>, job: crate::models::Job) -> JobExecution {
|
||||
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) => {
|
||||
@@ -52,16 +58,29 @@ impl JobHandler for IndexDocumentTextJob {
|
||||
let quickwit_endpoint = match &state.config.quickwit_endpoint {
|
||||
Some(endpoint) => endpoint.clone(),
|
||||
None => {
|
||||
warn!("quickwit endpoint missing; skipping indexing");
|
||||
return JobExecution::Success;
|
||||
return JobExecution::Failed {
|
||||
error: "quickwit endpoint missing".into(),
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
let quickwit_index = match &state.config.quickwit_index {
|
||||
Some(index) => index.clone(),
|
||||
let tenant = match state.tenants.get_by_id(job.tenant_id) {
|
||||
Ok(tenant) => tenant,
|
||||
Err(err) => {
|
||||
warn!(job_id = %job.id, error = ?err, "failed to load tenant for indexing");
|
||||
return JobExecution::Retry {
|
||||
delay: Duration::from_secs(30),
|
||||
error: format!("failed to load tenant: {err:?}"),
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
let quickwit_index = match tenant.quickwit_index.clone() {
|
||||
Some(index) => index,
|
||||
None => {
|
||||
warn!("quickwit index missing; skipping indexing");
|
||||
return JobExecution::Success;
|
||||
return JobExecution::Failed {
|
||||
error: "tenant quickwit index not configured".into(),
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
@@ -95,7 +114,7 @@ impl JobHandler for IndexDocumentTextJob {
|
||||
}
|
||||
|
||||
let s3_key = context.text_s3_key.unwrap();
|
||||
let text = match state.storage.get_object(&s3_key).await {
|
||||
let text = match storage.get_object(&s3_key).await {
|
||||
Ok(bytes) => match String::from_utf8(bytes) {
|
||||
Ok(text) => text,
|
||||
Err(err) => {
|
||||
@@ -129,6 +148,7 @@ impl JobHandler for IndexDocumentTextJob {
|
||||
let payload = json!({
|
||||
"document_id": context.document.id,
|
||||
"version_id": context.version.id,
|
||||
"tenant_id": job.tenant_id,
|
||||
"title": context.document.title.to_lowercase(),
|
||||
"text": text.to_lowercase()
|
||||
});
|
||||
|
||||
@@ -8,6 +8,7 @@ use crate::{
|
||||
jobs::{mark_job_failed, mark_job_succeeded, reserve_job, retry_job_after, JobQueueError},
|
||||
models::Job,
|
||||
state::AppState,
|
||||
storage::TenantStorage,
|
||||
};
|
||||
|
||||
pub mod analyze;
|
||||
@@ -25,7 +26,7 @@ pub enum JobExecution {
|
||||
#[async_trait]
|
||||
pub trait JobHandler: Send + Sync {
|
||||
fn job_type(&self) -> &'static str;
|
||||
async fn handle(&self, state: Arc<AppState>, job: Job) -> JobExecution;
|
||||
async fn handle(&self, state: Arc<AppState>, job: Job, storage: TenantStorage) -> JobExecution;
|
||||
}
|
||||
|
||||
pub struct Worker {
|
||||
@@ -84,8 +85,20 @@ impl Worker {
|
||||
|
||||
if let Some(job) = job_opt {
|
||||
if let Some(handler) = self.handlers.get(job.job_type.as_str()) {
|
||||
let result = handler.handle(self.state.clone(), job.clone()).await;
|
||||
match result {
|
||||
let execution = match self.state.storage_for_tenant(job.tenant_id) {
|
||||
Ok(storage) => {
|
||||
handler
|
||||
.handle(self.state.clone(), job.clone(), storage)
|
||||
.await
|
||||
}
|
||||
Err(err) => {
|
||||
error!(job_id = %job.id, error = ?err, "failed to load tenant storage for job");
|
||||
JobExecution::Failed {
|
||||
error: format!("tenant storage unavailable: {err:?}"),
|
||||
}
|
||||
}
|
||||
};
|
||||
match execution {
|
||||
JobExecution::Success => {
|
||||
if let Ok(mut conn) = self.state.db_unscoped() {
|
||||
mark_job_succeeded(&mut conn, job.id)?;
|
||||
|
||||
+12
-10
@@ -25,6 +25,7 @@ use crate::{
|
||||
},
|
||||
schema::{document_asset_objects, document_assets, document_versions, documents},
|
||||
state::AppState,
|
||||
storage::TenantStorage,
|
||||
utils::storage_paths::document_asset_object_prefix,
|
||||
};
|
||||
|
||||
@@ -55,7 +56,12 @@ impl JobHandler for GenerateOcrTextJob {
|
||||
JOB_GENERATE_OCR_TEXT
|
||||
}
|
||||
|
||||
async fn handle(&self, state: Arc<AppState>, job: crate::models::Job) -> JobExecution {
|
||||
async fn handle(
|
||||
&self,
|
||||
state: Arc<AppState>,
|
||||
job: crate::models::Job,
|
||||
storage: TenantStorage,
|
||||
) -> JobExecution {
|
||||
let payload: OcrPayload = match serde_json::from_value(job.payload.clone()) {
|
||||
Ok(payload) => payload,
|
||||
Err(err) => {
|
||||
@@ -92,7 +98,7 @@ impl JobHandler for GenerateOcrTextJob {
|
||||
return JobExecution::Success;
|
||||
}
|
||||
|
||||
let bytes = match state.storage.get_object(&context.version.s3_key).await {
|
||||
let bytes = match storage.get_object(&context.version.s3_key).await {
|
||||
Ok(bytes) => bytes,
|
||||
Err(err) => {
|
||||
warn!(job_id = %job.id, error = %err, "failed to fetch document for ocr");
|
||||
@@ -129,7 +135,7 @@ impl JobHandler for GenerateOcrTextJob {
|
||||
|
||||
if context.existing_asset.is_some() {
|
||||
for object in &context.existing_objects {
|
||||
if let Err(err) = state.storage.delete_object(&object.s3_key).await {
|
||||
if let Err(err) = storage.delete_object(&object.s3_key).await {
|
||||
warn!(job_id = %job.id, error = %err, s3_key = %object.s3_key, "failed to delete existing ocr asset object");
|
||||
}
|
||||
}
|
||||
@@ -144,8 +150,7 @@ impl JobHandler for GenerateOcrTextJob {
|
||||
asset_id,
|
||||
);
|
||||
|
||||
if let Err(err) = state
|
||||
.storage
|
||||
if let Err(err) = storage
|
||||
.put_object(
|
||||
&s3_key,
|
||||
generation.text.into_bytes(),
|
||||
@@ -168,11 +173,8 @@ impl JobHandler for GenerateOcrTextJob {
|
||||
.await
|
||||
{
|
||||
Ok(Ok(())) => {
|
||||
if state.config.quickwit_endpoint.is_some() && state.config.quickwit_index.is_some()
|
||||
{
|
||||
if let Err(err) = enqueue_index_job(&state, &payload) {
|
||||
warn!(job_id = %job.id, error = %err, "failed to enqueue index job");
|
||||
}
|
||||
if let Err(err) = enqueue_index_job(&state, &payload) {
|
||||
warn!(job_id = %job.id, error = %err, "failed to enqueue index job");
|
||||
}
|
||||
JobExecution::Success
|
||||
}
|
||||
|
||||
@@ -19,6 +19,7 @@ use crate::{
|
||||
},
|
||||
schema::{document_asset_objects, document_assets, document_versions, documents},
|
||||
state::AppState,
|
||||
storage::TenantStorage,
|
||||
utils::storage_paths::document_asset_object_key,
|
||||
};
|
||||
|
||||
@@ -53,7 +54,12 @@ impl JobHandler for GenerateThumbnailsJob {
|
||||
JOB_GENERATE_THUMBNAILS
|
||||
}
|
||||
|
||||
async fn handle(&self, state: Arc<AppState>, job: crate::models::Job) -> JobExecution {
|
||||
async fn handle(
|
||||
&self,
|
||||
state: Arc<AppState>,
|
||||
job: crate::models::Job,
|
||||
storage: TenantStorage,
|
||||
) -> JobExecution {
|
||||
let payload: ThumbnailPayload = match serde_json::from_value(job.payload.clone()) {
|
||||
Ok(p) => p,
|
||||
Err(err) => {
|
||||
@@ -89,7 +95,7 @@ impl JobHandler for GenerateThumbnailsJob {
|
||||
return JobExecution::Success;
|
||||
}
|
||||
|
||||
let bytes = match state.storage.get_object(&initial.version.s3_key).await {
|
||||
let bytes = match storage.get_object(&initial.version.s3_key).await {
|
||||
Ok(bytes) => bytes,
|
||||
Err(err) => {
|
||||
warn!(job_id = %job.id, error = %err, "thumbnail fetch failed; will retry");
|
||||
@@ -148,7 +154,7 @@ impl JobHandler for GenerateThumbnailsJob {
|
||||
|
||||
if initial.existing_preview.is_some() {
|
||||
for object in &initial.existing_preview_objects {
|
||||
if let Err(err) = state.storage.delete_object(&object.s3_key).await {
|
||||
if let Err(err) = storage.delete_object(&object.s3_key).await {
|
||||
warn!(
|
||||
job_id = %job.id,
|
||||
error = %err,
|
||||
@@ -161,7 +167,7 @@ impl JobHandler for GenerateThumbnailsJob {
|
||||
|
||||
if initial.existing_thumbnail.is_some() {
|
||||
for object in &initial.existing_thumbnail_objects {
|
||||
if let Err(err) = state.storage.delete_object(&object.s3_key).await {
|
||||
if let Err(err) = storage.delete_object(&object.s3_key).await {
|
||||
warn!(
|
||||
job_id = %job.id,
|
||||
error = %err,
|
||||
@@ -193,8 +199,7 @@ impl JobHandler for GenerateThumbnailsJob {
|
||||
ordinal,
|
||||
);
|
||||
|
||||
if let Err(err) = state
|
||||
.storage
|
||||
if let Err(err) = storage
|
||||
.put_object(
|
||||
&s3_key,
|
||||
image.image_bytes.clone(),
|
||||
@@ -235,8 +240,7 @@ impl JobHandler for GenerateThumbnailsJob {
|
||||
ordinal,
|
||||
);
|
||||
|
||||
if let Err(err) = state
|
||||
.storage
|
||||
if let Err(err) = storage
|
||||
.put_object(
|
||||
&s3_key,
|
||||
image.image_bytes.clone(),
|
||||
|
||||
Reference in New Issue
Block a user