Skip to main content

tuwunel_service/migrations/
retroactively_fix_bad_data_from_roomuserid_joined.rs

1use 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					// A member with no resolved member event is left untouched.
32					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}