Move have_notify into the peer
This commit is contained in:
parent
8c67127a85
commit
99dfb14895
2 changed files with 9 additions and 16 deletions
|
|
@ -227,9 +227,6 @@ pub struct LivePeerState {
|
||||||
// This is used to limit the number of chunk requests we send to a peer at a time.
|
// This is used to limit the number of chunk requests we send to a peer at a time.
|
||||||
pub requests_sem: Arc<Semaphore>,
|
pub requests_sem: Arc<Semaphore>,
|
||||||
|
|
||||||
// This is used to unpause processes after we were choked.
|
|
||||||
pub have_notify: Arc<Notify>,
|
|
||||||
|
|
||||||
// This is used to track the pieces the peer has.
|
// This is used to track the pieces the peer has.
|
||||||
pub bitfield: BF,
|
pub bitfield: BF,
|
||||||
|
|
||||||
|
|
@ -252,7 +249,6 @@ impl LivePeerState {
|
||||||
peer_interested: false,
|
peer_interested: false,
|
||||||
bitfield: BF::new(),
|
bitfield: BF::new(),
|
||||||
previously_requested_pieces: BF::new(),
|
previously_requested_pieces: BF::new(),
|
||||||
have_notify: Arc::new(Notify::new()),
|
|
||||||
requests_sem: Arc::new(Semaphore::new(0)),
|
requests_sem: Arc::new(Semaphore::new(0)),
|
||||||
inflight_requests: Default::default(),
|
inflight_requests: Default::default(),
|
||||||
tx,
|
tx,
|
||||||
|
|
|
||||||
|
|
@ -478,7 +478,8 @@ impl TorrentState {
|
||||||
|
|
||||||
let handler = PeerHandler {
|
let handler = PeerHandler {
|
||||||
addr,
|
addr,
|
||||||
on_bitfield_notify: Arc::new(Default::default()),
|
on_bitfield_notify: Default::default(),
|
||||||
|
have_notify: Default::default(),
|
||||||
state: state.clone(),
|
state: state.clone(),
|
||||||
spawner,
|
spawner,
|
||||||
};
|
};
|
||||||
|
|
@ -902,10 +903,13 @@ impl TorrentState {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Clone)]
|
|
||||||
struct PeerHandler {
|
struct PeerHandler {
|
||||||
state: Arc<TorrentState>,
|
state: Arc<TorrentState>,
|
||||||
on_bitfield_notify: Arc<Notify>,
|
// This is used to unpause chunk requester once the bitfield
|
||||||
|
// is received.
|
||||||
|
on_bitfield_notify: Notify,
|
||||||
|
// This is used to unpause after we were choked.
|
||||||
|
have_notify: Notify,
|
||||||
addr: SocketAddr,
|
addr: SocketAddr,
|
||||||
spawner: BlockingSpawner,
|
spawner: BlockingSpawner,
|
||||||
}
|
}
|
||||||
|
|
@ -1069,14 +1073,7 @@ impl PeerHandler {
|
||||||
WriterRequest::Message(MessageOwned::Interested),
|
WriterRequest::Message(MessageOwned::Interested),
|
||||||
])?;
|
])?;
|
||||||
|
|
||||||
let notify = match self
|
let notify = &self.have_notify;
|
||||||
.state
|
|
||||||
.peers
|
|
||||||
.with_live(handle, |l| l.have_notify.clone())
|
|
||||||
{
|
|
||||||
Some(notify) => notify,
|
|
||||||
None => return Ok(()),
|
|
||||||
};
|
|
||||||
#[allow(unused_must_use)]
|
#[allow(unused_must_use)]
|
||||||
{
|
{
|
||||||
timeout(Duration::from_secs(60), notify.notified()).await;
|
timeout(Duration::from_secs(60), notify.notified()).await;
|
||||||
|
|
@ -1224,11 +1221,11 @@ impl PeerHandler {
|
||||||
|
|
||||||
fn on_i_am_unchoked(&self, handle: PeerHandle) {
|
fn on_i_am_unchoked(&self, handle: PeerHandle) {
|
||||||
debug!("we are unchoked");
|
debug!("we are unchoked");
|
||||||
|
self.have_notify.notify_waiters();
|
||||||
self.state
|
self.state
|
||||||
.peers
|
.peers
|
||||||
.with_live_mut(handle, "on_i_am_unchoked", |live| {
|
.with_live_mut(handle, "on_i_am_unchoked", |live| {
|
||||||
live.i_am_choked = false;
|
live.i_am_choked = false;
|
||||||
live.have_notify.notify_waiters();
|
|
||||||
live.requests_sem.add_permits(16);
|
live.requests_sem.add_permits(16);
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue