1use 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 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 (_, 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 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 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 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, )
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 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 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 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 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 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 let database_id = self
966 .metadata_manager
967 .catalog_controller
968 .get_object_database_id(SourceId::from(source_id))
969 .await?;
970
971 tracing::info!(source_id = source_id, "Pausing stream");
973 self.barrier_scheduler
974 .run_command(database_id, Command::Pause)
975 .await?;
976
977 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, )
990 .await?;
991
992 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 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 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 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 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 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 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 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}