2
0
Fork 0
mirror of https://git.asonix.dog/asonix/pict-rs synced 2024-11-09 22:14:59 +00:00

Fix dangling unprocessed uploads

Adds error boundary around backgrounded ingest
This commit is contained in:
asonix 2023-07-17 22:30:10 -05:00
parent 42c801f0fe
commit 558605381d

View file

@ -1,5 +1,5 @@
use crate::{
error::Error,
error::{Error, UploadError},
formats::InputProcessableFormat,
ingest::Session,
queue::{Base64Bytes, LocalBoxFuture, Process},
@ -71,27 +71,37 @@ async fn process_ingest<R, S>(
unprocessed_identifier: Vec<u8>,
upload_id: UploadId,
declared_alias: Option<Alias>,
media: &crate::config::Media,
media: &'static crate::config::Media,
) -> Result<(), Error>
where
R: FullRepo + 'static,
S: Store,
S: Store + 'static,
{
let fut = async {
let unprocessed_identifier = S::Identifier::from_bytes(unprocessed_identifier)?;
let stream = store
.to_stream(&unprocessed_identifier, None, None)
let ident = unprocessed_identifier.clone();
let store2 = store.clone();
let repo = repo.clone();
let error_boundary = actix_rt::spawn(async move {
let stream = store2
.to_stream(&ident, None, None)
.await?
.map_err(Error::from);
let session = crate::ingest::ingest(repo, store, stream, declared_alias, media).await?;
let session =
crate::ingest::ingest(&repo, &store2, stream, declared_alias, media).await?;
let token = session.delete_token().await?;
Ok((session, token)) as Result<(Session<R, S>, DeleteToken), Error>
})
.await;
store.remove(&unprocessed_identifier).await?;
Ok((session, token)) as Result<(Session<R, S>, DeleteToken), Error>
error_boundary.map_err(|_| UploadError::Canceled)?
};
let result = match fut.await {