Skip to main content

tuwunel_api/client/sync/v5/extensions/
account_data.rs

1use std::collections::BTreeMap;
2
3use futures::TryStreamExt;
4use ruma::{
5	OwnedRoomId,
6	api::client::sync::sync_events::v5::response::AccountData,
7	events::{AnyRawAccountDataEvent, AnyRoomAccountDataEvent},
8	serde::Raw,
9};
10use tuwunel_core::{Result, extract_variant, utils::TryReadyExt};
11
12use super::{Connection, SyncInfo, Window, selector};
13use crate::client::{is_empty_account_data_event, sync::v5::range::Results};
14
15#[tracing::instrument(name = "account_data", level = "trace", skip_all)]
16pub(super) async fn collect(
17	SyncInfo { services, sender_user, .. }: SyncInfo<'_>,
18	conn: &Connection,
19) -> Result<AccountData> {
20	let globalsince = conn.globalsince;
21	let global = services
22		.account_data
23		.changes_since_fallible(None, sender_user, globalsince, Some(conn.next_batch))
24		.ready_try_filter_map(|event| Ok(extract_variant!(event, AnyRawAccountDataEvent::Global)))
25		.ready_try_filter(move |event| globalsince != 0 || !is_empty_account_data_event(event))
26		.try_collect()
27		.await?;
28
29	Ok(AccountData { global, rooms: Default::default() })
30}
31
32pub(super) fn collect_ranges(
33	conn: &Connection,
34	window: &Window,
35	ranges: &mut Results,
36) -> BTreeMap<OwnedRoomId, Vec<Raw<AnyRoomAccountDataEvent>>> {
37	let implicit = conn
38		.extensions
39		.account_data
40		.lists
41		.as_deref()
42		.map(<[_]>::iter);
43
44	let explicit = conn
45		.extensions
46		.account_data
47		.rooms
48		.as_deref()
49		.map(<[_]>::iter);
50
51	selector(conn, window, implicit, explicit)
52		.filter_map(|room_id| {
53			ranges
54				.take_account_data(room_id)
55				.map(|events| (room_id.to_owned(), events))
56		})
57		.collect()
58}