use faucet_core::DEFAULT_BATCH_SIZE;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct CsvSinkConfig {
pub path: String,
#[serde(default = "default_delimiter")]
pub delimiter: u8,
#[serde(default = "default_true")]
pub write_headers: bool,
#[serde(default)]
pub append: bool,
#[serde(default = "default_batch_size")]
pub batch_size: usize,
#[serde(default)]
pub on_unknown_field: OnUnknownField,
#[cfg(feature = "compression")]
#[serde(default)]
pub compression: faucet_core::CompressionConfig,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum OnUnknownField {
#[default]
Warn,
Error,
}
fn default_delimiter() -> u8 {
b','
}
fn default_true() -> bool {
true
}
fn default_batch_size() -> usize {
DEFAULT_BATCH_SIZE
}
impl CsvSinkConfig {
pub fn new(path: impl Into<String>) -> Self {
Self {
path: path.into(),
delimiter: b',',
write_headers: true,
append: false,
batch_size: DEFAULT_BATCH_SIZE,
on_unknown_field: OnUnknownField::Warn,
#[cfg(feature = "compression")]
compression: faucet_core::CompressionConfig::Auto,
}
}
pub fn delimiter(mut self, d: u8) -> Self {
self.delimiter = d;
self
}
pub fn write_headers(mut self, v: bool) -> Self {
self.write_headers = v;
self
}
pub fn append(mut self, v: bool) -> Self {
self.append = v;
self
}
pub fn with_batch_size(mut self, batch_size: usize) -> Self {
self.batch_size = batch_size;
self
}
pub fn on_unknown_field(mut self, policy: OnUnknownField) -> Self {
self.on_unknown_field = policy;
self
}
#[cfg(feature = "compression")]
pub fn compression(mut self, c: faucet_core::CompressionConfig) -> Self {
self.compression = c;
self
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn default_config() {
let config = CsvSinkConfig::new("/tmp/out.csv");
assert_eq!(config.path, "/tmp/out.csv");
assert_eq!(config.delimiter, b',');
assert!(config.write_headers);
assert!(!config.append);
}
#[test]
fn builder_methods() {
let config = CsvSinkConfig::new("/tmp/out.tsv")
.delimiter(b'\t')
.write_headers(false)
.append(true);
assert_eq!(config.delimiter, b'\t');
assert!(!config.write_headers);
assert!(config.append);
}
#[test]
fn batch_size_defaults_to_default_batch_size() {
let config = CsvSinkConfig::new("/tmp/out.csv");
assert_eq!(config.batch_size, faucet_core::DEFAULT_BATCH_SIZE);
}
#[test]
fn with_batch_size_overrides_default() {
let config = CsvSinkConfig::new("/tmp/out.csv").with_batch_size(250);
assert_eq!(config.batch_size, 250);
}
#[test]
fn batch_size_zero_is_accepted_as_no_batching_sentinel() {
let config = CsvSinkConfig::new("/tmp/out.csv").with_batch_size(0);
assert_eq!(config.batch_size, 0);
assert!(faucet_core::validate_batch_size(config.batch_size).is_ok());
}
#[test]
fn batch_size_above_max_is_rejected_by_validate_batch_size() {
let config =
CsvSinkConfig::new("/tmp/out.csv").with_batch_size(faucet_core::MAX_BATCH_SIZE + 1);
assert!(faucet_core::validate_batch_size(config.batch_size).is_err());
}
#[test]
fn batch_size_deserializes_from_json() {
let json = r#"{
"path": "/tmp/out.csv",
"delimiter": 44,
"write_headers": true,
"append": false,
"batch_size": 500
}"#;
let config: CsvSinkConfig = serde_json::from_str(json).unwrap();
assert_eq!(config.batch_size, 500);
}
#[test]
fn batch_size_defaults_when_missing_in_json() {
let json = r#"{"path": "/tmp/out.csv"}"#;
let config: CsvSinkConfig = serde_json::from_str(json).unwrap();
assert_eq!(config.batch_size, faucet_core::DEFAULT_BATCH_SIZE);
}
#[test]
fn on_unknown_field_defaults_to_warn() {
let config = CsvSinkConfig::new("/tmp/out.csv");
assert_eq!(config.on_unknown_field, OnUnknownField::Warn);
let json = r#"{"path": "/tmp/out.csv"}"#;
let config: CsvSinkConfig = serde_json::from_str(json).unwrap();
assert_eq!(config.on_unknown_field, OnUnknownField::Warn);
}
#[test]
fn on_unknown_field_deserializes_snake_case() {
let json = r#"{"path": "/tmp/out.csv", "on_unknown_field": "error"}"#;
let config: CsvSinkConfig = serde_json::from_str(json).unwrap();
assert_eq!(config.on_unknown_field, OnUnknownField::Error);
}
#[test]
fn on_unknown_field_builder_sets_policy() {
let config = CsvSinkConfig::new("/tmp/out.csv").on_unknown_field(OnUnknownField::Error);
assert_eq!(config.on_unknown_field, OnUnknownField::Error);
}
}