From 337c3626b8fecac7edc23327d44bf1ee58df103b Mon Sep 17 00:00:00 2001 From: Borja Castellano Date: Tue, 15 Sep 2026 13:27:01 +0000 Subject: [PATCH 1/6] perf(dash-spv): size storage segments per item type Every segment cache held 50 000 items, whatever an item costs. Near the tip a block segment holds tens of MB of decoded blocks, while a header segment spanning the same heights is 5.6 MB and a filter header one 1.6 MB. `Persistable::ITEMS_PER_SEGMENT` lets each type choose: headers 10 000 (~1.1 MB per segment), filter headers 50 000 (~1.6 MB), filters 2 000 (~2 MB near the tip) and blocks 1 000. This changes the on-disk layout of the header, filter and block segments. Storage written with 50 000-item segments is misread by this layout and has to be deleted until a migration or a versioned folder lands. Mainnet restore, mainnet.100mbi.100ms, #1015/#1016/#1014 applied, two resident segments, jemalloc heap profiling, wallet identical in every run (14114383 sat, 7112 records, 13389 addresses): segment items time peak RSS segment caches 50 000 for every type 7.0 min 818 MiB 332 MiB 5 000 for every type 7.3-8.3 min 554-687 MiB 23-106 MiB 1 000 for every type 9.2-12.2 min 453-577 MiB 2-66 MiB per type (this commit) 7.5 min 538 MiB 38 MiB With 1 000 items everywhere, the header and filter header phases paid an fsync per evicted segment (105-141 s and 254-332 s instead of ~60 s and ~150 s). Blocks at 500 items saved ~18 MiB of cache but reloaded 50 % more block segments. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_017DruChNTWXwJoWPartZwCf --- dash-spv/src/storage/segments.rs | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/dash-spv/src/storage/segments.rs b/dash-spv/src/storage/segments.rs index 2717857b1..f0c98dcc7 100644 --- a/dash-spv/src/storage/segments.rs +++ b/dash-spv/src/storage/segments.rs @@ -27,6 +27,7 @@ use crate::{ pub trait Persistable: Sized + Encodable + Decodable + PartialEq + Clone { const SEGMENT_PREFIX: &'static str = "segment"; const DATA_FILE_EXTENSION: &'static str = "dat"; + const ITEMS_PER_SEGMENT: u32; fn segment_file_name(segment_id: u32) -> String { format!("{}_{:04}.{}", Self::SEGMENT_PREFIX, segment_id, Self::DATA_FILE_EXTENSION) @@ -36,12 +37,16 @@ pub trait Persistable: Sized + Encodable + Decodable + PartialEq + Clone { } impl Persistable for Vec { + const ITEMS_PER_SEGMENT: u32 = 2_000; + fn sentinel() -> Self { vec![] } } impl Persistable for HashedBlockHeader { + const ITEMS_PER_SEGMENT: u32 = 10_000; + fn sentinel() -> Self { let header = BlockHeader { version: Version::from_consensus(i32::MAX), // Invalid version @@ -57,12 +62,16 @@ impl Persistable for HashedBlockHeader { } impl Persistable for FilterHeader { + const ITEMS_PER_SEGMENT: u32 = 50_000; + fn sentinel() -> Self { FilterHeader::from_byte_array([0u8; 32]) } } impl Persistable for HashedBlock { + const ITEMS_PER_SEGMENT: u32 = 1_000; + fn sentinel() -> Self { let block = Block { header: *HashedBlockHeader::sentinel().header(), @@ -562,7 +571,7 @@ pub struct Segment { } impl Segment { - const ITEMS_PER_SEGMENT: u32 = 50_000; + const ITEMS_PER_SEGMENT: u32 = I::ITEMS_PER_SEGMENT; fn new(segment_id: u32, mut items: Vec, state: SegmentState) -> Self { debug_assert!(items.len() <= Self::ITEMS_PER_SEGMENT as usize); From 25ffe0828730a8227bcc3abb19dbab6dab264d88 Mon Sep 17 00:00:00 2001 From: Borja Castellano Date: Tue, 15 Sep 2026 13:41:07 +0000 Subject: [PATCH 2/6] refactor(dash-spv): keep Persistable inside the storage module Nothing outside `storage` implements or names the trait, and the `segments` module it lives in is private to `storage` already. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_017DruChNTWXwJoWPartZwCf --- dash-spv/src/storage/segments.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dash-spv/src/storage/segments.rs b/dash-spv/src/storage/segments.rs index f0c98dcc7..4788bb0c7 100644 --- a/dash-spv/src/storage/segments.rs +++ b/dash-spv/src/storage/segments.rs @@ -24,7 +24,7 @@ use crate::{ StorageError, }; -pub trait Persistable: Sized + Encodable + Decodable + PartialEq + Clone { +pub(super) trait Persistable: Sized + Encodable + Decodable + PartialEq + Clone { const SEGMENT_PREFIX: &'static str = "segment"; const DATA_FILE_EXTENSION: &'static str = "dat"; const ITEMS_PER_SEGMENT: u32; From 8bf329f8692c721a46e86c9f7bd4cbd79ffbbcab Mon Sep 17 00:00:00 2001 From: Borja Castellano Date: Wed, 30 Sep 2026 21:51:40 +0000 Subject: [PATCH 3/6] feat(dash-spv): migrate segment storage to per-type segment sizes Storage version 2. Segment files move from 50_000 items each to the per-type sizes (block headers 10_000, filter headers 50_000, filters 2_000, blocks 1_000), and segment ids in file names are padded to 6 digits instead of 4, since the smaller segments push ids past 9999. The migration keeps its own frozen copy of the legacy and new layouts (names, items per segment, item encodings and sentinels) and copies raw item bytes, so it does not depend on the storage code as it evolves. Each folder is rebuilt in /tmp/, swapped in through tmp/.old, and tmp/ is removed before the next folder, so the extra disk space at any time is one folder. An interrupted run either finishes the swap or rebuilds the folder from the untouched legacy files. New segments holding only sentinels are not written: an all-sentinel segment with the highest id would make load_or_new report no tip. Co-Authored-By: Claude Opus 5.5 --- dash-spv/src/storage/filters.rs | 2 +- dash-spv/src/storage/migrator/mod.rs | 6 +- dash-spv/src/storage/migrator/v2.rs | 250 +++++++++++++++++++++++++++ dash-spv/src/storage/segments.rs | 20 +-- 4 files changed, 265 insertions(+), 13 deletions(-) create mode 100644 dash-spv/src/storage/migrator/v2.rs diff --git a/dash-spv/src/storage/filters.rs b/dash-spv/src/storage/filters.rs index 3578a26ea..b691e1684 100644 --- a/dash-spv/src/storage/filters.rs +++ b/dash-spv/src/storage/filters.rs @@ -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. diff --git a/dash-spv/src/storage/migrator/mod.rs b/dash-spv/src/storage/migrator/mod.rs index 0fab0f69a..fc10ecd67 100644 --- a/dash-spv/src/storage/migrator/mod.rs +++ b/dash-spv/src/storage/migrator/mod.rs @@ -1,4 +1,5 @@ mod v1; +mod v2; use std::io; use std::path::{Path, PathBuf}; @@ -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 { @@ -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?; diff --git a/dash-spv/src/storage/migrator/v2.rs b/dash-spv/src/storage/migrator/v2.rs new file mode 100644 index 000000000..e7e7aaddb --- /dev/null +++ b/dash-spv/src/storage/migrator/v2.rs @@ -0,0 +1,250 @@ +use std::fs::{self, File}; +use std::io::{self, BufReader, BufWriter, Read, Write}; +use std::path::Path; + +use async_trait::async_trait; +use dashcore::consensus::{encode, Decodable}; +use dashcore::Block; + +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 BLOCK_HEADER_LEN: usize = 80; +const HASH_LEN: usize = 32; + +const SENTINEL_BLOCK_HEADER: [u8; BLOCK_HEADER_LEN] = { + let mut header = [0xFF; BLOCK_HEADER_LEN]; + header[3] = 0x7F; + header +}; + +#[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> { + for (folder, kind, items_per_segment) in FOLDERS { + migrate_folder(storage_path, folder, kind, items_per_segment)?; + } + + let tmp = storage_path.join(TMP_DIR); + if tmp.exists() { + fs::remove_dir_all(tmp)?; + } + + Ok(()) + } +} + +fn migrate_folder( + storage_path: &Path, + name: &str, + kind: ItemKind, + items_per_segment: u32, +) -> Result<(), MigratorError> { + let folder = storage_path.join(name); + let tmp = storage_path.join(TMP_DIR); + let staged = tmp.join(name); + let replaced = tmp.join(format!("{name}.old")); + + if replaced.exists() { + if !folder.exists() { + fs::rename(&staged, &folder)?; + } + fs::remove_dir_all(&tmp)?; + return Ok(()); + } + + let legacy_ids = legacy_segment_ids(&folder)?; + if legacy_ids.is_empty() { + return Ok(()); + } + + if tmp.exists() { + fs::remove_dir_all(&tmp)?; + } + fs::create_dir_all(&staged)?; + + for legacy_id in legacy_ids { + split_segment(&folder, &staged, legacy_id, kind, items_per_segment)?; + } + + fs::rename(&folder, &replaced)?; + fs::rename(&staged, &folder)?; + fs::remove_dir_all(&tmp)?; + + Ok(()) +} + +fn legacy_segment_ids(folder: &Path) -> Result, MigratorError> { + let entries = match fs::read_dir(folder) { + Ok(entries) => entries, + Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()), + Err(e) => return Err(e.into()), + }; + + let prefix = format!("{SEGMENT_PREFIX}_"); + let suffix = format!(".{SEGMENT_EXTENSION}"); + + let mut ids = Vec::new(); + for entry in entries { + 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)) + else { + continue; + }; + if digits.len() == LEGACY_SEGMENT_ID_DIGITS && digits.bytes().all(|b| b.is_ascii_digit()) { + ids.push(digits.parse().expect("four ascii digits fit in a u32")); + } + } + ids.sort_unstable(); + + Ok(ids) +} + +fn split_segment( + folder: &Path, + staged: &Path, + legacy_id: u32, + kind: ItemKind, + items_per_segment: u32, +) -> Result<(), MigratorError> { + let legacy_path = folder.join(format!( + "{SEGMENT_PREFIX}_{legacy_id:0width$}.{SEGMENT_EXTENSION}", + width = LEGACY_SEGMENT_ID_DIGITS + )); + 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.iter().any(|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 segment_file_name(id: u32) -> String { + format!("{SEGMENT_PREFIX}_{id:0width$}.{SEGMENT_EXTENSION}", width = SEGMENT_ID_DIGITS) +} + +fn write_segment(path: &Path, items: &[Vec]) -> 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(reader: &mut R, kind: ItemKind) -> Result>, 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::::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(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[..BLOCK_HEADER_LEN] == SENTINEL_BLOCK_HEADER, + ItemKind::FilterHeader => item.iter().all(|b| *b == 0), + ItemKind::Filter => item == [0], + ItemKind::Block => { + item[HASH_LEN..BLOCK_HEADER_LEN + HASH_LEN] == SENTINEL_BLOCK_HEADER + && item[BLOCK_HEADER_LEN + HASH_LEN..] == [0] + } + } +} + +struct Tee<'a, R> { + inner: &'a mut R, + bytes: Vec, +} + +impl Read for Tee<'_, R> { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + let n = self.inner.read(buf)?; + self.bytes.extend_from_slice(&buf[..n]); + Ok(n) + } +} diff --git a/dash-spv/src/storage/segments.rs b/dash-spv/src/storage/segments.rs index 4788bb0c7..e3979fb90 100644 --- a/dash-spv/src/storage/segments.rs +++ b/dash-spv/src/storage/segments.rs @@ -30,7 +30,7 @@ pub(super) trait Persistable: Sized + Encodable + Decodable + PartialEq + Clone const ITEMS_PER_SEGMENT: u32; fn segment_file_name(segment_id: u32) -> String { - format!("{}_{:04}.{}", Self::SEGMENT_PREFIX, segment_id, Self::DATA_FILE_EXTENSION) + format!("{}_{:06}.{}", Self::SEGMENT_PREFIX, segment_id, Self::DATA_FILE_EXTENSION) } fn sentinel() -> Self; @@ -142,13 +142,13 @@ impl SegmentCache { } /// Parse the segment id out of a segment file name of the form - /// `{SEGMENT_PREFIX}_{id:04}.{DATA_FILE_EXTENSION}` (see + /// `{SEGMENT_PREFIX}_{id:06}.{DATA_FILE_EXTENSION}` (see /// [`Persistable::segment_file_name`]). /// /// The entire remaining component between the `{prefix}_` and /// `.{extension}` fixtures must parse as a `u32`, so trailing junk - /// (`segment_0000junk.dat`) is rejected and ids longer than the - /// zero-padding width (`segment_100000.dat`) are accepted, not truncated. + /// (`segment_000000junk.dat`) is rejected and ids longer than the + /// zero-padding width (`segment_1000000.dat`) are accepted, not truncated. fn parse_segment_id(file_name: &str) -> Option { let separator = format!("{}_", I::SEGMENT_PREFIX); let suffix = format!(".{}", I::DATA_FILE_EXTENSION); @@ -1273,19 +1273,19 @@ mod tests { fn test_parse_segment_id() { type Cache = SegmentCache; - // Round-trips the writer's `{prefix}_{id:04}.{ext}` format for a range + // Round-trips the writer's `{prefix}_{id:06}.{ext}` format for a range // of ids, including ones wider than the zero-padding. - for id in [0u32, 1, 42, 9999, 10_000, 100_000, u32::MAX] { + for id in [0u32, 1, 42, 999_999, 1_000_000, u32::MAX] { assert_eq!(Cache::parse_segment_id(&FilterHeader::segment_file_name(id)), Some(id)); } // Zero-padded low ids parse to their numeric value, not truncated. - assert_eq!(Cache::parse_segment_id("segment_0000.dat"), Some(0)); - assert_eq!(Cache::parse_segment_id("segment_0005.dat"), Some(5)); + assert_eq!(Cache::parse_segment_id("segment_000000.dat"), Some(0)); + assert_eq!(Cache::parse_segment_id("segment_000005.dat"), Some(5)); // Ids wider than the padding are accepted whole, never truncated to - // the first four digits. - assert_eq!(Cache::parse_segment_id("segment_100000.dat"), Some(100_000)); + // the first six digits. + assert_eq!(Cache::parse_segment_id("segment_1000000.dat"), Some(1_000_000)); // Trailing junk between the id and the extension is rejected rather // than silently parsed as a prefix of the component. From 96369195b698aadc10662bd615cc4610dd45d07e Mon Sep 17 00:00:00 2001 From: Borja Castellano Date: Wed, 30 Sep 2026 23:14:33 +0000 Subject: [PATCH 4/6] refactor(dash-spv): simplify segment migration recovery and sentinel check Recovery no longer tracks a swap state: the migration deletes /tmp on start, migrates every folder that still holds 4-digit segment files and skips folders that only hold 6-digit ones. Staged segments are moved into the folder file by file and the legacy files are deleted last, so every 6-digit file in a folder is complete and rebuilding from the remaining legacy files reproduces the same files. The sentinel check looks at one field per item type (version i32::MAX for headers and blocks, an all-zero filter header, an empty filter). It is still needed: on a dev V1 mainnet storage the last new filters segment is all sentinels, and writing it made the client load a filter tip of 0 and download every filter again from height 200000. Checked on that storage (synced by dev to 2547714): the client resumes at the tip, every new file equals its slice of the legacy file and every skipped slice is only sentinels, and a SIGKILL mid-migration followed by a restart yields identical files. Blocks skip 771 of 2350 segments (523 -> 443 MiB). Co-Authored-By: Claude Opus 5.5 --- dash-spv/src/storage/migrator/v2.rs | 77 ++++++++++++++--------------- 1 file changed, 38 insertions(+), 39 deletions(-) diff --git a/dash-spv/src/storage/migrator/v2.rs b/dash-spv/src/storage/migrator/v2.rs index e7e7aaddb..7c0212dc9 100644 --- a/dash-spv/src/storage/migrator/v2.rs +++ b/dash-spv/src/storage/migrator/v2.rs @@ -24,12 +24,7 @@ const TMP_DIR: &str = "tmp"; const BLOCK_HEADER_LEN: usize = 80; const HASH_LEN: usize = 32; - -const SENTINEL_BLOCK_HEADER: [u8; BLOCK_HEADER_LEN] = { - let mut header = [0xFF; BLOCK_HEADER_LEN]; - header[3] = 0x7F; - header -}; +const SENTINEL_VERSION: [u8; 4] = i32::MAX.to_le_bytes(); #[derive(Clone, Copy)] enum ItemKind { @@ -59,19 +54,26 @@ 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)?; + for (folder, kind, items_per_segment) in FOLDERS { migrate_folder(storage_path, folder, kind, items_per_segment)?; } - let tmp = storage_path.join(TMP_DIR); - if tmp.exists() { - fs::remove_dir_all(tmp)?; - } + remove_dir_if_exists(&tmp)?; Ok(()) } } +fn remove_dir_if_exists(path: &Path) -> io::Result<()> { + match fs::remove_dir_all(path) { + Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(()), + result => result, + } +} + fn migrate_folder( storage_path: &Path, name: &str, @@ -79,35 +81,30 @@ fn migrate_folder( items_per_segment: u32, ) -> Result<(), MigratorError> { let folder = storage_path.join(name); - let tmp = storage_path.join(TMP_DIR); - let staged = tmp.join(name); - let replaced = tmp.join(format!("{name}.old")); - - if replaced.exists() { - if !folder.exists() { - fs::rename(&staged, &folder)?; - } - fs::remove_dir_all(&tmp)?; - return Ok(()); - } + let staged = storage_path.join(TMP_DIR).join(name); let legacy_ids = legacy_segment_ids(&folder)?; if legacy_ids.is_empty() { return Ok(()); } - if tmp.exists() { - fs::remove_dir_all(&tmp)?; - } + remove_dir_if_exists(&staged)?; fs::create_dir_all(&staged)?; - for legacy_id in legacy_ids { + for &legacy_id in &legacy_ids { split_segment(&folder, &staged, legacy_id, kind, items_per_segment)?; } - fs::rename(&folder, &replaced)?; - fs::rename(&staged, &folder)?; - fs::remove_dir_all(&tmp)?; + for entry in fs::read_dir(&staged)? { + let entry = entry?; + fs::rename(entry.path(), folder.join(entry.file_name()))?; + } + + for legacy_id in legacy_ids { + fs::remove_file(folder.join(legacy_segment_file_name(legacy_id)))?; + } + + remove_dir_if_exists(&staged)?; Ok(()) } @@ -148,10 +145,7 @@ fn split_segment( kind: ItemKind, items_per_segment: u32, ) -> Result<(), MigratorError> { - let legacy_path = folder.join(format!( - "{SEGMENT_PREFIX}_{legacy_id:0width$}.{SEGMENT_EXTENSION}", - width = LEGACY_SEGMENT_ID_DIGITS - )); + 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; @@ -165,7 +159,11 @@ fn split_segment( items.push(item); } - if items.iter().any(|item| !is_sentinel(item, kind)) { + 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)?; } @@ -184,6 +182,10 @@ fn split_segment( 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) } @@ -226,13 +228,10 @@ fn read_fixed(reader: &mut R, len: usize) -> Result<(), encode::Error> fn is_sentinel(item: &[u8], kind: ItemKind) -> bool { match kind { - ItemKind::BlockHeader => item[..BLOCK_HEADER_LEN] == SENTINEL_BLOCK_HEADER, - ItemKind::FilterHeader => item.iter().all(|b| *b == 0), + ItemKind::BlockHeader => item[..4] == SENTINEL_VERSION, + ItemKind::FilterHeader => item == [0; HASH_LEN], ItemKind::Filter => item == [0], - ItemKind::Block => { - item[HASH_LEN..BLOCK_HEADER_LEN + HASH_LEN] == SENTINEL_BLOCK_HEADER - && item[BLOCK_HEADER_LEN + HASH_LEN..] == [0] - } + ItemKind::Block => item[HASH_LEN..HASH_LEN + 4] == SENTINEL_VERSION, } } From 44cc960568bde036c6d828a171abe0a52ab89e84 Mon Sep 17 00:00:00 2001 From: Borja Castellano Date: Thu, 1 Oct 2026 00:18:23 +0000 Subject: [PATCH 5/6] perf(dash-spv): migrate segment folders in parallel and treat mixed folders as corruption Each folder is migrated by its own tokio task, and inside it WORKERS_PER_FOLDER spawn_blocking workers split the legacy segments (worker n takes legacy ids n, n + 4, ...). The split stays the streaming, blocking code: reading whole legacy files through tokio::fs would hold up to ~600 MB at once (filters files reach 61 MB, blocks 88 MB). Recovery is now: delete /tmp on start, skip folders that only hold 6-digit segment files, and return Corruption for a folder holding both 4-digit and 6-digit files, which wipes the storage. On a dev V1 mainnet storage (cold page cache, this server) the migration took 41-46 s sequentially and 4-17 s now. Almost all of it is the per-file fsync, which stays: without it the whole migration takes ~2 s because writes only reach the page cache, and the legacy files are deleted right after. Deferring the fsyncs until all files are written was slower. Bytes were verified against the legacy files, a SIGKILL mid-migration followed by a restart migrates correctly, and a mixed folder wipes the storage and syncs again from the checkpoint. Co-Authored-By: Claude Opus 5.5 --- dash-spv/src/storage/migrator/v2.rs | 98 +++++++++++++++++++---------- 1 file changed, 64 insertions(+), 34 deletions(-) diff --git a/dash-spv/src/storage/migrator/v2.rs b/dash-spv/src/storage/migrator/v2.rs index 7c0212dc9..6c9a1faae 100644 --- a/dash-spv/src/storage/migrator/v2.rs +++ b/dash-spv/src/storage/migrator/v2.rs @@ -1,10 +1,11 @@ -use std::fs::{self, File}; +use std::fs::File; use std::io::{self, BufReader, BufWriter, Read, Write}; -use std::path::Path; +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}; @@ -21,6 +22,7 @@ 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; @@ -55,87 +57,115 @@ pub struct V2Migrator; 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)?; - - for (folder, kind, items_per_segment) in FOLDERS { - migrate_folder(storage_path, folder, kind, items_per_segment)?; + 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)?; + remove_dir_if_exists(&tmp).await?; Ok(()) } } -fn remove_dir_if_exists(path: &Path) -> io::Result<()> { - match fs::remove_dir_all(path) { +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, } } -fn migrate_folder( - storage_path: &Path, - name: &str, +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 = legacy_segment_ids(&folder)?; + 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" + ))); + } - remove_dir_if_exists(&staged)?; - fs::create_dir_all(&staged)?; - - for &legacy_id in &legacy_ids { - split_segment(&folder, &staged, legacy_id, kind, items_per_segment)?; + tokio::fs::create_dir_all(&staged).await?; + + let mut workers = JoinSet::new(); + for worker in 0..WORKERS_PER_FOLDER { + let ids: Vec = + 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")?; } - for entry in fs::read_dir(&staged)? { - let entry = entry?; - fs::rename(entry.path(), folder.join(entry.file_name()))?; + 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 { - fs::remove_file(folder.join(legacy_segment_file_name(legacy_id)))?; + tokio::fs::remove_file(folder.join(legacy_segment_file_name(legacy_id))).await?; } - remove_dir_if_exists(&staged)?; - Ok(()) } -fn legacy_segment_ids(folder: &Path) -> Result, MigratorError> { - let entries = match fs::read_dir(folder) { +async fn segment_ids(folder: &Path) -> Result<(Vec, 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()), + 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 ids = Vec::new(); - for entry in entries { - let name = entry?.file_name(); + 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; }; - if digits.len() == LEGACY_SEGMENT_ID_DIGITS && digits.bytes().all(|b| b.is_ascii_digit()) { - ids.push(digits.parse().expect("four ascii digits fit in a u32")); + 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, + _ => {} } } - ids.sort_unstable(); + legacy_ids.sort_unstable(); - Ok(ids) + Ok((legacy_ids, has_current)) } fn split_segment( From a6a2818c7660a11f3b58f1af22c7acc492358812 Mon Sep 17 00:00:00 2001 From: Borja Castellano Date: Thu, 1 Oct 2026 10:41:11 +0000 Subject: [PATCH 6/6] fix(dash-spv): treat a legacy segment ending mid-item as corruption `read_item` took every `UnexpectedEof` as the clean end of a legacy segment, including one hit after part of an item was read. The truncated segment was then migrated and its legacy file deleted, losing the tail silently. Only an EOF before the first byte of an item ends the segment now; anything else is `MigratorError::Corruption`, which clears the storage for a resync. Co-Authored-By: Claude Opus 5.5 --- dash-spv/src/storage/migrator/v2.rs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/dash-spv/src/storage/migrator/v2.rs b/dash-spv/src/storage/migrator/v2.rs index 6c9a1faae..f5a13f712 100644 --- a/dash-spv/src/storage/migrator/v2.rs +++ b/dash-spv/src/storage/migrator/v2.rs @@ -245,7 +245,11 @@ fn read_item(reader: &mut R, kind: ItemKind) -> Result>, match result { Ok(()) => Ok(Some(tee.bytes)), - Err(encode::Error::Io(e)) if e.kind() == io::ErrorKind::UnexpectedEof => Ok(None), + Err(encode::Error::Io(e)) + if e.kind() == io::ErrorKind::UnexpectedEof && tee.bytes.is_empty() => + { + Ok(None) + } Err(e) => Err(MigratorError::Corruption(format!("Failed to decode legacy item: {e}"))), } }