Skip to main content

risingwave_meta/manager/
metadata.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{BTreeMap, HashMap, HashSet};
16use std::fmt::{Debug, Formatter};
17
18use anyhow::anyhow;
19use itertools::Itertools;
20use risingwave_common::catalog::{DatabaseId, TableId, TableOption};
21use risingwave_common::id::JobId;
22use risingwave_meta_model::refresh_job::{self, RefreshState};
23use risingwave_meta_model::{SinkId, SourceId, WorkerId};
24use risingwave_pb::catalog::{PbSource, PbTable};
25use risingwave_pb::common::worker_node::{PbResource, Property as AddNodeProperty, State};
26use risingwave_pb::common::{HostAddress, PbWorkerNode, PbWorkerType, WorkerNode, WorkerType};
27use risingwave_pb::meta::list_rate_limits_response::RateLimitInfo;
28use risingwave_pb::stream_plan::{PbDispatcherType, PbStreamNode, PbStreamScanType};
29use sea_orm::TransactionTrait;
30use sea_orm::prelude::DateTime;
31use tokio::sync::mpsc::{UnboundedReceiver, unbounded_channel};
32use tokio::sync::oneshot;
33use tracing::warn;
34
35use crate::MetaResult;
36use crate::controller::catalog::CatalogControllerRef;
37use crate::controller::cluster::{ClusterControllerRef, StreamingClusterInfo, WorkerExtraInfo};
38use crate::controller::scale::find_fragment_no_shuffle_dags_detailed;
39use crate::manager::{LocalNotification, NotificationVersion};
40use crate::model::{ActorId, ClusterId, Fragment, FragmentId, StreamJobFragments, SubscriptionId};
41use crate::stream::SplitAssignment;
42use crate::telemetry::MetaTelemetryJobDesc;
43
44#[derive(Clone)]
45pub struct MetadataManager {
46    pub cluster_controller: ClusterControllerRef,
47    pub catalog_controller: CatalogControllerRef,
48}
49
50#[derive(Debug)]
51pub(crate) enum ActiveStreamingWorkerChange {
52    Add(WorkerNode),
53    Remove(WorkerNode),
54    Update(WorkerNode),
55}
56
57pub struct ActiveStreamingWorkerNodes {
58    worker_nodes: HashMap<WorkerId, WorkerNode>,
59    rx: UnboundedReceiver<LocalNotification>,
60    #[cfg_attr(not(debug_assertions), expect(dead_code))]
61    meta_manager: Option<MetadataManager>,
62}
63
64impl Debug for ActiveStreamingWorkerNodes {
65    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
66        f.debug_struct("ActiveStreamingWorkerNodes")
67            .field("worker_nodes", &self.worker_nodes)
68            .finish()
69    }
70}
71
72impl ActiveStreamingWorkerNodes {
73    pub(crate) fn uninitialized() -> Self {
74        Self {
75            worker_nodes: Default::default(),
76            rx: unbounded_channel().1,
77            meta_manager: None,
78        }
79    }
80
81    #[cfg(test)]
82    pub(crate) fn for_test(worker_nodes: HashMap<WorkerId, WorkerNode>) -> Self {
83        let (tx, rx) = unbounded_channel();
84        let _join_handle = tokio::spawn(async move {
85            let _tx = tx;
86            std::future::pending::<()>().await
87        });
88        Self {
89            worker_nodes,
90            rx,
91            meta_manager: None,
92        }
93    }
94
95    /// Return an uninitialized one as a placeholder for future initialized
96    pub(crate) async fn new_snapshot(meta_manager: MetadataManager) -> MetaResult<Self> {
97        let (nodes, rx) = meta_manager
98            .subscribe_active_streaming_compute_nodes()
99            .await?;
100        Ok(Self {
101            worker_nodes: nodes
102                .into_iter()
103                .filter_map(|node| {
104                    let is_streaming = node.property.as_ref().is_some_and(|p| p.is_streaming);
105                    is_streaming.then_some((node.id, node))
106                })
107                .collect(),
108            rx,
109            meta_manager: Some(meta_manager),
110        })
111    }
112
113    pub(crate) fn current(&self) -> &HashMap<WorkerId, WorkerNode> {
114        &self.worker_nodes
115    }
116
117    pub(crate) async fn changed(&mut self) -> ActiveStreamingWorkerChange {
118        loop {
119            let notification = self
120                .rx
121                .recv()
122                .await
123                .expect("notification stopped or uninitialized");
124            fn is_target_worker_node(worker: &WorkerNode) -> bool {
125                worker.r#type == WorkerType::ComputeNode as i32
126                    && worker.property.as_ref().unwrap().is_streaming
127            }
128            match notification {
129                LocalNotification::WorkerNodeDeleted(worker) => {
130                    let is_target_worker_node = is_target_worker_node(&worker);
131                    let Some(prev_worker) = self.worker_nodes.remove(&worker.id) else {
132                        if is_target_worker_node {
133                            warn!(
134                                ?worker,
135                                "notify to delete an non-existing streaming compute worker"
136                            );
137                        }
138                        continue;
139                    };
140                    if !is_target_worker_node {
141                        warn!(
142                            ?worker,
143                            ?prev_worker,
144                            "deleted worker has a different recent type"
145                        );
146                    }
147                    if worker.state == State::Starting as i32 {
148                        warn!(
149                            id = %worker.id,
150                            host = ?worker.host,
151                            state = worker.state,
152                            "a starting streaming worker is deleted"
153                        );
154                    }
155                    break ActiveStreamingWorkerChange::Remove(prev_worker);
156                }
157                LocalNotification::WorkerNodeActivated(worker) => {
158                    if !is_target_worker_node(&worker) {
159                        if let Some(prev_worker) = self.worker_nodes.remove(&worker.id) {
160                            warn!(
161                                ?worker,
162                                ?prev_worker,
163                                "the type of a streaming worker is changed"
164                            );
165                            break ActiveStreamingWorkerChange::Remove(prev_worker);
166                        } else {
167                            continue;
168                        }
169                    }
170                    assert_eq!(
171                        worker.state,
172                        State::Running as i32,
173                        "not started worker added: {:?}",
174                        worker
175                    );
176                    if let Some(prev_worker) = self.worker_nodes.insert(worker.id, worker.clone()) {
177                        assert_eq!(prev_worker.host, worker.host);
178                        assert_eq!(prev_worker.r#type, worker.r#type);
179                        warn!(
180                            ?prev_worker,
181                            ?worker,
182                            eq = prev_worker == worker,
183                            "notify to update an existing active worker"
184                        );
185                        if prev_worker == worker {
186                            continue;
187                        } else {
188                            break ActiveStreamingWorkerChange::Update(worker);
189                        }
190                    } else {
191                        break ActiveStreamingWorkerChange::Add(worker);
192                    }
193                }
194                _ => {
195                    continue;
196                }
197            }
198        }
199    }
200
201    #[cfg(debug_assertions)]
202    pub(crate) async fn validate_change(&self) {
203        use risingwave_pb::common::WorkerNode;
204        use thiserror_ext::AsReport;
205        let Some(meta_manager) = &self.meta_manager else {
206            return;
207        };
208        match meta_manager.list_active_streaming_compute_nodes().await {
209            Ok(worker_nodes) => {
210                let ignore_irrelevant_info = |node: &WorkerNode| {
211                    (
212                        node.id,
213                        WorkerNode {
214                            id: node.id,
215                            r#type: node.r#type,
216                            host: node.host.clone(),
217                            property: node.property.clone(),
218                            resource: node.resource.clone(),
219                            ..Default::default()
220                        },
221                    )
222                };
223                let worker_nodes: HashMap<_, _> =
224                    worker_nodes.iter().map(ignore_irrelevant_info).collect();
225                let curr_worker_nodes: HashMap<_, _> = self
226                    .current()
227                    .values()
228                    .map(ignore_irrelevant_info)
229                    .collect();
230                if worker_nodes != curr_worker_nodes {
231                    warn!(
232                        ?worker_nodes,
233                        ?curr_worker_nodes,
234                        "different to global snapshot"
235                    );
236                }
237            }
238            Err(e) => {
239                warn!(
240                    e = ?e.as_report(),
241                    "failed to list active streaming compute nodes for comparison with the local snapshot",
242                );
243            }
244        }
245    }
246}
247
248impl MetadataManager {
249    pub fn new(
250        cluster_controller: ClusterControllerRef,
251        catalog_controller: CatalogControllerRef,
252    ) -> Self {
253        Self {
254            cluster_controller,
255            catalog_controller,
256        }
257    }
258
259    pub async fn get_worker_by_id(&self, worker_id: WorkerId) -> MetaResult<Option<PbWorkerNode>> {
260        self.cluster_controller.get_worker_by_id(worker_id).await
261    }
262
263    pub async fn count_worker_node(&self) -> MetaResult<HashMap<WorkerType, u64>> {
264        let node_map = self.cluster_controller.count_worker_by_type().await?;
265        Ok(node_map
266            .into_iter()
267            .map(|(ty, cnt)| (ty.into(), cnt as u64))
268            .collect())
269    }
270
271    pub async fn get_worker_info_by_id(&self, worker_id: WorkerId) -> Option<WorkerExtraInfo> {
272        self.cluster_controller
273            .get_worker_info_by_id(worker_id as _)
274            .await
275    }
276
277    pub async fn add_worker_node(
278        &self,
279        r#type: PbWorkerType,
280        host_address: HostAddress,
281        property: AddNodeProperty,
282        resource: PbResource,
283    ) -> MetaResult<WorkerId> {
284        self.cluster_controller
285            .add_worker(r#type, host_address, property, resource)
286            .await
287            .map(|id| id as WorkerId)
288    }
289
290    pub async fn list_worker_node(
291        &self,
292        worker_type: Option<WorkerType>,
293        worker_state: Option<State>,
294    ) -> MetaResult<Vec<PbWorkerNode>> {
295        self.cluster_controller
296            .list_workers(worker_type.map(Into::into), worker_state.map(Into::into))
297            .await
298    }
299
300    pub async fn subscribe_active_streaming_compute_nodes(
301        &self,
302    ) -> MetaResult<(Vec<WorkerNode>, UnboundedReceiver<LocalNotification>)> {
303        self.cluster_controller
304            .subscribe_active_streaming_compute_nodes()
305            .await
306    }
307
308    pub async fn list_active_streaming_compute_nodes(&self) -> MetaResult<Vec<PbWorkerNode>> {
309        self.cluster_controller
310            .list_active_streaming_workers()
311            .await
312    }
313
314    pub async fn list_active_serving_compute_nodes(&self) -> MetaResult<Vec<PbWorkerNode>> {
315        self.cluster_controller.list_active_serving_workers().await
316    }
317
318    pub async fn list_active_database_ids(&self) -> MetaResult<HashSet<DatabaseId>> {
319        Ok(self
320            .catalog_controller
321            .list_fragment_database_ids(None)
322            .await?
323            .into_iter()
324            .map(|(_, database_id)| database_id)
325            .collect())
326    }
327
328    pub async fn split_fragment_map_by_database<T: Debug>(
329        &self,
330        fragment_map: HashMap<FragmentId, T>,
331    ) -> MetaResult<HashMap<DatabaseId, HashMap<FragmentId, T>>> {
332        let fragment_to_database_map: HashMap<_, _> = self
333            .catalog_controller
334            .list_fragment_database_ids(Some(
335                fragment_map
336                    .keys()
337                    .map(|fragment_id| *fragment_id as _)
338                    .collect(),
339            ))
340            .await?
341            .into_iter()
342            .map(|(fragment_id, database_id)| (fragment_id as FragmentId, database_id))
343            .collect();
344        let mut ret: HashMap<_, HashMap<_, _>> = HashMap::new();
345        for (fragment_id, value) in fragment_map {
346            let database_id = *fragment_to_database_map
347                .get(&fragment_id)
348                .ok_or_else(|| anyhow!("cannot get database_id of fragment {fragment_id}"))?;
349            ret.entry(database_id)
350                .or_default()
351                .try_insert(fragment_id, value)
352                .expect("non duplicate");
353        }
354        Ok(ret)
355    }
356
357    pub async fn list_creating_jobs(&self) -> MetaResult<HashSet<JobId>> {
358        Ok(self
359            .catalog_controller
360            .list_creating_jobs(false, None)
361            .await?
362            .into_iter()
363            .map(|(job_id, _, _, _, _)| job_id)
364            .collect())
365    }
366
367    pub async fn list_sources(&self) -> MetaResult<Vec<PbSource>> {
368        self.catalog_controller.list_sources().await
369    }
370
371    /// Get and filter the "**root**" fragments of the specified relations.
372    /// The root fragment is the bottom-most fragment of its fragment graph, and can be a `MView` or a `Source`.
373    ///
374    /// See also [`crate::controller::catalog::CatalogController::get_root_fragments`].
375    pub async fn get_upstream_root_fragments(
376        &self,
377        upstream_table_ids: &HashSet<TableId>,
378    ) -> MetaResult<HashMap<JobId, Fragment>> {
379        let upstream_root_fragments = self
380            .catalog_controller
381            .get_root_fragments(upstream_table_ids.iter().map(|id| id.as_job_id()).collect())
382            .await?;
383
384        Ok(upstream_root_fragments)
385    }
386
387    pub async fn get_streaming_cluster_info(&self) -> MetaResult<StreamingClusterInfo> {
388        self.cluster_controller.get_streaming_cluster_info().await
389    }
390
391    pub async fn get_all_table_options(&self) -> MetaResult<HashMap<TableId, TableOption>> {
392        self.catalog_controller.get_all_table_options().await
393    }
394
395    pub async fn get_table_name_type_mapping(
396        &self,
397    ) -> MetaResult<HashMap<TableId, (String, String)>> {
398        self.catalog_controller.get_table_name_type_mapping().await
399    }
400
401    pub async fn get_created_table_ids(&self) -> MetaResult<Vec<TableId>> {
402        self.catalog_controller.get_created_table_ids().await
403    }
404
405    pub async fn get_table_associated_source_id(
406        &self,
407        table_id: TableId,
408    ) -> MetaResult<Option<SourceId>> {
409        self.catalog_controller
410            .get_table_associated_source_id(table_id)
411            .await
412    }
413
414    pub async fn get_table_catalog_by_ids(&self, ids: &[TableId]) -> MetaResult<Vec<PbTable>> {
415        self.catalog_controller
416            .get_table_by_ids(ids.to_vec(), false)
417            .await
418    }
419
420    pub async fn list_refresh_jobs(&self) -> MetaResult<Vec<refresh_job::Model>> {
421        self.catalog_controller.list_refresh_jobs().await
422    }
423
424    pub async fn list_refreshable_table_ids(&self) -> MetaResult<Vec<TableId>> {
425        self.catalog_controller.list_refreshable_table_ids().await
426    }
427
428    pub async fn ensure_refresh_job(&self, table_id: TableId) -> MetaResult<()> {
429        self.catalog_controller.ensure_refresh_job(table_id).await
430    }
431
432    pub async fn update_refresh_job_status(
433        &self,
434        table_id: TableId,
435        status: RefreshState,
436        trigger_time: Option<DateTime>,
437        is_success: bool,
438    ) -> MetaResult<()> {
439        self.catalog_controller
440            .update_refresh_job_status(table_id, status, trigger_time, is_success)
441            .await
442    }
443
444    pub async fn reset_all_refresh_jobs_to_idle(&self) -> MetaResult<()> {
445        self.catalog_controller
446            .reset_all_refresh_jobs_to_idle()
447            .await
448    }
449
450    pub async fn update_refresh_job_interval(
451        &self,
452        table_id: TableId,
453        trigger_interval_secs: Option<i64>,
454    ) -> MetaResult<()> {
455        self.catalog_controller
456            .update_refresh_job_interval(table_id, trigger_interval_secs)
457            .await
458    }
459
460    pub async fn get_sink_state_table_ids(&self, sink_id: SinkId) -> MetaResult<Vec<TableId>> {
461        self.catalog_controller
462            .get_sink_state_table_ids(sink_id)
463            .await
464    }
465
466    pub async fn get_table_catalog_by_cdc_table_id(
467        &self,
468        cdc_table_id: &String,
469    ) -> MetaResult<Vec<PbTable>> {
470        self.catalog_controller
471            .get_table_by_cdc_table_id(cdc_table_id)
472            .await
473    }
474
475    pub async fn get_downstream_fragments(
476        &self,
477        job_id: JobId,
478    ) -> MetaResult<Vec<(PbDispatcherType, Fragment)>> {
479        self.catalog_controller
480            .get_downstream_fragments(job_id)
481            .await
482    }
483
484    pub async fn get_job_id_to_internal_table_ids_mapping(
485        &self,
486    ) -> Option<Vec<(JobId, Vec<TableId>)>> {
487        self.catalog_controller.get_job_internal_table_ids().await
488    }
489
490    pub async fn get_job_fragments_by_id(&self, job_id: JobId) -> MetaResult<StreamJobFragments> {
491        Ok(self
492            .catalog_controller
493            .get_job_fragments_by_id(job_id)
494            .await?
495            .0)
496    }
497
498    pub fn get_running_actors_of_fragment(&self, id: FragmentId) -> MetaResult<HashSet<ActorId>> {
499        self.catalog_controller
500            .get_running_actors_of_fragment(id as _)
501    }
502
503    // (backfill_actor_id, upstream_source_actor_id)
504    pub async fn get_running_actors_for_source_backfill(
505        &self,
506        source_backfill_fragment_id: FragmentId,
507        source_fragment_id: FragmentId,
508    ) -> MetaResult<HashSet<(ActorId, ActorId)>> {
509        let actor_ids = self
510            .catalog_controller
511            .get_running_actors_for_source_backfill(
512                source_backfill_fragment_id as _,
513                source_fragment_id as _,
514            )
515            .await?;
516        Ok(actor_ids
517            .into_iter()
518            .map(|(id, upstream)| (id as ActorId, upstream as ActorId))
519            .collect())
520    }
521
522    pub fn worker_actor_count(&self) -> MetaResult<HashMap<WorkerId, usize>> {
523        let actor_cnt = self.catalog_controller.worker_actor_count()?;
524        Ok(actor_cnt
525            .into_iter()
526            .map(|(id, cnt)| (id as WorkerId, cnt))
527            .collect())
528    }
529
530    pub async fn count_streaming_job(&self) -> MetaResult<usize> {
531        self.catalog_controller.count_streaming_jobs().await
532    }
533
534    pub async fn list_stream_job_desc(&self) -> MetaResult<Vec<MetaTelemetryJobDesc>> {
535        self.catalog_controller
536            .list_stream_job_desc_for_telemetry()
537            .await
538    }
539
540    pub async fn update_source_rate_limit_by_source_id(
541        &self,
542        source_id: SourceId,
543        rate_limit: Option<u32>,
544    ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
545        self.catalog_controller
546            .update_source_rate_limit_by_source_id(source_id as _, rate_limit)
547            .await
548    }
549
550    pub async fn update_backfill_rate_limit_by_job_id(
551        &self,
552        job_id: JobId,
553        rate_limit: Option<u32>,
554    ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
555        self.catalog_controller
556            .update_backfill_rate_limit_by_job_id(job_id, rate_limit)
557            .await
558    }
559
560    pub async fn update_sink_rate_limit_by_sink_id(
561        &self,
562        sink_id: SinkId,
563        rate_limit: Option<u32>,
564    ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
565        self.catalog_controller
566            .update_sink_rate_limit_by_job_id(sink_id, rate_limit)
567            .await
568    }
569
570    pub async fn update_dml_rate_limit_by_job_id(
571        &self,
572        job_id: JobId,
573        rate_limit: Option<u32>,
574    ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
575        self.catalog_controller
576            .update_dml_rate_limit_by_job_id(job_id, rate_limit)
577            .await
578    }
579
580    pub async fn update_sink_props_by_sink_id(
581        &self,
582        sink_id: SinkId,
583        props: BTreeMap<String, String>,
584    ) -> MetaResult<HashMap<String, String>> {
585        let new_props = self
586            .catalog_controller
587            .update_sink_props_by_sink_id(sink_id, props)
588            .await?;
589        Ok(new_props)
590    }
591
592    pub async fn update_iceberg_table_props_by_table_id(
593        &self,
594        table_id: TableId,
595        props: BTreeMap<String, String>,
596        alter_iceberg_table_props: Option<
597            risingwave_pb::meta::alter_connector_props_request::PbExtraOptions,
598        >,
599    ) -> MetaResult<(HashMap<String, String>, SinkId)> {
600        let (new_props, sink_id) = self
601            .catalog_controller
602            .update_iceberg_table_props_by_table_id(table_id, props, alter_iceberg_table_props)
603            .await?;
604        Ok((new_props, sink_id))
605    }
606
607    pub async fn update_fragment_rate_limit_by_fragment_id(
608        &self,
609        fragment_id: FragmentId,
610        throttle_type: risingwave_pb::common::ThrottleType,
611        rate_limit: Option<u32>,
612    ) -> MetaResult<PbStreamNode> {
613        self.catalog_controller
614            .update_fragment_rate_limit_by_fragment_id(fragment_id as _, throttle_type, rate_limit)
615            .await
616    }
617
618    #[await_tree::instrument]
619    pub async fn update_fragment_splits(
620        &self,
621        split_assignment: &SplitAssignment,
622    ) -> MetaResult<()> {
623        let fragment_splits = split_assignment
624            .iter()
625            .map(|(fragment_id, splits)| {
626                (
627                    *fragment_id as _,
628                    splits.values().flatten().cloned().collect_vec(),
629                )
630            })
631            .collect();
632
633        let inner = self.catalog_controller.inner.write().await;
634
635        self.catalog_controller
636            .update_fragment_splits(&inner.db, &fragment_splits)
637            .await
638    }
639
640    pub async fn get_mv_depended_subscriptions(
641        &self,
642        database_id: Option<DatabaseId>,
643    ) -> MetaResult<HashMap<TableId, HashMap<SubscriptionId, u64>>> {
644        Ok(self
645            .catalog_controller
646            .get_mv_depended_subscriptions(database_id)
647            .await?
648            .into_iter()
649            .map(|(table_id, subscriptions)| {
650                (
651                    table_id,
652                    subscriptions
653                        .into_iter()
654                        .map(|(subscription_id, retention_time)| {
655                            (subscription_id as SubscriptionId, retention_time)
656                        })
657                        .collect(),
658                )
659            })
660            .collect())
661    }
662
663    pub async fn get_job_max_parallelism(&self, job_id: JobId) -> MetaResult<usize> {
664        self.catalog_controller
665            .get_max_parallelism_by_id(job_id)
666            .await
667    }
668
669    pub async fn get_existing_job_resource_group(
670        &self,
671        streaming_job_id: JobId,
672    ) -> MetaResult<String> {
673        self.catalog_controller
674            .get_existing_job_resource_group(streaming_job_id)
675            .await
676    }
677
678    pub async fn get_database_resource_group(&self, database_id: DatabaseId) -> MetaResult<String> {
679        self.catalog_controller
680            .get_database_resource_group(database_id)
681            .await
682    }
683
684    pub fn cluster_id(&self) -> &ClusterId {
685        self.cluster_controller.cluster_id()
686    }
687
688    pub async fn list_rate_limits(&self) -> MetaResult<Vec<RateLimitInfo>> {
689        let rate_limits = self.catalog_controller.list_rate_limits().await?;
690        Ok(rate_limits)
691    }
692
693    pub async fn get_job_backfill_scan_types(
694        &self,
695        job_id: JobId,
696    ) -> MetaResult<HashMap<FragmentId, PbStreamScanType>> {
697        let backfill_types = self
698            .catalog_controller
699            .get_job_fragment_backfill_scan_type(job_id)
700            .await?;
701        Ok(backfill_types)
702    }
703
704    /// Returns jobs containing a scan that cannot be rescheduled online.
705    pub async fn collect_online_unreschedulable_backfill_jobs(
706        &self,
707        job_ids: impl IntoIterator<Item = &JobId>,
708    ) -> MetaResult<HashSet<JobId>> {
709        let mut unreschedulable = HashSet::new();
710
711        for job_id in job_ids {
712            let scan_types = self
713                .catalog_controller
714                .get_job_fragment_backfill_scan_type(*job_id)
715                .await?;
716            if scan_types
717                .values()
718                .any(|scan_type| !scan_type.is_reschedulable(true))
719            {
720                unreschedulable.insert(*job_id);
721            }
722        }
723
724        Ok(unreschedulable)
725    }
726
727    pub async fn collect_reschedule_blocked_jobs_for_creating_jobs(
728        &self,
729        creating_job_ids: impl IntoIterator<Item = &JobId>,
730        is_online: bool,
731    ) -> MetaResult<HashSet<JobId>> {
732        let creating_job_ids: HashSet<_> = creating_job_ids.into_iter().copied().collect();
733        if creating_job_ids.is_empty() {
734            return Ok(HashSet::new());
735        }
736
737        let inner = self.catalog_controller.inner.read().await;
738        let txn = inner.db.begin().await?;
739
740        let mut initial_fragment_ids = HashSet::new();
741        for job_id in &creating_job_ids {
742            let scan_types = self
743                .catalog_controller
744                .get_job_fragment_backfill_scan_type_in_txn(&txn, *job_id)
745                .await?;
746            initial_fragment_ids.extend(scan_types.into_iter().filter_map(
747                |(fragment_id, scan_type)| {
748                    (!scan_type.is_reschedulable(is_online)).then_some(fragment_id)
749                },
750            ));
751        }
752
753        if !initial_fragment_ids.is_empty() {
754            let upstream_fragments = self
755                .catalog_controller
756                .upstream_fragments_in_txn(&txn, initial_fragment_ids.iter().copied())
757                .await?;
758            initial_fragment_ids.extend(upstream_fragments.into_values().flatten());
759        }
760
761        let mut blocked_fragment_ids = initial_fragment_ids.clone();
762        if !initial_fragment_ids.is_empty() {
763            let initial_fragment_ids = initial_fragment_ids.into_iter().collect_vec();
764            let ensembles =
765                find_fragment_no_shuffle_dags_detailed(&txn, &initial_fragment_ids).await?;
766            for ensemble in ensembles {
767                blocked_fragment_ids.extend(ensemble.fragments());
768            }
769        }
770
771        let mut blocked_job_ids = HashSet::new();
772        if !blocked_fragment_ids.is_empty() {
773            let fragment_ids = blocked_fragment_ids.into_iter().collect_vec();
774            let fragment_job_ids = self
775                .catalog_controller
776                .get_fragment_job_id_in_txn(&txn, fragment_ids)
777                .await?;
778            blocked_job_ids.extend(
779                fragment_job_ids
780                    .into_iter()
781                    .map(|job_id| job_id.as_job_id()),
782            );
783        }
784
785        txn.commit().await?;
786
787        Ok(blocked_job_ids)
788    }
789}
790
791impl MetadataManager {
792    /// Wait for job finishing notification in `TrackingJob::finish`.
793    /// The progress is updated per barrier.
794    #[await_tree::instrument]
795    pub async fn wait_streaming_job_finished(
796        &self,
797        database_id: DatabaseId,
798        id: JobId,
799    ) -> MetaResult<NotificationVersion> {
800        tracing::debug!("wait_streaming_job_finished: {id:?}");
801        let mut mgr = self.catalog_controller.get_inner_write_guard().await;
802        if mgr.streaming_job_is_finished(id).await? {
803            return Ok(self.catalog_controller.notify_frontend_trivial().await);
804        }
805        let (tx, rx) = oneshot::channel();
806
807        mgr.register_finish_notifier(database_id, id, tx);
808        drop(mgr);
809        rx.await
810            .map_err(|_| "no received reason".to_owned())
811            .and_then(|result| result)
812            .map_err(|reason| anyhow!("failed to wait streaming job finish: {}", reason).into())
813    }
814
815    pub(crate) async fn notify_finish_failed(&self, database_id: Option<DatabaseId>, err: String) {
816        let mut mgr = self.catalog_controller.get_inner_write_guard().await;
817        mgr.notify_finish_failed(database_id, err);
818    }
819}