use serde::{Deserialize, Serialize};
use crate::streaming::{ExtractionConfig, StreamMode};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ZenohConfig {
pub enabled: bool,
pub mode: ZenohMode,
pub connect: Vec<String>,
pub listen: Vec<String>,
pub prefix: String,
pub auto_topics: Vec<AutoTopic>,
#[serde(skip_serializing)]
pub api_key: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "lowercase")]
pub enum ZenohMode {
#[default]
Peer,
Client,
Router,
}
impl ZenohMode {
pub fn as_str(&self) -> &'static str {
match self {
Self::Peer => "peer",
Self::Client => "client",
Self::Router => "router",
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AutoTopic {
pub key_expr: String,
pub user_id: String,
#[serde(default)]
pub mode: StreamMode,
#[serde(default)]
pub payload_mode: PayloadMode,
#[serde(default)]
pub extraction_config: ExtractionConfig,
#[serde(default)]
pub tags: Vec<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "lowercase")]
pub enum PayloadMode {
#[default]
Passthrough,
Structured,
}
impl Default for ZenohConfig {
fn default() -> Self {
Self {
enabled: false,
mode: ZenohMode::default(),
connect: Vec::new(),
listen: Vec::new(),
prefix: "shodh".to_string(),
auto_topics: Vec::new(),
api_key: None,
}
}
}
impl ZenohConfig {
pub fn from_env() -> Self {
let mut config = Self::default();
if let Ok(val) = std::env::var("SHODH_ZENOH_ENABLED") {
config.enabled = val == "true" || val == "1";
}
if let Ok(val) = std::env::var("SHODH_ZENOH_MODE") {
config.mode = match val.to_lowercase().as_str() {
"client" => ZenohMode::Client,
"router" => ZenohMode::Router,
_ => ZenohMode::Peer,
};
}
if let Ok(val) = std::env::var("SHODH_ZENOH_CONNECT") {
config.connect = val
.split(',')
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.collect();
}
if let Ok(val) = std::env::var("SHODH_ZENOH_LISTEN") {
config.listen = val
.split(',')
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.collect();
}
if let Ok(val) = std::env::var("SHODH_ZENOH_PREFIX") {
let trimmed = val.trim().to_string();
if !trimmed.is_empty() {
config.prefix = trimmed;
}
}
if let Ok(val) = std::env::var("SHODH_ZENOH_AUTO_TOPICS") {
match serde_json::from_str::<Vec<AutoTopic>>(&val) {
Ok(topics) => config.auto_topics = topics,
Err(e) => {
tracing::warn!(
"Failed to parse SHODH_ZENOH_AUTO_TOPICS: {}. Expected JSON array.",
e
);
}
}
}
if let Ok(val) = std::env::var("SHODH_ZENOH_API_KEY") {
let trimmed = val.trim().to_string();
if !trimmed.is_empty() {
config.api_key = Some(trimmed);
}
}
config
}
pub fn validate(&self) {
if self.prefix.contains('/') {
tracing::warn!(
"SHODH_ZENOH_PREFIX contains '/' — this may cause unexpected key expression nesting"
);
}
for (i, topic) in self.auto_topics.iter().enumerate() {
if topic.key_expr.is_empty() {
tracing::warn!("Auto-topic [{}] has empty key_expr — will be skipped", i);
}
if topic.user_id.is_empty() {
tracing::warn!("Auto-topic [{}] has empty user_id — will be skipped", i);
}
}
if !self.connect.is_empty() && self.mode == ZenohMode::Router {
tracing::warn!(
"Zenoh mode is 'router' but connect endpoints are set — routers typically only listen"
);
}
let binds_all_interfaces = self
.listen
.iter()
.any(|ep| ep.contains("0.0.0.0") || ep.contains("[::]"));
if binds_all_interfaces && self.api_key.is_none() {
tracing::warn!(
"Zenoh listen endpoints include 0.0.0.0 but no SHODH_ZENOH_API_KEY is set — \
any network peer can invoke memory operations. Set SHODH_ZENOH_API_KEY or \
bind to 127.0.0.1 for local-only deployments."
);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_default_config() {
let config = ZenohConfig::default();
assert!(!config.enabled);
assert_eq!(config.mode, ZenohMode::Peer);
assert_eq!(config.prefix, "shodh");
assert!(config.connect.is_empty());
assert!(config.listen.is_empty());
assert!(config.auto_topics.is_empty());
}
#[test]
fn test_zenoh_mode_serde() {
let json = r#""client""#;
let mode: ZenohMode = serde_json::from_str(json).unwrap();
assert_eq!(mode, ZenohMode::Client);
}
#[test]
fn test_auto_topic_deserialization() {
let json = r#"{
"key_expr": "rt/spot1/status",
"user_id": "spot-1",
"mode": "sensor",
"payload_mode": "passthrough",
"tags": ["spot", "sensor"]
}"#;
let topic: AutoTopic = serde_json::from_str(json).unwrap();
assert_eq!(topic.key_expr, "rt/spot1/status");
assert_eq!(topic.user_id, "spot-1");
assert_eq!(topic.mode, StreamMode::Sensor);
assert_eq!(topic.payload_mode, PayloadMode::Passthrough);
assert_eq!(topic.tags, vec!["spot", "sensor"]);
}
#[test]
fn test_auto_topic_defaults() {
let json = r#"{"key_expr": "test/topic", "user_id": "u1"}"#;
let topic: AutoTopic = serde_json::from_str(json).unwrap();
assert_eq!(topic.mode, StreamMode::Conversation);
assert_eq!(topic.payload_mode, PayloadMode::Passthrough);
assert!(topic.tags.is_empty());
}
#[test]
fn test_payload_mode_serde() {
let json = r#""structured""#;
let mode: PayloadMode = serde_json::from_str(json).unwrap();
assert_eq!(mode, PayloadMode::Structured);
}
}