Skip to main content

tmdb_rs/
stream.rs

1//! `stream` feature: drain a paginated endpoint as a stream of items
2
3use std::future::Future;
4use std::pin::Pin;
5use std::task::{Context, Poll};
6
7use futures_core::Stream;
8
9use crate::{Page, Result};
10
11/// a response the stream can walk page by page; implemented by [`Page`]
12pub trait Paginate {
13    type Item;
14    /// one item off the current page; pages are ~20 items, so a front
15    /// removal is fine
16    fn next_item(&mut self) -> Option<Self::Item>;
17    fn has_next_page(&self) -> bool;
18}
19
20impl<T> Paginate for Page<T> {
21    type Item = T;
22
23    fn next_item(&mut self) -> Option<T> {
24        if self.results.is_empty() {
25            None
26        } else {
27            Some(self.results.remove(0))
28        }
29    }
30
31    fn has_next_page(&self) -> bool {
32        self.has_next()
33    }
34}
35
36type Fetch<R> = Box<dyn FnMut(u32) -> Pin<Box<dyn Future<Output = Result<R>>>>>;
37
38/// the stream returned by `into_stream` on paginated request builders
39pub struct PageStream<R> {
40    fetch: Fetch<R>,
41    next_page: u32,
42    pending: Option<Pin<Box<dyn Future<Output = Result<R>>>>>,
43    page: Option<R>,
44    done: bool,
45}
46
47pub(crate) fn page_stream<R, F, Fut>(mut fetch: F) -> PageStream<R>
48where
49    F: FnMut(u32) -> Fut + 'static,
50    Fut: Future<Output = Result<R>> + 'static,
51{
52    PageStream {
53        fetch: Box::new(move |page| Box::pin(fetch(page))),
54        next_page: 1,
55        pending: None,
56        page: None,
57        done: false,
58    }
59}
60
61// everything pinned lives behind a Pin<Box>, so moving is safe
62impl<R> Unpin for PageStream<R> {}
63
64impl<R: Paginate> Stream for PageStream<R> {
65    type Item = Result<R::Item>;
66
67    fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
68        let this = self.get_mut();
69        loop {
70            if let Some(item) = this.page.as_mut().and_then(Paginate::next_item) {
71                return Poll::Ready(Some(Ok(item)));
72            }
73            if this.done {
74                return Poll::Ready(None);
75            }
76            if this.pending.is_none() {
77                this.pending = Some((this.fetch)(this.next_page));
78            }
79            // just set above, or on a previous poll
80            let pending = this.pending.as_mut().expect("pending future");
81            let page = match pending.as_mut().poll(cx) {
82                Poll::Pending => return Poll::Pending,
83                Poll::Ready(Err(error)) => {
84                    this.done = true;
85                    return Poll::Ready(Some(Err(error)));
86                }
87                Poll::Ready(Ok(page)) => page,
88            };
89            this.pending = None;
90            this.done = !page.has_next_page();
91            this.next_page += 1;
92            this.page = Some(page);
93        }
94    }
95}