Skip to main content

camel_component_grpc/
component.rs

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        // C2: direction-aware fail-closed. transport=tls declared but the
93        // server has no certs → refuse (no silent plaintext on a TLS port).
94        // Check this BEFORE proto compilation to fail fast.
95        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        // NOTE: direction-aware outbound intent is already enforced at parse time
133        // (parse_grpc_query_params rejects conflicting transport+cert combos) and
134        // in GrpcProducer::new (ClientTransport::Tls hard-errors on bad config).
135        // No extra validation needed here.
136        // route_id may not be set in test scenarios (e.g., standalone endpoint tests)
137        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    /// Regression: C2 — inbound transport=tls without server certs must error at create_consumer.
166    /// Parse succeeds (defers to create_consumer), but create_consumer must hard-error with
167    /// "serverCertPath" in the message.
168    #[test]
169    fn test_c2_inbound_tls_without_server_certs_errors() {
170        // Construct a GrpcEndpoint directly with transport_intent=Tls but server_transport=Plaintext
171        // This simulates the state after parsing "grpc://...?transport=tls" without server certs
172        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, // No server certs
183            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        // Now call create_consumer — this MUST fail because transport_intent=Tls but
208        // server_transport=Plaintext means no server certs were provided (C2 violation)
209        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}