tuwunel_service/migrations/
rebuild_roomid_tscount_pducount.rs1use 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}