Skip to main content

azure_core/http/
pager.rs

1// Copyright (c) Microsoft Corporation. All rights reserved.
2// Licensed under the MIT License.
3
4//! Types and methods for pageable responses.
5
6use crate::{
7    error::ErrorKind,
8    http::{
9        headers::HeaderName, policies::create_public_api_span, response::Response, Context,
10        DeserializeWith, Format, JsonFormat, Url,
11    },
12    tracing::{Span, SpanStatus},
13};
14use async_trait::async_trait;
15use futures::{stream::FusedStream, FutureExt, Stream};
16use pin_project::pin_project;
17use std::{fmt, future::Future, pin::Pin, sync::Arc, task};
18use typespec::error::ResultExt;
19
20/// Represents the state of a [`Pager`] or [`PageIterator`].
21#[derive(Clone, Debug, Default, PartialEq, Eq)]
22pub enum PagerState {
23    /// The pager should fetch the initial page.
24    #[default]
25    Initial,
26    /// The pager should fetch a subsequent page using the next link/continuation token `C`.
27    More(PagerContinuation),
28}
29
30/// The result of fetching a single page from a [`Pager`], whether there are more pages or paging is done.
31pub enum PagerResult<P> {
32    /// There are more pages the [`Pager`] may fetch using the `continuation` token.
33    More {
34        /// The response for the current page.
35        response: P,
36        /// Continuation state for the next page.
37        continuation: PagerContinuation,
38    },
39    /// The [`Pager`] is done and there are no additional pages to fetch.
40    Done {
41        /// The response for the current page.
42        response: P,
43    },
44}
45
46impl<P, F> PagerResult<Response<P, F>> {
47    /// Creates a [`PagerResult`] from the provided response, extracting the continuation value from the provided header.
48    ///
49    /// If the provided response has a header with the matching name, this returns [`PagerResult::More`], using the value from the header as the continuation.
50    /// If the provided response does not have a header with the matching name, this returns [`PagerResult::Done`].
51    pub fn from_response_header(response: Response<P, F>, header_name: &HeaderName) -> Self {
52        match response.headers().get_optional_string(header_name) {
53            Some(continuation) => PagerResult::More {
54                response,
55                continuation: PagerContinuation::Token(continuation),
56            },
57            None => PagerResult::Done { response },
58        }
59    }
60}
61
62impl<P> fmt::Debug for PagerResult<P> {
63    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
64        match self {
65            Self::More { continuation, .. } => f
66                .debug_struct("More")
67                .field("continuation", continuation)
68                .finish_non_exhaustive(),
69            Self::Done { .. } => f.debug_struct("Done").finish_non_exhaustive(),
70        }
71    }
72}
73
74/// Information returned by the server to continue to the next page.
75#[derive(Clone, Debug, PartialEq, Eq)]
76#[non_exhaustive]
77pub enum PagerContinuation {
78    /// The continuation is a next link for client paging.
79    Link(Url),
80
81    /// The continuation is a token for server paging.
82    Token(String),
83}
84
85impl AsRef<str> for PagerContinuation {
86    fn as_ref(&self) -> &str {
87        match self {
88            Self::Link(next_link) => next_link.as_str(),
89            Self::Token(continuation_token) => continuation_token.as_str(),
90        }
91    }
92}
93
94impl fmt::Display for PagerContinuation {
95    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
96        f.write_str(self.as_ref())
97    }
98}
99
100impl From<PagerContinuation> for String {
101    fn from(value: PagerContinuation) -> Self {
102        match value {
103            PagerContinuation::Link(next_link) => String::from(next_link.as_str()),
104            PagerContinuation::Token(continuation_token) => continuation_token,
105        }
106    }
107}
108
109impl TryFrom<PagerContinuation> for Url {
110    type Error = crate::Error;
111    fn try_from(value: PagerContinuation) -> Result<Self, Self::Error> {
112        match value {
113            PagerContinuation::Link(next_link) => Ok(next_link),
114            PagerContinuation::Token(continuation_token) => continuation_token
115                .parse()
116                .with_kind(ErrorKind::DataConversion),
117        }
118    }
119}
120
121/// Represents a single page of items returned by a collection request to a service.
122#[async_trait]
123pub trait Page {
124    /// The type of items in the collection.
125    type Item;
126    /// The type containing items in the collection e.g., [`Vec<Self::Item>`](Vec).
127    type IntoIter: Iterator<Item = Self::Item>;
128
129    /// Gets a single page of items returned by a collection request to a service.
130    async fn into_items(self) -> crate::Result<Self::IntoIter>;
131}
132
133#[async_trait]
134impl<P, F> Page for Response<P, F>
135where
136    P: DeserializeWith<F> + Page + Send,
137    F: Format + Send,
138{
139    type Item = P::Item;
140    type IntoIter = P::IntoIter;
141    async fn into_items(self) -> crate::Result<Self::IntoIter> {
142        let page: P = self.into_model()?;
143        page.into_items().await
144    }
145}
146
147/// Represents a paginated stream of items returned by a collection request to a service.
148///
149/// Specifically, this is a [`ItemIterator`] that yields [`Response`] items.
150///
151/// # Examples
152///
153/// For clients that return a `Pager`, you can iterate over items across one or more pages:
154///
155/// ```no_run
156/// # use azure_core::{credentials::TokenCredential, http::Transport};
157/// # use azure_core_examples::secrets::{ResourceExt, SecretClient, SecretClientOptions};
158/// # use futures::TryStreamExt;
159/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
160/// # let credential: std::sync::Arc<dyn TokenCredential> = unimplemented!();
161/// let client = SecretClient::new(
162///     "https://my-vault.vault.azure.net",
163///     credential.clone(),
164///     None,
165/// )?;
166///
167/// // List secret properties using a Pager.
168/// let mut pager = client.list_secret_properties(None)?;
169/// while let Some(secret) = pager.try_next().await? {
170///     println!("{}", secret.resource_id()?.name);
171/// }
172/// # Ok(()) }
173/// ```
174///
175/// If you want to iterate each page of items, you can call [`Pager::into_pages`] to get a [`PageIterator`]:
176///
177/// ```no_run
178/// # use azure_core::{credentials::TokenCredential, http::Transport};
179/// # use azure_core_examples::secrets::{ResourceExt, SecretClient, SecretClientOptions};
180/// # use futures::TryStreamExt;
181/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
182/// # let credential: std::sync::Arc<dyn TokenCredential> = unimplemented!();
183/// let client = SecretClient::new(
184///     "https://my-vault.vault.azure.net",
185///     credential.clone(),
186///     None,
187/// )?;
188///
189/// // Iterate each page of secrets using a PageIterator.
190/// let mut pager = client.list_secret_properties(None)?.into_pages();
191/// while let Some(page) = pager.try_next().await? {
192///     let page = page.into_model()?;
193///     for secret in page.value {
194///         println!("{}", secret.resource_id()?.name);
195///     }
196/// }
197/// # Ok(()) }
198/// ```
199pub type Pager<P, F = JsonFormat> = ItemIterator<Response<P, F>>;
200
201/// A pinned boxed [`Future`] that can be stored and called dynamically.
202///
203/// Intended only for [`ItemIterator`] and [`PageIterator`].
204pub type PagerResultFuture<P> =
205    Pin<Box<dyn Future<Output = crate::Result<PagerResult<P>>> + Send + 'static>>;
206
207type PagerFn<P> = Box<dyn Fn(PagerState, PagerOptions<'static>) -> PagerResultFuture<P> + Send>;
208
209/// Options for configuring the behavior of a [`Pager`].
210#[derive(Clone)]
211pub struct PagerOptions<'a> {
212    /// Context for HTTP requests made by the [`Pager`].
213    pub context: Context<'a>,
214
215    /// Optional continuation token or next link to resume paging.
216    ///
217    /// # Examples
218    ///
219    /// ``` no_run
220    /// # use azure_core_examples::identity as azure_identity;
221    /// # use azure_core_examples::secrets as azure_security_keyvault_secrets;
222    /// use azure_core::http::pager::PagerOptions;
223    /// use azure_identity::DeveloperToolsCredential;
224    /// use azure_security_keyvault_secrets::{
225    ///     models::SecretClientListSecretPropertiesOptions,
226    ///     SecretClient,
227    /// };
228    /// use futures::stream::TryStreamExt as _;
229    ///
230    /// # #[tokio::main]
231    /// # async fn main() -> azure_core::Result<()> {
232    /// let client = SecretClient::new("https://my-vault.vault.azure.net", DeveloperToolsCredential::new(None)?, None)?;
233    ///
234    /// // Start the first pager at the first page.
235    /// let mut pager = client.list_secret_properties(None)?;
236    ///
237    /// // Continue the second pager from where the first pager left off,
238    /// // which is the first page in this example.
239    /// let options = SecretClientListSecretPropertiesOptions {
240    ///     method_options: PagerOptions {
241    ///         continuation: pager.into_continuation(),
242    ///         ..Default::default()
243    ///     },
244    ///     ..Default::default()
245    /// };
246    /// let mut pager = client.list_secret_properties(Some(options))?;
247    /// while let Some(secret) = pager.try_next().await? {
248    ///     println!("{:?}", secret.id);
249    /// }
250    /// # Ok(())
251    /// # }
252    /// ```
253    pub continuation: Option<PagerContinuation>,
254}
255
256impl<'a> fmt::Debug for PagerOptions<'a> {
257    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
258        f.debug_struct("PagerOptions")
259            .field("context", &self.context)
260            .field("continuation", &self.continuation)
261            .finish()
262    }
263}
264
265impl<'a> Default for PagerOptions<'a> {
266    fn default() -> Self {
267        PagerOptions {
268            context: Context::new(),
269            continuation: None,
270        }
271    }
272}
273
274/// Iterates over a collection of items or individual pages of items from a service.
275///
276/// You can asynchronously iterate over items returned by a collection request to a service,
277/// or asynchronously fetch pages of items if preferred.
278///
279/// # Examples
280///
281/// For clients that return a `Pager`, you can iterate over items across one or more pages:
282///
283/// ```no_run
284/// # use azure_core::{credentials::TokenCredential, http::Transport};
285/// # use azure_core_examples::secrets::{ResourceExt, SecretClient, SecretClientOptions};
286/// # use futures::TryStreamExt;
287/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
288/// # let credential: std::sync::Arc<dyn TokenCredential> = unimplemented!();
289/// let client = SecretClient::new(
290///     "https://my-vault.vault.azure.net",
291///     credential.clone(),
292///     None,
293/// )?;
294///
295/// // List secret properties using a Pager.
296/// let mut pager = client.list_secret_properties(None)?;
297/// while let Some(secret) = pager.try_next().await? {
298///     println!("{}", secret.resource_id()?.name);
299/// }
300/// # Ok(()) }
301/// ```
302///
303/// If you want to iterate each page of items, you can call [`Pager::into_pages`] to get a [`PageIterator`]:
304///
305/// ```no_run
306/// # use azure_core::{credentials::TokenCredential, http::Transport};
307/// # use azure_core_examples::secrets::{ResourceExt, SecretClient, SecretClientOptions};
308/// # use futures::TryStreamExt;
309/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
310/// # let credential: std::sync::Arc<dyn TokenCredential> = unimplemented!();
311/// let client = SecretClient::new(
312///     "https://my-vault.vault.azure.net",
313///     credential.clone(),
314///     None,
315/// )?;
316///
317/// // Iterate each page of secrets using a PageIterator.
318/// let mut pager = client.list_secret_properties(None)?.into_pages();
319/// while let Some(page) = pager.try_next().await? {
320///     let page = page.into_model()?;
321///     for secret in page.value {
322///         println!("{}", secret.resource_id()?.name);
323///     }
324/// }
325/// # Ok(()) }
326/// ```
327#[must_use = "streams do nothing unless you poll them"]
328#[pin_project(project = ItemIteratorProjection, project_replace = ItemIteratorProjectionOwned)]
329pub struct ItemIterator<P>
330where
331    P: Page + Send,
332{
333    #[pin]
334    iter: PageIterator<P>,
335    /// The continuation token or next link for the current page.
336    ///
337    /// Unlike the inner [`PageIterator::continuation`], this continuation token might be a page behind.
338    /// To help avoid skipping items (barring underlying changes in the remote collection), we don't continue with the next page
339    /// until after we've iterated all items on the current page.
340    continuation: Option<PagerContinuation>,
341    current: Option<P::IntoIter>,
342}
343
344impl<P> ItemIterator<P>
345where
346    P: Page + Send,
347{
348    /// Creates a [`ItemIterator`] from a callback that will be called repeatedly to request each page.
349    ///
350    /// This method expect a callback that accepts a single [`PagerState`] parameter, and returns a [`PagerResult`] value asynchronously.
351    /// The result will be an asynchronous stream of [`Result`](crate::Result) values.
352    ///
353    /// The first time your callback is called, it will be called with [`Option::None`], indicating no next link/continuation token is present.
354    ///
355    /// Your callback must return one of:
356    /// * `Ok(result)` - The request succeeded, and the provided [`PagerResult`] indicates the value to return and if there are more pages.
357    /// * `Err(..)` - The request failed. The error will be yielded to the stream, the stream will end, and the callback will not be called again.
358    ///
359    /// ## Examples
360    ///
361    /// To page results using a next link:
362    ///
363    /// ```rust,no_run
364    /// # use azure_core::{Result, http::{RawResponse, ItemIterator, pager::{Page, PagerContinuation, PagerOptions, PagerResult, PagerState}, Pipeline, Request, Response, Method, Url}, json};
365    /// # let api_version = "2025-06-04".to_string();
366    /// # let pipeline: Pipeline = panic!("Not a runnable example");
367    /// #[derive(serde::Deserialize)]
368    /// struct ListItemsResult {
369    ///     items: Vec<String>,
370    ///     next_link: Option<String>,
371    /// }
372    /// #[async_trait::async_trait]
373    /// impl Page for ListItemsResult {
374    ///     type Item = String;
375    ///     type IntoIter = <Vec<String> as IntoIterator>::IntoIter;
376    ///     async fn into_items(self) -> Result<Self::IntoIter> {
377    ///         Ok(self.items.into_iter())
378    ///     }
379    /// }
380    /// let url = "https://example.com/my_paginated_api".parse().unwrap();
381    /// let mut base_req = Request::new(url, Method::Get);
382    /// let pager = ItemIterator::new(move |next_link: PagerState, options: PagerOptions<'static>| {
383    ///     // The callback must be 'static, so you have to clone and move any values you want to use.
384    ///     let pipeline = pipeline.clone();
385    ///     let api_version = api_version.clone();
386    ///     let mut req = base_req.clone();
387    ///     Box::pin(async move {
388    ///         if let PagerState::More(next_link) = next_link {
389    ///             let next_link: Url = next_link.try_into().expect("expected Url");
390    ///             // Ensure the api-version from the client is appended.
391    ///             let qp = next_link
392    ///                 .query_pairs()
393    ///                 .filter(|(name, _)| name.ne("api-version"));
394    ///             req
395    ///                 .url_mut()
396    ///                 .query_pairs_mut()
397    ///                 .clear()
398    ///                 .extend_pairs(qp)
399    ///                 .append_pair("api-version", &api_version);
400    ///         }
401    ///         let resp = pipeline
402    ///           .send(&options.context, &mut req, None)
403    ///           .await?;
404    ///         let (status, headers, body) = resp.deconstruct();
405    ///         let result: ListItemsResult = json::from_json(&body)?;
406    ///         let resp: Response<ListItemsResult> = RawResponse::from_bytes(status, headers, body).into();
407    ///         Ok(match result.next_link {
408    ///             Some(next_link) => PagerResult::More {
409    ///                 response: resp,
410    ///                 continuation: PagerContinuation::Link(next_link.parse()?),
411    ///             },
412    ///             None => PagerResult::Done { response: resp }
413    ///         })
414    ///     })
415    /// }, None);
416    /// ```
417    ///
418    /// To page results using headers:
419    ///
420    /// ```rust,no_run
421    /// # use azure_core::{Result, http::{Context, ItemIterator, pager::{Page, PagerResult, PagerState}, Pipeline, Request, Response, Method, headers::HeaderName}};
422    /// # let pipeline: Pipeline = panic!("Not a runnable example");
423    /// #[derive(serde::Deserialize)]
424    /// struct ListItemsResult {
425    ///     items: Vec<String>,
426    /// }
427    /// #[async_trait::async_trait]
428    /// impl Page for ListItemsResult {
429    ///     type Item = String;
430    ///     type IntoIter = <Vec<String> as IntoIterator>::IntoIter;
431    ///     async fn into_items(self) -> Result<Self::IntoIter> {
432    ///         Ok(self.items.into_iter())
433    ///     }
434    /// }
435    /// let url = "https://example.com/my_paginated_api".parse().unwrap();
436    /// let mut base_req = Request::new(url, Method::Get);
437    /// let pager = ItemIterator::new(move |continuation, options| {
438    ///     // The callback must be 'static, so you have to clone and move any values you want to use.
439    ///     let pipeline = pipeline.clone();
440    ///     let mut req = base_req.clone();
441    ///     Box::pin(async move {
442    ///         if let PagerState::More(continuation) = continuation {
443    ///             req.insert_header("x-ms-continuation", continuation.as_ref().to_string());
444    ///         }
445    ///         let resp: Response<ListItemsResult> = pipeline
446    ///           .send(&options.context, &mut req, None)
447    ///           .await?
448    ///           .into();
449    ///         Ok(PagerResult::from_response_header(resp, &HeaderName::from_static("x-next-continuation")))
450    ///     })
451    /// }, None);
452    /// ```
453    pub fn new<
454        F: Fn(PagerState, PagerOptions<'static>) -> PagerResultFuture<P> + Send + 'static,
455    >(
456        make_request: F,
457        options: Option<PagerOptions<'static>>,
458    ) -> Self {
459        let options = options.unwrap_or_default();
460
461        // Start from the optional `PagerOptions::continuation_token`.
462        let continuation_token = options.continuation.clone();
463
464        Self {
465            iter: PageIterator::new(make_request, Some(options)),
466            continuation: continuation_token,
467            current: None,
468        }
469    }
470
471    /// Gets the continuation token to pass to [`PagerOptions`] to resume paging in another iterator.
472    pub fn continuation(&self) -> Option<&PagerContinuation> {
473        self.continuation.as_ref()
474    }
475
476    /// Gets the continuation token to pass to [`PagerOptions`] to resume paging in another iterator.
477    ///
478    /// # Examples
479    ///
480    /// This takes ownership of the iterator and can be useful when constructing a new iterator.
481    ///
482    /// ```no_run
483    /// # use azure_core::http::pager::PagerOptions;
484    /// # use azure_core_examples::secrets::{models::SecretClientListSecretPropertiesOptions, SecretClient};
485    /// # use futures::stream::TryStreamExt;
486    /// # #[tokio::main] async fn main() -> azure_core::Result<()> {
487    /// # let client: SecretClient = unimplemented!();
488    /// let pager1 = client.list_secret_properties(None)?;
489    /// assert!(pager1.try_next().await?.is_some());
490    /// let options = SecretClientListSecretPropertiesOptions {
491    ///     method_options: PagerOptions {
492    ///         continuation: pager1.into_continuation(),
493    ///         ..Default::default()
494    ///     },
495    ///     ..Default::default()
496    /// };
497    /// let pager2 = client.list_secret_properties(Some(options))?;
498    /// assert!(pager2.try_next().await?.is_some());
499    /// # Ok(()) }
500    /// ```
501    pub fn into_continuation(self) -> Option<PagerContinuation> {
502        self.continuation
503    }
504
505    /// Gets a [`PageIterator`] to iterate over pages instead of items.
506    ///
507    /// Resumes from the current page of items until after all items in the current page have been iterated
508    /// to avoid skipping items in the current page.
509    pub fn into_pages(self) -> PageIterator<P> {
510        let mut iter = self.iter;
511
512        // Start with the current page until after all items are iterated.
513        iter.options.continuation = self.continuation;
514        iter.state = iter
515            .options
516            .continuation
517            .as_ref()
518            .map_or_else(|| State::Init, |_| State::More);
519
520        iter
521    }
522}
523
524impl<P> Stream for ItemIterator<P>
525where
526    P: Page + Send,
527{
528    type Item = crate::Result<P::Item>;
529
530    fn poll_next(
531        self: Pin<&mut Self>,
532        cx: &mut std::task::Context<'_>,
533    ) -> std::task::Poll<Option<Self::Item>> {
534        let mut this = self.project();
535        let mut iter = this.iter.as_mut();
536        loop {
537            if let Some(current) = this.current.as_mut() {
538                if let Some(item) = current.next() {
539                    return task::Poll::Ready(Some(Ok(item)));
540                }
541
542                // Reset the iterator and poll for the next page.
543                *this.current = None;
544            }
545
546            // Set the current to the next page only after iterating through all items.
547            tracing::trace!(
548                "updating continuation from {:?} to {:?}",
549                &this.continuation,
550                iter.continuation(),
551            );
552            *this.continuation = iter.options.continuation.clone();
553
554            match iter.as_mut().poll_next(cx) {
555                task::Poll::Ready(page) => match page {
556                    Some(Ok(page)) => match page.into_items().poll_unpin(cx) {
557                        task::Poll::Ready(Ok(iter)) => {
558                            *this.current = Some(iter);
559                            continue;
560                        }
561                        task::Poll::Ready(Err(err)) => return task::Poll::Ready(Some(Err(err))),
562                        task::Poll::Pending => return task::Poll::Pending,
563                    },
564                    Some(Err(err)) => return task::Poll::Ready(Some(Err(err))),
565                    None => return task::Poll::Ready(None),
566                },
567                task::Poll::Pending => return task::Poll::Pending,
568            }
569        }
570    }
571}
572
573impl<P> fmt::Debug for ItemIterator<P>
574where
575    P: Page + Send,
576{
577    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
578        f.debug_struct("ItemIterator")
579            .field("iter", &self.iter)
580            .field("continuation", &self.continuation)
581            .finish_non_exhaustive()
582    }
583}
584
585impl<P> FusedStream for ItemIterator<P>
586where
587    P: Page + Send,
588{
589    fn is_terminated(&self) -> bool {
590        self.iter.is_terminated()
591    }
592}
593
594/// Iterates over a collection pages of items from a service.
595///
596/// # Examples
597///
598/// Some clients may return a `PageIterator` if there are no items to iterate or multiple items to iterate.
599/// The following example shows how you can also get a `PageIterator` from a [`Pager`] to iterate over pages instead of items.
600/// The pattern for iterating pages is otherwise the same:
601///
602/// ```no_run
603/// # use azure_core::{credentials::TokenCredential, http::Transport};
604/// # use azure_core_examples::secrets::{ResourceExt, SecretClient, SecretClientOptions};
605/// # use futures::TryStreamExt;
606/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
607/// # let credential: std::sync::Arc<dyn TokenCredential> = unimplemented!();
608/// let client = SecretClient::new(
609///     "https://my-vault.vault.azure.net",
610///     credential.clone(),
611///     None,
612/// )?;
613///
614/// // Iterate each page of secrets using a PageIterator.
615/// let mut pager = client.list_secret_properties(None)?.into_pages();
616/// while let Some(page) = pager.try_next().await? {
617///     let page = page.into_model()?;
618///     for secret in page.value {
619///         println!("{}", secret.resource_id()?.name);
620///     }
621/// }
622/// # Ok(()) }
623/// ```
624#[must_use = "streams do nothing unless you poll them"]
625#[pin_project(project = PageIteratorProjection, project_replace = PageIteratorProjectionOwned)]
626pub struct PageIterator<P>
627where
628    P: Send,
629{
630    #[pin]
631    make_request: PagerFn<P>,
632    options: PagerOptions<'static>,
633    state: State<P>,
634    added_span: bool,
635}
636
637impl<P> PageIterator<P>
638where
639    P: Send,
640{
641    /// Creates a [`PageIterator`] from a callback that will be called repeatedly to request each page.
642    ///
643    /// This method expect a callback that accepts a single [`PagerState`] parameter, and returns a [`PagerResult`] value asynchronously.
644    /// The result will be an asynchronous stream of [`Result`](crate::Result) values.
645    ///
646    /// The first time your callback is called, it will be called with [`PagerState::Initial`], indicating no next link/continuation token is present.
647    ///
648    /// Your callback must return one of:
649    /// * `Ok(result)` - The request succeeded, and the provided [`PagerResult`] indicates the value to return and if there are more pages.
650    /// * `Err(..)` - The request failed. The error will be yielded to the stream, the stream will end, and the callback will not be called again.
651    ///
652    /// ## Examples
653    ///
654    /// To page results using a next link:
655    ///
656    /// ```rust,no_run
657    /// # use azure_core::{Result, http::{RawResponse, pager::{PageIterator, PagerContinuation, PagerOptions, PagerResult, PagerState}, Pipeline, Request, Response, Method, Url}, json};
658    /// # let api_version = "2025-06-04".to_string();
659    /// # let pipeline: Pipeline = panic!("Not a runnable example");
660    /// #[derive(serde::Deserialize)]
661    /// struct ListItemsResult {
662    ///     items: Vec<String>,
663    ///     next_link: Option<String>,
664    /// }
665    /// let url = "https://example.com/my_paginated_api".parse().unwrap();
666    /// let mut base_req = Request::new(url, Method::Get);
667    /// let pager = PageIterator::new(move |next_link: PagerState, options: PagerOptions<'static>| {
668    ///     // The callback must be 'static, so you have to clone and move any values you want to use.
669    ///     let pipeline = pipeline.clone();
670    ///     let api_version = api_version.clone();
671    ///     let mut req = base_req.clone();
672    ///     Box::pin(async move {
673    ///         if let PagerState::More(next_link) = next_link {
674    ///             let next_link: Url = next_link.try_into().expect("expected Url");
675    ///             // Ensure the api-version from the client is appended.
676    ///             let qp = next_link
677    ///                 .query_pairs()
678    ///                 .filter(|(name, _)| name.ne("api-version"));
679    ///             req
680    ///                 .url_mut()
681    ///                 .query_pairs_mut()
682    ///                 .clear()
683    ///                 .extend_pairs(qp)
684    ///                 .append_pair("api-version", &api_version);
685    ///         }
686    ///         let resp = pipeline
687    ///           .send(&options.context, &mut req, None)
688    ///           .await?;
689    ///         let (status, headers, body) = resp.deconstruct();
690    ///         let result: ListItemsResult = json::from_json(&body)?;
691    ///         let resp: Response<ListItemsResult> = RawResponse::from_bytes(status, headers, body).into();
692    ///         Ok(match result.next_link {
693    ///             Some(next_link) => PagerResult::More {
694    ///                 response: resp,
695    ///                 continuation: PagerContinuation::Link(next_link.parse()?),
696    ///             },
697    ///             None => PagerResult::Done { response: resp }
698    ///         })
699    ///     })
700    /// }, None);
701    /// ```
702    ///
703    /// To page results using headers:
704    ///
705    /// ```rust,no_run
706    /// # use azure_core::{Result, http::{Context, pager::{PageIterator, PagerResult, PagerState}, Pipeline, Request, Response, Method, headers::HeaderName}};
707    /// # let pipeline: Pipeline = panic!("Not a runnable example");
708    /// #[derive(serde::Deserialize)]
709    /// struct ListItemsResult {
710    ///     items: Vec<String>,
711    /// }
712    /// let url = "https://example.com/my_paginated_api".parse().unwrap();
713    /// let mut base_req = Request::new(url, Method::Get);
714    /// let pager = PageIterator::new(move |continuation, options| {
715    ///     // The callback must be 'static, so you have to clone and move any values you want to use.
716    ///     let pipeline = pipeline.clone();
717    ///     let mut req = base_req.clone();
718    ///     Box::pin(async move {
719    ///         if let PagerState::More(continuation) = continuation {
720    ///             req.insert_header("x-ms-continuation", continuation.as_ref().to_string());
721    ///         }
722    ///         let resp: Response<ListItemsResult> = pipeline
723    ///           .send(&options.context, &mut req, None)
724    ///           .await?
725    ///           .into();
726    ///         Ok(PagerResult::from_response_header(resp, &HeaderName::from_static("x-ms-continuation")))
727    ///     })
728    /// }, None);
729    /// ```
730    pub fn new<
731        F: Fn(PagerState, PagerOptions<'static>) -> PagerResultFuture<P> + Send + 'static,
732    >(
733        make_request: F,
734        options: Option<PagerOptions<'static>>,
735    ) -> Self {
736        let options = options.unwrap_or_default();
737        let state = options
738            .continuation
739            .as_ref()
740            .map_or_else(|| State::Init, |_| State::More);
741
742        Self {
743            make_request: Box::new(make_request),
744            options,
745            state,
746            added_span: false,
747        }
748    }
749
750    /// Gets the continuation token to pass to [`PagerOptions`] to resume paging in another iterator.
751    pub fn continuation(&self) -> Option<&PagerContinuation> {
752        self.options.continuation.as_ref()
753    }
754
755    /// Gets the continuation token to pass to [`PagerOptions`] to resume paging in another iterator.
756    ///
757    /// # Examples
758    ///
759    /// This takes ownership of the iterator and can be useful when constructing a new iterator.
760    ///
761    /// ```no_run
762    /// # use azure_core::http::pager::PagerOptions;
763    /// # use azure_core_examples::secrets::{models::SecretClientListSecretPropertiesOptions, SecretClient};
764    /// # use futures::stream::TryStreamExt;
765    /// # #[tokio::main] async fn main() -> azure_core::Result<()> {
766    /// # let client: SecretClient = unimplemented!();
767    /// let pager1 = client.list_secret_properties(None)?.into_pages();
768    /// assert!(pager1.try_next().await?.is_some());
769    /// let options = SecretClientListSecretPropertiesOptions {
770    ///     method_options: PagerOptions {
771    ///         continuation: pager1.into_continuation(),
772    ///         ..Default::default()
773    ///     },
774    ///     ..Default::default()
775    /// };
776    /// let pager2 = client.list_secret_properties(Some(options))?.into_pages();
777    /// assert!(pager2.try_next().await?.is_some());
778    /// # Ok(()) }
779    /// ```
780    pub fn into_continuation(self) -> Option<PagerContinuation> {
781        self.options.continuation
782    }
783}
784
785impl<P> Stream for PageIterator<P>
786where
787    P: Send,
788{
789    type Item = crate::Result<P>;
790
791    fn poll_next(
792        self: Pin<&mut Self>,
793        cx: &mut std::task::Context<'_>,
794    ) -> std::task::Poll<Option<Self::Item>> {
795        let this = self.project();
796
797        // When in the initial state or resuming from a continuation token,
798        // attach a span to the context for the entire paging operation.
799        if *this.state == State::Init || this.options.continuation.is_some() {
800            tracing::debug!("establish a public API span for new pager.");
801
802            // At the very start of polling, create a span for the entire request, and attach it to the context
803            let span = create_public_api_span(&this.options.context, None, None);
804            if let Some(ref s) = span {
805                *this.added_span = true;
806                this.options.context.insert(s.clone());
807            }
808        }
809
810        let result = match *this.state {
811            State::Init => {
812                tracing::debug!("initial page request");
813                let options = this.options.clone();
814                let mut fut = (this.make_request)(PagerState::Initial, options);
815
816                match fut.poll_unpin(cx) {
817                    task::Poll::Ready(result) => result,
818                    task::Poll::Pending => {
819                        *this.state = State::Pending(fut);
820                        return task::Poll::Pending;
821                    }
822                }
823            }
824            State::Pending(ref mut fut) => task::ready!(fut.poll_unpin(cx)),
825            State::More => {
826                let options = this.options.clone();
827                let continuation_token = options
828                    .continuation
829                    .clone()
830                    // We should always have a continuation_token with `State::More`.
831                    .expect("expected continuation_token");
832                tracing::debug!("subsequent page request to {:?}", &continuation_token,);
833
834                let mut fut = (this.make_request)(PagerState::More(continuation_token), options);
835
836                match fut.poll_unpin(cx) {
837                    task::Poll::Ready(result) => result,
838                    task::Poll::Pending => {
839                        *this.state = State::Pending(fut);
840                        return task::Poll::Pending;
841                    }
842                }
843            }
844            State::Done => {
845                tracing::debug!("done");
846                // Set the `continuation_token` to None now that we are done.
847                this.options.continuation = None;
848                return task::Poll::Ready(None);
849            }
850        };
851
852        // Update continuation token and instrumentation.
853        match result {
854            Err(e) => {
855                if *this.added_span {
856                    if let Some(span) = this.options.context.value::<Arc<dyn Span>>() {
857                        // Mark the span as an error with an appropriate description.
858                        span.set_status(SpanStatus::Error {
859                            description: e.to_string(),
860                        });
861                        span.set_attribute("error.type", e.kind().to_string().into());
862                        span.end();
863                    }
864                }
865
866                *this.state = State::Done;
867                task::Poll::Ready(Some(Err(e)))
868            }
869
870            Ok(PagerResult::More {
871                response,
872                continuation: continuation_token,
873            }) => {
874                // Set the `continuation_token` to the next page.
875                this.options.continuation = Some(continuation_token);
876                *this.state = State::More;
877                task::Poll::Ready(Some(Ok(response)))
878            }
879
880            Ok(PagerResult::Done { response }) => {
881                // Set the `continuation_token` to None now that we are done.
882                this.options.continuation = None;
883                *this.state = State::Done;
884
885                // When the result is done, finalize the span. Note that we only do that if we created the span in the first place;
886                // otherwise, it is the responsibility of the caller to end their span.
887                if *this.added_span {
888                    if let Some(span) = this.options.context.value::<Arc<dyn Span>>() {
889                        span.end();
890                    }
891                }
892
893                task::Poll::Ready(Some(Ok(response)))
894            }
895        }
896    }
897}
898
899impl<P> fmt::Debug for PageIterator<P>
900where
901    P: Send,
902{
903    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
904        f.debug_struct("PageIterator")
905            .field("continuation_token", &self.options.continuation)
906            .field("options", &self.options)
907            .field("state", &self.state)
908            .field("added_span", &self.added_span)
909            .finish_non_exhaustive()
910    }
911}
912
913impl<P> FusedStream for PageIterator<P>
914where
915    P: Send,
916{
917    fn is_terminated(&self) -> bool {
918        self.state == State::Done
919    }
920}
921
922enum State<P> {
923    Init,
924    Pending(PagerResultFuture<P>),
925    More,
926    Done,
927}
928
929impl<P> fmt::Debug for State<P> {
930    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
931        match self {
932            State::Init => f.write_str("Init"),
933            State::Pending(..) => f.debug_tuple("Pending").finish_non_exhaustive(),
934            State::More => f.write_str("More"),
935            State::Done => f.write_str("Done"),
936        }
937    }
938}
939
940impl<P> PartialEq for State<P> {
941    fn eq(&self, other: &Self) -> bool {
942        // Only needs to compare if both states are Init or Done; internally, we don't care about any other states.
943        matches!(
944            (self, other),
945            (State::Init, State::Init) | (State::Done, State::Done)
946        )
947    }
948}
949
950#[cfg(test)]
951mod tests {
952    use super::{
953        ItemIterator, PageIterator, Pager, PagerContinuation, PagerOptions, PagerResult, PagerState,
954    };
955    use crate::http::{
956        headers::{HeaderName, HeaderValue},
957        pager::PagerResultFuture,
958        RawResponse, Response, StatusCode,
959    };
960    use async_trait::async_trait;
961    use futures::{StreamExt as _, TryStreamExt};
962    use serde::Deserialize;
963    use std::collections::HashMap;
964
965    #[derive(Deserialize, Debug, PartialEq, Eq)]
966    struct Page {
967        pub items: Vec<i32>,
968        pub page: Option<i32>,
969    }
970
971    #[async_trait]
972    impl super::Page for Page {
973        type Item = i32;
974        type IntoIter = <Vec<i32> as IntoIterator>::IntoIter;
975
976        async fn into_items(self) -> crate::Result<Self::IntoIter> {
977            Ok(self.items.into_iter())
978        }
979    }
980
981    #[tokio::test]
982    async fn callback_item_pagination() {
983        let pager: Pager<Page> = Pager::new(
984            |continuation: PagerState, _ctx| {
985                Box::pin(async move {
986                    match continuation {
987                        PagerState::Initial => Ok(PagerResult::More {
988                            response: RawResponse::from_bytes(
989                                StatusCode::Ok,
990                                HashMap::from([(
991                                    HeaderName::from_static("x-test-header"),
992                                    HeaderValue::from_static("page-1"),
993                                )])
994                                .into(),
995                                r#"{"items":[1],"page":1}"#,
996                            )
997                            .into(),
998                            continuation: PagerContinuation::Token("1".into()),
999                        }),
1000                        PagerState::More(ref i) if i.as_ref() == "1" => Ok(PagerResult::More {
1001                            response: RawResponse::from_bytes(
1002                                StatusCode::Ok,
1003                                HashMap::from([(
1004                                    HeaderName::from_static("x-test-header"),
1005                                    HeaderValue::from_static("page-2"),
1006                                )])
1007                                .into(),
1008                                r#"{"items":[2],"page":2}"#,
1009                            )
1010                            .into(),
1011                            continuation: PagerContinuation::Token("2".into()),
1012                        }),
1013                        PagerState::More(ref i) if i.as_ref() == "2" => Ok(PagerResult::Done {
1014                            response: RawResponse::from_bytes(
1015                                StatusCode::Ok,
1016                                HashMap::from([(
1017                                    HeaderName::from_static("x-test-header"),
1018                                    HeaderValue::from_static("page-3"),
1019                                )])
1020                                .into(),
1021                                r#"{"items":[3],"page":3}"#,
1022                            )
1023                            .into(),
1024                        }),
1025                        _ => {
1026                            panic!("Unexpected continuation value")
1027                        }
1028                    }
1029                })
1030            },
1031            None,
1032        );
1033        let items: Vec<i32> = pager.try_collect().await.unwrap();
1034        assert_eq!(vec![1, 2, 3], items.as_slice())
1035    }
1036
1037    #[tokio::test]
1038    async fn callback_item_pagination_error() {
1039        let pager: Pager<Page> = ItemIterator::new(
1040            |continuation: PagerState, _options| {
1041                Box::pin(async move {
1042                    match continuation {
1043                        PagerState::Initial => Ok(PagerResult::More {
1044                            response: RawResponse::from_bytes(
1045                                StatusCode::Ok,
1046                                HashMap::from([(
1047                                    HeaderName::from_static("x-test-header"),
1048                                    HeaderValue::from_static("page-1"),
1049                                )])
1050                                .into(),
1051                                r#"{"items":[1],"page":1}"#,
1052                            )
1053                            .into(),
1054                            continuation: PagerContinuation::Token("1".into()),
1055                        }),
1056                        PagerState::More(ref i) if i.as_ref() == "1" => {
1057                            Err(crate::Error::with_message(
1058                                crate::error::ErrorKind::Other,
1059                                "yon request didst fail",
1060                            ))
1061                        }
1062                        _ => {
1063                            panic!("Unexpected continuation value")
1064                        }
1065                    }
1066                })
1067            },
1068            None,
1069        );
1070        let pages: Vec<Result<(String, Page), crate::Error>> = pager
1071            .into_pages()
1072            .then(|r| async move {
1073                let r = r?;
1074                let header = r
1075                    .headers()
1076                    .get_optional_string(&HeaderName::from_static("x-test-header"))
1077                    .unwrap();
1078                let body = r.into_model()?;
1079                Ok((header, body))
1080            })
1081            .collect()
1082            .await;
1083        assert_eq!(2, pages.len());
1084        assert_eq!(
1085            &(
1086                "page-1".to_string(),
1087                Page {
1088                    items: vec![1],
1089                    page: Some(1)
1090                }
1091            ),
1092            pages[0].as_ref().unwrap()
1093        );
1094
1095        let err = pages[1].as_ref().unwrap_err();
1096        assert_eq!(&crate::error::ErrorKind::Other, err.kind());
1097        assert_eq!("yon request didst fail", format!("{}", err));
1098    }
1099
1100    #[tokio::test]
1101    async fn page_iterator_iterate_all_pages() {
1102        // Create a PageIterator and iterate through all three pages.
1103        let mut pager = PageIterator::new(make_three_page_callback(), None);
1104
1105        // Should start with no continuation_token.
1106        assert_eq!(pager.continuation(), None);
1107
1108        // Get first page.
1109        let first_page = pager
1110            .next()
1111            .await
1112            .expect("expected first page")
1113            .expect("expected successful first page")
1114            .into_model()
1115            .expect("expected page");
1116        assert_eq!(first_page.page, Some(1));
1117        assert_eq!(first_page.items, vec![1, 2, 3]);
1118
1119        // continuation_token should now point to second page.
1120        assert_eq!(
1121            pager.continuation().map(AsRef::as_ref),
1122            Some("next-token-1")
1123        );
1124
1125        // Get second page.
1126        let second_page = pager
1127            .next()
1128            .await
1129            .expect("expected second page")
1130            .expect("expected successful second page")
1131            .into_model()
1132            .expect("expected page");
1133        assert_eq!(second_page.page, Some(2));
1134        assert_eq!(second_page.items, vec![4, 5, 6]);
1135
1136        // continuation_token should now point to third page.
1137        assert_eq!(
1138            pager.continuation().map(AsRef::as_ref),
1139            Some("next-token-2")
1140        );
1141
1142        // Get third page.
1143        let third_page = pager
1144            .next()
1145            .await
1146            .expect("expected third page")
1147            .expect("expected successful third page")
1148            .into_model()
1149            .expect("expected page");
1150        assert_eq!(third_page.page, None);
1151        assert_eq!(third_page.items, vec![7, 8, 9]);
1152
1153        // continuation_token should now be None (done).
1154        assert_eq!(pager.continuation(), None);
1155
1156        // Verify stream is exhausted.
1157        assert!(pager.next().await.is_none());
1158    }
1159
1160    #[tokio::test]
1161    async fn page_iterator_with_continuation_token() {
1162        // Create the first PageIterator.
1163        let mut first_pager = PageIterator::new(make_three_page_callback(), None);
1164
1165        // Should start with no continuation_token.
1166        assert_eq!(first_pager.continuation(), None);
1167
1168        // Advance to the first page.
1169        let first_page = first_pager
1170            .next()
1171            .await
1172            .expect("expected first page")
1173            .expect("expected successful first page")
1174            .into_model()
1175            .expect("expected page");
1176        assert_eq!(first_page.page, Some(1));
1177        assert_eq!(first_page.items, vec![1, 2, 3]);
1178
1179        // continuation_token should point to second page.
1180        let continuation_token = first_pager
1181            .continuation()
1182            .expect("expected continuation_token from first page");
1183        assert_eq!(continuation_token.as_ref(), "next-token-1");
1184
1185        // Create the second PageIterator.
1186        let mut second_pager = PageIterator::new(
1187            make_three_page_callback(),
1188            Some(PagerOptions {
1189                continuation: Some(continuation_token.clone()),
1190                ..Default::default()
1191            }),
1192        );
1193
1194        // Should start with link to second page.
1195        assert_eq!(
1196            second_pager.continuation().map(AsRef::as_ref),
1197            Some("next-token-1"),
1198        );
1199
1200        // Advance to second page.
1201        let second_page = second_pager
1202            .next()
1203            .await
1204            .expect("expected second page")
1205            .expect("expected successful second page")
1206            .into_model()
1207            .expect("expected page");
1208        assert_eq!(second_page.page, Some(2));
1209        assert_eq!(second_page.items, vec![4, 5, 6]);
1210        assert_eq!(
1211            second_pager.continuation().map(AsRef::as_ref),
1212            Some("next-token-2")
1213        );
1214
1215        // Advance to last page.
1216        let last_page = second_pager
1217            .next()
1218            .await
1219            .expect("expected last page")
1220            .expect("expected successful last page")
1221            .into_model()
1222            .expect("expected page");
1223        assert_eq!(last_page.page, None);
1224        assert_eq!(last_page.items, vec![7, 8, 9]);
1225        assert_eq!(second_pager.continuation(), None);
1226    }
1227
1228    #[tokio::test]
1229    async fn page_iterator_from_item_iterator_after_first_page() {
1230        // Create an ItemIterator and consume all items from first page.
1231        let mut item_pager = ItemIterator::new(make_three_page_callback(), None);
1232
1233        // Should start with no continuation_token.
1234        assert_eq!(item_pager.continuation(), None);
1235
1236        // Consume all three items from the first page.
1237        let first_item = item_pager
1238            .next()
1239            .await
1240            .expect("expected first item")
1241            .expect("expected successful first item");
1242        assert_eq!(first_item, 1);
1243
1244        let second_item = item_pager
1245            .next()
1246            .await
1247            .expect("expected second item")
1248            .expect("expected successful second item");
1249        assert_eq!(second_item, 2);
1250
1251        let third_item = item_pager
1252            .next()
1253            .await
1254            .expect("expected third item")
1255            .expect("expected successful third item");
1256        assert_eq!(third_item, 3);
1257
1258        // Convert to PageIterator after consuming first page.
1259        let mut page_pager = item_pager.into_pages();
1260
1261        // Should start with None initially.
1262        assert_eq!(page_pager.continuation(), None);
1263
1264        // Verify we start over with the first page again (ItemIterator.continuation_token() was None).
1265        let first_page = page_pager
1266            .next()
1267            .await
1268            .expect("expected first page")
1269            .expect("expected successful first page")
1270            .into_model()
1271            .expect("expected page");
1272        assert_eq!(first_page.page, Some(1));
1273        assert_eq!(first_page.items, vec![1, 2, 3]);
1274
1275        // continuation_token should now point to second page.
1276        let continuation_token = page_pager
1277            .continuation()
1278            .expect("expected continuation_token from first page");
1279        assert_eq!(continuation_token.as_ref(), "next-token-1");
1280    }
1281
1282    #[tokio::test]
1283    async fn page_iterator_from_item_iterator_second_page_first_item() {
1284        // Create an ItemIterator and consume items up to first item of second page.
1285        let mut item_pager = ItemIterator::new(make_three_page_callback(), None);
1286
1287        // Should start with no continuation_token.
1288        assert_eq!(item_pager.continuation(), None);
1289
1290        // Consume all three items from the first page.
1291        let first_item = item_pager
1292            .next()
1293            .await
1294            .expect("expected first item")
1295            .expect("expected successful first item");
1296        assert_eq!(first_item, 1);
1297
1298        let second_item = item_pager
1299            .next()
1300            .await
1301            .expect("expected second item")
1302            .expect("expected successful second item");
1303        assert_eq!(second_item, 2);
1304
1305        let third_item = item_pager
1306            .next()
1307            .await
1308            .expect("expected third item")
1309            .expect("expected successful third item");
1310        assert_eq!(third_item, 3);
1311
1312        // Get first item from second page.
1313        let fourth_item = item_pager
1314            .next()
1315            .await
1316            .expect("expected fourth item")
1317            .expect("expected successful fourth item");
1318        assert_eq!(fourth_item, 4);
1319
1320        // Convert to PageIterator after consuming first item of second page.
1321        let mut page_pager = item_pager.into_pages();
1322
1323        // Should start with second page since that's where we were.
1324        assert_eq!(
1325            page_pager.continuation().map(AsRef::as_ref),
1326            Some("next-token-1")
1327        );
1328
1329        // Get second page - should be the second page.
1330        let second_page = page_pager
1331            .next()
1332            .await
1333            .expect("expected second page")
1334            .expect("expected successful second page")
1335            .into_model()
1336            .expect("expected page");
1337        assert_eq!(second_page.page, Some(2));
1338        assert_eq!(second_page.items, vec![4, 5, 6]);
1339
1340        // continuation_token should now point to third page.
1341        let continuation_token = page_pager
1342            .continuation()
1343            .expect("expected continuation_token from second page");
1344        assert_eq!(continuation_token.as_ref(), "next-token-2");
1345    }
1346
1347    #[tokio::test]
1348    async fn item_iterator_with_continuation_token() {
1349        // Create the first ItemIterator.
1350        let mut first_pager = ItemIterator::new(make_three_page_callback(), None);
1351
1352        // Should start with no continuation_token.
1353        assert_eq!(first_pager.continuation(), None);
1354
1355        // Get first item from first page.
1356        let first_item = first_pager
1357            .next()
1358            .await
1359            .expect("expected first item")
1360            .expect("expected successful first item");
1361        assert_eq!(first_item, 1);
1362
1363        // Get second item from first page.
1364        let second_item = first_pager
1365            .next()
1366            .await
1367            .expect("expected second item")
1368            .expect("expected successful second item");
1369        assert_eq!(second_item, 2);
1370
1371        // continuation_token should point to current page after processing some, but not all, items.
1372        let continuation_token = first_pager.continuation();
1373        assert_eq!(continuation_token, None);
1374
1375        // Create the second ItemIterator with continuation token.
1376        let mut second_pager = ItemIterator::new(
1377            make_three_page_callback(),
1378            Some(PagerOptions {
1379                continuation: continuation_token.cloned(),
1380                ..Default::default()
1381            }),
1382        );
1383
1384        // Should start with link to first page.
1385        assert_eq!(second_pager.continuation(), None);
1386
1387        // When continuing with a continuation token, we should start over from the
1388        // beginning of the page, not where we left off in the item stream.
1389        // This means we should get the first item of the first page (1), not the
1390        // third item of the first page (3).
1391        let first_item_second_pager = second_pager
1392            .next()
1393            .await
1394            .expect("expected first item from second pager")
1395            .expect("expected successful first item from second pager");
1396        assert_eq!(first_item_second_pager, 1);
1397
1398        // Get remaining items.
1399        let items: Vec<i32> = second_pager.try_collect().await.unwrap();
1400        assert_eq!(items.as_slice(), vec![2, 3, 4, 5, 6, 7, 8, 9]);
1401    }
1402
1403    #[tokio::test]
1404    async fn item_iterator_continuation_second_page_second_item() {
1405        // Create the first ItemIterator.
1406        let mut first_pager = ItemIterator::new(make_three_page_callback(), None);
1407
1408        // Should start with no continuation_token.
1409        assert_eq!(first_pager.continuation(), None);
1410
1411        // Iterate to the second item of the second page.
1412        // First page: items 1, 2, 3
1413        let first_item = first_pager
1414            .next()
1415            .await
1416            .expect("expected first item")
1417            .expect("expected successful first item");
1418        assert_eq!(first_item, 1);
1419
1420        let second_item = first_pager
1421            .next()
1422            .await
1423            .expect("expected second item")
1424            .expect("expected successful second item");
1425        assert_eq!(second_item, 2);
1426
1427        let third_item = first_pager
1428            .next()
1429            .await
1430            .expect("expected third item")
1431            .expect("expected successful third item");
1432        assert_eq!(third_item, 3);
1433
1434        // Second page: item 4 (first of second page)
1435        let fourth_item = first_pager
1436            .next()
1437            .await
1438            .expect("expected fourth item")
1439            .expect("expected successful fourth item");
1440        assert_eq!(fourth_item, 4);
1441
1442        // Second page: item 5 (second of second page)
1443        let fifth_item = first_pager
1444            .next()
1445            .await
1446            .expect("expected fifth item")
1447            .expect("expected successful fifth item");
1448        assert_eq!(fifth_item, 5);
1449
1450        // Get continuation token - should point to current page (second page).
1451        let continuation_token = first_pager.into_continuation();
1452        assert_eq!(
1453            continuation_token.as_ref().map(AsRef::as_ref),
1454            Some("next-token-1")
1455        );
1456
1457        // Create the second ItemIterator with continuation token.
1458        let mut second_pager = ItemIterator::new(
1459            make_three_page_callback(),
1460            Some(PagerOptions {
1461                continuation: continuation_token,
1462                ..Default::default()
1463            }),
1464        );
1465
1466        // When continuing with a continuation token, we should start over from the
1467        // beginning of the current page (second page), not where we left off.
1468        // This means we should get the first item of the second page (4).
1469        let first_item_second_pager = second_pager
1470            .next()
1471            .await
1472            .expect("expected first item from second pager")
1473            .expect("expected successful first item from second pager");
1474        assert_eq!(first_item_second_pager, 4);
1475
1476        // Get remaining items.
1477        let items: Vec<i32> = second_pager.try_collect().await.unwrap();
1478        assert_eq!(items.as_slice(), vec![5, 6, 7, 8, 9]);
1479    }
1480
1481    #[tokio::test]
1482    async fn item_iterator_continuation_after_first_page() {
1483        // Create the first ItemIterator.
1484        let mut first_pager = ItemIterator::new(make_three_page_callback(), None);
1485
1486        // Should start with no continuation_token.
1487        assert_eq!(first_pager.continuation(), None);
1488
1489        // Iterate past the third item of the first page (all items of first page).
1490        let first_item = first_pager
1491            .next()
1492            .await
1493            .expect("expected first item")
1494            .expect("expected successful first item");
1495        assert_eq!(first_item, 1);
1496
1497        let second_item = first_pager
1498            .next()
1499            .await
1500            .expect("expected second item")
1501            .expect("expected successful second item");
1502        assert_eq!(second_item, 2);
1503
1504        let third_item = first_pager
1505            .next()
1506            .await
1507            .expect("expected third item")
1508            .expect("expected successful third item");
1509        assert_eq!(third_item, 3);
1510
1511        // Get continuation token after finishing the first page - should still point to current page (first page).
1512        let continuation_token = first_pager.continuation();
1513        assert_eq!(continuation_token, None);
1514
1515        // Create the second ItemIterator with continuation token.
1516        let mut second_pager = ItemIterator::new(
1517            make_three_page_callback(),
1518            Some(PagerOptions {
1519                continuation: continuation_token.cloned(),
1520                ..Default::default()
1521            }),
1522        );
1523
1524        // When continuing with a continuation token after finishing a page, we should
1525        // start from the beginning of the current page.
1526        // This means we should get the first item of the second page (4).
1527        let first_item_second_pager = second_pager
1528            .next()
1529            .await
1530            .expect("expected first item from first pager")
1531            .expect("expected successful first item from first pager");
1532        assert_eq!(first_item_second_pager, 1);
1533
1534        // Get remaining items.
1535        let items: Vec<i32> = second_pager.try_collect().await.unwrap();
1536        assert_eq!(items.as_slice(), vec![2, 3, 4, 5, 6, 7, 8, 9]);
1537    }
1538
1539    #[tokio::test]
1540    #[deny(unfulfilled_lint_expectations)]
1541    async fn must_use_item_iterator() {
1542        fn create() -> crate::Result<Pager<Page>> {
1543            Ok(ItemIterator::new(make_three_page_callback(), None))
1544        }
1545
1546        #[expect(unused_must_use)]
1547        create();
1548
1549        #[expect(unused_must_use)]
1550        create().unwrap();
1551
1552        // We cannot use #[must_use] to force a stream to be polled.
1553        #[deny(unused_must_use)]
1554        let _pager = create().unwrap();
1555
1556        #[deny(unused_must_use)]
1557        let pager = create().unwrap();
1558        let items: Vec<i32> = pager.try_collect().await.unwrap();
1559        assert_eq!(items.as_slice(), vec![1, 2, 3, 4, 5, 6, 7, 8, 9]);
1560    }
1561
1562    #[tokio::test]
1563    #[deny(unfulfilled_lint_expectations)]
1564    async fn must_use_page_iterator() {
1565        fn create() -> crate::Result<PageIterator<Response<Page>>> {
1566            Ok(PageIterator::new(make_three_page_callback(), None))
1567        }
1568
1569        #[expect(unused_must_use)]
1570        create();
1571
1572        #[expect(unused_must_use)]
1573        create().unwrap();
1574
1575        // We cannot use #[must_use] to force a stream to be polled.
1576        #[deny(unused_must_use)]
1577        let _pager = create().unwrap();
1578
1579        #[deny(unused_must_use)]
1580        let mut pager = create().unwrap();
1581        let i = pager
1582            .try_next()
1583            .await
1584            .expect("expected successful response")
1585            .expect("expected some page response")
1586            .into_model()
1587            .expect("expected page");
1588        assert_eq!(i.items.as_slice(), vec![1, 2, 3]);
1589    }
1590
1591    #[allow(clippy::type_complexity)]
1592    fn make_three_page_callback(
1593    ) -> impl Fn(PagerState, PagerOptions<'_>) -> PagerResultFuture<Response<Page>> {
1594        |continuation: PagerState, _options| {
1595            Box::pin(async move {
1596                match continuation {
1597                    PagerState::Initial => Ok(PagerResult::More {
1598                        response: RawResponse::from_bytes(
1599                            StatusCode::Ok,
1600                            Default::default(),
1601                            r#"{"items":[1,2,3],"page":1}"#,
1602                        )
1603                        .into(),
1604                        continuation: PagerContinuation::Token("next-token-1".into()),
1605                    }),
1606                    PagerState::More(continuation) if continuation.as_ref() == "next-token-1" => {
1607                        Ok(PagerResult::More {
1608                            response: RawResponse::from_bytes(
1609                                StatusCode::Ok,
1610                                HashMap::from([(
1611                                    HeaderName::from_static("x-test-header"),
1612                                    HeaderValue::from_static("page-2"),
1613                                )])
1614                                .into(),
1615                                r#"{"items":[4,5,6],"page":2}"#,
1616                            )
1617                            .into(),
1618                            continuation: PagerContinuation::Token("next-token-2".into()),
1619                        })
1620                    }
1621                    PagerState::More(continuation) if continuation.as_ref() == "next-token-2" => {
1622                        Ok(PagerResult::Done {
1623                            response: RawResponse::from_bytes(
1624                                StatusCode::Ok,
1625                                HashMap::from([(
1626                                    HeaderName::from_static("x-test-header"),
1627                                    HeaderValue::from_static("page-3"),
1628                                )])
1629                                .into(),
1630                                r#"{"items":[7,8,9]}"#,
1631                            )
1632                            .into(),
1633                        })
1634                    }
1635                    _ => {
1636                        panic!("Unexpected continuation value: {:?}", continuation)
1637                    }
1638                }
1639            })
1640        }
1641    }
1642}