Fix shutting down of resolver.

This commit is contained in:
Martin Hoffmann
2016-05-10 12:12:37 +02:00
parent 27661713d1
commit e00a2fac92
7 changed files with 97 additions and 37 deletions
+4 -1
View File
@@ -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();
}
+29 -21
View File
@@ -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<X> {
/// The query queue.
///
/// Resolvers will place their queries here and transports will return
/// unanswered queries here, too.
queries: RotorReceiver<Query>,
/// Resolvers will place their queries here. If this queue disconnects,
/// the dispatcher will close.
queries: mpsc::Receiver<Query>,
/// 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<Query>,
/// Are we in bootstrap state?
bootstrap: Option<DispatcherBootstrap>,
@@ -50,10 +58,12 @@ pub struct Dispatcher<X> {
impl<X> Dispatcher<X> {
/// Creates a new dispatcher from the given configuration.
pub fn new<S: GenericScope>(conf: ResolvConf, scope: &mut S)
-> Dispatcher<X> {
-> (Self, RotorSender<Query>) {
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<X> Dispatcher<X> {
};
res.configure();
scope.notifier().wakeup().unwrap();
res
}
/// Returns a sender for the query queue
pub fn query_sender(&self) -> RotorSender<Query> {
self.queries.sender()
(res, tx)
}
/*
@@ -108,17 +113,17 @@ impl<X> Dispatcher<X> {
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<X> Machine for Dispatcher<X> {
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);
}
}
}
+2 -2
View File
@@ -36,8 +36,8 @@ impl<X> DnsTransport<X> {
/// Returns the transport and a resolver.
pub fn new<S: GenericScope>(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)
}
+2 -2
View File
@@ -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.
+37 -7
View File
@@ -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<T> {
receiver: mpsc::Receiver<T>,
sender: mpsc::Sender<T>,
@@ -79,27 +82,54 @@ impl<T> RotorReceiver<T> {
/// A sender to a mpsc channel waking up the receiver.
#[derive(Debug)]
pub struct RotorSender<T> {
sender: mpsc::Sender<T>,
sender: Option<mpsc::Sender<T>>,
notifier: Option<Notifier>
}
impl<T> RotorSender<T> {
pub fn new(sender: mpsc::Sender<T>, notifier: Option<Notifier>) -> Self {
RotorSender { sender: sender, notifier: notifier }
RotorSender { sender: Some(sender), notifier: notifier }
}
pub fn send(&self, t: T) -> Result<(), mpsc::SendError<T>> {
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<T> Clone for RotorSender<T> {
fn clone(&self) -> Self {
RotorSender::new(self.sender.clone(), self.notifier.clone())
RotorSender {
sender: self.sender.clone(),
notifier: self.notifier.clone(),
}
}
}
impl<T> Drop for RotorSender<T> {
fn drop(&mut self) {
self.sender = None;
match self.notifier {
None => None,
Some(ref notifier) => notifier.wakeup().ok()
};
}
}
//------------ Functions ----------------------------------------------------
pub fn channel<T>(notifier: Option<Notifier>)
-> (RotorSender<T>, mpsc::Receiver<T>) {
let (tx, rx) = mpsc::channel();
(RotorSender::new(tx, notifier), rx)
}
+23 -3
View File
@@ -16,7 +16,8 @@ pub struct TcpTransport<X>(State<X>);
impl<X> TcpTransport<X> {
pub fn new(seed: ConnTransportSeed, _scope: &mut Scope<X>) -> Self {
pub fn new(seed: ConnTransportSeed, scope: &mut Scope<X>) -> Self {
seed.notifier.set(scope.notifier());
Idle::new(StreamTransportInfo::new(seed))
}
}
@@ -128,12 +129,28 @@ impl<X> Idle<X> {
unreachable!()
}
fn wakeup(self, scope: &mut Scope<X>)
fn wakeup(mut self, scope: &mut Scope<X>)
-> Response<TcpTransport<X>, 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<X> From<Idle<X>> for TcpTransport<X> {
fn from(idle: Idle<X>) -> Self { TcpTransport(State::Idle(idle)) }
}
//------------ Connecting ---------------------------------------------------
@@ -367,6 +384,7 @@ impl<X> Closing<X> {
-> Response<TcpTransport<X>, ConnTransportSeed> {
self.info.flush_timeouts();
scope.deregister(&self.sock).ok();
println!("TCP transport done");
Response::done()
}
@@ -403,6 +421,7 @@ impl<X> Closing<X> {
}
}
else {
println!("TCP transport done");
Response::done()
}
}
@@ -439,6 +458,7 @@ impl<X> Failed<X> {
fn wakeup(mut self, _scope: &mut Scope<X>)
-> Response<TcpTransport<X>, ConnTransportSeed> {
if self.info.reject_commands() {
println!("TCP transport done");
Response::done()
}
else {
-1
View File
@@ -265,7 +265,6 @@ impl<X> UdpTransport<X> {
match self.timeouts.next_timeout() {
None => Response::ok(self),
Some(timeout) => {
println!("TIMEOUT: {:?}", timeout);
Response::ok(self).deadline(timeout)
}
}