Skip to main content

risingwave_rpc_client/
meta_client.rs

1// Copyright 2022 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{BTreeMap, HashMap, HashSet};
16use std::fmt::{Debug, Display};
17use std::sync::Arc;
18use std::sync::atomic::AtomicBool;
19use std::sync::atomic::Ordering::Relaxed;
20use std::thread;
21use std::time::{Duration, SystemTime};
22
23use anyhow::{Context, anyhow};
24use async_trait::async_trait;
25use cluster_limit_service_client::ClusterLimitServiceClient;
26use either::Either;
27use futures::stream::BoxStream;
28use list_rate_limits_response::RateLimitInfo;
29use lru::LruCache;
30use replace_job_plan::ReplaceJob;
31use risingwave_common::catalog::{
32    AlterDatabaseParam, FunctionId, IndexId, ObjectId, SecretId, TableId,
33};
34use risingwave_common::config::{MAX_CONNECTION_WINDOW_SIZE, MetaConfig};
35use risingwave_common::hash::WorkerSlotMapping;
36use risingwave_common::id::{
37    ConnectionId, DatabaseId, JobId, SchemaId, SinkId, SubscriptionId, UserId, ViewId, WorkerId,
38};
39use risingwave_common::monitor::EndpointExt;
40use risingwave_common::system_param::AdaptiveParallelismStrategy;
41use risingwave_common::system_param::reader::SystemParamsReader;
42use risingwave_common::telemetry::report::TelemetryInfoFetcher;
43use risingwave_common::util::addr::HostAddr;
44use risingwave_common::util::meta_addr::MetaAddressStrategy;
45use risingwave_common::util::resource_util::cpu::total_cpu_available;
46use risingwave_common::util::resource_util::hostname;
47use risingwave_common::util::resource_util::memory::system_memory_available_bytes;
48use risingwave_common::util::retry::exponential_backoff;
49use risingwave_common::util::version::current_rw_version;
50use risingwave_error::bail;
51use risingwave_error::tonic::ErrorIsFromTonicServerImpl;
52use risingwave_hummock_sdk::change_log::{TableChangeLog, TableChangeLogs};
53use risingwave_hummock_sdk::version::{HummockVersion, HummockVersionDelta};
54use risingwave_hummock_sdk::{
55    CompactionGroupId, HummockEpoch, HummockVersionId, ObjectIdRange, SyncResult,
56};
57use risingwave_pb::backup_service::backup_service_client::BackupServiceClient;
58use risingwave_pb::backup_service::*;
59use risingwave_pb::catalog::{
60    Connection, PbComment, PbDatabase, PbFunction, PbIndex, PbSchema, PbSink, PbSource,
61    PbSubscription, PbTable, PbView, Table,
62};
63use risingwave_pb::cloud_service::cloud_service_client::CloudServiceClient;
64use risingwave_pb::cloud_service::*;
65use risingwave_pb::common::worker_node::Property;
66use risingwave_pb::common::{HostAddress, OptionalUint32, OptionalUint64, WorkerNode, WorkerType};
67use risingwave_pb::configured_monitor_service_client;
68use risingwave_pb::connector_service::sink_coordination_service_client::SinkCoordinationServiceClient;
69use risingwave_pb::ddl_service::alter_owner_request::Object;
70use risingwave_pb::ddl_service::create_iceberg_table_request::{PbSinkJobInfo, PbTableJobInfo};
71use risingwave_pb::ddl_service::ddl_service_client::DdlServiceClient;
72use risingwave_pb::ddl_service::*;
73use risingwave_pb::hummock::compact_task::TaskStatus;
74use risingwave_pb::hummock::get_compaction_score_response::PickerInfo;
75use risingwave_pb::hummock::get_table_change_logs_request::PbTableFilter;
76use risingwave_pb::hummock::hummock_manager_service_client::HummockManagerServiceClient;
77use risingwave_pb::hummock::rise_ctl_update_compaction_config_request::mutable_config::MutableConfig;
78use risingwave_pb::hummock::subscribe_compaction_event_request::Register;
79use risingwave_pb::hummock::write_limits::WriteLimit;
80use risingwave_pb::hummock::*;
81use risingwave_pb::iceberg_compaction::subscribe_iceberg_compaction_event_request::Register as IcebergRegister;
82use risingwave_pb::iceberg_compaction::{
83    SubscribeIcebergCompactionEventRequest, SubscribeIcebergCompactionEventResponse,
84    subscribe_iceberg_compaction_event_request,
85};
86use risingwave_pb::id::{ActorId, FragmentId, HummockSstableId, IcebergCompactionTaskId, SourceId};
87use risingwave_pb::meta::alter_connector_props_request::{
88    AlterConnectorPropsObject, AlterIcebergTableIds, ExtraOptions,
89};
90use risingwave_pb::meta::cancel_creating_jobs_request::PbJobs;
91use risingwave_pb::meta::cluster_service_client::ClusterServiceClient;
92use risingwave_pb::meta::event_log_service_client::EventLogServiceClient;
93use risingwave_pb::meta::heartbeat_service_client::HeartbeatServiceClient;
94use risingwave_pb::meta::hosted_iceberg_catalog_service_client::HostedIcebergCatalogServiceClient;
95use risingwave_pb::meta::list_actor_splits_response::ActorSplit;
96use risingwave_pb::meta::list_actor_states_response::ActorState;
97use risingwave_pb::meta::list_cdc_progress_response::PbCdcProgress;
98use risingwave_pb::meta::list_iceberg_tables_response::IcebergTable;
99use risingwave_pb::meta::list_refresh_table_states_response::RefreshTableState;
100use risingwave_pb::meta::list_streaming_job_states_response::StreamingJobState;
101use risingwave_pb::meta::list_table_fragments_response::TableFragmentInfo;
102use risingwave_pb::meta::meta_member_service_client::MetaMemberServiceClient;
103use risingwave_pb::meta::notification_service_client::NotificationServiceClient;
104use risingwave_pb::meta::scale_service_client::ScaleServiceClient;
105use risingwave_pb::meta::serving_service_client::ServingServiceClient;
106use risingwave_pb::meta::session_param_service_client::SessionParamServiceClient;
107use risingwave_pb::meta::stream_manager_service_client::StreamManagerServiceClient;
108use risingwave_pb::meta::system_params_service_client::SystemParamsServiceClient;
109use risingwave_pb::meta::telemetry_info_service_client::TelemetryInfoServiceClient;
110use risingwave_pb::meta::{FragmentDistribution, *};
111use risingwave_pb::monitor_service::monitor_service_client::MonitorServiceClient;
112use risingwave_pb::monitor_service::stack_trace_request::ActorTracesFormat;
113use risingwave_pb::monitor_service::{StackTraceRequest, StackTraceResponse};
114use risingwave_pb::secret::PbSecretRef;
115use risingwave_pb::stream_plan::StreamFragmentGraph;
116use risingwave_pb::user::alter_default_privilege_request::Operation as AlterDefaultPrivilegeOperation;
117use risingwave_pb::user::update_user_request::UpdateField;
118use risingwave_pb::user::user_service_client::UserServiceClient;
119use risingwave_pb::user::*;
120use thiserror_ext::AsReport;
121use tokio::sync::mpsc::{Receiver, UnboundedSender, unbounded_channel};
122use tokio::sync::oneshot::Sender;
123use tokio::sync::{RwLock, mpsc, oneshot};
124use tokio::task::JoinHandle;
125use tokio::time::{self};
126use tokio_retry::strategy::jitter;
127use tokio_stream::wrappers::UnboundedReceiverStream;
128use tonic::transport::Endpoint;
129use tonic::{Code, Request, Streaming};
130
131use crate::channel::{Channel, WrappedChannelExt};
132use crate::error::{Result, RpcError};
133use crate::hummock_meta_client::{
134    CompactionEventItem, HummockMetaClient, HummockMetaClientChangeLogInfo,
135    IcebergCompactionEventItem,
136};
137use crate::meta_rpc_client_method_impl;
138
139/// Client to meta server. Cloning the instance is lightweight.
140#[derive(Clone, Debug)]
141pub struct MetaClient {
142    worker_id: WorkerId,
143    worker_type: WorkerType,
144    host_addr: HostAddr,
145    inner: GrpcMetaClient,
146    meta_config: Arc<MetaConfig>,
147    cluster_id: String,
148    shutting_down: Arc<AtomicBool>,
149}
150
151impl MetaClient {
152    pub fn worker_id(&self) -> WorkerId {
153        self.worker_id
154    }
155
156    pub fn host_addr(&self) -> &HostAddr {
157        &self.host_addr
158    }
159
160    pub fn worker_type(&self) -> WorkerType {
161        self.worker_type
162    }
163
164    pub fn cluster_id(&self) -> &str {
165        &self.cluster_id
166    }
167
168    /// Subscribe to notification from meta.
169    pub async fn subscribe(
170        &self,
171        subscribe_type: SubscribeType,
172    ) -> Result<Streaming<SubscribeResponse>> {
173        let request = SubscribeRequest {
174            subscribe_type: subscribe_type as i32,
175            host: Some(self.host_addr.to_protobuf()),
176            worker_id: self.worker_id(),
177        };
178
179        let retry_strategy = GrpcMetaClient::retry_strategy_to_bound(
180            Duration::from_secs(self.meta_config.max_heartbeat_interval_secs as u64),
181            true,
182        );
183
184        tokio_retry::Retry::spawn(retry_strategy, || async {
185            let request = request.clone();
186            self.inner.subscribe(request).await
187        })
188        .await
189    }
190
191    pub async fn create_connection(
192        &self,
193        connection_name: String,
194        database_id: DatabaseId,
195        schema_id: SchemaId,
196        owner_id: UserId,
197        req: create_connection_request::Payload,
198    ) -> Result<WaitVersion> {
199        let request = CreateConnectionRequest {
200            name: connection_name,
201            database_id,
202            schema_id,
203            owner_id,
204            payload: Some(req),
205        };
206        let resp = self.inner.create_connection(request).await?;
207        Ok(resp
208            .version
209            .ok_or_else(|| anyhow!("wait version not set"))?)
210    }
211
212    pub async fn create_secret(
213        &self,
214        secret_name: String,
215        database_id: DatabaseId,
216        schema_id: SchemaId,
217        owner_id: UserId,
218        value: Vec<u8>,
219    ) -> Result<WaitVersion> {
220        let request = CreateSecretRequest {
221            name: secret_name,
222            database_id,
223            schema_id,
224            owner_id,
225            value,
226        };
227        let resp = self.inner.create_secret(request).await?;
228        Ok(resp
229            .version
230            .ok_or_else(|| anyhow!("wait version not set"))?)
231    }
232
233    pub async fn list_connections(&self, _name: Option<&str>) -> Result<Vec<Connection>> {
234        let request = ListConnectionsRequest {};
235        let resp = self.inner.list_connections(request).await?;
236        Ok(resp.connections)
237    }
238
239    pub async fn drop_connection(
240        &self,
241        connection_id: ConnectionId,
242        cascade: bool,
243    ) -> Result<WaitVersion> {
244        let request = DropConnectionRequest {
245            connection_id,
246            cascade,
247        };
248        let resp = self.inner.drop_connection(request).await?;
249        Ok(resp
250            .version
251            .ok_or_else(|| anyhow!("wait version not set"))?)
252    }
253
254    pub async fn drop_secret(&self, secret_id: SecretId, cascade: bool) -> Result<WaitVersion> {
255        let request = DropSecretRequest { secret_id, cascade };
256        let resp = self.inner.drop_secret(request).await?;
257        Ok(resp
258            .version
259            .ok_or_else(|| anyhow!("wait version not set"))?)
260    }
261
262    /// Register the current node to the cluster and set the corresponding worker id.
263    ///
264    /// Retry if there's connection issue with the meta node. Exit the process if the registration fails.
265    pub async fn register_new(
266        addr_strategy: MetaAddressStrategy,
267        worker_type: WorkerType,
268        addr: &HostAddr,
269        property: Property,
270        meta_config: Arc<MetaConfig>,
271    ) -> (Self, SystemParamsReader) {
272        let ret =
273            Self::register_new_inner(addr_strategy, worker_type, addr, property, meta_config).await;
274
275        match ret {
276            Ok(ret) => ret,
277            Err(err) => {
278                tracing::error!(error = %err.as_report(), "failed to register worker, exiting...");
279                std::process::exit(1);
280            }
281        }
282    }
283
284    async fn register_new_inner(
285        addr_strategy: MetaAddressStrategy,
286        worker_type: WorkerType,
287        addr: &HostAddr,
288        property: Property,
289        meta_config: Arc<MetaConfig>,
290    ) -> Result<(Self, SystemParamsReader)> {
291        tracing::info!("register meta client using strategy: {}", addr_strategy);
292
293        // Retry until reaching `max_heartbeat_interval_secs`
294        let retry_strategy = GrpcMetaClient::retry_strategy_to_bound(
295            Duration::from_secs(meta_config.max_heartbeat_interval_secs as u64),
296            true,
297        );
298
299        let init_result: Result<_> = tokio_retry::RetryIf::spawn(
300            retry_strategy,
301            || async {
302                let grpc_meta_client =
303                    GrpcMetaClient::new(&addr_strategy, meta_config.clone()).await?;
304
305                let add_worker_resp = grpc_meta_client
306                    .add_worker_node(AddWorkerNodeRequest {
307                        worker_type: worker_type as i32,
308                        host: Some(addr.to_protobuf()),
309                        property: Some(property.clone()),
310                        resource: Some(risingwave_pb::common::worker_node::Resource {
311                            rw_version: current_rw_version(),
312                            total_memory_bytes: system_memory_available_bytes() as _,
313                            total_cpu_cores: total_cpu_available() as _,
314                            hostname: hostname(),
315                        }),
316                    })
317                    .await
318                    .context("failed to add worker node")?;
319
320                let system_params_resp = grpc_meta_client
321                    .get_system_params(GetSystemParamsRequest {})
322                    .await
323                    .context("failed to get initial system params")?;
324
325                Ok((add_worker_resp, system_params_resp, grpc_meta_client))
326            },
327            // Only retry if there's any transient connection issue.
328            // If the error is from our implementation or business, do not retry it.
329            |e: &RpcError| !e.is_from_tonic_server_impl(),
330        )
331        .await;
332
333        let (add_worker_resp, system_params_resp, grpc_meta_client) = init_result?;
334        let worker_id = add_worker_resp
335            .node_id
336            .expect("AddWorkerNodeResponse::node_id is empty");
337
338        let meta_client = Self {
339            worker_id,
340            worker_type,
341            host_addr: addr.clone(),
342            inner: grpc_meta_client,
343            meta_config: meta_config.clone(),
344            cluster_id: add_worker_resp.cluster_id,
345            shutting_down: Arc::new(false.into()),
346        };
347
348        static REPORT_PANIC: std::sync::Once = std::sync::Once::new();
349        REPORT_PANIC.call_once(|| {
350            let meta_client_clone = meta_client.clone();
351            std::panic::update_hook(move |default_hook, info| {
352                // Try to report panic event to meta node.
353                meta_client_clone.try_add_panic_event_blocking(info, None);
354                default_hook(info);
355            });
356        });
357
358        Ok((meta_client, system_params_resp.params.unwrap().into()))
359    }
360
361    /// Activate the current node in cluster to confirm it's ready to serve.
362    pub async fn activate(&self, addr: &HostAddr) -> Result<()> {
363        let request = ActivateWorkerNodeRequest {
364            host: Some(addr.to_protobuf()),
365            node_id: self.worker_id,
366        };
367        let retry_strategy = GrpcMetaClient::retry_strategy_to_bound(
368            Duration::from_secs(self.meta_config.max_heartbeat_interval_secs as u64),
369            true,
370        );
371        tokio_retry::Retry::spawn(retry_strategy, || async {
372            let request = request.clone();
373            self.inner.activate_worker_node(request).await
374        })
375        .await?;
376
377        Ok(())
378    }
379
380    /// Send heartbeat signal to meta service.
381    pub async fn send_heartbeat(&self) -> Result<()> {
382        let request = HeartbeatRequest {
383            node_id: self.worker_id,
384        };
385        let resp = self.inner.heartbeat(request).await?;
386        if let Some(status) = resp.status
387            && status.code() == risingwave_pb::common::status::Code::UnknownWorker
388        {
389            // Ignore the error if we're already shutting down.
390            // Otherwise, exit the process.
391            if !self.shutting_down.load(Relaxed) {
392                tracing::error!(message = status.message, "worker expired");
393                std::process::exit(1);
394            }
395        }
396        Ok(())
397    }
398
399    pub async fn create_database(&self, db: PbDatabase) -> Result<WaitVersion> {
400        let request = CreateDatabaseRequest { db: Some(db) };
401        let resp = self.inner.create_database(request).await?;
402        // TODO: handle error in `resp.status` here
403        Ok(resp
404            .version
405            .ok_or_else(|| anyhow!("wait version not set"))?)
406    }
407
408    pub async fn create_schema(&self, schema: PbSchema) -> Result<WaitVersion> {
409        let request = CreateSchemaRequest {
410            schema: Some(schema),
411        };
412        let resp = self.inner.create_schema(request).await?;
413        // TODO: handle error in `resp.status` here
414        Ok(resp
415            .version
416            .ok_or_else(|| anyhow!("wait version not set"))?)
417    }
418
419    pub async fn create_materialized_view(
420        &self,
421        table: PbTable,
422        graph: StreamFragmentGraph,
423        dependencies: HashSet<ObjectId>,
424        resource_type: streaming_job_resource_type::ResourceType,
425        if_not_exists: bool,
426        refresh_interval_sec: Option<u64>,
427    ) -> Result<WaitVersion> {
428        let request = CreateMaterializedViewRequest {
429            materialized_view: Some(table),
430            fragment_graph: Some(graph),
431            resource_type: Some(PbStreamingJobResourceType {
432                resource_type: Some(resource_type),
433            }),
434            dependencies: dependencies.into_iter().collect(),
435            if_not_exists,
436            refresh_interval_sec,
437        };
438        let resp = self.inner.create_materialized_view(request).await?;
439        // TODO: handle error in `resp.status` here
440        Ok(resp
441            .version
442            .ok_or_else(|| anyhow!("wait version not set"))?)
443    }
444
445    pub async fn drop_materialized_view(
446        &self,
447        table_id: TableId,
448        cascade: bool,
449    ) -> Result<WaitVersion> {
450        let request = DropMaterializedViewRequest { table_id, cascade };
451
452        let resp = self.inner.drop_materialized_view(request).await?;
453        Ok(resp
454            .version
455            .ok_or_else(|| anyhow!("wait version not set"))?)
456    }
457
458    pub async fn create_source(
459        &self,
460        source: PbSource,
461        graph: Option<StreamFragmentGraph>,
462        if_not_exists: bool,
463    ) -> Result<WaitVersion> {
464        let request = CreateSourceRequest {
465            source: Some(source),
466            fragment_graph: graph,
467            if_not_exists,
468        };
469
470        let resp = self.inner.create_source(request).await?;
471        Ok(resp
472            .version
473            .ok_or_else(|| anyhow!("wait version not set"))?)
474    }
475
476    pub async fn create_sink(
477        &self,
478        sink: PbSink,
479        graph: StreamFragmentGraph,
480        dependencies: HashSet<ObjectId>,
481        resource_type: streaming_job_resource_type::ResourceType,
482        if_not_exists: bool,
483        since_timestamp_epoch: Option<u64>,
484    ) -> Result<WaitVersion> {
485        let request = CreateSinkRequest {
486            sink: Some(sink),
487            fragment_graph: Some(graph),
488            dependencies: dependencies.into_iter().collect(),
489            if_not_exists,
490            resource_type: Some(PbStreamingJobResourceType {
491                resource_type: Some(resource_type),
492            }),
493            since_timestamp_epoch,
494        };
495
496        let resp = self.inner.create_sink(request).await?;
497        Ok(resp
498            .version
499            .ok_or_else(|| anyhow!("wait version not set"))?)
500    }
501
502    pub async fn create_subscription(&self, subscription: PbSubscription) -> Result<WaitVersion> {
503        let request = CreateSubscriptionRequest {
504            subscription: Some(subscription),
505        };
506
507        let resp = self.inner.create_subscription(request).await?;
508        Ok(resp
509            .version
510            .ok_or_else(|| anyhow!("wait version not set"))?)
511    }
512
513    pub async fn create_function(&self, function: PbFunction) -> Result<WaitVersion> {
514        let request = CreateFunctionRequest {
515            function: Some(function),
516        };
517        let resp = self.inner.create_function(request).await?;
518        Ok(resp
519            .version
520            .ok_or_else(|| anyhow!("wait version not set"))?)
521    }
522
523    pub async fn create_table(
524        &self,
525        source: Option<PbSource>,
526        table: PbTable,
527        graph: StreamFragmentGraph,
528        job_type: PbTableJobType,
529        if_not_exists: bool,
530        dependencies: HashSet<ObjectId>,
531    ) -> Result<WaitVersion> {
532        let request = CreateTableRequest {
533            materialized_view: Some(table),
534            fragment_graph: Some(graph),
535            source,
536            job_type: job_type as _,
537            if_not_exists,
538            dependencies: dependencies.into_iter().collect(),
539        };
540        let resp = self.inner.create_table(request).await?;
541        // TODO: handle error in `resp.status` here
542        Ok(resp
543            .version
544            .ok_or_else(|| anyhow!("wait version not set"))?)
545    }
546
547    pub async fn comment_on(&self, comment: PbComment) -> Result<WaitVersion> {
548        let request = CommentOnRequest {
549            comment: Some(comment),
550        };
551        let resp = self.inner.comment_on(request).await?;
552        Ok(resp
553            .version
554            .ok_or_else(|| anyhow!("wait version not set"))?)
555    }
556
557    pub async fn alter_name(
558        &self,
559        object: alter_name_request::Object,
560        name: &str,
561    ) -> Result<WaitVersion> {
562        let request = AlterNameRequest {
563            object: Some(object),
564            new_name: name.to_owned(),
565        };
566        let resp = self.inner.alter_name(request).await?;
567        Ok(resp
568            .version
569            .ok_or_else(|| anyhow!("wait version not set"))?)
570    }
571
572    pub async fn alter_owner(&self, object: Object, owner_id: UserId) -> Result<WaitVersion> {
573        let request = AlterOwnerRequest {
574            object: Some(object),
575            owner_id,
576        };
577        let resp = self.inner.alter_owner(request).await?;
578        Ok(resp
579            .version
580            .ok_or_else(|| anyhow!("wait version not set"))?)
581    }
582
583    pub async fn alter_subscription_retention(
584        &self,
585        subscription_id: SubscriptionId,
586        retention_seconds: u64,
587        definition: String,
588    ) -> Result<WaitVersion> {
589        let request = AlterSubscriptionRetentionRequest {
590            subscription_id,
591            retention_seconds,
592            definition,
593        };
594        let resp = self.inner.alter_subscription_retention(request).await?;
595        Ok(resp
596            .version
597            .ok_or_else(|| anyhow!("wait version not set"))?)
598    }
599
600    pub async fn alter_database_param(
601        &self,
602        database_id: DatabaseId,
603        param: AlterDatabaseParam,
604    ) -> Result<WaitVersion> {
605        let request = match param {
606            AlterDatabaseParam::BarrierIntervalMs(barrier_interval_ms) => {
607                let barrier_interval_ms = OptionalUint32 {
608                    value: barrier_interval_ms,
609                };
610                AlterDatabaseParamRequest {
611                    database_id,
612                    param: Some(alter_database_param_request::Param::BarrierIntervalMs(
613                        barrier_interval_ms,
614                    )),
615                }
616            }
617            AlterDatabaseParam::CheckpointFrequency(checkpoint_frequency) => {
618                let checkpoint_frequency = OptionalUint64 {
619                    value: checkpoint_frequency,
620                };
621                AlterDatabaseParamRequest {
622                    database_id,
623                    param: Some(alter_database_param_request::Param::CheckpointFrequency(
624                        checkpoint_frequency,
625                    )),
626                }
627            }
628        };
629        let resp = self.inner.alter_database_param(request).await?;
630        Ok(resp
631            .version
632            .ok_or_else(|| anyhow!("wait version not set"))?)
633    }
634
635    pub async fn alter_set_schema(
636        &self,
637        object: alter_set_schema_request::Object,
638        new_schema_id: SchemaId,
639    ) -> Result<WaitVersion> {
640        let request = AlterSetSchemaRequest {
641            new_schema_id,
642            object: Some(object),
643        };
644        let resp = self.inner.alter_set_schema(request).await?;
645        Ok(resp
646            .version
647            .ok_or_else(|| anyhow!("wait version not set"))?)
648    }
649
650    pub async fn alter_source(&self, source: PbSource) -> Result<WaitVersion> {
651        let request = AlterSourceRequest {
652            source: Some(source),
653        };
654        let resp = self.inner.alter_source(request).await?;
655        Ok(resp
656            .version
657            .ok_or_else(|| anyhow!("wait version not set"))?)
658    }
659
660    pub async fn alter_parallelism(
661        &self,
662        job_id: JobId,
663        parallelism: PbTableParallelism,
664        adaptive_parallelism_strategy: Option<AdaptiveParallelismStrategy>,
665        deferred: bool,
666    ) -> Result<()> {
667        let request = AlterParallelismRequest {
668            table_id: job_id,
669            parallelism: Some(parallelism),
670            deferred,
671            adaptive_parallelism_strategy: adaptive_parallelism_strategy
672                .map(|strategy| strategy.to_string()),
673        };
674
675        self.inner.alter_parallelism(request).await?;
676        Ok(())
677    }
678
679    pub async fn alter_backfill_parallelism(
680        &self,
681        job_id: JobId,
682        parallelism: Option<PbTableParallelism>,
683        adaptive_parallelism_strategy: Option<AdaptiveParallelismStrategy>,
684        deferred: bool,
685    ) -> Result<()> {
686        let request = AlterBackfillParallelismRequest {
687            table_id: job_id,
688            parallelism,
689            deferred,
690            adaptive_parallelism_strategy: adaptive_parallelism_strategy
691                .map(|strategy| strategy.to_string()),
692        };
693
694        self.inner.alter_backfill_parallelism(request).await?;
695        Ok(())
696    }
697
698    pub async fn alter_streaming_job_config(
699        &self,
700        job_id: JobId,
701        entries_to_add: HashMap<String, String>,
702        keys_to_remove: Vec<String>,
703    ) -> Result<()> {
704        let request = AlterStreamingJobConfigRequest {
705            job_id,
706            entries_to_add,
707            keys_to_remove,
708        };
709
710        self.inner.alter_streaming_job_config(request).await?;
711        Ok(())
712    }
713
714    pub async fn alter_fragment_parallelism(
715        &self,
716        fragment_ids: Vec<FragmentId>,
717        parallelism: Option<PbTableParallelism>,
718    ) -> Result<()> {
719        let request = AlterFragmentParallelismRequest {
720            fragment_ids,
721            parallelism,
722        };
723
724        self.inner.alter_fragment_parallelism(request).await?;
725        Ok(())
726    }
727
728    pub async fn alter_cdc_table_backfill_parallelism(
729        &self,
730        table_id: JobId,
731        parallelism: PbTableParallelism,
732    ) -> Result<()> {
733        let request = AlterCdcTableBackfillParallelismRequest {
734            table_id,
735            parallelism: Some(parallelism),
736        };
737        self.inner
738            .alter_cdc_table_backfill_parallelism(request)
739            .await?;
740        Ok(())
741    }
742
743    pub async fn alter_resource_group(
744        &self,
745        job_id: JobId,
746        resource_group: Option<String>,
747        deferred: bool,
748    ) -> Result<()> {
749        let request = AlterResourceGroupRequest {
750            job_id,
751            resource_group,
752            deferred,
753        };
754
755        self.inner.alter_resource_group(request).await?;
756        Ok(())
757    }
758
759    pub async fn alter_database_resource_group(
760        &self,
761        database_id: DatabaseId,
762        resource_group: Option<String>,
763        deferred: bool,
764    ) -> Result<WaitVersion> {
765        let request = AlterDatabaseResourceGroupRequest {
766            database_id,
767            resource_group,
768            deferred,
769        };
770
771        let resp = self.inner.alter_database_resource_group(request).await?;
772        Ok(resp
773            .version
774            .ok_or_else(|| anyhow!("wait version not set"))?)
775    }
776
777    pub async fn alter_swap_rename(
778        &self,
779        object: alter_swap_rename_request::Object,
780    ) -> Result<WaitVersion> {
781        let request = AlterSwapRenameRequest {
782            object: Some(object),
783        };
784        let resp = self.inner.alter_swap_rename(request).await?;
785        Ok(resp
786            .version
787            .ok_or_else(|| anyhow!("wait version not set"))?)
788    }
789
790    pub async fn alter_secret(
791        &self,
792        secret_id: SecretId,
793        secret_name: String,
794        database_id: DatabaseId,
795        schema_id: SchemaId,
796        owner_id: UserId,
797        value: Vec<u8>,
798    ) -> Result<WaitVersion> {
799        let request = AlterSecretRequest {
800            secret_id,
801            name: secret_name,
802            database_id,
803            schema_id,
804            owner_id,
805            value,
806        };
807        let resp = self.inner.alter_secret(request).await?;
808        Ok(resp
809            .version
810            .ok_or_else(|| anyhow!("wait version not set"))?)
811    }
812
813    pub async fn replace_job(
814        &self,
815        graph: StreamFragmentGraph,
816        replace_job: ReplaceJob,
817    ) -> Result<WaitVersion> {
818        let request = ReplaceJobPlanRequest {
819            plan: Some(ReplaceJobPlan {
820                fragment_graph: Some(graph),
821                replace_job: Some(replace_job),
822            }),
823        };
824        let resp = self.inner.replace_job_plan(request).await?;
825        // TODO: handle error in `resp.status` here
826        Ok(resp
827            .version
828            .ok_or_else(|| anyhow!("wait version not set"))?)
829    }
830
831    pub async fn auto_schema_change(&self, schema_change: SchemaChangeEnvelope) -> Result<()> {
832        let request = AutoSchemaChangeRequest {
833            schema_change: Some(schema_change),
834        };
835        let _ = self.inner.auto_schema_change(request).await?;
836        Ok(())
837    }
838
839    /// Block until the pk-index sink `sink_id`'s database has committed through `epoch`, returning
840    /// the coordinator's committed iceberg snapshot id (`None` if no snapshot committed yet).
841    pub async fn wait_iceberg_pk_index_sink_epoch(
842        &self,
843        sink_id: SinkId,
844        epoch: u64,
845    ) -> Result<Option<i64>> {
846        let request = WaitIcebergPkIndexSinkEpochRequest { sink_id, epoch };
847        let resp = self.inner.wait_iceberg_pk_index_sink_epoch(request).await?;
848        Ok(resp.snapshot_id)
849    }
850
851    pub async fn create_view(
852        &self,
853        view: PbView,
854        dependencies: HashSet<ObjectId>,
855    ) -> Result<WaitVersion> {
856        let request = CreateViewRequest {
857            view: Some(view),
858            dependencies: dependencies.into_iter().collect(),
859        };
860        let resp = self.inner.create_view(request).await?;
861        // TODO: handle error in `resp.status` here
862        Ok(resp
863            .version
864            .ok_or_else(|| anyhow!("wait version not set"))?)
865    }
866
867    pub async fn create_index(
868        &self,
869        index: PbIndex,
870        table: PbTable,
871        graph: StreamFragmentGraph,
872        resource_type: streaming_job_resource_type::ResourceType,
873        if_not_exists: bool,
874    ) -> Result<WaitVersion> {
875        let request = CreateIndexRequest {
876            index: Some(index),
877            index_table: Some(table),
878            fragment_graph: Some(graph),
879            if_not_exists,
880            resource_type: Some(PbStreamingJobResourceType {
881                resource_type: Some(resource_type),
882            }),
883        };
884        let resp = self.inner.create_index(request).await?;
885        // TODO: handle error in `resp.status` here
886        Ok(resp
887            .version
888            .ok_or_else(|| anyhow!("wait version not set"))?)
889    }
890
891    pub async fn drop_table(
892        &self,
893        source_id: Option<SourceId>,
894        table_id: TableId,
895        cascade: bool,
896    ) -> Result<WaitVersion> {
897        let request = DropTableRequest {
898            source_id: source_id.map(risingwave_pb::ddl_service::drop_table_request::SourceId::Id),
899            table_id,
900            cascade,
901        };
902
903        let resp = self.inner.drop_table(request).await?;
904        Ok(resp
905            .version
906            .ok_or_else(|| anyhow!("wait version not set"))?)
907    }
908
909    pub async fn compact_iceberg_table(&self, sink_id: SinkId) -> Result<IcebergCompactionTaskId> {
910        let request = CompactIcebergTableRequest { sink_id };
911        let resp = self.inner.compact_iceberg_table(request).await?;
912        Ok(resp.task_id)
913    }
914
915    pub async fn rewrite_iceberg_table_manifests(&self, sink_id: SinkId) -> Result<()> {
916        let request = RewriteIcebergTableManifestsRequest { sink_id };
917        let _resp = self.inner.rewrite_iceberg_table_manifests(request).await?;
918        Ok(())
919    }
920
921    pub async fn expire_iceberg_table_snapshots(&self, sink_id: SinkId) -> Result<()> {
922        let request = ExpireIcebergTableSnapshotsRequest { sink_id };
923        let _resp = self.inner.expire_iceberg_table_snapshots(request).await?;
924        Ok(())
925    }
926
927    pub async fn drop_view(&self, view_id: ViewId, cascade: bool) -> Result<WaitVersion> {
928        let request = DropViewRequest { view_id, cascade };
929        let resp = self.inner.drop_view(request).await?;
930        Ok(resp
931            .version
932            .ok_or_else(|| anyhow!("wait version not set"))?)
933    }
934
935    pub async fn drop_source(&self, source_id: SourceId, cascade: bool) -> Result<WaitVersion> {
936        let request = DropSourceRequest { source_id, cascade };
937        let resp = self.inner.drop_source(request).await?;
938        Ok(resp
939            .version
940            .ok_or_else(|| anyhow!("wait version not set"))?)
941    }
942
943    pub async fn reset_source(&self, source_id: SourceId) -> Result<WaitVersion> {
944        let request = ResetSourceRequest { source_id };
945        let resp = self.inner.reset_source(request).await?;
946        Ok(resp
947            .version
948            .ok_or_else(|| anyhow!("wait version not set"))?)
949    }
950
951    pub async fn drop_sink(&self, sink_id: SinkId, cascade: bool) -> Result<WaitVersion> {
952        let request = DropSinkRequest { sink_id, cascade };
953        let resp = self.inner.drop_sink(request).await?;
954        Ok(resp
955            .version
956            .ok_or_else(|| anyhow!("wait version not set"))?)
957    }
958
959    pub async fn drop_subscription(
960        &self,
961        subscription_id: SubscriptionId,
962        cascade: bool,
963    ) -> Result<WaitVersion> {
964        let request = DropSubscriptionRequest {
965            subscription_id,
966            cascade,
967        };
968        let resp = self.inner.drop_subscription(request).await?;
969        Ok(resp
970            .version
971            .ok_or_else(|| anyhow!("wait version not set"))?)
972    }
973
974    pub async fn drop_index(&self, index_id: IndexId, cascade: bool) -> Result<WaitVersion> {
975        let request = DropIndexRequest { index_id, cascade };
976        let resp = self.inner.drop_index(request).await?;
977        Ok(resp
978            .version
979            .ok_or_else(|| anyhow!("wait version not set"))?)
980    }
981
982    pub async fn drop_function(
983        &self,
984        function_id: FunctionId,
985        cascade: bool,
986    ) -> Result<WaitVersion> {
987        let request = DropFunctionRequest {
988            function_id,
989            cascade,
990        };
991        let resp = self.inner.drop_function(request).await?;
992        Ok(resp
993            .version
994            .ok_or_else(|| anyhow!("wait version not set"))?)
995    }
996
997    pub async fn drop_database(&self, database_id: DatabaseId) -> Result<WaitVersion> {
998        let request = DropDatabaseRequest { database_id };
999        let resp = self.inner.drop_database(request).await?;
1000        Ok(resp
1001            .version
1002            .ok_or_else(|| anyhow!("wait version not set"))?)
1003    }
1004
1005    pub async fn drop_schema(&self, schema_id: SchemaId, cascade: bool) -> Result<WaitVersion> {
1006        let request = DropSchemaRequest { schema_id, cascade };
1007        let resp = self.inner.drop_schema(request).await?;
1008        Ok(resp
1009            .version
1010            .ok_or_else(|| anyhow!("wait version not set"))?)
1011    }
1012
1013    // TODO: using UserInfoVersion instead as return type.
1014    pub async fn create_user(&self, user: UserInfo) -> Result<u64> {
1015        let request = CreateUserRequest { user: Some(user) };
1016        let resp = self.inner.create_user(request).await?;
1017        Ok(resp.version)
1018    }
1019
1020    pub async fn drop_user(&self, user_id: UserId) -> Result<u64> {
1021        let request = DropUserRequest { user_id };
1022        let resp = self.inner.drop_user(request).await?;
1023        Ok(resp.version)
1024    }
1025
1026    pub async fn update_user(
1027        &self,
1028        user: UserInfo,
1029        update_fields: Vec<UpdateField>,
1030    ) -> Result<u64> {
1031        let request = UpdateUserRequest {
1032            user: Some(user),
1033            update_fields: update_fields
1034                .into_iter()
1035                .map(|field| field as i32)
1036                .collect::<Vec<_>>(),
1037        };
1038        let resp = self.inner.update_user(request).await?;
1039        Ok(resp.version)
1040    }
1041
1042    pub async fn grant_privilege(
1043        &self,
1044        user_ids: Vec<UserId>,
1045        privileges: Vec<GrantPrivilege>,
1046        with_grant_option: bool,
1047        granted_by: UserId,
1048    ) -> Result<u64> {
1049        let request = GrantPrivilegeRequest {
1050            user_ids,
1051            privileges,
1052            with_grant_option,
1053            granted_by,
1054        };
1055        let resp = self.inner.grant_privilege(request).await?;
1056        Ok(resp.version)
1057    }
1058
1059    pub async fn revoke_privilege(
1060        &self,
1061        user_ids: Vec<UserId>,
1062        privileges: Vec<GrantPrivilege>,
1063        granted_by: UserId,
1064        revoke_by: UserId,
1065        revoke_grant_option: bool,
1066        cascade: bool,
1067    ) -> Result<u64> {
1068        let request = RevokePrivilegeRequest {
1069            user_ids,
1070            privileges,
1071            granted_by,
1072            revoke_by,
1073            revoke_grant_option,
1074            cascade,
1075        };
1076        let resp = self.inner.revoke_privilege(request).await?;
1077        Ok(resp.version)
1078    }
1079
1080    pub async fn alter_default_privilege(
1081        &self,
1082        user_ids: Vec<UserId>,
1083        database_id: DatabaseId,
1084        schema_ids: Vec<SchemaId>,
1085        operation: AlterDefaultPrivilegeOperation,
1086        granted_by: UserId,
1087    ) -> Result<()> {
1088        let request = AlterDefaultPrivilegeRequest {
1089            user_ids,
1090            database_id,
1091            schema_ids,
1092            operation: Some(operation),
1093            granted_by,
1094        };
1095        self.inner.alter_default_privilege(request).await?;
1096        Ok(())
1097    }
1098
1099    /// Unregister the current node from the cluster.
1100    pub async fn unregister(&self) -> Result<()> {
1101        let request = DeleteWorkerNodeRequest {
1102            host: Some(self.host_addr.to_protobuf()),
1103        };
1104        self.inner.delete_worker_node(request).await?;
1105        self.shutting_down.store(true, Relaxed);
1106        Ok(())
1107    }
1108
1109    /// Try to unregister the current worker from the cluster with best effort. Log the result.
1110    pub async fn try_unregister(&self) {
1111        match self.unregister().await {
1112            Ok(_) => {
1113                tracing::info!(
1114                    worker_id = %self.worker_id(),
1115                    "successfully unregistered from meta service",
1116                )
1117            }
1118            Err(e) => {
1119                tracing::warn!(
1120                    error = %e.as_report(),
1121                    worker_id = %self.worker_id(),
1122                    "failed to unregister from meta service",
1123                );
1124            }
1125        }
1126    }
1127
1128    pub async fn list_worker_nodes(
1129        &self,
1130        worker_type: Option<WorkerType>,
1131    ) -> Result<Vec<WorkerNode>> {
1132        let request = ListAllNodesRequest {
1133            worker_type: worker_type.map(Into::into),
1134            include_starting_nodes: true,
1135        };
1136        let resp = self.inner.list_all_nodes(request).await?;
1137        Ok(resp.nodes)
1138    }
1139
1140    /// Starts a heartbeat worker.
1141    pub fn start_heartbeat_loop(
1142        meta_client: MetaClient,
1143        min_interval: Duration,
1144    ) -> (JoinHandle<()>, Sender<()>) {
1145        let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
1146        let join_handle = tokio::spawn(async move {
1147            let mut min_interval_ticker = tokio::time::interval(min_interval);
1148            loop {
1149                tokio::select! {
1150                    biased;
1151                    // Shutdown
1152                    _ = &mut shutdown_rx => {
1153                        tracing::info!("Heartbeat loop is stopped");
1154                        return;
1155                    }
1156                    // Wait for interval
1157                    _ = min_interval_ticker.tick() => {},
1158                }
1159                tracing::debug!(target: "events::meta::client_heartbeat", "heartbeat");
1160                match tokio::time::timeout(
1161                    // TODO: decide better min_interval for timeout
1162                    min_interval * 3,
1163                    meta_client.send_heartbeat(),
1164                )
1165                .await
1166                {
1167                    Ok(Ok(_)) => {}
1168                    Ok(Err(err)) => {
1169                        tracing::warn!(error = %err.as_report(), "Failed to send_heartbeat");
1170                    }
1171                    Err(_) => {
1172                        tracing::warn!("Failed to send_heartbeat: timeout");
1173                    }
1174                }
1175            }
1176        });
1177        (join_handle, shutdown_tx)
1178    }
1179
1180    pub async fn risectl_list_state_tables(&self) -> Result<Vec<PbTable>> {
1181        let request = RisectlListStateTablesRequest {};
1182        let resp = self.inner.risectl_list_state_tables(request).await?;
1183        Ok(resp.tables)
1184    }
1185
1186    pub async fn risectl_resume_backfill(
1187        &self,
1188        request: RisectlResumeBackfillRequest,
1189    ) -> Result<()> {
1190        self.inner.risectl_resume_backfill(request).await?;
1191        Ok(())
1192    }
1193
1194    pub async fn flush(&self, database_id: DatabaseId) -> Result<HummockVersionId> {
1195        let request = FlushRequest { database_id };
1196        let resp = self.inner.flush(request).await?;
1197        Ok(resp.hummock_version_id)
1198    }
1199
1200    pub async fn wait(&self, job_id: Option<JobId>) -> Result<WaitVersion> {
1201        let request = WaitRequest { job_id };
1202        let resp = self.inner.wait(request).await?;
1203        Ok(resp
1204            .version
1205            .ok_or_else(|| anyhow!("wait version not set"))?)
1206    }
1207
1208    pub async fn recover(&self) -> Result<()> {
1209        let request = RecoverRequest {};
1210        self.inner.recover(request).await?;
1211        Ok(())
1212    }
1213
1214    pub async fn cancel_creating_jobs(&self, jobs: PbJobs) -> Result<Vec<u32>> {
1215        let request = CancelCreatingJobsRequest { jobs: Some(jobs) };
1216        let resp = self.inner.cancel_creating_jobs(request).await?;
1217        Ok(resp.canceled_jobs)
1218    }
1219
1220    pub async fn list_table_fragments(
1221        &self,
1222        job_ids: &[JobId],
1223    ) -> Result<HashMap<JobId, TableFragmentInfo>> {
1224        let request = ListTableFragmentsRequest {
1225            table_ids: job_ids.to_vec(),
1226        };
1227        let resp = self.inner.list_table_fragments(request).await?;
1228        Ok(resp.table_fragments)
1229    }
1230
1231    pub async fn list_streaming_job_states(&self) -> Result<Vec<StreamingJobState>> {
1232        let resp = self
1233            .inner
1234            .list_streaming_job_states(ListStreamingJobStatesRequest {})
1235            .await?;
1236        Ok(resp.states)
1237    }
1238
1239    pub async fn list_fragment_distributions(
1240        &self,
1241        include_node: bool,
1242    ) -> Result<Vec<FragmentDistribution>> {
1243        let resp = self
1244            .inner
1245            .list_fragment_distribution(ListFragmentDistributionRequest {
1246                include_node: Some(include_node),
1247            })
1248            .await?;
1249        Ok(resp.distributions)
1250    }
1251
1252    pub async fn list_creating_fragment_distribution(&self) -> Result<Vec<FragmentDistribution>> {
1253        let resp = self
1254            .inner
1255            .list_creating_fragment_distribution(ListCreatingFragmentDistributionRequest {
1256                include_node: Some(true),
1257            })
1258            .await?;
1259        Ok(resp.distributions)
1260    }
1261
1262    pub async fn get_fragment_by_id(
1263        &self,
1264        fragment_id: FragmentId,
1265    ) -> Result<Option<FragmentDistribution>> {
1266        let resp = self
1267            .inner
1268            .get_fragment_by_id(GetFragmentByIdRequest { fragment_id })
1269            .await?;
1270        Ok(resp.distribution)
1271    }
1272
1273    pub async fn get_fragment_vnodes(
1274        &self,
1275        fragment_id: FragmentId,
1276    ) -> Result<Vec<(ActorId, Vec<u32>)>> {
1277        let resp = self
1278            .inner
1279            .get_fragment_vnodes(GetFragmentVnodesRequest { fragment_id })
1280            .await?;
1281        Ok(resp
1282            .actor_vnodes
1283            .into_iter()
1284            .map(|actor| (actor.actor_id, actor.vnode_indices))
1285            .collect())
1286    }
1287
1288    pub async fn get_actor_vnodes(&self, actor_id: ActorId) -> Result<Vec<u32>> {
1289        let resp = self
1290            .inner
1291            .get_actor_vnodes(GetActorVnodesRequest { actor_id })
1292            .await?;
1293        Ok(resp.vnode_indices)
1294    }
1295
1296    pub async fn list_actor_states(&self) -> Result<Vec<ActorState>> {
1297        let resp = self
1298            .inner
1299            .list_actor_states(ListActorStatesRequest {})
1300            .await?;
1301        Ok(resp.states)
1302    }
1303
1304    pub async fn list_actor_splits(&self) -> Result<Vec<ActorSplit>> {
1305        let resp = self
1306            .inner
1307            .list_actor_splits(ListActorSplitsRequest {})
1308            .await?;
1309
1310        Ok(resp.actor_splits)
1311    }
1312
1313    pub async fn pause(&self) -> Result<PauseResponse> {
1314        let request = PauseRequest {};
1315        let resp = self.inner.pause(request).await?;
1316        Ok(resp)
1317    }
1318
1319    pub async fn resume(&self) -> Result<ResumeResponse> {
1320        let request = ResumeRequest {};
1321        let resp = self.inner.resume(request).await?;
1322        Ok(resp)
1323    }
1324
1325    pub async fn apply_throttle(
1326        &self,
1327        throttle_target: PbThrottleTarget,
1328        throttle_type: risingwave_pb::common::PbThrottleType,
1329        id: u32,
1330        rate: Option<u32>,
1331    ) -> Result<ApplyThrottleResponse> {
1332        let request = ApplyThrottleRequest {
1333            throttle_target: throttle_target as i32,
1334            throttle_type: throttle_type as i32,
1335            id,
1336            rate,
1337        };
1338        let resp = self.inner.apply_throttle(request).await?;
1339        Ok(resp)
1340    }
1341
1342    pub async fn get_cluster_recovery_status(&self) -> Result<RecoveryStatus> {
1343        let resp = self
1344            .inner
1345            .get_cluster_recovery_status(GetClusterRecoveryStatusRequest {})
1346            .await?;
1347        Ok(resp.get_status().unwrap())
1348    }
1349
1350    pub async fn get_cluster_info(&self) -> Result<GetClusterInfoResponse> {
1351        let request = GetClusterInfoRequest {};
1352        let resp = self.inner.get_cluster_info(request).await?;
1353        Ok(resp)
1354    }
1355
1356    pub async fn reschedule(
1357        &self,
1358        worker_reschedules: HashMap<u32, PbWorkerReschedule>,
1359        revision: u64,
1360        resolve_no_shuffle_upstream: bool,
1361    ) -> Result<(bool, u64)> {
1362        let request = RescheduleRequest {
1363            revision,
1364            resolve_no_shuffle_upstream,
1365            worker_reschedules,
1366        };
1367        let resp = self.inner.reschedule(request).await?;
1368        Ok((resp.success, resp.revision))
1369    }
1370
1371    pub async fn risectl_get_pinned_versions_summary(
1372        &self,
1373    ) -> Result<RiseCtlGetPinnedVersionsSummaryResponse> {
1374        let request = RiseCtlGetPinnedVersionsSummaryRequest {};
1375        self.inner
1376            .rise_ctl_get_pinned_versions_summary(request)
1377            .await
1378    }
1379
1380    pub async fn risectl_get_checkpoint_hummock_version(
1381        &self,
1382    ) -> Result<RiseCtlGetCheckpointVersionResponse> {
1383        let request = RiseCtlGetCheckpointVersionRequest {};
1384        self.inner.rise_ctl_get_checkpoint_version(request).await
1385    }
1386
1387    pub async fn risectl_pause_hummock_version_checkpoint(
1388        &self,
1389    ) -> Result<RiseCtlPauseVersionCheckpointResponse> {
1390        let request = RiseCtlPauseVersionCheckpointRequest {};
1391        self.inner.rise_ctl_pause_version_checkpoint(request).await
1392    }
1393
1394    pub async fn risectl_resume_hummock_version_checkpoint(
1395        &self,
1396    ) -> Result<RiseCtlResumeVersionCheckpointResponse> {
1397        let request = RiseCtlResumeVersionCheckpointRequest {};
1398        self.inner.rise_ctl_resume_version_checkpoint(request).await
1399    }
1400
1401    pub async fn init_metadata_for_replay(
1402        &self,
1403        tables: Vec<PbTable>,
1404        compaction_groups: Vec<CompactionGroupInfo>,
1405    ) -> Result<()> {
1406        let req = InitMetadataForReplayRequest {
1407            tables,
1408            compaction_groups,
1409        };
1410        let _resp = self.inner.init_metadata_for_replay(req).await?;
1411        Ok(())
1412    }
1413
1414    pub async fn replay_version_delta(
1415        &self,
1416        version_delta: HummockVersionDelta,
1417    ) -> Result<(HummockVersion, Vec<CompactionGroupId>)> {
1418        let req = ReplayVersionDeltaRequest {
1419            version_delta: Some(version_delta.into()),
1420        };
1421        let resp = self.inner.replay_version_delta(req).await?;
1422        Ok((
1423            HummockVersion::from_rpc_protobuf(&resp.version.unwrap()),
1424            resp.modified_compaction_groups,
1425        ))
1426    }
1427
1428    pub async fn list_version_deltas(
1429        &self,
1430        start_id: HummockVersionId,
1431        num_limit: u32,
1432        committed_epoch_limit: HummockEpoch,
1433    ) -> Result<Vec<HummockVersionDelta>> {
1434        let req = ListVersionDeltasRequest {
1435            start_id,
1436            num_limit,
1437            committed_epoch_limit,
1438        };
1439        Ok(self
1440            .inner
1441            .list_version_deltas(req)
1442            .await?
1443            .version_deltas
1444            .unwrap()
1445            .version_deltas
1446            .iter()
1447            .map(HummockVersionDelta::from_rpc_protobuf)
1448            .collect())
1449    }
1450
1451    pub async fn trigger_compaction_deterministic(
1452        &self,
1453        version_id: HummockVersionId,
1454        compaction_groups: Vec<CompactionGroupId>,
1455    ) -> Result<()> {
1456        let req = TriggerCompactionDeterministicRequest {
1457            version_id,
1458            compaction_groups,
1459        };
1460        self.inner.trigger_compaction_deterministic(req).await?;
1461        Ok(())
1462    }
1463
1464    pub async fn disable_commit_epoch(&self) -> Result<HummockVersion> {
1465        let req = DisableCommitEpochRequest {};
1466        Ok(HummockVersion::from_rpc_protobuf(
1467            &self
1468                .inner
1469                .disable_commit_epoch(req)
1470                .await?
1471                .current_version
1472                .unwrap(),
1473        ))
1474    }
1475
1476    pub async fn get_assigned_compact_task_num(&self) -> Result<usize> {
1477        let req = GetAssignedCompactTaskNumRequest {};
1478        let resp = self.inner.get_assigned_compact_task_num(req).await?;
1479        Ok(resp.num_tasks as usize)
1480    }
1481
1482    pub async fn risectl_list_compaction_group(&self) -> Result<Vec<CompactionGroupInfo>> {
1483        let req = RiseCtlListCompactionGroupRequest {};
1484        let resp = self.inner.rise_ctl_list_compaction_group(req).await?;
1485        Ok(resp.compaction_groups)
1486    }
1487
1488    pub async fn risectl_update_compaction_config(
1489        &self,
1490        compaction_groups: &[CompactionGroupId],
1491        configs: &[MutableConfig],
1492    ) -> Result<()> {
1493        let req = RiseCtlUpdateCompactionConfigRequest {
1494            compaction_group_ids: compaction_groups.to_vec(),
1495            configs: configs
1496                .iter()
1497                .map(
1498                    |c| rise_ctl_update_compaction_config_request::MutableConfig {
1499                        mutable_config: Some(c.clone()),
1500                    },
1501                )
1502                .collect(),
1503        };
1504        let _resp = self.inner.rise_ctl_update_compaction_config(req).await?;
1505        Ok(())
1506    }
1507
1508    pub async fn backup_meta(&self, remarks: Option<String>) -> Result<u64> {
1509        let req = BackupMetaRequest { remarks };
1510        let resp = self.inner.backup_meta(req).await?;
1511        Ok(resp.job_id)
1512    }
1513
1514    pub async fn get_backup_job_status(&self, job_id: u64) -> Result<(BackupJobStatus, String)> {
1515        let req = GetBackupJobStatusRequest { job_id };
1516        let resp = self.inner.get_backup_job_status(req).await?;
1517        Ok((resp.job_status(), resp.message))
1518    }
1519
1520    pub async fn delete_meta_snapshot(&self, snapshot_ids: &[u64]) -> Result<()> {
1521        let req = DeleteMetaSnapshotRequest {
1522            snapshot_ids: snapshot_ids.to_vec(),
1523        };
1524        let _resp = self.inner.delete_meta_snapshot(req).await?;
1525        Ok(())
1526    }
1527
1528    pub async fn get_meta_snapshot_manifest(&self) -> Result<MetaSnapshotManifest> {
1529        let req = GetMetaSnapshotManifestRequest {};
1530        let resp = self.inner.get_meta_snapshot_manifest(req).await?;
1531        Ok(resp.manifest.expect("should exist"))
1532    }
1533
1534    pub async fn get_telemetry_info(&self) -> Result<TelemetryInfoResponse> {
1535        let req = GetTelemetryInfoRequest {};
1536        let resp = self.inner.get_telemetry_info(req).await?;
1537        Ok(resp)
1538    }
1539
1540    pub async fn get_meta_store_endpoint(&self) -> Result<String> {
1541        let req = GetMetaStoreInfoRequest {};
1542        let resp = self.inner.get_meta_store_info(req).await?;
1543        Ok(resp.meta_store_endpoint)
1544    }
1545
1546    pub async fn alter_sink_props(
1547        &self,
1548        sink_id: SinkId,
1549        changed_props: BTreeMap<String, String>,
1550        changed_secret_refs: BTreeMap<String, PbSecretRef>,
1551        connector_conn_ref: Option<ConnectionId>,
1552    ) -> Result<()> {
1553        let req = AlterConnectorPropsRequest {
1554            object_id: sink_id.as_raw_id(),
1555            changed_props: changed_props.into_iter().collect(),
1556            changed_secret_refs: changed_secret_refs.into_iter().collect(),
1557            connector_conn_ref,
1558            object_type: AlterConnectorPropsObject::Sink as i32,
1559            extra_options: None,
1560        };
1561        let _resp = self.inner.alter_connector_props(req).await?;
1562        Ok(())
1563    }
1564
1565    pub async fn alter_iceberg_table_props(
1566        &self,
1567        table_id: TableId,
1568        sink_id: SinkId,
1569        source_id: SourceId,
1570        changed_props: BTreeMap<String, String>,
1571        changed_secret_refs: BTreeMap<String, PbSecretRef>,
1572        connector_conn_ref: Option<ConnectionId>,
1573    ) -> Result<()> {
1574        let req = AlterConnectorPropsRequest {
1575            object_id: table_id.as_raw_id(),
1576            changed_props: changed_props.into_iter().collect(),
1577            changed_secret_refs: changed_secret_refs.into_iter().collect(),
1578            connector_conn_ref,
1579            object_type: AlterConnectorPropsObject::IcebergTable as i32,
1580            extra_options: Some(ExtraOptions::AlterIcebergTableIds(AlterIcebergTableIds {
1581                sink_id,
1582                source_id,
1583            })),
1584        };
1585        let _resp = self.inner.alter_connector_props(req).await?;
1586        Ok(())
1587    }
1588
1589    pub async fn alter_source_connector_props(
1590        &self,
1591        source_id: SourceId,
1592        changed_props: BTreeMap<String, String>,
1593        changed_secret_refs: BTreeMap<String, PbSecretRef>,
1594        connector_conn_ref: Option<ConnectionId>,
1595    ) -> Result<()> {
1596        let req = AlterConnectorPropsRequest {
1597            object_id: source_id.as_raw_id(),
1598            changed_props: changed_props.into_iter().collect(),
1599            changed_secret_refs: changed_secret_refs.into_iter().collect(),
1600            connector_conn_ref,
1601            object_type: AlterConnectorPropsObject::Source as i32,
1602            extra_options: None,
1603        };
1604        let _resp = self.inner.alter_connector_props(req).await?;
1605        Ok(())
1606    }
1607
1608    pub async fn alter_connection_connector_props(
1609        &self,
1610        connection_id: u32,
1611        changed_props: BTreeMap<String, String>,
1612        changed_secret_refs: BTreeMap<String, PbSecretRef>,
1613    ) -> Result<()> {
1614        let req = AlterConnectorPropsRequest {
1615            object_id: connection_id,
1616            changed_props: changed_props.into_iter().collect(),
1617            changed_secret_refs: changed_secret_refs.into_iter().collect(),
1618            connector_conn_ref: None, // Connections don't reference other connections
1619            object_type: AlterConnectorPropsObject::Connection as i32,
1620            extra_options: None,
1621        };
1622        let _resp = self.inner.alter_connector_props(req).await?;
1623        Ok(())
1624    }
1625
1626    /// Orchestrated source property update with pause/update/resume workflow.
1627    /// This is the "safe" version that pauses sources before updating and resumes after.
1628    pub async fn alter_source_properties_safe(
1629        &self,
1630        source_id: SourceId,
1631        changed_props: BTreeMap<String, String>,
1632        changed_secret_refs: BTreeMap<String, PbSecretRef>,
1633        reset_splits: bool,
1634    ) -> Result<()> {
1635        let req = AlterSourcePropertiesSafeRequest {
1636            source_id: source_id.as_raw_id(),
1637            changed_props: changed_props.into_iter().collect(),
1638            changed_secret_refs: changed_secret_refs.into_iter().collect(),
1639            options: Some(PropertyUpdateOptions { reset_splits }),
1640        };
1641        let _resp = self.inner.alter_source_properties_safe(req).await?;
1642        Ok(())
1643    }
1644
1645    /// Reset source split assignments (UNSAFE - admin only).
1646    /// This clears persisted split metadata and triggers re-discovery.
1647    pub async fn reset_source_splits(&self, source_id: SourceId) -> Result<()> {
1648        let req = ResetSourceSplitsRequest {
1649            source_id: source_id.as_raw_id(),
1650        };
1651        let _resp = self.inner.reset_source_splits(req).await?;
1652        Ok(())
1653    }
1654
1655    /// Inject specific offsets into source splits (UNSAFE - admin only).
1656    /// This can cause data duplication or loss depending on the correctness of the provided offsets.
1657    pub async fn inject_source_offsets(
1658        &self,
1659        source_id: SourceId,
1660        split_offsets: HashMap<String, String>,
1661    ) -> Result<Vec<String>> {
1662        let req = InjectSourceOffsetsRequest {
1663            source_id: source_id.as_raw_id(),
1664            split_offsets,
1665        };
1666        let resp = self.inner.inject_source_offsets(req).await?;
1667        Ok(resp.applied_split_ids)
1668    }
1669
1670    pub async fn set_system_param(
1671        &self,
1672        param: String,
1673        value: Option<String>,
1674    ) -> Result<Option<SystemParamsReader>> {
1675        let req = SetSystemParamRequest { param, value };
1676        let resp = self.inner.set_system_param(req).await?;
1677        Ok(resp.params.map(SystemParamsReader::from))
1678    }
1679
1680    pub async fn clear_file_cache(
1681        &self,
1682        clear_meta_cache: bool,
1683        clear_data_cache: bool,
1684    ) -> Result<()> {
1685        self.inner
1686            .clear_file_cache(ClearFileCacheRequest {
1687                clear_meta_cache,
1688                clear_data_cache,
1689            })
1690            .await?;
1691        Ok(())
1692    }
1693
1694    pub async fn get_session_params(&self) -> Result<String> {
1695        let req = GetSessionParamsRequest {};
1696        let resp = self.inner.get_session_params(req).await?;
1697        Ok(resp.params)
1698    }
1699
1700    pub async fn set_session_param(&self, param: String, value: Option<String>) -> Result<String> {
1701        let req = SetSessionParamRequest { param, value };
1702        let resp = self.inner.set_session_param(req).await?;
1703        Ok(resp.param)
1704    }
1705
1706    pub async fn get_ddl_progress(&self) -> Result<Vec<DdlProgress>> {
1707        let req = GetDdlProgressRequest {};
1708        let resp = self.inner.get_ddl_progress(req).await?;
1709        Ok(resp.ddl_progress)
1710    }
1711
1712    pub async fn split_compaction_group(
1713        &self,
1714        group_id: CompactionGroupId,
1715        table_ids_to_new_group: &[TableId],
1716        partition_vnode_count: u32,
1717    ) -> Result<CompactionGroupId> {
1718        let req = SplitCompactionGroupRequest {
1719            group_id,
1720            table_ids: table_ids_to_new_group.to_vec(),
1721            partition_vnode_count,
1722        };
1723        let resp = self.inner.split_compaction_group(req).await?;
1724        Ok(resp.new_group_id)
1725    }
1726
1727    pub async fn get_tables(
1728        &self,
1729        table_ids: Vec<TableId>,
1730        include_dropped_tables: bool,
1731    ) -> Result<HashMap<TableId, Table>> {
1732        let req = GetTablesRequest {
1733            table_ids,
1734            include_dropped_tables,
1735        };
1736        let resp = self.inner.get_tables(req).await?;
1737        Ok(resp.tables)
1738    }
1739
1740    pub async fn list_serving_vnode_mappings(
1741        &self,
1742    ) -> Result<HashMap<FragmentId, (JobId, WorkerSlotMapping)>> {
1743        let req = GetServingVnodeMappingsRequest {};
1744        let resp = self.inner.get_serving_vnode_mappings(req).await?;
1745        let mappings = resp
1746            .worker_slot_mappings
1747            .into_iter()
1748            .map(|p| {
1749                (
1750                    p.fragment_id,
1751                    (
1752                        resp.fragment_to_table
1753                            .get(&p.fragment_id)
1754                            .cloned()
1755                            .unwrap_or_default(),
1756                        WorkerSlotMapping::from_protobuf(p.mapping.as_ref().unwrap()),
1757                    ),
1758                )
1759            })
1760            .collect();
1761        Ok(mappings)
1762    }
1763
1764    pub async fn risectl_list_compaction_status(
1765        &self,
1766    ) -> Result<(
1767        Vec<CompactStatus>,
1768        Vec<CompactTaskAssignment>,
1769        Vec<CompactTaskProgress>,
1770    )> {
1771        let req = RiseCtlListCompactionStatusRequest {};
1772        let resp = self.inner.rise_ctl_list_compaction_status(req).await?;
1773        Ok((
1774            resp.compaction_statuses,
1775            resp.task_assignment,
1776            resp.task_progress,
1777        ))
1778    }
1779
1780    pub async fn get_compaction_score(
1781        &self,
1782        compaction_group_id: CompactionGroupId,
1783    ) -> Result<Vec<PickerInfo>> {
1784        let req = GetCompactionScoreRequest {
1785            compaction_group_id,
1786        };
1787        let resp = self.inner.get_compaction_score(req).await?;
1788        Ok(resp.scores)
1789    }
1790
1791    pub async fn risectl_rebuild_table_stats(&self) -> Result<()> {
1792        let req = RiseCtlRebuildTableStatsRequest {};
1793        let _resp = self.inner.rise_ctl_rebuild_table_stats(req).await?;
1794        Ok(())
1795    }
1796
1797    pub async fn list_branched_object(&self) -> Result<Vec<BranchedObject>> {
1798        let req = ListBranchedObjectRequest {};
1799        let resp = self.inner.list_branched_object(req).await?;
1800        Ok(resp.branched_objects)
1801    }
1802
1803    pub async fn list_active_write_limit(&self) -> Result<HashMap<CompactionGroupId, WriteLimit>> {
1804        let req = ListActiveWriteLimitRequest {};
1805        let resp = self.inner.list_active_write_limit(req).await?;
1806        Ok(resp.write_limits)
1807    }
1808
1809    pub async fn list_hummock_meta_config(&self) -> Result<HashMap<String, String>> {
1810        let req = ListHummockMetaConfigRequest {};
1811        let resp = self.inner.list_hummock_meta_config(req).await?;
1812        Ok(resp.configs)
1813    }
1814
1815    pub async fn get_table_change_logs(
1816        &self,
1817        epoch_only: bool,
1818        start_epoch_inclusive: Option<u64>,
1819        end_epoch_inclusive: Option<u64>,
1820        table_ids: Option<HashSet<TableId>>,
1821        exclude_empty: bool,
1822        limit: Option<u32>,
1823    ) -> Result<TableChangeLogs> {
1824        let req = GetTableChangeLogsRequest {
1825            epoch_only,
1826            start_epoch_inclusive,
1827            end_epoch_inclusive,
1828            table_ids: table_ids.map(|iter| PbTableFilter {
1829                table_ids: iter.into_iter().collect::<Vec<_>>(),
1830            }),
1831            exclude_empty,
1832            limit,
1833        };
1834        let resp = self.inner.get_table_change_logs(req).await?;
1835        Ok(resp
1836            .table_change_logs
1837            .into_iter()
1838            .map(|(id, change_log)| (TableId::new(id), TableChangeLog::from_protobuf(&change_log)))
1839            .collect())
1840    }
1841
1842    pub async fn delete_worker_node(&self, worker: HostAddress) -> Result<()> {
1843        let _resp = self
1844            .inner
1845            .delete_worker_node(DeleteWorkerNodeRequest { host: Some(worker) })
1846            .await?;
1847
1848        Ok(())
1849    }
1850
1851    pub async fn rw_cloud_validate_source(
1852        &self,
1853        source_type: SourceType,
1854        source_config: HashMap<String, String>,
1855    ) -> Result<RwCloudValidateSourceResponse> {
1856        let req = RwCloudValidateSourceRequest {
1857            source_type: source_type.into(),
1858            source_config,
1859        };
1860        let resp = self.inner.rw_cloud_validate_source(req).await?;
1861        Ok(resp)
1862    }
1863
1864    pub async fn sink_coordinate_client(&self) -> SinkCoordinationRpcClient {
1865        self.inner.core.read().await.sink_coordinate_client.clone()
1866    }
1867
1868    pub async fn list_compact_task_assignment(&self) -> Result<Vec<CompactTaskAssignment>> {
1869        let req = ListCompactTaskAssignmentRequest {};
1870        let resp = self.inner.list_compact_task_assignment(req).await?;
1871        Ok(resp.task_assignment)
1872    }
1873
1874    pub async fn list_event_log(&self) -> Result<Vec<EventLog>> {
1875        let req = ListEventLogRequest::default();
1876        let resp = self.inner.list_event_log(req).await?;
1877        Ok(resp.event_logs)
1878    }
1879
1880    pub async fn list_compact_task_progress(&self) -> Result<Vec<CompactTaskProgress>> {
1881        let req = ListCompactTaskProgressRequest {};
1882        let resp = self.inner.list_compact_task_progress(req).await?;
1883        Ok(resp.task_progress)
1884    }
1885
1886    #[cfg(madsim)]
1887    pub fn try_add_panic_event_blocking(
1888        &self,
1889        panic_info: impl Display,
1890        timeout_millis: Option<u64>,
1891    ) {
1892    }
1893
1894    /// If `timeout_millis` is None, default is used.
1895    #[cfg(not(madsim))]
1896    pub fn try_add_panic_event_blocking(
1897        &self,
1898        panic_info: impl Display,
1899        timeout_millis: Option<u64>,
1900    ) {
1901        let event = event_log::EventWorkerNodePanic {
1902            worker_id: self.worker_id,
1903            worker_type: self.worker_type.into(),
1904            host_addr: Some(self.host_addr.to_protobuf()),
1905            panic_info: format!("{panic_info}"),
1906        };
1907        let grpc_meta_client = self.inner.clone();
1908        let _ = thread::spawn(move || {
1909            let rt = tokio::runtime::Builder::new_current_thread()
1910                .enable_all()
1911                .build()
1912                .unwrap();
1913            let req = AddEventLogRequest {
1914                event: Some(add_event_log_request::Event::WorkerNodePanic(event)),
1915            };
1916            rt.block_on(async {
1917                let _ = tokio::time::timeout(
1918                    Duration::from_millis(timeout_millis.unwrap_or(1000)),
1919                    grpc_meta_client.add_event_log(req),
1920                )
1921                .await;
1922            });
1923        })
1924        .join();
1925    }
1926
1927    pub async fn add_sink_fail_evet(
1928        &self,
1929        sink_id: SinkId,
1930        sink_name: String,
1931        connector: String,
1932        error: String,
1933    ) -> Result<()> {
1934        let event = event_log::EventSinkFail {
1935            sink_id,
1936            sink_name,
1937            connector,
1938            error,
1939        };
1940        let req = AddEventLogRequest {
1941            event: Some(add_event_log_request::Event::SinkFail(event)),
1942        };
1943        self.inner.add_event_log(req).await?;
1944        Ok(())
1945    }
1946
1947    pub async fn add_cdc_auto_schema_change_fail_event(
1948        &self,
1949        source_id: SourceId,
1950        table_name: String,
1951        cdc_table_id: String,
1952        upstream_ddl: String,
1953        fail_info: String,
1954    ) -> Result<()> {
1955        let event = event_log::EventAutoSchemaChangeFail {
1956            table_id: source_id.as_cdc_table_id(),
1957            table_name,
1958            cdc_table_id,
1959            upstream_ddl,
1960            fail_info,
1961        };
1962        let req = AddEventLogRequest {
1963            event: Some(add_event_log_request::Event::AutoSchemaChangeFail(event)),
1964        };
1965        self.inner.add_event_log(req).await?;
1966        Ok(())
1967    }
1968
1969    pub async fn cancel_compact_task(&self, task_id: u64, task_status: TaskStatus) -> Result<bool> {
1970        let req = CancelCompactTaskRequest {
1971            task_id,
1972            task_status: task_status as _,
1973        };
1974        let resp = self.inner.cancel_compact_task(req).await?;
1975        Ok(resp.ret)
1976    }
1977
1978    pub async fn get_version_by_epoch(
1979        &self,
1980        epoch: HummockEpoch,
1981        table_id: TableId,
1982    ) -> Result<PbHummockVersion> {
1983        let req = GetVersionByEpochRequest { epoch, table_id };
1984        let resp = self.inner.get_version_by_epoch(req).await?;
1985        Ok(resp.version.unwrap())
1986    }
1987
1988    pub async fn get_cluster_limits(
1989        &self,
1990    ) -> Result<Vec<risingwave_common::util::cluster_limit::ClusterLimit>> {
1991        let req = GetClusterLimitsRequest {};
1992        let resp = self.inner.get_cluster_limits(req).await?;
1993        Ok(resp.active_limits.into_iter().map(|l| l.into()).collect())
1994    }
1995
1996    pub async fn merge_compaction_group(
1997        &self,
1998        left_group_id: CompactionGroupId,
1999        right_group_id: CompactionGroupId,
2000    ) -> Result<()> {
2001        let req = MergeCompactionGroupRequest {
2002            left_group_id,
2003            right_group_id,
2004        };
2005        self.inner.merge_compaction_group(req).await?;
2006        Ok(())
2007    }
2008
2009    /// List all rate limits for sources and backfills
2010    pub async fn list_rate_limits(&self) -> Result<Vec<RateLimitInfo>> {
2011        let request = ListRateLimitsRequest {};
2012        let resp = self.inner.list_rate_limits(request).await?;
2013        Ok(resp.rate_limits)
2014    }
2015
2016    pub async fn list_cdc_progress(&self) -> Result<HashMap<JobId, PbCdcProgress>> {
2017        let request = ListCdcProgressRequest {};
2018        let resp = self.inner.list_cdc_progress(request).await?;
2019        Ok(resp.cdc_progress)
2020    }
2021
2022    pub async fn list_refresh_table_states(&self) -> Result<Vec<RefreshTableState>> {
2023        let request = ListRefreshTableStatesRequest {};
2024        let resp = self.inner.list_refresh_table_states(request).await?;
2025        Ok(resp.states)
2026    }
2027
2028    pub async fn list_iceberg_compaction_status(
2029        &self,
2030    ) -> Result<Vec<list_iceberg_compaction_status_response::IcebergCompactionStatus>> {
2031        let request = ListIcebergCompactionStatusRequest {};
2032        let resp = self.inner.list_iceberg_compaction_status(request).await?;
2033        Ok(resp.statuses)
2034    }
2035
2036    pub async fn list_sink_log_store_tables(
2037        &self,
2038    ) -> Result<Vec<list_sink_log_store_tables_response::SinkLogStoreTable>> {
2039        let request = ListSinkLogStoreTablesRequest {};
2040        let resp = self.inner.list_sink_log_store_tables(request).await?;
2041        Ok(resp.tables)
2042    }
2043
2044    pub async fn create_iceberg_table(
2045        &self,
2046        table_job_info: PbTableJobInfo,
2047        sink_job_info: PbSinkJobInfo,
2048        iceberg_source: PbSource,
2049        if_not_exists: bool,
2050    ) -> Result<WaitVersion> {
2051        let request = CreateIcebergTableRequest {
2052            table_info: Some(table_job_info),
2053            sink_info: Some(sink_job_info),
2054            iceberg_source: Some(iceberg_source),
2055            if_not_exists,
2056        };
2057
2058        let resp = Box::pin(self.inner.create_iceberg_table(request)).await?;
2059        Ok(resp
2060            .version
2061            .ok_or_else(|| anyhow!("wait version not set"))?)
2062    }
2063
2064    pub async fn list_hosted_iceberg_tables(&self) -> Result<Vec<IcebergTable>> {
2065        let request = ListIcebergTablesRequest {};
2066        let resp = self.inner.list_iceberg_tables(request).await?;
2067        Ok(resp.iceberg_tables)
2068    }
2069
2070    /// Get the await tree of all nodes in the cluster.
2071    pub async fn get_cluster_stack_trace(
2072        &self,
2073        actor_traces_format: ActorTracesFormat,
2074    ) -> Result<StackTraceResponse> {
2075        let request = StackTraceRequest {
2076            actor_traces_format: actor_traces_format as i32,
2077        };
2078        let resp = self.inner.stack_trace(request).await?;
2079        Ok(resp)
2080    }
2081
2082    pub async fn set_sync_log_store_aligned(&self, job_id: JobId, aligned: bool) -> Result<()> {
2083        let request = SetSyncLogStoreAlignedRequest { job_id, aligned };
2084        self.inner.set_sync_log_store_aligned(request).await?;
2085        Ok(())
2086    }
2087
2088    pub async fn refresh(&self, request: RefreshRequest) -> Result<RefreshResponse> {
2089        self.inner.refresh(request).await
2090    }
2091
2092    pub async fn list_unmigrated_tables(&self) -> Result<HashMap<TableId, String>> {
2093        let request = ListUnmigratedTablesRequest {};
2094        let resp = self.inner.list_unmigrated_tables(request).await?;
2095        Ok(resp
2096            .tables
2097            .into_iter()
2098            .map(|table| (table.table_id, table.table_name))
2099            .collect())
2100    }
2101}
2102
2103#[async_trait]
2104impl HummockMetaClient for MetaClient {
2105    async fn unpin_version_before(&self, unpin_version_before: HummockVersionId) -> Result<()> {
2106        let req = UnpinVersionBeforeRequest {
2107            context_id: self.worker_id(),
2108            unpin_version_before,
2109        };
2110        self.inner.unpin_version_before(req).await?;
2111        Ok(())
2112    }
2113
2114    async fn get_current_version(&self) -> Result<HummockVersion> {
2115        let req = GetCurrentVersionRequest::default();
2116        Ok(HummockVersion::from_rpc_protobuf(
2117            &self
2118                .inner
2119                .get_current_version(req)
2120                .await?
2121                .current_version
2122                .unwrap(),
2123        ))
2124    }
2125
2126    async fn get_new_object_ids(&self, number: u32) -> Result<ObjectIdRange> {
2127        let resp = self
2128            .inner
2129            .get_new_object_ids(GetNewObjectIdsRequest { number })
2130            .await?;
2131        Ok(ObjectIdRange::new(resp.start_id, resp.end_id))
2132    }
2133
2134    async fn commit_epoch_with_change_log(
2135        &self,
2136        _epoch: HummockEpoch,
2137        _sync_result: SyncResult,
2138        _change_log_info: Option<HummockMetaClientChangeLogInfo>,
2139    ) -> Result<()> {
2140        panic!("Only meta service can commit_epoch in production.")
2141    }
2142
2143    async fn trigger_manual_compaction(
2144        &self,
2145        compaction_group_id: CompactionGroupId,
2146        table_id: JobId,
2147        level: u32,
2148        target_level: Option<u32>,
2149        sst_ids: Vec<HummockSstableId>,
2150        exclusive: bool,
2151    ) -> Result<bool> {
2152        // TODO: support key_range parameter
2153        let req = TriggerManualCompactionRequest {
2154            compaction_group_id,
2155            table_id,
2156            // if table_id not exist, manual_compaction will include all the sst
2157            // without check internal_table_id
2158            level,
2159            target_level,
2160            sst_ids,
2161            exclusive: Some(exclusive),
2162            ..Default::default()
2163        };
2164
2165        let resp = self.inner.trigger_manual_compaction(req).await?;
2166        Ok(resp.should_retry.unwrap_or(false))
2167    }
2168
2169    async fn trigger_full_gc(
2170        &self,
2171        sst_retention_time_sec: u64,
2172        prefix: Option<String>,
2173    ) -> Result<()> {
2174        self.inner
2175            .trigger_full_gc(TriggerFullGcRequest {
2176                sst_retention_time_sec,
2177                prefix,
2178            })
2179            .await?;
2180        Ok(())
2181    }
2182
2183    async fn subscribe_compaction_event(
2184        &self,
2185    ) -> Result<(
2186        UnboundedSender<SubscribeCompactionEventRequest>,
2187        BoxStream<'static, CompactionEventItem>,
2188    )> {
2189        let (request_sender, request_receiver) =
2190            unbounded_channel::<SubscribeCompactionEventRequest>();
2191        request_sender
2192            .send(SubscribeCompactionEventRequest {
2193                event: Some(subscribe_compaction_event_request::Event::Register(
2194                    Register {
2195                        context_id: self.worker_id(),
2196                    },
2197                )),
2198                create_at: SystemTime::now()
2199                    .duration_since(SystemTime::UNIX_EPOCH)
2200                    .expect("Clock may have gone backwards")
2201                    .as_millis() as u64,
2202            })
2203            .context("Failed to subscribe compaction event")?;
2204
2205        let stream = self
2206            .inner
2207            .subscribe_compaction_event(Request::new(UnboundedReceiverStream::new(
2208                request_receiver,
2209            )))
2210            .await?;
2211
2212        Ok((request_sender, Box::pin(stream)))
2213    }
2214
2215    async fn get_version_by_epoch(
2216        &self,
2217        epoch: HummockEpoch,
2218        table_id: TableId,
2219    ) -> Result<PbHummockVersion> {
2220        self.get_version_by_epoch(epoch, table_id).await
2221    }
2222
2223    async fn subscribe_iceberg_compaction_event(
2224        &self,
2225    ) -> Result<(
2226        UnboundedSender<SubscribeIcebergCompactionEventRequest>,
2227        BoxStream<'static, IcebergCompactionEventItem>,
2228    )> {
2229        let (request_sender, request_receiver) =
2230            unbounded_channel::<SubscribeIcebergCompactionEventRequest>();
2231        request_sender
2232            .send(SubscribeIcebergCompactionEventRequest {
2233                event: Some(subscribe_iceberg_compaction_event_request::Event::Register(
2234                    IcebergRegister {
2235                        context_id: self.worker_id(),
2236                    },
2237                )),
2238                create_at: SystemTime::now()
2239                    .duration_since(SystemTime::UNIX_EPOCH)
2240                    .expect("Clock may have gone backwards")
2241                    .as_millis() as u64,
2242            })
2243            .context("Failed to subscribe compaction event")?;
2244
2245        let stream = self
2246            .inner
2247            .subscribe_iceberg_compaction_event(Request::new(UnboundedReceiverStream::new(
2248                request_receiver,
2249            )))
2250            .await?;
2251
2252        Ok((request_sender, Box::pin(stream)))
2253    }
2254
2255    async fn get_table_change_logs(
2256        &self,
2257        epoch_only: bool,
2258        start_epoch_inclusive: Option<u64>,
2259        end_epoch_inclusive: Option<u64>,
2260        table_ids: Option<HashSet<TableId>>,
2261        exclude_empty: bool,
2262        limit: Option<u32>,
2263    ) -> Result<TableChangeLogs> {
2264        self.get_table_change_logs(
2265            epoch_only,
2266            start_epoch_inclusive,
2267            end_epoch_inclusive,
2268            table_ids,
2269            exclude_empty,
2270            limit,
2271        )
2272        .await
2273    }
2274}
2275
2276#[async_trait]
2277impl TelemetryInfoFetcher for MetaClient {
2278    async fn fetch_telemetry_info(&self) -> std::result::Result<Option<String>, String> {
2279        let resp = self
2280            .get_telemetry_info()
2281            .await
2282            .map_err(|e| e.to_report_string())?;
2283        let tracking_id = resp.get_tracking_id().ok();
2284        Ok(tracking_id.map(|id| id.to_owned()))
2285    }
2286}
2287
2288pub type SinkCoordinationRpcClient = SinkCoordinationServiceClient<Channel>;
2289
2290#[derive(Debug, Clone)]
2291struct GrpcMetaClientCore {
2292    cluster_client: ClusterServiceClient<Channel>,
2293    meta_member_client: MetaMemberServiceClient<Channel>,
2294    heartbeat_client: HeartbeatServiceClient<Channel>,
2295    ddl_client: DdlServiceClient<Channel>,
2296    hummock_client: HummockManagerServiceClient<Channel>,
2297    notification_client: NotificationServiceClient<Channel>,
2298    stream_client: StreamManagerServiceClient<Channel>,
2299    user_client: UserServiceClient<Channel>,
2300    scale_client: ScaleServiceClient<Channel>,
2301    backup_client: BackupServiceClient<Channel>,
2302    telemetry_client: TelemetryInfoServiceClient<Channel>,
2303    system_params_client: SystemParamsServiceClient<Channel>,
2304    session_params_client: SessionParamServiceClient<Channel>,
2305    serving_client: ServingServiceClient<Channel>,
2306    cloud_client: CloudServiceClient<Channel>,
2307    sink_coordinate_client: SinkCoordinationRpcClient,
2308    event_log_client: EventLogServiceClient<Channel>,
2309    cluster_limit_client: ClusterLimitServiceClient<Channel>,
2310    hosted_iceberg_catalog_service_client: HostedIcebergCatalogServiceClient<Channel>,
2311    monitor_client: MonitorServiceClient<Channel>,
2312}
2313
2314impl GrpcMetaClientCore {
2315    pub(crate) fn new(channel: Channel) -> Self {
2316        let cluster_client = ClusterServiceClient::new(channel.clone());
2317        let meta_member_client = MetaMemberClient::new(channel.clone());
2318        let heartbeat_client = HeartbeatServiceClient::new(channel.clone());
2319        let ddl_client =
2320            DdlServiceClient::new(channel.clone()).max_decoding_message_size(usize::MAX);
2321        let hummock_client =
2322            HummockManagerServiceClient::new(channel.clone()).max_decoding_message_size(usize::MAX);
2323        let notification_client =
2324            NotificationServiceClient::new(channel.clone()).max_decoding_message_size(usize::MAX);
2325        let stream_client =
2326            StreamManagerServiceClient::new(channel.clone()).max_decoding_message_size(usize::MAX);
2327        let user_client = UserServiceClient::new(channel.clone());
2328        let scale_client =
2329            ScaleServiceClient::new(channel.clone()).max_decoding_message_size(usize::MAX);
2330        let backup_client = BackupServiceClient::new(channel.clone());
2331        let telemetry_client =
2332            TelemetryInfoServiceClient::new(channel.clone()).max_decoding_message_size(usize::MAX);
2333        let system_params_client = SystemParamsServiceClient::new(channel.clone());
2334        let session_params_client = SessionParamServiceClient::new(channel.clone());
2335        let serving_client = ServingServiceClient::new(channel.clone());
2336        let cloud_client = CloudServiceClient::new(channel.clone());
2337        let sink_coordinate_client = SinkCoordinationServiceClient::new(channel.clone())
2338            .max_decoding_message_size(usize::MAX);
2339        let event_log_client = EventLogServiceClient::new(channel.clone());
2340        let cluster_limit_client = ClusterLimitServiceClient::new(channel.clone());
2341        let hosted_iceberg_catalog_service_client =
2342            HostedIcebergCatalogServiceClient::new(channel.clone());
2343        let monitor_client = configured_monitor_service_client(MonitorServiceClient::new(channel));
2344
2345        GrpcMetaClientCore {
2346            cluster_client,
2347            meta_member_client,
2348            heartbeat_client,
2349            ddl_client,
2350            hummock_client,
2351            notification_client,
2352            stream_client,
2353            user_client,
2354            scale_client,
2355            backup_client,
2356            telemetry_client,
2357            system_params_client,
2358            session_params_client,
2359            serving_client,
2360            cloud_client,
2361            sink_coordinate_client,
2362            event_log_client,
2363            cluster_limit_client,
2364            hosted_iceberg_catalog_service_client,
2365            monitor_client,
2366        }
2367    }
2368}
2369
2370/// Client to meta server. Cloning the instance is lightweight.
2371///
2372/// It is a wrapper of tonic client. See [`crate::meta_rpc_client_method_impl`].
2373#[derive(Debug, Clone)]
2374struct GrpcMetaClient {
2375    member_monitor_event_sender: mpsc::Sender<Sender<Result<()>>>,
2376    core: Arc<RwLock<GrpcMetaClientCore>>,
2377}
2378
2379type MetaMemberClient = MetaMemberServiceClient<Channel>;
2380
2381struct MetaMemberGroup {
2382    members: LruCache<http::Uri, Option<MetaMemberClient>>,
2383}
2384
2385struct MetaMemberManagement {
2386    core_ref: Arc<RwLock<GrpcMetaClientCore>>,
2387    members: Either<MetaMemberClient, MetaMemberGroup>,
2388    current_leader: http::Uri,
2389    meta_config: Arc<MetaConfig>,
2390}
2391
2392impl MetaMemberManagement {
2393    const META_MEMBER_REFRESH_PERIOD: Duration = Duration::from_secs(5);
2394
2395    fn host_address_to_uri(addr: HostAddress) -> http::Uri {
2396        format!("http://{}:{}", addr.host, addr.port)
2397            .parse()
2398            .unwrap()
2399    }
2400
2401    async fn recreate_core(&self, channel: Channel) {
2402        let mut core = self.core_ref.write().await;
2403        *core = GrpcMetaClientCore::new(channel);
2404    }
2405
2406    async fn refresh_members(&mut self) -> Result<()> {
2407        let leader_addr = match self.members.as_mut() {
2408            Either::Left(client) => {
2409                let resp = client
2410                    .to_owned()
2411                    .members(MembersRequest {})
2412                    .await
2413                    .map_err(RpcError::from_meta_status)?;
2414                let resp = resp.into_inner();
2415                resp.members.into_iter().find(|member| member.is_leader)
2416            }
2417            Either::Right(member_group) => {
2418                let mut fetched_members = None;
2419
2420                for (addr, client) in &mut member_group.members {
2421                    let members: Result<_> = try {
2422                        let mut client = match client {
2423                            Some(cached_client) => cached_client.to_owned(),
2424                            None => {
2425                                let endpoint = GrpcMetaClient::addr_to_endpoint(addr.clone());
2426                                let channel = GrpcMetaClient::connect_to_endpoint(endpoint)
2427                                    .await
2428                                    .context("failed to create client")
2429                                    .map_err(RpcError::from)?;
2430                                let new_client: MetaMemberClient =
2431                                    MetaMemberServiceClient::new(channel);
2432                                *client = Some(new_client.clone());
2433
2434                                new_client
2435                            }
2436                        };
2437
2438                        let resp = client
2439                            .members(MembersRequest {})
2440                            .await
2441                            .context("failed to fetch members")
2442                            .map_err(RpcError::from)?;
2443
2444                        resp.into_inner().members
2445                    };
2446
2447                    let fetched = members.is_ok();
2448                    fetched_members = Some(members);
2449                    if fetched {
2450                        break;
2451                    }
2452                }
2453
2454                let members = fetched_members
2455                    .context("no member available in the list")?
2456                    .context("could not refresh members")?;
2457
2458                // find new leader
2459                let mut leader = None;
2460                for member in members {
2461                    if member.is_leader {
2462                        leader = Some(member.clone());
2463                    }
2464
2465                    let addr = Self::host_address_to_uri(member.address.unwrap());
2466                    // We don't clean any expired addrs here to deal with some extreme situations.
2467                    if !member_group.members.contains(&addr) {
2468                        tracing::info!("new meta member joined: {}", addr);
2469                        member_group.members.put(addr, None);
2470                    }
2471                }
2472
2473                leader
2474            }
2475        };
2476
2477        if let Some(leader) = leader_addr {
2478            let discovered_leader = Self::host_address_to_uri(leader.address.unwrap());
2479
2480            if discovered_leader != self.current_leader {
2481                tracing::info!("new meta leader {} discovered", discovered_leader);
2482
2483                let retry_strategy = GrpcMetaClient::retry_strategy_to_bound(
2484                    Duration::from_secs(self.meta_config.meta_leader_lease_secs),
2485                    false,
2486                );
2487
2488                let channel = tokio_retry::Retry::spawn(retry_strategy, || async {
2489                    let endpoint = GrpcMetaClient::addr_to_endpoint(discovered_leader.clone());
2490                    GrpcMetaClient::connect_to_endpoint(endpoint).await
2491                })
2492                .await?;
2493
2494                self.recreate_core(channel).await;
2495                self.current_leader = discovered_leader;
2496            }
2497        }
2498
2499        Ok(())
2500    }
2501}
2502
2503impl GrpcMetaClient {
2504    // See `Endpoint::http2_keep_alive_interval`
2505    const ENDPOINT_KEEP_ALIVE_INTERVAL_SEC: u64 = 60;
2506    // See `Endpoint::keep_alive_timeout`
2507    const ENDPOINT_KEEP_ALIVE_TIMEOUT_SEC: u64 = 60;
2508    // Retry base interval in ms for connecting to meta server.
2509    const INIT_RETRY_BASE_INTERVAL_MS: u64 = 10;
2510    // Max retry times for connecting to meta server.
2511    const INIT_RETRY_MAX_INTERVAL_MS: u64 = 2000;
2512
2513    fn start_meta_member_monitor(
2514        &self,
2515        init_leader_addr: http::Uri,
2516        members: Either<MetaMemberClient, MetaMemberGroup>,
2517        force_refresh_receiver: Receiver<Sender<Result<()>>>,
2518        meta_config: Arc<MetaConfig>,
2519    ) -> Result<()> {
2520        let core_ref: Arc<RwLock<GrpcMetaClientCore>> = self.core.clone();
2521        let current_leader = init_leader_addr;
2522
2523        let enable_period_tick = matches!(members, Either::Right(_));
2524
2525        let member_management = MetaMemberManagement {
2526            core_ref,
2527            members,
2528            current_leader,
2529            meta_config,
2530        };
2531
2532        let mut force_refresh_receiver = force_refresh_receiver;
2533
2534        tokio::spawn(async move {
2535            let mut member_management = member_management;
2536            let mut ticker = time::interval(MetaMemberManagement::META_MEMBER_REFRESH_PERIOD);
2537
2538            loop {
2539                let event: Option<Sender<Result<()>>> = if enable_period_tick {
2540                    tokio::select! {
2541                        _ = ticker.tick() => None,
2542                        result_sender = force_refresh_receiver.recv() => {
2543                            if result_sender.is_none() {
2544                                break;
2545                            }
2546
2547                            result_sender
2548                        },
2549                    }
2550                } else {
2551                    let result_sender = force_refresh_receiver.recv().await;
2552
2553                    if result_sender.is_none() {
2554                        break;
2555                    }
2556
2557                    result_sender
2558                };
2559
2560                let tick_result = member_management.refresh_members().await;
2561                if let Err(e) = tick_result.as_ref() {
2562                    tracing::warn!(error = %e.as_report(),  "refresh meta member client failed");
2563                }
2564
2565                if let Some(sender) = event {
2566                    // ignore resp
2567                    let _resp = sender.send(tick_result);
2568                }
2569            }
2570        });
2571
2572        Ok(())
2573    }
2574
2575    async fn force_refresh_leader(&self) -> Result<()> {
2576        let (sender, receiver) = oneshot::channel();
2577
2578        self.member_monitor_event_sender
2579            .send(sender)
2580            .await
2581            .map_err(|e| anyhow!(e))?;
2582
2583        receiver.await.map_err(|e| anyhow!(e))?
2584    }
2585
2586    /// Connect to the meta server from `addrs`.
2587    pub async fn new(strategy: &MetaAddressStrategy, config: Arc<MetaConfig>) -> Result<Self> {
2588        let (channel, addr) = match strategy {
2589            MetaAddressStrategy::LoadBalance(addr) => {
2590                Self::try_build_rpc_channel(vec![addr.clone()]).await
2591            }
2592            MetaAddressStrategy::List(addrs) => Self::try_build_rpc_channel(addrs.clone()).await,
2593        }?;
2594        let (force_refresh_sender, force_refresh_receiver) = mpsc::channel(1);
2595        let client = GrpcMetaClient {
2596            member_monitor_event_sender: force_refresh_sender,
2597            core: Arc::new(RwLock::new(GrpcMetaClientCore::new(channel))),
2598        };
2599
2600        let meta_member_client = client.core.read().await.meta_member_client.clone();
2601        let members = match strategy {
2602            MetaAddressStrategy::LoadBalance(_) => Either::Left(meta_member_client),
2603            MetaAddressStrategy::List(addrs) => {
2604                let mut members = LruCache::new(20.try_into().unwrap());
2605                for addr in addrs {
2606                    members.put(addr.clone(), None);
2607                }
2608                members.put(addr.clone(), Some(meta_member_client));
2609
2610                Either::Right(MetaMemberGroup { members })
2611            }
2612        };
2613
2614        client.start_meta_member_monitor(addr, members, force_refresh_receiver, config.clone())?;
2615
2616        client.force_refresh_leader().await?;
2617
2618        Ok(client)
2619    }
2620
2621    fn addr_to_endpoint(addr: http::Uri) -> Endpoint {
2622        Endpoint::from(addr).initial_connection_window_size(MAX_CONNECTION_WINDOW_SIZE)
2623    }
2624
2625    pub(crate) async fn try_build_rpc_channel(
2626        addrs: impl IntoIterator<Item = http::Uri>,
2627    ) -> Result<(Channel, http::Uri)> {
2628        let endpoints: Vec<_> = addrs
2629            .into_iter()
2630            .map(|addr| (Self::addr_to_endpoint(addr.clone()), addr))
2631            .collect();
2632
2633        let mut last_error = None;
2634
2635        for (endpoint, addr) in endpoints {
2636            match Self::connect_to_endpoint(endpoint).await {
2637                Ok(channel) => {
2638                    tracing::info!("Connect to meta server {} successfully", addr);
2639                    return Ok((channel, addr));
2640                }
2641                Err(e) => {
2642                    tracing::warn!(
2643                        error = %e.as_report(),
2644                        "Failed to connect to meta server {}, trying again",
2645                        addr,
2646                    );
2647                    last_error = Some(e);
2648                }
2649            }
2650        }
2651
2652        if let Some(last_error) = last_error {
2653            Err(anyhow::anyhow!(last_error)
2654                .context("failed to connect to all meta servers")
2655                .into())
2656        } else {
2657            bail!("no meta server address provided")
2658        }
2659    }
2660
2661    async fn connect_to_endpoint(endpoint: Endpoint) -> Result<Channel> {
2662        let channel = endpoint
2663            .http2_keep_alive_interval(Duration::from_secs(Self::ENDPOINT_KEEP_ALIVE_INTERVAL_SEC))
2664            .keep_alive_timeout(Duration::from_secs(Self::ENDPOINT_KEEP_ALIVE_TIMEOUT_SEC))
2665            .connect_timeout(Duration::from_secs(5))
2666            .monitored_connect("grpc-meta-client", Default::default())
2667            .await?
2668            .wrapped();
2669
2670        Ok(channel)
2671    }
2672
2673    pub(crate) fn retry_strategy_to_bound(
2674        high_bound: Duration,
2675        exceed: bool,
2676    ) -> impl Iterator<Item = Duration> {
2677        let iter = exponential_backoff(
2678            Duration::from_millis(Self::INIT_RETRY_BASE_INTERVAL_MS),
2679            Self::INIT_RETRY_BASE_INTERVAL_MS,
2680            Duration::from_millis(Self::INIT_RETRY_MAX_INTERVAL_MS),
2681        )
2682        .map(jitter);
2683
2684        let mut sum = Duration::default();
2685
2686        iter.take_while(move |duration| {
2687            sum += *duration;
2688
2689            if exceed {
2690                sum < high_bound + *duration
2691            } else {
2692                sum < high_bound
2693            }
2694        })
2695    }
2696}
2697
2698macro_rules! for_all_meta_rpc {
2699    ($macro:ident) => {
2700        $macro! {
2701             { cluster_client, add_worker_node, AddWorkerNodeRequest, AddWorkerNodeResponse }
2702            ,{ cluster_client, activate_worker_node, ActivateWorkerNodeRequest, ActivateWorkerNodeResponse }
2703            ,{ cluster_client, delete_worker_node, DeleteWorkerNodeRequest, DeleteWorkerNodeResponse }
2704            ,{ cluster_client, list_all_nodes, ListAllNodesRequest, ListAllNodesResponse }
2705            ,{ cluster_client, get_cluster_recovery_status, GetClusterRecoveryStatusRequest, GetClusterRecoveryStatusResponse }
2706            ,{ cluster_client, get_meta_store_info, GetMetaStoreInfoRequest, GetMetaStoreInfoResponse }
2707            ,{ heartbeat_client, heartbeat, HeartbeatRequest, HeartbeatResponse }
2708            ,{ stream_client, flush, FlushRequest, FlushResponse }
2709            ,{ stream_client, pause, PauseRequest, PauseResponse }
2710            ,{ stream_client, resume, ResumeRequest, ResumeResponse }
2711            ,{ stream_client, apply_throttle, ApplyThrottleRequest, ApplyThrottleResponse }
2712            ,{ stream_client, cancel_creating_jobs, CancelCreatingJobsRequest, CancelCreatingJobsResponse }
2713            ,{ stream_client, list_table_fragments, ListTableFragmentsRequest, ListTableFragmentsResponse }
2714            ,{ stream_client, list_streaming_job_states, ListStreamingJobStatesRequest, ListStreamingJobStatesResponse }
2715            ,{ stream_client, list_fragment_distribution, ListFragmentDistributionRequest, ListFragmentDistributionResponse }
2716            ,{ stream_client, list_creating_fragment_distribution, ListCreatingFragmentDistributionRequest, ListCreatingFragmentDistributionResponse }
2717            ,{ stream_client, list_sink_log_store_tables, ListSinkLogStoreTablesRequest, ListSinkLogStoreTablesResponse }
2718            ,{ stream_client, list_actor_states, ListActorStatesRequest, ListActorStatesResponse }
2719            ,{ stream_client, list_actor_splits, ListActorSplitsRequest, ListActorSplitsResponse }
2720            ,{ stream_client, recover, RecoverRequest, RecoverResponse }
2721            ,{ stream_client, list_rate_limits, ListRateLimitsRequest, ListRateLimitsResponse }
2722            ,{ stream_client, list_cdc_progress, ListCdcProgressRequest, ListCdcProgressResponse }
2723            ,{ stream_client, list_refresh_table_states, ListRefreshTableStatesRequest, ListRefreshTableStatesResponse }
2724            ,{ stream_client, list_iceberg_compaction_status, ListIcebergCompactionStatusRequest, ListIcebergCompactionStatusResponse }
2725            ,{ stream_client, alter_connector_props, AlterConnectorPropsRequest, AlterConnectorPropsResponse }
2726            ,{ stream_client, alter_source_properties_safe, AlterSourcePropertiesSafeRequest, AlterSourcePropertiesSafeResponse }
2727            ,{ stream_client, reset_source_splits, ResetSourceSplitsRequest, ResetSourceSplitsResponse }
2728            ,{ stream_client, inject_source_offsets, InjectSourceOffsetsRequest, InjectSourceOffsetsResponse }
2729            ,{ stream_client, get_fragment_by_id, GetFragmentByIdRequest, GetFragmentByIdResponse }
2730            ,{ stream_client, get_fragment_vnodes, GetFragmentVnodesRequest, GetFragmentVnodesResponse }
2731            ,{ stream_client, get_actor_vnodes, GetActorVnodesRequest, GetActorVnodesResponse }
2732            ,{ stream_client, set_sync_log_store_aligned, SetSyncLogStoreAlignedRequest, SetSyncLogStoreAlignedResponse }
2733            ,{ stream_client, refresh, RefreshRequest, RefreshResponse }
2734            ,{ stream_client, list_unmigrated_tables, ListUnmigratedTablesRequest, ListUnmigratedTablesResponse }
2735            ,{ ddl_client, create_table, CreateTableRequest, CreateTableResponse }
2736            ,{ ddl_client, alter_name, AlterNameRequest, AlterNameResponse }
2737            ,{ ddl_client, alter_owner, AlterOwnerRequest, AlterOwnerResponse }
2738            ,{ ddl_client, alter_subscription_retention, AlterSubscriptionRetentionRequest, AlterSubscriptionRetentionResponse }
2739            ,{ ddl_client, alter_set_schema, AlterSetSchemaRequest, AlterSetSchemaResponse }
2740            ,{ ddl_client, alter_parallelism, AlterParallelismRequest, AlterParallelismResponse }
2741            ,{ ddl_client, alter_backfill_parallelism, AlterBackfillParallelismRequest, AlterBackfillParallelismResponse }
2742            ,{ ddl_client, alter_streaming_job_config, AlterStreamingJobConfigRequest, AlterStreamingJobConfigResponse }
2743            ,{ ddl_client, alter_fragment_parallelism, AlterFragmentParallelismRequest, AlterFragmentParallelismResponse }
2744            ,{ ddl_client, alter_cdc_table_backfill_parallelism, AlterCdcTableBackfillParallelismRequest, AlterCdcTableBackfillParallelismResponse }
2745            ,{ ddl_client, alter_resource_group, AlterResourceGroupRequest, AlterResourceGroupResponse }
2746            ,{ ddl_client, alter_database_resource_group, AlterDatabaseResourceGroupRequest, AlterDatabaseResourceGroupResponse }
2747            ,{ ddl_client, alter_database_param, AlterDatabaseParamRequest, AlterDatabaseParamResponse }
2748            ,{ ddl_client, create_materialized_view, CreateMaterializedViewRequest, CreateMaterializedViewResponse }
2749            ,{ ddl_client, create_view, CreateViewRequest, CreateViewResponse }
2750            ,{ ddl_client, create_source, CreateSourceRequest, CreateSourceResponse }
2751            ,{ ddl_client, create_sink, CreateSinkRequest, CreateSinkResponse }
2752            ,{ ddl_client, create_subscription, CreateSubscriptionRequest, CreateSubscriptionResponse }
2753            ,{ ddl_client, create_schema, CreateSchemaRequest, CreateSchemaResponse }
2754            ,{ ddl_client, create_database, CreateDatabaseRequest, CreateDatabaseResponse }
2755            ,{ ddl_client, create_secret, CreateSecretRequest, CreateSecretResponse }
2756            ,{ ddl_client, create_index, CreateIndexRequest, CreateIndexResponse }
2757            ,{ ddl_client, create_function, CreateFunctionRequest, CreateFunctionResponse }
2758            ,{ ddl_client, drop_table, DropTableRequest, DropTableResponse }
2759            ,{ ddl_client, drop_materialized_view, DropMaterializedViewRequest, DropMaterializedViewResponse }
2760            ,{ ddl_client, drop_view, DropViewRequest, DropViewResponse }
2761            ,{ ddl_client, drop_source, DropSourceRequest, DropSourceResponse }
2762            ,{ ddl_client, reset_source, ResetSourceRequest, ResetSourceResponse }
2763            ,{ ddl_client, drop_secret, DropSecretRequest, DropSecretResponse}
2764            ,{ ddl_client, drop_sink, DropSinkRequest, DropSinkResponse }
2765            ,{ ddl_client, drop_subscription, DropSubscriptionRequest, DropSubscriptionResponse }
2766            ,{ ddl_client, drop_database, DropDatabaseRequest, DropDatabaseResponse }
2767            ,{ ddl_client, drop_schema, DropSchemaRequest, DropSchemaResponse }
2768            ,{ ddl_client, drop_index, DropIndexRequest, DropIndexResponse }
2769            ,{ ddl_client, drop_function, DropFunctionRequest, DropFunctionResponse }
2770            ,{ ddl_client, replace_job_plan, ReplaceJobPlanRequest, ReplaceJobPlanResponse }
2771            ,{ ddl_client, alter_source, AlterSourceRequest, AlterSourceResponse }
2772            ,{ ddl_client, risectl_list_state_tables, RisectlListStateTablesRequest, RisectlListStateTablesResponse }
2773            ,{ ddl_client, risectl_resume_backfill, RisectlResumeBackfillRequest, RisectlResumeBackfillResponse }
2774            ,{ ddl_client, get_ddl_progress, GetDdlProgressRequest, GetDdlProgressResponse }
2775            ,{ ddl_client, create_connection, CreateConnectionRequest, CreateConnectionResponse }
2776            ,{ ddl_client, list_connections, ListConnectionsRequest, ListConnectionsResponse }
2777            ,{ ddl_client, drop_connection, DropConnectionRequest, DropConnectionResponse }
2778            ,{ ddl_client, comment_on, CommentOnRequest, CommentOnResponse }
2779            ,{ ddl_client, get_tables, GetTablesRequest, GetTablesResponse }
2780            ,{ ddl_client, wait, WaitRequest, WaitResponse }
2781            ,{ ddl_client, auto_schema_change, AutoSchemaChangeRequest, AutoSchemaChangeResponse }
2782            ,{ ddl_client, alter_swap_rename, AlterSwapRenameRequest, AlterSwapRenameResponse }
2783            ,{ ddl_client, alter_secret, AlterSecretRequest, AlterSecretResponse }
2784            ,{ ddl_client, compact_iceberg_table, CompactIcebergTableRequest, CompactIcebergTableResponse }
2785            ,{ ddl_client, rewrite_iceberg_table_manifests, RewriteIcebergTableManifestsRequest, RewriteIcebergTableManifestsResponse }
2786            ,{ ddl_client, expire_iceberg_table_snapshots, ExpireIcebergTableSnapshotsRequest, ExpireIcebergTableSnapshotsResponse }
2787            ,{ ddl_client, create_iceberg_table, CreateIcebergTableRequest, CreateIcebergTableResponse }
2788            ,{ ddl_client, wait_iceberg_pk_index_sink_epoch, WaitIcebergPkIndexSinkEpochRequest, WaitIcebergPkIndexSinkEpochResponse }
2789            ,{ hummock_client, unpin_version_before, UnpinVersionBeforeRequest, UnpinVersionBeforeResponse }
2790            ,{ hummock_client, get_current_version, GetCurrentVersionRequest, GetCurrentVersionResponse }
2791            ,{ hummock_client, replay_version_delta, ReplayVersionDeltaRequest, ReplayVersionDeltaResponse }
2792            ,{ hummock_client, list_version_deltas, ListVersionDeltasRequest, ListVersionDeltasResponse }
2793            ,{ hummock_client, get_assigned_compact_task_num, GetAssignedCompactTaskNumRequest, GetAssignedCompactTaskNumResponse }
2794            ,{ hummock_client, trigger_compaction_deterministic, TriggerCompactionDeterministicRequest, TriggerCompactionDeterministicResponse }
2795            ,{ hummock_client, disable_commit_epoch, DisableCommitEpochRequest, DisableCommitEpochResponse }
2796            ,{ hummock_client, get_new_object_ids, GetNewObjectIdsRequest, GetNewObjectIdsResponse }
2797            ,{ hummock_client, trigger_manual_compaction, TriggerManualCompactionRequest, TriggerManualCompactionResponse }
2798            ,{ hummock_client, trigger_full_gc, TriggerFullGcRequest, TriggerFullGcResponse }
2799            ,{ hummock_client, rise_ctl_get_pinned_versions_summary, RiseCtlGetPinnedVersionsSummaryRequest, RiseCtlGetPinnedVersionsSummaryResponse }
2800            ,{ hummock_client, rise_ctl_list_compaction_group, RiseCtlListCompactionGroupRequest, RiseCtlListCompactionGroupResponse }
2801            ,{ hummock_client, rise_ctl_update_compaction_config, RiseCtlUpdateCompactionConfigRequest, RiseCtlUpdateCompactionConfigResponse }
2802            ,{ hummock_client, rise_ctl_get_checkpoint_version, RiseCtlGetCheckpointVersionRequest, RiseCtlGetCheckpointVersionResponse }
2803            ,{ hummock_client, rise_ctl_pause_version_checkpoint, RiseCtlPauseVersionCheckpointRequest, RiseCtlPauseVersionCheckpointResponse }
2804            ,{ hummock_client, rise_ctl_resume_version_checkpoint, RiseCtlResumeVersionCheckpointRequest, RiseCtlResumeVersionCheckpointResponse }
2805            ,{ hummock_client, init_metadata_for_replay, InitMetadataForReplayRequest, InitMetadataForReplayResponse }
2806            ,{ hummock_client, split_compaction_group, SplitCompactionGroupRequest, SplitCompactionGroupResponse }
2807            ,{ hummock_client, rise_ctl_list_compaction_status, RiseCtlListCompactionStatusRequest, RiseCtlListCompactionStatusResponse }
2808            ,{ hummock_client, get_compaction_score, GetCompactionScoreRequest, GetCompactionScoreResponse }
2809            ,{ hummock_client, rise_ctl_rebuild_table_stats, RiseCtlRebuildTableStatsRequest, RiseCtlRebuildTableStatsResponse }
2810            ,{ hummock_client, subscribe_compaction_event, impl tonic::IntoStreamingRequest<Message = SubscribeCompactionEventRequest>, Streaming<SubscribeCompactionEventResponse> }
2811            ,{ hummock_client, subscribe_iceberg_compaction_event, impl tonic::IntoStreamingRequest<Message = SubscribeIcebergCompactionEventRequest>, Streaming<SubscribeIcebergCompactionEventResponse> }
2812            ,{ hummock_client, list_branched_object, ListBranchedObjectRequest, ListBranchedObjectResponse }
2813            ,{ hummock_client, list_active_write_limit, ListActiveWriteLimitRequest, ListActiveWriteLimitResponse }
2814            ,{ hummock_client, list_hummock_meta_config, ListHummockMetaConfigRequest, ListHummockMetaConfigResponse }
2815            ,{ hummock_client, list_compact_task_assignment, ListCompactTaskAssignmentRequest, ListCompactTaskAssignmentResponse }
2816            ,{ hummock_client, list_compact_task_progress, ListCompactTaskProgressRequest, ListCompactTaskProgressResponse }
2817            ,{ hummock_client, cancel_compact_task, CancelCompactTaskRequest, CancelCompactTaskResponse}
2818            ,{ hummock_client, get_version_by_epoch, GetVersionByEpochRequest, GetVersionByEpochResponse }
2819            ,{ hummock_client, merge_compaction_group, MergeCompactionGroupRequest, MergeCompactionGroupResponse }
2820            ,{ hummock_client, get_table_change_logs, GetTableChangeLogsRequest, GetTableChangeLogsResponse }
2821            ,{ user_client, create_user, CreateUserRequest, CreateUserResponse }
2822            ,{ user_client, update_user, UpdateUserRequest, UpdateUserResponse }
2823            ,{ user_client, drop_user, DropUserRequest, DropUserResponse }
2824            ,{ user_client, grant_privilege, GrantPrivilegeRequest, GrantPrivilegeResponse }
2825            ,{ user_client, revoke_privilege, RevokePrivilegeRequest, RevokePrivilegeResponse }
2826            ,{ user_client, alter_default_privilege, AlterDefaultPrivilegeRequest, AlterDefaultPrivilegeResponse }
2827            ,{ scale_client, get_cluster_info, GetClusterInfoRequest, GetClusterInfoResponse }
2828            ,{ scale_client, reschedule, RescheduleRequest, RescheduleResponse }
2829            ,{ notification_client, subscribe, SubscribeRequest, Streaming<SubscribeResponse> }
2830            ,{ backup_client, backup_meta, BackupMetaRequest, BackupMetaResponse }
2831            ,{ backup_client, get_backup_job_status, GetBackupJobStatusRequest, GetBackupJobStatusResponse }
2832            ,{ backup_client, delete_meta_snapshot, DeleteMetaSnapshotRequest, DeleteMetaSnapshotResponse}
2833            ,{ backup_client, get_meta_snapshot_manifest, GetMetaSnapshotManifestRequest, GetMetaSnapshotManifestResponse}
2834            ,{ telemetry_client, get_telemetry_info, GetTelemetryInfoRequest, TelemetryInfoResponse}
2835            ,{ system_params_client, get_system_params, GetSystemParamsRequest, GetSystemParamsResponse }
2836            ,{ system_params_client, set_system_param, SetSystemParamRequest, SetSystemParamResponse }
2837            ,{ system_params_client, clear_file_cache, ClearFileCacheRequest, ClearFileCacheResponse }
2838            ,{ session_params_client, get_session_params, GetSessionParamsRequest, GetSessionParamsResponse }
2839            ,{ session_params_client, set_session_param, SetSessionParamRequest, SetSessionParamResponse }
2840            ,{ serving_client, get_serving_vnode_mappings, GetServingVnodeMappingsRequest, GetServingVnodeMappingsResponse }
2841            ,{ cloud_client, rw_cloud_validate_source, RwCloudValidateSourceRequest, RwCloudValidateSourceResponse }
2842            ,{ event_log_client, list_event_log, ListEventLogRequest, ListEventLogResponse }
2843            ,{ event_log_client, add_event_log, AddEventLogRequest, AddEventLogResponse }
2844            ,{ cluster_limit_client, get_cluster_limits, GetClusterLimitsRequest, GetClusterLimitsResponse }
2845            ,{ hosted_iceberg_catalog_service_client, list_iceberg_tables, ListIcebergTablesRequest, ListIcebergTablesResponse }
2846            ,{ monitor_client, stack_trace, StackTraceRequest, StackTraceResponse }
2847        }
2848    };
2849}
2850
2851impl GrpcMetaClient {
2852    async fn refresh_client_if_needed(&self, code: Code) {
2853        if matches!(
2854            code,
2855            Code::Unknown | Code::Unimplemented | Code::Unavailable
2856        ) {
2857            tracing::debug!("matching tonic code {}", code);
2858            let (result_sender, result_receiver) = oneshot::channel();
2859            if self
2860                .member_monitor_event_sender
2861                .try_send(result_sender)
2862                .is_ok()
2863            {
2864                if let Ok(Err(e)) = result_receiver.await {
2865                    tracing::warn!(error = %e.as_report(), "force refresh meta client failed");
2866                }
2867            } else {
2868                tracing::debug!("skipping the current refresh, somewhere else is already doing it")
2869            }
2870        }
2871    }
2872}
2873
2874impl GrpcMetaClient {
2875    for_all_meta_rpc! { meta_rpc_client_method_impl }
2876}