Skip to main content

function_sdk_rust/
server.rs

1//! A spec-compliant gRPC server runtime for composition functions.
2
3use 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/// CLI arguments required by the Crossplane composition function spec.
14#[derive(clap::Parser, Debug)]
15#[command(version, about = "A Crossplane composition function")]
16pub struct Args {
17    /// Emit debug logs.
18    #[arg(short, long, env = "DEBUG", default_value_t = false)]
19    pub debug: bool,
20
21    /// Address at which to listen for gRPC connections.
22    #[arg(long, default_value = "0.0.0.0:9443")]
23    pub address: String,
24
25    /// Directory containing tls.crt, tls.key, and ca.crt; serve using mTLS.
26    #[arg(long, env = "TLS_SERVER_CERTS_DIR")]
27    pub tls_certs_dir: Option<PathBuf>,
28
29    /// Run without mTLS credentials. If set, --tls-certs-dir is ignored.
30    #[arg(long, default_value_t = false)]
31    pub insecure: bool,
32
33    /// Maximum size in bytes of gRPC messages the function accepts.
34    /// Defaults to the gRPC default of 4MB.
35    #[arg(long)]
36    pub max_recv_message_size: Option<usize>,
37
38    /// Address at which to serve metrics at /metrics (OpenMetrics 1.0, or
39    /// the classic Prometheus text format for an Accept header that asks
40    /// for it) - the Go SDK's default; empty disables them.
41    #[arg(long, env = "METRICS_ADDRESS", default_value = ":8080")]
42    pub metrics_address: String,
43}
44
45/// Starts a gRPC server and serves RunFunctionRequests until SIGTERM or
46/// SIGINT, then shuts down gracefully.
47///
48/// Serves with mTLS from `--tls-certs-dir` (tls.crt and tls.key must be the
49/// function's PEM-encoded certificate and key; ca.crt a PEM-encoded CA used
50/// to authenticate Crossplane) unless `--insecure` is set. gRPC server
51/// reflection is enabled for both the v1 and v1alpha reflection APIs, and
52/// the gRPC health service reports the function as serving.
53///
54/// Unless `--metrics-address` is empty, the gRPC server metrics the Go SDK
55/// serves (`grpc_server_started_total` and friends, same names and labels)
56/// are served there at /metrics - OpenMetrics 1.0 as the main format, the
57/// classic Prometheus text format for scrapers that ask for it.
58pub 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    // The metrics endpoint runs beside the gRPC server, like the Go SDK's;
88    // a failure to serve it is logged, never the function's failure. The
89    // counting layer below is always installed - with the endpoint disabled
90    // the counts are simply never served, which is indistinguishable from
91    // the Go SDK installing no interceptor.
92    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}