Skip to main content

risingwave_meta_service/
stream_service.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{HashMap, HashSet};
16
17use chrono::DateTime;
18use itertools::Itertools;
19use risingwave_common::id::JobId;
20use risingwave_common::secret::LocalSecretManager;
21use risingwave_common::util::stream_graph_visitor::visit_stream_node_mut;
22use risingwave_connector::source::SplitMetaData;
23use risingwave_meta::barrier::BarrierManagerRef;
24use risingwave_meta::controller::fragment::StreamingJobInfo;
25use risingwave_meta::controller::utils::FragmentDesc;
26use risingwave_meta::manager::MetadataManager;
27use risingwave_meta::manager::iceberg_compaction::IcebergCompactionManagerRef;
28use risingwave_meta::stream::{GlobalRefreshManagerRef, SourceManagerRunningInfo};
29use risingwave_meta::{MetaError, model};
30use risingwave_meta_model::{ConnectionId, FragmentId, JobStatus, SourceId, StreamingParallelism};
31use risingwave_pb::common::ThrottleType;
32use risingwave_pb::meta::alter_connector_props_request::AlterConnectorPropsObject;
33use risingwave_pb::meta::cancel_creating_jobs_request::Jobs;
34use risingwave_pb::meta::list_actor_splits_response::FragmentType;
35use risingwave_pb::meta::list_cdc_progress_response::PbCdcProgress;
36use risingwave_pb::meta::list_refresh_table_states_response::RefreshTableState;
37use risingwave_pb::meta::list_table_fragments_response::{
38    ActorInfo, FragmentInfo, TableFragmentInfo,
39};
40use risingwave_pb::meta::stream_manager_service_server::StreamManagerService;
41use risingwave_pb::meta::table_fragments::PbState;
42use risingwave_pb::meta::table_fragments::fragment::PbFragmentDistributionType;
43use risingwave_pb::meta::*;
44use risingwave_pb::stream_plan::stream_node::NodeBody;
45use risingwave_pb::stream_plan::throttle_mutation::ThrottleConfig;
46use tonic::{Request, Response, Status};
47
48use crate::barrier::{BarrierScheduler, Command};
49use crate::manager::MetaSrvEnv;
50use crate::stream::GlobalStreamManagerRef;
51
52pub type TonicResponse<T> = Result<Response<T>, Status>;
53
54#[derive(Clone)]
55pub struct StreamServiceImpl {
56    env: MetaSrvEnv,
57    barrier_scheduler: BarrierScheduler,
58    barrier_manager: BarrierManagerRef,
59    stream_manager: GlobalStreamManagerRef,
60    metadata_manager: MetadataManager,
61    refresh_manager: GlobalRefreshManagerRef,
62    iceberg_compaction_manager: IcebergCompactionManagerRef,
63}
64
65impl StreamServiceImpl {
66    pub fn new(
67        env: MetaSrvEnv,
68        barrier_scheduler: BarrierScheduler,
69        barrier_manager: BarrierManagerRef,
70        stream_manager: GlobalStreamManagerRef,
71        metadata_manager: MetadataManager,
72        refresh_manager: GlobalRefreshManagerRef,
73        iceberg_compaction_manager: IcebergCompactionManagerRef,
74    ) -> Self {
75        StreamServiceImpl {
76            env,
77            barrier_scheduler,
78            barrier_manager,
79            stream_manager,
80            metadata_manager,
81            refresh_manager,
82            iceberg_compaction_manager,
83        }
84    }
85}
86
87fn effective_streaming_job_parallelism(
88    job_status: JobStatus,
89    parallelism: StreamingParallelism,
90    adaptive_parallelism_strategy: Option<String>,
91    backfill_parallelism: Option<StreamingParallelism>,
92    backfill_adaptive_parallelism_strategy: Option<String>,
93) -> (StreamingParallelism, Option<String>) {
94    if job_status != JobStatus::Created {
95        (
96            backfill_parallelism.unwrap_or(parallelism),
97            backfill_adaptive_parallelism_strategy.or(adaptive_parallelism_strategy),
98        )
99    } else {
100        (parallelism, adaptive_parallelism_strategy)
101    }
102}
103
104#[async_trait::async_trait]
105impl StreamManagerService for StreamServiceImpl {
106    async fn flush(&self, request: Request<FlushRequest>) -> TonicResponse<FlushResponse> {
107        self.env.idle_manager().record_activity();
108        let req = request.into_inner();
109
110        let version_id = self.barrier_scheduler.flush(req.database_id).await?;
111        Ok(Response::new(FlushResponse {
112            status: None,
113            hummock_version_id: version_id,
114        }))
115    }
116
117    async fn list_refresh_table_states(
118        &self,
119        _request: Request<ListRefreshTableStatesRequest>,
120    ) -> TonicResponse<ListRefreshTableStatesResponse> {
121        let refresh_jobs = self.metadata_manager.list_refresh_jobs().await?;
122        let refresh_table_states = refresh_jobs
123            .into_iter()
124            .map(|job| RefreshTableState {
125                table_id: job.table_id,
126                current_status: job.current_status.to_string(),
127                last_trigger_time: job
128                    .last_trigger_time
129                    .map(|time| DateTime::from_timestamp_millis(time).unwrap().to_string()),
130                trigger_interval_secs: job.trigger_interval_secs,
131                last_success_time: job
132                    .last_success_time
133                    .map(|time| DateTime::from_timestamp_millis(time).unwrap().to_string()),
134            })
135            .collect();
136        Ok(Response::new(ListRefreshTableStatesResponse {
137            states: refresh_table_states,
138        }))
139    }
140
141    async fn list_iceberg_compaction_status(
142        &self,
143        _request: Request<ListIcebergCompactionStatusRequest>,
144    ) -> TonicResponse<ListIcebergCompactionStatusResponse> {
145        let statuses = self
146            .iceberg_compaction_manager
147            .list_compaction_statuses()
148            .into_iter()
149            .map(
150                |status| list_iceberg_compaction_status_response::IcebergCompactionStatus {
151                    sink_id: status.sink_id.as_raw_id(),
152                    task_type: status.task_type,
153                    trigger_interval_sec: status.trigger_interval_sec,
154                    trigger_snapshot_count: status.trigger_snapshot_count as u64,
155                    schedule_state: status.schedule_state,
156                    next_compaction_after_sec: status.next_compaction_after_sec,
157                    pending_snapshot_count: status.pending_snapshot_count.map(|count| count as u64),
158                    is_triggerable: status.is_triggerable,
159                },
160            )
161            .collect();
162
163        Ok(Response::new(ListIcebergCompactionStatusResponse {
164            statuses,
165        }))
166    }
167
168    async fn pause(&self, _: Request<PauseRequest>) -> Result<Response<PauseResponse>, Status> {
169        for database_id in self.metadata_manager.list_active_database_ids().await? {
170            self.barrier_scheduler
171                .run_command(database_id, Command::pause())
172                .await?;
173        }
174        Ok(Response::new(PauseResponse {}))
175    }
176
177    async fn resume(&self, _: Request<ResumeRequest>) -> Result<Response<ResumeResponse>, Status> {
178        for database_id in self.metadata_manager.list_active_database_ids().await? {
179            self.barrier_scheduler
180                .run_command(database_id, Command::resume())
181                .await?;
182        }
183        Ok(Response::new(ResumeResponse {}))
184    }
185
186    async fn apply_throttle(
187        &self,
188        request: Request<ApplyThrottleRequest>,
189    ) -> Result<Response<ApplyThrottleResponse>, Status> {
190        let request = request.into_inner();
191
192        // Decode enums from raw i32 fields to handle decoupled target/type.
193        let throttle_target = request.throttle_target();
194        let throttle_type = request.throttle_type();
195
196        let raw_object_id: u32;
197        let fragments: HashMap<FragmentId, _>;
198
199        match (throttle_type, throttle_target) {
200            (ThrottleType::Source, ThrottleTarget::Source | ThrottleTarget::Table) => {
201                fragments = self
202                    .metadata_manager
203                    .update_source_rate_limit_by_source_id(request.id.into(), request.rate)
204                    .await?;
205                raw_object_id = request.id;
206            }
207            (ThrottleType::Backfill, ThrottleTarget::Mv)
208            | (ThrottleType::Backfill, ThrottleTarget::Sink)
209            | (ThrottleType::Backfill, ThrottleTarget::Table) => {
210                fragments = self
211                    .metadata_manager
212                    .update_backfill_rate_limit_by_job_id(JobId::from(request.id), request.rate)
213                    .await?;
214                raw_object_id = request.id;
215            }
216            (ThrottleType::Dml, ThrottleTarget::Table) => {
217                fragments = self
218                    .metadata_manager
219                    .update_dml_rate_limit_by_job_id(JobId::from(request.id), request.rate)
220                    .await?;
221                raw_object_id = request.id;
222            }
223            (ThrottleType::Sink, ThrottleTarget::Sink) => {
224                fragments = self
225                    .metadata_manager
226                    .update_sink_rate_limit_by_sink_id(request.id.into(), request.rate)
227                    .await?;
228                raw_object_id = request.id;
229            }
230            // FIXME(kwannoel): specialize for throttle type x target
231            (_, ThrottleTarget::Fragment) => {
232                let fragment_id = request.id.into();
233                let stream_node = self
234                    .metadata_manager
235                    .update_fragment_rate_limit_by_fragment_id(
236                        fragment_id,
237                        throttle_type,
238                        request.rate,
239                    )
240                    .await?;
241                fragments = [(fragment_id, stream_node)].into_iter().collect();
242                let job_id = self
243                    .metadata_manager
244                    .catalog_controller
245                    .get_fragment_streaming_job_id(fragment_id)
246                    .await?;
247                raw_object_id = job_id.as_raw_id();
248            }
249            _ => {
250                return Err(Status::invalid_argument(format!(
251                    "unsupported throttle target/type: {:?}/{:?}",
252                    throttle_target, throttle_type
253                )));
254            }
255        };
256
257        let database_id = self
258            .metadata_manager
259            .catalog_controller
260            .get_object_database_id(raw_object_id)
261            .await?;
262
263        let throttle_config = ThrottleConfig {
264            rate_limit: request.rate,
265            throttle_type: throttle_type.into(),
266        };
267        let config = fragments
268            .into_iter()
269            .map(|(fragment_id, stream_node)| (fragment_id, (throttle_config, stream_node)))
270            .collect();
271        let _i = self
272            .barrier_scheduler
273            .run_command(database_id, Command::Throttle { config })
274            .await?;
275
276        Ok(Response::new(ApplyThrottleResponse { status: None }))
277    }
278
279    async fn cancel_creating_jobs(
280        &self,
281        request: Request<CancelCreatingJobsRequest>,
282    ) -> TonicResponse<CancelCreatingJobsResponse> {
283        let req = request.into_inner();
284        let job_ids = match req.jobs.unwrap() {
285            Jobs::Infos(infos) => self
286                .metadata_manager
287                .catalog_controller
288                .find_creating_streaming_job_ids(infos.infos)
289                .await?
290                .into_iter()
291                .map(|id| id.as_job_id())
292                .collect(),
293            Jobs::Ids(jobs) => jobs.job_ids,
294        };
295
296        let canceled_jobs = self
297            .stream_manager
298            .cancel_streaming_jobs(job_ids)
299            .await?
300            .into_iter()
301            .map(|id| id.as_raw_id())
302            .collect_vec();
303        Ok(Response::new(CancelCreatingJobsResponse {
304            status: None,
305            canceled_jobs,
306        }))
307    }
308
309    async fn list_table_fragments(
310        &self,
311        request: Request<ListTableFragmentsRequest>,
312    ) -> Result<Response<ListTableFragmentsResponse>, Status> {
313        let req = request.into_inner();
314        let table_ids = HashSet::<JobId>::from_iter(req.table_ids);
315
316        let mut info = HashMap::new();
317        for job_id in table_ids {
318            let (table_fragments, fragment_actors, _actor_status) = self
319                .metadata_manager
320                .catalog_controller
321                .get_job_fragments_by_id(job_id)
322                .await?;
323            let mut dispatchers = self
324                .metadata_manager
325                .catalog_controller
326                .get_fragment_actor_dispatchers(
327                    table_fragments.fragment_ids().map(|id| id as _).collect(),
328                )
329                .await?;
330            let ctx = table_fragments.ctx.to_protobuf();
331            info.insert(
332                table_fragments.stream_job_id(),
333                TableFragmentInfo {
334                    fragments: table_fragments
335                        .fragments
336                        .into_iter()
337                        .map(|(id, fragment)| FragmentInfo {
338                            id,
339                            actors: fragment_actors
340                                .get(&id)
341                                .into_iter()
342                                .flat_map(|actors| actors.iter().map(|actor| actor.actor_id))
343                                .map(|actor_id| ActorInfo {
344                                    id: actor_id,
345                                    node: Some(fragment.nodes.clone()),
346                                    dispatcher: dispatchers
347                                        .get_mut(&fragment.fragment_id)
348                                        .and_then(|dispatchers| dispatchers.remove(&actor_id))
349                                        .unwrap_or_default(),
350                                })
351                                .collect_vec(),
352                        })
353                        .collect_vec(),
354                    ctx: Some(ctx),
355                },
356            );
357        }
358
359        Ok(Response::new(ListTableFragmentsResponse {
360            table_fragments: info,
361        }))
362    }
363
364    async fn list_streaming_job_states(
365        &self,
366        _request: Request<ListStreamingJobStatesRequest>,
367    ) -> Result<Response<ListStreamingJobStatesResponse>, Status> {
368        let job_infos = self
369            .metadata_manager
370            .catalog_controller
371            .list_streaming_job_infos()
372            .await?;
373        let states = job_infos
374            .into_iter()
375            .map(
376                |StreamingJobInfo {
377                     job_id,
378                     job_status,
379                     name,
380                     parallelism,
381                     adaptive_parallelism_strategy,
382                     backfill_parallelism,
383                     backfill_adaptive_parallelism_strategy,
384                     max_parallelism,
385                     resource_group,
386                     database_id,
387                     schema_id,
388                     config_override,
389                     ..
390                 }| {
391                    // While a job is still being created, system tables should surface the
392                    // temporary backfill override instead of the post-backfill target value.
393                    let (parallelism, adaptive_parallelism_strategy) =
394                        effective_streaming_job_parallelism(
395                            job_status,
396                            parallelism,
397                            adaptive_parallelism_strategy,
398                            backfill_parallelism,
399                            backfill_adaptive_parallelism_strategy,
400                        );
401                    let parallelism = match parallelism {
402                        StreamingParallelism::Adaptive => model::TableParallelism::Adaptive,
403                        StreamingParallelism::Custom => model::TableParallelism::Custom,
404                        StreamingParallelism::Fixed(n) => model::TableParallelism::Fixed(n as _),
405                    };
406
407                    list_streaming_job_states_response::StreamingJobState {
408                        table_id: job_id,
409                        name,
410                        state: PbState::from(job_status) as _,
411                        parallelism: Some(parallelism.into()),
412                        max_parallelism: max_parallelism as _,
413                        resource_group,
414                        database_id,
415                        schema_id,
416                        config_override,
417                        adaptive_parallelism_strategy,
418                    }
419                },
420            )
421            .collect_vec();
422
423        Ok(Response::new(ListStreamingJobStatesResponse { states }))
424    }
425
426    async fn list_fragment_distribution(
427        &self,
428        _request: Request<ListFragmentDistributionRequest>,
429    ) -> Result<Response<ListFragmentDistributionResponse>, Status> {
430        let include_node = _request.into_inner().include_node.unwrap_or(true);
431        let distributions = if include_node {
432            self.metadata_manager
433                .catalog_controller
434                .list_fragment_descs_with_node(false)
435                .await?
436        } else {
437            self.metadata_manager
438                .catalog_controller
439                .list_fragment_descs_without_node(false)
440                .await?
441        }
442        .into_iter()
443        .map(|(dist, _)| dist)
444        .collect();
445
446        Ok(Response::new(ListFragmentDistributionResponse {
447            distributions,
448        }))
449    }
450
451    async fn list_creating_fragment_distribution(
452        &self,
453        _request: Request<ListCreatingFragmentDistributionRequest>,
454    ) -> Result<Response<ListCreatingFragmentDistributionResponse>, Status> {
455        let include_node = _request.into_inner().include_node.unwrap_or(true);
456        let distributions = if include_node {
457            self.metadata_manager
458                .catalog_controller
459                .list_fragment_descs_with_node(true)
460                .await?
461        } else {
462            self.metadata_manager
463                .catalog_controller
464                .list_fragment_descs_without_node(true)
465                .await?
466        }
467        .into_iter()
468        .map(|(dist, _)| dist)
469        .collect();
470
471        Ok(Response::new(ListCreatingFragmentDistributionResponse {
472            distributions,
473        }))
474    }
475
476    async fn list_sink_log_store_tables(
477        &self,
478        _request: Request<ListSinkLogStoreTablesRequest>,
479    ) -> Result<Response<ListSinkLogStoreTablesResponse>, Status> {
480        let tables = self
481            .metadata_manager
482            .catalog_controller
483            .list_sink_log_store_tables()
484            .await?
485            .into_iter()
486            .map(|(sink_id, internal_table_id)| {
487                list_sink_log_store_tables_response::SinkLogStoreTable {
488                    sink_id: sink_id.as_raw_id(),
489                    internal_table_id: internal_table_id.as_raw_id(),
490                }
491            })
492            .collect();
493
494        Ok(Response::new(ListSinkLogStoreTablesResponse { tables }))
495    }
496
497    async fn get_fragment_by_id(
498        &self,
499        request: Request<GetFragmentByIdRequest>,
500    ) -> Result<Response<GetFragmentByIdResponse>, Status> {
501        let req = request.into_inner();
502        let fragment_desc = self
503            .metadata_manager
504            .catalog_controller
505            .get_fragment_desc_by_id(req.fragment_id)
506            .await?;
507        let distribution = fragment_desc
508            .map(|(desc, upstreams)| fragment_desc_to_distribution(desc, upstreams, true));
509        Ok(Response::new(GetFragmentByIdResponse { distribution }))
510    }
511
512    async fn get_fragment_vnodes(
513        &self,
514        request: Request<GetFragmentVnodesRequest>,
515    ) -> Result<Response<GetFragmentVnodesResponse>, Status> {
516        let req = request.into_inner();
517        let fragment_id = req.fragment_id;
518
519        let shared_actor_infos = self.env.shared_actor_infos();
520        let guard = shared_actor_infos.read_guard();
521
522        let fragment_info = guard
523            .get_fragment(fragment_id)
524            .ok_or_else(|| Status::not_found(format!("Fragment {} not found", fragment_id)))?;
525
526        let actor_vnodes = fragment_info
527            .actors
528            .iter()
529            .map(|(actor_id, actor_info)| {
530                let vnode_indices = if let Some(ref vnode_bitmap) = actor_info.vnode_bitmap {
531                    vnode_bitmap.iter_ones().map(|v| v as u32).collect()
532                } else {
533                    vec![]
534                };
535
536                get_fragment_vnodes_response::ActorVnodes {
537                    actor_id: *actor_id,
538                    vnode_indices,
539                }
540            })
541            .collect();
542
543        Ok(Response::new(GetFragmentVnodesResponse { actor_vnodes }))
544    }
545
546    async fn get_actor_vnodes(
547        &self,
548        request: Request<GetActorVnodesRequest>,
549    ) -> Result<Response<GetActorVnodesResponse>, Status> {
550        let req = request.into_inner();
551        let actor_id = req.actor_id;
552
553        let shared_actor_infos = self.env.shared_actor_infos();
554        let guard = shared_actor_infos.read_guard();
555
556        // Find the actor across all fragments
557        let actor_info = guard
558            .iter_over_fragments()
559            .find_map(|(_, fragment_info)| fragment_info.actors.get(&actor_id))
560            .ok_or_else(|| Status::not_found(format!("Actor {} not found", actor_id)))?;
561
562        let vnode_indices = if let Some(ref vnode_bitmap) = actor_info.vnode_bitmap {
563            vnode_bitmap.iter_ones().map(|v| v as u32).collect()
564        } else {
565            vec![]
566        };
567
568        Ok(Response::new(GetActorVnodesResponse { vnode_indices }))
569    }
570
571    async fn list_actor_states(
572        &self,
573        _request: Request<ListActorStatesRequest>,
574    ) -> Result<Response<ListActorStatesResponse>, Status> {
575        let actor_locations = self
576            .metadata_manager
577            .catalog_controller
578            .list_actor_locations()?;
579        let states = actor_locations
580            .into_iter()
581            .map(|actor_location| list_actor_states_response::ActorState {
582                actor_id: actor_location.actor_id,
583                fragment_id: actor_location.fragment_id,
584                worker_id: actor_location.worker_id,
585            })
586            .collect_vec();
587
588        Ok(Response::new(ListActorStatesResponse { states }))
589    }
590
591    async fn recover(
592        &self,
593        _request: Request<RecoverRequest>,
594    ) -> Result<Response<RecoverResponse>, Status> {
595        self.barrier_manager.adhoc_recovery().await?;
596        Ok(Response::new(RecoverResponse {}))
597    }
598
599    async fn list_actor_splits(
600        &self,
601        _request: Request<ListActorSplitsRequest>,
602    ) -> Result<Response<ListActorSplitsResponse>, Status> {
603        let SourceManagerRunningInfo {
604            source_fragments,
605            backfill_fragments,
606        } = self.stream_manager.source_manager.get_running_info().await;
607
608        let mut actor_splits = self.env.shared_actor_infos().list_assignments();
609
610        let source_actors: HashMap<_, _> = {
611            let all_fragment_ids: HashSet<_> = backfill_fragments
612                .values()
613                .flat_map(|set| set.iter().flat_map(|&(id1, id2)| [id1, id2]))
614                .chain(source_fragments.values().flatten().copied())
615                .collect();
616
617            let guard = self.env.shared_actor_infos().read_guard();
618            guard
619                .iter_over_fragments()
620                .filter(|(frag_id, _)| all_fragment_ids.contains(*frag_id))
621                .flat_map(|(fragment_id, fragment_info)| {
622                    fragment_info
623                        .actors
624                        .keys()
625                        .copied()
626                        .map(|actor_id| (actor_id, *fragment_id))
627                })
628                .collect()
629        };
630
631        let is_shared_source = self
632            .metadata_manager
633            .catalog_controller
634            .list_source_id_with_shared_types()
635            .await?;
636
637        let fragment_to_source: HashMap<_, _> = source_fragments
638            .into_iter()
639            .flat_map(|(source_id, fragment_ids)| {
640                let source_type = if is_shared_source.get(&source_id).copied().unwrap_or(false) {
641                    FragmentType::SharedSource
642                } else {
643                    FragmentType::NonSharedSource
644                };
645
646                fragment_ids
647                    .into_iter()
648                    .map(move |fragment_id| (fragment_id, (source_id, source_type)))
649            })
650            .chain(
651                backfill_fragments
652                    .into_iter()
653                    .flat_map(|(source_id, fragment_ids)| {
654                        fragment_ids.into_iter().flat_map(
655                            move |(fragment_id, upstream_fragment_id)| {
656                                [
657                                    (fragment_id, (source_id, FragmentType::SharedSourceBackfill)),
658                                    (
659                                        upstream_fragment_id,
660                                        (source_id, FragmentType::SharedSource),
661                                    ),
662                                ]
663                            },
664                        )
665                    }),
666            )
667            .collect();
668
669        let actor_splits = source_actors
670            .into_iter()
671            .flat_map(|(actor_id, fragment_id)| {
672                let (source_id, fragment_type) = fragment_to_source
673                    .get(&fragment_id)
674                    .copied()
675                    .unwrap_or_default();
676
677                actor_splits
678                    .remove(&actor_id)
679                    .unwrap_or_default()
680                    .into_iter()
681                    .map(move |split| list_actor_splits_response::ActorSplit {
682                        actor_id,
683                        source_id,
684                        fragment_id,
685                        split_id: split.id().to_string(),
686                        fragment_type: fragment_type.into(),
687                    })
688            })
689            .collect_vec();
690
691        Ok(Response::new(ListActorSplitsResponse { actor_splits }))
692    }
693
694    async fn list_rate_limits(
695        &self,
696        _request: Request<ListRateLimitsRequest>,
697    ) -> Result<Response<ListRateLimitsResponse>, Status> {
698        let rate_limits = self
699            .metadata_manager
700            .catalog_controller
701            .list_rate_limits()
702            .await?;
703        Ok(Response::new(ListRateLimitsResponse { rate_limits }))
704    }
705
706    #[cfg_attr(coverage, coverage(off))]
707    async fn refresh(
708        &self,
709        request: Request<RefreshRequest>,
710    ) -> Result<Response<RefreshResponse>, Status> {
711        let req = request.into_inner();
712
713        tracing::info!("Refreshing table with id: {}", req.table_id);
714
715        let response = self
716            .refresh_manager
717            .trigger_manual_refresh(req, self.env.shared_actor_infos())
718            .await?;
719
720        Ok(Response::new(response))
721    }
722
723    async fn alter_connector_props(
724        &self,
725        request: Request<AlterConnectorPropsRequest>,
726    ) -> Result<Response<AlterConnectorPropsResponse>, Status> {
727        let request = request.into_inner();
728        let secret_manager = LocalSecretManager::global();
729        let (new_props_plaintext, object_id) = match AlterConnectorPropsObject::try_from(
730            request.object_type,
731        ) {
732            Ok(AlterConnectorPropsObject::Sink) => (
733                self.metadata_manager
734                    .update_sink_props_by_sink_id(
735                        request.object_id.into(),
736                        request.changed_props.clone().into_iter().collect(),
737                    )
738                    .await?,
739                request.object_id.into(),
740            ),
741            Ok(AlterConnectorPropsObject::IcebergTable) => {
742                let (prop, sink_id) = self
743                    .metadata_manager
744                    .update_iceberg_table_props_by_table_id(
745                        request.object_id.into(),
746                        request.changed_props.clone().into_iter().collect(),
747                        request.extra_options,
748                    )
749                    .await?;
750                (prop, sink_id.as_object_id())
751            }
752
753            Ok(AlterConnectorPropsObject::Source) => {
754                // alter source and table's associated source
755                if request.connector_conn_ref.is_some() {
756                    return Err(Status::invalid_argument(
757                        "alter connector_conn_ref is not supported",
758                    ));
759                }
760                let options_with_secret = self
761                    .metadata_manager
762                    .catalog_controller
763                    .update_source_props_by_source_id(
764                        request.object_id.into(),
765                        request.changed_props.clone().into_iter().collect(),
766                        request.changed_secret_refs.clone().into_iter().collect(),
767                        false, // SQL ALTER SOURCE enforces alter-on-fly check
768                    )
769                    .await?;
770
771                self.stream_manager
772                    .source_manager
773                    .validate_source_once(request.object_id.into(), options_with_secret.clone())
774                    .await?;
775
776                let (options, secret_refs) = options_with_secret.into_parts();
777                (
778                    secret_manager
779                        .fill_secrets(options, secret_refs)
780                        .map_err(MetaError::from)?
781                        .into_iter()
782                        .collect(),
783                    request.object_id.into(),
784                )
785            }
786            Ok(AlterConnectorPropsObject::Connection) => {
787                // Update the connection and all dependent sources/sinks atomically, and later broadcast
788                // the complete plaintext properties to dependents via barrier.
789                let (
790                    connection_options_with_secret,
791                    updated_sources_with_props,
792                    updated_sinks_with_props,
793                ) = self
794                    .metadata_manager
795                    .catalog_controller
796                    .update_connection_and_dependent_objects_props(
797                        ConnectionId::from(request.object_id),
798                        request.changed_props.clone().into_iter().collect(),
799                        request.changed_secret_refs.clone().into_iter().collect(),
800                    )
801                    .await?;
802
803                // Materialize connection plaintext for observability/debugging (not broadcast directly).
804                let (options, secret_refs) = connection_options_with_secret.into_parts();
805                let new_props_plaintext = secret_manager
806                    .fill_secrets(options, secret_refs)
807                    .map_err(MetaError::from)?
808                    .into_iter()
809                    .collect::<HashMap<String, String>>();
810
811                // Broadcast changes to dependent sources and sinks if any exist.
812                let mut dependent_mutation = HashMap::default();
813                for (source_id, complete_source_props) in updated_sources_with_props {
814                    dependent_mutation.insert(source_id.as_object_id(), complete_source_props);
815                }
816                for (sink_id, complete_sink_props) in updated_sinks_with_props {
817                    dependent_mutation.insert(sink_id.as_object_id(), complete_sink_props);
818                }
819
820                if !dependent_mutation.is_empty() {
821                    let database_id = self
822                        .metadata_manager
823                        .catalog_controller
824                        .get_object_database_id(ConnectionId::from(request.object_id))
825                        .await?;
826                    tracing::info!(
827                        "broadcasting connection {} property changes to dependent object ids: {:?}",
828                        request.object_id,
829                        dependent_mutation.keys().collect_vec()
830                    );
831                    let _version = self
832                        .barrier_scheduler
833                        .run_command(
834                            database_id,
835                            Command::ConnectorPropsChange(dependent_mutation),
836                        )
837                        .await?;
838                }
839
840                (
841                    new_props_plaintext,
842                    ConnectionId::from(request.object_id).as_object_id(),
843                )
844            }
845
846            _ => {
847                unimplemented!(
848                    "Unsupported object type for AlterConnectorProps: {:?}",
849                    request.object_type
850                );
851            }
852        };
853
854        let database_id = self
855            .metadata_manager
856            .catalog_controller
857            .get_object_database_id(object_id)
858            .await?;
859        // Connection updates are broadcast to dependent sources/sinks inside the `Connection` branch above.
860        // For sources/sinks/iceberg-table updates, broadcast the change to the object itself.
861        if AlterConnectorPropsObject::try_from(request.object_type)
862            .is_ok_and(|t| t != AlterConnectorPropsObject::Connection)
863        {
864            let mut mutation = HashMap::default();
865            mutation.insert(object_id, new_props_plaintext);
866            let _version = self
867                .barrier_scheduler
868                .run_command(database_id, Command::ConnectorPropsChange(mutation))
869                .await?;
870        }
871
872        Ok(Response::new(AlterConnectorPropsResponse {}))
873    }
874
875    async fn set_sync_log_store_aligned(
876        &self,
877        request: Request<SetSyncLogStoreAlignedRequest>,
878    ) -> Result<Response<SetSyncLogStoreAlignedResponse>, Status> {
879        let req = request.into_inner();
880        let job_id = req.job_id;
881        let aligned = req.aligned;
882
883        self.metadata_manager
884            .catalog_controller
885            .mutate_fragments_by_job_id(
886                job_id,
887                |_mask, stream_node| {
888                    let mut visited = false;
889                    visit_stream_node_mut(stream_node, |body| {
890                        if let NodeBody::SyncLogStore(sync_log_store) = body {
891                            sync_log_store.aligned = aligned;
892                            visited = true
893                        }
894                    });
895                    Ok(visited)
896                },
897                "no fragments found with synced log store",
898            )
899            .await?;
900
901        Ok(Response::new(SetSyncLogStoreAlignedResponse {}))
902    }
903
904    async fn list_cdc_progress(
905        &self,
906        _request: Request<ListCdcProgressRequest>,
907    ) -> Result<Response<ListCdcProgressResponse>, Status> {
908        let cdc_progress = self
909            .barrier_manager
910            .get_cdc_progress()
911            .await?
912            .into_iter()
913            .map(|(job_id, p)| {
914                (
915                    job_id,
916                    PbCdcProgress {
917                        split_total_count: p.split_total_count,
918                        split_backfilled_count: p.split_backfilled_count,
919                        split_completed_count: p.split_completed_count,
920                    },
921                )
922            })
923            .collect();
924        Ok(Response::new(ListCdcProgressResponse { cdc_progress }))
925    }
926
927    async fn list_unmigrated_tables(
928        &self,
929        _request: Request<ListUnmigratedTablesRequest>,
930    ) -> Result<Response<ListUnmigratedTablesResponse>, Status> {
931        let unmigrated_tables = self
932            .metadata_manager
933            .catalog_controller
934            .list_unmigrated_tables()
935            .await?
936            .into_iter()
937            .map(|table| list_unmigrated_tables_response::UnmigratedTable {
938                table_id: table.id,
939                table_name: table.name,
940            })
941            .collect();
942
943        Ok(Response::new(ListUnmigratedTablesResponse {
944            tables: unmigrated_tables,
945        }))
946    }
947
948    /// Orchestrated source property update with pause/update/resume workflow.
949    /// This is the "safe" version that pauses sources before updating and resumes after.
950    async fn alter_source_properties_safe(
951        &self,
952        request: Request<AlterSourcePropertiesSafeRequest>,
953    ) -> Result<Response<AlterSourcePropertiesSafeResponse>, Status> {
954        let request = request.into_inner();
955        let source_id = request.source_id;
956        let options = request.options.unwrap_or_default();
957
958        tracing::info!(
959            source_id = source_id,
960            reset_splits = options.reset_splits,
961            "Starting orchestrated source property update"
962        );
963
964        // Get the database ID for the source
965        let database_id = self
966            .metadata_manager
967            .catalog_controller
968            .get_object_database_id(SourceId::from(source_id))
969            .await?;
970
971        // Step 1: Pause the stream (already commits state)
972        tracing::info!(source_id = source_id, "Pausing stream");
973        self.barrier_scheduler
974            .run_command(database_id, Command::Pause)
975            .await?;
976
977        // Step 2: Update catalog and get the new properties
978        let result = async {
979            let secret_manager = LocalSecretManager::global();
980
981            let options_with_secret = self
982                .metadata_manager
983                .catalog_controller
984                .update_source_props_by_source_id(
985                    source_id.into(),
986                    request.changed_props.clone().into_iter().collect(),
987                    request.changed_secret_refs.clone().into_iter().collect(),
988                    true, // risectl admin operation skips alter-on-fly check
989                )
990                .await?;
991
992            // Validate the source
993            self.stream_manager
994                .source_manager
995                .validate_source_once(source_id.into(), options_with_secret.clone())
996                .await?;
997
998            let (props, secret_refs) = options_with_secret.into_parts();
999            let new_props_plaintext: HashMap<String, String> = secret_manager
1000                .fill_secrets(props, secret_refs)
1001                .map_err(MetaError::from)?
1002                .into_iter()
1003                .collect();
1004
1005            // Step 3: Issue ConnectorPropsChange barrier
1006            tracing::info!(
1007                source_id = source_id,
1008                "Issuing ConnectorPropsChange barrier"
1009            );
1010            let mut mutation = HashMap::default();
1011            mutation.insert(source_id.into(), new_props_plaintext);
1012            self.barrier_scheduler
1013                .run_command(database_id, Command::ConnectorPropsChange(mutation))
1014                .await?;
1015
1016            // Step 4: Optional split reset
1017            if options.reset_splits {
1018                tracing::info!(source_id = source_id, "Resetting source splits");
1019                self.stream_manager
1020                    .source_manager
1021                    .reset_source_splits(source_id.into())
1022                    .await?;
1023            }
1024
1025            Ok::<_, MetaError>(())
1026        }
1027        .await;
1028
1029        // Step 5: Resume the stream (even if previous steps failed)
1030        tracing::info!(source_id = source_id, "Resuming stream");
1031        let resume_result = self
1032            .barrier_scheduler
1033            .run_command(database_id, Command::Resume)
1034            .await;
1035
1036        // Return the first error if any
1037        result?;
1038        resume_result?;
1039
1040        tracing::info!(
1041            source_id = source_id,
1042            "Orchestrated source property update completed successfully"
1043        );
1044
1045        Ok(Response::new(AlterSourcePropertiesSafeResponse {}))
1046    }
1047
1048    /// Reset source split assignments (UNSAFE - admin only).
1049    /// This clears persisted split metadata and triggers re-discovery.
1050    async fn reset_source_splits(
1051        &self,
1052        request: Request<ResetSourceSplitsRequest>,
1053    ) -> Result<Response<ResetSourceSplitsResponse>, Status> {
1054        let request = request.into_inner();
1055        let source_id = request.source_id;
1056
1057        tracing::warn!(
1058            source_id = source_id,
1059            "UNSAFE: Resetting source splits - this may cause data duplication or loss"
1060        );
1061
1062        self.stream_manager
1063            .source_manager
1064            .reset_source_splits(source_id.into())
1065            .await?;
1066
1067        Ok(Response::new(ResetSourceSplitsResponse {}))
1068    }
1069
1070    /// Inject specific offsets into source splits (UNSAFE - admin only).
1071    /// This can cause data duplication or loss depending on the correctness of the provided offsets.
1072    async fn inject_source_offsets(
1073        &self,
1074        request: Request<InjectSourceOffsetsRequest>,
1075    ) -> Result<Response<InjectSourceOffsetsResponse>, Status> {
1076        let request = request.into_inner();
1077        let source_id = request.source_id;
1078        let split_offsets = request.split_offsets;
1079
1080        // Validate split IDs exist before proceeding
1081        let applied_split_ids = self
1082            .stream_manager
1083            .source_manager
1084            .validate_inject_source_offsets(source_id.into(), &split_offsets)
1085            .await?;
1086
1087        tracing::warn!(
1088            source_id = source_id,
1089            num_offsets = split_offsets.len(),
1090            "UNSAFE: Injecting source offsets - this may cause data duplication or loss"
1091        );
1092
1093        let database_id = self
1094            .metadata_manager
1095            .catalog_controller
1096            .get_object_database_id(SourceId::from(source_id))
1097            .await?;
1098
1099        self.barrier_scheduler
1100            .run_command(
1101                database_id,
1102                Command::InjectSourceOffsets {
1103                    source_id: SourceId::from(source_id),
1104                    split_offsets,
1105                },
1106            )
1107            .await?;
1108
1109        Ok(Response::new(InjectSourceOffsetsResponse {
1110            applied_split_ids,
1111        }))
1112    }
1113}
1114
1115fn fragment_desc_to_distribution(
1116    fragment_desc: FragmentDesc,
1117    upstreams: Vec<FragmentId>,
1118    include_node: bool,
1119) -> FragmentDistribution {
1120    let node = include_node.then(|| fragment_desc.stream_node.to_protobuf());
1121    FragmentDistribution {
1122        fragment_id: fragment_desc.fragment_id,
1123        table_id: fragment_desc.job_id,
1124        distribution_type: PbFragmentDistributionType::from(fragment_desc.distribution_type) as _,
1125        state_table_ids: fragment_desc.state_table_ids.0,
1126        upstream_fragment_ids: upstreams,
1127        fragment_type_mask: fragment_desc.fragment_type_mask as _,
1128        parallelism: fragment_desc.parallelism as _,
1129        vnode_count: fragment_desc.vnode_count as _,
1130        node,
1131        parallelism_policy: fragment_desc.parallelism_policy,
1132    }
1133}
1134
1135#[cfg(test)]
1136mod tests {
1137    use risingwave_meta_model::{JobStatus, StreamingParallelism};
1138
1139    use super::effective_streaming_job_parallelism;
1140
1141    #[test]
1142    fn test_effective_streaming_job_parallelism_prefers_backfill_override() {
1143        let (parallelism, strategy) = effective_streaming_job_parallelism(
1144            JobStatus::Creating,
1145            StreamingParallelism::Adaptive,
1146            Some("BOUNDED(4)".to_owned()),
1147            Some(StreamingParallelism::Adaptive),
1148            Some("BOUNDED(2)".to_owned()),
1149        );
1150
1151        assert_eq!(parallelism, StreamingParallelism::Adaptive);
1152        assert_eq!(strategy.as_deref(), Some("BOUNDED(2)"));
1153    }
1154
1155    #[test]
1156    fn test_effective_streaming_job_parallelism_falls_back_to_job_strategy() {
1157        let (parallelism, strategy) = effective_streaming_job_parallelism(
1158            JobStatus::Initial,
1159            StreamingParallelism::Adaptive,
1160            Some("RATIO(0.5)".to_owned()),
1161            Some(StreamingParallelism::Adaptive),
1162            None,
1163        );
1164
1165        assert_eq!(parallelism, StreamingParallelism::Adaptive);
1166        assert_eq!(strategy.as_deref(), Some("RATIO(0.5)"));
1167    }
1168
1169    #[test]
1170    fn test_effective_streaming_job_parallelism_ignores_backfill_after_creation() {
1171        let (parallelism, strategy) = effective_streaming_job_parallelism(
1172            JobStatus::Created,
1173            StreamingParallelism::Fixed(4),
1174            None,
1175            Some(StreamingParallelism::Fixed(2)),
1176            Some("BOUNDED(2)".to_owned()),
1177        );
1178
1179        assert_eq!(parallelism, StreamingParallelism::Fixed(4));
1180        assert_eq!(strategy, None);
1181    }
1182}