use crate::client::{Client, ClientBuilder};
use crate::error::Result;
pub struct PageStream<T: Client> {
builder: ClientBuilder<T>,
current_page: u32,
exhausted: bool,
max_pages: Option<u32>,
}
impl<T: Client> PageStream<T> {
pub fn new(builder: ClientBuilder<T>) -> Self {
let current_page = builder.page;
Self {
builder,
current_page,
exhausted: false,
max_pages: None,
}
}
#[must_use]
pub fn max_pages(mut self, max: u32) -> Self {
self.max_pages = Some(max);
self
}
pub fn current_page(&self) -> u32 {
self.current_page
}
pub async fn next(&mut self) -> Option<Result<Vec<T::Post>>> {
if self.exhausted {
return None;
}
if let Some(max) = self.max_pages {
let pages_fetched = self.current_page.saturating_sub(self.builder.page);
if pages_fetched >= max {
self.exhausted = true;
return None;
}
}
let mut page_builder = self.builder.clone();
page_builder.page = self.current_page;
let client = page_builder.build();
match client.get().await {
Ok(posts) => {
if posts.is_empty() {
self.exhausted = true;
return Some(Ok(posts));
}
self.current_page += 1;
Some(Ok(posts))
}
Err(e) => {
self.exhausted = true;
Some(Err(e))
}
}
}
}
pub struct PostStream<T: Client> {
page_stream: PageStream<T>,
buffer: Vec<T::Post>,
buffer_index: usize,
posts_yielded: u32,
max_posts: Option<u32>,
}
impl<T: Client> PostStream<T> {
pub fn new(builder: ClientBuilder<T>) -> Self {
Self {
page_stream: PageStream::new(builder),
buffer: Vec::new(),
buffer_index: 0,
posts_yielded: 0,
max_posts: None,
}
}
#[must_use]
pub fn max_posts(mut self, max: u32) -> Self {
self.max_posts = Some(max);
self
}
#[must_use]
pub fn max_pages(mut self, max: u32) -> Self {
self.page_stream = self.page_stream.max_pages(max);
self
}
pub fn posts_yielded(&self) -> u32 {
self.posts_yielded
}
pub fn current_page(&self) -> u32 {
self.page_stream.current_page()
}
pub async fn next(&mut self) -> Option<Result<T::Post>> {
if let Some(max) = self.max_posts
&& self.posts_yielded >= max
{
return None;
}
if self.buffer_index < self.buffer.len() {
let post = self.buffer.swap_remove(self.buffer_index);
self.buffer_index = 0; self.posts_yielded += 1;
return Some(Ok(post));
}
match self.page_stream.next().await? {
Ok(posts) => {
if posts.is_empty() {
return None;
}
self.buffer = posts;
self.buffer_index = 1; self.posts_yielded += 1;
if self.buffer.is_empty() {
None
} else {
Some(Ok(self.buffer.swap_remove(0)))
}
}
Err(e) => Some(Err(e)),
}
}
pub async fn collect(mut self) -> Result<Vec<T::Post>> {
let mut all_posts = Vec::new();
while let Some(result) = self.next().await {
all_posts.push(result?);
}
Ok(all_posts)
}
}
impl<T: Client> ClientBuilder<T> {
#[must_use]
pub fn into_page_stream(self) -> PageStream<T> {
PageStream::new(self)
}
#[must_use]
pub fn into_post_stream(self) -> PostStream<T> {
PostStream::new(self)
}
}