imdb-async 0.1.0

Opinionated and unopinionated async wrappers to efficiently retrieve and parse IMDB's dataset
Documentation
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?;
		// TODO:  Does it make sense to allow episodes with no season/ep number for our purposes?
		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); // this is just a capacity hint to hopefully minimize reallocations
	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())
}