1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
76#[serde(rename_all = "snake_case")]
77pub enum FileFormat {
78 #[default]
81 JsonLines,
82 JsonArray,
84 Csv,
86 Xml,
88 Xlsx,
90 Parquet,
93 RawText,
95 Avro,
98 Orc,
101}
102
103impl FileFormat {
104 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 pub fn requires_whole_object(self) -> bool {
127 matches!(self, Self::JsonArray | Self::Xml | Self::Xlsx | Self::Orc)
128 }
129
130 pub fn is_container(self) -> bool {
133 matches!(self, Self::Avro | Self::Orc)
134 }
135
136 pub fn is_writable(self) -> bool {
138 !matches!(self, Self::Orc | Self::Parquet)
139 }
140
141 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 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
199#[serde(deny_unknown_fields)]
200pub struct CsvOptions {
201 #[serde(default = "default_delimiter")]
203 pub delimiter: String,
204 #[serde(default = "default_true")]
209 pub has_headers: bool,
210 #[serde(default = "default_quote")]
212 pub quote: String,
213 #[serde(default, skip_serializing_if = "Option::is_none")]
219 pub flexible: Option<bool>,
220 #[serde(default, skip_serializing_if = "Vec::is_empty")]
222 pub null_values: Vec<String>,
223 #[serde(default)]
227 pub on_unknown_field: CsvUnknownField,
228}
229
230#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
232#[serde(rename_all = "snake_case")]
233pub enum CsvUnknownField {
234 #[default]
236 Widen,
237 Warn,
240 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 pub fn delimiter_byte(&self) -> Result<u8, FaucetError> {
268 single_byte("delimiter", &self.delimiter)
269 }
270
271 pub fn quote_byte(&self) -> Result<u8, FaucetError> {
273 single_byte("quote", &self.quote)
274 }
275
276 pub fn flexible_or(&self, default: bool) -> bool {
278 self.flexible.unwrap_or(default)
279 }
280
281 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#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
304#[serde(deny_unknown_fields)]
305pub struct ExcelOptions {
306 #[serde(default, skip_serializing_if = "Option::is_none")]
308 pub sheet: Option<String>,
309 #[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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
330#[serde(deny_unknown_fields)]
331pub struct XmlOptions {
332 #[serde(default = "default_record_element")]
335 pub record_element: String,
336 #[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
352#[serde(rename_all = "snake_case")]
353pub enum AvroCodec {
354 #[default]
356 Null,
357 Deflate,
359 Snappy,
361 Zstd,
363}
364
365#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
367#[serde(deny_unknown_fields)]
368pub struct AvroOptions {
369 #[serde(default, skip_serializing_if = "Option::is_none")]
375 pub schema: Option<Value>,
376 #[serde(default)]
379 pub codec: AvroCodec,
380}
381
382#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
384#[serde(deny_unknown_fields)]
385pub struct OrcOptions {
386 #[serde(default, skip_serializing_if = "Option::is_none")]
389 pub columns: Option<Vec<String>>,
390}
391
392#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
399pub struct FormatOptions {
400 #[serde(default)]
402 pub csv: CsvOptions,
403 #[serde(default)]
405 pub excel: ExcelOptions,
406 #[serde(default)]
408 pub xml: XmlOptions,
409 #[serde(default)]
411 pub avro: AvroOptions,
412 #[serde(default)]
414 pub orc: OrcOptions,
415}
416
417pub 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
440pub 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
463fn 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#[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
586pub 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
611pub 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 #[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 #[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 #[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 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 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 let recs = vec![json!({"a": 1}), json!({"b": 2}), json!("not an object")];
779 assert_eq!(header_union(&recs), vec!["a", "b"]);
780 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}