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