use serde::Deserialize as _;
use std::collections::HashMap;
use std::path::Path;
use noyalib::compat::serde_yaml;
use crate::commands::test::document::TestDocError;
use crate::commands::test::runner::find_camel_toml_root;
const BODY_SCALAR_SENTINEL: &str = "unsupported body scalar: ";
pub(crate) const JOB_SAFE_CONSUMER_SCHEMES: [&str; 4] = ["direct", "seda", "log", "mock"];
const JOB_SEND_SCHEMES: [&str; 2] = ["direct", "seda"];
#[derive(Debug)]
pub(crate) struct JobDocument {
pub(crate) execute: ExecuteSection,
pub(crate) route_files: Option<Vec<String>>,
pub(crate) route_files_from_root: Option<Vec<String>>,
pub(crate) routes: Option<serde_yaml::Value>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum JobMode {
OneShot,
Batch,
}
impl JobMode {
pub(crate) fn as_str(&self) -> &'static str {
match self {
Self::OneShot => "one-shot",
Self::Batch => "batch",
}
}
}
#[derive(Debug)]
pub(crate) struct ExecuteSection {
pub(crate) mode: JobMode,
pub(crate) send: JobSendAction,
pub(crate) capture_reply: bool,
pub(crate) timeout: std::time::Duration,
}
#[derive(Debug)]
pub(crate) struct JobSendAction {
pub(crate) to: String,
pub(crate) body: Option<JobBody>,
pub(crate) headers: Option<HashMap<String, serde_json::Value>>,
}
#[derive(Debug)]
pub(crate) enum JobBody {
Text(String),
Json(serde_json::Value),
}
#[derive(Debug)]
pub(crate) enum JobDocError {
NotJobSuffix { path: String },
Yaml(String),
UnknownField(String),
MissingExecute,
ExclusiveWithScenario,
MixedVocabulary { sections: Vec<&'static str> },
MissingMode,
UnsupportedMode(String),
MissingTimeout,
InvalidTimeout(String),
UnsupportedSendScheme { to: String },
UnsupportedBodyScalar(String),
RouteSource(TestDocError),
}
impl std::fmt::Display for JobDocError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::NotJobSuffix { path } => write!(
f,
"job document {path} must use the reserved .job.yaml/.job.yml suffix (rename it — .test.yaml names a camel test document)"
),
Self::Yaml(raw) => write!(f, "invalid job document: {raw}"),
Self::UnknownField(raw) => write!(f, "unknown field in job document: {raw}"),
Self::MissingExecute => {
write!(f, "job document requires an execute: section")
}
Self::ExclusiveWithScenario => {
write!(f, "execute: and scenario: are mutually exclusive sections")
}
Self::MixedVocabulary { sections } => write!(
f,
"execute: is mutually exclusive with the test sections {}",
sections
.iter()
.map(|s| format!("`{s}`"))
.collect::<Vec<_>>()
.join(", ")
),
Self::MissingMode => write!(f, "execute.mode is required"),
Self::UnsupportedMode(mode) => write!(
f,
"unsupported execute.mode `{mode}`: expected `one-shot` or `batch`"
),
Self::MissingTimeout => write!(f, "execute.timeout is required"),
Self::InvalidTimeout(raw) => write!(
f,
"invalid execute.timeout `{raw}`: expected a positive duration (e.g. `30s`)"
),
Self::UnsupportedSendScheme { to } => {
write!(f, "send target `{to}` must start with `direct:` or `seda:`")
}
Self::UnsupportedBodyScalar(raw) => write!(
f,
"unsupported body scalar `{raw}`: only string, object, and array bodies are supported"
),
Self::RouteSource(err) => write!(f, "{err}"),
}
}
}
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
struct ExecuteSectionDoc {
mode: Option<String>,
#[serde(default)]
send: Option<JobSendActionDoc>,
#[serde(rename = "capture-reply")]
capture_reply: Option<bool>,
timeout: Option<String>,
}
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
struct JobSendActionDoc {
to: String,
#[serde(default, deserialize_with = "deserialize_option_job_body")]
body: Option<JobBody>,
headers: Option<HashMap<String, serde_json::Value>>,
}
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
struct JobDocumentDoc {
execute: ExecuteSectionDoc,
#[serde(default)]
#[expect(dead_code, reason = "admit-only under deny_unknown_fields")]
description: Option<String>,
route_files: Option<Vec<String>>,
route_files_from_root: Option<Vec<String>>,
routes: Option<serde_yaml::Value>,
}
fn job_body_from_value(value: serde_json::Value) -> Result<Option<JobBody>, String> {
match &value {
serde_json::Value::String(s) => Ok(Some(JobBody::Text(s.clone()))),
serde_json::Value::Object(_) | serde_json::Value::Array(_) => {
Ok(Some(JobBody::Json(value)))
}
scalar => Err(format!("{BODY_SCALAR_SENTINEL}{scalar}")),
}
}
fn deserialize_option_job_body<'de, D>(deserializer: D) -> Result<Option<JobBody>, D::Error>
where
D: serde::Deserializer<'de>,
{
let value = serde_json::Value::deserialize(deserializer)?;
job_body_from_value(value).map_err(serde::de::Error::custom)
}
const TEST_VOCABULARY_KEYS: [&str; 8] = [
"inputs",
"expects",
"intercepts",
"beans",
"repositories",
"sequence",
"settle",
"env",
];
pub(crate) fn parse_job_document(path: &Path, text: &str) -> Result<JobDocument, JobDocError> {
if !camel_dsl::discovery::is_job_document(path) {
return Err(JobDocError::NotJobSuffix {
path: path.display().to_string(),
});
}
let value = serde_yaml::from_str::<serde_yaml::Value>(text)
.map_err(|e| JobDocError::Yaml(e.to_string()))?;
let has = |key: &str| value.get(key).is_some();
if !has("execute") {
return Err(JobDocError::MissingExecute);
}
if has("scenario") {
return Err(JobDocError::ExclusiveWithScenario);
}
let mixed: Vec<&'static str> = TEST_VOCABULARY_KEYS
.iter()
.copied()
.filter(|key| has(key))
.collect();
if !mixed.is_empty() {
return Err(JobDocError::MixedVocabulary { sections: mixed });
}
let raw = serde_yaml::from_str::<JobDocumentDoc>(text).map_err(|e| classify(&e.to_string()))?;
let mut present: Vec<&'static str> = Vec::new();
if raw.route_files.is_some() {
present.push("routeFiles");
}
if raw.route_files_from_root.is_some() {
present.push("routeFilesFromRoot");
}
if raw.routes.is_some() {
present.push("routes");
}
if present.len() != 1 {
return Err(JobDocError::RouteSource(
TestDocError::RouteSourceConflict { present },
));
}
let execute = raw.execute;
let mode = match execute.mode.ok_or(JobDocError::MissingMode)?.as_str() {
"one-shot" => JobMode::OneShot,
"batch" => JobMode::Batch,
other => return Err(JobDocError::UnsupportedMode(other.to_string())),
};
let timeout_raw = execute.timeout.ok_or(JobDocError::MissingTimeout)?;
let timeout = humantime::parse_duration(&timeout_raw)
.ok()
.filter(|d| *d > std::time::Duration::ZERO)
.ok_or_else(|| JobDocError::InvalidTimeout(timeout_raw.clone()))?;
let send = execute.send.ok_or(JobDocError::Yaml(
"execute.send is required: exactly one send action".to_string(),
))?;
if !JOB_SEND_SCHEMES
.iter()
.any(|scheme| send.to.starts_with(&format!("{scheme}:")))
{
return Err(JobDocError::UnsupportedSendScheme { to: send.to });
}
Ok(JobDocument {
execute: ExecuteSection {
mode,
send: JobSendAction {
to: send.to,
body: send.body,
headers: send.headers,
},
capture_reply: execute.capture_reply.unwrap_or(false),
timeout,
},
route_files: raw.route_files,
route_files_from_root: raw.route_files_from_root,
routes: raw.routes,
})
}
fn classify(raw: &str) -> JobDocError {
if let Some((_, after)) = raw.split_once(BODY_SCALAR_SENTINEL) {
let scalar = after.split_whitespace().next().unwrap_or_default();
return JobDocError::UnsupportedBodyScalar(scalar.to_string());
}
if raw.contains("unknown field") {
return JobDocError::UnknownField(raw.to_string());
}
JobDocError::Yaml(raw.to_string())
}
pub(crate) enum JobRouteSource {
Patterns(Vec<String>),
Inline(String),
}
pub(crate) fn resolve_route_source(
doc: &JobDocument,
doc_dir: &Path,
) -> Result<JobRouteSource, JobDocError> {
if let Some(files) = &doc.route_files_from_root {
let root = find_camel_toml_root(doc_dir).ok_or_else(|| {
JobDocError::RouteSource(TestDocError::NoProjectRoot {
doc_dir: doc_dir.display().to_string(),
})
})?;
Ok(JobRouteSource::Patterns(
files
.iter()
.map(|p| root.join(p).display().to_string())
.collect(),
))
} else if let Some(files) = &doc.route_files {
Ok(JobRouteSource::Patterns(
files
.iter()
.map(|p| doc_dir.join(p).display().to_string())
.collect(),
))
} else if let Some(value) = &doc.routes {
let mut mapping = serde_yaml::Mapping::new();
mapping.insert("routes", value.clone());
let text = serde_yaml::to_string(&serde_yaml::Value::Mapping(mapping))
.map_err(|e| JobDocError::Yaml(format!("failed to serialize inline routes: {e}")))?;
Ok(JobRouteSource::Inline(text))
} else {
Err(JobDocError::RouteSource(
TestDocError::RouteSourceConflict {
present: Vec::new(),
},
))
}
}
pub(crate) fn scheme_of_uri(uri: &str) -> Option<&str> {
uri.split_once(':').map(|(scheme, _)| scheme)
}
pub(crate) fn uri_base(uri: &str) -> &str {
uri.split('?').next().unwrap_or(uri)
}
pub(crate) fn seda_send_uri(to: &str) -> String {
const PARAM: &str = "waitForTaskToComplete";
match to.split_once('?') {
None => format!("{to}?{PARAM}=Always"),
Some((base, query)) => {
let kept: Vec<&str> = query
.split('&')
.filter(|pair| !pair.starts_with(&format!("{PARAM}=")))
.collect();
let mut uri = String::from(base);
uri.push('?');
if !kept.is_empty() {
uri.push_str(&kept.join("&"));
uri.push('&');
}
uri.push_str(&format!("{PARAM}=Always"));
uri
}
}
}
pub(crate) fn target_route_ids(
defs: &[camel_core::RouteDefinition],
target_base: &str,
) -> Vec<String> {
defs.iter()
.filter(|def| uri_base(def.from_uri()) == target_base)
.map(|def| def.route_id().to_string())
.collect()
}
pub(crate) fn validate_consumer_uri(from_uri: &str) -> Result<(), String> {
let scheme = scheme_of_uri(from_uri).unwrap_or_default();
if JOB_SAFE_CONSUMER_SCHEMES.contains(&scheme) {
Ok(())
} else {
Err(format!(
"route consumes from `{from_uri}`; one-shot job documents allow only {} \
consumers (producers/sinks as to: URIs are unrestricted); scheme `{scheme}` \
is rejected",
JOB_SAFE_CONSUMER_SCHEMES
.iter()
.map(|s| format!("`{s}:`"))
.collect::<Vec<_>>()
.join(", ")
))
}
}