use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use async_trait::async_trait;
use camel_api::security_policy::{
AccessMode, AuthPrincipal, CredentialSource, PRINCIPAL_SUBJECT_KEY, Principal,
RouteSecurityPlan, TransportId,
};
use camel_api::{Body, CamelError, Exchange, Message};
use camel_auth::native_auth::{
NativeCredential, NativeCredentialSecret, NativeCredentialStore, StaticTokenAuthenticator,
};
use camel_auth::{ProviderEntry, ProviderRegistry, RolePolicy, TokenAuthenticator, read_carrier};
use camel_component_api::{
Consumer, ConsumerContext, ExchangeEnvelope, NoOpComponentContext, RuntimeObservability,
SecurityContext,
};
use camel_component_grpc::GrpcMode;
use camel_component_grpc::config::GrpcServerConfig;
use camel_component_grpc::consumer::GrpcConsumer;
use tokio::net::TcpListener;
use tokio::sync::mpsc;
use tokio_stream::StreamExt;
use tokio_util::sync::CancellationToken;
use zeroize::Zeroizing;
mod helloworld {
tonic::include_proto!("helloworld");
}
mod streaming {
tonic::include_proto!("streaming");
}
const PROVIDER_ID: &str = "idp-grpc";
const TOKEN: &str = "test-token-grpc";
fn test_rt() -> Arc<dyn RuntimeObservability> {
Arc::new(NoOpComponentContext)
}
fn fixture_registry() -> ProviderRegistry {
let store = NativeCredentialStore::try_new(vec![NativeCredential {
secret: NativeCredentialSecret::Plaintext {
value: Zeroizing::new(TOKEN.to_string()),
},
principal: camel_api::security_policy::Principal {
subject: "svc-grpc".to_string(),
issuer: "test".to_string(),
audience: vec![],
scopes: vec![],
roles: vec!["grpc-role".to_string()],
claims: serde_json::Value::Null,
},
}])
.expect("credential store");
let registry = ProviderRegistry::new();
registry.register(
PROVIDER_ID,
ProviderEntry {
authenticator: Arc::new(StaticTokenAuthenticator::new(store)),
audience_binding: None,
},
);
registry
}
fn authenticated_plan(sources: Vec<CredentialSource>) -> RouteSecurityPlan {
RouteSecurityPlan {
access_mode: AccessMode::Authenticated,
provider_ref: Some(PROVIDER_ID.to_string()),
transport: TransportId::Grpc,
credential_sources: sources,
audience_binding: None,
}
}
async fn start_consumer(
sec_ctx: Option<SecurityContext>,
) -> (u16, mpsc::Receiver<ExchangeEnvelope>) {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let port = listener.local_addr().expect("addr").port();
let proto_path = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/helloworld.proto");
let mut consumer = GrpcConsumer::new(
"127.0.0.1".to_string(),
port,
"/helloworld.Greeter/SayHello".to_string(),
proto_path,
"helloworld.Greeter".to_string(),
"SayHello".to_string(),
GrpcMode::Unary,
test_rt(),
GrpcServerConfig::default(),
);
if let Some(sec_ctx) = sec_ctx {
consumer.set_security_context(sec_ctx);
}
let (route_tx, route_rx) = mpsc::channel(16);
let cancel_token = CancellationToken::new();
let ctx = ConsumerContext::new(route_tx, cancel_token, "grpc-auth-test-route".to_string());
tokio::spawn(async move {
consumer
.start_with_listener(ctx, listener)
.await
.expect("consumer start");
});
tokio::time::sleep(std::time::Duration::from_millis(150)).await;
(port, route_rx)
}
async fn start_secured_consumer(
plan: RouteSecurityPlan,
providers: Arc<ProviderRegistry>,
) -> (u16, mpsc::Receiver<ExchangeEnvelope>) {
let policy = RolePolicy::new(vec!["grpc-role".to_string()], true);
let sec_ctx = SecurityContext::new(policy)
.with_credential_sources(plan.credential_sources.clone())
.with_plan(plan)
.with_providers(providers);
start_consumer(Some(sec_ctx)).await
}
async fn greeter_client(
port: u16,
) -> helloworld::greeter_client::GreeterClient<tonic::transport::Channel> {
let channel = tonic::transport::Endpoint::from_shared(format!("http://127.0.0.1:{port}"))
.expect("endpoint")
.connect_lazy();
helloworld::greeter_client::GreeterClient::new(channel)
}
async fn start_streaming_consumer(
sec_ctx: Option<SecurityContext>,
method: &str,
mode: GrpcMode,
) -> (u16, mpsc::Receiver<ExchangeEnvelope>) {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let port = listener.local_addr().expect("addr").port();
let proto_path = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/streaming.proto");
let mut consumer = GrpcConsumer::new(
"127.0.0.1".to_string(),
port,
format!("/streaming.StreamService/{method}"),
proto_path,
"streaming.StreamService".to_string(),
method.to_string(),
mode,
test_rt(),
GrpcServerConfig::default(),
);
if let Some(sec_ctx) = sec_ctx {
consumer.set_security_context(sec_ctx);
}
let (route_tx, route_rx) = mpsc::channel(16);
let cancel_token = CancellationToken::new();
let ctx = ConsumerContext::new(route_tx, cancel_token, "grpc-auth-test-route".to_string());
tokio::spawn(async move {
consumer
.start_with_listener(ctx, listener)
.await
.expect("consumer start");
});
tokio::time::sleep(std::time::Duration::from_millis(150)).await;
(port, route_rx)
}
async fn stream_service_client(
port: u16,
) -> streaming::stream_service_client::StreamServiceClient<tonic::transport::Channel> {
let channel = tonic::transport::Endpoint::from_shared(format!("http://127.0.0.1:{port}"))
.expect("endpoint")
.connect_lazy();
streaming::stream_service_client::StreamServiceClient::new(channel)
}
fn ok_reply() -> Exchange {
Exchange::new(Message::new(Body::Json(
serde_json::json!({"message": "ok"}),
)))
}
#[tokio::test]
async fn grpc_denies_without_credentials_under_kernel() {
let providers = Arc::new(fixture_registry());
let plan = authenticated_plan(vec![CredentialSource::AuthorizationHeader]);
let (port, mut route_rx) = start_secured_consumer(plan, providers).await;
let mut client = greeter_client(port).await;
let err = client
.say_hello(helloworld::HelloRequest {
name: "World".to_string(),
})
.await
.expect_err("request without credentials must be denied");
assert_eq!(
err.code(),
tonic::Code::Unauthenticated,
"denial must be UNAUTHENTICATED, got: {err}"
);
assert!(
route_rx.try_recv().is_err(),
"denied request must not produce a route exchange"
);
}
#[tokio::test]
async fn grpc_named_header_credential_authenticates() {
let providers = Arc::new(fixture_registry());
let plan = authenticated_plan(vec![CredentialSource::Header {
name: "x-api-key".to_string(),
}]);
let (port, mut route_rx) = start_secured_consumer(plan, providers).await;
let pipeline = tokio::spawn(async move {
let envelope = route_rx.recv().await.expect("exchange reaches route");
let subject = envelope
.exchange
.property(PRINCIPAL_SUBJECT_KEY)
.and_then(|v| v.as_str().map(str::to_string));
if let Some(tx) = envelope.reply_tx {
let _ = tx.send(Ok(ok_reply()));
}
subject
});
let mut client = greeter_client(port).await;
let mut request = tonic::Request::new(helloworld::HelloRequest {
name: "World".to_string(),
});
request
.metadata_mut()
.insert("x-api-key", TOKEN.parse().expect("token metadata value"));
let response = client
.say_hello(request)
.await
.expect("valid named-header credential must authenticate");
assert_eq!(response.into_inner().message, "ok");
let subject = pipeline.await.expect("pipeline join");
assert_eq!(
subject.as_deref(),
Some("svc-grpc"),
"authenticated principal must reach the route"
);
}
#[tokio::test]
async fn grpc_carrier_present_on_second_request_fresh_exchange() {
let providers = Arc::new(fixture_registry());
let plan = authenticated_plan(vec![CredentialSource::AuthorizationHeader]);
let (port, mut route_rx) = start_secured_consumer(plan, providers).await;
let pipeline = tokio::spawn(async move {
let mut carried = 0usize;
for _ in 0..2 {
let envelope = route_rx.recv().await.expect("exchange reaches route");
if let Some(carrier) = read_carrier(&envelope.exchange) {
assert_eq!(carrier.provider_id(), PROVIDER_ID);
assert_eq!(carrier.principal().subject, "svc-grpc");
carried += 1;
}
if let Some(tx) = envelope.reply_tx {
let _ = tx.send(Ok(ok_reply()));
}
}
carried
});
let mut client = greeter_client(port).await;
for i in 0..2 {
let mut request = tonic::Request::new(helloworld::HelloRequest {
name: format!("World-{i}"),
});
request.metadata_mut().insert(
"authorization",
format!("Bearer {TOKEN}")
.parse()
.expect("bearer metadata value"),
);
let response = client
.say_hello(request)
.await
.expect("request {i} must succeed under the kernel path");
assert_eq!(response.into_inner().message, "ok");
}
let carried = pipeline.await.expect("pipeline join");
assert_eq!(
carried, 2,
"typed carrier must be present on EVERY fresh exchange, not just the first"
);
}
struct CountingProvider {
calls: Arc<AtomicUsize>,
}
#[async_trait]
impl TokenAuthenticator for CountingProvider {
async fn authenticate_bearer(&self, _token: &str) -> Result<Principal, CamelError> {
self.calls.fetch_add(1, Ordering::SeqCst);
Ok(Principal {
subject: "svc-grpc".to_string(),
issuer: "test".to_string(),
audience: vec![],
scopes: vec![],
roles: vec!["grpc-role".to_string()],
claims: serde_json::Value::Null,
})
}
}
#[tokio::test]
async fn grpc_public_no_extraction_counted_provider() {
let calls = Arc::new(AtomicUsize::new(0));
let registry = ProviderRegistry::new();
registry.register(
PROVIDER_ID,
ProviderEntry {
authenticator: Arc::new(CountingProvider {
calls: Arc::clone(&calls),
}),
audience_binding: None,
},
);
let sec_ctx = SecurityContext::from_arc(Arc::new(RolePolicy::new(vec![], true)))
.with_providers(Arc::new(registry));
let (port, mut route_rx) = start_consumer(Some(sec_ctx)).await;
let pipeline = tokio::spawn(async move {
let envelope = route_rx
.recv()
.await
.expect("public pass-through reaches route");
if let Some(tx) = envelope.reply_tx {
let _ = tx.send(Ok(ok_reply()));
}
});
let mut client = greeter_client(port).await;
let mut request = tonic::Request::new(helloworld::HelloRequest {
name: "World".to_string(),
});
request.metadata_mut().insert(
"authorization",
format!("Bearer {TOKEN}")
.parse()
.expect("bearer metadata value"),
);
let response = client
.say_hello(request)
.await
.expect("plan-less route is Public pass-through");
assert_eq!(response.into_inner().message, "ok");
pipeline.await.expect("pipeline join");
assert_eq!(
calls.load(Ordering::SeqCst),
0,
"plan-less route must never consult a provider"
);
}
#[tokio::test]
async fn grpc_kernel_auth_policy_denied_regression() {
let providers = Arc::new(fixture_registry());
let plan = authenticated_plan(vec![CredentialSource::AuthorizationHeader]);
let policy = RolePolicy::new(vec!["other-role".to_string()], true);
let sec_ctx = SecurityContext::new(policy)
.with_credential_sources(plan.credential_sources.clone())
.with_plan(plan)
.with_providers(providers);
let (port, mut route_rx) = start_consumer(Some(sec_ctx)).await;
let pipeline = tokio::spawn(async move {
let envelope = route_rx
.recv()
.await
.expect("kernel-authenticated request must reach the route pipeline");
if let Some(tx) = envelope.reply_tx {
let _ = tx.send(Err(CamelError::Unauthorized("missing role".to_string())));
}
});
let mut client = greeter_client(port).await;
let mut request = tonic::Request::new(helloworld::HelloRequest {
name: "World".to_string(),
});
request.metadata_mut().insert(
"authorization",
format!("Bearer {TOKEN}")
.parse()
.expect("bearer metadata value"),
);
let err = client
.say_hello(request)
.await
.expect_err("pipeline policy denial must deny the call");
assert_eq!(
err.code(),
tonic::Code::PermissionDenied,
"pipeline policy denial must surface as PERMISSION_DENIED, got: {err}"
);
pipeline.await.expect("pipeline join");
}
#[tokio::test]
async fn grpc_server_streaming_pipeline_denial_regression() {
let providers = Arc::new(fixture_registry());
let plan = authenticated_plan(vec![CredentialSource::AuthorizationHeader]);
let policy = RolePolicy::new(vec!["other-role".to_string()], true);
let sec_ctx = SecurityContext::new(policy)
.with_credential_sources(plan.credential_sources.clone())
.with_plan(plan)
.with_providers(providers);
let (port, mut route_rx) =
start_streaming_consumer(Some(sec_ctx), "ServerList", GrpcMode::ServerStreaming).await;
let pipeline = tokio::spawn(async move {
let envelope = route_rx
.recv()
.await
.expect("kernel-authenticated request must reach the route pipeline");
if let Some(tx) = envelope.reply_tx {
let _ = tx.send(Err(CamelError::Unauthorized("missing role".to_string())));
}
});
let mut client = stream_service_client(port).await;
let mut request = tonic::Request::new(streaming::ListRequest { count: 3 });
request.metadata_mut().insert(
"authorization",
format!("Bearer {TOKEN}")
.parse()
.expect("bearer metadata value"),
);
let mut stream = client
.server_list(request)
.await
.expect("stream must open before the denial verdict")
.into_inner();
let err = stream
.next()
.await
.expect("denial must produce a stream frame")
.expect_err("pipeline policy denial must surface on the stream");
assert_eq!(
err.code(),
tonic::Code::PermissionDenied,
"pipeline policy denial must surface as PERMISSION_DENIED, got: {err}"
);
pipeline.await.expect("pipeline join");
}
#[tokio::test]
async fn grpc_client_streaming_pipeline_denial_regression() {
let providers = Arc::new(fixture_registry());
let plan = authenticated_plan(vec![CredentialSource::AuthorizationHeader]);
let policy = RolePolicy::new(vec!["other-role".to_string()], true);
let sec_ctx = SecurityContext::new(policy)
.with_credential_sources(plan.credential_sources.clone())
.with_plan(plan)
.with_providers(providers);
let (port, mut route_rx) =
start_streaming_consumer(Some(sec_ctx), "ClientSum", GrpcMode::ClientStreaming).await;
let pipeline = tokio::spawn(async move {
loop {
let envelope = route_rx
.recv()
.await
.expect("kernel-authenticated request must reach the route pipeline");
if let Some(tx) = envelope.reply_tx {
let _ = tx.send(Err(CamelError::Unauthorized("missing role".to_string())));
}
let complete = matches!(
envelope
.exchange
.input
.header("CamelGrpcClientStreamComplete"),
Some(serde_json::Value::Bool(true))
);
if complete {
break;
}
}
});
let mut client = stream_service_client(port).await;
let mut request = tonic::Request::new(tokio_stream::iter(vec![
streaming::NumberRequest { value: 1 },
streaming::NumberRequest { value: 2 },
]));
request.metadata_mut().insert(
"authorization",
format!("Bearer {TOKEN}")
.parse()
.expect("bearer metadata value"),
);
let err = client
.client_sum(request)
.await
.expect_err("pipeline policy denial must fail the RPC");
assert_eq!(
err.code(),
tonic::Code::PermissionDenied,
"pipeline policy denial must surface as PERMISSION_DENIED, got: {err}"
);
pipeline.await.expect("pipeline join");
}
#[tokio::test]
async fn grpc_bidi_pipeline_denial_regression() {
let providers = Arc::new(fixture_registry());
let plan = authenticated_plan(vec![CredentialSource::AuthorizationHeader]);
let policy = RolePolicy::new(vec!["other-role".to_string()], true);
let sec_ctx = SecurityContext::new(policy)
.with_credential_sources(plan.credential_sources.clone())
.with_plan(plan)
.with_providers(providers);
let (port, mut route_rx) =
start_streaming_consumer(Some(sec_ctx), "BidiEcho", GrpcMode::Bidi).await;
let pipeline = tokio::spawn(async move {
let envelope = route_rx
.recv()
.await
.expect("kernel-authenticated request must reach the route pipeline");
if let Some(tx) = envelope.reply_tx {
let _ = tx.send(Err(CamelError::Unauthorized("missing role".to_string())));
}
});
let request_stream = tokio_stream::iter(vec![streaming::EchoRequest {
message: "hello".to_string(),
}])
.chain(tokio_stream::pending());
let mut request = tonic::Request::new(request_stream);
request.metadata_mut().insert(
"authorization",
format!("Bearer {TOKEN}")
.parse()
.expect("bearer metadata value"),
);
let mut client = stream_service_client(port).await;
let mut stream = client
.bidi_echo(request)
.await
.expect("stream must open before the denial verdict")
.into_inner();
let err = stream
.next()
.await
.expect("denial must produce a stream frame")
.expect_err("pipeline policy denial must surface on the stream");
assert_eq!(
err.code(),
tonic::Code::PermissionDenied,
"pipeline policy denial must surface as PERMISSION_DENIED, got: {err}"
);
pipeline.await.expect("pipeline join");
}
fn collect_rs_files(dir: &Path, out: &mut Vec<PathBuf>) {
for entry in std::fs::read_dir(dir).expect("read source dir") {
let path = entry.expect("dir entry").path();
if path.is_dir() {
collect_rs_files(&path, out);
} else if path.extension().is_some_and(|ext| ext == "rs") {
out.push(path);
}
}
}
#[test]
fn grpc_legacy_arm_deleted_source_scan() {
let needles = [
["Grpc", "Principal"].concat(),
["extract", "_principal"].concat(),
["legacy", "_authenticator"].concat(),
];
let src = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("src");
let mut files = Vec::new();
collect_rs_files(&src, &mut files);
assert!(
!files.is_empty(),
"source scan must find production sources"
);
for file in files {
let content = std::fs::read_to_string(&file).expect("read source file");
for needle in &needles {
assert!(
!content.contains(needle.as_str()),
"legacy marker `{needle}` must stay deleted: {}",
file.display()
);
}
}
}