use faucet_core::{FaucetError, Source, Stream, StreamPage};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::HashMap;
use std::pin::Pin;
pub const KIND: &str = "jsonl";
fn default_batch_size() -> usize {
faucet_core::DEFAULT_BATCH_SIZE
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct JsonLinesConfig {
pub path: String,
#[serde(default = "default_batch_size")]
pub batch_size: usize,
#[serde(default)]
pub limit: Option<usize>,
#[cfg(feature = "compression")]
#[serde(default)]
pub compression: faucet_core::CompressionConfig,
}
impl JsonLinesConfig {
pub fn new(path: impl Into<String>) -> Self {
Self {
path: path.into(),
batch_size: default_batch_size(),
limit: None,
#[cfg(feature = "compression")]
compression: faucet_core::CompressionConfig::Auto,
}
}
pub fn limit(mut self, limit: usize) -> Self {
self.limit = Some(limit);
self
}
pub fn batch_size(mut self, batch_size: usize) -> Self {
self.batch_size = batch_size;
self
}
}
pub struct JsonLinesSource {
config: JsonLinesConfig,
}
impl JsonLinesSource {
pub fn new(config: JsonLinesConfig) -> Self {
Self { config }
}
}
#[async_trait::async_trait]
impl Source for JsonLinesSource {
async fn fetch_with_context(
&self,
context: &HashMap<String, Value>,
) -> Result<Vec<Value>, FaucetError> {
use futures::StreamExt;
let mut out = Vec::new();
let mut pages = self.stream_pages(context, self.config.batch_size);
while let Some(page) = pages.next().await {
out.extend(page?.records);
}
Ok(out)
}
fn stream_pages<'a>(
&'a self,
_context: &'a HashMap<String, Value>,
_batch_size: usize,
) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>> {
let path = self.config.path.clone();
let chunk = if self.config.batch_size == 0 {
usize::MAX
} else {
self.config.batch_size
};
let limit = match self.config.limit {
None | Some(0) => usize::MAX,
Some(n) => n,
};
#[cfg(feature = "compression")]
let codec = self.config.compression.resolve(&path);
Box::pin(async_stream::try_stream! {
use tokio::io::AsyncBufReadExt;
let file = tokio::fs::File::open(&path).await.map_err(|e| {
FaucetError::Source(format!("failed to open JSON Lines file '{path}': {e}"))
})?;
let reader = tokio::io::BufReader::new(file);
#[cfg(feature = "compression")]
let reader = faucet_core::compression::wrap_async_reader(reader, codec);
let mut lines = reader.lines();
let mut batch: Vec<Value> = Vec::new();
let mut emitted = 0usize;
let mut lineno = 0usize;
while emitted < limit {
let Some(line) = lines.next_line().await.map_err(|e| {
FaucetError::Source(format!("reading '{path}': {e}"))
})? else {
break;
};
lineno += 1;
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
let value: Value = serde_json::from_str(trimmed).map_err(|e| {
FaucetError::Source(format!(
"malformed JSON at '{path}' line {lineno}: {e}"
))
})?;
batch.push(value);
emitted += 1;
if batch.len() >= chunk {
yield StreamPage { records: std::mem::take(&mut batch), bookmark: None };
}
}
if !batch.is_empty() {
yield StreamPage { records: batch, bookmark: None };
}
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use futures::StreamExt;
fn write(dir: &std::path::Path, name: &str, body: &str) -> String {
let p = dir.join(name);
std::fs::write(&p, body).unwrap();
p.to_string_lossy().to_string()
}
#[tokio::test]
async fn reads_one_record_per_line() {
let dir = tempfile::tempdir().unwrap();
let path = write(dir.path(), "out.jsonl", "{\"a\":1}\n{\"a\":2}\n{\"a\":3}\n");
let src = JsonLinesSource::new(JsonLinesConfig::new(path));
let rows = src.fetch_all().await.unwrap();
assert_eq!(rows.len(), 3);
assert_eq!(rows[2]["a"], 3);
}
#[tokio::test]
async fn blank_lines_are_skipped_not_records() {
let dir = tempfile::tempdir().unwrap();
let path = write(dir.path(), "out.jsonl", "{\"a\":1}\n\n \n{\"a\":2}\n");
let src = JsonLinesSource::new(JsonLinesConfig::new(path));
assert_eq!(src.fetch_all().await.unwrap().len(), 2);
}
#[tokio::test]
async fn an_empty_file_yields_no_pages() {
let dir = tempfile::tempdir().unwrap();
let path = write(dir.path(), "out.jsonl", "");
let src = JsonLinesSource::new(JsonLinesConfig::new(path));
assert!(src.fetch_all().await.unwrap().is_empty());
let ctx = HashMap::new();
let mut pages = src.stream_pages(&ctx, 0);
assert!(pages.next().await.is_none(), "no empty trailing page");
}
#[tokio::test]
async fn a_malformed_line_names_the_file_and_line() {
let dir = tempfile::tempdir().unwrap();
let path = write(dir.path(), "out.jsonl", "{\"a\":1}\n{\"a\":\n");
let src = JsonLinesSource::new(JsonLinesConfig::new(path));
let err = src.fetch_all().await.unwrap_err();
let msg = err.to_string();
assert!(msg.contains("line 2"), "{msg}");
assert!(msg.contains("out.jsonl"), "{msg}");
}
#[tokio::test]
async fn a_missing_file_is_a_typed_source_error() {
let src = JsonLinesSource::new(JsonLinesConfig::new("/definitely/not/here.jsonl"));
let err = src.fetch_all().await.unwrap_err();
assert!(matches!(err, FaucetError::Source(_)), "{err:?}");
assert!(err.to_string().contains("failed to open"), "{err}");
}
#[tokio::test]
async fn the_limit_stops_the_read_before_eof() {
let dir = tempfile::tempdir().unwrap();
let mut body = String::new();
for i in 0..1000 {
body.push_str(&format!("{{\"i\":{i}}}\n"));
}
body.push_str("this is not json\n");
let path = write(dir.path(), "big.jsonl", &body);
let src = JsonLinesSource::new(JsonLinesConfig::new(path).limit(3));
let rows = src.fetch_all().await.unwrap();
assert_eq!(rows.len(), 3);
assert_eq!(rows[0]["i"], 0);
}
#[tokio::test]
async fn a_zero_limit_means_no_limit_not_no_rows() {
let dir = tempfile::tempdir().unwrap();
let path = write(dir.path(), "out.jsonl", "{\"a\":1}\n{\"a\":2}\n");
let src = JsonLinesSource::new(JsonLinesConfig::new(path).limit(0));
assert_eq!(src.fetch_all().await.unwrap().len(), 2);
}
#[tokio::test]
async fn batch_size_controls_the_page_cadence() {
let dir = tempfile::tempdir().unwrap();
let path = write(
dir.path(),
"out.jsonl",
"{\"a\":1}\n{\"a\":2}\n{\"a\":3}\n{\"a\":4}\n{\"a\":5}\n",
);
let src = JsonLinesSource::new(JsonLinesConfig::new(path).batch_size(2));
let ctx = HashMap::new();
let mut pages = src.stream_pages(&ctx, 999);
let mut sizes = Vec::new();
while let Some(p) = pages.next().await {
sizes.push(p.unwrap().records.len());
}
assert_eq!(sizes, vec![2, 2, 1], "trailing partial page is emitted");
}
#[tokio::test]
async fn batch_size_zero_drains_into_one_page() {
let dir = tempfile::tempdir().unwrap();
let path = write(dir.path(), "out.jsonl", "{\"a\":1}\n{\"a\":2}\n");
let src = JsonLinesSource::new(JsonLinesConfig::new(path).batch_size(0));
let ctx = HashMap::new();
let mut pages = src.stream_pages(&ctx, 1);
let first = pages.next().await.unwrap().unwrap();
assert_eq!(first.records.len(), 2);
assert!(pages.next().await.is_none());
}
#[tokio::test]
async fn non_object_lines_are_read_as_values() {
let dir = tempfile::tempdir().unwrap();
let path = write(dir.path(), "out.jsonl", "1\n[2,3]\n\"x\"\n");
let src = JsonLinesSource::new(JsonLinesConfig::new(path));
let rows = src.fetch_all().await.unwrap();
assert_eq!(
rows,
vec![Value::from(1), serde_json::json!([2, 3]), Value::from("x")]
);
}
#[test]
fn config_round_trips_through_json() {
let cfg: JsonLinesConfig =
serde_json::from_value(serde_json::json!({"path": "/tmp/a.jsonl", "limit": 5}))
.unwrap();
assert_eq!(cfg.path, "/tmp/a.jsonl");
assert_eq!(cfg.limit, Some(5));
assert_eq!(cfg.batch_size, faucet_core::DEFAULT_BATCH_SIZE);
}
}