Skip to main content

ocre/runtime/
push.rs

1//! [`send`]: one web push message to one subscription.
2
3use worker::{Fetch, Headers, Method, Request, RequestInit, send::SendFuture};
4
5use super::Ctx;
6use crate::{
7    Error, Result,
8    push::{Subscription, VAPID_PRIVATE_KEY, VAPID_PUBLIC_KEY, VAPID_SUBJECT, VapidKeys, encrypt, vapid_authorization},
9};
10
11/// What the push service did with a message.
12#[derive(Debug, Clone, Copy, PartialEq, Eq)]
13pub enum Sent {
14    /// Accepted (201): the browser gets it when it is online, within the TTL.
15    Delivered,
16    /// The subscription no longer exists (404 or 410: the user unsubscribed
17    /// or the browser dropped it): delete it.
18    Gone,
19}
20
21/// Encrypts `message` (JSON, e.g. [`message`](crate::push::message)) for
22/// `subscription` and posts it to its push service, signed with the app's
23/// VAPID keys. The push service keeps it up to `ttl` seconds while the
24/// browser is offline. One subrequest.
25///
26/// # Errors
27///
28/// - [`Error::BadRequest`] when the subscription's keys are invalid or the
29///   message is over [`MAX_PAYLOAD`](crate::push::MAX_PAYLOAD) bytes.
30/// - [`Error::Internal`] when a `VAPID_*` setting is missing or invalid,
31///   or the push service refuses the message (other than 404/410), with its
32///   status and answer.
33///
34/// # Examples
35///
36/// ```no_run
37/// use ocre::push::{self, Sent, Subscription};
38/// use ocre::{Ctx, Result};
39///
40/// async fn notify(ctx: &Ctx, subscription: &Subscription) -> Result<()> {
41///     let message = push::message("Your video is ready", "Download it now.", "/videos/42");
42///     if push::send(ctx, subscription, &message, 24 * 3600).await? == Sent::Gone {
43///         // delete the subscription
44///     }
45///     Ok(())
46/// }
47/// # let _ = notify;
48/// ```
49pub fn send(
50    ctx: &Ctx,
51    subscription: &Subscription,
52    message: &serde_json::Value,
53    ttl: u32,
54) -> impl Future<Output = Result<Sent>> + Send + use<> {
55    let env = ctx.env().clone();
56    let subscription = subscription.clone();
57    let payload = message.to_string();
58    SendFuture::new(async move {
59        let var = |name: &str| env.var(name).ok().map(|value| value.to_string()).filter(|v| !v.trim().is_empty());
60        let missing = |name: &str| {
61            Error::internal(format!(
62                "web push needs {name} (not set). Fix: `ocre g push` writes a VAPID key pair to .dev.vars; in \
63                 production set VAPID_PUBLIC_KEY and VAPID_SUBJECT in worker.env and push VAPID_PRIVATE_KEY with \
64                 `ocre secrets push VAPID_PRIVATE_KEY --file .prod.vars`"
65            ))
66        };
67        let private_key = super::secrets::lookup(&env, VAPID_PRIVATE_KEY).await?.filter(|v| !v.trim().is_empty());
68        let keys = VapidKeys {
69            public_key: var(VAPID_PUBLIC_KEY).ok_or_else(|| missing(VAPID_PUBLIC_KEY))?,
70            private_key: private_key.ok_or_else(|| missing(VAPID_PRIVATE_KEY))?,
71        };
72        let subject = var(VAPID_SUBJECT).ok_or_else(|| missing(VAPID_SUBJECT))?;
73        let body = encrypt(&subscription.keys, payload.as_bytes())?;
74        let authorization = vapid_authorization(&subscription.endpoint, &subject, &keys, crate::now())?;
75        let headers = Headers::new();
76        headers.set("Authorization", &authorization)?;
77        headers.set("Content-Encoding", "aes128gcm")?;
78        headers.set("Content-Type", "application/octet-stream")?;
79        headers.set("TTL", &ttl.to_string())?;
80        let mut init = RequestInit::new();
81        let array = worker::js_sys::Uint8Array::from(body.as_slice());
82        init.with_method(Method::Post).with_headers(headers).with_body(Some(array.into()));
83        let mut response = Fetch::Request(Request::new_with_init(&subscription.endpoint, &init)?).send().await?;
84        match response.status_code() {
85            200..=299 => Ok(Sent::Delivered),
86            404 | 410 => Ok(Sent::Gone),
87            status => {
88                let text = response.text().await.unwrap_or_default();
89                Err(Error::internal(format!(
90                    "the push service answered {status}: {}",
91                    text.chars().take(200).collect::<String>()
92                )))
93            }
94        }
95    })
96}