Skip to main content

vantage_api_client/rest/
api.rs

1use std::sync::Arc;
2
3use ciborium::Value as CborValue;
4use indexmap::IndexMap;
5use reqwest::header::{AUTHORIZATION, HeaderMap, HeaderName, HeaderValue};
6use vantage_api_pool::resilient::{ResilientClient, TransportEvent, TransportObserver};
7use vantage_core::{Priority, error};
8use vantage_dataset::traits::Result;
9use vantage_expressions::Expression;
10use vantage_table::pagination::Pagination;
11use vantage_types::Record;
12
13use crate::transport::AuthHeader;
14
15/// How the API wraps its row array in the response body.
16///
17/// Most public APIs use one of these three shapes; the legacy vantage
18/// "wrapped under `data`" shape is `Wrapped { array_key: "data" }`.
19#[derive(Clone, Debug)]
20pub enum ResponseShape {
21    /// Body is a bare JSON array of records.
22    /// Example: `GET /users` → `[ {…}, {…} ]`. JSONPlaceholder, GitHub, etc.
23    BareArray,
24
25    /// Body is a JSON object with the array under a fixed key.
26    /// Example: `GET /users` → `{ "data": [ … ] }`.
27    Wrapped { array_key: String },
28
29    /// Body is a JSON object with the array under a key matching the
30    /// table name. Example (DummyJSON):
31    /// `GET /products` → `{ "products": [ … ], "total": …, "skip": …, "limit": … }`.
32    WrappedByTableName,
33}
34
35impl Default for ResponseShape {
36    /// Default matches the legacy 0.1.x shape: `{ "data": [...] }`.
37    fn default() -> Self {
38        ResponseShape::Wrapped {
39            array_key: "data".to_string(),
40        }
41    }
42}
43
44/// Names of the page/limit query parameters the API expects.
45///
46/// Defaults to `("_page", "_limit")` — the JSON Server convention used
47/// by JSONPlaceholder. DummyJSON uses `("skip", "limit")` (in items not
48/// pages). Customise via `RestApiBuilder::pagination_params`.
49#[derive(Clone, Debug)]
50pub struct PaginationParams {
51    pub page: String,
52    pub limit: String,
53    /// If true, the page parameter is sent as a *0-based item offset*
54    /// (`skip`) instead of a 1-based page index. DummyJSON-style.
55    pub skip_based: bool,
56}
57
58impl PaginationParams {
59    pub fn page_limit(page: impl Into<String>, limit: impl Into<String>) -> Self {
60        Self {
61            page: page.into(),
62            limit: limit.into(),
63            skip_based: false,
64        }
65    }
66
67    pub fn skip_limit(skip: impl Into<String>, limit: impl Into<String>) -> Self {
68        Self {
69            page: skip.into(),
70            limit: limit.into(),
71            skip_based: true,
72        }
73    }
74}
75
76impl Default for PaginationParams {
77    fn default() -> Self {
78        Self::page_limit("_page", "_limit")
79    }
80}
81
82/// REST API backend for Vantage — reads data from HTTP JSON endpoints.
83///
84/// Each table maps to an API endpoint: `{base_url}/{table_name}`.
85/// Response shape is configurable via [`RestApi::builder`]; see
86/// [`ResponseShape`] for the supported variants.
87///
88/// Currently read-only — write operations return errors.
89/// How a table's conditions are applied to a request.
90///
91/// URL `{placeholder}` path segments are always filled from matching
92/// eq-conditions regardless of strategy; this governs what happens to
93/// the *remaining* (non-path) eq-conditions.
94#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
95pub enum FilterStrategy {
96    /// Append remaining eq-conditions as `?field=value` query params
97    /// (JSON-Server semantics). The default.
98    #[default]
99    Query,
100    /// Apply remaining eq-conditions as in-memory row filters after the
101    /// fetch, never as query params. For APIs whose only server-side
102    /// filters are path segments and that reject (or ignore) unknown
103    /// query params — e.g. the Mercury control-API, whose CLI likewise
104    /// filters version/env client-side after fetching by product path.
105    Client,
106}
107
108/// How the API takes a sort: one query param naming the column, with a
109/// prefix that flips it to descending (`?ordering=-net`, the Django REST
110/// Framework and Launch Library convention). Configuring it makes the REST
111/// vista orderable, so a paged grid asks the server for its sort instead of
112/// re-sorting the rows it happens to hold.
113#[derive(Clone, Debug, PartialEq, Eq)]
114pub struct OrderingParams {
115    pub param: String,
116    pub desc_prefix: String,
117}
118
119impl OrderingParams {
120    pub fn new(param: impl Into<String>, desc_prefix: impl Into<String>) -> Self {
121        Self {
122            param: param.into(),
123            desc_prefix: desc_prefix.into(),
124        }
125    }
126
127    /// The value for `param`: the column, prefixed when descending.
128    pub fn value(&self, field: &str, dir: vantage_vista::SortDirection) -> String {
129        match dir {
130            vantage_vista::SortDirection::Ascending => field.to_string(),
131            vantage_vista::SortDirection::Descending => format!("{}{field}", self.desc_prefix),
132        }
133    }
134}
135
136/// A sort pushed down to the API: column and direction.
137pub(crate) type Order<'a> = Option<(&'a str, vantage_vista::SortDirection)>;
138
139/// One response's rows, see [`RestApi::fetch_page`].
140struct Page {
141    rows: IndexMap<String, Record<CborValue>>,
142    total: Option<i64>,
143    server_rows: usize,
144    lossy: bool,
145}
146
147#[derive(Clone, Debug)]
148pub struct RestApi {
149    base_url: String,
150    client: ResilientClient,
151    pub(crate) auth_header: AuthHeader,
152    response_shape: ResponseShape,
153    pagination: PaginationParams,
154    /// Paging params were configured explicitly, so windows can be fetched
155    /// even with no `total_key` — the end of the set shows as a short page.
156    paged: bool,
157    /// When true, no `_page`/`_limit` query params are appended and
158    /// list endpoints are assumed to return the full result set in
159    /// one shot. Caller-side requests for page > 1 short-circuit to
160    /// an empty result so a perpetual-grid stops paging after the
161    /// first chunk. Useful for FastAPI/Pydantic services that treat
162    /// unknown query params as strict filters.
163    no_pagination: bool,
164    /// How non-path eq-conditions are applied — query params vs.
165    /// in-memory post-fetch filtering. See [`FilterStrategy`].
166    filter_strategy: FilterStrategy,
167    /// Response-envelope key carrying the grand total of matching rows
168    /// (e.g. `count`). When set, the shell reports an exact count and
169    /// advertises `can_fetch_window` for lazy/scroll loading; when `None`
170    /// it falls back to counting fetched rows.
171    total_key: Option<String>,
172    /// The sort query param, when the API has one. Unset means the API
173    /// cannot sort and consumers order fetched rows themselves.
174    ordering: Option<OrderingParams>,
175    /// Emit `tracing` events for window/count requests.
176    debug: bool,
177}
178
179impl RestApi {
180    /// Create a new REST API pointing at `base_url`. Uses the legacy
181    /// default response shape (`{ "data": [...] }`). For other shapes
182    /// (bare array, wrapped-by-table-name) use [`RestApi::builder`].
183    pub fn new(base_url: impl Into<String>) -> Self {
184        RestApi::builder(base_url).build()
185    }
186
187    /// Start configuring a [`RestApi`] via the builder.
188    pub fn builder(base_url: impl Into<String>) -> RestApiBuilder {
189        RestApiBuilder::new(base_url.into())
190    }
191
192    /// Set the Authorization header value (e.g. "Bearer `<token>`").
193    /// Provided for backwards compatibility — prefer
194    /// `RestApi::builder(...).auth(...)`.
195    pub fn with_auth(mut self, auth: impl Into<String>) -> Self {
196        self.auth_header = AuthHeader::new(auth);
197        self
198    }
199
200    /// The configured response-envelope total key, if any. When set, the
201    /// REST shell can report an exact count and serve `fetch_window`.
202    pub fn total_key(&self) -> Option<&str> {
203        self.total_key.as_deref()
204    }
205
206    /// Whether the shell can serve absolute-offset windows: the API pages
207    /// (explicit `pagination_params`, or a `total_key`) and pagination isn't
208    /// switched off. Without a total the loader learns the end of the set
209    /// from a short page.
210    pub fn serves_windows(&self) -> bool {
211        !self.no_pagination && (self.paged || self.total_key.is_some())
212    }
213
214    /// The sort query param, when configured.
215    pub fn ordering(&self) -> Option<&OrderingParams> {
216        self.ordering.as_ref()
217    }
218
219    /// What the circuit breaker is doing right now. `None` when the
220    /// underlying client has no breaker configured.
221    pub fn breaker_state(&self) -> Option<crate::BreakerState> {
222        self.client.breaker_state()
223    }
224
225    /// The resilient client backing this API — the pool, breaker and
226    /// observer a caller outside the read path (e.g. an outbox replaying a
227    /// queued write) should share rather than build its own.
228    pub fn client(&self) -> &ResilientClient {
229        &self.client
230    }
231
232    /// The configured base URL requests are joined against.
233    pub fn base_url(&self) -> &str {
234        &self.base_url
235    }
236
237    /// Issue an arbitrary HTTP request against `base_url`/`path`, through
238    /// the same resilient client and auth header as reads.
239    ///
240    /// `headers` are applied after the configured `Authorization` header, so
241    /// a caller-supplied header of the same name wins. `body`, when given,
242    /// is sent as the JSON request body. The call policy follows
243    /// [`vantage_core::Priority::current`], same as a table read. On success,
244    /// a non-`GET` method reports [`TransportEvent::WritePushed`] to the
245    /// observer.
246    pub async fn http_request(
247        &self,
248        method: reqwest::Method,
249        path: &str,
250        headers: &[(&str, &str)],
251        body: Option<&serde_json::Value>,
252    ) -> vantage_core::Result<reqwest::Response> {
253        let url = join_base_path(&self.base_url, path);
254
255        // `HeaderMap::insert` replaces a same-named entry rather than
256        // appending — unlike `RequestBuilder::header` — so a caller header
257        // actually overrides the configured auth instead of riding alongside
258        // it on the wire.
259        let mut header_map = HeaderMap::new();
260        if let Some(auth) = self.auth_header.value() {
261            header_map.insert(
262                AUTHORIZATION,
263                HeaderValue::from_str(auth)
264                    .expect("configured auth header must be a valid header value"),
265            );
266        }
267        for (name, value) in headers {
268            header_map.insert(
269                HeaderName::from_bytes(name.as_bytes()).expect("header name must be valid"),
270                HeaderValue::from_str(value).expect("header value must be valid"),
271            );
272        }
273
274        let policy = crate::transport::policy_for(Priority::current());
275        let response = self
276            .client
277            .execute_with(&policy, |http| {
278                let req = http
279                    .request(method.clone(), &url)
280                    .headers(header_map.clone());
281                match body {
282                    Some(body) => req.json(body),
283                    None => req,
284                }
285            })
286            .await
287            .map_err(|e| crate::transport::client_error(e, "API request failed", &url))?;
288
289        if method != reqwest::Method::GET {
290            self.client.report(TransportEvent::WritePushed);
291        }
292
293        Ok(response)
294    }
295
296    /// Build the endpoint path for `table_name`, substituting any
297    /// `{placeholder}` segments from matching eq-conditions.
298    ///
299    /// Returns the absolute URL up to (but excluding) the query string,
300    /// alongside the indices of conditions consumed by the substitution
301    /// — those are dropped from the query string by `build_query_string`.
302    ///
303    /// Tables that don't use templates (no `{}` in the name) pass
304    /// through unchanged and consume no conditions.
305    fn endpoint_url(
306        &self,
307        table_name: &str,
308        conditions: &[&Expression<CborValue>],
309    ) -> Result<(String, Vec<usize>)> {
310        let mut consumed = Vec::new();
311        let mut path = String::with_capacity(table_name.len());
312        let mut rest = table_name;
313        while let Some(open) = rest.find('{') {
314            path.push_str(&rest[..open]);
315            let after = &rest[open + 1..];
316            let close = after.find('}').ok_or_else(|| {
317                error!(
318                    "Unclosed `{` in table name URI template",
319                    table_name = table_name
320                )
321            })?;
322            let placeholder = &after[..close];
323            let (idx, value) = conditions
324                .iter()
325                .enumerate()
326                .find_map(|(i, cond)| {
327                    if consumed.contains(&i) {
328                        return None;
329                    }
330                    let (field, value) = crate::condition_to_query_param(cond)?;
331                    (field == placeholder).then_some((i, value))
332                })
333                .ok_or_else(|| {
334                    error!(
335                        "No eq-condition provided for URI placeholder",
336                        placeholder = placeholder,
337                        table_name = table_name
338                    )
339                })?;
340            consumed.push(idx);
341            path.push_str(&urlencode(&value));
342            rest = &after[close + 1..];
343        }
344        path.push_str(rest);
345        Ok((format!("{}/{}", self.base_url, path), consumed))
346    }
347
348    /// Decide which conditions go in the query string and which are applied to
349    /// the rows after they arrive.
350    ///
351    /// Under [`FilterStrategy::Client`] non-path eq-conditions are *not* sent —
352    /// the API rejects or ignores unknown params — so they come back as
353    /// client-side filters and every condition is marked consumed to keep it out
354    /// of the URL. Otherwise nothing is filtered locally and only the path
355    /// placeholders are consumed.
356    ///
357    /// Shared by the real fetch and by [`preview_request`](Self::preview_request)
358    /// so a previewed URL cannot claim a filter the fetch would have applied in
359    /// memory, or vice versa.
360    fn split_filters(
361        &self,
362        conds: &[&Expression<CborValue>],
363        consumed: Vec<usize>,
364    ) -> (Vec<usize>, Vec<(String, String)>) {
365        if self.filter_strategy == FilterStrategy::Client {
366            let filters = conds
367                .iter()
368                .enumerate()
369                .filter(|(i, _)| !consumed.contains(i))
370                .filter_map(|(_, c)| crate::condition_to_query_param(c))
371                .collect();
372            ((0..conds.len()).collect(), filters)
373        } else {
374            (consumed, Vec::new())
375        }
376    }
377
378    /// Build the combined query-string from pagination + conditions.
379    /// `consumed` lists condition indices already baked into the URI
380    /// path; those don't appear in the query string. Conditions that
381    /// don't peel cleanly into eq pairs are skipped — same "best effort"
382    /// stance as before.
383    fn build_query_string(
384        &self,
385        window: Option<(i64, i64)>,
386        conditions: &[&Expression<CborValue>],
387        consumed: &[usize],
388        order: Order<'_>,
389    ) -> String {
390        let mut params: Vec<(String, String)> = Vec::new();
391
392        // Pagination first — matches the order users see in the URL bar.
393        // When `no_pagination` is set the API doesn't accept page/limit
394        // query params (and may treat them as strict filters that
395        // return empty), so we leave them off.
396        //
397        // `window` is a half-open `[offset, offset+limit)` band. Skip-based
398        // APIs take the offset verbatim; page-based APIs are addressed by
399        // 1-based page, derived from the offset (the loader may hand
400        // non-page-aligned windows, so it rounds down to the containing page).
401        if !self.no_pagination
402            && let Some((offset, limit)) = window
403        {
404            let offset = offset.max(0);
405            let limit = limit.max(1);
406            let page_value = if self.pagination.skip_based {
407                offset.to_string()
408            } else {
409                (offset / limit + 1).to_string()
410            };
411            params.push((self.pagination.page.clone(), page_value));
412            params.push((self.pagination.limit.clone(), limit.to_string()));
413        }
414
415        // A sort reaches the query only when the API declared how it takes
416        // one; the vista never offers `add_order` otherwise.
417        if let (Some(spec), Some((field, dir))) = (&self.ordering, order) {
418            params.push((spec.param.clone(), spec.value(field, dir)));
419        }
420
421        // Conditions: each `eq` becomes `?field=value`. Multiple
422        // conditions AND together (JSON Server semantics).
423        for (i, cond) in conditions.iter().enumerate() {
424            if consumed.contains(&i) {
425                continue;
426            }
427            if let Some((field, value)) = crate::condition_to_query_param(cond) {
428                params.push((field, value));
429            }
430        }
431
432        if params.is_empty() {
433            return String::new();
434        }
435        let mut s = String::from("?");
436        for (i, (k, v)) in params.iter().enumerate() {
437            if i > 0 {
438                s.push('&');
439            }
440            // Minimal URL encoding — we encode `&` and `=` and spaces
441            // because those break the query format. Anything else
442            // passes through; the JSON Server convention is permissive.
443            s.push_str(&urlencode(k));
444            s.push('=');
445            s.push_str(&urlencode(v));
446        }
447        s
448    }
449
450    /// Render the request a read would issue, without issuing it.
451    ///
452    /// Shares [`endpoint_url`](Self::endpoint_url) and
453    /// [`build_query_string`](Self::build_query_string) with the real fetch
454    /// path, so a previewed URL cannot drift from the one that gets sent.
455    ///
456    /// One difference is deliberate: `fetch_raw_body` first *awaits* any
457    /// deferred condition (a foreign key whose value arrives with the parent
458    /// row), and awaiting is what a preview must not do. Those are counted
459    /// under `deferred_conditions` and left out of the URL instead.
460    pub(crate) fn preview_request<'a>(
461        &self,
462        table_name: &str,
463        window: Option<(i64, i64)>,
464        order: Order<'_>,
465        conditions: impl IntoIterator<Item = &'a Expression<CborValue>>,
466    ) -> serde_json::Value {
467        let conds: Vec<&Expression<CborValue>> = conditions.into_iter().collect();
468
469        // Conditions the query-param lowering cannot peel into an eq pair —
470        // deferred foreign keys, and anything else not shaped `field = value`.
471        let unresolved = conds
472            .iter()
473            .filter(|c| crate::condition_to_query_param(c).is_none())
474            .count();
475
476        let (endpoint, consumed) = match self.endpoint_url(table_name, &conds) {
477            Ok(pair) => pair,
478            // A URI template placeholder went unfilled. If some condition is
479            // still unresolved, the real fetch would have awaited it *before*
480            // building the path, so this is a preview limitation and the
481            // template is the honest answer. With nothing outstanding, the
482            // fetch would fail here too — report that.
483            Err(_) if unresolved > 0 => {
484                return serde_json::json!({
485                    "driver": "rest-api",
486                    "method": "GET",
487                    "url": format!("{}/{}", self.base_url, table_name),
488                    "unresolved_conditions": unresolved,
489                    "note": "path placeholders are filled from conditions resolved \
490                             at fetch time; the template is shown unfilled",
491                });
492            }
493            Err(e) => {
494                return serde_json::json!({
495                    "driver": "rest-api",
496                    "base_url": self.base_url,
497                    "error": e.to_string(),
498                });
499            }
500        };
501
502        let (query_consumed, client_filters) = self.split_filters(&conds, consumed);
503        let query = self.build_query_string(window, &conds, &query_consumed, order);
504
505        serde_json::json!({
506            "driver": "rest-api",
507            "method": "GET",
508            "url": join_query(&endpoint, &query),
509            "auth_header": self.auth_header.masked(),
510            // Under `FilterStrategy::Client` these never reach the server: the
511            // rows come back unfiltered and are narrowed in memory.
512            "client_side_filters": client_filters
513                .into_iter()
514                .map(|(k, v)| format!("{k}={v}"))
515                .collect::<Vec<_>>(),
516            "unresolved_conditions": unresolved,
517        })
518    }
519
520    /// Fetch data from the API endpoint and return parsed records.
521    ///
522    /// `id_field` selects which JSON field is treated as the record ID;
523    /// if `None`, row indices are used. The page-based `pagination` is
524    /// lowered to a `[offset, offset+limit)` window; `conditions` are
525    /// pushed into the URL query string — eq-conditions become
526    /// `?field=value`. Conditions that can't be peeled into a simple
527    /// eq are silently skipped (caller-side filtering still applies if
528    /// needed).
529    pub(crate) async fn fetch_records<'a>(
530        &self,
531        table_name: &str,
532        id_field: Option<&str>,
533        pagination: Option<&Pagination>,
534        conditions: impl IntoIterator<Item = &'a Expression<CborValue>>,
535    ) -> Result<IndexMap<String, Record<CborValue>>> {
536        let window = pagination.map(|p| (p.skip(), p.limit()));
537        self.fetch_windowed(table_name, id_field, window, None, conditions)
538            .await
539            .map(|(records, _total)| records)
540    }
541
542    /// Fetch a single half-open row window `[offset, offset+limit)` — the
543    /// primitive a paged, lazily-loaded grid drives on scroll (offset is
544    /// an absolute row index, not a page number).
545    /// Fetch one half-open row window, plus the envelope's `total_key` when
546    /// the response carries one.
547    ///
548    /// The total comes out of the **same response as the rows**. Every paged
549    /// endpoint reports it on every reply, so a caller wanting both a window
550    /// and a grand total takes them together here rather than pairing a
551    /// window fetch with [`Self::fetch_total`] — that pairing costs a second
552    /// round trip for a number already in hand.
553    pub(crate) async fn fetch_window_records_counted<'a>(
554        &self,
555        table_name: &str,
556        id_field: Option<&str>,
557        offset: i64,
558        limit: i64,
559        order: Order<'_>,
560        conditions: impl IntoIterator<Item = &'a Expression<CborValue>>,
561    ) -> Result<(IndexMap<String, Record<CborValue>>, Option<i64>)> {
562        self.fetch_windowed(
563            table_name,
564            id_field,
565            Some((offset, limit)),
566            order,
567            conditions,
568        )
569        .await
570    }
571
572    /// Read the grand total of matching rows from the response envelope's
573    /// configured `total_key` (e.g. `count`). Returns `None` when no
574    /// `total_key` is set — the caller then falls back to counting fetched
575    /// rows. Issues a cheap `limit=1` request so the body carries the count
576    /// without paying for the rows.
577    pub(crate) async fn fetch_total<'a>(
578        &self,
579        table_name: &str,
580        conditions: impl IntoIterator<Item = &'a Expression<CborValue>>,
581    ) -> Result<Option<i64>> {
582        let Some(total_key) = self.total_key.clone() else {
583            return Ok(None);
584        };
585        let (body, _client_filters) = self
586            .fetch_raw_body(table_name, Some((0, 1)), None, conditions)
587            .await?;
588        let total = body
589            .get(total_key.as_str())
590            .and_then(|v| v.as_i64())
591            .ok_or_else(|| {
592                error!(
593                    "total_key missing or not an integer in API response",
594                    total_key = total_key.as_str()
595                )
596            })?;
597        if self.debug {
598            tracing::debug!(target: "vantage_api_client::rest", total, "REST count");
599        }
600        Ok(Some(total))
601    }
602
603    /// Resolve conditions, build the windowed request URL, GET it (with the
604    /// auth header if configured), and return the parsed JSON body together
605    /// with any client-side filters that still need applying (under
606    /// [`FilterStrategy::Client`]).
607    async fn fetch_raw_body<'a>(
608        &self,
609        table_name: &str,
610        window: Option<(i64, i64)>,
611        order: Order<'_>,
612        conditions: impl IntoIterator<Item = &'a Expression<CborValue>>,
613    ) -> Result<(serde_json::Value, Vec<(String, String)>)> {
614        // Conditions may carry `DeferredFn` values — typically from
615        // `related_in_condition` for `with_one`-style traversals where the FK
616        // lives in a parent record we haven't fetched yet. Resolve them once,
617        // up front, so the rest of the pipeline sees only sync scalars.
618        let raw: Vec<&Expression<CborValue>> = conditions.into_iter().collect();
619        let mut resolved: Vec<Expression<CborValue>> = Vec::with_capacity(raw.len());
620        for cond in raw {
621            resolved.push(cond.resolve_deferred().await?);
622        }
623        let conds: Vec<&Expression<CborValue>> = resolved.iter().collect();
624        let (endpoint, consumed) = self.endpoint_url(table_name, &conds)?;
625
626        let (query_consumed, client_filters) = self.split_filters(&conds, consumed);
627        let query = self.build_query_string(window, &conds, &query_consumed, order);
628        let url = join_query(&endpoint, &query);
629
630        // The `(0, 1)` window is the count probe (reads only the envelope's
631        // total); log it at debug so it doesn't drown the real data fetches.
632        if self.debug {
633            if window == Some((0, 1)) {
634                tracing::debug!(target: "vantage_api_client::rest", table = table_name, url = %url, "REST GET (count probe)");
635            } else {
636                tracing::info!(target: "vantage_api_client::rest", table = table_name, url = %url, "REST GET");
637            }
638        }
639
640        // Time every round trip, unconditionally — a remote API is the one part
641        // of a read the process cannot bound, and a slow page is far more often
642        // one slow GET than anything local. Reported regardless of `debug` so
643        // the cost is attributable from a default log; a request over a second
644        // is worth an operator's attention, hence `info` at that point.
645        let started = std::time::Instant::now();
646        let policy = crate::transport::policy_for(Priority::current());
647        let response = self
648            .client
649            .execute_with(&policy, |http| {
650                let req = http.get(&url);
651                match self.auth_header.value() {
652                    Some(auth) => req.header(AUTHORIZATION, auth),
653                    None => req,
654                }
655            })
656            .await
657            .map_err(|e| {
658                let ms = started.elapsed().as_millis() as u64;
659                // A transient failure (5xx, network, open breaker) is visible
660                // in the datasource's health and will be tried again by a
661                // later call; only an answer no retry can change is worth a
662                // warning.
663                if e.is_final() {
664                    tracing::warn!(
665                        target: "vantage_api_client::rest",
666                        table = table_name,
667                        url = %url,
668                        ms,
669                        attempts = e.attempts,
670                        "REST GET failed",
671                    );
672                } else {
673                    tracing::debug!(
674                        target: "vantage_api_client::rest",
675                        table = table_name,
676                        url = %url,
677                        ms,
678                        attempts = e.attempts,
679                        kind = e.kind_name(),
680                        "REST GET gave up on a transient failure",
681                    );
682                }
683                crate::transport::client_error(e, "API request failed", &url)
684            })?;
685
686        let body: serde_json::Value = response
687            .json()
688            .await
689            .map_err(|e| error!("Failed to parse API response as JSON", detail = e))?;
690
691        let ms = started.elapsed().as_millis() as u64;
692        let probe = window == Some((0, 1));
693        if ms >= 1000 {
694            tracing::info!(
695                target: "vantage_api_client::rest",
696                table = table_name,
697                url = %url,
698                ms,
699                count_probe = probe,
700                "slow REST GET",
701            );
702        } else {
703            tracing::debug!(
704                target: "vantage_api_client::rest",
705                table = table_name,
706                url = %url,
707                ms,
708                count_probe = probe,
709                "REST GET done",
710            );
711        }
712
713        Ok((body, client_filters))
714    }
715
716    /// Also reports the envelope total when `total_key` is configured and the
717    /// body carries it. Unlike [`Self::fetch_total`] this never errors on a
718    /// missing total: the rows are the point here, and a caller that needs a
719    /// definitive count can still ask for one.
720    ///
721    /// A window comes back short only when the server ran out: callers take a
722    /// short window as the end of the set.
723    async fn fetch_windowed<'a>(
724        &self,
725        table_name: &str,
726        id_field: Option<&str>,
727        window: Option<(i64, i64)>,
728        order: Order<'_>,
729        conditions: impl IntoIterator<Item = &'a Expression<CborValue>>,
730    ) -> Result<(IndexMap<String, Record<CborValue>>, Option<i64>)> {
731        // Non-paginating endpoints return the whole list on the first
732        // window; a later window would just re-deliver the same rows and the
733        // perpetual grid would never mark itself exhausted. Short-circuit any
734        // window past the start to empty so the grid sees the chunk shrink
735        // and stops asking for more.
736        if self.no_pagination && window.is_some_and(|(offset, _)| offset > 0) {
737            return Ok((IndexMap::new(), None));
738        }
739        let conditions: Vec<&Expression<CborValue>> = conditions.into_iter().collect();
740        let Some((offset, limit)) = window.filter(|_| !self.no_pagination) else {
741            let page = self
742                .fetch_page(table_name, id_field, window, order, &conditions)
743                .await?;
744            return Ok((page.rows, page.total));
745        };
746        let (offset, limit) = (offset.max(0), limit.max(1));
747
748        // Page-based APIs are addressed by whole pages: start at the page
749        // holding `offset`.
750        let start = if self.pagination.skip_based {
751            offset
752        } else {
753            offset - offset % limit
754        };
755        let first = self
756            .fetch_page(
757                table_name,
758                id_field,
759                Some((start, limit)),
760                order,
761                &conditions,
762            )
763            .await?;
764        if start == offset && !first.lossy {
765            return Ok((first.rows, first.total));
766        }
767
768        // Assemble the window page by page until it is full or the server
769        // runs out. When rows were dropped (client filters, repeated ids) a
770        // row's place in the result is only known by walking from the start.
771        let (mut server_offset, mut skip) = if first.lossy {
772            (0, offset)
773        } else {
774            (start, offset - start)
775        };
776        let mut page = (server_offset == start).then_some(first);
777        let mut rows = IndexMap::new();
778        let mut total = None;
779        let mut lossy = false;
780        loop {
781            let current = match page.take() {
782                Some(p) => p,
783                None => {
784                    self.fetch_page(
785                        table_name,
786                        id_field,
787                        Some((server_offset, limit)),
788                        order,
789                        &conditions,
790                    )
791                    .await?
792                }
793            };
794            total = total.or(current.total);
795            lossy |= current.lossy;
796            for (id, row) in current.rows {
797                if skip > 0 {
798                    skip -= 1;
799                } else if (rows.len() as i64) < limit {
800                    rows.insert(id, row);
801                }
802            }
803            if rows.len() as i64 >= limit || (current.server_rows as i64) < limit {
804                break;
805            }
806            server_offset += limit;
807        }
808        Ok((rows, if lossy { None } else { total }))
809    }
810
811    /// One request: the page's rows keyed by id, the envelope total, and how
812    /// many rows the server sent — more than `rows` holds when client filters
813    /// or repeated ids dropped some (`lossy`).
814    async fn fetch_page(
815        &self,
816        table_name: &str,
817        id_field: Option<&str>,
818        window: Option<(i64, i64)>,
819        order: Order<'_>,
820        conditions: &[&Expression<CborValue>],
821    ) -> Result<Page> {
822        let (body, client_filters) = self
823            .fetch_raw_body(table_name, window, order, conditions.iter().copied())
824            .await?;
825        let total = self
826            .total_key
827            .as_deref()
828            .and_then(|key| body.get(key))
829            .and_then(|v| v.as_i64());
830        let data = self.extract_array(&body, table_name)?;
831
832        let mut records = IndexMap::new();
833        for (row_idx, item) in data.iter().enumerate() {
834            let obj = item
835                .as_object()
836                .ok_or_else(|| error!("API data item is not an object", index = row_idx))?;
837
838            // Extract ID from the configured id_field, or use row index
839            let id = id_field
840                .and_then(|field| obj.get(field))
841                .and_then(|v| match v {
842                    serde_json::Value::String(s) => Some(s.clone()),
843                    serde_json::Value::Number(n) => Some(n.to_string()),
844                    _ => None,
845                })
846                .unwrap_or_else(|| row_idx.to_string());
847
848            // The HTTP body parses as JSON for free; convert to CBOR
849            // at this single boundary so the rest of the pipeline
850            // (Table, Vista) sees the universal carrier. json_to_cbor
851            // is total — JSON is a strict subset of CBOR.
852            let mut record: Record<CborValue> = Record::new();
853            for (k, v) in obj {
854                record.insert(k.clone(), vantage_types::json_to_cbor(v.clone()));
855            }
856
857            records.insert(id, record);
858        }
859
860        // Client-side filtering (FilterStrategy::Client): drop rows that
861        // don't match the non-path eq-conditions. A condition whose field
862        // is absent from a row is treated as a pass (it was a path/request
863        // param, not a record field) — mirroring the AWS connector and the
864        // Mercury CLI's own post-fetch `_filter_deployments`.
865        if !client_filters.is_empty() {
866            records.retain(|_id, record| {
867                client_filters
868                    .iter()
869                    .all(|(field, want)| match record.get(field) {
870                        Some(v) => crate::cbor_to_query_string(v).as_deref() == Some(want.as_str()),
871                        None => true,
872                    })
873            });
874        }
875
876        // Counts the rows this call hands back, so a client-side filter shows
877        // up as fewer rows pulled than the server sent.
878        self.client
879            .report(TransportEvent::RowsPulled { n: records.len() });
880
881        // The envelope counted what the SERVER matched, before any rows were
882        // dropped above — reporting that total now would size a grid to rows it
883        // will never be given. No total is better than a wrong one.
884        let total = if client_filters.is_empty() {
885            total
886        } else {
887            None
888        };
889        Ok(Page {
890            lossy: records.len() < data.len(),
891            server_rows: data.len(),
892            rows: records,
893            total,
894        })
895    }
896}
897
898fn urlencode(s: &str) -> String {
899    urlencoding::encode(s).into_owned()
900}
901
902/// Join a base URL and a path with exactly one slash between them,
903/// regardless of whether either side already carries one.
904fn join_base_path(base: &str, path: &str) -> String {
905    format!(
906        "{}/{}",
907        base.trim_end_matches('/'),
908        path.trim_start_matches('/')
909    )
910}
911
912/// Append a `build_query_string` result (always opening with `?`, or empty)
913/// to an endpoint URL. The table path may itself carry a query string (e.g.
914/// `launches/?mode=detailed`), in which case the appended params must join
915/// with `&` — otherwise the URL gets two `?` and the API rejects it.
916fn join_query(endpoint: &str, query: &str) -> String {
917    match query.strip_prefix('?') {
918        Some(rest) if endpoint.contains('?') => format!("{endpoint}&{rest}"),
919        _ => format!("{endpoint}{query}"),
920    }
921}
922
923impl RestApi {
924    /// Pull the row array out of the response body, according to the
925    /// configured `ResponseShape`. A bare-array API answers a single-record
926    /// endpoint (`activities/{id}`) with one object: that reads as one row.
927    fn extract_array<'a>(
928        &self,
929        body: &'a serde_json::Value,
930        table_name: &str,
931    ) -> Result<&'a [serde_json::Value]> {
932        match &self.response_shape {
933            ResponseShape::BareArray => match body {
934                serde_json::Value::Array(rows) => Ok(rows),
935                serde_json::Value::Object(_) => Ok(std::slice::from_ref(body)),
936                _ => Err(error!(
937                    "Expected response body to be a JSON array or object (BareArray shape)"
938                )),
939            },
940            ResponseShape::Wrapped { array_key } => body[array_key]
941                .as_array()
942                .map(Vec::as_slice)
943                .ok_or_else(|| {
944                    error!(
945                        "Response missing array under wrapper key",
946                        array_key = array_key
947                    )
948                }),
949            ResponseShape::WrappedByTableName => body[table_name]
950                .as_array()
951                .map(Vec::as_slice)
952                .ok_or_else(|| {
953                    error!(
954                        "Response missing array under table-name key",
955                        table_name = table_name
956                    )
957                }),
958        }
959    }
960}
961
962/// Builder for [`RestApi`]. Lets callers pick a [`ResponseShape`] and
963/// override the pagination parameter names.
964///
965/// ```no_run
966/// use vantage_api_client::{RestApi, ResponseShape, PaginationParams};
967///
968/// // JSONPlaceholder: bare arrays, JSON-Server pagination conventions.
969/// let api = RestApi::builder("https://jsonplaceholder.typicode.com")
970///     .response_shape(ResponseShape::BareArray)
971///     .build();
972///
973/// // DummyJSON: wrapped-by-table-name, skip-based pagination.
974/// let api = RestApi::builder("https://dummyjson.com")
975///     .response_shape(ResponseShape::WrappedByTableName)
976///     .pagination_params(PaginationParams::skip_limit("skip", "limit"))
977///     .build();
978/// ```
979#[derive(Clone, Debug)]
980pub struct RestApiBuilder {
981    base_url: String,
982    auth_header: AuthHeader,
983    response_shape: ResponseShape,
984    pagination: PaginationParams,
985    /// `pagination_params` was called: the API is known to page.
986    paged: bool,
987    no_pagination: bool,
988    filter_strategy: FilterStrategy,
989    total_key: Option<String>,
990    ordering: Option<OrderingParams>,
991    debug: bool,
992    transport: crate::transport::ClientConfig,
993}
994
995impl RestApiBuilder {
996    fn new(base_url: String) -> Self {
997        Self {
998            base_url,
999            auth_header: AuthHeader::default(),
1000            response_shape: ResponseShape::default(),
1001            pagination: PaginationParams::default(),
1002            paged: false,
1003            no_pagination: false,
1004            filter_strategy: FilterStrategy::default(),
1005            total_key: None,
1006            ordering: None,
1007            debug: false,
1008            transport: crate::transport::ClientConfig::default(),
1009        }
1010    }
1011
1012    /// Concurrent requests to this API at most (default 4).
1013    pub fn max_parallel(mut self, n: usize) -> Self {
1014        self.transport.max_parallel = n.max(1);
1015        self
1016    }
1017
1018    /// Requests per second to this API at most.
1019    pub fn rate_limit(mut self, per_second: f64) -> Self {
1020        self.transport.rate_limit = Some(per_second);
1021        self
1022    }
1023
1024    /// Report every attempt, retry and breaker transition under `key`
1025    /// (the datasource name).
1026    pub fn observer(
1027        mut self,
1028        key: impl Into<Arc<str>>,
1029        observer: Arc<dyn TransportObserver>,
1030    ) -> Self {
1031        self.transport.observer = Some((key.into(), observer));
1032        self
1033    }
1034
1035    /// A pre-configured `reqwest::Client` (timeouts, proxies, TLS).
1036    pub fn http_client(mut self, client: reqwest::Client) -> Self {
1037        self.transport.http = Some(client);
1038        self
1039    }
1040
1041    /// Set the Authorization header value (e.g. "Bearer `<token>`").
1042    pub fn auth(mut self, auth: impl Into<String>) -> Self {
1043        self.auth_header = AuthHeader::new(auth);
1044        self.transport.auth_refresher = None;
1045        self
1046    }
1047
1048    /// Get the bearer token from `refresher` instead of a fixed header: it
1049    /// is asked on the first request, and again when the API answers `401`.
1050    /// A request waits for it, so a refresher may run an interactive
1051    /// sign-in. Replaces [`auth`](Self::auth).
1052    pub fn auth_refresher(mut self, refresher: crate::AuthRefresher) -> Self {
1053        self.auth_header = AuthHeader::default();
1054        self.transport.auth_refresher = Some(refresher);
1055        self
1056    }
1057
1058    /// Choose how the API wraps its row array. Defaults to
1059    /// `Wrapped { array_key: "data" }` for backwards compat.
1060    pub fn response_shape(mut self, shape: ResponseShape) -> Self {
1061        self.response_shape = shape;
1062        self
1063    }
1064
1065    /// Override the page/limit query parameter names. Default is
1066    /// `("_page", "_limit")` (JSON Server convention).
1067    pub fn pagination_params(mut self, pagination: PaginationParams) -> Self {
1068        self.pagination = pagination;
1069        self.paged = true;
1070        self
1071    }
1072
1073    /// Disable pagination entirely — no `_page`/`_limit` query
1074    /// params are appended, and a request for page > 1 is short-
1075    /// circuited to an empty result. Use this for APIs that don't
1076    /// paginate (return the full list every call) or that treat
1077    /// unknown query params as strict filters.
1078    pub fn no_pagination(mut self) -> Self {
1079        self.no_pagination = true;
1080        self
1081    }
1082
1083    /// Choose how non-path eq-conditions are applied. Default is
1084    /// [`FilterStrategy::Query`]; use [`FilterStrategy::Client`] for
1085    /// APIs that only filter via path segments and reject/ignore unknown
1086    /// query params (the conditions are then applied in memory).
1087    pub fn filter_strategy(mut self, strategy: FilterStrategy) -> Self {
1088        self.filter_strategy = strategy;
1089        self
1090    }
1091
1092    /// Name the response-envelope key carrying the grand total of matching
1093    /// rows (e.g. `count`). Setting it lets the shell report an exact count
1094    /// and advertise `can_fetch_window` for lazy/scroll loading.
1095    pub fn total_key(mut self, key: impl Into<String>) -> Self {
1096        self.total_key = Some(key.into());
1097        self
1098    }
1099
1100    /// The API's sort query param (see [`OrderingParams`]). With it set the
1101    /// vista reports `can_order` and every column is orderable.
1102    pub fn ordering(mut self, ordering: OrderingParams) -> Self {
1103        self.ordering = Some(ordering);
1104        self
1105    }
1106
1107    /// Emit `tracing` events for window/count requests.
1108    pub fn debug(mut self, debug: bool) -> Self {
1109        self.debug = debug;
1110        self
1111    }
1112
1113    pub fn build(self) -> RestApi {
1114        RestApi {
1115            base_url: self.base_url,
1116            client: crate::transport::build_client(self.transport),
1117            auth_header: self.auth_header,
1118            response_shape: self.response_shape,
1119            pagination: self.pagination,
1120            paged: self.paged,
1121            no_pagination: self.no_pagination,
1122            filter_strategy: self.filter_strategy,
1123            total_key: self.total_key,
1124            ordering: self.ordering,
1125            debug: self.debug,
1126        }
1127    }
1128}
1129
1130#[cfg(test)]
1131mod tests {
1132    use super::*;
1133
1134    /// `build_query_string` with no conditions, exercising only the
1135    /// window → pagination-param mapping.
1136    fn qs(api: &RestApi, window: Option<(i64, i64)>) -> String {
1137        api.build_query_string(window, &[], &[], None)
1138    }
1139
1140    #[test]
1141    fn debug_masks_auth_header() {
1142        let api = RestApi::builder("http://x")
1143            .auth("Bearer secret-token")
1144            .build();
1145        let text = format!("{api:?}");
1146        assert!(!text.contains("secret-token"), "{text}");
1147        assert!(text.contains("<set>"), "{text}");
1148    }
1149
1150    #[test]
1151    fn skip_based_window_uses_offset_verbatim() {
1152        let api = RestApi::builder("http://x")
1153            .pagination_params(PaginationParams::skip_limit("skip", "limit"))
1154            .build();
1155        assert_eq!(qs(&api, Some((20, 10))), "?skip=20&limit=10");
1156    }
1157
1158    #[test]
1159    fn page_based_window_derives_one_based_page() {
1160        let api = RestApi::builder("http://x").build(); // default _page/_limit
1161        // offset 20 / limit 10 → page 3 (1-based).
1162        assert_eq!(qs(&api, Some((20, 10))), "?_page=3&_limit=10");
1163    }
1164
1165    #[test]
1166    fn no_window_emits_no_pagination_params() {
1167        let api = RestApi::builder("http://x").build();
1168        assert_eq!(qs(&api, None), "");
1169    }
1170
1171    #[test]
1172    fn no_pagination_suppresses_window_params() {
1173        let api = RestApi::builder("http://x").no_pagination().build();
1174        assert_eq!(qs(&api, Some((20, 10))), "");
1175    }
1176
1177    #[test]
1178    fn query_string_joins_plain_endpoint_with_question_mark() {
1179        assert_eq!(
1180            join_query("http://x/launches/", "?_page=1&_limit=10"),
1181            "http://x/launches/?_page=1&_limit=10"
1182        );
1183    }
1184
1185    #[test]
1186    fn query_string_joins_templated_endpoint_with_ampersand() {
1187        // Endpoint already carries `?mode=detailed`; pagination must append
1188        // with `&`, not a second `?`.
1189        assert_eq!(
1190            join_query("http://x/launches/?mode=detailed", "?offset=0&limit=1"),
1191            "http://x/launches/?mode=detailed&offset=0&limit=1"
1192        );
1193    }
1194
1195    #[test]
1196    fn empty_query_string_leaves_endpoint_untouched() {
1197        assert_eq!(
1198            join_query("http://x/launches/?mode=detailed", ""),
1199            "http://x/launches/?mode=detailed"
1200        );
1201    }
1202
1203    /// Live regression for the double-`?` bug: a real fetch against the
1204    /// Launch Library 2 dev API using a table path that already carries a
1205    /// query string (`launches/?mode=detailed`). Before the `join_query`
1206    /// fix the request URL was `…/launches/?mode=detailed?offset=0&limit=1`
1207    /// and the server answered 500. Network-gated, so `#[ignore]`d:
1208    /// `cargo test -p vantage-api-client -- --ignored query_string`.
1209    #[tokio::test]
1210    #[ignore = "hits the live Launch Library 2 dev API"]
1211    async fn live_templated_table_path_fetches_rows() {
1212        let api = RestApi::builder("https://lldev.thespacedevs.com/2.3.0")
1213            .pagination_params(PaginationParams::skip_limit("offset", "limit"))
1214            .response_shape(ResponseShape::Wrapped {
1215                array_key: "results".into(),
1216            })
1217            .total_key("count")
1218            .build();
1219
1220        let total = api
1221            .fetch_total("launches/?mode=detailed", [])
1222            .await
1223            .expect("fetch_total");
1224        assert!(total.is_some_and(|n| n > 0), "expected a positive count");
1225
1226        let (rows, window_total) = api
1227            .fetch_window_records_counted("launches/?mode=detailed", Some("id"), 0, 3, None, [])
1228            .await
1229            .expect("fetch_window_records_counted");
1230        assert_eq!(rows.len(), 3, "expected the requested 3-row window");
1231        // The whole point of the counted window: the same response that
1232        // carried the rows also carried the count, so `fetch_total`'s extra
1233        // round trip buys nothing a caller couldn't already have.
1234        assert_eq!(
1235            window_total, total,
1236            "the window's envelope total should match the dedicated count",
1237        );
1238    }
1239}