pub mod config;
pub mod document;
pub mod es;
pub mod offsets;
pub mod parse;
pub mod runner;
pub mod sink;
pub mod tail;
use serde::{Deserialize, Serialize};
pub use config::{LogForwardConfig, LogLevel, DEFAULT_ENDPOINT, DEFAULT_INDEX_PREFIX};
pub use document::{ForwardDocument, NodeTags};
pub use es::ElasticsearchSink;
pub use offsets::OffsetStore;
pub use parse::{parse_line, LogEvent};
pub use runner::{classify_nodes, spawn_log_forwarder, ForwarderHandle, DEFAULT_POLL_INTERVAL};
pub use sink::{BatchOutcome, DocumentOutcome, DocumentQueue, LogSink, RetryPolicy};
pub use tail::{LogTailer, TailedEvent};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, utoipa::ToSchema)]
pub struct ForwardingNode {
pub node_id: u32,
pub service: String,
pub log_dir: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, utoipa::ToSchema)]
pub struct SkippedNode {
pub node_id: u32,
pub service: String,
pub reason: String,
}
impl SkippedNode {
#[must_use]
pub fn no_logging(node_id: u32, service: impl Into<String>) -> Self {
Self {
node_id,
service: service.into(),
reason: "logging is not enabled for this node — re-add it with --log-dir-path to \
forward its logs"
.to_string(),
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, utoipa::ToSchema)]
pub struct ForwardStats {
pub events_forwarded: u64,
pub events_dropped_by_level: u64,
pub events_dropped_by_overflow: u64,
pub batches_sent: u64,
pub batches_failed: u64,
pub last_success_unix: Option<u64>,
pub last_error: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, utoipa::ToSchema)]
pub struct LogForwardStatus {
pub enabled: bool,
pub endpoint: String,
pub index_prefix: String,
pub min_level: LogLevel,
pub token_fingerprint: Option<String>,
pub active: bool,
pub nodes_forwarding: Vec<ForwardingNode>,
pub nodes_skipped: Vec<SkippedNode>,
pub stats: ForwardStats,
}
impl LogForwardStatus {
#[must_use]
pub fn inactive(config: &LogForwardConfig) -> Self {
Self {
enabled: config.enabled,
endpoint: config.endpoint.clone(),
index_prefix: config.index_prefix.clone(),
min_level: config.min_level,
token_fingerprint: config.token_fingerprint(),
active: false,
nodes_forwarding: Vec::new(),
nodes_skipped: Vec::new(),
stats: ForwardStats::default(),
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, utoipa::ToSchema)]
pub struct LogForwardEnableRequest {
#[serde(default)]
pub token: Option<String>,
#[serde(default)]
pub endpoint: Option<String>,
#[serde(default)]
pub min_level: Option<LogLevel>,
}
pub fn apply_enable(
stored: &LogForwardConfig,
request: &LogForwardEnableRequest,
) -> crate::error::Result<LogForwardConfig> {
let mut config = stored.clone();
config.enabled = true;
if let Some(token) = &request.token {
config.token = token.trim().to_string();
}
if let Some(endpoint) = &request.endpoint {
config.endpoint = endpoint.trim().to_string();
}
if let Some(level) = request.min_level {
config.min_level = level;
}
config.ensure_installation_id();
config.validate()?;
Ok(config)
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, utoipa::ToSchema)]
pub struct LogForwardResult {
pub enabled: bool,
pub already_in_state: bool,
pub endpoint: String,
pub min_level: LogLevel,
pub nodes_forwarding: Vec<ForwardingNode>,
pub nodes_skipped: Vec<SkippedNode>,
pub pending_daemon_start: bool,
}
#[cfg(test)]
mod tests {
use super::*;
fn stored_with_token() -> LogForwardConfig {
LogForwardConfig {
enabled: false,
token: "stored-key".to_string(),
installation_id: "0123456789abcdef".to_string(),
..LogForwardConfig::disabled()
}
}
#[test]
fn enabling_for_the_first_time_requires_a_token() {
let error = apply_enable(
&LogForwardConfig::disabled(),
&LogForwardEnableRequest::default(),
)
.unwrap_err()
.to_string();
assert!(error.contains("token"), "{error}");
}
#[test]
fn re_enabling_reuses_the_stored_token_and_settings() {
let stored = LogForwardConfig {
endpoint: "http://127.0.0.1:9999".to_string(),
min_level: LogLevel::Warn,
..stored_with_token()
};
let config = apply_enable(&stored, &LogForwardEnableRequest::default()).unwrap();
assert!(config.enabled);
assert_eq!(config.token, "stored-key");
assert_eq!(config.endpoint, "http://127.0.0.1:9999");
assert_eq!(config.min_level, LogLevel::Warn);
}
#[test]
fn enabling_mints_an_installation_id_and_later_enables_keep_it() {
let config = apply_enable(
&LogForwardConfig::disabled(),
&LogForwardEnableRequest {
token: Some("first-key".to_string()),
..LogForwardEnableRequest::default()
},
)
.unwrap();
assert_eq!(config.installation_id.len(), 16);
let re_enabled = apply_enable(&config, &LogForwardEnableRequest::default()).unwrap();
assert_eq!(re_enabled.installation_id, config.installation_id);
let rotated = apply_enable(
&config,
&LogForwardEnableRequest {
token: Some("rotated-key".to_string()),
..LogForwardEnableRequest::default()
},
)
.unwrap();
assert_eq!(rotated.installation_id, config.installation_id);
}
#[test]
fn a_supplied_token_endpoint_and_level_override_what_was_stored() {
let config = apply_enable(
&stored_with_token(),
&LogForwardEnableRequest {
token: Some(" rotated-key ".to_string()),
endpoint: Some("http://localhost:8080".to_string()),
min_level: Some(LogLevel::Error),
},
)
.unwrap();
assert_eq!(config.token, "rotated-key", "surrounding space is trimmed");
assert_eq!(config.endpoint, "http://localhost:8080");
assert_eq!(config.min_level, LogLevel::Error);
}
#[test]
fn an_invalid_endpoint_is_rejected_before_anything_is_persisted() {
let error = apply_enable(
&stored_with_token(),
&LogForwardEnableRequest {
endpoint: Some("logs.autonomi.com".to_string()),
..LogForwardEnableRequest::default()
},
)
.unwrap_err()
.to_string();
assert!(error.contains("http(s) URL"), "{error}");
}
#[test]
fn inactive_status_mirrors_the_config_without_exposing_the_token() {
let config = LogForwardConfig {
enabled: true,
token: "secret-api-key".to_string(),
..LogForwardConfig::disabled()
};
let status = LogForwardStatus::inactive(&config);
assert!(status.enabled);
assert!(!status.active);
assert_eq!(status.endpoint, DEFAULT_ENDPOINT);
assert_eq!(status.min_level, LogLevel::Info);
assert_eq!(status.token_fingerprint, config.token_fingerprint());
let json = serde_json::to_string(&status).unwrap();
assert!(
!json.contains("secret-api-key"),
"status must never carry the token: {json}"
);
}
#[test]
fn inactive_status_of_a_disabled_config_has_no_fingerprint() {
let status = LogForwardStatus::inactive(&LogForwardConfig::disabled());
assert!(!status.enabled);
assert_eq!(status.token_fingerprint, None);
}
#[test]
fn skip_reason_points_at_the_flag_that_fixes_it() {
let skipped = SkippedNode::no_logging(3, "node3");
assert_eq!(skipped.node_id, 3);
assert_eq!(skipped.service, "node3");
assert!(skipped.reason.contains("--log-dir-path"));
}
#[test]
fn stats_start_at_zero() {
let stats = ForwardStats::default();
assert_eq!(stats.events_forwarded, 0);
assert_eq!(stats.batches_failed, 0);
assert_eq!(stats.last_success_unix, None);
assert_eq!(stats.last_error, None);
}
}