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
#![forbid(unsafe_code)]
use anyhow::Result;
use diem_config::config::NodeConfig;
use diem_logger::prelude::*;
use diem_secure_net::NetworkServer;
use diem_types::{account_state_blob::AccountStateBlob, proof::SparseMerkleProof};
use diemdb::DiemDB;
use std::{
sync::Arc,
thread::{self, JoinHandle},
};
use storage_interface::{DbReader, DbWriter, Error, StartupInfo};
pub fn start_storage_service_with_db(config: &NodeConfig, diem_db: Arc<DiemDB>) -> JoinHandle<()> {
let storage_service = StorageService { db: diem_db };
storage_service.run(config)
}
#[derive(Clone)]
pub struct StorageService {
db: Arc<DiemDB>,
}
impl StorageService {
fn handle_message(&self, input_message: Vec<u8>) -> Result<Vec<u8>, Error> {
let input = bcs::from_bytes(&input_message)?;
let output = match input {
storage_interface::StorageRequest::GetAccountStateWithProofByVersionRequest(req) => {
bcs::to_bytes(&self.get_account_state_with_proof_by_version(&req))
}
storage_interface::StorageRequest::GetStartupInfoRequest => {
bcs::to_bytes(&self.get_startup_info())
}
storage_interface::StorageRequest::SaveTransactionsRequest(req) => {
bcs::to_bytes(&self.save_transactions(&req))
}
};
Ok(output?)
}
fn get_account_state_with_proof_by_version(
&self,
req: &storage_interface::GetAccountStateWithProofByVersionRequest,
) -> Result<
(
Option<AccountStateBlob>,
SparseMerkleProof<AccountStateBlob>,
),
Error,
> {
Ok(self
.db
.get_account_state_with_proof_by_version(req.address, req.version)?)
}
fn get_startup_info(&self) -> Result<Option<StartupInfo>, Error> {
Ok(self.db.get_startup_info()?)
}
fn save_transactions(
&self,
req: &storage_interface::SaveTransactionsRequest,
) -> Result<(), Error> {
Ok(self.db.save_transactions(
&req.txns_to_commit,
req.first_version,
req.ledger_info_with_signatures.as_ref(),
)?)
}
fn run(self, config: &NodeConfig) -> JoinHandle<()> {
let mut network_server =
NetworkServer::new("storage", config.storage.address, config.storage.timeout_ms);
let ret = thread::spawn(move || loop {
if let Err(e) = self.process_one_message(&mut network_server) {
warn!(
error = ?e,
"Failed to process message.",
);
}
});
info!("StorageService spawned.");
ret
}
fn process_one_message(&self, network_server: &mut NetworkServer) -> Result<(), Error> {
let request = network_server.read()?;
let response = self.handle_message(request)?;
network_server.write(&response)?;
Ok(())
}
}
#[cfg(test)]
mod storage_service_test;