From 03b167d3774068b43c2cff98b89d2913fdbe4a4b Mon Sep 17 00:00:00 2001 From: Brandon Thomas Date: Thu, 14 Jan 2021 00:22:52 -0500 Subject: [PATCH] This should be a more robust client that self-heals --- src/main.rs | 117 +++++++++++++++++++++++++++++++++++++++++----------- 1 file changed, 92 insertions(+), 25 deletions(-) diff --git a/src/main.rs b/src/main.rs index 16afbd4..8bea6e6 100644 --- a/src/main.rs +++ b/src/main.rs @@ -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; +} \ No newline at end of file