backend dry
This commit is contained in:
@@ -18,7 +18,13 @@ use crate::{
|
||||
storage::TenantStorage,
|
||||
};
|
||||
|
||||
use super::{ocr::OCR_TEXT_ASSET_TYPE, JobExecution, JobHandler};
|
||||
use super::{
|
||||
fetch_version_object,
|
||||
handle_fetch_error,
|
||||
ocr::OCR_TEXT_ASSET_TYPE,
|
||||
JobExecution,
|
||||
JobHandler,
|
||||
};
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct IndexPayload {
|
||||
@@ -114,21 +120,23 @@ impl JobHandler for IndexDocumentTextJob {
|
||||
}
|
||||
|
||||
let s3_key = context.text_s3_key.unwrap();
|
||||
let text = match storage.get_object(&s3_key).await {
|
||||
Ok(bytes) => match String::from_utf8(bytes) {
|
||||
Ok(text) => text,
|
||||
Err(err) => {
|
||||
warn!(job_id = %job.id, error = %err, "ocr text not valid UTF-8");
|
||||
return JobExecution::Failed {
|
||||
error: "ocr text not valid UTF-8".into(),
|
||||
};
|
||||
}
|
||||
},
|
||||
let bytes = match fetch_version_object(
|
||||
&context.version,
|
||||
&storage,
|
||||
&s3_key,
|
||||
state.config.worker_max_document_bytes,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(bytes) => bytes,
|
||||
Err(err) => return handle_fetch_error(job, err, "failed to download ocr text"),
|
||||
};
|
||||
let text = match String::from_utf8(bytes) {
|
||||
Ok(text) => text,
|
||||
Err(err) => {
|
||||
warn!(job_id = %job.id, error = %err, "failed to download ocr text");
|
||||
return JobExecution::Retry {
|
||||
delay: Duration::from_secs(30),
|
||||
error: err.to_string(),
|
||||
warn!(job_id = %job.id, error = %err, "ocr text not valid UTF-8");
|
||||
return JobExecution::Failed {
|
||||
error: "ocr text not valid UTF-8".into(),
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
use std::{collections::HashMap, sync::Arc, time::Duration};
|
||||
|
||||
use anyhow::Error as AnyhowError;
|
||||
use async_trait::async_trait;
|
||||
use tokio::time::sleep;
|
||||
use tracing::{error, info, warn};
|
||||
@@ -147,3 +148,68 @@ pub fn default_handlers() -> Vec<Arc<dyn JobHandler>> {
|
||||
Arc::new(index::IndexDocumentTextJob::new()),
|
||||
]
|
||||
}
|
||||
|
||||
pub(crate) fn check_worker_document_limit(
|
||||
size_bytes: i64,
|
||||
limit_bytes: u64,
|
||||
) -> Result<(), (u64, u64)> {
|
||||
let size = size_bytes.max(0) as u64;
|
||||
if size > limit_bytes {
|
||||
Err((size, limit_bytes))
|
||||
} else {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) enum FetchVersionError {
|
||||
TooLarge { size: u64, limit: u64 },
|
||||
Storage(AnyhowError),
|
||||
}
|
||||
|
||||
pub(crate) async fn fetch_version_object(
|
||||
version: &crate::models::DocumentVersion,
|
||||
storage: &TenantStorage,
|
||||
s3_key: &str,
|
||||
limit_bytes: u64,
|
||||
) -> Result<Vec<u8>, FetchVersionError> {
|
||||
check_worker_document_limit(version.size_bytes, limit_bytes).map_err(|(size, limit)| {
|
||||
FetchVersionError::TooLarge {
|
||||
size,
|
||||
limit,
|
||||
}
|
||||
})?;
|
||||
|
||||
storage
|
||||
.get_object(s3_key)
|
||||
.await
|
||||
.map_err(FetchVersionError::Storage)
|
||||
}
|
||||
|
||||
pub(crate) fn handle_fetch_error(
|
||||
job_id: crate::models::Job,
|
||||
err: FetchVersionError,
|
||||
message: &str,
|
||||
) -> JobExecution {
|
||||
match err {
|
||||
FetchVersionError::TooLarge { size, limit } => {
|
||||
warn!(
|
||||
job_id = %job_id.id,
|
||||
size_bytes = size,
|
||||
limit_bytes = limit,
|
||||
"document exceeds worker size limit"
|
||||
);
|
||||
JobExecution::Failed {
|
||||
error: format!(
|
||||
"document size {size} bytes exceeds worker limit of {limit} bytes"
|
||||
),
|
||||
}
|
||||
}
|
||||
FetchVersionError::Storage(err) => {
|
||||
warn!(job_id = %job_id.id, error = %err, "{message}");
|
||||
JobExecution::Retry {
|
||||
delay: Duration::from_secs(30),
|
||||
error: err.to_string(),
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -30,7 +30,7 @@ use crate::{
|
||||
utils::storage_paths::document_asset_object_prefix,
|
||||
};
|
||||
|
||||
use super::{JobExecution, JobHandler};
|
||||
use super::{fetch_version_object, handle_fetch_error, JobExecution, JobHandler};
|
||||
|
||||
pub const OCR_TEXT_ASSET_TYPE: &str = "ocr-text";
|
||||
const MIN_TEXT_LENGTH: usize = 50;
|
||||
@@ -99,15 +99,16 @@ impl JobHandler for GenerateOcrTextJob {
|
||||
return JobExecution::Success;
|
||||
}
|
||||
|
||||
let bytes = match storage.get_object(&context.version.s3_key).await {
|
||||
let bytes = match fetch_version_object(
|
||||
&context.version,
|
||||
&storage,
|
||||
&context.version.s3_key,
|
||||
state.config.worker_max_document_bytes,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(bytes) => bytes,
|
||||
Err(err) => {
|
||||
warn!(job_id = %job.id, error = %err, "failed to fetch document for ocr");
|
||||
return JobExecution::Retry {
|
||||
delay: Duration::from_secs(30),
|
||||
error: err.to_string(),
|
||||
};
|
||||
}
|
||||
Err(err) => return handle_fetch_error(job, err, "failed to fetch document for ocr"),
|
||||
};
|
||||
|
||||
let doc_meta = PdfDocumentMeta {
|
||||
|
||||
@@ -24,7 +24,13 @@ use crate::{
|
||||
utils::storage_paths::document_asset_object_key,
|
||||
};
|
||||
|
||||
use super::{analyze::determine_thumbnail_support, JobExecution, JobHandler};
|
||||
use super::{
|
||||
analyze::determine_thumbnail_support,
|
||||
fetch_version_object,
|
||||
handle_fetch_error,
|
||||
JobExecution,
|
||||
JobHandler,
|
||||
};
|
||||
|
||||
const THUMBNAIL_WIDTH: u32 = 512;
|
||||
const THUMBNAIL_HEIGHT: u32 = 512;
|
||||
@@ -96,15 +102,16 @@ impl JobHandler for GenerateThumbnailsJob {
|
||||
return JobExecution::Success;
|
||||
}
|
||||
|
||||
let bytes = match storage.get_object(&initial.version.s3_key).await {
|
||||
let bytes = match fetch_version_object(
|
||||
&initial.version,
|
||||
&storage,
|
||||
&initial.version.s3_key,
|
||||
state.config.worker_max_document_bytes,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(bytes) => bytes,
|
||||
Err(err) => {
|
||||
warn!(job_id = %job.id, error = %err, "thumbnail fetch failed; will retry");
|
||||
return JobExecution::Retry {
|
||||
delay: Duration::from_secs(30),
|
||||
error: err.to_string(),
|
||||
};
|
||||
}
|
||||
Err(err) => return handle_fetch_error(job, err, "thumbnail fetch failed; will retry"),
|
||||
};
|
||||
|
||||
let generation = match generate_preview_and_thumbnail(&initial.document, &bytes) {
|
||||
|
||||
Reference in New Issue
Block a user