use crate::actor::ActorHandle;
use crate::actor::server_actor::{
PeerAddress, ServerActor, ServerMessage, UserMessage,
};
use crate::download_store::{DownloadStore, collect_failed_tokens};
use crate::types::{
ClientVersion, DownloadMetadata, DownloadStatus, RoomEvent, RoomInfo,
SessionLoss, SessionWatch, UserInfo, UserPresence, UserStats, UserStatus,
};
use crate::utils::logger;
use crate::{
Transfer,
actor::{ActorSystem, peer_registry::PeerRegistry},
error::{Result, SoulseekRs},
message::peer::{FileEntry, SharedDirectory, build_file_search_response},
peer::{ConnectionType, DownloadPeer, Peer, PeerMessage, listen::Listen},
shares::Shares,
types::{Download, Search, SearchResult},
utils::{lock::RwLockExt, md5},
};
use std::{
collections::{HashMap, HashSet},
net::TcpStream,
sync::{
Arc, RwLock,
atomic::{AtomicBool, AtomicU32, Ordering},
mpsc::{self, Receiver, Sender},
},
thread::{self, sleep},
time::{Duration, Instant},
};
use upload_queue::QueuedUpload;
use crate::{debug, error, info, trace, warn};
const DEFAULT_LISTEN_PORT: u16 = 2234;
pub const DEFAULT_WISHLIST_INTERVAL: Duration = Duration::from_mins(12);
pub const DEFAULT_UPLOAD_SLOTS: usize = 10;
const BROKER_CONNECT_TIMEOUT: Duration = Duration::from_secs(20);
static NEXT_CONNECT_TOKEN: AtomicU32 = AtomicU32::new(1);
fn next_connect_token() -> u32 {
NEXT_CONNECT_TOKEN.fetch_add(1, Ordering::Relaxed).max(1)
}
static NEXT_UPLOAD_TOKEN: AtomicU32 = AtomicU32::new(0x8000_0000);
fn next_upload_token() -> u32 {
NEXT_UPLOAD_TOKEN.fetch_add(1, Ordering::Relaxed)
}
static NEXT_DOWNLOAD_TOKEN: AtomicU32 = AtomicU32::new(1);
fn next_download_token() -> u32 {
NEXT_DOWNLOAD_TOKEN.fetch_add(1, Ordering::Relaxed) % 0x8000_0000
}
struct UploadJob {
downloader: String,
real_path: std::path::PathBuf,
virtual_path: String,
size: u64,
}
struct ActiveUpload {
username: String,
filename: String,
size: u64,
bytes_sent: Arc<std::sync::atomic::AtomicU64>,
cancel: Arc<std::sync::atomic::AtomicBool>,
status: crate::types::UploadStatus,
started: Instant,
}
fn upload_speed(
status: &crate::types::UploadStatus,
bytes_sent: u64,
started: Instant,
) -> f64 {
if !matches!(status, crate::types::UploadStatus::InProgress) {
return 0.0;
}
let elapsed = started.elapsed().as_secs_f64();
if elapsed <= 0.0 {
return 0.0;
}
bytes_sent as f64 / elapsed
}
fn build_search_response(
shares: &Shares,
own_username: &str,
token: u32,
query: &str,
) -> Option<crate::message::Message> {
let matches = shares.search(query);
if matches.is_empty() {
return None;
}
let entries: Vec<FileEntry> = matches
.iter()
.map(|f| FileEntry {
name: &f.virtual_path,
size: f.size,
attribs: &f.attributes,
})
.collect();
Some(build_file_search_response(
own_username,
token,
&entries,
1,
0,
))
}
#[derive(Debug, Clone)]
pub struct ClientSettings {
pub username: String,
pub password: String,
pub server_address: PeerAddress,
pub enable_listen: bool,
pub listen_port: u16,
pub shared_directories: Vec<String>,
pub version: ClientVersion,
}
impl ClientSettings {
pub fn new(
username: impl Into<String>,
password: impl Into<String>,
) -> Self {
Self {
username: username.into(),
password: password.into(),
..Default::default()
}
}
}
impl Default for ClientSettings {
fn default() -> Self {
Self {
username: String::new(),
password: String::new(),
server_address: PeerAddress::new(
"server.slsknet.org".to_string(),
2416,
),
enable_listen: true,
listen_port: DEFAULT_LISTEN_PORT,
shared_directories: Vec::new(),
version: ClientVersion::default(),
}
}
}
#[derive(Debug)]
#[non_exhaustive]
pub enum ClientOperation {
ConnectToPeer(Peer),
SearchResult(SearchResult),
PeerDisconnected(u64, String, Option<SoulseekRs>),
DownloadFromPeer(u32, Peer, bool),
UpdateDownloadTokens(Transfer, String),
GetPeerAddressResponse {
username: String,
host: String,
port: u32,
obfuscation_type: u32,
obfuscated_port: u16,
},
UploadFailed(String, String),
PlaceInQueueUpdate {
username: String,
filename: String,
place: u32,
},
SetServerSender(Sender<ServerMessage>),
PrivateMessageReceived(UserMessage),
UserStatusReceived {
username: String,
status: u32,
privileged: bool,
},
UserStatsReceived {
username: String,
average_speed: u32,
shared_files: u32,
shared_folders: u32,
},
PeerConnected(String),
IncomingSearch {
username: String,
token: u32,
query: String,
},
QueueUpload {
requester_key: String,
filename: String,
},
StartUpload {
token: u32,
},
ShareListRequested {
requester_key: String,
},
BrowseResult {
username: String,
directories: Vec<SharedDirectory>,
},
PeerConnectFailed(u64, String),
RoomEvent(RoomEvent),
WishlistInterval(u32),
PrivilegedUsers(Vec<String>),
OwnPrivileges(u32),
PlaceInQueueRequested {
requester_key: String,
filename: String,
},
}
pub struct ClientContext {
pub peer_registry: Option<PeerRegistry>,
pub downloads: DownloadStore,
server_sender: Option<Sender<ServerMessage>>,
searches: HashMap<String, Search>,
private_messages: Vec<UserMessage>,
pending_connect_tokens: HashMap<u32, String>,
pub shares: Arc<Shares>,
pub shared_directories: Vec<String>,
peer_addresses: HashMap<String, (String, u32)>,
pending_peer_messages: HashMap<String, Vec<crate::message::Message>>,
uploads: HashMap<u32, UploadJob>,
active_uploads: HashMap<u32, ActiveUpload>,
pending_serves: HashMap<String, Vec<u32>>,
browse_results: HashMap<String, Vec<SharedDirectory>>,
room_list: Vec<RoomInfo>,
room_events: Vec<RoomEvent>,
room_members: HashMap<String, Vec<String>>,
user_info: HashMap<String, UserInfo>,
wishlist_interval: Option<u32>,
privileged_users: HashSet<String>,
own_privileges: Option<u32>,
upload_queue: Vec<QueuedUpload>,
upload_seq: u64,
upload_slots: usize,
upload_events: Vec<crate::types::UploadInfo>,
actor_system: Arc<ActorSystem>,
}
impl Default for ClientContext {
fn default() -> Self {
Self::new()
}
}
impl ClientContext {
pub fn add_download(&mut self, download: Download) {
self.downloads.add(download);
}
pub fn remove_download(&mut self, token: u32) {
self.downloads.remove(token);
}
#[must_use]
pub fn get_download_by_token(&self, token: u32) -> Option<&Download> {
self.downloads.get_by_token(token)
}
pub fn get_download_by_token_mut(
&mut self,
token: u32,
) -> Option<&mut Download> {
self.downloads.get_by_token_mut(token)
}
#[must_use]
pub fn get_download_tokens(&self) -> Vec<u32> {
self.downloads.tokens()
}
#[must_use]
pub const fn get_downloads(&self) -> &Vec<Download> {
self.downloads.list()
}
pub fn update_download_with_status(
&mut self,
token: u32,
status: DownloadStatus,
) {
self.downloads.update_status(token, status);
}
}
impl ClientContext {
#[must_use]
pub fn new() -> Self {
let actor_system = Arc::new(ActorSystem::new());
Self {
peer_registry: None,
server_sender: None,
searches: HashMap::new(),
private_messages: Vec::new(),
pending_connect_tokens: HashMap::new(),
shares: Arc::new(Shares::empty()),
shared_directories: Vec::new(),
peer_addresses: HashMap::new(),
pending_peer_messages: HashMap::new(),
uploads: HashMap::new(),
active_uploads: HashMap::new(),
pending_serves: HashMap::new(),
browse_results: HashMap::new(),
room_list: Vec::new(),
room_events: Vec::new(),
room_members: HashMap::new(),
user_info: HashMap::new(),
wishlist_interval: None,
privileged_users: HashSet::new(),
own_privileges: None,
upload_queue: Vec::new(),
upload_seq: 0,
upload_slots: DEFAULT_UPLOAD_SLOTS,
upload_events: Vec::new(),
downloads: DownloadStore::new(),
actor_system,
}
}
pub fn apply_room_event(&mut self, event: RoomEvent) {
match &event {
RoomEvent::List(rooms) => self.room_list.clone_from(rooms),
RoomEvent::Joined { room, users } => {
let mut members = users.clone();
members.sort();
members.dedup();
self.room_members.insert(room.clone(), members);
}
RoomEvent::Left { room } => {
self.room_members.remove(room);
}
RoomEvent::UserJoined { room, username } => {
let members =
self.room_members.entry(room.clone()).or_default();
if let Err(at) = members.binary_search(username) {
members.insert(at, username.clone());
}
}
RoomEvent::UserLeft { room, username } => {
if let Some(members) = self.room_members.get_mut(room)
&& let Ok(at) = members.binary_search(username)
{
members.remove(at);
}
}
RoomEvent::Message { .. } => {}
}
self.room_events.push(event);
}
pub fn apply_user_status(
&mut self,
username: String,
status: u32,
privileged: bool,
) {
self.user_info
.entry(username.clone())
.or_insert_with(|| UserInfo::pending(username))
.presence = Some(UserPresence {
status: UserStatus::from_code(status),
privileged,
});
}
pub fn invalidate_user_info(&mut self, username: &str) {
self.user_info.remove(username);
}
pub fn apply_user_stats(
&mut self,
username: String,
average_speed: u32,
shared_files: u32,
shared_folders: u32,
) {
self.user_info
.entry(username.clone())
.or_insert_with(|| UserInfo::pending(username))
.stats = Some(UserStats {
average_speed,
shared_files,
shared_folders,
});
}
#[must_use]
pub fn user_info(&self, username: &str) -> Option<UserInfo> {
self.user_info.get(username).cloned()
}
#[must_use]
pub fn room_members(&self, room: &str) -> Vec<String> {
self.room_members.get(room).cloned().unwrap_or_default()
}
#[must_use]
pub fn room_list(&self) -> Vec<RoomInfo> {
self.room_list.clone()
}
#[must_use]
pub fn take_room_events(&mut self) -> Vec<RoomEvent> {
std::mem::take(&mut self.room_events)
}
pub fn cache_peer_address(
&mut self,
username: &str,
host: String,
port: u32,
) {
self.peer_addresses
.insert(username.to_string(), (host, port));
}
#[must_use]
pub fn peer_address(&self, username: &str) -> Option<(String, u32)> {
self.peer_addresses.get(username).cloned()
}
pub fn queue_peer_message(
&mut self,
username: &str,
message: crate::message::Message,
) {
self.pending_peer_messages
.entry(username.to_string())
.or_default()
.push(message);
}
pub fn take_peer_messages(
&mut self,
username: &str,
) -> Vec<crate::message::Message> {
self.pending_peer_messages
.remove(username)
.unwrap_or_default()
}
pub fn store_browse_result(
&mut self,
username: String,
directories: Vec<SharedDirectory>,
) {
self.browse_results.insert(username, directories);
}
pub fn take_browse_result(
&mut self,
username: &str,
) -> Option<Vec<SharedDirectory>> {
self.browse_results.remove(username)
}
pub fn add_pending_connect(&mut self, token: u32, username: String) {
self.pending_connect_tokens.insert(token, username);
}
pub fn take_pending_connect(&mut self, token: u32) -> Option<String> {
self.pending_connect_tokens.remove(&token)
}
pub fn push_private_message(&mut self, message: UserMessage) {
self.private_messages.push(message);
}
pub fn take_private_messages(&mut self) -> Vec<UserMessage> {
std::mem::take(&mut self.private_messages)
}
}
pub struct Client {
enable_listen: bool,
listen_port: u16,
bound_port: Option<u16>,
address: PeerAddress,
username: String,
password: String,
version: ClientVersion,
shared_directories: Vec<String>,
server_handle: Option<ActorHandle<ServerMessage>>,
context: Arc<RwLock<ClientContext>>,
session: SessionWatch,
}
impl Client {
pub fn new(
username: impl Into<String>,
password: impl Into<String>,
) -> Self {
Self::with_settings(ClientSettings::new(username, password))
}
#[must_use]
pub fn with_settings(settings: ClientSettings) -> Self {
logger::init();
Self {
enable_listen: settings.enable_listen,
listen_port: settings.listen_port,
bound_port: None,
address: settings.server_address,
username: settings.username,
password: settings.password,
version: settings.version,
shared_directories: settings.shared_directories,
context: Arc::new(RwLock::new(ClientContext::new())),
server_handle: None,
session: SessionWatch::default(),
}
}
#[must_use]
pub fn username(&self) -> &str {
&self.username
}
#[must_use]
pub const fn listen_port(&self) -> Option<u16> {
self.bound_port
}
#[must_use]
pub fn session_loss(&self) -> Option<SessionLoss> {
self.session.loss()
}
#[must_use]
pub fn shared_directories(&self) -> Vec<String> {
self.context
.read_safe()
.map(|ctx| ctx.shared_directories.clone())
.unwrap_or_default()
}
#[must_use]
pub fn shared_counts(&self) -> (u32, u32) {
self.context.read_safe().map_or((0, 0), |ctx| {
(ctx.shares.folder_count(), ctx.shares.file_count())
})
}
#[must_use]
pub fn uploads(&self) -> Vec<crate::types::UploadInfo> {
self.context.read_safe().map_or_else(
|_| Vec::new(),
|ctx| {
let mut tokens: Vec<&u32> = ctx.active_uploads.keys().collect();
tokens.sort_unstable();
let mut all: Vec<crate::types::UploadInfo> = tokens
.into_iter()
.map(|token| {
let upload = &ctx.active_uploads[token];
let bytes_sent = upload
.bytes_sent
.load(std::sync::atomic::Ordering::Relaxed);
crate::types::UploadInfo {
username: upload.username.clone(),
filename: upload.filename.clone(),
size: upload.size,
bytes_sent,
speed_bytes_per_sec: upload_speed(
&upload.status,
bytes_sent,
upload.started,
),
status: upload.status.clone(),
}
})
.collect();
all.extend(ctx.queued_uploads());
all
},
)
}
#[must_use = "returns whether a matching upload was found"]
pub fn cancel_upload(&self, username: &str, filename: &str) -> bool {
self.context.read_safe().is_ok_and(|ctx| {
let mut found = false;
for upload in ctx.active_uploads.values() {
if upload.username == username
&& upload.filename == filename
&& upload.status == crate::types::UploadStatus::InProgress
{
upload
.cancel
.store(true, std::sync::atomic::Ordering::Relaxed);
found = true;
}
}
found
})
}
pub fn set_shared_directories(&self, dirs: Vec<String>) -> Result<()> {
let roots: Vec<std::path::PathBuf> = dirs
.iter()
.filter(|dir| !dir.trim().is_empty())
.map(std::path::PathBuf::from)
.collect();
let shares = if roots.is_empty() {
Shares::empty()
} else {
Shares::scan_many(&roots)
};
info!(
"Now sharing {} files in {} folders from {} directories",
shares.file_count(),
shares.folder_count(),
roots.len()
);
let folder_count = shares.folder_count();
let file_count = shares.file_count();
{
let mut ctx = self.context.write_safe()?;
ctx.shares = Arc::new(shares);
ctx.shared_directories = dirs;
}
self.send_server_message(
crate::message::server::MessageFactory::build_shared_folders_message(
folder_count,
file_count,
),
)
}
}
mod connection;
mod downloads;
mod operations;
mod rooms;
mod search;
mod upload_queue;
mod uploads;
#[cfg(test)]
mod tests;