rpc.rs 11.3 KB
Newer Older
1
use bitcoin::consensus::encode::serialize;
Roman Zeyde's avatar
Roman Zeyde committed
2 3
use bitcoin_hashes::hex::{FromHex, ToHex};
use bitcoin_hashes::sha256d::Hash as Sha256dHash;
4
use error_chain::ChainedError;
5
use serde_json::{from_str, Value};
6
use std::collections::HashMap;
7
use std::io::{BufRead, BufReader, Write};
8
use std::net::{Shutdown, SocketAddr, TcpListener, TcpStream};
9
use std::sync::mpsc::SyncSender;
10
use std::sync::{Arc, Mutex};
11
use std::thread;
Roman Zeyde's avatar
WiP  
Roman Zeyde committed
12

Roman Zeyde's avatar
Roman Zeyde committed
13
use crate::errors::*;
kenshin-samourai's avatar
kenshin-samourai committed
14
use crate::query::Query;
15
use crate::util::{spawn_thread, Channel, SyncChannel};
Roman Zeyde's avatar
WiP  
Roman Zeyde committed
16

kenshin-samourai's avatar
kenshin-samourai committed
17
// Indexer version
18
const ADDRINDEXRS_VERSION: &str = env!("CARGO_PKG_VERSION");
kenshin-samourai's avatar
kenshin-samourai committed
19
// Version of the simulated electrum protocol
20
const PROTOCOL_VERSION: &str = "1.4";
21

kenshin-samourai's avatar
kenshin-samourai committed
22 23 24
//
// Get a script hash from a given value
//
Roman Zeyde's avatar
Roman Zeyde committed
25
fn hash_from_value(val: Option<&Value>) -> Result<Sha256dHash> {
kenshin-samourai's avatar
kenshin-samourai committed
26
    // TODO: Sha256dHash should be a generic hash-container (since script hash is single SHA256)
Roman Zeyde's avatar
Roman Zeyde committed
27
    let script_hash = val.chain_err(|| "missing hash")?;
Roman Zeyde's avatar
Roman Zeyde committed
28 29 30 31 32
    let script_hash = script_hash.as_str().chain_err(|| "non-string hash")?;
    let script_hash = Sha256dHash::from_hex(script_hash).chain_err(|| "non-hex hash")?;
    Ok(script_hash)
}

kenshin-samourai's avatar
kenshin-samourai committed
33 34 35
//
// Connection with a RPC client
//
36 37
struct Connection {
    query: Arc<Query>,
38 39
    stream: TcpStream,
    addr: SocketAddr,
40
    chan: SyncChannel<Message>,
41 42
}

