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(¤t_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}