Skip to main content

faucet_core/file_format/
csv.rs

1//! CSV read/write for the file connectors (#604).
2//!
3//! The reader is the one the REST source has used since #497 — streaming
4//! RFC-4180 via `csv-async`, `flexible(true)`, header-derived keys with a
5//! `column_<i>` fallback, all-`String` values. Keeping those semantics exactly
6//! is the point of moving it here rather than writing a second one: a pipeline
7//! that reads a CSV through the REST source and one that reads the same file
8//! from S3 must produce the same records.
9
10use super::CsvOptions;
11use crate::error::FaucetError;
12use serde_json::{Map, Value};
13
14/// Parse CSV bytes into records (lenient about ragged rows).
15pub async fn decode(
16    bytes: &[u8],
17    delimiter: u8,
18    has_headers: bool,
19) -> Result<Vec<Value>, FaucetError> {
20    let opts = CsvOptions {
21        delimiter: (delimiter as char).to_string(),
22        has_headers,
23        ..CsvOptions::default()
24    };
25    let mut rdr = CsvRowReader::with_bytes(bytes, delimiter, &opts, true)?;
26    let mut out = Vec::new();
27    while let Some(r) = rdr.next_record().await? {
28        out.push(r);
29    }
30    Ok(out)
31}
32
33/// Parse CSV bytes into records with the full dialect in `opts`;
34/// `default_flexible` applies when `opts.flexible` is unset.
35pub async fn decode_with(
36    bytes: &[u8],
37    opts: &CsvOptions,
38    default_flexible: bool,
39) -> Result<Vec<Value>, FaucetError> {
40    let mut rdr = CsvRowReader::new(bytes, opts, default_flexible)?;
41    let mut out = Vec::new();
42    while let Some(r) = rdr.next_record().await? {
43        out.push(r);
44    }
45    Ok(out)
46}
47
48/// A streaming CSV reader yielding one JSON object per row.
49///
50/// Header names must be unique: rows are keyed by header, so a repeated name
51/// would silently drop a column. A ragged row fails naming its line unless
52/// the dialect is flexible.
53pub struct CsvRowReader<R> {
54    inner: csv_async::AsyncReader<R>,
55    record: csv_async::StringRecord,
56    headers: Option<Vec<String>>,
57    has_headers: bool,
58    flexible: bool,
59    null_values: Vec<String>,
60    line: usize,
61}
62
63impl<R: tokio::io::AsyncRead + Unpin + Send> CsvRowReader<R> {
64    /// Build a reader over `reader` with the dialect in `opts`.
65    pub fn new(reader: R, opts: &CsvOptions, default_flexible: bool) -> Result<Self, FaucetError> {
66        Self::with_bytes(reader, opts.delimiter_byte()?, opts, default_flexible)
67    }
68
69    fn with_bytes(
70        reader: R,
71        delimiter: u8,
72        opts: &CsvOptions,
73        default_flexible: bool,
74    ) -> Result<Self, FaucetError> {
75        let flexible = opts.flexible_or(default_flexible);
76        let inner = csv_async::AsyncReaderBuilder::new()
77            .has_headers(false)
78            .delimiter(delimiter)
79            .quote(opts.quote_byte()?)
80            .flexible(flexible)
81            .create_reader(reader);
82        Ok(Self {
83            inner,
84            record: csv_async::StringRecord::new(),
85            headers: None,
86            has_headers: opts.has_headers,
87            flexible,
88            null_values: opts.null_values.clone(),
89            line: 0,
90        })
91    }
92
93    /// The next row, or `None` at the end of the input.
94    pub async fn next_record(&mut self) -> Result<Option<Value>, FaucetError> {
95        loop {
96            self.line += 1;
97            let more = self
98                .inner
99                .read_record(&mut self.record)
100                .await
101                .map_err(|e| self.parse_error(e))?;
102            if !more {
103                return Ok(None);
104            }
105            if self.has_headers && self.headers.is_none() {
106                self.headers = Some(unique_headers(&self.record)?);
107                continue;
108            }
109            return Ok(Some(Value::Object(record_to_object(
110                &self.record,
111                self.headers.as_deref(),
112                &self.null_values,
113            ))));
114        }
115    }
116
117    fn parse_error(&self, e: csv_async::Error) -> FaucetError {
118        if !self.flexible && matches!(e.kind(), csv_async::ErrorKind::UnequalLengths { .. }) {
119            FaucetError::Source(format!(
120                "csv: ragged row at line {}: {e} — a short or long row is a structural defect \
121                 that would silently misalign fields; fix the file or set `csv.flexible: true` \
122                 to accept uneven rows",
123                self.line
124            ))
125        } else {
126            FaucetError::Source(format!("csv: parse error at line {}: {e}", self.line))
127        }
128    }
129}
130
131fn unique_headers(rec: &csv_async::StringRecord) -> Result<Vec<String>, FaucetError> {
132    let headers: Vec<String> = rec.iter().map(str::to_string).collect();
133    let mut seen = std::collections::HashMap::with_capacity(headers.len());
134    for (i, name) in headers.iter().enumerate() {
135        if let Some(first) = seen.insert(name.as_str(), i) {
136            let shown = if name.is_empty() { "(empty)" } else { name };
137            return Err(FaucetError::Source(format!(
138                "csv: duplicate header {shown} at columns {first} and {i}; rows are keyed by \
139                 header name, so a repeated header would silently drop a column — rename it \
140                 or set `csv.has_headers: false`"
141            )));
142        }
143    }
144    Ok(headers)
145}
146
147/// How one CSV record becomes a JSON object — the single definition, so the
148/// `Value` path and any byte path can never drift apart.
149fn record_to_object(
150    rec: &csv_async::StringRecord,
151    headers: Option<&[String]>,
152    null_values: &[String],
153) -> Map<String, Value> {
154    let mut obj = Map::new();
155    for (i, field) in rec.iter().enumerate() {
156        let key = headers
157            .and_then(|h| h.get(i).cloned())
158            .unwrap_or_else(|| format!("column_{i}"));
159        let value = if null_values.iter().any(|n| n == field) {
160            Value::Null
161        } else {
162            Value::String(field.to_string())
163        };
164        obj.insert(key, value);
165    }
166    obj
167}
168
169/// Write records as CSV.
170///
171/// Columns are the first-seen union of every record's keys
172/// ([`header_union`](super::header_union)), so a record that gains a field
173/// mid-page widens the file instead of losing the field. A record missing a
174/// column writes an empty cell.
175pub fn encode(records: &[Value], delimiter: u8, has_headers: bool) -> Result<Vec<u8>, FaucetError> {
176    encode_rows(records, delimiter, b'"', has_headers)
177}
178
179/// [`encode`] with the full dialect in `opts`.
180pub fn encode_with(records: &[Value], opts: &CsvOptions) -> Result<Vec<u8>, FaucetError> {
181    encode_rows(
182        records,
183        opts.delimiter_byte()?,
184        opts.quote_byte()?,
185        opts.has_headers,
186    )
187}
188
189fn encode_rows(
190    records: &[Value],
191    delimiter: u8,
192    quote: u8,
193    has_headers: bool,
194) -> Result<Vec<u8>, FaucetError> {
195    let headers = super::header_union(records);
196    let mut wtr = csv::WriterBuilder::new()
197        .delimiter(delimiter)
198        .quote(quote)
199        .from_writer(Vec::new());
200    if has_headers && !headers.is_empty() {
201        wtr.write_record(&headers)
202            .map_err(|e| FaucetError::Sink(format!("csv: writing header: {e}")))?;
203    }
204    for r in records {
205        let row: Vec<String> = headers
206            .iter()
207            .map(|h| r.get(h).map(super::cell_text).unwrap_or_default())
208            .collect();
209        wtr.write_record(&row)
210            .map_err(|e| FaucetError::Sink(format!("csv: writing row: {e}")))?;
211    }
212    wtr.into_inner()
213        .map_err(|e| FaucetError::Sink(format!("csv: finishing: {e}")))
214}
215
216#[cfg(test)]
217mod tests {
218    use super::*;
219    use serde_json::json;
220
221    #[tokio::test]
222    async fn headers_name_the_fields() {
223        let recs = decode(b"a,b\n1,2\n3,4\n", b',', true)
224            .await
225            .expect("decode");
226        assert_eq!(
227            recs,
228            vec![json!({"a": "1", "b": "2"}), json!({"a": "3", "b": "4"})]
229        );
230    }
231
232    #[tokio::test]
233    async fn without_headers_fields_fall_back_to_column_index() {
234        let recs = decode(b"1,2\n", b',', false).await.expect("decode");
235        assert_eq!(recs, vec![json!({"column_0": "1", "column_1": "2"})]);
236    }
237
238    #[tokio::test]
239    async fn a_quoted_embedded_newline_stays_one_record() {
240        let recs = decode(b"a\n\"x\ny\"\n", b',', true).await.expect("decode");
241        assert_eq!(recs, vec![json!({"a": "x\ny"})]);
242    }
243
244    #[tokio::test]
245    async fn a_ragged_row_is_kept_not_rejected() {
246        // `flexible(true)`: a short row yields the fields it has rather than
247        // failing the whole object, which is what the REST source does today.
248        let recs = decode(b"a,b\n1\n", b',', true).await.expect("decode");
249        assert_eq!(recs, vec![json!({"a": "1"})]);
250    }
251
252    #[tokio::test]
253    async fn header_only_and_empty_bodies_yield_nothing() {
254        assert!(
255            decode(b"a,b\n", b',', true)
256                .await
257                .expect("header")
258                .is_empty()
259        );
260        assert!(decode(b"", b',', true).await.expect("empty").is_empty());
261    }
262
263    #[test]
264    fn encode_uses_the_union_of_keys_and_blanks_the_missing_ones() {
265        let out = encode(
266            &[json!({"a": 1, "b": 2}), json!({"a": 3, "c": 4})],
267            b',',
268            true,
269        )
270        .expect("encode");
271        assert_eq!(String::from_utf8(out).unwrap(), "a,b,c\n1,2,\n3,,4\n");
272    }
273
274    #[test]
275    fn encode_can_omit_the_header() {
276        let out = encode(&[json!({"a": 1})], b',', false).expect("encode");
277        assert_eq!(String::from_utf8(out).unwrap(), "1\n");
278    }
279
280    #[tokio::test]
281    async fn a_tab_delimiter_round_trips() {
282        let out = encode(&[json!({"a": "x", "b": "y"})], b'\t', true).expect("encode");
283        assert_eq!(String::from_utf8(out.clone()).unwrap(), "a\tb\nx\ty\n");
284        let back = decode(&out, b'\t', true).await.expect("decode");
285        assert_eq!(back, vec![json!({"a": "x", "b": "y"})]);
286    }
287
288    #[test]
289    fn nested_structure_survives_as_json_rather_than_being_dropped() {
290        let out = encode(&[json!({"a": {"k": 1}})], b',', true).expect("encode");
291        assert_eq!(String::from_utf8(out).unwrap(), "a\n\"{\"\"k\"\":1}\"\n");
292    }
293    fn opts(f: impl FnOnce(&mut CsvOptions)) -> CsvOptions {
294        let mut o = CsvOptions::default();
295        f(&mut o);
296        o
297    }
298
299    #[tokio::test]
300    async fn a_custom_quote_character_is_honoured_on_read_and_write() {
301        let o = opts(|o| o.quote = "'".into());
302        let recs = decode_with(b"a,b\n'x,y',2\n", &o, false)
303            .await
304            .expect("decode");
305        assert_eq!(recs, vec![json!({"a": "x,y", "b": "2"})]);
306        let out = encode_with(&recs, &o).expect("encode");
307        assert_eq!(String::from_utf8(out).unwrap(), "a,b\n'x,y',2\n");
308    }
309
310    #[tokio::test]
311    async fn a_strict_dialect_rejects_a_ragged_row_naming_its_line() {
312        let o = opts(|o| o.flexible = Some(false));
313        let err = decode_with(b"a,b\n1,2\n3\n", &o, true).await.unwrap_err();
314        let msg = err.to_string();
315        assert!(msg.contains("ragged row at line 3"), "{msg}");
316        assert!(msg.contains("csv.flexible"), "{msg}");
317        let lenient = decode_with(b"a,b\n1,2\n3\n", &CsvOptions::default(), false)
318            .await
319            .unwrap_err();
320        assert!(lenient.to_string().contains("ragged"));
321        let ok = decode_with(b"a,b\n3\n", &CsvOptions::default(), true)
322            .await
323            .expect("default flexible");
324        assert_eq!(ok, vec![json!({"a": "3"})]);
325    }
326
327    #[tokio::test]
328    async fn null_values_read_as_null() {
329        let o = opts(|o| o.null_values = vec!["".into(), "NULL".into()]);
330        let recs = decode_with(b"a,b,c\n,NULL,x\n", &o, false)
331            .await
332            .expect("decode");
333        assert_eq!(recs, vec![json!({"a": null, "b": null, "c": "x"})]);
334    }
335
336    #[tokio::test]
337    async fn a_duplicate_header_is_refused_rather_than_dropping_a_column() {
338        let err = decode(b"a,a\n1,2\n", b',', true).await.unwrap_err();
339        let msg = err.to_string();
340        assert!(
341            msg.contains("duplicate header a at columns 0 and 1"),
342            "{msg}"
343        );
344        let err = decode(b",\n1,2\n", b',', true).await.unwrap_err();
345        assert!(err.to_string().contains("(empty)"));
346        let ok = decode(b"a,a\n1,2\n", b',', false)
347            .await
348            .expect("no headers");
349        assert_eq!(
350            ok,
351            vec![
352                json!({"column_0": "a", "column_1": "a"}),
353                json!({"column_0": "1", "column_1": "2"})
354            ]
355        );
356    }
357
358    #[tokio::test]
359    async fn a_malformed_quote_is_a_typed_parse_error() {
360        let o = opts(|o| o.flexible = Some(false));
361        let bad: &[u8] = b"a\n\xff\xfe\n";
362        let err = decode_with(bad, &o, false).await.unwrap_err();
363        assert!(err.to_string().contains("parse error at line"), "{err}");
364    }
365
366    #[test]
367    fn dialect_bytes_are_validated() {
368        assert!(opts(|o| o.quote = "''".into()).validate().is_err());
369        assert!(opts(|o| o.delimiter = ";;".into()).validate().is_err());
370        assert!(opts(|o| o.delimiter = "\\t".into()).validate().is_ok());
371        assert!(encode_with(&[json!({"a": 1})], &opts(|o| o.quote = "".into())).is_err());
372    }
373
374    #[test]
375    fn write_headers_is_another_name_for_has_headers() {
376        let o: CsvOptions = serde_json::from_value(json!({"write_headers": false})).unwrap();
377        assert!(!o.has_headers);
378        assert_eq!(o.on_unknown_field, super::super::CsvUnknownField::Widen);
379    }
380}