This commit is contained in:
Igor Katson 2025-01-30 11:24:43 +00:00
parent 62d5288398
commit 3e6c2eae79
No known key found for this signature in database
GPG key ID: B4EC22B66D61A3F5
5 changed files with 36 additions and 24 deletions

View file

@ -26,7 +26,7 @@ async fn h_api_root(parts: Parts) -> impl IntoResponse {
.headers .headers
.get("Accept") .get("Accept")
.and_then(|h| h.to_str().ok()) .and_then(|h| h.to_str().ok())
.map_or(false, |h| h.contains("text/html")) .is_some_and(|h| h.contains("text/html"))
{ {
return Redirect::temporary("./web/").into_response(); return Redirect::temporary("./web/").into_response();
} }

View file

@ -88,6 +88,16 @@ where
} }
} }
struct ManagePeerArgs<R, W> {
handshake_supports_extended: bool,
read_buf: ReadBuf,
write_buf: Vec<u8>,
read: R,
write: W,
outgoing_chan: tokio::sync::mpsc::UnboundedReceiver<WriterRequest>,
have_broadcast: tokio::sync::broadcast::Receiver<ValidPieceIndex>,
}
impl<H: PeerConnectionHandler> PeerConnection<H> { impl<H: PeerConnectionHandler> PeerConnection<H> {
pub fn new( pub fn new(
addr: SocketAddr, addr: SocketAddr,
@ -147,21 +157,21 @@ impl<H: PeerConnectionHandler> PeerConnection<H> {
.context("error writing handshake")?; .context("error writing handshake")?;
write_buf.clear(); write_buf.clear();
let h_supports_extended = handshake.supports_extended(); let handshake_supports_extended = handshake.supports_extended();
self.handler.on_handshake(handshake)?; self.handler.on_handshake(handshake)?;
let (read, write) = conn.into_split(); let (read, write) = conn.into_split();
self.manage_peer( self.manage_peer(ManagePeerArgs {
h_supports_extended, handshake_supports_extended,
read_buf, read_buf,
write_buf, write_buf,
read, read,
write, write,
outgoing_chan, outgoing_chan,
have_broadcast, have_broadcast,
) })
.await .await
} }
@ -201,7 +211,7 @@ impl<H: PeerConnectionHandler> PeerConnection<H> {
.read_handshake(&mut read, rwtimeout) .read_handshake(&mut read, rwtimeout)
.await .await
.context("error reading handshake")?; .context("error reading handshake")?;
let h_supports_extended = h.supports_extended(); let handshake_supports_extended = h.supports_extended();
trace!( trace!(
peer_id=?Id20::new(h.peer_id), peer_id=?Id20::new(h.peer_id),
decoded_id=?try_decode_peer_id(Id20::new(h.peer_id)), decoded_id=?try_decode_peer_id(Id20::new(h.peer_id)),
@ -217,28 +227,35 @@ impl<H: PeerConnectionHandler> PeerConnection<H> {
self.handler.on_handshake(h)?; self.handler.on_handshake(h)?;
self.manage_peer( self.manage_peer(ManagePeerArgs {
h_supports_extended, handshake_supports_extended,
read_buf, read_buf,
write_buf, write_buf,
read, read,
write, write,
outgoing_chan, outgoing_chan,
have_broadcast, have_broadcast,
) })
.await .await
} }
async fn manage_peer( async fn manage_peer(
&self, &self,
handshake_supports_extended: bool, args: ManagePeerArgs<
mut read_buf: ReadBuf, impl tokio::io::AsyncRead + Send + Unpin,
mut write_buf: Vec<u8>, impl tokio::io::AsyncWrite + Send + Unpin,
mut read: impl tokio::io::AsyncRead + Unpin + Send, >,
mut write: impl tokio::io::AsyncWrite + Unpin + Send,
mut outgoing_chan: tokio::sync::mpsc::UnboundedReceiver<WriterRequest>,
mut have_broadcast: tokio::sync::broadcast::Receiver<ValidPieceIndex>,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
let ManagePeerArgs {
handshake_supports_extended,
mut read_buf,
mut write_buf,
mut read,
mut write,
mut outgoing_chan,
mut have_broadcast,
} = args;
use tokio::io::AsyncWriteExt; use tokio::io::AsyncWriteExt;
let rwtimeout = self let rwtimeout = self

View file

@ -996,7 +996,7 @@ impl Session {
name, name,
} = add_res; } = add_res;
let private = metadata.as_ref().map_or(false, |m| m.info.private); let private = metadata.as_ref().is_some_and(|m| m.info.private);
let make_peer_rx = || { let make_peer_rx = || {
self.make_peer_rx( self.make_peer_rx(

View file

@ -1253,10 +1253,7 @@ impl PeerHandler {
/// ///
/// If this returns, an existing in-flight piece was marked to be ours. /// If this returns, an existing in-flight piece was marked to be ours.
fn try_steal_old_slow_piece(&self, threshold: f64) -> Option<ValidPieceIndex> { fn try_steal_old_slow_piece(&self, threshold: f64) -> Option<ValidPieceIndex> {
let my_avg_time = match self.counters.average_piece_download_time() { let my_avg_time = self.counters.average_piece_download_time()?;
Some(t) => t,
None => return None,
};
let (stolen_idx, from_peer) = { let (stolen_idx, from_peer) = {
let mut g = self.state.lock_write("try_steal_old_slow_piece"); let mut g = self.state.lock_write("try_steal_old_slow_piece");

View file

@ -263,8 +263,6 @@ impl LivePeerState {
} }
pub fn has_full_torrent(&self, total_pieces: usize) -> bool { pub fn has_full_torrent(&self, total_pieces: usize) -> bool {
self.bitfield self.bitfield.get(0..total_pieces).is_some_and(|s| s.all())
.get(0..total_pieces)
.map_or(false, |s| s.all())
} }
} }