Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 9 additions & 3 deletions src/source_api/cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<SourceProductList, ProxyError> {
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,
Expand Down
33 changes: 19 additions & 14 deletions src/source_api/registry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Vec<String>, 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),
}
}
}
}

Expand Down
8 changes: 7 additions & 1 deletion src/source_api/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -137,5 +137,11 @@ impl DataConnectionDetails {

#[derive(Debug, Clone, Deserialize)]
pub struct SourceProductList {
pub products: Vec<SourceProduct>,
/// 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<SourceProduct>,
/// The next page's `cursor`; absent on the last page.
#[serde(default)]
pub next_cursor: Option<String>,
}
22 changes: 14 additions & 8 deletions tests/fixtures.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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) ──────────
Expand Down
4 changes: 2 additions & 2 deletions tests/stub_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
11 changes: 7 additions & 4 deletions tests/test_contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -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():
Expand Down
Loading