Skip to main content

datui_lib/formats/
journal.rs

1//! The systemd journal as `journalctl -o json` writes it: NDJSON whose records carry
2//! `__CURSOR` and `__REALTIME_TIMESTAMP`.
3//!
4//! journalctl has already parsed the journal, so this is the NDJSON reader with a few
5//! expressions on top: `time` from `__REALTIME_TIMESTAMP`, `level` from `PRIORITY` in
6//! order of severity, `MESSAGE` as text where it came as bytes, and the columns put in
7//! the order a reader of logs looks for them. Every field stays. A file's records are
8//! read whole into memory, up to `limits.journal_bytes`, with the schema inferred from
9//! all of them, so a field first seen late is a column too. A pipe or a followed file is scanned as it grows
10//! (`crate::loading::follow::lines`); a pipe's fields first seen after the open join once it
11//! ends.
12
13use 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
22/// What datui does with journal JSON: see [`crate::formats::readers`].
23pub(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            // A listing does not parse text.
30            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
53/// The levels `PRIORITY` 0 to 7 names, most severe first: the order a sort and
54/// `level <= "err"` go by.
55pub const LEVELS: [&str; 8] = [
56    "emerg", "alert", "crit", "err", "warning", "notice", "info", "debug",
57];
58
59/// Whether `head` begins journal JSON: its first line is an object with `__CURSOR`
60/// and `__REALTIME_TIMESTAMP`. A first line longer than the head is judged on the keys
61/// it names so far.
62pub 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
85/// The records of `paths`, read whole up to `most` bytes in all, every record's fields
86/// among the columns; and the bytes left out past `most`, which end at a whole record.
87fn 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    // One file under the cap goes to the reader as it is, mapped rather than copied.
96    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    // One read of them all, so the schema is every file's fields.
103    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        // Once one file is cut short, the rest are left out whole.
109        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            // The record cut through is left out whole.
118            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    // Each entry's place in the journal, for `#`. Read whole, so the index costs no
182    // pushdown; a followed journal is scanned, and goes without.
183    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
191/// The journal's columns: `time` and `level` derived, `MESSAGE` as text, in the order
192/// a reader of logs looks for them, bookkeeping last. Also the same as Python method
193/// calls, for Copy as Python.
194pub(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
263/// `PRIORITY` 0 to 7 as its level, an enum in order of severity; anything else null.
264fn 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
276/// The columns in reading order: time, level, the unit (or the identifier when there
277/// is no unit), the PID and the message, then the rest as they came, then bookkeeping.
278fn 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
298/// A message journalctl wrote as bytes: an array of numbers, which Polars reads into a
299/// text column as `[115, 116]`, or a list column when every message was one.
300fn 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
308/// `MESSAGE` as text: bytes decoded, lossily.
309fn 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
340/// How many messages came as bytes rather than text.
341fn 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
358/// The Info panel's tab: the span, the entries, units, boots and hosts.
359pub(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
452/// Copy as Python: every record's fields, as the open read them.
453fn 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    /// Past `most` bytes the records are left out whole, and counted; a later file is
465    /// left out entirely.
466    #[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        // Cut before its end: by the keys it names.
489        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}