tuwunel_service/migrations/
retroactively_fix_bad_data_from_roomuserid_joined.rs1use futures::StreamExt;
2use ruma::events::room::member::MembershipState;
3use tuwunel_core::{
4 Result, debug_info, info,
5 matrix::PduCount,
6 utils::{ReadyExt, stream::BroadbandExt},
7 warn,
8};
9
10use crate::Services;
11
12pub(super) async fn retroactively_fix_bad_data_from_roomuserid_joined(
13 services: &Services,
14) -> Result {
15 warn!("Retroactively fixing bad data from broken roomuserid_joined");
16
17 let db = &services.db;
18 let _cork = db.cork_and_sync();
19
20 services
21 .metadata
22 .iter_ids()
23 .for_each(async |room_id| {
24 debug_info!(%room_id, "Fixing room");
25
26 services
27 .state_cache
28 .room_members(room_id)
29 .map(ToOwned::to_owned)
30 .broad_filter_map(async |user_id| {
31 let member = services
33 .state_accessor
34 .get_member(room_id, &user_id)
35 .await
36 .ok()?;
37
38 Some((user_id, member.membership))
39 })
40 .ready_for_each(|(user_id, membership)| {
41 let count = services.globals.next_count();
42
43 match membership {
44 | MembershipState::Join => services.state_cache.mark_as_joined(
45 &user_id,
46 room_id,
47 PduCount::Normal(*count),
48 ),
49 | _ => services.state_cache.mark_as_left(
50 &user_id,
51 room_id,
52 PduCount::Normal(*count),
53 ),
54 }
55 })
56 .await;
57
58 services
59 .state_cache
60 .update_joined_count(room_id)
61 .await;
62 })
63 .await;
64
65 info!("Finished fixing");
66
67 db["global"].insert(b"retroactively_fix_bad_data_from_roomuserid_joined", []);
68 Ok(())
69}