use crate::{ error::Error, repo::{FullRepo, UploadId, UploadRepo}, store::Store, }; use actix_web::web::Bytes; use futures_util::{Stream, TryStreamExt}; use tokio_util::io::StreamReader; pub(crate) struct Backgrounded where R: FullRepo + 'static, S: Store, { repo: R, identifier: Option, upload_id: Option, } impl Backgrounded where R: FullRepo + 'static, S: Store, { pub(crate) fn disarm(mut self) { let _ = self.identifier.take(); let _ = self.upload_id.take(); } pub(crate) fn upload_id(&self) -> Option { self.upload_id } pub(crate) fn identifier(&self) -> Option<&S::Identifier> { self.identifier.as_ref() } pub(crate) async fn proxy

(repo: R, store: S, stream: P) -> Result where P: Stream>, { let mut this = Self { repo, identifier: None, upload_id: Some(UploadId::generate()), }; this.do_proxy(store, stream).await?; Ok(this) } async fn do_proxy

(&mut self, store: S, stream: P) -> Result<(), Error> where P: Stream>, { UploadRepo::create(&self.repo, self.upload_id.expect("Upload id exists")).await?; let stream = stream.map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e)); let mut reader = StreamReader::new(Box::pin(stream)); let identifier = store.save_async_read(&mut reader).await?; self.identifier = Some(identifier.clone()); Ok(()) } } impl Drop for Backgrounded where R: FullRepo + 'static, S: Store, { fn drop(&mut self) { if let Some(identifier) = self.identifier.take() { let repo = self.repo.clone(); actix_rt::spawn(async move { let _ = crate::queue::cleanup_identifier(&repo, identifier).await; }); } if let Some(upload_id) = self.upload_id { let repo = self.repo.clone(); actix_rt::spawn(async move { let _ = repo.claim(upload_id).await; }); } } }