Skip to main content

sui_rpc_api/
lib.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use std::convert::Infallible;
5use std::sync::Arc;
6
7use reader::StateReader;
8use subscription::SubscriptionServiceHandle;
9use sui_http::middleware::callback::CallbackLayer;
10use sui_types::storage::RpcStateReader;
11use sui_types::transaction_executor::ProposerSelector;
12use sui_types::transaction_executor::TransactionExecutor;
13use tap::Pipe;
14use tonic::server::NamedService;
15use tower::Service;
16
17pub mod client;
18mod client_protocol_version;
19mod config;
20mod error;
21pub mod grpc;
22pub mod ledger_history;
23mod metrics;
24pub mod read_mask_defaults;
25mod reader;
26mod response;
27mod service;
28pub mod subscription;
29
30pub use client::Client;
31pub use client_protocol_version::{X_SUI_CLIENT_PROTOCOL_VERSION, client_protocol_version};
32pub use config::Config;
33pub use error::{
34    CheckpointNotFoundError, ErrorDetails, ErrorReason, ObjectNotFoundError, Result, RpcError,
35};
36pub use metrics::{
37    GrpcMethodAllowlist, RpcMetrics, RpcMetricsMakeCallbackHandler,
38    grpc_method_paths_from_file_descriptor_sets,
39};
40pub use reader::TransactionNotFoundError;
41pub use sui_rpc::proto;
42
43#[derive(Clone)]
44pub struct ServerVersion {
45    pub bin: &'static str,
46    pub version: &'static str,
47}
48
49impl ServerVersion {
50    pub fn new(bin: &'static str, version: &'static str) -> Self {
51        Self { bin, version }
52    }
53}
54
55impl std::fmt::Display for ServerVersion {
56    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
57        f.write_str(self.bin)?;
58        f.write_str("/")?;
59        f.write_str(self.version)
60    }
61}
62
63#[derive(Clone)]
64pub struct RpcService {
65    reader: StateReader,
66    executor: Option<Arc<dyn TransactionExecutor>>,
67    proposer_selector: Option<Arc<dyn ProposerSelector>>,
68    subscription_service_handle: Option<SubscriptionServiceHandle>,
69    chain_id: sui_types::digests::ChainIdentifier,
70    server_version: Option<ServerVersion>,
71    metrics: Option<Arc<RpcMetrics>>,
72    pub(crate) list_metrics: Option<Arc<metrics::ListApiMetrics>>,
73    config: Config,
74    extra_routes: axum::Router,
75    extra_service_names: Vec<&'static str>,
76    extra_file_descriptor_sets: Vec<&'static [u8]>,
77}
78
79impl RpcService {
80    pub fn new(reader: Arc<dyn RpcStateReader>) -> Self {
81        let chain_id = reader.get_chain_identifier().unwrap();
82        Self {
83            reader: StateReader::new(reader),
84            executor: None,
85            proposer_selector: None,
86            subscription_service_handle: None,
87            chain_id,
88            server_version: None,
89            metrics: None,
90            list_metrics: None,
91            config: Config::default(),
92            extra_routes: axum::Router::new(),
93            extra_service_names: Vec::new(),
94            extra_file_descriptor_sets: Vec::new(),
95        }
96    }
97
98    pub fn with_server_version(&mut self, server_version: ServerVersion) -> &mut Self {
99        self.server_version = Some(server_version);
100        self
101    }
102
103    pub fn with_config(&mut self, config: Config) {
104        self.config = config;
105    }
106
107    pub fn with_executor(&mut self, executor: Arc<dyn TransactionExecutor + Send + Sync>) {
108        self.executor = Some(executor);
109    }
110
111    pub fn with_proposer_selector(&mut self, proposer_selector: Arc<dyn ProposerSelector>) {
112        self.proposer_selector = Some(proposer_selector);
113    }
114
115    pub fn with_subscription_service(
116        &mut self,
117        subscription_service_handle: SubscriptionServiceHandle,
118    ) {
119        self.subscription_service_handle = Some(subscription_service_handle);
120    }
121
122    pub fn with_metrics(&mut self, registry: &prometheus::Registry) {
123        self.metrics = Some(Arc::new(RpcMetrics::new(registry)));
124        self.list_metrics = Some(Arc::new(metrics::ListApiMetrics::new(registry)));
125    }
126
127    pub fn with_custom_service<S>(&mut self, svc: S)
128    where
129        S: Service<
130                axum::extract::Request,
131                Response: axum::response::IntoResponse,
132                Error = Infallible,
133            > + NamedService
134            + Clone
135            + Send
136            + Sync
137            + 'static,
138        S::Future: Send + 'static,
139        S::Error: Into<grpc::BoxError> + Send,
140    {
141        self.extra_service_names.push(S::NAME);
142        self.extra_routes = std::mem::take(&mut self.extra_routes)
143            .route_service(&format!("/{}/{{*rest}}", S::NAME), svc);
144    }
145
146    pub fn with_file_descriptor_set(&mut self, encoded_fds: &'static [u8]) {
147        self.extra_file_descriptor_sets.push(encoded_fds);
148    }
149
150    pub fn chain_id(&self) -> sui_types::digests::ChainIdentifier {
151        self.chain_id
152    }
153
154    pub fn server_version(&self) -> Option<&ServerVersion> {
155        self.server_version.as_ref()
156    }
157
158    pub async fn into_router(mut self) -> axum::Router {
159        let metrics = self.metrics.clone();
160        let extra_routes = std::mem::take(&mut self.extra_routes);
161        let extra_service_names = std::mem::take(&mut self.extra_service_names);
162
163        // Single source of truth for every encoded FileDescriptorSet that
164        // backs a gRPC service mounted below. Consumed by the reflection
165        // services, the metrics allowlist, and the request-log layer so they
166        // cannot drift out of sync.
167        let built_in_file_descriptor_sets: [&[u8]; 5] = [
168            sui_rpc::proto::google::protobuf::FILE_DESCRIPTOR_SET,
169            sui_rpc::proto::google::rpc::FILE_DESCRIPTOR_SET,
170            sui_rpc::proto::sui::rpc::v2::FILE_DESCRIPTOR_SET,
171            sui_rpc::proto::sui::rpc::v2alpha::FILE_DESCRIPTOR_SET,
172            tonic_health::pb::FILE_DESCRIPTOR_SET,
173        ];
174        let file_descriptor_sets: Vec<&[u8]> = built_in_file_descriptor_sets
175            .into_iter()
176            .chain(std::mem::take(&mut self.extra_file_descriptor_sets))
177            .collect();
178
179        // Allowlist of `/Service/Method` paths used by the metrics middleware
180        // to bound prometheus label cardinality.
181        let grpc_method_allowlist = Arc::new(
182            metrics::grpc_method_paths_from_file_descriptor_sets(&file_descriptor_sets)
183                .expect("registered FileDescriptorSet bytes must be valid protobuf"),
184        );
185
186        let request_log =
187            mysten_network::request_log::GrpcRequestLogLayer::from_encoded_file_descriptor_sets(
188                file_descriptor_sets.iter().copied(),
189            )
190            .unwrap_or_else(|e| {
191                // Extra sets registered by embedders may not merge cleanly (e.g. missing
192                // imports). Reflection and metrics tolerate that, so don't fail startup —
193                // capture just won't decode those extra services.
194                tracing::warn!(
195                    "request-log descriptor pool falling back to built-in file descriptor sets: {e}"
196                );
197                mysten_network::request_log::GrpcRequestLogLayer::from_encoded_file_descriptor_sets(
198                    built_in_file_descriptor_sets,
199                )
200                .expect("built-in FileDescriptorSet bytes must be valid protobuf")
201            });
202
203        let router = {
204            let ledger_service =
205                sui_rpc::proto::sui::rpc::v2::ledger_service_server::LedgerServiceServer::new(
206                    self.clone(),
207                )
208                .send_compressed(tonic::codec::CompressionEncoding::Zstd);
209            let proof_service_v2alpha =
210                sui_rpc::proto::sui::rpc::v2alpha::proof_service_server::ProofServiceServer::new(
211                    self.clone(),
212                )
213                .send_compressed(tonic::codec::CompressionEncoding::Zstd);
214            let transaction_execution_service = sui_rpc::proto::sui::rpc::v2::transaction_execution_service_server::TransactionExecutionServiceServer::new(self.clone())
215                .send_compressed(tonic::codec::CompressionEncoding::Zstd);
216            let state_service =
217                sui_rpc::proto::sui::rpc::v2::state_service_server::StateServiceServer::new(
218                    self.clone(),
219                )
220                .send_compressed(tonic::codec::CompressionEncoding::Zstd);
221            let signature_verification_service = sui_rpc::proto::sui::rpc::v2::signature_verification_service_server::SignatureVerificationServiceServer::new(self.clone())
222                .send_compressed(tonic::codec::CompressionEncoding::Zstd);
223            let move_package_service = sui_rpc::proto::sui::rpc::v2::move_package_service_server::MovePackageServiceServer::new(self.clone())
224                .send_compressed(tonic::codec::CompressionEncoding::Zstd);
225            let name_service =
226                sui_rpc::proto::sui::rpc::v2::name_service_server::NameServiceServer::new(
227                    self.clone(),
228                )
229                .send_compressed(tonic::codec::CompressionEncoding::Zstd);
230
231            let (health_reporter, health_service) = tonic_health::server::health_reporter();
232
233            let mut reflection_v1_builder = tonic_reflection::server::Builder::configure();
234            let mut reflection_v1alpha_builder = tonic_reflection::server::Builder::configure();
235            for fds in &file_descriptor_sets {
236                reflection_v1_builder =
237                    reflection_v1_builder.register_encoded_file_descriptor_set(fds);
238                reflection_v1alpha_builder =
239                    reflection_v1alpha_builder.register_encoded_file_descriptor_set(fds);
240            }
241
242            let reflection_v1 = reflection_v1_builder.build_v1().unwrap();
243            let reflection_v1alpha = reflection_v1alpha_builder.build_v1alpha().unwrap();
244
245            fn service_name<S: tonic::server::NamedService>(_service: &S) -> &'static str {
246                S::NAME
247            }
248
249            for service_name in [
250                service_name(&ledger_service),
251                service_name(&transaction_execution_service),
252                service_name(&state_service),
253                service_name(&signature_verification_service),
254                service_name(&move_package_service),
255                service_name(&name_service),
256                service_name(&proof_service_v2alpha),
257                service_name(&reflection_v1),
258                service_name(&reflection_v1alpha),
259            ] {
260                health_reporter
261                    .set_service_status(service_name, tonic_health::ServingStatus::Serving)
262                    .await;
263            }
264
265            let mut services = grpc::Services::new()
266                .timeout(self.config.grpc_timeout())
267                // V2
268                .add_service(ledger_service)
269                .add_service(transaction_execution_service)
270                .add_service(state_service)
271                .add_service(signature_verification_service)
272                .add_service(move_package_service)
273                .add_service(name_service)
274                // V2alpha
275                .add_service(proof_service_v2alpha)
276                // Reflection
277                .add_service(reflection_v1)
278                .add_service(reflection_v1alpha);
279
280            if self.subscription_service_handle.is_some() {
281                let subscription_service =
282sui_rpc::proto::sui::rpc::v2::subscription_service_server::SubscriptionServiceServer::new(self.clone());
283                health_reporter
284                    .set_service_status(
285                        service_name(&subscription_service),
286                        tonic_health::ServingStatus::Serving,
287                    )
288                    .await;
289
290                services = services.add_service(subscription_service);
291            }
292
293            for name in &extra_service_names {
294                health_reporter
295                    .set_service_status(*name, tonic_health::ServingStatus::Serving)
296                    .await;
297            }
298
299            services
300                .merge_router(extra_routes)
301                .add_service(health_service)
302                .into_router(request_log)
303        };
304
305        let health_endpoint = axum::Router::new()
306            .route("/health", axum::routing::get(service::health::health))
307            .with_state(self.clone());
308
309        router
310            .merge(health_endpoint)
311            .layer(axum::middleware::map_response_with_state(
312                self,
313                response::append_info_headers,
314            ))
315            .pipe(|router| {
316                if let Some(metrics) = metrics {
317                    router.layer(CallbackLayer::new(
318                        metrics::RpcMetricsMakeCallbackHandler::with_grpc_method_allowlist(
319                            metrics,
320                            grpc_method_allowlist,
321                        ),
322                    ))
323                } else {
324                    router
325                }
326            })
327    }
328
329    pub async fn start_service(self, socket_address: std::net::SocketAddr) {
330        let listener = tokio::net::TcpListener::bind(socket_address).await.unwrap();
331        axum::serve(listener, self.into_router().await)
332            .await
333            .unwrap();
334    }
335}
336
337#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize)]
338#[serde(rename_all = "lowercase")]
339pub enum Direction {
340    Ascending,
341    Descending,
342}
343
344impl Direction {
345    pub fn is_descending(self) -> bool {
346        matches!(self, Self::Descending)
347    }
348}
349
350#[cfg(test)]
351mod tests {
352    /// The request-log layer's descriptor pool is built from these sets at server startup with an
353    /// `expect`, so they must always merge into one valid pool.
354    #[test]
355    fn request_log_pool_builds_from_registered_file_descriptor_sets() {
356        mysten_network::request_log::GrpcRequestLogLayer::from_encoded_file_descriptor_sets([
357            sui_rpc::proto::google::protobuf::FILE_DESCRIPTOR_SET,
358            sui_rpc::proto::google::rpc::FILE_DESCRIPTOR_SET,
359            sui_rpc::proto::sui::rpc::v2::FILE_DESCRIPTOR_SET,
360            sui_rpc::proto::sui::rpc::v2alpha::FILE_DESCRIPTOR_SET,
361            tonic_health::pb::FILE_DESCRIPTOR_SET,
362        ])
363        .unwrap();
364    }
365}