#![cfg(all(
feature = "source-csv",
feature = "sink-sqlite",
feature = "source-mongodb-cdc"
))]
use clap::Parser;
use faucet_cli::cli::{Cli, StateLoadArgs, StatusArgs};
use faucet_cli::commands::migrate::{StateKeyMigration, StateMigrationReport, render_state_report};
use faucet_cli::config::{ConnectorSpec, StateStoreSpec};
use faucet_cli::error::CliError;
use faucet_cli::status::Health;
use faucet_core::check::ProbeStatus;
use faucet_core::state_version::{StoredState, wrap_versioned};
use faucet_core::{FileStateStore, StateStore};
use serde_json::{Value, json};
use std::path::{Path, PathBuf};
async fn run(args: &[&str]) -> Result<(), CliError> {
let mut argv = vec!["faucet"];
argv.extend_from_slice(args);
let cli = Cli::try_parse_from(argv).expect("argv parses");
Box::pin(faucet_cli::run_command(cli)).await
}
fn s(p: &Path) -> String {
p.display().to_string()
}
fn write_config(dir: &Path) -> PathBuf {
std::fs::write(dir.join("in.csv"), "id,name\n1,x\n").unwrap();
let text = format!(
r#"version: 1
name: orders
pipeline:
sources:
csv: {{ type: csv, config: {{ path: {input} }} }}
cdc:
type: mongodb-cdc
config:
connection_uri: "mongodb://127.0.0.1:1/?directConnection=true"
scope: {{ type: collection, database: app, collection: users }}
sink: {{ type: sqlite, config: {{ database_url: "sqlite://{db}?mode=rwc", table_name: t, column_mapping: auto_map, create_table: true }} }}
state: {{ type: file, config: {{ path: {state} }} }}
matrix:
- id: a
source: {{ ref: csv }}
- id: cdc
source: {{ ref: cdc }}
"#,
input = s(&dir.join("in.csv")),
db = s(&dir.join("out.db")),
state = s(&dir.join("state")),
);
let path = dir.join("orders.yaml");
std::fs::write(&path, text).unwrap();
path
}
fn status_args(cfg: &Path) -> StatusArgs {
StatusArgs {
config: Some(cfg.to_path_buf()),
row: None,
probe: false,
load: StateLoadArgs {
json: true,
env_file: None,
no_env_file: true,
profile: None,
},
}
}
#[tokio::test]
async fn migrate_state_envelopes_migrates_and_refuses() {
let dir = tempfile::tempdir().unwrap();
let cfg = write_config(dir.path());
let cfgs = s(&cfg);
let store = FileStateStore::new(dir.path().join("state"));
let legacy_cdc = json!({ "resume_token": { "_data": "AA" } });
store.put("orders::a", &json!({ "id": 2 })).await.unwrap();
store.put("orders::cdc", &legacy_cdc).await.unwrap();
let load = StateLoadArgs {
json: true,
env_file: None,
no_env_file: true,
profile: None,
};
let (_, target, _) = faucet_cli::commands::state::load(Some(&cfg), &load)
.await
.unwrap();
let stores = faucet_cli::pipeline_state::ops::Stores::build(&target, None)
.await
.unwrap();
let show = faucet_cli::pipeline_state::ops::show(&target, &stores, None, chrono::Utc::now())
.await
.unwrap();
let fmt = |row: &str| {
show.rows
.iter()
.find(|r| r.row == row)
.and_then(|r| r.state_format.clone())
.expect("a state format")
};
assert_eq!(fmt("a").status, "legacy");
assert_eq!(fmt("cdc").status, "migrate");
assert_eq!(fmt("cdc").expected_schema, 1);
run(&["state", "show", &cfgs]).await.unwrap();
let err = run(&["migrate", "--state", &cfgs, "--check"])
.await
.unwrap_err();
assert!(err.to_string().contains("not current"), "{err}");
assert_eq!(store.get("orders::cdc").await.unwrap(), Some(legacy_cdc));
let report = faucet_cli::commands::migrate::migrate_state(&target, &stores, None, true)
.await
.unwrap();
assert_eq!(report.pending(), 2);
run(&["migrate", "--state", &cfgs, "--row", "cdc", "--json"])
.await
.unwrap();
let cdc = StoredState::parse(&store.get("orders::cdc").await.unwrap().unwrap());
assert_eq!(
(cdc.format, cdc.owner.as_deref(), cdc.schema),
(1, Some("mongodb-cdc"), 1)
);
assert_eq!(cdc.data["invalidate"], json!(false));
assert_eq!(
store.get("orders::a").await.unwrap(),
Some(json!({ "id": 2 })),
"--row leaves the other row alone"
);
run(&["migrate", "--state", &cfgs]).await.unwrap();
assert_eq!(
store.get("orders::a").await.unwrap(),
Some(wrap_versioned("csv", 0, &json!({ "id": 2 })))
);
run(&["migrate", "--state", &cfgs, "--check"])
.await
.expect("everything is current now");
let show = faucet_cli::pipeline_state::ops::show(&target, &stores, None, chrono::Utc::now())
.await
.unwrap();
assert!(
show.rows.iter().all(|r| r
.state_format
.as_ref()
.is_some_and(|f| f.status == "current")),
"{:?}",
show.rows
);
let foreign = wrap_versioned("mysql-cdc", 0, &json!({ "file": "b.1", "pos": 4 }));
store.put("orders::a", &foreign).await.unwrap();
let err = run(&["migrate", "--state", &cfgs]).await.unwrap_err();
assert!(err.to_string().contains("cannot be read"), "{err}");
assert_eq!(store.get("orders::a").await.unwrap(), Some(foreign));
let status = faucet_cli::commands::status::build(&status_args(&cfg))
.await
.unwrap();
let row = status.rows.iter().find(|r| r.row == "a").unwrap();
assert_eq!(row.health, Health::Degraded, "{:?}", row.reasons);
assert!(
row.reasons.iter().any(|r| r.contains("cannot be read")),
"{:?}",
row.reasons
);
let text = faucet_cli::status::render::render(&status);
assert!(text.contains("state: incompatible"), "{text}");
assert!(
run(&["status", &cfgs]).await.is_err(),
"a degraded row makes `faucet status` exit non-zero"
);
run(&[
"state",
"set",
&cfgs,
"--row",
"a",
"--bookmark",
r#"{"id":5}"#,
"--yes",
])
.await
.unwrap();
assert_eq!(
store.get("orders::a").await.unwrap(),
Some(wrap_versioned("csv", 0, &json!({ "id": 5 })))
);
}
#[test]
fn the_state_report_renders_every_action() {
let key = |action: &'static str, detail: Option<&str>| StateKeyMigration {
row: "r".into(),
key: "p::r".into(),
owner: "mongodb-cdc".into(),
action,
from_schema: 0,
to_schema: 1,
detail: detail.map(str::to_owned),
};
let keys = vec![
key("current", None),
key("enveloped", None),
key("migrated", None),
key("refused", Some("found 'x' state")),
];
let applied = StateMigrationReport {
pipeline: "p".into(),
check: false,
keys: keys.clone(),
};
let text = render_state_report(&applied);
assert!(text.contains("state migration"));
assert!(text.contains("current (mongodb-cdc schema 1)"));
assert!(text.contains("rewritten in the versioned envelope"));
assert!(text.contains("migrated: mongodb-cdc schema 0 → 1"));
assert!(text.contains("REFUSED — found 'x' state"));
assert_eq!(applied.pending(), 3);
let check = StateMigrationReport {
pipeline: "p".into(),
check: true,
keys,
};
let text = render_state_report(&check);
assert!(text.contains("state check"));
assert!(text.contains("would be rewritten in the envelope"));
assert!(text.contains("needs migration: mongodb-cdc schema 0 → 1"));
let empty = StateMigrationReport {
pipeline: "p".into(),
check: false,
keys: vec![],
};
assert!(render_state_report(&empty).contains("no stored bookmarks"));
}
fn spec<T: serde::de::DeserializeOwned>(v: Value) -> T {
serde_json::from_value(v).unwrap()
}
#[tokio::test]
async fn the_doctor_state_probe_reports_each_format() {
use faucet_cli::commands::doctor::state_format_probe;
let dir = tempfile::tempdir().unwrap();
let state: StateStoreSpec =
spec(json!({ "type": "file", "config": { "path": s(&dir.path().join("st")) } }));
let store = FileStateStore::new(dir.path().join("st"));
let csv: ConnectorSpec = spec(json!({ "type": "csv", "config": { "path": "in.csv" } }));
let cdc: ConnectorSpec = spec(json!({
"type": "mongodb-cdc",
"config": { "connection_uri": "mongodb://127.0.0.1:1", "scope": { "type": "collection", "database": "a", "collection": "b" } }
}));
assert!(
state_format_probe(&csv, Some(&state), "p::r")
.await
.is_none()
);
let memory: StateStoreSpec = spec(json!({ "type": "memory" }));
store.put("p::r", &json!({ "id": 1 })).await.unwrap();
assert!(
state_format_probe(&csv, Some(&memory), "p::r")
.await
.is_none()
);
assert!(state_format_probe(&csv, None, "p::r").await.is_none());
let probe = state_format_probe(&csv, Some(&state), "p::r")
.await
.unwrap();
assert_eq!((probe.role, probe.name), ("state", "format"));
assert!(matches!(probe.status, ProbeStatus::Pass), "{probe:?}");
assert!(probe.hint.unwrap().contains("before versioning"));
store
.put("p::r", &wrap_versioned("csv", 0, &json!({ "id": 1 })))
.await
.unwrap();
let probe = state_format_probe(&csv, Some(&state), "p::r")
.await
.unwrap();
assert!(matches!(probe.status, ProbeStatus::Pass));
assert!(probe.hint.is_none());
store
.put("p::r", &json!({ "resume_token": { "_data": "AA" } }))
.await
.unwrap();
let probe = state_format_probe(&cdc, Some(&state), "p::r")
.await
.unwrap();
assert!(
matches!(&probe.status, ProbeStatus::Skip { reason } if reason.contains("0 → 1")),
"{probe:?}"
);
store
.put("p::r", &wrap_versioned("kafka", 0, &json!({})))
.await
.unwrap();
let probe = state_format_probe(&csv, Some(&state), "p::r")
.await
.unwrap();
assert!(
matches!(&probe.status, ProbeStatus::Fail { reason } if reason.contains("kafka")),
"{probe:?}"
);
std::fs::write(dir.path().join("not-a-dir"), "x").unwrap();
let broken: StateStoreSpec =
spec(json!({ "type": "file", "config": { "path": s(&dir.path().join("not-a-dir")) } }));
let probe = state_format_probe(&csv, Some(&broken), "p::r")
.await
.unwrap();
assert!(
matches!(probe.status, ProbeStatus::Fail { .. }),
"{probe:?}"
);
}