use std::collections::HashSet;
use std::collections::HashMap;
use csv_async::AsyncReaderBuilder;
use futures::future::ready;
use futures::stream::Stream;
use futures::stream::StreamExt;
use futures::stream::TryStreamExt;
use serde::Deserialize;
use crate::Error;
use crate::EPISODES_URL;
use crate::start_stream;
use crate::start_stream_lines;
#[derive(Debug, Deserialize)]
pub struct EpisodeLink {
#[serde(rename = "tconst", deserialize_with = "crate::util::parse_imdb_id")]
pub imdb_id: u64,
#[serde(rename = "parentTconst", deserialize_with = "crate::util::parse_imdb_id")]
pub series_imdb_id: u64,
#[serde(rename = "seasonNumber", deserialize_with = "crate::util::parse_janky_tsv_option")]
pub season: Option<u16>,
#[serde(rename = "episodeNumber", deserialize_with = "crate::util::parse_janky_tsv_option")]
pub episode: Option<u16>
}
pub async fn get_episodes_filtered_by_show(show_ids: &[u64]) -> Result<HashMap<u64, Vec<EpisodeLink>>, Error> {
let mut stream = stream_episodes_filtered(show_ids).await?;
let mut episodes_by_show = HashMap::with_capacity(show_ids.len());
while let Some(result) = stream.next().await {
let episode = result?;
if(episode.season.is_none() || episode.episode.is_none()) {
continue;
}
let episodes = episodes_by_show.entry(episode.series_imdb_id).or_insert_with(|| Vec::new());
episodes.push(episode);
}
Ok(episodes_by_show)
}
pub async fn get_episodes_filtered(show_ids: &[u64]) -> Result<HashMap<u64, EpisodeLink>, Error> {
let mut stream = stream_episodes_filtered(show_ids).await?;
let mut episodes = HashMap::with_capacity(show_ids.len() * 8); while let Some(result) = stream.next().await {
let episode = result?;
if(episode.season.is_none() || episode.episode.is_none()) {
continue;
}
episodes.insert(episode.imdb_id, episode);
}
Ok(episodes)
}
pub async fn stream_episodes_filtered(ids: &[u64]) -> Result<impl Stream<Item = Result<EpisodeLink, Error>> + '_, Error> {
let ids: HashSet<String> = ids.iter().map(|id| format!("tt{:07}", id)).collect();
let reader = start_stream_lines(EPISODES_URL).await?;
let stream = AsyncReaderBuilder::new()
.delimiter(b'\t')
.has_headers(false)
.create_deserializer(reader
.filter(move |result| match result {
Ok(line) => ready(ids.contains(line.split('\t').nth(1).unwrap())),
Err(_) => ready(true)
})
.map_ok(|l| l + "\n")
.into_async_read()
)
.into_deserialize::<EpisodeLink>();
Ok(stream.err_into())
}
pub async fn stream_episodes() -> Result<impl Stream<Item = Result<EpisodeLink, Error>>, Error> {
let reader = start_stream(EPISODES_URL).await?;
let stream = AsyncReaderBuilder::new()
.delimiter(b'\t')
.create_deserializer(reader)
.into_deserialize::<EpisodeLink>();
Ok(stream.err_into())
}