#![allow(clippy::unwrap_used, clippy::panic, clippy::print_stderr)]
use std::collections::{BTreeMap, HashMap};
use fraiseql_functions::{IngestSource, PushSource, RawDelivery, Source, Transport};
use sqlx::PgPool;
use super::{
WebhookInboundState, WebhookRoutes, WebhookSource, webhook_router, webhook_routes_check,
};
use crate::config::WebhookRouteConfig;
fn built(routes: &HashMap<String, WebhookRouteConfig>) -> WebhookRoutes {
webhook_routes_check(routes, |_| Some("configured".to_string()), false)
.expect("these fixtures are valid configurations")
}
fn lazy_pool() -> PgPool {
PgPool::connect_lazy("postgres://test:test@localhost/test").unwrap()
}
fn timestamp() -> chrono::DateTime<chrono::Utc> {
chrono::DateTime::parse_from_rfc3339("2026-07-03T12:00:00Z")
.unwrap()
.with_timezone(&chrono::Utc)
}
mod router_construction {
use super::{WebhookInboundState, WebhookRoutes, lazy_pool, webhook_router};
#[tokio::test]
async fn webhook_router_constructs() {
let state = WebhookInboundState::new(lazy_pool(), &WebhookRoutes::default(), |_| None);
let _ = webhook_router(state);
}
}
mod after_ingest_bridge {
use std::{future::Future, pin::Pin, sync::Arc};
use fraiseql_functions::host::live::QueryExecutor;
use serde_json::Value;
use super::{WebhookInboundState, WebhookRoutes, lazy_pool};
use crate::routes::after_mutation::QueryExecutorFactory;
struct MockExec;
impl QueryExecutor for MockExec {
fn execute_query(
&self,
_query: &str,
_variables: Option<&Value>,
) -> Pin<Box<dyn Future<Output = fraiseql_error::Result<Value>> + Send + '_>> {
Box::pin(async { Ok(Value::Null) })
}
}
fn factory() -> QueryExecutorFactory {
Arc::new(|_identity| Arc::new(MockExec) as Arc<dyn QueryExecutor>)
}
#[tokio::test] async fn without_a_factory_the_bridge_is_unwired() {
let state = WebhookInboundState::new(lazy_pool(), &WebhookRoutes::default(), |_| None);
assert!(state.query_executor_factory().is_none());
}
#[tokio::test]
async fn with_a_factory_the_state_carries_the_after_ingest_bridge() {
let state = WebhookInboundState::new(lazy_pool(), &WebhookRoutes::default(), |_| None)
.with_query_executor_factory(factory());
assert!(
state.query_executor_factory().is_some(),
"#594: the webhook state must carry the query bridge so after:ingest can write back"
);
}
}
#[test]
fn webhook_source_declares_push_transport() {
let source = WebhookSource::new("stripe", "partner-a");
assert_eq!(
source.source(),
IngestSource::Webhook {
provider: "stripe".to_string(),
}
);
assert_eq!(source.transport(), Transport::Push);
}
#[test]
fn webhook_source_normalizes_delivery_and_carries_payload() {
let source = WebhookSource::new("stripe", "partner-a");
let payload = serde_json::json!({ "id": "evt_1", "type": "charge.succeeded" });
let mut headers = BTreeMap::new();
headers.insert("webhook-id".to_string(), "evt_1".to_string());
let raw = RawDelivery {
event_id: "evt_1",
event_type: "charge.succeeded",
payload: &payload,
headers: &headers,
received_at: timestamp(),
};
let message = source.normalize(&raw).unwrap();
assert_eq!(
message.source,
IngestSource::Webhook {
provider: "stripe".to_string(),
}
);
assert_eq!(
message.idempotency_key, "9:partner-a:evt_1",
"#1046: the spine dedup key is namespaced by the receiving route — `source` \
cannot carry it, being the after:ingest discriminant"
);
assert_eq!(message.subject.as_deref(), Some("charge.succeeded"));
assert_eq!(message.payload.as_ref(), Some(&payload));
assert_eq!(message.headers.get("webhook-id").map(String::as_str), Some("evt_1"));
assert_eq!(message.trigger_type(), "after:ingest:webhook:stripe");
}
#[test]
fn two_routes_on_one_provider_share_a_trigger_but_not_a_dedup_key() {
let payload = serde_json::json!({ "id": "1001" });
let headers = BTreeMap::new();
let raw = RawDelivery {
event_id: "1001",
event_type: "order.created",
payload: &payload,
headers: &headers,
received_at: timestamp(),
};
let a = WebhookSource::new("hmac-sha256", "partner-a").normalize(&raw).unwrap();
let b = WebhookSource::new("hmac-sha256", "partner-b").normalize(&raw).unwrap();
assert_ne!(
a.idempotency_key, b.idempotency_key,
"#1046: each sender numbers its own events, so one sender's `1001` must not \
deduplicate against another's"
);
assert_eq!(
a.trigger_type(),
b.trigger_type(),
"#1046: the after:ingest discriminant stays `webhook:<provider>` — changing it \
would break every declared trigger"
);
}
#[test]
fn a_colon_in_a_route_segment_cannot_forge_another_routes_dedup_key() {
let key = |route: &str, id: &str| {
let payload = serde_json::json!({ "id": id });
let headers = BTreeMap::new();
let raw = RawDelivery {
event_id: id,
event_type: "order.created",
payload: &payload,
headers: &headers,
received_at: timestamp(),
};
WebhookSource::new("hmac-sha256", route)
.normalize(&raw)
.unwrap()
.idempotency_key
};
assert_ne!(
key("a", "b:1"),
key("a:b", "1"),
"#1046: the route/event-id join must be unambiguous for every segment — the \
event id is chosen by whoever posts, so an ambiguous join hands one route's \
sender a way to suppress another route's message"
);
}
mod form_bodies {
use axum::http::{HeaderMap, HeaderValue, header::CONTENT_TYPE};
use super::super::{form_to_json, is_form_encoded};
fn headers(content_type: &str) -> HeaderMap {
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, HeaderValue::from_str(content_type).unwrap());
headers
}
#[test]
fn a_form_content_type_is_recognised_with_or_without_parameters() {
assert!(is_form_encoded(&headers("application/x-www-form-urlencoded")));
assert!(
is_form_encoded(&headers("application/x-www-form-urlencoded; charset=UTF-8")),
"a charset parameter is legal and must not change the reading"
);
assert!(
is_form_encoded(&headers("Application/X-WWW-Form-Urlencoded")),
"media types are case-insensitive"
);
}
#[test]
fn other_content_types_still_take_the_json_path() {
assert!(!is_form_encoded(&headers("application/json")));
assert!(!is_form_encoded(&headers("text/plain")));
assert!(
!is_form_encoded(&HeaderMap::new()),
"no declared type is not a form body — the sender's declaration decides"
);
}
#[test]
fn a_twilio_shaped_body_becomes_a_flat_object_with_values_decoded() {
assert_eq!(
form_to_json(b"CallSid=CA123&From=%2B15550001111&Body=hi+there"),
serde_json::json!({
"CallSid": "CA123",
"From": "+15550001111",
"Body": "hi there",
}),
"percent escapes and `+` must be decoded, or an after:ingest function \
reads a phone number as `%2B1…`"
);
}
#[test]
fn a_repeated_key_keeps_every_value_in_wire_order() {
assert_eq!(
form_to_json(b"Tag=a&CallSid=CA123&Tag=b"),
serde_json::json!({ "Tag": ["a", "b"], "CallSid": "CA123" }),
"form encoding permits repeats; collapsing them to one value would drop \
data with nothing on the wire to show for it"
);
}
#[test]
fn degenerate_bodies_parse_rather_than_fail() {
assert_eq!(form_to_json(b""), serde_json::json!({}));
assert_eq!(
form_to_json(b"novalue"),
serde_json::json!({ "novalue": "" }),
"form encoding has no invalid syntax to reject — a bare key is an empty value"
);
}
}
#[test]
fn webhook_source_rejects_delivery_without_event_id() {
let source = WebhookSource::new("stripe", "partner-a");
let payload = serde_json::json!({});
let headers = BTreeMap::new();
let raw = RawDelivery {
event_id: "",
event_type: "x",
payload: &payload,
headers: &headers,
received_at: timestamp(),
};
assert!(source.normalize(&raw).is_err());
}
mod error_body_sanitization {
use axum::{
Router,
body::Body,
http::{Request, StatusCode},
};
use tower::ServiceExt as _;
use super::{HashMap, WebhookInboundState, built, lazy_pool, webhook_router};
use crate::config::WebhookRouteConfig;
fn router_with_one_route() -> Router {
let mut routes = HashMap::new();
routes.insert(
"hooks".to_string(),
WebhookRouteConfig {
secret_env: Some("TEST_WEBHOOK_SECRET".to_string()),
provider: "hmac-sha256".to_string(),
path: None,
public_url: None,
credential: None,
encoding: None,
prefix: None,
..Default::default()
},
);
let state = WebhookInboundState::new(lazy_pool(), &built(&routes), |_| {
Some("s3cret-value".to_string())
});
webhook_router(state)
}
#[tokio::test]
async fn a_forged_signature_is_refused_without_echoing_the_internal_reason() {
let request = Request::builder()
.method("POST")
.uri("/webhooks/hooks")
.header("X-Signature", "deadbeef")
.body(Body::from(r#"{"id":"evt_1"}"#))
.unwrap();
let response = router_with_one_route().oneshot(request).await.unwrap();
assert_eq!(
response.status(),
StatusCode::UNAUTHORIZED,
"a forged signature is the sender's fault, so the status stays 401"
);
let bytes = axum::body::to_bytes(response.into_body(), usize::MAX).await.unwrap();
let text = String::from_utf8(bytes.to_vec()).unwrap();
let body: serde_json::Value = serde_json::from_str(&text).unwrap();
assert_eq!(
body["error_description"], "Authentication failed",
"the sanitizer's flat description is what an unauthenticated caller may see; \
got {text}"
);
assert_eq!(
body["error"], "authentication_error",
"the route must render through FraiseQLError's IntoResponse, not its own shape; \
got {text}"
);
assert!(
!text.contains("signature mismatch"),
"the verifier's internal reason must not reach the sender; got {text}"
);
}
}
mod key_material_is_not_the_senders_fault {
use axum::{
Router,
body::Body,
http::{Request, StatusCode},
};
use tower::ServiceExt as _;
use super::{HashMap, WebhookInboundState, built, lazy_pool, webhook_router};
use crate::config::WebhookRouteConfig;
const VALID_ED25519_PUBLIC_KEY: &str =
"d75a980182b10ab7d54bfed3c964073a0ee172f3daa62325af021a68f707511a";
fn discord_router(configured_key: &'static str) -> Router {
let mut routes = HashMap::new();
routes.insert(
"discord".to_string(),
WebhookRouteConfig {
secret_env: Some("TEST_DISCORD_KEY".to_string()),
provider: "discord".to_string(),
path: None,
public_url: None,
credential: None,
encoding: None,
prefix: None,
..Default::default()
},
);
let state = WebhookInboundState::new(lazy_pool(), &built(&routes), |_| {
Some(configured_key.to_string())
});
webhook_router(state)
}
fn request(signature: &str) -> Request<Body> {
Request::builder()
.method("POST")
.uri("/webhooks/discord")
.header("X-Signature-Ed25519", signature)
.header("X-Signature-Timestamp", chrono::Utc::now().timestamp().to_string())
.body(Body::from(r#"{"id":"evt_1","type":1}"#))
.unwrap()
}
#[tokio::test]
async fn a_misconfigured_signing_key_is_the_servers_fault_not_a_401() {
let response = discord_router("not-hex!").oneshot(request("aabb")).await.unwrap();
assert_eq!(
response.status(),
StatusCode::INTERNAL_SERVER_ERROR,
"unparseable configured key material is a server-side misconfiguration; \
answering 401 blames the sender and invites the provider to disable the endpoint"
);
let bytes = axum::body::to_bytes(response.into_body(), usize::MAX).await.unwrap();
let text = String::from_utf8(bytes.to_vec()).unwrap();
assert!(
!text.contains("hex") && !text.contains("public key"),
"the key-parsing detail belongs in the operator's log, not the response; got {text}"
);
}
#[tokio::test]
async fn a_malformed_sender_signature_is_still_a_401() {
let response =
discord_router(VALID_ED25519_PUBLIC_KEY).oneshot(request("zz")).await.unwrap();
assert_eq!(
response.status(),
StatusCode::UNAUTHORIZED,
"an unparseable signature header is the sender's fault and must stay 401, \
or an anonymous caller can produce a 5xx on demand"
);
}
}
mod empty_secret_is_not_configured {
use super::{super::webhook_routes_check, HashMap, WebhookInboundState, built, lazy_pool};
use crate::config::WebhookRouteConfig;
fn one_route() -> HashMap<String, WebhookRouteConfig> {
let mut routes = HashMap::new();
routes.insert(
"hooks".to_string(),
WebhookRouteConfig {
secret_env: Some("TEST_WEBHOOK_SECRET".to_string()),
provider: "hmac-sha256".to_string(),
path: None,
public_url: None,
credential: None,
encoding: None,
prefix: None,
..Default::default()
},
);
routes
}
#[test]
fn production_refuses_a_route_whose_secret_env_is_empty() {
let error = webhook_routes_check(&one_route(), |_| Some(String::new()), true)
.expect_err("an empty signing secret cannot verify a delivery, so it must not boot");
let message = error.to_string();
assert!(
message.contains("TEST_WEBHOOK_SECRET"),
"the refusal must name the variable the operator has to fix; got {message}"
);
}
#[test]
fn a_non_empty_secret_still_boots() {
assert!(
webhook_routes_check(&one_route(), |_| Some("s3cret".to_string()), true).is_ok(),
"a configured, non-empty secret is exactly the case that must boot"
);
}
#[tokio::test]
async fn a_route_with_an_empty_secret_is_not_mounted() {
let state =
WebhookInboundState::new(lazy_pool(), &built(&one_route()), |_| Some(String::new()));
assert!(
state.routes.is_empty(),
"a route whose secret is empty must be skipped, not mounted"
);
}
}
mod colliding_path_segments {
use super::{super::webhook_routes_check, HashMap, WebhookInboundState, built, lazy_pool};
use crate::config::WebhookRouteConfig;
fn route(provider: &str, secret_env: &str, path: Option<&str>) -> WebhookRouteConfig {
WebhookRouteConfig {
secret_env: Some(secret_env.to_string()),
provider: provider.to_string(),
path: path.map(str::to_string),
public_url: None,
credential: None,
encoding: None,
prefix: None,
..Default::default()
}
}
fn both_overridden() -> HashMap<String, WebhookRouteConfig> {
let mut routes = HashMap::new();
routes.insert("stripe_live".to_string(), route("stripe", "STRIPE_LIVE", Some("stripe")));
routes.insert("stripe_test".to_string(), route("stripe", "STRIPE_TEST", Some("stripe")));
routes
}
fn override_collides_with_a_name() -> HashMap<String, WebhookRouteConfig> {
let mut routes = HashMap::new();
routes.insert("github".to_string(), route("github", "GH_SECRET", None));
routes.insert("github_v2".to_string(), route("github", "GH_SECRET_2", Some("github")));
routes
}
#[test]
fn two_routes_on_one_segment_are_refused_at_boot() {
let error = webhook_routes_check(&both_overridden(), |_| Some("s".to_string()), true)
.expect_err("a shadowed route can never receive a delivery, so this must not boot");
let message = error.to_string();
assert!(
message.contains("stripe"),
"the refusal must name the colliding segment; got {message}"
);
assert!(
message.contains("stripe_live") && message.contains("stripe_test"),
"the refusal must name both colliding routes; got {message}"
);
}
#[test]
fn an_override_colliding_with_another_routes_name_is_refused() {
let error =
webhook_routes_check(&override_collides_with_a_name(), |_| Some("s".to_string()), true)
.expect_err("a collision does not require two explicit overrides");
let message = error.to_string();
assert!(
message.contains("github_v2") && message.contains("github"),
"the refusal must name both colliding routes; got {message}"
);
}
#[test]
fn the_refusal_is_identical_across_runs() {
let first = webhook_routes_check(&both_overridden(), |_| Some("s".to_string()), true)
.unwrap_err()
.to_string();
for _ in 0..64 {
let again = webhook_routes_check(&both_overridden(), |_| Some("s".to_string()), true)
.unwrap_err()
.to_string();
assert_eq!(again, first, "the collision refusal must not vary between runs");
}
}
#[tokio::test] async fn distinct_segments_still_boot() {
let mut routes = HashMap::new();
routes.insert("partner_a".to_string(), route("hmac-sha256", "A_SECRET", None));
routes.insert("partner_b".to_string(), route("hmac-sha256", "B_SECRET", None));
assert!(
webhook_routes_check(&routes, |_| Some("s".to_string()), true).is_ok(),
"distinct segments are not a collision, even on a shared provider"
);
let state =
WebhookInboundState::new(lazy_pool(), &built(&routes), |_| Some("s".to_string()));
assert_eq!(state.routes.len(), 2, "both routes must mount");
}
}
mod signed_dedup_key {
use std::collections::BTreeMap;
use super::super::{extract_event_id, extract_event_type};
fn headers(pairs: &[(&str, &str)]) -> BTreeMap<String, String> {
pairs.iter().map(|(k, v)| ((*k).to_string(), (*v).to_string())).collect()
}
#[test]
fn an_identical_body_always_yields_an_identical_key() {
let body = br#"{"action":"opened","number":1}"#;
let payload = serde_json::from_slice(body).unwrap();
assert_eq!(extract_event_id(&payload, body), extract_event_id(&payload, body));
}
#[test]
fn signed_payload_id_keys_the_delivery_when_present() {
let body = br#"{"id":"evt_1","type":"charge.succeeded"}"#;
let payload = serde_json::from_slice(body).unwrap();
assert_eq!(
extract_event_id(&payload, body),
"evt_1",
"a signed top-level id is the natural key and must be used verbatim"
);
}
#[test]
fn distinct_bodies_get_distinct_keys() {
let a = br#"{"action":"opened","number":1}"#;
let b = br#"{"action":"opened","number":2}"#;
let pa = serde_json::from_slice(a).unwrap();
let pb = serde_json::from_slice(b).unwrap();
assert_ne!(
extract_event_id(&pa, a),
extract_event_id(&pb, b),
"genuinely different deliveries must not collapse into one key"
);
}
#[test]
fn body_hash_key_is_sha256_not_the_unstable_default_hasher() {
let body = b"{}";
let payload = serde_json::from_slice(body).unwrap();
assert_eq!(
extract_event_id(&payload, body),
"body:44136fa355b3678a1146ad16f7e8649e94fb4fc21fe77e8310c060f61caaff8a"
);
}
#[test]
fn signed_payload_type_beats_an_injected_header() {
let payload = serde_json::json!({ "type": "charge.succeeded" });
let h = headers(&[("x-github-event", "push")]);
assert_eq!(
extract_event_type(&payload, &h),
"charge.succeeded",
"an unsigned header must not relabel a delivery whose signed body states its type"
);
}
#[test]
fn github_event_header_still_used_when_the_body_carries_no_type() {
let payload = serde_json::json!({ "action": "opened" });
let h = headers(&[("x-github-event", "pull_request")]);
assert_eq!(extract_event_type(&payload, &h), "pull_request");
}
}
mod spine_handler_disposition {
#![allow(clippy::panic, clippy::print_stderr)]
use fraiseql_functions::{InboundMessage, IngestSource};
use fraiseql_webhooks::EventHandler as _;
use sqlx::PgPool;
use crate::inbound::{
spine::{Emitted, PostgresInboundSpine, emit_in_tx},
webhook::SpineEventHandler,
};
async fn pool() -> Option<(PgPool, fraiseql_test_support::Service)> {
let svc = fraiseql_test_support::postgres().await?;
let pool = PgPool::connect(svc.url()).await.unwrap();
Some((pool, svc))
}
fn message(key: &str) -> InboundMessage {
InboundMessage::new(
IngestSource::Webhook {
provider: "stripe".to_string(),
},
key,
chrono::DateTime::parse_from_rfc3339("2026-07-03T12:00:00Z")
.unwrap()
.with_timezone(&chrono::Utc),
)
}
#[tokio::test]
async fn a_key_the_spine_already_owns_is_reported_duplicate() {
let Some((pool, _svc)) = pool().await else {
eprintln!("SKIP a_key_the_spine_already_owns_is_reported_duplicate: no postgres");
return;
};
let spine = PostgresInboundSpine::new(pool.clone());
spine.init().await.unwrap();
let key = uuid::Uuid::new_v4().to_string();
let msg = message(&key);
let mut tx = pool.begin().await.unwrap();
assert!(
matches!(emit_in_tx(&mut tx, &msg).await.unwrap(), Emitted::New(_)),
"the first emit must record the message"
);
tx.commit().await.unwrap();
let mut tx = pool.begin().await.unwrap();
let handled = SpineEventHandler
.handle("stripe", serde_json::to_value(&msg).unwrap(), &mut tx)
.await
.expect("a duplicate is an answer, not a failure");
tx.commit().await.unwrap();
assert!(
matches!(handled, fraiseql_webhooks::Handled::Duplicate),
"the spine refused the write as a duplicate; the handler reported {handled:?}, \
which the route renders as `processed` and dispatches after:ingest on"
);
}
#[tokio::test]
async fn a_fresh_key_is_reported_recorded() {
let Some((pool, _svc)) = pool().await else {
eprintln!("SKIP a_fresh_key_is_reported_recorded: no postgres");
return;
};
let spine = PostgresInboundSpine::new(pool.clone());
spine.init().await.unwrap();
let msg = message(&uuid::Uuid::new_v4().to_string());
let mut tx = pool.begin().await.unwrap();
let handled = SpineEventHandler
.handle("stripe", serde_json::to_value(&msg).unwrap(), &mut tx)
.await
.unwrap();
tx.commit().await.unwrap();
let fraiseql_webhooks::Handled::Recorded(value) = handled else {
panic!("a fresh key must be recorded, got {handled:?}");
};
assert_eq!(
value.get("idempotency_key").and_then(serde_json::Value::as_str),
Some(msg.idempotency_key.as_str()),
"the normalized message is what the route dispatches on: {value}"
);
}
}
mod the_event_is_what_the_scheme_authenticated {
use std::sync::Arc;
use axum::{
body::Body,
http::{Request, StatusCode},
};
use fraiseql_webhooks::{
InboundRequest, PostgresIdempotencyStore, SignatureError, SignatureVerifier,
StaticSecretProvider, Verified, WebhookPipeline,
};
use hmac::{Hmac, KeyInit, Mac as _};
use serde_json::{Value, json};
use sha2::Sha256;
use sqlx::{PgPool, postgres::PgPoolOptions};
use tower::ServiceExt as _;
use super::{BTreeMap, WebhookInboundState, webhook_router};
use crate::inbound::webhook::{BuiltRoute, SpineEventHandler};
const SECRET: &str = "whsec_1321_envelope";
const SECRET_ENV: &str = "FRAISEQL_TEST_ENVELOPE_SECRET";
const SEGMENT: &str = "envelope";
const SIGNATURE_HEADER: &str = "X-Envelope-Signature";
struct EnvelopeScheme;
impl SignatureVerifier for EnvelopeScheme {
fn name(&self) -> &'static str {
"test-envelope"
}
fn verify(
&self,
request: &InboundRequest<'_>,
secret: &str,
) -> Result<Verified, SignatureError> {
let signature = request.require_header(SIGNATURE_HEADER)?;
let payload = request.body();
let mut mac = Hmac::<Sha256>::new_from_slice(secret.as_bytes()).map_err(
|error: hmac::digest::InvalidLength| SignatureError::KeyMaterial(error.to_string()),
)?;
mac.update(payload);
let expected = hex::encode(mac.finalize().into_bytes());
if signature != expected {
return Err(SignatureError::Mismatch);
}
let envelope: Value =
serde_json::from_slice(payload).map_err(|_| SignatureError::InvalidFormat)?;
let event = envelope.get("event").ok_or(SignatureError::InvalidFormat)?;
let field = |name: &str| {
event
.get(name)
.and_then(Value::as_str)
.map(str::to_string)
.ok_or(SignatureError::InvalidFormat)
};
Ok(Verified::Event {
id: field("id")?,
event_type: field("type")?,
payload: event.clone(),
})
}
}
fn state(pool: PgPool) -> WebhookInboundState {
let mut routes = BTreeMap::new();
routes.insert(
SEGMENT.to_string(),
BuiltRoute {
name: SEGMENT.to_string(),
provider: "test-envelope".to_string(),
scheme: Arc::new(EnvelopeScheme),
secret_name: Some(SECRET_ENV.to_string()),
public_url: None,
},
);
let secrets =
StaticSecretProvider::new().with_secret(SECRET_ENV.to_string(), SECRET.to_string());
let store = PostgresIdempotencyStore::new(pool.clone());
WebhookInboundState {
pipeline: Arc::new(WebhookPipeline::new(
pool,
secrets,
store,
SpineEventHandler,
)),
routes: Arc::new(routes),
hooks: None,
query_executor_factory: None,
}
}
async fn setup() -> Option<PgPool> {
let url = fraiseql_test_support::try_database_url()?;
let pool = PgPoolOptions::new().max_connections(4).connect(&url).await.unwrap();
PostgresIdempotencyStore::new(pool.clone()).init().await.unwrap();
WebhookInboundState::init_spine(&pool).await.unwrap();
Some(pool)
}
fn unique() -> String {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
.to_string()
}
#[tokio::test]
async fn the_ledger_and_the_spine_record_the_authenticated_event() {
let Some(pool) = setup().await else {
eprintln!(
"skipping the_ledger_and_the_spine_record_the_authenticated_event: \
DATABASE_URL unset"
);
return;
};
let router = webhook_router(state(pool.clone()));
let run = unique();
let authenticated_id = format!("authenticated-{run}");
let envelope = json!({
"id": format!("envelope-{run}"),
"type": "envelope.received",
"note": "this field belongs to the envelope and to nothing else",
"event": {
"id": authenticated_id,
"type": "authenticated.event",
"amount": 4242,
},
});
let body = serde_json::to_vec(&envelope).unwrap();
let mut mac = Hmac::<Sha256>::new_from_slice(SECRET.as_bytes()).unwrap();
mac.update(&body);
let signature = hex::encode(mac.finalize().into_bytes());
let request = Request::builder()
.method("POST")
.uri(format!("/webhooks/{SEGMENT}"))
.header("X-Envelope-Signature", signature)
.header("content-type", "application/json")
.body(Body::from(body))
.unwrap();
let response = router.oneshot(request).await.unwrap();
assert_eq!(response.status(), StatusCode::OK, "the delivery is genuine and must be taken");
let (ledger_event_id, ledger_event_type): (String, String) = sqlx::query_as(
"SELECT event_id, event_type FROM webhooks.tb_inbound_delivery WHERE route = $1 \
ORDER BY pk_inbound_delivery DESC LIMIT 1",
)
.bind(SEGMENT)
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(
ledger_event_id, authenticated_id,
"#1321: the replay defence must be keyed on the id the scheme \
authenticated. Keyed on the envelope's own id, a sender that controls \
the envelope controls the dedup key of an event it did not sign (#751)."
);
assert_eq!(
ledger_event_type, "authenticated.event",
"the claim records the authenticated event's type, not the envelope's"
);
let stored: Value = sqlx::query_scalar(
"SELECT payload FROM _fraiseql_inbound_message WHERE idempotency_key = $1",
)
.bind(format!("{}:{SEGMENT}:{authenticated_id}", SEGMENT.len()))
.fetch_one(&pool)
.await
.unwrap_or_else(|error| {
panic!(
"#1321: the spine row must be keyed on the authenticated id \
({authenticated_id}); none found: {error}"
)
});
assert_eq!(
stored.get("subject").and_then(Value::as_str),
Some("authenticated.event"),
"the message subject is the authenticated event's type; got: {stored}"
);
assert_eq!(
stored.get("payload"),
Some(&json!({
"id": authenticated_id,
"type": "authenticated.event",
"amount": 4242,
})),
"the durable payload is what the scheme authenticated — the envelope's \
own `note` field must not be in it; got: {stored}"
);
}
}
mod the_credential_need_not_be_a_header_and_the_body_need_not_be_json {
use std::sync::Arc;
use axum::{
body::Body,
http::{Request, StatusCode},
};
use fraiseql_webhooks::{
CredentialLocation, InboundRequest, PostgresIdempotencyStore, SignatureError,
SignatureVerifier, StaticSecretProvider, Verified, WebhookPipeline,
};
use serde_json::{Value, json};
use sqlx::{PgPool, postgres::PgPoolOptions};
use tower::ServiceExt as _;
use super::{BTreeMap, WebhookInboundState, webhook_router};
use crate::inbound::webhook::{BuiltRoute, SpineEventHandler};
const SECRET: &str = "tok_1321_body";
const SECRET_ENV: &str = "FRAISEQL_TEST_BODY_CREDENTIAL_SECRET";
struct BodyFieldToken;
impl SignatureVerifier for BodyFieldToken {
fn name(&self) -> &'static str {
"test-body-field-token"
}
fn verify(
&self,
request: &InboundRequest<'_>,
secret: &str,
) -> Result<Verified, SignatureError> {
let presented =
request.credential(&CredentialLocation::BodyField("token".to_string()))?;
if presented == secret {
Ok(Verified::Body)
} else {
Err(SignatureError::Mismatch)
}
}
}
struct WholeBodyToken;
impl SignatureVerifier for WholeBodyToken {
fn name(&self) -> &'static str {
"test-whole-body-token"
}
fn verify(
&self,
request: &InboundRequest<'_>,
secret: &str,
) -> Result<Verified, SignatureError> {
let token = request.credential(&CredentialLocation::Body)?;
let (claims_hex, presented) =
token.split_once('.').ok_or(SignatureError::InvalidFormat)?;
if presented != secret {
return Err(SignatureError::Mismatch);
}
let claims: Value = serde_json::from_slice(
&hex::decode(claims_hex).map_err(|_| SignatureError::InvalidFormat)?,
)
.map_err(|_| SignatureError::InvalidFormat)?;
let field = |name: &str| {
claims
.get(name)
.and_then(Value::as_str)
.map(str::to_string)
.ok_or(SignatureError::InvalidFormat)
};
Ok(Verified::Event {
id: field("id")?,
event_type: field("type")?,
payload: claims.clone(),
})
}
}
fn state(
pool: PgPool,
segment: &str,
scheme: Arc<dyn SignatureVerifier>,
) -> WebhookInboundState {
let mut routes = BTreeMap::new();
routes.insert(
segment.to_string(),
BuiltRoute {
name: segment.to_string(),
provider: scheme.name().to_string(),
scheme,
secret_name: Some(SECRET_ENV.to_string()),
public_url: None,
},
);
let secrets =
StaticSecretProvider::new().with_secret(SECRET_ENV.to_string(), SECRET.to_string());
let store = PostgresIdempotencyStore::new(pool.clone());
WebhookInboundState {
pipeline: Arc::new(WebhookPipeline::new(
pool,
secrets,
store,
SpineEventHandler,
)),
routes: Arc::new(routes),
hooks: None,
query_executor_factory: None,
}
}
async fn setup() -> Option<PgPool> {
let url = fraiseql_test_support::try_database_url()?;
let pool = PgPoolOptions::new().max_connections(4).connect(&url).await.unwrap();
PostgresIdempotencyStore::new(pool.clone()).init().await.unwrap();
WebhookInboundState::init_spine(&pool).await.unwrap();
Some(pool)
}
fn unique() -> String {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
.to_string()
}
async fn ledger_event_id(pool: &PgPool, route: &str) -> Option<String> {
sqlx::query_scalar(
"SELECT event_id FROM webhooks.tb_inbound_delivery WHERE route = $1 \
ORDER BY pk_inbound_delivery DESC LIMIT 1",
)
.bind(route)
.fetch_optional(pool)
.await
.unwrap()
}
#[tokio::test]
async fn a_credential_in_the_body_reaches_verification() {
let Some(pool) = setup().await else {
eprintln!("skipping a_credential_in_the_body_reaches_verification: DATABASE_URL unset");
return;
};
let segment = "body-field";
let router = webhook_router(state(pool.clone(), segment, Arc::new(BodyFieldToken)));
let id = format!("body-field-{}", unique());
let body = json!({ "token": SECRET, "id": id, "type": "invoice.created" });
let request = Request::builder()
.method("POST")
.uri(format!("/webhooks/{segment}"))
.header("content-type", "application/json")
.body(Body::from(serde_json::to_vec(&body).unwrap()))
.unwrap();
let response = router.oneshot(request).await.unwrap();
let status = response.status();
let bytes = axum::body::to_bytes(response.into_body(), 64 * 1024).await.unwrap();
let text = String::from_utf8_lossy(&bytes).into_owned();
assert_eq!(
status,
StatusCode::OK,
"#1321: a scheme whose credential is a body field must reach verification; \
body: {text}"
);
assert!(text.contains("processed"), "expected processed, got: {text}");
assert_eq!(
ledger_event_id(&pool, segment).await.as_deref(),
Some(id.as_str()),
"the body is still the event for this scheme, so the claim is keyed on its id"
);
}
#[tokio::test]
async fn a_body_that_is_not_json_reaches_verification() {
let Some(pool) = setup().await else {
eprintln!("skipping a_body_that_is_not_json_reaches_verification: DATABASE_URL unset");
return;
};
let segment = "whole-body";
let router = webhook_router(state(pool.clone(), segment, Arc::new(WholeBodyToken)));
let id = format!("whole-body-{}", unique());
let claims = json!({ "id": id, "type": "session.created", "sub": "usr_42" });
let token = format!("{}.{SECRET}", hex::encode(serde_json::to_vec(&claims).unwrap()));
let request = Request::builder()
.method("POST")
.uri(format!("/webhooks/{segment}"))
.header("content-type", "application/jwt")
.body(Body::from(token))
.unwrap();
let response = router.oneshot(request).await.unwrap();
let status = response.status();
let bytes = axum::body::to_bytes(response.into_body(), 64 * 1024).await.unwrap();
let text = String::from_utf8_lossy(&bytes).into_owned();
assert_eq!(
status,
StatusCode::OK,
"#1321: a scheme whose body IS its credential must reach verification — the \
body-must-be-JSON refusal ran before the scheme and made the family \
unreachable; body: {text}"
);
assert!(text.contains("processed"), "expected processed, got: {text}");
assert_eq!(
ledger_event_id(&pool, segment).await.as_deref(),
Some(id.as_str()),
"the event comes out of the authenticated token, not out of the body bytes"
);
}
#[tokio::test]
async fn a_wrong_body_credential_is_still_refused() {
let Some(pool) = setup().await else {
eprintln!("skipping a_wrong_body_credential_is_still_refused: DATABASE_URL unset");
return;
};
let segment = "body-field-forged";
let router = webhook_router(state(pool.clone(), segment, Arc::new(BodyFieldToken)));
let body = json!({ "token": "not-the-secret", "id": unique(), "type": "invoice.created" });
let request = Request::builder()
.method("POST")
.uri(format!("/webhooks/{segment}"))
.header("content-type", "application/json")
.body(Body::from(serde_json::to_vec(&body).unwrap()))
.unwrap();
let response = router.oneshot(request).await.unwrap();
assert_eq!(response.status(), StatusCode::UNAUTHORIZED, "a wrong token must 401");
assert_eq!(
ledger_event_id(&pool, segment).await,
None,
"a delivery that failed verification must claim nothing"
);
}
}
mod token_scheme_boot {
use super::{super::webhook_routes_check, HashMap, WebhookInboundState, lazy_pool};
use crate::config::WebhookRouteConfig;
const JWKS: &str = "https://tenant.hanko.io/.well-known/jwks.json";
fn route(config: WebhookRouteConfig) -> HashMap<String, WebhookRouteConfig> {
let mut routes = HashMap::new();
routes.insert("idp".to_string(), config);
routes
}
fn hanko() -> WebhookRouteConfig {
WebhookRouteConfig {
provider: "hanko".to_string(),
jwks_uri: Some(JWKS.to_string()),
audience: Some("my-app".to_string()),
..Default::default()
}
}
#[tokio::test]
async fn a_token_route_with_no_secret_boots_and_is_mounted() {
for is_production in [false, true] {
let built = webhook_routes_check(&route(hanko()), |_| None, is_production)
.unwrap_or_else(|error| {
panic!(
"a hanko route needs no secret_env, so it must boot with \
is_production={is_production}: {error}"
)
});
let state = WebhookInboundState::new(lazy_pool(), &built, |_| None);
assert_eq!(
state.mounted_segments(),
vec!["idp".to_string()],
"and it must be MOUNTED, not skipped. The development skip is keyed on the \
scheme refusing this route's key material — keyed on `the secret is \
missing` it would swallow every jwt-jwks route (#787)"
);
}
}
#[test]
fn a_token_route_carrying_a_secret_is_refused_in_every_environment() {
let mut config = hanko();
config.secret_env = Some("HANKO_SECRET".to_string());
for is_production in [false, true] {
let error = webhook_routes_check(
&route(config.clone()),
|_| Some("s3cret".to_string()),
is_production,
)
.err()
.unwrap_or_else(|| {
panic!("unused key material must be refused (is_production={is_production})")
});
let message = error.to_string();
assert!(
message.contains("secret_env"),
"and the refusal must name the key to remove: {message}"
);
}
}
#[test]
fn a_token_route_with_no_key_set_is_refused_in_every_environment() {
let mut config = hanko();
config.jwks_uri = None;
for is_production in [false, true] {
let error = webhook_routes_check(&route(config.clone()), |_| None, is_production)
.err()
.unwrap_or_else(|| panic!("a route with nowhere to look a key up must be refused"));
assert!(
error.to_string().contains("jwks_uri"),
"and must name the key an operator has to set: {error}"
);
}
}
#[test]
fn a_key_set_that_could_never_be_fetched_is_refused_at_boot() {
for uri in ["http://tenant.hanko.io/jwks", "not a url", "ftp://idp/jwks"] {
let mut config = hanko();
config.jwks_uri = Some(uri.to_string());
let error = webhook_routes_check(&route(config), |_| None, true)
.err()
.unwrap_or_else(|| panic!("{uri} must be refused at boot"));
let _ = error;
}
let mut config = hanko();
config.jwks_uri = Some("http://localhost:9999/jwks".to_string());
webhook_routes_check(&route(config), |_| None, true)
.expect("a loopback key set is how a development IdP is reached");
}
#[test]
fn a_preset_refuses_the_claim_keys_it_fixes_itself() {
type Set = fn(&mut WebhookRouteConfig);
let cases: [(&str, Set); 5] = [
("event_type_claim", |c| c.event_type_claim = Some("evt".to_string())),
("payload_claim", |c| c.payload_claim = Some("data".to_string())),
("id_claim", |c| c.id_claim = Some("jti".to_string())),
("body_hash_claim", |c| c.body_hash_claim = Some("hash".to_string())),
("credential", |c| c.credential = Some("body".parse().expect("a location"))),
];
for (key, config) in cases {
let mut route_config = hanko();
config(&mut route_config);
let error =
webhook_routes_check(&route(route_config), |_| None, true).err().unwrap_or_else(
|| panic!("`{key}` is fixed by the hanko preset, so it must be refused"),
);
let message = error.to_string();
assert!(
message.contains(key),
"and the refusal must name the key that would have been ignored: {message}"
);
}
}
#[test]
fn the_generic_scheme_reads_what_the_presets_refuse() {
let config = WebhookRouteConfig {
provider: "jwt-jwks".to_string(),
jwks_uri: Some(JWKS.to_string()),
audience: Some("my-app".to_string()),
credential: Some("body:token".parse().expect("a location")),
event_type_claim: Some("kind".to_string()),
payload_claim: Some("event".to_string()),
id_claim: Some("ref".to_string()),
max_age_secs: Some(600),
..Default::default()
};
webhook_routes_check(&route(config), |_| None, true)
.expect("the operator owns the generic scheme's signing details");
}
#[test]
fn the_generic_scheme_refuses_a_body_digest_combined_with_claim_keys() {
let config = WebhookRouteConfig {
provider: "jwt-jwks".to_string(),
jwks_uri: Some(JWKS.to_string()),
body_hash_claim: Some("request_body_sha256".to_string()),
event_type_claim: Some("kind".to_string()),
..Default::default()
};
let Err(error) = webhook_routes_check(&route(config), |_| None, true) else {
panic!("the two are contradictory, not merely redundant, so this must be refused")
};
assert!(error.to_string().contains("event_type_claim"), "got {error}");
}
#[tokio::test]
async fn a_secret_scheme_with_no_secret_still_skips_in_development_and_refuses_in_production() {
let config = WebhookRouteConfig {
secret_env: Some("MISSING_SECRET".to_string()),
provider: "hmac-sha256".to_string(),
..Default::default()
};
let built = webhook_routes_check(&route(config.clone()), |_| None, false)
.expect("development downgrades this to a warning");
let state = WebhookInboundState::new(lazy_pool(), &built, |_| None);
assert!(
state.mounted_segments().is_empty(),
"an unconfigured secret route must be skipped, so it answers 404 rather than \
500ing with the variable's name in the body"
);
let Err(error) = webhook_routes_check(&route(config), |_| None, true) else {
panic!("production must refuse a route whose signing secret is unset")
};
assert!(
error.to_string().contains("FRAISEQL_ENV=development"),
"and says how a local setup may proceed: {error}"
);
}
}
#[test]
fn the_three_key_material_refusals_stay_distinguishable() {
use crate::config::WebhookRouteConfig;
let route = |config: WebhookRouteConfig| {
let mut routes = HashMap::new();
routes.insert("r".to_string(), config);
routes
};
let hmac_with_env = WebhookRouteConfig {
secret_env: Some("A_SECRET_ENV".to_string()),
provider: "hmac-sha256".to_string(),
..Default::default()
};
let Err(error) = super::webhook_routes_check(&route(hmac_with_env), |_| None, true) else {
panic!("production refuses a route whose secret is unset")
};
let message = error.to_string();
assert!(message.contains("is not set"), "names the missing variable: {message}");
assert!(
!message.contains("cannot use the key material"),
"and not the shape answer — a secret that is not there has no shape: {message}"
);
assert!(message.contains("A_SECRET_ENV"), "and says which variable: {message}");
let no_env = WebhookRouteConfig {
provider: "hmac-sha256".to_string(),
..Default::default()
};
let Err(error) = super::webhook_routes_check(&route(no_env), |_| None, true) else {
panic!("a shared-secret scheme with no secret_env cannot verify anything")
};
let message = error.to_string();
assert!(
message.contains("sets no `secret_env`"),
"there is nothing to look up, so the fix is to add the key rather than to set a \
variable: {message}"
);
let clerk = WebhookRouteConfig {
secret_env: Some("A_SECRET_ENV".to_string()),
provider: "clerk".to_string(),
..Default::default()
};
let Err(error) = super::webhook_routes_check(
&route(clerk),
|_| Some("whpk_C2FVsBQIhrscChlQIMV+b5sSYspob7oD".to_string()),
true,
) else {
panic!("asymmetric v1a material is not verifiable here")
};
let message = error.to_string();
assert!(
message.contains("cannot use the key material"),
"a key that IS there and is wrong gets the shape answer: {message}"
);
assert!(message.contains("whpk_"), "naming what was pasted: {message}");
}