tuwunel_api/client/push/
notifications.rs1use axum::extract::State;
2use futures::StreamExt;
3use ruma::{MilliSecondsSinceUnixEpoch, api::client::push::get_notifications, push::Action};
4use tuwunel_core::{
5 Result, at, err,
6 matrix::{Event, PduId},
7 utils::{
8 stream::{ReadyExt, WidebandExt},
9 string::to_small_string,
10 },
11};
12
13use crate::Ruma;
14
15pub(crate) async fn get_notifications_route(
20 State(services): State<crate::State>,
21 body: Ruma<get_notifications::v3::Request>,
22) -> Result<get_notifications::v3::Response> {
23 use get_notifications::v3::Notification;
24
25 let sender_user = body.sender_user();
26
27 let from = body
28 .body
29 .from
30 .as_deref()
31 .map(str::parse)
32 .transpose()
33 .map_err(|e| err!(Request(InvalidParam("Invalid `from' parameter: {e}"))))?;
34
35 let limit: usize = body
36 .body
37 .limit
38 .map(TryInto::try_into)
39 .transpose()?
40 .unwrap_or(50)
41 .clamp(1, 100);
42
43 let only_highlight = body
44 .body
45 .only
46 .as_deref()
47 .is_some_and(|only| only.contains("highlight"));
48
49 let mut next_token: Option<u64> = None;
50 let notifications = services
51 .pusher
52 .get_notifications(sender_user, from)
53 .ready_filter(|(_, notify)| {
54 !only_highlight || notify.actions.iter().any(Action::is_highlight)
55 })
56 .wide_filter_map(async |(count, notify)| {
57 let pdu_id = PduId {
58 shortroomid: notify.sroomid,
59 count: count.into(),
60 };
61
62 let event = services
63 .timeline
64 .get_pdu_from_id(&pdu_id.into())
65 .await
66 .ok()
67 .filter(|event| !event.is_redacted())?;
68
69 let read = services
70 .pusher
71 .last_notification_read(sender_user, event.room_id())
72 .await
73 .is_ok_and(|last_read| last_read.ge(&count));
74
75 let ts = notify
76 .ts
77 .try_into()
78 .map(MilliSecondsSinceUnixEpoch)
79 .ok()?;
80
81 let notification = Notification {
82 room_id: event.room_id().into(),
83 event: event.into_format(),
84 ts,
85 read,
86 profile_tag: notify.tag,
87 actions: notify.actions,
88 };
89
90 Some((count, notification))
91 })
92 .take(limit)
93 .inspect(|(count, _)| {
94 next_token.replace(*count);
95 })
96 .map(at!(1))
97 .collect::<Vec<_>>()
98 .await;
99
100 Ok(get_notifications::v3::Response {
101 next_token: next_token.map(to_small_string),
102 notifications,
103 })
104}