init
This commit is contained in:
+18
-5
@@ -255,6 +255,14 @@ impl ChannelService {
|
||||
) -> Result<(), anyhow::Error> {
|
||||
let db = &self.service_context.repositories.server.context.db;
|
||||
let event_bus = &self.service_context.event_bus;
|
||||
let server_id = self
|
||||
.service_context
|
||||
.repositories
|
||||
.channel
|
||||
.get_by_id(channel_id)
|
||||
.await?
|
||||
.and_then(|channel| channel.server_id)
|
||||
.ok_or_else(|| anyhow::anyhow!("Server channel not found"))?;
|
||||
|
||||
let txn = db.begin().await?;
|
||||
|
||||
@@ -279,10 +287,7 @@ impl ChannelService {
|
||||
|
||||
txn.commit().await?;
|
||||
|
||||
event_bus.emit(
|
||||
"channel_role_permission_created",
|
||||
(channel_id, role_id, permissions),
|
||||
);
|
||||
event_bus.emit("channel_role_permission_updated", (role_id, server_id));
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -294,6 +299,14 @@ impl ChannelService {
|
||||
) -> Result<(), anyhow::Error> {
|
||||
let db = &self.service_context.repositories.server.context.db;
|
||||
let event_bus = &self.service_context.event_bus;
|
||||
let server_id = self
|
||||
.service_context
|
||||
.repositories
|
||||
.channel
|
||||
.get_by_id(channel_id)
|
||||
.await?
|
||||
.and_then(|channel| channel.server_id)
|
||||
.ok_or_else(|| anyhow::anyhow!("Server channel not found"))?;
|
||||
|
||||
let txn = db.begin().await?;
|
||||
|
||||
@@ -305,7 +318,7 @@ impl ChannelService {
|
||||
|
||||
txn.commit().await?;
|
||||
|
||||
event_bus.emit("channel_role_permission_deleted", (channel_id, role_id));
|
||||
event_bus.emit("channel_role_permission_updated", (role_id, server_id));
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
+140
-68
@@ -1,7 +1,9 @@
|
||||
use crate::domain::events::channel::{ChannelCreatedEvent, ChannelDeletedEvent};
|
||||
use crate::domain::events::server_tree::ServerTreeInvalidatedEvent;
|
||||
use crate::models::server;
|
||||
use crate::repositories::Repositories;
|
||||
use crate::services::ServicesContext;
|
||||
use event_bus::EventBus;
|
||||
use std::sync::Arc;
|
||||
use uuid::Uuid;
|
||||
// list of all events :
|
||||
@@ -57,6 +59,16 @@ pub struct PermissionSyncService {
|
||||
}
|
||||
|
||||
impl PermissionSyncService {
|
||||
fn invalidate_tree(event_bus: &Arc<EventBus>, server_id: Uuid, user_ids: Option<Vec<Uuid>>) {
|
||||
event_bus.emit(
|
||||
"server_tree_invalidated",
|
||||
ServerTreeInvalidatedEvent {
|
||||
server_id,
|
||||
user_ids,
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
pub fn new(service_context: Arc<ServicesContext>) -> Self {
|
||||
Self { service_context }
|
||||
}
|
||||
@@ -70,6 +82,7 @@ impl PermissionSyncService {
|
||||
// ---------------------------------------------------------------------
|
||||
// Événements Serveur & Membres Serveur
|
||||
// ---------------------------------------------------------------------
|
||||
let notify = event_bus.clone();
|
||||
event_bus.on_async_with(
|
||||
"server_created",
|
||||
repositories.clone(),
|
||||
@@ -81,16 +94,25 @@ impl PermissionSyncService {
|
||||
event_bus.on_async_with(
|
||||
"server_user_created",
|
||||
repositories.clone(),
|
||||
move |repositories, (server_id, user_id): (Uuid, Uuid)| async move {
|
||||
Self::sync_user(repositories, user_id, server_id).await;
|
||||
move |repositories, (server_id, user_id): (Uuid, Uuid)| {
|
||||
let notify = notify.clone();
|
||||
async move {
|
||||
Self::sync_user(repositories, user_id, server_id).await;
|
||||
Self::invalidate_tree(¬ify, server_id, Some(vec![user_id]));
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
let notify = event_bus.clone();
|
||||
event_bus.on_async_with(
|
||||
"server_user_deleted",
|
||||
repositories.clone(),
|
||||
move |repositories, (server_id, user_id): (Uuid, Uuid)| async move {
|
||||
Self::sync_user(repositories, user_id, server_id).await;
|
||||
move |repositories, (server_id, user_id): (Uuid, Uuid)| {
|
||||
let notify = notify.clone();
|
||||
async move {
|
||||
Self::sync_user(repositories, user_id, server_id).await;
|
||||
Self::invalidate_tree(¬ify, server_id, Some(vec![user_id]));
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
@@ -98,19 +120,29 @@ impl PermissionSyncService {
|
||||
// Événements Rôles & Membres de Rôle
|
||||
// ---------------------------------------------------------------------
|
||||
|
||||
let notify = event_bus.clone();
|
||||
event_bus.on_async_with(
|
||||
"role_user_created",
|
||||
repositories.clone(),
|
||||
move |repositories, (_role_id, user_id, server_id): (Uuid, Uuid, Uuid)| async move {
|
||||
Self::sync_user(repositories, user_id, server_id).await;
|
||||
move |repositories, (_role_id, user_id, server_id): (Uuid, Uuid, Uuid)| {
|
||||
let notify = notify.clone();
|
||||
async move {
|
||||
Self::sync_user(repositories, user_id, server_id).await;
|
||||
Self::invalidate_tree(¬ify, server_id, Some(vec![user_id]));
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
let notify = event_bus.clone();
|
||||
event_bus.on_async_with(
|
||||
"role_user_deleted",
|
||||
repositories.clone(),
|
||||
move |repositories, (_role_id, user_id, server_id): (Uuid, Uuid, Uuid)| async move {
|
||||
Self::sync_user(repositories, user_id, server_id).await;
|
||||
move |repositories, (_role_id, user_id, server_id): (Uuid, Uuid, Uuid)| {
|
||||
let notify = notify.clone();
|
||||
async move {
|
||||
Self::sync_user(repositories, user_id, server_id).await;
|
||||
Self::invalidate_tree(¬ify, server_id, Some(vec![user_id]));
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
@@ -118,19 +150,29 @@ impl PermissionSyncService {
|
||||
// Overrides de Permissions Serveur
|
||||
// ---------------------------------------------------------------------
|
||||
|
||||
let notify = event_bus.clone();
|
||||
event_bus.on_async_with(
|
||||
"server_role_permission_updated",
|
||||
repositories.clone(),
|
||||
move |repositories, (role_id, server_id): (Uuid, Uuid)| async move {
|
||||
Self::sync_role_members(repositories, role_id, server_id).await;
|
||||
move |repositories, (role_id, server_id): (Uuid, Uuid)| {
|
||||
let notify = notify.clone();
|
||||
async move {
|
||||
Self::sync_role_members(repositories, role_id, server_id).await;
|
||||
Self::invalidate_tree(¬ify, server_id, None);
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
let notify = event_bus.clone();
|
||||
event_bus.on_async_with(
|
||||
"server_user_permission_updated",
|
||||
repositories.clone(),
|
||||
move |repositories, (server_id, user_id): (Uuid, Uuid)| async move {
|
||||
Self::sync_user(repositories, user_id, server_id).await;
|
||||
move |repositories, (server_id, user_id): (Uuid, Uuid)| {
|
||||
let notify = notify.clone();
|
||||
async move {
|
||||
Self::sync_user(repositories, user_id, server_id).await;
|
||||
Self::invalidate_tree(¬ify, server_id, Some(vec![user_id]));
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
@@ -138,73 +180,103 @@ impl PermissionSyncService {
|
||||
// Événements Canaux & Overrides de Permissions Canaux
|
||||
// ---------------------------------------------------------------------
|
||||
|
||||
let notify = event_bus.clone();
|
||||
event_bus.on_async_with(
|
||||
"channel_created",
|
||||
repositories.clone(),
|
||||
move |repositories, event: ChannelCreatedEvent| async move {
|
||||
if let Some(server_id) = event.channel.server_id {
|
||||
Self::sync_server(repositories, server_id).await;
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
event_bus.on_async_with(
|
||||
"channel_deleted",
|
||||
repositories.clone(),
|
||||
move |repositories, event: ChannelDeletedEvent| async move {
|
||||
if let Some(server_id) = event.channel.server_id {
|
||||
Self::sync_server(repositories, server_id).await;
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
event_bus.on_async_with(
|
||||
"channel_role_permission_updated",
|
||||
repositories.clone(),
|
||||
move |repositories, (role_id, server_id): (Uuid, Uuid)| async move {
|
||||
Self::sync_role_members(repositories, role_id, server_id).await;
|
||||
},
|
||||
);
|
||||
|
||||
event_bus.on_async_with(
|
||||
"channel_user_permission_updated",
|
||||
repositories.clone(),
|
||||
move |repositories, (server_id, user_id): (Uuid, Uuid)| async move {
|
||||
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;
|
||||
move |repositories, event: ChannelCreatedEvent| {
|
||||
let notify = notify.clone();
|
||||
async move {
|
||||
if let Some(server_id) = event.channel.server_id {
|
||||
Self::sync_server(repositories, server_id).await;
|
||||
Self::invalidate_tree(¬ify, server_id, None);
|
||||
}
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
let notify = event_bus.clone();
|
||||
event_bus.on_async_with(
|
||||
"channel_deleted",
|
||||
repositories.clone(),
|
||||
move |repositories, event: ChannelDeletedEvent| {
|
||||
let notify = notify.clone();
|
||||
async move {
|
||||
if let Some(server_id) = event.channel.server_id {
|
||||
Self::sync_server(repositories, server_id).await;
|
||||
Self::invalidate_tree(¬ify, server_id, None);
|
||||
}
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
let notify = event_bus.clone();
|
||||
event_bus.on_async_with(
|
||||
"channel_role_permission_updated",
|
||||
repositories.clone(),
|
||||
move |repositories, (role_id, server_id): (Uuid, Uuid)| {
|
||||
let notify = notify.clone();
|
||||
async move {
|
||||
Self::sync_role_members(repositories, role_id, server_id).await;
|
||||
Self::invalidate_tree(¬ify, server_id, None);
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
let notify = event_bus.clone();
|
||||
event_bus.on_async_with(
|
||||
"channel_user_permission_updated",
|
||||
repositories.clone(),
|
||||
move |repositories, (server_id, user_id): (Uuid, Uuid)| {
|
||||
let notify = notify.clone();
|
||||
async move {
|
||||
Self::sync_user(repositories, user_id, server_id).await;
|
||||
Self::invalidate_tree(¬ify, server_id, Some(vec![user_id]));
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
let notify = event_bus.clone();
|
||||
event_bus.on_async_with(
|
||||
"channel_user_permission_created",
|
||||
repositories.clone(),
|
||||
move |repositories, (channel_id, user_id, _permissions): (Uuid, Uuid, u64)| {
|
||||
let notify = notify.clone();
|
||||
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;
|
||||
Self::invalidate_tree(¬ify, server_id, Some(vec![user_id]));
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
let notify = event_bus.clone();
|
||||
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;
|
||||
move |repositories, (channel_id, user_id): (Uuid, Uuid)| {
|
||||
let notify = notify.clone();
|
||||
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;
|
||||
Self::invalidate_tree(¬ify, server_id, Some(vec![user_id]));
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
|
||||
Reference in New Issue
Block a user