fix clippy errs
This commit is contained in:
@@ -115,7 +115,11 @@ pub async fn get_supporter_messages(
|
|||||||
"Session not found. Chat history may have been cleared.",
|
"Session not found. Chat history may have been cleared.",
|
||||||
));
|
));
|
||||||
};
|
};
|
||||||
let html: String = session.messages.iter().map(render_supporter_message).collect();
|
let html: String = session
|
||||||
|
.messages
|
||||||
|
.iter()
|
||||||
|
.map(render_supporter_message)
|
||||||
|
.collect();
|
||||||
Html(html)
|
Html(html)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -171,7 +175,10 @@ pub async fn supporter_send_message(
|
|||||||
State(state): State<Arc<AppState>>,
|
State(state): State<Arc<AppState>>,
|
||||||
Form(form): Form<SendMessageForm>,
|
Form(form): Form<SendMessageForm>,
|
||||||
) -> Html<String> {
|
) -> Html<String> {
|
||||||
if state.is_rate_limited(limit_key("supporter-message", addr), MESSAGE_LIMIT_PER_WINDOW) {
|
if state.is_rate_limited(
|
||||||
|
limit_key("supporter-message", addr),
|
||||||
|
MESSAGE_LIMIT_PER_WINDOW,
|
||||||
|
) {
|
||||||
return Html(render_system_notice(
|
return Html(render_system_notice(
|
||||||
"Too many messages from this network. Please wait until the 3-minute window resets.",
|
"Too many messages from this network. Please wait until the 3-minute window resets.",
|
||||||
));
|
));
|
||||||
|
|||||||
+1
-1
@@ -1,3 +1,3 @@
|
|||||||
pub mod pages;
|
|
||||||
pub mod chat;
|
pub mod chat;
|
||||||
|
pub mod pages;
|
||||||
pub mod ws;
|
pub mod ws;
|
||||||
|
|||||||
+11
-14
@@ -1,7 +1,7 @@
|
|||||||
use axum::{
|
use axum::{
|
||||||
extract::{
|
extract::{
|
||||||
ws::{Message, WebSocket, WebSocketUpgrade},
|
|
||||||
Path, State,
|
Path, State,
|
||||||
|
ws::{Message, WebSocket, WebSocketUpgrade},
|
||||||
},
|
},
|
||||||
response::IntoResponse,
|
response::IntoResponse,
|
||||||
};
|
};
|
||||||
@@ -14,7 +14,6 @@ const USER_TYPING_STOP_PREFIX: &str = "typing:user:stop:";
|
|||||||
const SUPPORTER_TYPING_START_PREFIX: &str = "typing:supporter:start:";
|
const SUPPORTER_TYPING_START_PREFIX: &str = "typing:supporter:start:";
|
||||||
const SUPPORTER_TYPING_STOP_PREFIX: &str = "typing:supporter:stop:";
|
const SUPPORTER_TYPING_STOP_PREFIX: &str = "typing:supporter:stop:";
|
||||||
|
|
||||||
/// WebSocket for user — pushes notifications when supporter replies
|
|
||||||
pub async fn user_ws(
|
pub async fn user_ws(
|
||||||
ws: WebSocketUpgrade,
|
ws: WebSocketUpgrade,
|
||||||
Path(session_id): Path<String>,
|
Path(session_id): Path<String>,
|
||||||
@@ -36,15 +35,13 @@ async fn handle_user_socket(mut socket: WebSocket, session_id: String, state: Ar
|
|||||||
|
|
||||||
loop {
|
loop {
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
// Forward broadcast events to the user browser.
|
|
||||||
Ok(event) = rx.recv() => {
|
Ok(event) = rx.recv() => {
|
||||||
if event == session_id || event == RESET_EVENT || is_typing_event_for_session(&event, &session_id) {
|
if (event == session_id || event == RESET_EVENT || is_typing_event_for_session(&event, &session_id))
|
||||||
if socket.send(Message::Text(event)).await.is_err() {
|
&& socket.send(Message::Text(event)).await.is_err()
|
||||||
break;
|
{
|
||||||
}
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// Handle ping/close and typing events from client.
|
|
||||||
msg = socket.recv() => {
|
msg = socket.recv() => {
|
||||||
match msg {
|
match msg {
|
||||||
Some(Ok(Message::Text(text))) => {
|
Some(Ok(Message::Text(text))) => {
|
||||||
@@ -63,7 +60,6 @@ async fn handle_user_socket(mut socket: WebSocket, session_id: String, state: Ar
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// WebSocket for supporter dashboard — pushes notifications when any session updates
|
|
||||||
pub async fn supporter_ws(
|
pub async fn supporter_ws(
|
||||||
ws: WebSocketUpgrade,
|
ws: WebSocketUpgrade,
|
||||||
Path(_supporter_id): Path<String>,
|
Path(_supporter_id): Path<String>,
|
||||||
@@ -90,11 +86,12 @@ async fn handle_supporter_socket(mut socket: WebSocket, state: Arc<AppState>) {
|
|||||||
msg = socket.recv() => {
|
msg = socket.recv() => {
|
||||||
match msg {
|
match msg {
|
||||||
Some(Ok(Message::Text(text))) => {
|
Some(Ok(Message::Text(text))) => {
|
||||||
if let Some(session_id) = supporter_typing_session_id(&text) {
|
if let Some(session_id) = supporter_typing_session_id(&text)
|
||||||
if !session_id.is_empty() && state.sessions.contains_key(session_id) {
|
&& !session_id.is_empty()
|
||||||
let tx = state.get_or_create_notifier(session_id);
|
&& state.sessions.contains_key(session_id)
|
||||||
let _ = tx.send(text);
|
{
|
||||||
}
|
let tx = state.get_or_create_notifier(session_id);
|
||||||
|
let _ = tx.send(text);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Some(Ok(Message::Close(_))) | None => break,
|
Some(Ok(Message::Close(_))) | None => break,
|
||||||
|
|||||||
+11
-3
@@ -28,7 +28,9 @@ async fn main() {
|
|||||||
loop {
|
loop {
|
||||||
tokio::time::sleep(CHAT_RESET_INTERVAL).await;
|
tokio::time::sleep(CHAT_RESET_INTERVAL).await;
|
||||||
reset_state.clear_everything();
|
reset_state.clear_everything();
|
||||||
tracing::info!("Cleared all in-memory chat sessions, messages, notifiers, and rate limits");
|
tracing::info!(
|
||||||
|
"Cleared all in-memory chat sessions, messages, notifiers, and rate limits"
|
||||||
|
);
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -38,9 +40,15 @@ async fn main() {
|
|||||||
.route("/supporter", get(handlers::pages::supporter_page))
|
.route("/supporter", get(handlers::pages::supporter_page))
|
||||||
.route("/reset-countdown", get(handlers::pages::reset_countdown))
|
.route("/reset-countdown", get(handlers::pages::reset_countdown))
|
||||||
.route("/ws/user/:session_id", get(handlers::ws::user_ws))
|
.route("/ws/user/:session_id", get(handlers::ws::user_ws))
|
||||||
.route("/ws/supporter/:supporter_id", get(handlers::ws::supporter_ws))
|
.route(
|
||||||
|
"/ws/supporter/:supporter_id",
|
||||||
|
get(handlers::ws::supporter_ws),
|
||||||
|
)
|
||||||
.route("/messages/:session_id", get(handlers::chat::get_messages))
|
.route("/messages/:session_id", get(handlers::chat::get_messages))
|
||||||
.route("/supporter/messages/:session_id", get(handlers::chat::get_supporter_messages))
|
.route(
|
||||||
|
"/supporter/messages/:session_id",
|
||||||
|
get(handlers::chat::get_supporter_messages),
|
||||||
|
)
|
||||||
.route("/send/:session_id", post(handlers::chat::send_message))
|
.route("/send/:session_id", post(handlers::chat::send_message))
|
||||||
.route(
|
.route(
|
||||||
"/supporter/send/:session_id",
|
"/supporter/send/:session_id",
|
||||||
|
|||||||
+4
-1
@@ -1,6 +1,9 @@
|
|||||||
use crate::models::ChatSession;
|
use crate::models::ChatSession;
|
||||||
use dashmap::DashMap;
|
use dashmap::DashMap;
|
||||||
use std::{sync::Mutex, time::{Duration, Instant}};
|
use std::{
|
||||||
|
sync::Mutex,
|
||||||
|
time::{Duration, Instant},
|
||||||
|
};
|
||||||
use tokio::sync::broadcast;
|
use tokio::sync::broadcast;
|
||||||
|
|
||||||
pub const RESET_EVENT: &str = "__reset__";
|
pub const RESET_EVENT: &str = "__reset__";
|
||||||
|
|||||||
Reference in New Issue
Block a user