Skip to main content

self_hosted_node/
logs.rs

1//! Ships this node's warnings to Loki the same way metrics go to the push
2//! gateway: the node pushes out, so a fleet on someone else's machine needs no
3//! agent installed beside it.
4//!
5//! Without this the fleet's logs stayed on the box. `{service="self-hosted-node"}`
6//! returned nothing for production, only for the staging node that happens to
7//! run next to Loki, so a slow-decision warning naming its game was written and
8//! never seen.
9
10use std::sync::atomic::{AtomicU64, Ordering};
11use std::sync::Arc;
12use std::time::{Duration, SystemTime, UNIX_EPOCH};
13
14use tokio::sync::mpsc;
15use tracing::field::{Field, Visit};
16use tracing::{Event, Level, Subscriber};
17use tracing_subscriber::layer::Context;
18use tracing_subscriber::Layer;
19
20const ENV_URL: &str = "SELF_HOSTED_NODE_LOGS_PUSH_URL";
21const ENV_USERNAME: &str = "SELF_HOSTED_NODE_LOGS_PUSH_USERNAME";
22const ENV_PASSWORD: &str = "SELF_HOSTED_NODE_LOGS_PUSH_PASSWORD";
23const ENV_LEVEL: &str = "SELF_HOSTED_NODE_LOGS_PUSH_LEVEL";
24const ENV_INSTANCE: &str = "SELF_HOSTED_NODE_INSTANCE";
25const ENV_METRICS_URL: &str = "SELF_HOSTED_NODE_METRICS_PUSH_URL";
26
27const FLUSH_INTERVAL: Duration = Duration::from_secs(10);
28const QUEUE_CAPACITY: usize = 2048;
29const MAX_BATCH: usize = 512;
30
31pub struct LokiLayer {
32    tx: mpsc::Sender<(u128, String)>,
33    level: Level,
34    dropped: Arc<AtomicU64>,
35}
36
37#[derive(Default)]
38struct MessageVisitor {
39    message: String,
40    fields: String,
41}
42
43impl Visit for MessageVisitor {
44    fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
45        if field.name() == "message" {
46            self.message = format!("{value:?}");
47        } else {
48            self.fields
49                .push_str(&format!(" {}={:?}", field.name(), value));
50        }
51    }
52}
53
54impl<S: Subscriber> Layer<S> for LokiLayer {
55    fn on_event(&self, event: &Event<'_>, _ctx: Context<'_, S>) {
56        let meta = event.metadata();
57        if *meta.level() > self.level {
58            return;
59        }
60        // The flusher reports its own failures through this same layer, so
61        // without this an unreachable Loki would feed itself for ever.
62        if meta.target().starts_with(module_path!()) {
63            return;
64        }
65        let mut visitor = MessageVisitor::default();
66        event.record(&mut visitor);
67        let line = format!(
68            "{} {} {}{}",
69            meta.level(),
70            meta.target(),
71            visitor.message,
72            visitor.fields
73        );
74        let nanos = SystemTime::now()
75            .duration_since(UNIX_EPOCH)
76            .map(|d| d.as_nanos())
77            .unwrap_or_default();
78        if self.tx.try_send((nanos, line)).is_err() {
79            self.dropped.fetch_add(1, Ordering::Relaxed);
80        }
81    }
82}
83
84/// The instance label the push gateway already uses, so a node's logs and its
85/// metrics carry the same name without configuring it twice.
86fn instance_label() -> String {
87    if let Ok(name) = std::env::var(ENV_INSTANCE) {
88        if !name.is_empty() {
89            return name;
90        }
91    }
92    std::env::var(ENV_METRICS_URL)
93        .ok()
94        .and_then(|url| {
95            url.rsplit_once("/instance/")
96                .map(|(_, name)| name.trim_end_matches('/').to_string())
97        })
98        .filter(|name| !name.is_empty())
99        .unwrap_or_else(|| "unknown".to_string())
100}
101
102pub fn layer_from_env() -> Option<LokiLayer> {
103    let url = std::env::var(ENV_URL).ok().filter(|v| !v.is_empty())?;
104    let username = std::env::var(ENV_USERNAME).ok().filter(|v| !v.is_empty());
105    let password = std::env::var(ENV_PASSWORD).ok().filter(|v| !v.is_empty());
106    let level = match std::env::var(ENV_LEVEL)
107        .unwrap_or_default()
108        .to_lowercase()
109        .as_str()
110    {
111        "trace" => Level::TRACE,
112        "debug" => Level::DEBUG,
113        "info" => Level::INFO,
114        "error" => Level::ERROR,
115        _ => Level::WARN,
116    };
117    let (tx, rx) = mpsc::channel(QUEUE_CAPACITY);
118    let dropped = Arc::new(AtomicU64::new(0));
119    tokio::spawn(run(
120        rx,
121        url,
122        username,
123        password,
124        instance_label(),
125        dropped.clone(),
126    ));
127    Some(LokiLayer { tx, level, dropped })
128}
129
130async fn run(
131    mut rx: mpsc::Receiver<(u128, String)>,
132    url: String,
133    username: Option<String>,
134    password: Option<String>,
135    instance: String,
136    dropped: Arc<AtomicU64>,
137) {
138    let client = reqwest::Client::new();
139    let mut batch: Vec<(u128, String)> = Vec::new();
140    let mut ticker = tokio::time::interval(FLUSH_INTERVAL);
141    let mut complained = false;
142    loop {
143        tokio::select! {
144            line = rx.recv() => match line {
145                Some(line) => {
146                    batch.push(line);
147                    if batch.len() < MAX_BATCH {
148                        continue;
149                    }
150                }
151                None => break,
152            },
153            _ = ticker.tick() => {}
154        }
155        if batch.is_empty() {
156            continue;
157        }
158        let lost = dropped.swap(0, Ordering::Relaxed);
159        if lost > 0 {
160            let nanos = SystemTime::now()
161                .duration_since(UNIX_EPOCH)
162                .map(|d| d.as_nanos())
163                .unwrap_or_default();
164            batch.push((
165                nanos,
166                format!("WARN {} dropped {lost} log lines", module_path!()),
167            ));
168        }
169        let values: Vec<[String; 2]> = batch
170            .drain(..)
171            .map(|(nanos, line)| [nanos.to_string(), line])
172            .collect();
173        let payload = serde_json::json!({
174            "streams": [{
175                "stream": { "service": "self-hosted-node", "instance": instance },
176                "values": values,
177            }]
178        });
179        let mut request = client.post(&url).json(&payload);
180        if let Some(user) = &username {
181            request = request.basic_auth(user, password.clone());
182        }
183        match request.send().await {
184            Ok(response) if response.status().is_success() => complained = false,
185            // Reported to stderr, never through `tracing`: this task is what
186            // drains the layer's queue.
187            Ok(response) => {
188                if !complained {
189                    complained = true;
190                    eprintln!("[logs] loki push rejected: {}", response.status());
191                }
192            }
193            Err(error) => {
194                if !complained {
195                    complained = true;
196                    eprintln!("[logs] loki push failed: {error}");
197                }
198            }
199        }
200    }
201}