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 services.server.progress.advance();
31
32 let Ok(pdu) = serde_json::from_slice::<PduRoomTs>(value) else {
33 return count;
34 };
35
36 let ts = u64::from(pdu.origin_server_ts.get());
37 let pdu_id = RawPduId::from(key);
38 let count_key = bias_count(pdu_id.count());
39 let room_id: &RoomId = &pdu.room_id;
40
41 roomid_tscount_pducount.put_raw((room_id, ts, count_key), pdu_id.count());
42
43 count.saturating_add(1)
44 })
45 .await;
46
47 drop(cork);
48 info!(%count, "Rebuilt roomid_tscount_pducount index");
49
50 db["global"].insert(b"rebuild_roomid_tscount_pducount", []);
51 roomid_tscount_pducount.sort()
52}