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