use super::{Edits, next_migration_path, register_routes};
use crate::{
CliResult,
names::{RUST_KEYWORDS, is_identifier},
output::CliError,
project::Project,
secret,
};
const IGNORE: &str = "#[ignore = \"request test: run with `ocre test --e2e`\"]";
pub fn webhook(project: &Project, name: &str, standard: bool) -> CliResult {
let name = name.strip_suffix("_webhook").unwrap_or(name);
if !is_identifier(name) || name.chars().any(|c| c.is_ascii_uppercase()) || RUST_KEYWORDS.contains(&name) {
return Err(CliError::new(format!("invalid webhook name `{name}`"))
.hint("use snake_case, e.g. `ocre g webhook payments`"));
}
let secret_name = format!("{}_WEBHOOK_SECRET", name.to_uppercase());
let command = format!("ocre g webhook {name}{}", if standard { " --standard" } else { "" });
let mut edits = Edits::new(project);
edits.create(&format!("src/{name}_webhook.rs"), handler(name, &secret_name, standard, &command))?;
edits.create(&format!("tests/{name}_webhook.rs"), tests(name, &secret_name, standard, &command))?;
let new_table = !edits.has_create_migration("webhook_events")?;
if new_table {
let path = next_migration_path(&edits, "create_webhook_events")?;
let sql =
format!("-- Generated by `{command}`: one row per webhook event received.\n{}", ocre::webhooks::TABLE_SQL);
edits.create(&path, sql)?;
}
if let Some(vars) = edits.read(".dev.vars")? {
let value = if standard { format!("whsec_{}", secret::generate()) } else { secret::generate() };
let vars = if vars.is_empty() || vars.ends_with('\n') { vars } else { format!("{vars}\n") };
edits.update(".dev.vars", format!("{vars}{secret_name}={value}\n"));
}
register_routes(&mut edits, &format!("{name}_webhook"))?;
let mut report = edits.apply("generate webhook")?;
report.next = vec![
format!("write the effect in `handle` (src/{name}_webhook.rs), then `ocre test --e2e`"),
format!(
"give the sender https://<your host>/webhooks/{name} and the secret; production: put {secret_name} in .prod.vars, then `ocre secrets push {secret_name} --file .prod.vars`"
),
];
if new_table {
report.next.insert(0, "ocre migrate".to_owned());
}
Ok(report)
}
fn handler(name: &str, secret_name: &str, standard: bool, command: &str) -> String {
let (doc, verify) = if standard {
(
"//! The sender follows Standard Webhooks (https://www.standardwebhooks.com):\n\
//! `webhook-id`, `webhook-timestamp` and `webhook-signature` headers, a\n\
//! `whsec_...` secret, deliveries older than 5 minutes refused. The\n\
//! `webhook-id` identifies the event.",
" let id = webhooks::verify_standard(&secret, &headers, &body, 300, ocre::now())?;\n \
let event: Value = serde_json::from_slice(&body).map_err(|_| Error::bad_request(\"the body is not JSON\"))?;",
)
} else {
(
"//! The sender signs the raw body with HMAC-SHA256 and the shared secret,\n\
//! and sends it in the `X-Signature` header (hex, `sha256=<hex>` or\n\
//! base64). Each event is JSON with a unique `id`.",
" let signature = headers.get(\"x-signature\").and_then(|value| value.to_str().ok()).unwrap_or_default();\n \
webhooks::verify(secret.as_bytes(), &body, signature)?;\n \
let event: Value = serde_json::from_slice(&body).map_err(|_| Error::bad_request(\"the body is not JSON\"))?;\n \
let id = match &event[\"id\"] {\n \
Value::String(id) => id.clone(),\n \
Value::Number(id) => id.to_string(),\n \
_ => return Err(Error::bad_request(\"the event has no `id`\")),\n \
};",
)
};
format!(
r#"//! The `{name}` webhook: `POST /webhooks/{name}`. Generated by `{command}`.
//!
{doc}
//!
//! Every event is recorded in `webhook_events` (payload, status, attempts,
//! error) and `handle` runs once per event: a delivery of an event already
//! processed is answered `{{"status": "duplicate"}}`, and a failed one 500,
//! so the sender retries it. The secret is `{secret_name}` (.dev.vars;
//! production: `ocre secrets push {secret_name} --file .prod.vars`).
use axum::{{Json, Router, body::Bytes, extract::State, http::HeaderMap, routing::post}};
use ocre::{{
Ctx, Error, Result,
serde_json::{{self, Value, json}},
webhooks::{{self, Delivery}},
}};
/// The source of the events in `webhook_events`.
const SOURCE: &str = "{name}";
pub fn routes() -> Router<Ctx> {{
Router::new().route("/webhooks/{name}", post(receive))
}}
async fn receive(State(ctx): State<Ctx>, headers: HeaderMap, body: Bytes) -> Result<Json<Value>> {{
let secret = ctx.secret("{secret_name}").await?;
{verify}
let db = ctx.db()?;
let delivery = webhooks::once(&db, SOURCE, &id, &body, || handle(&ctx, &event)).await?;
let status = if matches!(delivery, Delivery::Duplicate) {{ "duplicate" }} else {{ "processed" }};
Ok(Json(json!({{ "status": status }})))
}}
/// The effect of one event, run once per event id: update records with
/// `ctx.db()?`, enqueue a job, broadcast a change. An `Err` marks the event
/// failed and answers 500, so the sender delivers it again.
async fn handle(_ctx: &Ctx, _event: &Value) -> Result<()> {{
Ok(())
}}
"#
)
}
fn tests(name: &str, secret_name: &str, standard: bool, command: &str) -> String {
let deliver = if standard {
format!(
r#"/// Delivers `body` as event `id`, signed with `secret` (Standard Webhooks).
fn deliver(id: &str, body: &str, secret: &str) -> Response {{
let timestamp = ocre::now();
let signature = webhooks::sign_standard(secret, id, timestamp, body.as_bytes()).expect("a whsec_ secret");
let mut client = Client::new()
.header("webhook-id", id)
.header("webhook-timestamp", ×tamp.to_string())
.header("webhook-signature", &signature);
client.request("POST", "/webhooks/{name}", Some(("application/json", body.as_bytes().to_vec())))
}}
fn wrong_secret() -> String {{
format!("whsec_{{}}", "d3Jvbmc")
}}"#
)
} else {
format!(
r#"/// Delivers `body` (event `_id`), signed with `secret` in `X-Signature`.
fn deliver(_id: &str, body: &str, secret: &str) -> Response {{
let mut client = Client::new().header("X-Signature", &webhooks::sign(secret.as_bytes(), body.as_bytes()));
client.request("POST", "/webhooks/{name}", Some(("application/json", body.as_bytes().to_vec())))
}}
fn wrong_secret() -> String {{
"wrong".to_owned()
}}"#
)
};
format!(
r##"//! Request tests for the `{name}` webhook. Generated by `{command}`.
//! They run against the app in workerd with `ocre test --e2e`; plain
//! `cargo test` skips them (`#[ignore]`).
use ocre::{{
serde_json::Value,
testing::{{self, Client, Response}},
webhooks,
}};
fn secret() -> String {{
testing::var("{secret_name}").expect("{secret_name} is in .dev.vars")
}}
{deliver}
/// An event id no earlier run used.
fn unique_id() -> String {{
let nanos = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos();
format!("evt_{{nanos}}")
}}
#[test]
{IGNORE}
fn refuses_unsigned_events() {{
let id = unique_id();
let body = format!(r#"{{{{"id":"{{id}}"}}}}"#);
deliver(&id, &body, &wrong_secret()).assert_status(401);
}}
#[test]
{IGNORE}
fn processes_each_event_once() {{
let id = unique_id();
let body = format!(r#"{{{{"id":"{{id}}"}}}}"#);
let first: Value = deliver(&id, &body, &secret()).assert_status(200).json();
assert_eq!(first["status"], "processed");
let again: Value = deliver(&id, &body, &secret()).assert_status(200).json();
assert_eq!(again["status"], "duplicate", "the effect runs once per event");
}}
"##
)
}