use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct SubscriptionKafkaConfig {
pub endpoint: String,
pub default_topic: String,
pub client_id: String,
pub timeout_ms: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub compression: Option<String>,
}
impl Default for SubscriptionKafkaConfig {
fn default() -> Self {
Self {
endpoint: String::new(),
default_topic: "fraiseql.subscriptions".to_owned(),
client_id: "fraiseql-subscriptions".to_owned(),
timeout_ms: 5_000,
compression: None,
}
}
}
impl SubscriptionKafkaConfig {
pub fn validate(&self) -> Result<(), String> {
if self.endpoint.trim().is_empty() {
return Err("[subscription_kafka] endpoint is empty. Omit the section to \
disable the transport; an empty endpoint is a mistake, not a \
way to switch it off."
.to_owned());
}
if self.default_topic.trim().is_empty() {
return Err("[subscription_kafka] default_topic is empty.".to_owned());
}
if self.timeout_ms == 0 {
return Err("[subscription_kafka] timeout_ms must be at least 1.".to_owned());
}
Ok(())
}
}
#[cfg(test)]
mod tests;