From 7f90cc6abc79288b20f7baf682217dee332e0369 Mon Sep 17 00:00:00 2001 From: mdecimus Date: Sat, 7 Jun 2025 09:20:49 +0200 Subject: [PATCH] Migrate messages not queued for delivery --- crates/migration/src/queue.rs | 20 +++++++++++++++++--- 1 file changed, 17 insertions(+), 3 deletions(-) diff --git a/crates/migration/src/queue.rs b/crates/migration/src/queue.rs index e5568c09..8dfb902b 100644 --- a/crates/migration/src/queue.rs +++ b/crates/migration/src/queue.rs @@ -35,7 +35,6 @@ pub(crate) async fn migrate_queue(server: &Server) -> trc::Result<()> { ))); let mut queue_ids = AHashSet::new(); - server .store() .iterate( @@ -49,7 +48,22 @@ pub(crate) async fn migrate_queue(server: &Server) -> trc::Result<()> { .await .caused_by(trc::location!())?; - let count = queue_ids.len(); + let from_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(0))); + let to_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(u64::MAX))); + server + .store() + .iterate( + IterateParams::new(from_key, to_key).ascending().no_values(), + |key, _| { + queue_ids.insert(key.deserialize_be_u64(0)?); + + Ok(true) + }, + ) + .await + .caused_by(trc::location!())?; + + let mut count = 0; for queue_id in queue_ids { match server @@ -96,7 +110,7 @@ pub(crate) async fn migrate_queue(server: &Server) -> trc::Result<()> { .serialize() .caused_by(trc::location!())?, ); - + count += 1; server .store() .write(batch.build_all())