gitoxide-core 0.63.0

The library implementing all capabilities of the gitoxide CLI
use std::path::Path;

use gix::{
    Result,
    error::{ResultExt, message},
};
use gix_trace::forest::Tree;
use parking_lot::Mutex;
use rusqlite::params;

use crate::trace::Output;

pub fn subscriber(db_path: impl AsRef<Path>, trace: u8, output: Output) -> Result<tracing::Dispatch> {
    let con = Mutex::new(
        rusqlite::Connection::open(db_path).or_raise(|| message("Could not open the corpus trace database"))?,
    );
    crate::trace::subscriber(
        trace,
        output,
        Some(Box::new(move |tree| {
            let run_id = tree
                .span()
                .ok()
                .filter(|span| span.name() == "run")
                .and_then(|span| span.fields().iter().find(|field| field.key() == "run_id"))
                .and_then(|field| field.value().parse::<super::db::Id>().ok());
            if let Some(run_id) = run_id {
                let json =
                    serde_json::to_string_pretty(&tree_json(tree)).expect("serialization to string always works");
                con.lock()
                    .execute("UPDATE run SET spans_json = ?1 WHERE id = ?2", params![json, run_id])
                    .or_raise(|| message!("Could not store the trace for corpus run {run_id}"))?;
            }
            Ok(())
        })),
    )
}

fn tree_json(tree: &Tree) -> serde_json::Value {
    let fields = |fields: &[gix_trace::forest::tree::Field]| {
        fields
            .iter()
            .map(|field| (field.key().to_owned(), serde_json::Value::from(field.value())))
            .collect::<serde_json::Map<_, _>>()
    };
    match tree {
        Tree::Event(event) => serde_json::json!({
            "Event": {
                "level": event.level().as_str(),
                "fields": fields(event.fields()),
                "message": event.message(),
                "tag": event.tag().map(|tag| tag.to_string()),
            }
        }),
        Tree::Span(span) => serde_json::json!({
            "Span": {
                "level": span.level().as_str(),
                "fields": fields(span.fields()),
                "name": span.name(),
                "nanos_total": span.total_duration().as_nanos(),
                "nanos_nested": span.inner_duration().as_nanos(),
                "nodes": span.nodes().iter().map(tree_json).collect::<Vec<_>>(),
            }
        }),
    }
}

#[cfg(test)]
mod tests {
    use gix::error::TestResult;
    use std::path::Path;

    use crate::{corpus::db, trace::Output};

    #[test]
    fn requested_trace_mode_controls_deferred_format_and_level() -> TestResult<()> {
        let fixture = tempfile::tempdir()?;

        let forest_info = messages(fixture.path(), 1)?;
        assert_eq!(forest_info.len(), 2);
        assert!(
            forest_info.iter().all(|line| line.contains("INFO")),
            "forest INFO output retains only INFO lines"
        );
        assert!(
            forest_info.iter().all(|line| line.contains("\x1b[")),
            "forest output includes ANSI styling"
        );

        let forest_debug = messages(fixture.path(), 2)?;
        assert_eq!(forest_debug.len(), 3);
        assert!(
            forest_debug.iter().any(|line| line.contains("DEBUG")),
            "forest DEBUG output includes DEBUG lines"
        );

        let flat_debug = messages(fixture.path(), 3)?;
        assert_eq!(flat_debug.len(), 3);
        assert!(
            flat_debug.iter().any(|line| line.contains("DEBUG")),
            "flat DEBUG output includes DEBUG lines"
        );
        assert!(
            flat_debug.iter().all(|line| line.contains("\x1b[")),
            "flat output includes ANSI styling"
        );
        assert!(flat_debug.iter().any(|line| line.contains("close")));
        assert!(
            !flat_debug.iter().any(|line| line.contains("TRACE")),
            "flat DEBUG output excludes TRACE lines"
        );

        let flat_trace = messages(fixture.path(), 4)?;
        assert_eq!(flat_trace.len(), 4);
        assert!(
            flat_trace.iter().any(|line| line.contains("TRACE")),
            "flat TRACE output includes TRACE lines"
        );
        Ok(())
    }

