oxedyne/fe2o3/fe2o3_steel/src/srv/webhook.rs
19.4 KiB, 109 runs
created by r1870400018:10053, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 2 | //! Anthropic Claude |
| 3 | |
| 4 | /// Webhook handler infrastructure. |
| 5 | /// |
| 6 | /// Steel provides the trait and registry; apps implement their own |
| 7 | /// handlers and register them before starting the server. |
| 8 | |
| 9 | use crate::srv::cfg::WebhookRoute; |
| 10 | |
| 11 | use oxedyne_fe2o3_core::prelude::*; |
| 12 | use oxedyne_fe2o3_net::{ |
| 13 | hmac::verify_hmac_sha256, |
| 14 | http::{ |
| 15 | fields::HeaderFields, |
| 16 | msg::HttpMessage, |
| 17 | status::HttpStatus, |
| 18 | }, |
| 19 | }; |
| 20 | |
| 21 | use std::{ |
| 22 | collections::HashMap, |
| 23 | future::Future, |
| 24 | pin::Pin, |
| 25 | sync::Arc, |
| 26 | }; |
| 27 | use tokio_rustls::rustls::ClientConfig; |
| 28 | |
| 29 | |
| 30 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 31 | // │ WEBHOOK HANDLER TRAIT │ |
| 32 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 33 | |
| 34 | /// Apps implement this trait for each webhook integration they need -- a payment |
| 35 | /// provider forwarding a purchase confirmation to a fulfilment upstream, a |
| 36 | /// notification service relaying an event -- and register instances with a |
| 37 | /// [`WebhookRegistry`] before starting Steel. |
| 38 | pub trait WebhookHandler: Send + Sync + 'static { |
| 39 | fn handle<'a>( |
| 40 | &'a self, |
| 41 | route: &'a WebhookRoute, |
| 42 | body: &'a [u8], |
| 43 | req_headers: &'a HeaderFields, |
| 44 | tls_client: &'a Option<Arc<ClientConfig>>, |
| 45 | id: &'a str, |
| 46 | ) -> Pin<Box<dyn Future<Output = Outcome<Option<HttpMessage>>> + Send + 'a>>; |
| 47 | } |
| 48 | |
| 49 | |
| 50 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 51 | // │ WEBHOOK REGISTRY │ |
| 52 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 53 | |
| 54 | /// Maps handler names, as written in config, to handler implementations. Built |
| 55 | /// by the app before server startup; the stock `steel` binary creates an empty |
| 56 | /// one. |
| 57 | #[derive(Default)] |
| 58 | pub struct WebhookRegistry { |
| 59 | handlers: HashMap<String, Box<dyn WebhookHandler>>, |
| 60 | } |
| 61 | |
| 62 | impl WebhookRegistry { |
| 63 | pub fn new() -> Self { |
| 64 | Self { |
| 65 | handlers: HashMap::new(), |
| 66 | } |
| 67 | } |
| 68 | |
| 69 | /// The name must match the `handler` field in the webhook route config. |
| 70 | pub fn register<H: WebhookHandler>(&mut self, name: &str, handler: H) { |
| 71 | self.handlers.insert(name.to_string(), Box::new(handler)); |
| 72 | } |
| 73 | |
| 74 | pub fn insert_boxed(&mut self, name: String, handler: Box<dyn WebhookHandler>) { |
| 75 | self.handlers.insert(name, handler); |
| 76 | } |
| 77 | |
| 78 | pub fn has(&self, name: &str) -> bool { |
| 79 | self.handlers.contains_key(name) |
| 80 | } |
| 81 | } |
| 82 | |
| 83 | // Manual Debug impl because Box<dyn WebhookHandler> is not Debug. |
| 84 | impl std::fmt::Debug for WebhookRegistry { |
| 85 | fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { |
| 86 | f.debug_struct("WebhookRegistry") |
| 87 | .field("handlers", &self.handlers.keys().collect::<Vec<_>>()) |
| 88 | .finish() |
| 89 | } |
| 90 | } |
| 91 | |
| 92 | |
| 93 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 94 | // │ DISPATCH │ |
| 95 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 96 | |
| 97 | /// Called only for handler-mode webhook routes; upstream-mode routes are |
| 98 | /// forwarded directly from the HTTPS dispatcher. |
| 99 | pub async fn dispatch( |
| 100 | registry: &WebhookRegistry, |
| 101 | route: &WebhookRoute, |
| 102 | body: &[u8], |
| 103 | req_headers: &HeaderFields, |
| 104 | tls_client: &Option<Arc<ClientConfig>>, |
| 105 | id: &str, |
| 106 | ) |
| 107 | -> Outcome<Option<HttpMessage>> |
| 108 | { |
| 109 | let name = match &route.handler { |
| 110 | Some(n) => n, |
| 111 | None => { |
| 112 | warn!("{}: webhook::dispatch called on an upstream-mode route.", id); |
| 113 | return Ok(Some(HttpMessage::respond_with_text( |
| 114 | HttpStatus::InternalServerError, |
| 115 | "Webhook route misconfigured.", |
| 116 | ))); |
| 117 | }, |
| 118 | }; |
| 119 | match registry.handlers.get(name) { |
| 120 | Some(handler) => handler.handle( |
| 121 | route, body, req_headers, tls_client, id, |
| 122 | ).await, |
| 123 | None => { |
| 124 | warn!("{}: No registered webhook handler '{}'.", id, name); |
| 125 | Ok(Some(HttpMessage::respond_with_text( |
| 126 | HttpStatus::NotFound, |
| 127 | "Unknown webhook handler.", |
| 128 | ))) |
| 129 | } |
| 130 | } |
| 131 | } |
| 132 | |
| 133 | |
| 134 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 135 | // │ HELPER UTILITIES (re-exported for app handlers) │ |
| 136 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 137 | |
| 138 | /// Percent-encodes for `application/x-www-form-urlencoded`. |
| 139 | pub fn url_encode(s: &str) -> String { |
| 140 | let mut out = String::with_capacity(s.len() * 2); |
| 141 | for b in s.bytes() { |
| 142 | match b { |
| 143 | b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => { |
| 144 | out.push(b as char); |
| 145 | } |
| 146 | b' ' => out.push('+'), |
| 147 | _ => { |
| 148 | out.push('%'); |
| 149 | out.push(char::from(b"0123456789ABCDEF"[(b >> 4) as usize])); |
| 150 | out.push(char::from(b"0123456789ABCDEF"[(b & 0x0F) as usize])); |
| 151 | } |
| 152 | } |
| 153 | } |
| 154 | out |
| 155 | } |
| 156 | |
| 157 | /// Looks for `"key": "value"` and returns the value, treating `null` as an empty |
| 158 | /// string. No full JSON parser -- just enough for a webhook payload. |
| 159 | pub fn extract_value(block: &str, key: &str) -> Option<String> { |
| 160 | let key_pos = ok!(block.find(key)); |
| 161 | let after_key = &block[key_pos + key.len()..]; |
| 162 | let colon_pos = ok!(after_key.find(':')); |
| 163 | let after_colon = after_key[colon_pos + 1..].trim_start(); |
| 164 | if after_colon.starts_with('"') { |
| 165 | let content = &after_colon[1..]; |
| 166 | let end = content.find('"').unwrap_or(content.len()); |
| 167 | Some(content[..end].to_string()) |
| 168 | } else if after_colon.starts_with("null") { |
| 169 | Some(String::new()) |
| 170 | } else { |
| 171 | let end = after_colon.find(|c: char| c == ',' || c == '}' || c == '\n') |
| 172 | .unwrap_or(after_colon.len()); |
| 173 | Some(after_colon[..end].trim().to_string()) |
| 174 | } |
| 175 | } |
| 176 | |
| 177 | /// As `extract_value`, but only within the section starting at `section_key`. |
| 178 | pub fn extract_json_string(json: &str, section_key: &str, key: &str) -> Option<String> { |
| 179 | let section_start = ok!(json.find(section_key)); |
| 180 | let block = &json[section_start..json.len().min(section_start + 3000)]; |
| 181 | extract_value(block, key) |
| 182 | } |
| 183 | |
| 184 | /// Extract a JSON string value for a key at the TOP level of the object -- brace |
| 185 | /// depth one -- ignoring a same-named key nested deeper. |
| 186 | /// |
| 187 | /// [`extract_value`] returns the first match anywhere, which is wrong for a |
| 188 | /// top-level field that a nested object also carries. A Stripe event is the |
| 189 | /// case this exists for: it orders the top-level `"type"` (e.g. |
| 190 | /// `"charge.refunded"`) *after* `"data"`, and the nested charge object holds an |
| 191 | /// `"outcome": { "type": "authorized" }`, so the first `"type"` in the payload |
| 192 | /// is the wrong one. Matching only at depth one returns the event type. |
| 193 | /// |
| 194 | /// `key` includes its surrounding quotes, as [`extract_value`] expects (e.g. |
| 195 | /// `"\"type\""`). Strings and their escapes are respected while tracking depth, |
| 196 | /// so a brace inside a string value does not shift the count. |
| 197 | pub fn extract_top_level_value(block: &str, key: &str) -> Option<String> { |
| 198 | let bytes = block.as_bytes(); |
| 199 | let mut depth: i32 = 0; |
| 200 | let mut in_str = false; |
| 201 | let mut esc = false; |
| 202 | let mut i = 0; |
| 203 | while i < bytes.len() { |
| 204 | let c = bytes[i]; |
| 205 | if in_str { |
| 206 | // Inside a string: consume until the unescaped closing quote, so a |
| 207 | // brace or a quote in the text cannot move the depth or end the |
| 208 | // string early. |
| 209 | if esc { esc = false; } |
| 210 | else if c == b'\\' { esc = true; } |
| 211 | else if c == b'"' { in_str = false; } |
| 212 | i += 1; |
| 213 | continue; |
| 214 | } |
| 215 | match c { |
| 216 | b'"' => { |
| 217 | // A key at depth one is `"key" :` -- a string, at the object's |
| 218 | // own level, followed by a colon. A value that happens to equal |
| 219 | // the key text is followed by a comma or a brace, not a colon, |
| 220 | // so the colon check keeps this to keys. |
| 221 | if depth == 1 && block[i..].starts_with(key) { |
| 222 | let after = block[i + key.len()..].trim_start(); |
| 223 | if after.starts_with(':') { |
| 224 | return extract_value(&block[i..], key); |
| 225 | } |
| 226 | } |
| 227 | in_str = true; |
| 228 | } |
| 229 | b'{' | b'[' => depth += 1, |
| 230 | b'}' | b']' => depth -= 1, |
| 231 | _ => {} |
| 232 | } |
| 233 | i += 1; |
| 234 | } |
| 235 | None |
| 236 | } |
| 237 | |
| 238 | |
| 239 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 240 | // │ STRIPE WEBHOOK SIGNATURE VERIFICATION │ |
| 241 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 242 | |
| 243 | // Stripe's own libraries default to five minutes. |
| 244 | pub const STRIPE_SIG_TOLERANCE_SECS: u64 = 300; |
| 245 | |
| 246 | /// Stripe signs each webhook with the endpoint's signing secret (`whsec_...`) |
| 247 | /// and sends the result in a `Stripe-Signature` header of the form |
| 248 | /// `t=<unix_ts>,v1=<hex_hmac>[,v1=<hex_hmac>...]` -- there may be several `v1` |
| 249 | /// schemes during a secret rotation, and other schemes such as `v0`, which are |
| 250 | /// ignored. The signed payload is the ASCII string `"<t>.<raw request body>"`, |
| 251 | /// and the tag is `HMAC-SHA256` of that payload under the signing secret, taken |
| 252 | /// over the whole `whsec_...` string as configured. |
| 253 | /// |
| 254 | /// Verification succeeds when the recomputed HMAC matches any supplied `v1` |
| 255 | /// value, compared in constant time, **and** the timestamp is within |
| 256 | /// `tolerance_secs` of `now_secs`. The current time is passed in rather than |
| 257 | /// read from a clock, so the check is deterministic and testable. The error |
| 258 | /// never carries the signing secret. |
| 259 | pub fn verify_stripe_signature( |
| 260 | secret: &str, |
| 261 | body: &[u8], |
| 262 | sig_header: &str, |
| 263 | now_secs: u64, |
| 264 | tolerance_secs: u64, |
| 265 | ) |
| 266 | -> Outcome<()> |
| 267 | { |
| 268 | // Parse the comma-separated scheme=value pairs, collecting the |
| 269 | // timestamp and every v1 tag. |
| 270 | let mut ts: Option<&str> = None; |
| 271 | let mut v1s: Vec<&str> = Vec::new(); |
| 272 | for part in sig_header.split(',') { |
| 273 | let mut kv = part.splitn(2, '='); |
| 274 | let scheme = match kv.next() { |
| 275 | Some(s) => s.trim(), |
| 276 | None => continue, |
| 277 | }; |
| 278 | let value = match kv.next() { |
| 279 | Some(v) => v.trim(), |
| 280 | None => continue, |
| 281 | }; |
| 282 | match scheme { |
| 283 | "t" => ts = Some(value), |
| 284 | "v1" => v1s.push(value), |
| 285 | _ => {}, // Ignore v0 and any future schemes. |
| 286 | } |
| 287 | } |
| 288 | |
| 289 | let ts = match ts { |
| 290 | Some(t) => t, |
| 291 | None => return Err(err!( |
| 292 | "Stripe-Signature header has no timestamp (t=)."; |
| 293 | Invalid, Input, Missing)), |
| 294 | }; |
| 295 | if v1s.is_empty() { |
| 296 | return Err(err!( |
| 297 | "Stripe-Signature header has no v1 signature."; |
| 298 | Invalid, Input, Missing)); |
| 299 | } |
| 300 | |
| 301 | // Enforce the timestamp tolerance before any HMAC work. |
| 302 | let t_secs = match ts.parse::<u64>() { |
| 303 | Ok(n) => n, |
| 304 | Err(_) => return Err(err!( |
| 305 | "Stripe-Signature timestamp '{}' is not an integer.", ts; |
| 306 | Invalid, Input, Mismatch)), |
| 307 | }; |
| 308 | let skew = if now_secs >= t_secs { |
| 309 | now_secs - t_secs |
| 310 | } else { |
| 311 | t_secs - now_secs |
| 312 | }; |
| 313 | if skew > tolerance_secs { |
| 314 | return Err(err!( |
| 315 | "Stripe-Signature timestamp outside tolerance ({}s > {}s).", |
| 316 | skew, tolerance_secs; |
| 317 | Invalid, Input, Range)); |
| 318 | } |
| 319 | |
| 320 | // Signed payload is "<t>.<body>". Build the byte string. |
| 321 | let mut signed = Vec::with_capacity(ts.len() + 1 + body.len()); |
| 322 | signed.extend_from_slice(ts.as_bytes()); |
| 323 | signed.push(b'.'); |
| 324 | signed.extend_from_slice(body); |
| 325 | |
| 326 | // Constant-time compare the recomputed HMAC against each v1 tag. |
| 327 | for v1 in &v1s { |
| 328 | let tag = match hex_decode(v1) { |
| 329 | Some(bytes) => bytes, |
| 330 | None => continue, // Malformed hex cannot match. |
| 331 | }; |
| 332 | if verify_hmac_sha256(secret.as_bytes(), &signed, &tag) { |
| 333 | return Ok(()); |
| 334 | } |
| 335 | } |
| 336 | |
| 337 | Err(err!( |
| 338 | "Stripe-Signature verification failed: no v1 tag matched."; |
| 339 | Invalid, Input, Security)) |
| 340 | } |
| 341 | |
| 342 | /// `None` on an odd length or a non-hex character. |
| 343 | fn hex_decode(s: &str) -> Option<Vec<u8>> { |
| 344 | let bytes = s.as_bytes(); |
| 345 | if bytes.len() % 2 != 0 { |
| 346 | return None; |
| 347 | } |
| 348 | let nibble = |c: u8| -> Option<u8> { |
| 349 | match c { |
| 350 | b'0'..=b'9' => Some(c - b'0'), |
| 351 | b'a'..=b'f' => Some(c - b'a' + 10), |
| 352 | b'A'..=b'F' => Some(c - b'A' + 10), |
| 353 | _ => None, |
| 354 | } |
| 355 | }; |
| 356 | let mut out = Vec::with_capacity(bytes.len() / 2); |
| 357 | let mut i = 0; |
| 358 | while i + 1 < bytes.len() { |
| 359 | let hi = ok!(nibble(bytes[i])); |
| 360 | let lo = ok!(nibble(bytes[i + 1])); |
| 361 | out.push((hi << 4) | lo); |
| 362 | i += 2; |
| 363 | } |
| 364 | Some(out) |
| 365 | } |
| 366 | |
| 367 | |
| 368 | #[cfg(test)] |
| 369 | mod tests { |
| 370 | use super::*; |
| 371 | use oxedyne_fe2o3_net::hmac::hmac_sha256; |
| 372 | |
| 373 | fn hex_encode(bytes: &[u8]) -> String { |
| 374 | let mut out = String::with_capacity(bytes.len() * 2); |
| 375 | for b in bytes { |
| 376 | out.push(char::from(b"0123456789abcdef"[(b >> 4) as usize])); |
| 377 | out.push(char::from(b"0123456789abcdef"[(b & 0x0F) as usize])); |
| 378 | } |
| 379 | out |
| 380 | } |
| 381 | |
| 382 | #[test] |
| 383 | fn top_level_value_ignores_a_nested_key_of_the_same_name() { |
| 384 | // A Stripe event, ordered as Stripe orders it: the top-level `type` |
| 385 | // comes AFTER `data`, and the nested charge carries `outcome.type`. |
| 386 | let event = r#"{"id":"evt_1","object":"event","data":{"object":{"id":"ch_1", |
| 387 | "object":"charge","outcome":{"network_status":"approved","type":"authorized"}, |
| 388 | "payment_method_details":{"type":"card"},"amount_refunded":1000}}, |
| 389 | "type":"charge.refunded"}"#; |
| 390 | // The naive scan grabs the first `type` -- the wrong one. |
| 391 | assert_eq!(extract_value(event, "\"type\""), Some(fmt!("authorized"))); |
| 392 | // The depth-aware one returns the event type. |
| 393 | assert_eq!(extract_top_level_value(event, "\"type\""), Some(fmt!("charge.refunded"))); |
| 394 | // The top-level `id` is returned even though a nested `id` exists. |
| 395 | assert_eq!(extract_top_level_value(event, "\"id\""), Some(fmt!("evt_1"))); |
| 396 | // A brace inside a string value does not shift the depth count. |
| 397 | let tricky = r#"{"note":"a } brace { in text","type":"top"}"#; |
| 398 | assert_eq!(extract_top_level_value(tricky, "\"type\""), Some(fmt!("top"))); |
| 399 | // A key present only deeper than the top level is not found at all. |
| 400 | let nested_only = r#"{"data":{"object":{"kind":"deep"}}}"#; |
| 401 | assert_eq!(extract_top_level_value(nested_only, "\"kind\""), None); |
| 402 | // A top-level value that equals the key text is not mistaken for the key. |
| 403 | let valueish = r#"{"a":"type","type":"real"}"#; |
| 404 | assert_eq!(extract_top_level_value(valueish, "\"type\""), Some(fmt!("real"))); |
| 405 | } |
| 406 | |
| 407 | fn sign(secret: &str, body: &[u8], t: u64) -> String { |
| 408 | let mut signed = Vec::new(); |
| 409 | signed.extend_from_slice(t.to_string().as_bytes()); |
| 410 | signed.push(b'.'); |
| 411 | signed.extend_from_slice(body); |
| 412 | let tag = hmac_sha256(secret.as_bytes(), &signed); |
| 413 | fmt!("t={},v1={}", t, hex_encode(&tag)) |
| 414 | } |
| 415 | |
| 416 | #[test] |
| 417 | fn test_valid_signature_verifies() { |
| 418 | let secret = "whsec_test_secret"; |
| 419 | let body = br#"{"id":"evt_1","type":"checkout.session.completed"}"#; |
| 420 | let t = 1_700_000_000u64; |
| 421 | let header = sign(secret, body, t); |
| 422 | // Within tolerance of the signing time. |
| 423 | assert!(verify_stripe_signature( |
| 424 | secret, body, &header, t + 10, STRIPE_SIG_TOLERANCE_SECS).is_ok()); |
| 425 | } |
| 426 | |
| 427 | #[test] |
| 428 | fn test_tampered_body_rejected() { |
| 429 | let secret = "whsec_test_secret"; |
| 430 | let body = br#"{"amount":100}"#; |
| 431 | let t = 1_700_000_000u64; |
| 432 | let header = sign(secret, body, t); |
| 433 | let tampered = br#"{"amount":999}"#; |
| 434 | assert!(verify_stripe_signature( |
| 435 | secret, tampered, &header, t, STRIPE_SIG_TOLERANCE_SECS).is_err()); |
| 436 | } |
| 437 | |
| 438 | #[test] |
| 439 | fn test_wrong_secret_rejected() { |
| 440 | let body = br#"{"amount":100}"#; |
| 441 | let t = 1_700_000_000u64; |
| 442 | let header = sign("whsec_real", body, t); |
| 443 | assert!(verify_stripe_signature( |
| 444 | "whsec_forged", body, &header, t, STRIPE_SIG_TOLERANCE_SECS).is_err()); |
| 445 | } |
| 446 | |
| 447 | #[test] |
| 448 | fn test_expired_timestamp_rejected() { |
| 449 | let secret = "whsec_test_secret"; |
| 450 | let body = br#"{"ok":true}"#; |
| 451 | let t = 1_700_000_000u64; |
| 452 | let header = sign(secret, body, t); |
| 453 | // now is well past t + tolerance. |
| 454 | let now = t + STRIPE_SIG_TOLERANCE_SECS + 1; |
| 455 | assert!(verify_stripe_signature( |
| 456 | secret, body, &header, now, STRIPE_SIG_TOLERANCE_SECS).is_err()); |
| 457 | } |
| 458 | |
| 459 | #[test] |
| 460 | fn test_future_timestamp_rejected() { |
| 461 | let secret = "whsec_test_secret"; |
| 462 | let body = br#"{"ok":true}"#; |
| 463 | let t = 1_700_000_000u64 + STRIPE_SIG_TOLERANCE_SECS + 5; |
| 464 | let header = sign(secret, body, t); |
| 465 | // now is before the signing time by more than tolerance. |
| 466 | let now = 1_700_000_000u64; |
| 467 | assert!(verify_stripe_signature( |
| 468 | secret, body, &header, now, STRIPE_SIG_TOLERANCE_SECS).is_err()); |
| 469 | } |
| 470 | |
| 471 | #[test] |
| 472 | fn test_multiple_v1_one_matches() { |
| 473 | let secret = "whsec_test_secret"; |
| 474 | let body = br#"{"ok":true}"#; |
| 475 | let t = 1_700_000_000u64; |
| 476 | let good = sign(secret, body, t); |
| 477 | // Prepend a bogus v1 to simulate a rotation window; the real |
| 478 | // one still matches. |
| 479 | let header = fmt!("{},v1=deadbeef", good); |
| 480 | assert!(verify_stripe_signature( |
| 481 | secret, body, &header, t, STRIPE_SIG_TOLERANCE_SECS).is_ok()); |
| 482 | } |
| 483 | |
| 484 | #[test] |
| 485 | fn test_missing_v1_rejected() { |
| 486 | let secret = "whsec_test_secret"; |
| 487 | let body = br#"{"ok":true}"#; |
| 488 | let header = "t=1700000000"; |
| 489 | assert!(verify_stripe_signature( |
| 490 | secret, body, header, 1_700_000_000u64, STRIPE_SIG_TOLERANCE_SECS).is_err()); |
| 491 | } |
| 492 | } |