mirror of
https://github.com/YGGverse/aquatic.git
synced 2026-04-01 18:25:30 +00:00
Merge pull request #56 from greatest-ape/2022-03-23
Small ws and http improvements
This commit is contained in:
commit
bdc719b755
5 changed files with 175 additions and 179 deletions
|
|
@ -1,19 +1,10 @@
|
||||||
use std::net::{Ipv4Addr, Ipv6Addr};
|
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::time::Instant;
|
|
||||||
|
|
||||||
use aquatic_common::access_list::{create_access_list_cache, AccessListArcSwap, AccessListCache};
|
use aquatic_common::access_list::AccessListArcSwap;
|
||||||
use aquatic_common::{AmortizedIndexMap, CanonicalSocketAddr};
|
use aquatic_common::CanonicalSocketAddr;
|
||||||
use either::Either;
|
|
||||||
use smartstring::{LazyCompact, SmartString};
|
|
||||||
|
|
||||||
pub use aquatic_common::ValidUntil;
|
pub use aquatic_common::ValidUntil;
|
||||||
|
|
||||||
use aquatic_http_protocol::common::*;
|
|
||||||
use aquatic_http_protocol::response::ResponsePeer;
|
|
||||||
|
|
||||||
use crate::config::Config;
|
|
||||||
|
|
||||||
use aquatic_http_protocol::{
|
use aquatic_http_protocol::{
|
||||||
request::{AnnounceRequest, ScrapeRequest},
|
request::{AnnounceRequest, ScrapeRequest},
|
||||||
response::{AnnounceResponse, ScrapeResponse},
|
response::{AnnounceResponse, ScrapeResponse},
|
||||||
|
|
@ -72,152 +63,6 @@ impl ChannelResponse {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub trait Ip: ::std::fmt::Debug + Copy + Eq + ::std::hash::Hash {}
|
|
||||||
|
|
||||||
impl Ip for Ipv4Addr {}
|
|
||||||
impl Ip for Ipv6Addr {}
|
|
||||||
|
|
||||||
#[derive(Clone, Copy, Debug)]
|
|
||||||
pub struct ConnectionMeta {
|
|
||||||
/// Index of socket worker responsible for this connection. Required for
|
|
||||||
/// sending back response through correct channel to correct worker.
|
|
||||||
pub response_consumer_id: ConsumerId,
|
|
||||||
pub peer_addr: CanonicalSocketAddr,
|
|
||||||
/// Connection id local to socket worker
|
|
||||||
pub connection_id: ConnectionId,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Clone, Copy, Debug)]
|
|
||||||
pub struct PeerConnectionMeta<I: Ip> {
|
|
||||||
pub response_consumer_id: ConsumerId,
|
|
||||||
pub connection_id: ConnectionId,
|
|
||||||
pub peer_ip_address: I,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(PartialEq, Eq, Clone, Copy, Debug)]
|
|
||||||
pub enum PeerStatus {
|
|
||||||
Seeding,
|
|
||||||
Leeching,
|
|
||||||
Stopped,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl PeerStatus {
|
|
||||||
/// Determine peer status from announce event and number of bytes left.
|
|
||||||
///
|
|
||||||
/// Likely, the last branch will be taken most of the time.
|
|
||||||
#[inline]
|
|
||||||
pub fn from_event_and_bytes_left(event: AnnounceEvent, opt_bytes_left: Option<usize>) -> Self {
|
|
||||||
if let AnnounceEvent::Stopped = event {
|
|
||||||
Self::Stopped
|
|
||||||
} else if let Some(0) = opt_bytes_left {
|
|
||||||
Self::Seeding
|
|
||||||
} else {
|
|
||||||
Self::Leeching
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Debug, Clone, Copy)]
|
|
||||||
pub struct Peer<I: Ip> {
|
|
||||||
pub connection_meta: PeerConnectionMeta<I>,
|
|
||||||
pub port: u16,
|
|
||||||
pub status: PeerStatus,
|
|
||||||
pub valid_until: ValidUntil,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl<I: Ip> Peer<I> {
|
|
||||||
pub fn to_response_peer(&self) -> ResponsePeer<I> {
|
|
||||||
ResponsePeer {
|
|
||||||
ip_address: self.connection_meta.peer_ip_address,
|
|
||||||
port: self.port,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
|
|
||||||
pub struct PeerMapKey<I: Ip> {
|
|
||||||
pub peer_id: PeerId,
|
|
||||||
pub ip_or_key: Either<I, SmartString<LazyCompact>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
pub type PeerMap<I> = AmortizedIndexMap<PeerMapKey<I>, Peer<I>>;
|
|
||||||
|
|
||||||
pub struct TorrentData<I: Ip> {
|
|
||||||
pub peers: PeerMap<I>,
|
|
||||||
pub num_seeders: usize,
|
|
||||||
pub num_leechers: usize,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl<I: Ip> Default for TorrentData<I> {
|
|
||||||
#[inline]
|
|
||||||
fn default() -> Self {
|
|
||||||
Self {
|
|
||||||
peers: Default::default(),
|
|
||||||
num_seeders: 0,
|
|
||||||
num_leechers: 0,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
pub type TorrentMap<I> = AmortizedIndexMap<InfoHash, TorrentData<I>>;
|
|
||||||
|
|
||||||
#[derive(Default)]
|
|
||||||
pub struct TorrentMaps {
|
|
||||||
pub ipv4: TorrentMap<Ipv4Addr>,
|
|
||||||
pub ipv6: TorrentMap<Ipv6Addr>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl TorrentMaps {
|
|
||||||
pub fn clean(&mut self, config: &Config, access_list: &Arc<AccessListArcSwap>) {
|
|
||||||
let mut access_list_cache = create_access_list_cache(access_list);
|
|
||||||
|
|
||||||
Self::clean_torrent_map(config, &mut access_list_cache, &mut self.ipv4);
|
|
||||||
Self::clean_torrent_map(config, &mut access_list_cache, &mut self.ipv6);
|
|
||||||
}
|
|
||||||
|
|
||||||
fn clean_torrent_map<I: Ip>(
|
|
||||||
config: &Config,
|
|
||||||
access_list_cache: &mut AccessListCache,
|
|
||||||
torrent_map: &mut TorrentMap<I>,
|
|
||||||
) {
|
|
||||||
let now = Instant::now();
|
|
||||||
|
|
||||||
torrent_map.retain(|info_hash, torrent_data| {
|
|
||||||
if !access_list_cache
|
|
||||||
.load()
|
|
||||||
.allows(config.access_list.mode, &info_hash.0)
|
|
||||||
{
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
|
|
||||||
let num_seeders = &mut torrent_data.num_seeders;
|
|
||||||
let num_leechers = &mut torrent_data.num_leechers;
|
|
||||||
|
|
||||||
torrent_data.peers.retain(|_, peer| {
|
|
||||||
let keep = peer.valid_until.0 >= now;
|
|
||||||
|
|
||||||
if !keep {
|
|
||||||
match peer.status {
|
|
||||||
PeerStatus::Seeding => {
|
|
||||||
*num_seeders -= 1;
|
|
||||||
}
|
|
||||||
PeerStatus::Leeching => {
|
|
||||||
*num_leechers -= 1;
|
|
||||||
}
|
|
||||||
_ => (),
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
keep
|
|
||||||
});
|
|
||||||
|
|
||||||
!torrent_data.peers.is_empty()
|
|
||||||
});
|
|
||||||
|
|
||||||
torrent_map.shrink_to_fit();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Default, Clone)]
|
#[derive(Default, Clone)]
|
||||||
pub struct State {
|
pub struct State {
|
||||||
pub access_list: Arc<AccessListArcSwap>,
|
pub access_list: Arc<AccessListArcSwap>,
|
||||||
|
|
|
||||||
|
|
@ -18,7 +18,7 @@ mod common;
|
||||||
pub mod config;
|
pub mod config;
|
||||||
mod workers;
|
mod workers;
|
||||||
|
|
||||||
pub const APP_NAME: &str = "aquatic_http: HTTP/TLS BitTorrent tracker";
|
pub const APP_NAME: &str = "aquatic_http: BitTorrent tracker (HTTP over TLS)";
|
||||||
|
|
||||||
const SHARED_CHANNEL_SIZE: usize = 1024;
|
const SHARED_CHANNEL_SIZE: usize = 1024;
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,27 +1,179 @@
|
||||||
|
use std::cell::RefCell;
|
||||||
use std::collections::BTreeMap;
|
use std::collections::BTreeMap;
|
||||||
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr};
|
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr};
|
||||||
|
use std::rc::Rc;
|
||||||
|
use std::sync::Arc;
|
||||||
|
use std::time::Duration;
|
||||||
|
use std::time::Instant;
|
||||||
|
|
||||||
use either::Either;
|
use either::Either;
|
||||||
use rand::Rng;
|
|
||||||
|
|
||||||
use aquatic_common::extract_response_peers;
|
|
||||||
use aquatic_http_protocol::request::*;
|
|
||||||
use aquatic_http_protocol::response::*;
|
|
||||||
|
|
||||||
use std::cell::RefCell;
|
|
||||||
use std::rc::Rc;
|
|
||||||
use std::time::Duration;
|
|
||||||
|
|
||||||
use futures_lite::{Stream, StreamExt};
|
use futures_lite::{Stream, StreamExt};
|
||||||
use glommio::channels::channel_mesh::{MeshBuilder, Partial, Role, Senders};
|
use glommio::channels::channel_mesh::{MeshBuilder, Partial, Role, Senders};
|
||||||
use glommio::timer::TimerActionRepeat;
|
use glommio::timer::TimerActionRepeat;
|
||||||
use glommio::{enclose, prelude::*};
|
use glommio::{enclose, prelude::*};
|
||||||
use rand::prelude::SmallRng;
|
use rand::prelude::SmallRng;
|
||||||
|
use rand::Rng;
|
||||||
use rand::SeedableRng;
|
use rand::SeedableRng;
|
||||||
|
use smartstring::{LazyCompact, SmartString};
|
||||||
|
|
||||||
|
use aquatic_common::access_list::{create_access_list_cache, AccessListArcSwap, AccessListCache};
|
||||||
|
use aquatic_common::extract_response_peers;
|
||||||
|
use aquatic_common::ValidUntil;
|
||||||
|
use aquatic_common::{AmortizedIndexMap, CanonicalSocketAddr};
|
||||||
|
use aquatic_http_protocol::common::*;
|
||||||
|
use aquatic_http_protocol::request::*;
|
||||||
|
use aquatic_http_protocol::response::ResponsePeer;
|
||||||
|
use aquatic_http_protocol::response::*;
|
||||||
|
|
||||||
use crate::common::*;
|
use crate::common::*;
|
||||||
use crate::config::Config;
|
use crate::config::Config;
|
||||||
|
|
||||||
|
pub trait Ip: ::std::fmt::Debug + Copy + Eq + ::std::hash::Hash {}
|
||||||
|
|
||||||
|
impl Ip for Ipv4Addr {}
|
||||||
|
impl Ip for Ipv6Addr {}
|
||||||
|
|
||||||
|
#[derive(Clone, Copy, Debug)]
|
||||||
|
pub struct ConnectionMeta {
|
||||||
|
/// Index of socket worker responsible for this connection. Required for
|
||||||
|
/// sending back response through correct channel to correct worker.
|
||||||
|
pub response_consumer_id: ConsumerId,
|
||||||
|
pub peer_addr: CanonicalSocketAddr,
|
||||||
|
/// Connection id local to socket worker
|
||||||
|
pub connection_id: ConnectionId,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Copy, Debug)]
|
||||||
|
pub struct PeerConnectionMeta<I: Ip> {
|
||||||
|
pub response_consumer_id: ConsumerId,
|
||||||
|
pub connection_id: ConnectionId,
|
||||||
|
pub peer_ip_address: I,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(PartialEq, Eq, Clone, Copy, Debug)]
|
||||||
|
pub enum PeerStatus {
|
||||||
|
Seeding,
|
||||||
|
Leeching,
|
||||||
|
Stopped,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PeerStatus {
|
||||||
|
/// Determine peer status from announce event and number of bytes left.
|
||||||
|
///
|
||||||
|
/// Likely, the last branch will be taken most of the time.
|
||||||
|
#[inline]
|
||||||
|
pub fn from_event_and_bytes_left(event: AnnounceEvent, opt_bytes_left: Option<usize>) -> Self {
|
||||||
|
if let AnnounceEvent::Stopped = event {
|
||||||
|
Self::Stopped
|
||||||
|
} else if let Some(0) = opt_bytes_left {
|
||||||
|
Self::Seeding
|
||||||
|
} else {
|
||||||
|
Self::Leeching
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy)]
|
||||||
|
pub struct Peer<I: Ip> {
|
||||||
|
pub connection_meta: PeerConnectionMeta<I>,
|
||||||
|
pub port: u16,
|
||||||
|
pub status: PeerStatus,
|
||||||
|
pub valid_until: ValidUntil,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<I: Ip> Peer<I> {
|
||||||
|
pub fn to_response_peer(&self) -> ResponsePeer<I> {
|
||||||
|
ResponsePeer {
|
||||||
|
ip_address: self.connection_meta.peer_ip_address,
|
||||||
|
port: self.port,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
|
||||||
|
pub struct PeerMapKey<I: Ip> {
|
||||||
|
pub peer_id: PeerId,
|
||||||
|
pub ip_or_key: Either<I, SmartString<LazyCompact>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
pub type PeerMap<I> = AmortizedIndexMap<PeerMapKey<I>, Peer<I>>;
|
||||||
|
|
||||||
|
pub struct TorrentData<I: Ip> {
|
||||||
|
pub peers: PeerMap<I>,
|
||||||
|
pub num_seeders: usize,
|
||||||
|
pub num_leechers: usize,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<I: Ip> Default for TorrentData<I> {
|
||||||
|
#[inline]
|
||||||
|
fn default() -> Self {
|
||||||
|
Self {
|
||||||
|
peers: Default::default(),
|
||||||
|
num_seeders: 0,
|
||||||
|
num_leechers: 0,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub type TorrentMap<I> = AmortizedIndexMap<InfoHash, TorrentData<I>>;
|
||||||
|
|
||||||
|
#[derive(Default)]
|
||||||
|
pub struct TorrentMaps {
|
||||||
|
pub ipv4: TorrentMap<Ipv4Addr>,
|
||||||
|
pub ipv6: TorrentMap<Ipv6Addr>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl TorrentMaps {
|
||||||
|
pub fn clean(&mut self, config: &Config, access_list: &Arc<AccessListArcSwap>) {
|
||||||
|
let mut access_list_cache = create_access_list_cache(access_list);
|
||||||
|
|
||||||
|
Self::clean_torrent_map(config, &mut access_list_cache, &mut self.ipv4);
|
||||||
|
Self::clean_torrent_map(config, &mut access_list_cache, &mut self.ipv6);
|
||||||
|
}
|
||||||
|
|
||||||
|
fn clean_torrent_map<I: Ip>(
|
||||||
|
config: &Config,
|
||||||
|
access_list_cache: &mut AccessListCache,
|
||||||
|
torrent_map: &mut TorrentMap<I>,
|
||||||
|
) {
|
||||||
|
let now = Instant::now();
|
||||||
|
|
||||||
|
torrent_map.retain(|info_hash, torrent_data| {
|
||||||
|
if !access_list_cache
|
||||||
|
.load()
|
||||||
|
.allows(config.access_list.mode, &info_hash.0)
|
||||||
|
{
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
let num_seeders = &mut torrent_data.num_seeders;
|
||||||
|
let num_leechers = &mut torrent_data.num_leechers;
|
||||||
|
|
||||||
|
torrent_data.peers.retain(|_, peer| {
|
||||||
|
let keep = peer.valid_until.0 >= now;
|
||||||
|
|
||||||
|
if !keep {
|
||||||
|
match peer.status {
|
||||||
|
PeerStatus::Seeding => {
|
||||||
|
*num_seeders -= 1;
|
||||||
|
}
|
||||||
|
PeerStatus::Leeching => {
|
||||||
|
*num_leechers -= 1;
|
||||||
|
}
|
||||||
|
_ => (),
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
keep
|
||||||
|
});
|
||||||
|
|
||||||
|
!torrent_data.peers.is_empty()
|
||||||
|
});
|
||||||
|
|
||||||
|
torrent_map.shrink_to_fit();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub async fn run_request_worker(
|
pub async fn run_request_worker(
|
||||||
config: Config,
|
config: Config,
|
||||||
state: State,
|
state: State,
|
||||||
|
|
|
||||||
|
|
@ -73,7 +73,6 @@ pub async fn run_socket_worker(
|
||||||
|
|
||||||
let connection_slab = Rc::new(RefCell::new(Slab::new()));
|
let connection_slab = Rc::new(RefCell::new(Slab::new()));
|
||||||
|
|
||||||
// Periodically remove closed connections
|
|
||||||
TimerActionRepeat::repeat(enclose!((config, connection_slab) move || {
|
TimerActionRepeat::repeat(enclose!((config, connection_slab) move || {
|
||||||
clean_connections(
|
clean_connections(
|
||||||
config.clone(),
|
config.clone(),
|
||||||
|
|
@ -139,15 +138,15 @@ async fn clean_connections(
|
||||||
let now = Instant::now();
|
let now = Instant::now();
|
||||||
|
|
||||||
connection_slab.borrow_mut().retain(|_, reference| {
|
connection_slab.borrow_mut().retain(|_, reference| {
|
||||||
let keep = reference.valid_until.0 > now;
|
if reference.valid_until.0 > now {
|
||||||
|
true
|
||||||
if !keep {
|
} else {
|
||||||
if let Some(ref handle) = reference.task_handle {
|
if let Some(ref handle) = reference.task_handle {
|
||||||
handle.cancel();
|
handle.cancel();
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
keep
|
false
|
||||||
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
connection_slab.borrow_mut().shrink_to_fit();
|
connection_slab.borrow_mut().shrink_to_fit();
|
||||||
|
|
|
||||||
|
|
@ -156,15 +156,15 @@ async fn clean_connections(
|
||||||
let now = Instant::now();
|
let now = Instant::now();
|
||||||
|
|
||||||
connection_slab.borrow_mut().retain(|_, reference| {
|
connection_slab.borrow_mut().retain(|_, reference| {
|
||||||
let keep = reference.valid_until.0 > now;
|
if reference.valid_until.0 > now {
|
||||||
|
true
|
||||||
if !keep {
|
} else {
|
||||||
if let Some(ref handle) = reference.task_handle {
|
if let Some(ref handle) = reference.task_handle {
|
||||||
handle.cancel();
|
handle.cancel();
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
keep
|
false
|
||||||
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
connection_slab.borrow_mut().shrink_to_fit();
|
connection_slab.borrow_mut().shrink_to_fit();
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue