Unverified Commit 06d0c8d7 authored by Roman Zeyde's avatar Roman Zeyde
Browse files

Split Store trait into its R/W parts

parent 4747e380
......@@ -12,7 +12,7 @@ use std::sync::{Arc, RwLock};
use std::time::{Duration, Instant};
use daemon::Daemon;
use store::{Row, Store};
use store::{DBStore, ReadStore, Row, WriteStore};
use util::{full_hash, hash_prefix, Bytes, FullHash, HashPrefix, HeaderEntry, HeaderList,
HeaderMap, HASH_PREFIX_LEN};
......@@ -198,7 +198,7 @@ fn index_block(block: &Block, height: usize) -> Vec<Row> {
rows
}
fn read_indexed_headers(store: &Store) -> HeaderMap {
fn read_indexed_headers(store: &ReadStore) -> HeaderMap {
let mut headers = HeaderMap::new();
for row in store.scan(b"B") {
let key: BlockKey = bincode::deserialize(&row.key).unwrap();
......@@ -208,7 +208,7 @@ fn read_indexed_headers(store: &Store) -> HeaderMap {
headers
}
fn read_last_indexed_blockhash(store: &Store) -> Option<Sha256dHash> {
fn read_last_indexed_blockhash(store: &ReadStore) -> Option<Sha256dHash> {
let row = store.get(b"L")?;
let blockhash: Sha256dHash = deserialize(&row).unwrap();
Some(blockhash)
......@@ -365,7 +365,7 @@ impl Index {
missing_headers
}
pub fn update(&self, store: &Store, daemon: &Daemon) -> Sha256dHash {
pub fn update(&self, store: &DBStore, daemon: &Daemon) -> Sha256dHash {
let mut indexed_headers: Arc<HeaderList> = self.headers_list();
let no_indexed_headers = indexed_headers.headers().is_empty();
if no_indexed_headers {
......@@ -383,7 +383,7 @@ impl Index {
/*use_progress_bar=*/ no_indexed_headers,
)) {
// TODO: add timing
store.persist(rows);
store.write(rows);
}
let tip = current_headers.tip();
*(self.headers.write().unwrap()) = Arc::new(current_headers);
......
......@@ -9,7 +9,7 @@ use std::time::{Duration, Instant};
use daemon::{Daemon, MempoolEntry};
use index::index_transaction;
use store::{Row, Store};
use store::{ReadStore, Row};
use util::Bytes;
error_chain!{}
......@@ -21,18 +21,19 @@ struct MempoolStore {
}
impl MempoolStore {
fn new() -> MempoolStore {
fn new(rows: Vec<Row>) -> MempoolStore {
let mut map = BTreeMap::new();
for row in rows {
let (key, value) = row.into_pair();
map.insert(key, value);
}
MempoolStore {
map: RwLock::new(BTreeMap::new()),
map: RwLock::new(map),
}
}
fn clear(&self) {
self.map.write().unwrap().clear()
}
}
impl Store for MempoolStore {
impl ReadStore for MempoolStore {
fn get(&self, key: &[u8]) -> Option<Bytes> {
self.map.read().unwrap().get(key).map(|v| v.to_vec())
}
......@@ -51,13 +52,6 @@ impl Store for MempoolStore {
}
rows
}
fn persist(&self, rows: Vec<Row>) {
let mut map = self.map.write().unwrap();
for row in rows {
let (key, value) = row.into_pair();
map.insert(key, value);
}
}
}
pub struct Stats {
......@@ -80,7 +74,7 @@ impl Tracker {
pub fn new() -> Tracker {
Tracker {
stats: HashMap::new(),
index: MempoolStore::new(),
index: MempoolStore::new(vec![]),
}
}
......@@ -134,14 +128,11 @@ impl Tracker {
for txid in old_txids.difference(&new_txids) {
self.remove(txid);
}
let mut rows = Vec::new();
for stats in self.stats.values() {
index_transaction(&stats.tx, 0, &mut rows)
}
self.index.clear();
self.index.persist(rows);
self.index = MempoolStore::new(rows);
debug!(
"mempool update took {:.1} ms ({} txns)",
t.elapsed().in_seconds() * 1e3,
......@@ -158,7 +149,7 @@ impl Tracker {
self.stats.remove(txid);
}
pub fn index(&self) -> &Store {
pub fn index(&self) -> &ReadStore {
&self.index
}
}
......
......@@ -10,7 +10,7 @@ use std::sync::RwLock;
use daemon::Daemon;
use index::{compute_script_hash, Index, TxInRow, TxOutRow, TxRow};
use mempool::Tracker;
use store::Store;
use store::ReadStore;
use util::{FullHash, HashPrefix, HeaderEntry};
error_chain!{}
......@@ -84,8 +84,8 @@ fn merklize(left: Sha256dHash, right: Sha256dHash) -> Sha256dHash {
Sha256dHash::from_data(&data)
}
// TODO: the 3 functions below can be part of Store.
fn txrows_by_prefix(store: &Store, txid_prefix: &HashPrefix) -> Vec<TxRow> {
// TODO: the 3 functions below can be part of ReadStore.
fn txrows_by_prefix(store: &ReadStore, txid_prefix: &HashPrefix) -> Vec<TxRow> {
store
.scan(&TxRow::filter(&txid_prefix))
.iter()
......@@ -93,7 +93,7 @@ fn txrows_by_prefix(store: &Store, txid_prefix: &HashPrefix) -> Vec<TxRow> {
.collect()
}
fn txids_by_script_hash(store: &Store, script_hash: &[u8]) -> Vec<HashPrefix> {
fn txids_by_script_hash(store: &ReadStore, script_hash: &[u8]) -> Vec<HashPrefix> {
store
.scan(&TxOutRow::filter(script_hash))
.iter()
......@@ -102,7 +102,7 @@ fn txids_by_script_hash(store: &Store, script_hash: &[u8]) -> Vec<HashPrefix> {
}
fn txids_by_funding_output(
store: &Store,
store: &ReadStore,
txn_id: &Sha256dHash,
output_index: usize,
) -> Vec<HashPrefix> {
......@@ -116,13 +116,13 @@ fn txids_by_funding_output(
pub struct Query<'a> {
daemon: &'a Daemon,
index: &'a Index,
index_store: &'a Store, // TODO: should be a part of index
index_store: &'a ReadStore, // TODO: should be a part of index
tracker: RwLock<Tracker>,
}
// TODO: return errors instead of panics
impl<'a> Query<'a> {
pub fn new(index_store: &'a Store, daemon: &'a Daemon, index: &'a Index) -> Query<'a> {
pub fn new(index_store: &'a ReadStore, daemon: &'a Daemon, index: &'a Index) -> Query<'a> {
Query {
daemon,
index,
......@@ -135,7 +135,7 @@ impl<'a> Query<'a> {
self.daemon
}
fn load_txns(&self, store: &Store, prefixes: Vec<HashPrefix>) -> Vec<TxnHeight> {
fn load_txns(&self, store: &ReadStore, prefixes: Vec<HashPrefix>) -> Vec<TxnHeight> {
let mut txns = Vec::new();
for txid_prefix in prefixes {
for tx_row in txrows_by_prefix(store, &txid_prefix) {
......@@ -150,7 +150,11 @@ impl<'a> Query<'a> {
txns
}
fn find_spending_input(&self, store: &Store, funding: &FundingOutput) -> Option<SpendingInput> {
fn find_spending_input(
&self,
store: &ReadStore,
funding: &FundingOutput,
) -> Option<SpendingInput> {
let spending_txns: Vec<TxnHeight> = self.load_txns(
store,
txids_by_funding_output(store, &funding.txn_id, funding.output_index),
......
......@@ -15,10 +15,13 @@ impl Row {
}
}
pub trait Store: Sync {
pub trait ReadStore: Sync {
fn get(&self, key: &[u8]) -> Option<Bytes>;
fn scan(&self, prefix: &[u8]) -> Vec<Row>;
fn persist(&self, rows: Vec<Row>);
}
pub trait WriteStore: Sync {
fn write(&self, rows: Vec<Row>);
}
pub struct DBStore {
......@@ -66,7 +69,7 @@ impl DBStore {
}
}
impl Store for DBStore {
impl ReadStore for DBStore {
fn get(&self, key: &[u8]) -> Option<Bytes> {
self.db.get(key).unwrap().map(|v| v.to_vec())
}
......@@ -89,8 +92,10 @@ impl Store for DBStore {
}
rows
}
}
fn persist(&self, rows: Vec<Row>) {
impl WriteStore for DBStore {
fn write(&self, rows: Vec<Row>) {
let mut batch = rocksdb::WriteBatch::default();
for row in rows {
batch.put(row.key.as_slice(), row.value.as_slice()).unwrap();
......
Markdown is supported
0% or .
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment