Skip to main content

ocre/runtime/
webhooks.rs

1//! [`once`]: a webhook event's effect runs once, however many times it is delivered.
2
3use crate::{Db, Result, params};
4
5/// A delivery left `processing` this long (seconds) is taken over by the
6/// next one: the invocation that claimed it died before finishing.
7pub const STALE_AFTER: i64 = 300;
8
9/// What [`once`] did with a delivery.
10#[derive(Debug, Clone, PartialEq, Eq)]
11pub enum Delivery<T> {
12    /// The effect ran (first delivery, or a retry after a failure) and returned this.
13    Processed(T),
14    /// Already processed, or being processed by another invocation: nothing ran.
15    Duplicate,
16}
17
18/// Runs `effect` for the event `event_id` of `source` unless it already ran.
19///
20/// The delivery is claimed with one insert into `webhook_events` (unique on
21/// `source, event_id`), with the raw `payload` kept for the log. When the
22/// effect succeeds the event is marked `processed`, and later deliveries of
23/// the same id get [`Delivery::Duplicate`]: answer them 200 so the provider
24/// stops retrying. When it fails, the event is marked `failed` with the
25/// error and the error is returned, so the handler answers 500 and the
26/// provider's retry runs the effect again. A delivery still `processing`
27/// after [`STALE_AFTER`] seconds is taken over.
28///
29/// Two to three D1 writes per delivery. Recording and effect are separate
30/// statements: an effect made of D1 writes should be idempotent itself
31/// (`INSERT ... ON CONFLICT DO NOTHING`, `UPDATE ... WHERE status = 'pending'`)
32/// to stay correct if the invocation dies between them.
33///
34/// # Errors
35///
36/// The effect's error, or a D1 error (the table is missing: run
37/// `ocre g webhook` and `ocre migrate`).
38pub async fn once<T, F, Fut>(db: &Db, source: &str, event_id: &str, payload: &[u8], effect: F) -> Result<Delivery<T>>
39where
40    F: FnOnce() -> Fut,
41    Fut: Future<Output = Result<T>>,
42{
43    let now = crate::now();
44    let claimed = db
45        .execute(
46            "INSERT INTO webhook_events (source, event_id, payload, status, received_at) \
47             VALUES (?1, ?2, ?3, 'processing', ?4) \
48             ON CONFLICT (source, event_id) DO UPDATE SET status = 'processing', attempts = attempts + 1, \
49             received_at = ?4 \
50             WHERE webhook_events.status = 'failed' \
51             OR (webhook_events.status = 'processing' AND webhook_events.received_at < ?5)",
52            params![source, event_id, String::from_utf8_lossy(payload).into_owned(), now, now - STALE_AFTER],
53        )
54        .await?;
55    if claimed == 0 {
56        return Ok(Delivery::Duplicate);
57    }
58    match effect().await {
59        Ok(value) => {
60            db.execute(
61                "UPDATE webhook_events SET status = 'processed', error = NULL, processed_at = ?3 \
62                 WHERE source = ?1 AND event_id = ?2",
63                params![source, event_id, crate::now()],
64            )
65            .await?;
66            Ok(Delivery::Processed(value))
67        }
68        Err(err) => {
69            db.execute(
70                "UPDATE webhook_events SET status = 'failed', error = ?3 WHERE source = ?1 AND event_id = ?2",
71                params![source, event_id, err.to_string()],
72            )
73            .await?;
74            Err(err)
75        }
76    }
77}
78
79/// The answer of [`post_signed`]: status and body.
80#[derive(Debug, Clone, PartialEq, Eq)]
81pub struct Answer {
82    /// HTTP status.
83    pub status: u16,
84    /// Body, as text.
85    pub text: String,
86}
87
88/// POSTs `body` as JSON to `url`, signed: the HMAC-SHA256 of the body with
89/// `secret` in `X-Signature` (lowercase hex, see [`sign`](crate::webhooks::sign)),
90/// and `Authorization: Bearer <bearer>` when given (API keys of services
91/// such as RunPod). One subrequest. Any status is an answer: check
92/// `status` (a service's 4xx is not an error of the call).
93///
94/// # Errors
95///
96/// [`Error::Internal`](crate::Error::Internal) when the request cannot be
97/// built or the service cannot be reached.
98pub fn post_signed(
99    url: &str,
100    secret: &[u8],
101    bearer: Option<&str>,
102    body: &serde_json::Value,
103) -> impl Future<Output = Result<Answer>> + Send + use<> {
104    let url = url.to_owned();
105    let body = body.to_string();
106    let signature = crate::webhooks::sign(secret, body.as_bytes());
107    let bearer = bearer.map(str::to_owned);
108    worker::send::SendFuture::new(async move {
109        let headers = worker::Headers::new();
110        headers.set("Content-Type", "application/json")?;
111        headers.set("X-Signature", &signature)?;
112        if let Some(bearer) = bearer {
113            headers.set("Authorization", &format!("Bearer {bearer}"))?;
114        }
115        let mut init = worker::RequestInit::new();
116        init.with_method(worker::Method::Post).with_headers(headers).with_body(Some(body.into()));
117        let mut response = worker::Fetch::Request(worker::Request::new_with_init(&url, &init)?).send().await?;
118        Ok(Answer { status: response.status_code(), text: response.text().await? })
119    })
120}