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(¤t_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, ¶ms)
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, ¶ms)
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}