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)]
200#[schemars(extend("x-faucet-aliases" = ["write_headers"]))]
201pub struct CsvOptions {
202 #[serde(default = "default_delimiter")]
204 pub delimiter: String,
205 #[serde(default = "default_true", alias = "write_headers")]
210 pub has_headers: bool,
211 #[serde(default = "default_quote")]
213 pub quote: String,
214 #[serde(default, skip_serializing_if = "Option::is_none")]
220 pub flexible: Option<bool>,
221 #[serde(default, skip_serializing_if = "Vec::is_empty")]
223 pub null_values: Vec<String>,
224 #[serde(default)]
228 pub on_unknown_field: CsvUnknownField,
229}
230
231#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
233#[serde(rename_all = "snake_case")]
234pub enum CsvUnknownField {
235 #[default]
237 Widen,
238 Warn,
241 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 pub fn delimiter_byte(&self) -> Result<u8, FaucetError> {
269 single_byte("delimiter", &self.delimiter)
270 }
271
272 pub fn quote_byte(&self) -> Result<u8, FaucetError> {
274 single_byte("quote", &self.quote)
275 }
276
277 pub fn flexible_or(&self, default: bool) -> bool {
279 self.flexible.unwrap_or(default)
280 }
281
282 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#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
305#[serde(deny_unknown_fields)]
306pub struct ExcelOptions {
307 #[serde(default, skip_serializing_if = "Option::is_none")]
309 pub sheet: Option<String>,
310 #[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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
331#[serde(deny_unknown_fields)]
332pub struct XmlOptions {
333 #[serde(default = "default_record_element")]
336 pub record_element: String,
337 #[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
353#[serde(rename_all = "snake_case")]
354pub enum AvroCodec {
355 #[default]
357 Null,
358 Deflate,
360 Snappy,
362 Zstd,
364}
365
366#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
368#[serde(deny_unknown_fields)]
369pub struct AvroOptions {
370 #[serde(default, skip_serializing_if = "Option::is_none")]
376 pub schema: Option<Value>,
377 #[serde(default)]
380 pub codec: AvroCodec,
381}
382
383#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
385#[serde(deny_unknown_fields)]
386pub struct OrcOptions {
387 #[serde(default, skip_serializing_if = "Option::is_none")]
390 pub columns: Option<Vec<String>>,
391}
392
393#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
400pub struct FormatOptions {
401 #[serde(default)]
403 pub csv: CsvOptions,
404 #[serde(default)]
406 pub excel: ExcelOptions,
407 #[serde(default)]
409 pub xml: XmlOptions,
410 #[serde(default)]
412 pub avro: AvroOptions,
413 #[serde(default)]
415 pub orc: OrcOptions,
416}
417
418pub 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
441pub 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
464fn 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#[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
587pub 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
612pub 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 #[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 #[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 #[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 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 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 let recs = vec![json!({"a": 1}), json!({"b": 2}), json!("not an object")];
780 assert_eq!(header_union(&recs), vec!["a", "b"]);
781 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}