    #[test]
    fn display_modes_do_not_filter_the_stored_trace() -> TestResult<()> {
        let fixture = tempfile::tempdir()?;
        for trace in 0..=4 {
            let db_path = fixture.path().join(format!("stored-{trace}.db"));
            let connection = db::create(&db_path)?;
            connection.execute("INSERT INTO run (insertion_time) VALUES (0)", [])?;
            let run_id = u32::try_from(connection.last_insert_rowid()).expect("test run id fits in u32");

            let output = Output::default();
            let dispatch = super::subscriber(&db_path, trace, output.clone())?;
            tracing::dispatcher::with_default(&dispatch, || {
                tracing::info_span!("run", run_id).in_scope(|| {
                    tracing::debug!("stored debug event\nprivate continuation");
                    tracing::trace!("stored trace event");
                    tracing::info!("visible info event");
                });
            });

            let stored: String =
                connection.query_row("SELECT spans_json FROM run WHERE id = ?1", [run_id], |row| row.get(0))?;
            for message in [
                "stored debug event",
                "private continuation",
                "stored trace event",
                "visible info event",
            ] {
                assert!(
                    stored.contains(message),
                    "display mode {trace} leaves storage complete: {message}"
                );
            }
            let output = output.lock().expect("trace output lock is not poisoned");
            let rendered = String::from_utf8_lossy(&output);
            assert_eq!(rendered.is_empty(), trace == 0, "disabled display remains silent");
            for (message, visible) in [
                ("visible info event", trace != 0),
                ("stored debug event", trace >= 2),
                ("private continuation", trace >= 2),
                ("stored trace event", trace == 4),
            ] {
                assert_eq!(
                    rendered.contains(message),
                    visible,
                    "display mode {trace} filters complete events: {message}"
                );
            }
        }
        Ok(())
    }

    #[test]
    fn delayed_run_spans_keep_their_ids_and_unrelated_roots_only_display() -> TestResult<()> {
        let fixture = tempfile::tempdir()?;
        let db_path = fixture.path().join("run-ids.db");
        let connection = db::create(&db_path)?;
        connection.execute("INSERT INTO run (insertion_time) VALUES (0), (0)", [])?;
        let second_run_id = u32::try_from(connection.last_insert_rowid()).expect("test run id fits in u32");
        let first_run_id = second_run_id - 1;
        let output = Output::default();
        let dispatch = super::subscriber(&db_path, 1, output.clone())?;
        {
            let _guard = tracing::dispatcher::set_default(&dispatch);
            let first = tracing::info_span!("run", run_id = first_run_id);
            first.in_scope(|| tracing::info!("first run"));
            tracing::info_span!("run", run_id = second_run_id).in_scope(|| tracing::info!("second run"));
            drop(first);
            tracing::info!("unrelated event");
            tracing::info_span!("unrelated span", run_id = second_run_id)
                .in_scope(|| tracing::info!("unrelated span event"));
            tracing::info_span!("run").in_scope(|| tracing::info!("run without id"));
        }

        for (run_id, message) in [(first_run_id, "first run"), (second_run_id, "second run")] {
            let stored: String =
                connection.query_row("SELECT spans_json FROM run WHERE id = ?1", [run_id], |row| row.get(0))?;
            let stored: serde_json::Value = serde_json::from_str(&stored)?;
            assert_eq!(
                stored["Span"]["fields"]["run_id"],
                run_id.to_string(),
                "completed run spans retain their own ID regardless of closure order"
            );
            assert_eq!(
                stored["Span"]["nodes"][0]["Event"]["message"], message,
                "unrelated roots do not overwrite a run's stored tree"
            );
        }
        let output = output.lock().expect("trace output lock is not poisoned");
        let rendered = String::from_utf8_lossy(&output);
        for message in [
            "first run",
            "second run",
            "unrelated event",
            "unrelated span event",
            "run without id",
        ] {
            assert!(
                rendered.contains(message),
                "all requested trace output remains visible: {message}"
            );
        }
        Ok(())
    }

