Export/import: keep archived items, spam samples and the spam model
--export skipped three things, so a move from one database to another
(RocksDB to PostgreSQL, say) lost them without a word:
- archived items (subspace j), the records behind undelete;
- spam training samples (subspace w);
- the trained spam classifier and its trainer state, blobs stored under
fixed names that no blob link points at, so the walk over links never
reached them.
j and w now travel with the registry family, where their indexes and id
counters already were, so EXPORT_TYPES=registry keeps them consistent.
The two named blobs travel with the blob family. The file format is
unchanged and import reads any subspace it is given, so an export made
by an older binary still imports.
The full-text index (subspace z) stays out, on purpose. It belongs to one
search backend: PostgreSQL and MySQL index into their own tables and have
no z table at all, and external engines keep the index themselves. So
--import now returns the subspaces it wrote, and boot queues the
reindexAccounts and reindexTelemetry store maintenance tasks, the same
ones an administrator can queue by hand, to rebuild the index for
whichever search store the server runs with once it starts.
The round trip also turned up a loss in import itself: the SQL stores
add a negative amount with an UPDATE, which does nothing to a row that
isn't there yet, so every negative counter or quota vanished on import
into PostgreSQL, MySQL or SQLite. Import now creates the row first.
The in-memory subspaces (m, y) stay out: rate limits, locks, greylisting,
ACME challenge tokens and OAuth codes, all short-lived. Issued
certificates are registry objects and travel.
The store test now writes archived items, spam samples, directory
entries, the fork's own subspace and the named blobs, checks they come
back in place, then imports the same export into a fresh store of the
other local backend (RocksDB to SQLite, or SQLite to RocksDB), compares
it key for key and counter for counter, and checks the queued reindex.
It fails on the old export code ("Subspace j was not exported").
--help now says what an export holds.
This commit is contained in:
@@ -2,6 +2,8 @@
|
||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||
*
|
||||
* 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<u8>)],
|
||||
) {
|
||||
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::<Task>(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<u8>)]) {
|
||||
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<KeyValue>,
|
||||
@@ -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],
|
||||
|
||||
Reference in New Issue
Block a user