From 5916de9a91e937c86115d141267bb5d0380deb4a Mon Sep 17 00:00:00 2001 From: Wind-Explorer Date: Wed, 26 Aug 2026 21:44:53 +0800 Subject: [PATCH] friends connection status system --- src-common/src/lib.rs | 8 + src-server/src/network/interactions.rs | 7 +- src-server/src/network/mod.rs | 212 +++++++++++++++++----- src-server/src/network/presence.rs | 179 ++++++++++++++++++ src-tauri/src/cursor/mod.rs | 34 ++++ src-tauri/src/lib.rs | 2 + src-tauri/src/network/mod.rs | 73 +++++++- src-tauri/src/network/presence.rs | 165 +++++++++++++++++ src/lib/bindings.ts | 6 + src/lib/components/friend-statuses.svelte | 29 +++ src/lib/listeners/friend-statuses.ts | 31 ++++ src/lib/listeners/index.ts | 2 + src/lib/listeners/live-metadata.ts | 11 ++ src/routes/+page.svelte | 4 + src/routes/scene/+page.svelte | 15 +- 15 files changed, 722 insertions(+), 56 deletions(-) create mode 100644 src-server/src/network/presence.rs create mode 100644 src-tauri/src/network/presence.rs create mode 100644 src/lib/components/friend-statuses.svelte create mode 100644 src/lib/listeners/friend-statuses.ts diff --git a/src-common/src/lib.rs b/src-common/src/lib.rs index 03f5d8f..2f65640 100644 --- a/src-common/src/lib.rs +++ b/src-common/src/lib.rs @@ -31,6 +31,7 @@ pub enum ClientMessage { signature: String, }, SyncFriendProfiles, + SyncFriendStatuses, Signed { payload: String, signature: String, @@ -106,6 +107,13 @@ pub enum ServerMessage { FriendProfiles { profiles: Vec, }, + FriendStatusChanged { + friend_id: String, + online: bool, + }, + FriendStatuses { + friend_ids: Vec, + }, FriendLiveData { friend_id: String, payload: String, diff --git a/src-server/src/network/interactions.rs b/src-server/src/network/interactions.rs index 98f3dbe..ee9458b 100644 --- a/src-server/src/network/interactions.rs +++ b/src-server/src/network/interactions.rs @@ -10,7 +10,7 @@ use wyd_common::{ MAX_INTERACTION_PAYLOAD_BYTES, ServerMessage, interaction_bytes, }; -use super::{Clients, verify}; +use super::{Clients, are_mutual_friends, verify}; pub(super) async fn relay( clients: &Clients, @@ -28,10 +28,7 @@ pub(super) async fn relay( .filter(|client| client.connection_id == connection_id)?; let recipient = clients .get(recipient_id) - .filter(|recipient| { - source.friends.iter().any(|friend| friend == recipient_id) - && recipient.friends.iter().any(|friend| friend == public_key) - }) + .filter(|recipient| are_mutual_friends(public_key, source, recipient_id, recipient)) .map(|recipient| recipient.sender.clone()); (source.key, recipient) }; diff --git a/src-server/src/network/mod.rs b/src-server/src/network/mod.rs index 63aa55b..479245b 100644 --- a/src-server/src/network/mod.rs +++ b/src-server/src/network/mod.rs @@ -14,11 +14,11 @@ use futures_util::{SinkExt, StreamExt}; use tokio::sync::{Mutex, mpsc}; use uuid::Uuid; use wyd_common::{ - ClientMessage, Profile, ServerMessage, friends_bytes, message_bytes, profile_bytes, - register_bytes, + ClientMessage, Profile, ServerMessage, message_bytes, profile_bytes, register_bytes, }; mod interactions; +mod presence; type Clients = Arc>>; @@ -31,6 +31,11 @@ struct Client { sender: mpsc::Sender, } +fn are_mutual_friends(source_id: &str, source: &Client, friend_id: &str, friend: &Client) -> bool { + source.friends.iter().any(|id| id == friend_id) + && friend.friends.iter().any(|id| id == source_id) +} + pub fn routes() -> Router { Router::new() .route("/v1/ws", get(upgrade)) @@ -79,20 +84,25 @@ async fn connected(mut socket: WebSocket, clients: Clients) { let connection_id = Uuid::new_v4(); let public_key = profile.id.clone(); let (sender, mut outgoing) = mpsc::channel(32); - clients.lock().await.insert( - public_key.clone(), - Client { - connection_id, - key, - profile, - friends, - sender, - }, - ); + let previous_friends = clients + .lock() + .await + .insert( + public_key.clone(), + Client { + connection_id, + key, + profile, + friends, + sender, + }, + ) + .map(|client| client.friends); if send(&mut socket, &ServerMessage::Registered).await.is_err() { - remove(&clients, &public_key, connection_id).await; + presence::disconnected(&clients, &public_key, connection_id).await; return; } + presence::connected(&clients, &public_key, previous_friends).await; let (mut writer, mut reader) = socket.split(); let mut ping = tokio::time::interval(Duration::from_secs(20)); @@ -143,7 +153,17 @@ async fn connected(mut socket: WebSocket, clients: Clients) { } } Ok(ClientMessage::FriendsUpdated { friends, signature }) => { - if !update_friends(&clients, &public_key, connection_id, friends, &signature).await { + let Some(friend_ids) = presence::update_friends( + &clients, + &public_key, + connection_id, + friends, + &signature, + ).await else { + break; + }; + let message = ServerMessage::FriendStatuses { friend_ids }; + if writer.send(Message::Text(serde_json::to_string(&message).unwrap().into())).await.is_err() { break; } } @@ -156,6 +176,15 @@ async fn connected(mut socket: WebSocket, clients: Clients) { break; } } + Ok(ClientMessage::SyncFriendStatuses) => { + let Some(friend_ids) = presence::snapshot(&clients, &public_key, connection_id).await else { + break; + }; + let message = ServerMessage::FriendStatuses { friend_ids }; + if writer.send(Message::Text(serde_json::to_string(&message).unwrap().into())).await.is_err() { + break; + } + } _ => break, } Some(Ok(Message::Ping(data))) => { @@ -167,7 +196,7 @@ async fn connected(mut socket: WebSocket, clients: Clients) { } } - remove(&clients, &public_key, connection_id).await; + presence::disconnected(&clients, &public_key, connection_id).await; } async fn update_profile( @@ -209,27 +238,6 @@ async fn update_profile( true } -async fn update_friends( - clients: &Clients, - public_key: &str, - connection_id: Uuid, - friends: Vec, - signature: &str, -) -> bool { - let mut clients = clients.lock().await; - let Some(client) = clients - .get_mut(public_key) - .filter(|client| client.connection_id == connection_id) - else { - return false; - }; - if !verify(&client.key, &friends_bytes(&friends), signature) { - return false; - } - client.friends = friends; - true -} - async fn friend_profiles( clients: &Clients, public_key: &str, @@ -270,8 +278,7 @@ async fn relay_live_data( .iter() .filter(|(recipient_id, recipient)| { recipient_id.as_str() != public_key - && source.friends.iter().any(|friend| friend == *recipient_id) - && recipient.friends.iter().any(|friend| friend == public_key) + && are_mutual_friends(public_key, source, recipient_id, recipient) }) .map(|(_, recipient)| recipient.sender.clone()) .collect(); @@ -291,16 +298,6 @@ async fn relay_live_data( true } -async fn remove(clients: &Clients, public_key: &str, connection_id: Uuid) { - let mut clients = clients.lock().await; - if clients - .get(public_key) - .is_some_and(|client| client.connection_id == connection_id) - { - clients.remove(public_key); - } -} - async fn send(socket: &mut WebSocket, message: &ServerMessage) -> Result<(), axum::Error> { socket .send(Message::Text( @@ -328,6 +325,7 @@ fn verify(key: &VerifyingKey, bytes: &[u8], signature: &str) -> bool { #[cfg(test)] mod tests { use ed25519_dalek::{Signer, SigningKey}; + use wyd_common::friends_bytes; use super::*; @@ -492,7 +490,7 @@ mod tests { URL_SAFE_NO_PAD.encode(signing_key.sign(&friends_bytes(&friends)).to_bytes()); assert!( - update_friends( + presence::update_friends( &clients, &public_key, connection_id, @@ -500,6 +498,7 @@ mod tests { &signature, ) .await + .is_some() ); assert_eq!( clients.lock().await.get(&public_key).unwrap().friends, @@ -565,6 +564,108 @@ mod tests { ); } + #[tokio::test] + async fn presence_sync_and_disconnect_broadcast_require_mutual_friendship() { + let signing_key = SigningKey::from_bytes(&[7; 32]); + let connection_id = Uuid::new_v4(); + let (source_sender, _source_receiver) = mpsc::channel(4); + let (mutual_sender, mut mutual_receiver) = mpsc::channel(4); + let (one_way_sender, mut one_way_receiver) = mpsc::channel(4); + let clients = Clients::default(); + let mut registry = clients.lock().await; + registry.insert( + "source".to_owned(), + Client { + connection_id, + key: signing_key.verifying_key(), + profile: Profile { + id: "source".to_owned(), + display_name: "Source".to_owned(), + }, + friends: vec!["mutual".to_owned()], + sender: source_sender, + }, + ); + registry.insert( + "mutual".to_owned(), + Client { + connection_id: Uuid::new_v4(), + key: signing_key.verifying_key(), + profile: Profile { + id: "mutual".to_owned(), + display_name: "Mutual".to_owned(), + }, + friends: vec!["source".to_owned()], + sender: mutual_sender, + }, + ); + registry.insert( + "one-way".to_owned(), + Client { + connection_id: Uuid::new_v4(), + key: signing_key.verifying_key(), + profile: Profile { + id: "one-way".to_owned(), + display_name: "One way".to_owned(), + }, + friends: vec!["source".to_owned()], + sender: one_way_sender, + }, + ); + drop(registry); + + assert_eq!( + presence::snapshot(&clients, "source", connection_id) + .await + .expect("current session"), + ["mutual"] + ); + assert!( + presence::snapshot(&clients, "source", Uuid::new_v4()) + .await + .is_none() + ); + + presence::broadcast_status(&clients, "source", true).await; + assert_friend_status(mutual_receiver.recv().await, "source", true); + assert!(one_way_receiver.try_recv().is_err()); + + let no_friends = Vec::new(); + let signature = + URL_SAFE_NO_PAD.encode(signing_key.sign(&friends_bytes(&no_friends)).to_bytes()); + assert_eq!( + presence::update_friends(&clients, "source", connection_id, no_friends, &signature,) + .await, + Some(Vec::new()) + ); + assert_friend_status(mutual_receiver.recv().await, "source", false); + + let mutual_friends = vec!["mutual".to_owned()]; + let signature = + URL_SAFE_NO_PAD.encode(signing_key.sign(&friends_bytes(&mutual_friends)).to_bytes()); + assert_eq!( + presence::update_friends( + &clients, + "source", + connection_id, + mutual_friends, + &signature, + ) + .await, + Some(vec!["mutual".to_owned()]) + ); + assert_friend_status(mutual_receiver.recv().await, "source", true); + + presence::disconnected(&clients, "source", Uuid::new_v4()).await; + assert!(clients.lock().await.contains_key("source")); + assert!(mutual_receiver.try_recv().is_err()); + + presence::disconnected(&clients, "source", connection_id).await; + assert_friend_status(mutual_receiver.recv().await, "source", false); + assert!(one_way_receiver.try_recv().is_err()); + assert!(!clients.lock().await.contains_key("source")); + } + #[tokio::test] async fn live_data_relay_requires_valid_session_signature_and_mutual_friendship() { let signing_key = SigningKey::from_bytes(&[7; 32]); @@ -669,4 +770,17 @@ mod tests { ); assert!(mutual_receiver.try_recv().is_err()); } + + fn assert_friend_status(message: Option, expected_id: &str, expected_online: bool) { + let Message::Text(message) = message.expect("friend receives status") else { + panic!("expected text status"); + }; + let ServerMessage::FriendStatusChanged { friend_id, online } = + serde_json::from_str(&message).expect("decode friend status") + else { + panic!("expected friend status"); + }; + assert_eq!(friend_id, expected_id); + assert_eq!(online, expected_online); + } } diff --git a/src-server/src/network/presence.rs b/src-server/src/network/presence.rs new file mode 100644 index 0000000..e73d1be --- /dev/null +++ b/src-server/src/network/presence.rs @@ -0,0 +1,179 @@ +use axum::extract::ws::Message; +use tokio::sync::mpsc; +use uuid::Uuid; +use wyd_common::{ServerMessage, friends_bytes}; + +use super::{Clients, are_mutual_friends, verify}; + +pub(super) async fn connected( + clients: &Clients, + public_key: &str, + previous_friends: Option>, +) { + if let Some(previous_friends) = previous_friends { + broadcast_friendship_changes(clients, public_key, &previous_friends).await; + } else { + broadcast_status(clients, public_key, true).await; + } +} + +pub(super) async fn update_friends( + clients: &Clients, + public_key: &str, + connection_id: Uuid, + friends: Vec, + signature: &str, +) -> Option> { + let mut clients = clients.lock().await; + let previous_friends = { + let client = clients + .get_mut(public_key) + .filter(|client| client.connection_id == connection_id)?; + if !verify(&client.key, &friends_bytes(&friends), signature) { + return None; + } + std::mem::replace(&mut client.friends, friends.clone()) + }; + + let mut notifications = Vec::new(); + let mut online_friend_ids = Vec::new(); + for (friend_id, friend) in clients.iter().filter(|(id, _)| id.as_str() != public_key) { + let was_visible = mutual_friends(public_key, &previous_friends, friend_id, &friend.friends); + let is_visible = mutual_friends(public_key, &friends, friend_id, &friend.friends); + if is_visible { + online_friend_ids.push(friend_id.clone()); + } + if was_visible != is_visible { + notifications.push((friend.sender.clone(), public_key.to_owned(), is_visible)); + } + } + drop(clients); + + send_notifications(notifications).await; + online_friend_ids.sort_unstable(); + Some(online_friend_ids) +} + +pub(super) async fn snapshot( + clients: &Clients, + public_key: &str, + connection_id: Uuid, +) -> Option> { + let clients = clients.lock().await; + let source = clients + .get(public_key) + .filter(|client| client.connection_id == connection_id)?; + let mut friend_ids = clients + .iter() + .filter(|(friend_id, friend)| { + friend_id.as_str() != public_key && visible_to(public_key, source, friend_id, friend) + }) + .map(|(friend_id, _)| friend_id.clone()) + .collect::>(); + friend_ids.sort_unstable(); + Some(friend_ids) +} + +pub(super) async fn broadcast_status(clients: &Clients, public_key: &str, online: bool) { + let recipients = { + let clients = clients.lock().await; + let Some(source) = clients.get(public_key) else { + return; + }; + clients + .iter() + .filter(|(friend_id, friend)| { + friend_id.as_str() != public_key + && visible_to(public_key, source, friend_id, friend) + }) + .map(|(_, friend)| friend.sender.clone()) + .collect::>() + }; + let message = status_message(public_key.to_owned(), online); + for recipient in recipients { + let _ = recipient.send(message.clone()).await; + } +} + +pub(super) async fn disconnected(clients: &Clients, public_key: &str, connection_id: Uuid) { + let recipients = { + let mut clients = clients.lock().await; + if !clients + .get(public_key) + .is_some_and(|client| client.connection_id == connection_id) + { + return; + } + let source = clients + .remove(public_key) + .expect("checked registered client"); + clients + .iter() + .filter(|(friend_id, friend)| visible_to(public_key, &source, friend_id, friend)) + .map(|(_, friend)| friend.sender.clone()) + .collect::>() + }; + let message = status_message(public_key.to_owned(), false); + for recipient in recipients { + let _ = recipient.send(message.clone()).await; + } +} + +async fn broadcast_friendship_changes( + clients: &Clients, + public_key: &str, + previous_friends: &[String], +) { + let notifications = { + let clients = clients.lock().await; + let Some(source) = clients.get(public_key) else { + return; + }; + clients + .iter() + .filter(|(friend_id, _)| friend_id.as_str() != public_key) + .filter_map(|(friend_id, friend)| { + let was_visible = + mutual_friends(public_key, previous_friends, friend_id, &friend.friends); + let is_visible = visible_to(public_key, source, friend_id, friend); + (was_visible != is_visible) + .then(|| (friend.sender.clone(), public_key.to_owned(), is_visible)) + }) + .collect::>() + }; + send_notifications(notifications).await; +} + +async fn send_notifications(notifications: Vec<(mpsc::Sender, String, bool)>) { + for (recipient, friend_id, online) in notifications { + let _ = recipient.send(status_message(friend_id, online)).await; + } +} + +fn visible_to( + source_id: &str, + source: &super::Client, + viewer_id: &str, + viewer: &super::Client, +) -> bool { + // Apply the source profile's per-friend online-visibility preference here. + are_mutual_friends(source_id, source, viewer_id, viewer) +} + +fn mutual_friends( + source_id: &str, + source_friends: &[String], + viewer_id: &str, + viewer_friends: &[String], +) -> bool { + source_friends.iter().any(|friend| friend == viewer_id) + && viewer_friends.iter().any(|friend| friend == source_id) +} + +fn status_message(friend_id: String, online: bool) -> Message { + Message::Text( + serde_json::to_string(&ServerMessage::FriendStatusChanged { friend_id, online }) + .expect("friend status serializes") + .into(), + ) +} diff --git a/src-tauri/src/cursor/mod.rs b/src-tauri/src/cursor/mod.rs index 4d65315..fe059c2 100644 --- a/src-tauri/src/cursor/mod.rs +++ b/src-tauri/src/cursor/mod.rs @@ -46,6 +46,16 @@ impl CursorState { positions_by_user.insert(user_id, positions); Ok(positions_by_user.clone()) } + + fn remove( + &self, + user_ids: &[String], + ) -> Result>, String> { + let mut positions_by_user = self.0.write().map_err(|error| error.to_string())?; + let previous_len = positions_by_user.len(); + positions_by_user.retain(|user_id, _| !user_ids.contains(user_id)); + Ok((positions_by_user.len() != previous_len).then(|| positions_by_user.clone())) + } } // Was private, but for some reason LSP @@ -125,6 +135,23 @@ pub(crate) fn emit_position(handle: &AppHandle, user_id: String, positions: Curs } } +pub(crate) fn remove_positions(handle: &AppHandle, user_ids: &[String]) { + if user_ids.is_empty() { + return; + } + let positions = match handle.state::().remove(user_ids) { + Ok(Some(positions)) => positions, + Ok(None) => return, + Err(error) => { + eprintln!("Failed to remove offline cursor positions: {error}"); + return; + } + }; + if let Err(error) = (CursorPositionChanged { positions }).emit(handle) { + eprintln!("Failed to emit cursor position removal: {error}"); + } +} + /// Convert absolute to normalized coordinates (0.12, 0.78), or normalized to absolute (1234, 567) pub fn transform_cursor_pos( pos: &CursorPosition, @@ -181,6 +208,13 @@ mod tests { let snapshot = state.update("local".to_owned(), positions(3.0)).unwrap(); assert_eq!(snapshot.len(), 2); assert_eq!(snapshot["local"].raw.x, 3.0); + + let snapshot = state + .remove(&["friend".to_owned()]) + .unwrap() + .expect("friend cursor removed"); + assert_eq!(snapshot.keys().collect::>(), ["local"]); + assert!(state.remove(&["missing".to_owned()]).unwrap().is_none()); } #[test] diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index a7e940c..5827795 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -71,6 +71,7 @@ fn specta_builder() -> tauri_specta::Builder { profile::update_profile, keypair::get_public_key, network::list_statuses, + network::list_friend_statuses, images::pick_and_send_image, images::send_image_bytes, interactions::send_interaction, @@ -81,6 +82,7 @@ fn specta_builder() -> tauri_specta::Builder { remotes::RemotesChanged, profile::ProfileChanged, network::NetworkStatusChanged, + network::FriendStatusesChanged, cursor::CursorPositionChanged, ufa::ForegroundAppChanged, ufa::FriendForegroundAppChanged, diff --git a/src-tauri/src/network/mod.rs b/src-tauri/src/network/mod.rs index 0332a13..b3e5b4b 100644 --- a/src-tauri/src/network/mod.rs +++ b/src-tauri/src/network/mod.rs @@ -23,6 +23,10 @@ use crate::keypair::AppKeypair; use crate::live_data::LiveData; use crate::remotes::{self, Remote, RemotesChanged}; +mod presence; + +use presence::{Change as FriendPresenceChange, FriendPresence}; + type Statuses = Arc>>; #[derive(Debug, Clone, Serialize, Deserialize, Type)] @@ -48,6 +52,12 @@ pub struct NetworkStatusChanged { pub statuses: Vec, } +#[derive(Debug, Clone, Serialize, Deserialize, Type, Event)] +#[serde(rename_all = "camelCase")] +pub struct FriendStatusesChanged { + pub friend_ids: Vec, +} + struct Connection { remote: Remote, sender: mpsc::Sender, @@ -66,6 +76,7 @@ struct InteractionRequest { pub struct Network { connections: Mutex>, statuses: Statuses, + friend_presence: Arc, profile: watch::Sender, friends: watch::Sender>, keypair: AppKeypair, @@ -173,6 +184,7 @@ impl Network { if let Some(connection) = connections.remove(&id) { connection.task.abort(); remove_status(&self.statuses, &id, connection.generation); + apply_friend_presence_change(handle, self.friend_presence.remove(&id)); } } @@ -193,6 +205,7 @@ impl Network { let task = tauri::async_runtime::spawn(run( handle.clone(), self.statuses.clone(), + self.friend_presence.clone(), remote.clone(), generation, self.profile.subscribe(), @@ -228,6 +241,7 @@ pub async fn init(handle: &AppHandle) -> Result<(), Box> let network = Network { connections: Mutex::new(HashMap::new()), statuses: Statuses::default(), + friend_presence: Arc::default(), profile, friends, keypair, @@ -252,9 +266,10 @@ pub async fn init(handle: &AppHandle) -> Result<(), Box> if *current == ids { return false; } - *current = ids; + *current = ids.clone(); true }); + apply_friend_presence_change(&listener_handle, network.friend_presence.retain(&ids)); }); Ok(()) } @@ -262,6 +277,7 @@ pub async fn init(handle: &AppHandle) -> Result<(), Box> async fn run( handle: AppHandle, statuses: Statuses, + friend_presence: Arc, remote: Remote, generation: u64, mut profiles: watch::Receiver, @@ -281,6 +297,7 @@ async fn run( if let Err(error) = connect( &handle, &statuses, + &friend_presence, &remote, generation, &mut profiles, @@ -293,6 +310,7 @@ async fn run( { eprintln!("remote {} disconnected: {error}", remote.id); } + apply_friend_presence_change(&handle, friend_presence.remove(&remote.id)); changed( &handle, &statuses, @@ -307,6 +325,7 @@ async fn run( async fn connect( handle: &AppHandle, statuses: &Statuses, + friend_presence: &FriendPresence, remote: &Remote, generation: u64, profiles: &mut watch::Receiver, @@ -354,6 +373,7 @@ async fn connect( .send(InteractionDeliveryStatus::Unavailable); } send(&mut writer, &ClientMessage::SyncFriendProfiles).await?; + send(&mut writer, &ClientMessage::SyncFriendStatuses).await?; changed( handle, statuses, @@ -421,6 +441,21 @@ async fn connect( eprintln!("failed to synchronize friend profiles: {error}"); } } + ServerMessage::FriendStatusChanged { friend_id, online } => { + update_friend_presence( + handle, + friend_presence, + &remote.id, + friend_id, + online, + ); + } + ServerMessage::FriendStatuses { friend_ids } => { + apply_friend_presence_change( + handle, + friend_presence.replace(&remote.id, friend_ids), + ); + } ServerMessage::FriendLiveData { friend_id, payload } => { match serde_json::from_str(&payload) { Ok(LiveData::Cursor { positions }) => { @@ -471,6 +506,42 @@ pub fn list_statuses( Ok(statuses) } +#[tauri::command] +#[specta::specta] +pub fn list_friend_statuses(network: State<'_, Network>) -> Result, String> { + network.friend_presence.snapshot() +} + +fn update_friend_presence( + handle: &AppHandle, + presence: &FriendPresence, + remote_id: &str, + friend_id: String, + online: bool, +) { + apply_friend_presence_change(handle, presence.update(remote_id, friend_id, online)); +} + +fn apply_friend_presence_change( + handle: &AppHandle, + result: Result, String>, +) { + match result { + Ok(Some(change)) => { + crate::cursor::remove_positions(handle, &change.went_offline); + if let Err(error) = (FriendStatusesChanged { + friend_ids: change.online, + }) + .emit(handle) + { + eprintln!("failed to emit friend statuses: {error}"); + } + } + Ok(None) => {} + Err(error) => eprintln!("failed to update friend presence: {error}"), + } +} + fn changed( handle: &AppHandle, statuses: &Statuses, diff --git a/src-tauri/src/network/presence.rs b/src-tauri/src/network/presence.rs new file mode 100644 index 0000000..a74c6d1 --- /dev/null +++ b/src-tauri/src/network/presence.rs @@ -0,0 +1,165 @@ +use std::collections::{HashMap, HashSet}; +use std::sync::Mutex; + +#[derive(Debug, PartialEq, Eq)] +pub(super) struct Change { + pub(super) online: Vec, + pub(super) went_offline: Vec, +} + +#[derive(Default)] +pub(super) struct FriendPresence(Mutex>>); + +impl FriendPresence { + pub(super) fn replace( + &self, + remote_id: &str, + friend_ids: Vec, + ) -> Result, String> { + let mut by_remote = self.0.lock().map_err(|error| error.to_string())?; + let before = aggregate(&by_remote); + let friend_ids = friend_ids.into_iter().collect::>(); + if friend_ids.is_empty() { + by_remote.remove(remote_id); + } else { + by_remote.insert(remote_id.to_owned(), friend_ids); + } + Ok(diff(before, &by_remote)) + } + + pub(super) fn update( + &self, + remote_id: &str, + friend_id: String, + online: bool, + ) -> Result, String> { + let mut by_remote = self.0.lock().map_err(|error| error.to_string())?; + let before = aggregate(&by_remote); + if online { + by_remote + .entry(remote_id.to_owned()) + .or_default() + .insert(friend_id); + } else if let Some(friend_ids) = by_remote.get_mut(remote_id) { + friend_ids.remove(&friend_id); + if friend_ids.is_empty() { + by_remote.remove(remote_id); + } + } + Ok(diff(before, &by_remote)) + } + + pub(super) fn remove(&self, remote_id: &str) -> Result, String> { + let mut by_remote = self.0.lock().map_err(|error| error.to_string())?; + let before = aggregate(&by_remote); + by_remote.remove(remote_id); + Ok(diff(before, &by_remote)) + } + + pub(super) fn retain(&self, known_friend_ids: &[String]) -> Result, String> { + let known_friend_ids = known_friend_ids.iter().collect::>(); + let mut by_remote = self.0.lock().map_err(|error| error.to_string())?; + let before = aggregate(&by_remote); + by_remote.retain(|_, friend_ids| { + friend_ids.retain(|friend_id| known_friend_ids.contains(friend_id)); + !friend_ids.is_empty() + }); + Ok(diff(before, &by_remote)) + } + + pub(super) fn snapshot(&self) -> Result, String> { + let by_remote = self.0.lock().map_err(|error| error.to_string())?; + let mut friend_ids = aggregate(&by_remote).into_iter().collect::>(); + friend_ids.sort_unstable(); + Ok(friend_ids) + } +} + +fn aggregate(by_remote: &HashMap>) -> HashSet { + by_remote + .values() + .flat_map(|friend_ids| friend_ids.iter().cloned()) + .collect() +} + +fn diff(before: HashSet, by_remote: &HashMap>) -> Option { + let after = aggregate(by_remote); + if before == after { + return None; + } + let mut went_offline = before.difference(&after).cloned().collect::>(); + let mut online = after.into_iter().collect::>(); + went_offline.sort_unstable(); + online.sort_unstable(); + Some(Change { + online, + went_offline, + }) +} + +#[cfg(test)] +mod tests { + use super::{Change, FriendPresence}; + + #[test] + fn friend_stays_online_while_any_remote_reports_presence() { + let presence = FriendPresence::default(); + assert_eq!( + presence + .replace("remote-a", vec!["friend".to_owned()]) + .unwrap(), + Some(Change { + online: vec!["friend".to_owned()], + went_offline: Vec::new(), + }) + ); + assert!( + presence + .replace("remote-b", vec!["friend".to_owned()]) + .unwrap() + .is_none() + ); + assert!(presence.remove("remote-a").unwrap().is_none()); + assert_eq!( + presence.remove("remote-b").unwrap(), + Some(Change { + online: Vec::new(), + went_offline: vec!["friend".to_owned()], + }) + ); + } + + #[test] + fn incremental_updates_and_friend_removal_share_the_same_snapshot() { + let presence = FriendPresence::default(); + presence + .replace("remote", vec!["kept".to_owned(), "removed".to_owned()]) + .unwrap(); + + assert_eq!( + presence + .update("remote", "added".to_owned(), true) + .unwrap() + .expect("friend came online") + .online, + ["added", "kept", "removed"] + ); + assert_eq!( + presence + .update("remote", "added".to_owned(), false) + .unwrap() + .expect("friend went offline") + .went_offline, + ["added"] + ); + assert_eq!( + presence + .retain(&["kept".to_owned()]) + .unwrap() + .expect("unknown friend was removed") + .went_offline, + ["removed"] + ); + assert_eq!(presence.snapshot().unwrap(), ["kept"]); + } +} diff --git a/src/lib/bindings.ts b/src/lib/bindings.ts index 26fec50..ecdb713 100644 --- a/src/lib/bindings.ts +++ b/src/lib/bindings.ts @@ -44,6 +44,9 @@ async getPublicKey() : Promise { async listStatuses() : Promise { return await TAURI_INVOKE("list_statuses"); }, +async listFriendStatuses() : Promise { + return await TAURI_INVOKE("list_friend_statuses"); +}, async pickAndSendImage(recipientId: string) : Promise { return await TAURI_INVOKE("pick_and_send_image", { recipientId }); }, @@ -66,6 +69,7 @@ cursorPositionChanged: CursorPositionChanged, foregroundAppChanged: ForegroundAppChanged, friendForegroundAppChanged: FriendForegroundAppChanged, friendInteractionReceived: FriendInteractionReceived, +friendStatusesChanged: FriendStatusesChanged, friendsChanged: FriendsChanged, networkStatusChanged: NetworkStatusChanged, profileChanged: ProfileChanged, @@ -75,6 +79,7 @@ cursorPositionChanged: "cursor-position-changed", foregroundAppChanged: "foreground-app-changed", friendForegroundAppChanged: "friend-foreground-app-changed", friendInteractionReceived: "friend-interaction-received", +friendStatusesChanged: "friend-statuses-changed", friendsChanged: "friends-changed", networkStatusChanged: "network-status-changed", profileChanged: "profile-changed", @@ -107,6 +112,7 @@ mapped: CursorPosition } export type ForegroundAppChanged = { meta: AppMeta } export type FriendForegroundAppChanged = { friendId: string; meta: AppMeta } export type FriendInteractionReceived = { interactionId: string; friendId: string; content: InteractionContent } +export type FriendStatusesChanged = { friendIds: string[] } export type FriendsChanged = { friends: User[] } export type InteractionContent = { type: "text"; text: string } | { type: "wave" } | { type: "image"; mediaType: string; data: string } export type NetworkStatusChanged = { statuses: ConnectionStatus[] } diff --git a/src/lib/components/friend-statuses.svelte b/src/lib/components/friend-statuses.svelte new file mode 100644 index 0000000..ce3600a --- /dev/null +++ b/src/lib/components/friend-statuses.svelte @@ -0,0 +1,29 @@ + + +
+

Friend statuses

+ + {#if $friendsListenerError || $friendStatusesListenerError} +

+ {$friendsListenerError || $friendStatusesListenerError} +

+ {:else if $friends.length === 0} +

No friends configured.

+ {:else} +
    + {#each $friends as friend (friend.id)} + {@const online = $onlineFriendIds.has(friend.id)} +
  • + {friend.displayName}: {online ? "Online" : "Offline"} + ({friend.id}) +
  • + {/each} +
+ {/if} +
diff --git a/src/lib/listeners/friend-statuses.ts b/src/lib/listeners/friend-statuses.ts new file mode 100644 index 0000000..4fb29d1 --- /dev/null +++ b/src/lib/listeners/friend-statuses.ts @@ -0,0 +1,31 @@ +import { writable } from "svelte/store"; +import { commands, events } from "$lib/bindings"; +import { removeFriendForegroundApps } from "./live-metadata"; + +export const onlineFriendIds = writable>(new Set()); +export const friendStatusesListenerError = writable(""); + +export async function initFriendStatusesListener() { + let current = new Set(); + + const apply = (friendIds: string[]) => { + const next = new Set(friendIds); + const wentOffline = [...current].filter((friendId) => !next.has(friendId)); + current = next; + onlineFriendIds.set(next); + + removeFriendForegroundApps(wentOffline); + }; + + const unlisten = await events.friendStatusesChanged.listen((event) => { + apply(event.payload.friendIds); + }); + + try { + apply(await commands.listFriendStatuses()); + } catch (error) { + friendStatusesListenerError.set(String(error)); + } + + return unlisten; +} diff --git a/src/lib/listeners/index.ts b/src/lib/listeners/index.ts index 5d90fb7..7b89e25 100644 --- a/src/lib/listeners/index.ts +++ b/src/lib/listeners/index.ts @@ -1,4 +1,5 @@ import { initConnectionStatusesListener } from "./connection-status"; +import { initFriendStatusesListener } from "./friend-statuses"; import { initFriendsListener } from "./friends"; import { initInteractionListener } from "./interactions"; import { initLiveMetadataListeners } from "./live-metadata"; @@ -16,6 +17,7 @@ export function initAppListeners(): Unlisten { initRemotesListener(), initProfileListener(), initConnectionStatusesListener(), + initFriendStatusesListener(), initLiveMetadataListeners(), initInteractionListener(), ]) diff --git a/src/lib/listeners/live-metadata.ts b/src/lib/listeners/live-metadata.ts index e8de1aa..3d51657 100644 --- a/src/lib/listeners/live-metadata.ts +++ b/src/lib/listeners/live-metadata.ts @@ -19,6 +19,17 @@ export const liveMetadata = writable({ }); export const liveMetadataListenerError = writable(""); +export function removeFriendForegroundApps(friendIds: string[]) { + if (friendIds.length === 0) return; + const removed = new Set(friendIds); + liveMetadata.update((metadata) => ({ + ...metadata, + foregroundApps: new Map( + [...metadata.foregroundApps].filter(([userId]) => !removed.has(userId)), + ), + })); +} + export async function initLiveMetadataListeners() { let localId = ""; let pendingLocalForegroundApp: AppMeta | null = null; diff --git a/src/routes/+page.svelte b/src/routes/+page.svelte index e4bbe0b..06ea9fb 100644 --- a/src/routes/+page.svelte +++ b/src/routes/+page.svelte @@ -1,5 +1,6 @@