Skip to main content

renox_core/auth/
notifications.rs

1//! Notifications: one message to a user (or to an address), delivered by
2//! email, stored for an in-app list, sent through the app's own channels
3//! (WhatsApp, SMS, Slack…), or all of these.
4//!
5//! ```
6//! # use renox::prelude::*;
7//! use renox::auth::{Channel, DatabaseMessage, Notification, Recipient};
8//! use renox::mail::Mail;
9//!
10//! struct OrderShipped { order_id: i64 }
11//!
12//! impl Notification for OrderShipped {
13//!     fn kind(&self) -> &'static str { "order-shipped" }
14//!     // Per recipient, like Laravel's `via`: WhatsApp only for those with a number.
15//!     fn channels(&self, to: &Recipient) -> Vec<Channel> {
16//!         let mut channels = vec![Channel::Mail, Channel::Database];
17//!         if to.address("whatsapp").is_some() {
18//!             channels.push(Channel::Custom("whatsapp"));
19//!         }
20//!         channels
21//!     }
22//!
23//!     fn to_mail(&self, to: &Recipient, state: &AppState) -> Result<Mail> {
24//!         let email = to.email().unwrap_or_default();
25//!         state.mail_view(&email, "Your order has shipped", "mail/shipped", context! { id => self.order_id })
26//!     }
27//!
28//!     // What the in-app list (the UI kit's `notification_bell`) shows.
29//!     fn to_database(&self, _: &Recipient, _: &AppState) -> Result<renox::serde_json::Value> {
30//!         Ok(DatabaseMessage::success(format!("Order #{} shipped", self.order_id))
31//!             .body("It arrives in 2–3 days.")
32//!             .url(format!("/orders/{}", self.order_id))
33//!             .with("order_id", self.order_id) // any other keys the app reads back
34//!             .into())
35//!     }
36//!
37//!     fn to_channel(&self, _channel: &str, _: &Recipient, _: &AppState) -> Result<renox::serde_json::Value> {
38//!         Ok(json!({ "text": format!("Order #{} has shipped", self.order_id) }))
39//!     }
40//! }
41//!
42//! # async fn demo(state: AppState, user: User, db: Db, order_id: i64) -> Result {
43//! state.notify(&user, &OrderShipped { order_id }).await?;          // now
44//! state.notify_later(&user, &OrderShipped { order_id }).await?;    // through the queue
45//! // Someone without an account:
46//! let guest = Recipient::to("mail", "guest@example.com").and("whatsapp", "+6281234567890");
47//! state.notify(&guest, &OrderShipped { order_id }).await?;
48//! let unread = user.unread_notifications(&db).await?;
49//! # let _ = unread; Ok(()) }
50//!
51//! // The app's own channel: `message` is what `to_channel` returned.
52//! # let _ =
53//! App::new().channel("whatsapp", |to: Recipient, message, state| async move {
54//!     let phone = to.address("whatsapp").or_else(|| to.user()?.get("phone"));
55//!     let _ = (state, phone, message); // call your provider's API here
56//!     Ok(())
57//! })
58//! # ;
59//! ```
60
61use std::collections::BTreeMap;
62use std::sync::Arc;
63
64use anyhow::anyhow;
65use serde::{Deserialize, Serialize};
66use serde_json::Value;
67
68use super::User;
69use crate::db::{DateTime, Db, now};
70use crate::i18n::with_locale;
71use crate::mail::Mail;
72use crate::queue::{Job, JobContext};
73use crate::toast::{ToastAction, ToastKind};
74use crate::{AppState, Result};
75
76/// Where a notification goes.
77#[derive(Debug, Clone, Copy, PartialEq, Eq)]
78#[non_exhaustive]
79pub enum Channel {
80    /// Sent with `state.mailer` (`to_mail`).
81    Mail,
82    /// Stored in the `notifications` table for an in-app list
83    /// (`to_database`); only for users.
84    Database,
85    /// One of the app's own channels, registered with `App::channel`
86    /// (`to_channel`).
87    Custom(&'static str),
88}
89
90/// Who a notification goes to: a user, addresses for someone without an
91/// account, or a user with an extra address.
92#[derive(Debug, Clone, Default, Serialize, Deserialize)]
93#[non_exhaustive]
94pub struct Recipient {
95    user: Option<User>,
96    /// An address per channel, e.g. `mail` → `a@b.c`, `whatsapp` → `+62…`.
97    routes: BTreeMap<String, String>,
98    /// The language to write in; see [`Recipient::locale`].
99    #[serde(default)]
100    language: Option<String>,
101}
102
103impl Recipient {
104    /// The user, at their email address (and any addresses added with `and`).
105    pub fn for_user(user: &User) -> Self {
106        Self {
107            user: Some(user.clone()),
108            routes: BTreeMap::new(),
109            language: None,
110        }
111    }
112
113    /// The user, if the recipient has an account.
114    pub fn user(&self) -> Option<&User> {
115        self.user.as_ref()
116    }
117
118    /// Someone without an account: `Recipient::to("mail", "a@b.c")`.
119    pub fn to(channel: &str, address: impl Into<String>) -> Self {
120        Self::default().and(channel, address)
121    }
122
123    /// Adds (or replaces) the address for `channel`.
124    pub fn and(mut self, channel: &str, address: impl Into<String>) -> Self {
125        self.routes.insert(channel.to_owned(), address.into());
126        self
127    }
128
129    /// The address for `channel`; for `mail`, the user's email when no
130    /// other address was given.
131    pub fn address(&self, channel: &str) -> Option<String> {
132        self.routes.get(channel).cloned().or_else(|| match channel {
133            "mail" => self.user.as_ref().map(|u| u.email.clone()),
134            _ => None,
135        })
136    }
137
138    /// Whether `address` is one of the recipient's addresses, on any channel.
139    pub(crate) fn has_address(&self, address: &str) -> bool {
140        self.routes.values().any(|a| a == address) || self.email().as_deref() == Some(address)
141    }
142
143    /// The address for the `mail` channel (see [`Recipient::address`]).
144    pub fn email(&self) -> Option<String> {
145        self.address("mail")
146    }
147
148    /// Writes to this recipient in `locale` (e.g. `"es"`).
149    pub fn in_locale(mut self, locale: impl Into<String>) -> Self {
150        self.language = Some(locale.into());
151        self
152    }
153
154    /// The recipient's language: the one given with [`Recipient::in_locale`],
155    /// else the user's `locale` column if the `users` table has one.
156    /// Notifications build their messages in it (`t()` in mail views,
157    /// `state.current_lang()` in code).
158    pub fn locale(&self) -> Option<String> {
159        self.language.clone().or_else(|| {
160            self.user
161                .as_ref()?
162                .get::<String>("locale")
163                .filter(|l| !l.is_empty())
164        })
165    }
166}
167
168impl From<&User> for Recipient {
169    fn from(user: &User) -> Self {
170        Self::for_user(user)
171    }
172}
173
174impl From<&crate::AuthUser> for Recipient {
175    fn from(user: &crate::AuthUser) -> Self {
176        Self::for_user(user)
177    }
178}
179
180impl From<&Recipient> for Recipient {
181    fn from(recipient: &Recipient) -> Self {
182        recipient.clone()
183    }
184}
185
186/// A message to a [`Recipient`], with a version for each channel it goes out on.
187pub trait Notification: Send + Sync {
188    /// Stored with database notifications, e.g. to pick an icon.
189    fn kind(&self) -> &'static str;
190
191    /// The channels to deliver on for this recipient (Laravel's `via`),
192    /// e.g. WhatsApp only for users who turned it on; only `Channel::Mail`
193    /// unless overridden.
194    fn channels(&self, to: &Recipient) -> Vec<Channel> {
195        let _ = to;
196        vec![Channel::Mail]
197    }
198
199    /// The mail for `Channel::Mail`; an error unless overridden.
200    fn to_mail(&self, _to: &Recipient, _state: &AppState) -> Result<Mail> {
201        Err(anyhow!("notification `{}` has no mail version", self.kind()).into())
202    }
203
204    /// The JSON stored for `Channel::Database` (`DatabaseNotification::data`);
205    /// `null` by default. Return a [`DatabaseMessage`] for the UI kit's
206    /// `notification_bell`. It runs in the recipient's language, like
207    /// `to_mail`.
208    fn to_database(&self, to: &Recipient, state: &AppState) -> Result<Value> {
209        let _ = (to, state);
210        Ok(Value::Null)
211    }
212
213    /// The message for one of the app's own channels; the channel's handler
214    /// gets it.
215    fn to_channel(&self, channel: &str, to: &Recipient, state: &AppState) -> Result<Value> {
216        let _ = (to, state);
217        Err(anyhow!(
218            "notification `{}` has no version for the `{channel}` channel",
219            self.kind()
220        )
221        .into())
222    }
223}
224
225/// What a notification stores for the in-app list, in the shape the UI
226/// kit's `notification_bell` shows (and pushes as a toast when it arrives):
227/// a status (its icon), a title, a line of text, a link and buttons. Return
228/// it from [`Notification::to_database`] with `.into()`; `with` adds the
229/// app's own keys next to them.
230#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
231#[non_exhaustive]
232pub struct DatabaseMessage {
233    /// Its icon and color: success, info, warning or error.
234    pub status: ToastKind,
235    /// The first line.
236    pub title: String,
237    /// A second, lighter line.
238    #[serde(default, skip_serializing_if = "Option::is_none")]
239    pub body: Option<String>,
240    /// Where opening it goes (it is marked read on the way).
241    #[serde(default, skip_serializing_if = "Option::is_none")]
242    pub url: Option<String>,
243    /// Links or buttons under the text.
244    #[serde(default, skip_serializing_if = "Vec::is_empty")]
245    pub actions: Vec<ToastAction>,
246    /// The app's own keys (`with`), stored next to the others.
247    #[serde(flatten)]
248    pub extra: serde_json::Map<String, Value>,
249}
250
251impl DatabaseMessage {
252    /// A message with `status` and `title`.
253    pub fn new(status: ToastKind, title: impl Into<String>) -> Self {
254        Self {
255            status,
256            title: title.into(),
257            body: None,
258            url: None,
259            actions: Vec::new(),
260            extra: serde_json::Map::new(),
261        }
262    }
263
264    /// A success message (a green check).
265    pub fn success(title: impl Into<String>) -> Self {
266        Self::new(ToastKind::Success, title)
267    }
268
269    /// An info message.
270    pub fn info(title: impl Into<String>) -> Self {
271        Self::new(ToastKind::Info, title)
272    }
273
274    /// A warning.
275    pub fn warning(title: impl Into<String>) -> Self {
276        Self::new(ToastKind::Warning, title)
277    }
278
279    /// An error.
280    pub fn error(title: impl Into<String>) -> Self {
281        Self::new(ToastKind::Error, title)
282    }
283
284    /// Adds a second line.
285    pub fn body(mut self, body: impl Into<String>) -> Self {
286        self.body = Some(body.into());
287        self
288    }
289
290    /// Where opening the notification goes.
291    pub fn url(mut self, url: impl Into<String>) -> Self {
292        self.url = Some(url.into());
293        self
294    }
295
296    /// Adds a link or button.
297    pub fn action(mut self, action: ToastAction) -> Self {
298        self.actions.push(action);
299        self
300    }
301
302    /// Adds a link: `.link("Invoice", "/invoices/7")`.
303    pub fn link(self, label: impl Into<String>, url: impl Into<String>) -> Self {
304        self.action(ToastAction::link(label, url))
305    }
306
307    /// Stores `key` next to the message, for the app's own pages
308    /// (`notification.data.order_id`). The message's own keys win.
309    pub fn with(mut self, key: &str, value: impl Serialize) -> Self {
310        self.extra.insert(
311            key.to_owned(),
312            serde_json::to_value(value).unwrap_or(Value::Null),
313        );
314        self
315    }
316}
317
318impl From<DatabaseMessage> for Value {
319    fn from(message: DatabaseMessage) -> Self {
320        serde_json::to_value(message).unwrap_or(Value::Null)
321    }
322}
323
324/// A stored notification.
325#[derive(Debug, Clone, Serialize)]
326#[non_exhaustive]
327pub struct DatabaseNotification {
328    /// The `notifications` row id.
329    pub id: i64,
330    /// The notification's [`Notification::kind`].
331    pub kind: String,
332    /// What [`Notification::to_database`] returned.
333    pub data: Value,
334    /// When it was marked read; `None` while unread.
335    pub read_at: Option<DateTime>,
336    /// When it was stored.
337    pub created_at: DateTime,
338}
339
340impl DatabaseNotification {
341    /// The stored [`DatabaseMessage`], when `to_database` returned one (its
342    /// data has a `title`).
343    pub fn message(&self) -> Option<DatabaseMessage> {
344        serde_json::from_value(self.data.clone()).ok()
345    }
346}
347
348/// Wakes the open notification streams (`/notifications/stream`) of one
349/// user when something changed for them in this process; streams also look
350/// at the table every few seconds, for changes made by other servers or by
351/// `queue:work`. It also carries the app's own events
352/// ([`AppState::broadcast`]) to the streams open in this process.
353pub(crate) struct Hub {
354    tx: tokio::sync::broadcast::Sender<Signal>,
355}
356
357#[derive(Debug, Clone, PartialEq, Eq)]
358pub(crate) enum Signal {
359    /// Something changed for this user.
360    User(i64),
361    /// An app event for one user's streams, or (`None`) for every stream.
362    Event(Option<i64>, Arc<Broadcast>),
363    /// The server is shutting down.
364    Stop,
365}
366
367/// An app event on its way to the open pages: its DOM event name and its
368/// data, already JSON.
369#[derive(Debug, PartialEq, Eq)]
370pub(crate) struct Broadcast {
371    pub(crate) event: String,
372    pub(crate) data: String,
373}
374
375impl Hub {
376    pub(crate) fn new() -> Self {
377        Self {
378            tx: tokio::sync::broadcast::channel(256).0,
379        }
380    }
381
382    /// Tells `user_id`'s streams to look now.
383    pub(crate) fn touch(&self, user_id: i64) {
384        let _ = self.tx.send(Signal::User(user_id));
385    }
386
387    /// Sends an app event to `user_id`'s streams, or to every stream.
388    pub(crate) fn event(&self, user_id: Option<i64>, event: Broadcast) {
389        let _ = self.tx.send(Signal::Event(user_id, Arc::new(event)));
390    }
391
392    /// Ends every stream, so a graceful shutdown doesn't wait for them.
393    pub(crate) fn stop(&self) {
394        let _ = self.tx.send(Signal::Stop);
395    }
396
397    pub(crate) fn subscribe(&self) -> tokio::sync::broadcast::Receiver<Signal> {
398        self.tx.subscribe()
399    }
400}
401
402pub(crate) type ChannelFn = Arc<
403    dyn Fn(
404            AppState,
405            Recipient,
406            Value,
407        ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result> + Send>>
408        + Send
409        + Sync,
410>;
411
412pub(crate) fn channel_fn<F, Fut>(send: F) -> ChannelFn
413where
414    F: Fn(Recipient, Value, AppState) -> Fut + Send + Sync + 'static,
415    Fut: std::future::Future<Output = Result> + Send + 'static,
416{
417    Arc::new(move |state, to, message| Box::pin(send(to, message, state)))
418}
419
420/// Delivers a message to one of the app's channels from the queue.
421#[derive(Serialize, Deserialize)]
422pub(crate) struct SendToChannel {
423    channel: String,
424    to: Recipient,
425    message: Value,
426}
427
428impl Job for SendToChannel {
429    const NAME: &'static str = "renox.send-to-channel";
430    const MAX_ATTEMPTS: u32 = 5;
431
432    async fn handle(self, ctx: JobContext) -> Result {
433        let send = ctx.state.channel(&self.channel)?;
434        send(ctx.state.clone(), self.to, self.message).await
435    }
436}
437
438/// The channels a notification is delivered to, in delivery order: the
439/// database row first, so a failure there doesn't leave a sent message
440/// behind that a retry would send again; mail last.
441fn ordered(notification: &impl Notification, to: &Recipient) -> Vec<Channel> {
442    let mut channels = notification.channels(to);
443    channels.sort_by_key(|channel| match channel {
444        Channel::Database => 0,
445        Channel::Custom(_) => 1,
446        _ => 2,
447    });
448    channels
449}
450
451/// Whether `name` can be a DOM event name sent to pages: letters, digits
452/// and `-`, `_`, `:`, `.` (what `hx-trigger` and Alpine's `x-on` can name).
453fn event_name(name: &str) -> Result<&str> {
454    let valid = !name.is_empty()
455        && name.len() <= 100
456        && name
457            .chars()
458            .all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | ':' | '.'));
459    if valid {
460        Ok(name)
461    } else {
462        Err(anyhow!(
463            "`{name}` can't be a broadcast event's name: use letters, digits, `-`, `_`, `:` and `.`"
464        )
465        .into())
466    }
467}
468
469impl AppState {
470    /// Sends the DOM event `event`, with `data` as its `detail`, to every
471    /// page open in this process with a notification stream (the UI kit's
472    /// `notification_bell` or `event_stream()`, which need
473    /// `Auth::new().notifications()`), whoever is logged in there.
474    ///
475    /// Pages listen as to any DOM event on `document`: htmx with
476    /// `hx-trigger="order-updated from:document"`, Alpine with
477    /// `x-on:order-updated.document="…"`, a script with
478    /// `document.addEventListener("order-updated", e => e.detail)`. The
479    /// event `renox:toast` with `{"toasts": [toast]}` shows toasts.
480    ///
481    /// It is fire-and-forget: nothing is stored, and pages that aren't
482    /// connected at that moment (closed, reconnecting, open on another
483    /// server, or a `queue:work` process sent it) never get it. Use a
484    /// database notification ([`AppState::notify`]) for what must arrive.
485    ///
486    /// ```
487    /// # use renox::prelude::*;
488    /// # fn demo(state: AppState, user: User) -> Result {
489    /// state.broadcast("order-updated", json!({ "id": 7, "status": "paid" }))?;
490    /// // Only Ana's pages (any tab, any device on this server):
491    /// state.broadcast_to(user.id, "renox:toast", json!({ "toasts": [Toast::info("Export ready")] }))?;
492    /// # Ok(()) }
493    /// ```
494    pub fn broadcast(&self, event: &str, data: impl Serialize) -> Result {
495        self.send_broadcast(None, event, data)
496    }
497
498    /// [`AppState::broadcast`] to the pages of one user (by id) only.
499    pub fn broadcast_to(&self, user_id: i64, event: &str, data: impl Serialize) -> Result {
500        self.send_broadcast(Some(user_id), event, data)
501    }
502
503    fn send_broadcast(&self, user_id: Option<i64>, event: &str, data: impl Serialize) -> Result {
504        let event = event_name(event)?.to_owned();
505        let data = serde_json::to_value(&data).map_err(anyhow::Error::from)?;
506        let sent = crate::SentBroadcast {
507            user_id,
508            event,
509            data,
510        };
511        if self.fakes.record_broadcast(sent.clone()) {
512            return Ok(());
513        }
514        self.notification_hub.event(
515            user_id,
516            Broadcast {
517                event: sent.event,
518                data: sent.data.to_string(),
519            },
520        );
521        Ok(())
522    }
523
524    fn channel(&self, name: &str) -> Result<ChannelFn> {
525        self.channels.get(name).cloned().ok_or_else(|| {
526            anyhow!("no `{name}` notification channel: register it with `App::channel`").into()
527        })
528    }
529
530    async fn store_notification(&self, to: &Recipient, notification: &impl Notification) -> Result {
531        // Only users have an in-app list.
532        let Some(user) = &to.user else { return Ok(()) };
533        crate::db::sql(
534            "INSERT INTO notifications (user_id, kind, data, created_at) VALUES (?, ?, ?, ?)",
535        )
536        .bind(user.id)
537        .bind(notification.kind())
538        .bind(
539            with_locale(to.locale().as_deref(), || {
540                notification.to_database(to, self)
541            })?
542            .to_string(),
543        )
544        .bind(now())
545        .execute(&self.db)
546        .await?;
547        self.notification_hub.touch(user.id);
548        Ok(())
549    }
550
551    /// Delivers `notification` now to a user (`&user`) or to addresses (a
552    /// [`Recipient`]), on each of its channels. The database row is written
553    /// first and mail sent last, so a failure doesn't leave a message behind
554    /// that a retry would send again.
555    pub async fn notify(
556        &self,
557        to: impl Into<Recipient>,
558        notification: &impl Notification,
559    ) -> Result {
560        let to = &to.into();
561        if self.fakes.record_notification(notification.kind(), to) {
562            return Ok(());
563        }
564        let locale = to.locale();
565        let locale = locale.as_deref();
566        for channel in ordered(notification, to) {
567            match channel {
568                Channel::Database => self.store_notification(to, notification).await?,
569                Channel::Custom(name) => {
570                    let send = self.channel(name)?;
571                    let message = with_locale(locale, || notification.to_channel(name, to, self))?;
572                    send(self.clone(), to.clone(), message).await?;
573                }
574                Channel::Mail => {
575                    let mail = with_locale(locale, || notification.to_mail(to, self))?;
576                    self.mailer.send(mail).await?;
577                }
578            }
579        }
580        Ok(())
581    }
582
583    /// Like `notify`, but mail and the app's channels are sent by queue
584    /// workers, each channel as its own job with its own retries. The
585    /// messages are built now; the database row is written now.
586    pub async fn notify_later(
587        &self,
588        to: impl Into<Recipient>,
589        notification: &impl Notification,
590    ) -> Result {
591        let to = to.into();
592        if self.fakes.record_notification(notification.kind(), &to) {
593            return Ok(());
594        }
595        let locale = to.locale();
596        let locale = locale.as_deref();
597        for channel in ordered(notification, &to) {
598            match channel {
599                Channel::Database => self.store_notification(&to, notification).await?,
600                Channel::Custom(name) => {
601                    self.channel(name)?; // fail now for an unknown channel
602                    let message = with_locale(locale, || notification.to_channel(name, &to, self))?;
603                    self.dispatch(SendToChannel {
604                        channel: name.to_owned(),
605                        to: to.clone(),
606                        message,
607                    })
608                    .await?;
609                }
610                Channel::Mail => {
611                    let mail = with_locale(locale, || notification.to_mail(&to, self))?;
612                    self.queue_mail(mail).await?;
613                }
614            }
615        }
616        Ok(())
617    }
618}
619
620pub(super) fn from_row(row: &crate::db::Row) -> Result<DatabaseNotification> {
621    let data: String = row.try_get("data")?;
622    Ok(DatabaseNotification {
623        id: row.try_get("id")?,
624        kind: row.try_get("kind")?,
625        data: serde_json::from_str(&data).unwrap_or(Value::Null),
626        read_at: row.try_get("read_at")?,
627        created_at: row.try_get("created_at")?,
628    })
629}
630
631impl User {
632    /// The user's notifications, newest first.
633    pub async fn notifications(&self, db: &Db, limit: u32) -> Result<Vec<DatabaseNotification>> {
634        let rows = crate::db::sql(
635            "SELECT id, kind, data, read_at, created_at FROM notifications \
636             WHERE user_id = ? ORDER BY id DESC LIMIT ?",
637        )
638        .bind(self.id)
639        .bind(i64::from(limit))
640        .fetch_all(db)
641        .await?;
642        rows.iter().map(from_row).collect()
643    }
644
645    /// The user's notifications older than the one with id `before`,
646    /// newest first: the next page after a list ending at `before`.
647    pub async fn notifications_before(
648        &self,
649        db: &Db,
650        before: i64,
651        limit: u32,
652    ) -> Result<Vec<DatabaseNotification>> {
653        let rows = crate::db::sql(
654            "SELECT id, kind, data, read_at, created_at FROM notifications \
655             WHERE user_id = ? AND id < ? ORDER BY id DESC LIMIT ?",
656        )
657        .bind(self.id)
658        .bind(before)
659        .bind(i64::from(limit))
660        .fetch_all(db)
661        .await?;
662        rows.iter().map(from_row).collect()
663    }
664
665    /// One of the user's notifications, or `None` if it isn't theirs.
666    pub async fn notification(&self, db: &Db, id: i64) -> Result<Option<DatabaseNotification>> {
667        let row = crate::db::sql(
668            "SELECT id, kind, data, read_at, created_at FROM notifications \
669             WHERE id = ? AND user_id = ?",
670        )
671        .bind(id)
672        .bind(self.id)
673        .fetch_optional(db)
674        .await?;
675        row.as_ref().map(from_row).transpose()
676    }
677
678    /// The user's unread notifications, newest first.
679    pub async fn unread_notifications(&self, db: &Db) -> Result<Vec<DatabaseNotification>> {
680        let rows = crate::db::sql(
681            "SELECT id, kind, data, read_at, created_at FROM notifications \
682             WHERE user_id = ? AND read_at IS NULL ORDER BY id DESC",
683        )
684        .bind(self.id)
685        .fetch_all(db)
686        .await?;
687        rows.iter().map(from_row).collect()
688    }
689
690    /// How many of the user's notifications are unread.
691    pub async fn unread_notification_count(&self, db: &Db) -> Result<i64> {
692        Ok(crate::db::sql(
693            "SELECT COUNT(*) FROM notifications WHERE user_id = ? AND read_at IS NULL",
694        )
695        .bind(self.id)
696        .scalar(db)
697        .await?)
698    }
699
700    /// Marks one of the user's notifications read; returns whether it was theirs.
701    pub async fn mark_notification_read(&self, db: &Db, id: i64) -> Result<bool> {
702        let done = crate::db::sql(
703            "UPDATE notifications SET read_at = COALESCE(read_at, ?) WHERE id = ? AND user_id = ?",
704        )
705        .bind(now())
706        .bind(id)
707        .bind(self.id)
708        .execute(db)
709        .await?;
710        Ok(done > 0)
711    }
712
713    /// Marks one of the user's notifications unread again; returns whether
714    /// it was theirs.
715    pub async fn mark_notification_unread(&self, db: &Db, id: i64) -> Result<bool> {
716        let done =
717            crate::db::sql("UPDATE notifications SET read_at = NULL WHERE id = ? AND user_id = ?")
718                .bind(id)
719                .bind(self.id)
720                .execute(db)
721                .await?;
722        Ok(done > 0)
723    }
724
725    /// Deletes one of the user's notifications; returns whether it was theirs.
726    pub async fn delete_notification(&self, db: &Db, id: i64) -> Result<bool> {
727        let done = crate::db::sql("DELETE FROM notifications WHERE id = ? AND user_id = ?")
728            .bind(id)
729            .bind(self.id)
730            .execute(db)
731            .await?;
732        Ok(done > 0)
733    }
734
735    /// Deletes all the user's notifications; returns how many there were.
736    pub async fn delete_notifications(&self, db: &Db) -> Result<u64> {
737        Ok(
738            crate::db::sql("DELETE FROM notifications WHERE user_id = ?")
739                .bind(self.id)
740                .execute(db)
741                .await?,
742        )
743    }
744
745    /// Marks all the user's unread notifications read; returns how many there were.
746    pub async fn mark_all_notifications_read(&self, db: &Db) -> Result<u64> {
747        let done = crate::db::sql(
748            "UPDATE notifications SET read_at = ? WHERE user_id = ? AND read_at IS NULL",
749        )
750        .bind(now())
751        .bind(self.id)
752        .execute(db)
753        .await?;
754        Ok(done)
755    }
756}
757
758/// Deletes notifications read more than `age` ago; returns how many. Unread
759/// ones stay, however old. `rnx notifications:prune` (from the `Auth` module)
760/// runs it with `--days` (30 by default); schedule it, e.g. daily, since
761/// nothing else removes read notifications while their user exists.
762pub async fn prune_read_notifications(db: &Db, age: std::time::Duration) -> Result<u64> {
763    let before = now() - chrono::Duration::from_std(age).unwrap_or_default();
764    Ok(
765        crate::db::sql("DELETE FROM notifications WHERE read_at < ?")
766            .bind(before)
767            .execute(db)
768            .await?,
769    )
770}