function_sdk_rust/
server.rs1use std::net::SocketAddr;
4use std::path::{Path, PathBuf};
5
6use tonic::transport::{Certificate, Identity, Server, ServerTlsConfig};
7
8use crate::Error;
9use crate::proto::v1::function_runner_service_server::{
10 FunctionRunnerService, FunctionRunnerServiceServer,
11};
12
13#[derive(clap::Parser, Debug)]
15#[command(version, about = "A Crossplane composition function")]
16pub struct Args {
17 #[arg(short, long, env = "DEBUG", default_value_t = false)]
19 pub debug: bool,
20
21 #[arg(long, default_value = "0.0.0.0:9443")]
23 pub address: String,
24
25 #[arg(long, env = "TLS_SERVER_CERTS_DIR")]
27 pub tls_certs_dir: Option<PathBuf>,
28
29 #[arg(long, default_value_t = false)]
31 pub insecure: bool,
32
33 #[arg(long)]
36 pub max_recv_message_size: Option<usize>,
37
38 #[arg(long, env = "METRICS_ADDRESS", default_value = ":8080")]
42 pub metrics_address: String,
43}
44
45pub async fn serve<F: FunctionRunnerService>(function: F, args: &Args) -> Result<(), Error> {
59 let address: SocketAddr = args.address.parse()?;
60
61 let mut builder = Server::builder();
62 if !args.insecure {
63 let dir = args
64 .tls_certs_dir
65 .as_deref()
66 .ok_or(Error::MissingTlsCertsDir)?;
67 builder = builder.tls_config(tls_config(dir)?)?;
68 }
69
70 let mut function_service = FunctionRunnerServiceServer::new(function);
71 if let Some(size) = args.max_recv_message_size {
72 function_service = function_service.max_decoding_message_size(size);
73 }
74
75 let reflection_v1 = tonic_reflection::server::Builder::configure()
76 .register_encoded_file_descriptor_set(crate::proto::FILE_DESCRIPTOR_SET)
77 .build_v1()?;
78 let reflection_v1alpha = tonic_reflection::server::Builder::configure()
79 .register_encoded_file_descriptor_set(crate::proto::FILE_DESCRIPTOR_SET)
80 .build_v1alpha()?;
81
82 let (health_reporter, health_service) = tonic_health::server::health_reporter();
83 health_reporter
84 .set_serving::<FunctionRunnerServiceServer<F>>()
85 .await;
86
87 if !args.metrics_address.is_empty() {
93 crate::metrics::initialize();
94 let metrics_address = args.metrics_address.clone();
95 tokio::spawn(async move {
96 if let Err(e) = crate::metrics::serve(&metrics_address).await {
97 tracing::error!(error = %e, "cannot serve metrics");
98 }
99 });
100 }
101
102 tracing::info!(%address, insecure = args.insecure, "serving FunctionRunnerService");
103
104 builder
105 .layer(crate::metrics::MetricsLayer)
106 .add_service(function_service)
107 .add_service(health_service)
108 .add_service(reflection_v1)
109 .add_service(reflection_v1alpha)
110 .serve_with_shutdown(address, shutdown_signal())
111 .await?;
112
113 Ok(())
114}
115
116fn tls_config(dir: &Path) -> Result<ServerTlsConfig, Error> {
117 let read = |name: &str| {
118 let path = dir.join(name);
119 std::fs::read(&path).map_err(|source| Error::ReadCertificate { path, source })
120 };
121 let cert = read("tls.crt")?;
122 let key = read("tls.key")?;
123 let ca = read("ca.crt")?;
124
125 Ok(ServerTlsConfig::new()
126 .identity(Identity::from_pem(cert, key))
127 .client_ca_root(Certificate::from_pem(ca))
128 .client_auth_optional(false))
129}
130
131#[cfg(unix)]
132async fn shutdown_signal() {
133 let mut sigterm = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
134 .expect("cannot install SIGTERM handler");
135 tokio::select! {
136 _ = sigterm.recv() => {}
137 _ = tokio::signal::ctrl_c() => {}
138 }
139 tracing::info!("shutting down");
140}
141
142#[cfg(not(unix))]
143async fn shutdown_signal() {
144 let _ = tokio::signal::ctrl_c().await;
145 tracing::info!("shutting down");
146}