Skip to content

Commit

Permalink
adding branching wakers + fixing the arc construction
Browse files Browse the repository at this point in the history
  • Loading branch information
jakubDoka committed Dec 22, 2023
1 parent d6b026c commit 59dfdf8
Show file tree
Hide file tree
Showing 2 changed files with 250 additions and 208 deletions.
32 changes: 15 additions & 17 deletions muxers/yamux/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -129,25 +129,23 @@ where
) -> Poll<Result<StreamMuxerEvent, Self::Error>> {
let this = self.get_mut();

let inbound_stream = ready!(this.poll_inner(cx))?;

if this.inbound_stream_buffer.len() >= MAX_BUFFERED_INBOUND_STREAMS {
tracing::warn!(
stream=%inbound_stream.0,
"dropping stream because buffer is full"
);
drop(inbound_stream);
} else {
this.inbound_stream_buffer.push_back(inbound_stream);

if let Some(waker) = this.inbound_stream_waker.take() {
waker.wake()
loop {
let inbound_stream = ready!(this.poll_inner(cx))?;

if this.inbound_stream_buffer.len() >= MAX_BUFFERED_INBOUND_STREAMS {
tracing::warn!(
stream=%inbound_stream.0,
"dropping stream because buffer is full"
);
drop(inbound_stream);
} else {
this.inbound_stream_buffer.push_back(inbound_stream);

if let Some(waker) = this.inbound_stream_waker.take() {
waker.wake()
}
}
}

// Schedule an immediate wake-up, allowing other code to run.
// cx.waker().wake_by_ref();
Poll::Pending
}
}

Expand Down
Loading

0 comments on commit 59dfdf8

Please sign in to comment.