-
Notifications
You must be signed in to change notification settings - Fork 355
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Implement persistence with the new structures #965
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,89 @@ | ||
use core::convert::Infallible; | ||
|
||
use crate::Append; | ||
|
||
/// `Persist` wraps a [`PersistBackend`] (`B`) to create a convenient staging area for changes (`C`) | ||
/// before they are persisted. | ||
/// | ||
/// Not all changes to the in-memory representation needs to be written to disk right away, so | ||
/// [`Persist::stage`] can be used to *stage* changes first and then [`Persist::commit`] can be used | ||
/// to write changes to disk. | ||
#[derive(Debug)] | ||
pub struct Persist<B, C> { | ||
backend: B, | ||
stage: C, | ||
} | ||
|
||
impl<B, C> Persist<B, C> | ||
where | ||
B: PersistBackend<C>, | ||
C: Default + Append, | ||
{ | ||
/// Create a new [`Persist`] from [`PersistBackend`]. | ||
pub fn new(backend: B) -> Self { | ||
Self { | ||
backend, | ||
stage: Default::default(), | ||
} | ||
} | ||
|
||
/// Stage a `changeset` to be commited later with [`commit`]. | ||
/// | ||
/// [`commit`]: Self::commit | ||
pub fn stage(&mut self, changeset: C) { | ||
self.stage.append(changeset) | ||
} | ||
|
||
/// Get the changes that have not been commited yet. | ||
pub fn staged(&self) -> &C { | ||
&self.stage | ||
} | ||
|
||
/// Commit the staged changes to the underlying persistance backend. | ||
/// | ||
/// Returns a backend-defined error if this fails. | ||
pub fn commit(&mut self) -> Result<(), B::WriteError> { | ||
let mut temp = C::default(); | ||
core::mem::swap(&mut temp, &mut self.stage); | ||
self.backend.write_changes(&temp) | ||
Comment on lines
+46
to
+48
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Mhhhh maybe put a small comment here, it took me a second to figure out what you were trying to do (commit changes and clean self.stage at the same time) Also, a shorter way of writing this, but not necessarily a cleaner one, might be:
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. actually shouldn't we only clear it in the case that it succeeds. There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. That's a good point There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. /// Commit the staged changes to the underlying persistance backend.
///
/// Changes that are committed (if any) are returned.
///
/// # Error
///
/// Returns a backend-defined error if this fails.
pub fn commit(&mut self) -> Result<Option<C>, B::WriteError> {
if self.stage.is_empty() {
return Ok(None);
}
self.backend
.write_changes(&self.stage)
// if written successfully, take and return `self.stage`
.map(|_| Some(core::mem::take(&mut self.stage)))
} There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. |
||
} | ||
} | ||
|
||
/// A persistence backend for [`Persist`]. | ||
/// | ||
/// `C` represents the changeset; a datatype that records changes made to in-memory data structures | ||
/// that are to be persisted, or retrieved from persistence. | ||
pub trait PersistBackend<C> { | ||
/// The error the backend returns when it fails to write. | ||
type WriteError: core::fmt::Debug; | ||
|
||
/// The error the backend returns when it fails to load changesets `C`. | ||
type LoadError: core::fmt::Debug; | ||
|
||
/// Writes a changeset to the persistence backend. | ||
/// | ||
/// It is up to the backend what it does with this. It could store every changeset in a list or | ||
/// it inserts the actual changes into a more structured database. All it needs to guarantee is | ||
/// that [`load_from_persistence`] restores a keychain tracker to what it should be if all | ||
/// changesets had been applied sequentially. | ||
/// | ||
/// [`load_from_persistence`]: Self::load_from_persistence | ||
fn write_changes(&mut self, changeset: &C) -> Result<(), Self::WriteError>; | ||
|
||
/// Return the aggregate changeset `C` from persistence. | ||
fn load_from_persistence(&mut self) -> Result<C, Self::LoadError>; | ||
} | ||
|
||
impl<C: Default> PersistBackend<C> for () { | ||
type WriteError = Infallible; | ||
|
||
type LoadError = Infallible; | ||
|
||
fn write_changes(&mut self, _changeset: &C) -> Result<(), Self::WriteError> { | ||
Ok(()) | ||
} | ||
|
||
fn load_from_persistence(&mut self) -> Result<C, Self::LoadError> { | ||
Ok(C::default()) | ||
} | ||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,6 +1,7 @@ | ||
use crate::collections::BTreeMap; | ||
use crate::collections::BTreeSet; | ||
use crate::BlockId; | ||
use alloc::vec::Vec; | ||
use bitcoin::{Block, OutPoint, Transaction, TxOut}; | ||
|
||
/// Trait to do something with every txout contained in a structure. | ||
|
@@ -64,20 +65,56 @@ impl<A: Anchor> Anchor for &'static A { | |
pub trait Append { | ||
/// Append another object of the same type onto `self`. | ||
fn append(&mut self, other: Self); | ||
|
||
/// Returns whether the structure is considered empty. | ||
fn is_empty(&self) -> bool; | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Let's say I wanted to implement There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I think implementing |
||
} | ||
|
||
impl Append for () { | ||
fn append(&mut self, _other: Self) {} | ||
|
||
fn is_empty(&self) -> bool { | ||
true | ||
} | ||
} | ||
|
||
impl<K: Ord, V> Append for BTreeMap<K, V> { | ||
fn append(&mut self, mut other: Self) { | ||
BTreeMap::append(self, &mut other) | ||
} | ||
|
||
fn is_empty(&self) -> bool { | ||
BTreeMap::is_empty(self) | ||
} | ||
} | ||
|
||
impl<T: Ord> Append for BTreeSet<T> { | ||
fn append(&mut self, mut other: Self) { | ||
BTreeSet::append(self, &mut other) | ||
} | ||
|
||
fn is_empty(&self) -> bool { | ||
BTreeSet::is_empty(self) | ||
} | ||
} | ||
|
||
impl<T> Append for Vec<T> { | ||
fn append(&mut self, mut other: Self) { | ||
Vec::append(self, &mut other) | ||
} | ||
|
||
fn is_empty(&self) -> bool { | ||
Vec::is_empty(self) | ||
} | ||
} | ||
|
||
impl<A: Append, B: Append> Append for (A, B) { | ||
fn append(&mut self, other: Self) { | ||
Append::append(&mut self.0, other.0); | ||
Append::append(&mut self.1, other.1); | ||
} | ||
|
||
fn is_empty(&self) -> bool { | ||
Append::is_empty(&self.0) && Append::is_empty(&self.1) | ||
} | ||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,100 @@ | ||
use bincode::Options; | ||
use std::{ | ||
fs::File, | ||
io::{self, 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`. | ||
/// | ||
/// [`next`]: Self::next | ||
pub struct EntryIter<'t, T> { | ||
db_file: Option<&'t mut File>, | ||
|
||
/// The file position for the first read of `db_file`. | ||
start_pos: Option<u64>, | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. nit: maybe expand in the docs why it's an Option and what happens when this is None |
||
types: PhantomData<T>, | ||
} | ||
|
||
impl<'t, T> EntryIter<'t, T> { | ||
pub fn new(start_pos: u64, db_file: &'t mut File) -> Self { | ||
Self { | ||
db_file: Some(db_file), | ||
start_pos: Some(start_pos), | ||
types: PhantomData, | ||
} | ||
} | ||
} | ||
|
||
impl<'t, T> Iterator for EntryIter<'t, T> | ||
where | ||
T: serde::de::DeserializeOwned, | ||
{ | ||
type Item = Result<T, IterError>; | ||
|
||
fn next(&mut self) -> Option<Self::Item> { | ||
evanlinjin marked this conversation as resolved.
Show resolved
Hide resolved
|
||
// closure which reads a single entry starting from `self.pos` | ||
let read_one = |f: &mut File, start_pos: Option<u64>| -> Result<Option<T>, IterError> { | ||
let pos = match start_pos { | ||
Some(pos) => f.seek(io::SeekFrom::Start(pos))?, | ||
None => f.stream_position()?, | ||
}; | ||
|
||
match bincode_options().deserialize_from(&*f) { | ||
Ok(changeset) => { | ||
f.stream_position()?; | ||
Ok(Some(changeset)) | ||
} | ||
Err(e) => { | ||
if let bincode::ErrorKind::Io(inner) = &*e { | ||
if inner.kind() == io::ErrorKind::UnexpectedEof { | ||
let eof = f.seek(io::SeekFrom::End(0))?; | ||
if pos == eof { | ||
return Ok(None); | ||
} | ||
} | ||
} | ||
f.seek(io::SeekFrom::Start(pos))?; | ||
Err(IterError::Bincode(*e)) | ||
} | ||
} | ||
}; | ||
|
||
let result = read_one(self.db_file.as_mut()?, self.start_pos.take()); | ||
if result.is_err() { | ||
self.db_file = None; | ||
} | ||
result.transpose() | ||
} | ||
} | ||
|
||
impl From<io::Error> for IterError { | ||
fn from(value: io::Error) -> Self { | ||
IterError::Io(value) | ||
} | ||
} | ||
|
||
/// Error type for [`EntryIter`]. | ||
#[derive(Debug)] | ||
pub enum IterError { | ||
/// Failure to read from the file. | ||
Io(io::Error), | ||
/// Failure to decode data from the file. | ||
Bincode(bincode::ErrorKind), | ||
} | ||
|
||
impl core::fmt::Display for IterError { | ||
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { | ||
match self { | ||
IterError::Io(e) => write!(f, "io error trying to read entry {}", e), | ||
IterError::Bincode(e) => write!(f, "bincode error while reading entry {}", e), | ||
} | ||
} | ||
} | ||
|
||
impl std::error::Error for IterError {} |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
nit: maybe it should be S for stage, instead of C?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
C
for changeset though 😮There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Then do
changeset: C
:PThere was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I support
changeset: C
🙂