From e00a2fac9245f43ad660d33f8dffeb53570644de Mon Sep 17 00:00:00 2001 From: Martin Hoffmann Date: Tue, 10 May 2016 12:12:37 +0200 Subject: [PATCH] Fix shutting down of resolver. --- examples/dig.rs | 5 +++- src/resolv/stub/dispatcher.rs | 50 ++++++++++++++++++++--------------- src/resolv/stub/mod.rs | 4 +-- src/resolv/stub/stream.rs | 4 +-- src/resolv/stub/sync.rs | 44 +++++++++++++++++++++++++----- src/resolv/stub/tcp.rs | 26 +++++++++++++++--- src/resolv/stub/udp.rs | 1 - 7 files changed, 97 insertions(+), 37 deletions(-) diff --git a/examples/dig.rs b/examples/dig.rs index fb62c38c..297de2e3 100644 --- a/examples/dig.rs +++ b/examples/dig.rs @@ -230,6 +230,8 @@ fn main() { let query = Query::new(&name, options.qtype().unwrap(), options.qclass().unwrap()); let response = resolver.sync_task(query).unwrap(); + println!("Done with task"); + /* let len = response.len(); print_result(response); println!(";; Query time: not yet available."); @@ -237,7 +239,8 @@ fn main() { println!(";; WHEN: not yet available."); println!(";; MSG SIZE rcvd: {} bytes", len); println!(""); - + */ + println!("Dropping resolver"); drop(resolver); join.join().unwrap(); } diff --git a/src/resolv/stub/dispatcher.rs b/src/resolv/stub/dispatcher.rs index 549cff7d..12f38483 100644 --- a/src/resolv/stub/dispatcher.rs +++ b/src/resolv/stub/dispatcher.rs @@ -8,7 +8,7 @@ use resolv::conf::ResolvConf; use resolv::error::Error; use super::conn::{ConnCommand, ConnTransportSeed}; use super::query::Query; -use super::sync::{RotorReceiver, RotorSender, SharedNotifier}; +use super::sync::{RotorReceiver, RotorSender, SharedNotifier, channel}; use super::udp::{UdpCommand, UdpTransportSeed}; @@ -21,9 +21,17 @@ pub struct Dispatcher { /// The query queue. /// - /// Resolvers will place their queries here and transports will return - /// unanswered queries here, too. - queries: RotorReceiver, + /// Resolvers will place their queries here. If this queue disconnects, + /// the dispatcher will close. + queries: mpsc::Receiver, + + /// The failed query queue. + /// + /// Transports will return failed queries into here. We need a + /// separate queue for this since we need to be able to spawn new + /// transports and, subsequently, need to hold a sender to clone in + /// this case. Because of that, the queue will never disconnect. + failed: RotorReceiver, /// Are we in bootstrap state? bootstrap: Option, @@ -50,10 +58,12 @@ pub struct Dispatcher { impl Dispatcher { /// Creates a new dispatcher from the given configuration. pub fn new(conf: ResolvConf, scope: &mut S) - -> Dispatcher { + -> (Self, RotorSender) { + let (tx, rx) = channel(Some(scope.notifier())); let mut res = Dispatcher { conf: conf, - queries: RotorReceiver::new(Some(scope.notifier())), + queries: rx, + failed: RotorReceiver::new(Some(scope.notifier())), bootstrap: None, dgram_servers: Vec::new(), dgram_start: 0, @@ -63,12 +73,7 @@ impl Dispatcher { }; res.configure(); scope.notifier().wakeup().unwrap(); - res - } - - /// Returns a sender for the query queue - pub fn query_sender(&self) -> RotorSender { - self.queries.sender() + (res, tx) } /* @@ -108,17 +113,17 @@ impl Dispatcher { fn configure(&mut self) { let mut bs = DispatcherBootstrap::new(); let (udp4, udp4_tx) = UdpTransportSeed::new(self.conf.clone(), - self.queries.sender(), + self.failed.sender(), false); let mut use_udp4 = false; let (udp6, udp6_tx) = UdpTransportSeed::new(self.conf.clone(), - self.queries.sender(), + self.failed.sender(), true); let mut use_udp6 = false; for addr in self.conf.servers.iter() { let (tcp, tcp_tx) = ConnTransportSeed::new(self.conf.clone(), addr.clone(), - self.queries.sender()); + self.failed.sender()); self.stream_servers.push(Server::Other(tcp_tx.clone(), tcp.notifier())); let server = match addr { @@ -265,18 +270,21 @@ impl Machine for Dispatcher { else { loop { match self.queries.try_recv() { - Ok(query) => { - self.dispatch(query); - } + Ok(query) => self.dispatch(query), + Err(mpsc::TryRecvError::Empty) => break, Err(mpsc::TryRecvError::Disconnected) => { self.close(); return Response::done(); } - Err(mpsc::TryRecvError::Empty) => { - return Response::ok(self); - } } } + loop { + match self.failed.try_recv() { + Ok(query) => self.dispatch(query), + _ => break, + } + } + return Response::ok(self); } } } diff --git a/src/resolv/stub/mod.rs b/src/resolv/stub/mod.rs index f13a496d..d8b7b594 100644 --- a/src/resolv/stub/mod.rs +++ b/src/resolv/stub/mod.rs @@ -36,8 +36,8 @@ impl DnsTransport { /// Returns the transport and a resolver. pub fn new(conf: ResolvConf, scope: &mut S) -> (Self, Resolver) { - let dispatcher = Dispatcher::new(conf, scope); - let resolver = Resolver::new(dispatcher.query_sender()); + let (dispatcher, tx) = Dispatcher::new(conf, scope); + let resolver = Resolver::new(tx); (DnsTransport(Composition::Dispatcher(dispatcher)), resolver) } diff --git a/src/resolv/stub/stream.rs b/src/resolv/stub/stream.rs index 3c710082..be16191f 100644 --- a/src/resolv/stub/stream.rs +++ b/src/resolv/stub/stream.rs @@ -66,8 +66,8 @@ impl StreamTransportInfo { pub fn conf(&self) -> &ResolvConf { &self.conf } pub fn can_read(&self) -> bool { !self.timeouts.is_empty() } - pub fn can_write(&self) -> bool { !self.send_queue.can_write() } - + pub fn can_write(&self) -> bool { self.send_queue.can_write() } + /// Processes the command queue. /// /// Returns whether a close command was received. diff --git a/src/resolv/stub/sync.rs b/src/resolv/stub/sync.rs index 882f3de1..4b74056b 100644 --- a/src/resolv/stub/sync.rs +++ b/src/resolv/stub/sync.rs @@ -48,6 +48,9 @@ impl SharedNotifier { //------------ RotorReceiver ------------------------------------------------ /// An `mpsc::Receiver` that can produce a `RotorSender` if necessary. +/// +/// Note that because the receiver internally holds a sender for later +/// cloning, it will never disconnect. pub struct RotorReceiver { receiver: mpsc::Receiver, sender: mpsc::Sender, @@ -79,27 +82,54 @@ impl RotorReceiver { /// A sender to a mpsc channel waking up the receiver. #[derive(Debug)] pub struct RotorSender { - sender: mpsc::Sender, + sender: Option>, notifier: Option } impl RotorSender { pub fn new(sender: mpsc::Sender, notifier: Option) -> Self { - RotorSender { sender: sender, notifier: notifier } + RotorSender { sender: Some(sender), notifier: notifier } } pub fn send(&self, t: T) -> Result<(), mpsc::SendError> { - try!(self.sender.send(t)); - if let Some(ref notifier) = self.notifier { - let _ = notifier.wakeup(); + match self.sender { + None => Err(mpsc::SendError(t)), + Some(ref sender) => { + try!(sender.send(t)); + if let Some(ref notifier) = self.notifier { + let _ = notifier.wakeup(); + } + Ok(()) + } } - Ok(()) } } impl Clone for RotorSender { fn clone(&self) -> Self { - RotorSender::new(self.sender.clone(), self.notifier.clone()) + RotorSender { + sender: self.sender.clone(), + notifier: self.notifier.clone(), + } } } +impl Drop for RotorSender { + fn drop(&mut self) { + self.sender = None; + match self.notifier { + None => None, + Some(ref notifier) => notifier.wakeup().ok() + }; + } +} + + +//------------ Functions ---------------------------------------------------- + +pub fn channel(notifier: Option) + -> (RotorSender, mpsc::Receiver) { + let (tx, rx) = mpsc::channel(); + (RotorSender::new(tx, notifier), rx) +} + diff --git a/src/resolv/stub/tcp.rs b/src/resolv/stub/tcp.rs index dbf3b756..0a899915 100644 --- a/src/resolv/stub/tcp.rs +++ b/src/resolv/stub/tcp.rs @@ -16,7 +16,8 @@ pub struct TcpTransport(State); impl TcpTransport { - pub fn new(seed: ConnTransportSeed, _scope: &mut Scope) -> Self { + pub fn new(seed: ConnTransportSeed, scope: &mut Scope) -> Self { + seed.notifier.set(scope.notifier()); Idle::new(StreamTransportInfo::new(seed)) } } @@ -128,12 +129,28 @@ impl Idle { unreachable!() } - fn wakeup(self, scope: &mut Scope) + fn wakeup(mut self, scope: &mut Scope) -> Response, ConnTransportSeed> { - Response::ok(Connecting::new(self.info, scope)) + if self.info.process_commands() { + // We got a close command and can shut down right away. + println!("TCP transport done"); + Response::done() + } + else { + if self.info.can_write() { + Response::ok(Connecting::new(self.info, scope)) + } + else { + Response::ok(self.into()) + } + } } } +impl From> for TcpTransport { + fn from(idle: Idle) -> Self { TcpTransport(State::Idle(idle)) } +} + //------------ Connecting --------------------------------------------------- @@ -367,6 +384,7 @@ impl Closing { -> Response, ConnTransportSeed> { self.info.flush_timeouts(); scope.deregister(&self.sock).ok(); + println!("TCP transport done"); Response::done() } @@ -403,6 +421,7 @@ impl Closing { } } else { + println!("TCP transport done"); Response::done() } } @@ -439,6 +458,7 @@ impl Failed { fn wakeup(mut self, _scope: &mut Scope) -> Response, ConnTransportSeed> { if self.info.reject_commands() { + println!("TCP transport done"); Response::done() } else { diff --git a/src/resolv/stub/udp.rs b/src/resolv/stub/udp.rs index 9cf04ed8..ae7f6aab 100644 --- a/src/resolv/stub/udp.rs +++ b/src/resolv/stub/udp.rs @@ -265,7 +265,6 @@ impl UdpTransport { match self.timeouts.next_timeout() { None => Response::ok(self), Some(timeout) => { - println!("TIMEOUT: {:?}", timeout); Response::ok(self).deadline(timeout) } }