Skip to content

Commit

Permalink
Added clean stop logic on server and peers. Broken unit test.
Browse files Browse the repository at this point in the history
  • Loading branch information
ignopeverell committed Oct 30, 2016
1 parent 42769c3 commit a23308d
Show file tree
Hide file tree
Showing 5 changed files with 82 additions and 20 deletions.
4 changes: 4 additions & 0 deletions p2p/src/peer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,4 +48,8 @@ impl Peer {
pub fn run(&self, na: &NetAdapter) -> Option<Error> {
self.proto.handle(na)
}

pub fn stop(&self) {
self.proto.as_ref().close()
}
}
38 changes: 23 additions & 15 deletions p2p/src/protocol.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ use std::ops::DerefMut;
use std::rc::Rc;

use mioco;
use mioco::sync::mpsc::sync_channel;
use mioco::sync::mpsc::{sync_channel, SyncSender};
use mioco::tcp::{TcpStream, Shutdown};

use core::ser;
Expand All @@ -27,12 +27,19 @@ use types::*;

pub struct ProtocolV1 {
conn: RefCell<TcpStream>,
//msg_send: Option<SyncSender<ser::Writeable>>,
stop_send: RefCell<Option<SyncSender<u8>>>,
}

impl Protocol for ProtocolV1 {
fn handle(&self, server: &NetAdapter) -> Option<ser::Error> {
// setup a channel so we can switch between reads and writes
let (send, recv) = sync_channel(10);
// setup channels so we can switch between reads, writes and close
let (msg_send, msg_recv) = sync_channel(10);
let (stop_send, stop_recv) = sync_channel(1);

//self.msg_send = Some(msg_send);
let mut stop_mut = self.stop_send.borrow_mut();
*stop_mut = Some(stop_send);

let mut conn = self.conn.borrow_mut();
loop {
Expand All @@ -43,25 +50,26 @@ impl Protocol for ProtocolV1 {
continue;
}
},
r:recv => {
ser::serialize(conn.deref_mut(), recv.recv().unwrap());
r:msg_recv => {
ser::serialize(conn.deref_mut(), msg_recv.recv().unwrap());
},
r:stop_recv => {
stop_recv.recv();
conn.shutdown(Shutdown::Both);
return None;;
}
);
}
}

fn close(&self) {
let stop_send = self.stop_send.borrow();
stop_send.as_ref().unwrap().send(0);
}
}

impl ProtocolV1 {
pub fn new(conn: TcpStream) -> ProtocolV1 {
ProtocolV1 { conn: RefCell::new(conn) }
ProtocolV1 { conn: RefCell::new(conn), /* msg_send: None, */ stop_send: RefCell::new(None) }
}

// fn close(&mut self, err_code: u32, explanation: &'static str) {
// ser::serialize(self.conn,
// &PeerError {
// code: err_code,
// message: explanation.to_string(),
// });
// self.conn.shutdown(Shutdown::Both);
// }
}
18 changes: 13 additions & 5 deletions p2p/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
use std::cell::RefCell;
use std::io;
use std::net::SocketAddr;
use std::ops::Deref;
use std::str::FromStr;
use std::sync::{Arc, RwLock};

Expand Down Expand Up @@ -52,7 +53,7 @@ impl Server {
}
/// Creates a new p2p server. Opens a TCP port to allow incoming
/// connections and starts the bootstrapping process to find peers.
pub fn start(&'static self) -> Result<(), Error> {
pub fn start(& self) -> Result<(), Error> {
mioco::spawn(move || -> io::Result<()> {
let addr = DEFAULT_LISTEN_ADDR.parse().unwrap();
let listener = try!(TcpListener::bind(&addr));
Expand All @@ -67,10 +68,10 @@ impl Server {
mioco::spawn(move || -> io::Result<()> {
let peer = try!(Peer::connect(conn, &hs).map_err(|_| io::Error::last_os_error()));
let wpeer = Arc::new(peer);
{
let mut peers = self.peers.write().unwrap();
peers.push(wpeer.clone());
}
// {
// let mut peers = self.peers.write().unwrap();
// peers.push(wpeer.clone());
// }
if let Some(err) = wpeer.run(&DummyAdapter{}) {
error!("{:?}", err);
}
Expand All @@ -82,6 +83,13 @@ impl Server {
Ok(())
}

pub fn stop(&self) {
let peers = self.peers.write().unwrap();
for p in peers.deref() {
p.stop();
}
}

/// Simulates an unrelated client connecting to our server. Mostly used for
/// tests.
pub fn connect_as_client(addr: SocketAddr) -> Result<Peer, Error> {
Expand Down
2 changes: 2 additions & 0 deletions p2p/src/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,8 @@ pub trait Protocol {
/// Starts handling protocol communication, the peer(s) is expected to be
/// known already, usually passed during construction.
fn handle(&self, na: &NetAdapter) -> Option<Error>;

fn close(&self);
}

/// Bridge between the networking layer and the rest of the system. Handles the
Expand Down
40 changes: 40 additions & 0 deletions p2p/tests/network_conn.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
// Copyright 2016 The Grin Developers
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

extern crate grin_p2p as p2p;
extern crate mioco;
extern crate env_logger;

use std::io;
use std::thread;
use std::time;

#[test]
fn peer_handshake() {
env_logger::init().unwrap();

mioco::start(|| -> io::Result<()> {
let server = p2p::Server::new();
server.start();

// given server a little time to start
mioco::sleep(time::Duration::from_millis(200));

let addr = p2p::DEFAULT_LISTEN_ADDR.parse().unwrap();
try!(p2p::Server::connect_as_client(addr).map_err(|_| io::Error::last_os_error()));

server.stop();
mioco::shutdown();
}).unwrap().unwrap();
}

0 comments on commit a23308d

Please sign in to comment.