resolve deadlock, fix component.rs

This commit is contained in:
lumi 2017-05-27 18:41:54 +02:00
commit 1550c52552
5 changed files with 55 additions and 50 deletions

View file

@ -77,7 +77,7 @@ impl ClientBuilder {
let host = &self.host.unwrap_or(self.jid.domain.clone()); let host = &self.host.unwrap_or(self.jid.domain.clone());
let mut transport = SslTransport::connect(host, self.port)?; let mut transport = SslTransport::connect(host, self.port)?;
C2S::init(&mut transport, &self.jid.domain, "before_sasl")?; C2S::init(&mut transport, &self.jid.domain, "before_sasl")?;
let dispatcher = Arc::new(Mutex::new(Dispatcher::new())); let dispatcher = Arc::new(Dispatcher::new());
let mut credentials = self.credentials; let mut credentials = self.credentials;
credentials.channel_binding = transport.channel_bind(); credentials.channel_binding = transport.channel_bind();
let transport = Arc::new(Mutex::new(transport)); let transport = Arc::new(Mutex::new(transport));
@ -88,7 +88,7 @@ impl ClientBuilder {
binding: PluginProxyBinding::new(dispatcher.clone()), binding: PluginProxyBinding::new(dispatcher.clone()),
dispatcher: dispatcher, dispatcher: dispatcher,
}; };
client.dispatcher.lock().unwrap().register(Priority::Default, move |evt: &SendElement| { client.dispatcher.register(Priority::Default, move |evt: &SendElement| {
let mut t = transport.lock().unwrap(); let mut t = transport.lock().unwrap();
t.write_element(&evt.0).unwrap(); t.write_element(&evt.0).unwrap();
Propagation::Continue Propagation::Continue
@ -105,7 +105,7 @@ pub struct Client {
transport: Arc<Mutex<SslTransport>>, transport: Arc<Mutex<SslTransport>>,
plugins: HashMap<TypeId, Arc<Box<Plugin>>>, plugins: HashMap<TypeId, Arc<Box<Plugin>>>,
binding: PluginProxyBinding, binding: PluginProxyBinding,
dispatcher: Arc<Mutex<Dispatcher>>, dispatcher: Arc<Dispatcher>,
} }
impl Client { impl Client {
@ -119,10 +119,7 @@ impl Client {
let binding = self.binding.clone(); let binding = self.binding.clone();
plugin.bind(binding); plugin.bind(binding);
let p = Arc::new(Box::new(plugin) as Box<Plugin>); let p = Arc::new(Box::new(plugin) as Box<Plugin>);
{ P::init(&self.dispatcher, p.clone());
let mut disp = self.dispatcher.lock().unwrap();
P::init(&mut disp, p.clone());
}
if self.plugins.insert(TypeId::of::<P>(), p).is_some() { if self.plugins.insert(TypeId::of::<P>(), p).is_some() {
panic!("registering a plugin that's already registered"); panic!("registering a plugin that's already registered");
} }
@ -132,7 +129,7 @@ impl Client {
where where
E: Event, E: Event,
F: Fn(&E) -> Propagation + 'static { F: Fn(&E) -> Propagation + 'static {
self.dispatcher.lock().unwrap().register(pri, func); self.dispatcher.register(pri, func);
} }
/// Returns the plugin given by the type parameter, if it exists, else panics. /// Returns the plugin given by the type parameter, if it exists, else panics.
@ -146,14 +143,11 @@ impl Client {
/// Returns the next event and flush the send queue. /// Returns the next event and flush the send queue.
pub fn main(&mut self) -> Result<(), Error> { pub fn main(&mut self) -> Result<(), Error> {
self.dispatcher.lock().unwrap().flush_all(); self.dispatcher.flush_all();
loop { loop {
let elem = self.read_element()?; let elem = self.read_element()?;
{ self.dispatcher.dispatch(ReceiveElement(elem));
let mut disp = self.dispatcher.lock().unwrap(); self.dispatcher.flush_all();
disp.dispatch(ReceiveElement(elem));
disp.flush_all();
}
} }
} }

View file

@ -4,7 +4,7 @@ use transport::{Transport, PlainTransport};
use error::Error; use error::Error;
use ns; use ns;
use plugin::{Plugin, PluginInit, PluginProxyBinding}; use plugin::{Plugin, PluginInit, PluginProxyBinding};
use event::{Dispatcher, ReceiveElement}; use event::{Dispatcher, ReceiveElement, SendElement, Propagation, Priority, Event};
use connection::{Connection, Component2S}; use connection::{Connection, Component2S};
use sha_1::{Sha1, Digest}; use sha_1::{Sha1, Digest};
@ -61,15 +61,20 @@ impl ComponentBuilder {
let host = &self.host.unwrap_or(self.jid.domain.clone()); let host = &self.host.unwrap_or(self.jid.domain.clone());
let mut transport = PlainTransport::connect(host, self.port)?; let mut transport = PlainTransport::connect(host, self.port)?;
Component2S::init(&mut transport, &self.jid.domain, "stream_opening")?; Component2S::init(&mut transport, &self.jid.domain, "stream_opening")?;
let dispatcher = Arc::new(Mutex::new(Dispatcher::new())); let dispatcher = Arc::new(Dispatcher::new());
let transport = Arc::new(Mutex::new(transport)); let transport = Arc::new(Mutex::new(transport));
let mut component = Component { let mut component = Component {
jid: self.jid, jid: self.jid,
transport: transport, transport: transport.clone(),
plugins: HashMap::new(), plugins: HashMap::new(),
binding: PluginProxyBinding::new(dispatcher.clone()), binding: PluginProxyBinding::new(dispatcher.clone()),
dispatcher: dispatcher, dispatcher: dispatcher,
}; };
component.dispatcher.register(Priority::Default, move |evt: &SendElement| {
let mut t = transport.lock().unwrap();
t.write_element(&evt.0).unwrap();
Propagation::Continue
});
component.connect(self.secret)?; component.connect(self.secret)?;
Ok(component) Ok(component)
} }
@ -81,7 +86,7 @@ pub struct Component {
transport: Arc<Mutex<PlainTransport>>, transport: Arc<Mutex<PlainTransport>>,
plugins: HashMap<TypeId, Arc<Box<Plugin>>>, plugins: HashMap<TypeId, Arc<Box<Plugin>>>,
binding: PluginProxyBinding, binding: PluginProxyBinding,
dispatcher: Arc<Mutex<Dispatcher>>, dispatcher: Arc<Dispatcher>,
} }
impl Component { impl Component {
@ -95,10 +100,7 @@ impl Component {
let binding = self.binding.clone(); let binding = self.binding.clone();
plugin.bind(binding); plugin.bind(binding);
let p = Arc::new(Box::new(plugin) as Box<Plugin>); let p = Arc::new(Box::new(plugin) as Box<Plugin>);
{ P::init(&self.dispatcher, p.clone());
let mut disp = self.dispatcher.lock().unwrap();
P::init(&mut disp, p.clone());
}
if self.plugins.insert(TypeId::of::<P>(), p).is_some() { if self.plugins.insert(TypeId::of::<P>(), p).is_some() {
panic!("registering a plugin that's already registered"); panic!("registering a plugin that's already registered");
} }
@ -113,16 +115,20 @@ impl Component {
.expect("plugin downcast failure (should not happen!!)") .expect("plugin downcast failure (should not happen!!)")
} }
pub fn register_handler<E, F>(&mut self, pri: Priority, func: F)
where
E: Event,
F: Fn(&E) -> Propagation + 'static {
self.dispatcher.register(pri, func);
}
/// Returns the next event and flush the send queue. /// Returns the next event and flush the send queue.
pub fn main(&mut self) -> Result<(), Error> { pub fn main(&mut self) -> Result<(), Error> {
self.dispatcher.lock().unwrap().flush_all(); self.dispatcher.flush_all();
loop { loop {
let elem = self.read_element()?; let elem = self.read_element()?;
{ self.dispatcher.dispatch(ReceiveElement(elem));
let mut disp = self.dispatcher.lock().unwrap(); self.dispatcher.flush_all();
disp.dispatch(ReceiveElement(elem));
disp.flush_all();
}
} }
} }

View file

@ -4,6 +4,7 @@ use std::fmt::Debug;
use std::collections::BTreeMap; use std::collections::BTreeMap;
use std::cmp::Ordering; use std::cmp::Ordering;
use std::mem; use std::mem;
use std::sync::Mutex;
use minidom::Element; use minidom::Element;
@ -72,21 +73,21 @@ pub enum Propagation {
/// An event dispatcher, this takes care of dispatching events to their respective handlers. /// An event dispatcher, this takes care of dispatching events to their respective handlers.
pub struct Dispatcher { pub struct Dispatcher {
handlers: BTreeMap<TypeId, Vec<Record<Priority, Box<EventHandler>>>>, handlers: Mutex<BTreeMap<TypeId, Vec<Record<Priority, Box<EventHandler>>>>>,
queue: Vec<(TypeId, AbstractEvent)>, queue: Mutex<Vec<(TypeId, AbstractEvent)>>,
} }
impl Dispatcher { impl Dispatcher {
/// Create a new `Dispatcher`. /// Create a new `Dispatcher`.
pub fn new() -> Dispatcher { pub fn new() -> Dispatcher {
Dispatcher { Dispatcher {
handlers: BTreeMap::new(), handlers: Mutex::new(BTreeMap::new()),
queue: Vec::new(), queue: Mutex::new(Vec::new()),
} }
} }
/// Register an event handler. /// Register an event handler.
pub fn register<E, F>(&mut self, priority: Priority, func: F) pub fn register<E, F>(&self, priority: Priority, func: F)
where where
E: Event, E: Event,
F: Fn(&E) -> Propagation + 'static { F: Fn(&E) -> Propagation + 'static {
@ -110,24 +111,28 @@ impl Dispatcher {
func: func, func: func,
_marker: PhantomData, _marker: PhantomData,
}) as Box<EventHandler>; }) as Box<EventHandler>;
let ent = self.handlers.entry(TypeId::of::<E>()) let mut guard = self.handlers.lock().unwrap();
.or_insert_with(|| Vec::new()); let ent = guard.entry(TypeId::of::<E>())
.or_insert_with(|| Vec::new());
ent.push(Record(priority, handler)); ent.push(Record(priority, handler));
ent.sort(); ent.sort();
} }
/// Append an event to the queue. /// Append an event to the queue.
pub fn dispatch<E>(&mut self, event: E) where E: Event { pub fn dispatch<E>(&self, event: E) where E: Event {
self.queue.push((TypeId::of::<E>(), AbstractEvent::new(event))); self.queue.lock().unwrap().push((TypeId::of::<E>(), AbstractEvent::new(event)));
} }
/// Flush all events in the queue so they can be handled by their respective handlers. /// Flush all events in the queue so they can be handled by their respective handlers.
/// Returns whether there are still pending events. /// Returns whether there are still pending events.
pub fn flush(&mut self) -> bool { pub fn flush(&self) -> bool {
let mut q = Vec::new(); let mut q = Vec::new();
mem::swap(&mut self.queue, &mut q); {
let mut my_q = self.queue.lock().unwrap();
mem::swap(my_q.as_mut(), &mut q);
}
'evts: for (t, evt) in q { 'evts: for (t, evt) in q {
if let Some(handlers) = self.handlers.get_mut(&t) { if let Some(handlers) = self.handlers.lock().unwrap().get_mut(&t) {
for &mut Record(_, ref mut handler) in handlers { for &mut Record(_, ref mut handler) in handlers {
match handler.handle(&evt) { match handler.handle(&evt) {
Propagation::Stop => { continue 'evts; }, Propagation::Stop => { continue 'evts; },
@ -136,12 +141,12 @@ impl Dispatcher {
} }
} }
} }
!self.queue.is_empty() !self.queue.lock().unwrap().is_empty()
} }
/// Flushes all events, like `flush`, but keeps doing this until there is nothing left in the /// Flushes all events, like `flush`, but keeps doing this until there is nothing left in the
/// queue. /// queue.
pub fn flush_all(&mut self) { pub fn flush_all(&self) {
while self.flush() {} while self.flush() {}
} }
} }
@ -176,7 +181,7 @@ mod tests {
#[test] #[test]
#[should_panic(expected = "success")] #[should_panic(expected = "success")]
fn test() { fn test() {
let mut disp = Dispatcher::new(); let disp = Dispatcher::new();
#[derive(Debug)] #[derive(Debug)]
struct MyEvent { struct MyEvent {

View file

@ -4,7 +4,7 @@ use event::{Event, Dispatcher, SendElement, Priority, Propagation};
use std::any::Any; use std::any::Any;
use std::sync::{Arc, Mutex}; use std::sync::Arc;
use std::mem; use std::mem;
@ -12,11 +12,11 @@ use minidom::Element;
#[derive(Clone)] #[derive(Clone)]
pub struct PluginProxyBinding { pub struct PluginProxyBinding {
dispatcher: Arc<Mutex<Dispatcher>>, dispatcher: Arc<Dispatcher>,
} }
impl PluginProxyBinding { impl PluginProxyBinding {
pub fn new(dispatcher: Arc<Mutex<Dispatcher>>) -> PluginProxyBinding { pub fn new(dispatcher: Arc<Dispatcher>) -> PluginProxyBinding {
PluginProxyBinding { PluginProxyBinding {
dispatcher: dispatcher, dispatcher: dispatcher,
} }
@ -57,7 +57,7 @@ impl PluginProxy {
pub fn dispatch<E: Event>(&self, event: E) { pub fn dispatch<E: Event>(&self, event: E) {
self.with_binding(move |binding| { self.with_binding(move |binding| {
// TODO: proper error handling // TODO: proper error handling
binding.dispatcher.lock().unwrap().dispatch(event); binding.dispatcher.dispatch(event);
}); });
} }
@ -68,7 +68,7 @@ impl PluginProxy {
F: Fn(&E) -> Propagation + 'static { F: Fn(&E) -> Propagation + 'static {
self.with_binding(move |binding| { self.with_binding(move |binding| {
// TODO: proper error handling // TODO: proper error handling
binding.dispatcher.lock().unwrap().register(priority, func); binding.dispatcher.register(priority, func);
}); });
} }
@ -90,7 +90,7 @@ pub trait Plugin: Any + PluginAny {
} }
pub trait PluginInit { pub trait PluginInit {
fn init(dispatcher: &mut Dispatcher, me: Arc<Box<Plugin>>); fn init(dispatcher: &Dispatcher, me: Arc<Box<Plugin>>);
} }
pub trait PluginAny { pub trait PluginAny {

View file

@ -9,7 +9,7 @@ macro_rules! impl_plugin {
#[allow(unused_variables)] #[allow(unused_variables)]
impl $crate::plugin::PluginInit for $plugin { impl $crate::plugin::PluginInit for $plugin {
fn init( dispatcher: &mut $crate::event::Dispatcher fn init( dispatcher: &$crate::event::Dispatcher
, me: ::std::sync::Arc<Box<$crate::plugin::Plugin>>) { , me: ::std::sync::Arc<Box<$crate::plugin::Plugin>>) {
$( $(
let new_arc = me.clone(); let new_arc = me.clone();