commit
afb830c91f
|
@ -1,7 +1,7 @@
|
||||||
[package]
|
[package]
|
||||||
name = "silk"
|
name = "silk"
|
||||||
description = "A silky smooth implementation of the Loom architecture"
|
description = "A silky smooth implementation of the Loom architecture"
|
||||||
version = "0.2.1"
|
version = "0.2.2"
|
||||||
documentation = "https://docs.rs/silk"
|
documentation = "https://docs.rs/silk"
|
||||||
homepage = "http://loomprotocol.com/"
|
homepage = "http://loomprotocol.com/"
|
||||||
repository = "https://github.com/loomprotocol/silk"
|
repository = "https://github.com/loomprotocol/silk"
|
||||||
|
|
13
README.md
13
README.md
|
@ -24,28 +24,27 @@ Create a *Historian* and send it *events* to generate an *event log*, where each
|
||||||
is tagged with the historian's latest *hash*. Then ensure the order of events was not tampered
|
is tagged with the historian's latest *hash*. Then ensure the order of events was not tampered
|
||||||
with by verifying each entry's hash can be generated from the hash in the previous entry:
|
with by verifying each entry's hash can be generated from the hash in the previous entry:
|
||||||
|
|
||||||
![historian](https://user-images.githubusercontent.com/55449/36492930-97a572be-16eb-11e8-8289-358e9507189e.png)
|
![historian](https://user-images.githubusercontent.com/55449/36499105-7c8db6a0-16fd-11e8-8b88-c6e0f52d7a50.png)
|
||||||
|
|
||||||
```rust
|
```rust
|
||||||
extern crate silk;
|
extern crate silk;
|
||||||
|
|
||||||
use silk::historian::Historian;
|
use silk::historian::Historian;
|
||||||
use silk::log::{verify_slice, Entry, Event, Sha256Hash};
|
use silk::log::{verify_slice, Entry, Event, Sha256Hash};
|
||||||
use std::{thread, time};
|
use std::thread::sleep;
|
||||||
|
use std::time::Duration;
|
||||||
use std::sync::mpsc::SendError;
|
use std::sync::mpsc::SendError;
|
||||||
|
|
||||||
fn create_log(hist: &Historian) -> Result<(), SendError<Event>> {
|
fn create_log(hist: &Historian) -> Result<(), SendError<Event>> {
|
||||||
hist.sender.send(Event::Tick)?;
|
sleep(Duration::from_millis(15));
|
||||||
thread::sleep(time::Duration::new(0, 100_000));
|
|
||||||
hist.sender.send(Event::UserDataKey(Sha256Hash::default()))?;
|
hist.sender.send(Event::UserDataKey(Sha256Hash::default()))?;
|
||||||
thread::sleep(time::Duration::new(0, 100_000));
|
sleep(Duration::from_millis(10));
|
||||||
hist.sender.send(Event::Tick)?;
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn main() {
|
fn main() {
|
||||||
let seed = Sha256Hash::default();
|
let seed = Sha256Hash::default();
|
||||||
let hist = Historian::new(&seed);
|
let hist = Historian::new(&seed, Some(10));
|
||||||
create_log(&hist).expect("send error");
|
create_log(&hist).expect("send error");
|
||||||
drop(hist.sender);
|
drop(hist.sender);
|
||||||
let entries: Vec<Entry> = hist.receiver.iter().collect();
|
let entries: Vec<Entry> = hist.receiver.iter().collect();
|
||||||
|
|
|
@ -1,8 +1,6 @@
|
||||||
msc {
|
msc {
|
||||||
client,historian,logger;
|
client,historian,logger;
|
||||||
|
|
||||||
client=>historian [ label = "Tick" ] ;
|
|
||||||
historian=>logger [ label = "Tick" ] ;
|
|
||||||
logger=>historian [ label = "e0 = Entry{hash: h0, n: 0, event: Tick}" ] ;
|
logger=>historian [ label = "e0 = Entry{hash: h0, n: 0, event: Tick}" ] ;
|
||||||
logger=>logger [ label = "h1 = hash(h0)" ] ;
|
logger=>logger [ label = "h1 = hash(h0)" ] ;
|
||||||
logger=>logger [ label = "h2 = hash(h1)" ] ;
|
logger=>logger [ label = "h2 = hash(h1)" ] ;
|
||||||
|
@ -13,8 +11,6 @@ msc {
|
||||||
logger=>logger [ label = "h4 = hash(h3)" ] ;
|
logger=>logger [ label = "h4 = hash(h3)" ] ;
|
||||||
logger=>logger [ label = "h5 = hash(h4)" ] ;
|
logger=>logger [ label = "h5 = hash(h4)" ] ;
|
||||||
logger=>logger [ label = "h6 = hash(h5)" ] ;
|
logger=>logger [ label = "h6 = hash(h5)" ] ;
|
||||||
client=>historian [ label = "Tick" ] ;
|
|
||||||
historian=>logger [ label = "Tick" ] ;
|
|
||||||
logger=>historian [ label = "e2 = Entry{hash: h6, n: 3, event: Tick}" ] ;
|
logger=>historian [ label = "e2 = Entry{hash: h6, n: 3, event: Tick}" ] ;
|
||||||
client=>historian [ label = "collect()" ] ;
|
client=>historian [ label = "collect()" ] ;
|
||||||
historian=>client [ label = "entries = [e0, e1, e2]" ] ;
|
historian=>client [ label = "entries = [e0, e1, e2]" ] ;
|
||||||
|
|
|
@ -2,21 +2,20 @@ extern crate silk;
|
||||||
|
|
||||||
use silk::historian::Historian;
|
use silk::historian::Historian;
|
||||||
use silk::log::{verify_slice, Entry, Event, Sha256Hash};
|
use silk::log::{verify_slice, Entry, Event, Sha256Hash};
|
||||||
use std::{thread, time};
|
use std::thread::sleep;
|
||||||
|
use std::time::Duration;
|
||||||
use std::sync::mpsc::SendError;
|
use std::sync::mpsc::SendError;
|
||||||
|
|
||||||
fn create_log(hist: &Historian) -> Result<(), SendError<Event>> {
|
fn create_log(hist: &Historian) -> Result<(), SendError<Event>> {
|
||||||
hist.sender.send(Event::Tick)?;
|
sleep(Duration::from_millis(15));
|
||||||
thread::sleep(time::Duration::new(0, 100_000));
|
|
||||||
hist.sender.send(Event::UserDataKey(Sha256Hash::default()))?;
|
hist.sender.send(Event::UserDataKey(Sha256Hash::default()))?;
|
||||||
thread::sleep(time::Duration::new(0, 100_000));
|
sleep(Duration::from_millis(10));
|
||||||
hist.sender.send(Event::Tick)?;
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn main() {
|
fn main() {
|
||||||
let seed = Sha256Hash::default();
|
let seed = Sha256Hash::default();
|
||||||
let hist = Historian::new(&seed);
|
let hist = Historian::new(&seed, Some(10));
|
||||||
create_log(&hist).expect("send error");
|
create_log(&hist).expect("send error");
|
||||||
drop(hist.sender);
|
drop(hist.sender);
|
||||||
let entries: Vec<Entry> = hist.receiver.iter().collect();
|
let entries: Vec<Entry> = hist.receiver.iter().collect();
|
||||||
|
|
|
@ -7,6 +7,7 @@
|
||||||
|
|
||||||
use std::thread::JoinHandle;
|
use std::thread::JoinHandle;
|
||||||
use std::sync::mpsc::{Receiver, Sender};
|
use std::sync::mpsc::{Receiver, Sender};
|
||||||
|
use std::time::{Duration, SystemTime};
|
||||||
use log::{extend_and_hash, hash, Entry, Event, Sha256Hash};
|
use log::{extend_and_hash, hash, Entry, Event, Sha256Hash};
|
||||||
|
|
||||||
pub struct Historian {
|
pub struct Historian {
|
||||||
|
@ -20,17 +21,12 @@ pub enum ExitReason {
|
||||||
RecvDisconnected,
|
RecvDisconnected,
|
||||||
SendDisconnected,
|
SendDisconnected,
|
||||||
}
|
}
|
||||||
|
fn log_event(
|
||||||
fn log_events(
|
|
||||||
receiver: &Receiver<Event>,
|
|
||||||
sender: &Sender<Entry>,
|
sender: &Sender<Entry>,
|
||||||
num_hashes: &mut u64,
|
num_hashes: &mut u64,
|
||||||
end_hash: &mut Sha256Hash,
|
end_hash: &mut Sha256Hash,
|
||||||
|
event: Event,
|
||||||
) -> Result<(), (Entry, ExitReason)> {
|
) -> Result<(), (Entry, ExitReason)> {
|
||||||
use std::sync::mpsc::TryRecvError;
|
|
||||||
loop {
|
|
||||||
match receiver.try_recv() {
|
|
||||||
Ok(event) => {
|
|
||||||
if let Event::UserDataKey(key) = event {
|
if let Event::UserDataKey(key) = event {
|
||||||
*end_hash = extend_and_hash(end_hash, &key);
|
*end_hash = extend_and_hash(end_hash, &key);
|
||||||
}
|
}
|
||||||
|
@ -43,6 +39,30 @@ fn log_events(
|
||||||
return Err((entry, ExitReason::SendDisconnected));
|
return Err((entry, ExitReason::SendDisconnected));
|
||||||
}
|
}
|
||||||
*num_hashes = 0;
|
*num_hashes = 0;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn log_events(
|
||||||
|
receiver: &Receiver<Event>,
|
||||||
|
sender: &Sender<Entry>,
|
||||||
|
num_hashes: &mut u64,
|
||||||
|
end_hash: &mut Sha256Hash,
|
||||||
|
epoch: SystemTime,
|
||||||
|
num_ticks: &mut u64,
|
||||||
|
ms_per_tick: Option<u64>,
|
||||||
|
) -> Result<(), (Entry, ExitReason)> {
|
||||||
|
use std::sync::mpsc::TryRecvError;
|
||||||
|
loop {
|
||||||
|
if let Some(ms) = ms_per_tick {
|
||||||
|
let now = SystemTime::now();
|
||||||
|
if now > epoch + Duration::from_millis((*num_ticks + 1) * ms) {
|
||||||
|
log_event(sender, num_hashes, end_hash, Event::Tick)?;
|
||||||
|
*num_ticks += 1;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
match receiver.try_recv() {
|
||||||
|
Ok(event) => {
|
||||||
|
log_event(sender, num_hashes, end_hash, event)?;
|
||||||
}
|
}
|
||||||
Err(TryRecvError::Empty) => {
|
Err(TryRecvError::Empty) => {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
|
@ -63,6 +83,7 @@ fn log_events(
|
||||||
/// sending back Entry messages until either the receiver or sender channel is closed.
|
/// sending back Entry messages until either the receiver or sender channel is closed.
|
||||||
pub fn create_logger(
|
pub fn create_logger(
|
||||||
start_hash: Sha256Hash,
|
start_hash: Sha256Hash,
|
||||||
|
ms_per_tick: Option<u64>,
|
||||||
receiver: Receiver<Event>,
|
receiver: Receiver<Event>,
|
||||||
sender: Sender<Entry>,
|
sender: Sender<Entry>,
|
||||||
) -> JoinHandle<(Entry, ExitReason)> {
|
) -> JoinHandle<(Entry, ExitReason)> {
|
||||||
|
@ -70,8 +91,18 @@ pub fn create_logger(
|
||||||
thread::spawn(move || {
|
thread::spawn(move || {
|
||||||
let mut end_hash = start_hash;
|
let mut end_hash = start_hash;
|
||||||
let mut num_hashes = 0;
|
let mut num_hashes = 0;
|
||||||
|
let mut num_ticks = 0;
|
||||||
|
let epoch = SystemTime::now();
|
||||||
loop {
|
loop {
|
||||||
if let Err(err) = log_events(&receiver, &sender, &mut num_hashes, &mut end_hash) {
|
if let Err(err) = log_events(
|
||||||
|
&receiver,
|
||||||
|
&sender,
|
||||||
|
&mut num_hashes,
|
||||||
|
&mut end_hash,
|
||||||
|
epoch,
|
||||||
|
&mut num_ticks,
|
||||||
|
ms_per_tick,
|
||||||
|
) {
|
||||||
return err;
|
return err;
|
||||||
}
|
}
|
||||||
end_hash = hash(&end_hash);
|
end_hash = hash(&end_hash);
|
||||||
|
@ -81,11 +112,11 @@ pub fn create_logger(
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Historian {
|
impl Historian {
|
||||||
pub fn new(start_hash: &Sha256Hash) -> Self {
|
pub fn new(start_hash: &Sha256Hash, ms_per_tick: Option<u64>) -> Self {
|
||||||
use std::sync::mpsc::channel;
|
use std::sync::mpsc::channel;
|
||||||
let (sender, event_receiver) = channel();
|
let (sender, event_receiver) = channel();
|
||||||
let (entry_sender, receiver) = channel();
|
let (entry_sender, receiver) = channel();
|
||||||
let thread_hdl = create_logger(*start_hash, event_receiver, entry_sender);
|
let thread_hdl = create_logger(*start_hash, ms_per_tick, event_receiver, entry_sender);
|
||||||
Historian {
|
Historian {
|
||||||
sender,
|
sender,
|
||||||
receiver,
|
receiver,
|
||||||
|
@ -98,14 +129,13 @@ impl Historian {
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
use log::*;
|
use log::*;
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_historian() {
|
|
||||||
use std::thread::sleep;
|
use std::thread::sleep;
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_historian() {
|
||||||
let zero = Sha256Hash::default();
|
let zero = Sha256Hash::default();
|
||||||
let hist = Historian::new(&zero);
|
let hist = Historian::new(&zero, None);
|
||||||
|
|
||||||
hist.sender.send(Event::Tick).unwrap();
|
hist.sender.send(Event::Tick).unwrap();
|
||||||
sleep(Duration::new(0, 1_000_000));
|
sleep(Duration::new(0, 1_000_000));
|
||||||
|
@ -129,7 +159,7 @@ mod tests {
|
||||||
#[test]
|
#[test]
|
||||||
fn test_historian_closed_sender() {
|
fn test_historian_closed_sender() {
|
||||||
let zero = Sha256Hash::default();
|
let zero = Sha256Hash::default();
|
||||||
let hist = Historian::new(&zero);
|
let hist = Historian::new(&zero, None);
|
||||||
drop(hist.receiver);
|
drop(hist.receiver);
|
||||||
hist.sender.send(Event::Tick).unwrap();
|
hist.sender.send(Event::Tick).unwrap();
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
|
@ -137,4 +167,22 @@ mod tests {
|
||||||
ExitReason::SendDisconnected
|
ExitReason::SendDisconnected
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_ticking_historian() {
|
||||||
|
let zero = Sha256Hash::default();
|
||||||
|
let hist = Historian::new(&zero, Some(20));
|
||||||
|
sleep(Duration::from_millis(30));
|
||||||
|
hist.sender.send(Event::UserDataKey(zero)).unwrap();
|
||||||
|
sleep(Duration::from_millis(15));
|
||||||
|
drop(hist.sender);
|
||||||
|
assert_eq!(
|
||||||
|
hist.thread_hdl.join().unwrap().1,
|
||||||
|
ExitReason::RecvDisconnected
|
||||||
|
);
|
||||||
|
|
||||||
|
let entries: Vec<Entry> = hist.receiver.iter().collect();
|
||||||
|
assert!(entries.len() > 1);
|
||||||
|
assert!(verify_slice(&entries, &zero));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
Loading…
Reference in New Issue