Skip to main content

camel_master/
endpoint.rs

1use std::sync::Arc;
2use std::time::Duration;
3
4use camel_api::{CamelError, MetricsCollector, PlatformService};
5use camel_component_api::{
6    BoxProcessor, Component, ComponentContext, Consumer, Endpoint, NetworkRetryPolicy,
7    ProducerContext,
8};
9use camel_language_api::Language;
10
11use crate::consumer::MasterConsumer;
12
13pub(crate) struct MasterEndpoint {
14    pub(crate) uri: String,
15    pub(crate) lock_name: String,
16    pub(crate) delegate_uri: String,
17    pub(crate) delegate_component: Arc<dyn Component>,
18    // TODO(MST-001): MetricsCollector is wired through from ComponentContext but never called.
19    // Should emit metrics on leader acquisition (increment_exchanges), leader loss
20    // (increment_errors), and delegate start/stop events (record_circuit_breaker_change).
21    pub(crate) metrics: Arc<dyn MetricsCollector>,
22    pub(crate) platform_service: Arc<dyn PlatformService>,
23    pub(crate) drain_timeout: Duration,
24    pub(crate) reconnect: NetworkRetryPolicy,
25}
26
27impl Endpoint for MasterEndpoint {
28    fn uri(&self) -> &str {
29        &self.uri
30    }
31
32    fn create_consumer(
33        &self,
34        rt: Arc<dyn camel_component_api::RuntimeObservability>,
35    ) -> Result<Box<dyn Consumer>, CamelError> {
36        Ok(Box::new(MasterConsumer::new(
37            self.lock_name.clone(),
38            self.delegate_uri.clone(),
39            Arc::clone(&self.delegate_component),
40            Arc::clone(&self.metrics),
41            Arc::clone(&self.platform_service),
42            self.drain_timeout,
43            self.reconnect.clone(),
44            rt,
45        )))
46    }
47
48    fn create_producer(
49        &self,
50        rt: Arc<dyn camel_component_api::RuntimeObservability>,
51        ctx: &ProducerContext,
52    ) -> Result<BoxProcessor, CamelError> {
53        let delegate_ctx = MasterDelegateContext {
54            delegate_component: Arc::clone(&self.delegate_component),
55            metrics: Arc::clone(&self.metrics),
56            platform_service: Arc::clone(&self.platform_service),
57        };
58
59        self.delegate_component
60            .create_endpoint(&self.delegate_uri, &delegate_ctx)?
61            .create_producer(rt, ctx)
62    }
63}
64
65pub(crate) struct MasterDelegateContext {
66    // Fields are pub(crate) during Tasks 1-2 because reconcile_event (still in
67    // lib.rs) constructs this struct via field-init syntax. After Task 3 (Commit 2
68    // — leadership.rs extraction), they CAN be narrowed to private since only
69    // endpoint.rs constructs them (in create_producer).
70    pub(crate) delegate_component: Arc<dyn Component>,
71    pub(crate) metrics: Arc<dyn MetricsCollector>,
72    pub(crate) platform_service: Arc<dyn PlatformService>,
73}
74
75impl ComponentContext for MasterDelegateContext {
76    fn resolve_component(&self, scheme: &str) -> Option<Arc<dyn Component>> {
77        if self.delegate_component.scheme() == scheme {
78            Some(Arc::clone(&self.delegate_component))
79        } else {
80            None
81        }
82    }
83
84    fn resolve_language(&self, _name: &str) -> Option<Arc<dyn Language>> {
85        None
86    }
87
88    fn metrics(&self) -> Arc<dyn MetricsCollector> {
89        Arc::clone(&self.metrics)
90    }
91
92    fn platform_service(&self) -> Arc<dyn PlatformService> {
93        Arc::clone(&self.platform_service)
94    }
95
96    fn register_route_health_check(
97        &self,
98        _route_id: &str,
99        _check: Arc<dyn camel_api::AsyncHealthCheck>,
100    ) {
101    }
102
103    fn unregister_route_health_check(&self, _route_id: &str) {}
104}