1use std::io::Write;
4
5use serde_json::Value;
6use tracing::field::{Field, Visit};
7use tracing_subscriber::fmt::MakeWriter;
8use tracing_subscriber::fmt::format::Writer;
9use tracing_subscriber::fmt::time::{FormatTime, SystemTime};
10use tracing_subscriber::layer::Context;
11
12use crate::sink::is_point_target;
13
14const RESERVED: [&str; 4] = ["ts", "level", "target", "cause"];
15
16pub fn json<W>(writer: W) -> Json<W>
29where
30 W: for<'w> MakeWriter<'w> + 'static,
31{
32 Json { writer }
33}
34
35pub struct Json<W> {
36 writer: W,
37}
38
39impl<S, W> tracing_subscriber::Layer<S> for Json<W>
40where
41 S: tracing::Subscriber,
42 W: for<'w> MakeWriter<'w> + 'static,
43{
44 fn on_event(&self, event: &tracing::Event<'_>, _: Context<'_, S>) {
45 let meta = event.metadata();
46
47 let mut ts = String::new();
48 if SystemTime.format_time(&mut Writer::new(&mut ts)).is_err() {
49 ts.clear();
50 }
51
52 let mut line = Line(vec![
53 ("ts", Value::from(ts)),
54 ("level", Value::from(meta.level().as_str())),
55 ("target", Value::from(meta.target())),
56 ]);
57 if !is_point_target(meta.target())
58 && let Some(cause) = crate::current()
59 {
60 line.0.push(("cause", Value::from(cause.get())));
61 }
62
63 event.record(&mut line);
64
65 let Ok(bytes) = line.into_bytes() else {
66 return;
67 };
68
69 let _ = self.writer.make_writer_for(meta).write_all(&bytes);
70 }
71}
72
73struct Line(Vec<(&'static str, Value)>);
74
75impl Line {
76 fn put(&mut self, field: &Field, value: Value) {
77 let name = field.name();
78 if RESERVED.contains(&name) || self.0.iter().any(|(taken, _)| *taken == name) {
79 return;
80 }
81
82 self.0.push((name, value));
83 }
84
85 fn into_bytes(self) -> serde_json::Result<Vec<u8>> {
87 let mut bytes = vec![b'{'];
88 for (at, (name, value)) in self.0.iter().enumerate() {
89 if at > 0 {
90 bytes.push(b',');
91 }
92 serde_json::to_writer(&mut bytes, name)?;
93 bytes.push(b':');
94 serde_json::to_writer(&mut bytes, value)?;
95 }
96 bytes.extend_from_slice(b"}\n");
97
98 Ok(bytes)
99 }
100}
101
102impl Visit for Line {
103 fn record_f64(&mut self, field: &Field, value: f64) {
104 self.put(field, Value::from(value));
105 }
106
107 fn record_i64(&mut self, field: &Field, value: i64) {
108 self.put(field, Value::from(value));
109 }
110
111 fn record_u64(&mut self, field: &Field, value: u64) {
112 self.put(field, Value::from(value));
113 }
114
115 fn record_i128(&mut self, field: &Field, value: i128) {
116 let value = i64::try_from(value).map_or_else(|_| Value::from(value.to_string()), Value::from);
117 self.put(field, value);
118 }
119
120 fn record_u128(&mut self, field: &Field, value: u128) {
121 let value = u64::try_from(value).map_or_else(|_| Value::from(value.to_string()), Value::from);
122 self.put(field, value);
123 }
124
125 fn record_bool(&mut self, field: &Field, value: bool) {
126 self.put(field, Value::from(value));
127 }
128
129 fn record_str(&mut self, field: &Field, value: &str) {
130 self.put(field, Value::from(value));
131 }
132
133 fn record_error(&mut self, field: &Field, value: &(dyn std::error::Error + 'static)) {
134 self.put(field, Value::from(value.to_string()));
135 }
136
137 fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
138 self.put(field, Value::from(format!("{value:?}")));
139 }
140}
141
142#[cfg(test)]
143mod tests {
144 use std::sync::{Arc, Mutex};
145
146 use serde_json::json as value;
147 use tracing_subscriber::layer::SubscriberExt;
148
149 use crate::{Point, mark, resume};
150
151 #[derive(Clone, Default)]
152 struct Buffer(Arc<Mutex<Vec<u8>>>);
153
154 impl std::io::Write for Buffer {
155 fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
156 self.0.lock().unwrap().extend_from_slice(bytes);
157 Ok(bytes.len())
158 }
159
160 fn flush(&mut self) -> std::io::Result<()> {
161 Ok(())
162 }
163 }
164
165 impl<'w> tracing_subscriber::fmt::MakeWriter<'w> for Buffer {
166 type Writer = Buffer;
167
168 fn make_writer(&'w self) -> Buffer {
169 self.clone()
170 }
171 }
172
173 fn lines(buffer: &Buffer) -> Vec<(String, serde_json::Value)> {
174 let bytes = buffer.0.lock().unwrap();
175 String::from_utf8(bytes.clone())
176 .unwrap()
177 .lines()
178 .map(|line| (line.to_string(), serde_json::from_str(line).unwrap()))
179 .collect()
180 }
181
182 #[test]
183 fn an_event_is_a_line_of_json_with_its_fields_typed_and_its_cause() {
184 let buffer = Buffer::default();
185 let subscriber = tracing_subscriber::registry().with(super::json(buffer.clone()));
186
187 tracing::subscriber::with_default(subscriber, || {
188 let action = mark(|| Point::Action { message: "Kill" });
189 let _resumed = resume(Some(action));
190 tracing::info!(target: "app", pid = 42, alive = false, name = %"init", "process killed");
191 });
192
193 let lines = lines(&buffer);
194 assert_eq!(lines.len(), 2, "{lines:?}");
195
196 let (_, point) = &lines[0];
197 assert_eq!(point["target"], "guinea::action");
198 assert_eq!(point["action"], "Kill");
199 assert_eq!(point.get("cause"), None, "a point has its parent already");
200
201 let (line, event) = &lines[1];
202 let order: Vec<usize> = ["ts", "level", "target", "cause", "message", "pid", "alive", "name"]
203 .iter()
204 .map(|key| line.find(&format!("\"{key}\":")).expect(key))
205 .collect();
206 assert!(order.is_sorted(), "{line}");
207
208 assert_eq!(event["level"], "INFO");
209 assert_eq!(event["cause"], point["id"]);
210 assert_eq!(event["message"], "process killed");
211 assert_eq!(event["pid"], value!(42));
212 assert_eq!(event["alive"], value!(false));
213 assert_eq!(event["name"], "init");
214 }
215}