mirror of
https://github.com/n0-computer/noq.git
synced 2026-09-24 20:25:00 +00:00
Move idle Notify handle into endpoint Shared
This commit is contained in:
committed by
Dirkjan Ochtman
parent
7d51d2e57e
commit
00fb34759d
+6
-14
@@ -266,21 +266,13 @@ impl Endpoint {
|
||||
/// [`close()`]: Endpoint::close
|
||||
pub async fn wait_idle(&self) {
|
||||
loop {
|
||||
let idle;
|
||||
{
|
||||
let endpoint = &mut *self.inner.state.lock().unwrap();
|
||||
if endpoint.connections.is_empty() {
|
||||
break;
|
||||
}
|
||||
// Clone the `Arc<Notify>` so we can wait on the underlying `Notify` without holding
|
||||
// the lock. Store it in the outer scope to ensure it outlives the lock guard.
|
||||
idle = endpoint.idle.clone();
|
||||
// Construct the future while the lock is held to ensure we can't miss a wakeup if
|
||||
// the `Notify` is signaled immediately after we release the lock. `await` it after
|
||||
// the lock guard is out of scope.
|
||||
idle.notified()
|
||||
}
|
||||
.await;
|
||||
self.inner.shared.idle.notified().await;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -312,7 +304,7 @@ impl Future for EndpointDriver {
|
||||
let now = Instant::now();
|
||||
let mut keep_going = false;
|
||||
keep_going |= endpoint.drive_recv(cx, now)?;
|
||||
keep_going |= endpoint.handle_events(cx);
|
||||
keep_going |= endpoint.handle_events(cx, &self.0.shared);
|
||||
keep_going |= endpoint.drive_send(cx)?;
|
||||
|
||||
if !endpoint.incoming.is_empty() {
|
||||
@@ -368,13 +360,13 @@ pub(crate) struct State {
|
||||
recv_limiter: WorkLimiter,
|
||||
recv_buf: Box<[u8]>,
|
||||
send_limiter: WorkLimiter,
|
||||
idle: Arc<Notify>,
|
||||
runtime: Arc<dyn Runtime>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct Shared {
|
||||
incoming: Notify,
|
||||
idle: Notify,
|
||||
}
|
||||
|
||||
impl State {
|
||||
@@ -491,7 +483,7 @@ impl State {
|
||||
result
|
||||
}
|
||||
|
||||
fn handle_events(&mut self, cx: &mut Context) -> bool {
|
||||
fn handle_events(&mut self, cx: &mut Context, shared: &Shared) -> bool {
|
||||
use EndpointEvent::*;
|
||||
|
||||
for _ in 0..IO_LOOP_BOUND {
|
||||
@@ -501,7 +493,7 @@ impl State {
|
||||
if e.is_drained() {
|
||||
self.connections.senders.remove(&ch);
|
||||
if self.connections.is_empty() {
|
||||
self.idle.notify_waiters();
|
||||
shared.idle.notify_waiters();
|
||||
}
|
||||
}
|
||||
if let Some(event) = self.inner.handle_event(ch, e) {
|
||||
@@ -626,6 +618,7 @@ impl EndpointRef {
|
||||
Self(Arc::new(EndpointInner {
|
||||
shared: Shared {
|
||||
incoming: Notify::new(),
|
||||
idle: Notify::new(),
|
||||
},
|
||||
state: Mutex::new(State {
|
||||
socket,
|
||||
@@ -646,7 +639,6 @@ impl EndpointRef {
|
||||
recv_buf: recv_buf.into(),
|
||||
recv_limiter: WorkLimiter::new(RECV_TIME_BOUND),
|
||||
send_limiter: WorkLimiter::new(SEND_TIME_BOUND),
|
||||
idle: Arc::new(Notify::new()),
|
||||
runtime,
|
||||
}),
|
||||
}))
|
||||
|
||||
Reference in New Issue
Block a user