Merge pull request #4 from recurse-chat/working-messages
Working messages
This commit is contained in:
@@ -27,3 +27,4 @@ target
|
|||||||
|
|
||||||
private.key
|
private.key
|
||||||
config.json
|
config.json
|
||||||
|
messages
|
||||||
|
|||||||
Generated
+200
-1
@@ -2,6 +2,18 @@
|
|||||||
# It is not intended for manual editing.
|
# It is not intended for manual editing.
|
||||||
version = 4
|
version = 4
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "ahash"
|
||||||
|
version = "0.8.12"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "5a15f179cd60c4584b8a8c596927aadc462e27f2ca70c04e0071964a73ba7a75"
|
||||||
|
dependencies = [
|
||||||
|
"cfg-if",
|
||||||
|
"once_cell",
|
||||||
|
"version_check",
|
||||||
|
"zerocopy",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "anyhow"
|
name = "anyhow"
|
||||||
version = "1.0.104"
|
version = "1.0.104"
|
||||||
@@ -105,12 +117,28 @@ dependencies = [
|
|||||||
"tinyvec",
|
"tinyvec",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "bumpalo"
|
||||||
|
version = "3.20.3"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "72f5acc6cb2ba439de613abc23857ec3d78374d8ed5ac84e9d11336e87da8649"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "bytes"
|
name = "bytes"
|
||||||
version = "1.12.1"
|
version = "1.12.1"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04"
|
checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "cc"
|
||||||
|
version = "1.4.3"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "509591b7bcd67f4ef775afad7662703b4935daaa6ec0e5605cfb1090b32a2b6d"
|
||||||
|
dependencies = [
|
||||||
|
"find-msvc-tools",
|
||||||
|
"shlex",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "cfg-if"
|
name = "cfg-if"
|
||||||
version = "1.0.4"
|
version = "1.0.4"
|
||||||
@@ -229,18 +257,38 @@ dependencies = [
|
|||||||
"bs58",
|
"bs58",
|
||||||
"ed25519-dalek",
|
"ed25519-dalek",
|
||||||
"rand 0.8.7",
|
"rand 0.8.7",
|
||||||
|
"rusqlite",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
"tokio",
|
"tokio",
|
||||||
"tower-http",
|
"tower-http",
|
||||||
|
"uuid",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "fallible-iterator"
|
||||||
|
version = "0.3.0"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "2acce4a10f12dc2fb14a218589d4f1f62ef011b2d0cc4b3cb1bba8e94da14649"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "fallible-streaming-iterator"
|
||||||
|
version = "0.1.9"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "7360491ce676a36bf9bb3c56c1aa791658183a54d2744120f27285738d90465a"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "fiat-crypto"
|
name = "fiat-crypto"
|
||||||
version = "0.2.9"
|
version = "0.2.9"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "28dea519a9695b9977216879a3ebfddf92f1c08c05d984f8996aecd6ecdc811d"
|
checksum = "28dea519a9695b9977216879a3ebfddf92f1c08c05d984f8996aecd6ecdc811d"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "find-msvc-tools"
|
||||||
|
version = "0.1.11"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "d45db016d36b838f563236e9193d0ee6ce38f3f68b6c94e914b4929c96bbb890"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "form_urlencoded"
|
name = "form_urlencoded"
|
||||||
version = "1.2.2"
|
version = "1.2.2"
|
||||||
@@ -319,10 +367,39 @@ checksum = "899def5c37c4fd7b2664648c28120ecec138e4d395b459e5ca34f9cce2dd77fd"
|
|||||||
dependencies = [
|
dependencies = [
|
||||||
"cfg-if",
|
"cfg-if",
|
||||||
"libc",
|
"libc",
|
||||||
"r-efi",
|
"r-efi 5.3.0",
|
||||||
"wasip2",
|
"wasip2",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "getrandom"
|
||||||
|
version = "0.4.3"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "300e883d756b2e4ec94e02791f39b04b522276138852cfc41d9fb7e904106099"
|
||||||
|
dependencies = [
|
||||||
|
"cfg-if",
|
||||||
|
"libc",
|
||||||
|
"r-efi 6.0.0",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "hashbrown"
|
||||||
|
version = "0.14.5"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1"
|
||||||
|
dependencies = [
|
||||||
|
"ahash",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "hashlink"
|
||||||
|
version = "0.9.1"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "6ba4ff7128dee98c7dc9794b6a411377e1404dba1c97deb8d1a55297bd25d8af"
|
||||||
|
dependencies = [
|
||||||
|
"hashbrown",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "http"
|
name = "http"
|
||||||
version = "1.5.0"
|
version = "1.5.0"
|
||||||
@@ -415,12 +492,34 @@ version = "1.0.18"
|
|||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682"
|
checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "js-sys"
|
||||||
|
version = "0.3.104"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "0e0c1080212aad755ea003d18543e8768dd432c48819efd73a7bf1e39b7a5a3a"
|
||||||
|
dependencies = [
|
||||||
|
"cfg-if",
|
||||||
|
"futures-util",
|
||||||
|
"wasm-bindgen",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "libc"
|
name = "libc"
|
||||||
version = "0.2.189"
|
version = "0.2.189"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "3eaf3ede3fee6db1a4c2ee091bf8a8b4dccdc6d17f656fb07896ee72867612f2"
|
checksum = "3eaf3ede3fee6db1a4c2ee091bf8a8b4dccdc6d17f656fb07896ee72867612f2"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "libsqlite3-sys"
|
||||||
|
version = "0.28.0"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "0c10584274047cb335c23d3e61bcef8e323adae7c5c8c760540f73610177fc3f"
|
||||||
|
dependencies = [
|
||||||
|
"cc",
|
||||||
|
"pkg-config",
|
||||||
|
"vcpkg",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "log"
|
name = "log"
|
||||||
version = "0.4.33"
|
version = "0.4.33"
|
||||||
@@ -494,6 +593,12 @@ dependencies = [
|
|||||||
"spki",
|
"spki",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "pkg-config"
|
||||||
|
version = "0.3.34"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "f6b464fbc74e149a392436b17d523f769e057cb6877f6a5c4618bc6f11800548"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "ppv-lite86"
|
name = "ppv-lite86"
|
||||||
version = "0.2.21"
|
version = "0.2.21"
|
||||||
@@ -527,6 +632,12 @@ version = "5.3.0"
|
|||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f"
|
checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "r-efi"
|
||||||
|
version = "6.0.0"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "rand"
|
name = "rand"
|
||||||
version = "0.8.7"
|
version = "0.8.7"
|
||||||
@@ -586,6 +697,20 @@ dependencies = [
|
|||||||
"getrandom 0.3.4",
|
"getrandom 0.3.4",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "rusqlite"
|
||||||
|
version = "0.31.0"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "b838eba278d213a8beaf485bd313fd580ca4505a00d5871caeb1457c55322cae"
|
||||||
|
dependencies = [
|
||||||
|
"bitflags",
|
||||||
|
"fallible-iterator",
|
||||||
|
"fallible-streaming-iterator",
|
||||||
|
"hashlink",
|
||||||
|
"libsqlite3-sys",
|
||||||
|
"smallvec",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "rustc_version"
|
name = "rustc_version"
|
||||||
version = "0.4.1"
|
version = "0.4.1"
|
||||||
@@ -595,6 +720,12 @@ dependencies = [
|
|||||||
"semver",
|
"semver",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "rustversion"
|
||||||
|
version = "1.0.23"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "cf54715a573b99ac80df0bc206da022bcd442c974952c7b9720069370852e21f"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "ryu"
|
name = "ryu"
|
||||||
version = "1.0.23"
|
version = "1.0.23"
|
||||||
@@ -695,6 +826,12 @@ dependencies = [
|
|||||||
"digest",
|
"digest",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "shlex"
|
||||||
|
version = "2.0.1"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "f8fadd59c855ef2080decdef8ff161eb6661b86933c9d82e5ba29dc602a55aba"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "signature"
|
name = "signature"
|
||||||
version = "2.2.0"
|
version = "2.2.0"
|
||||||
@@ -963,6 +1100,23 @@ version = "1.0.24"
|
|||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75"
|
checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "uuid"
|
||||||
|
version = "1.24.1"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "2cefc03fd367c0c6d4305de1b312cf00248c4114f4a0418ce6a6af769e3b0bd9"
|
||||||
|
dependencies = [
|
||||||
|
"getrandom 0.4.3",
|
||||||
|
"js-sys",
|
||||||
|
"wasm-bindgen",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "vcpkg"
|
||||||
|
version = "0.2.15"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "version_check"
|
name = "version_check"
|
||||||
version = "0.9.5"
|
version = "0.9.5"
|
||||||
@@ -984,6 +1138,51 @@ dependencies = [
|
|||||||
"wit-bindgen",
|
"wit-bindgen",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "wasm-bindgen"
|
||||||
|
version = "0.2.127"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "1b70935747edd64d89de3efa29d73789b806c15798f8e7dca4d8ac356b50ce70"
|
||||||
|
dependencies = [
|
||||||
|
"cfg-if",
|
||||||
|
"once_cell",
|
||||||
|
"rustversion",
|
||||||
|
"wasm-bindgen-macro",
|
||||||
|
"wasm-bindgen-shared",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "wasm-bindgen-macro"
|
||||||
|
version = "0.2.127"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "77775f8f3f7217702089053b94958f8f54061a3f663417df76e19cbdcca29bc1"
|
||||||
|
dependencies = [
|
||||||
|
"quote",
|
||||||
|
"wasm-bindgen-macro-support",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "wasm-bindgen-macro-support"
|
||||||
|
version = "0.2.127"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "e11d33f857dc2fb11b8bc75aee111aa9cbeb12cd9f25efd3d4c2a3dd4e235284"
|
||||||
|
dependencies = [
|
||||||
|
"bumpalo",
|
||||||
|
"proc-macro2",
|
||||||
|
"quote",
|
||||||
|
"syn 2.0.119",
|
||||||
|
"wasm-bindgen-shared",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "wasm-bindgen-shared"
|
||||||
|
version = "0.2.127"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "7ef64dbcc55df09c7e5a46182d181c2cfa3e925f3da937ea764728b4bbb9dcbf"
|
||||||
|
dependencies = [
|
||||||
|
"unicode-ident",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "windows-link"
|
name = "windows-link"
|
||||||
version = "0.2.1"
|
version = "0.2.1"
|
||||||
|
|||||||
@@ -13,3 +13,5 @@ ed25519-dalek = { version = "2", features = ["rand_core"] }
|
|||||||
rand = "0.8"
|
rand = "0.8"
|
||||||
bs58 = "0.5.1"
|
bs58 = "0.5.1"
|
||||||
tower-http = { version = "0.7.0", features = ["fs", "cors"] }
|
tower-http = { version = "0.7.0", features = ["fs", "cors"] }
|
||||||
|
rusqlite = { version = "0.31", features = ["bundled"] }
|
||||||
|
uuid = { version = "1.24.1", features = ["v4"] }
|
||||||
|
|||||||
@@ -0,0 +1,124 @@
|
|||||||
|
use anyhow::Result;
|
||||||
|
use rusqlite::{Connection, params};
|
||||||
|
use serde::{Deserialize, Serialize};
|
||||||
|
use std::collections::HashMap;
|
||||||
|
use std::path::PathBuf;
|
||||||
|
use std::sync::Mutex;
|
||||||
|
|
||||||
|
pub struct MessageStore {
|
||||||
|
data_dir: PathBuf,
|
||||||
|
connections: Mutex<HashMap<String, Connection>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
|
pub struct MessageData {
|
||||||
|
pub content: String,
|
||||||
|
pub timestamp: u64,
|
||||||
|
pub signature: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
|
pub struct StoredMessage {
|
||||||
|
pub id: String,
|
||||||
|
pub author: String,
|
||||||
|
#[serde(flatten)]
|
||||||
|
pub data: MessageData,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl MessageStore {
|
||||||
|
pub fn new(data_dir: PathBuf) -> Result<Self> {
|
||||||
|
std::fs::create_dir_all(&data_dir)?;
|
||||||
|
Ok(Self {
|
||||||
|
data_dir,
|
||||||
|
connections: Mutex::new(HashMap::new()),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Opens (or reuses an already-open) connection for a channel,
|
||||||
|
/// creating the schema if this is the first time.
|
||||||
|
fn with_channel<T>(
|
||||||
|
&self,
|
||||||
|
channel_id: &str,
|
||||||
|
f: impl FnOnce(&Connection) -> Result<T>,
|
||||||
|
) -> Result<T> {
|
||||||
|
let mut conns = self.connections.lock().unwrap();
|
||||||
|
|
||||||
|
if !conns.contains_key(channel_id) {
|
||||||
|
let path = self.data_dir.join(format!("{channel_id}.db"));
|
||||||
|
let conn = Connection::open(&path)?;
|
||||||
|
|
||||||
|
conn.execute(
|
||||||
|
"CREATE TABLE IF NOT EXISTS messages (
|
||||||
|
id TEXT PRIMARY KEY,
|
||||||
|
author_pubkey TEXT NOT NULL,
|
||||||
|
content TEXT NOT NULL,
|
||||||
|
timestamp INTEGER NOT NULL,
|
||||||
|
signature TEXT NOT NULL
|
||||||
|
)",
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
|
||||||
|
conn.execute(
|
||||||
|
"CREATE INDEX IF NOT EXISTS idx_messages_timestamp ON messages(timestamp)",
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
|
||||||
|
conns.insert(channel_id.to_string(), conn);
|
||||||
|
}
|
||||||
|
|
||||||
|
let conn = conns.get(channel_id).unwrap();
|
||||||
|
f(conn)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn insert_message(&self, channel_id: &str, msg: &StoredMessage) -> Result<()> {
|
||||||
|
self.with_channel(channel_id, |conn| {
|
||||||
|
conn.execute(
|
||||||
|
"INSERT INTO messages (id, author_pubkey, content, timestamp, signature)
|
||||||
|
VALUES (?1, ?2, ?3, ?4, ?5)",
|
||||||
|
params![
|
||||||
|
msg.id,
|
||||||
|
msg.author,
|
||||||
|
msg.data.content,
|
||||||
|
msg.data.timestamp,
|
||||||
|
msg.data.signature,
|
||||||
|
],
|
||||||
|
)?;
|
||||||
|
Ok(())
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Fetches the most recent `limit` messages by the `offset`.
|
||||||
|
pub fn get_recent_messages(
|
||||||
|
&self,
|
||||||
|
channel_id: &str,
|
||||||
|
limit: u32,
|
||||||
|
chunk: u32,
|
||||||
|
) -> Result<Vec<StoredMessage>> {
|
||||||
|
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 OFFSET ?2",
|
||||||
|
)?;
|
||||||
|
|
||||||
|
let rows = stmt.query_map(params![limit, offset], |row| {
|
||||||
|
Ok(StoredMessage {
|
||||||
|
id: row.get(0)?,
|
||||||
|
author: row.get(1)?,
|
||||||
|
data: MessageData {
|
||||||
|
content: row.get(2)?,
|
||||||
|
timestamp: row.get(3)?,
|
||||||
|
signature: row.get(4)?,
|
||||||
|
},
|
||||||
|
})
|
||||||
|
})?;
|
||||||
|
|
||||||
|
let mut messages: Vec<StoredMessage> = rows.collect::<Result<_, _>>()?;
|
||||||
|
messages.reverse(); // DESC query, then flip to oldest-first for display
|
||||||
|
Ok(messages)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,2 @@
|
|||||||
|
pub mod config;
|
||||||
|
pub mod messages;
|
||||||
+1
-1
@@ -1,4 +1,4 @@
|
|||||||
pub mod config;
|
pub mod data;
|
||||||
pub mod protocol;
|
pub mod protocol;
|
||||||
pub mod server;
|
pub mod server;
|
||||||
pub mod signature;
|
pub mod signature;
|
||||||
|
|||||||
@@ -0,0 +1,79 @@
|
|||||||
|
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;
|
||||||
|
|
||||||
|
pub async fn send_message(
|
||||||
|
server: &Arc<Server>,
|
||||||
|
verifying_key: VerifyingKey,
|
||||||
|
_socket: &Arc<Mutex<WebSocket>>,
|
||||||
|
message: MessageData,
|
||||||
|
channel_id: String,
|
||||||
|
) -> anyhow::Result<()> {
|
||||||
|
let server_timestamp = SystemTime::now().duration_since(UNIX_EPOCH)?.as_millis() as u64;
|
||||||
|
|
||||||
|
if server_timestamp.saturating_sub(message.timestamp) > 2000 {
|
||||||
|
anyhow::bail!(
|
||||||
|
"Message timestamp out of range ({server_timestamp} - {}) < 2000",
|
||||||
|
message.timestamp
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
let server_pubkey_string = crate::signature::to_string(&server.key.verifying_key());
|
||||||
|
let signed_string = format!(
|
||||||
|
"{}@{}@{}",
|
||||||
|
message.timestamp, server_pubkey_string, message.content
|
||||||
|
);
|
||||||
|
|
||||||
|
let signature = crate::signature::from_string_sig(&message.signature)
|
||||||
|
.map_err(|_| anyhow::anyhow!("Invalid signature encoding"))?;
|
||||||
|
|
||||||
|
verifying_key
|
||||||
|
.verify(signed_string.as_bytes(), &signature)
|
||||||
|
.map_err(|_| anyhow::anyhow!("Signature verification failed"))?;
|
||||||
|
|
||||||
|
let stored = StoredMessage {
|
||||||
|
id: uuid::Uuid::new_v4().to_string(),
|
||||||
|
author: crate::signature::to_string(&verifying_key),
|
||||||
|
data: message,
|
||||||
|
};
|
||||||
|
|
||||||
|
server.message_store.insert_message(&channel_id, &stored)?;
|
||||||
|
|
||||||
|
server
|
||||||
|
.broadcast(&ClientMethod::Messages {
|
||||||
|
messages: HashMap::from([(channel_id, vec![stored])]),
|
||||||
|
})
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn get_messages(
|
||||||
|
server: &Arc<Server>,
|
||||||
|
_verifying_key: VerifyingKey,
|
||||||
|
socket: &Arc<Mutex<WebSocket>>,
|
||||||
|
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(())
|
||||||
|
}
|
||||||
+38
-4
@@ -1,12 +1,14 @@
|
|||||||
use std::{borrow::Cow, sync::Arc};
|
use std::{borrow::Cow, collections::HashMap, sync::Arc};
|
||||||
|
|
||||||
use axum::extract::ws::{Message, Utf8Bytes, WebSocket};
|
use axum::extract::ws::{Message, Utf8Bytes, WebSocket};
|
||||||
|
use ed25519_dalek::VerifyingKey;
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use tokio::sync::Mutex;
|
use tokio::sync::Mutex;
|
||||||
|
|
||||||
use crate::types::ClientMeta;
|
use crate::{data::messages::StoredMessage, server::Server, types::ClientMeta};
|
||||||
|
|
||||||
pub mod initialize;
|
pub mod initialize;
|
||||||
|
pub mod message;
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
#[serde(tag = "method")]
|
#[serde(tag = "method")]
|
||||||
@@ -19,6 +21,10 @@ pub enum ClientMethod {
|
|||||||
hostname: String,
|
hostname: String,
|
||||||
},
|
},
|
||||||
|
|
||||||
|
Messages {
|
||||||
|
messages: HashMap<String, Vec<StoredMessage>>,
|
||||||
|
},
|
||||||
|
|
||||||
Error {
|
Error {
|
||||||
error: Cow<'static, str>,
|
error: Cow<'static, str>,
|
||||||
},
|
},
|
||||||
@@ -35,6 +41,16 @@ pub enum ServerMethod {
|
|||||||
hostname: String,
|
hostname: String,
|
||||||
},
|
},
|
||||||
|
|
||||||
|
SendMessage {
|
||||||
|
channel_id: String,
|
||||||
|
data: crate::data::messages::MessageData,
|
||||||
|
},
|
||||||
|
|
||||||
|
GetMessages {
|
||||||
|
channel_id: String,
|
||||||
|
chunk: u32,
|
||||||
|
},
|
||||||
|
|
||||||
Meta(ClientMeta),
|
Meta(ClientMeta),
|
||||||
|
|
||||||
Error {
|
Error {
|
||||||
@@ -42,8 +58,16 @@ pub enum ServerMethod {
|
|||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn read_loop(socket: &Arc<Mutex<WebSocket>>) -> anyhow::Result<()> {
|
pub async fn read_loop(
|
||||||
while let Some(message) = read_socket(&mut *socket.lock().await).await? {
|
server: &Arc<Server>,
|
||||||
|
verifying_key: VerifyingKey,
|
||||||
|
socket: &Arc<Mutex<WebSocket>>,
|
||||||
|
) -> anyhow::Result<()> {
|
||||||
|
let mut socket_lock = socket.lock().await;
|
||||||
|
|
||||||
|
while let Some(message) = read_socket(&mut *socket_lock).await? {
|
||||||
|
drop(socket_lock);
|
||||||
|
|
||||||
match message {
|
match message {
|
||||||
ServerMethod::Initialize { .. } => {
|
ServerMethod::Initialize { .. } => {
|
||||||
send_socket(
|
send_socket(
|
||||||
@@ -61,7 +85,17 @@ pub async fn read_loop(socket: &Arc<Mutex<WebSocket>>) -> anyhow::Result<()> {
|
|||||||
ServerMethod::Error { error } => {
|
ServerMethod::Error { error } => {
|
||||||
eprintln!("Client error: {error}");
|
eprintln!("Client error: {error}");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
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(())
|
Ok(())
|
||||||
|
|||||||
+51
-15
@@ -1,5 +1,6 @@
|
|||||||
use std::{
|
use std::{
|
||||||
collections::HashMap,
|
collections::HashMap,
|
||||||
|
path::PathBuf,
|
||||||
sync::{Arc, atomic::AtomicU16},
|
sync::{Arc, atomic::AtomicU16},
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -8,10 +9,10 @@ use axum::{
|
|||||||
response::Response,
|
response::Response,
|
||||||
};
|
};
|
||||||
use ed25519_dalek::{SigningKey, VerifyingKey};
|
use ed25519_dalek::{SigningKey, VerifyingKey};
|
||||||
use tokio::sync::Mutex;
|
use tokio::{sync::Mutex, task::JoinSet};
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
config::Config,
|
data::{config::Config, messages::MessageStore},
|
||||||
protocol::{ClientMethod, read_loop, send_socket},
|
protocol::{ClientMethod, read_loop, send_socket},
|
||||||
types::ClientMeta,
|
types::ClientMeta,
|
||||||
};
|
};
|
||||||
@@ -20,13 +21,14 @@ pub struct UserConnections {
|
|||||||
pub meta: ClientMeta,
|
pub meta: ClientMeta,
|
||||||
pub counter: AtomicU16,
|
pub counter: AtomicU16,
|
||||||
pub public_key: VerifyingKey,
|
pub public_key: VerifyingKey,
|
||||||
pub connections: HashMap<u16, Arc<Mutex<WebSocket>>>,
|
pub connections: Mutex<HashMap<u16, Arc<Mutex<WebSocket>>>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
pub struct Server {
|
pub struct Server {
|
||||||
pub key: SigningKey,
|
pub key: SigningKey,
|
||||||
pub config: Config,
|
pub config: Config,
|
||||||
pub clients: Mutex<HashMap<VerifyingKey, UserConnections>>,
|
pub clients: Mutex<HashMap<VerifyingKey, Arc<UserConnections>>>,
|
||||||
|
pub message_store: MessageStore,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Server {
|
impl Server {
|
||||||
@@ -35,6 +37,7 @@ impl Server {
|
|||||||
key: crate::signature::get().await?,
|
key: crate::signature::get().await?,
|
||||||
config: Config::get().await?,
|
config: Config::get().await?,
|
||||||
clients: Mutex::new(HashMap::new()),
|
clients: Mutex::new(HashMap::new()),
|
||||||
|
message_store: MessageStore::new(PathBuf::from("messages"))?,
|
||||||
}))
|
}))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -50,31 +53,43 @@ impl Server {
|
|||||||
|
|
||||||
let mut clients_meta = s.clients.lock().await;
|
let mut clients_meta = s.clients.lock().await;
|
||||||
|
|
||||||
let clients =
|
let clients = clients_meta
|
||||||
clients_meta
|
|
||||||
.entry(public_key)
|
.entry(public_key)
|
||||||
.or_insert_with(|| UserConnections {
|
.or_insert_with(|| {
|
||||||
|
Arc::new(UserConnections {
|
||||||
meta,
|
meta,
|
||||||
public_key: public_key,
|
public_key: public_key,
|
||||||
counter: AtomicU16::new(0),
|
counter: AtomicU16::new(0),
|
||||||
connections: HashMap::new(),
|
connections: Mutex::new(HashMap::new()),
|
||||||
});
|
})
|
||||||
|
})
|
||||||
|
.clone();
|
||||||
|
|
||||||
let conid = clients
|
let conid = clients
|
||||||
.counter
|
.counter
|
||||||
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
||||||
|
|
||||||
clients.connections.insert(conid, client.clone());
|
clients
|
||||||
|
.connections
|
||||||
|
.lock()
|
||||||
|
.await
|
||||||
|
.insert(conid, client.clone());
|
||||||
|
|
||||||
if let Err(e) = read_loop(&client).await {
|
drop(clients_meta);
|
||||||
|
|
||||||
|
if let Err(e) = read_loop(&s, public_key, &client).await {
|
||||||
eprintln!("Failed to handle client: {e}");
|
eprintln!("Failed to handle client: {e}");
|
||||||
} else {
|
} else {
|
||||||
println!("Client connection closed")
|
println!("Client connection closed")
|
||||||
}
|
}
|
||||||
|
|
||||||
clients.connections.remove(&conid);
|
let mut clients_meta = s.clients.lock().await;
|
||||||
|
|
||||||
if clients.connections.len() == 0 {
|
let mut connections = clients.connections.lock().await;
|
||||||
|
|
||||||
|
connections.remove(&conid);
|
||||||
|
|
||||||
|
if connections.len() == 0 {
|
||||||
clients_meta.remove(&public_key);
|
clients_meta.remove(&public_key);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -85,11 +100,32 @@ impl Server {
|
|||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub async fn broadcast(self: &Arc<Self>, message: &ClientMethod) -> anyhow::Result<()> {
|
||||||
|
let mut set = JoinSet::new();
|
||||||
|
|
||||||
|
for (_, client) in self.clients.lock().await.iter() {
|
||||||
|
let msg = message.clone();
|
||||||
|
let client = client.clone();
|
||||||
|
|
||||||
|
set.spawn(async move { client.send(&msg).await });
|
||||||
|
}
|
||||||
|
|
||||||
|
// Await all spawned tasks to finish
|
||||||
|
while let Some(res) = set.join_next().await {
|
||||||
|
// handle task panic or errors if necessary
|
||||||
|
if let Ok(Err(e)) = res {
|
||||||
|
eprintln!("Failed to send to a client: {:?}", e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl UserConnections {
|
impl UserConnections {
|
||||||
pub async fn send(&self, message: &ClientMethod) -> anyhow::Result<()> {
|
pub async fn send(&self, message: &ClientMethod) -> anyhow::Result<()> {
|
||||||
for (_, conn) in &self.connections {
|
for (_, conn) in self.connections.lock().await.iter() {
|
||||||
send_socket(&mut *conn.lock().await, message).await?;
|
send_socket(&mut *conn.lock().await, message).await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -97,7 +133,7 @@ impl UserConnections {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub async fn send_to(&self, id: u16, message: &ClientMethod) -> anyhow::Result<bool> {
|
pub async fn send_to(&self, id: u16, message: &ClientMethod) -> anyhow::Result<bool> {
|
||||||
if let Some(conn) = self.connections.get(&id) {
|
if let Some(conn) = self.connections.lock().await.get(&id) {
|
||||||
send_socket(&mut *conn.lock().await, message).await?;
|
send_socket(&mut *conn.lock().await, message).await?;
|
||||||
|
|
||||||
Ok(true)
|
Ok(true)
|
||||||
|
|||||||
+6
-1
@@ -1,9 +1,14 @@
|
|||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
pub struct ClientMeta {}
|
#[serde(rename_all = "camelCase")]
|
||||||
|
pub struct ClientMeta {
|
||||||
|
pub display_name: String,
|
||||||
|
pub avatar: Option<String>,
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
|
#[serde(rename_all = "camelCase")]
|
||||||
pub struct ServerMeta {
|
pub struct ServerMeta {
|
||||||
pub name: String,
|
pub name: String,
|
||||||
pub description: String,
|
pub description: String,
|
||||||
|
|||||||
Reference in New Issue
Block a user