Skip to main content

stmo_cli/
api.rs

1#![allow(clippy::missing_errors_doc)]
2
3use crate::models::{
4    CreateDashboard, CreateQuery, CreateWidget, Dashboard, DashboardSummary, DashboardsResponse,
5    DataSource, DataSourceSchema, QueriesResponse, Query,
6};
7use anyhow::{Context, Result};
8use reqwest::{Client, header};
9use serde::Serialize;
10use serde::de::DeserializeOwned;
11
12pub struct RedashClient {
13    client: Client,
14    base_url: String,
15}
16
17impl RedashClient {
18    pub fn new(base_url: String, api_key: &str) -> Result<Self> {
19        let mut headers = header::HeaderMap::new();
20        headers.insert(
21            "Authorization",
22            header::HeaderValue::from_str(&format!("Key {api_key}"))
23                .context("Invalid API key format")?,
24        );
25
26        let client = Client::builder()
27            .default_headers(headers)
28            .build()
29            .context("Failed to build HTTP client")?;
30
31        Ok(Self { client, base_url })
32    }
33
34    #[must_use]
35    pub fn base_url(&self) -> &str {
36        &self.base_url
37    }
38
39    async fn get_json<T: DeserializeOwned>(&self, url: &str, ctx: &str) -> Result<T> {
40        let response = self
41            .client
42            .get(url)
43            .send()
44            .await
45            .with_context(|| format!("Failed to request {ctx}"))?;
46
47        let response = ensure_success(response).await?;
48
49        response
50            .json()
51            .await
52            .with_context(|| format!("Failed to parse {ctx} response"))
53    }
54
55    async fn post_json<T: DeserializeOwned, B: Serialize + ?Sized>(
56        &self,
57        url: &str,
58        body: &B,
59        ctx: &str,
60    ) -> Result<T> {
61        let response = self
62            .client
63            .post(url)
64            .json(body)
65            .send()
66            .await
67            .with_context(|| format!("Failed to request {ctx}"))?;
68
69        let response = ensure_success(response).await?;
70
71        response
72            .json()
73            .await
74            .with_context(|| format!("Failed to parse {ctx} response"))
75    }
76
77    pub async fn list_my_queries(&self, page: u32, page_size: u32) -> Result<QueriesResponse> {
78        let url = format!(
79            "{}/api/queries/my?page={page}&page_size={page_size}",
80            self.base_url
81        );
82        self.get_json(&url, "my queries").await
83    }
84
85    pub async fn get_query(&self, query_id: u64) -> Result<Query> {
86        let url = format!("{}/api/queries/{query_id}", self.base_url);
87        self.get_json(&url, &format!("query {query_id}")).await
88    }
89
90    pub async fn list_data_sources(&self) -> Result<Vec<DataSource>> {
91        let url = format!("{}/api/data_sources", self.base_url);
92        self.get_json(&url, "data sources").await
93    }
94
95    pub async fn get_data_source(&self, data_source_id: u64) -> Result<DataSource> {
96        let url = format!("{}/api/data_sources/{data_source_id}", self.base_url);
97        self.get_json(&url, &format!("data source {data_source_id}"))
98            .await
99    }
100
101    pub async fn get_data_source_schema(
102        &self,
103        data_source_id: u64,
104        refresh: bool,
105    ) -> Result<DataSourceSchema> {
106        let url = if refresh {
107            format!(
108                "{}/api/data_sources/{data_source_id}/schema?refresh=true",
109                self.base_url
110            )
111        } else {
112            format!("{}/api/data_sources/{data_source_id}/schema", self.base_url)
113        };
114
115        self.get_json(&url, &format!("schema for data source {data_source_id}"))
116            .await
117    }
118
119    pub async fn create_query(&self, create_query: &CreateQuery) -> Result<Query> {
120        let url = format!("{}/api/queries", self.base_url);
121        self.post_json(&url, create_query, "new query").await
122    }
123
124    pub async fn create_or_update_query(&self, query: &Query) -> Result<Query> {
125        let url = format!("{}/api/queries/{}", self.base_url, query.id);
126        self.post_json(&url, query, &format!("query {} update", query.id))
127            .await
128    }
129
130    pub async fn create_visualization(
131        &self,
132        query_id: u64,
133        viz: &crate::models::CreateVisualization,
134    ) -> Result<crate::models::Visualization> {
135        let url = format!("{}/api/visualizations", self.base_url);
136        self.post_json(&url, viz, &format!("visualization for query {query_id}"))
137            .await
138    }
139
140    pub async fn update_visualization(
141        &self,
142        viz: &crate::models::Visualization,
143    ) -> Result<crate::models::Visualization> {
144        let url = format!("{}/api/visualizations/{}", self.base_url, viz.id);
145        self.post_json(&url, viz, &format!("visualization {} update", viz.id))
146            .await
147    }
148
149    pub async fn fetch_all_queries(&self) -> Result<Vec<Query>> {
150        let mut all_queries = Vec::new();
151        let mut page = 1;
152        let page_size = 100;
153
154        loop {
155            let response = self.list_my_queries(page, page_size).await?;
156
157            if response.results.is_empty() {
158                break;
159            }
160
161            all_queries.extend(response.results);
162            eprintln!(
163                "Fetched {} / {} queries...",
164                all_queries.len(),
165                response.count
166            );
167
168            #[allow(clippy::cast_possible_truncation)]
169            if all_queries.len() >= response.count as usize {
170                break;
171            }
172
173            page += 1;
174        }
175
176        Ok(all_queries)
177    }
178
179    pub async fn refresh_query(
180        &self,
181        query_id: u64,
182        parameters: Option<std::collections::HashMap<String, serde_json::Value>>,
183    ) -> Result<crate::models::Job> {
184        let url = format!("{}/api/queries/{query_id}/results", self.base_url);
185
186        let request = crate::models::RefreshRequest {
187            max_age: 0,
188            parameters,
189        };
190
191        let job_response: crate::models::JobResponse = self
192            .post_json(&url, &request, &format!("query {query_id} refresh"))
193            .await?;
194
195        Ok(job_response.job)
196    }
197
198    pub async fn poll_job(&self, job_id: &str) -> Result<crate::models::Job> {
199        let url = format!("{}/api/jobs/{job_id}", self.base_url);
200
201        let job_response: crate::models::JobResponse =
202            self.get_json(&url, &format!("job {job_id}")).await?;
203
204        Ok(job_response.job)
205    }
206
207    pub async fn get_query_result(
208        &self,
209        query_id: u64,
210        result_id: u64,
211    ) -> Result<crate::models::QueryResult> {
212        let url = format!(
213            "{}/api/queries/{query_id}/results/{result_id}.json",
214            self.base_url
215        );
216
217        let result_response: crate::models::QueryResultResponse = self
218            .get_json(&url, &format!("result {result_id} for query {query_id}"))
219            .await?;
220
221        Ok(result_response.query_result)
222    }
223
224    pub async fn execute_query_with_polling(
225        &self,
226        query_id: u64,
227        parameters: Option<std::collections::HashMap<String, serde_json::Value>>,
228        timeout_secs: u64,
229        poll_interval_ms: u64,
230    ) -> Result<crate::models::QueryResult> {
231        use crate::models::JobStatus;
232        use tokio::time::{Duration, sleep};
233
234        eprintln!("Executing query {query_id}...");
235        let job = self.refresh_query(query_id, parameters).await?;
236
237        let start = std::time::Instant::now();
238        let timeout = Duration::from_secs(timeout_secs);
239        let poll_interval = Duration::from_millis(poll_interval_ms);
240
241        let mut current_job = job;
242        loop {
243            if start.elapsed() > timeout {
244                anyhow::bail!("Query execution timed out after {timeout_secs} seconds");
245            }
246
247            let status = JobStatus::from_u8(current_job.status)?;
248
249            match status {
250                JobStatus::Success => {
251                    let result_id = current_job
252                        .query_result_id
253                        .context("Job succeeded but no result_id returned")?;
254
255                    eprintln!("Query completed, fetching results...");
256                    return self.get_query_result(query_id, result_id).await;
257                }
258                JobStatus::Failure => {
259                    let error = current_job
260                        .error
261                        .unwrap_or_else(|| "Unknown error".to_string());
262                    anyhow::bail!("Query execution failed: {error}");
263                }
264                JobStatus::Cancelled => {
265                    anyhow::bail!("Query execution was cancelled");
266                }
267                JobStatus::Pending | JobStatus::Started => {
268                    eprint!(".");
269                    sleep(poll_interval).await;
270                    current_job = self.poll_job(&current_job.id).await?;
271                }
272            }
273        }
274    }
275
276    pub async fn archive_query(&self, query_id: u64) -> Result<Query> {
277        let url = format!("{}/api/queries/{query_id}", self.base_url);
278        let payload = serde_json::json!({"is_archived": true});
279        self.post_json(&url, &payload, &format!("query {query_id} archive"))
280            .await
281    }
282
283    pub async fn unarchive_query(&self, query_id: u64) -> Result<Query> {
284        let url = format!("{}/api/queries/{query_id}", self.base_url);
285        let payload = serde_json::json!({"is_archived": false});
286        self.post_json(&url, &payload, &format!("query {query_id} unarchive"))
287            .await
288    }
289
290    pub async fn create_dashboard(&self, dashboard: &CreateDashboard) -> Result<Dashboard> {
291        let url = format!("{}/api/dashboards", self.base_url);
292        self.post_json(&url, dashboard, "new dashboard").await
293    }
294
295    pub async fn list_favorite_dashboards(
296        &self,
297        page: u32,
298        page_size: u32,
299    ) -> Result<DashboardsResponse> {
300        let url = format!(
301            "{}/api/dashboards/favorites?page={page}&page_size={page_size}",
302            self.base_url
303        );
304        self.get_json(&url, "favorite dashboards").await
305    }
306
307    pub async fn get_dashboard(&self, slug_or_id: &str) -> Result<Dashboard> {
308        let url = format!("{}/api/dashboards/{slug_or_id}", self.base_url);
309        self.get_json(&url, &format!("dashboard {slug_or_id}"))
310            .await
311    }
312
313    pub async fn update_dashboard(&self, dashboard: &Dashboard) -> Result<Dashboard> {
314        let url = format!("{}/api/dashboards/{}", self.base_url, dashboard.id);
315        self.post_json(
316            &url,
317            dashboard,
318            &format!("dashboard {} update", dashboard.id),
319        )
320        .await
321    }
322
323    pub async fn archive_dashboard(&self, dashboard_id: u64) -> Result<()> {
324        let url = format!("{}/api/dashboards/{dashboard_id}", self.base_url);
325        let payload = serde_json::json!({"is_archived": true});
326        let response = self
327            .client
328            .post(&url)
329            .json(&payload)
330            .send()
331            .await
332            .context(format!("Failed to archive dashboard {dashboard_id}"))?;
333
334        ensure_success(response).await?;
335
336        Ok(())
337    }
338
339    pub async fn unarchive_dashboard(&self, dashboard_id: u64) -> Result<Dashboard> {
340        let url = format!("{}/api/dashboards/{dashboard_id}", self.base_url);
341        let payload = serde_json::json!({"is_archived": false});
342        self.post_json(
343            &url,
344            &payload,
345            &format!("dashboard {dashboard_id} unarchive"),
346        )
347        .await
348    }
349
350    pub async fn create_widget(&self, widget: &CreateWidget) -> Result<crate::models::Widget> {
351        let url = format!("{}/api/widgets", self.base_url);
352        self.post_json(&url, widget, "new widget").await
353    }
354
355    pub async fn update_widget(
356        &self,
357        widget_id: u64,
358        widget: &CreateWidget,
359    ) -> Result<crate::models::Widget> {
360        let url = format!("{}/api/widgets/{widget_id}", self.base_url);
361        self.post_json(&url, widget, &format!("widget {widget_id} update"))
362            .await
363    }
364
365    pub async fn delete_widget(&self, widget_id: u64) -> Result<()> {
366        let url = format!("{}/api/widgets/{widget_id}", self.base_url);
367        let response = self
368            .client
369            .delete(&url)
370            .send()
371            .await
372            .context(format!("Failed to delete widget {widget_id}"))?;
373
374        ensure_success(response).await?;
375
376        Ok(())
377    }
378
379    pub async fn favorite_dashboard(&self, slug: &str) -> Result<()> {
380        let url = format!("{}/api/dashboards/{slug}/favorite", self.base_url);
381        let response = self
382            .client
383            .post(&url)
384            .json(&serde_json::json!({}))
385            .send()
386            .await
387            .context(format!("Failed to favorite dashboard {slug}"))?;
388
389        ensure_success(response).await?;
390
391        Ok(())
392    }
393
394    pub async fn fetch_favorite_dashboards(&self) -> Result<Vec<DashboardSummary>> {
395        let mut all_dashboards = Vec::new();
396        let mut page = 1;
397        let page_size = 100;
398
399        loop {
400            let response = self.list_favorite_dashboards(page, page_size).await?;
401
402            if response.results.is_empty() {
403                break;
404            }
405
406            all_dashboards.extend(response.results);
407            eprintln!(
408                "Fetched {} / {} dashboards...",
409                all_dashboards.len(),
410                response.count
411            );
412
413            #[allow(clippy::cast_possible_truncation)]
414            if all_dashboards.len() >= response.count as usize {
415                break;
416            }
417
418            page += 1;
419        }
420
421        Ok(all_dashboards)
422    }
423
424    async fn get_with_retry(
425        &self,
426        url: &str,
427        params: &[(&str, String)],
428    ) -> Result<reqwest::Response> {
429        use tokio::time::{Duration, sleep};
430
431        const MAX_ATTEMPTS: u32 = 4;
432        let base_delays = [250u64, 500, 1000, 2000];
433
434        let mut last_error = anyhow::anyhow!("No attempts made");
435        for attempt in 0..MAX_ATTEMPTS {
436            let response = self
437                .client
438                .get(url)
439                .query(params)
440                .send()
441                .await
442                .with_context(|| format!("Failed to GET {url}"))?;
443
444            let status = response.status();
445            let should_retry = status.as_u16() == 429 || status.is_server_error();
446
447            match ensure_success(response).await {
448                Ok(response) => return Ok(response),
449                Err(err) if !should_retry || attempt + 1 == MAX_ATTEMPTS => return Err(err),
450                Err(err) => last_error = err,
451            }
452
453            let delay_ms = base_delays[attempt as usize];
454            sleep(Duration::from_millis(delay_ms)).await;
455        }
456
457        Err(last_error)
458    }
459
460    async fn list_queries(&self, q: &str, page: u32, page_size: u32) -> Result<QueriesResponse> {
461        let url = format!("{}/api/queries", self.base_url);
462        let params = [
463            ("page", page.to_string()),
464            ("page_size", page_size.to_string()),
465            ("q", q.to_string()),
466        ];
467        self.get_with_retry(&url, &params)
468            .await?
469            .json()
470            .await
471            .context("Failed to parse queries response")
472    }
473
474    async fn list_dashboards(
475        &self,
476        q: &str,
477        page: u32,
478        page_size: u32,
479    ) -> Result<DashboardsResponse> {
480        let url = format!("{}/api/dashboards", self.base_url);
481        let params = [
482            ("page", page.to_string()),
483            ("page_size", page_size.to_string()),
484            ("q", q.to_string()),
485        ];
486        self.get_with_retry(&url, &params)
487            .await?
488            .json()
489            .await
490            .context("Failed to parse dashboards response")
491    }
492
493    pub async fn search_queries(&self, q: &str, limit: usize) -> Result<Vec<Query>> {
494        const PAGE_SIZE: usize = 250;
495
496        let mut results: Vec<Query> = Vec::new();
497        let mut page = 1u32;
498
499        loop {
500            let remaining = limit - results.len();
501            #[allow(clippy::cast_possible_truncation)]
502            let page_size = remaining.min(PAGE_SIZE) as u32;
503            let response = self.list_queries(q, page, page_size).await?;
504
505            results.extend(response.results);
506
507            #[allow(clippy::cast_possible_truncation)]
508            if results.len() >= limit || results.len() >= response.count as usize {
509                break;
510            }
511
512            page += 1;
513        }
514
515        results.truncate(limit);
516        Ok(results)
517    }
518
519    pub async fn search_dashboards(&self, q: &str, limit: usize) -> Result<Vec<DashboardSummary>> {
520        const PAGE_SIZE: usize = 250;
521
522        let mut results: Vec<DashboardSummary> = Vec::new();
523        let mut page = 1u32;
524
525        loop {
526            let remaining = limit - results.len();
527            #[allow(clippy::cast_possible_truncation)]
528            let page_size = remaining.min(PAGE_SIZE) as u32;
529            let response = self.list_dashboards(q, page, page_size).await?;
530
531            results.extend(response.results);
532
533            #[allow(clippy::cast_possible_truncation)]
534            if results.len() >= limit || results.len() >= response.count as usize {
535                break;
536            }
537
538            page += 1;
539        }
540
541        results.truncate(limit);
542        Ok(results)
543    }
544}
545
546async fn ensure_success(response: reqwest::Response) -> Result<reqwest::Response> {
547    let status = response.status();
548    if !status.is_success() {
549        let body = response.text().await.unwrap_or_default();
550        anyhow::bail!("API error {status}: {body}");
551    }
552    Ok(response)
553}