Skip to main content

ironflow_engine/notify/
betterstack.rs

1//! [`BetterStackSubscriber`] -- forwards error events to BetterStack Logs.
2
3use reqwest::Client;
4use serde::Serialize;
5
6use super::retry::{RetryConfig, deliver_with_retry, is_accepted_202};
7use super::{Event, EventSubscriber, SubscriberFuture};
8
9/// Default BetterStack Logs ingestion endpoint.
10const DEFAULT_INGEST_URL: &str = "https://in.logs.betterstack.com";
11
12/// Payload sent to BetterStack Logs API.
13#[derive(Debug, Serialize)]
14struct LogPayload {
15    /// ISO-8601 timestamp.
16    dt: String,
17    /// Log level (always "error" for this subscriber).
18    level: &'static str,
19    /// Human-readable message.
20    message: String,
21    /// Structured event data.
22    event: serde_json::Value,
23}
24
25/// Subscriber that forwards error events to [BetterStack Logs](https://betterstack.com/docs/logs/).
26///
27/// Only acts on events that represent failures:
28/// - [`Event::StepFailed`]
29/// - [`Event::RunFailed`]
30///
31/// All other events are silently ignored (filtering by event type
32/// should already be done at subscription time, but this subscriber
33/// adds an extra safety check).
34///
35/// Retries failed deliveries with exponential backoff (up to 3 attempts,
36/// 5 s timeout per attempt).
37///
38/// # Examples
39///
40/// ```no_run
41/// use ironflow_engine::notify::{BetterStackSubscriber, Event, EventPublisher};
42///
43/// let mut publisher = EventPublisher::new();
44/// publisher.subscribe(
45///     BetterStackSubscriber::new("my-source-token"),
46///     &[Event::STEP_FAILED, Event::RUN_FAILED],
47/// );
48/// ```
49pub 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    /// Create a new subscriber with the given BetterStack source token.
59    ///
60    /// Uses the default ingestion endpoint (`https://in.logs.betterstack.com`)
61    /// and default [`RetryConfig`].
62    ///
63    /// # Panics
64    ///
65    /// Panics if the HTTP client cannot be built (TLS backend unavailable).
66    ///
67    /// # Examples
68    ///
69    /// ```
70    /// use ironflow_engine::notify::BetterStackSubscriber;
71    ///
72    /// let subscriber = BetterStackSubscriber::new("my-source-token");
73    /// assert_eq!(subscriber.source_token(), "my-source-token");
74    /// ```
75    pub fn new(source_token: &str) -> Self {
76        Self::with_url(source_token, DEFAULT_INGEST_URL)
77    }
78
79    /// Create a subscriber with a custom ingestion URL.
80    ///
81    /// Useful for testing or self-hosted BetterStack instances.
82    ///
83    /// # Panics
84    ///
85    /// Panics if the HTTP client cannot be built (TLS backend unavailable).
86    ///
87    /// # Examples
88    ///
89    /// ```
90    /// use ironflow_engine::notify::BetterStackSubscriber;
91    ///
92    /// let subscriber = BetterStackSubscriber::with_url(
93    ///     "my-source-token",
94    ///     "https://custom.logs.example.com",
95    /// );
96    /// assert_eq!(subscriber.ingest_url(), "https://custom.logs.example.com");
97    /// ```
98    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    /// Create a subscriber with a custom ingestion URL and retry configuration.
103    ///
104    /// # Panics
105    ///
106    /// Panics if the HTTP client cannot be built (TLS backend unavailable).
107    ///
108    /// # Examples
109    ///
110    /// ```
111    /// use ironflow_engine::notify::{BetterStackSubscriber, RetryConfig};
112    ///
113    /// let config = RetryConfig::new(
114    ///     5,
115    ///     std::time::Duration::from_secs(10),
116    ///     std::time::Duration::from_secs(1),
117    /// );
118    /// let subscriber = BetterStackSubscriber::with_url_and_retry(
119    ///     "my-source-token",
120    ///     "https://custom.logs.example.com",
121    ///     config,
122    /// );
123    /// ```
124    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    /// Returns the source token.
140    pub fn source_token(&self) -> &str {
141        &self.source_token
142    }
143
144    /// Returns the ingestion URL.
145    pub fn ingest_url(&self) -> &str {
146        &self.ingest_url
147    }
148
149    /// Build a log payload from an error event. Returns `None` for non-error events.
150    #[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            requirement: None,
369            at: Utc::now(),
370        });
371
372        assert!(BetterStackSubscriber::build_payload(&event).is_none());
373    }
374
375    #[test]
376    fn build_payload_user_signed_in_returns_none() {
377        let event = Event::UserSignedIn(UserSignedInEvent {
378            user_id: Uuid::now_v7(),
379            username: "alice".to_string(),
380            at: Utc::now(),
381        });
382
383        assert!(BetterStackSubscriber::build_payload(&event).is_none());
384    }
385
386    #[tokio::test]
387    async fn handle_ignores_non_error_events() {
388        let sub = BetterStackSubscriber::with_url("token", "http://127.0.0.1:1");
389        let event = Event::RunCreated(RunCreatedEvent {
390            run_id: Uuid::now_v7(),
391            workflow_name: "deploy".to_string(),
392            at: Utc::now(),
393        });
394        // Should return immediately without attempting HTTP
395        sub.handle(&event).await;
396    }
397
398    #[tokio::test]
399    async fn deliver_to_real_endpoint_returns_202() {
400        use axum::Router;
401        use axum::http::StatusCode;
402        use axum::routing::post;
403        use tokio::net::TcpListener;
404
405        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
406        let addr = listener.local_addr().unwrap();
407
408        let app = Router::new().route("/", post(|| async { StatusCode::ACCEPTED }));
409
410        tokio::spawn(async move {
411            axum::serve(listener, app).await.unwrap();
412        });
413
414        let sub = BetterStackSubscriber::with_url("test-token", &format!("http://{}", addr));
415        let event = Event::StepFailed(StepFailedEvent {
416            run_id: Uuid::now_v7(),
417            step_id: Uuid::now_v7(),
418            step_name: "build".to_string(),
419            kind: StepKind::Shell,
420            error: "exit code 1".to_string(),
421            at: Utc::now(),
422        });
423
424        sub.handle(&event).await;
425    }
426}