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 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 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}