rename verified signal message to admin message, reformat

This commit is contained in:
Chris Beck
2025-12-05 22:07:22 -07:00
parent 010a5a9336
commit 7a163f3962
8 changed files with 47 additions and 22 deletions
+2 -2
View File
@@ -6,7 +6,7 @@
use async_trait::async_trait; use async_trait::async_trait;
use conf::Conf; use conf::Conf;
use signal_gateway::{ use signal_gateway::{
AdminMessageResponse, Context, MessageHandler, MessageHandlerResult, VerifiedSignalMessage, AdminMessage, AdminMessageResponse, Context, MessageHandler, MessageHandlerResult,
}; };
use std::time::Duration; use std::time::Duration;
@@ -49,7 +49,7 @@ struct AdminHttpHandler {
impl MessageHandler for AdminHttpHandler { impl MessageHandler for AdminHttpHandler {
async fn handle_verified_signal_message( async fn handle_verified_signal_message(
&self, &self,
msg: VerifiedSignalMessage, msg: AdminMessage,
_context: &dyn Context, _context: &dyn Context,
) -> MessageHandlerResult { ) -> MessageHandlerResult {
let response = self let response = self
+2 -2
View File
@@ -6,7 +6,7 @@
use async_trait::async_trait; use async_trait::async_trait;
use conf::Conf; use conf::Conf;
use signal_gateway::{ use signal_gateway::{
AdminMessageResponse, Context, MessageHandler, MessageHandlerResult, VerifiedSignalMessage, AdminMessage, AdminMessageResponse, Context, MessageHandler, MessageHandlerResult,
}; };
use std::time::Duration; use std::time::Duration;
use tokio::{ use tokio::{
@@ -45,7 +45,7 @@ struct AdminNetcatHandler {
impl MessageHandler for AdminNetcatHandler { impl MessageHandler for AdminNetcatHandler {
async fn handle_verified_signal_message( async fn handle_verified_signal_message(
&self, &self,
msg: VerifiedSignalMessage, msg: AdminMessage,
_context: &dyn Context, _context: &dyn Context,
) -> MessageHandlerResult { ) -> MessageHandlerResult {
// Connect to server // Connect to server
+4 -1
View File
@@ -51,7 +51,10 @@ impl LogBuffer {
// caller-determined, but we need to pass our concrete iterator type. HRTB with the // caller-determined, but we need to pass our concrete iterator type. HRTB with the
// concrete type (`F: for<'a> FnOnce(Rev<vec_deque::Iter<'a, T>>)`) works but leaks // concrete type (`F: for<'a> FnOnce(Rev<vec_deque::Iter<'a, T>>)`) works but leaks
// implementation details. // implementation details.
pub fn with_iter<R>(&self, f: impl FnOnce(&mut dyn ExactSizeIterator<Item = &LogMessage>) -> R) -> R { pub fn with_iter<R>(
&self,
f: impl FnOnce(&mut dyn ExactSizeIterator<Item = &LogMessage>) -> R,
) -> R {
let buf = self.buf.lock().unwrap(); let buf = self.buf.lock().unwrap();
let mut iter = buf.iter().rev(); let mut iter = buf.iter().rev();
f(&mut iter) f(&mut iter)
+11 -4
View File
@@ -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::{SignalAlertMessage, LimitResult, Limiter, LimiterSet}; use super::{LimitResult, Limiter, LimiterSet, SignalAlertMessage};
use crate::{ use crate::{
concurrent_map::ConcurrentMap, concurrent_map::ConcurrentMap,
log_format::LogFormatConfig, log_format::LogFormatConfig,
@@ -85,7 +85,10 @@ pub struct LogHandler {
impl LogHandler { impl LogHandler {
/// Initialize a new log handler /// Initialize a new log handler
pub fn new(config: LogHandlerConfig, signal_alert_mq_tx: UnboundedSender<SignalAlertMessage>) -> Self { pub fn new(
config: LogHandlerConfig,
signal_alert_mq_tx: UnboundedSender<SignalAlertMessage>,
) -> Self {
let routes = config let routes = config
.routes .routes
.iter() .iter()
@@ -133,7 +136,9 @@ impl LogHandler {
// Guess at how much to reserve // Guess at how much to reserve
text.reserve(iter.len() * 128); text.reserve(iter.len() * 128);
for log_msg in iter { for log_msg in iter {
self.config.log_format.write_log_msg(&mut text, log_msg, now); self.config
.log_format
.write_log_msg(&mut text, log_msg, now);
} }
}); });
text.push('\n'); text.push('\n');
@@ -180,7 +185,9 @@ impl LogHandler {
let now = Utc::now(); let now = Utc::now();
buffer.push_back_and_drain(log_msg, |log_msg| { buffer.push_back_and_drain(log_msg, |log_msg| {
self.config.log_format.write_log_msg(&mut text, log_msg, now); self.config
.log_format
.write_log_msg(&mut text, log_msg, now);
}); });
Some(text) Some(text)
+6 -6
View File
@@ -4,14 +4,14 @@
use crate::signal_jsonrpc::connect_ipc; use crate::signal_jsonrpc::connect_ipc;
use crate::{ use crate::{
alertmanager::AlertPost, alertmanager::AlertPost,
log_message::{LogMessage, Origin},
message_handler::{
AdminMessage, AdminMessageResponse, Context, MessageHandler, MessageHandlerResult,
},
prometheus::{Prometheus, PrometheusConfig},
signal_jsonrpc::{ signal_jsonrpc::{
Envelope, Identity, MessageTarget, RpcClient, RpcClientError, SignalMessage, connect_tcp, Envelope, Identity, MessageTarget, RpcClient, RpcClientError, SignalMessage, connect_tcp,
}, },
log_message::{LogMessage, Origin},
message_handler::{
AdminMessageResponse, Context, MessageHandler, MessageHandlerResult, VerifiedSignalMessage,
},
prometheus::{Prometheus, PrometheusConfig},
}; };
use chrono::Utc; use chrono::Utc;
use conf::{Conf, Subcommands}; use conf::{Conf, Subcommands};
@@ -494,7 +494,7 @@ impl Gateway {
self.handle_gateway_command(cmd).await self.handle_gateway_command(cmd).await
} else if let Some(handler) = &self.message_handler { } else if let Some(handler) = &self.message_handler {
let msg = VerifiedSignalMessage { let msg = AdminMessage {
message: data.message.clone(), message: data.message.clone(),
timestamp: data.timestamp, timestamp: data.timestamp,
sender_uuid: msg.source_uuid.clone(), sender_uuid: msg.source_uuid.clone(),
+2 -2
View File
@@ -13,12 +13,12 @@ pub(crate) mod circular_buffer;
pub(crate) mod concurrent_map; pub(crate) mod concurrent_map;
pub(crate) mod log_format; pub(crate) mod log_format;
pub(crate) mod log_message; pub(crate) mod log_message;
pub(crate) mod signal_jsonrpc;
pub(crate) mod prometheus; pub(crate) mod prometheus;
pub(crate) mod signal_jsonrpc;
pub(crate) mod transports; pub(crate) mod transports;
pub use gateway::{Gateway, GatewayConfig}; pub use gateway::{Gateway, GatewayConfig};
pub use log_message::{Level, LogFilter, LogMessage, LogMessageBuilder}; pub use log_message::{Level, LogFilter, LogMessage, LogMessageBuilder};
pub use message_handler::{ pub use message_handler::{
AdminMessageResponse, Context, MessageHandler, MessageHandlerResult, VerifiedSignalMessage, AdminMessage, AdminMessageResponse, Context, MessageHandler, MessageHandlerResult,
}; };
+18 -3
View File
@@ -25,11 +25,20 @@ impl LogFormatConfig {
/// ///
/// The `now` parameter is the current time, used to calculate relative /// The `now` parameter is the current time, used to calculate relative
/// timestamps (e.g., "T-10s"). /// timestamps (e.g., "T-10s").
pub fn write_log_msg(&self, mut writer: impl std::fmt::Write, log_msg: &LogMessage, now: DateTime<Utc>) { pub fn write_log_msg(
&self,
mut writer: impl std::fmt::Write,
log_msg: &LogMessage,
now: DateTime<Utc>,
) {
// Format: "ERROR T-10s [foo bar.rs:42]: message" // Format: "ERROR T-10s [foo bar.rs:42]: message"
// Pad severity to 5 chars (left-aligned), time to 8 chars (right-aligned) // Pad severity to 5 chars (left-aligned), time to 8 chars (right-aligned)
if self.write_log_msg_inner(&mut writer, log_msg, now).is_err() { if self.write_log_msg_inner(&mut writer, log_msg, now).is_err() {
error!("Couldn't write log message: {}: {}", log_msg.level.to_str(), log_msg.msg); error!(
"Couldn't write log message: {}: {}",
log_msg.level.to_str(),
log_msg.msg
);
} }
} }
@@ -132,7 +141,13 @@ fn write_t_minus(writer: &mut impl Write, delta: TimeDelta, align: usize) -> std
} }
// For durations < 1 minute, show hundredths; < 10 minutes, show tenths // For durations < 1 minute, show hundredths; < 10 minutes, show tenths
let decimal_places = if total_secs < 60 { 2 } else if total_secs < 600 { 1 } else { 0 }; let decimal_places = if total_secs < 60 {
2
} else if total_secs < 600 {
1
} else {
0
};
let frac = match decimal_places { let frac = match decimal_places {
2 => nanos / 10_000_000, // hundredths 2 => nanos / 10_000_000, // hundredths
1 => nanos / 100_000_000, // tenths 1 => nanos / 100_000_000, // tenths
+2 -2
View File
@@ -9,7 +9,7 @@ use std::{error::Error, path::PathBuf};
/// that has been verified as coming from a trusted admin. /// that has been verified as coming from a trusted admin.
#[non_exhaustive] #[non_exhaustive]
#[derive(Clone, Debug)] #[derive(Clone, Debug)]
pub struct VerifiedSignalMessage { pub struct AdminMessage {
/// The text content of the message. /// The text content of the message.
pub message: String, pub message: String,
/// The timestamp of the message (milliseconds since Unix epoch). /// The timestamp of the message (milliseconds since Unix epoch).
@@ -111,7 +111,7 @@ pub trait MessageHandler: Send + Sync {
/// handled as gateway commands). /// handled as gateway commands).
async fn handle_verified_signal_message( async fn handle_verified_signal_message(
&self, &self,
msg: VerifiedSignalMessage, msg: AdminMessage,
context: &dyn Context, context: &dyn Context,
) -> MessageHandlerResult; ) -> MessageHandlerResult;
} }