1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
use std::sync::Arc;
use client_traits::{BlockInfo, ChainNotify};
use common_types::{
ids::BlockId,
io_message::ClientIoMessage,
chain_notify::NewBlocks,
};
use vapory_types::H256;
use vapcore_io::IoChannel;
use log::{trace, warn};
use parking_lot::Mutex;
use crate::traits::{Broadcast, Oracle};
struct StandardOracle<F> where F: 'static + Send + Sync + Fn() -> bool {
client: Arc<dyn BlockInfo>,
sync_status: F,
}
impl<F> Oracle for StandardOracle<F>
where F: Send + Sync + Fn() -> bool
{
fn to_number(&self, hash: H256) -> Option<u64> {
self.client.block_header(BlockId::Hash(hash)).map(|h| h.number())
}
fn is_major_importing(&self) -> bool {
(self.sync_status)()
}
}
impl<C: 'static> Broadcast for Mutex<IoChannel<ClientIoMessage<C>>> {
fn request_snapshot_at(&self, num: u64) {
if let Err(e) = self.lock().send(ClientIoMessage::TakeSnapshot(num)) {
warn!(target: "snapshot_watcher", "Snapshot watcher disconnected from IoService: {}", e);
} else {
trace!(target: "snapshot_watcher", "Snapshot requested at block #{}", num);
}
}
}
pub struct Watcher {
oracle: Box<dyn Oracle>,
broadcast: Box<dyn Broadcast>,
period: u64,
history: u64,
}
impl Watcher {
pub fn new<F, C>(
client: Arc<dyn BlockInfo>,
sync_status: F,
channel: IoChannel<ClientIoMessage<C>>,
period: u64,
history: u64
) -> Self
where
F: 'static + Send + Sync + Fn() -> bool,
C: 'static + Send + Sync,
{
Watcher {
oracle: Box::new(StandardOracle { client, sync_status }),
broadcast: Box::new(Mutex::new(channel)),
period,
history,
}
}
#[cfg(any(test, feature = "test-helpers"))]
pub fn new_test(oracle: Box<dyn Oracle>, broadcast: Box<dyn Broadcast>, period: u64, history: u64) -> Self {
Watcher { oracle, broadcast, period, history }
}
}
impl ChainNotify for Watcher {
fn new_blocks(&self, new_blocks: NewBlocks) {
if self.oracle.is_major_importing() || new_blocks.has_more_blocks_to_import { return }
let highest = new_blocks.imported.into_iter()
.filter_map(|h| self.oracle.to_number(h))
.map(|num| num.saturating_sub(self.history) )
.filter(|num| num % self.period == 0 )
.fold(0, ::std::cmp::max);
if highest > 0 {
self.broadcast.request_snapshot_at(highest);
}
}
}