Skip to main content

tuwunel_api/client/
events.rs

1use 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
27/// One user's listen on one room's event stream.
28///
29/// A non-member's listen is a peek, which sees only what a room preview may
30/// show.
31struct Listen<'a> {
32	services: &'a Services,
33	sender_user: &'a UserId,
34	peeking: bool,
35}
36
37/// GET `/_matrix/client/v3/events`
38pub(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	// The endpoint listens for new events, so a stream without a token starts now.
74	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		// Any new event answers, hidden or not, so no later wake rescans the window.
116		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
136/// The page for a window of new events, or `None` while the window is empty.
137///
138/// The peek keeps the event it looked at, so the page scans the window once.
139async 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
158/// The events of a window the user may see, as one page of the stream.
159///
160/// A full page ends at its last event. A short one scanned the whole window,
161/// hidden events included, so it ends at `next_batch`. An empty page starts
162/// where the stream did.
163async 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
192/// Keeps an event the listener may see.
193///
194/// A member gets the general history-visibility rule. A peek gets only what a
195/// room preview may show: events sent while the room was world-readable, and the
196/// event that made it so.
197async 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}