init
This commit is contained in:
@@ -27,8 +27,14 @@ impl RealtimeRegistry {
|
||||
{
|
||||
continue;
|
||||
}
|
||||
channel_users.entry(permission.resource_id).or_default().insert(permission.user_id);
|
||||
user_channels.entry(permission.user_id).or_default().insert(permission.resource_id);
|
||||
channel_users
|
||||
.entry(permission.resource_id)
|
||||
.or_default()
|
||||
.insert(permission.user_id);
|
||||
user_channels
|
||||
.entry(permission.user_id)
|
||||
.or_default()
|
||||
.insert(permission.resource_id);
|
||||
}
|
||||
|
||||
*self.channel_users.write() = channel_users;
|
||||
@@ -37,20 +43,32 @@ impl RealtimeRegistry {
|
||||
}
|
||||
|
||||
pub fn users_for_channel(&self, channel_id: Uuid) -> HashSet<Uuid> {
|
||||
self.channel_users.read().get(&channel_id).cloned().unwrap_or_default()
|
||||
self.channel_users
|
||||
.read()
|
||||
.get(&channel_id)
|
||||
.cloned()
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
pub fn set_user_channels(&self, user_id: Uuid, channels: impl IntoIterator<Item = Uuid>) {
|
||||
let channels: HashSet<_> = channels.into_iter().collect();
|
||||
let old = self.user_channels.write().insert(user_id, channels.clone()).unwrap_or_default();
|
||||
let old = self
|
||||
.user_channels
|
||||
.write()
|
||||
.insert(user_id, channels.clone())
|
||||
.unwrap_or_default();
|
||||
let mut by_channel = self.channel_users.write();
|
||||
for channel_id in old.difference(&channels) {
|
||||
if let Some(users) = by_channel.get_mut(channel_id) {
|
||||
users.remove(&user_id);
|
||||
if users.is_empty() { by_channel.remove(channel_id); }
|
||||
if users.is_empty() {
|
||||
by_channel.remove(channel_id);
|
||||
}
|
||||
}
|
||||
}
|
||||
for channel_id in channels { by_channel.entry(channel_id).or_default().insert(user_id); }
|
||||
for channel_id in channels {
|
||||
by_channel.entry(channel_id).or_default().insert(user_id);
|
||||
}
|
||||
}
|
||||
|
||||
pub fn remove_user(&self, user_id: Uuid) {
|
||||
@@ -59,7 +77,9 @@ impl RealtimeRegistry {
|
||||
for channel_id in channels {
|
||||
if let Some(users) = by_channel.get_mut(&channel_id) {
|
||||
users.remove(&user_id);
|
||||
if users.is_empty() { by_channel.remove(&channel_id); }
|
||||
if users.is_empty() {
|
||||
by_channel.remove(&channel_id);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -71,33 +91,69 @@ impl RealtimeRegistry {
|
||||
for user_id in users {
|
||||
if let Some(channels) = by_user.get_mut(&user_id) {
|
||||
channels.remove(&channel_id);
|
||||
if channels.is_empty() { by_user.remove(&user_id); }
|
||||
if channels.is_empty() {
|
||||
by_user.remove(&user_id);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub fn start_listening(self: &Arc<Self>, repositories: Arc<Repositories>, event_bus: Arc<EventBus>) {
|
||||
pub fn start_listening(
|
||||
self: &Arc<Self>,
|
||||
repositories: Arc<Repositories>,
|
||||
event_bus: Arc<EventBus>,
|
||||
) {
|
||||
let registry = Arc::clone(self);
|
||||
event_bus.on_async_with("channel_user_permission_updated", repositories.clone(), move |repositories, (_channel_id, user_id, _permissions): (Uuid, Uuid, u64)| {
|
||||
let registry = Arc::clone(®istry);
|
||||
async move {
|
||||
match repositories.computed_permission.get_all().await {
|
||||
Ok(all) => registry.set_user_channels(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)),
|
||||
Err(error) => tracing::error!(%user_id, ?error, "Unable to refresh realtime registry"),
|
||||
event_bus.on_async_with(
|
||||
"channel_user_permission_updated",
|
||||
repositories.clone(),
|
||||
move |repositories, (_channel_id, user_id, _permissions): (Uuid, Uuid, u64)| {
|
||||
let registry = Arc::clone(®istry);
|
||||
async move {
|
||||
match repositories.computed_permission.get_all().await {
|
||||
Ok(all) => registry.set_user_channels(
|
||||
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),
|
||||
),
|
||||
Err(error) => {
|
||||
tracing::error!(%user_id, ?error, "Unable to refresh realtime registry")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
},
|
||||
);
|
||||
|
||||
let registry = Arc::clone(self);
|
||||
let repositories = repositories.clone();
|
||||
event_bus.on_async_with("server_user_permission_updated", repositories, move |repositories, (_server_id, user_id): (Uuid, Uuid)| {
|
||||
let registry = Arc::clone(®istry);
|
||||
async move {
|
||||
if let Ok(all) = repositories.computed_permission.get_all().await {
|
||||
registry.set_user_channels(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));
|
||||
event_bus.on_async_with(
|
||||
"server_user_permission_updated",
|
||||
repositories,
|
||||
move |repositories, (_server_id, user_id): (Uuid, Uuid)| {
|
||||
let registry = Arc::clone(®istry);
|
||||
async move {
|
||||
if let Ok(all) = repositories.computed_permission.get_all().await {
|
||||
registry.set_user_channels(
|
||||
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),
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
},
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user