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}