Webhooks: send one sample event to a saved webhook
This commit is contained in:
@@ -156,15 +156,7 @@ async fn post_webhook_events(
|
||||
|
||||
// Add HMAC-SHA256 signature
|
||||
let mut headers = settings.headers.clone();
|
||||
if !settings.key.is_empty() {
|
||||
let key = hmac::Key::new(hmac::HMAC_SHA256, settings.key.as_bytes());
|
||||
let tag = hmac::sign(&key, body.as_bytes());
|
||||
|
||||
headers.insert(
|
||||
"X-Signature",
|
||||
STANDARD.encode(tag.as_ref()).parse().unwrap(),
|
||||
);
|
||||
}
|
||||
sign(&mut headers, &settings.key, &body);
|
||||
|
||||
// Send request
|
||||
let response = settings
|
||||
@@ -188,3 +180,150 @@ async fn post_webhook_events(
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
/// Adds the HMAC-SHA256 `X-Signature` a receiver checks, when the webhook has a key.
|
||||
fn sign(headers: &mut hyper::HeaderMap, key: &str, body: &str) {
|
||||
if !key.is_empty() {
|
||||
let key = hmac::Key::new(hmac::HMAC_SHA256, key.as_bytes());
|
||||
let tag = hmac::sign(&key, body.as_bytes());
|
||||
|
||||
headers.insert(
|
||||
"X-Signature",
|
||||
STANDARD.encode(tag.as_ref()).parse().unwrap(),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// inbuxa: "Send test" for a saved webhook (settings-reorg, Webhooks). One
|
||||
/// sample event, sent the way a real batch is: the same URL, headers, sign-in,
|
||||
/// signature, timeout and certificate checks. The event's type,
|
||||
/// `webhook.test`, is none the server raises, and an `X-Inbuxa-Test` header
|
||||
/// marks it, so a receiver can tell it apart. Answers the HTTP status, or why
|
||||
/// nothing came back.
|
||||
pub async fn send_test(hook: ®istry::schema::structs::WebHook) -> Result<u16, String> {
|
||||
let mut headers = hook
|
||||
.http_auth
|
||||
.build_headers(hook.http_headers.clone(), "application/json".into())
|
||||
.await
|
||||
.map_err(|err| format!("Unable to build HTTP headers: {err}"))?;
|
||||
let key = hook
|
||||
.signature_key
|
||||
.secret()
|
||||
.await
|
||||
.map_err(|err| format!("Unable to retrieve signature key: {err}"))?
|
||||
.unwrap_or_default()
|
||||
.into_owned();
|
||||
|
||||
let created = now();
|
||||
let body = serde_json::json!({
|
||||
"events": [{
|
||||
"id": format!("test-{created}"),
|
||||
"createdAt": mail_parser::DateTime::from_timestamp(created as i64).to_rfc3339(),
|
||||
"type": "webhook.test",
|
||||
"data": { "details": "A test from inbuxa Admin. Nothing happened on the server." },
|
||||
}]
|
||||
})
|
||||
.to_string();
|
||||
sign(&mut headers, &key, &body);
|
||||
headers.insert("X-Inbuxa-Test", "true".parse().unwrap());
|
||||
|
||||
let response = utils::http::http_client_builder(hook.allow_invalid_certs)
|
||||
.build()
|
||||
.map_err(|err| format!("Unable to build an HTTP client: {err}"))?
|
||||
.post(&hook.url)
|
||||
.timeout(hook.timeout.into_inner())
|
||||
.headers(headers)
|
||||
.body(body)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|err| format!("Webhook request to {} failed: {err}", hook.url))?;
|
||||
Ok(response.status().as_u16())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use registry::schema::structs::{SecretKeyOptional, SecretKeyValue, WebHook};
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
|
||||
/// One request in, the given status out; hands back what was received.
|
||||
async fn receiver(status: &'static str) -> (String, tokio::task::JoinHandle<String>) {
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let url = format!("http://{}/hook", listener.local_addr().unwrap());
|
||||
let task = tokio::spawn(async move {
|
||||
let (mut socket, _) = listener.accept().await.unwrap();
|
||||
let mut buf = Vec::new();
|
||||
let mut chunk = [0u8; 4096];
|
||||
loop {
|
||||
let n = socket.read(&mut chunk).await.unwrap();
|
||||
buf.extend_from_slice(&chunk[..n]);
|
||||
let text = String::from_utf8_lossy(&buf);
|
||||
if let Some(end) = text.find("\r\n\r\n") {
|
||||
let length = text[..end]
|
||||
.lines()
|
||||
.find_map(|l| {
|
||||
l.to_ascii_lowercase()
|
||||
.strip_prefix("content-length:")
|
||||
.map(|v| v.trim().parse::<usize>().unwrap())
|
||||
})
|
||||
.unwrap_or(0);
|
||||
if buf.len() >= end + 4 + length || n == 0 {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
socket
|
||||
.write_all(
|
||||
format!("HTTP/1.1 {status}\r\ncontent-length: 0\r\nconnection: close\r\n\r\n")
|
||||
.as_bytes(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
String::from_utf8_lossy(&buf).into_owned()
|
||||
});
|
||||
(url, task)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn send_test_signs_and_marks_the_sample() {
|
||||
let (url, task) = receiver("204 No Content").await;
|
||||
let hook = WebHook {
|
||||
url,
|
||||
enable: false,
|
||||
signature_key: SecretKeyOptional::Value(SecretKeyValue { secret: "k".into() }),
|
||||
..Default::default()
|
||||
};
|
||||
assert_eq!(send_test(&hook).await, Ok(204));
|
||||
|
||||
let request = task.await.unwrap();
|
||||
let (head, body) = request.split_once("\r\n\r\n").unwrap();
|
||||
let head = head.to_ascii_lowercase();
|
||||
assert!(head.contains("x-inbuxa-test: true"), "{head}");
|
||||
let parsed: serde_json::Value = serde_json::from_str(body).unwrap();
|
||||
assert_eq!(parsed["events"][0]["type"], "webhook.test");
|
||||
let tag = hmac::sign(&hmac::Key::new(hmac::HMAC_SHA256, b"k"), body.as_bytes());
|
||||
assert!(
|
||||
head.contains(&format!(
|
||||
"x-signature: {}",
|
||||
STANDARD.encode(tag.as_ref()).to_ascii_lowercase()
|
||||
)),
|
||||
"{head}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn send_test_reports_what_came_back() {
|
||||
let (url, _task) = receiver("403 Forbidden").await;
|
||||
let hook = WebHook {
|
||||
url,
|
||||
..Default::default()
|
||||
};
|
||||
assert_eq!(send_test(&hook).await, Ok(403));
|
||||
|
||||
let hook = WebHook {
|
||||
url: "http://127.0.0.1:9/hook".into(),
|
||||
..Default::default()
|
||||
};
|
||||
assert!(send_test(&hook).await.unwrap_err().contains("failed"));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -131,6 +131,17 @@ impl ManagementApi for Server {
|
||||
let answer = jmap::inbuxa::directory_test::test(self, &request).await?;
|
||||
Ok(JsonResponse::new(answer).no_cache().into_http_response())
|
||||
}
|
||||
// inbuxa: send one sample event to a saved webhook
|
||||
"webhook" if is_post && path.get(1).copied() == Some("test") => {
|
||||
let (_in_flight, access_token) = self.authenticate_headers(req, session).await?;
|
||||
jmap::inbuxa::webhook_test::assert_allowed(&access_token)?;
|
||||
let request = body
|
||||
.as_deref()
|
||||
.and_then(|body| serde_json::from_slice::<serde_json::Value>(body).ok())
|
||||
.unwrap_or_default();
|
||||
let answer = jmap::inbuxa::webhook_test::test(self, &request).await?;
|
||||
Ok(JsonResponse::new(answer).no_cache().into_http_response())
|
||||
}
|
||||
"account" => {
|
||||
// Authenticate request
|
||||
let (_in_flight, access_token) = self.authenticate_headers(req, session).await?;
|
||||
|
||||
@@ -18,6 +18,7 @@ pub mod ai_limits;
|
||||
pub mod log_settings;
|
||||
pub mod data_inventory;
|
||||
pub mod directory_test;
|
||||
pub mod webhook_test;
|
||||
pub mod explanation;
|
||||
pub mod protocol_policy;
|
||||
pub mod tenant_protocol_policy;
|
||||
|
||||
@@ -0,0 +1,61 @@
|
||||
/*
|
||||
* SPDX-FileCopyrightText: 2026 Coffey Labs
|
||||
*
|
||||
* SPDX-License-Identifier: AGPL-3.0-only
|
||||
*/
|
||||
|
||||
//! `POST /api/webhook/test`: send one sample event to a saved webhook
|
||||
//! (settings-reorg, Webhooks "Send test").
|
||||
//!
|
||||
//! ```json
|
||||
//! {"webhookId": "b"}
|
||||
//! ```
|
||||
//!
|
||||
//! The answer is `{"sent": true, "status": 200, "ms": 84}` when the receiver
|
||||
//! answered 2xx, `{"sent": false, "status": 403, …}` when it answered
|
||||
//! otherwise, and `{"sent": false, "error": "…"}` when nothing came back. The
|
||||
//! webhook is used as saved, even when it's off, so it can be tried before
|
||||
//! it's switched on. The request goes where the saved webhook already sends,
|
||||
//! so this gives nobody a reach they didn't have.
|
||||
//!
|
||||
//! For server-level administrators who may change webhooks.
|
||||
|
||||
use common::{Server, auth::AccessToken};
|
||||
use registry::schema::{enums::Permission, structs::WebHook};
|
||||
use serde_json::{Value, json};
|
||||
use std::{str::FromStr, time::Instant};
|
||||
use types::id::Id;
|
||||
|
||||
pub fn assert_allowed(access_token: &AccessToken) -> trc::Result<()> {
|
||||
if access_token.tenant_id().is_some() {
|
||||
return Err(trc::JmapEvent::Forbidden
|
||||
.into_err()
|
||||
.details("Webhook tests are for server-level administrators."));
|
||||
}
|
||||
access_token.enforce_permission(Permission::SysWebHookUpdate)
|
||||
}
|
||||
|
||||
pub async fn test(server: &Server, body: &Value) -> trc::Result<Value> {
|
||||
let webhook_id = body
|
||||
.get("webhookId")
|
||||
.and_then(Value::as_str)
|
||||
.and_then(|id| Id::from_str(id).ok())
|
||||
.ok_or_else(|| {
|
||||
trc::ResourceEvent::BadParameters
|
||||
.into_err()
|
||||
.details("Expected {\"webhookId\": …}")
|
||||
})?;
|
||||
let Some(hook) = server.registry().object::<WebHook>(webhook_id).await? else {
|
||||
return Ok(json!({ "sent": false, "error": "There's no such webhook. Save it first." }));
|
||||
};
|
||||
|
||||
let started = Instant::now();
|
||||
Ok(match common::telemetry::webhooks::send_test(&hook).await {
|
||||
Ok(status) => json!({
|
||||
"sent": (200..300).contains(&status),
|
||||
"status": status,
|
||||
"ms": started.elapsed().as_millis() as u64,
|
||||
}),
|
||||
Err(error) => json!({ "sent": false, "error": error }),
|
||||
})
|
||||
}
|
||||
Reference in New Issue
Block a user