diff --git a/README.md b/README.md index bbe9402..c583d19 100644 --- a/README.md +++ b/README.md @@ -1,8 +1,9 @@ # Things not covered - WebSocket recovery when connection is lost + - Exponential backoff function could solve this with upper limit after which manual refresh is required - User authentication. Currently you just need to know the user name to sign in -- Virtualization of chat message lists. Too many messages will currently have a performance impact +- Virtualization of lists like chat messages. Too many messages will currently have a performance impact - Possible solution: [TanStack Virtual](https://tanstack.com/virtual/latest/docs/introduction) - Message time is controlled by client and not server and client can write in it whatever they want - Removing users when they disappear on the server diff --git a/client/src/context.tsx b/client/src/context.tsx index 362f83f..239c925 100644 --- a/client/src/context.tsx +++ b/client/src/context.tsx @@ -3,11 +3,26 @@ import { Signal, createContext, createEffect, + createResource, createSignal, useContext, } from "solid-js"; import { ChatMessage } from "./routes/Chat"; +// #[derive(Serialize)] +// #[serde(tag = "type")] +// enum ClientMessage { +// ChatMessage{ message: Arc}, +// AddUser { name: Arc}, +// RemoveUser{name: Arc}, +// SynchronizeMessage { message: Arc }, +// } +type Message = + | { type: "ChatMessage"; message: ChatMessage } + | { type: "AddUser"; name: string } + | { type: "RemoveUser"; name: string } + | { type: "SynchronizeMessage"; message: ChatMessage }; + const [name, setName] = createSignal( localStorage.getItem("name") ); @@ -26,15 +41,14 @@ const [socket, setSocket] = createSignal(); // Not sure if using a map is better const messagesByUser = new Map>(); -async function handleMessage(event: MessageEvent) { - if (typeof event.data !== "string") - throw new Error("Message is not a string"); - - console.debug("Received message", event.data); - // For now we just pray it's right 🙂 - const message = JSON.parse(event.data) as ChatMessage; - - let signal = messagesByUser.get(message.sender); +/** + * + * @param message + * @param contact Contact needs to be set separately and is not provided by the message + * because the message could be from the current user but a different client for syncing + */ +function addChatMessage(message: ChatMessage, contact: string) { + let signal = messagesByUser.get(contact); if (signal === undefined) { signal = createSignal([]); @@ -46,7 +60,67 @@ async function handleMessage(event: MessageEvent) { setMessages((messages) => [message, ...messages]); } -const state = { socket, name, setName, messagesByUser }; +const isSecureRequired = + window.location.protocol === "https:" || + import.meta.env.MODE !== "development"; + +export const backendUrl = new URL( + (isSecureRequired ? "https://" : "http://") + + (import.meta.env.VITE_DEV_BACKEND_HOST ?? window.location.host) +); +console.debug("Backend at", backendUrl); +const socketUrl = new URL(backendUrl.href); +socketUrl.protocol = isSecureRequired ? "wss:" : "ws:"; + +async function fetchUsers() { + const response = await fetch(backendUrl + "users"); + const data = await response.json(); + return data; +} + +const [users, { mutate }] = createResource(fetchUsers); + +function addUser(name: string) { + mutate((previous) => { + if (previous === undefined) return [name]; + const index = previous.indexOf(name); + if (index === -1) return [...previous, name]; + return previous; + }); +} + +function removeUser(name: string) { + mutate((previous) => previous?.filter((user) => user !== name)); +} + +async function handleMessage(event: MessageEvent) { + if (typeof event.data !== "string") + throw new Error("Message is not a string"); + + console.debug("Received message", event.data); + // For now we just pray it's the right type 🙂 + const message = JSON.parse(event.data) as Message; + + switch (message.type) { + case "ChatMessage": + // If we get a message from a different user, we need to use the sender to find the chat + addChatMessage(message.message, message.message.sender); + break; + case "AddUser": + addUser(message.name); + break; + case "RemoveUser": + removeUser(message.name); + break; + case "SynchronizeMessage": + // If we get a message from us but a different client, we need to use the find the intended recipient + // of the message to fin the chat partner + addChatMessage(message.message, message.message.recipient); + break; + } +} + +const state = { socket, name, setName, messagesByUser, users }; const Context = createContext(state); // Close socket if name goes to null. Meaningif the user signs out @@ -64,18 +138,6 @@ createEffect>((previous) => { return value; }, name()); -const isSecureRequired = - window.location.protocol === "https:" || - import.meta.env.MODE !== "development"; - -export const backendUrl = new URL( - (isSecureRequired ? "https://" : "http://") + - (import.meta.env.VITE_DEV_BACKEND_HOST ?? window.location.host) -); -console.debug("Backend at", backendUrl); -const socketUrl = new URL(backendUrl.href); -socketUrl.protocol = isSecureRequired ? "wss:" : "ws:"; - // Use new socket if name changes createEffect((previous) => { const id = name(); diff --git a/client/src/routes/Index.tsx b/client/src/routes/Index.tsx index 5c3b58a..f41218f 100644 --- a/client/src/routes/Index.tsx +++ b/client/src/routes/Index.tsx @@ -1,15 +1,10 @@ -import { For, Show, createEffect, createResource } from "solid-js"; -import { backendUrl, useAppContext } from "../context"; +import { For, Show, createEffect } from "solid-js"; +import { useAppContext } from "../context"; import { Navigate, useNavigate } from "@solidjs/router"; import TopAppBar from "../components/TopAppBar"; -async function fetchUsers() { - const response = await fetch(backendUrl + "users"); - const data = await response.json(); - return data; -} export default function Index() { - const { socket, name, setName } = useAppContext(); + const { socket, name, setName, users } = useAppContext(); const navigate = useNavigate(); if (socket() === undefined || name() === null) @@ -20,9 +15,6 @@ export default function Index() { }); // Available chats - - const [users] = createResource(fetchUsers); - const usersWithoutSelf = () => users()?.filter((user) => user !== name()); return ( @@ -51,7 +43,7 @@ export default function Index() { when={usersWithoutSelf() && usersWithoutSelf()!.length > 0} fallback={