43
impl Connection {
44 45 46 47 48
    pub fn new(
        query: Arc<Query>,
        stream: TcpStream,
        addr: SocketAddr,
    ) -> Connection {
49
        Connection {
50
            query,
51 52
            stream,
            addr,
53
            chan: SyncChannel::new(10),
54 55 56
        }
    }

57
    fn server_version(&self) -> Result<Value> {
58
        Ok(json!([
59
            format!("addrindexrs {}", ADDRINDEXRS_VERSION),
60 61
            PROTOCOL_VERSION
        ]))
62
    }
Roman Zeyde's avatar
WiP  
Roman Zeyde committed
63

64 65 66 67 68 69 70
    fn blockchain_headers_subscribe(&mut self) -> Result<Value> {
        let entry = self.query.get_best_header()?;
        let hex_header = hex::encode(serialize(entry.header()));
        let result = json!({"hex": hex_header, "height": entry.height()});
        Ok(result)
    }

71 72 73 74 75 76
    fn blockchain_scripthash_get_balance(&self, _params: &[Value]) -> Result<Value> {
        Ok(
            json!({ "confirmed": null, "unconfirmed": null }),
        )
    }

Roman Zeyde's avatar
Roman Zeyde committed
77
    fn blockchain_scripthash_get_history(&self, params: &[Value]) -> Result<Value> {
Roman Zeyde's avatar
Roman Zeyde committed
78
        let script_hash = hash_from_value(params.get(0)).chain_err(|| "bad script_hash")?;
79
        let status = self.query.status(&script_hash[..])?;
Roman Zeyde's avatar
Roman Zeyde committed
80
        Ok(json!(Value::Array(
81 82
            status
                .history()
Roman Zeyde's avatar
Roman Zeyde committed
83
                .into_iter()
kenshin-samourai's avatar
kenshin-samourai committed
84
                .map(|item| json!({"tx_hash": item.to_hex()}))
Roman Zeyde's avatar
Roman Zeyde committed
85 86
                .collect()
        )))
87
    }
Roman Zeyde's avatar
WiP  
Roman Zeyde committed
88

Chiguireitor's avatar
Chiguireitor committed
89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109
    fn blockchain_scripthash_get_utxos(&self, params: &[Value]) -> Result<Value> {
        let script_hash = hash_from_value(params.get(0)).chain_err(|| "bad script_hash")?;
        let status = self.query.status(&script_hash[..])?;

        let mut dict = HashMap::new();
        for item in status.funding().into_iter() {
            dict.insert(item.txid.to_hex() + ":" + &item.vout.to_string(), "");
        }

        for item in status.spending().into_iter() {
            dict.remove( &(item.outpoint.0.to_hex() + ":" + &item.outpoint.1.to_string()) );
        }

        let mut utxos = vec![];
        for (outpoint, _drop) in &dict {
            utxos.push(outpoint)
        }

        Ok(json!(utxos))
    }

110
    fn handle_command(&mut self, method: &str, params: &[Value], id: &Value) -> Result<Value> {
111
        let result = match method {
112
            "blockchain.headers.subscribe" => self.blockchain_headers_subscribe(),
113
            "blockchain.scripthash.get_balance" => self.blockchain_scripthash_get_balance(&params),
114
            "blockchain.scripthash.get_history" => self.blockchain_scripthash_get_history(&params),
Chiguireitor's avatar
Chiguireitor committed
115
            "blockchain.scripthash.get_utxos" => self.blockchain_scripthash_get_utxos(&params),
116
            "server.ping" => Ok(Value::Null),
117
            "server.version" => self.server_version(),
118
            &_ => bail!("unknown method {} {:?}", method, params),
119
        };
120
        // TODO: return application errors should be sent to the client
121 122 123
        Ok(match result {
            Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
            Err(e) => {
Roman Zeyde's avatar
Roman Zeyde committed
124 125 126 127 128 129 130
                warn!(
                    "rpc #{} {} {:?} failed: {}",
                    id,
                    method,
                    params,
                    e.display_chain()
                );
131 132 133
                json!({"jsonrpc": "2.0", "id": id, "error": format!("{}", e)})
            }
        })
134 135
    }

136 137 138 139 140 141 142 143
    fn send_values(&mut self, values: &[Value]) -> Result<()> {
        for value in values {
            let line = value.to_string() + "\n";
            self.stream
                .write_all(line.as_bytes())
                .chain_err(|| format!("failed to send {}", value))?;
        }
        Ok(())
144 145
    }

146
    fn handle_replies(&mut self) -> Result<()> {
147
        let empty_params = json!([]);
Roman Zeyde's avatar
Roman Zeyde committed
148
        loop {
149
            let msg = self.chan.receiver().recv().chain_err(|| "channel closed")?;
150
            trace!("RPC {:?}", msg);
Roman Zeyde's avatar
Roman Zeyde committed
151 152 153
            match msg {
                Message::Request(line) => {
                    let cmd: Value = from_str(&line).chain_err(|| "invalid JSON format")?;
154 155 156 157 158
                    let reply = match (
                        cmd.get("method"),
                        cmd.get("params").unwrap_or_else(|| &empty_params),
                        cmd.get("id"),
                    ) {
Roman Zeyde's avatar
Roman Zeyde committed
159 160
                        (
                            Some(&Value::String(ref method)),
161
                            &Value::Array(ref params),
162
                            Some(ref id),
Roman Zeyde's avatar
Roman Zeyde committed
163 164 165
                        ) => self.handle_command(method, params, id)?,
                        _ => bail!("invalid command: {}", cmd),
                    };
166
                    self.send_values(&[reply])?
167
                }
168
                Message::Done => return Ok(()),
169
            }
Roman Zeyde's avatar
Roman Zeyde committed
170 171 172
        }
    }

173
    fn handle_requests(mut reader: BufReader<TcpStream>, tx: SyncSender<Message>) -> Result<()> {
Roman Zeyde's avatar
Roman Zeyde committed
174
        loop {
175
            let mut line = Vec::<u8>::new();
Roman Zeyde's avatar
Roman Zeyde committed
176
            reader
177 178
                .read_until(b'\n', &mut line)
                .chain_err(|| "failed to read a request")?;
Roman Zeyde's avatar
Roman Zeyde committed
179
            if line.is_empty() {
180 181
                tx.send(Message::Done).chain_err(|| "channel closed")?;
                return Ok(());
Roman Zeyde's avatar
Roman Zeyde committed
182
            } else {
183 184 185 186 187
                if line.starts_with(&[22, 3, 1]) {
                    // (very) naive SSL handshake detection
                    let _ = tx.send(Message::Done);
                    bail!("invalid request - maybe SSL-encrypted data?: {:?}", line)
                }
188
                match String::from_utf8(line) {
Roman Zeyde's avatar
Roman Zeyde committed
189 190
                    Ok(req) => tx
                        .send(Message::Request(req))
191 192 193
                        .chain_err(|| "channel closed")?,
                    Err(err) => {
                        let _ = tx.send(Message::Done);
194
                        bail!("invalid UTF8: {}", err)
195 196
                    }
                }
Roman Zeyde's avatar
Roman Zeyde committed
197 198 199 200
            }
        }
    }

201
    pub fn run(mut self) {
202
        let reader = BufReader::new(self.stream.try_clone().expect("failed to clone TcpStream"));
203
        let tx = self.chan.sender();
204
        let child = spawn_thread("reader", || Connection::handle_requests(reader, tx));
205
        if let Err(e) = self.handle_replies() {
206 207 208 209 210 211
            error!(
                "[{}] connection handling failed: {}",
                self.addr,
                e.display_chain().to_string()
            );
        }
212
        debug!("[{}] shutting down connection", self.addr);
213
        let _ = self.stream.shutdown(Shutdown::Both);
214 215
        if let Err(err) = child.join().expect("receiver panicked") {
            error!("[{}] receiver failed: {}", self.addr, err);
216
        }
217 218 219
    }
}

