This commit is contained in:
2025-11-11 01:51:42 +01:00
parent 4ce4bd1a4b
commit 425e03650f
3 changed files with 3 additions and 174 deletions
-2
View File
@@ -151,8 +151,6 @@ impl Worker {
pub fn default_handlers() -> Vec<Arc<dyn JobHandler>> {
vec![
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()),
+2 -92
View File
@@ -10,7 +10,6 @@ use async_trait::async_trait;
use chrono::Utc;
use diesel::{pg::upsert::excluded, prelude::*};
use pdfium_render::prelude::*;
use serde::Deserialize;
use serde_json::json;
use tempfile::NamedTempFile;
use tokio::task;
@@ -20,111 +19,22 @@ use uuid::Uuid;
use crate::{
documents::asset::delete_asset,
error::AppResult,
jobs::JOB_GENERATE_OCR_TEXT,
models::{
Document, DocumentAsset, DocumentAssetObject, DocumentVersion, NewDocumentAsset,
NewDocumentAssetObject,
},
schema::{document_asset_objects, document_assets},
state::AppState,
storage::TenantStorage,
utils::storage_paths::document_asset_object_prefix,
};
use super::{
index::IndexDocumentTask,
job_execution_from_task_error,
taskflow::{
document::DocumentVersionTaskContext, BoxedTask, Task, TaskContext, TaskError,
TaskExecutor, TaskPlanner, TaskResult,
},
JobExecution, JobHandler,
use super::taskflow::{
document::DocumentVersionTaskContext, Task, TaskContext, TaskError, TaskResult,
};
pub const OCR_TEXT_ASSET_TYPE: &str = "ocr-text";
const MIN_TEXT_LENGTH: usize = 50;
#[derive(Clone, Debug, Deserialize)]
struct OcrPayload {
document_id: Uuid,
document_version_id: Uuid,
#[serde(default)]
force: bool,
}
pub struct GenerateOcrTextJob;
impl GenerateOcrTextJob {
pub fn new() -> Self {
Self
}
}
#[async_trait]
impl JobHandler for GenerateOcrTextJob {
fn job_type(&self) -> &'static str {
JOB_GENERATE_OCR_TEXT
}
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) => {
return JobExecution::Failed {
error: format!("invalid OCR payload: {err}"),
};
}
};
let mut context = DocumentVersionTaskContext::new(
job.id,
JOB_GENERATE_OCR_TEXT,
job.tenant_id,
payload.document_id,
payload.document_version_id,
payload.force,
state.config.worker_max_document_bytes,
state.clone(),
storage,
);
let planner = OcrPlanner::new(payload.force, state.clone());
match TaskExecutor::run(&planner, &mut context).await {
Ok(()) => JobExecution::Success,
Err(err) => job_execution_from_task_error(err),
}
}
}
struct OcrPlanner {
force: bool,
state: Arc<AppState>,
}
impl OcrPlanner {
fn new(force: bool, state: Arc<AppState>) -> Self {
Self { force, state }
}
}
#[async_trait]
impl TaskPlanner<DocumentVersionTaskContext> for OcrPlanner {
async fn plan(
&self,
_ctx: &mut DocumentVersionTaskContext,
) -> TaskResult<Vec<BoxedTask<DocumentVersionTaskContext>>> {
Ok(vec![
Box::new(GenerateOcrTask::new(self.force, self.state.clone())),
Box::new(IndexDocumentTask::new()),
])
}
}
pub struct GenerateOcrTask {
force: bool,
state: Arc<AppState>,
+1 -80
View File
@@ -5,7 +5,6 @@ use chrono::Utc;
use diesel::{pg::upsert::excluded, prelude::*};
use image::{GenericImageView, ImageFormat, ImageReader};
use pdfium_render::prelude::*;
use serde::Deserialize;
use serde_json::{json, Map, Value};
use tokio::task;
use tracing::{info, warn};
@@ -14,25 +13,20 @@ use uuid::Uuid;
use crate::{
documents::asset::delete_asset,
error::AppResult,
jobs::JOB_GENERATE_THUMBNAILS,
models::{
Document, DocumentAsset, DocumentAssetObject, DocumentVersion, NewDocumentAsset,
NewDocumentAssetObject,
},
schema::{document_asset_objects, document_assets, document_versions},
state::AppState,
storage::TenantStorage,
utils::storage_paths::document_asset_object_key,
};
use super::{
analyze::determine_thumbnail_support,
job_execution_from_task_error,
taskflow::{
document::DocumentVersionTaskContext, BoxedTask, Task, TaskContext, TaskError,
TaskExecutor, TaskPlanner, TaskResult,
document::DocumentVersionTaskContext, Task, TaskContext, TaskError, TaskResult,
},
JobExecution, JobHandler,
};
pub const THUMBNAIL_WIDTH: u32 = 512;
@@ -42,79 +36,6 @@ const PREVIEW_HEIGHT: u32 = THUMBNAIL_HEIGHT * 4;
pub const THUMBNAIL_ASSET_TYPE: &str = "thumbnail";
pub const PREVIEW_ASSET_TYPE: &str = "preview";
#[derive(Debug, Deserialize, Clone)]
struct ThumbnailPayload {
document_id: Uuid,
document_version_id: Uuid,
#[serde(default)]
force: bool,
}
pub struct GenerateThumbnailsJob;
impl GenerateThumbnailsJob {
pub fn new() -> Self {
Self
}
}
#[async_trait]
impl JobHandler for GenerateThumbnailsJob {
fn job_type(&self) -> &'static str {
JOB_GENERATE_THUMBNAILS
}
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) => {
return JobExecution::Failed {
error: format!("invalid thumbnail payload: {err}"),
}
}
};
let mut context = DocumentVersionTaskContext::new(
job.id,
JOB_GENERATE_THUMBNAILS,
job.tenant_id,
payload.document_id,
payload.document_version_id,
payload.force,
state.config.worker_max_document_bytes,
state.clone(),
storage,
);
let planner = ThumbnailPlanner {
force: payload.force,
};
match TaskExecutor::run(&planner, &mut context).await {
Ok(()) => JobExecution::Success,
Err(err) => job_execution_from_task_error(err),
}
}
}
struct ThumbnailPlanner {
force: bool,
}
#[async_trait]
impl TaskPlanner<DocumentVersionTaskContext> for ThumbnailPlanner {
async fn plan(
&self,
_ctx: &mut DocumentVersionTaskContext,
) -> TaskResult<Vec<BoxedTask<DocumentVersionTaskContext>>> {
Ok(vec![Box::new(GenerateThumbnailsTask::new(self.force))])
}
}
pub struct GenerateThumbnailsTask {
force: bool,
}