1use std::net::IpAddr;
18use std::path::Path;
19
20use tonic::codegen::{Body, Bytes, StdError};
21use tonic::service::interceptor::InterceptedService;
22use tonic::service::Interceptor;
23use tonic::transport::{Certificate, Channel, ClientTlsConfig, Endpoint};
24use tonic::{Request, Status};
25
26use crate::admin::v1::admin_service_client::AdminServiceClient;
27use crate::admin::v1::git_ops_service_client::GitOpsServiceClient;
28
29#[derive(Debug, thiserror::Error)]
31pub enum ClientError {
32 #[error("token must not be empty")]
34 EmptyToken,
35
36 #[error("invalid CA PEM bundle: no certificates found")]
38 InvalidCaPem,
39
40 #[error("force_tls and plaintext are mutually exclusive")]
45 ForceTlsPlaintextConflict,
46
47 #[error("plaintext and ca_pem are mutually exclusive")]
51 PlaintextCaPemConflict,
52
53 #[error("invalid CA PEM bundle: {0}")]
55 CaPemParse(String),
56
57 #[error("invalid address {addr:?}: {source}")]
59 InvalidEndpoint {
60 addr: String,
61 #[source]
62 source: tonic::transport::Error,
63 },
64
65 #[error("connect to {addr:?}: {source}")]
67 Connect {
68 addr: String,
69 #[source]
70 source: tonic::transport::Error,
71 },
72
73 #[error("read CA file {path:?}: {source}")]
75 ReadCaFile {
76 path: String,
77 #[source]
78 source: std::io::Error,
79 },
80
81 #[cfg(feature = "spiffe-workload")]
83 #[error("connect to SPIFFE Workload API at {socket:?}: {source}")]
84 WorkloadApi {
85 socket: String,
86 #[source]
87 source: spiffe::x509_source::X509SourceError,
88 },
89
90 #[cfg(feature = "spiffe-workload")]
96 #[error("connect to SPIFFE workload API at {socket:?}: no identity issued after {attempts} attempts: {source}")]
97 WorkloadNoIdentityIssued {
98 socket: String,
99 attempts: usize,
100 #[source]
101 source: spiffe::WorkloadApiError,
102 },
103
104 #[cfg(feature = "spiffe-workload")]
111 #[error("connect to SPIFFE workload API at {socket:?}: {source}")]
112 WorkloadProbe {
113 socket: String,
114 #[source]
115 source: spiffe::WorkloadApiError,
116 },
117
118 #[cfg(feature = "spiffe-workload")]
120 #[error("invalid trust domain {0:?}: {1}")]
121 InvalidTrustDomain(String, String),
122
123 #[cfg(feature = "spiffe-workload")]
125 #[error("build SPIFFE mTLS client config: {0}")]
126 SpiffeTls(String),
127
128 #[cfg(feature = "spiffe-workload")]
130 #[error("invalid workload address {0:?}: expected host:port")]
131 InvalidWorkloadAddress(String),
132
133 #[cfg(feature = "spiffe-workload")]
136 #[error("connect to {addr:?}: {source}")]
137 WorkloadConnect {
138 addr: String,
139 #[source]
140 source: std::io::Error,
141 },
142}
143
144#[derive(Clone)]
148pub struct TokenInterceptor {
149 header_value: tonic::metadata::MetadataValue<tonic::metadata::Ascii>,
150}
151
152impl Interceptor for TokenInterceptor {
153 fn call(&mut self, mut req: Request<()>) -> Result<Request<()>, Status> {
154 req.metadata_mut()
155 .insert("authorization", self.header_value.clone());
156 Ok(req)
157 }
158}
159
160pub type AdminChannel = InterceptedService<Channel, TokenInterceptor>;
164
165pub async fn dial_admin(
201 addr: impl AsRef<str>,
202 token: impl AsRef<str>,
203 ca_pem: Option<&[u8]>,
204 force_tls: bool,
205 plaintext: bool,
206) -> Result<AdminChannel, ClientError> {
207 let addr = addr.as_ref();
208 let token = token.as_ref().trim();
209 if token.is_empty() {
210 return Err(ClientError::EmptyToken);
211 }
212
213 let decision = admin_transport_decision(addr, ca_pem, force_tls, plaintext)?;
214 let uri = format!(
215 "{}://{addr}",
216 if decision.requires_tls() { "https" } else { "http" }
217 );
218
219 let mut endpoint = Endpoint::from_shared(uri).map_err(|source| ClientError::InvalidEndpoint {
220 addr: addr.to_string(),
221 source,
222 })?;
223 if let TransportDecision::Tls(tls_config) = decision {
224 endpoint = endpoint
225 .tls_config(tls_config)
226 .map_err(|source| ClientError::InvalidEndpoint {
227 addr: addr.to_string(),
228 source,
229 })?;
230 }
231
232 let channel = endpoint
233 .connect()
234 .await
235 .map_err(|source| ClientError::Connect {
236 addr: addr.to_string(),
237 source,
238 })?;
239
240 let header_value = format!("Bearer {token}")
241 .parse()
242 .expect("Bearer <token> is always valid ASCII metadata once token is trimmed non-empty");
243
244 Ok(InterceptedService::new(
245 channel,
246 TokenInterceptor { header_value },
247 ))
248}
249
250pub fn admin_client<T>(channel: T) -> AdminServiceClient<T>
261where
262 T: tonic::client::GrpcService<tonic::body::Body>,
263 T::Error: Into<StdError>,
264 T::ResponseBody: Body<Data = Bytes> + Send + 'static,
265 <T::ResponseBody as Body>::Error: Into<StdError> + Send,
266{
267 AdminServiceClient::new(channel)
268}
269
270pub fn gitops_client<T>(channel: T) -> GitOpsServiceClient<T>
289where
290 T: tonic::client::GrpcService<tonic::body::Body>,
291 T::Error: Into<StdError>,
292 T::ResponseBody: Body<Data = Bytes> + Send + 'static,
293 <T::ResponseBody as Body>::Error: Into<StdError> + Send,
294{
295 GitOpsServiceClient::new(channel)
296}
297
298pub fn read_ca_file(path: impl AsRef<Path>) -> Result<Vec<u8>, ClientError> {
300 let path_ref = path.as_ref();
301 std::fs::read(path_ref).map_err(|source| ClientError::ReadCaFile {
302 path: path_ref.display().to_string(),
303 source,
304 })
305}
306
307#[derive(Debug)]
308pub(crate) enum TransportDecision {
309 Plaintext,
310 Tls(ClientTlsConfig),
311}
312
313impl TransportDecision {
314 pub(crate) fn requires_tls(&self) -> bool {
315 matches!(self, TransportDecision::Tls(_))
316 }
317}
318
319pub(crate) fn admin_transport_decision(
320 addr: &str,
321 ca_pem: Option<&[u8]>,
322 force_tls: bool,
323 plaintext: bool,
324) -> Result<TransportDecision, ClientError> {
325 let ca_pem_non_empty = ca_pem.map(|pem| !pem.is_empty()).unwrap_or(false);
326
327 if plaintext && force_tls {
328 return Err(ClientError::ForceTlsPlaintextConflict);
329 }
330 if plaintext && ca_pem_non_empty {
331 return Err(ClientError::PlaintextCaPemConflict);
332 }
333 if plaintext {
334 return Ok(TransportDecision::Plaintext);
335 }
336
337 let host = host_of(addr);
338 let use_tls = force_tls || ca_pem_non_empty || !is_loopback_host(&host);
339
340 if !use_tls {
341 return Ok(TransportDecision::Plaintext);
342 }
343
344 let mut tls = ClientTlsConfig::new();
345 if let Some(pem) = ca_pem {
346 if !pem.is_empty() {
347 validate_ca_pem(pem)?;
348 tls = tls.ca_certificate(Certificate::from_pem(pem));
349 }
350 }
351 Ok(TransportDecision::Tls(tls))
352}
353
354pub(crate) fn host_of(addr: &str) -> String {
359 if let Ok(sock) = addr.parse::<std::net::SocketAddr>() {
360 return sock.ip().to_string();
361 }
362 if let Some(idx) = addr.rfind(':') {
363 let (host_part, port_part) = (&addr[..idx], &addr[idx + 1..]);
364 if !host_part.is_empty() && !port_part.is_empty() && port_part.bytes().all(|b| b.is_ascii_digit()) {
365 return host_part.trim_start_matches('[').trim_end_matches(']').to_string();
366 }
367 }
368 addr.to_string()
369}
370
371pub(crate) fn is_loopback_host(host: &str) -> bool {
372 if host.eq_ignore_ascii_case("localhost") {
373 return true;
374 }
375 host.parse::<IpAddr>()
376 .map(|ip| ip.is_loopback())
377 .unwrap_or(false)
378}
379
380fn validate_ca_pem(pem: &[u8]) -> Result<(), ClientError> {
381 let mut reader = std::io::BufReader::new(pem);
382 let mut count = 0usize;
383 for item in rustls_pemfile::certs(&mut reader) {
384 match item {
385 Ok(_) => count += 1,
386 Err(e) => return Err(ClientError::CaPemParse(e.to_string())),
387 }
388 }
389 if count == 0 {
390 return Err(ClientError::InvalidCaPem);
391 }
392 Ok(())
393}
394
395#[cfg(feature = "spiffe-workload")]
396mod workload {
397 use super::ClientError;
398 use std::future::Future;
399 use std::net::ToSocketAddrs;
400 use std::pin::Pin;
401 use std::sync::Arc;
402 use std::task::{Context, Poll};
403 use std::time::Duration;
404
405 use tokio::net::TcpStream;
406 use tonic::codegen::http::Uri;
407 use tonic::codegen::Service;
408 use tonic::transport::{Channel, Endpoint};
409
410 const WORKLOAD_DIAL_BACKOFF: [Duration; 4] = [
422 Duration::from_secs(1),
423 Duration::from_secs(2),
424 Duration::from_secs(4),
425 Duration::from_secs(8),
426 ];
427
428 const WORKLOAD_DIAL_MAX_ATTEMPTS: usize = WORKLOAD_DIAL_BACKOFF.len() + 1;
432
433 async fn retry_until_identity_issued<F, Fut>(
443 socket_path: &str,
444 mut probe: F,
445 ) -> Result<(), ClientError>
446 where
447 F: FnMut() -> Fut,
448 Fut: Future<Output = Result<(), spiffe::WorkloadApiError>>,
449 {
450 let mut last_err: Option<spiffe::WorkloadApiError> = None;
451 let delays_before_each_attempt =
458 std::iter::once(None).chain(WORKLOAD_DIAL_BACKOFF.into_iter().map(Some));
459
460 for delay in delays_before_each_attempt {
461 if let Some(delay) = delay {
462 tokio::time::sleep(delay).await;
463 }
464 match probe().await {
465 Ok(()) => return Ok(()),
466 Err(e) => {
467 if !matches!(e, spiffe::WorkloadApiError::NoIdentityIssued) {
468 return Err(ClientError::WorkloadProbe {
469 socket: socket_path.to_string(),
470 source: e,
471 });
472 }
473 last_err = Some(e);
474 }
475 }
476 }
477 Err(ClientError::WorkloadNoIdentityIssued {
478 socket: socket_path.to_string(),
479 attempts: WORKLOAD_DIAL_MAX_ATTEMPTS,
480 source: last_err
481 .expect("loop always records last_err before exhausting WORKLOAD_DIAL_MAX_ATTEMPTS"),
482 })
483 }
484
485 async fn probe_identity_issued(socket_path: &str) -> Result<(), spiffe::WorkloadApiError> {
490 let client = spiffe::WorkloadApiClient::connect_to(socket_path).await?;
491 client.fetch_x509_context().await?;
492 Ok(())
493 }
494
495 pub async fn dial_workload(
559 addr: impl AsRef<str>,
560 socket_path: impl AsRef<str>,
561 trust_domain: impl AsRef<str>,
562 ) -> Result<Channel, ClientError> {
563 let addr = addr.as_ref().to_string();
564 let socket_path = socket_path.as_ref().to_string();
565 let trust_domain = trust_domain.as_ref().to_string();
566
567 retry_until_identity_issued(&socket_path, || probe_identity_issued(&socket_path)).await?;
568
569 let source = spiffe::X509Source::builder()
570 .endpoint(&socket_path)
571 .build()
572 .await
573 .map_err(|source| ClientError::WorkloadApi {
574 socket: socket_path.clone(),
575 source,
576 })?;
577
578 let td = spiffe::TrustDomain::try_from(trust_domain.as_str())
579 .map_err(|e| ClientError::InvalidTrustDomain(trust_domain.clone(), e.to_string()))?;
580
581 let authorizer = spiffe_rustls::authorizer::trust_domains([td.clone()])
582 .map_err(|e| ClientError::SpiffeTls(e.to_string()))?;
583
584 let tls_config = spiffe_rustls::mtls_client(source)
585 .authorize(authorizer)
586 .trust_domain_policy(spiffe_rustls::TrustDomainPolicy::LocalOnly(td))
587 .with_alpn_protocols([b"h2".to_vec()])
588 .build()
589 .map_err(|e| ClientError::SpiffeTls(e.to_string()))?;
590
591 let host = super::host_of(&addr);
596 let server_name = rustls::pki_types::ServerName::try_from(host.clone())
597 .map_err(|_| ClientError::InvalidWorkloadAddress(addr.clone()))?;
598
599 let connector = SpiffeConnector {
600 target_addr: addr.clone(),
601 tls_config: Arc::new(tls_config),
602 server_name,
603 };
604
605 let endpoint =
613 Endpoint::from_shared(format!("http://{addr}")).map_err(|source| {
614 ClientError::InvalidEndpoint {
615 addr: addr.clone(),
616 source,
617 }
618 })?;
619
620 endpoint
621 .connect_with_connector(connector)
622 .await
623 .map_err(|source| ClientError::Connect { addr, source })
624 }
625
626 #[derive(Clone)]
627 struct SpiffeConnector {
628 target_addr: String,
629 tls_config: Arc<rustls::ClientConfig>,
630 server_name: rustls::pki_types::ServerName<'static>,
631 }
632
633 impl Service<Uri> for SpiffeConnector {
634 type Response = hyper_util::rt::TokioIo<tokio_rustls::client::TlsStream<TcpStream>>;
635 type Error = ClientError;
636 type Future = Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send>>;
637
638 fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
639 Poll::Ready(Ok(()))
640 }
641
642 fn call(&mut self, _uri: Uri) -> Self::Future {
643 let target_addr = self.target_addr.clone();
644 let tls_config = self.tls_config.clone();
645 let server_name = self.server_name.clone();
646
647 Box::pin(async move {
648 let socket_addr = target_addr
649 .to_socket_addrs()
650 .map_err(|source| ClientError::WorkloadConnect {
651 addr: target_addr.clone(),
652 source,
653 })?
654 .next()
655 .ok_or_else(|| ClientError::InvalidWorkloadAddress(target_addr.clone()))?;
656
657 let tcp = TcpStream::connect(socket_addr).await.map_err(|source| {
658 ClientError::WorkloadConnect {
659 addr: target_addr.clone(),
660 source,
661 }
662 })?;
663 let _ = tcp.set_nodelay(true);
664
665 let connector = tokio_rustls::TlsConnector::from(tls_config);
666 let tls_stream = connector
667 .connect(server_name, tcp)
668 .await
669 .map_err(|source| ClientError::WorkloadConnect {
670 addr: target_addr,
671 source,
672 })?;
673
674 Ok(hyper_util::rt::TokioIo::new(tls_stream))
675 })
676 }
677 }
678
679 #[cfg(test)]
680 mod tests {
681 use super::*;
682 use std::sync::atomic::{AtomicUsize, Ordering};
683
684 #[allow(dead_code)]
693 fn _dial_workload_channel_satisfies_gitops_and_admin_client_bound(channel: Channel) {
694 let _ = super::super::gitops_client(channel.clone());
695 let _ = super::super::admin_client(channel);
696 }
697
698 fn no_identity_issued() -> spiffe::WorkloadApiError {
699 spiffe::WorkloadApiError::NoIdentityIssued
700 }
701
702 fn permission_denied(msg: &str) -> spiffe::WorkloadApiError {
703 spiffe::WorkloadApiError::PermissionDenied(msg.to_string())
704 }
705
706 #[test]
707 fn workload_dial_backoff_matches_verified_kluster_schedule() {
708 assert_eq!(
712 WORKLOAD_DIAL_BACKOFF,
713 [
714 Duration::from_secs(1),
715 Duration::from_secs(2),
716 Duration::from_secs(4),
717 Duration::from_secs(8),
718 ]
719 );
720 assert_eq!(WORKLOAD_DIAL_MAX_ATTEMPTS, 5);
721 }
722
723 #[tokio::test]
724 async fn retry_until_identity_issued_succeeds_immediately() {
725 let calls = Arc::new(AtomicUsize::new(0));
726 let calls_probe = calls.clone();
727
728 let result = retry_until_identity_issued("unix:///test.sock", move || {
729 let calls = calls_probe.clone();
730 async move {
731 calls.fetch_add(1, Ordering::SeqCst);
732 Ok(())
733 }
734 })
735 .await;
736
737 assert!(result.is_ok());
738 assert_eq!(
739 calls.load(Ordering::SeqCst),
740 1,
741 "a successful first probe must not retry"
742 );
743 }
744
745 #[tokio::test(start_paused = true)]
746 async fn retry_until_identity_issued_retries_no_identity_issued_then_succeeds() {
747 let calls = Arc::new(AtomicUsize::new(0));
748 let calls_probe = calls.clone();
749 let start = tokio::time::Instant::now();
750
751 let result = retry_until_identity_issued("unix:///test.sock", move || {
752 let calls = calls_probe.clone();
753 async move {
754 let attempt = calls.fetch_add(1, Ordering::SeqCst);
755 if attempt < 2 {
756 Err(no_identity_issued())
757 } else {
758 Ok(())
759 }
760 }
761 })
762 .await;
763
764 assert!(result.is_ok());
765 assert_eq!(
766 calls.load(Ordering::SeqCst),
767 3,
768 "expected 2 failed probes then 1 succeeding probe"
769 );
770 assert_eq!(start.elapsed(), Duration::from_secs(1 + 2));
773 }
774
775 #[tokio::test]
776 async fn retry_until_identity_issued_returns_other_error_immediately_unretried() {
777 let calls = Arc::new(AtomicUsize::new(0));
778 let calls_probe = calls.clone();
779
780 let err = retry_until_identity_issued("unix:///test.sock", move || {
781 let calls = calls_probe.clone();
782 async move {
783 calls.fetch_add(1, Ordering::SeqCst);
784 Err(permission_denied("selectors do not match"))
785 }
786 })
787 .await
788 .unwrap_err();
789
790 assert_eq!(
791 calls.load(Ordering::SeqCst),
792 1,
793 "a non-'no identity issued' failure must never be retried"
794 );
795 match err {
796 ClientError::WorkloadProbe { socket, source } => {
797 assert_eq!(socket, "unix:///test.sock");
798 assert!(matches!(
799 source,
800 spiffe::WorkloadApiError::PermissionDenied(_)
801 ));
802 }
803 other => panic!("expected ClientError::WorkloadProbe, got {other:?}"),
804 }
805 }
806
807 #[tokio::test(start_paused = true)]
808 async fn retry_until_identity_issued_exhausts_after_five_attempts_with_full_backoff() {
809 let calls = Arc::new(AtomicUsize::new(0));
810 let calls_probe = calls.clone();
811 let start = tokio::time::Instant::now();
812
813 let err = retry_until_identity_issued("unix:///test.sock", move || {
814 let calls = calls_probe.clone();
815 async move {
816 calls.fetch_add(1, Ordering::SeqCst);
817 Err(no_identity_issued())
818 }
819 })
820 .await
821 .unwrap_err();
822
823 assert_eq!(
824 calls.load(Ordering::SeqCst),
825 WORKLOAD_DIAL_MAX_ATTEMPTS,
826 "expected exactly WORKLOAD_DIAL_MAX_ATTEMPTS probes when every one fails with no identity issued"
827 );
828 assert_eq!(start.elapsed(), Duration::from_secs(1 + 2 + 4 + 8));
830
831 match err {
832 ClientError::WorkloadNoIdentityIssued {
833 socket,
834 attempts,
835 source,
836 } => {
837 assert_eq!(socket, "unix:///test.sock");
838 assert_eq!(attempts, WORKLOAD_DIAL_MAX_ATTEMPTS);
839 assert!(matches!(source, spiffe::WorkloadApiError::NoIdentityIssued));
840 }
841 other => panic!("expected ClientError::WorkloadNoIdentityIssued, got {other:?}"),
842 }
843 }
844 }
845}
846
847#[cfg(feature = "spiffe-workload")]
848pub use workload::dial_workload;
849
850#[cfg(test)]
851mod tests {
852 use super::*;
853
854 #[test]
855 fn is_loopback_host_matches_go_client_table() {
856 let cases: &[(&str, bool)] = &[
857 ("localhost", true),
858 ("127.0.0.1", true),
859 ("::1", true),
860 ("10.0.0.5", false),
861 ("signet.internal", false),
862 ];
863 for (host, want) in cases {
864 assert_eq!(is_loopback_host(host), *want, "is_loopback_host({host:?})");
865 }
866 }
867
868 #[test]
869 fn admin_transport_decision_loopback_defaults_to_plaintext() {
870 let decision = admin_transport_decision("localhost:8444", None, false, false).unwrap();
871 assert!(!decision.requires_tls());
872 }
873
874 #[test]
875 fn admin_transport_decision_non_loopback_requires_tls() {
876 let decision =
877 admin_transport_decision("signet.internal:8444", None, false, false).unwrap();
878 assert!(decision.requires_tls());
879 }
880
881 #[test]
882 fn admin_transport_decision_force_tls_on_loopback() {
883 let decision = admin_transport_decision("localhost:8444", None, true, false).unwrap();
884 assert!(decision.requires_tls());
885 }
886
887 #[test]
888 fn admin_transport_decision_plaintext_overrides_non_loopback() {
889 let decision =
894 admin_transport_decision("signet.internal:8444", None, false, true).unwrap();
895 assert!(!decision.requires_tls());
896 }
897
898 #[test]
899 fn admin_transport_decision_plaintext_leaves_loopback_unaffected() {
900 let decision = admin_transport_decision("localhost:8444", None, false, true).unwrap();
903 assert!(!decision.requires_tls());
904 }
905
906 #[test]
907 fn admin_transport_decision_rejects_force_tls_and_plaintext_together() {
908 let err = admin_transport_decision("localhost:8444", None, true, true).unwrap_err();
909 assert!(
910 matches!(err, ClientError::ForceTlsPlaintextConflict),
911 "expected ClientError::ForceTlsPlaintextConflict, got {err:?}"
912 );
913 assert_eq!(
914 err.to_string(),
915 "force_tls and plaintext are mutually exclusive"
916 );
917 }
918
919 #[test]
920 fn admin_transport_decision_rejects_plaintext_with_ca_pem() {
921 let pem = b"-----BEGIN CERTIFICATE-----\nnot validated at this layer\n-----END CERTIFICATE-----\n";
922 let err = admin_transport_decision("localhost:8444", Some(pem), false, true).unwrap_err();
923 assert!(
924 matches!(err, ClientError::PlaintextCaPemConflict),
925 "expected ClientError::PlaintextCaPemConflict, got {err:?}"
926 );
927 assert_eq!(
928 err.to_string(),
929 "plaintext and ca_pem are mutually exclusive"
930 );
931 }
932
933 #[test]
934 fn admin_transport_decision_plaintext_with_empty_ca_pem_is_not_a_conflict() {
935 let decision = admin_transport_decision("localhost:8444", Some(&[]), false, true).unwrap();
940 assert!(!decision.requires_tls());
941 }
942
943 #[test]
944 fn admin_transport_decision_ca_pem_forces_tls_even_on_loopback() {
945 let pem = b"-----BEGIN CERTIFICATE-----\n\
951MIICoDCCAYgCCQDLsJN6ayvwqTANBgkqhkiG9w0BAQsFADASMRAwDgYDVQQDDAd0\n\
952ZXN0LWNhMB4XDTI2MDcxMjIxMzM1MFoXDTI2MDcxMzIxMzM1MFowEjEQMA4GA1UE\n\
953AwwHdGVzdC1jYTCCASIwDQYJKoZIhvcNAQEBBQADggEPADCCAQoCggEBALKQT/9e\n\
954HkJlnufQ8dCzc0JdZRO1gHDMY6stgfZljK1dEj2SaANpP3MDVIyDcmKq/6Gwbj4K\n\
955fexqB+1VGLn7CKmopYBvAIwiMDHsQ/R8xDOLwVJRCwnxzAbUUsBF9LvRkDqV4U0/\n\
956i7jizdwtxDHLoB9qEkDKWo3flgIGQtgJ6Vsj7YM9CPq369fby5ZBsPCR3itvEsiZ\n\
957BoM13D3A2RFywYWFpvAvzlzR6LoFd4OnH/8QMh9KTTtxNYw2K8C/a2Cv3GZRROhN\n\
958g5vcQbXLSyYVBUSwdEBT50/pl97KLStN54XEE2YQvoBZCU/kUBrOP888wn+ljafk\n\
959XEMVrZiAKRDnZokCAwEAATANBgkqhkiG9w0BAQsFAAOCAQEAQpRqAdDsxNm+1qFf\n\
9603IW8jJnfMrwdIUukE4c/ms7v3+n6QkdQYidfnZSXCrd0TAzXkRGonrFUDWAfRoGX\n\
961ty0EN/hiU/wmDEvmsNgg9PS5KW3qqoIFRGYdwxn97hjJ0GdgUrbBLg0BweeaP+WW\n\
9620Q7Jive55TT4W+Hwl5KETWOGi2FnvrlrDQGHWY1XKQKQn9J/tEQDMd+COyM9BHez\n\
963oWg4npa5Q/5SdfJs3i4GyGRU4NWYxGfgFi7JiHOZx8t2Nv0RJkYqQu1SMNq97IDo\n\
964ezQtmgLYbjPG41WWrdNT76h1mJgtlCzH0DfI7lQTBIi9AuE5poxPQiBoaC7flMsV\n\
965w8cAzA==\n\
966-----END CERTIFICATE-----\n";
967 let decision = admin_transport_decision("localhost:8444", Some(pem), false, false);
968 assert!(decision.is_ok());
969 assert!(decision.unwrap().requires_tls());
970 }
971
972 #[test]
973 fn admin_transport_decision_rejects_invalid_ca_pem() {
974 let err =
975 admin_transport_decision("signet.internal:8444", Some(b"not a cert"), false, false)
976 .unwrap_err();
977 assert!(
978 matches!(err, ClientError::InvalidCaPem),
979 "expected ClientError::InvalidCaPem, got {err:?}"
980 );
981 assert_eq!(err.to_string(), "invalid CA PEM bundle: no certificates found");
982 }
983
984 #[tokio::test]
985 async fn dial_admin_rejects_empty_token() {
986 let err = dial_admin("localhost:8444", " ", None, false, false)
987 .await
988 .unwrap_err();
989 assert!(matches!(err, ClientError::EmptyToken));
990 assert_eq!(err.to_string(), "token must not be empty");
991 }
992
993 #[tokio::test]
994 async fn gitops_client_and_admin_client_accept_both_channel_kinds() {
995 let plain_channel: Channel = Endpoint::from_static("http://localhost:1").connect_lazy();
1003 let _gitops_over_plain_channel = gitops_client(plain_channel.clone());
1004 let _admin_over_plain_channel = admin_client(plain_channel);
1005
1006 let header_value: tonic::metadata::MetadataValue<tonic::metadata::Ascii> =
1007 "Bearer test-token".parse().unwrap();
1008 let admin_channel: AdminChannel = InterceptedService::new(
1009 Endpoint::from_static("http://localhost:1").connect_lazy(),
1010 TokenInterceptor { header_value },
1011 );
1012 let _gitops_over_admin_channel = gitops_client(admin_channel.clone());
1013 let _admin_over_admin_channel = admin_client(admin_channel);
1014 }
1015
1016 #[test]
1017 fn host_of_handles_bracketed_ipv6_and_bare_hosts() {
1018 assert_eq!(host_of("localhost:8444"), "localhost");
1019 assert_eq!(host_of("127.0.0.1:8444"), "127.0.0.1");
1020 assert_eq!(host_of("[::1]:8444"), "::1");
1021 assert_eq!(host_of("signet.internal"), "signet.internal");
1022 }
1023}