Skip to main content

camel_component_grpc/
component.rs

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        // C2: direction-aware fail-closed. transport=tls declared but the
87        // server has no certs → refuse (no silent plaintext on a TLS port).
88        // Check this BEFORE proto compilation to fail fast.
89        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        // NOTE: direction-aware outbound intent is already enforced at parse time
127        // (parse_grpc_query_params rejects conflicting transport+cert combos) and
128        // in GrpcProducer::new (ClientTransport::Tls hard-errors on bad config).
129        // No extra validation needed here.
130        // route_id may not be set in test scenarios (e.g., standalone endpoint tests)
131        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    /// Regression: C2 — inbound transport=tls without server certs must error at create_consumer.
160    /// Parse succeeds (defers to create_consumer), but create_consumer must hard-error with
161    /// "serverCertPath" in the message.
162    #[test]
163    fn test_c2_inbound_tls_without_server_certs_errors() {
164        // Construct a GrpcEndpoint directly with transport_intent=Tls but server_transport=Plaintext
165        // This simulates the state after parsing "grpc://...?transport=tls" without server certs
166        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, // No server certs
177            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        // Now call create_consumer — this MUST fail because transport_intent=Tls but
202        // server_transport=Plaintext means no server certs were provided (C2 violation)
203        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}