Skip to main content

faucet_core/file_format/
mod.rs

1//! One file-format vocabulary for every file and object-store connector (#604).
2//!
3//! S3, GCS, Azure Blob and SFTP all move *files*, and the files they move are
4//! CSV exports, JSON dumps, XML feeds and Excel workbooks as often as they are
5//! JSON Lines. Before this module each connector carried its own small format
6//! enum — `JsonLines | JsonArray | RawText | Parquet` on the source side, and
7//! only `JsonLines | Parquet` on the sink side — so what you could read depended
8//! on which bucket it was in, and what you could write was a strict subset of
9//! what you could read. The gap was filled by pre- and post-processing outside
10//! the pipeline, which is exactly the work a movement engine exists to absorb.
11//!
12//! The parsers themselves are not new: the REST source has decoded
13//! `json | csv | xlsx | xml` since #497/#515. This module is where they move so
14//! that every connector — including a third-party one, which depends only on
15//! `faucet-core` — gets the same set, with the same option names and the same
16//! edge-case behaviour.
17//!
18//! # Shape
19//!
20//! ```yaml
21//! source:
22//!   type: s3
23//!   config:
24//!     bucket: exports
25//!     prefix: daily/
26//!     format: csv                  # FileFormat
27//!     compression: auto            # faucet_core::compression
28//!     csv: { has_headers: true, delimiter: "," }
29//! sink:
30//!   type: s3
31//!   config:
32//!     bucket: reports
33//!     format: xlsx
34//!     compression: gzip
35//! ```
36//!
37//! # What is deliberately not here
38//!
39//! **Parquet.** It is columnar, self-describing, and already has a dedicated
40//! Arrow read/write path per connector that must not be routed through a
41//! `Vec<Value>`. [`FileFormat::Parquet`] therefore names the format for config
42//! purposes but [`decode`]/[`encode`] refuse it, so a caller cannot silently
43//! lose the columnar fast path by going through the generic helper.
44//!
45//! **Compression.** [`crate::compression`] already owns codec selection,
46//! extension-based `auto` resolution, and the mismatch warning. Format and
47//! codec compose; neither needs to know about the other.
48
49use crate::error::FaucetError;
50use schemars::JsonSchema;
51use serde::{Deserialize, Serialize};
52use serde_json::Value;
53
54#[cfg(feature = "file-format-avro")]
55pub mod avro;
56pub mod container;
57#[cfg(feature = "file-format-csv")]
58pub mod csv;
59#[cfg(feature = "file-format-excel")]
60pub mod excel;
61pub mod json;
62#[cfg(feature = "file-format-orc")]
63pub mod orc;
64pub mod parquet_io;
65#[cfg(feature = "file-format-xml")]
66pub mod xml;
67
68pub use container::{ContainerDecoder, FileInput};
69
70/// How the bytes of one object are laid out.
71///
72/// The variants are the union of what the file connectors could previously read
73/// *or* write, so adopting this enum never removes a format from a connector
74/// that had it.
75#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
76#[serde(rename_all = "snake_case")]
77pub enum FileFormat {
78    /// One JSON value per line (NDJSON). The default everywhere, because it is
79    /// the only format that both streams and round-trips arbitrary JSON.
80    #[default]
81    JsonLines,
82    /// A single JSON array holding every record.
83    JsonArray,
84    /// Delimited text. See [`CsvOptions`].
85    Csv,
86    /// XML. See [`XmlOptions`].
87    Xml,
88    /// An Excel workbook (`.xlsx`). See [`ExcelOptions`].
89    Xlsx,
90    /// Apache Parquet — named for config, handled by each connector's own
91    /// Arrow path, never by [`decode`]/[`encode`].
92    Parquet,
93    /// Unparsed text: one record per object, `{"text": "<whole body>"}`.
94    RawText,
95    /// Apache Avro Object Container File. See [`AvroOptions`]. Read and
96    /// write; decodes to records or, with `arrow`, straight to Arrow batches.
97    Avro,
98    /// Apache ORC. See [`OrcOptions`]. **Read-only** — see the `orc` module
99    /// for why there is no writer.
100    Orc,
101}
102
103impl FileFormat {
104    /// The conventional file extension, used to name objects a sink creates
105    /// when the user did not pick one.
106    pub fn extension(self) -> &'static str {
107        match self {
108            Self::JsonLines => ".jsonl",
109            Self::JsonArray => ".json",
110            Self::Csv => ".csv",
111            Self::Xml => ".xml",
112            Self::Xlsx => ".xlsx",
113            Self::Parquet => ".parquet",
114            Self::RawText => ".txt",
115            Self::Avro => ".avro",
116            Self::Orc => ".orc",
117        }
118    }
119
120    /// Whether the whole object must be in memory to produce the first record.
121    ///
122    /// Callers use this to decide between a streaming read and a buffered one.
123    /// It is a property of the format, not of the source: a JSON array has no
124    /// record boundary until it is parsed, a zip-container workbook has its
125    /// directory at the end, and an XML document is a tree.
126    pub fn requires_whole_object(self) -> bool {
127        matches!(self, Self::JsonArray | Self::Xml | Self::Xlsx | Self::Orc)
128    }
129
130    /// Whether the format is a self-describing binary container decoded by
131    /// [`ContainerDecoder`] (Avro, ORC) rather than by [`decode`] per object.
132    pub fn is_container(self) -> bool {
133        matches!(self, Self::Avro | Self::Orc)
134    }
135
136    /// Whether a sink can write this format through [`encode`].
137    pub fn is_writable(self) -> bool {
138        !matches!(self, Self::Orc | Self::Parquet)
139    }
140
141    /// The format a file name's extension names, looking through a trailing
142    /// compression suffix (`.gz`, `.gzip`, `.zst`, `.zstd`), so
143    /// `export.csv.gz` is CSV. `None` for an extension no format claims.
144    ///
145    /// `.json` is a JSON array (a lone object is one record); `.jsonl` and
146    /// `.ndjson` are JSON Lines; `.tsv` is not claimed, because its delimiter
147    /// is an option rather than a format.
148    pub fn from_path(path: &str) -> Option<Self> {
149        let name = path
150            .rsplit(['/', '\\'])
151            .next()
152            .unwrap_or(path)
153            .to_ascii_lowercase();
154        let stem = ["gz", "gzip", "zst", "zstd"]
155            .iter()
156            .find_map(|c| name.strip_suffix(&format!(".{c}")))
157            .unwrap_or(&name);
158        let ext = stem.rsplit_once('.')?.1;
159        Some(match ext {
160            "jsonl" | "ndjson" => Self::JsonLines,
161            "json" => Self::JsonArray,
162            "csv" => Self::Csv,
163            "xml" => Self::Xml,
164            "xlsx" => Self::Xlsx,
165            "parquet" => Self::Parquet,
166            "txt" => Self::RawText,
167            "avro" => Self::Avro,
168            "orc" => Self::Orc,
169            _ => return None,
170        })
171    }
172
173    /// The name used in config and error messages.
174    pub fn as_str(self) -> &'static str {
175        match self {
176            Self::JsonLines => "json_lines",
177            Self::JsonArray => "json_array",
178            Self::Csv => "csv",
179            Self::Xml => "xml",
180            Self::Xlsx => "xlsx",
181            Self::Parquet => "parquet",
182            Self::RawText => "raw_text",
183            Self::Avro => "avro",
184            Self::Orc => "orc",
185        }
186    }
187}
188
189fn default_true() -> bool {
190    true
191}
192
193fn default_delimiter() -> String {
194    ",".into()
195}
196
197/// CSV dialect.
198#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
199#[serde(deny_unknown_fields)]
200#[schemars(extend("x-faucet-aliases" = ["write_headers"]))]
201pub struct CsvOptions {
202    /// Field separator. A single character; `"\t"` is accepted for tabs.
203    #[serde(default = "default_delimiter")]
204    pub delimiter: String,
205    /// Whether the first row names the fields. When false, fields are named
206    /// `column_0`, `column_1`, … — the same fallback the REST source and the
207    /// `csv` connector already use. On write it decides whether a header row
208    /// is written; `write_headers` is accepted as another name for it.
209    #[serde(default = "default_true", alias = "write_headers")]
210    pub has_headers: bool,
211    /// Quote character. A single byte; default `"`.
212    #[serde(default = "default_quote")]
213    pub quote: String,
214    /// Whether rows may have a different number of fields than the header.
215    /// `false` fails the read on the first ragged row, naming its line; `true`
216    /// keeps the fields the row has. Unset, each connector picks: the `file`
217    /// source is strict, the object-store sources and the REST source are
218    /// lenient.
219    #[serde(default, skip_serializing_if = "Option::is_none")]
220    pub flexible: Option<bool>,
221    /// Cell values read as `null` instead of a string (e.g. `["", "NULL"]`).
222    #[serde(default, skip_serializing_if = "Vec::is_empty")]
223    pub null_values: Vec<String>,
224    /// What a streaming CSV writer does with a field the header does not
225    /// have. See [`CsvUnknownField`]. Whole-object encoders always write the
226    /// union of every record's fields, so it only affects the `file` sink.
227    #[serde(default)]
228    pub on_unknown_field: CsvUnknownField,
229}
230
231/// What a CSV writer does with a field that is not in the header.
232#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
233#[serde(rename_all = "snake_case")]
234pub enum CsvUnknownField {
235    /// Add it as a new column; earlier rows get an empty cell for it.
236    #[default]
237    Widen,
238    /// Fix the header from the first page and drop the field, logging a
239    /// warning once per field.
240    Warn,
241    /// Fix the header from the first page and fail the write naming the field.
242    Error,
243}
244
245fn default_quote() -> String {
246    "\"".into()
247}
248
249impl Default for CsvOptions {
250    fn default() -> Self {
251        Self {
252            delimiter: default_delimiter(),
253            has_headers: true,
254            quote: default_quote(),
255            flexible: None,
256            null_values: Vec::new(),
257            on_unknown_field: CsvUnknownField::Widen,
258        }
259    }
260}
261
262impl CsvOptions {
263    /// The delimiter as the single byte the CSV readers want.
264    ///
265    /// Rejected rather than truncated: a multi-byte delimiter silently becoming
266    /// its first byte would split every row in the wrong place and produce
267    /// plausible-looking garbage.
268    pub fn delimiter_byte(&self) -> Result<u8, FaucetError> {
269        single_byte("delimiter", &self.delimiter)
270    }
271
272    /// The quote character as a single byte.
273    pub fn quote_byte(&self) -> Result<u8, FaucetError> {
274        single_byte("quote", &self.quote)
275    }
276
277    /// Whether ragged rows are accepted, with `default` when unset.
278    pub fn flexible_or(&self, default: bool) -> bool {
279        self.flexible.unwrap_or(default)
280    }
281
282    /// Check the single-byte fields, so a bad dialect fails at load time.
283    pub fn validate(&self) -> Result<(), FaucetError> {
284        self.delimiter_byte()?;
285        self.quote_byte()?;
286        Ok(())
287    }
288}
289
290fn single_byte(name: &str, value: &str) -> Result<u8, FaucetError> {
291    let d = match value {
292        "\\t" => "\t",
293        other => other,
294    };
295    match d.as_bytes() {
296        [b] => Ok(*b),
297        _ => Err(FaucetError::Config(format!(
298            "csv.{name} must be exactly one byte, got {value:?}"
299        ))),
300    }
301}
302
303/// Excel worksheet selection.
304#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
305#[serde(deny_unknown_fields)]
306pub struct ExcelOptions {
307    /// Worksheet name, or an index as a string. Default: the first sheet.
308    #[serde(default, skip_serializing_if = "Option::is_none")]
309    pub sheet: Option<String>,
310    /// 0-based row supplying the field names.
311    #[serde(default)]
312    pub header_row: usize,
313}
314
315fn default_record_element() -> String {
316    "record".into()
317}
318
319fn default_root_element() -> String {
320    "records".into()
321}
322
323/// XML record framing.
324///
325/// XML has no canonical record boundary, so one must be declared. On read,
326/// `record_element` names the repeated element; without it the decoder falls
327/// back to "the document's root children", which is right for the common
328/// `<rows><row/>…</rows>` shape and wrong often enough that naming the element
329/// is recommended in the docs.
330#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
331#[serde(deny_unknown_fields)]
332pub struct XmlOptions {
333    /// Element that delimits one record. Read: select these. Write: wrap each
334    /// record in one.
335    #[serde(default = "default_record_element")]
336    pub record_element: String,
337    /// Document element wrapping the records. Write only.
338    #[serde(default = "default_root_element")]
339    pub root_element: String,
340}
341
342impl Default for XmlOptions {
343    fn default() -> Self {
344        Self {
345            record_element: default_record_element(),
346            root_element: default_root_element(),
347        }
348    }
349}
350
351/// Block compression for written files.
352#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
353#[serde(rename_all = "snake_case")]
354pub enum AvroCodec {
355    /// Uncompressed blocks — the Avro default and the most portable.
356    #[default]
357    Null,
358    /// Raw deflate (RFC 1951), readable by every Avro implementation.
359    Deflate,
360    /// Snappy with the per-block CRC-32 the spec requires.
361    Snappy,
362    /// Zstandard.
363    Zstd,
364}
365
366/// Avro read/write options.
367#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
368#[serde(deny_unknown_fields)]
369pub struct AvroOptions {
370    /// An Avro schema (the JSON form: an object, an array for a union, or a
371    /// primitive name). On a **source** it is the *reader* schema every file
372    /// is resolved against — projection, aliases, defaults for added fields.
373    /// On a **sink** it is the *writer* schema; without it the schema is
374    /// inferred from the records.
375    #[serde(default, skip_serializing_if = "Option::is_none")]
376    pub schema: Option<Value>,
377    /// Block codec for written files: `null` (default), `deflate`, `snappy`,
378    /// `zstd`. Reading detects the codec from the file header.
379    #[serde(default)]
380    pub codec: AvroCodec,
381}
382
383/// ORC read options.
384#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
385#[serde(deny_unknown_fields)]
386pub struct OrcOptions {
387    /// Top-level columns to read, in file order. Default: every column. A
388    /// name the file does not have is an error, not an empty column.
389    #[serde(default, skip_serializing_if = "Option::is_none")]
390    pub columns: Option<Vec<String>>,
391}
392
393/// The per-format option blocks, flattened into a connector's config so they
394/// appear at its top level (`csv: {...}`, `excel: {...}`, `xml: {...}`).
395///
396/// Every block is defaulted, so a connector that adopts this struct adds no
397/// required field and the change stays a minor bump under the project's
398/// [API-evolution policy](https://github.com/faucet-hq/faucet-stream).
399#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
400pub struct FormatOptions {
401    /// CSV dialect, used when `format: csv`.
402    #[serde(default)]
403    pub csv: CsvOptions,
404    /// Worksheet selection, used when `format: xlsx`.
405    #[serde(default)]
406    pub excel: ExcelOptions,
407    /// Record framing, used when `format: xml`.
408    #[serde(default)]
409    pub xml: XmlOptions,
410    /// Reader / writer schema and codec, used when `format: avro`.
411    #[serde(default)]
412    pub avro: AvroOptions,
413    /// Column projection, used when `format: orc`.
414    #[serde(default)]
415    pub orc: OrcOptions,
416}
417
418/// Decode one object's bytes into records.
419///
420/// Async because the CSV reader is (`csv-async`), and because keeping one
421/// signature for every format is worth more than saving an `.await` on the
422/// three that are synchronous.
423pub async fn decode(
424    bytes: &[u8],
425    format: FileFormat,
426    opts: &FormatOptions,
427) -> Result<Vec<Value>, FaucetError> {
428    match format {
429        FileFormat::JsonLines => json::decode_lines(bytes),
430        FileFormat::JsonArray => json::decode_array(bytes),
431        FileFormat::RawText => json::decode_raw_text(bytes),
432        FileFormat::Csv => decode_csv(bytes, opts).await,
433        FileFormat::Xml => decode_xml(bytes, opts),
434        FileFormat::Xlsx => decode_xlsx(bytes, opts),
435        FileFormat::Parquet => Err(columnar_refusal("decode")),
436        FileFormat::Avro => decode_avro(bytes, opts),
437        FileFormat::Orc => decode_orc(bytes, opts),
438    }
439}
440
441/// Encode records into one object's bytes.
442pub fn encode(
443    records: &[Value],
444    format: FileFormat,
445    opts: &FormatOptions,
446) -> Result<Vec<u8>, FaucetError> {
447    match format {
448        FileFormat::JsonLines => json::encode_lines(records),
449        FileFormat::JsonArray => json::encode_array(records),
450        FileFormat::RawText => json::encode_raw_text(records),
451        FileFormat::Csv => encode_csv(records, opts),
452        FileFormat::Xml => encode_xml(records, opts),
453        FileFormat::Xlsx => encode_xlsx(records, opts),
454        FileFormat::Parquet => Err(columnar_refusal("encode")),
455        FileFormat::Avro => encode_avro(records, opts),
456        FileFormat::Orc => Err(FaucetError::Config(
457            "`format: orc` is read-only — there is no ORC writer; write Parquet for a columnar \
458             output"
459                .into(),
460        )),
461    }
462}
463
464/// Parquet never goes through the generic helper.
465///
466/// Routing it here would work and would silently cost the Arrow fast path and
467/// the columnar memory profile, so it is refused rather than accommodated.
468fn columnar_refusal(verb: &str) -> FaucetError {
469    FaucetError::Config(format!(
470        "file_format::{verb}: `parquet` is columnar and is handled by each connector's Arrow \
471         path, not the generic record helper"
472    ))
473}
474
475/// The error a build without the feature reports.
476///
477/// Named rather than mis-parsed: decoding an Excel workbook as CSV produces
478/// records, which is the worst possible outcome.
479#[cfg_attr(
480    all(
481        feature = "file-format-csv",
482        feature = "file-format-xml",
483        feature = "file-format-excel",
484        feature = "file-format-avro",
485        feature = "file-format-orc"
486    ),
487    allow(dead_code)
488)]
489fn missing_feature(format: FileFormat, feature: &str) -> FaucetError {
490    FaucetError::Config(format!(
491        "`format: {}` requires the `{feature}` build feature — rebuild with \
492         `--features {feature}`",
493        format.as_str()
494    ))
495}
496
497#[cfg(feature = "file-format-csv")]
498async fn decode_csv(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
499    csv::decode_with(bytes, &opts.csv, true).await
500}
501
502#[cfg(not(feature = "file-format-csv"))]
503async fn decode_csv(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
504    Err(missing_feature(FileFormat::Csv, "file-format-csv"))
505}
506
507#[cfg(feature = "file-format-csv")]
508fn encode_csv(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
509    csv::encode_with(records, &opts.csv)
510}
511
512#[cfg(not(feature = "file-format-csv"))]
513fn encode_csv(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
514    Err(missing_feature(FileFormat::Csv, "file-format-csv"))
515}
516
517#[cfg(feature = "file-format-xml")]
518fn decode_xml(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
519    xml::decode(bytes, &opts.xml.record_element)
520}
521
522#[cfg(not(feature = "file-format-xml"))]
523fn decode_xml(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
524    Err(missing_feature(FileFormat::Xml, "file-format-xml"))
525}
526
527#[cfg(feature = "file-format-xml")]
528fn encode_xml(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
529    xml::encode(records, &opts.xml.root_element, &opts.xml.record_element)
530}
531
532#[cfg(not(feature = "file-format-xml"))]
533fn encode_xml(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
534    Err(missing_feature(FileFormat::Xml, "file-format-xml"))
535}
536
537#[cfg(feature = "file-format-excel")]
538fn decode_xlsx(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
539    excel::decode(bytes, opts.excel.sheet.as_deref(), opts.excel.header_row)
540}
541
542#[cfg(not(feature = "file-format-excel"))]
543fn decode_xlsx(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
544    Err(missing_feature(FileFormat::Xlsx, "file-format-excel"))
545}
546
547#[cfg(feature = "file-format-excel")]
548fn encode_xlsx(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
549    excel::encode(records, opts.excel.sheet.as_deref())
550}
551
552#[cfg(not(feature = "file-format-excel"))]
553fn encode_xlsx(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
554    Err(missing_feature(FileFormat::Xlsx, "file-format-excel"))
555}
556
557#[cfg(feature = "file-format-avro")]
558fn decode_avro(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
559    avro::decode(bytes, &opts.avro)
560}
561
562#[cfg(not(feature = "file-format-avro"))]
563fn decode_avro(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
564    Err(missing_feature(FileFormat::Avro, "file-format-avro"))
565}
566
567#[cfg(feature = "file-format-avro")]
568fn encode_avro(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
569    avro::encode(records, &opts.avro)
570}
571
572#[cfg(not(feature = "file-format-avro"))]
573fn encode_avro(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
574    Err(missing_feature(FileFormat::Avro, "file-format-avro"))
575}
576
577#[cfg(feature = "file-format-orc")]
578fn decode_orc(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
579    orc::decode(bytes, &opts.orc)
580}
581
582#[cfg(not(feature = "file-format-orc"))]
583fn decode_orc(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
584    Err(missing_feature(FileFormat::Orc, "file-format-orc"))
585}
586
587/// Field names across `records`, in first-seen order.
588///
589/// A union rather than the first record's keys, so a later record carrying an
590/// extra field widens the file instead of losing the field silently.
591///
592/// "First-seen" orders *records* deterministically; within one record the order
593/// is whatever `serde_json::Map` iterates, which is insertion order only when
594/// the `preserve_order` feature is unified into the build and sorted otherwise.
595/// Column order is therefore a build-time property, not a guarantee — do not
596/// write a test that pins it.
597pub fn header_union(records: &[Value]) -> Vec<String> {
598    let mut seen = std::collections::HashSet::new();
599    let mut out = Vec::new();
600    for r in records {
601        if let Value::Object(map) = r {
602            for k in map.keys() {
603                if seen.insert(k.clone()) {
604                    out.push(k.clone());
605                }
606            }
607        }
608    }
609    out
610}
611
612/// One cell's text for the tabular formats.
613///
614/// A nested object or array is re-serialized as JSON rather than dropped or
615/// `Debug`-printed: a spreadsheet cell cannot hold structure, and the round
616/// trip back through `json_parse` is at least lossless.
617pub fn cell_text(v: &Value) -> String {
618    match v {
619        Value::Null => String::new(),
620        Value::String(s) => s.clone(),
621        Value::Bool(b) => b.to_string(),
622        Value::Number(n) => n.to_string(),
623        other => serde_json::to_string(other).unwrap_or_default(),
624    }
625}
626
627#[cfg(test)]
628mod tests {
629    use super::*;
630    use serde_json::json;
631
632    /// Every variant's extension and wire name, so adding a format without
633    /// giving it both is a test failure rather than a `.txt` file called
634    /// "raw_text".
635    #[test]
636    fn every_variant_has_an_extension_and_a_wire_name() {
637        let all = [
638            (FileFormat::JsonLines, ".jsonl", "json_lines"),
639            (FileFormat::JsonArray, ".json", "json_array"),
640            (FileFormat::Csv, ".csv", "csv"),
641            (FileFormat::Xml, ".xml", "xml"),
642            (FileFormat::Xlsx, ".xlsx", "xlsx"),
643            (FileFormat::Parquet, ".parquet", "parquet"),
644            (FileFormat::RawText, ".txt", "raw_text"),
645            (FileFormat::Avro, ".avro", "avro"),
646            (FileFormat::Orc, ".orc", "orc"),
647        ];
648        for (f, ext, name) in all {
649            assert_eq!(f.extension(), ext, "{f:?}");
650            assert_eq!(f.as_str(), name, "{f:?}");
651        }
652    }
653
654    /// `raw_text` is source-only in the connectors, but the shared helper
655    /// still round-trips it — the sink side is what `encode_raw_text` exists
656    /// for, and routing it through the top-level dispatch is how a connector
657    /// reaches it.
658    #[tokio::test]
659    async fn raw_text_round_trips_through_the_top_level_dispatch() {
660        let opts = FormatOptions::default();
661        let recs = decode(b"hello", FileFormat::RawText, &opts)
662            .await
663            .expect("decode");
664        assert_eq!(recs, vec![json!({"text": "hello"})]);
665        let bytes = encode(&recs, FileFormat::RawText, &opts).expect("encode");
666        assert_eq!(bytes, b"hello\n");
667    }
668
669    /// The refusal a build without the feature emits — it must name the
670    /// format and the feature, because mis-parsing an Excel blob as CSV is
671    /// the outcome this exists to prevent.
672    #[test]
673    fn a_missing_feature_is_named_not_guessed() {
674        let err = missing_feature(FileFormat::Xlsx, "file-format-excel");
675        let msg = err.to_string();
676        assert!(msg.contains("xlsx"), "{msg}");
677        assert!(msg.contains("file-format-excel"), "{msg}");
678    }
679
680    #[test]
681    fn extensions_and_names_are_stable() {
682        assert_eq!(FileFormat::default(), FileFormat::JsonLines);
683        assert_eq!(FileFormat::Csv.extension(), ".csv");
684        assert_eq!(FileFormat::Xlsx.as_str(), "xlsx");
685    }
686
687    #[test]
688    fn whole_object_formats_are_the_ones_without_a_record_boundary() {
689        for f in [
690            FileFormat::JsonArray,
691            FileFormat::Xml,
692            FileFormat::Xlsx,
693            FileFormat::Orc,
694        ] {
695            assert!(f.requires_whole_object(), "{f:?}");
696        }
697        for f in [
698            FileFormat::JsonLines,
699            FileFormat::Csv,
700            FileFormat::RawText,
701            FileFormat::Parquet,
702            FileFormat::Avro,
703        ] {
704            assert!(!f.requires_whole_object(), "{f:?}");
705        }
706    }
707
708    #[test]
709    fn extensions_resolve_through_compression_suffixes() {
710        let cases = [
711            ("a.jsonl", Some(FileFormat::JsonLines)),
712            ("dir/b.NDJSON", Some(FileFormat::JsonLines)),
713            ("c.json.gz", Some(FileFormat::JsonArray)),
714            ("d.csv.gz", Some(FileFormat::Csv)),
715            ("e.xml.zst", Some(FileFormat::Xml)),
716            ("f.xlsx", Some(FileFormat::Xlsx)),
717            ("g.parquet", Some(FileFormat::Parquet)),
718            ("h.txt", Some(FileFormat::RawText)),
719            ("i.avro", Some(FileFormat::Avro)),
720            ("j.orc", Some(FileFormat::Orc)),
721            ("k.csv.gzip", Some(FileFormat::Csv)),
722            ("l.jsonl.zstd", Some(FileFormat::JsonLines)),
723            ("m.tsv", None),
724            ("noext", None),
725            ("n.gz", None),
726        ];
727        for (path, want) in cases {
728            assert_eq!(FileFormat::from_path(path), want, "{path}");
729        }
730    }
731
732    #[test]
733    fn container_and_writable_classification() {
734        assert!(FileFormat::Avro.is_container() && FileFormat::Orc.is_container());
735        assert!(!FileFormat::Csv.is_container());
736        assert!(FileFormat::Avro.is_writable());
737        assert!(!FileFormat::Orc.is_writable() && !FileFormat::Parquet.is_writable());
738    }
739
740    #[tokio::test]
741    async fn orc_is_refused_for_writing() {
742        let err = encode(&[], FileFormat::Orc, &FormatOptions::default()).expect_err("orc");
743        assert!(err.to_string().contains("read-only"), "{err}");
744    }
745
746    #[tokio::test]
747    async fn parquet_is_refused_on_both_directions() {
748        // Going through the generic helper would work and would silently cost
749        // the Arrow path, so it must be an error, not a fallback.
750        let err = decode(b"", FileFormat::Parquet, &FormatOptions::default())
751            .await
752            .expect_err("decode refuses parquet");
753        assert!(err.to_string().contains("columnar"), "{err}");
754        let err = encode(&[], FileFormat::Parquet, &FormatOptions::default())
755            .expect_err("encode refuses parquet");
756        assert!(err.to_string().contains("columnar"), "{err}");
757    }
758
759    #[test]
760    fn delimiter_must_be_one_byte() {
761        let mut o = CsvOptions::default();
762        assert_eq!(o.delimiter_byte().expect("default"), b',');
763        o.delimiter = "\\t".into();
764        assert_eq!(o.delimiter_byte().expect("tab"), b'\t');
765        o.delimiter = "||".into();
766        let err = o.delimiter_byte().expect_err("multi-byte");
767        assert!(err.to_string().contains("exactly one byte"), "{err}");
768        // A multi-byte character is rejected too — truncating it to its first
769        // UTF-8 byte would split rows at a byte that never appears alone.
770        o.delimiter = "§".into();
771        assert!(o.delimiter_byte().is_err());
772    }
773
774    #[test]
775    fn header_union_widens_across_records_and_ignores_non_objects() {
776        // Records are visited in order, so a key only this record has lands
777        // after the earlier ones. Intra-record order is a build-time property
778        // of `serde_json::Map`, so it is deliberately not asserted.
779        let recs = vec![json!({"a": 1}), json!({"b": 2}), json!("not an object")];
780        assert_eq!(header_union(&recs), vec!["a", "b"]);
781        // A key seen twice appears once.
782        assert_eq!(header_union(&[json!({"a": 1}), json!({"a": 2})]), vec!["a"]);
783        assert!(header_union(&[json!([1, 2])]).is_empty());
784    }
785
786    #[test]
787    fn cells_render_scalars_plainly_and_structure_as_json() {
788        assert_eq!(cell_text(&Value::Null), "");
789        assert_eq!(cell_text(&json!("x")), "x");
790        assert_eq!(cell_text(&json!(3)), "3");
791        assert_eq!(cell_text(&json!(true)), "true");
792        assert_eq!(cell_text(&json!({"k": 1})), r#"{"k":1}"#);
793    }
794}