Skip to main content

datui_lib/
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, with the schema inferred from all of them, so a field first
9//! seen late is a column too. A pipe or a followed file is scanned as it grows
10//! ([`crate::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::model_files::MetaValue;
20use crate::text_formats::Detail;
21
22/// What datui does with journal JSON: see [`crate::readers`].
23pub(crate) const READER: crate::readers::Reader = crate::readers::Reader {
24    scan,
25    signatures: &[crate::readers::Signature {
26        says: |head, _| looks_like(head),
27        kind: crate::readers::Kind::Text,
28        trusted: crate::readers::Trusted {
29            // A listing does not parse text.
30            listing: false,
31            ..crate::readers::EVERYWHERE
32        },
33    }],
34    refines: &[crate::FileFormat::Jsonl, crate::FileFormat::Json],
35    python: Some(crate::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_modal::ExportFormat::Ndjson),
42    ..crate::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, every record's fields among the columns.
86fn read_all(paths: &[PathBuf]) -> Result<DataFrame> {
87    fn read<R: polars::io::mmap::MmapBytesReader>(reader: R) -> PolarsResult<DataFrame> {
88        JsonReader::new(reader)
89            .with_json_format(JsonFormat::JsonLines)
90            .infer_schema_len(None)
91            .finish()
92    }
93    let df = match paths {
94        [one] => read(std::fs::File::open(one)?)?,
95        many => {
96            // One read of them all, so the schema is every file's fields.
97            let mut bytes = Vec::new();
98            for path in many {
99                bytes.extend(std::fs::read(path)?);
100                if bytes.last().is_some_and(|&b| b != b'\n') {
101                    bytes.push(b'\n');
102                }
103            }
104            read(std::io::Cursor::new(bytes))?
105        }
106    };
107    Ok(df)
108}
109
110fn scan(input: crate::readers::ScanIn<'_>) -> Result<crate::scan::Scan> {
111    let path = input.paths[0].clone();
112    let failed = |e: &dyn std::fmt::Display| -> color_eyre::Report {
113        crate::error_display::FileError::new(&path, format!("not journal JSON: {e}")).into()
114    };
115    let raw = if input.options.follow {
116        crate::follow::scan_lines(&path, input.options, true, &mut input.report.read_python)?
117    } else {
118        read_all(input.paths).map_err(|e| failed(&e))?.lazy()
119    };
120    let schema = raw.clone().collect_schema().map_err(|e| failed(&e))?;
121    let bytes = came_as_bytes(&raw, &schema).map_err(|e| failed(&e))?;
122    let (lf, python) = derive(raw, &schema);
123    input.report.read_python.extend(python);
124    let detail = summary(&lf).map_err(|e| failed(&e))?;
125    let mut notes = Vec::new();
126    if bytes > 0 {
127        notes.push(crate::text_formats::note(
128            format!(
129                "{} stored as bytes {} shown as text, invalid UTF-8 as \u{fffd}",
130                crate::text_formats::count(bytes as u64, "message", "messages"),
131                crate::glyphs::get().middot
132            ),
133            "the journal".to_string(),
134        ));
135    }
136    input.report.opened = Some(Arc::new(crate::members::Opened {
137        detail: Some(Arc::new(detail)),
138        notes,
139        ..Default::default()
140    }));
141    // Each entry's place in the journal, for `#`. Read whole, so the index costs no
142    // pushdown; a followed journal is scanned, and goes without.
143    let lf = if input.options.follow {
144        lf
145    } else {
146        lf.with_row_index(crate::row_index::INDEX, None)
147    };
148    Ok(lf.into())
149}
150
151/// The journal's columns: `time` and `level` derived, `MESSAGE` as text, in the order
152/// a reader of logs looks for them, bookkeeping last. Also the same as Python method
153/// calls, for Copy as Python.
154pub(crate) fn derive(lf: LazyFrame, schema: &Schema) -> (LazyFrame, Vec<String>) {
155    let has = |name: &str| schema.contains(name);
156    let mut exprs = Vec::new();
157    let mut python = Vec::new();
158    if has(REALTIME) {
159        exprs.push(
160            col(REALTIME)
161                .cast(DataType::Int64)
162                .cast(DataType::Datetime(
163                    TimeUnit::Microseconds,
164                    Some(polars::prelude::TimeZone::UTC),
165                ))
166                .alias(TIME),
167        );
168        python.push(format!(
169            "pl.col(\"{REALTIME}\").cast(pl.Int64).cast(pl.Datetime(\"us\", \"UTC\")).alias(\"{TIME}\")"
170        ));
171    }
172    if has(PRIORITY) {
173        exprs.push(level(col(PRIORITY).cast(DataType::String)));
174        let names: Vec<String> = LEVELS.iter().map(|l| format!("\"{l}\"")).collect();
175        let numbers: Vec<String> = (0..LEVELS.len()).map(|i| format!("\"{i}\"")).collect();
176        python.push(format!(
177            "pl.col(\"{PRIORITY}\").cast(pl.String).replace_strict([{}], [{}], default=None, return_dtype=pl.Enum([{}])).alias(\"{LEVEL}\")",
178            numbers.join(", "),
179            names.join(", "),
180            names.join(", ")
181        ));
182    }
183    if let Some(dtype) = schema.get(MESSAGE)
184        && (dtype.is_string() || matches!(dtype, DataType::List(_)))
185    {
186        exprs.push(
187            col(MESSAGE)
188                .map(
189                    |c: Column| as_text(&c).map(|s| s.into_column()),
190                    |_, field| Ok(Field::new(field.name().clone(), DataType::String)),
191                )
192                .alias(MESSAGE),
193        );
194        python.push(format!(
195            "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)"
196        ));
197    }
198    let lf = if exprs.is_empty() {
199        lf
200    } else {
201        lf.with_columns(exprs)
202    };
203    let mut names: Vec<String> = schema.iter_names().map(|n| n.to_string()).collect();
204    for derived in [TIME, LEVEL] {
205        if (derived == TIME && has(REALTIME)) || (derived == LEVEL && has(PRIORITY)) {
206            names.retain(|n| n != derived);
207            names.push(derived.to_string());
208        }
209    }
210    let order = order(&names);
211    let quoted: Vec<String> = order.iter().map(|n| format!("\"{n}\"")).collect();
212    let mut calls = Vec::new();
213    if !python.is_empty() {
214        calls.push(format!(".with_columns({})", python.join(", ")));
215    }
216    calls.push(format!(".select([{}])", quoted.join(", ")));
217    (
218        lf.select(order.iter().map(|n| col(n.as_str())).collect::<Vec<_>>()),
219        calls,
220    )
221}
222
223/// `PRIORITY` 0 to 7 as its level, an enum in order of severity; anything else null.
224fn level(priority: Expr) -> Expr {
225    let levels = FrozenCategories::new(LEVELS).expect("eight distinct names");
226    let mut expr = lit(NULL).cast(DataType::String);
227    for (i, name) in LEVELS.iter().enumerate().rev() {
228        expr = when(priority.clone().eq(lit(i.to_string())))
229            .then(lit(*name))
230            .otherwise(expr);
231    }
232    expr.cast(DataType::from_frozen_categories(levels))
233        .alias(LEVEL)
234}
235
236/// The columns in reading order: time, level, the unit (or the identifier when there
237/// is no unit), the PID and the message, then the rest as they came, then bookkeeping.
238fn order(names: &[String]) -> Vec<String> {
239    let has = |n: &str| names.iter().any(|m| m == n);
240    let unit = if has(UNIT) { UNIT } else { IDENTIFIER };
241    let first: Vec<&str> = [TIME, LEVEL, unit, "_PID", MESSAGE]
242        .into_iter()
243        .filter(|n| has(n))
244        .collect();
245    let bookkeeping =
246        |n: &str| n.starts_with("__") || matches!(n, "_BOOT_ID" | "_MACHINE_ID" | "_RUNTIME_SCOPE");
247    let mut out: Vec<String> = first.iter().map(|n| n.to_string()).collect();
248    out.extend(
249        names
250            .iter()
251            .filter(|n| !first.contains(&n.as_str()) && !bookkeeping(n))
252            .cloned(),
253    );
254    out.extend(names.iter().filter(|n| bookkeeping(n)).cloned());
255    out
256}
257
258/// A message journalctl wrote as bytes: an array of numbers, which Polars reads into a
259/// text column as `[115, 116]`, or a list column when every message was one.
260fn bytes_of(text: &str) -> Option<Vec<u8>> {
261    let inner = text.strip_prefix('[')?.strip_suffix(']')?;
262    inner
263        .split(',')
264        .map(|b| b.trim().parse::<u8>().ok())
265        .collect()
266}
267
268/// `MESSAGE` as text: bytes decoded, lossily.
269fn as_text(column: &Column) -> PolarsResult<Series> {
270    let name = column.name().clone();
271    match column.dtype() {
272        DataType::List(_) => {
273            let lists = column.list()?;
274            let out: StringChunked = (0..lists.len())
275                .map(|i| {
276                    let item = lists.get_as_series(i)?;
277                    let item = item.cast(&DataType::UInt8).ok()?;
278                    let bytes: Vec<u8> = item.u8().ok()?.iter().map(|b| b.unwrap_or(0)).collect();
279                    Some(String::from_utf8_lossy(&bytes).into_owned())
280                })
281                .collect();
282            Ok(out.with_name(name).into_series())
283        }
284        _ => {
285            let text = column.str()?;
286            let out: StringChunked = text
287                .iter()
288                .map(|v| {
289                    v.map(|v| match bytes_of(v) {
290                        Some(bytes) => String::from_utf8_lossy(&bytes).into_owned(),
291                        None => v.to_string(),
292                    })
293                })
294                .collect();
295            Ok(out.with_name(name).into_series())
296        }
297    }
298}
299
300/// How many messages came as bytes rather than text.
301fn came_as_bytes(raw: &LazyFrame, schema: &Schema) -> PolarsResult<usize> {
302    let Some(dtype) = schema.get(MESSAGE) else {
303        return Ok(0);
304    };
305    let df = raw.clone().select([col(MESSAGE)]).collect()?;
306    let column = df.column(MESSAGE)?;
307    Ok(match dtype {
308        DataType::List(_) => column.len() - column.null_count(),
309        DataType::String => column
310            .str()?
311            .iter()
312            .filter(|v| v.is_some_and(|v| bytes_of(v).is_some()))
313            .count(),
314        _ => 0,
315    })
316}
317
318/// The Info panel's tab: the span, the entries, units, boots and hosts.
319pub(crate) fn summary(lf: &LazyFrame) -> PolarsResult<Detail> {
320    let schema = lf.clone().collect_schema()?;
321    let has = |n: &str| schema.contains(n);
322    let distinct = |n: &str| {
323        if has(n) {
324            col(n)
325                .drop_nulls()
326                .n_unique()
327                .cast(DataType::UInt64)
328                .alias(n)
329        } else {
330            lit(0u64).alias(n)
331        }
332    };
333    let mut aggs = vec![
334        len().cast(DataType::UInt64).alias("entries"),
335        distinct(UNIT),
336        distinct("_BOOT_ID"),
337        distinct("_HOSTNAME"),
338    ];
339    if has(TIME) {
340        aggs.push(col(TIME).min().alias("from"));
341        aggs.push(col(TIME).max().alias("to"));
342    }
343    let totals = lf.clone().select(aggs).collect()?;
344    let number = |n: &str| -> u64 {
345        totals
346            .column(n)
347            .ok()
348            .and_then(|c| c.u64().ok()?.get(0))
349            .unwrap_or(0)
350    };
351    let group = |n: u64| crate::numfmt::group_chrome(usize::try_from(n).unwrap_or(usize::MAX));
352    let mut lines = vec![format!("Entries: {}", group(number("entries")))];
353    if has(TIME) {
354        let at = |n: &str| {
355            totals
356                .column(n)
357                .ok()
358                .and_then(|c| c.get(0).ok())
359                .map(|v| v.to_string())
360                .filter(|v| v != "null")
361        };
362        if let (Some(from), Some(to)) = (at("from"), at("to")) {
363            lines.push(format!("From: {from}"));
364            lines.push(format!("To: {to}"));
365        }
366    }
367    lines.push(format!("Units: {}", group(number(UNIT))));
368    lines.push(format!("Boots: {}", group(number("_BOOT_ID"))));
369    lines.push(format!("Hosts: {}", group(number("_HOSTNAME"))));
370    let mut list = Vec::new();
371    let mut units = 0;
372    if has(UNIT) {
373        let counts = lf
374            .clone()
375            .filter(col(UNIT).is_not_null())
376            .group_by([col(UNIT)])
377            .agg([len().alias("n")])
378            .sort(
379                ["n"],
380                SortMultipleOptions::default().with_order_descending(true),
381            )
382            .collect()?;
383        units = counts.height();
384        let names = counts.column(UNIT)?.str()?.clone();
385        let n = counts.column("n")?.cast(&DataType::UInt64)?;
386        let n = n.u64()?;
387        list = names
388            .iter()
389            .zip(n.iter())
390            .map(|(name, n)| {
391                (
392                    name.unwrap_or_default().to_string(),
393                    MetaValue::Text(crate::text_formats::count(
394                        n.unwrap_or(0),
395                        "entry",
396                        "entries",
397                    )),
398                )
399            })
400            .collect();
401    }
402    let detail = Detail {
403        tab: crate::text_formats::tab(crate::FileFormat::Journal),
404        lines,
405        list_title: "Units",
406        list: crate::text_formats::capped_list(list.into_iter(), units),
407        ..Default::default()
408    };
409    Ok(detail)
410}
411
412/// Copy as Python: every record's fields, as the open read them.
413fn python_arguments(
414    call: &mut crate::python_script::Call<'_>,
415) -> Option<crate::python_script::Source> {
416    call.args.push("infer_schema_length=None".to_string());
417    None
418}
419
420#[cfg(test)]
421mod tests {
422    use super::*;
423
424    #[test]
425    fn journal_json_by_its_first_record() {
426        let entry = br#"{"__CURSOR":"s=1;i=1","__REALTIME_TIMESTAMP":"1","MESSAGE":"hi"}"#;
427        let mut head = entry.to_vec();
428        head.push(b'\n');
429        assert!(looks_like(&head));
430        // Cut before its end: by the keys it names.
431        assert!(looks_like(&entry[..50]));
432        assert!(!looks_like(b"{\"a\": 1}\n"));
433        assert!(!looks_like(b"__CURSOR __REALTIME_TIMESTAMP\n"));
434        let piped = crate::readers::sniff(&head, None, crate::readers::Asked::Pipe, |_| true);
435        assert_eq!(piped, Some(crate::FileFormat::Journal));
436    }
437
438    #[test]
439    fn columns_in_reading_order() {
440        let names: Vec<String> = [
441            "__CURSOR",
442            "_BOOT_ID",
443            "SYSLOG_IDENTIFIER",
444            "MESSAGE",
445            "_PID",
446            "_HOSTNAME",
447            "time",
448            "level",
449        ]
450        .map(String::from)
451        .to_vec();
452        assert_eq!(
453            order(&names),
454            [
455                "time",
456                "level",
457                "SYSLOG_IDENTIFIER",
458                "_PID",
459                "MESSAGE",
460                "_HOSTNAME",
461                "__CURSOR",
462                "_BOOT_ID"
463            ]
464        );
465    }
466
467    #[test]
468    fn bytes_are_numbers_in_brackets() {
469        assert_eq!(bytes_of("[104, 105]"), Some(b"hi".to_vec()));
470        assert_eq!(bytes_of("[104,105]"), Some(b"hi".to_vec()));
471        assert_eq!(bytes_of("[256]"), None);
472        assert_eq!(bytes_of("[]"), None);
473        assert_eq!(bytes_of("hi [1]"), None);
474    }
475
476    #[test]
477    fn the_info_tab_sums_the_journal() {
478        let df = df!(
479            REALTIME => ["1767225600000000", "1767225660000000", "1767225720000000"],
480            PRIORITY => ["6", "3", "9"],
481            UNIT => [Some("a.service"), Some("a.service"), None],
482            "_BOOT_ID" => ["b1", "b2", "b2"],
483            "_HOSTNAME" => ["h", "h", "h"],
484            MESSAGE => ["x", "[104, 105]", "y"],
485        )
486        .unwrap();
487        let lf = df.lazy();
488        let schema = lf.clone().collect_schema().unwrap();
489        assert_eq!(came_as_bytes(&lf, &schema).unwrap(), 1);
490        let (lf, python) = derive(lf, &schema);
491        assert!(python.last().unwrap().starts_with(".select("));
492        let out = lf.clone().collect().unwrap();
493        let level: Vec<Option<String>> = out
494            .column(LEVEL)
495            .unwrap()
496            .cast(&DataType::String)
497            .unwrap()
498            .str()
499            .unwrap()
500            .iter()
501            .map(|v| v.map(str::to_string))
502            .collect();
503        assert_eq!(level, [Some("info".into()), Some("err".into()), None]);
504        let detail = summary(&lf).unwrap();
505        assert_eq!(detail.tab, "Journal");
506        for line in ["Entries: 3", "Units: 1", "Boots: 2", "Hosts: 1"] {
507            assert!(detail.lines.iter().any(|l| l == line), "{:?}", detail.lines);
508        }
509        assert!(
510            detail
511                .lines
512                .iter()
513                .any(|l| l.starts_with("From: 2026-01-01 00:00:00"))
514        );
515        assert_eq!(detail.list.len(), 1);
516    }
517}