Skip to main content

tuwunel_service/migrations/
rebuild_roomid_tscount_pducount.rs

1use ruma::{MilliSecondsSinceUnixEpoch, OwnedRoomId, RoomId};
2use serde::Deserialize;
3use tuwunel_core::{
4	Result, info,
5	matrix::pdu::RawPduId,
6	utils::{ReadyExt, stream::TryIgnore},
7	warn,
8};
9
10use crate::{Services, rooms::timeline::bias_count};
11
12#[derive(Deserialize)]
13struct PduRoomTs {
14	room_id: OwnedRoomId,
15	origin_server_ts: MilliSecondsSinceUnixEpoch,
16}
17
18pub(super) async fn rebuild_roomid_tscount_pducount(services: &Services) -> Result {
19	let db = &services.db;
20	let cork = db.cork_and_sync();
21	let pduid_pdu = db["pduid_pdu"].clone();
22	let roomid_tscount_pducount = db["roomid_tscount_pducount"].clone();
23
24	warn!("Rebuilding roomid_tscount_pducount index for same-timestamp event ordering");
25
26	let count = pduid_pdu
27		.raw_stream()
28		.ignore_err()
29		.ready_fold(0_usize, |count, (key, value)| {
30			let Ok(pdu) = serde_json::from_slice::<PduRoomTs>(value) else {
31				return count;
32			};
33
34			let ts = u64::from(pdu.origin_server_ts.get());
35			let pdu_id = RawPduId::from(key);
36			let count_key = bias_count(pdu_id.count());
37			let room_id: &RoomId = &pdu.room_id;
38
39			roomid_tscount_pducount.put_raw((room_id, ts, count_key), pdu_id.count());
40
41			count.saturating_add(1)
42		})
43		.await;
44
45	drop(cork);
46	info!(%count, "Rebuilt roomid_tscount_pducount index");
47
48	db["global"].insert(b"rebuild_roomid_tscount_pducount", []);
49	roomid_tscount_pducount.sort()
50}