1use std::path::PathBuf;
2use std::sync::Arc;
3
4use camel_api::component_metadata::ComponentMetadata;
5use camel_component_api::{
6 BoxProcessor, CamelError, Component, ComponentContext, Consumer, Endpoint, ProducerContext,
7 RuntimeObservability,
8};
9
10use crate::config::{
11 GrpcConfig, GrpcServerConfig, ServerTransport, TransportIntent, parse_grpc_uri,
12};
13use crate::consumer::{GrpcConsumer, resolve_grpc_mode};
14use crate::health::GrpcHealthCheck;
15use crate::metadata::GrpcMetadataDescriptor;
16use crate::producer::GrpcProducer;
17
18pub struct GrpcComponent;
19
20impl GrpcComponent {
21 pub fn new() -> Self {
22 Self
23 }
24}
25
26impl Default for GrpcComponent {
27 fn default() -> Self {
28 Self::new()
29 }
30}
31
32impl Component for GrpcComponent {
33 fn scheme(&self) -> &str {
34 "grpc"
35 }
36
37 fn metadata(&self) -> ComponentMetadata {
38 GrpcMetadataDescriptor::metadata()
39 }
40
41 fn create_endpoint(
42 &self,
43 uri: &str,
44 ctx: &dyn ComponentContext,
45 ) -> Result<Box<dyn Endpoint>, CamelError> {
46 let (host, port, service_name, method_name, cfg) = parse_grpc_uri(uri)?;
47 let proto_file = cfg.proto_file.clone().ok_or_else(|| {
48 CamelError::EndpointCreationFailed(
49 "missing required query parameter: protoFile".to_string(),
50 )
51 })?;
52 let addr = format!("http://{host}:{port}");
53
54 let health_check = GrpcHealthCheck::new(host.clone(), port);
55 ctx.register_current_route_health_check(Arc::new(health_check));
56
57 Ok(Box::new(GrpcEndpoint {
58 uri: uri.to_string(),
59 addr,
60 host,
61 port,
62 proto_path: PathBuf::from(proto_file),
63 service_name,
64 method_name,
65 deadline_ms: cfg.deadline_ms,
66 config: cfg,
67 }))
68 }
69}
70
71struct GrpcEndpoint {
72 uri: String,
73 addr: String,
74 host: String,
75 port: u16,
76 proto_path: PathBuf,
77 service_name: String,
78 method_name: String,
79 deadline_ms: Option<u64>,
80 config: GrpcConfig,
81}
82
83impl Endpoint for GrpcEndpoint {
84 fn uri(&self) -> &str {
85 &self.uri
86 }
87
88 fn create_consumer(
89 &self,
90 rt: Arc<dyn RuntimeObservability>,
91 ) -> Result<Box<dyn Consumer>, CamelError> {
92 if self.config.transport_intent == TransportIntent::Tls
96 && !matches!(self.config.server_transport, ServerTransport::Tls(_))
97 {
98 return Err(CamelError::EndpointCreationFailed(
99 "gRPC from grpc://...?transport=tls requires serverCertPath + serverKeyPath \
100 (inbound TLS termination). Refusing to serve plaintext on a TLS-intent endpoint."
101 .to_string(),
102 ));
103 }
104
105 let path = format!("/{}/{}", self.service_name, self.method_name);
106 let mode = resolve_grpc_mode(&self.proto_path, &self.service_name, &self.method_name)?;
107
108 let server_config = GrpcServerConfig {
109 max_receive_message_len: Some(self.config.max_receive_message_length),
110 transport: self.config.server_transport.clone(),
111 };
112
113 Ok(Box::new(GrpcConsumer::new(
114 self.host.clone(),
115 self.port,
116 path,
117 self.proto_path.clone(),
118 self.service_name.clone(),
119 self.method_name.clone(),
120 mode,
121 rt,
122 server_config,
123 )))
124 }
125
126 fn create_producer(
127 &self,
128 rt: Arc<dyn RuntimeObservability>,
129 ctx: &ProducerContext,
130 ) -> Result<BoxProcessor, CamelError> {
131 let mode = resolve_grpc_mode(&self.proto_path, &self.service_name, &self.method_name)?;
132 let route_id = ctx.route_id().unwrap_or("unknown");
138 let producer = GrpcProducer::new(
139 self.addr.clone(),
140 self.proto_path.clone(),
141 self.service_name.clone(),
142 self.method_name.clone(),
143 mode,
144 self.deadline_ms,
145 &self.config,
146 rt,
147 route_id,
148 )?;
149 Ok(BoxProcessor::new(producer))
150 }
151}
152
153#[cfg(test)]
154mod tests {
155 use std::sync::Arc;
156
157 use camel_component_api::{NoOpComponentContext, RuntimeObservability};
158
159 use super::*;
160
161 fn test_runtime() -> Arc<dyn RuntimeObservability> {
162 Arc::new(NoOpComponentContext)
163 }
164
165 #[test]
169 fn test_c2_inbound_tls_without_server_certs_errors() {
170 let proto_path =
173 std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/helloworld.proto");
174
175 let config = GrpcConfig {
176 proto_file: None,
177 service: None,
178 method: None,
179 reflection: false,
180 transport_intent: TransportIntent::Tls,
181 client_transport: crate::config::ClientTransport::Plaintext,
182 server_transport: ServerTransport::Plaintext, max_receive_message_length: 4 * 1024 * 1024,
184 deadline_ms: None,
185 metadata: None,
186 connect_timeout_ms: 10_000,
187 default_deadline_ms: 30_000,
188 auth: crate::config::AuthConfig::None,
189 interceptors: crate::config::InterceptorConfig::default(),
190 consumer_strategy: crate::config::ConsumerStrategy::default(),
191 producer_strategy: crate::config::ProducerStrategy::default(),
192 retry: camel_component_api::NetworkRetryPolicy::default(),
193 };
194
195 let endpoint = GrpcEndpoint {
196 uri: "grpc://127.0.0.1:0/helloworld.Greeter/SayHello?transport=tls".to_string(),
197 addr: "http://127.0.0.1:0".to_string(),
198 host: "127.0.0.1".to_string(),
199 port: 0,
200 proto_path,
201 service_name: "helloworld.Greeter".to_string(),
202 method_name: "SayHello".to_string(),
203 deadline_ms: None,
204 config,
205 };
206
207 let result = endpoint.create_consumer(test_runtime());
210
211 assert!(
212 result.is_err(),
213 "C2: transport=tls without serverCertPath/serverKeyPath must fail at create_consumer"
214 );
215 match result {
216 Err(e) => {
217 let err = e.to_string();
218 assert!(
219 err.contains("serverCertPath") || err.contains("serverKeyPath"),
220 "error must mention the missing server certs: {err}"
221 );
222 }
223 Ok(_) => panic!("expected error, got Ok"),
224 }
225 }
226}