Skip to main content

camel_api/
runtime.rs

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    /// Stable minimal route representation for runtime command registration.
86    ///
87    /// Scope notes:
88    /// - This is intentionally a partial model and does not mirror every `BuilderStep`.
89    /// - Version 2 adds: auto_startup, startup_order, concurrency.
90    /// - Still excluded: error_handler, unit_of_work. These are set to defaults
91    ///   when compiling from canonical.
92    /// - Round-trip (YAML → Canonical → YAML) loses these fields.
93    /// - Advanced EIPs continue to use the existing RouteDefinition/BuilderStep path.
94    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        /// Threshold above which split fragments start new traces (0 = off).
144        /// Absent means the runtime default applies.
145        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        /// Coalesce concurrent misses on the same key into a single `on_miss`
162        /// run. Absent/null means `false`.
163        #[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        /// Exact key to invalidate (simple-language expression).
170        #[serde(default, skip_serializing_if = "Option::is_none")]
171        key: Option<String>,
172        /// Namespace prefix to invalidate (simple-language expression).
173        #[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        /// On-miss policy: `"stop"` (default) or `"continue"`. Absent/null means `"stop"`.
186        #[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    /// Per-bucket accumulation bound (audit 2026-08-31, F6-2).
287    /// Absent = builder default (10_000).
288    #[serde(default)]
289    pub max_bucket_size: Option<usize>,
290    #[ts(type = "number")]
291    pub bucket_ttl_ms: Option<u64>,
292    /// Language expression predicate; completes the bucket when it evaluates
293    /// true against the incoming exchange (runtime-resolved via the language
294    /// registry). Mirrors `CorrelationStrategy::Expression`.
295    #[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    /// Internal lifecycle command emitted by runtime adapters when a route crashes at runtime.
546    ///
547    /// This keeps aggregate/projection state aligned with controller-observed failures.
548    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    /// Reload template sources for a route (infrastructure command).
567    ///
568    /// Intercepted in `RuntimeBus::execute` BEFORE journal recovery + dedup,
569    /// exactly like `ReloadTlsCerts`: it is idempotent, NOT journaled, and does
570    /// not mutate `RouteStatus`. Dispatches to
571    /// `TemplateReloadRegistry::reload_route`.
572    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    /// **Note:** This variant is intercepted by `RuntimeBus::ask` *before* reaching
642    /// `execute_query`. Do not handle it in `execute_query` — it has no access to
643    /// the in-flight counter. See `runtime_bus.rs` for the intercept.
644    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        // Legacy wire form carrying only `key` still deserializes with
1123        // `key_prefix: None` (no key-iteration field).
1124        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        // A spec serialized without the fallback key deserializes with empty fallback.
1172        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        // CRITICAL: v1 specs must pass validation in v2 runtime (backward compat)
1321        spec.validate_contract().unwrap();
1322    }
1323}