1use 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 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
84fn 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 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}