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)]
200pub struct CsvOptions {
201    /// Field separator. A single character; `"\t"` is accepted for tabs.
202    #[serde(default = "default_delimiter")]
203    pub delimiter: String,
204    /// Whether the first row names the fields. When false, fields are named
205    /// `column_0`, `column_1`, … — the same fallback the REST source and the
206    /// `csv` connector already use. On write it decides whether a header row
207    /// is written.
208    #[serde(default = "default_true")]
209    pub has_headers: bool,
210    /// Quote character. A single byte; default `"`.
211    #[serde(default = "default_quote")]
212    pub quote: String,
213    /// Whether rows may have a different number of fields than the header.
214    /// `false` fails the read on the first ragged row, naming its line; `true`
215    /// keeps the fields the row has. Unset, each connector picks: the `file`
216    /// source is strict, the object-store sources and the REST source are
217    /// lenient.
218    #[serde(default, skip_serializing_if = "Option::is_none")]
219    pub flexible: Option<bool>,
220    /// Cell values read as `null` instead of a string (e.g. `["", "NULL"]`).
221    #[serde(default, skip_serializing_if = "Vec::is_empty")]
222    pub null_values: Vec<String>,
223    /// What a streaming CSV writer does with a field the header does not
224    /// have. See [`CsvUnknownField`]. Whole-object encoders always write the
225    /// union of every record's fields, so it only affects the `file` sink.
226    #[serde(default)]
227    pub on_unknown_field: CsvUnknownField,
228}
229
230/// What a CSV writer does with a field that is not in the header.
231#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
232#[serde(rename_all = "snake_case")]
233pub enum CsvUnknownField {
234    /// Add it as a new column; earlier rows get an empty cell for it.
235    #[default]
236    Widen,
237    /// Fix the header from the first page and drop the field, logging a
238    /// warning once per field.
239    Warn,
240    /// Fix the header from the first page and fail the write naming the field.
241    Error,
242}
243
244fn default_quote() -> String {
245    "\"".into()
246}
247
248impl Default for CsvOptions {
249    fn default() -> Self {
250        Self {
251            delimiter: default_delimiter(),
252            has_headers: true,
253            quote: default_quote(),
254            flexible: None,
255            null_values: Vec::new(),
256            on_unknown_field: CsvUnknownField::Widen,
257        }
258    }
259}
260
261impl CsvOptions {
262    /// The delimiter as the single byte the CSV readers want.
263    ///
264    /// Rejected rather than truncated: a multi-byte delimiter silently becoming
265    /// its first byte would split every row in the wrong place and produce
266    /// plausible-looking garbage.
267    pub fn delimiter_byte(&self) -> Result<u8, FaucetError> {
268        single_byte("delimiter", &self.delimiter)
269    }
270
271    /// The quote character as a single byte.
272    pub fn quote_byte(&self) -> Result<u8, FaucetError> {
273        single_byte("quote", &self.quote)
274    }
275
276    /// Whether ragged rows are accepted, with `default` when unset.
277    pub fn flexible_or(&self, default: bool) -> bool {
278        self.flexible.unwrap_or(default)
279    }
280
281    /// Check the single-byte fields, so a bad dialect fails at load time.
282    pub fn validate(&self) -> Result<(), FaucetError> {
283        self.delimiter_byte()?;
284        self.quote_byte()?;
285        Ok(())
286    }
287}
288
289fn single_byte(name: &str, value: &str) -> Result<u8, FaucetError> {
290    let d = match value {
291        "\\t" => "\t",
292        other => other,
293    };
294    match d.as_bytes() {
295        [b] => Ok(*b),
296        _ => Err(FaucetError::Config(format!(
297            "csv.{name} must be exactly one byte, got {value:?}"
298        ))),
299    }
300}
301
302/// Excel worksheet selection.
303#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
304#[serde(deny_unknown_fields)]
305pub struct ExcelOptions {
306    /// Worksheet name, or an index as a string. Default: the first sheet.
307    #[serde(default, skip_serializing_if = "Option::is_none")]
308    pub sheet: Option<String>,
309    /// 0-based row supplying the field names.
310    #[serde(default)]
311    pub header_row: usize,
312}
313
314fn default_record_element() -> String {
315    "record".into()
316}
317
318fn default_root_element() -> String {
319    "records".into()
320}
321
322/// XML record framing.
323///
324/// XML has no canonical record boundary, so one must be declared. On read,
325/// `record_element` names the repeated element; without it the decoder falls
326/// back to "the document's root children", which is right for the common
327/// `<rows><row/>…</rows>` shape and wrong often enough that naming the element
328/// is recommended in the docs.
329#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
330#[serde(deny_unknown_fields)]
331pub struct XmlOptions {
332    /// Element that delimits one record. Read: select these. Write: wrap each
333    /// record in one.
334    #[serde(default = "default_record_element")]
335    pub record_element: String,
336    /// Document element wrapping the records. Write only.
337    #[serde(default = "default_root_element")]
338    pub root_element: String,
339}
340
341impl Default for XmlOptions {
342    fn default() -> Self {
343        Self {
344            record_element: default_record_element(),
345            root_element: default_root_element(),
346        }
347    }
348}
349
350/// Block compression for written files.
351#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
352#[serde(rename_all = "snake_case")]
353pub enum AvroCodec {
354    /// Uncompressed blocks — the Avro default and the most portable.
355    #[default]
356    Null,
357    /// Raw deflate (RFC 1951), readable by every Avro implementation.
358    Deflate,
359    /// Snappy with the per-block CRC-32 the spec requires.
360    Snappy,
361    /// Zstandard.
362    Zstd,
363}
364
365/// Avro read/write options.
366#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
367#[serde(deny_unknown_fields)]
368pub struct AvroOptions {
369    /// An Avro schema (the JSON form: an object, an array for a union, or a
370    /// primitive name). On a **source** it is the *reader* schema every file
371    /// is resolved against — projection, aliases, defaults for added fields.
372    /// On a **sink** it is the *writer* schema; without it the schema is
373    /// inferred from the records.
374    #[serde(default, skip_serializing_if = "Option::is_none")]
375    pub schema: Option<Value>,
376    /// Block codec for written files: `null` (default), `deflate`, `snappy`,
377    /// `zstd`. Reading detects the codec from the file header.
378    #[serde(default)]
379    pub codec: AvroCodec,
380}
381
382/// ORC read options.
383#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
384#[serde(deny_unknown_fields)]
385pub struct OrcOptions {
386    /// Top-level columns to read, in file order. Default: every column. A
387    /// name the file does not have is an error, not an empty column.
388    #[serde(default, skip_serializing_if = "Option::is_none")]
389    pub columns: Option<Vec<String>>,
390}
391
392/// The per-format option blocks, flattened into a connector's config so they
393/// appear at its top level (`csv: {...}`, `excel: {...}`, `xml: {...}`).
394///
395/// Every block is defaulted, so a connector that adopts this struct adds no
396/// required field and the change stays a minor bump under the project's
397/// [API-evolution policy](https://github.com/faucet-hq/faucet-stream).
398#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
399pub struct FormatOptions {
400    /// CSV dialect, used when `format: csv`.
401    #[serde(default)]
402    pub csv: CsvOptions,
403    /// Worksheet selection, used when `format: xlsx`.
404    #[serde(default)]
405    pub excel: ExcelOptions,
406    /// Record framing, used when `format: xml`.
407    #[serde(default)]
408    pub xml: XmlOptions,
409    /// Reader / writer schema and codec, used when `format: avro`.
410    #[serde(default)]
411    pub avro: AvroOptions,
412    /// Column projection, used when `format: orc`.
413    #[serde(default)]
414    pub orc: OrcOptions,
415}
416
417/// Decode one object's bytes into records.
418///
419/// Async because the CSV reader is (`csv-async`), and because keeping one
420/// signature for every format is worth more than saving an `.await` on the
421/// three that are synchronous.
422pub async fn decode(
423    bytes: &[u8],
424    format: FileFormat,
425    opts: &FormatOptions,
426) -> Result<Vec<Value>, FaucetError> {
427    match format {
428        FileFormat::JsonLines => json::decode_lines(bytes),
429        FileFormat::JsonArray => json::decode_array(bytes),
430        FileFormat::RawText => json::decode_raw_text(bytes),
431        FileFormat::Csv => decode_csv(bytes, opts).await,
432        FileFormat::Xml => decode_xml(bytes, opts),
433        FileFormat::Xlsx => decode_xlsx(bytes, opts),
434        FileFormat::Parquet => Err(columnar_refusal("decode")),
435        FileFormat::Avro => decode_avro(bytes, opts),
436        FileFormat::Orc => decode_orc(bytes, opts),
437    }
438}
439
440/// Encode records into one object's bytes.
441pub fn encode(
442    records: &[Value],
443    format: FileFormat,
444    opts: &FormatOptions,
445) -> Result<Vec<u8>, FaucetError> {
446    match format {
447        FileFormat::JsonLines => json::encode_lines(records),
448        FileFormat::JsonArray => json::encode_array(records),
449        FileFormat::RawText => json::encode_raw_text(records),
450        FileFormat::Csv => encode_csv(records, opts),
451        FileFormat::Xml => encode_xml(records, opts),
452        FileFormat::Xlsx => encode_xlsx(records, opts),
453        FileFormat::Parquet => Err(columnar_refusal("encode")),
454        FileFormat::Avro => encode_avro(records, opts),
455        FileFormat::Orc => Err(FaucetError::Config(
456            "`format: orc` is read-only — there is no ORC writer; write Parquet for a columnar \
457             output"
458                .into(),
459        )),
460    }
461}
462
463/// Parquet never goes through the generic helper.
464///
465/// Routing it here would work and would silently cost the Arrow fast path and
466/// the columnar memory profile, so it is refused rather than accommodated.
467fn columnar_refusal(verb: &str) -> FaucetError {
468    FaucetError::Config(format!(
469        "file_format::{verb}: `parquet` is columnar and is handled by each connector's Arrow \
470         path, not the generic record helper"
471    ))
472}
473
474/// The error a build without the feature reports.
475///
476/// Named rather than mis-parsed: decoding an Excel workbook as CSV produces
477/// records, which is the worst possible outcome.
478#[cfg_attr(
479    all(
480        feature = "file-format-csv",
481        feature = "file-format-xml",
482        feature = "file-format-excel",
483        feature = "file-format-avro",
484        feature = "file-format-orc"
485    ),
486    allow(dead_code)
487)]
488fn missing_feature(format: FileFormat, feature: &str) -> FaucetError {
489    FaucetError::Config(format!(
490        "`format: {}` requires the `{feature}` build feature — rebuild with \
491         `--features {feature}`",
492        format.as_str()
493    ))
494}
495
496#[cfg(feature = "file-format-csv")]
497async fn decode_csv(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
498    csv::decode_with(bytes, &opts.csv, true).await
499}
500
501#[cfg(not(feature = "file-format-csv"))]
502async fn decode_csv(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
503    Err(missing_feature(FileFormat::Csv, "file-format-csv"))
504}
505
506#[cfg(feature = "file-format-csv")]
507fn encode_csv(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
508    csv::encode_with(records, &opts.csv)
509}
510
511#[cfg(not(feature = "file-format-csv"))]
512fn encode_csv(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
513    Err(missing_feature(FileFormat::Csv, "file-format-csv"))
514}
515
516#[cfg(feature = "file-format-xml")]
517fn decode_xml(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
518    xml::decode(bytes, &opts.xml.record_element)
519}
520
521#[cfg(not(feature = "file-format-xml"))]
522fn decode_xml(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
523    Err(missing_feature(FileFormat::Xml, "file-format-xml"))
524}
525
526#[cfg(feature = "file-format-xml")]
527fn encode_xml(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
528    xml::encode(records, &opts.xml.root_element, &opts.xml.record_element)
529}
530
531#[cfg(not(feature = "file-format-xml"))]
532fn encode_xml(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
533    Err(missing_feature(FileFormat::Xml, "file-format-xml"))
534}
535
536#[cfg(feature = "file-format-excel")]
537fn decode_xlsx(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
538    excel::decode(bytes, opts.excel.sheet.as_deref(), opts.excel.header_row)
539}
540
541#[cfg(not(feature = "file-format-excel"))]
542fn decode_xlsx(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
543    Err(missing_feature(FileFormat::Xlsx, "file-format-excel"))
544}
545
546#[cfg(feature = "file-format-excel")]
547fn encode_xlsx(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
548    excel::encode(records, opts.excel.sheet.as_deref())
549}
550
551#[cfg(not(feature = "file-format-excel"))]
552fn encode_xlsx(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
553    Err(missing_feature(FileFormat::Xlsx, "file-format-excel"))
554}
555
556#[cfg(feature = "file-format-avro")]
557fn decode_avro(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
558    avro::decode(bytes, &opts.avro)
559}
560
561#[cfg(not(feature = "file-format-avro"))]
562fn decode_avro(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
563    Err(missing_feature(FileFormat::Avro, "file-format-avro"))
564}
565
566#[cfg(feature = "file-format-avro")]
567fn encode_avro(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
568    avro::encode(records, &opts.avro)
569}
570
571#[cfg(not(feature = "file-format-avro"))]
572fn encode_avro(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
573    Err(missing_feature(FileFormat::Avro, "file-format-avro"))
574}
575
576#[cfg(feature = "file-format-orc")]
577fn decode_orc(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
578    orc::decode(bytes, &opts.orc)
579}
580
581#[cfg(not(feature = "file-format-orc"))]
582fn decode_orc(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
583    Err(missing_feature(FileFormat::Orc, "file-format-orc"))
584}
585
586/// Field names across `records`, in first-seen order.
587///
588/// A union rather than the first record's keys, so a later record carrying an
589/// extra field widens the file instead of losing the field silently.
590///
591/// "First-seen" orders *records* deterministically; within one record the order
592/// is whatever `serde_json::Map` iterates, which is insertion order only when
593/// the `preserve_order` feature is unified into the build and sorted otherwise.
594/// Column order is therefore a build-time property, not a guarantee — do not
595/// write a test that pins it.
596pub fn header_union(records: &[Value]) -> Vec<String> {
597    let mut seen = std::collections::HashSet::new();
598    let mut out = Vec::new();
599    for r in records {
600        if let Value::Object(map) = r {
601            for k in map.keys() {
602                if seen.insert(k.clone()) {
603                    out.push(k.clone());
604                }
605            }
606        }
607    }
608    out
609}
610
611/// One cell's text for the tabular formats.
612///
613/// A nested object or array is re-serialized as JSON rather than dropped or
614/// `Debug`-printed: a spreadsheet cell cannot hold structure, and the round
615/// trip back through `json_parse` is at least lossless.
616pub fn cell_text(v: &Value) -> String {
617    match v {
618        Value::Null => String::new(),
619        Value::String(s) => s.clone(),
620        Value::Bool(b) => b.to_string(),
621        Value::Number(n) => n.to_string(),
622        other => serde_json::to_string(other).unwrap_or_default(),
623    }
624}
625
626#[cfg(test)]
627mod tests {
628    use super::*;
629    use serde_json::json;
630
631    /// Every variant's extension and wire name, so adding a format without
632    /// giving it both is a test failure rather than a `.txt` file called
633    /// "raw_text".
634    #[test]
635    fn every_variant_has_an_extension_and_a_wire_name() {
636        let all = [
637            (FileFormat::JsonLines, ".jsonl", "json_lines"),
638            (FileFormat::JsonArray, ".json", "json_array"),
639            (FileFormat::Csv, ".csv", "csv"),
640            (FileFormat::Xml, ".xml", "xml"),
641            (FileFormat::Xlsx, ".xlsx", "xlsx"),
642            (FileFormat::Parquet, ".parquet", "parquet"),
643            (FileFormat::RawText, ".txt", "raw_text"),
644            (FileFormat::Avro, ".avro", "avro"),
645            (FileFormat::Orc, ".orc", "orc"),
646        ];
647        for (f, ext, name) in all {
648            assert_eq!(f.extension(), ext, "{f:?}");
649            assert_eq!(f.as_str(), name, "{f:?}");
650        }
651    }
652
653    /// `raw_text` is source-only in the connectors, but the shared helper
654    /// still round-trips it — the sink side is what `encode_raw_text` exists
655    /// for, and routing it through the top-level dispatch is how a connector
656    /// reaches it.
657    #[tokio::test]
658    async fn raw_text_round_trips_through_the_top_level_dispatch() {
659        let opts = FormatOptions::default();
660        let recs = decode(b"hello", FileFormat::RawText, &opts)
661            .await
662            .expect("decode");
663        assert_eq!(recs, vec![json!({"text": "hello"})]);
664        let bytes = encode(&recs, FileFormat::RawText, &opts).expect("encode");
665        assert_eq!(bytes, b"hello\n");
666    }
667
668    /// The refusal a build without the feature emits — it must name the
669    /// format and the feature, because mis-parsing an Excel blob as CSV is
670    /// the outcome this exists to prevent.
671    #[test]
672    fn a_missing_feature_is_named_not_guessed() {
673        let err = missing_feature(FileFormat::Xlsx, "file-format-excel");
674        let msg = err.to_string();
675        assert!(msg.contains("xlsx"), "{msg}");
676        assert!(msg.contains("file-format-excel"), "{msg}");
677    }
678
679    #[test]
680    fn extensions_and_names_are_stable() {
681        assert_eq!(FileFormat::default(), FileFormat::JsonLines);
682        assert_eq!(FileFormat::Csv.extension(), ".csv");
683        assert_eq!(FileFormat::Xlsx.as_str(), "xlsx");
684    }
685
686    #[test]
687    fn whole_object_formats_are_the_ones_without_a_record_boundary() {
688        for f in [
689            FileFormat::JsonArray,
690            FileFormat::Xml,
691            FileFormat::Xlsx,
692            FileFormat::Orc,
693        ] {
694            assert!(f.requires_whole_object(), "{f:?}");
695        }
696        for f in [
697            FileFormat::JsonLines,
698            FileFormat::Csv,
699            FileFormat::RawText,
700            FileFormat::Parquet,
701            FileFormat::Avro,
702        ] {
703            assert!(!f.requires_whole_object(), "{f:?}");
704        }
705    }
706
707    #[test]
708    fn extensions_resolve_through_compression_suffixes() {
709        let cases = [
710            ("a.jsonl", Some(FileFormat::JsonLines)),
711            ("dir/b.NDJSON", Some(FileFormat::JsonLines)),
712            ("c.json.gz", Some(FileFormat::JsonArray)),
713            ("d.csv.gz", Some(FileFormat::Csv)),
714            ("e.xml.zst", Some(FileFormat::Xml)),
715            ("f.xlsx", Some(FileFormat::Xlsx)),
716            ("g.parquet", Some(FileFormat::Parquet)),
717            ("h.txt", Some(FileFormat::RawText)),
718            ("i.avro", Some(FileFormat::Avro)),
719            ("j.orc", Some(FileFormat::Orc)),
720            ("k.csv.gzip", Some(FileFormat::Csv)),
721            ("l.jsonl.zstd", Some(FileFormat::JsonLines)),
722            ("m.tsv", None),
723            ("noext", None),
724            ("n.gz", None),
725        ];
726        for (path, want) in cases {
727            assert_eq!(FileFormat::from_path(path), want, "{path}");
728        }
729    }
730
731    #[test]
732    fn container_and_writable_classification() {
733        assert!(FileFormat::Avro.is_container() && FileFormat::Orc.is_container());
734        assert!(!FileFormat::Csv.is_container());
735        assert!(FileFormat::Avro.is_writable());
736        assert!(!FileFormat::Orc.is_writable() && !FileFormat::Parquet.is_writable());
737    }
738
739    #[tokio::test]
740    async fn orc_is_refused_for_writing() {
741        let err = encode(&[], FileFormat::Orc, &FormatOptions::default()).expect_err("orc");
742        assert!(err.to_string().contains("read-only"), "{err}");
743    }
744
745    #[tokio::test]
746    async fn parquet_is_refused_on_both_directions() {
747        // Going through the generic helper would work and would silently cost
748        // the Arrow path, so it must be an error, not a fallback.
749        let err = decode(b"", FileFormat::Parquet, &FormatOptions::default())
750            .await
751            .expect_err("decode refuses parquet");
752        assert!(err.to_string().contains("columnar"), "{err}");
753        let err = encode(&[], FileFormat::Parquet, &FormatOptions::default())
754            .expect_err("encode refuses parquet");
755        assert!(err.to_string().contains("columnar"), "{err}");
756    }
757
758    #[test]
759    fn delimiter_must_be_one_byte() {
760        let mut o = CsvOptions::default();
761        assert_eq!(o.delimiter_byte().expect("default"), b',');
762        o.delimiter = "\\t".into();
763        assert_eq!(o.delimiter_byte().expect("tab"), b'\t');
764        o.delimiter = "||".into();
765        let err = o.delimiter_byte().expect_err("multi-byte");
766        assert!(err.to_string().contains("exactly one byte"), "{err}");
767        // A multi-byte character is rejected too — truncating it to its first
768        // UTF-8 byte would split rows at a byte that never appears alone.
769        o.delimiter = "§".into();
770        assert!(o.delimiter_byte().is_err());
771    }
772
773    #[test]
774    fn header_union_widens_across_records_and_ignores_non_objects() {
775        // Records are visited in order, so a key only this record has lands
776        // after the earlier ones. Intra-record order is a build-time property
777        // of `serde_json::Map`, so it is deliberately not asserted.
778        let recs = vec![json!({"a": 1}), json!({"b": 2}), json!("not an object")];
779        assert_eq!(header_union(&recs), vec!["a", "b"]);
780        // A key seen twice appears once.
781        assert_eq!(header_union(&[json!({"a": 1}), json!({"a": 2})]), vec!["a"]);
782        assert!(header_union(&[json!([1, 2])]).is_empty());
783    }
784
785    #[test]
786    fn cells_render_scalars_plainly_and_structure_as_json() {
787        assert_eq!(cell_text(&Value::Null), "");
788        assert_eq!(cell_text(&json!("x")), "x");
789        assert_eq!(cell_text(&json!(3)), "3");
790        assert_eq!(cell_text(&json!(true)), "true");
791        assert_eq!(cell_text(&json!({"k": 1})), r#"{"k":1}"#);
792    }
793}