1use async_trait::async_trait;
2use serde::{Deserialize, Serialize};
3
4use crate::CamelError;
5use crate::declarative::LanguageExpressionDef;
6use crate::splitter::StreamSplitConfig;
7
8pub const CANONICAL_CONTRACT_NAME: &str = "canonical-v1";
9pub const CANONICAL_CONTRACT_VERSION: u32 = 2;
10pub const CANONICAL_CONTRACT_SUPPORTED_STEPS: &[&str] = &[
11 "to",
12 "log",
13 "wire_tap",
14 "script",
15 "filter",
16 "choice",
17 "split",
18 "aggregate",
19 "stop",
20 "delay",
21 "cache",
22 "cache_invalidate",
23 "cache_peek_stale",
24];
25pub const CANONICAL_CONTRACT_DECLARATIVE_ONLY_STEPS: &[&str] =
26 &["script", "filter", "choice", "split"];
27pub const CANONICAL_CONTRACT_EXCLUDED_DECLARATIVE_STEPS: &[&str] = &[
28 "set_header",
29 "set_property",
30 "set_body",
31 "multicast",
32 "convert_body_to",
33 "bean",
34 "marshal",
35 "unmarshal",
36];
37pub const CANONICAL_CONTRACT_RUST_ONLY_STEPS: &[&str] = &[
38 "processor",
39 "process",
40 "process_fn",
41 "map_body",
42 "set_body_fn",
43 "set_header_fn",
44];
45
46pub fn canonical_contract_supports_step(step: &str) -> bool {
47 CANONICAL_CONTRACT_SUPPORTED_STEPS.contains(&step)
48}
49
50pub fn canonical_contract_rejection_reason(step: &str) -> Option<&'static str> {
51 if CANONICAL_CONTRACT_EXCLUDED_DECLARATIVE_STEPS.contains(&step) {
52 return Some(
53 "declared out-of-scope for canonical v2; use declarative route compilation path outside CQRS canonical commands",
54 );
55 }
56
57 if CANONICAL_CONTRACT_RUST_ONLY_STEPS.contains(&step) {
58 return Some("rust-only programmable step; not representable in canonical v2 contract");
59 }
60
61 if canonical_contract_supports_step(step)
62 && CANONICAL_CONTRACT_DECLARATIVE_ONLY_STEPS.contains(&step)
63 {
64 return Some(
65 "supported only as declarative/serializable expression form; closure/processor variants are outside canonical v2",
66 );
67 }
68
69 None
70}
71
72#[derive(
73 Debug,
74 Clone,
75 PartialEq,
76 Eq,
77 serde::Serialize,
78 serde::Deserialize,
79 schemars::JsonSchema,
80 ts_rs::TS,
81)]
82#[serde(rename_all = "snake_case")]
83#[ts(rename_all = "snake_case")]
84pub struct CanonicalRouteSpec {
85 pub route_id: String,
95 pub from: String,
96 pub steps: Vec<CanonicalStepSpec>,
97 pub circuit_breaker: Option<CanonicalCircuitBreakerSpec>,
98 pub auto_startup: Option<bool>,
99 pub startup_order: Option<i32>,
100 pub concurrency: Option<CanonicalConcurrencySpec>,
101 pub version: u32,
102}
103
104#[derive(
105 Debug,
106 Clone,
107 PartialEq,
108 Eq,
109 serde::Serialize,
110 serde::Deserialize,
111 schemars::JsonSchema,
112 ts_rs::TS,
113)]
114#[serde(tag = "step", content = "config", rename_all = "snake_case")]
115#[ts(rename_all = "snake_case")]
116#[non_exhaustive]
117pub enum CanonicalStepSpec {
118 To {
119 uri: String,
120 },
121 Log {
122 message: String,
123 },
124 WireTap {
125 uri: String,
126 },
127 Script {
128 expression: LanguageExpressionDef,
129 },
130 Filter {
131 predicate: LanguageExpressionDef,
132 steps: Vec<CanonicalStepSpec>,
133 },
134 Choice {
135 whens: Vec<CanonicalWhenSpec>,
136 otherwise: Option<Vec<CanonicalStepSpec>>,
137 },
138 Split {
139 expression: CanonicalSplitExpressionSpec,
140 aggregation: CanonicalSplitAggregationSpec,
141 parallel: bool,
142 parallel_limit: Option<usize>,
143 trace_item_threshold: Option<usize>,
146 stop_on_exception: bool,
147 steps: Vec<CanonicalStepSpec>,
148 },
149 Aggregate(CanonicalAggregateSpec),
150 Stop,
151 Delay {
152 #[ts(type = "number")]
153 delay_ms: u64,
154 dynamic_header: Option<String>,
155 },
156 Cache {
157 repository: Option<String>,
158 key: String,
159 ttl: Option<String>,
160 max_entry_bytes: Option<usize>,
161 #[serde(default, skip_serializing_if = "Option::is_none")]
164 coalesce_misses: Option<bool>,
165 on_miss: Vec<CanonicalStepSpec>,
166 },
167 CacheInvalidate {
168 repository: Option<String>,
169 #[serde(default, skip_serializing_if = "Option::is_none")]
171 key: Option<String>,
172 #[serde(default, skip_serializing_if = "Option::is_none")]
174 key_prefix: Option<String>,
175 },
176 CacheClear {
177 repository: Option<String>,
178 },
179 CacheStats {
180 repository: Option<String>,
181 },
182 CachePeekStale {
183 repository: Option<String>,
184 key: String,
185 #[serde(default, skip_serializing_if = "Option::is_none")]
187 on_miss: Option<String>,
188 },
189}
190
191#[derive(
192 Debug,
193 Clone,
194 PartialEq,
195 Eq,
196 serde::Serialize,
197 serde::Deserialize,
198 schemars::JsonSchema,
199 ts_rs::TS,
200)]
201#[serde(rename_all = "snake_case")]
202#[ts(rename_all = "snake_case")]
203pub struct CanonicalWhenSpec {
204 pub predicate: LanguageExpressionDef,
205 pub steps: Vec<CanonicalStepSpec>,
206}
207
208#[derive(
209 Debug,
210 Clone,
211 PartialEq,
212 Eq,
213 serde::Serialize,
214 serde::Deserialize,
215 schemars::JsonSchema,
216 ts_rs::TS,
217)]
218#[serde(rename_all = "snake_case")]
219#[ts(rename_all = "snake_case")]
220#[non_exhaustive]
221pub enum CanonicalSplitExpressionSpec {
222 BodyLines,
223 BodyJsonArray,
224 Language(LanguageExpressionDef),
225 Stream(StreamSplitConfig),
226}
227
228#[derive(
229 Debug,
230 Clone,
231 PartialEq,
232 Eq,
233 serde::Serialize,
234 serde::Deserialize,
235 schemars::JsonSchema,
236 ts_rs::TS,
237)]
238#[serde(rename_all = "snake_case")]
239#[ts(rename_all = "snake_case")]
240#[non_exhaustive]
241pub enum CanonicalSplitAggregationSpec {
242 LastWins,
243 CollectAll,
244 Original,
245}
246
247#[derive(
248 Debug,
249 Clone,
250 PartialEq,
251 Eq,
252 serde::Serialize,
253 serde::Deserialize,
254 schemars::JsonSchema,
255 ts_rs::TS,
256)]
257#[serde(rename_all = "snake_case")]
258#[ts(rename_all = "snake_case")]
259#[non_exhaustive]
260pub enum CanonicalAggregateStrategySpec {
261 CollectAll,
262}
263
264#[derive(
265 Debug,
266 Clone,
267 PartialEq,
268 Eq,
269 serde::Serialize,
270 serde::Deserialize,
271 schemars::JsonSchema,
272 ts_rs::TS,
273)]
274#[serde(rename_all = "snake_case")]
275#[ts(rename_all = "snake_case")]
276pub struct CanonicalAggregateSpec {
277 pub header: String,
278 pub completion_size: Option<usize>,
279 #[ts(type = "number")]
280 pub completion_timeout_ms: Option<u64>,
281 pub correlation_key: Option<String>,
282 pub force_completion_on_stop: Option<bool>,
283 pub discard_on_timeout: Option<bool>,
284 pub strategy: CanonicalAggregateStrategySpec,
285 pub max_buckets: Option<usize>,
286 #[serde(default)]
289 pub max_bucket_size: Option<usize>,
290 #[ts(type = "number")]
291 pub bucket_ttl_ms: Option<u64>,
292 #[serde(default)]
296 pub completion_predicate: Option<LanguageExpressionDef>,
297}
298
299#[derive(
300 Debug,
301 Clone,
302 PartialEq,
303 Eq,
304 serde::Serialize,
305 serde::Deserialize,
306 schemars::JsonSchema,
307 ts_rs::TS,
308)]
309#[serde(rename_all = "snake_case")]
310#[ts(rename_all = "snake_case")]
311pub struct CanonicalCircuitBreakerSpec {
312 pub failure_threshold: u32,
313 #[ts(type = "number")]
314 pub open_duration_ms: u64,
315 #[serde(default, skip_serializing_if = "Vec::is_empty")]
316 pub fallback: Vec<CanonicalStepSpec>,
317}
318
319#[derive(
320 Debug,
321 Clone,
322 PartialEq,
323 Eq,
324 serde::Serialize,
325 serde::Deserialize,
326 schemars::JsonSchema,
327 ts_rs::TS,
328)]
329#[serde(tag = "mode", rename_all = "snake_case")]
330#[non_exhaustive]
331pub enum CanonicalConcurrencySpec {
332 Sequential,
333 Concurrent { max: usize },
334}
335
336impl CanonicalRouteSpec {
337 pub fn new(route_id: impl Into<String>, from: impl Into<String>) -> Self {
338 Self {
339 route_id: route_id.into(),
340 from: from.into(),
341 steps: Vec::new(),
342 circuit_breaker: None,
343 auto_startup: None,
344 startup_order: None,
345 concurrency: None,
346 version: CANONICAL_CONTRACT_VERSION,
347 }
348 }
349
350 pub fn with_auto_startup(mut self, auto: bool) -> Self {
351 self.auto_startup = Some(auto);
352 self
353 }
354
355 pub fn with_startup_order(mut self, order: i32) -> Self {
356 self.startup_order = Some(order);
357 self
358 }
359
360 pub fn with_concurrency(mut self, concurrency: CanonicalConcurrencySpec) -> Self {
361 self.concurrency = Some(concurrency);
362 self
363 }
364
365 pub fn validate_contract(&self) -> Result<(), CamelError> {
366 if self.route_id.trim().is_empty() {
367 return Err(CamelError::RouteError(
368 "canonical contract violation: route_id cannot be empty".to_string(),
369 ));
370 }
371 if self.from.trim().is_empty() {
372 return Err(CamelError::RouteError(
373 "canonical contract violation: from cannot be empty".to_string(),
374 ));
375 }
376 if self.version == 0 || self.version > CANONICAL_CONTRACT_VERSION {
377 return Err(CamelError::RouteError(format!(
378 "canonical contract violation: expected version {}, got {}",
379 CANONICAL_CONTRACT_VERSION, self.version
380 )));
381 }
382 validate_steps(&self.steps)?;
383 if let Some(cb) = &self.circuit_breaker {
384 if cb.failure_threshold == 0 {
385 return Err(CamelError::RouteError(
386 "canonical contract violation: circuit_breaker.failure_threshold must be > 0"
387 .to_string(),
388 ));
389 }
390 if cb.open_duration_ms == 0 {
391 return Err(CamelError::RouteError(
392 "canonical contract violation: circuit_breaker.open_duration_ms must be > 0"
393 .to_string(),
394 ));
395 }
396 validate_steps(&cb.fallback)?;
397 }
398 if let Some(CanonicalConcurrencySpec::Concurrent { max: 0 }) = &self.concurrency {
399 return Err(CamelError::RouteError(
400 "canonical contract violation: concurrency max must be > 0".to_string(),
401 ));
402 }
403 Ok(())
404 }
405}
406
407#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
408pub struct CanonicalFieldLoss {
409 pub field: &'static str,
410 pub reason: String,
411 pub target_version: u32,
412}
413
414#[derive(Debug, Clone, PartialEq, Eq, Default, serde::Serialize)]
415pub struct CanonicalLossReport {
416 pub dropped_fields: Vec<CanonicalFieldLoss>,
417}
418
419impl CanonicalLossReport {
420 pub fn from_field(field: &'static str, reason: &str, target_version: u32) -> Self {
421 Self {
422 dropped_fields: vec![CanonicalFieldLoss {
423 field,
424 reason: reason.to_string(),
425 target_version,
426 }],
427 }
428 }
429
430 pub fn is_empty(&self) -> bool {
431 self.dropped_fields.is_empty()
432 }
433}
434
435fn validate_steps(steps: &[CanonicalStepSpec]) -> Result<(), CamelError> {
436 for step in steps {
437 match step {
438 CanonicalStepSpec::To { uri } | CanonicalStepSpec::WireTap { uri } => {
439 if uri.trim().is_empty() {
440 return Err(CamelError::RouteError(
441 "canonical contract violation: endpoint uri cannot be empty".to_string(),
442 ));
443 }
444 }
445 CanonicalStepSpec::Filter { steps, .. } => validate_steps(steps)?,
446 CanonicalStepSpec::Choice { whens, otherwise } => {
447 for when in whens {
448 validate_steps(&when.steps)?;
449 }
450 if let Some(otherwise) = otherwise {
451 validate_steps(otherwise)?;
452 }
453 }
454 CanonicalStepSpec::Split {
455 parallel_limit,
456 steps,
457 ..
458 } => {
459 if let Some(limit) = parallel_limit
460 && *limit == 0
461 {
462 return Err(CamelError::RouteError(
463 "canonical contract violation: split.parallel_limit must be > 0"
464 .to_string(),
465 ));
466 }
467 validate_steps(steps)?;
468 }
469 CanonicalStepSpec::Aggregate(config) => {
470 let header_present = !config.header.trim().is_empty();
471 let key_present = config
472 .correlation_key
473 .as_deref()
474 .is_some_and(|k| !k.trim().is_empty());
475 if !key_present && config.correlation_key.is_some() {
476 return Err(CamelError::RouteError(
477 "canonical contract violation: aggregate.correlation_key cannot be empty"
478 .to_string(),
479 ));
480 }
481 if !header_present && !key_present {
482 return Err(CamelError::RouteError(
483 "canonical contract violation: aggregate requires a correlation source: header or correlation_key"
484 .to_string(),
485 ));
486 }
487 if let Some(size) = config.completion_size
488 && size == 0
489 {
490 return Err(CamelError::RouteError(
491 "canonical contract violation: aggregate.completion_size must be > 0"
492 .to_string(),
493 ));
494 }
495 }
496 CanonicalStepSpec::Cache { on_miss, .. } => {
497 validate_steps(on_miss)?;
498 }
499 CanonicalStepSpec::Log { .. }
500 | CanonicalStepSpec::Script { .. }
501 | CanonicalStepSpec::Stop
502 | CanonicalStepSpec::Delay { .. }
503 | CanonicalStepSpec::CacheInvalidate { .. }
504 | CanonicalStepSpec::CacheClear { .. }
505 | CanonicalStepSpec::CacheStats { .. }
506 | CanonicalStepSpec::CachePeekStale { .. } => {}
507 }
508 }
509 Ok(())
510}
511
512#[derive(Debug, Clone, PartialEq, Eq)]
513#[non_exhaustive]
514pub enum RuntimeCommand {
515 RegisterRoute {
516 spec: CanonicalRouteSpec,
517 command_id: String,
518 causation_id: Option<String>,
519 },
520 StartRoute {
521 route_id: String,
522 command_id: String,
523 causation_id: Option<String>,
524 },
525 StopRoute {
526 route_id: String,
527 command_id: String,
528 causation_id: Option<String>,
529 },
530 SuspendRoute {
531 route_id: String,
532 command_id: String,
533 causation_id: Option<String>,
534 },
535 ResumeRoute {
536 route_id: String,
537 command_id: String,
538 causation_id: Option<String>,
539 },
540 ReloadRoute {
541 route_id: String,
542 command_id: String,
543 causation_id: Option<String>,
544 },
545 FailRoute {
549 route_id: String,
550 error: String,
551 command_id: String,
552 causation_id: Option<String>,
553 },
554 RemoveRoute {
555 route_id: String,
556 command_id: String,
557 causation_id: Option<String>,
558 },
559 ReloadTlsCerts {
560 scheme: String,
561 host: String,
562 port: u16,
563 command_id: String,
564 causation_id: Option<String>,
565 },
566 ReloadTemplates {
573 route_id: String,
574 command_id: String,
575 causation_id: Option<String>,
576 },
577}
578
579impl RuntimeCommand {
580 pub fn command_id(&self) -> &str {
581 match self {
582 RuntimeCommand::RegisterRoute { command_id, .. }
583 | RuntimeCommand::StartRoute { command_id, .. }
584 | RuntimeCommand::StopRoute { command_id, .. }
585 | RuntimeCommand::SuspendRoute { command_id, .. }
586 | RuntimeCommand::ResumeRoute { command_id, .. }
587 | RuntimeCommand::ReloadRoute { command_id, .. }
588 | RuntimeCommand::FailRoute { command_id, .. }
589 | RuntimeCommand::RemoveRoute { command_id, .. }
590 | RuntimeCommand::ReloadTlsCerts { command_id, .. }
591 | RuntimeCommand::ReloadTemplates { command_id, .. } => command_id,
592 }
593 }
594
595 pub fn causation_id(&self) -> Option<&str> {
596 match self {
597 RuntimeCommand::RegisterRoute { causation_id, .. }
598 | RuntimeCommand::StartRoute { causation_id, .. }
599 | RuntimeCommand::StopRoute { causation_id, .. }
600 | RuntimeCommand::SuspendRoute { causation_id, .. }
601 | RuntimeCommand::ResumeRoute { causation_id, .. }
602 | RuntimeCommand::ReloadRoute { causation_id, .. }
603 | RuntimeCommand::FailRoute { causation_id, .. }
604 | RuntimeCommand::RemoveRoute { causation_id, .. }
605 | RuntimeCommand::ReloadTlsCerts { causation_id, .. }
606 | RuntimeCommand::ReloadTemplates { causation_id, .. } => causation_id.as_deref(),
607 }
608 }
609}
610
611#[derive(Debug, Clone, PartialEq, Eq)]
612#[non_exhaustive]
613pub enum RuntimeCommandResult {
614 Accepted,
615 Duplicate {
616 command_id: String,
617 },
618 RouteRegistered {
619 route_id: String,
620 },
621 RouteStateChanged {
622 route_id: String,
623 status: String,
624 },
625 TlsCertsReloaded {
626 scheme: String,
627 host: String,
628 port: u16,
629 },
630 TemplatesReloaded {
631 route_id: String,
632 },
633}
634
635#[derive(Debug, Clone, PartialEq, Eq)]
636#[non_exhaustive]
637pub enum RuntimeQuery {
638 GetRouteStatus {
639 route_id: String,
640 },
641 InFlightCount {
645 route_id: String,
646 },
647 ListRoutes,
648}
649
650#[derive(Debug, Clone, PartialEq, Eq)]
651#[non_exhaustive]
652pub enum RuntimeQueryResult {
653 InFlightCount { route_id: String, count: u64 },
654 RouteNotFound { route_id: String },
655 RouteStatus { route_id: String, status: String },
656 Routes { route_ids: Vec<String> },
657}
658
659#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
660#[non_exhaustive]
661pub enum RuntimeEvent {
662 RouteRegistered { route_id: String },
663 RouteStartRequested { route_id: String },
664 RouteStarted { route_id: String },
665 RouteFailed { route_id: String, error: String },
666 RouteStopped { route_id: String },
667 RouteSuspended { route_id: String },
668 RouteResumed { route_id: String },
669 RouteReloaded { route_id: String },
670 RouteRemoved { route_id: String },
671}
672
673#[async_trait]
674pub trait RuntimeCommandBus: Send + Sync {
675 async fn execute(&self, cmd: RuntimeCommand) -> Result<RuntimeCommandResult, CamelError>;
676}
677
678#[async_trait]
679pub trait RuntimeQueryBus: Send + Sync {
680 async fn ask(&self, query: RuntimeQuery) -> Result<RuntimeQueryResult, CamelError>;
681}
682
683pub trait RuntimeHandle: RuntimeCommandBus + RuntimeQueryBus {}
684
685impl<T> RuntimeHandle for T where T: RuntimeCommandBus + RuntimeQueryBus {}
686
687#[cfg(test)]
688mod tests {
689 use super::*;
690 use async_trait::async_trait;
691 use futures::executor::block_on;
692
693 struct NoopRuntime;
694
695 #[async_trait]
696 impl RuntimeCommandBus for NoopRuntime {
697 async fn execute(&self, cmd: RuntimeCommand) -> Result<RuntimeCommandResult, CamelError> {
698 Ok(match cmd {
699 RuntimeCommand::RegisterRoute { spec, .. } => {
700 RuntimeCommandResult::RouteRegistered {
701 route_id: spec.route_id,
702 }
703 }
704 RuntimeCommand::StartRoute { route_id, .. }
705 | RuntimeCommand::StopRoute { route_id, .. }
706 | RuntimeCommand::SuspendRoute { route_id, .. }
707 | RuntimeCommand::ResumeRoute { route_id, .. }
708 | RuntimeCommand::ReloadRoute { route_id, .. }
709 | RuntimeCommand::FailRoute { route_id, .. }
710 | RuntimeCommand::RemoveRoute { route_id, .. } => {
711 RuntimeCommandResult::RouteStateChanged {
712 route_id,
713 status: "ok".to_string(),
714 }
715 }
716 RuntimeCommand::ReloadTlsCerts {
717 scheme, host, port, ..
718 } => RuntimeCommandResult::TlsCertsReloaded { scheme, host, port },
719 RuntimeCommand::ReloadTemplates { route_id, .. } => {
720 RuntimeCommandResult::TemplatesReloaded { route_id }
721 }
722 })
723 }
724 }
725
726 #[async_trait]
727 impl RuntimeQueryBus for NoopRuntime {
728 async fn ask(&self, query: RuntimeQuery) -> Result<RuntimeQueryResult, CamelError> {
729 Ok(match query {
730 RuntimeQuery::GetRouteStatus { route_id } => RuntimeQueryResult::RouteStatus {
731 route_id,
732 status: "Started".to_string(),
733 },
734 RuntimeQuery::InFlightCount { route_id } => {
735 RuntimeQueryResult::InFlightCount { route_id, count: 0 }
736 }
737 RuntimeQuery::ListRoutes => RuntimeQueryResult::Routes {
738 route_ids: vec!["r1".to_string()],
739 },
740 })
741 }
742 }
743
744 #[test]
745 fn command_and_query_ids_are_exposed() {
746 let cmd = RuntimeCommand::StartRoute {
747 route_id: "r1".into(),
748 command_id: "c1".into(),
749 causation_id: None,
750 };
751 assert_eq!(cmd.command_id(), "c1");
752 }
753
754 #[test]
755 fn canonical_spec_requires_route_id_and_from() {
756 let spec = CanonicalRouteSpec::new("r1", "timer:tick");
757 assert_eq!(spec.route_id, "r1");
758 assert_eq!(spec.from, "timer:tick");
759 assert_eq!(spec.version, CANONICAL_CONTRACT_VERSION);
760 assert!(spec.steps.is_empty());
761 assert!(spec.circuit_breaker.is_none());
762 }
763
764 #[test]
765 fn canonical_contract_rejects_invalid_version() {
766 let mut spec = CanonicalRouteSpec::new("r1", "timer:tick");
767 spec.version = 3;
768 let err = spec.validate_contract().unwrap_err().to_string();
769 assert!(err.contains("expected version"));
770 }
771
772 #[test]
773 fn canonical_contract_declares_subset_scope() {
774 assert!(canonical_contract_supports_step("to"));
775 assert!(canonical_contract_supports_step("split"));
776 assert!(!canonical_contract_supports_step("set_header"));
777 assert!(!canonical_contract_supports_step("set_property"));
778
779 assert!(CANONICAL_CONTRACT_DECLARATIVE_ONLY_STEPS.contains(&"split"));
780 assert!(CANONICAL_CONTRACT_EXCLUDED_DECLARATIVE_STEPS.contains(&"set_header"));
781 assert!(CANONICAL_CONTRACT_EXCLUDED_DECLARATIVE_STEPS.contains(&"set_property"));
782 assert!(CANONICAL_CONTRACT_RUST_ONLY_STEPS.contains(&"processor"));
783 }
784
785 #[test]
786 fn canonical_contract_rejection_reason_is_explicit() {
787 let set_header_reason = canonical_contract_rejection_reason("set_header")
788 .expect("set_header should have explicit reason");
789 assert!(set_header_reason.contains("out-of-scope"));
790
791 let set_property_reason = canonical_contract_rejection_reason("set_property")
792 .expect("set_property should have explicit reason");
793 assert!(set_property_reason.contains("out-of-scope"));
794
795 let processor_reason = canonical_contract_rejection_reason("processor")
796 .expect("processor should be rust-only");
797 assert!(processor_reason.contains("rust-only"));
798
799 let split_reason = canonical_contract_rejection_reason("split")
800 .expect("split should require declarative form");
801 assert!(split_reason.contains("declarative"));
802 }
803
804 #[test]
805 fn command_causation_id_is_exposed() {
806 let cmd = RuntimeCommand::StopRoute {
807 route_id: "r1".into(),
808 command_id: "c2".into(),
809 causation_id: Some("c1".into()),
810 };
811 assert_eq!(cmd.command_id(), "c2");
812 assert_eq!(cmd.causation_id(), Some("c1"));
813 }
814
815 #[test]
816 fn canonical_contract_rejects_empty_route_id_and_from() {
817 let spec = CanonicalRouteSpec::new(" ", "timer:tick");
818 let err = spec.validate_contract().unwrap_err().to_string();
819 assert!(err.contains("route_id cannot be empty"));
820
821 let spec = CanonicalRouteSpec::new("r1", " ");
822 let err = spec.validate_contract().unwrap_err().to_string();
823 assert!(err.contains("from cannot be empty"));
824 }
825
826 #[test]
827 fn canonical_contract_rejects_invalid_nested_steps() {
828 let mut spec = CanonicalRouteSpec::new("r1", "timer:tick");
829 spec.steps = vec![CanonicalStepSpec::Split {
830 expression: CanonicalSplitExpressionSpec::BodyLines,
831 aggregation: CanonicalSplitAggregationSpec::CollectAll,
832 parallel: true,
833 parallel_limit: Some(0),
834 trace_item_threshold: None,
835 stop_on_exception: false,
836 steps: vec![CanonicalStepSpec::To {
837 uri: "log:ok".to_string(),
838 }],
839 }];
840 let err = spec.validate_contract().unwrap_err().to_string();
841 assert!(err.contains("split.parallel_limit must be > 0"));
842
843 spec.steps = vec![CanonicalStepSpec::To {
844 uri: " ".to_string(),
845 }];
846 let err = spec.validate_contract().unwrap_err().to_string();
847 assert!(err.contains("endpoint uri cannot be empty"));
848 }
849
850 #[test]
851 fn canonical_contract_rejects_invalid_aggregate_and_circuit_breaker() {
852 let mut spec = CanonicalRouteSpec::new("r1", "timer:tick");
853 spec.steps = vec![CanonicalStepSpec::Aggregate(CanonicalAggregateSpec {
854 header: " ".to_string(),
855 completion_size: Some(1),
856 completion_timeout_ms: None,
857 correlation_key: None,
858 force_completion_on_stop: None,
859 discard_on_timeout: None,
860 strategy: CanonicalAggregateStrategySpec::CollectAll,
861 max_buckets: None,
862 max_bucket_size: None,
863 bucket_ttl_ms: None,
864 completion_predicate: None,
865 })];
866 let err = spec.validate_contract().unwrap_err().to_string();
867 assert!(err.contains("correlation source"));
868
869 spec.steps = vec![CanonicalStepSpec::Aggregate(CanonicalAggregateSpec {
870 header: "k".to_string(),
871 completion_size: Some(0),
872 completion_timeout_ms: None,
873 correlation_key: None,
874 force_completion_on_stop: None,
875 discard_on_timeout: None,
876 strategy: CanonicalAggregateStrategySpec::CollectAll,
877 max_buckets: None,
878 max_bucket_size: None,
879 bucket_ttl_ms: None,
880 completion_predicate: None,
881 })];
882 let err = spec.validate_contract().unwrap_err().to_string();
883 assert!(err.contains("aggregate.completion_size must be > 0"));
884
885 spec.steps = vec![];
886 spec.circuit_breaker = Some(CanonicalCircuitBreakerSpec {
887 failure_threshold: 0,
888 open_duration_ms: 10,
889 fallback: vec![],
890 });
891 let err = spec.validate_contract().unwrap_err().to_string();
892 assert!(err.contains("failure_threshold must be > 0"));
893
894 spec.circuit_breaker = Some(CanonicalCircuitBreakerSpec {
895 failure_threshold: 1,
896 open_duration_ms: 0,
897 fallback: vec![],
898 });
899 let err = spec.validate_contract().unwrap_err().to_string();
900 assert!(err.contains("open_duration_ms must be > 0"));
901 }
902
903 #[test]
904 fn contract_accepts_expression_only_aggregate() {
905 let mut spec = CanonicalRouteSpec::new("r1", "timer:tick");
906 spec.steps = vec![CanonicalStepSpec::Aggregate(CanonicalAggregateSpec {
907 header: String::new(),
908 completion_size: None,
909 completion_timeout_ms: None,
910 correlation_key: Some("${header.orderId}".to_string()),
911 force_completion_on_stop: None,
912 discard_on_timeout: None,
913 strategy: CanonicalAggregateStrategySpec::CollectAll,
914 max_buckets: None,
915 max_bucket_size: None,
916 bucket_ttl_ms: None,
917 completion_predicate: None,
918 })];
919 assert!(spec.validate_contract().is_ok());
920 }
921
922 #[test]
923 fn contract_rejects_aggregate_missing_both_sources() {
924 let mut spec = CanonicalRouteSpec::new("r1", "timer:tick");
925 spec.steps = vec![CanonicalStepSpec::Aggregate(CanonicalAggregateSpec {
926 header: String::new(),
927 completion_size: None,
928 completion_timeout_ms: None,
929 correlation_key: None,
930 force_completion_on_stop: None,
931 discard_on_timeout: None,
932 strategy: CanonicalAggregateStrategySpec::CollectAll,
933 max_buckets: None,
934 max_bucket_size: None,
935 bucket_ttl_ms: None,
936 completion_predicate: None,
937 })];
938 let err = spec.validate_contract().unwrap_err().to_string();
939 assert!(err.contains("correlation source"));
940 }
941
942 #[test]
943 fn contract_rejects_empty_correlation_key() {
944 let mut spec = CanonicalRouteSpec::new("r1", "timer:tick");
945 spec.steps = vec![CanonicalStepSpec::Aggregate(CanonicalAggregateSpec {
946 header: "region".to_string(),
947 completion_size: None,
948 completion_timeout_ms: None,
949 correlation_key: Some(String::new()),
950 force_completion_on_stop: None,
951 discard_on_timeout: None,
952 strategy: CanonicalAggregateStrategySpec::CollectAll,
953 max_buckets: None,
954 max_bucket_size: None,
955 bucket_ttl_ms: None,
956 completion_predicate: None,
957 })];
958 let err = spec.validate_contract().unwrap_err().to_string();
959 assert!(err.contains("correlation_key cannot be empty"));
960 }
961
962 #[test]
963 fn canonical_contract_rejects_invalid_fallback_step() {
964 let mut spec = CanonicalRouteSpec::new("r1", "timer:tick");
965 spec.circuit_breaker = Some(CanonicalCircuitBreakerSpec {
966 failure_threshold: 1,
967 open_duration_ms: 10,
968 fallback: vec![CanonicalStepSpec::To {
969 uri: " ".to_string(),
970 }],
971 });
972 let err = spec.validate_contract().unwrap_err().to_string();
973 assert!(err.contains("endpoint uri cannot be empty"));
974 }
975
976 #[test]
977 fn canonical_contract_rejection_reason_none_for_regular_steps() {
978 assert!(canonical_contract_rejection_reason("to").is_none());
979 assert!(canonical_contract_rejection_reason("unknown-step").is_none());
980 }
981
982 #[test]
983 fn command_helpers_cover_all_variants() {
984 let spec = CanonicalRouteSpec::new("r1", "timer:tick");
985 let cmds = [
986 RuntimeCommand::RegisterRoute {
987 spec,
988 command_id: "c1".into(),
989 causation_id: Some("root".into()),
990 },
991 RuntimeCommand::StartRoute {
992 route_id: "r1".into(),
993 command_id: "c2".into(),
994 causation_id: None,
995 },
996 RuntimeCommand::StopRoute {
997 route_id: "r1".into(),
998 command_id: "c3".into(),
999 causation_id: None,
1000 },
1001 RuntimeCommand::SuspendRoute {
1002 route_id: "r1".into(),
1003 command_id: "c4".into(),
1004 causation_id: None,
1005 },
1006 RuntimeCommand::ResumeRoute {
1007 route_id: "r1".into(),
1008 command_id: "c5".into(),
1009 causation_id: None,
1010 },
1011 RuntimeCommand::ReloadRoute {
1012 route_id: "r1".into(),
1013 command_id: "c6".into(),
1014 causation_id: None,
1015 },
1016 RuntimeCommand::FailRoute {
1017 route_id: "r1".into(),
1018 error: "boom".into(),
1019 command_id: "c7".into(),
1020 causation_id: None,
1021 },
1022 RuntimeCommand::RemoveRoute {
1023 route_id: "r1".into(),
1024 command_id: "c8".into(),
1025 causation_id: None,
1026 },
1027 RuntimeCommand::ReloadTlsCerts {
1028 scheme: "https".into(),
1029 host: "example.com".into(),
1030 port: 8443,
1031 command_id: "c9".into(),
1032 causation_id: None,
1033 },
1034 ];
1035
1036 let ids: Vec<&str> = cmds.iter().map(RuntimeCommand::command_id).collect();
1037 assert_eq!(
1038 ids,
1039 vec!["c1", "c2", "c3", "c4", "c5", "c6", "c7", "c8", "c9"]
1040 );
1041 assert_eq!(cmds[0].causation_id(), Some("root"));
1042 assert_eq!(cmds[1].causation_id(), None);
1043 }
1044
1045 #[test]
1046 fn canonical_route_spec_serde_roundtrip() {
1047 let mut spec = CanonicalRouteSpec::new("test-route", "timer:tick?period=1000");
1048 spec.steps.push(CanonicalStepSpec::Log {
1049 message: "Hello".into(),
1050 });
1051 spec.steps.push(CanonicalStepSpec::To {
1052 uri: "log:info".into(),
1053 });
1054 spec.steps.push(CanonicalStepSpec::Stop);
1055
1056 let json = serde_json::to_string(&spec).unwrap();
1057 let deserialized: CanonicalRouteSpec = serde_json::from_str(&json).unwrap();
1058 assert_eq!(spec, deserialized);
1059 }
1060
1061 #[test]
1062 fn canonical_step_spec_serde_variants() {
1063 let steps = vec![
1064 CanonicalStepSpec::To {
1065 uri: "direct:a".into(),
1066 },
1067 CanonicalStepSpec::Log {
1068 message: "msg".into(),
1069 },
1070 CanonicalStepSpec::WireTap {
1071 uri: "direct:audit".into(),
1072 },
1073 CanonicalStepSpec::Stop,
1074 CanonicalStepSpec::Delay {
1075 delay_ms: 100,
1076 dynamic_header: None,
1077 },
1078 ];
1079 let json = serde_json::to_string_pretty(&steps).unwrap();
1080 let back: Vec<CanonicalStepSpec> = serde_json::from_str(&json).unwrap();
1081 assert_eq!(steps, back);
1082 }
1083
1084 #[test]
1085 fn canonical_cache_peek_stale_on_miss_round_trip() {
1086 let some = CanonicalStepSpec::CachePeekStale {
1087 repository: None,
1088 key: "k".into(),
1089 on_miss: Some("continue".into()),
1090 };
1091 let json = serde_json::to_string(&some).unwrap();
1092 assert!(json.contains("continue"));
1093 let back: CanonicalStepSpec = serde_json::from_str(&json).unwrap();
1094 assert_eq!(some, back);
1095
1096 let none = CanonicalStepSpec::CachePeekStale {
1097 repository: None,
1098 key: "k".into(),
1099 on_miss: None,
1100 };
1101 let json = serde_json::to_string(&none).unwrap();
1102 assert!(
1103 !json.contains("on_miss"),
1104 "on_miss: None must omit the key from JSON, got: {json}"
1105 );
1106 let back: CanonicalStepSpec = serde_json::from_str(&json).unwrap();
1107 assert_eq!(none, back);
1108 }
1109
1110 #[test]
1111 fn canonical_cache_invalidate_prefix_round_trip() {
1112 let prefixed = CanonicalStepSpec::CacheInvalidate {
1113 repository: Some("persistent".into()),
1114 key: None,
1115 key_prefix: Some("ns:".into()),
1116 };
1117 let json = serde_json::to_string(&prefixed).unwrap();
1118 assert!(json.contains("key_prefix"), "must emit key_prefix: {json}");
1119 let back: CanonicalStepSpec = serde_json::from_str(&json).unwrap();
1120 assert_eq!(prefixed, back);
1121
1122 let legacy = r#"{"step":"cache_invalidate","config":{"repository":null,"key":"k"}}"#;
1125 let parsed: CanonicalStepSpec = serde_json::from_str(legacy).unwrap();
1126 assert_eq!(
1127 parsed,
1128 CanonicalStepSpec::CacheInvalidate {
1129 repository: None,
1130 key: Some("k".into()),
1131 key_prefix: None,
1132 }
1133 );
1134 }
1135
1136 #[test]
1137 fn canonical_cache_clear_stats_round_trip() {
1138 let clear = CanonicalStepSpec::CacheClear {
1139 repository: Some("persistent".into()),
1140 };
1141 let json = serde_json::to_string(&clear).unwrap();
1142 assert!(json.contains("cache_clear"));
1143 let back: CanonicalStepSpec = serde_json::from_str(&json).unwrap();
1144 assert_eq!(clear, back);
1145
1146 let stats = CanonicalStepSpec::CacheStats { repository: None };
1147 let json = serde_json::to_string(&stats).unwrap();
1148 assert!(json.contains("cache_stats"));
1149 let back: CanonicalStepSpec = serde_json::from_str(&json).unwrap();
1150 assert_eq!(stats, back);
1151 }
1152
1153 #[test]
1154 fn canonical_circuit_breaker_fallback_roundtrip() {
1155 let mut spec = CanonicalRouteSpec::new("cb-fallback", "direct:start");
1156 spec.circuit_breaker = Some(CanonicalCircuitBreakerSpec {
1157 failure_threshold: 1,
1158 open_duration_ms: 60000,
1159 fallback: vec![CanonicalStepSpec::CachePeekStale {
1160 repository: Some("persistent".into()),
1161 key: "tile-xyz".into(),
1162 on_miss: None,
1163 }],
1164 });
1165
1166 let json = serde_json::to_string(&spec).unwrap();
1167 let back: CanonicalRouteSpec = serde_json::from_str(&json).unwrap();
1168 assert_eq!(spec, back);
1169 assert_eq!(back.circuit_breaker.as_ref().unwrap().fallback.len(), 1);
1170
1171 let without = CanonicalCircuitBreakerSpec {
1173 failure_threshold: 1,
1174 open_duration_ms: 60000,
1175 fallback: vec![],
1176 };
1177 let json = serde_json::to_string(&without).unwrap();
1178 assert!(
1179 !json.contains("fallback"),
1180 "empty fallback must omit the key from JSON, got: {json}"
1181 );
1182 let back: CanonicalCircuitBreakerSpec = serde_json::from_str(&json).unwrap();
1183 assert!(back.fallback.is_empty());
1184 }
1185
1186 #[test]
1187 fn canonical_route_spec_json_schema_generates() {
1188 let schema = schemars::schema_for!(CanonicalRouteSpec);
1189 let json = serde_json::to_string(&schema).unwrap();
1190 assert!(json.contains("CanonicalRouteSpec"));
1191 assert!(json.contains("route_id"));
1192 }
1193
1194 #[test]
1195 fn canonical_json_schema_has_no_function_step() {
1196 let schema = schemars::schema_for!(CanonicalRouteSpec);
1197 let json = serde_json::to_string(&schema).unwrap();
1198 assert!(
1199 !json.contains("\"function\""),
1200 "canonical JSON schema must not contain 'function' step"
1201 );
1202 }
1203
1204 #[test]
1205 fn canonical_contract_does_not_support_function() {
1206 assert!(
1207 !canonical_contract_supports_step("function"),
1208 "function must not be in CANONICAL_CONTRACT_SUPPORTED_STEPS"
1209 );
1210 }
1211
1212 #[test]
1213 fn runtime_command_result_all_variants_are_distinct() {
1214 let accepted = RuntimeCommandResult::Accepted;
1215 let dup = RuntimeCommandResult::Duplicate {
1216 command_id: "c1".into(),
1217 };
1218 let registered = RuntimeCommandResult::RouteRegistered {
1219 route_id: "r1".into(),
1220 };
1221 let changed = RuntimeCommandResult::RouteStateChanged {
1222 route_id: "r1".into(),
1223 status: "Started".into(),
1224 };
1225
1226 assert_ne!(accepted, dup);
1227 assert_ne!(dup, registered);
1228 assert_ne!(registered, changed);
1229
1230 let dup2 = RuntimeCommandResult::Duplicate {
1231 command_id: "c1".into(),
1232 };
1233 assert_eq!(dup, dup2);
1234 }
1235
1236 #[test]
1237 fn runtime_event_serialization_round_trip() {
1238 let event = RuntimeEvent::RouteFailed {
1239 route_id: "route-a".to_string(),
1240 error: "boom".to_string(),
1241 };
1242 let json = serde_json::to_string(&event).unwrap();
1243 let back: RuntimeEvent = serde_json::from_str(&json).unwrap();
1244 assert_eq!(event, back);
1245 }
1246
1247 #[test]
1248 fn noop_runtime_execute_and_ask_return_expected_shapes() {
1249 let rt = NoopRuntime;
1250 let cmd = RuntimeCommand::RegisterRoute {
1251 spec: CanonicalRouteSpec::new("r2", "timer:tick"),
1252 command_id: "c1".into(),
1253 causation_id: None,
1254 };
1255 let cmd_result = block_on(rt.execute(cmd)).unwrap();
1256 assert_eq!(
1257 cmd_result,
1258 RuntimeCommandResult::RouteRegistered {
1259 route_id: "r2".into()
1260 }
1261 );
1262
1263 let query_result = block_on(rt.ask(RuntimeQuery::GetRouteStatus {
1264 route_id: "r2".into(),
1265 }))
1266 .unwrap();
1267 assert_eq!(
1268 query_result,
1269 RuntimeQueryResult::RouteStatus {
1270 route_id: "r2".into(),
1271 status: "Started".into()
1272 }
1273 );
1274 }
1275
1276 #[test]
1277 fn canonical_contract_name_and_version_constants_match() {
1278 assert_eq!(CANONICAL_CONTRACT_NAME, "canonical-v1");
1279 assert_eq!(CANONICAL_CONTRACT_VERSION, 2);
1280 }
1281
1282 #[test]
1283 fn canonical_concurrency_spec_rejects_zero_max() {
1284 let spec = CanonicalRouteSpec::new("r1", "timer:tick")
1285 .with_concurrency(CanonicalConcurrencySpec::Concurrent { max: 0 });
1286 let err = spec.validate_contract().unwrap_err().to_string();
1287 assert!(err.contains("concurrency max must be > 0"), "{err}");
1288 }
1289
1290 #[test]
1291 fn canonical_v2_round_trip() {
1292 let spec = CanonicalRouteSpec::new("r1", "timer:tick")
1293 .with_auto_startup(false)
1294 .with_startup_order(42)
1295 .with_concurrency(CanonicalConcurrencySpec::Concurrent { max: 8 });
1296 spec.validate_contract().unwrap();
1297 }
1298
1299 #[test]
1300 fn canonical_v2_version_is_2() {
1301 assert_eq!(CANONICAL_CONTRACT_VERSION, 2);
1302 }
1303
1304 #[test]
1305 fn canonical_loss_report_builder() {
1306 let report =
1307 CanonicalLossReport::from_field("error_handler", "not supported by canonical path", 2);
1308 assert_eq!(report.dropped_fields.len(), 1);
1309 assert_eq!(report.dropped_fields[0].field, "error_handler");
1310 }
1311
1312 #[test]
1313 fn canonical_v1_json_deserializes_in_v2() {
1314 let json = r#"{"route_id":"r1","from":"timer:tick","steps":[],"version":1}"#;
1315 let spec: CanonicalRouteSpec = serde_json::from_str(json).unwrap();
1316 assert_eq!(spec.route_id, "r1");
1317 assert!(spec.auto_startup.is_none());
1318 assert!(spec.startup_order.is_none());
1319 assert!(spec.concurrency.is_none());
1320 spec.validate_contract().unwrap();
1322 }
1323}