tuwunel_api/client/sync/v5/extensions/
account_data.rs1use 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}