use crate::local_outputs::{LocalOutputRecord, LocalOutputState};
use crate::serve::error::ServeError;
use crate::serve::preview::engine::PREVIEW_UNLIMITED_PAGE_ROWS;
use crate::serve::preview::{Capped, PreviewPage, PreviewRequest, RowCap, RowRequest, read_capped};
use crate::serve::rbac::AuthContext;
use crate::serve::state::ServerState;
use axum::Json;
use axum::extract::{Extension, Path, Query, State};
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
#[derive(Debug, Default, Deserialize)]
pub struct PreviewQuery {
pub row_count_to_load: Option<String>,
}
impl PreviewQuery {
fn requested(&self) -> Result<Option<RowRequest>, ServeError> {
self.row_count_to_load
.as_deref()
.map(RowRequest::parse)
.transpose()
.map_err(ServeError::BadConfig)
}
}
#[derive(Debug, Serialize)]
pub struct PreviewResponse {
pub output_id: String,
pub path: String,
pub kind: String,
pub dataset_id: String,
pub pipeline: String,
pub run_id: String,
pub rows: Vec<Value>,
pub columns: Vec<String>,
pub row_count: usize,
pub row_limit: Option<usize>,
pub max_rows: Option<usize>,
pub truncated: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub capped_by: Option<Capped>,
pub elapsed_ms: u64,
}
pub async fn preview_output(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
Query(query): Query<PreviewQuery>,
) -> Result<Json<PreviewResponse>, ServeError> {
let policy = *state.preview();
if !policy.enabled {
return Err(ServeError::Forbidden(DISABLED.into()));
}
let record = state
.history()
.local_output_get(&id)
.await
.map_err(|e| ServeError::Internal(e.to_string()))?
.ok_or(ServeError::NotFound)?;
let rows = policy.resolve_rows(query.requested()?);
let request = source_spec(&record, rows)?;
let page = read_capped(&request, &crate::auth_catalog::AuthCatalog::new())
.await
.map_err(|e| read_error(&record, e))?;
audit(&state, &actor, &record, page.rows.len()).await;
Ok(Json(response(&record, page, rows, policy.max_rows())))
}
async fn audit(state: &ServerState, actor: &AuthContext, record: &LocalOutputRecord, rows: usize) {
crate::serve::audit::write(
state,
actor,
"local_output.preview",
Some(record.run_id.clone()),
None,
&format!("output={} kind={} rows={rows}", record.id, record.kind),
)
.await;
}
const DISABLED: &str = "local-output preview is disabled on this server — start \
`faucet serve` with `--preview-local-outputs` (or set \
FAUCET_SERVE_PREVIEW_LOCAL_OUTPUTS=true) to enable it. It reads \
the contents of files on the server's disk, so it is opt-in and \
intended for local testing.";
fn source_spec(record: &LocalOutputRecord, rows: RowCap) -> Result<PreviewRequest, ServeError> {
if record.state() == LocalOutputState::Expired {
return Err(ServeError::Conflict(format!(
"`{}` was cleaned up by local-output retention — the run record is kept, \
but the file is gone. Re-run the pipeline to regenerate it.",
record.path
)));
}
if !record.fs_path().exists() {
return Err(ServeError::Conflict(format!(
"`{}` is no longer on disk (removed outside faucet). The ledger row is \
kept; re-run the pipeline to regenerate the file.",
record.path
)));
}
if record.state() == LocalOutputState::External {
return Err(ServeError::Forbidden(format!(
"`{}` already existed when faucet first opened it, so faucet wrote to a file \
it does not own. Its contents are not previewed, for the same reason the \
retention GC will not delete it. The output is still listed, with state \
`external`.",
record.path
)));
}
let page = match rows {
RowCap::Rows(n) => n.saturating_add(1),
RowCap::Unlimited => PREVIEW_UNLIMITED_PAGE_ROWS,
};
let row_limit = match rows {
RowCap::Rows(n) => n.saturating_add(1),
RowCap::Unlimited => 0,
};
let config = match record.kind.as_str() {
"jsonl" => json!({ "path": record.path, "batch_size": page, "limit": row_limit }),
"csv" => json!({ "path": record.path, "batch_size": page }),
"file" => json!({ "path": record.path, "batch_size": page, "strict": true }),
"parquet" => json!({
"source": { "type": "local_path", "path": record.path },
"batch_size": page,
}),
other => {
return Err(ServeError::BadConfig(format!(
"preview is not supported for `{other}` outputs — only the local file \
sinks faucet can read back (jsonl, csv, parquet, file)"
)));
}
};
Ok(PreviewRequest {
kind: record.kind.clone(),
config,
rows,
})
}
fn read_error(record: &LocalOutputRecord, err: crate::error::CliError) -> ServeError {
use crate::error::CliError;
match err {
CliError::UnknownConnector { .. } => ServeError::BadConfig(format!(
"this build of faucet cannot read `{}` outputs: {err}",
record.kind
)),
CliError::Serve(m) => ServeError::Unavailable(m),
other => ServeError::Unprocessable {
message: format!("could not read `{}`: {other}", record.path),
details: None,
},
}
}
fn response(
record: &LocalOutputRecord,
page: PreviewPage,
rows: RowCap,
max_rows: Option<usize>,
) -> PreviewResponse {
PreviewResponse {
output_id: record.id.clone(),
path: record.path.clone(),
kind: record.kind.clone(),
dataset_id: record.dataset_id.clone(),
pipeline: record.pipeline.clone(),
run_id: record.run_id.clone(),
row_count: page.rows.len(),
truncated: page.truncated(),
capped_by: page.capped_by,
rows: page.rows,
columns: page.columns,
row_limit: rows.rows(),
max_rows,
elapsed_ms: page.elapsed_ms,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::local_outputs::{LocalOutputObservation, LocalOutputRecord};
use chrono::Utc;
use std::path::Path as FsPath;
fn record(path: &FsPath, kind: &str) -> LocalOutputRecord {
LocalOutputRecord::new(&LocalOutputObservation {
path: path.to_path_buf(),
dataset_uri: format!("file://{}", path.display()),
dataset_id: "ds-1".into(),
kind: kind.into(),
pipeline: "demo".into(),
row: "default".into(),
run_id: "run-1".into(),
pre_existing: false,
replaced: false,
retention_days: None,
observed_at: Utc::now(),
})
}
fn actor() -> AuthContext {
AuthContext {
principal: "bob".into(),
role: crate::serve::rbac::Role::Viewer,
source_ip: None,
tenant: None,
}
}
fn touch(dir: &FsPath, name: &str) -> std::path::PathBuf {
let p = dir.join(name);
std::fs::write(&p, "{\"a\":1}\n").unwrap();
p
}
#[test]
fn jsonl_spec_carries_the_ledger_path_and_a_bounded_page() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "out.jsonl");
let spec = source_spec(&record(&p, "jsonl"), RowCap::Rows(10)).unwrap();
assert_eq!(spec.kind, "jsonl");
assert_eq!(spec.config["path"], p.to_string_lossy().as_ref());
assert_eq!(spec.config["batch_size"], 11);
assert_eq!(spec.config["limit"], 11, "the reader stops on its own too");
assert_eq!(spec.rows, RowCap::Rows(10));
}
#[test]
fn csv_spec_targets_the_csv_source() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "rows.csv");
let spec = source_spec(&record(&p, "csv"), RowCap::Rows(5)).unwrap();
assert_eq!(spec.kind, "csv");
assert_eq!(spec.config["path"], p.to_string_lossy().as_ref());
assert_eq!(spec.config["batch_size"], 6);
}
#[test]
fn file_spec_targets_the_file_source_strictly() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "rows.jsonl");
let spec = source_spec(&record(&p, "file"), RowCap::Rows(3)).unwrap();
assert_eq!(spec.kind, "file");
assert_eq!(spec.config["path"], p.to_string_lossy().as_ref());
assert_eq!(spec.config["batch_size"], 4);
assert_eq!(spec.config["strict"], true);
}
#[test]
fn parquet_spec_uses_a_single_local_path_never_a_glob() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "part.parquet");
let spec = source_spec(&record(&p, "parquet"), RowCap::Rows(7)).unwrap();
assert_eq!(spec.kind, "parquet");
assert_eq!(spec.config["source"]["type"], "local_path");
assert_eq!(spec.config["source"]["path"], p.to_string_lossy().as_ref());
assert!(
spec.config["source"].get("pattern").is_none(),
"a glob would let one row expand to many files"
);
}
#[test]
fn an_unreadable_kind_is_rejected_by_name() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "thing.out");
let err = source_spec(&record(&p, "stdout"), RowCap::Rows(10)).unwrap_err();
match err {
ServeError::BadConfig(m) => assert!(m.contains("stdout"), "{m}"),
other => panic!("expected BadConfig, got {other:?}"),
}
}
#[test]
fn a_collected_output_is_a_conflict_not_a_failed_open() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "out.jsonl");
let mut r = record(&p, "jsonl");
r.deleted_at = Some(Utc::now());
let err = source_spec(&r, RowCap::Rows(10)).unwrap_err();
match err {
ServeError::Conflict(m) => assert!(m.contains("cleaned up"), "{m}"),
other => panic!("expected Conflict, got {other:?}"),
}
}
#[test]
fn a_file_removed_out_of_band_is_a_conflict_too() {
let dir = tempfile::tempdir().unwrap();
let r = record(&dir.path().join("never-written.jsonl"), "jsonl");
let err = source_spec(&r, RowCap::Rows(10)).unwrap_err();
match err {
ServeError::Conflict(m) => assert!(m.contains("no longer on disk"), "{m}"),
other => panic!("expected Conflict, got {other:?}"),
}
}
#[test]
fn an_external_file_is_not_previewable() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "theirs.jsonl");
let mut r = record(&p, "jsonl");
r.pre_existing = true;
assert_eq!(r.state(), LocalOutputState::External);
match source_spec(&r, RowCap::Rows(10)).unwrap_err() {
ServeError::Forbidden(m) => {
assert!(m.contains("does not own"), "{m}");
assert!(m.contains("external"), "the state stays discoverable: {m}");
}
other => panic!("expected Forbidden, got {other:?}"),
}
}
#[test]
fn a_replaced_file_is_previewable() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "overwritten.jsonl");
let mut r = record(&p, "jsonl");
r.pre_existing = true;
r.replaced = true;
assert_eq!(r.state(), LocalOutputState::Replaced);
assert!(source_spec(&r, RowCap::Rows(10)).is_ok());
}
#[test]
fn an_abandoned_read_is_unavailable_not_a_500() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "out.jsonl");
let err = read_error(
&record(&p, "jsonl"),
crate::error::CliError::Serve("preview abandoned after 60s".into()),
);
match err {
ServeError::Unavailable(m) => assert!(m.contains("abandoned"), "{m}"),
other => panic!("expected Unavailable, got {other:?}"),
}
assert_eq!(
ServeError::Unavailable(String::new()).status(),
axum::http::StatusCode::SERVICE_UNAVAILABLE
);
}
#[test]
fn an_unreadable_file_is_unprocessable_and_keeps_the_connectors_diagnostic() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "rows.csv");
let err = read_error(
&record(&p, "csv"),
crate::error::CliError::Config("ragged CSV row at line 7".into()),
);
match err {
ServeError::Unprocessable { message, .. } => {
assert!(message.contains("line 7"), "{message}");
assert!(message.contains("rows.csv"), "{message}");
}
other => panic!("expected Unprocessable, got {other:?}"),
}
}
#[test]
fn a_missing_connector_says_the_build_lacks_it() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "part.parquet");
let err = read_error(
&record(&p, "parquet"),
crate::error::CliError::UnknownConnector {
kind: "source",
name: "parquet".into(),
available: "csv".into(),
},
);
match err {
ServeError::BadConfig(m) => assert!(m.contains("parquet"), "{m}"),
other => panic!("expected BadConfig, got {other:?}"),
}
}
#[tokio::test]
async fn a_disabled_server_refuses_before_touching_the_ledger() {
let state = crate::serve::test_support::test_state();
assert!(!state.preview().enabled, "off by default");
let err = preview_output(
axum::extract::State(state),
axum::extract::Extension(actor()),
axum::extract::Path("whatever".into()),
axum::extract::Query(PreviewQuery::default()),
)
.await
.unwrap_err();
match err {
ServeError::Forbidden(m) => {
assert!(m.contains("--preview-local-outputs"), "{m}")
}
other => panic!("expected Forbidden, got {other:?}"),
}
}
#[tokio::test]
async fn an_unknown_id_is_not_found() {
let mut config = crate::serve::test_support::test_config();
config.preview = crate::serve::preview::PreviewConfig::new(true, 100, 1000);
let state = crate::serve::test_support::state_from(&config);
let err = preview_output(
axum::extract::State(state),
axum::extract::Extension(actor()),
axum::extract::Path("no-such-output".into()),
axum::extract::Query(PreviewQuery::default()),
)
.await
.unwrap_err();
assert!(matches!(err, ServeError::NotFound), "{err:?}");
}
#[tokio::test]
async fn previews_a_recorded_jsonl_output_end_to_end() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("out.jsonl");
std::fs::write(&path, "{\"a\":1,\"b\":\"x\"}\n{\"a\":2,\"b\":\"y\"}\n").unwrap();
let mut config = crate::serve::test_support::test_config();
config.preview = crate::serve::preview::PreviewConfig::new(true, 100, 1000);
let state = crate::serve::test_support::state_from(&config);
let obs = LocalOutputObservation {
path: path.clone(),
dataset_uri: "file:///out.jsonl".into(),
dataset_id: "ds-1".into(),
kind: "jsonl".into(),
pipeline: "demo".into(),
row: "default".into(),
run_id: "run-1".into(),
pre_existing: false,
replaced: false,
retention_days: None,
observed_at: Utc::now(),
};
state.history().local_output_record(&obs).await.unwrap();
let id = crate::local_outputs::ledger::output_id(&path);
let axum::Json(body) = preview_output(
axum::extract::State(state),
axum::extract::Extension(actor()),
axum::extract::Path(id.clone()),
axum::extract::Query(PreviewQuery::default()),
)
.await
.unwrap();
assert_eq!(body.output_id, id);
assert_eq!(body.rows.len(), 2);
assert_eq!(body.row_count, 2);
assert_eq!(body.columns, vec!["a".to_string(), "b".to_string()]);
assert_eq!(body.row_limit, Some(100));
assert_eq!(body.max_rows, Some(1000));
assert!(!body.truncated);
assert_eq!(body.capped_by, None);
assert_eq!(body.kind, "jsonl");
assert_eq!(body.run_id, "run-1");
}
#[tokio::test]
async fn row_count_to_load_is_clamped_to_the_hard_cap() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("out.jsonl");
let body: String = (0..20).map(|i| format!("{{\"i\":{i}}}\n")).collect();
std::fs::write(&path, body).unwrap();
let mut config = crate::serve::test_support::test_config();
config.preview = crate::serve::preview::PreviewConfig::new(true, 5, 3);
let state = crate::serve::test_support::state_from(&config);
state
.history()
.local_output_record(&LocalOutputObservation {
path: path.clone(),
dataset_uri: "file:///out.jsonl".into(),
dataset_id: "ds-1".into(),
kind: "jsonl".into(),
pipeline: "demo".into(),
row: "default".into(),
run_id: "run-1".into(),
pre_existing: false,
replaced: false,
retention_days: None,
observed_at: Utc::now(),
})
.await
.unwrap();
let axum::Json(body) = preview_output(
axum::extract::State(state),
axum::extract::Extension(actor()),
axum::extract::Path(crate::local_outputs::ledger::output_id(&path)),
axum::extract::Query(PreviewQuery {
row_count_to_load: Some("1000".into()),
}),
)
.await
.unwrap();
assert_eq!(body.row_limit, Some(3), "clamped to the hard cap");
assert_eq!(body.rows.len(), 3);
assert!(body.truncated, "17 rows were left unread");
assert_eq!(body.capped_by, Some(Capped::Rows));
}
async fn enabled_state_with(
path: &std::path::Path,
default_rows: usize,
max_rows: usize,
) -> (crate::serve::state::ServerState, String) {
let mut config = crate::serve::test_support::test_config();
config.preview = crate::serve::preview::PreviewConfig::new(true, default_rows, max_rows);
let state = crate::serve::test_support::state_from(&config);
state
.history()
.local_output_record(&LocalOutputObservation {
path: path.to_path_buf(),
dataset_uri: "file:///out.jsonl".into(),
dataset_id: "ds-1".into(),
kind: "jsonl".into(),
pipeline: "demo".into(),
row: "default".into(),
run_id: "run-1".into(),
pre_existing: false,
replaced: false,
retention_days: None,
observed_at: Utc::now(),
})
.await
.unwrap();
(state, crate::local_outputs::ledger::output_id(path))
}
fn write_rows(dir: &FsPath, rows: usize) -> std::path::PathBuf {
let p = dir.join("out.jsonl");
let body: String = (0..rows).map(|i| format!("{{\"i\":{i}}}\n")).collect();
std::fs::write(&p, body).unwrap();
p
}
#[test]
fn an_unlimited_read_is_unlimited_in_rows_but_still_paged() {
let dir = tempfile::tempdir().unwrap();
let p = touch(dir.path(), "out.jsonl");
let spec = source_spec(&record(&p, "jsonl"), RowCap::Unlimited).unwrap();
assert_eq!(
spec.config["batch_size"], PREVIEW_UNLIMITED_PAGE_ROWS,
"an unbounded page defeats every bound the engine has"
);
assert_ne!(spec.config["batch_size"], 0);
assert_eq!(spec.config["limit"], 0, "no *row* limit, though");
assert_eq!(spec.rows, RowCap::Unlimited);
}
#[test]
fn every_kind_is_paged_under_an_unlimited_read() {
let dir = tempfile::tempdir().unwrap();
for kind in ["jsonl", "csv", "parquet"] {
let p = touch(dir.path(), &format!("out.{kind}"));
let spec = source_spec(&record(&p, kind), RowCap::Unlimited).unwrap();
assert_eq!(
spec.config["batch_size"], PREVIEW_UNLIMITED_PAGE_ROWS,
"{kind} was handed an unbounded page"
);
}
}
#[tokio::test]
async fn row_count_to_load_all_returns_the_whole_dataset_when_no_ceiling_is_set() {
let dir = tempfile::tempdir().unwrap();
let path = write_rows(dir.path(), 1_200);
let (state, id) = enabled_state_with(&path, 500, 0).await;
let axum::Json(body) = preview_output(
axum::extract::State(state),
axum::extract::Extension(actor()),
axum::extract::Path(id),
axum::extract::Query(PreviewQuery {
row_count_to_load: Some("all".into()),
}),
)
.await
.unwrap();
assert_eq!(body.rows.len(), 1_200, "every row, not the soft cap");
assert_eq!(body.row_count, 1_200);
assert_eq!(body.row_limit, None, "null = unlimited");
assert_eq!(body.max_rows, None);
assert!(!body.truncated);
assert_eq!(body.capped_by, None);
}
#[tokio::test]
async fn row_count_to_load_all_is_still_clamped_by_a_configured_ceiling() {
let dir = tempfile::tempdir().unwrap();
let path = write_rows(dir.path(), 1_200);
let (state, id) = enabled_state_with(&path, 100, 250).await;
let axum::Json(body) = preview_output(
axum::extract::State(state),
axum::extract::Extension(actor()),
axum::extract::Path(id),
axum::extract::Query(PreviewQuery {
row_count_to_load: Some("all".into()),
}),
)
.await
.unwrap();
assert_eq!(body.rows.len(), 250);
assert_eq!(body.row_limit, Some(250));
assert!(body.truncated);
assert_eq!(body.capped_by, Some(Capped::Rows));
}
#[tokio::test]
async fn zero_is_the_same_request_as_all() {
let dir = tempfile::tempdir().unwrap();
let path = write_rows(dir.path(), 40);
let (state, id) = enabled_state_with(&path, 5, 0).await;
let axum::Json(body) = preview_output(
axum::extract::State(state),
axum::extract::Extension(actor()),
axum::extract::Path(id),
axum::extract::Query(PreviewQuery {
row_count_to_load: Some("0".into()),
}),
)
.await
.unwrap();
assert_eq!(body.rows.len(), 40);
assert_eq!(body.row_limit, None);
}
#[tokio::test]
async fn an_unparseable_row_count_is_a_400_naming_the_parameter() {
let dir = tempfile::tempdir().unwrap();
let path = write_rows(dir.path(), 3);
let (state, id) = enabled_state_with(&path, 100, 1000).await;
let err = preview_output(
axum::extract::State(state),
axum::extract::Extension(actor()),
axum::extract::Path(id),
axum::extract::Query(PreviewQuery {
row_count_to_load: Some("lots".into()),
}),
)
.await
.unwrap_err();
match err {
ServeError::BadConfig(m) => assert!(m.contains("row_count_to_load"), "{m}"),
other => panic!("expected BadConfig, got {other:?}"),
}
}
#[tokio::test]
async fn a_malformed_file_is_unprocessable_not_internal() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("out.jsonl");
std::fs::write(&path, "{\"a\":1}\n{oops\n").unwrap();
let mut config = crate::serve::test_support::test_config();
config.preview = crate::serve::preview::PreviewConfig::new(true, 100, 1000);
let state = crate::serve::test_support::state_from(&config);
state
.history()
.local_output_record(&LocalOutputObservation {
path: path.clone(),
dataset_uri: "file:///out.jsonl".into(),
dataset_id: "ds-1".into(),
kind: "jsonl".into(),
pipeline: "demo".into(),
row: "default".into(),
run_id: "run-1".into(),
pre_existing: false,
replaced: false,
retention_days: None,
observed_at: Utc::now(),
})
.await
.unwrap();
let err = preview_output(
axum::extract::State(state),
axum::extract::Extension(actor()),
axum::extract::Path(crate::local_outputs::ledger::output_id(&path)),
axum::extract::Query(PreviewQuery::default()),
)
.await
.unwrap_err();
match err {
ServeError::Unprocessable { message, .. } => {
assert!(message.contains("line 2"), "{message}")
}
other => panic!("expected Unprocessable, got {other:?}"),
}
}
}