Skip to content
Merged
Show file tree
Hide file tree
Changes from 6 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion dash-spv/src/storage/filters.rs
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,7 @@ mod tests {

storage.persist(tmp_dir.path()).await.unwrap();
let segment_file =
tmp_dir.path().join(PersistentFilterStorage::FOLDER_NAME).join("segment_0000.dat");
tmp_dir.path().join(PersistentFilterStorage::FOLDER_NAME).join("segment_000000.dat");
assert!(segment_file.exists());

// The start height survives a reload.
Expand Down
6 changes: 4 additions & 2 deletions dash-spv/src/storage/migrator/mod.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
mod v1;
mod v2;

use std::io;
use std::path::{Path, PathBuf};
Expand All @@ -9,8 +10,9 @@ use serde::{Deserialize, Serialize};
use crate::error::{StorageError, StorageResult};
use crate::storage::io::atomic_write;
use v1::V1Migrator;
use v2::V2Migrator;

pub const CURRENT_VERSION: u32 = 1;
pub const CURRENT_VERSION: u32 = 2;

#[async_trait]
trait Migrator {
Expand Down Expand Up @@ -116,7 +118,7 @@ async fn migration_loop(storage_path: &Path) -> Result<(), MigratorError> {
while version < CURRENT_VERSION {
tracing::info!("Migrating storage from version {} to {}", version, version + 1);

const MIGRATORS: &[&dyn Migrator; CURRENT_VERSION as usize] = &[&V1Migrator];
const MIGRATORS: &[&dyn Migrator; CURRENT_VERSION as usize] = &[&V1Migrator, &V2Migrator];

MIGRATORS[version as usize].apply_migration(storage_path).await?;

Expand Down
279 changes: 279 additions & 0 deletions dash-spv/src/storage/migrator/v2.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,279 @@
use std::fs::File;
use std::io::{self, BufReader, BufWriter, Read, Write};
use std::path::{Path, PathBuf};

use async_trait::async_trait;
use dashcore::consensus::{encode, Decodable};
use dashcore::Block;
use tokio::task::JoinSet;

use super::{Migrator, MigratorError};

const SEGMENT_PREFIX: &str = "segment";
const SEGMENT_EXTENSION: &str = "dat";

const LEGACY_SEGMENT_ID_DIGITS: usize = 4;
const LEGACY_ITEMS_PER_SEGMENT: u32 = 50_000;

const SEGMENT_ID_DIGITS: usize = 6;
const BLOCK_HEADERS_PER_SEGMENT: u32 = 10_000;
const FILTER_HEADERS_PER_SEGMENT: u32 = 50_000;
const FILTERS_PER_SEGMENT: u32 = 2_000;
const BLOCKS_PER_SEGMENT: u32 = 1_000;

const TMP_DIR: &str = "tmp";
const WORKERS_PER_FOLDER: usize = 4;

const BLOCK_HEADER_LEN: usize = 80;
const HASH_LEN: usize = 32;
const SENTINEL_VERSION: [u8; 4] = i32::MAX.to_le_bytes();

#[derive(Clone, Copy)]
enum ItemKind {
BlockHeader,
FilterHeader,
Filter,
Block,
}

const FOLDERS: [(&str, ItemKind, u32); 4] = [
("block_headers", ItemKind::BlockHeader, BLOCK_HEADERS_PER_SEGMENT),
("filter_headers", ItemKind::FilterHeader, FILTER_HEADERS_PER_SEGMENT),
("filters", ItemKind::Filter, FILTERS_PER_SEGMENT),
("blocks", ItemKind::Block, BLOCKS_PER_SEGMENT),
];

const _: () = {
let mut i = 0;
while i < FOLDERS.len() {
assert!(LEGACY_ITEMS_PER_SEGMENT.is_multiple_of(FOLDERS[i].2));
i += 1;
}
};

pub struct V2Migrator;

#[async_trait]
impl Migrator for V2Migrator {
async fn apply_migration(&self, storage_path: &Path) -> Result<(), MigratorError> {
let tmp = storage_path.join(TMP_DIR);
remove_dir_if_exists(&tmp).await?;

let mut folders = JoinSet::new();
for (name, kind, items_per_segment) in FOLDERS {
folders.spawn(migrate_folder(
storage_path.to_path_buf(),
name,
kind,
items_per_segment,
));
}
while let Some(result) = folders.join_next().await {
result.expect("folder migration panicked")?;
}

remove_dir_if_exists(&tmp).await?;

Ok(())
}
}

async fn remove_dir_if_exists(path: &Path) -> io::Result<()> {
match tokio::fs::remove_dir_all(path).await {
Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(()),
result => result,
}
}

async fn migrate_folder(
storage_path: PathBuf,
name: &'static str,
kind: ItemKind,
items_per_segment: u32,
) -> Result<(), MigratorError> {
let folder = storage_path.join(name);
let staged = storage_path.join(TMP_DIR).join(name);

let (legacy_ids, has_current) = segment_ids(&folder).await?;
if legacy_ids.is_empty() {
return Ok(());
}
if has_current {
return Err(MigratorError::Corruption(format!(
"{folder:?} holds both legacy and current segment files"
)));
}

tokio::fs::create_dir_all(&staged).await?;

let mut workers = JoinSet::new();
for worker in 0..WORKERS_PER_FOLDER {
let ids: Vec<u32> =
legacy_ids.iter().skip(worker).step_by(WORKERS_PER_FOLDER).copied().collect();
let (folder, staged) = (folder.clone(), staged.clone());
workers.spawn_blocking(move || {
ids.into_iter().try_for_each(|legacy_id| {
split_segment(&folder, &staged, legacy_id, kind, items_per_segment)
})
});
}
while let Some(result) = workers.join_next().await {
result.expect("migration worker panicked")?;
}

let mut staged_entries = tokio::fs::read_dir(&staged).await?;
while let Some(entry) = staged_entries.next_entry().await? {
tokio::fs::rename(entry.path(), folder.join(entry.file_name())).await?;
}

for legacy_id in legacy_ids {
tokio::fs::remove_file(folder.join(legacy_segment_file_name(legacy_id))).await?;
}

Ok(())
}

async fn segment_ids(folder: &Path) -> Result<(Vec<u32>, bool), MigratorError> {
let mut entries = match tokio::fs::read_dir(folder).await {
Ok(entries) => entries,
Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok((Vec::new(), false)),
Err(e) => return Err(e.into()),
};

let prefix = format!("{SEGMENT_PREFIX}_");
let suffix = format!(".{SEGMENT_EXTENSION}");

let mut legacy_ids = Vec::new();
let mut has_current = false;
while let Some(entry) = entries.next_entry().await? {
let name = entry.file_name();
let Some(digits) = name
.to_str()
.and_then(|name| name.strip_prefix(&prefix))
.and_then(|rest| rest.strip_suffix(&suffix))
.filter(|digits| digits.bytes().all(|b| b.is_ascii_digit()))
else {
continue;
};
match digits.len() {
LEGACY_SEGMENT_ID_DIGITS => {
legacy_ids.push(digits.parse().expect("four ascii digits fit in a u32"))
}
SEGMENT_ID_DIGITS => has_current = true,
_ => {}
}
}
legacy_ids.sort_unstable();

Ok((legacy_ids, has_current))
}

fn split_segment(
folder: &Path,
staged: &Path,
legacy_id: u32,
kind: ItemKind,
items_per_segment: u32,
) -> Result<(), MigratorError> {
let legacy_path = folder.join(legacy_segment_file_name(legacy_id));
let mut reader = BufReader::new(File::open(&legacy_path)?);

let segments_per_legacy = LEGACY_ITEMS_PER_SEGMENT / items_per_segment;

for chunk in 0..segments_per_legacy {
let mut items = Vec::with_capacity(items_per_segment as usize);
while items.len() < items_per_segment as usize {
let Some(item) = read_item(&mut reader, kind)? else {
break;
};
items.push(item);
}

if items.is_empty() {
return Ok(());
}

if !items.iter().all(|item| is_sentinel(item, kind)) {
let id = legacy_id * segments_per_legacy + chunk;
write_segment(&staged.join(segment_file_name(id)), &items)?;
}

if items.len() < items_per_segment as usize {
return Ok(());
}
}

if read_item(&mut reader, kind)?.is_some() {
return Err(MigratorError::Corruption(format!(
"{legacy_path:?} holds more than {LEGACY_ITEMS_PER_SEGMENT} items"
)));
}

Ok(())
}

fn legacy_segment_file_name(id: u32) -> String {
format!("{SEGMENT_PREFIX}_{id:0width$}.{SEGMENT_EXTENSION}", width = LEGACY_SEGMENT_ID_DIGITS)
}

fn segment_file_name(id: u32) -> String {
format!("{SEGMENT_PREFIX}_{id:0width$}.{SEGMENT_EXTENSION}", width = SEGMENT_ID_DIGITS)
}

fn write_segment(path: &Path, items: &[Vec<u8>]) -> Result<(), MigratorError> {
let mut writer = BufWriter::new(File::create(path)?);
for item in items {
writer.write_all(item)?;
}
writer.into_inner().map_err(|e| e.into_error())?.sync_all()?;
Ok(())
}

fn read_item<R: Read>(reader: &mut R, kind: ItemKind) -> Result<Option<Vec<u8>>, MigratorError> {
let mut tee = Tee {
inner: reader,
bytes: Vec::new(),
};

let result = match kind {
ItemKind::BlockHeader => read_fixed(&mut tee, BLOCK_HEADER_LEN + HASH_LEN),
ItemKind::FilterHeader => read_fixed(&mut tee, HASH_LEN),
ItemKind::Filter => Vec::<u8>::consensus_decode(&mut tee).map(drop),
ItemKind::Block => read_fixed(&mut tee, HASH_LEN)
.and_then(|()| Block::consensus_decode(&mut tee).map(drop)),
};

match result {
Ok(()) => Ok(Some(tee.bytes)),
Err(encode::Error::Io(e)) if e.kind() == io::ErrorKind::UnexpectedEof => Ok(None),
Err(e) => Err(MigratorError::Corruption(format!("Failed to decode legacy item: {e}"))),
}
}

fn read_fixed<R: Read>(reader: &mut R, len: usize) -> Result<(), encode::Error> {
let mut buf = vec![0; len];
reader.read_exact(&mut buf)?;
Ok(())
}

fn is_sentinel(item: &[u8], kind: ItemKind) -> bool {
match kind {
ItemKind::BlockHeader => item[..4] == SENTINEL_VERSION,
ItemKind::FilterHeader => item == [0; HASH_LEN],
ItemKind::Filter => item == [0],
ItemKind::Block => item[HASH_LEN..HASH_LEN + 4] == SENTINEL_VERSION,
}
}

struct Tee<'a, R> {
inner: &'a mut R,
bytes: Vec<u8>,
}

impl<R: Read> Read for Tee<'_, R> {
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
let n = self.inner.read(buf)?;
self.bytes.extend_from_slice(&buf[..n]);
Ok(n)
}
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Loading
Loading