#![warn(missing_docs)]
use crate::{
config::{ByteSize, CommunicationConfig, Input, NodeRunConfig},
id::{DataId, NodeId, OperatorId},
};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_with_expand_env::with_expand_envs;
use std::{
collections::{BTreeMap, BTreeSet},
fmt,
path::PathBuf,
};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "kebab-case")]
pub enum OutputFraming {
#[default]
Raw,
ArrowIpc,
}
pub const SHELL_SOURCE: &str = "shell";
pub const DYNAMIC_SOURCE: &str = "dynamic";
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
#[schemars(title = "dora-rs specification")]
pub struct Descriptor {
pub nodes: Vec<Node>,
#[schemars(skip)]
#[serde(default)]
pub communication: CommunicationConfig,
#[schemars(skip)]
#[serde(rename = "_unstable_deploy")]
pub deploy: Option<Deploy>,
#[schemars(skip)]
#[serde(default, rename = "_unstable_debug")]
pub debug: Debug,
#[serde(default)]
pub health_check_interval: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub strict_types: Option<bool>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub type_rules: Vec<TypeRuleDef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub env: Option<BTreeMap<String, EnvValue>>,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct TypeRuleDef {
pub from: String,
pub to: String,
}
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "kebab-case")]
pub enum RestartPolicy {
#[default]
Never,
OnFailure,
Always,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct Deploy {
pub machine: Option<String>,
pub working_dir: Option<PathBuf>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub labels: BTreeMap<String, String>,
#[serde(default)]
pub distribute: DistributeStrategy,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum DistributeStrategy {
#[default]
Local,
Scp,
Http,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)]
pub struct Debug {
#[serde(default, alias = "publish_all_messages_to_zenoh")]
pub enable_debug_inspection: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct Node {
pub id: NodeId,
pub name: Option<String>,
pub description: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub path: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub path_sha256: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub args: Option<String>,
pub env: Option<BTreeMap<String, EnvValue>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub operators: Option<RuntimeNode>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub operator: Option<SingleOperatorDefinition>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub ros2: Option<Ros2BridgeConfig>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub custom: Option<CustomNode>,
#[serde(default)]
pub outputs: BTreeSet<DataId>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub output_types: BTreeMap<DataId, String>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub output_framing: BTreeMap<DataId, OutputFraming>,
#[serde(default)]
pub inputs: BTreeMap<DataId, Input>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub input_types: BTreeMap<DataId, String>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub output_metadata: BTreeMap<DataId, Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub pattern: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub send_stdout_as: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub send_logs_as: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub min_log_level: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub max_log_size: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_rotated_files: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub build: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub git: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub hub: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub branch: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tag: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub rev: Option<String>,
#[serde(default)]
pub restart_policy: RestartPolicy,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub shared_memory_pool_size: Option<ByteSize>,
#[serde(default)]
pub max_restarts: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub restart_delay: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_restart_delay: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub restart_window: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub health_check_timeout: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub finish_grace_secs: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub module: Option<String>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub params: BTreeMap<String, String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cpu_affinity: Option<Vec<usize>>,
#[schemars(skip)]
#[serde(rename = "_unstable_deploy")]
pub deploy: Option<Deploy>,
}
#[allow(missing_docs)]
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ResolvedNode {
pub id: NodeId,
pub name: Option<String>,
pub description: Option<String>,
pub env: Option<BTreeMap<String, EnvValue>>,
#[serde(default)]
pub cpu_affinity: Option<Vec<usize>>,
#[serde(default)]
pub deploy: Option<Deploy>,
#[serde(flatten)]
pub kind: CoreNodeKind,
}
#[allow(missing_docs)]
impl ResolvedNode {
pub fn has_git_source(&self) -> bool {
self.kind
.as_custom()
.map(|n| n.source.is_git())
.unwrap_or_default()
}
}
#[allow(missing_docs)]
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
#[allow(clippy::large_enum_variant)]
pub enum CoreNodeKind {
#[serde(rename = "operators")]
Runtime(RuntimeNode),
Custom(CustomNode),
}
#[allow(missing_docs)]
impl CoreNodeKind {
pub fn as_custom(&self) -> Option<&CustomNode> {
match self {
CoreNodeKind::Runtime(_) => None,
CoreNodeKind::Custom(custom_node) => Some(custom_node),
}
}
}
#[allow(missing_docs)]
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(transparent)]
pub struct RuntimeNode {
pub operators: Vec<OperatorDefinition>,
}
#[allow(missing_docs)]
#[derive(Debug, Serialize, Deserialize, JsonSchema, Clone)]
pub struct OperatorDefinition {
pub id: OperatorId,
#[serde(flatten)]
pub config: OperatorConfig,
}
#[allow(missing_docs)]
#[derive(Debug, Serialize, Deserialize, JsonSchema, Clone)]
pub struct SingleOperatorDefinition {
pub id: Option<OperatorId>,
#[serde(flatten)]
pub config: OperatorConfig,
}
#[allow(missing_docs)]
#[derive(Debug, Serialize, Deserialize, JsonSchema, Clone)]
pub struct OperatorConfig {
pub name: Option<String>,
pub description: Option<String>,
#[serde(default)]
pub inputs: BTreeMap<DataId, Input>,
#[serde(default)]
pub outputs: BTreeSet<DataId>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub output_types: BTreeMap<DataId, String>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub output_framing: BTreeMap<DataId, OutputFraming>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub input_types: BTreeMap<DataId, String>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub output_metadata: BTreeMap<DataId, Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub pattern: Option<String>,
#[serde(flatten)]
pub source: OperatorSource,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub build: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub send_stdout_as: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub send_logs_as: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub min_log_level: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub max_log_size: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_rotated_files: Option<u32>,
}
#[allow(missing_docs)]
#[derive(Debug, Serialize, Deserialize, JsonSchema, Clone)]
#[serde(rename_all = "kebab-case")]
pub enum OperatorSource {
SharedLibrary(String),
Python(PythonSource),
#[schemars(skip)]
Wasm(String),
}
#[allow(missing_docs)]
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(from = "PythonSourceDef", into = "PythonSourceDef")]
pub struct PythonSource {
pub source: String,
pub conda_env: Option<String>,
}
#[allow(missing_docs)]
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(untagged)]
pub enum PythonSourceDef {
SourceOnly(String),
WithOptions {
source: String,
conda_env: Option<String>,
},
}
impl From<PythonSource> for PythonSourceDef {
fn from(input: PythonSource) -> Self {
match input {
PythonSource {
source,
conda_env: None,
} => Self::SourceOnly(source),
PythonSource { source, conda_env } => Self::WithOptions { source, conda_env },
}
}
}
impl From<PythonSourceDef> for PythonSource {
fn from(value: PythonSourceDef) -> Self {
match value {
PythonSourceDef::SourceOnly(source) => Self {
source,
conda_env: None,
},
PythonSourceDef::WithOptions { source, conda_env } => Self { source, conda_env },
}
}
}
#[allow(missing_docs)]
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct CustomNode {
pub path: String,
pub source: NodeSource,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub path_sha256: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub args: Option<String>,
pub envs: Option<BTreeMap<String, EnvValue>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub build: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub send_stdout_as: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub send_logs_as: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub min_log_level: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub max_log_size: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_rotated_files: Option<u32>,
#[serde(default)]
pub restart_policy: RestartPolicy,
#[serde(default)]
pub max_restarts: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub restart_delay: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_restart_delay: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub restart_window: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub health_check_timeout: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub finish_grace_secs: Option<f64>,
#[serde(flatten)]
pub run_config: NodeRunConfig,
}
#[allow(missing_docs)]
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub enum NodeSource {
Local,
GitBranch {
repo: String,
rev: Option<GitRepoRev>,
},
}
#[allow(missing_docs)]
impl NodeSource {
pub fn is_git(&self) -> bool {
matches!(self, Self::GitBranch { .. })
}
}
#[allow(missing_docs)]
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub enum GitRepoRev {
Branch(String),
Tag(String),
Rev(String),
}
#[allow(missing_docs)]
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(untagged)]
pub enum EnvValue {
#[serde(deserialize_with = "with_expand_envs")]
Bool(bool),
#[serde(deserialize_with = "with_expand_envs")]
Integer(i64),
#[serde(deserialize_with = "with_expand_envs")]
Float(f64),
#[serde(deserialize_with = "with_expand_envs")]
String(String),
}
impl fmt::Display for EnvValue {
fn fmt(&self, fmt: &mut fmt::Formatter) -> fmt::Result {
match self {
EnvValue::Bool(bool) => fmt.write_str(&bool.to_string()),
EnvValue::Integer(i64) => fmt.write_str(&i64.to_string()),
EnvValue::Float(f64) => fmt.write_str(&f64.to_string()),
EnvValue::String(str) => fmt.write_str(str),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct Ros2BridgeConfig {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub topic: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub message_type: Option<String>,
#[serde(default)]
pub direction: Ros2Direction,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub topics: Option<Vec<Ros2TopicConfig>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub service: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub service_type: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub action: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub action_type: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub role: Option<Ros2Role>,
#[serde(default)]
pub qos: Ros2QosConfig,
#[serde(default = "default_ros2_namespace")]
pub namespace: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub node_name: Option<String>,
}
impl Default for Ros2BridgeConfig {
fn default() -> Self {
Self {
topic: None,
message_type: None,
direction: Ros2Direction::default(),
topics: None,
service: None,
service_type: None,
action: None,
action_type: None,
role: None,
qos: Ros2QosConfig::default(),
namespace: default_ros2_namespace(),
node_name: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum Ros2Role {
Client,
Server,
}
fn default_ros2_namespace() -> String {
"/".to_string()
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct Ros2TopicConfig {
pub topic: String,
pub message_type: String,
#[serde(default)]
pub direction: Ros2Direction,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub input: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub qos: Option<Ros2QosConfig>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum Ros2Direction {
#[default]
Subscribe,
Publish,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct Ros2QosConfig {
#[serde(default)]
pub reliable: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub durability: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub liveliness: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub lease_duration: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_blocking_time: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub keep_last: Option<i32>,
#[serde(default)]
pub keep_all: bool,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn output_framing_defaults_to_raw() {
let yaml = r#"
nodes:
- id: test
path: test.py
outputs:
- data
"#;
let desc: Descriptor = serde_yaml::from_str(yaml).unwrap();
assert!(desc.nodes[0].output_framing.is_empty());
}
#[test]
fn output_framing_parses_arrow_ipc() {
let yaml = r#"
nodes:
- id: test
path: test.py
outputs:
- data
output_framing:
data: arrow-ipc
"#;
let desc: Descriptor = serde_yaml::from_str(yaml).unwrap();
assert_eq!(
desc.nodes[0].output_framing.get::<DataId>(&"data".into()),
Some(&OutputFraming::ArrowIpc)
);
}
#[test]
fn cpu_affinity_parses() {
let yaml = r#"
nodes:
- id: test
path: test.py
cpu_affinity: [0, 2, 4]
"#;
let desc: Descriptor = serde_yaml::from_str(yaml).unwrap();
assert_eq!(desc.nodes[0].cpu_affinity, Some(vec![0, 2, 4]));
}
#[test]
fn cpu_affinity_defaults_to_none() {
let yaml = r#"
nodes:
- id: test
path: test.py
"#;
let desc: Descriptor = serde_yaml::from_str(yaml).unwrap();
assert_eq!(desc.nodes[0].cpu_affinity, None);
}
#[test]
fn debug_flag_accepts_new_name() {
let yaml = r#"
nodes:
- id: test
path: test.py
_unstable_debug:
enable_debug_inspection: true
"#;
let desc: Descriptor = serde_yaml::from_str(yaml).unwrap();
assert!(desc.debug.enable_debug_inspection);
}
#[test]
fn debug_flag_accepts_legacy_alias() {
let yaml = r#"
nodes:
- id: test
path: test.py
_unstable_debug:
publish_all_messages_to_zenoh: true
"#;
let desc: Descriptor = serde_yaml::from_str(yaml).unwrap();
assert!(desc.debug.enable_debug_inspection);
}
}