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