1use std::path::PathBuf;
14use std::sync::Arc;
15
16use color_eyre::Result;
17use polars::prelude::*;
18
19use crate::formats::model_files::MetaValue;
20use crate::formats::text_formats::Detail;
21
22pub(crate) const READER: crate::formats::readers::Reader = crate::formats::readers::Reader {
24 scan,
25 signatures: &[crate::formats::readers::Signature {
26 says: |head, _| looks_like(head),
27 kind: crate::formats::readers::Kind::Text,
28 trusted: crate::formats::readers::Trusted {
29 listing: false,
31 ..crate::formats::readers::EVERYWHERE
32 },
33 }],
34 refines: &[crate::FileFormat::Jsonl, crate::FileFormat::Json],
35 python: Some(crate::export::python_script::Python {
36 call: "pl.scan_ndjson",
37 eager: false,
38 glob_flag: false,
39 arguments: Some(python_arguments),
40 }),
41 export: Some(crate::export::export_modal::ExportFormat::Ndjson),
42 ..crate::formats::readers::BASE
43};
44
45pub const TIME: &str = "time";
46pub const LEVEL: &str = "level";
47const REALTIME: &str = "__REALTIME_TIMESTAMP";
48const PRIORITY: &str = "PRIORITY";
49const MESSAGE: &str = "MESSAGE";
50const UNIT: &str = "_SYSTEMD_UNIT";
51const IDENTIFIER: &str = "SYSLOG_IDENTIFIER";
52
53pub const LEVELS: [&str; 8] = [
56 "emerg", "alert", "crit", "err", "warning", "notice", "info", "debug",
57];
58
59pub fn looks_like(head: &[u8]) -> bool {
63 let text = head.strip_prefix(b"\xef\xbb\xbf").unwrap_or(head);
64 let start = text
65 .iter()
66 .position(|b| !b.is_ascii_whitespace())
67 .unwrap_or(text.len());
68 let text = &text[start..];
69 if text.first() != Some(&b'{') {
70 return false;
71 }
72 let names = |keys: &dyn Fn(&str) -> bool| keys("__CURSOR") && keys(REALTIME);
73 match memchr::memchr(b'\n', text) {
74 Some(end) => match serde_json::from_slice::<serde_json::Value>(&text[..end]) {
75 Ok(serde_json::Value::Object(object)) => names(&|k| object.contains_key(k)),
76 _ => false,
77 },
78 None => {
79 let line = String::from_utf8_lossy(text);
80 names(&|k| line.contains(&format!("\"{k}\":")))
81 }
82 }
83}
84
85fn read_all(paths: &[PathBuf], most: u64) -> Result<(DataFrame, u64)> {
88 use std::io::Read;
89 fn read(reader: impl polars::io::mmap::MmapBytesReader) -> PolarsResult<DataFrame> {
90 JsonReader::new(reader)
91 .with_json_format(JsonFormat::JsonLines)
92 .infer_schema_len(None)
93 .finish()
94 }
95 if let [path] = paths {
97 let file = std::fs::File::open(path)?;
98 if file.metadata()?.len() <= most {
99 return Ok((read(file)?, 0));
100 }
101 }
102 let mut bytes = Vec::new();
104 let mut left_out = 0;
105 for path in paths {
106 let file = std::fs::File::open(path)?;
107 let len = file.metadata()?.len();
108 let room = if left_out > 0 {
110 0
111 } else {
112 most.saturating_sub(bytes.len() as u64)
113 };
114 file.take(room).read_to_end(&mut bytes)?;
115 if len > room {
116 left_out += len - room;
117 let whole = memchr::memrchr(b'\n', &bytes).map_or(0, |at| at + 1);
119 left_out += (bytes.len() - whole) as u64;
120 bytes.truncate(whole);
121 }
122 if bytes.last().is_some_and(|&b| b != b'\n') {
123 bytes.push(b'\n');
124 }
125 }
126 Ok((read(std::io::Cursor::new(bytes))?, left_out))
127}
128
129fn scan(input: crate::formats::readers::ScanIn<'_>) -> Result<crate::loading::scan::Scan> {
130 let path = input.paths[0].clone();
131 let failed = |e: &dyn std::fmt::Display| -> color_eyre::Report {
132 crate::error_display::FileError::new(&path, format!("not journal JSON: {e}")).into()
133 };
134 let mut left_out = 0;
135 let raw = if input.options.follow {
136 crate::loading::follow::scan_lines(
137 &path,
138 input.options,
139 true,
140 &mut input.report.read_python,
141 )?
142 } else {
143 let most = crate::limits::get().journal_bytes.bytes();
144 let (df, past) = read_all(input.paths, most).map_err(|e| failed(&e))?;
145 left_out = past;
146 df.lazy()
147 };
148 let schema = raw.clone().collect_schema().map_err(|e| failed(&e))?;
149 let bytes = came_as_bytes(&raw, &schema).map_err(|e| failed(&e))?;
150 let (lf, python) = derive(raw, &schema);
151 input.report.read_python.extend(python);
152 let detail = summary(&lf).map_err(|e| failed(&e))?;
153 let mut notes = Vec::new();
154 if left_out > 0 {
155 let size = crate::numfmt::bytes;
156 notes.push(crate::formats::text_formats::note(
157 format!(
158 "{} left out: past the first {} {} limits.journal_bytes raises it",
159 size(left_out),
160 size(crate::limits::get().journal_bytes.bytes()),
161 crate::glyphs::get().middot
162 ),
163 "the journal".to_string(),
164 ));
165 }
166 if bytes > 0 {
167 notes.push(crate::formats::text_formats::note(
168 format!(
169 "{} stored as bytes {} shown as text, invalid UTF-8 as \u{fffd}",
170 crate::formats::text_formats::count(bytes as u64, "message", "messages"),
171 crate::glyphs::get().middot
172 ),
173 "the journal".to_string(),
174 ));
175 }
176 input.report.opened = Some(Arc::new(crate::formats::members::Opened {
177 detail: Some(Arc::new(detail)),
178 notes,
179 ..Default::default()
180 }));
181 let lf = if input.options.follow {
184 lf
185 } else {
186 lf.with_row_index(crate::formats::row_index::INDEX, None)
187 };
188 Ok(lf.into())
189}
190
191pub(crate) fn derive(lf: LazyFrame, schema: &Schema) -> (LazyFrame, Vec<String>) {
195 let has = |name: &str| schema.contains(name);
196 let mut exprs = Vec::new();
197 let mut python = Vec::new();
198 if has(REALTIME) {
199 exprs.push(
200 col(REALTIME)
201 .cast(DataType::Int64)
202 .cast(DataType::Datetime(
203 TimeUnit::Microseconds,
204 Some(polars::prelude::TimeZone::UTC),
205 ))
206 .alias(TIME),
207 );
208 python.push(format!(
209 "pl.col(\"{REALTIME}\").cast(pl.Int64).cast(pl.Datetime(\"us\", \"UTC\")).alias(\"{TIME}\")"
210 ));
211 }
212 if has(PRIORITY) {
213 exprs.push(level(col(PRIORITY).cast(DataType::String)));
214 let names: Vec<String> = LEVELS.iter().map(|l| format!("\"{l}\"")).collect();
215 let numbers: Vec<String> = (0..LEVELS.len()).map(|i| format!("\"{i}\"")).collect();
216 python.push(format!(
217 "pl.col(\"{PRIORITY}\").cast(pl.String).replace_strict([{}], [{}], default=None, return_dtype=pl.Enum([{}])).alias(\"{LEVEL}\")",
218 numbers.join(", "),
219 names.join(", "),
220 names.join(", ")
221 ));
222 }
223 if let Some(dtype) = schema.get(MESSAGE)
224 && (dtype.is_string() || matches!(dtype, DataType::List(_)))
225 {
226 exprs.push(
227 col(MESSAGE)
228 .map(
229 |c: Column| as_text(&c).map(|s| s.into_column()),
230 |_, field| Ok(Field::new(field.name().clone(), DataType::String)),
231 )
232 .alias(MESSAGE),
233 );
234 python.push(format!(
235 "pl.col(\"{MESSAGE}\").map_elements(lambda m: bytes(int(b) for b in m.strip(\"[]\").split(\",\")).decode(\"utf-8\", \"replace\") if m.startswith(\"[\") and m.endswith(\"]\") else m, return_dtype=pl.String)"
236 ));
237 }
238 let lf = if exprs.is_empty() {
239 lf
240 } else {
241 lf.with_columns(exprs)
242 };
243 let mut names: Vec<String> = schema.iter_names().map(|n| n.to_string()).collect();
244 for derived in [TIME, LEVEL] {
245 if (derived == TIME && has(REALTIME)) || (derived == LEVEL && has(PRIORITY)) {
246 names.retain(|n| n != derived);
247 names.push(derived.to_string());
248 }
249 }
250 let order = order(&names);
251 let quoted: Vec<String> = order.iter().map(|n| format!("\"{n}\"")).collect();
252 let mut calls = Vec::new();
253 if !python.is_empty() {
254 calls.push(format!(".with_columns({})", python.join(", ")));
255 }
256 calls.push(format!(".select([{}])", quoted.join(", ")));
257 (
258 lf.select(order.iter().map(|n| col(n.as_str())).collect::<Vec<_>>()),
259 calls,
260 )
261}
262
263fn level(priority: Expr) -> Expr {
265 let levels = FrozenCategories::new(LEVELS).expect("eight distinct names");
266 let mut expr = lit(NULL).cast(DataType::String);
267 for (i, name) in LEVELS.iter().enumerate().rev() {
268 expr = when(priority.clone().eq(lit(i.to_string())))
269 .then(lit(*name))
270 .otherwise(expr);
271 }
272 expr.cast(DataType::from_frozen_categories(levels))
273 .alias(LEVEL)
274}
275
276fn order(names: &[String]) -> Vec<String> {
279 let has = |n: &str| names.iter().any(|m| m == n);
280 let unit = if has(UNIT) { UNIT } else { IDENTIFIER };
281 let first: Vec<&str> = [TIME, LEVEL, unit, "_PID", MESSAGE]
282 .into_iter()
283 .filter(|n| has(n))
284 .collect();
285 let bookkeeping =
286 |n: &str| n.starts_with("__") || matches!(n, "_BOOT_ID" | "_MACHINE_ID" | "_RUNTIME_SCOPE");
287 let mut out: Vec<String> = first.iter().map(|n| n.to_string()).collect();
288 out.extend(
289 names
290 .iter()
291 .filter(|n| !first.contains(&n.as_str()) && !bookkeeping(n))
292 .cloned(),
293 );
294 out.extend(names.iter().filter(|n| bookkeeping(n)).cloned());
295 out
296}
297
298fn bytes_of(text: &str) -> Option<Vec<u8>> {
301 let inner = text.strip_prefix('[')?.strip_suffix(']')?;
302 inner
303 .split(',')
304 .map(|b| b.trim().parse::<u8>().ok())
305 .collect()
306}
307
308fn as_text(column: &Column) -> PolarsResult<Series> {
310 let name = column.name().clone();
311 match column.dtype() {
312 DataType::List(_) => {
313 let lists = column.list()?;
314 let out: StringChunked = (0..lists.len())
315 .map(|i| {
316 let item = lists.get_as_series(i)?;
317 let item = item.cast(&DataType::UInt8).ok()?;
318 let bytes: Vec<u8> = item.u8().ok()?.iter().map(|b| b.unwrap_or(0)).collect();
319 Some(String::from_utf8_lossy(&bytes).into_owned())
320 })
321 .collect();
322 Ok(out.with_name(name).into_series())
323 }
324 _ => {
325 let text = column.str()?;
326 let out: StringChunked = text
327 .iter()
328 .map(|v| {
329 v.map(|v| match bytes_of(v) {
330 Some(bytes) => String::from_utf8_lossy(&bytes).into_owned(),
331 None => v.to_string(),
332 })
333 })
334 .collect();
335 Ok(out.with_name(name).into_series())
336 }
337 }
338}
339
340fn came_as_bytes(raw: &LazyFrame, schema: &Schema) -> PolarsResult<usize> {
342 let Some(dtype) = schema.get(MESSAGE) else {
343 return Ok(0);
344 };
345 let df = raw.clone().select([col(MESSAGE)]).collect()?;
346 let column = df.column(MESSAGE)?;
347 Ok(match dtype {
348 DataType::List(_) => column.len() - column.null_count(),
349 DataType::String => column
350 .str()?
351 .iter()
352 .filter(|v| v.is_some_and(|v| bytes_of(v).is_some()))
353 .count(),
354 _ => 0,
355 })
356}
357
358pub(crate) fn summary(lf: &LazyFrame) -> PolarsResult<Detail> {
360 let schema = lf.clone().collect_schema()?;
361 let has = |n: &str| schema.contains(n);
362 let distinct = |n: &str| {
363 if has(n) {
364 col(n)
365 .drop_nulls()
366 .n_unique()
367 .cast(DataType::UInt64)
368 .alias(n)
369 } else {
370 lit(0u64).alias(n)
371 }
372 };
373 let mut aggs = vec![
374 len().cast(DataType::UInt64).alias("entries"),
375 distinct(UNIT),
376 distinct("_BOOT_ID"),
377 distinct("_HOSTNAME"),
378 ];
379 if has(TIME) {
380 aggs.push(col(TIME).min().alias("from"));
381 aggs.push(col(TIME).max().alias("to"));
382 }
383 let totals = lf.clone().select(aggs).collect()?;
384 let number = |n: &str| -> u64 {
385 totals
386 .column(n)
387 .ok()
388 .and_then(|c| c.u64().ok()?.get(0))
389 .unwrap_or(0)
390 };
391 let group = |n: u64| crate::numfmt::group_chrome(usize::try_from(n).unwrap_or(usize::MAX));
392 let mut lines = vec![format!("Entries: {}", group(number("entries")))];
393 if has(TIME) {
394 let at = |n: &str| {
395 totals
396 .column(n)
397 .ok()
398 .and_then(|c| c.get(0).ok())
399 .map(|v| v.to_string())
400 .filter(|v| v != "null")
401 };
402 if let (Some(from), Some(to)) = (at("from"), at("to")) {
403 lines.push(format!("From: {from}"));
404 lines.push(format!("To: {to}"));
405 }
406 }
407 lines.push(format!("Units: {}", group(number(UNIT))));
408 lines.push(format!("Boots: {}", group(number("_BOOT_ID"))));
409 lines.push(format!("Hosts: {}", group(number("_HOSTNAME"))));
410 let mut list = Vec::new();
411 let mut units = 0;
412 if has(UNIT) {
413 let counts = lf
414 .clone()
415 .filter(col(UNIT).is_not_null())
416 .group_by([col(UNIT)])
417 .agg([len().alias("n")])
418 .sort(
419 ["n"],
420 SortMultipleOptions::default().with_order_descending(true),
421 )
422 .collect()?;
423 units = counts.height();
424 let names = counts.column(UNIT)?.str()?.clone();
425 let n = counts.column("n")?.cast(&DataType::UInt64)?;
426 let n = n.u64()?;
427 list = names
428 .iter()
429 .zip(n.iter())
430 .map(|(name, n)| {
431 (
432 name.unwrap_or_default().to_string(),
433 MetaValue::Text(crate::formats::text_formats::count(
434 n.unwrap_or(0),
435 "entry",
436 "entries",
437 )),
438 )
439 })
440 .collect();
441 }
442 let detail = Detail {
443 tab: crate::formats::text_formats::tab(crate::FileFormat::Journal),
444 lines,
445 list_title: "Units",
446 list: crate::formats::text_formats::capped_list(list.into_iter(), units),
447 ..Default::default()
448 };
449 Ok(detail)
450}
451
452fn python_arguments(
454 call: &mut crate::export::python_script::Call<'_>,
455) -> Option<crate::export::python_script::Source> {
456 call.args.push("infer_schema_length=None".to_string());
457 None
458}
459
460#[cfg(test)]
461mod tests {
462 use super::*;
463
464 #[test]
467 fn a_read_stops_at_its_cap_on_a_whole_record() {
468 let dir = tempfile::tempdir().unwrap();
469 let one = dir.path().join("one.json");
470 let two = dir.path().join("two.json");
471 let record = |n: u32| format!("{{\"MESSAGE\":\"m{n}\"}}\n");
472 std::fs::write(&one, (0..3).map(record).collect::<String>()).unwrap();
473 std::fs::write(&two, record(9)).unwrap();
474 let line = record(0).len() as u64;
475 let (df, left_out) = read_all(&[one.clone(), two.clone()], line + 3).unwrap();
476 assert_eq!(df.height(), 1);
477 assert_eq!(left_out, 3 * line);
478 let (df, left_out) = read_all(&[one, two], u64::MAX).unwrap();
479 assert_eq!((df.height(), left_out), (4, 0));
480 }
481
482 #[test]
483 fn journal_json_by_its_first_record() {
484 let entry = br#"{"__CURSOR":"s=1;i=1","__REALTIME_TIMESTAMP":"1","MESSAGE":"hi"}"#;
485 let mut head = entry.to_vec();
486 head.push(b'\n');
487 assert!(looks_like(&head));
488 assert!(looks_like(&entry[..50]));
490 assert!(!looks_like(b"{\"a\": 1}\n"));
491 assert!(!looks_like(b"__CURSOR __REALTIME_TIMESTAMP\n"));
492 let piped = crate::formats::readers::sniff(
493 &head,
494 None,
495 crate::formats::readers::Asked::Pipe,
496 |_| true,
497 );
498 assert_eq!(piped, Some(crate::FileFormat::Journal));
499 }
500
501 #[test]
502 fn columns_in_reading_order() {
503 let names: Vec<String> = [
504 "__CURSOR",
505 "_BOOT_ID",
506 "SYSLOG_IDENTIFIER",
507 "MESSAGE",
508 "_PID",
509 "_HOSTNAME",
510 "time",
511 "level",
512 ]
513 .map(String::from)
514 .to_vec();
515 assert_eq!(
516 order(&names),
517 [
518 "time",
519 "level",
520 "SYSLOG_IDENTIFIER",
521 "_PID",
522 "MESSAGE",
523 "_HOSTNAME",
524 "__CURSOR",
525 "_BOOT_ID"
526 ]
527 );
528 }
529
530 #[test]
531 fn bytes_are_numbers_in_brackets() {
532 assert_eq!(bytes_of("[104, 105]"), Some(b"hi".to_vec()));
533 assert_eq!(bytes_of("[104,105]"), Some(b"hi".to_vec()));
534 assert_eq!(bytes_of("[256]"), None);
535 assert_eq!(bytes_of("[]"), None);
536 assert_eq!(bytes_of("hi [1]"), None);
537 }
538
539 #[test]
540 fn the_info_tab_sums_the_journal() {
541 let df = df!(
542 REALTIME => ["1767225600000000", "1767225660000000", "1767225720000000"],
543 PRIORITY => ["6", "3", "9"],
544 UNIT => [Some("a.service"), Some("a.service"), None],
545 "_BOOT_ID" => ["b1", "b2", "b2"],
546 "_HOSTNAME" => ["h", "h", "h"],
547 MESSAGE => ["x", "[104, 105]", "y"],
548 )
549 .unwrap();
550 let lf = df.lazy();
551 let schema = lf.clone().collect_schema().unwrap();
552 assert_eq!(came_as_bytes(&lf, &schema).unwrap(), 1);
553 let (lf, python) = derive(lf, &schema);
554 assert!(python.last().unwrap().starts_with(".select("));
555 let out = lf.clone().collect().unwrap();
556 let level: Vec<Option<String>> = out
557 .column(LEVEL)
558 .unwrap()
559 .cast(&DataType::String)
560 .unwrap()
561 .str()
562 .unwrap()
563 .iter()
564 .map(|v| v.map(str::to_string))
565 .collect();
566 assert_eq!(level, [Some("info".into()), Some("err".into()), None]);
567 let detail = summary(&lf).unwrap();
568 assert_eq!(detail.tab, "Journal");
569 for line in ["Entries: 3", "Units: 1", "Boots: 2", "Hosts: 1"] {
570 assert!(detail.lines.iter().any(|l| l == line), "{:?}", detail.lines);
571 }
572 assert!(
573 detail
574 .lines
575 .iter()
576 .any(|l| l.starts_with("From: 2026-01-01 00:00:00"))
577 );
578 assert_eq!(detail.list.len(), 1);
579 }
580}