This should be a more robust client that self-heals

This commit is contained in:
Brandon Thomas
2021-01-14 00:22:52 -05:00
parent 3e4f7aa5e2
commit 03b167d377
+92 -25
View File
@@ -5,47 +5,114 @@ mod color;
use point::Point;
use std::thread;
use zmq::Error;
use zmq::{Error, Socket, Context, DONTWAIT};
use std::time::Duration;
use crate::color::Color;
const SOCKET_ADDRESS : &'static str = "tcp://127.0.0.1:8888";
enum MessagingState {
Sending,
Receiving,
}
fn main() {
let ctx = zmq::Context::new();
let context = zmq::Context::new();
//let socket = ctx.socket(zmq::PUSH).unwrap();
//let socket = ctx.socket(zmq::REQ).unwrap();
let socket = ctx.socket(zmq::PUSH).unwrap();
let mut socket = context.socket(zmq::REQ).unwrap();
socket.connect("tcp://127.0.0.1:5555").unwrap();
//socket.send("hello world!", 0).unwrap();
//socket.connect("tcp://127.0.0.1:5555").unwrap();
//socket.bind("tcp://127.0.0.1:5555").unwrap();
socket.connect(SOCKET_ADDRESS).unwrap();
let mut reconnect = false;
let mut fail_count = 0;
let mut messaging_state = MessagingState::Sending;
loop {
let point = Point::at_random_range(-1000.0f32, 1000.0f32, Color::White);
let bytes = point.to_bytes();
println!("Location: {}", point.location_string());
let bytes = point.to_bytes();
if reconnect {
socket = reconnect_socket(&context, socket, SOCKET_ADDRESS);
reconnect = false;
messaging_state = MessagingState::Sending;
}
//let result = socket.send("hello world!", 0);
let result = socket.send(&bytes, 0);
match messaging_state {
MessagingState::Sending => {
println!("Sending request...");
//let result = socket.send("hello world!", 0);
let result = socket.send(&bytes, 0);
match result {
Ok(_) => {
println!("Sent!");
thread::sleep(Duration::from_millis(250));
messaging_state = MessagingState::Receiving;
},
Err(e) => {
eprintln!("Send Error: {:?}", e);
thread::sleep(Duration::from_millis(2000));
fail_count += 1;
},
}
},
MessagingState::Receiving => {
println!("Awaiting response...");
let result = socket.recv_bytes(DONTWAIT);
match result {
Ok(_) => {
println!("Response received!");
messaging_state = MessagingState::Sending;
},
Err(e) => {
eprintln!("Recv Error: {:?}", e);
thread::sleep(Duration::from_millis(2000));
fail_count += 1;
},
}
match result {
Ok(_) => {},
Err(e) => {
eprintln!("Send Error: {:?}", e);
thread::sleep(Duration::from_millis(5_000));
continue;
},
}
/*let result = socket.recv_msg(0);
match result {
Ok(_) => {},
Err(e) => {
eprintln!("Recv Error: {:?}", e);
thread::sleep(Duration::from_millis(5_000));
continue;
},
}*/
if fail_count > 5 {
reconnect = true;
fail_count = 0;
thread::sleep(Duration::from_millis(5000));
}
}
}
fn reconnect_socket(context: &Context, socket: Socket, address: &str) -> Socket {
println!("[reconnect] Creating new socket...");
let mut socket = match context.socket(zmq::REQ) {
Ok(s) => {
println!("New socket created.");
s
},
Err(e) => {
println!("Error creating new socket: {:?}", e);
return socket;
},
};
println!("Connecting new socket...");
match socket.connect("tcp://127.0.0.1:5555") {
Ok(_) => {
println!("New socket connected.");
},
Err(err) => {
println!("Error connecting new socket: {:?}", err);
},
}
return socket;
}