tuwunel_service/migrations/
fix_referencedevents_missing_sep.rs1use std::{cmp::max, str};
2
3use futures::StreamExt;
4use tuwunel_core::{
5 Result, debug, info,
6 utils::{ReadyExt, stream::TryExpect},
7 warn,
8};
9use tuwunel_database::SEP;
10
11use crate::Services;
12
13pub(super) async fn fix_referencedevents_missing_sep(services: &Services) -> Result {
14 warn!("Fixing missing record separator between room_id and event_id in referencedevents");
15
16 let db = &services.db;
17 let cork = db.cork_and_sync();
18
19 let referencedevents = db["referencedevents"].clone();
20
21 let totals: (usize, usize) = (0, 0);
22 let (total, fixed) = referencedevents
23 .raw_stream()
24 .expect_ok()
25 .enumerate()
26 .ready_fold(totals, |mut a, (i, (key, val))| {
27 debug_assert!(val.is_empty(), "expected no value");
28
29 services.server.progress.advance();
30
31 let has_sep = key.contains(&SEP);
32
33 if !has_sep {
34 let key_str = str::from_utf8(key).expect("key not utf-8");
35 let room_id_len = key_str.find('$').expect("missing '$' in key");
36 let (room_id, event_id) = key_str.split_at(room_id_len);
37
38 debug!(?a, "fixing {room_id}, {event_id}");
39
40 let new_key = (room_id, event_id);
41 referencedevents.put_raw(new_key, val);
42 referencedevents.remove(key);
43 }
44
45 a.0 = max(i, a.0);
46 a.1 = a.1.saturating_add((!has_sep).into());
47 a
48 })
49 .await;
50
51 drop(cork);
52 info!(?total, ?fixed, "Fixed missing record separators in 'referencedevents'.");
53
54 db["global"].insert(b"fix_referencedevents_missing_sep", []);
55 referencedevents.sort()
56}