Skip to main content

cli/
export.rs

1//! Export writer for `mushroomdb export`.
2//!
3//! Supports two formats:
4//!   - `jsonl`: nodes.jsonl, edges.jsonl, rules.jsonl (always deterministic)
5//!   - `parquet`: nodes.parquet, edges.parquet, rules.parquet (byte-identical
6//!     output is NOT guaranteed across parquet-rs versions; use JSONL when
7//!     byte-stable reproducibility is required)
8//!
9//! Two runs on the same store state produce byte-identical JSONL output (sorted
10//! by key / (edge_type, src, dst) / name — callers must pre-sort the slices).
11//!
12//! # Float precision loss
13//!
14//! `Value::Float` values that are NaN or ±Inf are not representable in JSON.
15//! They are serialised as JSON `null` rather than causing the export to fail.
16//! This is a lossy mapping; the original value is irrecoverable from the
17//! export.  Stores that need exact float fidelity should use the binary backup
18//! command instead of export.
19
20use core_api::{ExportEdge, NodeInfo, RuleDef, Value};
21use std::io::Write as _;
22use std::path::Path;
23use std::sync::Arc;
24
25use arrow_array::{ArrayRef, BooleanArray, RecordBatch, StringArray};
26use arrow_schema::{DataType, Field, Schema};
27use parquet::arrow::ArrowWriter;
28use parquet::file::properties::WriterProperties;
29
30/// Accepted export formats.
31#[derive(Debug, Clone, PartialEq, Eq)]
32pub enum ExportFormat {
33    Jsonl,
34    Parquet,
35}
36
37impl ExportFormat {
38    pub fn parse(s: &str) -> Option<Self> {
39        match s {
40            "jsonl" => Some(Self::Jsonl),
41            "parquet" => Some(Self::Parquet),
42            _ => None,
43        }
44    }
45
46    pub fn name(&self) -> &'static str {
47        match self {
48            Self::Jsonl => "jsonl",
49            Self::Parquet => "parquet",
50        }
51    }
52}
53
54/// Serialize a `Value` to serde_json `Value`.
55///
56/// NaN and ±Inf floats are not representable in JSON and are mapped to `null`
57/// (see module-level docs for the loss semantics).
58fn value_to_json(v: &Value) -> serde_json::Value {
59    match v {
60        Value::Int(i) => serde_json::Value::Number((*i).into()),
61        // serde_json::Number::from_f64 returns None for NaN/±Inf; map those to null.
62        Value::Float(f) => serde_json::Number::from_f64(*f)
63            .map(serde_json::Value::Number)
64            .unwrap_or(serde_json::Value::Null),
65        Value::Str(s) => serde_json::Value::String(s.clone()),
66        Value::Bool(b) => serde_json::Value::Bool(*b),
67        Value::List(xs) => serde_json::Value::Array(xs.iter().map(value_to_json).collect()),
68        Value::Map(m) => serde_json::Value::Object(
69            m.iter()
70                .map(|(k, v)| (k.clone(), value_to_json(v)))
71                .collect(),
72        ),
73    }
74}
75
76/// Write JSONL files to `dest`: nodes.jsonl, edges.jsonl, rules.jsonl.
77///
78/// `nodes` must be sorted by key, `edges` by (edge_type, src, dst), `rules` by name.
79/// Deterministic: same sorted input → byte-identical output.
80pub fn write_jsonl(
81    nodes: &[NodeInfo],
82    edges: &[ExportEdge],
83    rules: &[RuleDef],
84    dest: &Path,
85) -> Result<(), crate::CliError> {
86    std::fs::create_dir_all(dest)?;
87
88    // ── nodes.jsonl ────────────────────────────────────────────────────────
89    {
90        let mut f = std::fs::File::create(dest.join("nodes.jsonl"))?;
91        for node in nodes {
92            let mut obj = serde_json::Map::new();
93            obj.insert("key".into(), serde_json::Value::String(node.key.clone()));
94            obj.insert(
95                "label".into(),
96                serde_json::Value::String(node.label.clone()),
97            );
98            for (k, v) in &node.props {
99                obj.insert(k.clone(), value_to_json(v));
100            }
101            let line = serde_json::to_string(&serde_json::Value::Object(obj))
102                .map_err(|e| crate::CliError(format!("json encode error: {e}")))?;
103            writeln!(f, "{line}")?;
104        }
105    }
106
107    // ── edges.jsonl ────────────────────────────────────────────────────────
108    {
109        let mut f = std::fs::File::create(dest.join("edges.jsonl"))?;
110        for edge in edges {
111            let obj = serde_json::json!({
112                "edge_type": edge.edge_type,
113                "src": edge.src,
114                "dst": edge.dst,
115                "derived": edge.derived,
116                "rule": edge.rule,
117            });
118            let line = serde_json::to_string(&obj)
119                .map_err(|e| crate::CliError(format!("json encode error: {e}")))?;
120            writeln!(f, "{line}")?;
121        }
122    }
123
124    // ── rules.jsonl ────────────────────────────────────────────────────────
125    {
126        let mut f = std::fs::File::create(dest.join("rules.jsonl"))?;
127        for rule in rules {
128            let line = serde_json::to_string(rule)
129                .map_err(|e| crate::CliError(format!("json encode rule: {e}")))?;
130            writeln!(f, "{line}")?;
131        }
132    }
133
134    Ok(())
135}
136
137/// Write parquet files to `dest`: nodes.parquet, edges.parquet, rules.parquet.
138///
139/// Schema for each file:
140/// - nodes.parquet: key (Utf8), label (Utf8), props (Utf8, JSON)
141/// - edges.parquet: edge_type (Utf8), src (Utf8), dst (Utf8), derived (Boolean), rule (Utf8, nullable)
142/// - rules.parquet: name (Utf8), definition (Utf8, JSON)
143///
144/// `nodes` must be sorted by key, `edges` by (edge_type, src, dst), `rules` by name.
145///
146/// Compression: Snappy (parquet-rs default).  This is a format detail, not a
147/// stability contract — the column encoding and page layout may differ across
148/// parquet-rs versions.  Use JSONL export when byte-stable reproducibility is
149/// required.
150pub fn write_parquet(
151    nodes: &[NodeInfo],
152    edges: &[ExportEdge],
153    rules: &[RuleDef],
154    dest: &Path,
155) -> Result<(), crate::CliError> {
156    std::fs::create_dir_all(dest)?;
157    let props = WriterProperties::builder().build();
158
159    // ── nodes.parquet ──────────────────────────────────────────────────────
160    {
161        let schema = Arc::new(Schema::new(vec![
162            Field::new("key", DataType::Utf8, false),
163            Field::new("label", DataType::Utf8, false),
164            Field::new("props", DataType::Utf8, false),
165        ]));
166        let mut keys = Vec::new();
167        let mut labels = Vec::new();
168        let mut props_json = Vec::new();
169        for node in nodes {
170            keys.push(node.key.clone());
171            labels.push(node.label.clone());
172            let p: serde_json::Map<String, serde_json::Value> = node
173                .props
174                .iter()
175                .map(|(k, v)| (k.clone(), value_to_json(v)))
176                .collect();
177            props_json.push(
178                serde_json::to_string(&serde_json::Value::Object(p))
179                    .map_err(|e| crate::CliError(format!("json encode: {e}")))?,
180            );
181        }
182        let batch = RecordBatch::try_new(
183            schema.clone(),
184            vec![
185                Arc::new(StringArray::from(keys)) as ArrayRef,
186                Arc::new(StringArray::from(labels)) as ArrayRef,
187                Arc::new(StringArray::from(props_json)) as ArrayRef,
188            ],
189        )
190        .map_err(|e| crate::CliError(format!("arrow error: {e}")))?;
191
192        let file = std::fs::File::create(dest.join("nodes.parquet"))?;
193        let mut writer = ArrowWriter::try_new(file, schema, Some(props.clone()))
194            .map_err(|e| crate::CliError(format!("parquet writer: {e}")))?;
195        writer
196            .write(&batch)
197            .map_err(|e| crate::CliError(format!("parquet write: {e}")))?;
198        writer
199            .close()
200            .map_err(|e| crate::CliError(format!("parquet close: {e}")))?;
201    }
202
203    // ── edges.parquet ──────────────────────────────────────────────────────
204    {
205        let schema = Arc::new(Schema::new(vec![
206            Field::new("edge_type", DataType::Utf8, false),
207            Field::new("src", DataType::Utf8, false),
208            Field::new("dst", DataType::Utf8, false),
209            Field::new("derived", DataType::Boolean, false),
210            Field::new("rule", DataType::Utf8, true),
211        ]));
212        let mut etypes = Vec::new();
213        let mut srcs = Vec::new();
214        let mut dsts = Vec::new();
215        let mut deriveds = Vec::new();
216        let mut rule_names: Vec<Option<String>> = Vec::new();
217        for edge in edges {
218            etypes.push(edge.edge_type.clone());
219            srcs.push(edge.src.clone());
220            dsts.push(edge.dst.clone());
221            deriveds.push(edge.derived);
222            rule_names.push(edge.rule.clone());
223        }
224        let batch = RecordBatch::try_new(
225            schema.clone(),
226            vec![
227                Arc::new(StringArray::from(etypes)) as ArrayRef,
228                Arc::new(StringArray::from(srcs)) as ArrayRef,
229                Arc::new(StringArray::from(dsts)) as ArrayRef,
230                Arc::new(BooleanArray::from(deriveds)) as ArrayRef,
231                Arc::new(StringArray::from(rule_names)) as ArrayRef,
232            ],
233        )
234        .map_err(|e| crate::CliError(format!("arrow error: {e}")))?;
235
236        let file = std::fs::File::create(dest.join("edges.parquet"))?;
237        let mut writer = ArrowWriter::try_new(file, schema, Some(props.clone()))
238            .map_err(|e| crate::CliError(format!("parquet writer: {e}")))?;
239        writer
240            .write(&batch)
241            .map_err(|e| crate::CliError(format!("parquet write: {e}")))?;
242        writer
243            .close()
244            .map_err(|e| crate::CliError(format!("parquet close: {e}")))?;
245    }
246
247    // ── rules.parquet ──────────────────────────────────────────────────────
248    {
249        let schema = Arc::new(Schema::new(vec![
250            Field::new("name", DataType::Utf8, false),
251            Field::new("definition", DataType::Utf8, false),
252        ]));
253        let mut names = Vec::new();
254        let mut definitions = Vec::new();
255        for rule in rules {
256            names.push(rule.name.clone());
257            definitions.push(
258                serde_json::to_string(rule)
259                    .map_err(|e| crate::CliError(format!("json encode: {e}")))?,
260            );
261        }
262        let batch = RecordBatch::try_new(
263            schema.clone(),
264            vec![
265                Arc::new(StringArray::from(names)) as ArrayRef,
266                Arc::new(StringArray::from(definitions)) as ArrayRef,
267            ],
268        )
269        .map_err(|e| crate::CliError(format!("arrow error: {e}")))?;
270
271        let file = std::fs::File::create(dest.join("rules.parquet"))?;
272        let mut writer = ArrowWriter::try_new(file, schema, Some(props))
273            .map_err(|e| crate::CliError(format!("parquet writer: {e}")))?;
274        writer
275            .write(&batch)
276            .map_err(|e| crate::CliError(format!("parquet write: {e}")))?;
277        writer
278            .close()
279            .map_err(|e| crate::CliError(format!("parquet close: {e}")))?;
280    }
281
282    Ok(())
283}