Skip to main content

systemprompt_agent/services/external_integrations/webhook/service/
delivery.rs

1use super::WebhookService;
2use super::types::{WebhookConfig, WebhookDeliveryResult, WebhookStats, WebhookTestResult};
3use crate::models::external_integrations::{IntegrationError, IntegrationResult};
4use serde_json::Value;
5use std::collections::HashMap;
6use systemprompt_identifiers::WebhookEndpointId;
7use systemprompt_models::net::validate_outbound_url;
8
9impl WebhookService {
10    pub async fn send_webhook(
11        &self,
12        url: &str,
13        payload: Value,
14        config: Option<WebhookConfig>,
15    ) -> IntegrationResult<WebhookDeliveryResult> {
16        validate_outbound_url(url)
17            .map_err(|e| IntegrationError::Webhook(format!("invalid webhook url: {e}")))?;
18        let config = config.unwrap_or_else(WebhookConfig::default);
19
20        let mut request_builder = self
21            .http_client
22            .post(url)
23            .json(&payload)
24            .header("Content-Type", "application/json")
25            .header(
26                "User-Agent",
27                concat!("systemprompt.io-Webhook/", env!("CARGO_PKG_VERSION")),
28            );
29
30        for (key, value) in &config.headers {
31            request_builder = request_builder.header(key, value);
32        }
33
34        if let Some(secret) = &config.secret {
35            let signature = Self::generate_signature(secret, &payload)?;
36            request_builder = request_builder.header("X-Webhook-Signature", signature);
37        }
38
39        if let Some(timeout) = config.timeout {
40            request_builder = request_builder.timeout(timeout);
41        }
42
43        let start_time = std::time::Instant::now();
44
45        match request_builder.send().await {
46            Ok(response) => {
47                let status = response.status().as_u16();
48                let headers: HashMap<String, String> = response
49                    .headers()
50                    .iter()
51                    .map(|(k, v)| (k.to_string(), v.to_str().unwrap_or("").to_owned()))
52                    .collect();
53
54                let body = response
55                    .text()
56                    .await
57                    .unwrap_or_else(|e| format!("<error reading response: {}>", e));
58                let duration = start_time.elapsed();
59
60                Ok(WebhookDeliveryResult {
61                    success: (200..300).contains(&status),
62                    status_code: status,
63                    response_body: body,
64                    response_headers: headers,
65                    duration_ms: duration.as_millis() as u64,
66                    error: None,
67                })
68            },
69            Err(e) => {
70                let duration = start_time.elapsed();
71                Ok(WebhookDeliveryResult {
72                    success: false,
73                    status_code: 0,
74                    response_body: String::new(),
75                    response_headers: HashMap::new(),
76                    duration_ms: duration.as_millis() as u64,
77                    error: Some(e.to_string()),
78                })
79            },
80        }
81    }
82
83    pub async fn get_endpoint_stats(
84        &self,
85        endpoint_id: &WebhookEndpointId,
86    ) -> IntegrationResult<WebhookStats> {
87        let endpoint = {
88            let endpoints = self.endpoints.read().await;
89            endpoints.get(endpoint_id).cloned().ok_or_else(|| {
90                IntegrationError::Webhook(format!("Endpoint not found: {endpoint_id}"))
91            })?
92        };
93
94        Ok(WebhookStats {
95            endpoint_id: endpoint.id,
96            total_requests: 0,
97            successful_requests: 0,
98            failed_requests: 0,
99            last_request_at: None,
100            average_response_time_ms: 0,
101        })
102    }
103
104    pub async fn test_endpoint(
105        &self,
106        endpoint_id: &WebhookEndpointId,
107    ) -> IntegrationResult<WebhookTestResult> {
108        let endpoint = {
109            let endpoints = self.endpoints.read().await;
110            endpoints.get(endpoint_id).cloned().ok_or_else(|| {
111                IntegrationError::Webhook(format!("Endpoint not found: {endpoint_id}"))
112            })?
113        };
114
115        let test_payload = serde_json::json!({
116            "test": true,
117            "timestamp": chrono::Utc::now().to_rfc3339(),
118            "endpoint_id": endpoint_id
119        });
120
121        let config = WebhookConfig {
122            secret: endpoint.secret.clone(),
123            headers: endpoint.headers.clone(),
124            timeout: Some(std::time::Duration::from_secs(10)),
125        };
126
127        let result = self
128            .send_webhook(&endpoint.url, test_payload, Some(config))
129            .await?;
130
131        Ok(WebhookTestResult {
132            endpoint_id: endpoint.id,
133            success: result.success,
134            status_code: result.status_code,
135            response_time_ms: result.duration_ms,
136            error: result.error,
137        })
138    }
139}