From 1fb5f2e4cf457a85fe13c15e781831a3caee3775 Mon Sep 17 00:00:00 2001 From: Brijesh-Thakkar Date: Wed, 30 Sep 2026 10:40:47 +0530 Subject: [PATCH 1/2] fix(file_store): don't overwrite entries appended via another handle Signed-off-by: Brijesh-Thakkar --- crates/file_store/src/store.rs | 126 ++++++++++++++++++- crates/file_store/tests/test_stale_handle.rs | 31 +++++ 2 files changed, 156 insertions(+), 1 deletion(-) create mode 100644 crates/file_store/tests/test_stale_handle.rs diff --git a/crates/file_store/src/store.rs b/crates/file_store/src/store.rs index 858b9d2cdf..ca4dbe471b 100644 --- a/crates/file_store/src/store.rs +++ b/crates/file_store/src/store.rs @@ -4,7 +4,7 @@ use bincode::Options; use std::{ fmt::{self, Debug}, fs::{File, OpenOptions}, - io::{self, Read, Write}, + io::{self, Read, Seek, Write}, marker::PhantomData, path::Path, }; @@ -230,6 +230,10 @@ where /// not needed because file pointer is always moved to the end of the last decodable data from /// beginning to end. /// + /// Before writing, valid changesets appended through other handles to the same file since this + /// handle last read or wrote are skipped over, so they are not overwritten. Concurrent writes + /// are still not synchronized. + /// /// If multiple garbage writes are produced on the file, the next load will only retrieve the /// first chunk of valid changesets. /// @@ -242,6 +246,26 @@ where return Ok(()); } + // Skip over valid entries that other handles appended since this handle's cursor was last + // synced, so we never overwrite them. `EntryIter` leaves the cursor right after the last + // decodable entry, or at the start of an undecodable one, so trailing garbage is handled + // as before. + let pos = self.db_file.stream_position()?; + for entry in EntryIter::::new(pos, &mut self.db_file) { + match entry { + Ok(_) => {} + // A real I/O failure means we cannot tell where the valid data ends. + Err(StoreError::Io(error)) => return Err(error), + Err(StoreError::Bincode(bincode::ErrorKind::Io(error))) + if error.kind() != io::ErrorKind::UnexpectedEof => + { + return Err(error) + } + // Undecodable or partial entry: write from here, as before. + Err(_) => break, + } + } + bincode_options() .serialize_into(&mut self.db_file, changeset) .map_err(|e| match *e { @@ -544,6 +568,106 @@ mod test { } } + /// Appending after a failed `dump` must still write at the start of the undecodable data, not + /// after it, so the new changeset stays reachable. + #[test] + fn append_after_failed_dump_writes_at_start_of_garbage() { + let temp_dir = tempfile::tempdir().unwrap(); + let file_path = temp_dir.path().join("db_file"); + let changeset1 = TestChangeSet::from(["a".to_string()]); + let changeset2 = TestChangeSet::from(["b".to_string()]); + + let mut store = + Store::::create(&TEST_MAGIC_BYTES, &file_path).expect("must create"); + store.append(&changeset1).expect("must append"); + store + .db_file + .write_all(&[255_u8; 20]) + .expect("should write"); + store.dump().expect_err("must fail on garbage"); + + store.append(&changeset2).expect("must append"); + drop(store); + + let err = Store::::load(&TEST_MAGIC_BYTES, &file_path) + .expect_err("leftover garbage must still fail to load"); + let mut expected = changeset1; + expected.extend(changeset2); + assert_eq!(err.changeset, Some(Box::new(expected))); + } + + /// A partial trailing entry written by another handle ends the catch-up scan; the append then + /// starts at that point, as it did before the scan existed. + #[test] + fn append_over_partial_trailing_entry_from_other_handle() { + let temp_dir = tempfile::tempdir().unwrap(); + let file_path = temp_dir.path().join("db_file"); + let changeset1 = TestChangeSet::from(["a".to_string()]); + let changeset2 = TestChangeSet::from(["b".to_string()]); + let torn = bincode_options() + .serialize(&TestChangeSet::from(["zzzzzz".to_string()])) + .unwrap(); + + let mut first = + Store::::create(&TEST_MAGIC_BYTES, &file_path).expect("must create"); + first.append(&changeset1).expect("must append"); + let (mut second, _) = + Store::::load(&TEST_MAGIC_BYTES, &file_path).expect("must load"); + first.db_file.write_all(&torn[..2]).expect("should write"); + + second.append(&changeset2).expect("must append"); + drop((first, second)); + + let (_, recovered) = + Store::::load(&TEST_MAGIC_BYTES, &file_path).expect("must load"); + let mut expected = changeset1; + expected.extend(changeset2); + assert_eq!(recovered, Some(expected)); + } + + /// After `append` the cursor must sit right after the new entry, both for a single handle and + /// for a handle that had to skip entries written by another one. + #[test] + fn append_leaves_cursor_after_new_entry() { + let temp_dir = tempfile::tempdir().unwrap(); + let file_path = temp_dir.path().join("db_file"); + let file_len = || fs::metadata(&file_path).unwrap().len(); + + let mut first = + Store::::create(&TEST_MAGIC_BYTES, &file_path).expect("must create"); + first.append(&TestChangeSet::from(["a".into()])).unwrap(); + assert_eq!(first.db_file.stream_position().unwrap(), file_len()); + + let (mut second, _) = + Store::::load(&TEST_MAGIC_BYTES, &file_path).expect("must load"); + first.append(&TestChangeSet::from(["b".into()])).unwrap(); + assert_eq!(first.db_file.stream_position().unwrap(), file_len()); + + second.append(&TestChangeSet::from(["c".into()])).unwrap(); + assert_eq!(second.db_file.stream_position().unwrap(), file_len()); + } + + /// An I/O error while scanning must be returned, not ignored, and nothing may be written. + #[test] + fn append_returns_io_error_from_scan() { + let temp_dir = tempfile::tempdir().unwrap(); + let file_path = temp_dir.path().join("db_file"); + let changeset = TestChangeSet::from(["a".to_string()]); + + let mut store = + Store::::create(&TEST_MAGIC_BYTES, &file_path).expect("must create"); + store.append(&changeset).expect("must append"); + let bytes_before = fs::read(&file_path).unwrap(); + + // A write-only handle cannot be read, so the scan fails with an I/O error. + store.db_file = OpenOptions::new().write(true).open(&file_path).unwrap(); + store + .append(&TestChangeSet::from(["b".to_string()])) + .expect_err("scan error must be returned"); + + assert_eq!(fs::read(&file_path).unwrap(), bytes_before); + } + #[test] fn test_load_recovers_state_after_last_write() { let temp_dir = tempfile::tempdir().unwrap(); diff --git a/crates/file_store/tests/test_stale_handle.rs b/crates/file_store/tests/test_stale_handle.rs new file mode 100644 index 0000000000..ad93c82fd0 --- /dev/null +++ b/crates/file_store/tests/test_stale_handle.rs @@ -0,0 +1,31 @@ +use bdk_file_store::Store; +use std::collections::BTreeSet; + +const MAGIC: &[u8] = b"bdk_test_magic"; +type ChangeSet = BTreeSet; + +#[test] +fn append_through_second_handle_keeps_earlier_append() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("db"); + + let mut first = Store::::create(MAGIC, &path).unwrap(); + first + .append(&ChangeSet::from(["initial".to_string()])) + .unwrap(); + + // Second handle opened while the file ends after "initial". + let (mut second, _) = Store::::load(MAGIC, &path).unwrap(); + + first + .append(&ChangeSet::from(["first".to_string()])) + .unwrap(); + second + .append(&ChangeSet::from(["other".to_string()])) + .unwrap(); + drop((first, second)); + + let (_, recovered) = Store::::load(MAGIC, &path).unwrap(); + let expected = ChangeSet::from(["initial".into(), "first".into(), "other".into()]); + assert_eq!(recovered, Some(expected)); +} From 956d277d75330fb085a1d43bf83a1ea16996b055 Mon Sep 17 00:00:00 2001 From: Brijesh-Thakkar Date: Wed, 30 Sep 2026 22:56:04 +0530 Subject: [PATCH 2/2] refactor(file_store)!: replace bincode with postcard `bincode` is deprecated. Migrate the file store's on-disk encoding to `postcard`, as part of the fix for #2308 (see also #2258). Each entry is now stored as a `postcard` `u64` varint length prefix followed by the `postcard`-encoded changeset. The length prefix frames entries, so a torn or corrupt entry can be told apart from a clean end of file, and a corrupt, huge length cannot trigger a huge allocation. Bytes left in a frame after decoding are rejected. Real I/O errors are reported as `StoreError::Io`, which `Store::append` now returns when scanning past entries written by another handle. BREAKING CHANGE: - The on-disk format changed. Files written by earlier versions cannot be read; use new magic bytes so old files fail with `StoreError::InvalidMagicBytes`. - `StoreError::Bincode(bincode::ErrorKind)` is replaced by `StoreError::Decode(postcard::Error)`. Co-Authored-By: Claude Sonnet 5.5 --- crates/file_store/Cargo.toml | 2 +- crates/file_store/src/entry_iter.rs | 104 ++++++++++++++++++------ crates/file_store/src/lib.rs | 9 +- crates/file_store/src/store.rs | 122 ++++++++++++++++++++++------ 4 files changed, 179 insertions(+), 58 deletions(-) diff --git a/crates/file_store/Cargo.toml b/crates/file_store/Cargo.toml index 8fbdc358de..131b9e9b04 100644 --- a/crates/file_store/Cargo.toml +++ b/crates/file_store/Cargo.toml @@ -16,7 +16,7 @@ workspace = true [dependencies] bdk_core = { path = "../core", version = "0.6.1", features = ["serde"]} -bincode = { version = "1" } +postcard = { version = "1.1", default-features = false, features = ["use-std"] } serde = { version = "1", features = ["derive"] } [dev-dependencies] diff --git a/crates/file_store/src/entry_iter.rs b/crates/file_store/src/entry_iter.rs index 8b284f1814..6e9015f815 100644 --- a/crates/file_store/src/entry_iter.rs +++ b/crates/file_store/src/entry_iter.rs @@ -1,18 +1,18 @@ use crate::StoreError; -use bincode::Options; use std::{ fs::File, - io::{self, BufReader, Seek}, + io::{self, BufRead, BufReader, Read, Seek}, marker::PhantomData, }; -use crate::bincode_options; - /// Iterator over entries in a file store. /// /// Reads and returns an entry each time [`next`] is called. If an error occurs while reading the /// iterator will yield a `Result::Err(_)` instead and then `None` for the next call to `next`. /// +/// Each entry is stored as a `postcard`-encoded `u64` varint length prefix followed by that many +/// bytes of `postcard`-encoded data. +/// /// [`next`]: Self::next pub struct EntryIter<'t, T> { /// Buffered reader around the file @@ -44,31 +44,81 @@ where if self.finished { return None; } - (|| { - if let Some(start) = self.start_pos.take() { - self.db_file.seek(io::SeekFrom::Start(start))?; - } + let entry = self.read_entry().transpose(); + // stop after the end of the file or the first error + if !matches!(entry, Some(Ok(_))) { + self.finished = true; + } + entry + } +} + +impl EntryIter<'_, T> +where + T: serde::de::DeserializeOwned, +{ + /// Reads the next entry, or `Ok(None)` on a clean end-of-file. + /// + /// If the entry cannot be read the file is rewound to where it started, so the position is + /// never left in the middle of an entry. + fn read_entry(&mut self) -> Result, StoreError> { + if let Some(start) = self.start_pos.take() { + self.db_file.seek(io::SeekFrom::Start(start))?; + } + let pos_before_read = self.db_file.stream_position()?; + + // no bytes left: this is the end of the file, not a torn entry + if self.db_file.fill_buf()?.is_empty() { + return Ok(None); + } + + let entry = self.read_frame(); + if entry.is_err() { + self.db_file.seek(io::SeekFrom::Start(pos_before_read))?; + } + entry.map(Some) + } + + /// Reads one length-prefixed frame and decodes it. + /// + /// A frame that ends early (end-of-file) or does not decode is [`StoreError::Decode`]. Any + /// other failure to read is reported as [`StoreError::Io`]. + fn read_frame(&mut self) -> Result { + let len = self.read_len_prefix()?; + + // `take` + `read_to_end` only allocates what the file actually holds, so a corrupt, huge + // length prefix cannot trigger a huge allocation. + let mut payload = Vec::new(); + (&mut self.db_file).take(len).read_to_end(&mut payload)?; + if payload.len() as u64 != len { + return Err(StoreError::Decode( + postcard::Error::DeserializeUnexpectedEnd, + )); + } + + // The length prefix is authoritative, so bytes left over after decoding are corruption. + match postcard::take_from_bytes(&payload) { + Ok((entry, [])) => Ok(entry), + Ok(_) => Err(StoreError::Decode(postcard::Error::SerdeDeCustom)), + Err(e) => Err(StoreError::Decode(e)), + } + } - let pos_before_read = self.db_file.stream_position()?; - match bincode_options().deserialize_from(&mut self.db_file) { - Ok(changeset) => Ok(Some(changeset)), - Err(e) => { - self.finished = true; - let pos_after_read = self.db_file.stream_position()?; - // allow unexpected EOF if 0 bytes were read - if let bincode::ErrorKind::Io(inner) = &*e { - if inner.kind() == io::ErrorKind::UnexpectedEof - && pos_after_read == pos_before_read - { - return Ok(None); - } - } - self.db_file.seek(io::SeekFrom::Start(pos_before_read))?; - Err(StoreError::Bincode(*e)) - } + /// Reads the `postcard` varint `u64` that holds the length of the frame's payload. + fn read_len_prefix(&mut self) -> Result { + // a `u64` varint is at most 10 bytes; the high bit of a byte flags that another follows + let mut buf = [0_u8; 10]; + for i in 0..buf.len() { + if self.db_file.read(&mut buf[i..=i])? == 0 { + return Err(StoreError::Decode( + postcard::Error::DeserializeUnexpectedEnd, + )); + } + if buf[i] & 0x80 == 0 { + return postcard::from_bytes(&buf[..=i]).map_err(StoreError::Decode); } - })() - .transpose() + } + Err(StoreError::Decode(postcard::Error::DeserializeBadVarint)) } } diff --git a/crates/file_store/src/lib.rs b/crates/file_store/src/lib.rs index 3731d50309..d061a7941c 100644 --- a/crates/file_store/src/lib.rs +++ b/crates/file_store/src/lib.rs @@ -4,14 +4,9 @@ mod entry_iter; mod store; use std::io; -use bincode::{DefaultOptions, Options}; pub use entry_iter::*; pub use store::*; -pub(crate) fn bincode_options() -> impl bincode::Options { - DefaultOptions::new().with_varint_encoding() -} - /// Error that occurs due to problems encountered with the file. #[derive(Debug)] pub enum StoreError { @@ -20,7 +15,7 @@ pub enum StoreError { /// Magic bytes do not match what is expected. InvalidMagicBytes { got: Vec, expected: Vec }, /// Failure to decode data from the file. - Bincode(bincode::ErrorKind), + Decode(postcard::Error), } impl core::fmt::Display for StoreError { @@ -34,7 +29,7 @@ impl core::fmt::Display for StoreError { match self { Self::Io(e) => write!(f, "io error while reading store file: {}", e), - Self::Bincode(e) => write!(f, "bincode error while decoding entry {}", e), + Self::Decode(e) => write!(f, "decode error while decoding entry: {}", e), Self::InvalidMagicBytes { got, expected } => { write!(f, "invalid magic bytes: ")?; write!(f, "expected 0x")?; diff --git a/crates/file_store/src/store.rs b/crates/file_store/src/store.rs index ca4dbe471b..244f78e263 100644 --- a/crates/file_store/src/store.rs +++ b/crates/file_store/src/store.rs @@ -1,6 +1,5 @@ -use crate::{bincode_options, EntryIter, StoreError}; +use crate::{EntryIter, StoreError}; use bdk_core::Merge; -use bincode::Options; use std::{ fmt::{self, Debug}, fs::{File, OpenOptions}, @@ -11,6 +10,10 @@ use std::{ /// Persists an append-only list of changesets (`C`) to a single file. /// +/// After the magic bytes, each changeset is stored as a `postcard`-encoded `u64` varint length +/// prefix followed by the `postcard`-encoded changeset. Files written by versions that used +/// `bincode` cannot be read, so use new magic bytes when upgrading. +/// /// > ⚠ This is a development/testing database. It does not natively support backwards compatible /// > BDK version upgrades so should not be used in production. #[derive(Debug)] @@ -60,7 +63,7 @@ where /// /// If there exist changesets in the file, [`load`] will try to aggregate them in /// a single changeset to verify their integrity. If aggregation fails - /// [`StoreErrorWithDump`] will be returned with the [`StoreError::Bincode`] error variant in + /// [`StoreErrorWithDump`] will be returned with the [`StoreError::Decode`] error variant in /// its error field and the aggregated changeset so far in the changeset field. /// /// To get a new working file store from this error use [`Store::create`] and [`Store::append`] @@ -178,7 +181,7 @@ where /// /// If there exist changesets in the file, [`dump`] will try to aggregate them in a single /// changeset. If aggregation fails [`StoreErrorWithDump`] will be returned with the - /// [`StoreError::Bincode`] error variant in its error field and the aggregated changeset so + /// [`StoreError::Decode`] error variant in its error field and the aggregated changeset so /// far in the changeset field. /// /// [`dump`]: Store::dump @@ -256,24 +259,17 @@ where Ok(_) => {} // A real I/O failure means we cannot tell where the valid data ends. Err(StoreError::Io(error)) => return Err(error), - Err(StoreError::Bincode(bincode::ErrorKind::Io(error))) - if error.kind() != io::ErrorKind::UnexpectedEof => - { - return Err(error) - } // Undecodable or partial entry: write from here, as before. Err(_) => break, } } - bincode_options() - .serialize_into(&mut self.db_file, changeset) - .map_err(|e| match *e { - bincode::ErrorKind::Io(error) => error, - unexpected_err => panic!("unexpected bincode error: {unexpected_err}"), - })?; - - Ok(()) + // Each entry is a `postcard` varint length prefix followed by the `postcard` payload, + // written in one call. + let payload = postcard::to_allocvec(changeset).map_err(io::Error::other)?; + let mut frame = postcard::to_allocvec(&(payload.len() as u64)).map_err(io::Error::other)?; + frame.extend_from_slice(&payload); + self.db_file.write_all(&frame) } } @@ -321,6 +317,14 @@ mod test { type TestChangeSet = BTreeSet; + /// The bytes [`Store::append`] writes for `changeset`: varint length prefix + payload. + fn frame(changeset: &TestChangeSet) -> Vec { + let payload = postcard::to_allocvec(changeset).unwrap(); + let mut frame = postcard::to_allocvec(&(payload.len() as u64)).unwrap(); + frame.extend_from_slice(&payload); + frame + } + /// Check behavior of [`Store::create`] and [`Store::load`]. #[test] fn construct_store() { @@ -393,7 +397,7 @@ mod test { match Store::::load(&TEST_MAGIC_BYTES, file_path) { Err(StoreErrorWithDump { changeset, - error: StoreError::Bincode(_), + error: StoreError::Decode(_), }) => { assert_eq!(changeset, Some(Box::new(test_changesets))) } @@ -421,7 +425,7 @@ mod test { match store.dump() { Err(StoreErrorWithDump { changeset, - error: StoreError::Bincode(_), + error: StoreError::Decode(_), }) => { assert_eq!(changeset, Some(Box::new(test_changesets))) } @@ -498,7 +502,7 @@ mod test { TestChangeSet::from(["4".into(), "5".into(), "6".into()]), ]; let last_changeset = TestChangeSet::from(["7".into(), "8".into(), "9".into()]); - let last_changeset_bytes = bincode_options().serialize(&last_changeset).unwrap(); + let last_changeset_bytes = frame(&last_changeset); for short_write_len in 1..last_changeset_bytes.len() - 1 { let file_path = temp_dir.path().join(format!("{short_write_len}.dat")); @@ -604,9 +608,7 @@ mod test { let file_path = temp_dir.path().join("db_file"); let changeset1 = TestChangeSet::from(["a".to_string()]); let changeset2 = TestChangeSet::from(["b".to_string()]); - let torn = bincode_options() - .serialize(&TestChangeSet::from(["zzzzzz".to_string()])) - .unwrap(); + let torn = frame(&TestChangeSet::from(["zzzzzz".to_string()])); let mut first = Store::::create(&TEST_MAGIC_BYTES, &file_path).expect("must create"); @@ -668,6 +670,80 @@ mod test { assert_eq!(fs::read(&file_path).unwrap(), bytes_before); } + /// Write `bytes` after the magic bytes of a fresh file and try to load it. + fn load_raw(bytes: &[u8]) -> Result<(Store, Option), StoreError> { + let temp_dir = tempfile::tempdir().unwrap(); + let file_path = temp_dir.path().join("db_file"); + fs::write(&file_path, [&TEST_MAGIC_BYTES[..], bytes].concat()).unwrap(); + Store::::load(&TEST_MAGIC_BYTES, &file_path).map_err(|e| e.error) + } + + /// The length prefix is 1 byte below 128 and 2 bytes from 128 on; both must round-trip. + #[test] + fn roundtrip_at_varint_length_boundaries() { + let temp_dir = tempfile::tempdir().unwrap(); + let file_path = temp_dir.path().join("db_file"); + // a one-entry set encodes as: count (1 byte) + string length (1 byte) + the string + let overhead = postcard::to_allocvec(&TestChangeSet::from([String::new()])) + .unwrap() + .len(); + + let mut store = Store::::create(&TEST_MAGIC_BYTES, &file_path).unwrap(); + let mut expected = TestChangeSet::new(); + for payload_len in [127_usize, 128] { + let changeset = TestChangeSet::from(["x".repeat(payload_len - overhead)]); + assert_eq!( + postcard::to_allocvec(&changeset).unwrap().len(), + payload_len + ); + store.append(&changeset).unwrap(); + expected.extend(changeset); + } + drop(store); + + let (_, recovered) = Store::::load(&TEST_MAGIC_BYTES, &file_path).unwrap(); + assert_eq!(recovered, Some(expected)); + } + + /// A corrupt, huge length prefix must be reported as a decode error, not panic or abort. + #[test] + fn load_fails_on_oversized_length_prefix() { + let bytes = postcard::to_allocvec(&u64::MAX).unwrap(); + assert!(matches!( + load_raw(&bytes), + Err(StoreError::Decode( + postcard::Error::DeserializeUnexpectedEnd + )) + )); + } + + /// Bytes left in a frame after the payload decoded are corruption. + #[test] + fn load_fails_on_trailing_bytes_in_frame() { + let mut payload = + postcard::to_allocvec(&TestChangeSet::from(["hello".to_string()])).unwrap(); + payload.extend_from_slice(&[0xaa, 0xbb]); + let mut bytes = postcard::to_allocvec(&(payload.len() as u64)).unwrap(); + bytes.extend_from_slice(&payload); + assert!(matches!(load_raw(&bytes), Err(StoreError::Decode(_)))); + } + + /// A length prefix cut short by the end of the file is a torn entry, and so is one that never + /// terminates. + #[test] + fn load_fails_on_torn_or_overlong_length_prefix() { + assert!(matches!( + load_raw(&[0x80]), + Err(StoreError::Decode( + postcard::Error::DeserializeUnexpectedEnd + )) + )); + assert!(matches!( + load_raw(&[0xff; 10]), + Err(StoreError::Decode(postcard::Error::DeserializeBadVarint)) + )); + } + #[test] fn test_load_recovers_state_after_last_write() { let temp_dir = tempfile::tempdir().unwrap();