systemprompt_agent/services/external_integrations/webhook/service/
delivery.rs1use 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}