use serde::{Deserialize, Serialize};
use tokio::io::{AsyncRead, AsyncWrite, AsyncWriteExt, BufReader};
use tokio::net::TcpListener;
use tokio::sync::mpsc;
use tracing::{debug, error, trace, warn};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum TelnetCommand {
Ping(PingSpec),
Traceroute(TracerouteSpec),
Dns(DnsSpec),
Http(HttpSpec),
Tls(TlsSpec),
Ntp(NtpSpec),
HostTelemetry(HostTelemetrySpec),
Status,
Stop(u64),
Ignored(String),
Unknown(String),
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
pub enum HostTelemetryKind {
Buddyinfo,
Rptaddrs,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HostTelemetrySpec {
pub kind: HostTelemetryKind,
pub interval: u64,
pub msm_id: Option<u32>,
pub lowmem: Option<u32>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct ScheduleSpec {
pub interval: u64,
pub start_time: i64,
pub stop_time: i64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PingSpec {
pub msm_id: u64,
pub target: String,
pub af: u8, pub packets: u32,
pub size: u16,
pub packet_interval: u32, pub spread: Option<u32>, pub schedule: ScheduleSpec,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TracerouteSpec {
pub msm_id: u64,
pub target: String,
pub af: u8,
pub protocol: String, pub paris: Option<u32>,
pub first_hop: u8,
pub max_hops: u8,
pub size: u16,
pub spread: Option<u32>,
pub schedule: ScheduleSpec,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DnsSpec {
pub msm_id: u64,
pub target: String, pub af: u8,
pub protocol: String, pub query_type: String,
pub query_class: String,
pub query_argument: String,
pub use_dnssec: bool,
pub recursion_desired: bool,
pub spread: Option<u32>,
pub schedule: ScheduleSpec,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HttpSpec {
pub msm_id: u64,
pub url: String,
pub method: String,
pub af: u8,
pub headers: Vec<String>,
pub body: Option<String>,
pub max_body_size: Option<u32>,
pub spread: Option<u32>,
pub schedule: ScheduleSpec,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TlsSpec {
pub msm_id: u64,
pub target: String,
pub port: u16,
pub af: u8,
pub hostname: Option<String>,
pub spread: Option<u32>,
pub schedule: ScheduleSpec,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct NtpSpec {
pub msm_id: u64,
pub target: String,
pub af: u8,
pub packets: u32,
pub spread: Option<u32>,
pub schedule: ScheduleSpec,
}
pub struct TelnetServer {
port: u16,
probe_id: u32,
session_id: std::sync::Arc<tokio::sync::RwLock<Option<String>>>,
command_tx: Option<mpsc::Sender<TelnetCommand>>,
}
const ATLAS_LOGIN: &str = "C_TO_P_TEST_V1";
const LOGIN_PREFIX: &str = "Atlas probe, see http://atlas.ripe.net/\r\n\r\n";
const LOGIN_PROMPT: &str = " login: ";
const PASSWORD_PROMPT: &str = "\r\nPassword: ";
const RESULT_OK: &str = "OK\r\n\r\n";
const BAD_PASSWORD: &str = "BAD_PASSWORD\r\n\r\n";
#[allow(dead_code)]
const BAD_COMMAND: &str = "BAD_COMMAND\r\n\r\n";
impl TelnetServer {
pub fn new(port: u16, probe_id: u32) -> Self {
Self {
port,
probe_id,
session_id: std::sync::Arc::new(tokio::sync::RwLock::new(None)),
command_tx: None,
}
}
pub fn with_channel(port: u16, probe_id: u32, tx: mpsc::Sender<TelnetCommand>) -> Self {
Self {
port,
probe_id,
session_id: std::sync::Arc::new(tokio::sync::RwLock::new(None)),
command_tx: Some(tx),
}
}
pub async fn set_session_id(&self, session_id: String) {
debug!("Setting session ID for telnet authentication");
*self.session_id.write().await = Some(session_id);
}
pub fn port(&self) -> u16 {
self.port
}
pub async fn run(&self) -> anyhow::Result<()> {
let addr = format!("127.0.0.1:{}", self.port);
let listener = TcpListener::bind(&addr).await?;
debug!("Telnet server listening on {}", addr);
loop {
let (socket, remote_addr) = listener.accept().await?;
debug!("Accepted telnet connection from {}", remote_addr);
let command_tx = self.command_tx.clone();
let probe_id = self.probe_id;
let session_id = self.session_id.clone();
tokio::spawn(async move {
if let Err(e) = handle_connection(socket, command_tx, probe_id, session_id).await {
error!("Error handling telnet connection: {}", e);
}
debug!("Telnet connection from {} closed", remote_addr);
});
}
}
}
const IAC: u8 = 255; const WILL: u8 = 251;
const WONT: u8 = 252;
const DO: u8 = 253;
const DONT: u8 = 254;
const SB: u8 = 250; const SE: u8 = 240;
const TELOPT_ECHO: u8 = 1;
const TELOPT_SGA: u8 = 3; const TELOPT_NAWS: u8 = 31;
#[derive(Debug, Clone, Copy, PartialEq)]
enum ConnectionState {
AwaitingLoginName,
AwaitingPassword,
Authenticated,
}
pub async fn handle_connection(
stream: impl AsyncRead + AsyncWrite + Unpin + Send + 'static,
command_tx: Option<mpsc::Sender<TelnetCommand>>,
probe_id: u32,
session_id: std::sync::Arc<tokio::sync::RwLock<Option<String>>>,
) -> anyhow::Result<()> {
let (reader, mut writer) = tokio::io::split(stream);
let initial_iacs: &[u8] = &[
IAC,
DO,
TELOPT_ECHO, IAC,
DO,
TELOPT_NAWS, IAC,
WILL,
TELOPT_ECHO, IAC,
WILL,
TELOPT_SGA, ];
if let Err(e) = writer.write_all(initial_iacs).await {
warn!("Failed to send telnet negotiation: {}", e);
}
let hostname = hostname::get()
.map(|h: std::ffi::OsString| h.to_string_lossy().to_string())
.unwrap_or_else(|_| "unknown".to_string());
let banner = format!(
"{}Probe {} ({}){}",
LOGIN_PREFIX, probe_id, hostname, LOGIN_PROMPT
);
if let Err(e) = writer.write_all(banner.as_bytes()).await {
warn!("Failed to send login banner: {}", e);
}
let _ = writer.flush().await;
trace!("Sent telnet login banner for probe {}", probe_id);
let mut reader = BufReader::new(reader);
let mut raw_buf = [0u8; 4096];
let mut state = ConnectionState::AwaitingLoginName;
let mut line_buffer = String::new();
const MAX_LINE_LEN: usize = 256 * 1024;
loop {
use tokio::io::AsyncReadExt;
match reader.read(&mut raw_buf).await {
Ok(0) => break, Ok(n) => {
let text = filter_telnet_commands(&raw_buf[..n]);
if text.is_empty() {
continue;
}
if line_buffer.len() + text.len() > MAX_LINE_LEN {
error!(
"Telnet line buffer exceeded {}KB, disconnecting",
MAX_LINE_LEN / 1024
);
anyhow::bail!("Line buffer overflow");
}
line_buffer.push_str(&text);
while let Some(newline_pos) = line_buffer.find(['\n', '\r']) {
let line: String = line_buffer.drain(..=newline_pos).collect();
let line = line.trim();
if line.is_empty() {
continue;
}
debug!("Telnet state {:?}, received: '{}'", state, line);
match state {
ConnectionState::AwaitingLoginName => {
if line == ATLAS_LOGIN {
trace!("Atlas login received, requesting password");
if let Err(e) = writer.write_all(PASSWORD_PROMPT.as_bytes()).await {
warn!("Failed to send password prompt: {}", e);
}
let _ = writer.flush().await;
state = ConnectionState::AwaitingPassword;
} else {
debug!("Unknown login '{}', closing connection", line);
break;
}
}
ConnectionState::AwaitingPassword => {
let valid = {
let sid = session_id.read().await;
sid.as_ref().map(|s| s == line).unwrap_or(false)
};
if valid {
debug!("Telnet session authenticated");
if let Err(e) = writer.write_all(RESULT_OK.as_bytes()).await {
warn!("Failed to send OK: {}", e);
}
let _ = writer.flush().await;
state = ConnectionState::Authenticated;
} else {
warn!("Bad password received");
if let Err(e) = writer.write_all(BAD_PASSWORD.as_bytes()).await {
warn!("Failed to send BAD_PASSWORD: {}", e);
}
let _ = writer.flush().await;
break; }
}
ConnectionState::Authenticated => {
trace!("Received command: {}", line);
let command = parse_command(line);
let response = match &command {
TelnetCommand::Unknown(s) => {
format!("ERROR: Unknown command: {}\r\n\r\n", s)
}
_ => RESULT_OK.to_string(),
};
if let Err(e) = writer.write_all(response.as_bytes()).await {
debug!("Failed to send response: {}", e);
}
let _ = writer.flush().await;
if let Some(ref tx) = command_tx {
if let Err(e) = tx.send(command).await {
debug!("Command channel closed (shutdown in progress): {}", e);
break; }
}
}
}
}
}
Err(e) => {
error!("Error reading from telnet: {}", e);
break;
}
}
}
Ok(())
}
fn filter_telnet_commands(data: &[u8]) -> String {
let mut result = Vec::new();
let mut i = 0;
while i < data.len() {
if data[i] == IAC && i + 1 < data.len() {
match data[i + 1] {
IAC => {
result.push(IAC);
i += 2;
}
WILL | WONT | DO | DONT => {
if i + 2 < data.len() {
debug!(
"Telnet: {:?} option {}",
match data[i + 1] {
WILL => "WILL",
WONT => "WONT",
DO => "DO",
DONT => "DONT",
_ => "?",
},
data[i + 2]
);
i += 3;
} else {
i += 2;
}
}
SB => {
i += 2;
while i + 1 < data.len() {
if data[i] == IAC && data[i + 1] == SE {
i += 2;
break;
}
i += 1;
}
}
240..=249 => {
i += 2;
}
_ => {
i += 2;
}
}
} else {
result.push(data[i]);
i += 1;
}
}
String::from_utf8_lossy(&result).to_string()
}
pub fn parse_command(cmd: &str) -> TelnetCommand {
if cmd.starts_with('{') {
if let Ok(spec) = serde_json::from_str::<serde_json::Value>(cmd) {
return parse_json_command(&spec);
}
}
let parts: Vec<&str> = cmd.split_whitespace().collect();
if parts.is_empty() {
return TelnetCommand::Unknown(cmd.to_string());
}
match parts[0].to_uppercase().as_str() {
"STATUS" => TelnetCommand::Status,
"STOP" if parts.len() >= 2 => {
if let Ok(msm_id) = parts[1].parse() {
TelnetCommand::Stop(msm_id)
} else {
TelnetCommand::Unknown(cmd.to_string())
}
}
"CRONLINE" => parse_cronline(cmd),
"ONEOFF" => {
if parts.len() >= 3 {
let measurement_parts = &parts[2..]; let fake_cronline =
format!("CRONLINE 0 0 0 UNIFORM 0 {}", measurement_parts.join(" "));
parse_cronline(&fake_cronline)
} else {
TelnetCommand::Unknown(cmd.to_string())
}
}
"CRONTAB" => {
trace!("CRONTAB command: {}", cmd);
TelnetCommand::Ignored(cmd.to_string())
}
_ => TelnetCommand::Unknown(cmd.to_string()),
}
}
fn parse_cronline(cmd: &str) -> TelnetCommand {
let tokens = tokenize_cronline(cmd);
if tokens.len() < 7 {
warn!("CRONLINE too short: {}", cmd);
return TelnetCommand::Unknown(cmd.to_string());
}
let interval: u64 = match tokens[1].parse() {
Ok(v) => v,
Err(_) => {
warn!("CRONLINE invalid interval '{}': {}", tokens[1], cmd);
return TelnetCommand::Unknown(cmd.to_string());
}
};
let _offset: u64 = tokens[2].parse().unwrap_or(0);
let end_time: i64 = tokens[3].parse().unwrap_or(0);
let _spread_type = &tokens[4]; let spread: u32 = tokens[5].parse().unwrap_or(0);
let measurement_cmd = &tokens[6];
let measurement_args = &tokens[7..];
trace!(
"CRONLINE: interval={}, end_time={}, spread={}, cmd={}, args={:?}",
interval,
end_time,
spread,
measurement_cmd,
measurement_args
);
match measurement_cmd.as_str() {
"evping" => parse_evping(measurement_args, interval, end_time, spread),
"evtraceroute" => parse_evtraceroute(measurement_args, interval, end_time, spread),
"evtdig" => parse_evtdig(measurement_args, interval, end_time, spread),
"evhttpget" => parse_evhttpget(measurement_args, interval, end_time, spread),
"evsslgetcert" => parse_evsslgetcert(measurement_args, interval, end_time, spread),
"evntp" => parse_evntp(measurement_args, interval, end_time, spread),
"buddyinfo" => parse_buddyinfo(measurement_args, interval),
"rptaddrs" => parse_rptaddrs(measurement_args, interval),
"httppost" => {
trace!("Ignoring httppost CRONLINE (uploader handles uploads natively)");
TelnetCommand::Ignored(cmd.to_string())
}
"condmv" | "dfrm" => {
trace!(
"Ignoring {} CRONLINE (no on-disk spool to manage)",
measurement_cmd
);
TelnetCommand::Ignored(cmd.to_string())
}
"conntrack" => {
trace!("Ignoring conntrack CRONLINE (no reference impl available)");
TelnetCommand::Ignored(cmd.to_string())
}
_ => {
warn!("Unknown measurement command: {}", measurement_cmd);
TelnetCommand::Unknown(cmd.to_string())
}
}
}
fn parse_buddyinfo(args: &[String], interval: u64) -> TelnetCommand {
let lowmem = args.first().and_then(|s| s.parse::<u32>().ok());
TelnetCommand::HostTelemetry(HostTelemetrySpec {
kind: HostTelemetryKind::Buddyinfo,
interval,
msm_id: None,
lowmem,
})
}
fn parse_rptaddrs(args: &[String], interval: u64) -> TelnetCommand {
let mut msm_id: Option<u32> = None;
let mut it = args.iter();
while let Some(arg) = it.next() {
match arg.as_str() {
"-A" => {
msm_id = it.next().and_then(|s| s.parse().ok());
}
"-c" | "-O" => {
let _ = it.next();
}
_ => {}
}
}
TelnetCommand::HostTelemetry(HostTelemetrySpec {
kind: HostTelemetryKind::Rptaddrs,
interval,
msm_id,
lowmem: None,
})
}
fn tokenize_cronline(cmd: &str) -> Vec<String> {
let mut tokens = Vec::new();
let mut current = String::new();
let mut in_quotes = false;
for c in cmd.chars() {
match c {
'"' => {
in_quotes = !in_quotes;
}
' ' | '\t' if !in_quotes => {
if !current.is_empty() {
tokens.push(current.clone());
current.clear();
}
}
_ => {
current.push(c);
}
}
}
if !current.is_empty() {
tokens.push(current);
}
tokens
}
fn parse_evping(args: &[String], interval: u64, end_time: i64, spread: u32) -> TelnetCommand {
let mut af: u8 = 4;
let mut count: u32 = 3;
let mut size: u16 = 64;
let mut msm_id: u64 = 0;
let mut target = String::new();
let mut i = 0;
while i < args.len() {
let arg = &args[i];
match arg.as_str() {
"-4" => af = 4,
"-6" => af = 6,
"-c" => {
if i + 1 < args.len() {
count = args[i + 1].parse().unwrap_or(3);
i += 1;
}
}
"-s" => {
if i + 1 < args.len() {
size = args[i + 1].parse().unwrap_or(64);
i += 1;
}
}
"-A" => {
if i + 1 < args.len() {
msm_id = args[i + 1].parse().unwrap_or(0);
i += 1;
}
}
"-O" => {
if i + 1 < args.len() {
i += 1;
}
}
"-I" => {
if i + 1 < args.len() {
i += 1;
}
}
"-R" => {
}
_ if !arg.starts_with('-') => {
target = arg.clone();
}
_ => {
debug!("evping: ignoring unknown option {}", arg);
}
}
i += 1;
}
if target.is_empty() {
warn!("evping: no target specified");
return TelnetCommand::Unknown(format!("evping {:?}", args));
}
trace!(
"Parsed evping: msm_id={}, target={}, af={}, count={}, size={}, interval={}",
msm_id,
target,
af,
count,
size,
interval
);
TelnetCommand::Ping(PingSpec {
msm_id,
target,
af,
packets: count,
size,
packet_interval: 1000, spread: Some(spread),
schedule: ScheduleSpec {
interval,
start_time: 0, stop_time: end_time,
},
})
}
fn parse_evtraceroute(args: &[String], interval: u64, end_time: i64, spread: u32) -> TelnetCommand {
let mut af: u8 = 4;
let mut protocol = "ICMP".to_string();
let mut first_hop: u8 = 1;
let mut max_hops: u8 = 32;
let mut paris: Option<u32> = None;
let mut msm_id: u64 = 0;
let mut target = String::new();
let mut size: u16 = 40;
let mut i = 0;
while i < args.len() {
let arg = &args[i];
match arg.as_str() {
"-4" => af = 4,
"-6" => af = 6,
"-I" => protocol = "ICMP".to_string(),
"-U" => protocol = "UDP".to_string(),
"-T" => protocol = "TCP".to_string(),
"-a" => {
if i + 1 < args.len() {
i += 1;
}
}
"-c" => {
if i + 1 < args.len() {
i += 1;
}
}
"-w" => {
if i + 1 < args.len() {
i += 1;
}
}
"-f" => {
if i + 1 < args.len() {
first_hop = args[i + 1].parse().unwrap_or(1);
i += 1;
}
}
"-m" => {
if i + 1 < args.len() {
max_hops = args[i + 1].parse().unwrap_or(32);
i += 1;
}
}
"-p" => {
if i + 1 < args.len() {
paris = Some(args[i + 1].parse().unwrap_or(0));
i += 1;
}
}
"-s" | "-S" => {
if i + 1 < args.len() {
size = args[i + 1].parse().unwrap_or(40);
i += 1;
}
}
"-A" => {
if i + 1 < args.len() {
msm_id = args[i + 1].parse().unwrap_or(0);
i += 1;
}
}
"-O" => {
if i + 1 < args.len() {
i += 1;
}
}
_ if !arg.starts_with('-') => {
target = arg.clone();
}
_ => {
debug!("evtraceroute: ignoring unknown option {}", arg);
}
}
i += 1;
}
if target.is_empty() {
warn!("evtraceroute: no target specified");
return TelnetCommand::Unknown(format!("evtraceroute {:?}", args));
}
trace!(
"Parsed evtraceroute: msm_id={}, target={}, af={}, protocol={}, hops={}-{}, interval={}",
msm_id,
target,
af,
protocol,
first_hop,
max_hops,
interval
);
TelnetCommand::Traceroute(TracerouteSpec {
msm_id,
target,
af,
protocol,
paris,
first_hop,
max_hops,
size,
spread: Some(spread),
schedule: ScheduleSpec {
interval,
start_time: 0,
stop_time: end_time,
},
})
}
fn parse_evtdig(args: &[String], interval: u64, end_time: i64, spread: u32) -> TelnetCommand {
let mut af: u8 = 4;
let mut protocol = "UDP".to_string();
let mut query_type = "A".to_string();
let mut query_class = "IN".to_string();
let mut use_dnssec = false;
let mut recursion_desired = true;
let mut msm_id: u64 = 0;
let mut target = String::new(); let mut query_argument = String::new(); let mut use_resolver = false;
let mut i = 0;
while i < args.len() {
let arg = &args[i];
match arg.as_str() {
"-4" => af = 4,
"-6" => af = 6,
"-p" => {
if i + 1 < args.len() {
i += 1;
}
}
"--soa" => {
query_type = "SOA".to_string();
if query_argument.is_empty() {
query_argument = ".".to_string();
}
}
"--resolv" => {
use_resolver = true;
}
"-h" => {
query_type = "TXT".to_string();
query_class = "CH".to_string();
query_argument = "hostname.bind".to_string();
}
"-b" => {
query_type = "TXT".to_string();
query_class = "CH".to_string();
query_argument = "version.bind".to_string();
}
"-i" => {
query_type = "TXT".to_string();
query_class = "CH".to_string();
query_argument = "id.server".to_string();
}
"-r" => {
use_dnssec = true;
}
"-d" | "-D" | "+dnssec" => {
use_dnssec = true;
}
"-t" => {
protocol = "TCP".to_string();
}
"--type" | "-type" => {
if i + 1 < args.len() {
let type_num: u16 = args[i + 1].parse().unwrap_or(1);
query_type = match type_num {
1 => "A",
2 => "NS",
5 => "CNAME",
6 => "SOA",
12 => "PTR",
15 => "MX",
16 => "TXT",
28 => "AAAA",
33 => "SRV",
43 => "DS",
46 => "RRSIG",
47 => "NSEC",
48 => "DNSKEY",
_ => "A",
}
.to_string();
i += 1;
}
}
"--class" | "-class" => {
if i + 1 < args.len() {
let class_num: u16 = args[i + 1].parse().unwrap_or(1);
query_class = match class_num {
1 => "IN",
3 => "CH",
_ => "IN",
}
.to_string();
i += 1;
}
}
"--query" | "-query" => {
if i + 1 < args.len() {
query_argument = args[i + 1].clone();
i += 1;
}
}
"--a" | "-a" => {
query_type = "A".to_string();
if i + 1 < args.len() && !args[i + 1].starts_with('-') {
query_argument = args[i + 1].clone();
i += 1;
}
}
"--aaaa" | "-aaaa" => {
query_type = "AAAA".to_string();
if i + 1 < args.len() && !args[i + 1].starts_with('-') {
query_argument = args[i + 1].clone();
i += 1;
}
}
"-e" => {
if i + 1 < args.len() {
i += 1;
}
}
"-R" => {
recursion_desired = false;
}
"--retry" => {
if i + 1 < args.len() {
i += 1;
}
}
"--qbuf" => {
}
"-T" => {
protocol = "TCP".to_string();
}
"-A" => {
if i + 1 < args.len() {
msm_id = args[i + 1].parse().unwrap_or(0);
i += 1;
}
}
"-O" => {
if i + 1 < args.len() {
i += 1;
}
}
_ if arg.starts_with('@') => {
target = arg[1..].to_string();
}
_ if arg.starts_with('+') => {
match arg.as_str() {
"+dnssec" | "+do" => use_dnssec = true,
"+tcp" => protocol = "TCP".to_string(),
"+nord" | "+norecurse" => recursion_desired = false,
_ => debug!("evtdig: ignoring option {}", arg),
}
}
_ if arg.starts_with("--") => {
debug!("evtdig: ignoring unknown option {}", arg);
}
_ if !arg.starts_with('-') => {
if target.is_empty() && arg.parse::<std::net::IpAddr>().is_ok() {
target = arg.clone();
} else if query_argument.is_empty() {
query_argument = arg.clone();
}
}
_ => {
debug!("evtdig: ignoring unknown option {}", arg);
}
}
i += 1;
}
if use_resolver && target.is_empty() {
target = if af == 4 {
"8.8.8.8"
} else {
"2001:4860:4860::8888"
}
.to_string();
}
if target.is_empty()
&& !query_argument.is_empty()
&& query_argument.parse::<std::net::IpAddr>().is_ok()
{
target = query_argument.clone();
query_argument = ".".to_string();
}
if target.is_empty() {
warn!("evtdig: no DNS server specified, args={:?}", args);
return TelnetCommand::Unknown(format!("evtdig {:?}", args));
}
if query_argument.is_empty() {
query_argument = ".".to_string();
}
if query_argument.ends_with('.') && query_argument.len() > 1 {
query_argument = query_argument[..query_argument.len() - 1].to_string();
}
trace!(
"Parsed evtdig: msm_id={}, server={}, query={} {} {}, protocol={}, interval={}",
msm_id,
target,
query_argument,
query_type,
query_class,
protocol,
interval
);
TelnetCommand::Dns(DnsSpec {
msm_id,
target,
af,
protocol,
query_type,
query_class,
query_argument,
use_dnssec,
recursion_desired,
spread: Some(spread),
schedule: ScheduleSpec {
interval,
start_time: 0,
stop_time: end_time,
},
})
}
fn parse_evhttpget(args: &[String], interval: u64, end_time: i64, spread: u32) -> TelnetCommand {
let mut af: u8 = 4;
let mut msm_id: u64 = 0;
let mut url = String::new();
let mut max_body_size: Option<u32> = None;
let mut i = 0;
while i < args.len() {
let arg = &args[i];
match arg.as_str() {
"-4" => af = 4,
"-6" => af = 6,
"-1" => {
}
"-A" => {
if i + 1 < args.len() {
msm_id = args[i + 1].parse().unwrap_or(0);
i += 1;
}
}
"-O" => {
if i + 1 < args.len() {
i += 1;
}
}
"-M" | "--max-body" => {
if i + 1 < args.len() {
max_body_size = args[i + 1].parse().ok();
i += 1;
}
}
"--store-headers" => {
if i + 1 < args.len() {
i += 1;
}
}
_ if !arg.starts_with('-') => {
url = arg.clone();
}
_ => {
debug!("evhttpget: ignoring unknown option {}", arg);
}
}
i += 1;
}
if url.is_empty() {
warn!("evhttpget: no URL specified");
return TelnetCommand::Unknown(format!("evhttpget {:?}", args));
}
trace!(
"Parsed evhttpget: msm_id={}, url={}, af={}, interval={}",
msm_id,
url,
af,
interval
);
TelnetCommand::Http(HttpSpec {
msm_id,
url,
method: "GET".to_string(),
af,
headers: vec![],
body: None,
max_body_size,
spread: Some(spread),
schedule: ScheduleSpec {
interval,
start_time: 0,
stop_time: end_time,
},
})
}
fn parse_evsslgetcert(args: &[String], interval: u64, end_time: i64, spread: u32) -> TelnetCommand {
let mut af: u8 = 4;
let mut port: u16 = 443;
let mut msm_id: u64 = 0;
let mut target = String::new();
let mut hostname: Option<String> = None;
let mut i = 0;
while i < args.len() {
let arg = &args[i];
match arg.as_str() {
"-4" => af = 4,
"-6" => af = 6,
"-p" => {
if i + 1 < args.len() {
port = args[i + 1].parse().unwrap_or(443);
i += 1;
}
}
"-A" => {
if i + 1 < args.len() {
msm_id = args[i + 1].parse().unwrap_or(0);
i += 1;
}
}
"-O" => {
if i + 1 < args.len() {
i += 1;
}
}
"-h" | "-H" | "--hostname" => {
if i + 1 < args.len() {
hostname = Some(args[i + 1].clone());
i += 1;
}
}
_ if !arg.starts_with('-') => {
target = arg.clone();
}
_ => {
debug!("evsslgetcert: ignoring unknown option {}", arg);
}
}
i += 1;
}
if target.is_empty() {
warn!("evsslgetcert: no target specified");
return TelnetCommand::Unknown(format!("evsslgetcert {:?}", args));
}
trace!(
"Parsed evsslgetcert: msm_id={}, target={}, hostname={:?}, port={}, af={}, interval={}",
msm_id,
target,
hostname,
port,
af,
interval
);
TelnetCommand::Tls(TlsSpec {
msm_id,
target: target.clone(),
port,
af,
hostname: hostname.or(Some(target)),
spread: Some(spread),
schedule: ScheduleSpec {
interval,
start_time: 0,
stop_time: end_time,
},
})
}
fn parse_evntp(args: &[String], interval: u64, end_time: i64, spread: u32) -> TelnetCommand {
let mut af: u8 = 4;
let mut packets: u32 = 3;
let mut msm_id: u64 = 0;
let mut target = String::new();
let mut i = 0;
while i < args.len() {
let arg = &args[i];
match arg.as_str() {
"-4" => af = 4,
"-6" => af = 6,
"-c" => {
if i + 1 < args.len() {
packets = args[i + 1].parse().unwrap_or(3);
i += 1;
}
}
"-A" => {
if i + 1 < args.len() {
msm_id = args[i + 1].parse().unwrap_or(0);
i += 1;
}
}
"-O" => {
if i + 1 < args.len() {
i += 1;
}
}
_ if !arg.starts_with('-') => {
target = arg.clone();
}
_ => {
debug!("evntp: ignoring unknown option {}", arg);
}
}
i += 1;
}
if target.is_empty() {
warn!("evntp: no target specified");
return TelnetCommand::Unknown(format!("evntp {:?}", args));
}
trace!(
"Parsed evntp: msm_id={}, target={}, af={}, packets={}, interval={}",
msm_id,
target,
af,
packets,
interval
);
TelnetCommand::Ntp(NtpSpec {
msm_id,
target,
af,
packets,
spread: Some(spread),
schedule: ScheduleSpec {
interval,
start_time: 0,
stop_time: end_time,
},
})
}
fn parse_schedule(spec: &serde_json::Value) -> ScheduleSpec {
ScheduleSpec {
interval: spec.get("interval").and_then(|v| v.as_u64()).unwrap_or(0),
start_time: spec.get("start_time").and_then(|v| v.as_i64()).unwrap_or(0),
stop_time: spec.get("stop_time").and_then(|v| v.as_i64()).unwrap_or(0),
}
}
fn parse_json_command(spec: &serde_json::Value) -> TelnetCommand {
let msm_type = spec
.get("type")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_lowercase();
let msm_id = spec.get("msm_id").and_then(|v| v.as_u64()).unwrap_or(0);
let target = spec
.get("target")
.or_else(|| spec.get("dst_name"))
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let af = spec.get("af").and_then(|v| v.as_u64()).unwrap_or(4) as u8;
let spread = spec
.get("spread")
.and_then(|v| v.as_u64())
.map(|v| v as u32);
let schedule = parse_schedule(spec);
match msm_type.as_str() {
"ping" => TelnetCommand::Ping(PingSpec {
msm_id,
target,
af,
packets: spec.get("packets").and_then(|v| v.as_u64()).unwrap_or(3) as u32,
size: spec.get("size").and_then(|v| v.as_u64()).unwrap_or(64) as u16,
packet_interval: spec
.get("packet_interval")
.and_then(|v| v.as_u64())
.unwrap_or(1000) as u32,
spread,
schedule,
}),
"traceroute" => TelnetCommand::Traceroute(TracerouteSpec {
msm_id,
target,
af,
protocol: spec
.get("protocol")
.and_then(|v| v.as_str())
.unwrap_or("ICMP")
.to_string(),
paris: spec.get("paris").and_then(|v| v.as_u64()).map(|v| v as u32),
first_hop: spec.get("first_hop").and_then(|v| v.as_u64()).unwrap_or(1) as u8,
max_hops: spec.get("max_hops").and_then(|v| v.as_u64()).unwrap_or(32) as u8,
size: spec.get("size").and_then(|v| v.as_u64()).unwrap_or(40) as u16,
spread,
schedule,
}),
"dns" => TelnetCommand::Dns(DnsSpec {
msm_id,
target,
af,
protocol: spec
.get("protocol")
.and_then(|v| v.as_str())
.unwrap_or("UDP")
.to_string(),
query_type: spec
.get("query_type")
.and_then(|v| v.as_str())
.unwrap_or("A")
.to_string(),
query_class: spec
.get("query_class")
.and_then(|v| v.as_str())
.unwrap_or("IN")
.to_string(),
query_argument: spec
.get("query_argument")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string(),
use_dnssec: spec
.get("use_dnssec")
.and_then(|v| v.as_bool())
.unwrap_or(false),
recursion_desired: spec
.get("recursion_desired")
.and_then(|v| v.as_bool())
.unwrap_or(true),
spread,
schedule,
}),
"http" | "https" => TelnetCommand::Http(HttpSpec {
msm_id,
url: spec
.get("url")
.or_else(|| spec.get("target"))
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string(),
method: spec
.get("method")
.and_then(|v| v.as_str())
.unwrap_or("GET")
.to_string(),
af,
headers: spec
.get("headers")
.and_then(|v| v.as_array())
.map(|arr| {
arr.iter()
.filter_map(|v| v.as_str().map(String::from))
.collect()
})
.unwrap_or_default(),
body: spec.get("body").and_then(|v| v.as_str()).map(String::from),
max_body_size: spec
.get("max_body_size")
.and_then(|v| v.as_u64())
.map(|v| v as u32),
spread,
schedule,
}),
"sslcert" | "tls" => TelnetCommand::Tls(TlsSpec {
msm_id,
target: target.clone(),
port: spec.get("port").and_then(|v| v.as_u64()).unwrap_or(443) as u16,
af,
hostname: spec
.get("hostname")
.and_then(|v| v.as_str())
.map(String::from)
.or(Some(target)),
spread,
schedule,
}),
"ntp" => TelnetCommand::Ntp(NtpSpec {
msm_id,
target,
af,
packets: spec.get("packets").and_then(|v| v.as_u64()).unwrap_or(3) as u32,
spread,
schedule,
}),
_ => TelnetCommand::Unknown(format!("Unknown measurement type: {}", msm_type)),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_parse_ping_command() {
let json = r#"{"type":"ping","msm_id":1001,"target":"8.8.8.8","af":4,"packets":3}"#;
let cmd = parse_command(json);
match cmd {
TelnetCommand::Ping(spec) => {
assert_eq!(spec.msm_id, 1001);
assert_eq!(spec.target, "8.8.8.8");
assert_eq!(spec.packets, 3);
assert_eq!(spec.schedule.interval, 0); }
_ => panic!("Expected Ping command"),
}
}
#[test]
fn test_parse_recurring_ping_command() {
let json = r#"{"type":"ping","msm_id":1001,"target":"8.8.8.8","af":4,"packets":3,"interval":300,"start_time":1700000000,"stop_time":1700086400}"#;
let cmd = parse_command(json);
match cmd {
TelnetCommand::Ping(spec) => {
assert_eq!(spec.msm_id, 1001);
assert_eq!(spec.schedule.interval, 300);
assert_eq!(spec.schedule.start_time, 1700000000);
assert_eq!(spec.schedule.stop_time, 1700086400);
}
_ => panic!("Expected Ping command"),
}
}
#[test]
fn test_parse_dns_command() {
let json = r#"{"type":"dns","msm_id":1002,"target":"9.9.9.9","query_argument":"example.com","query_type":"A"}"#;
let cmd = parse_command(json);
match cmd {
TelnetCommand::Dns(spec) => {
assert_eq!(spec.msm_id, 1002);
assert_eq!(spec.target, "9.9.9.9");
assert_eq!(spec.query_argument, "example.com");
}
_ => panic!("Expected Dns command"),
}
}
#[test]
fn test_parse_status_command() {
let cmd = parse_command("STATUS");
assert!(matches!(cmd, TelnetCommand::Status));
}
#[test]
fn test_parse_stop_command() {
let cmd = parse_command("STOP 12345");
match cmd {
TelnetCommand::Stop(id) => assert_eq!(id, 12345),
_ => panic!("Expected Stop command"),
}
}
#[test]
fn test_tokenize_cronline() {
let tokens = tokenize_cronline(
r#"CRONLINE 240 274 1770761451 UNIFORM 3 evping -4 -c 3 -A "1001" -O /home/atlas/data/new/7 193.0.14.129"#,
);
assert_eq!(tokens[0], "CRONLINE");
assert_eq!(tokens[1], "240");
assert_eq!(tokens[6], "evping");
assert_eq!(tokens[11], "1001"); assert_eq!(tokens[tokens.len() - 1], "193.0.14.129");
}
#[test]
fn test_parse_cronline_evping() {
let cmd = parse_command(
r#"CRONLINE 240 274 1770761451 UNIFORM 3 evping -4 -c 3 -A "1001" -O /home/atlas/data/new/7 193.0.14.129"#,
);
match cmd {
TelnetCommand::Ping(spec) => {
assert_eq!(spec.msm_id, 1001);
assert_eq!(spec.target, "193.0.14.129");
assert_eq!(spec.af, 4);
assert_eq!(spec.packets, 3);
assert_eq!(spec.schedule.interval, 240);
assert_eq!(spec.schedule.stop_time, 1770761451);
assert_eq!(spec.spread, Some(3));
}
_ => panic!("Expected Ping command, got {:?}", cmd),
}
}
#[test]
fn test_parse_cronline_evping_ipv6() {
let cmd = parse_command(
r#"CRONLINE 900 123 1800000000 UNIFORM 5 evping -6 -c 5 -s 128 -A "2001" 2001:4860:4860::8888"#,
);
match cmd {
TelnetCommand::Ping(spec) => {
assert_eq!(spec.msm_id, 2001);
assert_eq!(spec.target, "2001:4860:4860::8888");
assert_eq!(spec.af, 6);
assert_eq!(spec.packets, 5);
assert_eq!(spec.size, 128);
assert_eq!(spec.schedule.interval, 900);
}
_ => panic!("Expected Ping command"),
}
}
#[test]
fn test_parse_cronline_evtraceroute_udp() {
let cmd = parse_command(
r#"CRONLINE 1800 331 1925935851 UNIFORM 3 evtraceroute -4 -U -c 3 -w 1000 -A "5001" -O /home/atlas/data/new/7 193.0.14.129"#,
);
match cmd {
TelnetCommand::Traceroute(spec) => {
assert_eq!(spec.msm_id, 5001);
assert_eq!(spec.target, "193.0.14.129");
assert_eq!(spec.af, 4);
assert_eq!(spec.protocol, "UDP");
assert_eq!(spec.schedule.interval, 1800);
}
_ => panic!("Expected Traceroute command, got {:?}", cmd),
}
}
#[test]
fn test_parse_cronline_evtraceroute_icmp() {
let cmd = parse_command(
r#"CRONLINE 3600 0 0 UNIFORM 10 evtraceroute -4 -I -m 64 -f 1 -A "5002" 8.8.8.8"#,
);
match cmd {
TelnetCommand::Traceroute(spec) => {
assert_eq!(spec.msm_id, 5002);
assert_eq!(spec.target, "8.8.8.8");
assert_eq!(spec.protocol, "ICMP");
assert_eq!(spec.first_hop, 1);
assert_eq!(spec.max_hops, 64);
}
_ => panic!("Expected Traceroute command"),
}
}
#[test]
fn test_parse_cronline_evtdig() {
let cmd = parse_command(
r#"CRONLINE 900 0 0 UNIFORM 5 evtdig -4 --aaaa example.com. -A "6001" @9.9.9.9"#,
);
match cmd {
TelnetCommand::Dns(spec) => {
assert_eq!(spec.msm_id, 6001);
assert_eq!(spec.target, "9.9.9.9");
assert_eq!(spec.query_argument, "example.com"); assert_eq!(spec.query_type, "AAAA");
assert_eq!(spec.af, 4);
}
_ => panic!("Expected Dns command, got {:?}", cmd),
}
}
#[test]
fn test_parse_cronline_evtdig_dnssec() {
let cmd = parse_command(
r#"CRONLINE 1800 0 0 UNIFORM 3 evtdig -4 -d --a dnssec-test.org. -A "6002" @8.8.8.8"#,
);
match cmd {
TelnetCommand::Dns(spec) => {
assert_eq!(spec.msm_id, 6002);
assert!(spec.use_dnssec);
assert_eq!(spec.query_argument, "dnssec-test.org");
}
_ => panic!("Expected Dns command"),
}
}
#[test]
fn test_parse_cronline_evtdig_soa() {
let cmd = parse_command(
r#"CRONLINE 1800 795 1770762780 UNIFORM 3 evtdig -4 --soa . -A "10001" -O /home/atlas/data/new/7 193.0.14.129"#,
);
match cmd {
TelnetCommand::Dns(spec) => {
assert_eq!(spec.msm_id, 10001);
assert_eq!(spec.target, "193.0.14.129");
assert_eq!(spec.query_argument, ".");
assert_eq!(spec.query_type, "SOA");
assert_eq!(spec.af, 4);
assert_eq!(spec.protocol, "UDP");
}
_ => panic!("Expected Dns command, got {:?}", cmd),
}
}
#[test]
fn test_parse_cronline_evtdig_tcp_soa() {
let cmd = parse_command(
r#"CRONLINE 1800 1945 1770762780 UNIFORM 3 evtdig -4 -t --soa . -A "10101" -O /home/atlas/data/new/7 193.0.14.129"#,
);
match cmd {
TelnetCommand::Dns(spec) => {
assert_eq!(spec.msm_id, 10101);
assert_eq!(spec.target, "193.0.14.129");
assert_eq!(spec.query_type, "SOA");
assert_eq!(spec.protocol, "TCP");
}
_ => panic!("Expected Dns command, got {:?}", cmd),
}
}
#[test]
fn test_parse_cronline_evtdig_hostname() {
let cmd = parse_command(
r#"CRONLINE 240 659 1770762780 UNIFORM 3 evtdig -4 -h -A "10301" -O /home/atlas/data/new/7 193.0.14.129"#,
);
match cmd {
TelnetCommand::Dns(spec) => {
assert_eq!(spec.msm_id, 10301);
assert_eq!(spec.target, "193.0.14.129");
assert_eq!(spec.query_argument, "hostname.bind");
assert_eq!(spec.query_type, "TXT");
assert_eq!(spec.query_class, "CH");
}
_ => panic!("Expected Dns command, got {:?}", cmd),
}
}
#[test]
fn test_parse_cronline_evhttpget() {
let cmd = parse_command(
r#"CRONLINE 3600 0 0 UNIFORM 60 evhttpget -4 -A "7001" http://example.com/path"#,
);
match cmd {
TelnetCommand::Http(spec) => {
assert_eq!(spec.msm_id, 7001);
assert_eq!(spec.url, "http://example.com/path");
assert_eq!(spec.method, "GET");
assert_eq!(spec.af, 4);
}
_ => panic!("Expected Http command, got {:?}", cmd),
}
}
#[test]
fn test_parse_cronline_evsslgetcert() {
let cmd = parse_command(
r#"CRONLINE 1800 0 0 UNIFORM 10 evsslgetcert -4 -p 443 -A "8001" example.com"#,
);
match cmd {
TelnetCommand::Tls(spec) => {
assert_eq!(spec.msm_id, 8001);
assert_eq!(spec.target, "example.com");
assert_eq!(spec.port, 443);
assert_eq!(spec.af, 4);
}
_ => panic!("Expected Tls command, got {:?}", cmd),
}
}
#[test]
fn test_parse_cronline_evsslgetcert_custom_port() {
let cmd = parse_command(
r#"CRONLINE 900 0 0 UNIFORM 5 evsslgetcert -6 -p 8443 -A "8002" secure.example.com"#,
);
match cmd {
TelnetCommand::Tls(spec) => {
assert_eq!(spec.msm_id, 8002);
assert_eq!(spec.target, "secure.example.com");
assert_eq!(spec.port, 8443);
assert_eq!(spec.af, 6);
}
_ => panic!("Expected Tls command"),
}
}
#[test]
fn test_parse_cronline_evntp() {
let cmd =
parse_command(r#"CRONLINE 1800 0 0 UNIFORM 3 evntp -4 -c 5 -A "9001" pool.ntp.org"#);
match cmd {
TelnetCommand::Ntp(spec) => {
assert_eq!(spec.msm_id, 9001);
assert_eq!(spec.target, "pool.ntp.org");
assert_eq!(spec.packets, 5);
assert_eq!(spec.af, 4);
}
_ => panic!("Expected Ntp command, got {:?}", cmd),
}
}
#[test]
fn test_parse_cronline_buddyinfo() {
let cmd = parse_command("CRONLINE 600 0 0 UNIFORM 5 buddyinfo 2048");
match cmd {
TelnetCommand::HostTelemetry(spec) => {
assert_eq!(spec.kind, HostTelemetryKind::Buddyinfo);
assert_eq!(spec.interval, 600);
assert_eq!(spec.lowmem, Some(2048));
assert_eq!(spec.msm_id, None);
}
_ => panic!("Expected HostTelemetry/Buddyinfo, got {:?}", cmd),
}
}
#[test]
fn test_parse_cronline_rptaddrs() {
let cmd = parse_command(
r#"CRONLINE 14400 0 0 UNIFORM 30 rptaddrs -A 9104 -c /home/atlas/data/new/v6addr.vol -O /home/atlas/data/new/v6addr.txt"#,
);
match cmd {
TelnetCommand::HostTelemetry(spec) => {
assert_eq!(spec.kind, HostTelemetryKind::Rptaddrs);
assert_eq!(spec.interval, 14400);
assert_eq!(spec.msm_id, Some(9104));
}
_ => panic!("Expected HostTelemetry/Rptaddrs, got {:?}", cmd),
}
}
#[test]
fn test_parse_cronline_httppost_ignored() {
let cmd = parse_command(
"CRONLINE 60 0 0 UNIFORM 5 httppost -O /home/atlas/data/out/foo --post-file foo",
);
assert!(matches!(cmd, TelnetCommand::Ignored(_)));
}
#[test]
fn test_parse_cronline_condmv_ignored() {
let cmd = parse_command(
"CRONLINE 60 0 0 UNIFORM 5 condmv /home/atlas/data/new/foo /home/atlas/data/out/foo",
);
assert!(matches!(cmd, TelnetCommand::Ignored(_)));
}
#[test]
fn test_parse_cronline_conntrack_ignored() {
let cmd = parse_command("CRONLINE 600 0 0 UNIFORM 5 conntrack");
assert!(matches!(cmd, TelnetCommand::Ignored(_)));
}
}