Compare commits

...
2 Commits
Author SHA1 Message Date
Nell c96101ec3d init 2026-07-13 18:46:19 +02:00
Nell dc2be94a8d init 2026-07-13 16:16:00 +02:00
16 changed files with 157 additions and 574 deletions
Generated
+1 -1
View File
@@ -1384,10 +1384,10 @@ name = "event_bus"
version = "0.1.0"
dependencies = [
"criterion",
"glob",
"parking_lot",
"tokio",
"tracing",
"uuid",
]
[[package]]
+2 -2
View File
@@ -14,10 +14,10 @@ harness = false
[dependencies]
tokio = { version = "1.52.3", default-features = false, features = ["rt", "sync"] }
glob = "0.3.3"
parking_lot = "0.12.5"
tracing = "0.1"
uuid = { version = "1.23.5", features = ["v4"] }
[dev-dependencies]
tokio = { version = "1.52.3", default-features = false, features = ["rt", "rt-multi-thread", "macros", "time", "sync"] }
criterion = { version = "0.8.2", features = ["async_tokio"] }
criterion = { version = "0.8.2", features = ["async_tokio"] }
-134
View File
@@ -9,7 +9,6 @@ use event_bus::EventBus;
use tokio::runtime::Runtime;
const TOPIC: &str = "bench-topic";
const PATTERN: &str = "bench-*";
#[derive(Clone)]
struct SmallEvent {
@@ -473,82 +472,6 @@ fn bench_typed_callback(c: &mut Criterion) {
group.finish();
}
fn bench_pattern_callback(c: &mut Criterion) {
let rt = runtime();
let mut group = c.benchmark_group("event_bus/pattern_callback");
group.throughput(Throughput::Elements(1));
group.bench_function("small_struct", |b| {
b.to_async(&rt).iter_custom(|iters| async move {
let bus = Arc::new(EventBus::with_capacity(iters as usize + 1024));
let received = Arc::new(AtomicU64::new(0));
let handler_count = Arc::clone(&received);
let subscription = bus.on_pattern::<SmallEvent, _>(PATTERN, move |topic, event| {
let _ = topic;
let _ = event.value;
handler_count.fetch_add(1, Ordering::Relaxed);
});
let start = Instant::now();
for i in 0..iters {
bus.emit(TOPIC, SmallEvent { value: i });
}
wait_until_received(&received, iters).await;
let elapsed = start.elapsed();
subscription.abort();
elapsed
});
});
group.bench_function("arc_payload_1kb", |b| {
b.to_async(&rt).iter_custom(|iters| async move {
let bus = Arc::new(EventBus::with_capacity(iters as usize + 1024));
let payload: Arc<[u8]> = Arc::from(vec![7_u8; 1024].into_boxed_slice());
let received = Arc::new(AtomicU64::new(0));
let handler_count = Arc::clone(&received);
let subscription =
bus.on_pattern::<ArcPayloadEvent, _>(PATTERN, move |topic, event| {
let _ = topic;
let _ = event.id;
let _ = event.payload.len();
handler_count.fetch_add(1, Ordering::Relaxed);
});
let start = Instant::now();
for i in 0..iters {
bus.emit(
TOPIC,
ArcPayloadEvent {
id: i,
payload: Arc::clone(&payload),
},
);
}
wait_until_received(&received, iters).await;
let elapsed = start.elapsed();
subscription.abort();
elapsed
});
});
group.finish();
}
fn bench_multiple_subscribers(c: &mut Criterion) {
let rt = runtime();
@@ -597,69 +520,12 @@ fn bench_multiple_subscribers(c: &mut Criterion) {
group.finish();
}
fn bench_multiple_patterns(c: &mut Criterion) {
let rt = runtime();
let mut group = c.benchmark_group("event_bus/multiple_patterns");
group.throughput(Throughput::Elements(1));
for pattern_count in [1_u64, 4, 16, 64, 256] {
group.bench_function(format!("{pattern_count}_patterns_one_match"), |b| {
b.to_async(&rt).iter_custom(|iters| async move {
let bus = Arc::new(EventBus::with_capacity(iters as usize + 1024));
let received = Arc::new(AtomicU64::new(0));
let mut subscriptions = Vec::with_capacity(pattern_count as usize);
for index in 0..pattern_count {
let pattern = if index == 0 {
PATTERN.to_string()
} else {
format!("unused-{index}-*")
};
let handler_count = Arc::clone(&received);
let subscription =
bus.on_pattern::<SmallEvent, _>(&pattern, move |topic, event| {
let _ = topic;
let _ = event.value;
handler_count.fetch_add(1, Ordering::Relaxed);
});
subscriptions.push(subscription);
}
let start = Instant::now();
for i in 0..iters {
bus.emit(TOPIC, SmallEvent { value: i });
}
wait_until_received(&received, iters).await;
let elapsed = start.elapsed();
for subscription in subscriptions {
subscription.abort();
}
elapsed
});
});
}
group.finish();
}
criterion_group!(
benches,
bench_emit_no_subscriber,
bench_raw_subscriber,
bench_typed_callback,
bench_pattern_callback,
bench_multiple_subscribers,
bench_multiple_patterns,
);
criterion_main!(benches);
+49 -182
View File
@@ -2,12 +2,13 @@ use std::any::Any;
use std::future::Future;
use std::sync::Arc;
use glob::Pattern;
use parking_lot::RwLock;
use std::collections::HashMap;
use tokio::sync::broadcast;
use tokio::task::JoinHandle;
use tracing::log::kv::{Key, Value};
use tracing::{debug, trace, warn};
use uuid::Uuid;
/// Raw event type: an atomic reference-counted pointer to any value.
pub type AnyEvent = Arc<dyn Any + Send + Sync>;
@@ -15,6 +16,25 @@ pub type AnyEvent = Arc<dyn Any + Send + Sync>;
/// Default buffer capacity for each broadcast channel.
const DEFAULT_CAPACITY: usize = 64;
#[derive(Debug, Clone, PartialEq, Eq)]
enum ScopeValue {
String(String),
Uuid(Uuid),
}
impl ScopeValue {
fn into_string(self) -> String {
match self {
Self::String(value) => value,
Self::Uuid(value) => value.to_string(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct Scope {
key: String,
value: ScopeValue,
}
/// The central event bus.
///
/// Share it via `Arc<EventBus>` across modules.
@@ -60,37 +80,10 @@ const DEFAULT_CAPACITY: usize = 64;
/// # tokio::time::sleep(std::time::Duration::from_millis(10)).await;
/// # });
/// ```
///
/// # Example — glob pattern with topic
/// ```rust,no_run
/// use std::sync::Arc;
/// use oxspeak_server_lib::event_bus::EventBus;
///
/// #[derive(Clone, Debug)]
/// struct User { name: String }
///
/// # tokio_test::block_on(async {
/// let bus = Arc::new(EventBus::new());
///
/// bus.on_pattern::<User, _>("user-*", |topic, user| {
/// match topic.as_str() {
/// "user-created" => println!("Created : {:?}", user),
/// "user-deleted" => println!("Deleted : {:?}", user),
/// other => println!("{}: {:?}", other, user),
/// }
/// });
///
/// bus.emit("user-created", User { name: "Alice".into() });
/// bus.emit("user-deleted", User { name: "Bob".into() });
/// # tokio::time::sleep(std::time::Duration::from_millis(10)).await;
/// # });
/// ```
#[derive(Debug)]
pub struct EventBus {
/// Channels indexed by exact topic.
channels: RwLock<HashMap<String, broadcast::Sender<AnyEvent>>>,
/// Channels for glob-pattern subscriptions.
patterns: RwLock<Vec<(Pattern, broadcast::Sender<(String, AnyEvent)>)>>,
capacity: usize,
}
@@ -103,7 +96,6 @@ impl EventBus {
);
Self {
channels: RwLock::new(HashMap::new()),
patterns: RwLock::new(Vec::new()),
capacity: DEFAULT_CAPACITY,
}
}
@@ -113,7 +105,6 @@ impl EventBus {
debug!("EventBus created with capacity {}", capacity);
Self {
channels: RwLock::new(HashMap::new()),
patterns: RwLock::new(Vec::new()),
capacity,
}
}
@@ -151,7 +142,6 @@ impl EventBus {
/// Emits an event on a topic.
///
/// - Pushes the event into the exact-topic channel (if subscribers exist).
/// - Pushes the event into all glob-pattern channels that match the topic.
/// - If nobody is listening, the event is silently dropped.
///
/// # Example
@@ -167,7 +157,6 @@ impl EventBus {
trace!(topic, "Emitting event");
let event: AnyEvent = Arc::new(event);
// Exact-topic subscribers
if let Some(tx) = self.channels.read().get(topic) {
let receiver_count = tx.receiver_count();
let _ = tx.send(Arc::clone(&event));
@@ -176,20 +165,35 @@ impl EventBus {
receiver_count, "Event delivered to exact-topic channel"
);
}
}
// Glob-pattern subscribers
let patterns = self.patterns.read();
for (pattern, tx) in patterns.iter() {
if pattern.matches(topic) {
let receiver_count = tx.receiver_count();
let _ = tx.send((topic.to_string(), Arc::clone(&event)));
trace!(
topic,
pattern = pattern.as_str(),
receiver_count,
"Event delivered to pattern channel"
);
}
// todo : undocumented...
pub fn emit_scoped<T>(&self, topic: &str, scopes: impl IntoIterator<Item = Scope>, event: T)
where
T: Any + Send + Sync + 'static,
{
let event: AnyEvent = Arc::new(event);
// Émission sur le topic général.
self.emit_arc(topic, Arc::clone(&event));
// Émission sur chaque topic scoped.
for scope in scopes {
let scoped_topic = format!("{}:{}:{}", scope.key, scope.value.into_string(), topic);
self.emit_arc(&scoped_topic, Arc::clone(&event));
}
}
// todo : undocumented...
fn emit_arc(&self, topic: &str, event: AnyEvent) {
trace!(topic, "Emitting event");
if let Some(tx) = self.channels.read().get(topic) {
let receiver_count = tx.receiver_count();
let _ = tx.send(event);
trace!(topic, receiver_count, "Event delivered to channel");
}
}
@@ -307,143 +311,6 @@ impl EventBus {
})
}
/// Subscribes to all topics matching a glob pattern.
///
/// The handler receives `(topic, value)` — the topic name is included to
/// distinguish `user-created` from `user-deleted`, for example.
///
/// **No pre-registration required**: future topics are automatically covered.
/// Supports glob syntax: `*` (any sequence), `?` (one character),
/// `[abc]` (character class).
///
/// Returns a [`JoinHandle`] to cancel the subscription if needed.
///
/// # Example
/// ```rust,no_run
/// # use std::sync::Arc;
/// # use oxspeak_server_lib::event_bus::EventBus;
/// # #[derive(Clone, Debug)] struct User { name: String }
/// # let bus = Arc::new(EventBus::new());
/// bus.on_pattern::<User, _>("user-*", |topic, user| {
/// match topic.as_str() {
/// "user-created" => println!("Created : {:?}", user),
/// "user-deleted" => println!("Deleted : {:?}", user),
/// other => println!("{}: {:?}", other, user),
/// }
/// });
///
/// bus.emit("user-created", User { name: "Alice".into() });
/// bus.emit("user-deleted", User { name: "Bob".into() });
/// ```
pub fn on_pattern<T, F>(&self, pattern: &str, handler: F) -> JoinHandle<()>
where
T: Any + Send + Sync + Clone + 'static,
F: Fn(String, T) + Send + Sync + 'static,
{
let glob = Pattern::new(pattern).expect("invalid glob pattern");
let (tx, mut rx) = broadcast::channel(self.capacity);
self.patterns.write().push((glob, tx));
let pattern_owned = pattern.to_string();
debug!(pattern, "Sync pattern subscriber registered");
tokio::spawn(async move {
loop {
match rx.recv().await {
Ok((topic, evt)) => {
if let Some(typed) = evt.downcast_ref::<T>() {
trace!(
topic,
pattern = pattern_owned,
"Sync pattern handler invoked"
);
handler(topic, typed.clone());
}
}
Err(broadcast::error::RecvError::Lagged(n)) => {
warn!(
pattern = pattern_owned,
skipped = n,
"Pattern subscriber lagged, messages dropped"
);
}
Err(broadcast::error::RecvError::Closed) => {
debug!(
pattern = pattern_owned,
"Channel closed, sync pattern subscriber exiting"
);
break;
}
}
}
})
}
/// Subscribes to all topics matching a glob pattern, with an **async** handler.
///
/// Returns a [`JoinHandle`] to cancel the subscription if needed.
///
/// # Example
/// ```rust,no_run
/// # use std::sync::Arc;
/// # use oxspeak_server_lib::event_bus::EventBus;
/// # #[derive(Clone, Debug)] struct User { name: String }
/// # let bus = Arc::new(EventBus::new());
/// bus.on_pattern_async::<User, _, _>("user-*", |topic, user| async move {
/// match topic.as_str() {
/// "user-created" => println!("(async) Created : {:?}", user),
/// "user-deleted" => println!("(async) Deleted : {:?}", user),
/// other => println!("(async) {}: {:?}", other, user),
/// }
/// });
///
/// bus.emit("user-created", User { name: "Alice".into() });
/// ```
pub fn on_pattern_async<T, F, Fut>(&self, pattern: &str, handler: F) -> JoinHandle<()>
where
T: Any + Send + Sync + Clone + 'static,
F: Fn(String, T) -> Fut + Send + Sync + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
let glob = Pattern::new(pattern).expect("invalid glob pattern");
let (tx, mut rx) = broadcast::channel(self.capacity);
self.patterns.write().push((glob, tx));
let pattern_owned = pattern.to_string();
debug!(pattern, "Async pattern subscriber registered");
tokio::spawn(async move {
loop {
match rx.recv().await {
Ok((topic, evt)) => {
if let Some(typed) = evt.downcast_ref::<T>() {
trace!(
topic,
pattern = pattern_owned,
"Async pattern handler invoked"
);
handler(topic, typed.clone()).await;
}
}
Err(broadcast::error::RecvError::Lagged(n)) => {
warn!(
pattern = pattern_owned,
skipped = n,
"Pattern subscriber lagged, messages dropped"
);
}
Err(broadcast::error::RecvError::Closed) => {
debug!(
pattern = pattern_owned,
"Channel closed, async pattern subscriber exiting"
);
break;
}
}
}
})
}
// ─────────────────────────────────────────────────────────────────────────
// Subscription — low-level access (advanced use cases)
// ─────────────────────────────────────────────────────────────────────────
+2 -115
View File
@@ -1,115 +1,3 @@
//! # Event Bus
//!
//! An asynchronous event bus for routing typed messages between modules
//! without direct coupling.
//!
//! ## Features
//!
//! - **Key → event mapping**: each topic (`&str`) is independent
//! - **Type-unrestricted**: any `T: Any + Send + Sync + Clone`
//! - **Targeted wake-up**: only subscribers of the matching topic are woken up
//! - **Callback API**: JavaScript-style — `bus.on("topic", |payload| { ... })`
//! - **Glob pattern**: `on_pattern("user-*", |topic, payload| { ... })`
//! - **Async handlers**: `on_async` and `on_pattern_async`
//!
//! ## Example — sync callback
//!
//! ```rust,no_run
//! use std::sync::Arc;
//! use oxspeak_server_lib::event_bus::EventBus;
//!
//! #[derive(Clone, Debug)]
//! struct User { name: String }
//!
//! # tokio_test::block_on(async {
//! let bus = Arc::new(EventBus::new());
//!
//! bus.on::<User>("user-connected", |user| {
//! println!("Connected: {:?}", user);
//! });
//!
//! bus.emit("user-connected", User { name: "Alice".into() });
//! # tokio::time::sleep(std::time::Duration::from_millis(10)).await;
//! # });
//! ```
//!
//! ## Example — async callback
//!
//! ```rust,no_run
//! use std::sync::Arc;
//! use oxspeak_server_lib::event_bus::EventBus;
//!
//! #[derive(Clone, Debug)]
//! struct User { name: String }
//!
//! # tokio_test::block_on(async {
//! let bus = Arc::new(EventBus::new());
//!
//! bus.on_async::<User, _, _>("user-connected", |user| async move {
//! println!("(async) Connected: {:?}", user);
//! });
//!
//! bus.emit("user-connected", User { name: "Bob".into() });
//! # tokio::time::sleep(std::time::Duration::from_millis(10)).await;
//! # });
//! ```
//!
//! ## Example — glob pattern (topic included in the callback)
//!
//! ```rust,no_run
//! use std::sync::Arc;
//! use oxspeak_server_lib::event_bus::EventBus;
//!
//! #[derive(Clone, Debug)]
//! struct User { name: String }
//!
//! # tokio_test::block_on(async {
//! let bus = Arc::new(EventBus::new());
//!
//! bus.on_pattern::<User, _>("user-*", |topic, user| {
//! match topic.as_str() {
//! "user-created" => println!("Created : {:?}", user),
//! "user-deleted" => println!("Deleted : {:?}", user),
//! other => println!("{}: {:?}", other, user),
//! }
//! });
//!
//! bus.emit("user-created", User { name: "Alice".into() });
//! bus.emit("user-deleted", User { name: "Bob".into() });
//! # tokio::time::sleep(std::time::Duration::from_millis(10)).await;
//! # });
//! ```
//!
//! ## Example — multi-type with `match_event!` (advanced)
//!
//! ```rust,no_run
//! use std::sync::Arc;
//! use oxspeak_server_lib::event_bus::EventBus;
//! use oxspeak_server_lib::match_event;
//!
//! #[derive(Clone, Debug)] struct User { name: String }
//! #[derive(Clone, Debug)] struct UdpMetric { value: f32 }
//!
//! # tokio_test::block_on(async {
//! let bus = Arc::new(EventBus::new());
//! let mut rx = bus.on_raw("mixed-topic");
//!
//! bus.emit("mixed-topic", User { name: "Alice".into() });
//!
//! if let Ok(evt) = rx.recv().await {
//! match_event!(evt,
//! User => |u| println!("User: {:?}", u),
//! UdpMetric => |m| println!("Metric: {:?}", m),
//! );
//! }
//! # });
//! ```
mod bus;
// Public re-exports
pub use bus::{AnyEvent, EventBus};
/// Downcasts an [`AnyEvent`] to one or more concrete types and executes
/// the matching closure if the type matches.
///
@@ -154,9 +42,8 @@ macro_rules! match_event {
};
}
// ─────────────────────────────────────────────────────────────────────────────
// Tests
// ─────────────────────────────────────────────────────────────────────────────
mod bus;
pub use bus::{AnyEvent, EventBus};
#[cfg(test)]
mod tests;
+1 -102
View File
@@ -1,6 +1,7 @@
use crate::{match_event, EventBus};
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
use std::sync::Arc;
#[derive(Clone, Debug, PartialEq)]
struct User {
name: String,
@@ -126,108 +127,6 @@ async fn test_on_async_callback() {
assert!(received.load(Ordering::SeqCst));
}
// ── on_pattern ────────────────────────────────────────────────────────────
#[tokio::test]
async fn test_on_pattern_callback_receives_topic_and_payload() {
let bus = Arc::new(EventBus::new());
let created = Arc::new(AtomicBool::new(false));
let deleted = Arc::new(AtomicBool::new(false));
let c = Arc::clone(&created);
let d = Arc::clone(&deleted);
bus.on_pattern::<User, _>("user-*", move |topic, user| match topic.as_str() {
"user-created" if user.name == "Dave" => c.store(true, Ordering::SeqCst),
"user-deleted" if user.name == "Eve" => d.store(true, Ordering::SeqCst),
_ => {}
});
bus.emit(
"user-created",
User {
name: "Dave".into(),
},
);
bus.emit("user-deleted", User { name: "Eve".into() });
tokio::time::sleep(std::time::Duration::from_millis(30)).await;
assert!(created.load(Ordering::SeqCst));
assert!(deleted.load(Ordering::SeqCst));
}
#[tokio::test]
async fn test_on_pattern_does_not_match_other_topics() {
let bus = Arc::new(EventBus::new());
let called = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&called);
bus.on_pattern::<User, _>("user-*", move |_, _| {
flag.store(true, Ordering::SeqCst);
});
bus.emit(
"server-created",
User {
name: "Ghost".into(),
},
);
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
assert!(!called.load(Ordering::SeqCst));
}
#[tokio::test]
async fn test_on_pattern_no_pre_registration_needed() {
// Le pattern est enregistré avant que le topic n'existe
let bus = Arc::new(EventBus::new());
let received = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&received);
bus.on_pattern::<User, _>("user-*", move |topic, user| {
if topic == "user-new-topic" && user.name == "Frank" {
flag.store(true, Ordering::SeqCst);
}
});
bus.emit(
"user-new-topic",
User {
name: "Frank".into(),
},
);
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
assert!(received.load(Ordering::SeqCst));
}
// ── on_pattern_async ──────────────────────────────────────────────────────
#[tokio::test]
async fn test_on_pattern_async_callback() {
let bus = Arc::new(EventBus::new());
let received = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&received);
bus.on_pattern_async::<User, _, _>("user-*", move |topic, user| {
let f = Arc::clone(&flag);
async move {
if topic == "user-created" && user.name == "Hank" {
f.store(true, Ordering::SeqCst);
}
}
});
bus.emit(
"user-created",
User {
name: "Hank".into(),
},
);
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
assert!(received.load(Ordering::SeqCst));
}
// ── on_raw + match_event! ─────────────────────────────────────────────────
#[tokio::test]
+4 -16
View File
@@ -490,7 +490,7 @@ impl MigrationTrait for Migration {
.col(ColumnDef::new(Alias::new("server_id")).uuid().not_null())
.col(ColumnDef::new(Alias::new("role_id")).uuid().not_null())
.col(
ColumnDef::new(Alias::new("permission"))
ColumnDef::new(Alias::new("permissions"))
.big_integer()
.not_null()
.default(0),
@@ -537,7 +537,7 @@ impl MigrationTrait for Migration {
.col(ColumnDef::new(Alias::new("channel_id")).uuid().not_null())
.col(ColumnDef::new(Alias::new("role_id")).uuid().not_null())
.col(
ColumnDef::new(Alias::new("permission"))
ColumnDef::new(Alias::new("permissions"))
.big_integer()
.not_null()
.default(0),
@@ -584,7 +584,7 @@ impl MigrationTrait for Migration {
.col(ColumnDef::new(Alias::new("channel_id")).uuid().not_null())
.col(ColumnDef::new(Alias::new("user_id")).uuid().not_null())
.col(
ColumnDef::new(Alias::new("permission"))
ColumnDef::new(Alias::new("permissions"))
.big_integer()
.not_null()
.default(0),
@@ -635,19 +635,7 @@ impl MigrationTrait for Migration {
)
.col(ColumnDef::new(Alias::new("resource_id")).uuid().not_null())
.col(
ColumnDef::new(Alias::new("server_permissions"))
.big_integer()
.not_null()
.default(0),
)
.col(
ColumnDef::new(Alias::new("channel_permissions"))
.big_integer()
.not_null()
.default(0),
)
.col(
ColumnDef::new(Alias::new("voice_permissions"))
ColumnDef::new(Alias::new("permissions"))
.big_integer()
.not_null()
.default(0),
+12 -2
View File
@@ -1,6 +1,9 @@
mod permission_sync;
pub mod state;
use crate::config::AppConfig;
use crate::core::permission_sync::PermissionSyncService;
use crate::core::state::Services;
use crate::database::Database;
use crate::http::server::HttpServer;
use crate::metrics::{reporter, AppMetrics};
@@ -16,6 +19,7 @@ use uuid::Uuid;
pub struct App {
pub state: AppState,
pub services: Services,
}
impl App {
@@ -31,7 +35,7 @@ impl App {
let event_bus = Arc::new(EventBus::with_capacity(1024));
// Initialize shared repositories
let repositories = Repositories::new(db.clone(), event_bus.clone());
let repositories = Arc::new(Repositories::new(db.clone(), event_bus.clone()));
// Initialize gateway manager
let gateway = Arc::new(GatewayManager::default());
@@ -66,6 +70,8 @@ impl App {
let metrics = AppMetrics::new();
let permission_sync = PermissionSyncService::new(repositories.clone(), event_bus.clone());
let state = AppState {
db,
config: Arc::new(config),
@@ -77,7 +83,11 @@ impl App {
event_bus,
};
Ok(Self { state })
let services = Services {
permission_sync: Arc::new(permission_sync),
};
Ok(Self { state, services })
}
pub async fn run(self) -> Result<(), Box<dyn std::error::Error>> {
+44
View File
@@ -0,0 +1,44 @@
use crate::repositories::Repositories;
use event_bus::EventBus;
use std::sync::Arc;
// list of all events :
// server_user_created
// server_user_deleted
//
// role_user_created
// role_user_deleted
//
// server_role_permission_created
// server_role_permission_updated
// server_role_permission_deleted
//
// server_user_permission_created
// server_user_permission_updated
// server_user_permission_deleted
//
// channel_role_permission_created
// channel_role_permission_updated
// channel_role_permission_deleted
//
// channel_user_permission_created
// channel_user_permission_updated
// channel_user_permission_deleted
//
// channel_created
// channel_deleted
#[derive(Debug, Clone)]
pub struct PermissionSyncService {
repositories: Arc<Repositories>,
event_bus: Arc<EventBus>,
}
impl PermissionSyncService {
pub fn new(repositories: Arc<Repositories>, event_bus: Arc<EventBus>) -> Self {
Self {
repositories,
event_bus,
}
}
}
+7 -1
View File
@@ -1,4 +1,5 @@
use crate::config::AppConfig;
use crate::core::permission_sync::PermissionSyncService;
use crate::metrics::AppMetrics;
use crate::models::server;
use crate::repositories::Repositories;
@@ -11,7 +12,7 @@ use std::sync::{Arc, RwLock};
pub struct AppState {
pub db: DatabaseConnection,
pub config: Arc<AppConfig>,
pub repositories: Repositories,
pub repositories: Arc<Repositories>,
pub init_token: Arc<RwLock<Option<uuid::Uuid>>>,
pub default_server: Arc<server::Model>,
pub metrics: AppMetrics,
@@ -20,3 +21,8 @@ pub struct AppState {
}
impl AppState {}
#[derive(Debug, Clone)]
pub struct Services {
pub permission_sync: Arc<PermissionSyncService>,
}
+3 -2
View File
@@ -17,7 +17,8 @@ pub struct Model {
pub role_id: Uuid,
/// Bitmask des permissions accordées au rôle dans ce canal.
pub permission: i64,
pub permissions: i64,
#[sea_orm(
belongs_to,
from = "channel_id",
@@ -43,7 +44,7 @@ impl ActiveModelBehavior for ActiveModel {
id: Set(Uuid::new_v4()),
channel_id: NotSet,
role_id: NotSet,
permission: Set(ChannelPermission::empty().bits() as i64),
permissions: Set(ChannelPermission::empty().bits() as i64),
}
}
}
+3 -2
View File
@@ -17,7 +17,8 @@ pub struct Model {
pub user_id: Uuid,
/// Bitmask des permissions accordées directement à l'utilisateur dans ce canal.
pub permission: i64,
pub permissions: i64,
#[sea_orm(
belongs_to,
from = "channel_id",
@@ -43,7 +44,7 @@ impl ActiveModelBehavior for ActiveModel {
id: Set(Uuid::new_v4()),
channel_id: NotSet,
user_id: NotSet,
permission: Set(ChannelPermission::empty().bits() as i64),
permissions: Set(ChannelPermission::empty().bits() as i64),
}
}
}
+4 -5
View File
@@ -37,9 +37,9 @@ pub struct Model {
#[sea_orm(primary_key, auto_increment = false)]
pub resource_id: Uuid,
/// Cache des permissions serveur (stocké en i64 pour SQL, utilisé en u64)
pub server_permissions: i64,
pub channel_permissions: i64,
/// Cache des permissions (stocké en i64 pour SQL, utilisé en u64)
pub permissions: i64,
#[sea_orm(
belongs_to,
from = "user_id",
@@ -59,8 +59,7 @@ impl ActiveModelBehavior for ActiveModel {
server_id: NotSet,
scope_type: NotSet,
resource_id: NotSet,
server_permissions: Set(0),
channel_permissions: Set(0),
permissions: Set(0),
}
}
}
+3 -2
View File
@@ -17,7 +17,8 @@ pub struct Model {
pub role_id: Uuid,
/// Bitmask des permissions accordées directement à l'utilisateur.
pub permission: i64,
pub permissions: i64,
#[sea_orm(
belongs_to,
from = "server_id",
@@ -43,7 +44,7 @@ impl ActiveModelBehavior for ActiveModel {
id: Set(Uuid::new_v4()),
server_id: NotSet,
role_id: NotSet,
permission: Set(ServerPermission::empty().bits() as i64),
permissions: Set(ServerPermission::empty().bits() as i64),
}
}
}
+2 -2
View File
@@ -17,7 +17,7 @@ pub struct Model {
pub user_id: Uuid,
/// Bitmask des permissions accordées directement à l'utilisateur.
pub permission: i64,
pub permissions: i64,
#[sea_orm(
belongs_to,
from = "server_id",
@@ -43,7 +43,7 @@ impl ActiveModelBehavior for ActiveModel {
id: Set(Uuid::new_v4()),
server_id: NotSet,
user_id: NotSet,
permission: Set(ServerPermission::empty().bits() as i64),
permissions: Set(ServerPermission::empty().bits() as i64),
}
}
}
+20 -6
View File
@@ -1,6 +1,6 @@
use crate::models::{
channel, channel_role_permission, channel_user_permission, computed_permission, role_user,
server_role_permission, server_user,
server_role_permission, server_user, server_user_permission,
};
use crate::permissions::{ChannelPermission, ServerPermission};
use crate::repositories::{AnyResult, RepositoryContext};
@@ -46,6 +46,7 @@ impl ComputedPermissionRepository {
/// Les permissions effectives sont composées de :
///
/// - permissions serveur accordées aux rôles de l'utilisateur ;
/// - permissions serveur accordées directement à l'utilisateur ;
/// - permissions de canal accordées aux rôles de l'utilisateur ;
/// - permissions directes de l'utilisateur dans les canaux.
pub async fn full_sync_user(&self, user_id: Uuid, server_id: Uuid) -> AnyResult<()> {
@@ -76,10 +77,23 @@ impl ComputedPermissionRepository {
for permission in role_permissions {
server_permissions |=
ServerPermission::from_bits_retain(permission.permission as u64);
ServerPermission::from_bits_retain(permission.permissions as u64);
}
}
// ---------------------------------------------------------------------
// Permissions serveur directes de l'utilisateur
// ---------------------------------------------------------------------
if let Some(permission) = server_user_permission::Entity::find()
.filter(server_user_permission::Column::ServerId.eq(server_id))
.filter(server_user_permission::Column::UserId.eq(user_id))
.one(&self.context.db)
.await?
{
server_permissions |= ServerPermission::from_bits_retain(permission.permissions as u64);
}
// ---------------------------------------------------------------------
// Canaux du serveur
// ---------------------------------------------------------------------
@@ -112,7 +126,7 @@ impl ComputedPermissionRepository {
.entry(permission.channel_id)
.or_default()
.insert(ChannelPermission::from_bits_retain(
permission.permission as u64,
permission.permissions as u64,
));
}
@@ -135,7 +149,7 @@ impl ComputedPermissionRepository {
.entry(permission.channel_id)
.or_default()
.insert(ChannelPermission::from_bits_retain(
permission.permission as u64,
permission.permissions as u64,
));
}
@@ -151,7 +165,7 @@ impl ComputedPermissionRepository {
server_id: Set(server_id),
scope_type: Set(PermissionScopeType::Server),
resource_id: Set(server_id),
server_permissions: Set(server_permissions.bits() as i64),
permissions: Set(server_permissions.bits() as i64),
..Default::default()
});
@@ -166,7 +180,7 @@ impl ComputedPermissionRepository {
server_id: Set(server_id),
scope_type: Set(PermissionScopeType::Channel),
resource_id: Set(channel.id),
channel_permissions: Set(channel_permissions.bits() as i64),
permissions: Set(channel_permissions.bits() as i64),
..Default::default()
});
}