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