1use super::CsvOptions;
11use crate::error::FaucetError;
12use serde_json::{Map, Value};
13
14pub 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
33pub 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
48pub 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 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 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
147fn 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
169pub fn encode(records: &[Value], delimiter: u8, has_headers: bool) -> Result<Vec<u8>, FaucetError> {
176 encode_rows(records, delimiter, b'"', has_headers)
177}
178
179pub 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 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}