From 0bce209dd139e8ac3e0f6315c4012872cc3fd74f Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 20 Aug 2026 15:08:55 +0200 Subject: [PATCH] Get Messages --- .gitignore | 1 + src/data/messages.rs | 15 +++++++++++---- src/protocol/message.rs | 26 ++++++++++++++++++++++++++ src/protocol/mod.rs | 25 ++++++++++++++++++++++--- 4 files changed, 60 insertions(+), 7 deletions(-) diff --git a/.gitignore b/.gitignore index 57c39e9..4fc816a 100644 --- a/.gitignore +++ b/.gitignore @@ -27,3 +27,4 @@ target private.key config.json +messages diff --git a/src/data/messages.rs b/src/data/messages.rs index 6c613da..8434097 100644 --- a/src/data/messages.rs +++ b/src/data/messages.rs @@ -87,17 +87,24 @@ impl MessageStore { }) } - /// Fetches the most recent `limit` messages, oldest-first (ready to render top-to-bottom). - pub fn get_recent_messages(&self, channel_id: &str, limit: u32) -> Result> { + /// Fetches the most recent `limit` messages by the `offset`. + pub fn get_recent_messages( + &self, + channel_id: &str, + limit: u32, + chunk: u32, + ) -> Result> { self.with_channel(channel_id, |conn| { + let offset = chunk * limit; + let mut stmt = conn.prepare( "SELECT id, author_pubkey, content, timestamp, signature FROM messages ORDER BY timestamp DESC - LIMIT ?1", + LIMIT ?1 OFFSET ?2", )?; - let rows = stmt.query_map(params![limit], |row| { + let rows = stmt.query_map(params![limit, offset], |row| { Ok(StoredMessage { id: row.get(0)?, author: row.get(1)?, diff --git a/src/protocol/message.rs b/src/protocol/message.rs index b5a8458..0f3b76f 100644 --- a/src/protocol/message.rs +++ b/src/protocol/message.rs @@ -1,7 +1,9 @@ use crate::data::messages::{MessageData, StoredMessage}; +use crate::protocol::{ClientMethod, send_socket}; use crate::server::Server; use axum::extract::ws::WebSocket; use ed25519_dalek::{Verifier, VerifyingKey}; +use std::collections::HashMap; use std::sync::Arc; use std::time::{SystemTime, UNIX_EPOCH}; use tokio::sync::Mutex; @@ -45,3 +47,27 @@ pub async fn send_message( Ok(()) } + +pub async fn get_messages( + server: &Arc, + _verifying_key: VerifyingKey, + socket: &Arc>, + channel_id: String, + chunk: u32, +) -> anyhow::Result<()> { + const CHUNK_SIZE: u32 = 16; + + let messages = server + .message_store + .get_recent_messages(&channel_id, CHUNK_SIZE, chunk)?; + + send_socket( + &mut *socket.lock().await, + &ClientMethod::Messages { + messages: HashMap::from([(channel_id, messages)]), + }, + ) + .await?; + + Ok(()) +} diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index d3eead3..d15c575 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -1,11 +1,11 @@ -use std::{borrow::Cow, sync::Arc}; +use std::{borrow::Cow, collections::HashMap, sync::Arc}; use axum::extract::ws::{Message, Utf8Bytes, WebSocket}; use ed25519_dalek::VerifyingKey; use serde::{Deserialize, Serialize}; use tokio::sync::Mutex; -use crate::{server::Server, types::ClientMeta}; +use crate::{data::messages::StoredMessage, server::Server, types::ClientMeta}; pub mod initialize; pub mod message; @@ -21,6 +21,10 @@ pub enum ClientMethod { hostname: String, }, + Messages { + messages: HashMap>, + }, + Error { error: Cow<'static, str>, }, @@ -42,6 +46,11 @@ pub enum ServerMethod { data: crate::data::messages::MessageData, }, + GetMessages { + channel_id: String, + chunk: u32, + }, + Meta(ClientMeta), Error { @@ -54,7 +63,11 @@ pub async fn read_loop( verifying_key: VerifyingKey, socket: &Arc>, ) -> anyhow::Result<()> { - while let Some(message) = read_socket(&mut *socket.lock().await).await? { + let mut socket_lock = socket.lock().await; + + while let Some(message) = read_socket(&mut *socket_lock).await? { + drop(socket_lock); + match message { ServerMethod::Initialize { .. } => { send_socket( @@ -76,7 +89,13 @@ pub async fn read_loop( ServerMethod::SendMessage { channel_id, data } => { message::send_message(server, verifying_key, socket, data, channel_id).await?; } + + ServerMethod::GetMessages { channel_id, chunk } => { + message::get_messages(server, verifying_key, socket, channel_id, chunk).await?; + } } + + socket_lock = socket.lock().await; } Ok(())