use std::collections::BTreeMap;
use std::fs::create_dir;
use std::path::PathBuf;
use crate::db::{Hash, Table};
use crate::errors::AppError;
use crate::imdl::ImdlCommand;
use crate::options::CacheOptions;
use crate::queue::QueueItem;
use di::{inject, injectable, Ref};
use futures::stream::{iter, StreamExt};
use log::error;
#[injectable]
pub struct Queue {
table: Table<20, 1, QueueItem>,
}
#[allow(dead_code)]
impl Queue {
#[allow(dead_code)]
pub fn from_path(path: PathBuf) -> Self {
Self {
table: Table::new(path),
}
}
#[inject]
pub fn from_options(options: Ref<CacheOptions>) -> Self {
let path = options.cache.clone().expect("queue path should be set");
let path = path.join("queue");
if !path.exists() {
create_dir(&path)
.expect("should be able to create queue directory if it does not exist");
}
Self::from_path(path)
}
pub fn get(&self, hash: Hash<20>) -> Result<Option<QueueItem>, AppError> {
self.table.get(hash)
}
pub async fn get_unprocessed(
&mut self,
indexer: String,
transcode_enabled: bool,
upload_enabled: bool,
) -> Result<Vec<Hash<20>>, AppError> {
let items = self.table.get_all().await?;
let mut items: Vec<&QueueItem> = items
.values()
.filter(|item| {
item.indexer == indexer
&& exclude_verified_if_transcode_disabled(item, transcode_enabled)
&& exclude_transcoded_if_upload_disabled(item, upload_enabled)
&& exclude_verify_failures(item)
&& exclude_transcode_failures(item)
&& item.upload.is_none()
})
.collect();
items.sort_by_key(|x| &x.name);
let hashes = items.iter().map(|x| x.hash).collect();
Ok(hashes)
}
pub async fn get_all(&mut self) -> Result<BTreeMap<Hash<20>, QueueItem>, AppError> {
self.table.get_all().await
}
pub async fn set(&mut self, item: QueueItem) -> Result<(), AppError> {
self.table.set(item.hash, item).await
}
pub async fn set_many(
&self,
items: BTreeMap<Hash<20>, QueueItem>,
replace: bool,
) -> Result<usize, AppError> {
self.table.set_many(items, replace).await
}
pub async fn insert_new_torrent_files(
&mut self,
paths: Vec<PathBuf>,
) -> Result<usize, AppError> {
let stream = iter(paths.into_iter());
let items: BTreeMap<_, _> = stream
.filter_map(|path| async {
let torrent = match ImdlCommand::show(&path).await {
Ok(torrent) => Some(torrent),
Err(error) => {
error!("Failed to read torrent: {}\n{error}", path.display());
None
}
};
let item = QueueItem::from_torrent(path, torrent?);
Some((item.hash, item))
})
.collect()
.await;
self.table.set_many(items, false).await
}
}
fn exclude_verify_failures(item: &QueueItem) -> bool {
if let Some(verify) = &item.verify {
verify.verified
} else {
true
}
}
fn exclude_transcode_failures(item: &QueueItem) -> bool {
if let Some(transcode) = &item.transcode {
transcode.success
} else {
true
}
}
fn exclude_verified_if_transcode_disabled(item: &QueueItem, transcode_enabled: bool) -> bool {
transcode_enabled || item.verify.is_none()
}
fn exclude_transcoded_if_upload_disabled(x: &QueueItem, upload_enabled: bool) -> bool {
upload_enabled || x.transcode.is_none()
}