ping fixes
This commit is contained in:
+13
-7
@@ -1,7 +1,7 @@
|
||||
use failure::{Error, err_msg};
|
||||
use tungstenite::Message;
|
||||
use tungstenite::protocol::WebSocket;
|
||||
use tungstenite::server::accept;
|
||||
use tungstenite::error::Error;
|
||||
use tungstenite::Message::Binary;
|
||||
use std::net::{TcpListener, TcpStream};
|
||||
use serde_cbor::{to_vec};
|
||||
@@ -35,18 +35,19 @@ struct RpcErrorResponse {
|
||||
err: String
|
||||
}
|
||||
|
||||
fn receive(db: Db, rpc: &Rpc, msg: Message, client: &mut WebSocket<TcpStream>) -> Result<(), Error> {
|
||||
fn receive(db: Db, rpc: &Rpc, msg: Message, client: &mut WebSocket<TcpStream>) -> Result<String, Error> {
|
||||
match rpc.receive(msg, &db, client) {
|
||||
Ok(reply) => {
|
||||
let response = to_vec(&reply)
|
||||
.expect("failed to serialize response");
|
||||
client.write_message(Binary(response))
|
||||
client.write_message(Binary(response))?;
|
||||
return Ok(reply.method);
|
||||
},
|
||||
Err(e) => {
|
||||
info!("{:?}", e);
|
||||
let response = to_vec(&RpcErrorResponse { err: e.to_string() })
|
||||
.expect("failed to serialize error response");
|
||||
client.write_message(Binary(response))
|
||||
client.write_message(Binary(response))?;
|
||||
return Err(err_msg(e));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -120,8 +121,13 @@ pub fn start() {
|
||||
let begin = Instant::now();
|
||||
let db_connection = db.get().expect("unable to get db connection");
|
||||
match receive(db_connection, &rpc, msg, &mut websocket) {
|
||||
Ok(_r) => info!("response sent. total duration: {:?}", begin.elapsed()),
|
||||
Err(e) => info!("{:?}", e),
|
||||
Ok(method) => {
|
||||
match method.as_ref() {
|
||||
"pong" => (),
|
||||
_ => info!("response sent. total duration: {:?}", begin.elapsed()),
|
||||
}
|
||||
},
|
||||
Err(e) => warn!("{:?}", e),
|
||||
}
|
||||
},
|
||||
// connection is closed
|
||||
|
||||
+7
-1
@@ -35,6 +35,10 @@ impl Rpc {
|
||||
// cast the msg to this type to receive method name
|
||||
match from_slice::<RpcMessage>(&data) {
|
||||
Ok(v) => {
|
||||
if v.method == "ping" {
|
||||
return Ok(RpcResponse { method: "pong".to_string(), params: RpcResult::Pong(()) });
|
||||
}
|
||||
|
||||
info!("message method: {:?}", v.method);
|
||||
let mut tx = db.transaction()?;
|
||||
|
||||
@@ -382,7 +386,7 @@ impl Rpc {
|
||||
|
||||
#[derive(Debug,Clone,Serialize,Deserialize)]
|
||||
pub struct RpcResponse {
|
||||
method: String,
|
||||
pub method: String,
|
||||
params: RpcResult,
|
||||
}
|
||||
|
||||
@@ -402,6 +406,8 @@ pub enum RpcResult {
|
||||
|
||||
InstanceList(Vec<Instance>),
|
||||
InstanceState(Instance),
|
||||
|
||||
Pong(()),
|
||||
}
|
||||
|
||||
#[derive(Debug,Clone,Serialize,Deserialize)]
|
||||
|
||||
Reference in New Issue
Block a user