tuwunel_api/client/
events.rs1use std::iter::once;
2
3use axum::extract::State;
4use futures::{Stream, StreamExt, future::ok, pin_mut};
5use ruma::{
6 UserId,
7 api::client::peeking::listen_to_new_events::v3::{Request, Response},
8};
9use tokio::time::{Duration, Instant, timeout_at};
10use tuwunel_core::{
11 Err, Event, Result, err,
12 matrix::PduCount,
13 utils::{
14 BoolExt, OptionExt,
15 future::OptionFutureExt,
16 result::FlatOk,
17 stream::{IterStream, ReadyExt, WidebandExt},
18 },
19};
20use tuwunel_service::{Services, rooms::timeline::PdusIterItem};
21
22use super::visibility_filter;
23use crate::Ruma;
24
25const EVENT_LIMIT: usize = 50;
26
27struct Listen<'a> {
32 services: &'a Services,
33 sender_user: &'a UserId,
34 peeking: bool,
35}
36
37pub(crate) async fn events_route(
39 State(services): State<crate::State>,
40 body: Ruma<Request>,
41) -> Result<Response> {
42 let sender_user = body.sender_user();
43
44 let timeout = body
45 .body
46 .timeout
47 .as_ref()
48 .map(Duration::as_millis)
49 .map(TryInto::try_into)
50 .flat_ok()
51 .unwrap_or(services.config.client_sync_timeout_default)
52 .max(services.config.client_sync_timeout_min)
53 .min(services.config.client_sync_timeout_max);
54
55 let room_id = body.room_id.as_ref();
56
57 let peeking = services
58 .state_cache
59 .is_joined(sender_user, room_id)
60 .await
61 .is_false();
62
63 if peeking
64 && services
65 .state_accessor
66 .is_world_readable(room_id)
67 .await
68 .is_false()
69 {
70 return Err!(Request(Forbidden("No room preview available.")));
71 }
72
73 let from = body
75 .body
76 .from
77 .as_deref()
78 .map(str::parse)
79 .transpose()
80 .map_err(|_| err!(Request(InvalidParam("Invalid `from` token."))))?
81 .map_async(ok)
82 .unwrap_or_else_async(async || {
83 services
84 .globals
85 .wait_pending()
86 .await
87 .map(PduCount::Normal)
88 })
89 .await?;
90
91 let listen = Listen {
92 services: &services,
93 sender_user,
94 peeking,
95 };
96
97 let stop_at = Instant::now()
98 .checked_add(Duration::from_millis(timeout))
99 .expect("configuration must limit maximum timeout");
100
101 loop {
102 let watchers = services
103 .sync
104 .watch(sender_user, body.sender_device.as_deref(), once(room_id).stream())
105 .await;
106
107 let next_batch = services.globals.wait_pending().await?;
108
109 let window = services
110 .timeline
111 .pdus(Some(sender_user), room_id, Some(from))
112 .ready_filter_map(Result::ok)
113 .ready_take_while(|(count, _)| PduCount::Normal(next_batch).ge(count));
114
115 if let Some(response) = window_page(&listen, window, from, next_batch).await {
117 return Ok(response);
118 }
119
120 if timeout_at(stop_at, watchers).await.is_err() || services.server.is_stopping() {
121 return Ok(Response {
122 chunk: Default::default(),
123 start: from.to_string().into(),
124 end: services
125 .server
126 .is_stopping()
127 .is_false()
128 .then_some(next_batch)
129 .as_ref()
130 .map(ToString::to_string),
131 });
132 }
133 }
134}
135
136async fn window_page<Window>(
140 listen: &Listen<'_>,
141 window: Window,
142 from: PduCount,
143 next_batch: u64,
144) -> Option<Response>
145where
146 Window: Stream<Item = PdusIterItem> + Send,
147{
148 let window = window.peekable();
149
150 pin_mut!(window);
151 window.as_mut().peek().await?;
152
153 visible_page(listen, window, from, next_batch)
154 .await
155 .into()
156}
157
158async fn visible_page<Window>(
164 listen: &Listen<'_>,
165 window: Window,
166 from: PduCount,
167 next_batch: u64,
168) -> Response
169where
170 Window: Stream<Item = PdusIterItem> + Send,
171{
172 let (first, last, chunk) = window
173 .wide_filter_map(|item| listen_filter(listen, item))
174 .take(EVENT_LIMIT)
175 .ready_fold((None, None, Vec::new()), |(first, _, mut chunk), (count, pdu)| {
176 chunk.push(pdu.into_format());
177 (first.or(Some(count)), Some(count), chunk)
178 })
179 .await;
180
181 let start = first.unwrap_or(from).to_string().into();
182
183 let end = last
184 .filter(|_| chunk.len().eq(&EVENT_LIMIT))
185 .unwrap_or(PduCount::Normal(next_batch))
186 .to_string()
187 .into();
188
189 Response { start, end, chunk }
190}
191
192async fn listen_filter(listen: &Listen<'_>, item: PdusIterItem) -> Option<PdusIterItem> {
198 if !listen.peeking {
199 return visibility_filter(listen.services, item, listen.sender_user).await;
200 }
201
202 let (_, pdu) = &item;
203
204 listen
205 .services
206 .state_accessor
207 .is_world_readable_at(pdu)
208 .await
209 .then_some(item)
210}