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, math::u64_from_usize_saturating, 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 let total = services.metadata.iter_ids().count().await;
22
23 services
24 .server
25 .progress
26 .expect_total(u64_from_usize_saturating(total));
27
28 services
29 .metadata
30 .iter_ids()
31 .for_each(async |room_id| {
32 debug_info!(%room_id, "Fixing room");
33
34 services
35 .state_cache
36 .room_members(room_id)
37 .map(ToOwned::to_owned)
38 .broad_filter_map(async |user_id| {
39 let member = services
41 .state_accessor
42 .get_member(room_id, &user_id)
43 .await
44 .ok()?;
45
46 Some((user_id, member.membership))
47 })
48 .ready_for_each(|(user_id, membership)| {
49 let count = services.globals.next_count();
50
51 match membership {
52 | MembershipState::Join => services.state_cache.mark_as_joined(
53 &user_id,
54 room_id,
55 PduCount::Normal(*count),
56 ),
57 | _ => services.state_cache.mark_as_left(
58 &user_id,
59 room_id,
60 PduCount::Normal(*count),
61 ),
62 }
63 })
64 .await;
65
66 services
67 .state_cache
68 .update_joined_count(room_id)
69 .await;
70
71 services.server.progress.advance();
72 })
73 .await;
74
75 info!("Finished fixing");
76
77 db["global"].insert(b"retroactively_fix_bad_data_from_roomuserid_joined", []);
78 Ok(())
79}