-
Notifications
You must be signed in to change notification settings - Fork 1.2k
Fixes to the mplex implementation #360
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from 5 commits
Commits
Show all changes
14 commits
Select commit
Hold shift + click to select a range
c3b7103
Fixes to the mplex implementation
tomaka 3e99bef
Merge branch 'master' into mplex-fixes
tomaka 7cc95d4
Fix mem leak and wrong logging message
tomaka 236e4ea
Fix potential leak
tomaka e61a5b9
Immediately return an error if an error happened
tomaka d4c3951
Correctly handle Close and Reset
tomaka 68529ed
Check the even-ness of the substream id
tomaka aae820e
Fix concerns
tomaka 05d883b
Merge remote-tracking branch 'upstream/master' into mplex-fixes
tomaka 06ebefd
Fix tasks notification system
tomaka 4d0c55f
Merge branch 'master' into mplex-fixes
tomaka 6867565
Merge branch 'master' into mplex-fixes
tomaka e82aa7f
Revert "Fix tasks notification system"
tomaka d95a2c0
Merge remote-tracking branch 'upstream/master' into mplex-fixes
tomaka File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -23,6 +23,8 @@ extern crate fnv; | |
| #[macro_use] | ||
| extern crate futures; | ||
| extern crate libp2p_core as core; | ||
| #[macro_use] | ||
| extern crate log; | ||
| extern crate parking_lot; | ||
| extern crate tokio_codec; | ||
| extern crate tokio_io; | ||
|
|
@@ -32,12 +34,11 @@ mod codec; | |
|
|
||
| use std::{cmp, iter}; | ||
| use std::io::{Read, Write, Error as IoError, ErrorKind as IoErrorKind}; | ||
| use std::mem; | ||
| use std::sync::Arc; | ||
| use std::sync::{Arc, atomic::AtomicUsize, atomic::Ordering}; | ||
| use bytes::Bytes; | ||
| use core::{ConnectionUpgrade, Endpoint, StreamMuxer}; | ||
| use parking_lot::Mutex; | ||
| use fnv::FnvHashSet; | ||
| use fnv::{FnvHashMap, FnvHashSet}; | ||
| use futures::prelude::*; | ||
| use futures::{future, stream::Fuse, task}; | ||
| use tokio_codec::Framed; | ||
|
|
@@ -46,7 +47,7 @@ use tokio_io::{AsyncRead, AsyncWrite}; | |
| // Maximum number of simultaneously-open substreams. | ||
| const MAX_SUBSTREAMS: usize = 1024; | ||
| // Maximum number of elements in the internal buffer. | ||
| const MAX_BUFFER_LEN: usize = 256; | ||
| const MAX_BUFFER_LEN: usize = 1024; | ||
|
|
||
| /// Configuration for the multiplexer. | ||
| #[derive(Debug, Clone, Default)] | ||
|
|
@@ -74,11 +75,12 @@ where | |
| fn upgrade(self, i: C, _: (), endpoint: Endpoint, remote_addr: Maf) -> Self::Future { | ||
| let out = Multiplex { | ||
| inner: Arc::new(Mutex::new(MultiplexInner { | ||
| error: Ok(()), | ||
| inner: Framed::new(i, codec::Codec::new()).fuse(), | ||
| buffer: Vec::with_capacity(32), | ||
| opened_substreams: Default::default(), | ||
| next_outbound_stream_id: if endpoint == Endpoint::Dialer { 0 } else { 1 }, | ||
| to_notify: Vec::new(), | ||
| to_notify: Default::default(), | ||
| })) | ||
| }; | ||
|
|
||
|
|
@@ -107,27 +109,35 @@ impl<C> Clone for Multiplex<C> { | |
|
|
||
| // Struct shared throughout the implementation. | ||
| struct MultiplexInner<C> { | ||
| // Errored that happend earlier. Should poison any attempt to use this `MultiplexError`. | ||
| error: Result<(), IoError>, | ||
| // Underlying stream. | ||
| inner: Fuse<Framed<C, codec::Codec>>, | ||
| // Buffer of elements pulled from the stream but not processed yet. | ||
| buffer: Vec<codec::Elem>, | ||
| // List of Ids of opened substreams. Used to filter out messages that don't belong to any | ||
| // substream. | ||
| // substream. Note that this is handled exclusively by `next_match`. | ||
| opened_substreams: FnvHashSet<u32>, | ||
| // Id of the next outgoing substream. Should always increase by two. | ||
| next_outbound_stream_id: u32, | ||
| // List of tasks to notify when a new element is inserted in `buffer`. | ||
| to_notify: Vec<task::Task>, | ||
| // List of tasks to notify when a new element is inserted in `buffer` or an error happens. | ||
| to_notify: FnvHashMap<usize, task::Task>, | ||
| } | ||
|
|
||
| // Processes elements in `inner` until one matching `filter` is found. | ||
| // | ||
| // If `NotReady` is returned, the current task is scheduled for later, just like with any `Poll`. | ||
| // `Ready(Some())` is almost always returned. `Ready(None)` is returned if the stream is EOF. | ||
| /// Processes elements in `inner` until one matching `filter` is found. | ||
| /// | ||
| /// If `NotReady` is returned, the current task is scheduled for later, just like with any `Poll`. | ||
| /// `Ready(Some())` is almost always returned. `Ready(None)` is returned if the stream is EOF. | ||
| fn next_match<C, F, O>(inner: &mut MultiplexInner<C>, mut filter: F) -> Poll<Option<O>, IoError> | ||
| where C: AsyncRead + AsyncWrite, | ||
| F: FnMut(&codec::Elem) -> Option<O>, | ||
| { | ||
| // If an error happened earlier, immediately return it. | ||
| match inner.error { | ||
| Ok(()) => (), | ||
| Err(ref err) => return Err(IoError::new(err.kind(), err.to_string())), | ||
| }; | ||
|
|
||
| if let Some((offset, out)) = inner.buffer.iter().enumerate().filter_map(|(n, v)| filter(v).map(|v| (n, v))).next() { | ||
| inner.buffer.remove(offset); | ||
| return Ok(Async::Ready(Some(out))); | ||
|
|
@@ -137,27 +147,56 @@ where C: AsyncRead + AsyncWrite, | |
| let elem = match inner.inner.poll() { | ||
| Ok(Async::Ready(item)) => item, | ||
| Ok(Async::NotReady) => { | ||
| inner.to_notify.push(task::current()); | ||
| static NEXT_TASK_ID: AtomicUsize = AtomicUsize::new(0); | ||
| task_local!{ | ||
| static TASK_ID: usize = NEXT_TASK_ID.fetch_add(1, Ordering::Relaxed) | ||
| } | ||
| inner.to_notify.insert(TASK_ID.with(|&t| t), task::current()); | ||
| return Ok(Async::NotReady); | ||
| }, | ||
| Err(err) => { | ||
| return Err(err); | ||
| let err2 = IoError::new(err.kind(), err.to_string()); | ||
| inner.error = Err(err); | ||
| for task in inner.to_notify.drain() { | ||
| task.1.notify(); | ||
| } | ||
| return Err(err2); | ||
| }, | ||
| }; | ||
|
|
||
| if let Some(elem) = elem { | ||
| trace!("Received message: {:?}", elem); | ||
|
|
||
| // Handle substreams opening/closing. | ||
| match elem { | ||
| codec::Elem::Open { substream_id } => { // TODO: check even/uneven? | ||
| inner.opened_substreams.insert(substream_id); | ||
| }, | ||
| codec::Elem::Close { substream_id, .. } | codec::Elem::Reset { substream_id, .. } => { | ||
| inner.opened_substreams.remove(&substream_id); | ||
| }, | ||
| _ => () | ||
| } | ||
|
|
||
| if let Some(out) = filter(&elem) { | ||
| return Ok(Async::Ready(Some(out))); | ||
| } else { | ||
| if inner.buffer.len() >= MAX_BUFFER_LEN { | ||
| return Err(IoError::new(IoErrorKind::InvalidData, "reached maximum buffer length")); | ||
| debug!("Reached mplex maximum buffer length"); | ||
| inner.error = Err(IoError::new(IoErrorKind::Other, "reached maximum buffer length")); | ||
| for task in inner.to_notify.drain() { | ||
| task.1.notify(); | ||
| } | ||
| return Err(IoError::new(IoErrorKind::Other, "reached maximum buffer length")); | ||
| } | ||
|
|
||
| if inner.opened_substreams.contains(&elem.substream_id()) || elem.is_open_msg() { | ||
| inner.buffer.push(elem); | ||
| for task in inner.to_notify.drain(..) { | ||
| task.notify(); | ||
| for task in inner.to_notify.drain() { | ||
| task.1.notify(); | ||
| } | ||
| } else { | ||
| debug!("Ignored message because the substream wasn't open"); | ||
| } | ||
| } | ||
| } else { | ||
|
|
@@ -166,13 +205,6 @@ where C: AsyncRead + AsyncWrite, | |
| } | ||
| } | ||
|
|
||
| // Closes a substream in `inner`. | ||
| fn clean_out_substream<C>(inner: &mut MultiplexInner<C>, num: u32) { | ||
| let was_in = inner.opened_substreams.remove(&num); | ||
| debug_assert!(was_in, "Dropped substream which wasn't open ; programmer error"); | ||
| inner.buffer.retain(|elem| elem.substream_id() != num); | ||
| } | ||
|
|
||
| // Small convenience function that tries to write `elem` to the stream. | ||
| fn poll_send<C>(inner: &mut MultiplexInner<C>, elem: codec::Elem) -> Poll<(), IoError> | ||
| where C: AsyncRead + AsyncWrite | ||
|
|
@@ -212,12 +244,17 @@ where C: AsyncRead + AsyncWrite + 'static // TODO: 'static :-/ | |
| }; | ||
|
|
||
| // We use an RAII guard, so that we close the substream in case of an error. | ||
| struct OpenedSubstreamGuard<C>(Arc<Mutex<MultiplexInner<C>>>, u32); | ||
| struct OpenedSubstreamGuard<C>(Option<Arc<Mutex<MultiplexInner<C>>>>, u32); | ||
| impl<C> Drop for OpenedSubstreamGuard<C> { | ||
| fn drop(&mut self) { clean_out_substream(&mut self.0.lock(), self.1); } | ||
| fn drop(&mut self) { | ||
| if let Some(inner) = self.0.take() { | ||
| debug!("Failed to open outbound substream {}", self.1); | ||
| inner.lock().buffer.retain(|elem| elem.substream_id() != self.1); | ||
| } | ||
| } | ||
| } | ||
| inner.opened_substreams.insert(substream_id); | ||
| let guard = OpenedSubstreamGuard(self.inner.clone(), substream_id); | ||
| let mut guard = OpenedSubstreamGuard(Some(self.inner.clone()), substream_id); | ||
|
|
||
| // We send `Open { substream_id }`, then flush, then only produce the substream. | ||
| let future = { | ||
|
|
@@ -232,17 +269,14 @@ where C: AsyncRead + AsyncWrite + 'static // TODO: 'static :-/ | |
| move |()| { | ||
| future::poll_fn(move || inner.lock().inner.poll_complete()) | ||
| } | ||
| }).map({ | ||
| let inner = self.inner.clone(); | ||
| move |()| { | ||
| mem::forget(guard); | ||
| Some(Substream { | ||
| inner: inner.clone(), | ||
| num: substream_id, | ||
| current_data: Bytes::new(), | ||
| endpoint: Endpoint::Dialer, | ||
| }) | ||
| } | ||
| }).map(move |()| { | ||
| debug!("Successfully opened outbound substream {}", substream_id); | ||
| Some(Substream { | ||
| inner: guard.0.take().unwrap(), | ||
| num: substream_id, | ||
| current_data: Bytes::new(), | ||
| endpoint: Endpoint::Dialer, | ||
| }) | ||
| }) | ||
| }; | ||
|
|
||
|
|
@@ -265,19 +299,20 @@ where C: AsyncRead + AsyncWrite | |
| let mut inner = self.inner.lock(); | ||
|
|
||
| if inner.opened_substreams.len() >= MAX_SUBSTREAMS { | ||
| debug!("Refused substream ; reached maximum number of substreams {}", MAX_SUBSTREAMS); | ||
| return Err(IoError::new(IoErrorKind::ConnectionRefused, | ||
| "exceeded maximum number of open substreams")); | ||
| } | ||
|
|
||
| let num = try_ready!(next_match(&mut inner, |elem| { | ||
| match elem { | ||
| codec::Elem::Open { substream_id } => Some(*substream_id), // TODO: check even/uneven? | ||
| codec::Elem::Open { substream_id } => Some(*substream_id), | ||
| _ => None, | ||
| } | ||
| })); | ||
|
|
||
| if let Some(num) = num { | ||
| inner.opened_substreams.insert(num); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. whty don't we insert this any longer?
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. It's dealt with in |
||
| debug!("Successfully opened inbound substream {}", num); | ||
| Ok(Async::Ready(Some(Substream { | ||
| inner: self.inner.clone(), | ||
| current_data: Bytes::new(), | ||
|
|
@@ -306,13 +341,14 @@ where C: AsyncRead + AsyncWrite | |
| { | ||
| fn read(&mut self, buf: &mut [u8]) -> Result<usize, IoError> { | ||
| loop { | ||
| // First transfer from `current_data`. | ||
| // First, transfer from `current_data`. | ||
| if self.current_data.len() != 0 { | ||
| let len = cmp::min(self.current_data.len(), buf.len()); | ||
| buf[..len].copy_from_slice(&self.current_data.split_to(len)); | ||
| return Ok(len); | ||
| } | ||
|
|
||
| // Try to find a packet of data in the buffer. | ||
| let mut inner = self.inner.lock(); | ||
| let next_data_poll = next_match(&mut inner, |elem| { | ||
| match elem { | ||
|
|
@@ -328,7 +364,15 @@ where C: AsyncRead + AsyncWrite | |
| match next_data_poll { | ||
| Ok(Async::Ready(Some(data))) => self.current_data = data.freeze(), | ||
| Ok(Async::Ready(None)) => return Ok(0), | ||
| Ok(Async::NotReady) => return Err(IoErrorKind::WouldBlock.into()), | ||
| Ok(Async::NotReady) => { | ||
| // There was no data packet in the buffer about this substream ; maybe it's | ||
| // because it has been closed. | ||
| if inner.opened_substreams.contains(&self.num) { | ||
| return Err(IoErrorKind::WouldBlock.into()); | ||
| } else { | ||
| return Ok(0); | ||
| } | ||
| }, | ||
| Err(err) => return Err(err), | ||
| } | ||
| } | ||
|
|
@@ -386,7 +430,7 @@ impl<C> Drop for Substream<C> | |
| where C: AsyncRead + AsyncWrite | ||
| { | ||
| fn drop(&mut self) { | ||
| let _ = self.shutdown(); | ||
| clean_out_substream(&mut self.inner.lock(), self.num); | ||
| let _ = self.shutdown(); // TODO: this doesn't necessarily send the close message | ||
| self.inner.lock().buffer.retain(|elem| elem.substream_id() != self.num); | ||
| } | ||
| } | ||
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
maybe make it more explicit rather than having an empty match?