diff --git a/src/source_api/cache.rs b/src/source_api/cache.rs index e382499..f3c8351 100644 --- a/src/source_api/cache.rs +++ b/src/source_api/cache.rs @@ -157,19 +157,25 @@ pub async fn get_or_fetch_data_connection( .await } -/// Fetch an account's product list, cached for `PRODUCT_LIST_CACHE_SECS`. +/// Fetch one page of an account's product list, cached for +/// `PRODUCT_LIST_CACHE_SECS`: the first page, or the one `cursor` names. pub async fn get_or_fetch_product_list( api_base_url: &str, account: &str, api_auth: &crate::ApiAuth, request_id: &str, subject: Option<&str>, + cursor: Option<&str>, ) -> Result { - let api_url = format!( - "{}/api/v1/products/{}", + let mut api_url = format!( + "{}/api/v1/products/{}?limit=100", api_base_url, utf8_percent_encode(account, PATH_SEGMENT), ); + if let Some(cursor) = cursor { + api_url.push_str("&cursor="); + api_url.extend(utf8_percent_encode(cursor, NON_ALPHANUMERIC)); + } let cache_key = cache_key_with_subject(&api_url, subject); cached_fetch( &cache_key, diff --git a/src/source_api/registry.rs b/src/source_api/registry.rs index d23141c..9a048a9 100644 --- a/src/source_api/registry.rs +++ b/src/source_api/registry.rs @@ -25,21 +25,26 @@ impl SourceCoopRegistry { } } - /// List products for an account via the Source API. + /// List products for an account via the Source API, every page of it. pub async fn list_products(&self, account: &str) -> Result, ProxyError> { - let product_list = super::cache::get_or_fetch_product_list( - &self.api_base_url, - account, - &self.api_auth, - &self.request_id, - None, - ) - .await?; - Ok(product_list - .products - .into_iter() - .map(|p| p.product_id) - .collect()) + let mut ids = Vec::new(); + let mut cursor = None; + loop { + let page = super::cache::get_or_fetch_product_list( + &self.api_base_url, + account, + &self.api_auth, + &self.request_id, + None, + cursor.as_deref(), + ) + .await?; + ids.extend(page.items.into_iter().map(|p| p.product_id)); + match page.next_cursor { + Some(next) => cursor = Some(next), + None => return Ok(ids), + } + } } } diff --git a/src/source_api/types.rs b/src/source_api/types.rs index c341b1a..5a2ecd4 100644 --- a/src/source_api/types.rs +++ b/src/source_api/types.rs @@ -137,5 +137,11 @@ impl DataConnectionDetails { #[derive(Debug, Clone, Deserialize)] pub struct SourceProductList { - pub products: Vec, + /// One page of the account's products. `products`, unpaged, is the shape + /// the API had before source.coop#590. + #[serde(alias = "products")] + pub items: Vec, + /// The next page's `cursor`; absent on the last page. + #[serde(default)] + pub next_cursor: Option, } diff --git a/tests/fixtures.rs b/tests/fixtures.rs index cf5df36..6109848 100644 --- a/tests/fixtures.rs +++ b/tests/fixtures.rs @@ -71,14 +71,20 @@ fn restricted_product_fixture_is_not_public() { #[test] fn product_list_wrapper_parses() { - // The stub wraps the product fixture as {"products": [...]} for the - // account listing route; pin that wrapper shape too. - let json = format!( - r#"{{"products":[{}]}}"#, - include_str!("fixtures/product.json") - ); - let l: SourceProductList = serde_json::from_str(&json).unwrap(); - assert_eq!(l.products.len(), 1); + // The stub wraps the product fixture as a page, {"items": [...], + // "next_cursor": ...}, for the account listing route; pin that shape, and + // the unpaged {"products": [...]} the API answered with before it. + let product = include_str!("fixtures/product.json"); + let page: SourceProductList = + serde_json::from_str(&format!(r#"{{"items":[{product}],"next_cursor":"abc"}}"#)).unwrap(); + assert_eq!(page.items.len(), 1); + assert_eq!(page.next_cursor.as_deref(), Some("abc")); + let last: SourceProductList = + serde_json::from_str(&format!(r#"{{"items":[{product}],"next_cursor":null}}"#)).unwrap(); + assert_eq!(last.next_cursor, None); + let unpaged: SourceProductList = + serde_json::from_str(&format!(r#"{{"products":[{product}]}}"#)).unwrap(); + assert_eq!((unpaged.items.len(), unpaged.next_cursor), (1, None)); } // ── backend_options: provider → (backend_type, multistore options) ────────── diff --git a/tests/stub_api.py b/tests/stub_api.py index 457a22f..fed8699 100644 --- a/tests/stub_api.py +++ b/tests/stub_api.py @@ -90,12 +90,12 @@ def _fixture(name): RESTRICTED_PRODUCT_JSON = _fixture("product_restricted") ROUTES = { - f"/api/v1/products/{ACCOUNT}": {"products": [PRODUCT_JSON]}, + f"/api/v1/products/{ACCOUNT}": {"items": [PRODUCT_JSON], "next_cursor": None}, f"/api/v1/products/{ACCOUNT}/{PRODUCT}": PRODUCT_JSON, # No `authentication` field -> BackendAuth::Unsigned -> unsigned reads. f"/api/v1/data-connections/{CONNECTION}": _fixture("data_connection"), # Write probe (see above). - f"/api/v1/products/{WRITE_ACCOUNT}": {"products": [WRITE_PRODUCT_JSON]}, + f"/api/v1/products/{WRITE_ACCOUNT}": {"items": [WRITE_PRODUCT_JSON], "next_cursor": None}, f"/api/v1/products/{WRITE_ACCOUNT}/{WRITE_PRODUCT}": WRITE_PRODUCT_JSON, f"/api/v1/products/{WRITE_ACCOUNT}/{WRITE_PRODUCT}/permissions": ["read", "write"], f"/api/v1/data-connections/{WRITE_CONNECTION}": WRITE_CONNECTION_JSON, diff --git a/tests/test_contract.py b/tests/test_contract.py index d9ae10c..47d8fcb 100644 --- a/tests/test_contract.py +++ b/tests/test_contract.py @@ -72,12 +72,15 @@ def test_product_shape(): def test_product_list_shape(): - real = fetch(f"/api/v1/products/{ACCOUNT}") + # The proxy asks for pages of 100, as here. Until source.coop#590 deploys + # the API answers unpaged, as {"products": [...]}, which the proxy also reads. + real = fetch(f"/api/v1/products/{ACCOUNT}?limit=100") stub = ROUTES[f"/api/v1/products/{ACCOUNT}"] - assert json_type(real.get("products")) == "array", "products list missing" - matches = [p for p in real["products"] if p.get("product_id") == PRODUCT] + items = real.get("items", real.get("products")) + assert json_type(items) == "array", "items list missing" + matches = [p for p in items if p.get("product_id") == PRODUCT] assert matches, f"product {PRODUCT} missing from real list response" - assert_shape_subset(stub["products"][0], matches[0], "products[]") + assert_shape_subset(stub["items"][0], matches[0], "items[]") def test_data_connection_shape():