Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3b85fc2e17 |
@@ -34,7 +34,6 @@ default_exchange = "revolt.default"
|
||||
|
||||
[rabbit.queues]
|
||||
acks = "internal.ack"
|
||||
events = "internal.event"
|
||||
|
||||
[api]
|
||||
|
||||
|
||||
@@ -125,7 +125,6 @@ pub struct Database {
|
||||
#[derive(Deserialize, Debug, Clone)]
|
||||
pub struct RabbitQueues {
|
||||
pub acks: String,
|
||||
pub events: String,
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Debug, Clone)]
|
||||
|
||||
@@ -1,29 +1,20 @@
|
||||
use std::{
|
||||
collections::HashSet,
|
||||
sync::{Arc, OnceLock},
|
||||
};
|
||||
use std::collections::HashSet;
|
||||
use std::sync::Arc;
|
||||
|
||||
use crate::events::{client::EventV1, rabbit::*};
|
||||
use crate::events::rabbit::*;
|
||||
use crate::User;
|
||||
use lapin::{
|
||||
options::BasicPublishOptions,
|
||||
protocol::basic::AMQPProperties,
|
||||
types::{AMQPValue, FieldTable},
|
||||
BasicProperties, Channel, Connection, ConnectionProperties, Error as AMQPError,
|
||||
Channel, Connection, ConnectionProperties, Error as AMQPError,
|
||||
};
|
||||
use revolt_config::config;
|
||||
use revolt_models::v0::PushNotification;
|
||||
use revolt_presence::filter_online;
|
||||
use revolt_result::Result;
|
||||
|
||||
use serde_json::to_string;
|
||||
|
||||
static AMQP_INSTANCE: OnceLock<AMQP> = OnceLock::new();
|
||||
|
||||
pub fn get_amqp() -> &'static AMQP {
|
||||
AMQP_INSTANCE.get().expect("No AMQP instance set.")
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct AMQP {
|
||||
friend_request_accepted: Arc<Channel>,
|
||||
@@ -34,14 +25,13 @@ pub struct AMQP {
|
||||
ack_notification_message: Arc<Channel>,
|
||||
dm_call_updated: Arc<Channel>,
|
||||
process_ack: Arc<Channel>,
|
||||
publish_event: Arc<Channel>,
|
||||
#[allow(unused)]
|
||||
connection: Arc<Connection>,
|
||||
}
|
||||
|
||||
impl AMQP {
|
||||
pub async fn new(connection: Arc<Connection>) -> Self {
|
||||
let this = Self {
|
||||
Self {
|
||||
friend_request_accepted: Self::create_channel(&connection).await,
|
||||
friend_request_received: Self::create_channel(&connection).await,
|
||||
generic_message: Self::create_channel(&connection).await,
|
||||
@@ -50,13 +40,8 @@ impl AMQP {
|
||||
ack_notification_message: Self::create_channel(&connection).await,
|
||||
dm_call_updated: Self::create_channel(&connection).await,
|
||||
process_ack: Self::create_channel(&connection).await,
|
||||
publish_event: Self::create_channel(&connection).await,
|
||||
connection,
|
||||
};
|
||||
|
||||
let _ = AMQP_INSTANCE.set(this.clone());
|
||||
|
||||
this
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn new_auto() -> Self {
|
||||
@@ -394,23 +379,4 @@ impl AMQP {
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn publish_event(&self, channel: String, event: &EventV1) -> Result<(), AMQPError> {
|
||||
let mut headers = FieldTable::default();
|
||||
headers.insert("c".into(), AMQPValue::LongString(channel.into()));
|
||||
|
||||
let config = config().await;
|
||||
|
||||
self.publish_event
|
||||
.basic_publish(
|
||||
config.rabbit.default_exchange.clone().into(),
|
||||
config.rabbit.queues.events.into(),
|
||||
BasicPublishOptions::default(),
|
||||
&serde_json::to_vec(event).unwrap(),
|
||||
BasicProperties::default().with_headers(headers),
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,4 +1,2 @@
|
||||
#[allow(clippy::module_inception)]
|
||||
pub mod amqp;
|
||||
|
||||
pub use amqp::{AMQP, get_amqp};
|
||||
@@ -11,7 +11,7 @@ use revolt_models::v0::{
|
||||
UserVoiceState, Webhook,
|
||||
};
|
||||
|
||||
use crate::{Database, amqp::get_amqp};
|
||||
use crate::Database;
|
||||
|
||||
/// Ping Packet
|
||||
#[derive(Serialize, Deserialize, Debug, Clone)]
|
||||
@@ -372,16 +372,8 @@ impl EventV1 {
|
||||
#[cfg(debug_assertions)]
|
||||
info!("Publishing event to {channel}: {self:?}");
|
||||
|
||||
// #[cfg(debug_assertions)]
|
||||
// redis_kiss::publish(channel, self).await.unwrap();
|
||||
|
||||
if let Err(e) = get_amqp().publish_event(channel, &self).await {
|
||||
if cfg!(debug_assertions) {
|
||||
panic!("{e:?}");
|
||||
} else {
|
||||
log::error!("{e:?}");
|
||||
};
|
||||
};
|
||||
#[cfg(debug_assertions)]
|
||||
redis_kiss::publish(channel, self).await.unwrap();
|
||||
}
|
||||
|
||||
/// Publish user event
|
||||
|
||||
@@ -30,7 +30,7 @@ pub async fn create_session(user_id: &str, flags: u8) -> (bool, u32) {
|
||||
|
||||
if let Ok(mut conn) = get_connection().await {
|
||||
// Check whether this is the first session
|
||||
let was_empty = __get_set_size(&mut conn, user_id).await == 0;
|
||||
let was_empty = __get_set_size(&mut conn, &format!("sessions:{user_id}")).await == 0;
|
||||
|
||||
// A session ID is comprised of random data and any flags ORed to the end
|
||||
let session_id = {
|
||||
@@ -39,7 +39,7 @@ pub async fn create_session(user_id: &str, flags: u8) -> (bool, u32) {
|
||||
};
|
||||
|
||||
// Add session to user's sessions and to the region
|
||||
__add_to_set_u32(&mut conn, user_id, session_id).await;
|
||||
__add_to_set_u32(&mut conn, &format!("sessions:{user_id}"), session_id).await;
|
||||
__add_to_set_string(&mut conn, ONLINE_SET, user_id).await;
|
||||
__add_to_set_string(&mut conn, ®ION_KEY, &format!("{user_id}:{session_id}")).await;
|
||||
info!("Created session for {user_id}, assigned them a session ID of {session_id}.");
|
||||
@@ -62,7 +62,7 @@ async fn delete_session_internal(user_id: &str, session_id: u32, skip_region: bo
|
||||
|
||||
if let Ok(mut conn) = get_connection().await {
|
||||
// Remove the session
|
||||
__remove_from_set_u32(&mut conn, user_id, session_id).await;
|
||||
__remove_from_set_u32(&mut conn, &format!("sessions:{user_id}"), session_id).await;
|
||||
|
||||
// Remove from the region
|
||||
if !skip_region {
|
||||
|
||||
Reference in New Issue
Block a user