kenshin-samourai's avatar
kenshin-samourai committed
220 221 222
//
// Messages supported by the RPC API
//
223
#[derive(Debug)]
224 225 226 227 228
pub enum Message {
    Request(String),
    Done,
}

kenshin-samourai's avatar
kenshin-samourai committed
229 230 231
//
// RPC server
//
232
pub struct RPC {
Roman Zeyde's avatar
Roman Zeyde committed
233
    server: Option<thread::JoinHandle<()>>, // so we can join the server while dropping this ojbect
234 235 236
}

impl RPC {
237
    fn start_acceptor(addr: SocketAddr) -> Channel<Option<(TcpStream, SocketAddr)>> {
Roman Zeyde's avatar
Roman Zeyde committed
238
        let chan = Channel::unbounded();
239
        let acceptor = chan.sender();
240
        spawn_thread("acceptor", move || {
Roman Zeyde's avatar
Roman Zeyde committed
241 242
            let listener =
                TcpListener::bind(addr).unwrap_or_else(|e| panic!("bind({}) failed: {}", addr, e));
243
            info!(
kenshin-samourai's avatar
kenshin-samourai committed
244
                "Indexer RPC server running on {} (protocol {})",
245 246
                addr, PROTOCOL_VERSION
            );
247 248
            loop {
                let (stream, addr) = listener.accept().expect("accept failed");
Roman Zeyde's avatar
Roman Zeyde committed
249 250 251
                stream
                    .set_nonblocking(false)
                    .expect("failed to set connection as blocking");
252
                acceptor.send(Some((stream, addr))).expect("send failed");
253 254 255 256 257
            }
        });
        chan
    }

258
    pub fn start(addr: SocketAddr, query: Arc<Query>) -> RPC {
Roman Zeyde's avatar
Roman Zeyde committed
259
        RPC {
Roman Zeyde's avatar
Roman Zeyde committed
260
            server: Some(spawn_thread("rpc", move || {
261 262 263 264 265
                let senders = Arc::new(Mutex::new(HashMap::<i32, SyncSender<Message>>::new()));
                let handles = Arc::new(Mutex::new(
                    HashMap::<i32, std::thread::JoinHandle<()>>::new(),
                ));

266
                let acceptor = RPC::start_acceptor(addr);
267
                let mut handle_count = 0;
kenshin-samourai's avatar
kenshin-samourai committed
268

269
                while let Some((stream, addr)) = acceptor.receiver().recv().unwrap() {
270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292
                    let handle_id = handle_count;
                    handle_count += 1;
                    // explicitely scope the shadowed variables for the new thread
                    let handle: thread::JoinHandle<()> = {
                        let query = Arc::clone(&query);
                        let senders = Arc::clone(&senders);
                        let handles = Arc::clone(&handles);

                        spawn_thread("peer", move || {
                            info!("[{}] connected peer #{}", addr, handle_id);
                            let conn = Connection::new(query, stream, addr);
                            senders
                                .lock()
                                .unwrap()
                                .insert(handle_id, conn.chan.sender());
                            conn.run();
                            info!("[{}] disconnected peer #{}", addr, handle_id);
                            senders.lock().unwrap().remove(&handle_id);
                            handles.lock().unwrap().remove(&handle_id);
                        })
                    };

                    handles.lock().unwrap().insert(handle_id, handle);
293
                }
kenshin-samourai's avatar
kenshin-samourai committed
294

295
                trace!("closing {} RPC connections", senders.lock().unwrap().len());
296
                for sender in senders.lock().unwrap().values() {
297 298
                    let _ = sender.send(Message::Done);
                }
kenshin-samourai's avatar
kenshin-samourai committed
299

300 301 302 303 304
                trace!("waiting for {} RPC handling threads", handles.lock().unwrap().len());
                for (_, handle) in handles.lock().unwrap().drain() {
                    if let Err(e) = handle.join() {
                        warn!("failed to join thread: {:?}", e);
                    }
305
                }
kenshin-samourai's avatar
kenshin-samourai committed
306

Roman Zeyde's avatar
Roman Zeyde committed
307
                trace!("RPC connections are closed");
Roman Zeyde's avatar
Roman Zeyde committed
308
            })),
Roman Zeyde's avatar
Roman Zeyde committed
309
        }
310
    }
Roman Zeyde's avatar
Roman Zeyde committed
311
}
312

Roman Zeyde's avatar
Roman Zeyde committed
313 314
impl Drop for RPC {
    fn drop(&mut self) {
Roman Zeyde's avatar
Roman Zeyde committed
315
        trace!("stop accepting new RPCs");
Roman Zeyde's avatar
Roman Zeyde committed
316 317 318
        if let Some(handle) = self.server.take() {
            handle.join().unwrap();
        }
Roman Zeyde's avatar
Roman Zeyde committed
319
        trace!("RPC server is stopped");
320
    }
Roman Zeyde's avatar
Roman Zeyde committed
321
}