diff --git a/src/service/migrations.rs b/src/service/migrations.rs index 747da65ac..0a5881063 100644 --- a/src/service/migrations.rs +++ b/src/service/migrations.rs @@ -2,6 +2,8 @@ cmp, collections::{BTreeMap, HashMap}, future::ready, + sync::Arc, + time::Instant, }; use conduwuit::{ @@ -15,7 +17,7 @@ }, warn, }; -use database::Json; +use database::{Database, Json}; use futures::{FutureExt, StreamExt, TryStreamExt}; use itertools::Itertools; use ruma::{ @@ -30,7 +32,7 @@ }; use serde_json::value::to_raw_value; -use crate::{Services, media, rooms::short::ShortStateHash}; +use crate::{Services, media, rooms, rooms::short::ShortStateHash}; /// The current schema version. /// - If database is opened at greater version we reject with error. The @@ -125,7 +127,7 @@ async fn migrate(services: &Services) -> Result<()> { .map_err(|e| err!("Failed to run SHA256 media migration: {e}"))?; } else if config.media_startup_check { info!("Starting media startup integrity check."); - let now = std::time::Instant::now(); + let now = Instant::now(); media::migrations::checkup_sha256_media(services) .await .map_err(|e| err!("Failed to verify media integrity: {e}"))?; @@ -263,9 +265,11 @@ async fn migrate(services: &Services) -> Result<()> { } if db["global"].get(UNEMBED_UNSIGNED_INFO).await.is_not_found() { - unembed_unsigned_info(services) - .await - .map_err(|e| err!("Failed to run 'unembed_unsigned_info' migration: {e:?}"))?; + services.server.runtime().spawn(unembed_unsigned_info( + services.db.clone(), + services.rooms.metadata.clone(), + services.rooms.timeline.clone(), + )); } assert_eq!( @@ -916,12 +920,25 @@ async fn obliterate_roomsynctoken_shortstatehash_with_extreme_prejudice( } const UNEMBED_UNSIGNED_INFO: &str = "unembed_unsigned_info"; -async fn unembed_unsigned_info(services: &Services) -> Result { +#[tracing::instrument(name = "unembed_unsigned_info", skip_all)] +async fn unembed_unsigned_info( + db: Arc, + metadata: Arc, + timeline: Arc, +) -> Result { type Map = BTreeMap>; + info!( + "Starting background migration. This may take a long time, and may cause your server to \ + lag. Please avoid restarting until the migration is complete." + ); - let mut all_rooms = services.rooms.metadata.iter_ids(); + let start = Instant::now(); + let mut total_migrated = 0_usize; + let mut all_rooms = metadata.iter_ids(); while let Some(room_id) = all_rooms.next().await { - let mut pdus = std::pin::pin!(services.rooms.timeline.all_pdus(&room_id)); + info!("Migrating PDUs in {room_id}"); + let mut pdus = std::pin::pin!(timeline.all_pdus(&room_id)); + let mut migrated = 0_usize; while let Some((_, pdu)) = pdus.next().await { let mut pdu = pdu.clone(); let mut unsigned: Map = pdu @@ -951,20 +968,22 @@ async fn unembed_unsigned_info(services: &Services) -> Result { } } pdu.unsigned = Some(to_raw_value(&unsigned).expect("unsigned must be valid JSON")); - let pdu_id = services - .rooms - .timeline + let pdu_id = timeline .get_pdu_id(pdu.event_id()) .await .expect("PDU must have an id"); - services - .rooms - .timeline + timeline .replace_pdu(&pdu_id, &to_canonical_object(pdu).expect("PDU must be valid JSON")) .await .expect("must be able to replace PDU after migrating unsigned data"); + migrated = migrated.saturating_add(1); } + total_migrated = total_migrated.saturating_add(migrated); + info!(elapsed=?start.elapsed(), "Migrated {migrated} PDUs in {room_id}"); } + db["global"].insert(UNEMBED_UNSIGNED_INFO, []); + info!(elapsed=?start.elapsed(), total_migrated_pdus=total_migrated, "Finished migration."); + Ok(()) }