Skip to content

Commit 4357ccb

Browse files
authored
feat: Phase 2 - incidents, logs, monitor expansion, query compiler (#8)
feat: add incidents, logs, monitor expansion, and query compiler (Phase 2) Incidents (API v3): list, get, create, ack, resolve, escalate, delete, timeline with colored status output (started/acknowledged/resolved). Monitor expansion: update command with field-level edits, availability/SLA endpoint, response-times with nested region breakdown. Logs: sources listing via Telemetry API, raw SQL queries via ClickHouse Query API with Basic auth, simple filter query compiler (level:error, status:>=500, wildcards, AND/OR), and polling-based live tail. Query compiler translates human filters to ClickHouse SQL with duration parsing (1h, 30m, 7d), exact match, numeric comparison, wildcard, quoted contains, and boolean combinators. Includes 9 unit tests. Adds v3 HTTP helpers, shared parse_one/check_status, paginate_all_v3, telemetry client factory, SqlClient with HTTP Basic auth, and optional telemetry client in AppContext with env var override support.
1 parent 1ce883b commit 4357ccb

19 files changed

Lines changed: 1604 additions & 28 deletions

File tree

‎src/adapters/http/incidents.rs‎

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,66 @@
1+
use anyhow::Result;
2+
3+
use super::retry::with_retry;
4+
use super::{HttpClient, check_status, parse_one};
5+
use crate::types::{CreateIncidentRequest, IncidentFilters, IncidentResource, TimelineEvent};
6+
7+
impl HttpClient {
8+
pub async fn list_incidents(&self, filters: &IncidentFilters) -> Result<Vec<IncidentResource>> {
9+
let mut params: Vec<(&str, &str)> = Vec::new();
10+
if let Some(ref status) = filters.status {
11+
params.push(("status", status.as_str()));
12+
}
13+
if let Some(ref monitor_id) = filters.monitor_id {
14+
params.push(("monitor_id", monitor_id.as_str()));
15+
}
16+
if let Some(ref from) = filters.from {
17+
params.push(("from", from.as_str()));
18+
}
19+
if let Some(ref to) = filters.to {
20+
params.push(("to", to.as_str()));
21+
}
22+
23+
self.paginate_all_v3("/incidents", &params).await
24+
}
25+
26+
pub async fn get_incident(&self, id: &str) -> Result<IncidentResource> {
27+
let path = format!("/incidents/{id}");
28+
let resp = with_retry(|| async { Ok(self.get_v3(&path).send().await?) }).await?;
29+
parse_one(resp).await
30+
}
31+
32+
pub async fn create_incident(&self, req: &CreateIncidentRequest) -> Result<IncidentResource> {
33+
let resp =
34+
with_retry(|| async { Ok(self.post_v3("/incidents").json(req).send().await?) }).await?;
35+
parse_one(resp).await
36+
}
37+
38+
pub async fn acknowledge_incident(&self, id: &str) -> Result<IncidentResource> {
39+
let path = format!("/incidents/{id}/acknowledge");
40+
let resp = with_retry(|| async { Ok(self.post_v3(&path).send().await?) }).await?;
41+
parse_one(resp).await
42+
}
43+
44+
pub async fn resolve_incident(&self, id: &str) -> Result<IncidentResource> {
45+
let path = format!("/incidents/{id}/resolve");
46+
let resp = with_retry(|| async { Ok(self.post_v3(&path).send().await?) }).await?;
47+
parse_one(resp).await
48+
}
49+
50+
pub async fn escalate_incident(&self, id: &str) -> Result<IncidentResource> {
51+
let path = format!("/incidents/{id}/escalate");
52+
let resp = with_retry(|| async { Ok(self.post_v3(&path).send().await?) }).await?;
53+
parse_one(resp).await
54+
}
55+
56+
pub async fn delete_incident(&self, id: &str) -> Result<()> {
57+
let path = format!("/incidents/{id}");
58+
let resp = with_retry(|| async { Ok(self.delete_v3(&path).send().await?) }).await?;
59+
check_status(resp).await
60+
}
61+
62+
pub async fn incident_timeline(&self, id: &str) -> Result<Vec<TimelineEvent>> {
63+
let path = format!("/incidents/{id}/timeline");
64+
self.paginate_all_v3(&path, &[]).await
65+
}
66+
}

‎src/adapters/http/mod.rs‎

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,15 @@
1+
pub mod incidents;
12
pub mod pagination;
23
pub mod retry;
4+
pub mod sql;
5+
pub mod telemetry;
36
pub mod uptime;
47

8+
use anyhow::{Context, Result, bail};
59
use reqwest::header::{AUTHORIZATION, CONTENT_TYPE, HeaderMap, HeaderValue};
610

11+
use crate::types::SingleResponse;
12+
713
pub struct HttpClient {
814
client: reqwest::Client,
915
base_url: String,
@@ -39,14 +45,27 @@ impl HttpClient {
3945
format!("{}{}", self.base_url, path)
4046
}
4147

48+
pub fn url_v3(&self, path: &str) -> String {
49+
let v3_base = self.base_url.replace("/api/v2", "/api/v3");
50+
format!("{}{}", v3_base, path)
51+
}
52+
4253
pub fn get(&self, path: &str) -> reqwest::RequestBuilder {
4354
self.client.get(self.url(path)).headers(self.headers())
4455
}
4556

57+
pub fn get_v3(&self, path: &str) -> reqwest::RequestBuilder {
58+
self.client.get(self.url_v3(path)).headers(self.headers())
59+
}
60+
4661
pub fn post(&self, path: &str) -> reqwest::RequestBuilder {
4762
self.client.post(self.url(path)).headers(self.headers())
4863
}
4964

65+
pub fn post_v3(&self, path: &str) -> reqwest::RequestBuilder {
66+
self.client.post(self.url_v3(path)).headers(self.headers())
67+
}
68+
5069
pub fn patch(&self, path: &str) -> reqwest::RequestBuilder {
5170
self.client.patch(self.url(path)).headers(self.headers())
5271
}
@@ -55,7 +74,34 @@ impl HttpClient {
5574
self.client.delete(self.url(path)).headers(self.headers())
5675
}
5776

77+
pub fn delete_v3(&self, path: &str) -> reqwest::RequestBuilder {
78+
self.client
79+
.delete(self.url_v3(path))
80+
.headers(self.headers())
81+
}
82+
5883
pub fn get_absolute(&self, url: &str) -> reqwest::RequestBuilder {
5984
self.client.get(url).headers(self.headers())
6085
}
6186
}
87+
88+
/// Parse a single-resource response, or return an API error.
89+
pub async fn parse_one<T: serde::de::DeserializeOwned>(resp: reqwest::Response) -> Result<T> {
90+
let status = resp.status();
91+
if !status.is_success() {
92+
let body = resp.text().await.unwrap_or_default();
93+
bail!("API error ({}): {}", status, body);
94+
}
95+
let single: SingleResponse<T> = resp.json().await.context("Failed to parse response")?;
96+
Ok(single.data)
97+
}
98+
99+
/// Check response status and return an error with body if not successful.
100+
pub async fn check_status(resp: reqwest::Response) -> Result<()> {
101+
if !resp.status().is_success() {
102+
let status = resp.status();
103+
let body = resp.text().await.unwrap_or_default();
104+
bail!("API error ({}): {}", status, body);
105+
}
106+
Ok(())
107+
}

‎src/adapters/http/pagination.rs‎

Lines changed: 20 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,25 @@ impl HttpClient {
1212
&self,
1313
path: &str,
1414
query_params: &[(&str, &str)],
15+
) -> Result<Vec<T>> {
16+
self.paginate_all_with_builder(self.get(path), query_params)
17+
.await
18+
}
19+
20+
/// Fetches all pages from a v3 paginated endpoint, following `pagination.next` links.
21+
pub async fn paginate_all_v3<T: DeserializeOwned>(
22+
&self,
23+
path: &str,
24+
query_params: &[(&str, &str)],
25+
) -> Result<Vec<T>> {
26+
self.paginate_all_with_builder(self.get_v3(path), query_params)
27+
.await
28+
}
29+
30+
async fn paginate_all_with_builder<T: DeserializeOwned>(
31+
&self,
32+
initial_req: reqwest::RequestBuilder,
33+
query_params: &[(&str, &str)],
1534
) -> Result<Vec<T>> {
1635
let mut all = Vec::new();
1736
let mut next_url: Option<String> = None;
@@ -20,7 +39,7 @@ impl HttpClient {
2039
let resp = if let Some(ref url) = next_url {
2140
with_retry(|| async { Ok(self.get_absolute(url).send().await?) }).await?
2241
} else {
23-
let mut req = self.get(path);
42+
let mut req = initial_req.try_clone().expect("request clone failed");
2443
for (k, v) in query_params {
2544
req = req.query(&[(*k, *v)]);
2645
}

‎src/adapters/http/sql.rs‎

Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,59 @@
1+
use anyhow::{Context, Result, bail};
2+
3+
use super::retry::with_retry;
4+
5+
pub struct SqlClient {
6+
client: reqwest::Client,
7+
base_url: String,
8+
username: String,
9+
password: String,
10+
}
11+
12+
impl SqlClient {
13+
pub fn new(region: &str, username: &str, password: &str) -> Self {
14+
let base_url = format!("https://{region}-connect.betterstackdata.com");
15+
Self {
16+
client: reqwest::Client::new(),
17+
base_url,
18+
username: username.to_string(),
19+
password: password.to_string(),
20+
}
21+
}
22+
23+
pub async fn query(&self, sql: &str) -> Result<String> {
24+
let url = &self.base_url;
25+
let resp = with_retry(|| async {
26+
Ok(self
27+
.client
28+
.post(url)
29+
.basic_auth(&self.username, Some(&self.password))
30+
.body(format!("{sql} FORMAT JSONEachRow"))
31+
.send()
32+
.await?)
33+
})
34+
.await?;
35+
36+
let status = resp.status();
37+
if !status.is_success() {
38+
let body = resp.text().await.unwrap_or_default();
39+
bail!("SQL query error ({}): {}", status, body);
40+
}
41+
42+
resp.text().await.context("Failed to read SQL response")
43+
}
44+
45+
pub async fn query_json(&self, sql: &str) -> Result<Vec<serde_json::Value>> {
46+
let raw = self.query(sql).await?;
47+
let mut rows = Vec::new();
48+
for line in raw.lines() {
49+
let line = line.trim();
50+
if line.is_empty() {
51+
continue;
52+
}
53+
let value: serde_json::Value =
54+
serde_json::from_str(line).context("Failed to parse SQL result row")?;
55+
rows.push(value);
56+
}
57+
Ok(rows)
58+
}
59+
}

‎src/adapters/http/telemetry.rs‎

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,21 @@
1+
use anyhow::Result;
2+
3+
use super::retry::with_retry;
4+
use super::{HttpClient, parse_one};
5+
use crate::types::SourceResource;
6+
7+
impl HttpClient {
8+
pub fn telemetry(token: &str) -> Self {
9+
Self::new("https://telemetry.betterstack.com/api/v1", token)
10+
}
11+
12+
pub async fn list_sources(&self) -> Result<Vec<SourceResource>> {
13+
self.paginate_all("/sources", &[]).await
14+
}
15+
16+
pub async fn get_source(&self, id: &str) -> Result<SourceResource> {
17+
let path = format!("/sources/{id}");
18+
let resp = with_retry(|| async { Ok(self.get(&path).send().await?) }).await?;
19+
parse_one(resp).await
20+
}
21+
}

‎src/adapters/http/uptime.rs‎

Lines changed: 58 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -1,29 +1,11 @@
1-
use anyhow::{Context, Result, bail};
1+
use anyhow::Result;
22

3-
use super::HttpClient;
43
use super::retry::with_retry;
5-
use crate::types::{CreateMonitorRequest, MonitorFilters, MonitorResource, SingleResponse};
6-
7-
/// Parse a single-resource response, or return an API error.
8-
async fn parse_one<T: serde::de::DeserializeOwned>(resp: reqwest::Response) -> Result<T> {
9-
let status = resp.status();
10-
if !status.is_success() {
11-
let body = resp.text().await.unwrap_or_default();
12-
bail!("API error ({}): {}", status, body);
13-
}
14-
let single: SingleResponse<T> = resp.json().await.context("Failed to parse response")?;
15-
Ok(single.data)
16-
}
17-
18-
/// Check response status and return an error with body if not successful.
19-
async fn check_status(resp: reqwest::Response) -> Result<()> {
20-
if !resp.status().is_success() {
21-
let status = resp.status();
22-
let body = resp.text().await.unwrap_or_default();
23-
bail!("API error ({}): {}", status, body);
24-
}
25-
Ok(())
26-
}
4+
use super::{HttpClient, check_status, parse_one};
5+
use crate::types::{
6+
CreateMonitorRequest, MonitorFilters, MonitorResource, ResponseTimesResource, SlaResource,
7+
UpdateMonitorRequest,
8+
};
279

2810
impl HttpClient {
2911
pub async fn list_monitors(&self, filters: &MonitorFilters) -> Result<Vec<MonitorResource>> {
@@ -64,9 +46,61 @@ impl HttpClient {
6446
parse_one(resp).await
6547
}
6648

49+
pub async fn update_monitor(
50+
&self,
51+
id: &str,
52+
req: &UpdateMonitorRequest,
53+
) -> Result<MonitorResource> {
54+
let path = format!("/monitors/{id}");
55+
let resp = with_retry(|| async { Ok(self.patch(&path).json(req).send().await?) }).await?;
56+
parse_one(resp).await
57+
}
58+
6759
pub async fn delete_monitor(&self, id: &str) -> Result<()> {
6860
let path = format!("/monitors/{id}");
6961
let resp = with_retry(|| async { Ok(self.delete_req(&path).send().await?) }).await?;
7062
check_status(resp).await
7163
}
64+
65+
pub async fn monitor_sla(
66+
&self,
67+
id: &str,
68+
from: Option<&str>,
69+
to: Option<&str>,
70+
) -> Result<SlaResource> {
71+
let path = format!("/monitors/{id}/sla");
72+
let resp = with_retry(|| async {
73+
let mut req = self.get(&path);
74+
if let Some(f) = from {
75+
req = req.query(&[("from", f)]);
76+
}
77+
if let Some(t) = to {
78+
req = req.query(&[("to", t)]);
79+
}
80+
Ok(req.send().await?)
81+
})
82+
.await?;
83+
parse_one(resp).await
84+
}
85+
86+
pub async fn monitor_response_times(
87+
&self,
88+
id: &str,
89+
from: Option<&str>,
90+
to: Option<&str>,
91+
) -> Result<ResponseTimesResource> {
92+
let path = format!("/monitors/{id}/response-times");
93+
let resp = with_retry(|| async {
94+
let mut req = self.get(&path);
95+
if let Some(f) = from {
96+
req = req.query(&[("from", f)]);
97+
}
98+
if let Some(t) = to {
99+
req = req.query(&[("to", t)]);
100+
}
101+
Ok(req.send().await?)
102+
})
103+
.await?;
104+
parse_one(resp).await
105+
}
72106
}

0 commit comments

Comments
 (0)