pub mod connection;
pub mod handshake;
use std::
{
path::Path,
future::Future,
net::{ IpAddr, SocketAddr },
time::{ Instant, Duration },
collections::{ HashSet, HashMap },
sync::
{
Arc,
LazyLock,
Mutex as MutexSync,
},
};
use tokio::
{
fs,
time,
sync::{ Mutex, oneshot },
net::tcp::OwnedWriteHalf,
task::{ self, AbortHandle },
};
use rand::
{
TryRng,
rngs::SysRng,
};
use zeroize::Zeroizing;
use dashmap::DashMap;
use connection::
{
AvailableFile,
ConnectionType,
Connection,
};
use crate::
{
misc,
options,
role::Role,
crypto::password,
config::{ self, users },
consts::
{
self,
SharedKeys,
Streams,
},
network::
{
self,
voice::server as voice_server,
file::server as file,
codes::
{
PacketCode,
OnlineUser,
StoredMessage,
UserFile,
UserScreen,
},
},
};
pub static PENDING_TOKENS: LazyLock<DashMap<[u8; 32], (usize, ConnectionType, Instant)>> = LazyLock::new(|| DashMap::new());
pub static CONNECTIONS: LazyLock<DashMap<SocketAddr, Connection>> = LazyLock::new(|| DashMap::new()); pub static AVAILABLE_FILES: LazyLock<DashMap<String, Vec<AvailableFile>>> = LazyLock::new(|| DashMap::new());
async fn send_bans(write_stream: &Arc<Mutex<OwnedWriteHalf>>, keys: &SharedKeys) {
network::send(&mut *write_stream.lock().await, PacketCode::ServerBans
{
users: Some(config::bans::users()),
ips: Some(config::bans::ips()),
}, Some(keys)).await;
}
async fn send_history(write_stream: &Arc<Mutex<OwnedWriteHalf>>, keys: &SharedKeys) {
if !config::read_config::<bool>("persistent_messages") { return; }
let stored = config::messages::all();
if stored.is_empty() { return; }
let mut messages: Vec<StoredMessage> = Vec::new();
let mut budget = consts::MAX_HISTORY_SIZE;
for message in stored.into_iter().rev()
{
let size = message.username.len() + message.text.len();
if size > budget { break; }
budget -= size;
messages.push(message);
}
if messages.is_empty() { return; }
messages.reverse();
log::debug!("Replaying history ({} messages)", messages.len());
network::send(&mut *write_stream.lock().await, PacketCode::History { messages }, Some(keys)).await;
}
async fn remove_connections(addr: &IpAddr, grace: bool, info: Option<&str>) {
let addrs: Vec<SocketAddr> = CONNECTIONS.iter()
.filter(|conn| conn.key().ip() == *addr)
.map(|conn| *conn.key()).collect();
for addr in addrs
{
remove_connection(&addr, grace, info).await;
}
}
pub fn log_addr(id: &usize) -> String
{
CONNECTIONS.iter().find(|conn| conn.id() == Some(id))
.map(|conn| conn.peer_addr().to_string())
.unwrap_or_else(|| String::from("<gone>"))
}
pub fn spawn_with_abort<F, Fut>(f: F) -> AbortHandle where
F: FnOnce(AbortHandle) -> Fut + Send + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
let (tx, rx) = oneshot::channel::<AbortHandle>();
let handle = tokio::spawn(async move
{
let task = match rx.await
{
Ok(t) => t,
Err(_) => return
};
f(task).await;
});
let task = handle.abort_handle();
tx.send(task.clone()).ok();
task
}
pub fn send_to_all(code: PacketCode, filter_channel: bool, channel: Option<&str>) {
let entries: Vec<Connection> = CONNECTIONS.iter().filter_map(|entry|
{
match entry.value()
{
Connection::Authenticated { channel: c, .. } if !filter_channel || c.as_deref() == channel =>
{
Some(entry.value().clone())
},
_ => None,
}
}).collect();
for ref entry in entries
{
let write_stream = entry.write_stream().clone();
let code = code.clone();
let keys = entry.keys().cloned();
tokio::spawn(async move
{
network::send(&mut *write_stream.lock().await, code, keys.as_ref()).await;
});
}
}
pub async fn remove_connection(peer_addr: &SocketAddr, grace: bool, info: Option<&str>) {
let mut connection = match CONNECTIONS.remove(peer_addr)
{
Some((_, conn)) => conn,
None => return
};
if grace
{
network::send(&mut *connection.write_stream().lock().await, PacketCode::Disconnect, connection.keys()).await;
}
if let Some(streams) = connection.file_streams()
{
let uids: Vec<u64> = streams.lock().unwrap().keys().copied().collect();
for uid in uids
{
connection.remove_file_stream(uid);
}
}
if let Some(id) = connection.remove_screen_stream()
{
deattach(id, connection.username().unwrap()).await;
}
if connection.is_authenticated()
{
if let Some(target_id) = connection.attached_screen().as_ref().map(|a| a.target_id)
{
notify(target_id, PacketCode::Deattached
{
username: connection.username().unwrap().to_owned(),
}).await;
}
if options::voice_chat_enabled()
{
voice_server::remove_connection(connection.id().unwrap());
}
let username = connection.username().unwrap();
let _ = fs::remove_dir_all(misc::get_upload_dir(username)).await; file::ACTIVE_FILESHARES.retain(|_, u| u.client_id != *connection.id().unwrap());
AVAILABLE_FILES.remove(username);
send_to_all(PacketCode::Leave
{
username: connection.username().unwrap().to_string(),
id: *connection.id().unwrap(),
}, false, None);
}
log::info!
(
"Close connection{}: {peer_addr} ({}, {} left)",
if let Some(info) = info
{
format!(" ({info})")
} else { String::new() },
if connection.is_authenticated() { "authenticated" } else { "unauthenticated" },
CONNECTIONS.len(),
);
connection.task().abort();
}
fn user_connected(username: &str) -> bool {
CONNECTIONS.iter().any(|conn|
{
conn.username().map_or(false, |u| u == &username.to_string())
})
}
fn get_latest_id() -> usize
{
let ids: HashSet<usize> = CONNECTIONS.iter().filter_map(|conn|
{
if let Some(id) = conn.id()
{
Some(*id)
} else
{
None
}
}).collect();
for i in 0..
{
if !ids.contains(&i) {
return i;
}
}
unreachable!("what the fuck");
}
fn update_client_keys(peer_addr: &SocketAddr, keys: &SharedKeys) {
CONNECTIONS.alter(peer_addr, |_, old_connection|
{
match old_connection
{
Connection::NonAuthenticated { write_stream, task, seq, peer_addr, connect, obfuscation_key, .. } =>
{
Connection::NonAuthenticated
{
write_stream,
task,
peer_addr,
username: None,
keys: Some(keys.to_owned()),
obfuscation_key,
last_activity: Instant::now(),
seq,
connect,
}
},
Connection::Authenticated { write_stream, task, file_streams, screen_stream, username, role,
id, attached_screen, last_activity, last_image, channel, seq, server_seq, peer_addr, alive, muted,
credit, refill, throttles, .. } =>
{
Connection::Authenticated
{
write_stream,
task,
file_streams,
screen_stream,
peer_addr,
username,
role,
id,
keys: keys.to_owned(),
attached_screen,
last_activity,
last_key_exchange: Instant::now(),
last_image,
spam_violations: 0,
credit,
refill,
throttles,
channel,
seq,
server_seq,
alive,
muted,
}
}
}
});
}
fn authenticate_client(peer_addr: &SocketAddr, username: &str, role: Role, id: usize) {
CONNECTIONS.alter(&peer_addr, |_, old_connection|
{
Connection::Authenticated
{
write_stream: old_connection.write_stream().clone(),
task: old_connection.task().clone(),
file_streams: Arc::new(MutexSync::new(HashMap::new())),
screen_stream: None,
peer_addr: *old_connection.peer_addr(),
username: username.to_string(),
role,
id,
keys: old_connection.keys().unwrap().to_owned(),
attached_screen: None,
last_activity: Instant::now() - Duration::from_millis(config::read_config("min_message_delay")),
last_key_exchange: old_connection.last_key_exchange().copied().unwrap_or_else(Instant::now),
last_image: Instant::now() - consts::IMAGE_REQUEST_DELAY,
spam_violations: 0,
credit: config::read_config::<f32>("max_packet_burst"),
refill: Instant::now(),
throttles: 0,
channel: None,
seq: *old_connection.seq(),
server_seq: 0,
alive: true,
muted: false,
}
});
AVAILABLE_FILES.insert(username.to_string(), Vec::new());
log::info!("Authenticate connection: {peer_addr} (id {id}, role {role})");
}
fn update_client_channel(peer_addr: &SocketAddr, channel: &Option<String>) {
let mut old_channel = None;
CONNECTIONS.alter(&peer_addr, |_, old_connection|
{
old_channel = old_connection.channel().clone();
Connection::Authenticated
{
write_stream: old_connection.write_stream().clone(),
task: old_connection.task().clone(),
file_streams: old_connection.file_streams().unwrap().clone(),
screen_stream: old_connection.screen_stream().clone(),
peer_addr: *old_connection.peer_addr(),
username: old_connection.username().unwrap().clone(),
role: *old_connection.role().unwrap(),
id: *old_connection.id().unwrap(),
keys: old_connection.keys().unwrap().to_owned(),
attached_screen: old_connection.attached_screen().clone(),
last_activity: Instant::now(),
last_key_exchange: *old_connection.last_key_exchange().unwrap(),
last_image: *old_connection.last_image().unwrap(),
spam_violations: *old_connection.spam_violations().unwrap(),
credit: *old_connection.credit().unwrap(),
refill: *old_connection.refill().unwrap(),
throttles: *old_connection.throttles().unwrap(),
channel: channel.clone(),
seq: *old_connection.seq(),
server_seq: *old_connection.server_seq().unwrap(),
alive: true,
muted: *old_connection.muted(),
}
});
if old_channel == *channel { return; }
log::info!("Channel switch: {peer_addr} ({} -> {})",
if old_channel.is_some() { "channel" } else { "lobby" },
if channel.is_some() { "channel" } else { "lobby" });
if let Some(old_channel) = old_channel
{
if !CONNECTIONS.iter().any(|c| c.channel().as_ref() == Some(&old_channel))
{
log::info!("Channel destroyed (last client left): {peer_addr}");
send_to_all(PacketCode::ChannelDestroyed
{
name: old_channel,
}, false, None);
}
}
if let Some(channel) = channel
{
if CONNECTIONS.iter().filter(|c| c.channel().as_ref() == Some(channel)).count() == 1
{
log::info!("Channel created: {peer_addr}");
send_to_all(PacketCode::ChannelCreated
{
name: channel.clone(),
}, false, None);
}
}
}
async fn send_voice_clients(stream: &mut OwnedWriteHalf, keys: &SharedKeys, id: usize)
{
let sender_channel = match CONNECTIONS.iter().find(|e| e.value().id() == Some(&id))
{
Some(entry) => entry.value().channel().clone(),
None => return,
};
let mut clients: Vec<(usize, String)> = Vec::new();
for entry in CONNECTIONS.iter()
{
let conn = entry.value();
let uid = match conn.id()
{
Some(i) => *i,
None => continue
};
if uid == id { continue; } if conn.channel() != &sender_channel { continue; }
if voice_server::CONNECTIONS.contains_key(&uid)
{
if let Some(username) = conn.username()
{
clients.push((uid, username.clone()));
}
}
}
network::send(stream, PacketCode::VoiceClients { clients }, Some(keys)).await;
}
fn open_connection(id: usize, conn_type: ConnectionType) -> [u8; 32] {
let mut token = [0u8; 32];
SysRng.try_fill_bytes(&mut token).unwrap();
PENDING_TOKENS.insert(token, (id, conn_type, Instant::now()));
token
}
pub async fn notify(id: usize, code: PacketCode) {
let target = CONNECTIONS.iter().find(|conn| conn.id() == Some(&id))
.map(|conn| (conn.write_stream().clone(), conn.keys().cloned()));
if let Some((write_stream, keys)) = target
{
network::send(&mut *write_stream.lock().await, code, keys.as_ref()).await;
}
}
pub async fn deattach(sharer_id: usize, sharer_uname: &String) {
let mut to_notify = Vec::new();
for mut conn in CONNECTIONS.iter_mut()
.filter(|conn| conn.attached_screen().as_ref().is_some_and(|a| a.target_id == sharer_id))
{
conn.deattach_screen();
to_notify.push((conn.write_stream().clone(), conn.keys().cloned()));
}
for (stream_mutex, keys) in to_notify
{
network::send(&mut *stream_mutex.lock().await,
PacketCode::Deattach { username: Some(sharer_uname.to_owned()) }, keys.as_ref()).await;
}
}
pub async fn listen_client (
streams: &mut Streams<'_>,
peer_addr: SocketAddr,
obfuscation_key: [u8; 32],
task: AbortHandle,
)
{
log::info!("New connection: {peer_addr}");
CONNECTIONS.insert(peer_addr, Connection::NonAuthenticated
{
write_stream: streams.1.clone(),
task,
peer_addr: peer_addr,
username: None,
keys: None,
obfuscation_key,
last_activity: Instant::now(),
seq: 0,
connect: Instant::now(),
});
let mut keys = (Zeroizing::new(vec![]), Zeroizing::new(vec![]));
handshake::key_exchange(streams, &peer_addr, &obfuscation_key, &mut keys, None).await;
if keys.0.is_empty() || keys.1.is_empty()
{
log::warn!("Key exchange failed: {peer_addr}");
return remove_connection(&peer_addr, false, None).await
}
log::debug!("Key exchange done: {peer_addr}");
if config::read_config("check_client_version")
{
let version = handshake::ask_version(streams, &keys).await;
if version.is_none() || version != Some(misc::get_version().to_string())
{
log::warn!("Version mismatch: {peer_addr} (client {}, server {})",
version.as_deref().unwrap_or("none"), misc::get_version());
return remove_connection(&peer_addr, true, Some("version")).await;
}
}
handshake::send_welcome_packet(&mut *streams.1.lock().await, &keys).await;
let mut username: Option<String> = None;
let max_tries = config::read_config::<usize>("max_auth_tries"); let min_len = config::read_config::<usize>("min_username_length");
let max_len = config::read_config::<usize>("max_username_length");
let disabled_registration = !config::read_config::<bool>("allow_register");
if disabled_registration
{
network::send(&mut *streams.1.lock().await, PacketCode::RegisterDisabled, Some(&keys)).await;
}
for attempt in 0..max_tries
{
log::debug!("Asking for username: {peer_addr} (try {}/{max_tries})", attempt + 1);
network::send(&mut *streams.1.lock().await, PacketCode::Username { username: None }, Some(&keys)).await;
match network::receive(streams, Some(&keys), None).await
{
Some(PacketCode::Username { username: uname }) =>
{
if let Some(uname) = uname
{
if uname.len() >= min_len && uname.len() <= max_len &&
uname.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-') &&
!user_connected(&uname) && uname != options::get_server_username()
{
username = Some(uname);
break;
}
log::debug!("Username refused: {peer_addr}");
}
},
_ => return remove_connection(&peer_addr, false, Some("username")).await,
}
}
if username.is_none()
{
return remove_connection(&peer_addr, true, Some("username")).await;
}
let username = username.unwrap();
if let Some(mut conn) = CONNECTIONS.get_mut(&peer_addr)
{
if !conn.is_authenticated()
{
*conn.username_mut() = Some(username.clone());
}
}
let user_exists = config::users::contains(&username);
if !user_exists && !disabled_registration {
let max_tries = config::read_config::<usize>("max_auth_tries"); let mut password: Option<Zeroizing<String>> = None;
for _ in 0..max_tries
{
network::send(&mut *streams.1.lock().await, PacketCode::PasswordR { password: None }, Some(&keys)).await;
match network::receive(streams, Some(&keys), None).await
{
Some(PacketCode::PasswordR { password: pass }) =>
{
if let Some(pass) = pass
{
let pass = Zeroizing::new(pass);
if pass.len() >= config::read_config("min_password_length")
{
password = Some(pass);
break;
}
}
},
_ => return remove_connection(&peer_addr, false, Some("register")).await
};
}
if password.is_none()
{
return remove_connection(&peer_addr, true, Some("register")).await;
}
log::info!("Registering new user: {peer_addr}");
let hash = task::spawn_blocking(move || password::hash_password(password.as_ref().unwrap().as_str()))
.await.expect("Hashing password failed");
if config::users::add(&username, &hash)
{
log::info!("First user registered, granting {}: {peer_addr}", Role::Owner);
network::send(&mut *streams.1.lock().await, PacketCode::FirstUser, Some(&keys)).await;
}
} else {
network::send(&mut *streams.1.lock().await, PacketCode::PasswordL { password: None }, Some(&keys)).await;
let password = loop
{
match network::receive(streams, Some(&keys), None).await
{
Some(PacketCode::PasswordL { password: Some(password) }) => break Zeroizing::new(password),
_ => return remove_connection(&peer_addr, false, Some("login")).await,
}
};
let valid = if password.is_empty() || config::bans::banned(&username)
{
false
} else if let Some(hashed) = config::users::password(&username)
{
task::spawn_blocking(move || password::compare_password_hash(&hashed, &password))
.await.expect("Comparing password failed")
} else {
false
};
if !valid
{
log::warn!("Login refused: {peer_addr}");
return remove_connection(&peer_addr, true, Some("login")).await;
}
log::debug!("Login accepted: {peer_addr}");
}
let id = get_latest_id(); let role = config::users::role(&username).unwrap();
let mut channel: Option<String> = None;
authenticate_client(&peer_addr, &username, role, id);
send_history(&streams.1, &keys).await;
network::send(&mut *streams.1.lock().await, PacketCode::Accept { id, role }, Some(&keys)).await;
send_to_all(PacketCode::Join { username: username.clone() }, false, None);
if options::voice_chat_enabled()
{
send_voice_clients(&mut *streams.1.lock().await, &keys, id).await;
}
loop
{
let read = match network::receive(streams, Some(&keys), None).await
{
Some(r) => r,
None => return
};
if Instant::now().duration_since(*CONNECTIONS.get(&peer_addr).unwrap().last_key_exchange().unwrap()) >=
Duration::from_secs(consts::REKEY_INTERVAL) &&
!file::ACTIVE_FILESHARES.iter().any(|entry| entry.client_id == id) {
log::debug!("Rekeying: {peer_addr}");
let current_keys = keys.clone();
handshake::key_exchange(streams, &peer_addr, &obfuscation_key, &mut keys, Some(¤t_keys)).await; }
let role = CONNECTIONS.get(&peer_addr).and_then(|entry| entry.role().copied()).unwrap_or(role);
log::debug!("Packet {}: {peer_addr}", read.name());
match read
{
PacketCode::Message { text, colors, .. } =>
{
if *CONNECTIONS.get(&peer_addr).unwrap().muted()
{
log::debug!("Message dropped (muted): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::Muted, Some(&keys)).await;
continue;
}
let text = text.trim().to_owned();
log::info!("Message ({} chars) in {}: {peer_addr}", text.chars().count(),
if channel.is_some() { "channel" } else { "lobby" });
if channel.is_none() && config::read_config::<bool>("persistent_messages")
{
config::messages::store(&username, &text, &colors);
}
send_to_all(PacketCode::Message
{
text,
username: Some(username.clone()),
id: Some(id),
colors,
}, true, channel.as_deref());
}
PacketCode::Disconnect =>
{
return remove_connection(&peer_addr, true, None).await;
},
PacketCode::Voice { .. } =>
{
if !options::voice_chat_enabled()
{
log::warn!("Voice refused (disabled): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidFeature, Some(&keys)).await;
} else if !voice_server::CONNECTIONS.contains_key(&id) {
log::info!("Voice slot opened: {peer_addr}");
let token = voice_server::open_connection(id, username.clone());
network::send(&mut *streams.1.lock().await, PacketCode::Voice { token: Some(token) }, Some(&keys)).await;
send_to_all(PacketCode::VoiceJoin { username: username.clone(), id }, true, channel.as_deref());
send_voice_clients(&mut *streams.1.lock().await, &keys, id).await;
} else {
log::info!("Voice leave: {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::Voice { token: None }, Some(&keys)).await;
send_to_all(PacketCode::VoiceLeave { id }, true, channel.as_deref());
voice_server::remove_connection(&id);
}
},
PacketCode::Channel { channel: tchannel } =>
{
if tchannel.iter().all(|s| !s.is_empty() && s.len() <= config::read_config("max_channel_length") && s.chars().all(|c| c.is_ascii_alphanumeric() && c != ' '))
{
if options::voice_chat_enabled() && voice_server::CONNECTIONS.contains_key(&id)
{
send_to_all(PacketCode::VoiceLeave { id }, true, channel.as_deref());
}
update_client_channel(&peer_addr, &tchannel);
channel = tchannel.clone();
network::send(&mut *streams.1.lock().await, PacketCode::Channel { channel: tchannel }, Some(&keys)).await;
if options::voice_chat_enabled() && voice_server::CONNECTIONS.contains_key(&id)
{
send_to_all(PacketCode::VoiceJoin { username: username.clone(), id }, true, channel.as_deref());
}
send_voice_clients(&mut *streams.1.lock().await, &keys, id).await;
} else {
log::warn!("Channel refused (invalid name): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
}
},
PacketCode::List { .. } =>
{
let mut users = Vec::new();
for connection_enum in CONNECTIONS.iter()
{
if let Connection::Authenticated { username: uname, id: user_id, channel, .. } = connection_enum.value()
{
users.push(OnlineUser
{
username: uname.clone(),
id: *user_id,
channel: channel.clone(),
});
}
}
log::debug!("Sending online list ({} users): {peer_addr}", users.len());
network::send(&mut *streams.1.lock().await, PacketCode::List { users: Some(users) }, Some(&keys)).await;
},
PacketCode::Upload { hash, .. } | PacketCode::Image { hash, .. } =>
{
if *CONNECTIONS.get(&peer_addr).unwrap().muted()
{
network::send(&mut *streams.1.lock().await, PacketCode::Muted, Some(&keys)).await;
continue;
}
let image = matches!(read, PacketCode::Image { .. });
let username_color = match &read
{
PacketCode::Image { username_color, .. } => *username_color,
_ => None,
};
if image && config::messages::has_image(&hash)
{
let filename = match &read
{
PacketCode::Image { filename, .. } => Path::new(filename)
.file_name()
.and_then(|f| f.to_str())
.unwrap_or("unnamed_file")
.to_string(),
_ => String::from("unnamed_file"),
};
if channel.is_none() && config::read_config::<bool>("persistent_messages")
{
config::messages::store_image(&username, &filename, &hash, username_color);
}
send_to_all(PacketCode::ImageDisplay
{
username: username.clone(),
filename: filename.clone(),
hash,
data: None,
username_color,
}, true, channel.as_deref());
network::send(&mut *streams.1.lock().await, PacketCode::Image
{
hash,
filename,
token: None,
uid: None,
username_color,
}, Some(&keys)).await;
log::info!("Image already stored, upload skipped: {peer_addr}");
continue;
}
let active_count = file::ACTIVE_FILESHARES.iter().filter(|u| u.client_id == id).count();
if active_count >= config::read_config::<usize>("max_client_parallel_uploads")
{
log::warn!("Upload refused ({active_count} already running): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::UploadLimit, Some(&keys)).await;
continue;
}
let uid = rand::random::<u64>();
let token = open_connection(id, if image
{
ConnectionType::Image { uid, username_color }
} else
{
ConnectionType::FileUpload { uid }
});
log::info!("Upload request ({}): {peer_addr}", if image { "image" } else { "file" });
let packet = if image
{
PacketCode::Image
{
hash,
filename: String::new(), token: Some(token),
uid: Some(uid),
username_color,
}
} else
{
PacketCode::Upload
{
hash,
token: Some(token),
uid: Some(uid),
}
};
network::send(&mut *streams.1.lock().await, packet, Some(&keys)).await;
},
PacketCode::Download { id: owner_id, file_id, .. } =>
{
let username = CONNECTIONS.iter()
.find(|entry| entry.value().id() == owner_id.as_ref())
.and_then(|entry| entry.value().username().cloned());
if let Some(username) = username &&
let Some(file_id) = file_id &&
let Some(file) = AVAILABLE_FILES.get(&username).and_then(|f| f.value().get(file_id).cloned())
{
let uid = rand::random::<u64>();
let file_size = file.size;
let token = open_connection(id, ConnectionType::FileDownload { uid, file });
network::send(&mut *streams.1.lock().await, PacketCode::Download
{
token: Some(token),
file_id: None,
id: None,
}, Some(&keys)).await;
log::info!("Download request ({} bytes): {peer_addr}", file_size);
} else
{
log::warn!("Download refused (unknown file): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
}
},
PacketCode::Screen { .. } =>
{
if let Some((Some(removed_id), Some(username))) = CONNECTIONS.get_mut(&peer_addr)
.and_then(|mut conn| Some((conn.remove_screen_stream(), conn.username().cloned())))
{
deattach(removed_id, &username).await;
log::info!("Screen share stopped: {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::Screen { token: None }, Some(&keys)).await;
send_to_all(PacketCode::ScreenshareEnd { username: username.clone() }, false, None);
continue;
}
if config::read_config("enable_screenshare")
{
network::send(&mut *streams.1.lock().await, PacketCode::Screen
{
token: Some(open_connection(id, ConnectionType::Screen))
}, Some(&keys)).await;
send_to_all(PacketCode::Screenshare { username: username.clone() }, false, None);
log::info!("Screen share: {peer_addr}");
} else
{
log::warn!("Screen share refused (disabled): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidFeature, Some(&keys)).await;
}
},
PacketCode::Attach { id: sharer_id, .. } =>
{
let sharer_info = sharer_id.and_then(|sid|
{
CONNECTIONS.iter().find(|entry| entry.value().id() == Some(&sid) && entry.screen_stream().is_some())
.and_then(|conn| conn.username().map(|u| (sid, u.to_owned())))
});
if let Some((sharer_id, sharer_username)) = sharer_info &&
CONNECTIONS.get(&peer_addr).is_some_and(|c| c.attached_screen().as_ref()
.is_none_or(|s| s.target_id != sharer_id)) { let token = open_connection(id, ConnectionType::Attach
{
id: sharer_id,
});
network::send(&mut *streams.1.lock().await, PacketCode::Attach
{
id: None,
username: Some(sharer_username.to_owned()),
token: Some(token),
}, Some(&keys)).await;
notify(sharer_id, PacketCode::Attached
{
username: username.clone(),
}).await;
log::info!("Screen attach: {peer_addr} watching {}", log_addr(&sharer_id));
} else
{
log::warn!("Screen attach refused (no such share): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
}
},
PacketCode::Deattach { .. } =>
{
if let Some(sharer_id) = if let Some(mut conn) = CONNECTIONS.get_mut(&peer_addr)
{
if let Some(target_id) = conn.attached_screen().as_ref().map(|a| a.target_id)
{
conn.deattach_screen();
Some(target_id)
} else { None }
} else { None }
{
let sharer_uname = CONNECTIONS.iter()
.find(|c| c.value().id() == Some(&sharer_id))
.and_then(|c| c.value().username().cloned());
network::send(&mut *streams.1.lock().await, PacketCode::Deattach { username: sharer_uname }, Some(&keys)).await;
notify(sharer_id, PacketCode::Deattached
{
username: username.clone(),
}).await;
} else
{
log::warn!("Deattach refused (not attached): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
}
},
PacketCode::Files { .. } =>
{
let mut users = Vec::new();
for entry in AVAILABLE_FILES.iter()
{
let username = entry.key();
let uploads = entry.value();
if !uploads.is_empty() {
let id = CONNECTIONS.iter()
.find(|c| c.username() == Some(&username))
.and_then(|c| c.id().copied()).unwrap();
let upload: Vec<(String, usize)> = uploads.iter().enumerate()
.map(|(idx, u)| (u.filename.clone(), idx)).collect();
users.push(UserFile
{
username: username.to_string(),
id,
upload,
});
}
}
log::debug!("Sending file list ({} users): {peer_addr}", users.len());
network::send(&mut *streams.1.lock().await, PacketCode::Files { users: Some(users) }, Some(&keys)).await;
},
PacketCode::ImageData { hash, .. } =>
{
let wait = CONNECTIONS.get(&peer_addr)
.and_then(|conn| conn.last_image().map(|last| consts::IMAGE_REQUEST_DELAY
.saturating_sub(last.elapsed())))
.unwrap_or_default();
if !wait.is_zero()
{
log::debug!("Image fetch held for {}ms: {peer_addr}", wait.as_millis());
time::sleep(wait).await;
}
if let Some(mut conn) = CONNECTIONS.get_mut(&peer_addr)
&& let Some(last) = conn.last_image_mut() { *last = Instant::now(); }
let image = match config::messages::has_image(&hash)
{
true => file::read_image(&hash).await,
false => None,
};
match image.as_ref()
{
Some(data) => log::info!("Image fetch served ({} bytes): {peer_addr}", data.len()),
None => log::warn!("Image fetch refused (not in history): {peer_addr}"),
}
network::send(&mut *streams.1.lock().await, PacketCode::ImageData { hash, data: image }, Some(&keys)).await;
},
PacketCode::Screens { .. } =>
{
let mut users = Vec::new();
for connection_enum in CONNECTIONS.iter()
{
if let Connection::Authenticated { username: uname, id: user_id, screen_stream, .. } = connection_enum.value()
{
if screen_stream.is_none() { continue; }
users.push(UserScreen
{
username: uname.clone(),
id: *user_id,
});
}
}
log::debug!("Sending screenshare list ({} users): {peer_addr}", users.len());
network::send(&mut *streams.1.lock().await, PacketCode::Screens { users: Some(users) }, Some(&keys)).await;
},
PacketCode::PrivateMessage { text, id: recipient_id, .. } =>
{
if *CONNECTIONS.get(&peer_addr).unwrap().muted()
{
network::send(&mut *streams.1.lock().await, PacketCode::Muted, Some(&keys)).await;
continue;
}
let recipient_addr = CONNECTIONS.iter()
.find(|entry| entry.value().id() == Some(&recipient_id))
.map(|entry| *entry.key());
if let Some(recipient_addr) = recipient_addr
{
log::info!("Private message ({} chars): {peer_addr} -> {recipient_addr}", text.chars().count());
if recipient_id != id
{
let recipient_data = if let Some(recipient) =
CONNECTIONS.get(&recipient_addr)
{
Some((recipient.write_stream().clone(), recipient.keys().cloned()))
} else
{
None
};
if let Some((recipient_stream, recipient_keys)) = recipient_data
{
network::send(&mut *recipient_stream.lock().await, PacketCode::PrivateMessage
{
text: text.clone(),
username: Some(username.clone()),
id,
}, recipient_keys.as_ref()).await;
}
}
let recipient_uname = CONNECTIONS.get(&recipient_addr).and_then(|e| e.username().cloned()).unwrap();
network::send(&mut *streams.1.lock().await, PacketCode::PrivateMessageBack
{
text,
id: recipient_id,
username: recipient_uname,
}, Some(&keys)).await;
} else
{
log::warn!("Private message refused (no such recipient): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
}
},
PacketCode::ServerMute { id } =>
{
if role < Role::Moderator
{
log::warn!("Refused (permissions): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
let muted = if let Some(mut conn) = CONNECTIONS.iter_mut()
.find(|entry| entry.value().id() == Some(&id)) && conn.role() < Some(&role)
{
conn.toggle_mute();
Some((*conn.peer_addr(), *conn.muted()))
} else { None };
match muted
{
Some((target, muted)) => log::info!("{} by {peer_addr}: {target}",
if muted { "Mute" } else { "Unmute" }),
None =>
{
log::warn!("Mute refused (no such user, or a peer): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
}
}
},
PacketCode::ServerKick { id } =>
{
if role < Role::Moderator
{
log::warn!("Refused (permissions): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
let target = CONNECTIONS.iter()
.find(|entry| entry.value().id() == Some(&id))
.map(|entry| (*entry.key(), entry.role().cloned()));
if let Some((addr, Some(trole))) = target && trole < role
{
log::info!("Kick by {peer_addr}: {addr}");
remove_connection(&addr, true, Some("kick")).await;
} else {
log::warn!("Kick refused (no such user, or a peer): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
}
},
PacketCode::ServerBan { id: uid } =>
{
if role < Role::Owner || id == uid
{
log::warn!("Refused (permissions): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
let target = CONNECTIONS.iter()
.find(|entry| entry.value().id() == Some(&uid))
.map(|entry| (*entry.key(), entry.username().cloned()));
if let Some((addr, Some(username))) = target
{
log::info!("Ban by {peer_addr}: {addr}");
config::bans::ban(&username);
remove_connection(&addr, true, Some("ban")).await;
} else {
log::warn!("Ban refused (no such user): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
}
},
PacketCode::ServerBanIp { id: uid } =>
{
if role < Role::Owner || id == uid
{
log::warn!("Refused (permissions): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
let target = CONNECTIONS.iter()
.find(|entry| entry.value().id() == Some(&uid))
.map(|entry| *entry.key());
if let Some(addr) = target
{
log::info!("IP ban by {peer_addr}: {}", addr.ip());
config::bans::ban_ip(&addr.ip());
remove_connections(&addr.ip(), true, Some("ip ban")).await;
} else {
log::warn!("IP ban refused (no such user): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
}
},
PacketCode::ServerBans { .. } =>
{
if role < Role::Owner
{
log::warn!("Refused (permissions): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
send_bans(&streams.1, &keys).await;
},
PacketCode::ServerPardon { id: ban } =>
{
if role < Role::Owner
{
log::warn!("Refused (permissions): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
if config::bans::pardon(ban)
{
log::info!("Pardon by {peer_addr}");
send_bans(&streams.1, &keys).await;
} else
{
log::warn!("Pardon refused (no such ban): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
}
},
PacketCode::ServerPardonIp { id: ban } =>
{
if role < Role::Owner
{
log::warn!("Refused (permissions): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
if config::bans::pardon_ip(ban)
{
log::info!("IP pardon by {peer_addr}");
send_bans(&streams.1, &keys).await;
} else
{
log::warn!("IP pardon refused (no such ban): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
}
},
PacketCode::ServerSay { message } =>
{
if role < Role::Owner
{
log::warn!("Refused (permissions): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
log::info!("Server announcement ({} chars) by {peer_addr}", message.chars().count());
send_to_all(PacketCode::ServerSay { message }, false, None);
},
PacketCode::ServerRole { id: uid, role: new_role, .. } =>
{
if role < Role::Owner || id == uid || new_role > role
{
log::warn!("Refused (permissions): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
let target = CONNECTIONS.iter()
.find(|entry| entry.value().id() == Some(&uid))
.map(|entry| (entry.username().cloned(), entry.role().copied()));
let Some((Some(target_username), Some(target_role))) = target else
{
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
};
if target_role >= role || target_role == new_role
{
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
log::info!("Role change by {peer_addr}: {} is now {new_role} (was {target_role})", log_addr(&uid));
users::set_role(&target_username, new_role);
let sessions: Vec<(usize, Arc<Mutex<OwnedWriteHalf>>, SharedKeys)> =
{
let mut sessions = Vec::new();
for mut entry in CONNECTIONS.iter_mut()
.filter(|entry| entry.username() == Some(&target_username) && entry.role().is_some())
{
entry.set_role(new_role);
sessions.push((entry.id().copied().unwrap(), entry.write_stream().clone(), entry.keys().cloned().unwrap()));
}
sessions
};
for (sid, write_stream, target_keys) in sessions
{
network::send(&mut *write_stream.lock().await, PacketCode::ServerRole
{
id: sid,
role: new_role,
username: None,
}, Some(&target_keys)).await;
}
network::send(&mut *streams.1.lock().await, PacketCode::ServerRole
{
id: uid,
role: new_role,
username: Some(target_username),
}, Some(&keys)).await;
},
PacketCode::ServerSettings { settings, save } =>
{
if role < Role::Owner
{
log::warn!("Refused (permissions): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
if save && let Some(settings) = &settings
{
let accepted = config::settings::write(settings);
log::info!("Server settings saved by {peer_addr}: {accepted}/{} rows accepted", settings.len());
} else if !save
{
log::info!("Server settings read by {peer_addr}");
}
network::send(&mut *streams.1.lock().await, PacketCode::ServerSettings
{
settings: Some(config::settings::all()),
save,
}, Some(&keys)).await;
},
PacketCode::ServerRestart =>
{
if role < Role::Owner
{
log::warn!("Refused (permissions): {peer_addr}");
network::send(&mut *streams.1.lock().await, PacketCode::InvalidUsage, Some(&keys)).await;
continue;
}
log::info!("Restart requested by {peer_addr}");
tokio::spawn(async
{
disconnect_all().await;
misc::restart();
});
},
PacketCode::KeepAlive =>
{
if let Some(mut conn) = CONNECTIONS.get_mut(&peer_addr)
{
conn.set_alive(true);
}
},
_ => {}
}
}
}
pub async fn disconnect_all() {
let addrs: Vec<SocketAddr> = CONNECTIONS.iter().map(|conn| *conn.peer_addr()).collect();
log::info!("Disconnecting {} connections", addrs.len());
for addr in &addrs
{
remove_connection(addr, true, None).await; }
}
pub async fn disconnect_inactive() {
let now = Instant::now();
let inactive_addrs: Vec<SocketAddr> = CONNECTIONS.iter()
.filter(|conn| conn.is_inactive(Some(now)))
.map(|conn| *conn.peer_addr())
.collect();
if !inactive_addrs.is_empty() { log::debug!("Disconnecting {} inactive connections", inactive_addrs.len()); }
for addr in &inactive_addrs
{
remove_connection(addr, true, Some("inactive")).await;
}
}
pub async fn send_keepalive() {
let addresses: Vec<SocketAddr> = CONNECTIONS.iter()
.filter(|entry| entry.is_authenticated())
.map(|entry| *entry.key())
.collect();
let mut dead_clients = Vec::new();
for addr in addresses
{
let mut stream = None;
let mut keys = None;
if let Some(mut conn) = CONNECTIONS.get_mut(&addr)
{
if !conn.is_alive()
{
dead_clients.push(addr);
continue;
}
stream = Some(conn.write_stream().clone());
keys = conn.keys().cloned();
conn.set_alive(false);
}
if let Some(stream) = stream
{
network::send(&mut *stream.lock().await, PacketCode::KeepAlive, keys.as_ref()).await;
}
}
for dead in dead_clients
{
log::warn!("Missed keepalive: {dead}");
remove_connection(&dead, false, Some("dead")).await;
}
}