No one available to chat

} > -
    +
      {(user) => (
    • diff --git a/client/src/routes/SetUp.tsx b/client/src/routes/SetUp.tsx index 0672d48..3b7382f 100644 --- a/client/src/routes/SetUp.tsx +++ b/client/src/routes/SetUp.tsx @@ -24,9 +24,10 @@ function Alert() { Don't share sensitive information
      +

      Messages are not stored on the server.

      - Anyone using the same name as you can view your messages. Although - they are not permanently stored. + But anyone using the same name as you will also receive your + messages you send and receive.

      diff --git a/server/src/actor/delivery_service.rs b/server/src/actor/delivery_service.rs index f711d89..ff8c51d 100644 --- a/server/src/actor/delivery_service.rs +++ b/server/src/actor/delivery_service.rs @@ -6,7 +6,7 @@ use tokio::sync::{mpsc, oneshot}; enum Message { //TODO handle case where a user actor is not found - SendMessage(ChatMessage), + SendMessage(Arc), GetOrInsertUser(Arc, oneshot::Sender), GetUsers(oneshot::Sender]>>), RemoveUser(Arc), @@ -29,6 +29,19 @@ impl DeliveryService { fn get_handle(&self) -> Handle { self.sender.clone().into() } + + async fn add_available_contact(&self, user_name: Arc) { + // Notifying all users of new users kind of defeats the purpose of a start topology, but it + // will only get send to the client through websockets, so it shouldn't cause a server + // memory problem. Additionally, this is app is only limited in scope anyway. + // Also, this loop can be parallelized. + for user in self.users_by_name.values() { + let result = user.add_contact(user_name.clone()).await; + if let Err(error) = result { + tracing::error!("Error adding sending new contact to user: {}", error); + } + } + } } async fn run_actor(mut actor: DeliveryService) { @@ -51,25 +64,30 @@ async fn run_actor(mut actor: DeliveryService) { Message::GetOrInsertUser(user_name, respond) => { let entry = actor.users_by_name.get(user_name.as_ref()); - let user = entry.cloned().unwrap_or_else(|| { - let user = user::Handle::new(user_name.clone(), actor.get_handle()); - actor.users_by_name.insert(user_name.clone(), user.clone()); - user - }); + let user = match entry.cloned() { + Some(user) => user, + None => { + let user = user::Handle::new(user_name.clone(), actor.get_handle()); + actor.users_by_name.insert(user_name.clone(), user.clone()); + // Add new contact + actor.add_available_contact(user_name.clone()).await; + + user + } + }; let result = respond.send(user.clone()); + if result.is_ok() { - continue; + tracing::error!("Error sending user handle back"); } - tracing::error!("Error sending user handle back"); } Message::GetUsers(responder) => { let users = actor.users_by_name.keys().cloned().collect(); let result = responder.send(users); - if result.is_ok() { - continue; + if result.is_err() { + tracing::error!("Error sending users back"); } - tracing::error!("Error sending users back"); } Message::RemoveUser(name) => { let _ = actor.users_by_name.remove(&name); @@ -134,7 +152,7 @@ impl Handle { Ok(users) } - pub(super) async fn send_message(&self, message: ChatMessage) -> Result<(), impl Error> { + pub(super) async fn send_message(&self, message: Arc) -> Result<(), impl Error> { self.sender.send(Message::SendMessage(message)).await } diff --git a/server/src/actor/user.rs b/server/src/actor/user.rs index c92576a..77975bf 100644 --- a/server/src/actor/user.rs +++ b/server/src/actor/user.rs @@ -15,9 +15,14 @@ struct User { enum Message { AddSocket(websocket::Handle), - ProcessSocketMessage(ChatMessage), - ReceiveMessage(ChatMessage), + ProcessSocketMessage( + SocketId, + Arc, + ), + ReceiveMessage(Arc), RemoveSocket(SocketId), + AddContact(Arc), + RemoveContact(Arc), } async fn run_actor(mut actor: User) { @@ -26,8 +31,21 @@ async fn run_actor(mut actor: User) { Message::AddSocket(socket) => { actor.sockets.push(socket); } - Message::ProcessSocketMessage(message) => { - //TODO send message to all other connected sockets for this user + Message::ProcessSocketMessage(source, message) => { + // Synchronize message to all other connected sockets for this user + for socket in &actor.sockets { + if socket.id == source { + continue; + } + + + tracing::debug!("Syncing message"); + let result = socket.synchronize_message(message.clone()).await; + if let Err(error) = result { + tracing::error!("Error sending message to socket: {}", error); + } + } + // Send message to the user it is intended for through delivery service let result = actor.delivery_service.send_message(message).await; let Err(error) = result else { @@ -37,10 +55,9 @@ async fn run_actor(mut actor: User) { } Message::ReceiveMessage(message) => { // Creating a reference that is easier to clone for each loop iteration - let reference = Arc::new(message); //TODO this can easily be parallelized as it is fire and forget for socket in &actor.sockets { - let result = socket.send_message(reference.clone()).await; + let result = socket.send_message(message.clone()).await; if let Err(error) = result { tracing::error!("Error sending message to socket: {}", error); } @@ -64,6 +81,7 @@ async fn run_actor(mut actor: User) { // If the socket is the last one, we can remove the user from the delivery service if actor.sockets.is_empty() { + tracing::debug!("All sockets closed, removing user from delivery service"); let result = actor.delivery_service.remove_user(actor.name.clone()).await; // Shut down if let Err(error) = result { @@ -73,6 +91,22 @@ async fn run_actor(mut actor: User) { break; } } + Message::AddContact(user_name) => { + for socket in &actor.sockets { + let result = socket.add_contact(user_name.clone()).await; + if let Err(error) = result { + tracing::error!("Error adding user to socket: {}", error); + } + } + } + Message::RemoveContact(user_name) => { + for socket in &actor.sockets { + let result = socket.remove_contact(user_name.clone()).await; + if let Err(error) = result { + tracing::error!("Error removing user from socket: {}", error); + } + } + } } } } @@ -102,24 +136,33 @@ impl Handle { pub(crate) async fn add_socket( &self, socket: websocket::Handle, - ) -> Result<(), impl std::error::Error> { + ) -> Result<(), impl Error> { self.sender.send(Message::AddSocket(socket)).await } pub(super) async fn process_socket_message( &self, - message: ChatMessage, + source: SocketId, + message: Arc, ) -> Result<(), impl Error + Send + Sync> { self.sender - .send(Message::ProcessSocketMessage(message)) + .send(Message::ProcessSocketMessage(source, message)) .await } - pub(super) async fn receive_message(&self, message: ChatMessage) -> Result<(), impl Error> { + pub(super) async fn receive_message(&self, message: Arc) -> Result<(), impl Error> { self.sender.send(Message::ReceiveMessage(message)).await } pub(super) async fn remove_socket(&self, socket_id: SocketId) -> Result<(), impl Error> { self.sender.send(Message::RemoveSocket(socket_id)).await } + + pub(super) async fn add_contact(&self, user_name: Arc) -> Result<(), impl Error> { + self.sender.send(Message::AddContact(user_name)).await + } + + pub(super) async fn remove_contact(&self, user_name: Arc) -> Result<(), impl Error> { + self.sender.send(Message::RemoveContact(user_name)).await + } } diff --git a/server/src/actor/websocket.rs b/server/src/actor/websocket.rs index 056c46f..2b07b4a 100644 --- a/server/src/actor/websocket.rs +++ b/server/src/actor/websocket.rs @@ -2,12 +2,25 @@ use ::axum::extract::ws::Message as WebSocketMessage; use axum::extract::ws as axum; use nanoid::nanoid; use std::sync::Arc; +use serde::Serialize; use tokio::sync::mpsc; use super::{user, ChatMessage}; enum Message { SendMessage(Arc), + AddContact { name: Arc }, + RemoveContact { name: Arc }, + SynchronizeMessage { message: Arc }, +} + +#[derive(Serialize)] +#[serde(tag = "type")] +enum ClientMessage { + ChatMessage { message: Arc }, + AddUser { name: Arc }, + RemoveUser { name: Arc }, + SynchronizeMessage { message: Arc }, } #[derive(Clone, PartialEq, Eq)] @@ -21,7 +34,7 @@ struct WebSocket { user: user::Handle, } -async fn process_socket_message(user: &user::Handle, json: String) { +async fn process_socket_message(id: SocketId, user: &user::Handle, json: String) { //TODO validate message sender name and send back error if not let result = serde_json::from_str::(&json); let message = match result { @@ -32,27 +45,44 @@ async fn process_socket_message(user: &user::Handle, json: String) { } }; - let result = user.process_socket_message(message).await; + let result = user.process_socket_message(id, message.into()).await; if let Err(error) = result { tracing::error!("Error processing message: {:?}", error); } } +async fn send_to_socket(socket: &mut axum::WebSocket, message: ClientMessage) { + let json = serde_json::to_string(&message); + let result = match json { + Ok(json) => socket.send(WebSocketMessage::Text(json)).await, + Err(error) => { + tracing::error!("Error serializing message: {:?}", error); + return; + } + }; + + if let Err(error) = result { + tracing::error!("Error sending message through websocket: {:?}", error); + } +} + async fn process_actor_message(socket: &mut axum::WebSocket, message: Message) { match message { Message::SendMessage(message) => { - let json = serde_json::to_string(&message); - let result = match json { - Ok(json) => socket.send(WebSocketMessage::Text(json)).await, - Err(error) => { - tracing::error!("Error serializing message: {:?}", error); - return; - } - }; - - if let Err(error) = result { - tracing::error!("Error sending message through websocket: {:?}", error); - } + let message = ClientMessage::ChatMessage { message }; + send_to_socket(socket, message).await; + } + Message::AddContact { name } => { + let message = ClientMessage::AddUser { name }; + send_to_socket(socket, message).await; + } + Message::RemoveContact { name } => { + let message = ClientMessage::RemoveUser { name }; + send_to_socket(socket, message).await; + } + Message::SynchronizeMessage { message } => { + let message = ClientMessage::SynchronizeMessage { message }; + send_to_socket(socket, message).await; } } } @@ -63,8 +93,6 @@ async fn run_actor(mut actor: WebSocket) { Some(message) = actor.receiver.recv() => process_actor_message(&mut actor.socket, message).await, // Stop actor on error Some(Ok(message)) = actor.socket.recv() => { - tracing::info!("Received message: {:?}", message); - match message { WebSocketMessage::Close(_) => { //TODO remove socket from user to clean up or we have a memory leak @@ -74,7 +102,7 @@ async fn run_actor(mut actor: WebSocket) { tracing::error!("Error removing socket from user: {:?}", error); break; }, - WebSocketMessage::Text(text) => process_socket_message(&actor.user, text.into()).await, + WebSocketMessage::Text(text) => process_socket_message(actor.id.clone(), &actor.user, text.into()).await, other => tracing::error!("Unexpected message type: {:?}", other), } }, @@ -114,4 +142,16 @@ impl Handle { ) -> Result<(), impl std::error::Error> { self.sender.send(Message::SendMessage(message)).await } + + pub(super) async fn add_contact(&self, name: Arc) -> Result<(), impl std::error::Error> { + self.sender.send(Message::AddContact { name }).await + } + + pub(super) async fn remove_contact(&self, name: Arc) -> Result<(), impl std::error::Error> { + self.sender.send(Message::RemoveContact { name }).await + } + + pub(super) async fn synchronize_message(&self, message: Arc) -> Result<(), impl std::error::Error> { + self.sender.send(Message::SynchronizeMessage { message }).await + } }