remove admin-netcat thing
This commit is contained in:
@@ -1,84 +0,0 @@
|
|||||||
//! Admin netcat TCP client for forwarding messages to a TCP server
|
|
||||||
//!
|
|
||||||
//! This module handles admin messages not handled by the gateway by opening a TCP connection,
|
|
||||||
//! writing the message terminated with CRLF, and reading the response until CRLF.
|
|
||||||
|
|
||||||
use async_trait::async_trait;
|
|
||||||
use conf::Conf;
|
|
||||||
use signal_gateway::{
|
|
||||||
AdminMessage, AdminMessageResponse, Context, MessageHandler, MessageHandlerResult,
|
|
||||||
};
|
|
||||||
use std::time::Duration;
|
|
||||||
use tokio::{
|
|
||||||
io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
|
|
||||||
net::TcpStream,
|
|
||||||
time::timeout,
|
|
||||||
};
|
|
||||||
|
|
||||||
/// Configuration for the admin netcat TCP client
|
|
||||||
#[derive(Clone, Conf, Debug)]
|
|
||||||
#[conf(serde)]
|
|
||||||
pub struct AdminNetcatConfig {
|
|
||||||
/// TCP address to forward admin commands to
|
|
||||||
#[conf(long, env)]
|
|
||||||
pub tcp_addr: String,
|
|
||||||
/// Timeout for connecting, writing, and reading
|
|
||||||
#[conf(long, env, default_value = "5s", value_parser = conf_extra::parse_duration)]
|
|
||||||
pub timeout: Duration,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl AdminNetcatConfig {
|
|
||||||
/// Create a message handler from this config.
|
|
||||||
///
|
|
||||||
/// The returned handler opens a TCP connection to the configured address,
|
|
||||||
/// writes the message terminated with CRLF, and reads the response until CRLF.
|
|
||||||
pub fn into_handler(self) -> Box<dyn MessageHandler> {
|
|
||||||
Box::new(AdminNetcatHandler { config: self })
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Message handler that forwards messages to a TCP server.
|
|
||||||
struct AdminNetcatHandler {
|
|
||||||
config: AdminNetcatConfig,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[async_trait]
|
|
||||||
impl MessageHandler for AdminNetcatHandler {
|
|
||||||
async fn handle_verified_signal_message(
|
|
||||||
&self,
|
|
||||||
msg: AdminMessage,
|
|
||||||
_context: &dyn Context,
|
|
||||||
) -> MessageHandlerResult {
|
|
||||||
// Connect to server
|
|
||||||
let mut stream = timeout(
|
|
||||||
self.config.timeout,
|
|
||||||
TcpStream::connect(&self.config.tcp_addr),
|
|
||||||
)
|
|
||||||
.await
|
|
||||||
.map_err(|_| (504u16, "connecting: timeout".into()))?
|
|
||||||
.map_err(|err| (502u16, format!("connecting: {err}").into()))?;
|
|
||||||
|
|
||||||
// Write message with CRLF terminator
|
|
||||||
let message = format!("{}\r\n", msg.message);
|
|
||||||
timeout(self.config.timeout, stream.write_all(message.as_bytes()))
|
|
||||||
.await
|
|
||||||
.map_err(|_| (504u16, "writing: timeout".into()))?
|
|
||||||
.map_err(|err| (502u16, format!("writing: {err}").into()))?;
|
|
||||||
|
|
||||||
// Read response until CR
|
|
||||||
let mut reader = BufReader::new(stream);
|
|
||||||
let mut buf = Vec::new();
|
|
||||||
timeout(self.config.timeout, reader.read_until(b'\r', &mut buf))
|
|
||||||
.await
|
|
||||||
.map_err(|_| (504u16, "reading: timeout".into()))?
|
|
||||||
.map_err(|err| (502u16, format!("reading: {err}").into()))?;
|
|
||||||
|
|
||||||
// Convert to string and trim the trailing CR
|
|
||||||
let text = std::str::from_utf8(&buf)
|
|
||||||
.map_err(|err| (502u16, format!("utf8: {err}").into()))?
|
|
||||||
.trim_end_matches(['\r', '\n'])
|
|
||||||
.to_owned();
|
|
||||||
|
|
||||||
Ok(AdminMessageResponse::new(text))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -15,8 +15,13 @@ use tracing_subscriber::EnvFilter;
|
|||||||
mod admin_http;
|
mod admin_http;
|
||||||
use admin_http::AdminHttpConfig;
|
use admin_http::AdminHttpConfig;
|
||||||
|
|
||||||
mod admin_netcat;
|
/// Handler for admin messages that don't match built-in commands.
|
||||||
use admin_netcat::AdminNetcatConfig;
|
#[derive(Subcommands, Debug)]
|
||||||
|
#[conf(serde)]
|
||||||
|
pub enum AdminHandlerCommand {
|
||||||
|
/// Forward unhandled admin messages to an HTTP endpoint.
|
||||||
|
AdminHttp(AdminHttpConfig),
|
||||||
|
}
|
||||||
|
|
||||||
mod syslog;
|
mod syslog;
|
||||||
use syslog::SyslogConfig;
|
use syslog::SyslogConfig;
|
||||||
@@ -24,28 +29,6 @@ use syslog::SyslogConfig;
|
|||||||
pub mod json;
|
pub mod json;
|
||||||
use json::JsonConfig;
|
use json::JsonConfig;
|
||||||
|
|
||||||
/// Admin message handler configuration - select how non-command messages are handled
|
|
||||||
#[derive(Clone, Debug, Subcommands)]
|
|
||||||
#[conf(serde)]
|
|
||||||
enum AdminHandlerCommand {
|
|
||||||
/// Forward (unhandled) admin messages to a TCP endpoint (netcat-style)
|
|
||||||
/// Useful if an http server would be heavy in the target process
|
|
||||||
#[conf(name = "admin-netcat")]
|
|
||||||
Netcat(AdminNetcatConfig),
|
|
||||||
/// Forward (unhandled) admin messages to an HTTP endpoint via POST
|
|
||||||
#[conf(name = "admin-http")]
|
|
||||||
Http(AdminHttpConfig),
|
|
||||||
}
|
|
||||||
|
|
||||||
impl AdminHandlerCommand {
|
|
||||||
fn into_handler(self) -> Box<dyn signal_gateway::MessageHandler> {
|
|
||||||
match self {
|
|
||||||
AdminHandlerCommand::Netcat(config) => config.into_handler(),
|
|
||||||
AdminHandlerCommand::Http(config) => config.into_handler(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Top-level configuration for signal-gateway.
|
/// Top-level configuration for signal-gateway.
|
||||||
#[derive(Conf, Debug)]
|
#[derive(Conf, Debug)]
|
||||||
#[conf(serde, test)]
|
#[conf(serde, test)]
|
||||||
@@ -60,7 +43,7 @@ pub struct Config {
|
|||||||
syslog: Option<SyslogConfig>,
|
syslog: Option<SyslogConfig>,
|
||||||
#[conf(flatten, prefix)]
|
#[conf(flatten, prefix)]
|
||||||
json: Option<JsonConfig>,
|
json: Option<JsonConfig>,
|
||||||
/// Optional admin message handler (netcat or http)
|
/// Optional handler for admin messages that don't match built-in commands.
|
||||||
#[conf(subcommands)]
|
#[conf(subcommands)]
|
||||||
admin_handler: Option<AdminHandlerCommand>,
|
admin_handler: Option<AdminHandlerCommand>,
|
||||||
#[conf(flatten, serde(flatten))]
|
#[conf(flatten, serde(flatten))]
|
||||||
@@ -109,7 +92,9 @@ async fn main() {
|
|||||||
|
|
||||||
let token = CancellationToken::new();
|
let token = CancellationToken::new();
|
||||||
|
|
||||||
let message_handler = config.admin_handler.map(|c| c.into_handler());
|
let message_handler = config.admin_handler.map(|cmd| match cmd {
|
||||||
|
AdminHandlerCommand::AdminHttp(config) => config.into_handler(),
|
||||||
|
});
|
||||||
let gateway = Arc::new(Gateway::new(config.gateway, token.clone(), message_handler).await);
|
let gateway = Arc::new(Gateway::new(config.gateway, token.clone(), message_handler).await);
|
||||||
|
|
||||||
let listener = TcpListener::bind(config.http_listen_addr).await.unwrap();
|
let listener = TcpListener::bind(config.http_listen_addr).await.unwrap();
|
||||||
|
|||||||
Reference in New Issue
Block a user