use crate::auth_catalog::AuthCatalog;
use crate::error::{CliError, CliResult};
use faucet_core::Source;
use futures::StreamExt;
use serde::Serialize;
use serde_json::Value;
use std::collections::HashMap;
use std::time::{Duration, Instant};
use super::RowCap;
pub const PREVIEW_DEADLINE: Duration = Duration::from_secs(30);
pub const PREVIEW_HARD_TIMEOUT: Duration = Duration::from_secs(60);
pub const PREVIEW_MAX_BYTES: usize = 64 * 1024 * 1024;
pub const PREVIEW_UNLIMITED_PAGE_ROWS: usize = faucet_core::DEFAULT_BATCH_SIZE;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Bounds {
pub max_bytes: usize,
pub deadline: Duration,
pub hard_timeout: Duration,
}
impl Default for Bounds {
fn default() -> Self {
Self {
max_bytes: PREVIEW_MAX_BYTES,
deadline: PREVIEW_DEADLINE,
hard_timeout: PREVIEW_HARD_TIMEOUT,
}
}
}
#[derive(Debug, Clone)]
pub struct PreviewRequest {
pub kind: String,
pub config: Value,
pub rows: RowCap,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum Capped {
Rows,
Bytes,
Time,
}
impl Capped {
pub fn as_str(self) -> &'static str {
match self {
Self::Rows => "rows",
Self::Bytes => "bytes",
Self::Time => "time",
}
}
}
#[derive(Debug, Clone)]
pub struct PreviewPage {
pub rows: Vec<Value>,
pub columns: Vec<String>,
pub capped_by: Option<Capped>,
pub pages_read: usize,
pub elapsed_ms: u64,
}
impl PreviewPage {
pub fn truncated(&self) -> bool {
self.capped_by.is_some()
}
}
pub async fn read_capped(req: &PreviewRequest, auth: &AuthCatalog) -> CliResult<PreviewPage> {
read_capped_with(req, auth, Bounds::default()).await
}
pub async fn read_capped_with(
req: &PreviewRequest,
auth: &AuthCatalog,
bounds: Bounds,
) -> CliResult<PreviewPage> {
let started = Instant::now();
let page = tokio::time::timeout(bounds.hard_timeout, read_inner(req, auth, started, bounds))
.await
.map_err(|_| {
CliError::Serve(format!(
"preview abandoned after {:?} reading a `{}` source — a single page never \
returned",
bounds.hard_timeout, req.kind
))
})??;
Ok(PreviewPage {
elapsed_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
..page
})
}
async fn read_inner(
req: &PreviewRequest,
auth: &AuthCatalog,
started: Instant,
bounds: Bounds,
) -> CliResult<PreviewPage> {
let want = match req.rows {
RowCap::Rows(n) => Some(n.max(1).saturating_add(1)),
RowCap::Unlimited => None,
};
let source = build_source(&req.kind, req.config.clone(), auth).await?;
let mut rows: Vec<Value> = Vec::new();
if let Some(want) = want {
rows.reserve(want.min(4096));
}
let mut pages_read = 0usize;
let mut bytes = 0usize;
let mut capped_by = None;
{
let context = HashMap::new();
let hint = want.unwrap_or(PREVIEW_UNLIMITED_PAGE_ROWS);
let mut pages = source.stream_pages(&context, hint);
'outer: while let Some(page) = pages.next().await {
let page = page?;
pages_read += 1;
for record in page.records {
if want.is_some_and(|want| rows.len() >= want) {
break 'outer;
}
if bytes >= bounds.max_bytes {
capped_by = Some(Capped::Bytes);
break 'outer;
}
bytes += approx_bytes(&record);
rows.push(record);
}
if started.elapsed() >= bounds.deadline {
capped_by = Some(Capped::Time);
break;
}
}
}
if let Some(want) = want.filter(|want| rows.len() >= *want) {
capped_by = Some(Capped::Rows);
rows.truncate(want.saturating_sub(1));
}
Ok(PreviewPage {
columns: columns_of(&rows),
rows,
capped_by,
pages_read,
elapsed_ms: 0,
})
}
async fn build_source(kind: &str, config: Value, auth: &AuthCatalog) -> CliResult<Box<dyn Source>> {
if kind == super::jsonl::KIND {
let cfg: super::jsonl::JsonLinesConfig = serde_json::from_value(config)
.map_err(|e| CliError::Config(format!("invalid jsonl preview config: {e}")))?;
return Ok(Box::new(super::jsonl::JsonLinesSource::new(cfg)));
}
crate::registry::build_source(kind, config, auth, None).await
}
fn approx_bytes(value: &Value) -> usize {
match value {
Value::Null => 4,
Value::Bool(_) => 5,
Value::Number(_) => 8,
Value::String(s) => s.len() + 2,
Value::Array(items) => 2 + items.len() + items.iter().map(approx_bytes).sum::<usize>(),
Value::Object(map) => {
2 + map
.iter()
.map(|(k, v)| k.len() + 4 + approx_bytes(v))
.sum::<usize>()
}
}
}
fn columns_of(rows: &[Value]) -> Vec<String> {
let mut seen = std::collections::HashSet::new();
let mut columns = Vec::new();
for row in rows {
if let Value::Object(map) = row {
for key in map.keys() {
if seen.insert(key.as_str()) {
columns.push(key.clone());
}
}
}
}
columns
}
#[cfg(test)]
mod tests {
use super::*;
fn write_jsonl(dir: &std::path::Path, rows: usize) -> String {
let mut body = String::new();
for i in 0..rows {
body.push_str(&format!("{{\"i\":{i},\"name\":\"r{i}\"}}\n"));
}
let p = dir.join("out.jsonl");
std::fs::write(&p, body).unwrap();
p.to_string_lossy().to_string()
}
fn request(path: &str, rows: RowCap) -> PreviewRequest {
PreviewRequest {
kind: "jsonl".into(),
config: serde_json::json!({ "path": path, "batch_size": 2 }),
rows,
}
}
#[tokio::test]
async fn reads_up_to_the_cap_and_reports_truncation() {
let dir = tempfile::tempdir().unwrap();
let path = write_jsonl(dir.path(), 50);
let page = read_capped(&request(&path, RowCap::Rows(10)), &AuthCatalog::new())
.await
.unwrap();
assert_eq!(page.rows.len(), 10);
assert!(page.truncated(), "50 rows behind a cap of 10");
assert_eq!(page.capped_by, Some(Capped::Rows));
assert_eq!(page.columns, vec!["i".to_string(), "name".to_string()]);
assert!(page.pages_read <= 6, "read did not stop early: {page:?}");
}
#[tokio::test]
async fn a_file_shorter_than_the_cap_is_not_truncated() {
let dir = tempfile::tempdir().unwrap();
let path = write_jsonl(dir.path(), 3);
let page = read_capped(&request(&path, RowCap::Rows(100)), &AuthCatalog::new())
.await
.unwrap();
assert_eq!(page.rows.len(), 3);
assert!(!page.truncated());
assert_eq!(page.capped_by, None);
}
#[tokio::test]
async fn a_file_exactly_the_cap_is_not_truncated() {
let dir = tempfile::tempdir().unwrap();
let path = write_jsonl(dir.path(), 10);
let page = read_capped(&request(&path, RowCap::Rows(10)), &AuthCatalog::new())
.await
.unwrap();
assert_eq!(page.rows.len(), 10);
assert_eq!(
page.capped_by, None,
"exactly `cap` rows is a complete read"
);
}
#[tokio::test]
async fn a_file_of_exactly_one_more_row_than_the_cap_is_truncated() {
let dir = tempfile::tempdir().unwrap();
let path = write_jsonl(dir.path(), 11);
let req = PreviewRequest {
kind: "jsonl".into(),
config: serde_json::json!({ "path": path, "batch_size": 11, "limit": 11 }),
rows: RowCap::Rows(10),
};
let page = read_capped(&req, &AuthCatalog::new()).await.unwrap();
assert_eq!(page.rows.len(), 10, "the surplus row is evidence, not data");
assert_eq!(page.capped_by, Some(Capped::Rows));
assert!(page.truncated());
}
#[tokio::test]
async fn an_unlimited_read_returns_the_whole_dataset() {
let dir = tempfile::tempdir().unwrap();
let path = write_jsonl(dir.path(), 2_500);
let page = read_capped(&request(&path, RowCap::Unlimited), &AuthCatalog::new())
.await
.unwrap();
assert_eq!(page.rows.len(), 2_500);
assert_eq!(
page.capped_by, None,
"nothing was left behind, so nothing should claim it was"
);
assert_eq!(page.rows[2_499]["i"], 2_499);
}
#[tokio::test]
async fn an_empty_file_previews_as_zero_rows_not_an_error() {
let dir = tempfile::tempdir().unwrap();
let path = write_jsonl(dir.path(), 0);
let page = read_capped(&request(&path, RowCap::Rows(10)), &AuthCatalog::new())
.await
.unwrap();
assert!(page.rows.is_empty());
assert!(page.columns.is_empty());
assert!(!page.truncated());
}
#[tokio::test]
async fn a_one_row_cap_still_returns_a_row() {
let dir = tempfile::tempdir().unwrap();
let path = write_jsonl(dir.path(), 5);
let page = read_capped(&request(&path, RowCap::Rows(1)), &AuthCatalog::new())
.await
.unwrap();
assert_eq!(page.rows.len(), 1);
assert_eq!(page.capped_by, Some(Capped::Rows));
}
#[tokio::test]
async fn a_source_error_propagates_rather_than_returning_an_empty_page() {
let req = request("/definitely/not/here.jsonl", RowCap::Rows(10));
let err = read_capped(&req, &AuthCatalog::new()).await.unwrap_err();
assert!(err.to_string().contains("failed to open"), "{err}");
}
#[tokio::test]
async fn an_unknown_kind_is_an_error_not_a_panic() {
let req = PreviewRequest {
kind: "not-a-connector".into(),
config: serde_json::json!({}),
rows: RowCap::Rows(10),
};
assert!(read_capped(&req, &AuthCatalog::new()).await.is_err());
}
#[tokio::test]
async fn a_bad_jsonl_config_is_a_config_error() {
let req = PreviewRequest {
kind: "jsonl".into(),
config: serde_json::json!({ "nope": 1 }),
rows: RowCap::Rows(10),
};
let err = read_capped(&req, &AuthCatalog::new()).await.unwrap_err();
assert!(matches!(err, CliError::Config(_)), "{err:?}");
}
#[test]
fn columns_union_keeps_first_seen_order_and_ragged_keys() {
let rows = vec![
serde_json::json!({"a": 1, "b": 2}),
serde_json::json!({"a": 3, "z": 4}),
];
assert_eq!(columns_of(&rows), vec!["a", "b", "z"]);
}
#[test]
fn columns_are_empty_for_non_object_records() {
let rows = vec![serde_json::json!(1), serde_json::json!([1, 2])];
assert!(columns_of(&rows).is_empty());
}
#[tokio::test]
async fn non_object_records_survive_with_no_columns() {
let dir = tempfile::tempdir().unwrap();
let p = dir.path().join("scalars.jsonl");
std::fs::write(&p, "1\n2\n").unwrap();
let req = PreviewRequest {
kind: "jsonl".into(),
config: serde_json::json!({ "path": p.to_string_lossy() }),
rows: RowCap::Rows(10),
};
let page = read_capped(&req, &AuthCatalog::new()).await.unwrap();
assert_eq!(page.rows.len(), 2);
assert!(page.columns.is_empty());
}
#[test]
fn approx_bytes_scales_with_the_record() {
let small = serde_json::json!({"a": 1});
let big = serde_json::json!({"a": "x".repeat(10_000)});
assert!(approx_bytes(&small) < 40);
assert!(approx_bytes(&big) > 10_000);
let nested = serde_json::json!({"a": [{"b": "yy"}, {"b": "zz"}]});
assert!(approx_bytes(&nested) > approx_bytes(&small));
}
fn unlimited_request(path: &str) -> PreviewRequest {
PreviewRequest {
kind: "jsonl".into(),
config: serde_json::json!({
"path": path,
"batch_size": PREVIEW_UNLIMITED_PAGE_ROWS,
"limit": 0,
}),
rows: RowCap::Unlimited,
}
}
#[tokio::test]
async fn the_byte_budget_bounds_an_unlimited_read() {
let dir = tempfile::tempdir().unwrap();
let p = dir.path().join("fat.jsonl");
let blob = "x".repeat(4096);
let body: String = (0..500)
.map(|i| format!("{{\"i\":{i},\"blob\":\"{blob}\"}}\n"))
.collect();
std::fs::write(&p, body).unwrap();
let bounds = Bounds {
max_bytes: 64 * 1024,
..Bounds::default()
};
let page = read_capped_with(
&unlimited_request(&p.to_string_lossy()),
&AuthCatalog::new(),
bounds,
)
.await
.unwrap();
assert_eq!(page.capped_by, Some(Capped::Bytes), "{:?}", page.capped_by);
assert!(page.rows.len() < 500, "the read must stop before EOF");
assert!(!page.rows.is_empty(), "…but still return what it read");
assert!(page.truncated());
}
#[tokio::test]
async fn an_unbounded_page_would_defeat_the_byte_budget() {
let dir = tempfile::tempdir().unwrap();
let p = dir.path().join("all-at-once.jsonl");
let body: String = (0..500).map(|i| format!("{{\"i\":{i}}}\n")).collect();
std::fs::write(&p, body).unwrap();
let req = PreviewRequest {
kind: "jsonl".into(),
config: serde_json::json!({ "path": p.to_string_lossy(), "batch_size": 0 }),
rows: RowCap::Unlimited,
};
let page = read_capped_with(
&req,
&AuthCatalog::new(),
Bounds {
max_bytes: 64,
..Bounds::default()
},
)
.await
.unwrap();
assert_eq!(page.pages_read, 1, "one page: the whole file");
assert!(page.rows.len() < 500);
}
#[tokio::test]
async fn the_deadline_returns_a_partial_answer_rather_than_an_error() {
let dir = tempfile::tempdir().unwrap();
let path = write_jsonl(dir.path(), 5_000);
let bounds = Bounds {
deadline: Duration::ZERO,
..Bounds::default()
};
let page = read_capped_with(&unlimited_request(&path), &AuthCatalog::new(), bounds)
.await
.unwrap();
assert_eq!(page.capped_by, Some(Capped::Time));
assert!(page.truncated());
assert_eq!(page.rows.len(), PREVIEW_UNLIMITED_PAGE_ROWS);
}
#[test]
fn the_unlimited_page_size_is_never_the_drain_everything_sentinel() {
assert_ne!(PREVIEW_UNLIMITED_PAGE_ROWS, 0);
}
#[test]
fn capped_by_serializes_as_a_client_readable_word() {
assert_eq!(serde_json::to_string(&Capped::Rows).unwrap(), "\"rows\"");
assert_eq!(serde_json::to_string(&Capped::Bytes).unwrap(), "\"bytes\"");
assert_eq!(serde_json::to_string(&Capped::Time).unwrap(), "\"time\"");
assert_eq!(Capped::Time.as_str(), "time");
}
}