diff --git a/crates/common/src/manager/backup.rs b/crates/common/src/manager/backup.rs index 752b87e..4c38bc0 100644 --- a/crates/common/src/manager/backup.rs +++ b/crates/common/src/manager/backup.rs @@ -23,6 +23,13 @@ use utils::{UnwrapFailure, codec::leb128::Leb128_}; pub(super) const MAGIC_MARKER: u8 = 123; +// inbuxa: blobs kept under a fixed name instead of a content hash. Nothing +// links to them, so the export names them outright. +const NAMED_BLOBS: &[&[u8]] = &[ + crate::manager::SPAM_CLASSIFIER_KEY, + crate::manager::SPAM_TRAINER_KEY, +]; + #[derive(Debug, Clone, Copy, Hash, PartialEq, Eq)] pub(super) enum Family { Data = 0, @@ -143,15 +150,21 @@ impl Core { .await .failed("Failed to iterate over data store"); - for hash in blobs { + // inbuxa: the trained spam classifier and its trainer state are + // blobs stored under fixed names with no blob link, so the walk + // over links above never reaches them. + let named = NAMED_BLOBS.iter().map(|key| key.to_vec()); + for key in blobs + .into_iter() + .map(|hash| hash.as_slice().to_vec()) + .chain(named) + { if let Some(blob) = blob_store - .get_blob(hash.as_slice(), 0..usize::MAX) + .get_blob(&key, 0..usize::MAX) .await .failed("Failed to get blob") { - writer - .send((hash.as_slice().to_vec(), blob)) - .failed("Failed to send key"); + writer.send((key, blob)).failed("Failed to send key"); } } }), @@ -323,7 +336,13 @@ impl Family { SUBSPACE_REGISTRY_IDX, SUBSPACE_REGISTRY_PK, SUBSPACE_DIRECTORY, - store::SUBSPACE_INBUXA, // inbuxa: masked email + // inbuxa: registry objects the upstream list left out, so an + // export dropped them: archived items (undelete) and spam + // training samples. Their indexes and id counters already + // travel in this family and in `data`, so they ride along. + SUBSPACE_DELETED_ITEMS, + SUBSPACE_SPAM_SAMPLES, + store::SUBSPACE_INBUXA, // inbuxa: the fork's own data (masked email, undelete, policies) ], Family::Changelog => &[SUBSPACE_LOGS], Family::Queue => &[SUBSPACE_QUEUE_MESSAGE, SUBSPACE_QUEUE_EVENT], diff --git a/crates/common/src/manager/boot.rs b/crates/common/src/manager/boot.rs index 65fc2db..9f2f85a 100644 --- a/crates/common/src/manager/boot.rs +++ b/crates/common/src/manager/boot.rs @@ -54,6 +54,13 @@ Options: -o, --console Open the store console -h, --help Print help -V, --version Print version + +An export holds everything in the data and blob stores except short-lived +in-memory state (rate limits, locks, greylisting) and the full-text search +index, which belongs to one search backend. An import into an empty store +queues the index to be rebuilt when the server next starts. EXPORT_TYPES +limits an export to some of: data, registry, blob, changelog, queue, report, +telemetry, tasks. "# ); @@ -256,10 +263,10 @@ impl BootManager { telemetry.enable(); // Parse settings and restore - Box::pin(Core::parse(&mut bootstrap, storage)) - .await - .restore(path) - .await; + let core = Box::pin(Core::parse(&mut bootstrap, storage)).await; + let imported = core.restore(path).await; + // inbuxa: the search index isn't exported; rebuild it + core.queue_reindex(&imported).await; std::process::exit(0); } StoreOp::Console => { diff --git a/crates/common/src/manager/restore.rs b/crates/common/src/manager/restore.rs index b8cf35e..99b3780 100644 --- a/crates/common/src/manager/restore.rs +++ b/crates/common/src/manager/restore.rs @@ -9,15 +9,22 @@ use super::backup::MAGIC_MARKER; use crate::{Core, DATABASE_SCHEMA_VERSION}; use lz4_flex::frame::FrameDecoder; -use registry::schema::enums::CompressionAlgo; +use registry::{ + schema::{ + enums::{CompressionAlgo, TaskStoreMaintenanceType}, + structs::{Task, TaskStatus, TaskStoreMaintenance}, + }, + types::EnumImpl, +}; use std::{ fs::File, io::{BufReader, ErrorKind, Read}, path::{Path, PathBuf}, }; use store::{ - BlobStore, IterateParams, SUBSPACE_BLOBS, SUBSPACE_COUNTER, SUBSPACE_INDEXES, SUBSPACE_QUOTA, - SUBSPACE_REGISTRY_PK, Store, U32_LEN, + BlobStore, IterateParams, SUBSPACE_BLOBS, SUBSPACE_COUNTER, SUBSPACE_INDEXES, + SUBSPACE_PROPERTY, SUBSPACE_QUOTA, SUBSPACE_REGISTRY_PK, SUBSPACE_TELEMETRY_SPAN, Store, + U32_LEN, write::{ AnyClass, AnyKey, BatchBuilder, ValueClass, key::{DeserializeBigEndian, is_node_id_key}, @@ -27,7 +34,9 @@ use types::{collection::Collection, field::Field}; use utils::{UnwrapFailure, failed}; impl Core { - pub async fn restore(&self, src: PathBuf) { + /// Imports an export into an empty store and returns the subspaces it + /// wrote. inbuxa: the caller hands them to [`Core::queue_reindex`]. + pub async fn restore(&self, src: PathBuf) -> Vec { // Backup the core let paths = if src.is_dir() { let mut paths = Vec::new(); @@ -64,6 +73,13 @@ impl Core { std::process::exit(1); } + let mut imported = paths + .iter() + .map(|path| KeyValueReader::new(path).subspace) + .collect::>(); + imported.sort_unstable(); + imported.dedup(); + let mut tasks = Vec::new(); for path in paths { let storage = self.storage.clone(); @@ -76,6 +92,54 @@ impl Core { for task in tasks { task.await.failed("Failed to wait for task"); } + + imported + } + + /// inbuxa: an export never carries the full-text index. It is built by + /// and for one search backend (the SQL stores index into their own + /// tables, the key-value stores into a subspace, external engines keep it + /// themselves), so it would be wrong or unreadable after a move to + /// another one. Instead, an import queues the same reindex tasks an + /// administrator can queue by hand (`reindexAccounts` and + /// `reindexTelemetry` store maintenance), and the server rebuilds the + /// index for whatever search store it is configured with once it starts. + pub async fn queue_reindex(&self, imported: &[u8]) -> Vec { + let mut queued = Vec::new(); + if imported.contains(&SUBSPACE_PROPERTY) { + queued.push(TaskStoreMaintenanceType::ReindexAccounts); + } + if imported.contains(&SUBSPACE_TELEMETRY_SPAN) { + queued.push(TaskStoreMaintenanceType::ReindexTelemetry); + } + if queued.is_empty() { + return queued; + } + + let mut batch = BatchBuilder::new(); + for maintenance_type in &queued { + batch.schedule_task(Task::StoreMaintenance(TaskStoreMaintenance { + maintenance_type: *maintenance_type, + status: TaskStatus::now(), + shard_index: None, + })); + } + self.storage + .data + .write(batch.build_all()) + .await + .failed("Failed to queue the reindex tasks"); + + println!( + "Queued {} to rebuild the search index; it runs when the server starts.", + queued + .iter() + .map(|t| t.as_str()) + .collect::>() + .join(" and ") + ); + + queued } } @@ -125,17 +189,22 @@ async fn restore_file(store: Store, blob_store: BlobStore, path: &Path) { } SUBSPACE_COUNTER | SUBSPACE_QUOTA => { while let Some((key, value)) = reader.next() { - batch.add( - ValueClass::Any(AnyClass { - subspace: reader.subspace, - key, - }), - u64::from_le_bytes( - value - .try_into() - .expect("Failed to deserialize counter/quota"), - ) as i64, - ); + let class = ValueClass::Any(AnyClass { + subspace: reader.subspace, + key, + }); + let value = u64::from_le_bytes( + value + .try_into() + .expect("Failed to deserialize counter/quota"), + ) as i64; + // inbuxa: the SQL stores add a negative amount with an UPDATE, + // which does nothing to a row that isn't there yet, so a + // negative counter vanished on import. Create the row first. + if value < 0 { + batch.add(class.clone(), 0); + } + batch.add(class, value); if batch.is_large_batch() { store .write(batch.build_all()) diff --git a/tests/src/store/import_export.rs b/tests/src/store/import_export.rs index 326bab2..3e95472 100644 --- a/tests/src/store/import_export.rs +++ b/tests/src/store/import_export.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use crate::utils::{ @@ -9,14 +11,22 @@ use crate::utils::{ server::TestServer, temp_dir::TempDir, }; -use ::registry::schema::enums::CompressionAlgo; +use ::registry::schema::{ + enums::{CompressionAlgo, TaskStoreMaintenanceType}, + prelude::ObjectType, + structs::Task, +}; use ahash::AHashSet; -use common::{DATABASE_SCHEMA_VERSION, manager::backup::BackupParams}; +use common::{ + DATABASE_SCHEMA_VERSION, + manager::{SPAM_CLASSIFIER_KEY, SPAM_TRAINER_KEY, backup::BackupParams}, +}; use store::{ rand, write::{ AnyClass, AnyKey, BatchBuilder, BlobLink, BlobOp, Operation, QueueClass, QueueEvent, - RegistryClass, ValueClass, key::KeySerializer, + RegistryClass, TaskQueueClass, ValueClass, + key::{DeserializeBigEndian, KeySerializer}, }, *, }; @@ -167,6 +177,50 @@ pub async fn test(test: &TestServer) { } db.write(batch.build_all()).await.unwrap(); + // inbuxa: registry objects kept outside the registry subspace (archived + // items for undelete, spam training samples, directory entries) and the + // fork's own subspace. Exports used to leave the first two behind. + println!("Creating archived items, spam samples and fork data..."); + let mut batch = BatchBuilder::new(); + for item_id in [1u64, 2, 3] { + for object in [ + ObjectType::ArchivedItem, + ObjectType::SpamTrainingSample, + ObjectType::Account, + ] { + batch.set( + ValueClass::Registry(RegistryClass::Item { + object_id: object as u16, + item_id, + }), + random_bytes(item_id as usize * 64), + ); + } + batch.set( + ValueClass::Any(AnyClass { + subspace: SUBSPACE_INBUXA, + key: [b'U', b'x'] + .into_iter() + .chain(item_id.to_be_bytes()) + .collect(), + }), + random_bytes(32), + ); + } + db.write(batch.build_all()).await.unwrap(); + + // inbuxa: the trained spam classifier lives in blobs with fixed names + let mut named_blobs = Vec::new(); + for key in [SPAM_CLASSIFIER_KEY, SPAM_TRAINER_KEY] { + let data = random_bytes(4096); + test.server + .blob_store() + .put_blob(key, &data, CompressionAlgo::Lz4) + .await + .unwrap(); + named_blobs.push((key, data)); + } + // Create directory data println!("Creating directory data..."); let mut batch = BatchBuilder::new(); @@ -185,6 +239,17 @@ pub async fn test(test: &TestServer) { println!("Calculating store hash..."); let snapshot = Snapshot::new(&db).await; assert!(!snapshot.keys.is_empty(), "Store hash counts are empty",); + for subspace in [ + SUBSPACE_DELETED_ITEMS, + SUBSPACE_SPAM_SAMPLES, + SUBSPACE_INBUXA, + ] { + assert!( + snapshot.keys.iter().any(|k| k.subspace == subspace), + "No test data in subspace {}", + char::from(subspace) + ); + } // Export store println!("Exporting store..."); @@ -210,22 +275,188 @@ pub async fn test(test: &TestServer) { .finalize(), ); db.write(batch.build_all()).await.unwrap(); - test.server.core.restore(temp_dir.path.clone()).await; + for (key, _) in &named_blobs { + test.server.blob_store().delete_blob(key).await.unwrap(); + } + let imported = test.server.core.restore(temp_dir.path.clone()).await; let mut batch = BatchBuilder::new(); batch.clear(ValueClass::NodeId(0)); db.write(batch.build_all()).await.unwrap(); + for subspace in [ + SUBSPACE_DELETED_ITEMS, + SUBSPACE_SPAM_SAMPLES, + SUBSPACE_INBUXA, + ] { + assert!( + imported.contains(&subspace), + "Subspace {} was not exported", + char::from(subspace) + ); + } // Verify hash print!("Verifying store hash..."); snapshot.assert_is_eq(&Snapshot::new(&db).await); + assert_named_blobs(test.server.blob_store(), &named_blobs).await; println!(" GREAT SUCCESS!"); + // inbuxa: import the same export into a fresh store of another backend, + // the way a move from one database to another does it + #[cfg(all(feature = "rocks", feature = "sqlite"))] + cross_backend(test, &db, &temp_dir, &named_blobs).await; + // Destroy store + for (key, _) in &named_blobs { + test.server.blob_store().delete_blob(key).await.unwrap(); + } store_destroy(&db).await; store_assert_is_empty(&db, db.clone().into(), true).await; temp_dir.delete(); } +#[cfg(all(feature = "rocks", feature = "sqlite"))] +async fn cross_backend( + test: &TestServer, + source: &Store, + export: &TempDir, + named_blobs: &[(&[u8], Vec)], +) { + let source_type = std::env::var("STORE").unwrap(); + let target_type = if source_type.eq_ignore_ascii_case("sqlite") { + "RocksDb" + } else { + "Sqlite" + }; + println!("Importing the export into a fresh {target_type} store..."); + + let target_dir = TempDir::new("art_vandelay_cross_backend", true); + let target = Store::build( + crate::utils::storage::build_data_store(target_type, &target_dir.path.to_string_lossy()) + .await, + ) + .await + .unwrap(); + target.create_tables().await.unwrap(); + store_destroy(&target).await; + + let mut core = test.server.core.as_ref().clone(); + core.storage.data = target.clone(); + core.storage.blob = target.clone().into(); + let imported = core.restore(export.path.clone()).await; + + // Counters are stored differently by the SQL and key-value backends, so + // compare their keys here and their values through the counter API. + print!("Verifying {target_type} store hash..."); + Snapshot::new_portable(source) + .await + .assert_is_eq(&Snapshot::new_portable(&target).await); + for subspace in [SUBSPACE_COUNTER, SUBSPACE_QUOTA] { + let mut keys = Vec::new(); + source + .iterate( + IterateParams::new( + AnyKey { + subspace, + key: vec![0u8], + }, + AnyKey { + subspace, + key: vec![u8::MAX; 10], + }, + ) + .no_values(), + |key, _| { + keys.push(key.to_vec()); + Ok(true) + }, + ) + .await + .unwrap(); + for key in keys { + let class = || { + ValueClass::Any(AnyClass { + subspace, + key: key.clone(), + }) + }; + assert_eq!( + source.get_counter(class()).await.unwrap(), + target.get_counter(class()).await.unwrap(), + "Counter mismatch in {} for {key:?}", + char::from(subspace) + ); + } + } + assert_named_blobs(&core.storage.blob, named_blobs).await; + println!(" GREAT SUCCESS!"); + + // The search index isn't exported; the import queues its rebuild + let queued = core.queue_reindex(&imported).await; + let expected = [ + TaskStoreMaintenanceType::ReindexAccounts, + TaskStoreMaintenanceType::ReindexTelemetry, + ]; + assert_eq!(queued, expected); + let mut task_ids = Vec::new(); + target + .iterate( + IterateParams::new( + AnyKey { + subspace: SUBSPACE_TASK_QUEUE, + key: vec![0u8], + }, + AnyKey { + subspace: SUBSPACE_TASK_QUEUE, + key: vec![u8::MAX; 20], + }, + ) + .no_values(), + |key, _| { + if key.deserialize_be_u64(0)? == 0 { + task_ids.push(key.deserialize_be_u64(U64_LEN)?); + } + Ok(true) + }, + ) + .await + .unwrap(); + let mut found = Vec::new(); + for id in task_ids { + match target + .get_value::(ValueKey::from(ValueClass::TaskQueue( + TaskQueueClass::Task { id }, + ))) + .await + .unwrap() + { + Some(Task::StoreMaintenance(task)) => found.push(task.maintenance_type), + other => panic!("Unexpected task {other:?}"), + } + } + found.sort_by_key(|t| *t as u16); + assert_eq!(found, expected, "Queued tasks don't match"); + + store_destroy(&target).await; + drop(core); + drop(target); + target_dir.delete(); +} + +async fn assert_named_blobs(blob_store: &BlobStore, named_blobs: &[(&[u8], Vec)]) { + for (key, data) in named_blobs { + assert_eq!( + blob_store + .get_blob(key, 0..usize::MAX) + .await + .unwrap() + .as_ref(), + Some(data), + "Blob {} was not restored", + String::from_utf8_lossy(key) + ); + } +} + #[derive(Debug, PartialEq, Eq)] struct Snapshot { keys: AHashSet, @@ -240,7 +471,19 @@ struct KeyValue { impl Snapshot { async fn new(db: &Store) -> Self { - let is_sql = db.is_sql(); + Self::build(db, !db.is_sql(), true).await + } + + /// Comparable across backends: no counter values, which the SQL and + /// key-value stores encode differently, and no blobs, which only live in + /// the data store when it doubles as the blob store. + #[cfg(all(feature = "rocks", feature = "sqlite"))] + async fn new_portable(db: &Store) -> Self { + Self::build(db, false, false).await + } + + async fn build(db: &Store, counter_values: bool, with_blobs: bool) -> Self { + let is_sql = !counter_values; let mut keys = AHashSet::new(); @@ -265,7 +508,12 @@ impl Snapshot { (SUBSPACE_QUOTA, !is_sql), (SUBSPACE_REPORT_OUT, true), (SUBSPACE_REPORT_IN, true), + (SUBSPACE_DIRECTORY, true), + (SUBSPACE_INBUXA, true), ] { + if subspace == SUBSPACE_BLOBS && !with_blobs { + continue; + } let from_key = AnyKey { subspace, key: vec![0u8],