HUGE REFACTOR to suppor multiple states. Incomplete, broken

This commit is contained in:
Igor Katson 2023-11-24 09:30:21 +00:00
parent cc1ef9d0e4
commit 739666ff88
No known key found for this signature in database
GPG key ID: B4EC22B66D61A3F5
7 changed files with 242 additions and 244 deletions

View file

@ -63,6 +63,7 @@ use futures::{stream::FuturesUnordered, StreamExt};
use librqbit_core::{
id20::Id20,
lengths::{ChunkInfo, Lengths, ValidPieceIndex},
speed_estimator::{self, SpeedEstimator},
torrent_metainfo::TorrentMetaV1Info,
};
use parking_lot::{Mutex, RwLock, RwLockReadGuard, RwLockWriteGuard};
@ -144,6 +145,8 @@ pub struct TorrentStateLive {
peer_queue_tx: UnboundedSender<SocketAddr>,
finished_notify: Notify,
speed_estimator: SpeedEstimator,
}
impl TorrentStateLive {
@ -163,6 +166,9 @@ impl TorrentStateLive {
) -> Arc<Self> {
let options = options.unwrap_or_default();
let (peer_queue_tx, peer_queue_rx) = unbounded_channel();
let speed_estimator = SpeedEstimator::new(5);
let state = Arc::new(TorrentStateLive {
info_hash,
info,
@ -186,6 +192,7 @@ impl TorrentStateLive {
peer_semaphore: Semaphore::new(128),
peer_queue_tx,
finished_notify: Notify::new(),
speed_estimator,
});
spawn(
span!(Level::ERROR, "peer_adder"),
@ -194,6 +201,10 @@ impl TorrentStateLive {
state
}
pub fn speed_estimator(&self) -> &SpeedEstimator {
&self.speed_estimator
}
async fn task_manage_peer(
self: Arc<Self>,
addr: SocketAddr,