use crate::error::FaucetError;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::Value;
#[cfg(feature = "file-format-avro")]
pub mod avro;
pub mod container;
#[cfg(feature = "file-format-csv")]
pub mod csv;
#[cfg(feature = "file-format-excel")]
pub mod excel;
pub mod json;
#[cfg(feature = "file-format-orc")]
pub mod orc;
pub mod parquet_io;
#[cfg(feature = "file-format-xml")]
pub mod xml;
pub use container::{ContainerDecoder, FileInput};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum FileFormat {
#[default]
JsonLines,
JsonArray,
Csv,
Xml,
Xlsx,
Parquet,
RawText,
Avro,
Orc,
}
impl FileFormat {
pub fn extension(self) -> &'static str {
match self {
Self::JsonLines => ".jsonl",
Self::JsonArray => ".json",
Self::Csv => ".csv",
Self::Xml => ".xml",
Self::Xlsx => ".xlsx",
Self::Parquet => ".parquet",
Self::RawText => ".txt",
Self::Avro => ".avro",
Self::Orc => ".orc",
}
}
pub fn requires_whole_object(self) -> bool {
matches!(self, Self::JsonArray | Self::Xml | Self::Xlsx | Self::Orc)
}
pub fn is_container(self) -> bool {
matches!(self, Self::Avro | Self::Orc)
}
pub fn is_writable(self) -> bool {
!matches!(self, Self::Orc | Self::Parquet)
}
pub fn from_path(path: &str) -> Option<Self> {
let name = path
.rsplit(['/', '\\'])
.next()
.unwrap_or(path)
.to_ascii_lowercase();
let stem = ["gz", "gzip", "zst", "zstd"]
.iter()
.find_map(|c| name.strip_suffix(&format!(".{c}")))
.unwrap_or(&name);
let ext = stem.rsplit_once('.')?.1;
Some(match ext {
"jsonl" | "ndjson" => Self::JsonLines,
"json" => Self::JsonArray,
"csv" => Self::Csv,
"xml" => Self::Xml,
"xlsx" => Self::Xlsx,
"parquet" => Self::Parquet,
"txt" => Self::RawText,
"avro" => Self::Avro,
"orc" => Self::Orc,
_ => return None,
})
}
pub fn as_str(self) -> &'static str {
match self {
Self::JsonLines => "json_lines",
Self::JsonArray => "json_array",
Self::Csv => "csv",
Self::Xml => "xml",
Self::Xlsx => "xlsx",
Self::Parquet => "parquet",
Self::RawText => "raw_text",
Self::Avro => "avro",
Self::Orc => "orc",
}
}
}
fn default_true() -> bool {
true
}
fn default_delimiter() -> String {
",".into()
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct CsvOptions {
#[serde(default = "default_delimiter")]
pub delimiter: String,
#[serde(default = "default_true")]
pub has_headers: bool,
#[serde(default = "default_quote")]
pub quote: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub flexible: Option<bool>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub null_values: Vec<String>,
#[serde(default)]
pub on_unknown_field: CsvUnknownField,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum CsvUnknownField {
#[default]
Widen,
Warn,
Error,
}
fn default_quote() -> String {
"\"".into()
}
impl Default for CsvOptions {
fn default() -> Self {
Self {
delimiter: default_delimiter(),
has_headers: true,
quote: default_quote(),
flexible: None,
null_values: Vec::new(),
on_unknown_field: CsvUnknownField::Widen,
}
}
}
impl CsvOptions {
pub fn delimiter_byte(&self) -> Result<u8, FaucetError> {
single_byte("delimiter", &self.delimiter)
}
pub fn quote_byte(&self) -> Result<u8, FaucetError> {
single_byte("quote", &self.quote)
}
pub fn flexible_or(&self, default: bool) -> bool {
self.flexible.unwrap_or(default)
}
pub fn validate(&self) -> Result<(), FaucetError> {
self.delimiter_byte()?;
self.quote_byte()?;
Ok(())
}
}
fn single_byte(name: &str, value: &str) -> Result<u8, FaucetError> {
let d = match value {
"\\t" => "\t",
other => other,
};
match d.as_bytes() {
[b] => Ok(*b),
_ => Err(FaucetError::Config(format!(
"csv.{name} must be exactly one byte, got {value:?}"
))),
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct ExcelOptions {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub sheet: Option<String>,
#[serde(default)]
pub header_row: usize,
}
fn default_record_element() -> String {
"record".into()
}
fn default_root_element() -> String {
"records".into()
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct XmlOptions {
#[serde(default = "default_record_element")]
pub record_element: String,
#[serde(default = "default_root_element")]
pub root_element: String,
}
impl Default for XmlOptions {
fn default() -> Self {
Self {
record_element: default_record_element(),
root_element: default_root_element(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum AvroCodec {
#[default]
Null,
Deflate,
Snappy,
Zstd,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct AvroOptions {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub schema: Option<Value>,
#[serde(default)]
pub codec: AvroCodec,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct OrcOptions {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub columns: Option<Vec<String>>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
pub struct FormatOptions {
#[serde(default)]
pub csv: CsvOptions,
#[serde(default)]
pub excel: ExcelOptions,
#[serde(default)]
pub xml: XmlOptions,
#[serde(default)]
pub avro: AvroOptions,
#[serde(default)]
pub orc: OrcOptions,
}
pub async fn decode(
bytes: &[u8],
format: FileFormat,
opts: &FormatOptions,
) -> Result<Vec<Value>, FaucetError> {
match format {
FileFormat::JsonLines => json::decode_lines(bytes),
FileFormat::JsonArray => json::decode_array(bytes),
FileFormat::RawText => json::decode_raw_text(bytes),
FileFormat::Csv => decode_csv(bytes, opts).await,
FileFormat::Xml => decode_xml(bytes, opts),
FileFormat::Xlsx => decode_xlsx(bytes, opts),
FileFormat::Parquet => Err(columnar_refusal("decode")),
FileFormat::Avro => decode_avro(bytes, opts),
FileFormat::Orc => decode_orc(bytes, opts),
}
}
pub fn encode(
records: &[Value],
format: FileFormat,
opts: &FormatOptions,
) -> Result<Vec<u8>, FaucetError> {
match format {
FileFormat::JsonLines => json::encode_lines(records),
FileFormat::JsonArray => json::encode_array(records),
FileFormat::RawText => json::encode_raw_text(records),
FileFormat::Csv => encode_csv(records, opts),
FileFormat::Xml => encode_xml(records, opts),
FileFormat::Xlsx => encode_xlsx(records, opts),
FileFormat::Parquet => Err(columnar_refusal("encode")),
FileFormat::Avro => encode_avro(records, opts),
FileFormat::Orc => Err(FaucetError::Config(
"`format: orc` is read-only — there is no ORC writer; write Parquet for a columnar \
output"
.into(),
)),
}
}
fn columnar_refusal(verb: &str) -> FaucetError {
FaucetError::Config(format!(
"file_format::{verb}: `parquet` is columnar and is handled by each connector's Arrow \
path, not the generic record helper"
))
}
#[cfg_attr(
all(
feature = "file-format-csv",
feature = "file-format-xml",
feature = "file-format-excel",
feature = "file-format-avro",
feature = "file-format-orc"
),
allow(dead_code)
)]
fn missing_feature(format: FileFormat, feature: &str) -> FaucetError {
FaucetError::Config(format!(
"`format: {}` requires the `{feature}` build feature — rebuild with \
`--features {feature}`",
format.as_str()
))
}
#[cfg(feature = "file-format-csv")]
async fn decode_csv(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
csv::decode_with(bytes, &opts.csv, true).await
}
#[cfg(not(feature = "file-format-csv"))]
async fn decode_csv(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
Err(missing_feature(FileFormat::Csv, "file-format-csv"))
}
#[cfg(feature = "file-format-csv")]
fn encode_csv(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
csv::encode_with(records, &opts.csv)
}
#[cfg(not(feature = "file-format-csv"))]
fn encode_csv(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
Err(missing_feature(FileFormat::Csv, "file-format-csv"))
}
#[cfg(feature = "file-format-xml")]
fn decode_xml(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
xml::decode(bytes, &opts.xml.record_element)
}
#[cfg(not(feature = "file-format-xml"))]
fn decode_xml(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
Err(missing_feature(FileFormat::Xml, "file-format-xml"))
}
#[cfg(feature = "file-format-xml")]
fn encode_xml(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
xml::encode(records, &opts.xml.root_element, &opts.xml.record_element)
}
#[cfg(not(feature = "file-format-xml"))]
fn encode_xml(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
Err(missing_feature(FileFormat::Xml, "file-format-xml"))
}
#[cfg(feature = "file-format-excel")]
fn decode_xlsx(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
excel::decode(bytes, opts.excel.sheet.as_deref(), opts.excel.header_row)
}
#[cfg(not(feature = "file-format-excel"))]
fn decode_xlsx(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
Err(missing_feature(FileFormat::Xlsx, "file-format-excel"))
}
#[cfg(feature = "file-format-excel")]
fn encode_xlsx(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
excel::encode(records, opts.excel.sheet.as_deref())
}
#[cfg(not(feature = "file-format-excel"))]
fn encode_xlsx(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
Err(missing_feature(FileFormat::Xlsx, "file-format-excel"))
}
#[cfg(feature = "file-format-avro")]
fn decode_avro(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
avro::decode(bytes, &opts.avro)
}
#[cfg(not(feature = "file-format-avro"))]
fn decode_avro(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
Err(missing_feature(FileFormat::Avro, "file-format-avro"))
}
#[cfg(feature = "file-format-avro")]
fn encode_avro(records: &[Value], opts: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
avro::encode(records, &opts.avro)
}
#[cfg(not(feature = "file-format-avro"))]
fn encode_avro(_: &[Value], _: &FormatOptions) -> Result<Vec<u8>, FaucetError> {
Err(missing_feature(FileFormat::Avro, "file-format-avro"))
}
#[cfg(feature = "file-format-orc")]
fn decode_orc(bytes: &[u8], opts: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
orc::decode(bytes, &opts.orc)
}
#[cfg(not(feature = "file-format-orc"))]
fn decode_orc(_: &[u8], _: &FormatOptions) -> Result<Vec<Value>, FaucetError> {
Err(missing_feature(FileFormat::Orc, "file-format-orc"))
}
pub fn header_union(records: &[Value]) -> Vec<String> {
let mut seen = std::collections::HashSet::new();
let mut out = Vec::new();
for r in records {
if let Value::Object(map) = r {
for k in map.keys() {
if seen.insert(k.clone()) {
out.push(k.clone());
}
}
}
}
out
}
pub fn cell_text(v: &Value) -> String {
match v {
Value::Null => String::new(),
Value::String(s) => s.clone(),
Value::Bool(b) => b.to_string(),
Value::Number(n) => n.to_string(),
other => serde_json::to_string(other).unwrap_or_default(),
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn every_variant_has_an_extension_and_a_wire_name() {
let all = [
(FileFormat::JsonLines, ".jsonl", "json_lines"),
(FileFormat::JsonArray, ".json", "json_array"),
(FileFormat::Csv, ".csv", "csv"),
(FileFormat::Xml, ".xml", "xml"),
(FileFormat::Xlsx, ".xlsx", "xlsx"),
(FileFormat::Parquet, ".parquet", "parquet"),
(FileFormat::RawText, ".txt", "raw_text"),
(FileFormat::Avro, ".avro", "avro"),
(FileFormat::Orc, ".orc", "orc"),
];
for (f, ext, name) in all {
assert_eq!(f.extension(), ext, "{f:?}");
assert_eq!(f.as_str(), name, "{f:?}");
}
}
#[tokio::test]
async fn raw_text_round_trips_through_the_top_level_dispatch() {
let opts = FormatOptions::default();
let recs = decode(b"hello", FileFormat::RawText, &opts)
.await
.expect("decode");
assert_eq!(recs, vec![json!({"text": "hello"})]);
let bytes = encode(&recs, FileFormat::RawText, &opts).expect("encode");
assert_eq!(bytes, b"hello\n");
}
#[test]
fn a_missing_feature_is_named_not_guessed() {
let err = missing_feature(FileFormat::Xlsx, "file-format-excel");
let msg = err.to_string();
assert!(msg.contains("xlsx"), "{msg}");
assert!(msg.contains("file-format-excel"), "{msg}");
}
#[test]
fn extensions_and_names_are_stable() {
assert_eq!(FileFormat::default(), FileFormat::JsonLines);
assert_eq!(FileFormat::Csv.extension(), ".csv");
assert_eq!(FileFormat::Xlsx.as_str(), "xlsx");
}
#[test]
fn whole_object_formats_are_the_ones_without_a_record_boundary() {
for f in [
FileFormat::JsonArray,
FileFormat::Xml,
FileFormat::Xlsx,
FileFormat::Orc,
] {
assert!(f.requires_whole_object(), "{f:?}");
}
for f in [
FileFormat::JsonLines,
FileFormat::Csv,
FileFormat::RawText,
FileFormat::Parquet,
FileFormat::Avro,
] {
assert!(!f.requires_whole_object(), "{f:?}");
}
}
#[test]
fn extensions_resolve_through_compression_suffixes() {
let cases = [
("a.jsonl", Some(FileFormat::JsonLines)),
("dir/b.NDJSON", Some(FileFormat::JsonLines)),
("c.json.gz", Some(FileFormat::JsonArray)),
("d.csv.gz", Some(FileFormat::Csv)),
("e.xml.zst", Some(FileFormat::Xml)),
("f.xlsx", Some(FileFormat::Xlsx)),
("g.parquet", Some(FileFormat::Parquet)),
("h.txt", Some(FileFormat::RawText)),
("i.avro", Some(FileFormat::Avro)),
("j.orc", Some(FileFormat::Orc)),
("k.csv.gzip", Some(FileFormat::Csv)),
("l.jsonl.zstd", Some(FileFormat::JsonLines)),
("m.tsv", None),
("noext", None),
("n.gz", None),
];
for (path, want) in cases {
assert_eq!(FileFormat::from_path(path), want, "{path}");
}
}
#[test]
fn container_and_writable_classification() {
assert!(FileFormat::Avro.is_container() && FileFormat::Orc.is_container());
assert!(!FileFormat::Csv.is_container());
assert!(FileFormat::Avro.is_writable());
assert!(!FileFormat::Orc.is_writable() && !FileFormat::Parquet.is_writable());
}
#[tokio::test]
async fn orc_is_refused_for_writing() {
let err = encode(&[], FileFormat::Orc, &FormatOptions::default()).expect_err("orc");
assert!(err.to_string().contains("read-only"), "{err}");
}
#[tokio::test]
async fn parquet_is_refused_on_both_directions() {
let err = decode(b"", FileFormat::Parquet, &FormatOptions::default())
.await
.expect_err("decode refuses parquet");
assert!(err.to_string().contains("columnar"), "{err}");
let err = encode(&[], FileFormat::Parquet, &FormatOptions::default())
.expect_err("encode refuses parquet");
assert!(err.to_string().contains("columnar"), "{err}");
}
#[test]
fn delimiter_must_be_one_byte() {
let mut o = CsvOptions::default();
assert_eq!(o.delimiter_byte().expect("default"), b',');
o.delimiter = "\\t".into();
assert_eq!(o.delimiter_byte().expect("tab"), b'\t');
o.delimiter = "||".into();
let err = o.delimiter_byte().expect_err("multi-byte");
assert!(err.to_string().contains("exactly one byte"), "{err}");
o.delimiter = "§".into();
assert!(o.delimiter_byte().is_err());
}
#[test]
fn header_union_widens_across_records_and_ignores_non_objects() {
let recs = vec![json!({"a": 1}), json!({"b": 2}), json!("not an object")];
assert_eq!(header_union(&recs), vec!["a", "b"]);
assert_eq!(header_union(&[json!({"a": 1}), json!({"a": 2})]), vec!["a"]);
assert!(header_union(&[json!([1, 2])]).is_empty());
}
#[test]
fn cells_render_scalars_plainly_and_structure_as_json() {
assert_eq!(cell_text(&Value::Null), "");
assert_eq!(cell_text(&json!("x")), "x");
assert_eq!(cell_text(&json!(3)), "3");
assert_eq!(cell_text(&json!(true)), "true");
assert_eq!(cell_text(&json!({"k": 1})), r#"{"k":1}"#);
}
}