1use reqwest::Client;
4use serde::Serialize;
5
6use super::retry::{RetryConfig, deliver_with_retry, is_accepted_202};
7use super::{Event, EventSubscriber, SubscriberFuture};
8
9const DEFAULT_INGEST_URL: &str = "https://in.logs.betterstack.com";
11
12#[derive(Debug, Serialize)]
14struct LogPayload {
15 dt: String,
17 level: &'static str,
19 message: String,
21 event: serde_json::Value,
23}
24
25pub struct BetterStackSubscriber {
50 source_token: String,
51 authorization_header: String,
52 ingest_url: String,
53 client: Client,
54 retry_config: RetryConfig,
55}
56
57impl BetterStackSubscriber {
58 pub fn new(source_token: &str) -> Self {
76 Self::with_url(source_token, DEFAULT_INGEST_URL)
77 }
78
79 pub fn with_url(source_token: &str, ingest_url: &str) -> Self {
99 Self::with_url_and_retry(source_token, ingest_url, RetryConfig::default())
100 }
101
102 pub fn with_url_and_retry(
125 source_token: &str,
126 ingest_url: &str,
127 retry_config: RetryConfig,
128 ) -> Self {
129 let client = retry_config.build_client();
130 Self {
131 authorization_header: format!("Bearer {}", source_token),
132 source_token: source_token.to_string(),
133 ingest_url: ingest_url.to_string(),
134 client,
135 retry_config,
136 }
137 }
138
139 pub fn source_token(&self) -> &str {
141 &self.source_token
142 }
143
144 pub fn ingest_url(&self) -> &str {
146 &self.ingest_url
147 }
148
149 #[deny(unreachable_patterns)]
151 fn build_payload(event: &Event) -> Option<LogPayload> {
152 match event {
153 Event::StepFailed(e) => {
154 let message = format!(
155 "Step '{}' ({}) failed on run {}: {}",
156 e.step_name, e.kind, e.run_id, e.error
157 );
158 let event_json = serde_json::json!({
159 "type": "step_failed",
160 "run_id": e.run_id.to_string(),
161 "step_id": e.step_id.to_string(),
162 "step_name": e.step_name,
163 "kind": e.kind.to_string(),
164 "error": e.error,
165 });
166 Some(LogPayload {
167 dt: e.at.to_rfc3339(),
168 level: "error",
169 message,
170 event: event_json,
171 })
172 }
173 Event::RunFailed(e) => {
174 let error_detail = e.error.as_deref().unwrap_or("unknown error");
175 let message = format!(
176 "Run {} (workflow '{}') failed: {}",
177 e.run_id, e.workflow_name, error_detail
178 );
179 let event_json = serde_json::json!({
180 "type": "run_failed",
181 "run_id": e.run_id.to_string(),
182 "workflow_name": e.workflow_name,
183 "error": error_detail,
184 "cost_usd": e.cost_usd.to_string(),
185 "duration_ms": e.duration_ms,
186 });
187 Some(LogPayload {
188 dt: e.at.to_rfc3339(),
189 level: "error",
190 message,
191 event: event_json,
192 })
193 }
194 _ => None,
195 }
196 }
197}
198
199impl EventSubscriber for BetterStackSubscriber {
200 fn name(&self) -> &str {
201 "betterstack"
202 }
203
204 fn handle<'a>(&'a self, event: &'a Event) -> SubscriberFuture<'a> {
205 Box::pin(async move {
206 if let Some(payload) = Self::build_payload(event) {
207 deliver_with_retry(
208 &self.retry_config,
209 || {
210 self.client
211 .post(&self.ingest_url)
212 .header("Authorization", &self.authorization_header)
213 .json(&payload)
214 },
215 is_accepted_202,
216 "betterstack",
217 &payload.message,
218 )
219 .await;
220 }
221 })
222 }
223}
224
225#[cfg(test)]
226mod tests {
227 use std::collections::HashMap;
228
229 use super::*;
230 use crate::notify::{
231 ApprovalRequestedEvent, RunCreatedEvent, RunFailedEvent, RunStatusChangedEvent,
232 StepCompletedEvent, StepFailedEvent, UserSignedInEvent,
233 };
234 use chrono::Utc;
235 use ironflow_store::models::{RunStatus, StepKind};
236 use rust_decimal::Decimal;
237 use uuid::Uuid;
238
239 #[test]
240 fn new_sets_default_ingest_url() {
241 let sub = BetterStackSubscriber::new("token-123");
242 assert_eq!(sub.source_token(), "token-123");
243 assert_eq!(sub.ingest_url(), DEFAULT_INGEST_URL);
244 }
245
246 #[test]
247 fn with_url_sets_custom_ingest_url() {
248 let sub = BetterStackSubscriber::with_url("token-123", "https://custom.example.com");
249 assert_eq!(sub.source_token(), "token-123");
250 assert_eq!(sub.ingest_url(), "https://custom.example.com");
251 }
252
253 #[test]
254 fn name_is_betterstack() {
255 let sub = BetterStackSubscriber::new("token");
256 assert_eq!(sub.name(), "betterstack");
257 }
258
259 #[test]
260 fn build_payload_step_failed() {
261 let event = Event::StepFailed(StepFailedEvent {
262 run_id: Uuid::now_v7(),
263 step_id: Uuid::now_v7(),
264 step_name: "build".to_string(),
265 kind: StepKind::Shell,
266 error: "exit code 1".to_string(),
267 at: Utc::now(),
268 });
269
270 let payload = BetterStackSubscriber::build_payload(&event);
271 assert!(payload.is_some());
272 let payload = payload.unwrap();
273 assert_eq!(payload.level, "error");
274 assert!(payload.message.contains("build"));
275 assert!(payload.message.contains("exit code 1"));
276 assert_eq!(payload.event["type"], "step_failed");
277 assert_eq!(payload.event["error"], "exit code 1");
278 }
279
280 #[test]
281 fn build_payload_run_failed() {
282 let event = Event::RunFailed(RunFailedEvent {
283 run_id: Uuid::now_v7(),
284 workflow_name: "deploy".to_string(),
285 error: Some("step 'build' failed".to_string()),
286 cost_usd: Decimal::new(42, 2),
287 duration_ms: 5000,
288 labels: HashMap::new(),
289 at: Utc::now(),
290 });
291
292 let payload = BetterStackSubscriber::build_payload(&event);
293 assert!(payload.is_some());
294 let payload = payload.unwrap();
295 assert_eq!(payload.level, "error");
296 assert!(payload.message.contains("deploy"));
297 assert!(payload.message.contains("step 'build' failed"));
298 assert_eq!(payload.event["type"], "run_failed");
299 assert_eq!(payload.event["workflow_name"], "deploy");
300 }
301
302 #[test]
303 fn build_payload_run_failed_without_error_message() {
304 let event = Event::RunFailed(RunFailedEvent {
305 run_id: Uuid::now_v7(),
306 workflow_name: "deploy".to_string(),
307 error: None,
308 cost_usd: Decimal::ZERO,
309 duration_ms: 1000,
310 labels: HashMap::new(),
311 at: Utc::now(),
312 });
313
314 let payload = BetterStackSubscriber::build_payload(&event).unwrap();
315 assert!(payload.message.contains("unknown error"));
316 assert_eq!(payload.event["error"], "unknown error");
317 }
318
319 #[test]
320 fn build_payload_run_completed_returns_none() {
321 let event = Event::RunStatusChanged(RunStatusChangedEvent {
322 run_id: Uuid::now_v7(),
323 workflow_name: "deploy".to_string(),
324 from: RunStatus::Running,
325 to: RunStatus::Completed,
326 error: None,
327 cost_usd: Decimal::ZERO,
328 duration_ms: 1000,
329 labels: HashMap::new(),
330 at: Utc::now(),
331 });
332
333 assert!(BetterStackSubscriber::build_payload(&event).is_none());
334 }
335
336 #[test]
337 fn build_payload_run_created_returns_none() {
338 let event = Event::RunCreated(RunCreatedEvent {
339 run_id: Uuid::now_v7(),
340 workflow_name: "deploy".to_string(),
341 at: Utc::now(),
342 });
343
344 assert!(BetterStackSubscriber::build_payload(&event).is_none());
345 }
346
347 #[test]
348 fn build_payload_step_completed_returns_none() {
349 let event = Event::StepCompleted(StepCompletedEvent {
350 run_id: Uuid::now_v7(),
351 step_id: Uuid::now_v7(),
352 step_name: "build".to_string(),
353 kind: StepKind::Shell,
354 duration_ms: 500,
355 cost_usd: Decimal::ZERO,
356 at: Utc::now(),
357 });
358
359 assert!(BetterStackSubscriber::build_payload(&event).is_none());
360 }
361
362 #[test]
363 fn build_payload_approval_requested_returns_none() {
364 let event = Event::ApprovalRequested(ApprovalRequestedEvent {
365 run_id: Uuid::now_v7(),
366 step_id: Uuid::now_v7(),
367 message: "Deploy to prod?".to_string(),
368 at: Utc::now(),
369 });
370
371 assert!(BetterStackSubscriber::build_payload(&event).is_none());
372 }
373
374 #[test]
375 fn build_payload_user_signed_in_returns_none() {
376 let event = Event::UserSignedIn(UserSignedInEvent {
377 user_id: Uuid::now_v7(),
378 username: "alice".to_string(),
379 at: Utc::now(),
380 });
381
382 assert!(BetterStackSubscriber::build_payload(&event).is_none());
383 }
384
385 #[tokio::test]
386 async fn handle_ignores_non_error_events() {
387 let sub = BetterStackSubscriber::with_url("token", "http://127.0.0.1:1");
388 let event = Event::RunCreated(RunCreatedEvent {
389 run_id: Uuid::now_v7(),
390 workflow_name: "deploy".to_string(),
391 at: Utc::now(),
392 });
393 sub.handle(&event).await;
395 }
396
397 #[tokio::test]
398 async fn deliver_to_real_endpoint_returns_202() {
399 use axum::Router;
400 use axum::http::StatusCode;
401 use axum::routing::post;
402 use tokio::net::TcpListener;
403
404 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
405 let addr = listener.local_addr().unwrap();
406
407 let app = Router::new().route("/", post(|| async { StatusCode::ACCEPTED }));
408
409 tokio::spawn(async move {
410 axum::serve(listener, app).await.unwrap();
411 });
412
413 let sub = BetterStackSubscriber::with_url("test-token", &format!("http://{}", addr));
414 let event = Event::StepFailed(StepFailedEvent {
415 run_id: Uuid::now_v7(),
416 step_id: Uuid::now_v7(),
417 step_name: "build".to_string(),
418 kind: StepKind::Shell,
419 error: "exit code 1".to_string(),
420 at: Utc::now(),
421 });
422
423 sub.handle(&event).await;
424 }
425}