2
0
Fork 0
mirror of https://git.asonix.dog/asonix/pict-rs synced 2025-01-24 02:15:51 +00:00
pict-rs/src/queue/process.rs

202 lines
5.3 KiB
Rust
Raw Normal View History

use time::Instant;
use tracing::{Instrument, Span};
2023-09-06 01:45:07 +00:00
use crate::{
2023-07-22 16:15:30 +00:00
concurrent_processor::ProcessMap,
error::{Error, UploadError},
formats::InputProcessableFormat,
2023-09-05 02:51:27 +00:00
future::LocalBoxFuture,
ingest::Session,
2023-09-05 02:51:27 +00:00
queue::Process,
2024-02-04 00:19:14 +00:00
repo::{Alias, UploadId, UploadResult},
serde_str::Serde,
state::State,
store::Store,
};
use std::{path::PathBuf, sync::Arc};
pub(super) fn perform<'a, S>(
state: &'a State<S>,
2023-07-22 16:15:30 +00:00
process_map: &'a ProcessMap,
2023-09-03 17:47:06 +00:00
job: serde_json::Value,
) -> LocalBoxFuture<'a, Result<(), Error>>
where
S: Store + 'static,
{
Box::pin(async move {
2023-09-03 17:47:06 +00:00
match serde_json::from_value(job) {
Ok(job) => match job {
Process::Ingest {
identifier,
upload_id,
declared_alias,
} => {
process_ingest(
state,
Arc::from(identifier),
Serde::into_inner(upload_id),
declared_alias.map(Serde::into_inner),
)
.await?
}
Process::Generate {
target_format,
source,
process_path,
process_args,
} => {
generate(
state,
2023-07-22 16:15:30 +00:00
process_map,
target_format,
Serde::into_inner(source),
process_path,
process_args,
)
.await?
}
},
Err(e) => {
2023-01-29 17:57:59 +00:00
tracing::warn!("Invalid job: {}", format!("{e}"));
}
}
Ok(())
})
}
struct UploadGuard {
armed: bool,
start: Instant,
upload_id: UploadId,
}
impl UploadGuard {
fn guard(upload_id: UploadId) -> Self {
Self {
armed: true,
start: Instant::now(),
upload_id,
}
}
fn disarm(mut self) {
self.armed = false;
}
}
impl Drop for UploadGuard {
fn drop(&mut self) {
2024-02-04 21:45:47 +00:00
metrics::counter!(crate::init_metrics::BACKGROUND_UPLOAD_INGEST, "completed" => (!self.armed).to_string()).increment(1);
metrics::histogram!(crate::init_metrics::BACKGROUND_UPLOAD_INGEST_DURATION, "completed" => (!self.armed).to_string()).record(self.start.elapsed().as_seconds_f64());
if self.armed {
tracing::warn!(
"Upload future for {} dropped before completion! This can cause clients to wait forever",
self.upload_id,
);
}
}
}
#[tracing::instrument(skip(state))]
async fn process_ingest<S>(
state: &State<S>,
unprocessed_identifier: Arc<str>,
upload_id: UploadId,
declared_alias: Option<Alias>,
) -> Result<(), Error>
where
S: Store + 'static,
{
let guard = UploadGuard::guard(upload_id);
let fut = async {
2024-02-04 00:18:13 +00:00
let ident = unprocessed_identifier.clone();
let state2 = state.clone();
let current_span = Span::current();
let span = tracing::info_span!(parent: current_span, "error_boundary");
let error_boundary = crate::sync::abort_on_drop(crate::sync::spawn(
"ingest-media",
async move {
let stream =
crate::stream::from_err(state2.store.to_stream(&ident, None, None).await?);
let session = crate::ingest::ingest(&state2, stream, declared_alias).await?;
Ok(session) as Result<Session, Error>
}
.instrument(span),
))
.await;
state.store.remove(&unprocessed_identifier).await?;
2022-04-03 02:15:39 +00:00
error_boundary.map_err(|_| UploadError::Canceled)?
};
let result = match fut.await {
Ok(session) => {
let alias = session.alias().take().expect("Alias should exist").clone();
let token = session.disarm();
2023-07-27 03:53:41 +00:00
UploadResult::Success { alias, token }
}
Err(e) => {
2023-07-10 22:15:43 +00:00
tracing::warn!("Failed to ingest\n{}\n{}", format!("{e}"), format!("{e:?}"));
UploadResult::Failure {
2023-07-17 02:51:14 +00:00
message: e.root_cause().to_string(),
2023-09-02 01:50:10 +00:00
code: e.error_code().into_owned(),
}
}
};
state.repo.complete_upload(upload_id, result).await?;
guard.disarm();
Ok(())
}
#[tracing::instrument(skip(state, process_map, process_path, process_args))]
async fn generate<S: Store + 'static>(
state: &State<S>,
2023-07-22 16:15:30 +00:00
process_map: &ProcessMap,
target_format: InputProcessableFormat,
source: Alias,
process_path: PathBuf,
process_args: Vec<String>,
) -> Result<(), Error> {
let Some(hash) = state.repo.hash(&source).await? else {
// Nothing to do
return Ok(());
};
let path_string = process_path.to_string_lossy().to_string();
let identifier_opt = state
.repo
.variant_identifier(hash.clone(), path_string)
.await?;
if identifier_opt.is_some() {
return Ok(());
}
let original_details = crate::ensure_details(state, &source).await?;
crate::generate::generate(
state,
2023-07-22 16:15:30 +00:00
process_map,
target_format,
process_path,
process_args,
&original_details,
hash,
)
.await?;
Ok(())
}