#![allow(deprecated)]
use crate::auth_catalog::{self, AuthCatalog};
use crate::error::{CliError, CliResult};
use faucet_core::{FaucetError, Sink, Source};
use serde::de::DeserializeOwned;
use serde_json::Value;
use std::collections::{BTreeMap, BTreeSet};
use std::sync::{Arc, OnceLock};
pub type SourceFactory = Arc<dyn Fn(Value) -> CliResult<Box<dyn Source>> + Send + Sync>;
pub type SinkFactory = Arc<dyn Fn(Value) -> CliResult<Box<dyn Sink>> + Send + Sync>;
pub type SchemaFn = Arc<dyn Fn() -> Value + Send + Sync>;
struct SourceEntry {
factory: SourceFactory,
schema: SchemaFn,
description: &'static str,
}
struct SinkEntry {
factory: SinkFactory,
schema: SchemaFn,
description: &'static str,
}
#[derive(Default)]
pub struct PluginRegistry {
sources: BTreeMap<&'static str, SourceEntry>,
sinks: BTreeMap<&'static str, SinkEntry>,
errors: Vec<String>,
}
impl PluginRegistry {
pub fn new() -> Self {
Self::default()
}
pub fn with_builtins() -> Self {
Self::default()
}
#[must_use]
pub fn register_source<F>(self, name: &str, factory: F) -> Self
where
F: Fn(Value) -> CliResult<Box<dyn Source>> + Send + Sync + 'static,
{
self.register_source_with(name, factory, || serde_json::json!({"type": "object"}), "")
}
#[must_use]
pub fn register_source_with<F, S>(
mut self,
name: &str,
factory: F,
schema: S,
description: &str,
) -> Self
where
F: Fn(Value) -> CliResult<Box<dyn Source>> + Send + Sync + 'static,
S: Fn() -> Value + Send + Sync + 'static,
{
let key = leak_str(name);
if builtin_source_descriptions().iter().any(|(n, _)| *n == key) {
self.errors.push(format!(
"cannot register source `{name}`: a built-in source already uses that name"
));
return self;
}
if self.sources.contains_key(key) {
self.errors
.push(format!("source `{name}` is registered more than once"));
return self;
}
self.sources.insert(
key,
SourceEntry {
factory: Arc::new(factory),
schema: Arc::new(schema),
description: leak_str(description),
},
);
self
}
#[must_use]
pub fn register_sink<F>(self, name: &str, factory: F) -> Self
where
F: Fn(Value) -> CliResult<Box<dyn Sink>> + Send + Sync + 'static,
{
self.register_sink_with(name, factory, || serde_json::json!({"type": "object"}), "")
}
#[must_use]
pub fn register_sink_with<F, S>(
mut self,
name: &str,
factory: F,
schema: S,
description: &str,
) -> Self
where
F: Fn(Value) -> CliResult<Box<dyn Sink>> + Send + Sync + 'static,
S: Fn() -> Value + Send + Sync + 'static,
{
let key = leak_str(name);
if builtin_sink_descriptions().iter().any(|(n, _)| *n == key) {
self.errors.push(format!(
"cannot register sink `{name}`: a built-in sink already uses that name"
));
return self;
}
if self.sinks.contains_key(key) {
self.errors
.push(format!("sink `{name}` is registered more than once"));
return self;
}
self.sinks.insert(
key,
SinkEntry {
factory: Arc::new(factory),
schema: Arc::new(schema),
description: leak_str(description),
},
);
self
}
pub fn install(self) -> CliResult<()> {
if !self.errors.is_empty() {
return Err(CliError::Config(self.errors.join("; ")));
}
GLOBAL_REGISTRY
.set(self)
.map_err(|_| CliError::Config("connector registry already installed".to_owned()))
}
fn custom_source_descriptions(&self) -> Vec<(&'static str, &'static str)> {
self.sources
.iter()
.map(|(name, e)| {
(
*name,
if e.description.is_empty() {
"custom source connector"
} else {
e.description
},
)
})
.collect()
}
fn custom_sink_descriptions(&self) -> Vec<(&'static str, &'static str)> {
self.sinks
.iter()
.map(|(name, e)| {
(
*name,
if e.description.is_empty() {
"custom sink connector"
} else {
e.description
},
)
})
.collect()
}
}
fn leak_str(s: &str) -> &'static str {
Box::leak(s.to_owned().into_boxed_str())
}
static GLOBAL_REGISTRY: OnceLock<PluginRegistry> = OnceLock::new();
fn global() -> &'static PluginRegistry {
GLOBAL_REGISTRY.get_or_init(PluginRegistry::default)
}
type Built<T> = std::pin::Pin<Box<dyn std::future::Future<Output = CliResult<Box<T>>> + Send>>;
#[allow(dead_code)]
fn boxed_source<S, F>(fut: F) -> Built<dyn Source>
where
S: Source + 'static,
F: std::future::Future<Output = Result<S, faucet_core::FaucetError>> + Send + 'static,
{
Box::pin(async move { Ok(Box::new(fut.await?) as Box<dyn Source>) })
}
#[allow(dead_code)]
fn boxed_sink<S, F>(fut: F) -> Built<dyn Sink>
where
S: Sink + 'static,
F: std::future::Future<Output = Result<S, faucet_core::FaucetError>> + Send + 'static,
{
Box::pin(async move { Ok(Box::new(fut.await?) as Box<dyn Sink>) })
}
pub async fn build_source(
kind: &str,
config: Value,
auth: &AuthCatalog,
retry_policy: Option<&faucet_core::RetryPolicy>,
) -> CliResult<Box<dyn Source>> {
if let Some(entry) = global().sources.get(kind) {
reject_unknown_config_keys("source", kind, kind, &config, &(entry.schema)())?;
return (entry.factory)(config);
}
if let Ok(schema) = source_schema(kind) {
reject_unknown_config_keys("source", kind, kind, &config, &schema)?;
}
let auth_ref = auth_catalog::auth_ref(&config);
match kind {
#[cfg(feature = "source-rest")]
"rest" => {
let cfg = decode::<faucet_source_rest::RestStreamConfig>("source", "rest", config)?;
let mut s = faucet_source_rest::RestStream::new(cfg)?;
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
if let Some(rp) = retry_policy {
s = s.with_retry_policy(rp.clone());
}
Ok(Box::new(s))
}
#[cfg(feature = "source-graphql")]
"graphql" => {
let cfg =
decode::<faucet_source_graphql::GraphqlStreamConfig>("source", "graphql", config)?;
cfg.validate()?;
let mut s = faucet_source_graphql::GraphqlStream::try_new(cfg)?;
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
if let Some(rp) = retry_policy {
s = s.with_retry_policy(rp.clone());
}
Ok(Box::new(s))
}
#[cfg(feature = "source-xml")]
"xml" => {
let cfg = decode::<faucet_source_xml::XmlStreamConfig>("source", "xml", config)?;
let mut s = faucet_source_xml::XmlStream::try_new(cfg)?;
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
if let Some(rp) = retry_policy {
s = s.with_retry_policy(rp.clone());
}
Ok(Box::new(s))
}
#[cfg(feature = "source-grpc")]
"grpc" => {
let cfg = decode::<faucet_source_grpc::GrpcStreamConfig>("source", "grpc", config)?;
let mut s = faucet_source_grpc::GrpcStream::new(cfg)?;
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
Ok(Box::new(s))
}
#[cfg(feature = "source-postgres")]
"postgres" => {
let cfg = decode::<faucet_source_postgres::PostgresSourceConfig>(
"source", "postgres", config,
)?;
boxed_source(faucet_source_postgres::PostgresSource::new(cfg)).await
}
#[cfg(feature = "source-postgres-cdc")]
"postgres-cdc" => {
let cfg = decode::<faucet_source_postgres_cdc::PostgresCdcSourceConfig>(
"source",
"postgres-cdc",
config,
)?;
boxed_source(faucet_source_postgres_cdc::PostgresCdcSource::new(cfg)).await
}
#[cfg(feature = "source-mysql")]
"mysql" => {
let cfg = decode::<faucet_source_mysql::MysqlSourceConfig>("source", "mysql", config)?;
boxed_source(faucet_source_mysql::MysqlSource::new(cfg)).await
}
#[cfg(feature = "source-mssql")]
"mssql" => {
let cfg = decode::<faucet_source_mssql::MssqlSourceConfig>("source", "mssql", config)?;
boxed_source(faucet_source_mssql::MssqlSource::new(cfg)).await
}
#[cfg(feature = "source-sqlite")]
"sqlite" => {
let cfg =
decode::<faucet_source_sqlite::SqliteSourceConfig>("source", "sqlite", config)?;
boxed_source(faucet_source_sqlite::SqliteSource::new(cfg)).await
}
#[cfg(feature = "source-duckdb")]
"duckdb" => {
let cfg =
decode::<faucet_source_duckdb::DuckdbSourceConfig>("source", "duckdb", config)?;
boxed_source(faucet_source_duckdb::DuckdbSource::new(cfg)).await
}
#[cfg(feature = "source-sqs")]
"sqs" => {
let cfg = decode::<faucet_source_sqs::SqsSourceConfig>("source", "sqs", config)?;
boxed_source(faucet_source_sqs::SqsSource::new(cfg)).await
}
#[cfg(feature = "source-nats")]
"nats" => {
let cfg = decode::<faucet_source_nats::NatsSourceConfig>("source", "nats", config)?;
boxed_source(faucet_source_nats::NatsSource::new(cfg)).await
}
#[cfg(feature = "source-rabbitmq")]
"rabbitmq" => {
let cfg = decode::<faucet_source_rabbitmq::RabbitMqSourceConfig>(
"source", "rabbitmq", config,
)?;
boxed_source(faucet_source_rabbitmq::RabbitMqSource::new(cfg)).await
}
#[cfg(feature = "source-sftp")]
"sftp" => {
let cfg = decode::<faucet_source_sftp::SftpSourceConfig>("source", "sftp", config)?;
Ok(Box::new(faucet_source_sftp::SftpSource::new(cfg)?))
}
#[cfg(feature = "source-file")]
"file" => {
let cfg = decode::<faucet_source_file::FileSourceConfig>("source", "file", config)?;
Ok(Box::new(faucet_source_file::FileSource::new(cfg)?))
}
#[cfg(feature = "source-s3")]
"s3" => {
let cfg = decode::<faucet_source_s3::S3SourceConfig>("source", "s3", config)?;
boxed_source(faucet_source_s3::S3Source::new(cfg)).await
}
#[cfg(feature = "source-mongodb")]
"mongodb" => {
let cfg =
decode::<faucet_source_mongodb::MongoSourceConfig>("source", "mongodb", config)?;
boxed_source(faucet_source_mongodb::MongoSource::new(cfg)).await
}
#[cfg(feature = "source-mongodb-cdc")]
"mongodb-cdc" => {
let cfg = decode::<faucet_source_mongodb_cdc::MongoCdcSourceConfig>(
"source",
"mongodb-cdc",
config,
)?;
boxed_source(faucet_source_mongodb_cdc::MongoCdcSource::new(cfg)).await
}
#[cfg(feature = "source-mysql-cdc")]
"mysql-cdc" => {
let cfg = decode::<faucet_source_mysql_cdc::MysqlCdcSourceConfig>(
"source",
"mysql-cdc",
config,
)?;
boxed_source(faucet_source_mysql_cdc::MysqlCdcSource::new(cfg)).await
}
#[cfg(feature = "source-redis")]
"redis" => {
let cfg = decode::<faucet_source_redis::RedisSourceConfig>("source", "redis", config)?;
Ok(Box::new(faucet_source_redis::RedisSource::new(cfg)?))
}
#[cfg(feature = "source-webhook")]
"webhook" => {
let cfg =
decode::<faucet_source_webhook::WebhookSourceConfig>("source", "webhook", config)?;
Ok(Box::new(faucet_source_webhook::WebhookSource::new(cfg)))
}
#[cfg(feature = "source-websocket")]
"websocket" => {
let cfg = decode::<faucet_source_websocket::WebsocketSourceConfig>(
"source",
"websocket",
config,
)?;
let mut s = faucet_source_websocket::WebsocketSource::new(cfg)?;
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
Ok(Box::new(s))
}
#[cfg(feature = "source-csv")]
"csv" => {
let cfg = decode::<faucet_source_csv::CsvSourceConfig>("source", "csv", config)?;
cfg.validate()?;
Ok(Box::new(faucet_source_csv::CsvSource::new(cfg)))
}
#[cfg(feature = "source-singer")]
"singer" => {
let cfg =
decode::<faucet_source_singer::SingerSourceConfig>("source", "singer", config)?;
Ok(Box::new(faucet_source_singer::SingerSource::new(cfg)))
}
#[cfg(feature = "source-elasticsearch")]
"elasticsearch" => {
let cfg = decode::<faucet_source_elasticsearch::ElasticsearchSourceConfig>(
"source",
"elasticsearch",
config,
)?;
let mut s = faucet_source_elasticsearch::ElasticsearchSource::new(cfg)?;
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
Ok(Box::new(s))
}
#[cfg(feature = "source-kafka")]
"kafka" => {
let cfg = decode::<faucet_source_kafka::KafkaSourceConfig>("source", "kafka", config)?;
boxed_source(faucet_source_kafka::KafkaSource::new(cfg)).await
}
#[cfg(feature = "source-kinesis")]
"kinesis" => {
let cfg =
decode::<faucet_source_kinesis::KinesisSourceConfig>("source", "kinesis", config)?;
boxed_source(faucet_source_kinesis::KinesisSource::new(cfg)).await
}
#[cfg(feature = "source-spanner")]
"spanner" => {
let cfg =
decode::<faucet_source_spanner::SpannerSourceConfig>("source", "spanner", config)?;
boxed_source(faucet_source_spanner::SpannerSource::new(cfg)).await
}
#[cfg(feature = "source-parquet")]
"parquet" => {
let cfg =
decode::<faucet_source_parquet::ParquetSourceConfig>("source", "parquet", config)?;
boxed_source(faucet_source_parquet::ParquetSource::new(cfg)).await
}
#[cfg(feature = "source-delta")]
"delta" => {
let cfg = decode::<faucet_source_delta::DeltaSourceConfig>("source", "delta", config)?;
boxed_source(faucet_source_delta::DeltaSource::new(cfg)).await
}
#[cfg(feature = "source-databricks")]
"databricks" => {
let cfg = decode::<faucet_source_databricks::DatabricksSourceConfig>(
"source",
"databricks",
config,
)?;
let mut s = faucet_source_databricks::DatabricksSource::new(cfg)?;
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
Ok(Box::new(s))
}
#[cfg(feature = "source-iceberg")]
"iceberg" => {
let cfg =
decode::<faucet_source_iceberg::IcebergSourceConfig>("source", "iceberg", config)?;
boxed_source(faucet_source_iceberg::IcebergSource::new(cfg)).await
}
#[cfg(feature = "source-dynamodb")]
"dynamodb" => {
let cfg = decode::<faucet_source_dynamodb::DynamoDbSourceConfig>(
"source", "dynamodb", config,
)?;
boxed_source(faucet_source_dynamodb::DynamoDbSource::new(cfg)).await
}
#[cfg(feature = "source-oracle")]
"oracle" => {
let cfg =
decode::<faucet_source_oracle::OracleSourceConfig>("source", "oracle", config)?;
boxed_source(faucet_source_oracle::OracleSource::new(cfg)).await
}
#[cfg(feature = "source-oracle-cdc")]
"oracle-cdc" => {
let cfg = decode::<faucet_source_oracle_cdc::OracleCdcSourceConfig>(
"source",
"oracle-cdc",
config,
)?;
boxed_source(faucet_source_oracle_cdc::OracleCdcSource::new(cfg)).await
}
#[cfg(feature = "source-gcs")]
"gcs" => {
let cfg = decode::<faucet_source_gcs::GcsSourceConfig>("source", "gcs", config)?;
boxed_source(faucet_source_gcs::GcsSource::new(cfg)).await
}
#[cfg(feature = "source-bigquery")]
"bigquery" => {
let cfg = decode::<faucet_source_bigquery::BigQuerySourceConfig>(
"source", "bigquery", config,
)?;
boxed_source(faucet_source_bigquery::BigQuerySource::new(cfg)).await
}
#[cfg(feature = "source-snowflake")]
"snowflake" => {
let cfg = decode::<faucet_source_snowflake::SnowflakeSourceConfig>(
"source",
"snowflake",
config,
)?;
let mut s = faucet_source_snowflake::SnowflakeSource::new(cfg)?;
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
Ok(Box::new(s))
}
#[cfg(feature = "source-mssql-cdc")]
"mssql-cdc" => {
let cfg = decode::<faucet_source_mssql_cdc::MssqlCdcSourceConfig>(
"source",
"mssql-cdc",
config,
)?;
boxed_source(faucet_source_mssql_cdc::MssqlCdcSource::new(cfg)).await
}
#[cfg(feature = "source-redshift")]
"redshift" => {
let cfg = decode::<faucet_source_redshift::RedshiftSourceConfig>(
"source", "redshift", config,
)?;
Ok(Box::new(faucet_source_redshift::RedshiftSource::new(cfg)?))
}
#[cfg(feature = "source-pubsub")]
"pubsub" => {
let cfg =
decode::<faucet_source_pubsub::PubsubSourceConfig>("source", "pubsub", config)?;
boxed_source(faucet_source_pubsub::PubsubSource::new(cfg)).await
}
#[cfg(feature = "source-clickhouse")]
"clickhouse" => {
let cfg = decode::<faucet_source_clickhouse::ClickHouseSourceConfig>(
"source",
"clickhouse",
config,
)?;
Ok(Box::new(faucet_source_clickhouse::ClickHouseSource::new(
cfg,
)?))
}
#[cfg(feature = "source-azure-blob")]
"azure-blob" => {
let cfg = decode::<faucet_source_azure_blob::AzureBlobSourceConfig>(
"source",
"azure-blob",
config,
)?;
boxed_source(faucet_source_azure_blob::AzureBlobSource::new(cfg)).await
}
other => Err(unknown(other, "source", source_kinds())),
}
}
pub async fn build_sink(kind: &str, config: Value, auth: &AuthCatalog) -> CliResult<Box<dyn Sink>> {
if let Some(entry) = global().sinks.get(kind) {
reject_unknown_config_keys("sink", kind, kind, &config, &(entry.schema)())?;
return (entry.factory)(config);
}
if let Ok(schema) = sink_schema(kind) {
reject_unknown_config_keys("sink", kind, kind, &config, &schema)?;
}
let auth_ref = auth_catalog::auth_ref(&config);
match kind {
#[cfg(feature = "sink-bigquery")]
"bigquery" => {
let cfg =
decode::<faucet_sink_bigquery::BigQuerySinkConfig>("sink", "bigquery", config)?;
boxed_sink(faucet_sink_bigquery::BigQuerySink::new(cfg)).await
}
#[cfg(feature = "sink-iceberg")]
"iceberg" => {
let cfg = decode::<faucet_sink_iceberg::IcebergSinkConfig>("sink", "iceberg", config)?;
boxed_sink(faucet_sink_iceberg::IcebergSink::new(cfg)).await
}
#[cfg(feature = "sink-delta")]
"delta" => {
let cfg = decode::<faucet_sink_delta::DeltaSinkConfig>("sink", "delta", config)?;
boxed_sink(faucet_sink_delta::DeltaSink::new(cfg)).await
}
#[cfg(feature = "sink-postgres")]
"postgres" => {
let cfg =
decode::<faucet_sink_postgres::PostgresSinkConfig>("sink", "postgres", config)?;
boxed_sink(faucet_sink_postgres::PostgresSink::new(cfg)).await
}
#[cfg(feature = "sink-jsonl")]
"jsonl" => {
let cfg = decode::<faucet_sink_jsonl::JsonlSinkConfig>("sink", "jsonl", config)?;
Ok(Box::new(faucet_sink_jsonl::JsonlSink::new(cfg)))
}
#[cfg(feature = "sink-snowflake")]
"snowflake" => {
let cfg =
decode::<faucet_sink_snowflake::SnowflakeSinkConfig>("sink", "snowflake", config)?;
let mut s = faucet_sink_snowflake::SnowflakeSink::new(cfg)?;
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
Ok(Box::new(s))
}
#[cfg(feature = "sink-mysql")]
"mysql" => {
let cfg = decode::<faucet_sink_mysql::MysqlSinkConfig>("sink", "mysql", config)?;
boxed_sink(faucet_sink_mysql::MysqlSink::new(cfg)).await
}
#[cfg(feature = "sink-mssql")]
"mssql" => {
let cfg = decode::<faucet_sink_mssql::MssqlSinkConfig>("sink", "mssql", config)?;
boxed_sink(faucet_sink_mssql::MssqlSink::new(cfg)).await
}
#[cfg(feature = "sink-sqlite")]
"sqlite" => {
let cfg = decode::<faucet_sink_sqlite::SqliteSinkConfig>("sink", "sqlite", config)?;
boxed_sink(faucet_sink_sqlite::SqliteSink::new(cfg)).await
}
#[cfg(feature = "sink-duckdb")]
"duckdb" => {
let cfg = decode::<faucet_sink_duckdb::DuckdbSinkConfig>("sink", "duckdb", config)?;
boxed_sink(faucet_sink_duckdb::DuckdbSink::new(cfg)).await
}
#[cfg(feature = "sink-sqs")]
"sqs" => {
let cfg = decode::<faucet_sink_sqs::SqsSinkConfig>("sink", "sqs", config)?;
boxed_sink(faucet_sink_sqs::SqsSink::new(cfg)).await
}
#[cfg(feature = "sink-nats")]
"nats" => {
let cfg = decode::<faucet_sink_nats::NatsSinkConfig>("sink", "nats", config)?;
boxed_sink(faucet_sink_nats::NatsSink::new(cfg)).await
}
#[cfg(feature = "sink-rabbitmq")]
"rabbitmq" => {
let cfg =
decode::<faucet_sink_rabbitmq::RabbitMqSinkConfig>("sink", "rabbitmq", config)?;
boxed_sink(faucet_sink_rabbitmq::RabbitMqSink::new(cfg)).await
}
#[cfg(feature = "sink-sftp")]
"sftp" => {
let cfg = decode::<faucet_sink_sftp::SftpSinkConfig>("sink", "sftp", config)?;
Ok(Box::new(faucet_sink_sftp::SftpSink::new(cfg)?))
}
#[cfg(feature = "sink-singer")]
"singer" => {
let cfg = decode::<faucet_sink_singer::SingerSinkConfig>("sink", "singer", config)?;
for secret in faucet_sink_singer::secret_like_values(&cfg.target_config) {
crate::secrets::registry::register(&secret);
}
Ok(Box::new(faucet_sink_singer::SingerSink::new(cfg)?))
}
#[cfg(feature = "sink-s3")]
"s3" => {
let cfg = decode::<faucet_sink_s3::S3SinkConfig>("sink", "s3", config)?;
boxed_sink(faucet_sink_s3::S3Sink::new(cfg)).await
}
#[cfg(feature = "sink-mongodb")]
"mongodb" => {
let cfg = decode::<faucet_sink_mongodb::MongoSinkConfig>("sink", "mongodb", config)?;
boxed_sink(faucet_sink_mongodb::MongoSink::new(cfg)).await
}
#[cfg(feature = "sink-redis")]
"redis" => {
let cfg = decode::<faucet_sink_redis::RedisSinkConfig>("sink", "redis", config)?;
boxed_sink(faucet_sink_redis::RedisSink::new(cfg)).await
}
#[cfg(feature = "sink-csv")]
"csv" => {
let cfg = decode::<faucet_sink_csv::CsvSinkConfig>("sink", "csv", config)?;
Ok(Box::new(faucet_sink_csv::CsvSink::new(cfg)))
}
#[cfg(feature = "sink-elasticsearch")]
"elasticsearch" => {
let cfg = decode::<faucet_sink_elasticsearch::ElasticsearchSinkConfig>(
"sink",
"elasticsearch",
config,
)?;
let mut s = faucet_sink_elasticsearch::ElasticsearchSink::new(cfg)?;
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
Ok(Box::new(s))
}
#[cfg(feature = "sink-kafka")]
"kafka" => {
let cfg = decode::<faucet_sink_kafka::KafkaSinkConfig>("sink", "kafka", config)?;
boxed_sink(faucet_sink_kafka::KafkaSink::new(cfg)).await
}
#[cfg(feature = "sink-kinesis")]
"kinesis" => {
let cfg = decode::<faucet_sink_kinesis::KinesisSinkConfig>("sink", "kinesis", config)?;
boxed_sink(faucet_sink_kinesis::KinesisSink::new(cfg)).await
}
#[cfg(feature = "sink-spanner")]
"spanner" => {
let cfg = decode::<faucet_sink_spanner::SpannerSinkConfig>("sink", "spanner", config)?;
boxed_sink(faucet_sink_spanner::SpannerSink::new(cfg)).await
}
#[cfg(feature = "sink-http")]
"http" => {
let cfg = decode::<faucet_sink_http::HttpSinkConfig>("sink", "http", config)?;
let mut s = faucet_sink_http::HttpSink::new(cfg);
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
Ok(Box::new(s))
}
#[cfg(feature = "sink-stdout")]
"stdout" => {
let cfg = decode::<faucet_sink_stdout::StdoutSinkConfig>("sink", "stdout", config)?;
Ok(Box::new(faucet_sink_stdout::StdoutSink::new(cfg)))
}
#[cfg(feature = "sink-parquet")]
"parquet" => {
let cfg = decode::<faucet_sink_parquet::ParquetSinkConfig>("sink", "parquet", config)?;
boxed_sink(faucet_sink_parquet::ParquetSink::new(cfg)).await
}
#[cfg(feature = "sink-file")]
"file" => {
let cfg = decode::<faucet_sink_file::FileSinkConfig>("sink", "file", config)?;
Ok(Box::new(faucet_sink_file::FileSink::new(cfg)?))
}
#[cfg(feature = "sink-gcs")]
"gcs" => {
let cfg = decode::<faucet_sink_gcs::GcsSinkConfig>("sink", "gcs", config)?;
boxed_sink(faucet_sink_gcs::GcsSink::new(cfg)).await
}
#[cfg(feature = "sink-redshift")]
"redshift" => {
let cfg =
decode::<faucet_sink_redshift::RedshiftSinkConfig>("sink", "redshift", config)?;
boxed_sink(faucet_sink_redshift::RedshiftSink::new(cfg)).await
}
#[cfg(feature = "sink-pubsub")]
"pubsub" => {
let cfg = decode::<faucet_sink_pubsub::PubsubSinkConfig>("sink", "pubsub", config)?;
boxed_sink(faucet_sink_pubsub::PubsubSink::new(cfg)).await
}
#[cfg(feature = "sink-clickhouse")]
"clickhouse" => {
let cfg = decode::<faucet_sink_clickhouse::ClickHouseSinkConfig>(
"sink",
"clickhouse",
config,
)?;
Ok(Box::new(faucet_sink_clickhouse::ClickHouseSink::new(cfg)?))
}
#[cfg(feature = "sink-azure-blob")]
"azure-blob" => {
let cfg = decode::<faucet_sink_azure_blob::AzureBlobSinkConfig>(
"sink",
"azure-blob",
config,
)?;
boxed_sink(faucet_sink_azure_blob::AzureBlobSink::new(cfg)).await
}
#[cfg(feature = "sink-dynamodb")]
"dynamodb" => {
let cfg =
decode::<faucet_sink_dynamodb::DynamoDbSinkConfig>("sink", "dynamodb", config)?;
boxed_sink(faucet_sink_dynamodb::DynamoDbSink::new(cfg)).await
}
#[cfg(feature = "sink-databricks")]
"databricks" => {
let cfg = decode::<faucet_sink_databricks::DatabricksSinkConfig>(
"sink",
"databricks",
config,
)?;
let mut s = faucet_sink_databricks::DatabricksSink::new(cfg)?;
if let Some(name) = &auth_ref {
s = s.with_auth_provider(auth_catalog::resolve(auth, name)?);
}
Ok(Box::new(s))
}
#[cfg(feature = "sink-oracle")]
"oracle" => {
let cfg = decode::<faucet_sink_oracle::OracleSinkConfig>("sink", "oracle", config)?;
boxed_sink(faucet_sink_oracle::OracleSink::new(cfg)).await
}
other => Err(unknown(other, "sink", sink_kinds())),
}
}
pub const EXACTLY_ONCE_SOURCE_KINDS: &[&str] = &[
"postgres-cdc",
"mysql-cdc",
"mssql-cdc",
"mongodb-cdc",
"kafka",
"oracle-cdc",
];
pub const LAG_SOURCE_KINDS: &[&str] = &[
"postgres-cdc",
"mysql-cdc",
"mssql-cdc",
"mongodb-cdc",
"oracle-cdc",
"kafka",
"kinesis",
"dynamodb",
];
pub fn source_reports_lag(kind: &str) -> bool {
LAG_SOURCE_KINDS.contains(&kind)
}
pub fn source_state_schema(kind: &str) -> u32 {
match kind {
#[cfg(feature = "source-mongodb-cdc")]
"mongodb-cdc" => faucet_source_mongodb_cdc::STATE_SCHEMA,
_ => 0,
}
}
pub fn migrate_source_state(kind: &str, from: u32, data: Value) -> Result<Value, FaucetError> {
match kind {
#[cfg(feature = "source-mongodb-cdc")]
"mongodb-cdc" => faucet_source_mongodb_cdc::migrate_state(from, data),
_ if from == source_state_schema(kind) => Ok(data),
_ => Err(FaucetError::State(format!(
"{kind} has no migration from bookmark schema {from} to {}",
source_state_schema(kind)
))),
}
}
pub fn state_codec_for(source: Option<(&str, &Value)>) -> faucet_core::state_version::StateCodec {
match source {
Some((kind, config)) => {
if let Some(entry) = global().sources.get(kind)
&& let Ok(built) = (entry.factory)(config.clone())
{
return faucet_core::state_version::StateCodec::for_source(built.as_ref(), false);
}
faucet_core::state_version::StateCodec {
owner: kind.to_string(),
schema: source_state_schema(kind),
legacy: false,
}
}
None => faucet_core::state_version::StateCodec {
owner: "topology".into(),
schema: 0,
legacy: false,
},
}
}
pub const IDEMPOTENT_SINK_KINDS: &[&str] = &[
"sqlite",
"postgres",
"mysql",
"mssql",
"iceberg",
"bigquery",
"kafka",
"snowflake",
"redis",
"mongodb",
"spanner",
"databricks",
"oracle",
];
pub const SCHEMA_EVOLUTION_SINK_KINDS: &[&str] = &[
"postgres",
"mysql",
"mssql",
"sqlite",
"bigquery",
"elasticsearch",
"spanner",
"iceberg",
"databricks",
"oracle",
];
pub const UPSERT_SINK_KINDS: &[&str] = &[
"postgres",
"sqlite",
"mysql",
"mssql",
"mongodb",
"elasticsearch",
"bigquery",
"spanner",
"dynamodb",
"databricks",
"oracle",
];
pub const CLEANUP_SINK_KINDS: &[&str] = &[
"postgres",
"sqlite",
"mysql",
"mssql",
"mongodb",
"elasticsearch",
"bigquery",
"spanner",
];
pub const OVERWRITE_SINK_KINDS: &[&str] = &[
"sqlite",
"postgres",
"mysql",
"mssql",
"mongodb",
"bigquery",
"elasticsearch",
"databricks",
"oracle",
"file",
"s3",
"gcs",
"azure-blob",
"sftp",
];
pub fn sink_supports_overwrite(kind: &str) -> bool {
OVERWRITE_SINK_KINDS.contains(&kind)
}
pub const TRUNCATING_FILE_SINK_KINDS: &[&str] = &["jsonl", "csv", "parquet", "file"];
pub fn sink_truncating_path<'a>(kind: &str, cfg: &'a Value) -> Option<&'a str> {
let flag = |k: &str| cfg.get(k).and_then(Value::as_bool).unwrap_or(false);
let text = |v: Option<&'a Value>| v.and_then(Value::as_str);
match kind {
"jsonl" | "csv" if !flag("append") => text(cfg.get("path")),
"parquet" => {
let dest = cfg.get("destination")?;
let path = text(dest.get("path"))?;
let rolls = ["max_rows_per_file", "max_bytes_per_file"]
.iter()
.any(|k| cfg.get(*k).is_some_and(|v| !v.is_null()));
(dest.get("type").and_then(Value::as_str) == Some("local_path")
&& path.ends_with(".parquet")
&& !rolls)
.then_some(path)
}
"file" => {
let if_exists = cfg
.get("if_exists")
.or_else(|| cfg.get("mode"))
.and_then(Value::as_str)
.unwrap_or("replace");
let staged = cfg.get("write_mode").and_then(Value::as_str) == Some("overwrite");
if matches!(if_exists, "replace" | "overwrite") && !staged {
text(cfg.get("path"))
} else {
None
}
}
_ => None,
}
}
pub const SCOPED_OVERWRITE_SINK_KINDS: &[&str] = &["postgres", "bigquery"];
pub fn sink_supports_scoped_overwrite(kind: &str) -> bool {
SCOPED_OVERWRITE_SINK_KINDS.contains(&kind)
}
pub const DIRECT_OVERWRITE_SINK_KINDS: &[&str] = &["bigquery"];
pub fn sink_supports_direct_overwrite(kind: &str) -> bool {
DIRECT_OVERWRITE_SINK_KINDS.contains(&kind)
}
pub const STAGED_LOAD_SINK_KINDS: &[&str] = &[
"redshift",
"snowflake",
"bigquery",
"clickhouse",
"mssql",
"databricks",
];
pub fn sink_supports_staged_load(kind: &str) -> bool {
STAGED_LOAD_SINK_KINDS.contains(&kind)
}
pub const DISCOVER_SOURCE_KINDS: &[&str] = &[
"postgres",
"mysql",
"mssql",
"sqlite",
"mongodb",
"elasticsearch",
"bigquery",
"snowflake",
"spanner",
"s3",
"gcs",
"iceberg",
"dynamodb",
"oracle",
"file",
];
pub fn source_supports_discover(kind: &str) -> bool {
DISCOVER_SOURCE_KINDS.contains(&kind)
}
pub fn source_replay_guarantee(kind: &str) -> faucet_core::ReplayGuarantee {
if EXACTLY_ONCE_SOURCE_KINDS.contains(&kind) {
faucet_core::ReplayGuarantee::Deterministic
} else {
faucet_core::ReplayGuarantee::NonDeterministic
}
}
pub fn sink_guarantee(kind: &str) -> faucet_core::SinkGuarantee {
if IDEMPOTENT_SINK_KINDS.contains(&kind) {
faucet_core::SinkGuarantee::AtomicWatermark
} else if UPSERT_SINK_KINDS.contains(&kind) {
faucet_core::SinkGuarantee::KeyedUpsert
} else {
faucet_core::SinkGuarantee::AtLeastOnce
}
}
pub fn source_supports_exactly_once(kind: &str) -> bool {
source_replay_guarantee(kind) == faucet_core::ReplayGuarantee::Deterministic
}
pub fn sink_supports_idempotent_writes(kind: &str) -> bool {
sink_guarantee(kind) == faucet_core::SinkGuarantee::AtomicWatermark
}
pub fn sink_supports_schema_evolution(kind: &str) -> bool {
SCHEMA_EVOLUTION_SINK_KINDS.contains(&kind)
}
pub fn sink_supports_cleanup(kind: &str) -> bool {
CLEANUP_SINK_KINDS.contains(&kind)
}
pub const SINGER_WRITE_MODES: &[faucet_core::WriteMode] = &[
faucet_core::WriteMode::Append,
faucet_core::WriteMode::Upsert,
faucet_core::WriteMode::Overwrite,
];
pub fn sink_supported_write_modes(kind: &str) -> &'static [faucet_core::WriteMode] {
use faucet_core::WriteMode;
if kind == "singer" {
return SINGER_WRITE_MODES;
}
match (
UPSERT_SINK_KINDS.contains(&kind),
OVERWRITE_SINK_KINDS.contains(&kind),
) {
(true, true) => &[
WriteMode::Append,
WriteMode::Upsert,
WriteMode::Delete,
WriteMode::Overwrite,
],
(true, false) => &[WriteMode::Append, WriteMode::Upsert, WriteMode::Delete],
(false, true) => &[WriteMode::Append, WriteMode::Overwrite],
(false, false) => &[WriteMode::Append],
}
}
fn reject_unknown_config_keys(
role: &'static str,
kind: &str,
name: &str,
config: &Value,
schema: &Value,
) -> CliResult<()> {
let Some(given) = config.as_object() else {
return Ok(());
};
let mut known = BTreeSet::new();
if !collect_schema_keys(schema, schema, &mut known, 0) {
return Ok(());
}
if known.is_empty() {
return Ok(());
}
let unknown: Vec<&str> = given
.keys()
.map(String::as_str)
.filter(|k| !known.contains(*k))
.collect();
if unknown.is_empty() {
return Ok(());
}
let suggestions: Vec<String> = unknown
.iter()
.filter_map(|u| {
nearest_key(u, known.iter()).map(|k| format!("`{u}` — did you mean `{k}`?"))
})
.collect();
let hint = if suggestions.is_empty() {
String::new()
} else {
format!(" ({})", suggestions.join("; "))
};
Err(CliError::InvalidConnectorConfig {
kind: if role == "source" { "source" } else { "sink" },
name: name.to_owned(),
message: format!(
"unknown {role} `{kind}` config key(s): {}{hint}. A key the connector does not \
declare is silently ignored, so an integrity or batching knob would read as set \
while doing nothing — run `faucet schema {role} {kind}` for the full list.",
unknown
.iter()
.map(|k| format!("`{k}`"))
.collect::<Vec<_>>()
.join(", ")
),
})
}
fn collect_schema_keys(
schema: &Value,
root: &Value,
out: &mut BTreeSet<String>,
depth: usize,
) -> bool {
const MAX_DEPTH: usize = 8;
if depth > MAX_DEPTH {
return false;
}
let Some(obj) = schema.as_object() else {
return false;
};
match obj.get("additionalProperties") {
None | Some(Value::Bool(false)) => {}
Some(_) => return false,
}
if let Some(props) = obj.get("properties").and_then(Value::as_object) {
out.extend(props.keys().cloned());
}
if let Some(aliases) = obj.get("x-faucet-aliases").and_then(Value::as_array) {
out.extend(aliases.iter().filter_map(Value::as_str).map(str::to_string));
}
if let Some(r) = obj.get("$ref").and_then(Value::as_str) {
let Some(target) = resolve_local_ref(root, r) else {
return false;
};
if !collect_schema_keys(target, root, out, depth + 1) {
return false;
}
}
for key in ["allOf", "anyOf", "oneOf"] {
if let Some(branches) = obj.get(key).and_then(Value::as_array) {
for b in branches {
if !collect_schema_keys(b, root, out, depth + 1) {
return false;
}
}
}
}
true
}
fn resolve_local_ref<'a>(root: &'a Value, r: &str) -> Option<&'a Value> {
let name = r.strip_prefix("#/$defs/")?;
root.get("$defs")?.get(name)
}
fn nearest_key<'a, I: Iterator<Item = &'a String>>(given: &str, known: I) -> Option<&'a str> {
let budget = (given.len() / 3).max(2);
known
.map(|k| (edit_distance(given, k), k.as_str()))
.filter(|(d, _)| *d <= budget)
.min_by_key(|(d, k)| (*d, k.len()))
.map(|(_, k)| k)
}
fn edit_distance(a: &str, b: &str) -> usize {
let (a, b) = (a.as_bytes(), b.as_bytes());
let mut prev: Vec<usize> = (0..=b.len()).collect();
let mut cur = vec![0usize; b.len() + 1];
for (i, &ac) in a.iter().enumerate() {
cur[0] = i + 1;
for (j, &bc) in b.iter().enumerate() {
let cost = usize::from(ac != bc);
cur[j + 1] = (prev[j] + cost).min(prev[j + 1] + 1).min(cur[j] + 1);
}
std::mem::swap(&mut prev, &mut cur);
}
prev[b.len()]
}
const CONCURRENCY_KNOBS: [&str; 5] = [
"max_connections",
"request_concurrency",
"partition_concurrency",
"shard_concurrency",
"concurrency",
];
pub fn apply_concurrency_override(
schema: &Value,
config: &mut Value,
n: usize,
) -> Option<&'static str> {
let declared = schema.get("properties")?.as_object()?;
let knob = CONCURRENCY_KNOBS
.into_iter()
.find(|k| declared.contains_key(*k))?;
config.as_object_mut()?.insert(
knob.to_string(),
Value::Number(serde_json::Number::from(n as u64)),
);
Some(knob)
}
pub fn override_source_concurrency(
kind: &str,
config: &mut Value,
n: usize,
) -> Option<&'static str> {
let schema = source_schema(kind).ok()?;
apply_concurrency_override(&schema, config, n)
}
pub fn override_sink_concurrency(kind: &str, config: &mut Value, n: usize) -> Option<&'static str> {
let schema = sink_schema(kind).ok()?;
apply_concurrency_override(&schema, config, n)
}
fn check<T: DeserializeOwned>(kind: &'static str, name: &str, config: Value) -> CliResult<()> {
decode::<T>(kind, name, config).map(|_| ())
}
fn check_with<T, E, F>(kind: &'static str, name: &str, config: Value, validate: F) -> CliResult<()>
where
T: DeserializeOwned,
E: std::fmt::Display,
F: Fn(&T) -> Result<(), E>,
{
let cfg = decode::<T>(kind, name, config)?;
validate(&cfg).map_err(|e| CliError::InvalidConnectorConfig {
kind,
name: name.to_owned(),
message: e.to_string(),
})
}
pub fn validate_source_config(kind: &str, name: &str, config: Value) -> CliResult<()> {
if let Some(entry) = global().sources.get(kind) {
reject_unknown_config_keys("source", kind, name, &config, &(entry.schema)())?;
return (entry.factory)(config).map(|_| ());
}
if let Ok(schema) = source_schema(kind) {
reject_unknown_config_keys("source", kind, name, &config, &schema)?;
}
match kind {
#[cfg(feature = "source-rest")]
"rest" => {
check_with::<faucet_source_rest::RestStreamConfig, _, _>("rest", name, config, |c| {
c.validate()
})
}
#[cfg(feature = "source-graphql")]
"graphql" => check_with::<faucet_source_graphql::GraphqlStreamConfig, _, _>(
"graphql",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-xml")]
"xml" => check_with::<faucet_source_xml::XmlStreamConfig, _, _>("xml", name, config, |c| {
c.validate()
}),
#[cfg(feature = "source-grpc")]
"grpc" => check::<faucet_source_grpc::GrpcStreamConfig>("grpc", name, config),
#[cfg(feature = "source-postgres")]
"postgres" => {
check::<faucet_source_postgres::PostgresSourceConfig>("postgres", name, config)
}
#[cfg(feature = "source-postgres-cdc")]
"postgres-cdc" => check_with::<faucet_source_postgres_cdc::PostgresCdcSourceConfig, _, _>(
"postgres-cdc",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-mysql")]
"mysql" => check::<faucet_source_mysql::MysqlSourceConfig>("mysql", name, config),
#[cfg(feature = "source-mssql")]
"mssql" => {
check_with::<faucet_source_mssql::MssqlSourceConfig, _, _>("mssql", name, config, |c| {
c.validate()
})
}
#[cfg(feature = "source-sqlite")]
"sqlite" => check::<faucet_source_sqlite::SqliteSourceConfig>("sqlite", name, config),
#[cfg(feature = "source-duckdb")]
"duckdb" => check_with::<faucet_source_duckdb::DuckdbSourceConfig, _, _>(
"duckdb",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-sqs")]
"sqs" => check_with::<faucet_source_sqs::SqsSourceConfig, _, _>("sqs", name, config, |c| {
c.validate()
}),
#[cfg(feature = "source-nats")]
"nats" => {
check_with::<faucet_source_nats::NatsSourceConfig, _, _>("nats", name, config, |c| {
c.validate()
})
}
#[cfg(feature = "source-rabbitmq")]
"rabbitmq" => check_with::<faucet_source_rabbitmq::RabbitMqSourceConfig, _, _>(
"rabbitmq",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-sftp")]
"sftp" => check::<faucet_source_sftp::SftpSourceConfig>("sftp", name, config),
#[cfg(feature = "source-file")]
"file" => {
check_with::<faucet_source_file::FileSourceConfig, _, _>("file", name, config, |c| {
c.validate()
})
}
#[cfg(feature = "source-s3")]
"s3" => check::<faucet_source_s3::S3SourceConfig>("s3", name, config),
#[cfg(feature = "source-mongodb")]
"mongodb" => check::<faucet_source_mongodb::MongoSourceConfig>("mongodb", name, config),
#[cfg(feature = "source-mongodb-cdc")]
"mongodb-cdc" => check_with::<faucet_source_mongodb_cdc::MongoCdcSourceConfig, _, _>(
"mongodb-cdc",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-mysql-cdc")]
"mysql-cdc" => check_with::<faucet_source_mysql_cdc::MysqlCdcSourceConfig, _, _>(
"mysql-cdc",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-redis")]
"redis" => check::<faucet_source_redis::RedisSourceConfig>("redis", name, config),
#[cfg(feature = "source-webhook")]
"webhook" => check::<faucet_source_webhook::WebhookSourceConfig>("webhook", name, config),
#[cfg(feature = "source-websocket")]
"websocket" => check_with::<faucet_source_websocket::WebsocketSourceConfig, _, _>(
"websocket",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-csv")]
"csv" => check_with::<faucet_source_csv::CsvSourceConfig, _, _>("csv", name, config, |c| {
c.validate()
}),
#[cfg(feature = "source-singer")]
"singer" => check::<faucet_source_singer::SingerSourceConfig>("singer", name, config),
#[cfg(feature = "source-elasticsearch")]
"elasticsearch" => check::<faucet_source_elasticsearch::ElasticsearchSourceConfig>(
"elasticsearch",
name,
config,
),
#[cfg(feature = "source-kafka")]
"kafka" => {
check_with::<faucet_source_kafka::KafkaSourceConfig, _, _>("kafka", name, config, |c| {
c.validate()
})
}
#[cfg(feature = "source-kinesis")]
"kinesis" => check_with::<faucet_source_kinesis::KinesisSourceConfig, _, _>(
"kinesis",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-spanner")]
"spanner" => check_with::<faucet_source_spanner::SpannerSourceConfig, _, _>(
"spanner",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-parquet")]
"parquet" => check::<faucet_source_parquet::ParquetSourceConfig>("parquet", name, config),
#[cfg(feature = "source-delta")]
"delta" => {
check_with::<faucet_source_delta::DeltaSourceConfig, _, _>("delta", name, config, |c| {
c.validate()
})
}
#[cfg(feature = "source-databricks")]
"databricks" => check_with::<faucet_source_databricks::DatabricksSourceConfig, _, _>(
"databricks",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-iceberg")]
"iceberg" => check_with::<faucet_source_iceberg::IcebergSourceConfig, _, _>(
"iceberg",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-dynamodb")]
"dynamodb" => check_with::<faucet_source_dynamodb::DynamoDbSourceConfig, _, _>(
"dynamodb",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-oracle")]
"oracle" => check_with::<faucet_source_oracle::OracleSourceConfig, _, _>(
"oracle",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-oracle-cdc")]
"oracle-cdc" => check_with::<faucet_source_oracle_cdc::OracleCdcSourceConfig, _, _>(
"oracle-cdc",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-gcs")]
"gcs" => check_with::<faucet_source_gcs::GcsSourceConfig, _, _>("gcs", name, config, |c| {
c.validate()
}),
#[cfg(feature = "source-bigquery")]
"bigquery" => {
check::<faucet_source_bigquery::BigQuerySourceConfig>("bigquery", name, config)
}
#[cfg(feature = "source-snowflake")]
"snowflake" => {
check::<faucet_source_snowflake::SnowflakeSourceConfig>("snowflake", name, config)
}
#[cfg(feature = "source-mssql-cdc")]
"mssql-cdc" => check_with::<faucet_source_mssql_cdc::MssqlCdcSourceConfig, _, _>(
"mssql-cdc",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-redshift")]
"redshift" => check_with::<faucet_source_redshift::RedshiftSourceConfig, _, _>(
"redshift",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-pubsub")]
"pubsub" => check_with::<faucet_source_pubsub::PubsubSourceConfig, _, _>(
"pubsub",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-clickhouse")]
"clickhouse" => check_with::<faucet_source_clickhouse::ClickHouseSourceConfig, _, _>(
"clickhouse",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "source-azure-blob")]
"azure-blob" => check_with::<faucet_source_azure_blob::AzureBlobSourceConfig, _, _>(
"azure-blob",
name,
config,
|c| c.validate(),
),
other => Err(unknown(other, "source", source_kinds())),
}
}
pub fn sink_batch_atomicity(kind: &str, config: &Value) -> Option<faucet_core::BatchAtomicity> {
if let Some(entry) = global().sinks.get(kind) {
return (entry.factory)(config.clone())
.ok()
.map(|s| s.batch_atomicity());
}
#[allow(dead_code)]
fn atomicity_of<T: serde::de::DeserializeOwned>(
config: &Value,
f: impl Fn(&T) -> faucet_core::BatchAtomicity,
) -> Option<faucet_core::BatchAtomicity> {
serde_json::from_value::<T>(config.clone())
.ok()
.map(|c| f(&c))
}
match kind {
#[cfg(feature = "sink-bigquery")]
"bigquery" => atomicity_of::<faucet_sink_bigquery::BigQuerySinkConfig>(config, |c| {
c.batch_atomicity()
}),
#[cfg(feature = "sink-iceberg")]
"iceberg" => {
atomicity_of::<faucet_sink_iceberg::IcebergSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-delta")]
"delta" => {
atomicity_of::<faucet_sink_delta::DeltaSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-postgres")]
"postgres" => atomicity_of::<faucet_sink_postgres::PostgresSinkConfig>(config, |c| {
c.batch_atomicity()
}),
#[cfg(feature = "sink-jsonl")]
"jsonl" => {
atomicity_of::<faucet_sink_jsonl::JsonlSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-snowflake")]
"snowflake" => atomicity_of::<faucet_sink_snowflake::SnowflakeSinkConfig>(config, |c| {
c.batch_atomicity()
}),
#[cfg(feature = "sink-mysql")]
"mysql" => {
atomicity_of::<faucet_sink_mysql::MysqlSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-mssql")]
"mssql" => {
atomicity_of::<faucet_sink_mssql::MssqlSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-sqlite")]
"sqlite" => {
atomicity_of::<faucet_sink_sqlite::SqliteSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-duckdb")]
"duckdb" => {
atomicity_of::<faucet_sink_duckdb::DuckdbSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-sqs")]
"sqs" => atomicity_of::<faucet_sink_sqs::SqsSinkConfig>(config, |c| c.batch_atomicity()),
#[cfg(feature = "sink-nats")]
"nats" => atomicity_of::<faucet_sink_nats::NatsSinkConfig>(config, |c| c.batch_atomicity()),
#[cfg(feature = "sink-rabbitmq")]
"rabbitmq" => atomicity_of::<faucet_sink_rabbitmq::RabbitMqSinkConfig>(config, |c| {
c.batch_atomicity()
}),
#[cfg(feature = "sink-sftp")]
"sftp" => atomicity_of::<faucet_sink_sftp::SftpSinkConfig>(config, |c| c.batch_atomicity()),
#[cfg(feature = "sink-singer")]
"singer" => {
atomicity_of::<faucet_sink_singer::SingerSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-s3")]
"s3" => atomicity_of::<faucet_sink_s3::S3SinkConfig>(config, |c| c.batch_atomicity()),
#[cfg(feature = "sink-mongodb")]
"mongodb" => {
atomicity_of::<faucet_sink_mongodb::MongoSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-redis")]
"redis" => {
atomicity_of::<faucet_sink_redis::RedisSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-csv")]
"csv" => atomicity_of::<faucet_sink_csv::CsvSinkConfig>(config, |c| c.batch_atomicity()),
#[cfg(feature = "sink-elasticsearch")]
"elasticsearch" => {
atomicity_of::<faucet_sink_elasticsearch::ElasticsearchSinkConfig>(config, |c| {
c.batch_atomicity()
})
}
#[cfg(feature = "sink-kafka")]
"kafka" => {
atomicity_of::<faucet_sink_kafka::KafkaSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-kinesis")]
"kinesis" => {
atomicity_of::<faucet_sink_kinesis::KinesisSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-spanner")]
"spanner" => {
atomicity_of::<faucet_sink_spanner::SpannerSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-http")]
"http" => atomicity_of::<faucet_sink_http::HttpSinkConfig>(config, |c| c.batch_atomicity()),
#[cfg(feature = "sink-stdout")]
"stdout" => {
atomicity_of::<faucet_sink_stdout::StdoutSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-parquet")]
"parquet" => {
atomicity_of::<faucet_sink_parquet::ParquetSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-file")]
"file" => atomicity_of::<faucet_sink_file::FileSinkConfig>(config, |c| c.batch_atomicity()),
#[cfg(feature = "sink-gcs")]
"gcs" => atomicity_of::<faucet_sink_gcs::GcsSinkConfig>(config, |c| c.batch_atomicity()),
#[cfg(feature = "sink-redshift")]
"redshift" => atomicity_of::<faucet_sink_redshift::RedshiftSinkConfig>(config, |c| {
c.batch_atomicity()
}),
#[cfg(feature = "sink-pubsub")]
"pubsub" => {
atomicity_of::<faucet_sink_pubsub::PubsubSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-clickhouse")]
"clickhouse" => atomicity_of::<faucet_sink_clickhouse::ClickHouseSinkConfig>(config, |c| {
c.batch_atomicity()
}),
#[cfg(feature = "sink-azure-blob")]
"azure-blob" => atomicity_of::<faucet_sink_azure_blob::AzureBlobSinkConfig>(config, |c| {
c.batch_atomicity()
}),
#[cfg(feature = "sink-oracle")]
"oracle" => {
atomicity_of::<faucet_sink_oracle::OracleSinkConfig>(config, |c| c.batch_atomicity())
}
#[cfg(feature = "sink-dynamodb")]
"dynamodb" => atomicity_of::<faucet_sink_dynamodb::DynamoDbSinkConfig>(config, |c| {
c.batch_atomicity()
}),
#[cfg(feature = "sink-databricks")]
"databricks" => atomicity_of::<faucet_sink_databricks::DatabricksSinkConfig>(config, |c| {
c.batch_atomicity()
}),
_ => None,
}
}
pub fn validate_sink_config(kind: &str, name: &str, config: Value) -> CliResult<()> {
if let Some(entry) = global().sinks.get(kind) {
reject_unknown_config_keys("sink", kind, name, &config, &(entry.schema)())?;
return (entry.factory)(config).map(|_| ());
}
if let Ok(schema) = sink_schema(kind) {
reject_unknown_config_keys("sink", kind, name, &config, &schema)?;
}
match kind {
#[cfg(feature = "sink-bigquery")]
"bigquery" => check::<faucet_sink_bigquery::BigQuerySinkConfig>("bigquery", name, config),
#[cfg(feature = "sink-iceberg")]
"iceberg" => check_with::<faucet_sink_iceberg::IcebergSinkConfig, _, _>(
"iceberg",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "sink-delta")]
"delta" => {
check_with::<faucet_sink_delta::DeltaSinkConfig, _, _>("delta", name, config, |c| {
c.validate()
})
}
#[cfg(feature = "sink-postgres")]
"postgres" => check::<faucet_sink_postgres::PostgresSinkConfig>("postgres", name, config),
#[cfg(feature = "sink-jsonl")]
"jsonl" => check::<faucet_sink_jsonl::JsonlSinkConfig>("jsonl", name, config),
#[cfg(feature = "sink-snowflake")]
"snowflake" => {
check::<faucet_sink_snowflake::SnowflakeSinkConfig>("snowflake", name, config)
}
#[cfg(feature = "sink-mysql")]
"mysql" => check::<faucet_sink_mysql::MysqlSinkConfig>("mysql", name, config),
#[cfg(feature = "sink-mssql")]
"mssql" => {
check_with::<faucet_sink_mssql::MssqlSinkConfig, _, _>("mssql", name, config, |c| {
c.validate()
})
}
#[cfg(feature = "sink-sqlite")]
"sqlite" => check::<faucet_sink_sqlite::SqliteSinkConfig>("sqlite", name, config),
#[cfg(feature = "sink-duckdb")]
"duckdb" => check::<faucet_sink_duckdb::DuckdbSinkConfig>("duckdb", name, config),
#[cfg(feature = "sink-sqs")]
"sqs" => check_with::<faucet_sink_sqs::SqsSinkConfig, _, _>("sqs", name, config, |c| {
c.validate()
}),
#[cfg(feature = "sink-nats")]
"nats" => check_with::<faucet_sink_nats::NatsSinkConfig, _, _>("nats", name, config, |c| {
c.validate()
}),
#[cfg(feature = "sink-rabbitmq")]
"rabbitmq" => check_with::<faucet_sink_rabbitmq::RabbitMqSinkConfig, _, _>(
"rabbitmq",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "sink-sftp")]
"sftp" => check::<faucet_sink_sftp::SftpSinkConfig>("sftp", name, config),
#[cfg(feature = "sink-singer")]
"singer" => {
check_with::<faucet_sink_singer::SingerSinkConfig, _, _>("singer", name, config, |c| {
c.validate()
})
}
#[cfg(feature = "sink-s3")]
"s3" => {
check_with::<faucet_sink_s3::S3SinkConfig, _, _>("s3", name, config, |c| c.validate())
}
#[cfg(feature = "sink-mongodb")]
"mongodb" => check::<faucet_sink_mongodb::MongoSinkConfig>("mongodb", name, config),
#[cfg(feature = "sink-redis")]
"redis" => check::<faucet_sink_redis::RedisSinkConfig>("redis", name, config),
#[cfg(feature = "sink-csv")]
"csv" => check::<faucet_sink_csv::CsvSinkConfig>("csv", name, config),
#[cfg(feature = "sink-elasticsearch")]
"elasticsearch" => check::<faucet_sink_elasticsearch::ElasticsearchSinkConfig>(
"elasticsearch",
name,
config,
),
#[cfg(feature = "sink-kafka")]
"kafka" => {
check_with::<faucet_sink_kafka::KafkaSinkConfig, _, _>("kafka", name, config, |c| {
c.validate()
})
}
#[cfg(feature = "sink-kinesis")]
"kinesis" => check_with::<faucet_sink_kinesis::KinesisSinkConfig, _, _>(
"kinesis",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "sink-spanner")]
"spanner" => check_with::<faucet_sink_spanner::SpannerSinkConfig, _, _>(
"spanner",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "sink-http")]
"http" => check::<faucet_sink_http::HttpSinkConfig>("http", name, config),
#[cfg(feature = "sink-stdout")]
"stdout" => check::<faucet_sink_stdout::StdoutSinkConfig>("stdout", name, config),
#[cfg(feature = "sink-parquet")]
"parquet" => check_with::<faucet_sink_parquet::ParquetSinkConfig, _, _>(
"parquet",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "sink-file")]
"file" => check_with::<faucet_sink_file::FileSinkConfig, _, _>("file", name, config, |c| {
c.validate()
}),
#[cfg(feature = "sink-gcs")]
"gcs" => check_with::<faucet_sink_gcs::GcsSinkConfig, _, _>("gcs", name, config, |c| {
c.validate()
}),
#[cfg(feature = "sink-redshift")]
"redshift" => check_with::<faucet_sink_redshift::RedshiftSinkConfig, _, _>(
"redshift",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "sink-pubsub")]
"pubsub" => {
check_with::<faucet_sink_pubsub::PubsubSinkConfig, _, _>("pubsub", name, config, |c| {
c.validate()
})
}
#[cfg(feature = "sink-clickhouse")]
"clickhouse" => check_with::<faucet_sink_clickhouse::ClickHouseSinkConfig, _, _>(
"clickhouse",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "sink-azure-blob")]
"azure-blob" => {
check::<faucet_sink_azure_blob::AzureBlobSinkConfig>("azure-blob", name, config)
}
#[cfg(feature = "sink-dynamodb")]
"dynamodb" => check_with::<faucet_sink_dynamodb::DynamoDbSinkConfig, _, _>(
"dynamodb",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "sink-databricks")]
"databricks" => check_with::<faucet_sink_databricks::DatabricksSinkConfig, _, _>(
"databricks",
name,
config,
|c| c.validate(),
),
#[cfg(feature = "sink-oracle")]
"oracle" => {
check_with::<faucet_sink_oracle::OracleSinkConfig, _, _>("oracle", name, config, |c| {
c.validate()
})
}
other => Err(unknown(other, "sink", sink_kinds())),
}
}
pub fn source_schema(kind: &str) -> CliResult<Value> {
if let Some(entry) = global().sources.get(kind) {
return Ok((entry.schema)());
}
match kind {
#[cfg(feature = "source-rest")]
"rest" => Ok(schema::<faucet_source_rest::RestStreamConfig>()),
#[cfg(feature = "source-graphql")]
"graphql" => Ok(schema::<faucet_source_graphql::GraphqlStreamConfig>()),
#[cfg(feature = "source-xml")]
"xml" => Ok(schema::<faucet_source_xml::XmlStreamConfig>()),
#[cfg(feature = "source-grpc")]
"grpc" => Ok(schema::<faucet_source_grpc::GrpcStreamConfig>()),
#[cfg(feature = "source-postgres")]
"postgres" => Ok(schema::<faucet_source_postgres::PostgresSourceConfig>()),
#[cfg(feature = "source-postgres-cdc")]
"postgres-cdc" => Ok(schema::<faucet_source_postgres_cdc::PostgresCdcSourceConfig>()),
#[cfg(feature = "source-mysql")]
"mysql" => Ok(schema::<faucet_source_mysql::MysqlSourceConfig>()),
#[cfg(feature = "source-mssql")]
"mssql" => Ok(schema::<faucet_source_mssql::MssqlSourceConfig>()),
#[cfg(feature = "source-sqlite")]
"sqlite" => Ok(schema::<faucet_source_sqlite::SqliteSourceConfig>()),
#[cfg(feature = "source-duckdb")]
"duckdb" => Ok(schema::<faucet_source_duckdb::DuckdbSourceConfig>()),
#[cfg(feature = "source-sqs")]
"sqs" => Ok(schema::<faucet_source_sqs::SqsSourceConfig>()),
#[cfg(feature = "source-nats")]
"nats" => Ok(schema::<faucet_source_nats::NatsSourceConfig>()),
#[cfg(feature = "source-rabbitmq")]
"rabbitmq" => Ok(schema::<faucet_source_rabbitmq::RabbitMqSourceConfig>()),
#[cfg(feature = "source-sftp")]
"sftp" => Ok(schema::<faucet_source_sftp::SftpSourceConfig>()),
#[cfg(feature = "source-file")]
"file" => Ok(schema::<faucet_source_file::FileSourceConfig>()),
#[cfg(feature = "source-s3")]
"s3" => Ok(schema::<faucet_source_s3::S3SourceConfig>()),
#[cfg(feature = "source-mongodb")]
"mongodb" => Ok(schema::<faucet_source_mongodb::MongoSourceConfig>()),
#[cfg(feature = "source-mongodb-cdc")]
"mongodb-cdc" => Ok(schema::<faucet_source_mongodb_cdc::MongoCdcSourceConfig>()),
#[cfg(feature = "source-mysql-cdc")]
"mysql-cdc" => Ok(schema::<faucet_source_mysql_cdc::MysqlCdcSourceConfig>()),
#[cfg(feature = "source-redis")]
"redis" => Ok(schema::<faucet_source_redis::RedisSourceConfig>()),
#[cfg(feature = "source-webhook")]
"webhook" => Ok(schema::<faucet_source_webhook::WebhookSourceConfig>()),
#[cfg(feature = "source-websocket")]
"websocket" => Ok(schema::<faucet_source_websocket::WebsocketSourceConfig>()),
#[cfg(feature = "source-csv")]
"csv" => Ok(schema::<faucet_source_csv::CsvSourceConfig>()),
#[cfg(feature = "source-singer")]
"singer" => Ok(schema::<faucet_source_singer::SingerSourceConfig>()),
#[cfg(feature = "source-elasticsearch")]
"elasticsearch" => Ok(schema::<
faucet_source_elasticsearch::ElasticsearchSourceConfig,
>()),
#[cfg(feature = "source-kafka")]
"kafka" => Ok(schema::<faucet_source_kafka::KafkaSourceConfig>()),
#[cfg(feature = "source-kinesis")]
"kinesis" => Ok(schema::<faucet_source_kinesis::KinesisSourceConfig>()),
#[cfg(feature = "source-spanner")]
"spanner" => Ok(schema::<faucet_source_spanner::SpannerSourceConfig>()),
#[cfg(feature = "source-parquet")]
"parquet" => Ok(schema::<faucet_source_parquet::ParquetSourceConfig>()),
#[cfg(feature = "source-delta")]
"delta" => Ok(schema::<faucet_source_delta::DeltaSourceConfig>()),
#[cfg(feature = "source-databricks")]
"databricks" => Ok(schema::<faucet_source_databricks::DatabricksSourceConfig>()),
#[cfg(feature = "source-iceberg")]
"iceberg" => Ok(schema::<faucet_source_iceberg::IcebergSourceConfig>()),
#[cfg(feature = "source-dynamodb")]
"dynamodb" => Ok(schema::<faucet_source_dynamodb::DynamoDbSourceConfig>()),
#[cfg(feature = "source-oracle")]
"oracle" => Ok(schema::<faucet_source_oracle::OracleSourceConfig>()),
#[cfg(feature = "source-oracle-cdc")]
"oracle-cdc" => Ok(schema::<faucet_source_oracle_cdc::OracleCdcSourceConfig>()),
#[cfg(feature = "source-gcs")]
"gcs" => Ok(schema::<faucet_source_gcs::GcsSourceConfig>()),
#[cfg(feature = "source-bigquery")]
"bigquery" => Ok(schema::<faucet_source_bigquery::BigQuerySourceConfig>()),
#[cfg(feature = "source-snowflake")]
"snowflake" => Ok(schema::<faucet_source_snowflake::SnowflakeSourceConfig>()),
#[cfg(feature = "source-mssql-cdc")]
"mssql-cdc" => Ok(schema::<faucet_source_mssql_cdc::MssqlCdcSourceConfig>()),
#[cfg(feature = "source-redshift")]
"redshift" => Ok(schema::<faucet_source_redshift::RedshiftSourceConfig>()),
#[cfg(feature = "source-pubsub")]
"pubsub" => Ok(schema::<faucet_source_pubsub::PubsubSourceConfig>()),
#[cfg(feature = "source-clickhouse")]
"clickhouse" => Ok(schema::<faucet_source_clickhouse::ClickHouseSourceConfig>()),
#[cfg(feature = "source-azure-blob")]
"azure-blob" => Ok(schema::<faucet_source_azure_blob::AzureBlobSourceConfig>()),
other => Err(unknown(other, "source", source_kinds())),
}
}
pub fn source_exists(kind: &str) -> bool {
source_schema(kind).is_ok()
}
pub fn sink_exists(kind: &str) -> bool {
sink_schema(kind).is_ok()
}
pub fn sink_schema(kind: &str) -> CliResult<Value> {
if let Some(entry) = global().sinks.get(kind) {
return Ok((entry.schema)());
}
match kind {
#[cfg(feature = "sink-bigquery")]
"bigquery" => Ok(schema::<faucet_sink_bigquery::BigQuerySinkConfig>()),
#[cfg(feature = "sink-iceberg")]
"iceberg" => Ok(schema::<faucet_sink_iceberg::IcebergSinkConfig>()),
#[cfg(feature = "sink-delta")]
"delta" => Ok(schema::<faucet_sink_delta::DeltaSinkConfig>()),
#[cfg(feature = "sink-postgres")]
"postgres" => Ok(schema::<faucet_sink_postgres::PostgresSinkConfig>()),
#[cfg(feature = "sink-jsonl")]
"jsonl" => Ok(schema::<faucet_sink_jsonl::JsonlSinkConfig>()),
#[cfg(feature = "sink-snowflake")]
"snowflake" => Ok(schema::<faucet_sink_snowflake::SnowflakeSinkConfig>()),
#[cfg(feature = "sink-mysql")]
"mysql" => Ok(schema::<faucet_sink_mysql::MysqlSinkConfig>()),
#[cfg(feature = "sink-mssql")]
"mssql" => Ok(schema::<faucet_sink_mssql::MssqlSinkConfig>()),
#[cfg(feature = "sink-sqlite")]
"sqlite" => Ok(schema::<faucet_sink_sqlite::SqliteSinkConfig>()),
#[cfg(feature = "sink-duckdb")]
"duckdb" => Ok(schema::<faucet_sink_duckdb::DuckdbSinkConfig>()),
#[cfg(feature = "sink-sqs")]
"sqs" => Ok(schema::<faucet_sink_sqs::SqsSinkConfig>()),
#[cfg(feature = "sink-nats")]
"nats" => Ok(schema::<faucet_sink_nats::NatsSinkConfig>()),
#[cfg(feature = "sink-rabbitmq")]
"rabbitmq" => Ok(schema::<faucet_sink_rabbitmq::RabbitMqSinkConfig>()),
#[cfg(feature = "sink-sftp")]
"sftp" => Ok(schema::<faucet_sink_sftp::SftpSinkConfig>()),
#[cfg(feature = "sink-singer")]
"singer" => Ok(schema::<faucet_sink_singer::SingerSinkConfig>()),
#[cfg(feature = "sink-s3")]
"s3" => Ok(schema::<faucet_sink_s3::S3SinkConfig>()),
#[cfg(feature = "sink-mongodb")]
"mongodb" => Ok(schema::<faucet_sink_mongodb::MongoSinkConfig>()),
#[cfg(feature = "sink-redis")]
"redis" => Ok(schema::<faucet_sink_redis::RedisSinkConfig>()),
#[cfg(feature = "sink-csv")]
"csv" => Ok(schema::<faucet_sink_csv::CsvSinkConfig>()),
#[cfg(feature = "sink-elasticsearch")]
"elasticsearch" => Ok(schema::<faucet_sink_elasticsearch::ElasticsearchSinkConfig>()),
#[cfg(feature = "sink-kafka")]
"kafka" => Ok(schema::<faucet_sink_kafka::KafkaSinkConfig>()),
#[cfg(feature = "sink-kinesis")]
"kinesis" => Ok(schema::<faucet_sink_kinesis::KinesisSinkConfig>()),
#[cfg(feature = "sink-spanner")]
"spanner" => Ok(schema::<faucet_sink_spanner::SpannerSinkConfig>()),
#[cfg(feature = "sink-http")]
"http" => Ok(schema::<faucet_sink_http::HttpSinkConfig>()),
#[cfg(feature = "sink-stdout")]
"stdout" => Ok(schema::<faucet_sink_stdout::StdoutSinkConfig>()),
#[cfg(feature = "sink-parquet")]
"parquet" => Ok(schema::<faucet_sink_parquet::ParquetSinkConfig>()),
#[cfg(feature = "sink-file")]
"file" => Ok(schema::<faucet_sink_file::FileSinkConfig>()),
#[cfg(feature = "sink-gcs")]
"gcs" => Ok(schema::<faucet_sink_gcs::GcsSinkConfig>()),
#[cfg(feature = "sink-redshift")]
"redshift" => Ok(schema::<faucet_sink_redshift::RedshiftSinkConfig>()),
#[cfg(feature = "sink-pubsub")]
"pubsub" => Ok(schema::<faucet_sink_pubsub::PubsubSinkConfig>()),
#[cfg(feature = "sink-clickhouse")]
"clickhouse" => Ok(schema::<faucet_sink_clickhouse::ClickHouseSinkConfig>()),
#[cfg(feature = "sink-azure-blob")]
"azure-blob" => Ok(schema::<faucet_sink_azure_blob::AzureBlobSinkConfig>()),
#[cfg(feature = "sink-dynamodb")]
"dynamodb" => Ok(schema::<faucet_sink_dynamodb::DynamoDbSinkConfig>()),
#[cfg(feature = "sink-databricks")]
"databricks" => Ok(schema::<faucet_sink_databricks::DatabricksSinkConfig>()),
#[cfg(feature = "sink-oracle")]
"oracle" => Ok(schema::<faucet_sink_oracle::OracleSinkConfig>()),
other => Err(unknown(other, "sink", sink_kinds())),
}
}
#[allow(dead_code)]
fn deprecated(kind: &str) -> &'static str {
crate::vocabulary::deprecated_kind_description(kind).unwrap_or("Deprecated connector.")
}
pub fn source_descriptions() -> Vec<(&'static str, &'static str)> {
let mut v = builtin_source_descriptions();
v.extend(global().custom_source_descriptions());
v
}
#[allow(clippy::vec_init_then_push)]
fn builtin_source_descriptions() -> Vec<(&'static str, &'static str)> {
let mut v: Vec<(&'static str, &'static str)> = Vec::new();
#[cfg(feature = "source-rest")]
v.push(("rest", "REST API source with pagination, auth, transforms"));
#[cfg(feature = "source-graphql")]
v.push(("graphql", "GraphQL API source with cursor pagination"));
#[cfg(feature = "source-xml")]
v.push(("xml", "XML / SOAP API source with XML→JSON conversion"));
#[cfg(feature = "source-grpc")]
v.push(("grpc", "gRPC source with dynamic protobuf"));
#[cfg(feature = "source-postgres")]
v.push(("postgres", "PostgreSQL query source"));
#[cfg(feature = "source-postgres-cdc")]
v.push((
"postgres-cdc",
"PostgreSQL CDC source (logical replication)",
));
#[cfg(feature = "source-mysql")]
v.push(("mysql", "MySQL query source"));
#[cfg(feature = "source-mssql")]
v.push(("mssql", "Microsoft SQL Server query source"));
#[cfg(feature = "source-sqlite")]
v.push(("sqlite", "SQLite query source"));
#[cfg(feature = "source-duckdb")]
v.push((
"duckdb",
"DuckDB query source. Runs SQL against a DuckDB file or in-memory database and streams rows as JSON with bounded memory.",
));
#[cfg(feature = "source-sqs")]
v.push((
"sqs",
"AWS SQS source. Long-polls ReceiveMessage, deletes after the batch is emitted (at-least-once), with idle/max-messages termination.",
));
#[cfg(feature = "source-nats")]
v.push((
"nats",
"NATS source. Subscribes to a subject (or a JetStream durable consumer) and drains with idle/max-messages termination.",
));
#[cfg(feature = "source-rabbitmq")]
v.push((
"rabbitmq",
"RabbitMQ (AMQP 0.9.1) source. Consumes a queue (optional exchange bindings), acks each page only after the sink confirms it (at-least-once), with idle/max-messages termination.",
));
#[cfg(feature = "source-sftp")]
v.push((
"sftp",
"SFTP source. Lists/globs a remote directory and streams JSONL / JSON-array / raw-text files over SSH.",
));
#[cfg(feature = "source-file")]
v.push((
"file",
"Local file source. Reads JSONL / JSON / CSV / Excel / XML / Parquet / Avro / ORC from a path, directory, glob or http(s) URL, format and compression resolved per file; incremental by mtime or name.",
));
#[cfg(feature = "source-s3")]
v.push(("s3", "AWS S3 object source"));
#[cfg(feature = "source-mongodb")]
v.push(("mongodb", "MongoDB query source"));
#[cfg(feature = "source-mongodb-cdc")]
v.push(("mongodb-cdc", "MongoDB CDC source (Change Streams)"));
#[cfg(feature = "source-mysql-cdc")]
v.push(("mysql-cdc", "MySQL CDC source (binlog replication)"));
#[cfg(feature = "source-mssql-cdc")]
v.push((
"mssql-cdc",
"Microsoft SQL Server CDC source (change data capture, exactly-once capable)",
));
#[cfg(feature = "source-redshift")]
v.push((
"redshift",
"Amazon Redshift query source (PostgreSQL wire; streaming rows, incremental replication)",
));
#[cfg(feature = "source-pubsub")]
v.push((
"pubsub",
"Google Cloud Pub/Sub consumer — streaming pull with per-message records, attribute mapping, and ack at durable page boundaries (at-least-once)",
));
#[cfg(feature = "source-clickhouse")]
v.push((
"clickhouse",
"ClickHouse query source (HTTP interface, JSONEachRow streaming)",
));
#[cfg(feature = "source-azure-blob")]
v.push((
"azure-blob",
"Azure Blob Storage / ADLS Gen2 source — JSONL, JSON array, or raw text",
));
#[cfg(feature = "source-redis")]
v.push(("redis", "Redis (streams, lists, keys) source"));
#[cfg(feature = "source-webhook")]
v.push(("webhook", "Webhook HTTP receiver source"));
#[cfg(feature = "source-websocket")]
v.push((
"websocket",
"WebSocket streaming source — connects, subscribes, streams each message as a record",
));
#[cfg(feature = "source-csv")]
v.push(("csv", deprecated("csv")));
#[cfg(feature = "source-singer")]
v.push((
"singer",
"Singer tap bridge (runs an external Singer tap; single-stream v0, Tier-2/experimental)",
));
#[cfg(feature = "source-elasticsearch")]
v.push(("elasticsearch", "Elasticsearch search / scroll source"));
#[cfg(feature = "source-kafka")]
v.push(("kafka", "Apache Kafka consumer (rdkafka). Subscribes to topics and drains messages with idle/max-messages termination."));
#[cfg(feature = "source-kinesis")]
v.push(("kinesis", "AWS Kinesis Data Streams consumer. Per-shard workers with resumable sequence-number checkpoints and idle/max-messages termination."));
#[cfg(feature = "source-spanner")]
v.push(("spanner", "Google Cloud Spanner query source. Streaming SQL reads with incremental replication bookmarks, stale reads, and PK-range sharding."));
#[cfg(feature = "source-parquet")]
v.push(("parquet", deprecated("parquet")));
#[cfg(feature = "source-delta")]
v.push(("delta", "Apache Delta Lake source (local FS or S3/Azure/GCS). Streams active data files with time travel and projection pushdown."));
#[cfg(feature = "source-databricks")]
v.push(("databricks", "Databricks SQL query source (Statement Execution API). Streams typed query results with chunk pagination and incremental replication."));
#[cfg(feature = "source-iceberg")]
v.push(("iceberg", "Apache Iceberg table source (REST/Glue/SQL/HMS catalogs). Projection and filter pushdown, snapshot time travel, and incremental snapshot reads."));
#[cfg(feature = "source-dynamodb")]
v.push(("dynamodb", "Amazon DynamoDB source — parallel Scan, Query, or DynamoDB Streams change capture with resumable shard bookmarks."));
#[cfg(feature = "source-oracle")]
v.push(("oracle", "Oracle Database query source (ODPI-C session pool). Streams rows with incremental replication, PK-range sharding and discovery; needs Oracle Instant Client at runtime."));
#[cfg(feature = "source-oracle-cdc")]
v.push(("oracle-cdc", "Oracle CDC source via LogMiner. Committed-transaction change events with SCN bookmarks (exactly-once capable); needs Oracle Instant Client at runtime."));
#[cfg(feature = "source-gcs")]
v.push((
"gcs",
"Google Cloud Storage source — JSONL, JSON array, or raw text",
));
#[cfg(feature = "source-bigquery")]
v.push((
"bigquery",
"Google BigQuery query source (jobs.query + jobs.getQueryResults)",
));
#[cfg(feature = "source-snowflake")]
v.push((
"snowflake",
"Snowflake query source (SQL REST API with partition paging)",
));
v
}
pub fn sink_descriptions() -> Vec<(&'static str, &'static str)> {
let mut v = builtin_sink_descriptions();
v.extend(global().custom_sink_descriptions());
v
}
#[allow(clippy::vec_init_then_push)]
fn builtin_sink_descriptions() -> Vec<(&'static str, &'static str)> {
let mut v: Vec<(&'static str, &'static str)> = Vec::new();
#[cfg(feature = "sink-bigquery")]
v.push(("bigquery", "Google BigQuery streaming-insert sink"));
#[cfg(feature = "sink-iceberg")]
v.push((
"iceberg",
"Apache Iceberg sink (append, REST/Glue/SQL/HMS catalogs)",
));
#[cfg(feature = "sink-postgres")]
v.push(("postgres", "PostgreSQL sink (JSONB or auto-mapped columns)"));
#[cfg(feature = "sink-jsonl")]
v.push(("jsonl", deprecated("jsonl")));
#[cfg(feature = "sink-snowflake")]
v.push(("snowflake", "Snowflake SQL REST API sink"));
#[cfg(feature = "sink-mysql")]
v.push(("mysql", "MySQL sink"));
#[cfg(feature = "sink-mssql")]
v.push((
"mssql",
"Microsoft SQL Server sink (auto-mapped columns or JSON column)",
));
#[cfg(feature = "sink-sqlite")]
v.push(("sqlite", "SQLite sink"));
#[cfg(feature = "sink-duckdb")]
v.push((
"duckdb",
"DuckDB sink. Transaction-wrapped multi-row INSERT (JSON column or auto-mapped columns).",
));
#[cfg(feature = "sink-sqs")]
v.push((
"sqs",
"AWS SQS sink. Batched SendMessageBatch (10-message chunks) with per-entry partial-failure retry; FIFO group/dedup support.",
));
#[cfg(feature = "sink-nats")]
v.push((
"nats",
"NATS sink. Publishes records to a subject (optionally subject-per-record) and flushes per batch.",
));
#[cfg(feature = "sink-rabbitmq")]
v.push((
"rabbitmq",
"RabbitMQ (AMQP 0.9.1) sink. Publishes to an exchange with a static/per-record routing key, publisher confirms, and mandatory-return rows surfaced for the DLQ.",
));
#[cfg(feature = "sink-sftp")]
v.push((
"sftp",
"SFTP sink. Writes JSONL files over SSH with atomic temp-then-rename uploads.",
));
#[cfg(feature = "sink-singer")]
v.push((
"singer",
"Singer target bridge. Runs a Singer target executable and feeds it SCHEMA/RECORD/STATE messages; bookmarks advance only after the target confirms (echoed STATE or clean exit).",
));
#[cfg(feature = "sink-s3")]
v.push(("s3", "AWS S3 object sink"));
#[cfg(feature = "sink-mongodb")]
v.push(("mongodb", "MongoDB insert sink"));
#[cfg(feature = "sink-redis")]
v.push(("redis", "Redis (streams, lists, key-value) sink"));
#[cfg(feature = "sink-csv")]
v.push(("csv", deprecated("csv")));
#[cfg(feature = "sink-elasticsearch")]
v.push(("elasticsearch", "Elasticsearch bulk index sink"));
#[cfg(feature = "sink-kafka")]
v.push(("kafka", "Apache Kafka producer (rdkafka). FuturesUnordered batched sends with QueueFull retry; supports fixed or per-record topic routing."));
#[cfg(feature = "sink-kinesis")]
v.push(("kinesis", "AWS Kinesis Data Streams producer. Batched PutRecords with partition-key routing and partial-failure retry (DLQ-routable)."));
#[cfg(feature = "sink-spanner")]
v.push(("spanner", "Google Cloud Spanner sink. Batched mutations with upsert/delete write modes, exactly-once commit tokens, and schema evolution."));
#[cfg(feature = "sink-http")]
v.push(("http", "HTTP POST sink (individual or array batch)"));
#[cfg(feature = "sink-stdout")]
v.push(("stdout", "Stdout / stderr sink (JSON Lines, pretty, TSV)"));
#[cfg(feature = "sink-parquet")]
v.push(("parquet", deprecated("parquet")));
#[cfg(feature = "sink-file")]
v.push(("file", "Local file sink. JSONL, JSON, CSV, XML, Excel, Avro or Parquet by extension; rollover, compression, temp-then-rename finalisation, atomic overwrite."));
#[cfg(feature = "sink-delta")]
v.push(("delta", "Apache Delta Lake sink (local FS or S3/Azure/GCS). Append-only, schema-inferred table creation, one commit per flush."));
#[cfg(feature = "sink-gcs")]
v.push(("gcs", "Google Cloud Storage sink — JSONL files"));
#[cfg(feature = "sink-redshift")]
v.push((
"redshift",
"Amazon Redshift sink (COPY-from-S3 or multi-row INSERT)",
));
#[cfg(feature = "sink-pubsub")]
v.push((
"pubsub",
"Google Cloud Pub/Sub producer — batched publish with optional ordering keys, bounded concurrency, and partial-failure retry (DLQ-routable)",
));
#[cfg(feature = "sink-clickhouse")]
v.push((
"clickhouse",
"ClickHouse sink (HTTP INSERT … FORMAT JSONEachRow; optional async inserts)",
));
#[cfg(feature = "sink-azure-blob")]
v.push((
"azure-blob",
"Azure Blob Storage / ADLS Gen2 sink — JSONL files",
));
#[cfg(feature = "sink-dynamodb")]
v.push((
"dynamodb",
"Amazon DynamoDB sink — batched writes, upsert/delete by key, conditional writes",
));
#[cfg(feature = "sink-databricks")]
v.push((
"databricks",
"Databricks SQL warehouse sink — append/upsert/delete/overwrite into Delta tables, exactly-once, schema evolution, staged COPY INTO",
));
#[cfg(feature = "sink-oracle")]
v.push((
"oracle",
"Oracle Database sink — array-bound inserts, MERGE upsert/delete, overwrite, exactly-once, schema evolution (needs Oracle Instant Client at runtime)",
));
v
}
pub fn source_kinds() -> Vec<&'static str> {
source_descriptions().into_iter().map(|(k, _)| k).collect()
}
pub fn sink_kinds() -> Vec<&'static str> {
sink_descriptions().into_iter().map(|(k, _)| k).collect()
}
fn decode<T: DeserializeOwned>(kind: &'static str, name: &str, config: Value) -> CliResult<T> {
serde_json::from_value(config).map_err(|e| CliError::InvalidConnectorConfig {
kind,
name: name.to_owned(),
message: scrub_config_error(&e.to_string()),
})
}
fn scrub_config_error(msg: &str) -> String {
const MAX_CHARS: usize = 200;
let mut out = String::with_capacity(msg.len());
let mut in_quote = false;
for c in msg.chars() {
if c == '"' {
if !in_quote {
out.push_str("\"<redacted>\"");
}
in_quote = !in_quote;
continue;
}
if !in_quote {
out.push(c);
}
}
if out.chars().count() > MAX_CHARS {
let truncated: String = out.chars().take(MAX_CHARS).collect();
return format!("{truncated}…");
}
out
}
fn schema<T: faucet_core::JsonSchema>() -> Value {
serde_json::to_value(faucet_core::schema_for!(T))
.unwrap_or_else(|_| serde_json::json!({"type": "object"}))
}
fn unknown(name: &str, kind: &'static str, available: Vec<&'static str>) -> CliError {
CliError::UnknownConnector {
kind,
name: name.to_owned(),
available: if available.is_empty() {
"(none — rebuild faucet-cli with the relevant feature enabled)".to_owned()
} else {
available.join(", ")
},
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn bookmark_schema_and_migration_per_source_kind() {
assert_eq!(source_state_schema("csv"), 0);
assert_eq!(
migrate_source_state("csv", 0, json!({"a": 1})).unwrap(),
json!({"a": 1})
);
assert!(migrate_source_state("csv", 1, json!({})).is_err());
#[cfg(feature = "source-mongodb-cdc")]
{
assert_eq!(source_state_schema("mongodb-cdc"), 1);
let m = migrate_source_state("mongodb-cdc", 0, json!({"resume_token": {"_data": "A"}}))
.unwrap();
assert_eq!(m["invalidate"], json!(false));
assert_eq!(state_codec_for(Some(("mongodb-cdc", &json!({})))).schema, 1);
}
let topo = state_codec_for(None);
assert_eq!(
(topo.owner.as_str(), topo.schema, topo.legacy),
("topology", 0, false)
);
assert_eq!(state_codec_for(Some(("csv", &json!({})))).owner, "csv");
}
#[test]
fn every_compiled_source_kind_validates_its_own_config_type() {
for (kind, _) in source_descriptions() {
let err = validate_source_config(kind, "row-a", serde_json::json!(42))
.expect_err("a non-object config cannot deserialize into any connector config");
let msg = err.to_string();
assert!(
msg.contains(kind),
"the failure must name the connector kind so the operator knows which \
entry is wrong: {msg}"
);
assert!(
msg.contains("row-a"),
"the failure must name the row/template entry: {msg}"
);
}
}
#[test]
fn every_compiled_sink_kind_validates_its_own_config_type() {
for (kind, _) in sink_descriptions() {
let err = validate_sink_config(kind, "row-b", serde_json::json!("not an object"))
.expect_err("a non-object config cannot deserialize into any connector config");
let msg = err.to_string();
assert!(msg.contains(kind), "must name the kind: {msg}");
assert!(msg.contains("row-b"), "must name the entry: {msg}");
}
}
#[test]
fn an_empty_config_object_is_never_a_panic_and_always_actionable() {
for (kind, _) in source_descriptions() {
if let Err(e) = validate_source_config(kind, "row-c", serde_json::json!({})) {
let msg = e.to_string();
assert!(msg.contains(kind) && msg.contains("row-c"), "{msg}");
}
}
for (kind, _) in sink_descriptions() {
if let Err(e) = validate_sink_config(kind, "row-d", serde_json::json!({})) {
let msg = e.to_string();
assert!(msg.contains(kind) && msg.contains("row-d"), "{msg}");
}
}
}
#[test]
fn an_unknown_kind_is_refused_with_the_available_kinds() {
let src = validate_source_config("not-a-connector", "row", serde_json::json!({}))
.expect_err("an unknown source kind must not validate");
assert!(src.to_string().contains("not-a-connector"), "{src}");
let sink = validate_sink_config("not-a-connector", "row", serde_json::json!({}))
.expect_err("an unknown sink kind must not validate");
assert!(sink.to_string().contains("not-a-connector"), "{sink}");
}
#[test]
fn staged_load_allowlist_matches_capable_sinks() {
for k in ["redshift", "snowflake", "bigquery", "clickhouse", "mssql"] {
assert!(
sink_supports_staged_load(k),
"{k} should support staged load"
);
}
for k in ["jsonl", "postgres", "sqlite", "stdout"] {
assert!(
!sink_supports_staged_load(k),
"{k} should not (yet) support staged load"
);
}
}
#[derive(Clone)]
struct DummySource;
#[faucet_core::async_trait]
impl Source for DummySource {
async fn fetch_with_context(
&self,
_ctx: &std::collections::HashMap<String, Value>,
) -> Result<Vec<Value>, faucet_core::FaucetError> {
Ok(vec![serde_json::json!({"ok": true})])
}
fn config_schema(&self) -> Value {
serde_json::json!({"type": "object"})
}
}
#[test]
fn register_source_rejects_builtin_collision() {
let reg = PluginRegistry::with_builtins()
.register_source("csv", |_| Ok(Box::new(DummySource) as Box<dyn Source>));
let err = reg
.install()
.expect_err("built-in collision must be rejected");
match err {
CliError::Config(msg) => assert!(msg.contains("built-in source"), "{msg}"),
other => panic!("expected Config error, got {other:?}"),
}
}
#[test]
fn register_source_rejects_duplicate() {
let reg = PluginRegistry::new()
.register_source("dup", |_| Ok(Box::new(DummySource) as Box<dyn Source>))
.register_source("dup", |_| Ok(Box::new(DummySource) as Box<dyn Source>));
let err = reg
.install()
.expect_err("duplicate registration must be rejected");
match err {
CliError::Config(msg) => assert!(msg.contains("more than once"), "{msg}"),
other => panic!("expected Config error, got {other:?}"),
}
}
#[test]
fn register_sink_rejects_duplicate() {
let reg = PluginRegistry::new()
.register_sink("dupsink", |_| Err(CliError::Config("unused".into())))
.register_sink("dupsink", |_| Err(CliError::Config("unused".into())));
assert!(
reg.errors.iter().any(|e| e.contains("more than once")),
"{:?}",
reg.errors
);
}
#[test]
fn custom_descriptions_use_default_when_blank() {
let reg = PluginRegistry::new()
.register_source("acme", |_| Ok(Box::new(DummySource) as Box<dyn Source>));
let descs = reg.custom_source_descriptions();
assert_eq!(descs.len(), 1);
assert_eq!(descs[0].0, "acme");
assert_eq!(descs[0].1, "custom source connector");
}
#[test]
fn custom_descriptions_carry_explicit_summary() {
let reg = PluginRegistry::new().register_source_with(
"acme",
|_| Ok(Box::new(DummySource) as Box<dyn Source>),
|| serde_json::json!({"type": "object", "title": "acme"}),
"Acme widget source",
);
let descs = reg.custom_source_descriptions();
assert_eq!(descs[0], ("acme", "Acme widget source"));
assert_eq!(
(reg.sources.get("acme").unwrap().schema)()["title"],
serde_json::json!("acme")
);
}
#[test]
fn capability_constants_match_their_predicates() {
for &k in EXACTLY_ONCE_SOURCE_KINDS {
assert!(
source_supports_exactly_once(k),
"{k} should be exactly-once"
);
}
for &k in IDEMPOTENT_SINK_KINDS {
assert!(
sink_supports_idempotent_writes(k),
"{k} should be idempotent"
);
}
for &k in UPSERT_SINK_KINDS {
use faucet_core::WriteMode;
assert!(
sink_supported_write_modes(k).contains(&WriteMode::Upsert),
"{k} should support upsert"
);
}
assert!(IDEMPOTENT_SINK_KINDS.contains(&"bigquery"));
assert!(IDEMPOTENT_SINK_KINDS.contains(&"kafka"));
assert!(UPSERT_SINK_KINDS.contains(&"bigquery"));
}
#[cfg(feature = "source-rest")]
#[test]
fn rest_source_appears_in_listings() {
assert!(source_kinds().contains(&"rest"));
}
#[cfg(feature = "sink-jsonl")]
#[test]
fn jsonl_sink_appears_in_listings() {
assert!(sink_kinds().contains(&"jsonl"));
}
#[tokio::test]
async fn unknown_source_kind_errors() {
let err = build_source("nope", serde_json::json!({}), &AuthCatalog::new(), None)
.await
.err()
.expect("should fail");
match err {
CliError::UnknownConnector { kind, name, .. } => {
assert_eq!(kind, "source");
assert_eq!(name, "nope");
}
other => panic!("expected UnknownConnector, got {other:?}"),
}
}
#[tokio::test]
async fn unknown_sink_kind_errors() {
let err = build_sink("nope", serde_json::json!({}), &AuthCatalog::new())
.await
.err()
.expect("should fail");
assert!(matches!(
err,
CliError::UnknownConnector { kind: "sink", .. }
));
}
#[cfg(feature = "source-rest")]
#[test]
fn rest_schema_is_object() {
let s = source_schema("rest").unwrap();
assert!(s.is_object());
}
#[cfg(feature = "sink-jsonl")]
#[test]
fn jsonl_schema_is_object() {
let s = sink_schema("jsonl").unwrap();
assert!(s.is_object());
}
#[test]
fn scrub_config_error_redacts_quoted_values() {
let msg =
r#"invalid type: string "sk-super-secret-123", expected a sequence at line 1 column 9"#;
let scrubbed = scrub_config_error(msg);
assert!(!scrubbed.contains("sk-super-secret-123"), "{scrubbed}");
assert!(scrubbed.contains("<redacted>"), "{scrubbed}");
assert!(scrubbed.contains("invalid type"), "{scrubbed}");
assert!(scrubbed.contains("expected a sequence"), "{scrubbed}");
}
#[test]
fn scrub_config_error_truncates_long_messages() {
let msg = "x".repeat(500);
let scrubbed = scrub_config_error(&msg);
assert!(
scrubbed.chars().count() <= 201,
"len {}",
scrubbed.chars().count()
);
assert!(scrubbed.ends_with('…'));
}
#[cfg(feature = "source-csv")]
#[tokio::test]
async fn build_source_csv_succeeds_without_io() {
let src = build_source(
"csv",
serde_json::json!({ "path": "/tmp/does-not-need-to-exist.csv" }),
&AuthCatalog::new(),
None,
)
.await
.expect("csv source should build without I/O");
assert_eq!(src.connector_name(), "csv");
}
#[cfg(feature = "sink-jsonl")]
#[tokio::test]
async fn build_sink_jsonl_succeeds_without_io() {
let sink = build_sink(
"jsonl",
serde_json::json!({ "path": "/tmp/does-not-need-to-exist.jsonl" }),
&AuthCatalog::new(),
)
.await
.expect("jsonl sink should build without I/O");
assert_eq!(sink.connector_name(), "jsonl");
}
#[cfg(feature = "sink-stdout")]
#[tokio::test]
async fn build_sink_stdout_succeeds_without_io() {
let sink = build_sink("stdout", serde_json::json!({}), &AuthCatalog::new())
.await
.expect("stdout sink should build without I/O");
assert_eq!(sink.connector_name(), "stdout");
}
#[cfg(all(feature = "source-delta", feature = "sink-delta"))]
#[tokio::test]
async fn delta_registry_round_trip() {
let dir = tempfile::tempdir().unwrap();
let uri = dir.path().join("reg_delta").to_string_lossy().into_owned();
assert!(source_schema("delta").is_ok());
assert!(sink_schema("delta").is_ok());
assert!(source_descriptions().iter().any(|(n, _)| *n == "delta"));
assert!(sink_descriptions().iter().any(|(n, _)| *n == "delta"));
let sink = build_sink(
"delta",
serde_json::json!({ "table_uri": uri }),
&AuthCatalog::new(),
)
.await
.expect("delta sink builds");
assert_eq!(sink.connector_name(), "delta");
let n = sink
.write_batch(&[serde_json::json!({"id": 1}), serde_json::json!({"id": 2})])
.await
.expect("write");
assert_eq!(n, 2);
sink.flush().await.expect("flush");
let source = build_source(
"delta",
serde_json::json!({ "table_uri": uri }),
&AuthCatalog::new(),
None,
)
.await
.expect("delta source builds");
assert_eq!(source.connector_name(), "delta");
let rows = source
.fetch_with_context(&std::collections::HashMap::new())
.await
.expect("read");
assert_eq!(rows.len(), 2);
}
#[cfg(feature = "source-databricks")]
#[tokio::test]
async fn databricks_registry_source_builds() {
assert!(source_schema("databricks").is_ok());
assert!(
source_descriptions()
.iter()
.any(|(n, _)| *n == "databricks")
);
let cfg = serde_json::json!({
"workspace_url": "https://x.cloud.databricks.com",
"warehouse_id": "wh1",
"sql": "SELECT 1",
"auth": { "type": "pat", "config": { "token": "t" } }
});
let src = build_source("databricks", cfg, &AuthCatalog::new(), None)
.await
.expect("databricks source builds");
assert_eq!(src.connector_name(), "databricks");
}
#[cfg(feature = "source-csv")]
#[tokio::test]
async fn build_source_csv_invalid_config_errors() {
let res = build_source(
"csv",
serde_json::json!({ "path": 42 }),
&AuthCatalog::new(),
None,
)
.await;
match res {
Err(CliError::InvalidConnectorConfig { kind, name, .. }) => {
assert_eq!(kind, "source");
assert_eq!(name, "csv");
}
Ok(_) => panic!("expected InvalidConnectorConfig, got Ok"),
Err(other) => panic!("expected InvalidConnectorConfig, got {other:?}"),
}
}
#[cfg(feature = "source-csv")]
#[tokio::test]
async fn build_source_csv_rejects_oversized_batch_size() {
let res = build_source(
"csv",
serde_json::json!({
"path": "/tmp/x.csv",
"batch_size": faucet_core::MAX_BATCH_SIZE + 1
}),
&AuthCatalog::new(),
None,
)
.await;
assert!(matches!(
res,
Err(CliError::Faucet(faucet_core::FaucetError::Config(_)))
));
}
#[cfg(feature = "source-graphql")]
#[tokio::test]
async fn build_source_graphql_rejects_empty_endpoint() {
let res = build_source(
"graphql",
serde_json::json!({
"endpoint": "",
"query": "query { x }",
"variables": {},
"auth": { "type": "none" }
}),
&AuthCatalog::new(),
None,
)
.await;
assert!(matches!(
res,
Err(CliError::Faucet(faucet_core::FaucetError::Config(_)))
));
}
#[cfg(feature = "source-csv")]
#[test]
fn source_schema_csv_exposes_path_property() {
let schema = source_schema("csv").expect("csv schema");
let props = schema
.get("properties")
.and_then(Value::as_object)
.expect("schema should have a properties object");
assert!(props.contains_key("path"), "schema props: {props:?}");
}
#[cfg(feature = "sink-jsonl")]
#[test]
fn sink_schema_jsonl_exposes_path_property() {
let schema = sink_schema("jsonl").expect("jsonl schema");
let props = schema
.get("properties")
.and_then(Value::as_object)
.expect("schema should have a properties object");
assert!(props.contains_key("path"), "schema props: {props:?}");
}
#[test]
fn unknown_source_schema_errors_with_available_list() {
let err = source_schema("definitely-not-a-source").expect_err("unknown source");
match err {
CliError::UnknownConnector {
kind,
name,
available,
} => {
assert_eq!(kind, "source");
assert_eq!(name, "definitely-not-a-source");
assert!(!available.is_empty());
}
other => panic!("expected UnknownConnector, got {other:?}"),
}
}
#[test]
fn unknown_sink_schema_errors() {
let err = sink_schema("definitely-not-a-sink").expect_err("unknown sink");
assert!(matches!(
err,
CliError::UnknownConnector { kind: "sink", .. }
));
}
#[cfg(feature = "source-csv")]
#[test]
fn source_exists_is_true_for_known_and_false_for_unknown() {
assert!(source_exists("csv"));
assert!(!source_exists("definitely-not-a-source"));
}
#[cfg(feature = "sink-jsonl")]
#[test]
fn sink_exists_is_true_for_known_and_false_for_unknown() {
assert!(sink_exists("jsonl"));
assert!(!sink_exists("definitely-not-a-sink"));
}
#[test]
fn source_descriptions_are_non_empty_and_consistent() {
let descs = source_descriptions();
assert!(!descs.is_empty());
for (name, summary) in &descs {
assert!(!name.is_empty(), "empty connector name");
assert!(!summary.is_empty(), "empty summary for {name}");
assert!(
source_schema(name).is_ok(),
"listed source `{name}` has no schema"
);
}
}
#[test]
fn sink_descriptions_are_non_empty_and_consistent() {
let descs = sink_descriptions();
assert!(!descs.is_empty());
for (name, summary) in &descs {
assert!(!name.is_empty(), "empty connector name");
assert!(!summary.is_empty(), "empty summary for {name}");
assert!(
sink_schema(name).is_ok(),
"listed sink `{name}` has no schema"
);
}
}
#[cfg(all(feature = "source-csv", feature = "source-rest"))]
#[test]
fn source_kinds_contains_expected_builtins() {
let kinds = source_kinds();
assert!(kinds.contains(&"csv"));
assert!(kinds.contains(&"rest"));
}
#[cfg(all(feature = "sink-jsonl", feature = "sink-stdout"))]
#[test]
fn sink_kinds_contains_expected_builtins() {
let kinds = sink_kinds();
assert!(kinds.contains(&"jsonl"));
assert!(kinds.contains(&"stdout"));
}
#[cfg(feature = "source-rest")]
#[tokio::test]
async fn build_source_injects_referenced_auth_provider() {
let mut specs = std::collections::HashMap::new();
specs.insert(
"tok".to_string(),
serde_json::json!({"type": "static", "config": {"token": "abc"}}),
);
let catalog = auth_catalog::build_auth_catalog(Some(&specs)).expect("catalog");
let src = build_source("rest", rest_config_with_auth_ref("tok"), &catalog, None)
.await
.expect("rest source with a resolvable auth ref should build");
assert_eq!(src.connector_name(), "rest");
}
#[cfg(feature = "source-rest")]
fn rest_config_with_auth_ref(name: &str) -> Value {
let cfg = faucet_source_rest::RestStreamConfig::new("https://api.example.com", "/v1");
let mut v = serde_json::to_value(cfg).expect("serialize rest config");
v.as_object_mut()
.unwrap()
.insert("auth".to_string(), serde_json::json!({ "ref": name }));
v
}
#[cfg(feature = "source-rest")]
#[tokio::test]
async fn build_source_unknown_auth_ref_errors() {
let res = build_source(
"rest",
rest_config_with_auth_ref("missing"),
&AuthCatalog::new(),
None,
)
.await;
match res {
Err(CliError::UnknownAuthProvider { name, .. }) => assert_eq!(name, "missing"),
Ok(_) => panic!("expected UnknownAuthProvider, got Ok"),
Err(other) => panic!("expected UnknownAuthProvider, got {other:?}"),
}
}
#[test]
fn exactly_once_capability_allowlists() {
assert!(source_supports_exactly_once("postgres-cdc"));
assert!(source_supports_exactly_once("mysql-cdc"));
assert!(source_supports_exactly_once("mongodb-cdc"));
assert!(source_supports_exactly_once("kafka"));
assert!(!source_supports_exactly_once("rest"));
assert!(sink_supports_idempotent_writes("postgres"));
assert!(sink_supports_idempotent_writes("iceberg"));
assert!(sink_supports_idempotent_writes("bigquery"));
assert!(sink_supports_idempotent_writes("kafka"));
assert!(sink_supports_idempotent_writes("snowflake"));
assert!(sink_supports_idempotent_writes("redis"));
assert!(sink_supports_idempotent_writes("mongodb"));
assert!(!sink_supports_idempotent_writes("jsonl"));
}
#[test]
fn typed_delivery_capabilities_derive_from_kind_tables() {
use faucet_core::{ReplayGuarantee, SinkGuarantee};
assert_eq!(
source_replay_guarantee("kafka"),
ReplayGuarantee::Deterministic
);
assert_eq!(
source_replay_guarantee("rest"),
ReplayGuarantee::NonDeterministic
);
assert_eq!(sink_guarantee("postgres"), SinkGuarantee::AtomicWatermark);
assert_eq!(sink_guarantee("elasticsearch"), SinkGuarantee::KeyedUpsert);
assert_eq!(sink_guarantee("jsonl"), SinkGuarantee::AtLeastOnce);
}
#[test]
fn singer_write_modes_are_not_key_deduplicating() {
use faucet_core::WriteMode;
assert_eq!(sink_supported_write_modes("singer"), SINGER_WRITE_MODES);
assert!(!sink_supported_write_modes("singer").contains(&WriteMode::Delete));
assert!(!UPSERT_SINK_KINDS.contains(&"singer"));
assert!(!sink_supports_overwrite("singer"));
assert_eq!(
sink_guarantee("singer"),
faucet_core::SinkGuarantee::AtLeastOnce
);
}
#[cfg(feature = "sink-singer")]
#[tokio::test]
async fn singer_sink_builds_validates_and_registers_secrets() {
let auth = crate::auth_catalog::AuthCatalog::new();
let cfg = json!({
"target_command": "target-jsonl",
"target_config": {"api_token": "singer-reg-secret-1234"},
"stream": "s"
});
let sink = build_sink("singer", cfg.clone(), &auth).await.unwrap();
assert_eq!(sink.connector_name(), "singer");
assert_eq!(
crate::secrets::registry::redact("x singer-reg-secret-1234 y"),
"x *** y"
);
assert_eq!(
sink_batch_atomicity("singer", &cfg),
Some(faucet_core::BatchAtomicity::BestEffort)
);
assert!(sink_schema("singer").is_ok());
assert!(validate_sink_config("singer", "row", cfg).is_ok());
let bad = json!({"target_command": "t", "write_mode": "delete", "key": ["id"]});
assert!(validate_sink_config("singer", "row", bad).is_err());
assert!(sink_descriptions().iter().any(|(k, _)| *k == "singer"));
}
#[test]
fn sink_supported_write_modes_allowlist() {
use faucet_core::WriteMode;
assert!(sink_supported_write_modes("postgres").contains(&WriteMode::Upsert));
assert!(sink_supported_write_modes("elasticsearch").contains(&WriteMode::Delete));
assert!(sink_supported_write_modes("bigquery").contains(&WriteMode::Upsert));
assert_eq!(sink_supported_write_modes("jsonl"), &[WriteMode::Append]);
assert_eq!(sink_supported_write_modes("kafka"), &[WriteMode::Append]);
for k in [
"postgres",
"sqlite",
"mysql",
"mssql",
"mongodb",
"bigquery",
"elasticsearch",
] {
assert!(
sink_supported_write_modes(k).contains(&WriteMode::Overwrite),
"{k} should support overwrite"
);
assert!(sink_supports_overwrite(k), "{k} sink_supports_overwrite");
}
assert!(!sink_supports_overwrite("jsonl"));
}
#[test]
fn connectors_710_713_capability_allowlists() {
use faucet_core::{SinkGuarantee, WriteMode};
assert_eq!(
sink_supported_write_modes("dynamodb"),
&[WriteMode::Append, WriteMode::Upsert, WriteMode::Delete]
);
assert!(!sink_supports_overwrite("dynamodb"));
assert!(!sink_supports_idempotent_writes("dynamodb"));
assert!(!sink_supports_schema_evolution("dynamodb"));
assert!(!sink_supports_cleanup("dynamodb"));
assert_eq!(sink_guarantee("dynamodb"), SinkGuarantee::KeyedUpsert);
assert_eq!(
sink_supported_write_modes("databricks"),
&[
WriteMode::Append,
WriteMode::Upsert,
WriteMode::Delete,
WriteMode::Overwrite
]
);
assert!(sink_supports_idempotent_writes("databricks"));
assert!(sink_supports_schema_evolution("databricks"));
assert!(sink_supports_staged_load("databricks"));
assert!(!sink_supports_cleanup("databricks"));
for k in ["iceberg", "dynamodb"] {
assert!(source_supports_discover(k), "{k} discovers");
assert!(!source_supports_exactly_once(k), "{k} is not exactly-once");
}
assert!(source_supports_exactly_once("oracle-cdc"));
assert!(!source_supports_exactly_once("oracle"));
assert!(source_supports_discover("oracle"));
assert_eq!(
sink_supported_write_modes("oracle"),
&[
WriteMode::Append,
WriteMode::Upsert,
WriteMode::Delete,
WriteMode::Overwrite
]
);
assert!(sink_supports_idempotent_writes("oracle"));
assert!(sink_supports_schema_evolution("oracle"));
assert!(!sink_supports_cleanup("oracle"));
}
#[cfg(feature = "sink-databricks")]
#[tokio::test]
async fn databricks_registry_sink_builds_with_auth_ref() {
assert!(sink_schema("databricks").is_ok());
assert!(sink_kinds().contains(&"databricks"));
let mut specs = std::collections::HashMap::new();
specs.insert(
"dbx".to_string(),
serde_json::json!({"type": "static", "config": {"token": "abc"}}),
);
let catalog = auth_catalog::build_auth_catalog(Some(&specs)).expect("catalog");
let cfg = serde_json::json!({
"workspace_url": "https://x.cloud.databricks.com",
"warehouse_id": "w",
"auth": { "ref": "dbx" },
"schema": "s",
"table": "t",
});
validate_sink_config("databricks", "row", cfg.clone()).expect("valid databricks config");
let sink = build_sink("databricks", cfg, &catalog)
.await
.expect("databricks sink builds");
assert_eq!(sink.connector_name(), "databricks");
assert!(sink.supports_idempotent_writes());
}
#[cfg(all(
feature = "source-iceberg",
feature = "source-dynamodb",
feature = "sink-dynamodb"
))]
#[test]
fn iceberg_and_dynamodb_are_registered() {
assert!(source_schema("iceberg").is_ok());
assert!(source_schema("dynamodb").is_ok());
assert!(sink_schema("dynamodb").is_ok());
assert!(source_kinds().contains(&"iceberg"));
assert!(source_kinds().contains(&"dynamodb"));
assert!(sink_kinds().contains(&"dynamodb"));
let err = validate_sink_config(
"dynamodb",
"row",
serde_json::json!({"table_name": "t", "write_mode": "overwrite"}),
)
.expect_err("dynamodb rejects overwrite");
assert!(err.to_string().contains("overwrite"), "{err}");
validate_source_config(
"dynamodb",
"row",
serde_json::json!({"table_name": "t", "mode": "scan"}),
)
.expect("scan config is valid");
}
#[test]
fn sink_supports_schema_evolution_allowlist() {
assert!(sink_supports_schema_evolution("postgres"));
assert!(sink_supports_schema_evolution("mysql"));
assert!(sink_supports_schema_evolution("mssql"));
assert!(sink_supports_schema_evolution("sqlite"));
assert!(sink_supports_schema_evolution("bigquery"));
assert!(sink_supports_schema_evolution("elasticsearch"));
assert!(sink_supports_schema_evolution("iceberg"));
assert!(!sink_supports_schema_evolution("jsonl"));
assert!(!sink_supports_schema_evolution("kafka"));
}
#[test]
fn edit_distance_is_levenshtein() {
assert_eq!(edit_distance("", ""), 0);
assert_eq!(edit_distance("abc", "abc"), 0);
assert_eq!(edit_distance("", "abc"), 3);
assert_eq!(edit_distance("abc", ""), 3);
assert_eq!(edit_distance("batch_size", "batch_size"), 0);
assert_eq!(edit_distance("batch_sizes", "batch_size"), 1);
assert_eq!(edit_distance("kitten", "sitting"), 3);
}
#[test]
fn nearest_key_only_suggests_a_plausible_typo() {
let known: Vec<String> = ["batch_size", "verify_checksum", "prefix"]
.iter()
.map(|s| s.to_string())
.collect();
assert_eq!(
nearest_key("verify_checksums", known.iter()),
Some("verify_checksum")
);
assert_eq!(nearest_key("batch_sze", known.iter()), Some("batch_size"));
assert_eq!(nearest_key("region", known.iter()), None);
}
#[test]
fn unknown_connector_keys_are_rejected_with_the_declared_list() {
let schema = json!({
"type": "object",
"properties": { "bucket": {}, "prefix": {}, "verify_checksum": {} }
});
let ok = json!({ "bucket": "b", "verify_checksum": true });
assert!(reject_unknown_config_keys("source", "s3", "row", &ok, &schema).is_ok());
let typo = json!({ "bucket": "b", "verify_checksums": true });
let err = reject_unknown_config_keys("source", "s3", "row", &typo, &schema)
.expect_err("a typo'd key must not be silently ignored");
let msg = err.to_string();
assert!(msg.contains("verify_checksums"), "{msg}");
assert!(
msg.contains("did you mean `verify_checksum`"),
"the near-miss must be named, or the user re-reads the whole schema: {msg}"
);
}
#[test]
fn a_declared_catch_all_schema_accepts_any_key() {
let schema = json!({
"type": "object",
"properties": { "url": {} },
"additionalProperties": { "type": "string" }
});
let cfg = json!({ "url": "u", "anything_at_all": "x" });
assert!(reject_unknown_config_keys("sink", "http", "row", &cfg, &schema).is_ok());
let strict = json!({
"type": "object",
"properties": { "url": {} },
"additionalProperties": false
});
assert!(reject_unknown_config_keys("sink", "http", "row", &cfg, &strict).is_err());
}
#[test]
fn a_non_object_config_or_schema_is_left_alone() {
let schema = json!({ "type": "object", "properties": { "a": {} } });
assert!(
reject_unknown_config_keys("source", "k", "row", &json!("${vars.x}"), &schema).is_ok()
);
assert!(
reject_unknown_config_keys("source", "k", "row", &json!({"zzz": 1}), &json!({}))
.is_ok()
);
}
#[test]
fn a_flattened_tagged_enum_contributes_its_variant_keys() {
let schema = json!({
"type": "object",
"properties": { "host": {}, "username": {} },
"oneOf": [
{ "type": "object", "properties": { "type": {}, "config": {} } },
{ "type": "object", "properties": { "type": {}, "key_path": {} } }
]
});
let cfg = json!({ "host": "h", "username": "u", "type": "password", "config": {} });
assert!(reject_unknown_config_keys("source", "sftp", "row", &cfg, &schema).is_ok());
let bad = json!({ "host": "h", "usernme": "u" });
assert!(reject_unknown_config_keys("source", "sftp", "row", &bad, &schema).is_err());
}
#[test]
fn a_declared_serde_alias_is_accepted() {
let schema = json!({
"type": "object",
"properties": { "create_table": {} },
"x-faucet-aliases": ["create_if_missing"]
});
let cfg = json!({ "create_if_missing": true });
assert!(reject_unknown_config_keys("sink", "iceberg", "row", &cfg, &schema).is_ok());
assert!(
reject_unknown_config_keys("sink", "iceberg", "row", &json!({"nope": 1}), &schema)
.is_err()
);
}
#[test]
fn a_ref_into_defs_is_followed() {
let schema = json!({
"type": "object",
"properties": { "a": {} },
"allOf": [ { "$ref": "#/$defs/Extra" } ],
"$defs": { "Extra": { "type": "object", "properties": { "b": {} } } }
});
assert!(
reject_unknown_config_keys("sink", "k", "row", &json!({"a": 1, "b": 2}), &schema)
.is_ok()
);
assert!(reject_unknown_config_keys("sink", "k", "row", &json!({"c": 3}), &schema).is_err());
let dangling = json!({
"type": "object",
"properties": { "a": {} },
"allOf": [ { "$ref": "#/$defs/Missing" } ]
});
assert!(
reject_unknown_config_keys("sink", "k", "row", &json!({"c": 3}), &dangling).is_ok()
);
}
#[cfg(feature = "source-sftp")]
#[test]
fn the_shipped_sftp_shape_with_flattened_auth_is_accepted() {
let cfg = json!({
"host": "sftp.example.com", "port": 22, "username": "reporting",
"type": "password", "config": { "password": "x" },
"known_hosts": { "mode": "accept_new" },
"path": "/exports/daily", "glob": "*.jsonl",
"format": "jsonl", "batch_size": 1000
});
assert!(validate_source_config("sftp", "row", cfg).is_ok());
}
#[cfg(feature = "sink-bigquery")]
#[test]
fn a_flattened_config_still_declares_its_flattened_keys() {
let schema = sink_schema("bigquery").expect("bigquery schema");
let props = schema["properties"].as_object().expect("properties");
assert!(props.contains_key("write_mode"), "flattened WriteSpec key");
assert!(props.contains_key("project_id"), "own key");
let good = json!({
"project_id": "p", "dataset_id": "d", "table_id": "t",
"auth": {"type": "application_default"}, "write_mode": "append"
});
assert!(validate_sink_config("bigquery", "row", good).is_ok());
let bad = json!({
"project_id": "p", "dataset_id": "d", "table_id": "t",
"auth": {"type": "application_default"}, "write_modes": "append"
});
let err = validate_sink_config("bigquery", "row", bad)
.expect_err("a typo'd flattened key must be rejected");
assert!(err.to_string().contains("write_modes"), "{err}");
}
#[test]
fn concurrency_override_sets_the_first_declared_knob() {
let schema = json!({
"type": "object",
"properties": { "connection_url": {}, "max_connections": {}, "batch_size": {} }
});
let mut cfg = json!({ "connection_url": "postgres://u@h/db", "max_connections": 10 });
assert_eq!(
apply_concurrency_override(&schema, &mut cfg, 20),
Some("max_connections")
);
assert_eq!(
cfg["max_connections"], 20,
"the override must win over the config"
);
assert_eq!(
cfg["connection_url"], "postgres://u@h/db",
"nothing else touched"
);
}
#[test]
fn concurrency_override_is_a_no_op_without_a_knob() {
let schema = json!({ "type": "object", "properties": { "path": {} } });
let mut cfg = json!({ "path": "./out.jsonl" });
assert_eq!(apply_concurrency_override(&schema, &mut cfg, 20), None);
assert_eq!(cfg, json!({ "path": "./out.jsonl" }));
}
#[test]
fn concurrency_override_prefers_knobs_in_a_fixed_order() {
let schema = json!({
"type": "object",
"properties": { "concurrency": {}, "max_connections": {} }
});
let mut cfg = json!({});
assert_eq!(
apply_concurrency_override(&schema, &mut cfg, 3),
Some("max_connections")
);
assert!(cfg.get("concurrency").is_none());
}
#[cfg(feature = "source-postgres")]
#[test]
fn a_real_source_schema_resolves_to_its_pool_knob() {
let mut cfg = json!({ "connection_url": "postgres://u@h/db", "query": "SELECT 1" });
assert_eq!(
override_source_concurrency("postgres", &mut cfg, 20),
Some("max_connections")
);
assert_eq!(cfg["max_connections"], 20);
validate_source_config("postgres", "row", cfg).expect("overridden config stays valid");
}
#[cfg(feature = "source-rest")]
#[test]
fn the_rest_source_resolves_to_its_request_knob() {
let mut cfg = json!({ "base_url": "https://api.example.com", "path": "/x" });
assert_eq!(
override_source_concurrency("rest", &mut cfg, 6),
Some("request_concurrency")
);
assert_eq!(cfg["request_concurrency"], 6);
validate_source_config("rest", "row", cfg).expect("overridden config stays valid");
let old = json!({ "base_url": "https://a", "path": "/x", "partition_concurrency": 2, "partitions": [{"id": 1}] });
validate_source_config("rest", "row", old).expect("old spelling accepted");
}
#[test]
fn an_unknown_connector_kind_is_ignored_rather_than_erroring() {
let mut cfg = json!({ "a": 1 });
assert_eq!(
override_source_concurrency("not-a-connector", &mut cfg, 4),
None
);
assert_eq!(
override_sink_concurrency("not-a-connector", &mut cfg, 4),
None
);
}
}
#[cfg(test)]
mod truncating_tests {
use super::*;
use serde_json::json;
#[test]
fn truncating_path_by_kind() {
let p = |k: &str, v: Value| sink_truncating_path(k, &v).map(str::to_string);
assert_eq!(
p("jsonl", json!({"path": "a.jsonl"})).as_deref(),
Some("a.jsonl")
);
assert_eq!(
p("csv", json!({"path": "a.csv", "append": false})).as_deref(),
Some("a.csv")
);
assert_eq!(p("jsonl", json!({"path": "a.jsonl", "append": true})), None);
let pq = json!({"destination": {"type": "local_path", "path": "o/a.parquet"}});
assert_eq!(p("parquet", pq).as_deref(), Some("o/a.parquet"));
let rolled = json!({
"destination": {"type": "local_path", "path": "o/a.parquet"},
"max_rows_per_file": 10
});
assert_eq!(p("parquet", rolled), None);
let dir = json!({"destination": {"type": "local_path", "path": "o/"}});
assert_eq!(p("parquet", dir), None);
let s3 = json!({"destination": {"type": "s3", "bucket": "b"}});
assert_eq!(p("parquet", s3), None);
assert_eq!(p("parquet", json!({})), None);
assert_eq!(
p("file", json!({"path": "a.jsonl"})).as_deref(),
Some("a.jsonl")
);
assert_eq!(
p("file", json!({"path": "a.jsonl", "mode": "append"})),
None
);
assert_eq!(
p("file", json!({"path": "a.jsonl", "if_exists": "append"})),
None
);
assert_eq!(
p("file", json!({"path": "a.jsonl", "if_exists": "error"})),
None
);
for replace in [
json!({"if_exists": "replace"}),
json!({"mode": "overwrite"}),
] {
let mut cfg = replace;
cfg["path"] = json!("b.jsonl");
assert_eq!(p("file", cfg).as_deref(), Some("b.jsonl"));
}
assert_eq!(
p(
"file",
json!({"path": "a.jsonl", "write_mode": "overwrite"})
),
None
);
assert_eq!(p("postgres", json!({"path": "x"})), None);
assert!(TRUNCATING_FILE_SINK_KINDS.contains(&"jsonl"));
}
}