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};
9
10pub struct RedashClient {
11    client: Client,
12    base_url: String,
13}
14
15impl RedashClient {
16    pub fn new(base_url: String, api_key: &str) -> Result<Self> {
17        let mut headers = header::HeaderMap::new();
18        headers.insert(
19            "Authorization",
20            header::HeaderValue::from_str(&format!("Key {api_key}"))
21                .context("Invalid API key format")?,
22        );
23
24        let client = Client::builder()
25            .default_headers(headers)
26            .build()
27            .context("Failed to build HTTP client")?;
28
29        Ok(Self { client, base_url })
30    }
31
32    pub async fn list_my_queries(&self, page: u32, page_size: u32) -> Result<QueriesResponse> {
33        let url = format!(
34            "{}/api/queries/my?page={page}&page_size={page_size}",
35            self.base_url
36        );
37        let response = self
38            .client
39            .get(&url)
40            .send()
41            .await
42            .context("Failed to fetch my queries")?
43            .error_for_status()
44            .context("API returned error status")?;
45
46        response
47            .json()
48            .await
49            .context("Failed to parse queries response")
50    }
51
52    pub async fn get_query(&self, query_id: u64) -> Result<Query> {
53        let url = format!("{}/api/queries/{query_id}", self.base_url);
54        let response = self
55            .client
56            .get(&url)
57            .send()
58            .await
59            .context(format!("Failed to fetch query {query_id}"))?
60            .error_for_status()
61            .context("API returned error status")?;
62
63        response
64            .json()
65            .await
66            .context("Failed to parse query response")
67    }
68
69    pub async fn list_data_sources(&self) -> Result<Vec<DataSource>> {
70        let url = format!("{}/api/data_sources", self.base_url);
71        let response = self
72            .client
73            .get(&url)
74            .send()
75            .await
76            .context("Failed to fetch data sources")?
77            .error_for_status()
78            .context("API returned error status")?;
79
80        response
81            .json()
82            .await
83            .context("Failed to parse data sources response")
84    }
85
86    pub async fn get_data_source(&self, data_source_id: u64) -> Result<DataSource> {
87        let url = format!("{}/api/data_sources/{data_source_id}", self.base_url);
88        let response = self
89            .client
90            .get(&url)
91            .send()
92            .await
93            .context(format!("Failed to fetch data source {data_source_id}"))?
94            .error_for_status()
95            .context("API returned error status")?;
96
97        response
98            .json()
99            .await
100            .context("Failed to parse data source response")
101    }
102
103    pub async fn get_data_source_schema(
104        &self,
105        data_source_id: u64,
106        refresh: bool,
107    ) -> Result<DataSourceSchema> {
108        let url = if refresh {
109            format!(
110                "{}/api/data_sources/{data_source_id}/schema?refresh=true",
111                self.base_url
112            )
113        } else {
114            format!("{}/api/data_sources/{data_source_id}/schema", self.base_url)
115        };
116
117        let response = self
118            .client
119            .get(&url)
120            .send()
121            .await
122            .context(format!(
123                "Failed to fetch schema for data source {data_source_id}"
124            ))?
125            .error_for_status()
126            .context("API returned error status")?;
127
128        response
129            .json()
130            .await
131            .context("Failed to parse schema response")
132    }
133
134    pub async fn create_query(&self, create_query: &CreateQuery) -> Result<Query> {
135        let url = format!("{}/api/queries", self.base_url);
136        let response = self
137            .client
138            .post(&url)
139            .json(create_query)
140            .send()
141            .await
142            .context("Failed to create query")?
143            .error_for_status()
144            .context("API returned error status")?;
145
146        response
147            .json()
148            .await
149            .context("Failed to parse query create response")
150    }
151
152    pub async fn create_or_update_query(&self, query: &Query) -> Result<Query> {
153        let url = format!("{}/api/queries/{}", self.base_url, query.id);
154        let response = self
155            .client
156            .post(&url)
157            .json(query)
158            .send()
159            .await
160            .context(format!("Failed to update query {}", query.id))?
161            .error_for_status()
162            .context("API returned error status")?;
163
164        response
165            .json()
166            .await
167            .context("Failed to parse query update response")
168    }
169
170    pub async fn create_visualization(
171        &self,
172        query_id: u64,
173        viz: &crate::models::CreateVisualization,
174    ) -> Result<crate::models::Visualization> {
175        let url = format!("{}/api/visualizations", self.base_url);
176        let response = self
177            .client
178            .post(&url)
179            .json(viz)
180            .send()
181            .await
182            .context(format!(
183                "Failed to create visualization for query {query_id}"
184            ))?
185            .error_for_status()
186            .context("API returned error status")?;
187
188        response
189            .json()
190            .await
191            .context("Failed to parse visualization create response")
192    }
193
194    pub async fn update_visualization(
195        &self,
196        viz: &crate::models::Visualization,
197    ) -> Result<crate::models::Visualization> {
198        let url = format!("{}/api/visualizations/{}", self.base_url, viz.id);
199        let response = self
200            .client
201            .post(&url)
202            .json(viz)
203            .send()
204            .await
205            .context(format!("Failed to update visualization {}", viz.id))?
206            .error_for_status()
207            .context("API returned error status")?;
208
209        response
210            .json()
211            .await
212            .context("Failed to parse visualization update response")
213    }
214
215    pub async fn fetch_all_queries(&self) -> Result<Vec<Query>> {
216        let mut all_queries = Vec::new();
217        let mut page = 1;
218        let page_size = 100;
219
220        loop {
221            let response = self.list_my_queries(page, page_size).await?;
222
223            if response.results.is_empty() {
224                break;
225            }
226
227            all_queries.extend(response.results);
228            eprintln!(
229                "Fetched {} / {} queries...",
230                all_queries.len(),
231                response.count
232            );
233
234            #[allow(clippy::cast_possible_truncation)]
235            if all_queries.len() >= response.count as usize {
236                break;
237            }
238
239            page += 1;
240        }
241
242        Ok(all_queries)
243    }
244
245    pub async fn refresh_query(
246        &self,
247        query_id: u64,
248        parameters: Option<std::collections::HashMap<String, serde_json::Value>>,
249    ) -> Result<crate::models::Job> {
250        let url = format!("{}/api/queries/{query_id}/results", self.base_url);
251
252        let request = crate::models::RefreshRequest {
253            max_age: 0,
254            parameters,
255        };
256
257        let response = self
258            .client
259            .post(&url)
260            .json(&request)
261            .send()
262            .await
263            .context(format!("Failed to refresh query {query_id}"))?;
264
265        let status = response.status();
266        if !status.is_success() {
267            let error_body = response
268                .text()
269                .await
270                .unwrap_or_else(|_| "Unable to read error response".to_string());
271            anyhow::bail!("API returned error status {status}: {error_body}");
272        }
273
274        let job_response: crate::models::JobResponse = response
275            .json()
276            .await
277            .context("Failed to parse job response")?;
278
279        Ok(job_response.job)
280    }
281
282    pub async fn poll_job(&self, job_id: &str) -> Result<crate::models::Job> {
283        let url = format!("{}/api/jobs/{job_id}", self.base_url);
284
285        let response = self
286            .client
287            .get(&url)
288            .send()
289            .await
290            .context(format!("Failed to poll job {job_id}"))?
291            .error_for_status()
292            .context("API returned error status")?;
293
294        let job_response: crate::models::JobResponse = response
295            .json()
296            .await
297            .context("Failed to parse job response")?;
298
299        Ok(job_response.job)
300    }
301
302    pub async fn get_query_result(
303        &self,
304        query_id: u64,
305        result_id: u64,
306    ) -> Result<crate::models::QueryResult> {
307        let url = format!(
308            "{}/api/queries/{query_id}/results/{result_id}.json",
309            self.base_url
310        );
311
312        let response = self
313            .client
314            .get(&url)
315            .send()
316            .await
317            .context(format!(
318                "Failed to fetch result {result_id} for query {query_id}"
319            ))?
320            .error_for_status()
321            .context("API returned error status")?;
322
323        let result_response: crate::models::QueryResultResponse = response
324            .json()
325            .await
326            .context("Failed to parse query result response")?;
327
328        Ok(result_response.query_result)
329    }
330
331    pub async fn execute_query_with_polling(
332        &self,
333        query_id: u64,
334        parameters: Option<std::collections::HashMap<String, serde_json::Value>>,
335        timeout_secs: u64,
336        poll_interval_ms: u64,
337    ) -> Result<crate::models::QueryResult> {
338        use crate::models::JobStatus;
339        use tokio::time::{Duration, sleep};
340
341        eprintln!("Executing query {query_id}...");
342        let job = self.refresh_query(query_id, parameters).await?;
343
344        let start = std::time::Instant::now();
345        let timeout = Duration::from_secs(timeout_secs);
346        let poll_interval = Duration::from_millis(poll_interval_ms);
347
348        let mut current_job = job;
349        loop {
350            if start.elapsed() > timeout {
351                anyhow::bail!("Query execution timed out after {timeout_secs} seconds");
352            }
353
354            let status = JobStatus::from_u8(current_job.status)?;
355
356            match status {
357                JobStatus::Success => {
358                    let result_id = current_job
359                        .query_result_id
360                        .context("Job succeeded but no result_id returned")?;
361
362                    eprintln!("Query completed, fetching results...");
363                    return self.get_query_result(query_id, result_id).await;
364                }
365                JobStatus::Failure => {
366                    let error = current_job
367                        .error
368                        .unwrap_or_else(|| "Unknown error".to_string());
369                    anyhow::bail!("Query execution failed: {error}");
370                }
371                JobStatus::Cancelled => {
372                    anyhow::bail!("Query execution was cancelled");
373                }
374                JobStatus::Pending | JobStatus::Started => {
375                    eprint!(".");
376                    sleep(poll_interval).await;
377                    current_job = self.poll_job(&current_job.id).await?;
378                }
379            }
380        }
381    }
382
383    pub async fn archive_query(&self, query_id: u64) -> Result<Query> {
384        let url = format!("{}/api/queries/{query_id}", self.base_url);
385        let payload = serde_json::json!({"is_archived": true});
386
387        let response = self
388            .client
389            .post(&url)
390            .json(&payload)
391            .send()
392            .await
393            .context(format!("Failed to archive query {query_id}"))?
394            .error_for_status()
395            .context("API returned error status")?;
396
397        response
398            .json()
399            .await
400            .context("Failed to parse archive response")
401    }
402
403    pub async fn unarchive_query(&self, query_id: u64) -> Result<Query> {
404        let url = format!("{}/api/queries/{query_id}", self.base_url);
405        let payload = serde_json::json!({"is_archived": false});
406
407        let response = self
408            .client
409            .post(&url)
410            .json(&payload)
411            .send()
412            .await
413            .context(format!("Failed to unarchive query {query_id}"))?
414            .error_for_status()
415            .context("API returned error status")?;
416
417        response
418            .json()
419            .await
420            .context("Failed to parse unarchive response")
421    }
422
423    pub async fn create_dashboard(&self, dashboard: &CreateDashboard) -> Result<Dashboard> {
424        let url = format!("{}/api/dashboards", self.base_url);
425        let response = self
426            .client
427            .post(&url)
428            .json(dashboard)
429            .send()
430            .await
431            .context("Failed to create dashboard")?;
432
433        let status = response.status();
434        if !status.is_success() {
435            anyhow::bail!(
436                "HTTP {}: {}",
437                status.as_u16(),
438                status.canonical_reason().unwrap_or("Unknown error")
439            );
440        }
441
442        response
443            .json()
444            .await
445            .context("Failed to parse dashboard create response")
446    }
447
448    pub async fn list_favorite_dashboards(
449        &self,
450        page: u32,
451        page_size: u32,
452    ) -> Result<DashboardsResponse> {
453        let url = format!(
454            "{}/api/dashboards/favorites?page={page}&page_size={page_size}",
455            self.base_url
456        );
457        let response = self
458            .client
459            .get(&url)
460            .send()
461            .await
462            .context("Failed to fetch dashboards")?;
463
464        let status = response.status();
465        if !status.is_success() {
466            anyhow::bail!(
467                "HTTP {}: {}",
468                status.as_u16(),
469                status.canonical_reason().unwrap_or("Unknown error")
470            );
471        }
472
473        response
474            .json()
475            .await
476            .context("Failed to parse dashboards response")
477    }
478
479    pub async fn get_dashboard(&self, slug_or_id: &str) -> Result<Dashboard> {
480        let url = format!("{}/api/dashboards/{slug_or_id}", self.base_url);
481        let response = self
482            .client
483            .get(&url)
484            .send()
485            .await
486            .context(format!("Failed to fetch dashboard {slug_or_id}"))?;
487
488        let status = response.status();
489        if !status.is_success() {
490            let body = response.text().await.unwrap_or_default();
491            anyhow::bail!(
492                "HTTP {}: {} — {body}",
493                status.as_u16(),
494                status.canonical_reason().unwrap_or("Unknown error")
495            );
496        }
497
498        response
499            .json()
500            .await
501            .context("Failed to parse dashboard response")
502    }
503
504    pub async fn update_dashboard(&self, dashboard: &Dashboard) -> Result<Dashboard> {
505        let url = format!("{}/api/dashboards/{}", self.base_url, dashboard.id);
506        let response = self
507            .client
508            .post(&url)
509            .json(dashboard)
510            .send()
511            .await
512            .context(format!("Failed to update dashboard {}", dashboard.id))?;
513
514        let status = response.status();
515        if !status.is_success() {
516            let body = response.text().await.unwrap_or_default();
517            anyhow::bail!(
518                "HTTP {}: {} — {body}",
519                status.as_u16(),
520                status.canonical_reason().unwrap_or("Unknown error")
521            );
522        }
523
524        response
525            .json()
526            .await
527            .context("Failed to parse dashboard update response")
528    }
529
530    pub async fn archive_dashboard(&self, dashboard_id: u64) -> Result<()> {
531        let url = format!("{}/api/dashboards/{dashboard_id}", self.base_url);
532        let payload = serde_json::json!({"is_archived": true});
533        let response = self
534            .client
535            .post(&url)
536            .json(&payload)
537            .send()
538            .await
539            .context(format!("Failed to archive dashboard {dashboard_id}"))?;
540
541        let status = response.status();
542        if !status.is_success() {
543            anyhow::bail!(
544                "HTTP {}: {}",
545                status.as_u16(),
546                status.canonical_reason().unwrap_or("Unknown error")
547            );
548        }
549
550        Ok(())
551    }
552
553    pub async fn unarchive_dashboard(&self, dashboard_id: u64) -> Result<Dashboard> {
554        let url = format!("{}/api/dashboards/{dashboard_id}", self.base_url);
555        let payload = serde_json::json!({"is_archived": false});
556
557        let response = self
558            .client
559            .post(&url)
560            .json(&payload)
561            .send()
562            .await
563            .context(format!("Failed to unarchive dashboard {dashboard_id}"))?;
564
565        let status = response.status();
566        if !status.is_success() {
567            anyhow::bail!(
568                "HTTP {}: {}",
569                status.as_u16(),
570                status.canonical_reason().unwrap_or("Unknown error")
571            );
572        }
573
574        response
575            .json()
576            .await
577            .context("Failed to parse unarchive response")
578    }
579
580    pub async fn create_widget(&self, widget: &CreateWidget) -> Result<crate::models::Widget> {
581        let url = format!("{}/api/widgets", self.base_url);
582        let response = self
583            .client
584            .post(&url)
585            .json(widget)
586            .send()
587            .await
588            .context("Failed to create widget")?;
589
590        let status = response.status();
591        if !status.is_success() {
592            let body = response.text().await.unwrap_or_default();
593            anyhow::bail!(
594                "HTTP {}: {} — {body}",
595                status.as_u16(),
596                status.canonical_reason().unwrap_or("Unknown error")
597            );
598        }
599
600        response
601            .json()
602            .await
603            .context("Failed to parse widget create response")
604    }
605
606    pub async fn update_widget(
607        &self,
608        widget_id: u64,
609        widget: &CreateWidget,
610    ) -> Result<crate::models::Widget> {
611        let url = format!("{}/api/widgets/{widget_id}", self.base_url);
612        let response = self
613            .client
614            .post(&url)
615            .json(widget)
616            .send()
617            .await
618            .context(format!("Failed to update widget {widget_id}"))?;
619
620        let status = response.status();
621        if !status.is_success() {
622            let body = response.text().await.unwrap_or_default();
623            anyhow::bail!(
624                "HTTP {}: {} — {body}",
625                status.as_u16(),
626                status.canonical_reason().unwrap_or("Unknown error")
627            );
628        }
629
630        response
631            .json()
632            .await
633            .context("Failed to parse widget update response")
634    }
635
636    pub async fn delete_widget(&self, widget_id: u64) -> Result<()> {
637        let url = format!("{}/api/widgets/{widget_id}", self.base_url);
638        let response = self
639            .client
640            .delete(&url)
641            .send()
642            .await
643            .context(format!("Failed to delete widget {widget_id}"))?;
644
645        let status = response.status();
646        if !status.is_success() {
647            anyhow::bail!(
648                "HTTP {}: {}",
649                status.as_u16(),
650                status.canonical_reason().unwrap_or("Unknown error")
651            );
652        }
653
654        Ok(())
655    }
656
657    pub async fn favorite_dashboard(&self, slug: &str) -> Result<()> {
658        let url = format!("{}/api/dashboards/{slug}/favorite", self.base_url);
659        let response = self
660            .client
661            .post(&url)
662            .json(&serde_json::json!({}))
663            .send()
664            .await
665            .context(format!("Failed to favorite dashboard {slug}"))?;
666
667        let status = response.status();
668        if !status.is_success() {
669            anyhow::bail!(
670                "HTTP {}: {}",
671                status.as_u16(),
672                status.canonical_reason().unwrap_or("Unknown error")
673            );
674        }
675
676        Ok(())
677    }
678
679    pub async fn fetch_favorite_dashboards(&self) -> Result<Vec<DashboardSummary>> {
680        let mut all_dashboards = Vec::new();
681        let mut page = 1;
682        let page_size = 100;
683
684        loop {
685            let response = self.list_favorite_dashboards(page, page_size).await?;
686
687            if response.results.is_empty() {
688                break;
689            }
690
691            all_dashboards.extend(response.results);
692            eprintln!(
693                "Fetched {} / {} dashboards...",
694                all_dashboards.len(),
695                response.count
696            );
697
698            #[allow(clippy::cast_possible_truncation)]
699            if all_dashboards.len() >= response.count as usize {
700                break;
701            }
702
703            page += 1;
704        }
705
706        Ok(all_dashboards)
707    }
708}