    #[test]
    fn concurrent_run_roots_share_a_dispatch_without_mixing_storage_or_output() -> TestResult<()> {
        let fixture = tempfile::tempdir()?;
        for trace in [0, 1, 3, 4] {
            let db_path = fixture.path().join(format!("concurrent-{trace}.db"));
            let connection = db::create(&db_path)?;
            let mut run_ids = Vec::new();
            for _ in 0..4 {
                connection.execute("INSERT INTO run (insertion_time) VALUES (0)", [])?;
                run_ids.push(db::Id::try_from(connection.last_insert_rowid())?);
            }
            let output = Output::default();
            let dispatch = super::subscriber(&db_path, trace, output.clone())?;
            let barrier = std::sync::Barrier::new(run_ids.len());
            std::thread::scope(|scope| {
                for &run_id in &run_ids {
                    let (dispatch, barrier) = (&dispatch, &barrier);
                    scope.spawn(move || {
                        let _guard = tracing::dispatcher::set_default(dispatch);
                        tracing::info_span!("run", run_id).in_scope(|| {
                            barrier.wait();
                            tracing::info!("run event {run_id}\ncontinuation {run_id}");
                        });
                    });
                }
            });
            let output = output.lock().expect("trace output lock is not poisoned");
            let rendered = String::from_utf8_lossy(&output);
            let lines: Vec<_> = rendered.lines().collect();
            assert_eq!(
                rendered.is_empty(),
                trace == 0,
                "concurrent storage does not enable display"
            );
            for run_id in run_ids {
                let stored: String =
                    connection.query_row("SELECT spans_json FROM run WHERE id = ?1", [run_id], |row| row.get(0))?;
                let stored: serde_json::Value = serde_json::from_str(&stored)?;
                let event = format!("run event {run_id}");
                let continuation = format!("continuation {run_id}");
                assert_eq!(
                    stored["Span"]["fields"]["run_id"],
                    run_id.to_string(),
                    "concurrent roots keep their own IDs"
                );
                assert_eq!(
                    stored["Span"]["nodes"][0]["Event"]["message"],
                    format!("{event}\n{continuation}"),
                    "each stored tree contains only its own event",
                );
                if trace != 0 {
                    let index = lines
                        .iter()
                        .position(|line| line.contains(&event))
                        .expect("each event is displayed");
                    assert!(
                        lines.get(index + 1).is_some_and(|line| line.contains(&continuation)),
                        "each multiline event is written contiguously in display mode {trace}",
                    );
                }
            }
        }
        Ok(())
    }

    fn messages(root: &Path, trace: u8) -> TestResult<Vec<String>> {
        let db_path = root.join(format!("trace-{trace}.db"));
        drop(db::create(&db_path)?);
        let output = Output::default();
        let dispatch = super::subscriber(&db_path, trace, output.clone())?;
        {
            let _guard = tracing::dispatcher::set_default(&dispatch);
            tracing::info_span!("root").in_scope(|| {
                tracing::info!("info event");
                tracing::debug!("debug event");
                tracing::trace!("trace event");
            });
        }
        let output = output.lock().expect("trace output lock is not poisoned");
        Ok(String::from_utf8_lossy(&output)
            .lines()
            .map(ToOwned::to_owned)
            .collect())
    }
}

#[cfg(test)]
mod serialization_tests {
    use super::*;
    use gix::error::{ErrorExt, TestResult, message};
    use gix_trace::{
        ForestLayer,
        forest::{Tag, processor},
    };
    use tracing_subscriber::layer::SubscriberExt;

    #[test]
    fn serialization_preserves_the_corpus_json_shape() -> TestResult<()> {
        let (sender, receiver) = std::sync::mpsc::channel();
        let processor = processor::from_fn(move |tree| {
            sender
                .send(tree)
                .map_err(|err| processor::error(err.0, message("tree receiver dropped").raise()))
        });
        let layer = ForestLayer::new(processor, |event: &tracing::Event<'_>| {
            (event.metadata().target() == "tagged").then(|| {
                Tag::builder()
                    .prefix("request")
                    .level(*event.metadata().level())
                    .build()
            })
        });
        tracing::subscriber::with_default(tracing_subscriber::Registry::default().with(layer), || {
            let root = tracing::info_span!("root", initial = 1, deferred = tracing::field::Empty);
            root.record("initial", 2).record("deferred", "ready");
            root.in_scope(|| {
                tracing::info!(count = 2, "complete");
                let _child = tracing::debug_span!("child").entered();
                tracing::warn!(target: "tagged", "retrying");
            });
        });
        let tree = receiver.try_recv()?;
        let root = tree.span()?;
        let child = root.nodes()[1].span()?;
        let json = serde_json::to_string_pretty(&tree_json(&tree))?;
        assert_eq!(
            serde_json::from_str::<serde_json::Value>(&json)?,
            serde_json::json!({
                "Span": {
                    "level": "INFO",
                    "fields": {"initial": "2", "deferred": "\"ready\""},
                    "name": "root",
                    "nanos_total": root.total_duration().as_nanos(),
                    "nanos_nested": child.total_duration().as_nanos(),
                    "nodes": [
                        {"Event": {
                            "level": "INFO",
                            "fields": {"count": "2"},
                            "message": "complete",
                            "tag": null
                        }},
                        {"Span": {
                            "level": "DEBUG",
                            "fields": {},
                            "name": "child",
                            "nanos_total": child.total_duration().as_nanos(),
                            "nanos_nested": 0,
                            "nodes": [{"Event": {
                                "level": "WARN",
                                "fields": {},
                                "message": "retrying",
                                "tag": "request.warn"
                            }}]
                        }}
                    ]
                }
            }),
            "tree variants, final field values, tags, and nanosecond durations retain the stored schema"
        );
        Ok(())
    }
}