Small stuff
This commit is contained in:
@@ -1,21 +1,36 @@
|
|||||||
setInterval(() => {
|
import readline from 'readline';
|
||||||
const ws = new WebSocket('ws://localhost:7080');
|
|
||||||
|
|
||||||
function sendMessage(message) {
|
const rl = readline.createInterface({
|
||||||
ws.send(JSON.stringify({ type: 'send_message', params: {
|
input: process.stdin,
|
||||||
channel_id: 'Hello',
|
output: process.stdout
|
||||||
contents: message
|
});
|
||||||
}}));
|
|
||||||
}
|
|
||||||
|
|
||||||
ws.onopen = () => {
|
const ws = new WebSocket('ws://localhost:7080');
|
||||||
console.log('WebSocket connection established');
|
|
||||||
sendMessage('Hello, Server!');
|
|
||||||
};
|
|
||||||
|
|
||||||
ws.onmessage = (event) => {
|
ws.onopen = () => {
|
||||||
const message = JSON.parse(event.data);
|
console.log('WebSocket connection established');
|
||||||
console.log('Received:', message);
|
};
|
||||||
ws.close();
|
|
||||||
};
|
ws.onmessage = (event) => {
|
||||||
}, 100)
|
const message = JSON.parse(event.data);
|
||||||
|
console.log('Received:', message);
|
||||||
|
// ws.close();
|
||||||
|
};
|
||||||
|
|
||||||
|
function sendMessage(message) {
|
||||||
|
ws.send(JSON.stringify({ type: 'send_message', params: {
|
||||||
|
channel_id: 'Hello',
|
||||||
|
contents: message
|
||||||
|
}}));
|
||||||
|
}
|
||||||
|
|
||||||
|
function ask(question) {
|
||||||
|
return new Promise((resolve) => {
|
||||||
|
rl.question(question, resolve);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
while (true) {
|
||||||
|
const message = await ask('');
|
||||||
|
sendMessage(message);
|
||||||
|
}
|
||||||
+11
-3
@@ -7,7 +7,7 @@ use std::{
|
|||||||
use anyhow::Error;
|
use anyhow::Error;
|
||||||
use tungstenite::{Message, Utf8Bytes, WebSocket, accept};
|
use tungstenite::{Message, Utf8Bytes, WebSocket, accept};
|
||||||
|
|
||||||
use crate::types::{ClientMessage, ServerMessage, WsMessage};
|
use crate::types::{ClientMessage, ServerMessage, WsMessage, data::ResponseError};
|
||||||
|
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
pub struct Client(Arc<Mutex<WebSocket<TcpStream>>>);
|
pub struct Client(Arc<Mutex<WebSocket<TcpStream>>>);
|
||||||
@@ -18,7 +18,7 @@ impl Client {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub fn new_tcp(ws: TcpStream) -> crate::Result<Self> {
|
pub fn new_tcp(ws: TcpStream) -> crate::Result<Self> {
|
||||||
Ok(Self(Arc::new(Mutex::new(accept(ws)?))))
|
Ok(Self::new_ws(accept(ws)?))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -42,7 +42,7 @@ impl Client {
|
|||||||
Message::Text(t) => {
|
Message::Text(t) => {
|
||||||
let v = t.to_string();
|
let v = t.to_string();
|
||||||
match serde_json::from_str(&v) {
|
match serde_json::from_str(&v) {
|
||||||
Ok(f) => Ok(Some(WsMessage::FromClient(f))),
|
Ok(f) => Ok(Some(WsMessage::Message(f))),
|
||||||
Err(_) => Ok(Some(WsMessage::String(v))),
|
Err(_) => Ok(Some(WsMessage::String(v))),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -62,4 +62,12 @@ impl Client {
|
|||||||
.send(Message::Text(Utf8Bytes::from(serde_json::to_string(&m)?)))
|
.send(Message::Text(Utf8Bytes::from(serde_json::to_string(&m)?)))
|
||||||
.map_err(|e| e.into())
|
.map_err(|e| e.into())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn send_err(&self, m: ResponseError) -> crate::Result<()> {
|
||||||
|
self.0
|
||||||
|
.lock()
|
||||||
|
.unwrap()
|
||||||
|
.send(Message::Text(Utf8Bytes::from(serde_json::to_string(&m)?)))
|
||||||
|
.map_err(|e| e.into())
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+33
-23
@@ -1,4 +1,4 @@
|
|||||||
use crate::ServerConfig;
|
use crate::{ServerConfig, types::data::Message};
|
||||||
use rusqlite::{Connection, Result, params};
|
use rusqlite::{Connection, Result, params};
|
||||||
|
|
||||||
pub struct Database {
|
pub struct Database {
|
||||||
@@ -42,29 +42,33 @@ impl MessagesDb {
|
|||||||
user_id: &str,
|
user_id: &str,
|
||||||
contents: &str,
|
contents: &str,
|
||||||
timestamp: i64,
|
timestamp: i64,
|
||||||
) -> Result<usize> {
|
) -> Result<Message> {
|
||||||
self.0.execute(
|
self.0.execute(
|
||||||
"INSERT INTO chat (channel_id, user_id, contents, timestamp)
|
"INSERT INTO chat (channel_id, user_id, contents, timestamp)
|
||||||
VALUES (?1, ?2, ?3, ?4)",
|
VALUES (?1, ?2, ?3, ?4)",
|
||||||
params![channel_id, user_id, contents, timestamp],
|
params![channel_id, user_id, contents, timestamp],
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Fetch the latest N messages for a channel
|
|
||||||
pub fn fetch_recent(
|
|
||||||
&self,
|
|
||||||
channel_id: &str,
|
|
||||||
limit: usize,
|
|
||||||
) -> Result<Vec<(i64, String, String, String, i64)>> {
|
|
||||||
let mut stmt = self.0.prepare(
|
|
||||||
"SELECT id, channel_id, user_id, contents, timestamp
|
|
||||||
FROM chat
|
|
||||||
WHERE channel_id = ?1
|
|
||||||
ORDER BY timestamp DESC
|
|
||||||
LIMIT ?2",
|
|
||||||
)?;
|
)?;
|
||||||
|
|
||||||
let rows = stmt.query_map(params![channel_id, limit], |row| {
|
let id = self.0.last_insert_rowid();
|
||||||
|
|
||||||
|
Ok(Message {
|
||||||
|
id,
|
||||||
|
channel_id: channel_id.to_string(),
|
||||||
|
from: user_id.to_string(),
|
||||||
|
contents: contents.to_string(),
|
||||||
|
timestamp,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Get a message by its ID
|
||||||
|
pub fn get_by_id(&self, message_id: usize) -> Result<Option<Message>> {
|
||||||
|
let mut stmt = self.0.prepare(
|
||||||
|
"SELECT id, channel_id, user_id, contents, timestamp
|
||||||
|
FROM chat
|
||||||
|
WHERE id = ?1",
|
||||||
|
)?;
|
||||||
|
|
||||||
|
let mut rows = stmt.query_map(params![message_id], |row| {
|
||||||
Ok((
|
Ok((
|
||||||
row.get::<_, i64>(0)?, // id
|
row.get::<_, i64>(0)?, // id
|
||||||
row.get::<_, String>(1)?, // channel_id
|
row.get::<_, String>(1)?, // channel_id
|
||||||
@@ -74,10 +78,16 @@ impl MessagesDb {
|
|||||||
))
|
))
|
||||||
})?;
|
})?;
|
||||||
|
|
||||||
let mut results = Vec::new();
|
if let Some(row) = rows.next() {
|
||||||
for row in rows {
|
let (id, channel_id, user_id, contents, timestamp) = row?;
|
||||||
results.push(row?);
|
return Ok(Some(Message {
|
||||||
|
id,
|
||||||
|
channel_id,
|
||||||
|
from: user_id,
|
||||||
|
contents,
|
||||||
|
timestamp,
|
||||||
|
}));
|
||||||
}
|
}
|
||||||
Ok(results)
|
Ok(None)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+7
-2
@@ -124,13 +124,13 @@ impl Server {
|
|||||||
// The main req/res loop
|
// The main req/res loop
|
||||||
loop {
|
loop {
|
||||||
match client.read()? {
|
match client.read()? {
|
||||||
Some(WsMessage::FromClient(req)) => match req {
|
Some(WsMessage::Message(req)) => match req {
|
||||||
ClientMessage::SendMessage {
|
ClientMessage::SendMessage {
|
||||||
channel_id,
|
channel_id,
|
||||||
contents,
|
contents,
|
||||||
} => {
|
} => {
|
||||||
Self::LOGGER.info(format!("SendMessage to {channel_id}: {contents}"));
|
Self::LOGGER.info(format!("SendMessage to {channel_id}: {contents}"));
|
||||||
self.wrap_err(
|
let msg = self.wrap_err(
|
||||||
&client,
|
&client,
|
||||||
self.db.messages_db.insert(
|
self.db.messages_db.insert(
|
||||||
&channel_id,
|
&channel_id,
|
||||||
@@ -139,6 +139,11 @@ impl Server {
|
|||||||
chrono::Utc::now().timestamp(),
|
chrono::Utc::now().timestamp(),
|
||||||
),
|
),
|
||||||
)?;
|
)?;
|
||||||
|
|
||||||
|
self.wrap_err(
|
||||||
|
&client,
|
||||||
|
client.send(types::ServerMessage::MessageCreate(msg)),
|
||||||
|
)?;
|
||||||
}
|
}
|
||||||
|
|
||||||
ClientMessage::EditMessage {
|
ClientMessage::EditMessage {
|
||||||
|
|||||||
+13
-4
@@ -33,7 +33,7 @@ pub enum ServerMessage {
|
|||||||
Authenticated { user_id: String },
|
Authenticated { user_id: String },
|
||||||
|
|
||||||
/// Error responses
|
/// Error responses
|
||||||
Error { message: String },
|
Error(),
|
||||||
|
|
||||||
/// A new message in a channel
|
/// A new message in a channel
|
||||||
MessageCreate(data::Message),
|
MessageCreate(data::Message),
|
||||||
@@ -57,7 +57,7 @@ pub enum ServerMessage {
|
|||||||
/// WebSocket wrapper
|
/// WebSocket wrapper
|
||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
pub enum WsMessage<T: Serialize + for<'de> Deserialize<'de>> {
|
pub enum WsMessage<T: Serialize + for<'de> Deserialize<'de>> {
|
||||||
FromClient(T),
|
Message(T),
|
||||||
Binary(Bytes),
|
Binary(Bytes),
|
||||||
String(String),
|
String(String),
|
||||||
}
|
}
|
||||||
@@ -68,8 +68,8 @@ pub mod data {
|
|||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
pub struct Message {
|
pub struct Message {
|
||||||
pub id: String,
|
pub id: i64,
|
||||||
pub channel_id: u8,
|
pub channel_id: String,
|
||||||
pub from: String,
|
pub from: String,
|
||||||
pub contents: String,
|
pub contents: String,
|
||||||
pub timestamp: i64,
|
pub timestamp: i64,
|
||||||
@@ -88,4 +88,13 @@ pub mod data {
|
|||||||
Text,
|
Text,
|
||||||
Voice,
|
Voice,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
|
#[serde(tag = "error", rename_all = "snake_case")]
|
||||||
|
pub enum ResponseError {
|
||||||
|
InvalidRequest { message: String },
|
||||||
|
Unauthorized { message: String },
|
||||||
|
NotFound { message: String },
|
||||||
|
InternalError { message: String },
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user