This commit is contained in:
2026-08-29 18:39:28 +02:00
parent 88ad3c389f
commit 9551e90678
17 changed files with 642 additions and 269 deletions
+20
View File
@@ -1,5 +1,6 @@
import {defineStore} from "pinia"; import {defineStore} from "pinia";
import {useApi} from "@/composables/useApi.ts"; import {useApi} from "@/composables/useApi.ts";
import {onGatewayEvent} from "@/plugins/events.ts";
export interface Emoji { export interface Emoji {
id: string; id: string;
@@ -65,3 +66,22 @@ export const useEmojiStore = defineStore("emoji", {
}, },
}, },
}); });
onGatewayEvent("Emoji", (payload) => {
const store = useEmojiStore();
const emoji = payload.content as Emoji;
if (!emoji?.id) return;
// The store is scoped to the currently displayed server. Global emojis are
// always relevant; server emojis are relevant only for the active server.
if (emoji.server_id !== null && emoji.server_id !== store.activeServerId) return;
if (payload.action === "add" || payload.action === "update") {
store.emojis = [
...store.emojis.filter(item => item.id !== emoji.id),
emoji,
];
} else if (payload.action === "remove") {
store.emojis = store.emojis.filter(item => item.id !== emoji.id);
}
});
+15
View File
@@ -2,6 +2,7 @@ import {defineStore} from "pinia";
import {useApi} from "@/composables/useApi.ts"; import {useApi} from "@/composables/useApi.ts";
import {useChannelStore} from "@/stores/channel.ts"; import {useChannelStore} from "@/stores/channel.ts";
import {useCategoryStore} from "@/stores/category.ts"; import {useCategoryStore} from "@/stores/category.ts";
import {onGatewayEvent} from "@/plugins/events.ts";
export interface Server { export interface Server {
id: string id: string
@@ -153,3 +154,17 @@ export const useServerStore = defineStore("server", {
} }
} }
}); });
onGatewayEvent("Channel", (payload) => {
const channel = payload.content as { server_id?: string | null };
if (channel?.server_id) {
void useServerStore().fetchServerTree(channel.server_id);
}
});
onGatewayEvent("Category", (payload) => {
const category = payload.content as { server_id?: string | null };
if (category?.server_id) {
void useServerStore().fetchServerTree(category.server_id);
}
});
+8 -3
View File
@@ -5,7 +5,7 @@ use crate::database::Database;
use crate::http::server::HttpServer; use crate::http::server::HttpServer;
use crate::metrics::{AppMetrics, reporter}; use crate::metrics::{AppMetrics, reporter};
use crate::repositories::Repositories; use crate::repositories::Repositories;
use crate::routes::gateway::GatewayManager; use crate::routes::gateway::{GatewayManager, RealtimeRouter};
use crate::services::Services; use crate::services::Services;
use crate::udp::server::UdpServer; use crate::udp::server::UdpServer;
use event_bus::EventBus; use event_bus::EventBus;
@@ -72,8 +72,13 @@ impl App {
services services
.realtime_registry .realtime_registry
.start_listening(repositories.clone(), event_bus.clone()); .start_listening(repositories.clone(), event_bus.clone());
let gateway = Arc::new(GatewayManager::new(services.clone(), repositories.clone())); let gateway = Arc::new(GatewayManager::new());
gateway.start(event_bus.clone()); Arc::new(RealtimeRouter::new(
gateway.clone(),
services.clone(),
repositories.clone(),
))
.start(event_bus.clone());
let state = AppState { let state = AppState {
db, db,
+1 -4
View File
@@ -1,20 +1,17 @@
use crate::models::channel; use crate::models::channel;
use uuid::Uuid;
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct ChannelCreatedEvent { pub struct ChannelCreatedEvent {
pub server_id: Uuid,
pub channel: channel::Model, pub channel: channel::Model,
} }
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct ChannelUpdatedEvent { pub struct ChannelUpdatedEvent {
pub server_id: Uuid, pub previous: channel::Model,
pub channel: channel::Model, pub channel: channel::Model,
} }
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct ChannelDeletedEvent { pub struct ChannelDeletedEvent {
pub server_id: Uuid,
pub channel: channel::Model, pub channel: channel::Model,
} }
+17
View File
@@ -0,0 +1,17 @@
use crate::models::emoji;
#[derive(Debug, Clone)]
pub struct EmojiCreatedEvent {
pub emoji: emoji::Model,
}
#[derive(Debug, Clone)]
pub struct EmojiUpdatedEvent {
pub previous: emoji::Model,
pub emoji: emoji::Model,
}
#[derive(Debug, Clone)]
pub struct EmojiDeletedEvent {
pub emoji: emoji::Model,
}
+1
View File
@@ -1,3 +1,4 @@
pub mod channel; pub mod channel;
pub mod emoji;
pub mod message; pub mod message;
pub mod server; pub mod server;
+7 -1
View File
@@ -3,6 +3,7 @@ use crate::domain::dto::conversation::{
ConversationParticipantResponse, ConversationResponse, CreateConversationRequest, ConversationParticipantResponse, ConversationResponse, CreateConversationRequest,
ForkConversationRequest, ForkConversationRequest,
}; };
use crate::domain::events::channel::ChannelCreatedEvent;
use crate::http::context::CurrentUser; use crate::http::context::CurrentUser;
use crate::http::error::HTTPError; use crate::http::error::HTTPError;
use crate::models::{channel, channel_user, message, user}; use crate::models::{channel, channel_user, message, user};
@@ -107,7 +108,12 @@ async fn create_channel(state: &AppState, ids: &[Uuid]) -> Result<channel::Model
.services .services
.realtime_registry .realtime_registry
.set_channel_users(channel.id, ids.iter().copied()); .set_channel_users(channel.id, ids.iter().copied());
state.event_bus.emit("channel_created", channel.clone()); state.event_bus.emit(
"channel_created",
ChannelCreatedEvent {
channel: channel.clone(),
},
);
Ok(channel) Ok(channel)
} }
+26 -4
View File
@@ -1,3 +1,5 @@
use crate::domain::events::emoji::{EmojiCreatedEvent, EmojiDeletedEvent, EmojiUpdatedEvent};
use crate::services::media;
use crate::{ use crate::{
core::state::AppState, core::state::AppState,
domain::dto::emoji::{EmojiQueryParams, UpdateEmojiRequest}, domain::dto::emoji::{EmojiQueryParams, UpdateEmojiRequest},
@@ -6,7 +8,6 @@ use crate::{
routes::emoji::mapper, routes::emoji::mapper,
services::emoji::EmojiService, services::emoji::EmojiService,
}; };
use crate::services::media;
use axum::{ use axum::{
Json, Json,
body::Body, body::Body,
@@ -129,9 +130,13 @@ pub async fn create(
.ok_or_else(|| HTTPError::BadRequest("Unsupported or invalid image format".into()))?; .ok_or_else(|| HTTPError::BadRequest("Unsupported or invalid image format".into()))?;
let extension = media::extension_from_mime(&detected); let extension = media::extension_from_mime(&detected);
mime = Some(detected); mime = Some(detected);
let (p, h) = let (p, h) = EmojiService::save_asset(
EmojiService::save_asset(std::path::Path::new(&state.config.media.root), id, &bytes, extension) std::path::Path::new(&state.config.media.root),
.await?; id,
&bytes,
extension,
)
.await?;
path = Some(p); path = Some(p);
sha = Some(h); sha = Some(h);
} }
@@ -150,6 +155,12 @@ pub async fn create(
..Default::default() ..Default::default()
}; };
let created = state.services.emoji.create(model, name).await?; let created = state.services.emoji.create(model, name).await?;
state.event_bus.emit(
"emoji_created",
EmojiCreatedEvent {
emoji: created.clone(),
},
);
Ok((StatusCode::CREATED, Json(mapper::response(created)))) Ok((StatusCode::CREATED, Json(mapper::response(created))))
} }
@@ -172,6 +183,7 @@ pub async fn update(
.emoji .emoji
.name_available(target_name, target_server_id, Some(id)) .name_available(target_name, target_server_id, Some(id))
.await?; .await?;
let previous = existing.clone();
let mut active: emoji::ActiveModel = existing.into(); let mut active: emoji::ActiveModel = existing.into();
if let Some(server_id) = payload.server_id { if let Some(server_id) = payload.server_id {
state state
@@ -189,6 +201,13 @@ pub async fn update(
active.name = Set(EmojiService::normalize_name(&name)?); active.name = Set(EmojiService::normalize_name(&name)?);
} }
let updated = state.repositories.emoji.update(active).await?; let updated = state.repositories.emoji.update(active).await?;
state.event_bus.emit(
"emoji_updated",
EmojiUpdatedEvent {
previous,
emoji: updated.clone(),
},
);
Ok(Json(mapper::response(updated))) Ok(Json(mapper::response(updated)))
} }
@@ -222,6 +241,9 @@ pub async fn delete(
model.file_path.as_deref(), model.file_path.as_deref(),
) )
.await; .await;
state
.event_bus
.emit("emoji_deleted", EmojiDeletedEvent { emoji: model });
Ok(StatusCode::NO_CONTENT) Ok(StatusCode::NO_CONTENT)
} else { } else {
Err(HTTPError::NotFound) Err(HTTPError::NotFound)
+1 -1
View File
@@ -41,7 +41,7 @@ async fn handle_socket(socket: WebSocket, state: AppState, user: User) {
let (mut sender, mut receiver) = socket.split(); let (mut sender, mut receiver) = socket.split();
let (tx, mut rx) = mpsc::unbounded_channel::<Message>(); let (tx, mut rx) = mpsc::unbounded_channel::<Message>();
let mut client = GatewayClient::new(user, tx, state.event_bus.clone()); let mut client = GatewayClient::new(user, tx);
client.on_connect().await; client.on_connect().await;
state.gateway.add_client(client.clone()); state.gateway.add_client(client.clone());
+331 -205
View File
@@ -1,37 +1,40 @@
use crate::domain::events::channel::{
ChannelCreatedEvent, ChannelDeletedEvent, ChannelUpdatedEvent,
};
use crate::domain::events::emoji::{EmojiCreatedEvent, EmojiDeletedEvent, EmojiUpdatedEvent};
use crate::domain::events::message::{ use crate::domain::events::message::{
MessageCreatedEvent, MessageDeletedEvent, MessageReactionAddedEvent, MessageCreatedEvent, MessageDeletedEvent, MessageReactionAddedEvent,
MessageReactionRemovedEvent, MessageUpdatedEvent, MessageReactionRemovedEvent, MessageUpdatedEvent,
}; };
use crate::models::user::Model as User; use crate::models::{category, server, server_user, user};
use crate::repositories::Repositories;
use crate::routes::category::mapper::category_model_to_category_response; use crate::routes::category::mapper::category_model_to_category_response;
use crate::routes::channel::mapper::channel_model_to_channel_response; use crate::routes::channel::mapper::channel_model_to_channel_response;
use crate::routes::message::mapper::{ use crate::routes::message::mapper::{
message_model_to_message_response_with_data, message_model_to_message_response_with_data, message_model_to_message_response_with_reactions,
message_model_to_message_response_with_reactions,
reaction_model_to_response, reaction_model_to_response,
}; };
use crate::routes::server::mapper::server_model_to_server_response; use crate::routes::server::mapper::server_model_to_server_response;
use crate::services::Services; use crate::services::Services;
use crate::repositories::Repositories;
use axum::extract::ws::Message; use axum::extract::ws::Message;
use event_bus::EventBus; use event_bus::EventBus;
use events::GatewayEvent; use events::GatewayEvent;
use parking_lot::RwLock; use parking_lot::RwLock;
use std::collections::HashMap; use sea_orm::{ColumnTrait, EntityTrait, QueryFilter, QuerySelect};
use serde::Serialize;
use std::collections::{HashMap, HashSet};
use std::sync::Arc; use std::sync::Arc;
use tokio::sync::mpsc; use tokio::sync::mpsc;
use tokio::task::JoinHandle;
use uuid::Uuid; use uuid::Uuid;
pub mod events; pub mod events;
pub mod handlers; pub mod handlers;
pub mod routes; pub mod routes;
/// Couche de transport WebSocket. Elle ne connaît pas les permissions.
#[derive(Debug)] #[derive(Debug)]
pub struct GatewayManager { pub struct GatewayManager {
pub clients: RwLock<HashMap<ConnectionKey, GatewayClient>>, pub clients: RwLock<HashMap<ConnectionKey, GatewayClient>>,
services: Arc<Services>,
repositories: Arc<Repositories>,
} }
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
@@ -42,59 +45,119 @@ pub struct ConnectionKey {
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct GatewayClient { pub struct GatewayClient {
user: User, user: crate::models::user::Model,
connection_id: Uuid, connection_id: Uuid,
pub sender: mpsc::UnboundedSender<Message>, pub sender: mpsc::UnboundedSender<Message>,
pub event_bus: Arc<EventBus>,
_event_handles: Vec<Arc<JoinHandle<()>>>,
} }
impl GatewayManager { impl GatewayManager {
pub fn new(services: Arc<Services>, repositories: Arc<Repositories>) -> Self { pub fn new() -> Self {
Self { Self {
clients: RwLock::new(HashMap::new()), clients: RwLock::new(HashMap::new()),
}
}
pub(crate) fn add_client(&self, client: GatewayClient) {
self.clients.write().insert(client.key(), client);
}
pub(crate) fn remove_client(&self, client: &GatewayClient) {
self.clients.write().remove(&client.key());
}
pub(crate) fn send_to_users<T: Serialize>(
&self,
users: impl IntoIterator<Item = Uuid>,
namespace: &'static str,
action: &'static str,
content: T,
) {
let Ok(json) = serde_json::to_string(&GatewayEvent {
namespace,
action,
content,
}) else {
return;
};
let users: HashSet<_> = users.into_iter().collect();
for (key, client) in self.clients.read().iter() {
if users.contains(&key.user_id) {
let _ = client.sender.send(Message::Text(json.clone().into()));
}
}
}
}
/// Résout les audiences puis délègue l'envoi au transport WebSocket.
pub struct RealtimeRouter {
gateway: Arc<GatewayManager>,
services: Arc<Services>,
repositories: Arc<Repositories>,
}
impl RealtimeRouter {
pub fn new(
gateway: Arc<GatewayManager>,
services: Arc<Services>,
repositories: Arc<Repositories>,
) -> Self {
Self {
gateway,
services, services,
repositories, repositories,
} }
} }
/// Démarre les routeurs centraux des événements de messages.
pub fn start(self: &Arc<Self>, event_bus: Arc<EventBus>) { pub fn start(self: &Arc<Self>, event_bus: Arc<EventBus>) {
let manager = Arc::clone(self); let router = Arc::clone(self);
event_bus.on_async::<MessageCreatedEvent, _, _>("message_created", move |event| { event_bus.on_async::<MessageCreatedEvent, _, _>("message_created", move |event| {
let manager = Arc::clone(&manager); let router = Arc::clone(&router);
async move { async move {
let message_id = event.message.id; let attachments = router
let attachments = manager
.repositories .repositories
.message .message
.attachments_for_messages(&[message_id]) .attachments_for_messages(&[event.message.id])
.await .await
.ok() .ok()
.and_then(|mut items| items.remove(&message_id)) .and_then(|mut items| items.remove(&event.message.id))
.unwrap_or_default(); .unwrap_or_default();
manager.broadcast_message( let users = router
event.channel_id, .services
.realtime_registry
.users_for_channel(event.channel_id);
router.gateway.send_to_users(
users,
"Message",
"add", "add",
message_model_to_message_response_with_data(event.message, event.server_id, Vec::new(), attachments), message_model_to_message_response_with_data(
event.message,
event.server_id,
Vec::new(),
attachments,
),
); );
} }
}); });
let manager = Arc::clone(self); let router = Arc::clone(self);
event_bus.on_async::<MessageUpdatedEvent, _, _>("message_updated", move |event| { event_bus.on_async::<MessageUpdatedEvent, _, _>("message_updated", move |event| {
let manager = Arc::clone(&manager); let router = Arc::clone(&router);
async move { async move {
let message_id = event.message.id; let reactions = router
let reactions = manager
.services .services
.message_reaction .message_reaction
.grouped_for_messages(&[message_id]) .grouped_for_messages(&[event.message.id])
.await .await
.ok() .ok()
.and_then(|mut groups| groups.remove(&message_id)) .and_then(|mut groups| groups.remove(&event.message.id))
.unwrap_or_default(); .unwrap_or_default();
manager.broadcast_message( let users = router
event.channel_id, .services
.realtime_registry
.users_for_channel(event.channel_id);
router.gateway.send_to_users(
users,
"Message",
"update", "update",
message_model_to_message_response_with_reactions( message_model_to_message_response_with_reactions(
event.message, event.message,
@@ -105,22 +168,33 @@ impl GatewayManager {
} }
}); });
let manager = Arc::clone(self); let router = Arc::clone(self);
event_bus.on_async::<MessageDeletedEvent, _, _>("message_deleted", move |event| { event_bus.on_async::<MessageDeletedEvent, _, _>("message_deleted", move |event| {
let manager = Arc::clone(&manager); let router = Arc::clone(&router);
async move { async move {
manager.broadcast_message(event.channel_id, "remove", event.message.id); let users = router
.services
.realtime_registry
.users_for_channel(event.channel_id);
router
.gateway
.send_to_users(users, "Message", "remove", event.message.id);
} }
}); });
let manager = Arc::clone(self); let router = Arc::clone(self);
event_bus.on_async::<MessageReactionAddedEvent, _, _>( event_bus.on_async::<MessageReactionAddedEvent, _, _>(
"message_reaction_added", "message_reaction_added",
move |event| { move |event| {
let manager = Arc::clone(&manager); let router = Arc::clone(&router);
async move { async move {
manager.broadcast_reaction( let users = router
event.channel_id, .services
.realtime_registry
.users_for_channel(event.channel_id);
router.gateway.send_to_users(
users,
"Reaction",
"add", "add",
reaction_model_to_response(event.reaction), reaction_model_to_response(event.reaction),
); );
@@ -128,99 +202,248 @@ impl GatewayManager {
}, },
); );
let manager = Arc::clone(self); let router = Arc::clone(self);
event_bus.on_async::<MessageReactionRemovedEvent, _, _>( event_bus.on_async::<MessageReactionRemovedEvent, _, _>(
"message_reaction_removed", "message_reaction_removed",
move |event| { move |event| {
let manager = Arc::clone(&manager); let router = Arc::clone(&router);
async move { async move {
manager.broadcast_reaction( let users = router
event.channel_id, .services
.realtime_registry
.users_for_channel(event.channel_id);
router.gateway.send_to_users(
users,
"Reaction",
"remove", "remove",
reaction_model_to_response(event.reaction), reaction_model_to_response(event.reaction),
); );
} }
}, },
); );
}
pub(crate) fn add_client(&self, gateway_client: GatewayClient) { let router = Arc::clone(self);
let key = gateway_client.key(); event_bus.on_async::<ChannelCreatedEvent, _, _>("channel_created", move |event| {
self.clients.write().insert(key, gateway_client); let router = Arc::clone(&router);
} async move {
if let Err(error) = router
pub(crate) fn remove_client(&self, gateway_client: &GatewayClient) { .services
let key = gateway_client.key(); .realtime_registry
self.clients.write().remove(&key); .refresh_channel(&router.repositories, event.channel.id)
} .await
{
fn broadcast_message<T: serde::Serialize>( tracing::error!(channel_id = %event.channel.id, ?error, "Unable to refresh channel audience");
&self, return;
channel_id: Uuid, }
action: &'static str, let users = router
content: T, .services
) { .realtime_registry
let event = GatewayEvent { .users_for_channel(event.channel.id);
namespace: "Message", router.gateway.send_to_users(
action, users,
content, "Channel",
}; "add",
let Ok(json) = serde_json::to_string(&event) else { channel_model_to_channel_response(event.channel),
return; );
};
let users = self
.services
.realtime_registry
.users_for_channel(channel_id);
let clients = self.clients.read();
for (key, client) in clients.iter() {
if users.contains(&key.user_id) {
let _ = client.sender.send(Message::Text(json.clone().into()));
} }
} });
let router = Arc::clone(self);
event_bus.on_async::<ChannelUpdatedEvent, _, _>("channel_updated", move |event| {
let router = Arc::clone(&router);
async move {
let mut users = router
.services
.realtime_registry
.users_for_channel(event.channel.id);
users.extend(
router
.services
.realtime_registry
.users_for_channel(event.previous.id),
);
router.gateway.send_to_users(
users,
"Channel",
"update",
channel_model_to_channel_response(event.channel),
);
}
});
let router = Arc::clone(self);
event_bus.on_async::<ChannelDeletedEvent, _, _>("channel_deleted", move |event| {
let router = Arc::clone(&router);
async move {
let users = router
.services
.realtime_registry
.users_for_channel(event.channel.id);
router
.gateway
.send_to_users(users, "Channel", "remove", event.channel.id);
router
.services
.realtime_registry
.remove_channel(event.channel.id);
}
});
let router = Arc::clone(self);
event_bus.on_async::<category::Model, _, _>("category_created", move |category| {
let router = Arc::clone(&router);
async move {
let users = router.server_users(category.server_id).await;
router.gateway.send_to_users(
users,
"Category",
"add",
category_model_to_category_response(category),
);
}
});
let router = Arc::clone(self);
event_bus.on_async::<category::Model, _, _>("category_updated", move |category| {
let router = Arc::clone(&router);
async move {
let users = router.server_users(category.server_id).await;
router.gateway.send_to_users(
users,
"Category",
"update",
category_model_to_category_response(category),
);
}
});
let router = Arc::clone(self);
event_bus.on_async::<category::Model, _, _>("category_deleted", move |category| {
let router = Arc::clone(&router);
async move {
let users = router.server_users(category.server_id).await;
router
.gateway
.send_to_users(users, "Category", "remove", category.id);
}
});
let router = Arc::clone(self);
event_bus.on_async::<server::Model, _, _>("server_created", move |server| {
let router = Arc::clone(&router);
async move {
let users = router.server_users(server.id).await;
router.gateway.send_to_users(
users,
"Server",
"add",
server_model_to_server_response(server),
);
}
});
let router = Arc::clone(self);
event_bus.on_async::<server::Model, _, _>("server_updated", move |server| {
let router = Arc::clone(&router);
async move {
let users = router.server_users(server.id).await;
router.gateway.send_to_users(
users,
"Server",
"update",
server_model_to_server_response(server),
);
}
});
let router = Arc::clone(self);
event_bus.on_async::<(server::Model, Vec<Uuid>), _, _>(
"server_deleted",
move |(server, users)| {
let router = Arc::clone(&router);
async move {
router
.gateway
.send_to_users(users, "Server", "remove", server.id);
}
},
);
let router = Arc::clone(self);
event_bus.on_async::<EmojiCreatedEvent, _, _>("emoji_created", move |event| {
let router = Arc::clone(&router);
async move {
let users = router.emoji_users(event.emoji.server_id).await;
router.gateway.send_to_users(
users,
"Emoji",
"add",
crate::routes::emoji::mapper::response(event.emoji),
);
}
});
let router = Arc::clone(self);
event_bus.on_async::<EmojiUpdatedEvent, _, _>("emoji_updated", move |event| {
let router = Arc::clone(&router);
async move {
let mut users = router.emoji_users(event.emoji.server_id).await;
users.extend(router.emoji_users(event.previous.server_id).await);
router.gateway.send_to_users(
users,
"Emoji",
"update",
crate::routes::emoji::mapper::response(event.emoji),
);
}
});
let router = Arc::clone(self);
event_bus.on_async::<EmojiDeletedEvent, _, _>("emoji_deleted", move |event| {
let router = Arc::clone(&router);
async move {
let users = router.emoji_users(event.emoji.server_id).await;
router.gateway.send_to_users(
users,
"Emoji",
"remove",
crate::routes::emoji::mapper::response(event.emoji),
);
}
});
} }
fn broadcast_reaction<T: serde::Serialize>( async fn server_users(&self, server_id: Uuid) -> Vec<Uuid> {
&self, server_user::Entity::find()
channel_id: Uuid, .filter(server_user::Column::ServerId.eq(server_id))
action: &'static str, .select_only()
content: T, .column(server_user::Column::UserId)
) { .into_tuple::<Uuid>()
let event = GatewayEvent { .all(&self.repositories.server.context.db)
namespace: "Reaction", .await
action, .unwrap_or_default()
content, }
};
let Ok(json) = serde_json::to_string(&event) else {
return;
};
let users = self async fn emoji_users(&self, server_id: Option<Uuid>) -> Vec<Uuid> {
.services match server_id {
.realtime_registry Some(server_id) => self.server_users(server_id).await,
.users_for_channel(channel_id); None => user::Entity::find()
let clients = self.clients.read(); .select_only()
for (key, client) in clients.iter() { .column(user::Column::Id)
if users.contains(&key.user_id) { .into_tuple::<Uuid>()
let _ = client.sender.send(Message::Text(json.clone().into())); .all(&self.repositories.server.context.db)
} .await
.unwrap_or_default(),
} }
} }
} }
impl GatewayClient { impl GatewayClient {
pub fn new( pub fn new(user: crate::models::user::Model, sender: mpsc::UnboundedSender<Message>) -> Self {
user: User,
sender: mpsc::UnboundedSender<Message>,
event_bus: Arc<EventBus>,
) -> Self {
Self { Self {
user, user,
connection_id: Uuid::new_v4(), connection_id: Uuid::new_v4(),
sender, sender,
event_bus,
_event_handles: Vec::new(),
} }
} }
@@ -231,113 +454,16 @@ impl GatewayClient {
} }
} }
fn subscribe_event<T, F, R>(
&self,
event_name: &'static str,
namespace: &'static str,
action: &'static str,
mapper: F,
) -> Arc<JoinHandle<()>>
where
T: Clone + Send + Sync + 'static,
R: serde::Serialize + Send + 'static,
F: Fn(T) -> R + Send + Sync + 'static,
{
let sender = self.sender.clone();
Arc::new(
self.event_bus
.on_async::<T, _, _>(event_name, move |payload| {
let sender = sender.clone();
let content = mapper(payload);
async move {
let event = GatewayEvent {
namespace,
action,
content,
};
if let Ok(json) = serde_json::to_string(&event) {
let _ = sender.send(Message::Text(json.into()));
}
}
}),
)
}
pub fn subscribe_to_events(&mut self) {
let mut handles = Vec::new();
// Les messages sont routés par GatewayManager selon le channel_id.
handles.push(self.subscribe_event(
"channel_created",
"Channel",
"add",
channel_model_to_channel_response,
));
handles.push(self.subscribe_event(
"channel_updated",
"Channel",
"update",
channel_model_to_channel_response,
));
handles.push(self.subscribe_event("channel_deleted", "Channel", "remove", |id: Uuid| id));
handles.push(self.subscribe_event(
"category_created",
"Category",
"add",
category_model_to_category_response,
));
handles.push(self.subscribe_event(
"category_updated",
"Category",
"update",
category_model_to_category_response,
));
handles.push(self.subscribe_event("category_deleted", "Category", "remove", |id: Uuid| id));
handles.push(self.subscribe_event(
"server_created",
"Server",
"add",
server_model_to_server_response,
));
handles.push(self.subscribe_event(
"server_updated",
"Server",
"update",
server_model_to_server_response,
));
handles.push(self.subscribe_event("server_deleted", "Server", "remove", |id: Uuid| id));
self._event_handles = handles;
}
pub fn unsubscribe_all(&mut self) {
for handle in self._event_handles.drain(..) {
handle.abort();
}
}
async fn on_connect(&mut self) { async fn on_connect(&mut self) {
tracing::info!(user_id = %self.user.id, "Client connected"); tracing::info!(user_id = %self.user.id, "Client connected");
self.subscribe_to_events();
} }
async fn on_disconnect(&mut self) { async fn on_disconnect(&mut self) {
tracing::info!(user_id = %self.user.id, "Client disconnected"); tracing::info!(user_id = %self.user.id, "Client disconnected");
self.unsubscribe_all();
} }
async fn on_message(&self, message: Message) { async fn on_message(&self, message: Message) {
match message { if let Message::Text(content) = message {
Message::Binary(_) => {} tracing::info!(user_id = %self.user.id, "Received text message: {}", content);
Message::Text(content) => {
tracing::info!(user_id = %self.user.id, "Received text message: {}", content);
}
Message::Ping(_) => {}
Message::Pong(_) => {}
Message::Close(_) => {}
} }
} }
} }
+1 -1
View File
@@ -1,7 +1,7 @@
use super::handlers; use super::handlers;
use crate::core::AppState; use crate::core::AppState;
use axum::routing::get;
use axum::Router; use axum::Router;
use axum::routing::get;
pub fn router() -> Router<AppState> { pub fn router() -> Router<AppState> {
Router::new().route("/gateway", get(handlers::ws_handler)) Router::new().route("/gateway", get(handlers::ws_handler))
+6 -1
View File
@@ -87,6 +87,11 @@ impl CategoryService {
let txn = db.begin().await?; let txn = db.begin().await?;
let existing = category::Entity::find_by_id(id)
.one(&txn)
.await?
.ok_or_else(|| anyhow::anyhow!("Category not found"))?;
self.service_context self.service_context
.services .services
.get() .get()
@@ -102,7 +107,7 @@ impl CategoryService {
txn.commit().await?; txn.commit().await?;
if deleted { if deleted {
event_bus.emit("category_deleted", id); event_bus.emit("category_deleted", existing);
} }
Ok(deleted) Ok(deleted)
+23 -3
View File
@@ -1,4 +1,7 @@
use crate::domain::dto::channel::{CreateChannelRequest, UpdateChannelRequest}; use crate::domain::dto::channel::{CreateChannelRequest, UpdateChannelRequest};
use crate::domain::events::channel::{
ChannelCreatedEvent, ChannelDeletedEvent, ChannelUpdatedEvent,
};
use crate::models::server_item_order::OrderedResourceType; use crate::models::server_item_order::OrderedResourceType;
use crate::models::{channel, role}; use crate::models::{channel, role};
use crate::permissions::PermissionSet; use crate::permissions::PermissionSet;
@@ -88,7 +91,12 @@ impl ChannelService {
.await?; .await?;
// Post-commit event emission // Post-commit event emission
event_bus.emit("channel_created", channel.clone()); event_bus.emit(
"channel_created",
ChannelCreatedEvent {
channel: channel.clone(),
},
);
Ok(channel) Ok(channel)
} }
@@ -108,6 +116,7 @@ impl ChannelService {
.await? .await?
.ok_or_else(|| anyhow::anyhow!("Channel not found"))?; .ok_or_else(|| anyhow::anyhow!("Channel not found"))?;
let previous = existing.clone();
let mut active: channel::ActiveModel = existing.into(); let mut active: channel::ActiveModel = existing.into();
active.server_id = Set(payload.server_id); active.server_id = Set(payload.server_id);
active.category_id = Set(payload.category_id); active.category_id = Set(payload.category_id);
@@ -132,7 +141,13 @@ impl ChannelService {
txn.commit().await?; txn.commit().await?;
event_bus.emit("channel_updated", channel.clone()); event_bus.emit(
"channel_updated",
ChannelUpdatedEvent {
previous,
channel: channel.clone(),
},
);
Ok(channel) Ok(channel)
} }
@@ -143,6 +158,11 @@ impl ChannelService {
let txn = db.begin().await?; let txn = db.begin().await?;
let existing = channel::Entity::find_by_id(id)
.one(&txn)
.await?
.ok_or_else(|| anyhow::anyhow!("Channel not found"))?;
self.service_context self.service_context
.services .services
.get() .get()
@@ -158,7 +178,7 @@ impl ChannelService {
txn.commit().await?; txn.commit().await?;
if deleted { if deleted {
event_bus.emit("channel_deleted", id); event_bus.emit("channel_deleted", ChannelDeletedEvent { channel: existing });
} }
Ok(deleted) Ok(deleted)
+48 -6
View File
@@ -1,3 +1,5 @@
use crate::domain::events::channel::{ChannelCreatedEvent, ChannelDeletedEvent};
use crate::models::server;
use crate::repositories::Repositories; use crate::repositories::Repositories;
use crate::services::ServicesContext; use crate::services::ServicesContext;
use std::sync::Arc; use std::sync::Arc;
@@ -71,8 +73,8 @@ impl PermissionSyncService {
event_bus.on_async_with( event_bus.on_async_with(
"server_created", "server_created",
repositories.clone(), repositories.clone(),
move |repositories, server_id: Uuid| async move { move |repositories, server: server::Model| async move {
Self::sync_server(repositories, server_id).await; Self::sync_server(repositories, server.id).await;
}, },
); );
@@ -139,16 +141,20 @@ impl PermissionSyncService {
event_bus.on_async_with( event_bus.on_async_with(
"channel_created", "channel_created",
repositories.clone(), repositories.clone(),
move |repositories, server_id: Uuid| async move { move |repositories, event: ChannelCreatedEvent| async move {
Self::sync_server(repositories, server_id).await; if let Some(server_id) = event.channel.server_id {
Self::sync_server(repositories, server_id).await;
}
}, },
); );
event_bus.on_async_with( event_bus.on_async_with(
"channel_deleted", "channel_deleted",
repositories.clone(), repositories.clone(),
move |repositories, server_id: Uuid| async move { move |repositories, event: ChannelDeletedEvent| async move {
Self::sync_server(repositories, server_id).await; if let Some(server_id) = event.channel.server_id {
Self::sync_server(repositories, server_id).await;
}
}, },
); );
@@ -167,6 +173,42 @@ impl PermissionSyncService {
Self::sync_user(repositories, user_id, server_id).await; Self::sync_user(repositories, user_id, server_id).await;
}, },
); );
event_bus.on_async_with(
"channel_user_permission_created",
repositories.clone(),
move |repositories, (channel_id, user_id, _permissions): (Uuid, Uuid, u64)| async move {
if let Some(channel) = repositories
.channel
.get_by_id(channel_id)
.await
.ok()
.flatten()
{
if let Some(server_id) = channel.server_id {
Self::sync_user(repositories, user_id, server_id).await;
}
}
},
);
event_bus.on_async_with(
"channel_user_permission_deleted",
repositories,
move |repositories, (channel_id, user_id): (Uuid, Uuid)| async move {
if let Some(channel) = repositories
.channel
.get_by_id(channel_id)
.await
.ok()
.flatten()
{
if let Some(server_id) = channel.server_id {
Self::sync_user(repositories, user_id, server_id).await;
}
}
},
);
} }
// ------------------------------------------------------------------------- // -------------------------------------------------------------------------
+115 -34
View File
@@ -3,10 +3,10 @@ use crate::permissions::ChannelPermission;
use crate::repositories::Repositories; use crate::repositories::Repositories;
use event_bus::EventBus; use event_bus::EventBus;
use parking_lot::RwLock; use parking_lot::RwLock;
use sea_orm::{ColumnTrait, EntityTrait, QueryFilter};
use std::collections::{HashMap, HashSet}; use std::collections::{HashMap, HashSet};
use std::sync::Arc; use std::sync::Arc;
use uuid::Uuid; use uuid::Uuid;
use sea_orm::{ColumnTrait, EntityTrait, QueryFilter};
/// In-memory index of the users that can receive events for each channel. /// In-memory index of the users that can receive events for each channel.
#[derive(Debug, Default)] #[derive(Debug, Default)]
@@ -16,6 +16,58 @@ pub struct RealtimeRegistry {
} }
impl RealtimeRegistry { impl RealtimeRegistry {
/// Rebuilds one channel audience after a committed structural change.
pub async fn refresh_channel(
&self,
repositories: &Repositories,
channel_id: Uuid,
) -> anyhow::Result<()> {
let channel = channel::Entity::find_by_id(channel_id)
.one(&repositories.channel.context.db)
.await?
.ok_or_else(|| anyhow::anyhow!("Channel not found"))?;
if channel.channel_type == channel::ChannelType::DM {
let users = channel_user::Entity::find()
.filter(channel_user::Column::ChannelId.eq(channel_id))
.all(&repositories.channel.context.db)
.await?
.into_iter()
.map(|member| member.user_id);
self.set_channel_users(channel_id, users);
} else {
let permissions = repositories.computed_permission.get_all().await?;
self.set_channel_users(
channel_id,
permissions.into_iter().filter_map(|permission| {
(permission.scope_type == PermissionScopeType::Channel
&& permission.resource_id == channel_id
&& ChannelPermission::from_bits_retain(permission.permissions as u64)
.contains(ChannelPermission::READ_CHANNEL))
.then_some(permission.user_id)
}),
);
}
Ok(())
}
async fn refresh_user(&self, repositories: &Repositories, user_id: Uuid) -> anyhow::Result<()> {
let channels = repositories
.computed_permission
.get_all()
.await?
.into_iter()
.filter(|permission| {
permission.user_id == user_id
&& permission.scope_type == PermissionScopeType::Channel
&& ChannelPermission::from_bits_retain(permission.permissions as u64)
.contains(ChannelPermission::READ_CHANNEL)
})
.map(|permission| permission.resource_id);
self.set_user_channels(user_id, channels);
Ok(())
}
pub async fn initialize(&self, repositories: &Repositories) -> anyhow::Result<()> { pub async fn initialize(&self, repositories: &Repositories) -> anyhow::Result<()> {
let permissions = repositories.computed_permission.get_all().await?; let permissions = repositories.computed_permission.get_all().await?;
let mut channel_users = HashMap::<Uuid, HashSet<Uuid>>::new(); let mut channel_users = HashMap::<Uuid, HashSet<Uuid>>::new();
@@ -51,8 +103,14 @@ impl RealtimeRegistry {
.all(&repositories.channel.context.db) .all(&repositories.channel.context.db)
.await?; .await?;
for member in members { for member in members {
channel_users.entry(member.channel_id).or_default().insert(member.user_id); channel_users
user_channels.entry(member.user_id).or_default().insert(member.channel_id); .entry(member.channel_id)
.or_default()
.insert(member.user_id);
user_channels
.entry(member.user_id)
.or_default()
.insert(member.channel_id);
} }
} }
@@ -73,13 +131,17 @@ impl RealtimeRegistry {
let users: HashSet<_> = users.into_iter().collect(); let users: HashSet<_> = users.into_iter().collect();
let old = { let old = {
let mut by_channel = self.channel_users.write(); let mut by_channel = self.channel_users.write();
by_channel.insert(channel_id, users.clone()).unwrap_or_default() by_channel
.insert(channel_id, users.clone())
.unwrap_or_default()
}; };
let mut by_user = self.user_channels.write(); let mut by_user = self.user_channels.write();
for user_id in old.difference(&users) { for user_id in old.difference(&users) {
if let Some(channels) = by_user.get_mut(user_id) { if let Some(channels) = by_user.get_mut(user_id) {
channels.remove(&channel_id); channels.remove(&channel_id);
if channels.is_empty() { by_user.remove(user_id); } if channels.is_empty() {
by_user.remove(user_id);
}
} }
} }
for user_id in users { for user_id in users {
@@ -148,21 +210,36 @@ impl RealtimeRegistry {
move |repositories, (_channel_id, user_id, _permissions): (Uuid, Uuid, u64)| { move |repositories, (_channel_id, user_id, _permissions): (Uuid, Uuid, u64)| {
let registry = Arc::clone(&registry); let registry = Arc::clone(&registry);
async move { async move {
match repositories.computed_permission.get_all().await { if let Err(error) = registry.refresh_user(&repositories, user_id).await {
Ok(all) => registry.set_user_channels( tracing::error!(%user_id, ?error, "Unable to refresh realtime registry")
user_id, }
all.into_iter() }
.filter(|p| { },
p.user_id == user_id );
&& p.scope_type == PermissionScopeType::Channel
&& ChannelPermission::from_bits_retain(p.permissions as u64) let registry = Arc::clone(self);
.contains(ChannelPermission::READ_CHANNEL) event_bus.on_async_with(
}) "channel_user_permission_created",
.map(|p| p.resource_id), repositories.clone(),
), move |repositories, (_channel_id, user_id, _permissions): (Uuid, Uuid, u64)| {
Err(error) => { let registry = Arc::clone(&registry);
tracing::error!(%user_id, ?error, "Unable to refresh realtime registry") async move {
} if let Err(error) = registry.refresh_user(&repositories, user_id).await {
tracing::error!(%user_id, ?error, "Unable to refresh realtime registry")
}
}
},
);
let registry = Arc::clone(self);
event_bus.on_async_with(
"channel_user_permission_deleted",
repositories.clone(),
move |repositories, (_channel_id, user_id): (Uuid, Uuid)| {
let registry = Arc::clone(&registry);
async move {
if let Err(error) = registry.refresh_user(&repositories, user_id).await {
tracing::error!(%user_id, ?error, "Unable to refresh realtime registry")
} }
} }
}, },
@@ -176,18 +253,8 @@ impl RealtimeRegistry {
move |repositories, (_server_id, user_id): (Uuid, Uuid)| { move |repositories, (_server_id, user_id): (Uuid, Uuid)| {
let registry = Arc::clone(&registry); let registry = Arc::clone(&registry);
async move { async move {
if let Ok(all) = repositories.computed_permission.get_all().await { if let Err(error) = registry.refresh_user(&repositories, user_id).await {
registry.set_user_channels( tracing::error!(%user_id, ?error, "Unable to refresh realtime registry")
user_id,
all.into_iter()
.filter(|p| {
p.user_id == user_id
&& p.scope_type == PermissionScopeType::Channel
&& ChannelPermission::from_bits_retain(p.permissions as u64)
.contains(ChannelPermission::READ_CHANNEL)
})
.map(|p| p.resource_id),
);
} }
} }
}, },
@@ -210,8 +277,22 @@ mod tests {
registry.set_channel_users(channel_id, [first, second]); registry.set_channel_users(channel_id, [first, second]);
assert_eq!(registry.users_for_channel(channel_id).len(), 2); assert_eq!(registry.users_for_channel(channel_id).len(), 2);
assert!(registry.user_channels.read().get(&first).unwrap().contains(&channel_id)); assert!(
assert!(registry.user_channels.read().get(&second).unwrap().contains(&channel_id)); registry
.user_channels
.read()
.get(&first)
.unwrap()
.contains(&channel_id)
);
assert!(
registry
.user_channels
.read()
.get(&second)
.unwrap()
.contains(&channel_id)
);
registry.set_channel_users(channel_id, [second]); registry.set_channel_users(channel_id, [second]);
assert!(!registry.user_channels.read().contains_key(&first)); assert!(!registry.user_channels.read().contains_key(&first));
+14 -2
View File
@@ -1,4 +1,4 @@
use crate::models::{role, server}; use crate::models::{role, server, server_user};
use crate::repositories::Repositories; use crate::repositories::Repositories;
use crate::services::ServicesContext; use crate::services::ServicesContext;
use sea_orm::{ use sea_orm::{
@@ -85,6 +85,18 @@ impl ServerService {
let txn = db.begin().await?; let txn = db.begin().await?;
let existing = server::Entity::find_by_id(id)
.one(&txn)
.await?
.ok_or_else(|| anyhow::anyhow!("Server not found"))?;
let audience = server_user::Entity::find()
.filter(server_user::Column::ServerId.eq(id))
.all(&txn)
.await?
.into_iter()
.map(|member| member.user_id)
.collect::<Vec<_>>();
let res = server::Entity::delete_by_id(id).exec(&txn).await?; let res = server::Entity::delete_by_id(id).exec(&txn).await?;
let deleted = res.rows_affected > 0; let deleted = res.rows_affected > 0;
@@ -92,7 +104,7 @@ impl ServerService {
txn.commit().await?; txn.commit().await?;
if deleted { if deleted {
event_bus.emit("server_deleted", id); event_bus.emit("server_deleted", (existing, audience));
} }
Ok(deleted) Ok(deleted)
+8 -4
View File
@@ -13,10 +13,14 @@ pub struct ScopedGuard<'a, K: Eq + Hash + Clone> {
impl<'a, K: Eq + Hash + Clone> Drop for ScopedGuard<'a, K> { impl<'a, K: Eq + Hash + Clone> Drop for ScopedGuard<'a, K> {
fn drop(&mut self) { fn drop(&mut self) {
// Optionnel : Nettoyage des Mutex orphelins dans la HashMap scopée // `Drop` peut être exécuté sur une tâche Tokio. `blocking_lock` y
let mut map = self.manager.scopes.blocking_lock(); // panique ; le nettoyage est opportuniste et sera retenté au prochain
if let Some(weak) = map.get(&self.key) { // `lock_scope` si le mutex de la table est momentanément occupé.
if weak.strong_count() == 0 { if let Ok(mut map) = self.manager.scopes.try_lock() {
if map
.get(&self.key)
.is_some_and(|weak| weak.strong_count() <= 1)
{
map.remove(&self.key); map.remove(&self.key);
} }
} }