summaryrefslogtreecommitdiffstats
path: root/src/entities/itemsiter.rs
blob: 6f293912d8cfe4f917456fc1f114fef08e9d53f2 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
use futures::{stream::unfold, Stream};
use log::{as_debug, as_serde, debug, info, warn};

use crate::page::Page;
use serde::{Deserialize, Serialize};

/// Abstracts away the `next_page` logic into a single stream of items
///
/// ```no_run,async
/// use mastodon_async::prelude::*;
/// use futures::stream::StreamExt;
/// use futures_util::pin_mut;
///
/// tokio_test::block_on(async {
///     let data = Data::default();
///     let client = Mastodon::from(data);
///     let statuses = client.statuses(&AccountId::new("user-id"), Default::default()).await.unwrap().items_iter();
///     statuses.for_each(|status| async move {
///         // Do something with the status
///     }).await;
/// })
/// ```
///
/// See documentation for `futures::Stream::StreamExt` for available methods.
#[derive(Debug, Clone)]
pub(crate) struct ItemsIter<T: Clone + for<'de> Deserialize<'de> + Serialize> {
    page: Page<T>,
    buffer: Vec<T>,
    cur_idx: usize,
    use_initial: bool,
}

impl<'a, T: Clone + for<'de> Deserialize<'de> + Serialize> ItemsIter<T> {
    pub(crate) fn new(page: Page<T>) -> ItemsIter<T> {
        ItemsIter {
            page,
            buffer: vec![],
            cur_idx: 0,
            use_initial: true,
        }
    }

    fn need_next_page(&self) -> bool {
        if self.buffer.is_empty() || self.cur_idx == self.buffer.len() {
            debug!(idx = self.cur_idx, buffer_len = self.buffer.len(); "next page needed");
            true
        } else {
            false
        }
    }

    async fn fill_next_page(&mut self) -> Option<()> {
        match self.page.next_page().await {
            Ok(Some(items)) => {
                info!(item_count = items.len(); "next page received");
                if items.is_empty() {
                    return None;
                }
                self.buffer = items;
                self.cur_idx = 0;
                Some(())
            }
            Err(err) => {
                warn!(err = as_debug!(err); "error encountered filling next page");
                None
            }
            _ => None,
        }
    }

    pub(crate) fn stream(self) -> impl Stream<Item = T> {
        unfold(self, |mut this| async move {
            if this.use_initial {
                let idx = this.cur_idx;
                if this.page.initial_items.is_empty() || idx == this.page.initial_items.len() {
                    debug!(index = idx, n_initial_items = this.page.initial_items.len(); "exhausted initial items and no more pages are present");
                    return None;
                }
                if idx == this.page.initial_items.len() - 1 {
                    this.cur_idx = 0;
                    this.use_initial = false;
                    debug!(index = idx, n_initial_items = this.page.initial_items.len(); "exhausted initial items");
                } else {
                    this.cur_idx += 1;
                }
                let item = this.page.initial_items[idx].clone();
                debug!(item = as_serde!(item), index = idx; "yielding item from initial items");
                // let item = Box::pin(item);
                // pin_mut!(item);
                Some((item, this))
            } else {
                if this.need_next_page() {
                    this.fill_next_page().await?;
                }
                let idx = this.cur_idx;
                this.cur_idx += 1;
                let item = this.buffer[idx].clone();
                debug!(item = as_serde!(item), index = idx; "yielding item from initial stream");
                Some((item, this))
            }
        })
    }
}