rename AdminMessage to SignalAlert message which is less confusing
This commit is contained in:
@@ -1,6 +1,6 @@
|
|||||||
use super::log_buffer::LogBuffer;
|
use super::log_buffer::LogBuffer;
|
||||||
use super::route::{Destination, Limit, Route};
|
use super::route::{Destination, Limit, Route};
|
||||||
use super::{AdminMessage, LimitResult, Limiter, LimiterSet};
|
use super::{SignalAlertMessage, LimitResult, Limiter, LimiterSet};
|
||||||
use crate::{
|
use crate::{
|
||||||
concurrent_map::ConcurrentMap,
|
concurrent_map::ConcurrentMap,
|
||||||
log_format::LogFormatConfig,
|
log_format::LogFormatConfig,
|
||||||
@@ -74,7 +74,7 @@ pub struct LogHandlerConfig {
|
|||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub struct LogHandler {
|
pub struct LogHandler {
|
||||||
config: LogHandlerConfig,
|
config: LogHandlerConfig,
|
||||||
admin_mq_tx: UnboundedSender<AdminMessage>,
|
signal_alert_mq_tx: UnboundedSender<SignalAlertMessage>,
|
||||||
/// Log buffers keyed by origin (app + host). Lazily created.
|
/// Log buffers keyed by origin (app + host). Lazily created.
|
||||||
log_buffers: ConcurrentMap<Origin, LogBuffer>,
|
log_buffers: ConcurrentMap<Origin, LogBuffer>,
|
||||||
/// Routes with their associated limiter sets.
|
/// Routes with their associated limiter sets.
|
||||||
@@ -85,7 +85,7 @@ pub struct LogHandler {
|
|||||||
|
|
||||||
impl LogHandler {
|
impl LogHandler {
|
||||||
/// Initialize a new log handler
|
/// Initialize a new log handler
|
||||||
pub fn new(config: LogHandlerConfig, admin_mq_tx: UnboundedSender<AdminMessage>) -> Self {
|
pub fn new(config: LogHandlerConfig, signal_alert_mq_tx: UnboundedSender<SignalAlertMessage>) -> Self {
|
||||||
let routes = config
|
let routes = config
|
||||||
.routes
|
.routes
|
||||||
.iter()
|
.iter()
|
||||||
@@ -100,7 +100,7 @@ impl LogHandler {
|
|||||||
|
|
||||||
Self {
|
Self {
|
||||||
config,
|
config,
|
||||||
admin_mq_tx,
|
signal_alert_mq_tx,
|
||||||
log_buffers: ConcurrentMap::new(),
|
log_buffers: ConcurrentMap::new(),
|
||||||
routes,
|
routes,
|
||||||
overall_limits,
|
overall_limits,
|
||||||
@@ -192,7 +192,7 @@ impl LogHandler {
|
|||||||
// Send alert if we have formatted text
|
// Send alert if we have formatted text
|
||||||
if let Some(text) = formatted_text {
|
if let Some(text) = formatted_text {
|
||||||
let destination_override = rate_limit_result.ok().flatten();
|
let destination_override = rate_limit_result.ok().flatten();
|
||||||
if let Err(_err) = self.admin_mq_tx.send(AdminMessage {
|
if let Err(_err) = self.signal_alert_mq_tx.send(SignalAlertMessage {
|
||||||
origin: Some(origin),
|
origin: Some(origin),
|
||||||
text,
|
text,
|
||||||
attachment_paths: Default::default(),
|
attachment_paths: Default::default(),
|
||||||
|
|||||||
@@ -162,7 +162,7 @@ fn parse_gateway_command(s: &str) -> Result<GatewayCommand, String> {
|
|||||||
/// A message queued to be sent to all admins.
|
/// A message queued to be sent to all admins.
|
||||||
/// This is generally an alert message, which may have attached images.
|
/// This is generally an alert message, which may have attached images.
|
||||||
#[derive(Clone, Debug, Default)]
|
#[derive(Clone, Debug, Default)]
|
||||||
struct AdminMessage {
|
struct SignalAlertMessage {
|
||||||
/// The origin of the message (app + host), if from syslog
|
/// The origin of the message (app + host), if from syslog
|
||||||
origin: Option<Origin>,
|
origin: Option<Origin>,
|
||||||
text: String,
|
text: String,
|
||||||
@@ -188,8 +188,8 @@ struct AdminMessage {
|
|||||||
/// and then the buffer is purged. This serves as a minimal log-aggregation and alerting system.
|
/// and then the buffer is purged. This serves as a minimal log-aggregation and alerting system.
|
||||||
pub struct Gateway {
|
pub struct Gateway {
|
||||||
config: GatewayConfig,
|
config: GatewayConfig,
|
||||||
admin_mq_tx: UnboundedSender<AdminMessage>,
|
signal_alert_mq_tx: UnboundedSender<SignalAlertMessage>,
|
||||||
admin_mq_rx: Mutex<UnboundedReceiver<AdminMessage>>,
|
signal_alert_mq_rx: Mutex<UnboundedReceiver<SignalAlertMessage>>,
|
||||||
token: CancellationToken,
|
token: CancellationToken,
|
||||||
prometheus: Option<Prometheus>,
|
prometheus: Option<Prometheus>,
|
||||||
/// Log handler for processing log messages from all origins.
|
/// Log handler for processing log messages from all origins.
|
||||||
@@ -205,7 +205,7 @@ impl Gateway {
|
|||||||
token: CancellationToken,
|
token: CancellationToken,
|
||||||
message_handler: Option<Box<dyn MessageHandler>>,
|
message_handler: Option<Box<dyn MessageHandler>>,
|
||||||
) -> Self {
|
) -> Self {
|
||||||
let (admin_mq_tx, admin_mq_rx) = unbounded_channel();
|
let (signal_alert_mq_tx, signal_alert_mq_rx) = unbounded_channel();
|
||||||
|
|
||||||
let prometheus = config
|
let prometheus = config
|
||||||
.prometheus
|
.prometheus
|
||||||
@@ -214,12 +214,12 @@ impl Gateway {
|
|||||||
.transpose()
|
.transpose()
|
||||||
.expect("Invalid prometheus config");
|
.expect("Invalid prometheus config");
|
||||||
|
|
||||||
let log_handler = LogHandler::new(config.log_handler.clone(), admin_mq_tx.clone());
|
let log_handler = LogHandler::new(config.log_handler.clone(), signal_alert_mq_tx.clone());
|
||||||
|
|
||||||
Self {
|
Self {
|
||||||
config,
|
config,
|
||||||
admin_mq_tx,
|
signal_alert_mq_tx,
|
||||||
admin_mq_rx: Mutex::new(admin_mq_rx),
|
signal_alert_mq_rx: Mutex::new(signal_alert_mq_rx),
|
||||||
token,
|
token,
|
||||||
prometheus,
|
prometheus,
|
||||||
log_handler,
|
log_handler,
|
||||||
@@ -364,8 +364,8 @@ impl Gateway {
|
|||||||
async fn do_run(&self, signal_cli: &impl RpcClient) -> Result<(), RpcClientError> {
|
async fn do_run(&self, signal_cli: &impl RpcClient) -> Result<(), RpcClientError> {
|
||||||
self.update_trust(signal_cli).await;
|
self.update_trust(signal_cli).await;
|
||||||
|
|
||||||
let mut admin_mq_rx = self
|
let mut signal_alert_mq_rx = self
|
||||||
.admin_mq_rx
|
.signal_alert_mq_rx
|
||||||
.try_lock()
|
.try_lock()
|
||||||
.expect("Mutex should not be contended");
|
.expect("Mutex should not be contended");
|
||||||
let mut signal_rx = signal_cli
|
let mut signal_rx = signal_cli
|
||||||
@@ -378,7 +378,7 @@ impl Gateway {
|
|||||||
info!("Stop requested");
|
info!("Stop requested");
|
||||||
return Ok(());
|
return Ok(());
|
||||||
},
|
},
|
||||||
outbound_admin_msg = admin_mq_rx.recv() => {
|
outbound_admin_msg = signal_alert_mq_rx.recv() => {
|
||||||
if let Some(msg) = outbound_admin_msg {
|
if let Some(msg) = outbound_admin_msg {
|
||||||
// Log summary, or first 500 bytes of text if no summary provided
|
// Log summary, or first 500 bytes of text if no summary provided
|
||||||
let summary = msg.summary.as_deref().unwrap_or_else(|| {
|
let summary = msg.summary.as_deref().unwrap_or_else(|| {
|
||||||
@@ -413,7 +413,7 @@ impl Gateway {
|
|||||||
attachments,
|
attachments,
|
||||||
}.send(signal_cli).await?;
|
}.send(signal_cli).await?;
|
||||||
} else {
|
} else {
|
||||||
warn!("admin_mq_rx is closed, halting service");
|
warn!("signal_alert_mq_rx is closed, halting service");
|
||||||
self.token.cancel();
|
self.token.cancel();
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
@@ -762,8 +762,8 @@ impl Gateway {
|
|||||||
.collect::<Vec<_>>()
|
.collect::<Vec<_>>()
|
||||||
.join(" ");
|
.join(" ");
|
||||||
|
|
||||||
self.admin_mq_tx
|
self.signal_alert_mq_tx
|
||||||
.send(AdminMessage {
|
.send(SignalAlertMessage {
|
||||||
origin: None, // Prometheus alerts don't have a syslog origin
|
origin: None, // Prometheus alerts don't have a syslog origin
|
||||||
text,
|
text,
|
||||||
attachment_paths,
|
attachment_paths,
|
||||||
|
|||||||
Reference in New Issue
Block a user