Skip to main content

risingwave_meta/barrier/
info.rs

1// Copyright 2022 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::hash_map::Entry;
16use std::collections::{HashMap, HashSet};
17use std::mem::replace;
18use std::sync::Arc;
19
20use itertools::Itertools;
21use parking_lot::RawRwLock;
22use parking_lot::lock_api::RwLockReadGuard;
23use risingwave_common::bitmap::Bitmap;
24use risingwave_common::catalog::{DatabaseId, FragmentTypeFlag, FragmentTypeMask, TableId};
25use risingwave_common::id::JobId;
26use risingwave_common::util::epoch::EpochPair;
27use risingwave_common::util::stream_graph_visitor::visit_stream_node_mut;
28use risingwave_connector::source::{SplitImpl, SplitMetaData};
29use risingwave_meta_model::WorkerId;
30use risingwave_meta_model::fragment::DistributionType;
31use risingwave_pb::ddl_service::PbBackfillType;
32use risingwave_pb::hummock::HummockVersionStats;
33use risingwave_pb::id::SubscriberId;
34use risingwave_pb::meta::PbFragmentWorkerSlotMapping;
35use risingwave_pb::meta::subscribe_response::Operation;
36use risingwave_pb::source::PbCdcTableSnapshotSplits;
37use risingwave_pb::stream_plan::PbUpstreamSinkInfo;
38use risingwave_pb::stream_plan::barrier_mutation::Mutation;
39use risingwave_pb::stream_plan::stream_node::NodeBody;
40use risingwave_pb::stream_service::BarrierCompleteResponse;
41use tracing::{info, warn};
42
43use crate::barrier::cdc_progress::{CdcProgress, CdcTableBackfillTracker};
44use crate::barrier::command::{
45    CreateStreamingJobCommandInfo, PostCollectCommand, ReplaceStreamJobPlan, ThrottleConfigMap,
46    extract_throttle_config,
47};
48use crate::barrier::edge_builder::{FragmentEdgeBuildResult, FragmentEdgeBuilder};
49use crate::barrier::progress::{CreateMviewProgressTracker, StagingCommitInfo};
50use crate::barrier::rpc::{ControlStreamManager, to_partial_graph_id};
51use crate::barrier::{
52    BackfillProgress, BarrierKind, CreateStreamingJobType, FragmentBackfillProgress, TracedEpoch,
53};
54use crate::controller::fragment::{InflightActorInfo, InflightFragmentInfo};
55use crate::controller::utils::rebuild_fragment_mapping;
56use crate::manager::NotificationManagerRef;
57use crate::model::{
58    ActorId, ActorNewNoShuffle, BackfillUpstreamType, FragmentId, StreamActor, StreamJobFragments,
59};
60use crate::stream::UpstreamSinkInfo;
61use crate::{MetaError, MetaResult};
62
63#[derive(Debug, Clone)]
64pub struct SharedActorInfo {
65    pub worker_id: WorkerId,
66    pub vnode_bitmap: Option<Bitmap>,
67    pub splits: Vec<SplitImpl>,
68}
69
70impl From<&InflightActorInfo> for SharedActorInfo {
71    fn from(value: &InflightActorInfo) -> Self {
72        Self {
73            worker_id: value.worker_id,
74            vnode_bitmap: value.vnode_bitmap.clone(),
75            splits: value.splits.clone(),
76        }
77    }
78}
79
80#[derive(Debug, Clone)]
81pub struct SharedFragmentInfo {
82    pub fragment_id: FragmentId,
83    pub job_id: JobId,
84    pub distribution_type: DistributionType,
85    pub actors: HashMap<ActorId, SharedActorInfo>,
86    pub vnode_count: usize,
87    pub fragment_type_mask: FragmentTypeMask,
88    pub state_table_ids: HashSet<TableId>,
89}
90
91impl From<(&InflightFragmentInfo, JobId)> for SharedFragmentInfo {
92    fn from(pair: (&InflightFragmentInfo, JobId)) -> Self {
93        let (info, job_id) = pair;
94
95        let InflightFragmentInfo {
96            fragment_id,
97            distribution_type,
98            fragment_type_mask,
99            actors,
100            vnode_count,
101            state_table_ids,
102            ..
103        } = info;
104
105        Self {
106            fragment_id: *fragment_id,
107            job_id,
108            distribution_type: *distribution_type,
109            fragment_type_mask: *fragment_type_mask,
110            actors: actors
111                .iter()
112                .map(|(actor_id, actor)| (*actor_id, actor.into()))
113                .collect(),
114            vnode_count: *vnode_count,
115            state_table_ids: state_table_ids.clone(),
116        }
117    }
118}
119
120#[derive(Default, Debug)]
121pub struct SharedActorInfosInner {
122    info: HashMap<DatabaseId, HashMap<FragmentId, SharedFragmentInfo>>,
123}
124
125impl SharedActorInfosInner {
126    pub fn get_fragment(&self, fragment_id: FragmentId) -> Option<&SharedFragmentInfo> {
127        self.info
128            .values()
129            .find_map(|database| database.get(&fragment_id))
130    }
131
132    pub fn iter_over_fragments(&self) -> impl Iterator<Item = (&FragmentId, &SharedFragmentInfo)> {
133        self.info.values().flatten()
134    }
135}
136
137#[derive(Clone, educe::Educe)]
138#[educe(Debug)]
139pub struct SharedActorInfos {
140    inner: Arc<parking_lot::RwLock<SharedActorInfosInner>>,
141    #[educe(Debug(ignore))]
142    notification_manager: NotificationManagerRef,
143}
144
145impl SharedActorInfos {
146    pub fn read_guard(&self) -> RwLockReadGuard<'_, RawRwLock, SharedActorInfosInner> {
147        self.inner.read()
148    }
149
150    pub fn list_assignments(&self) -> HashMap<ActorId, Vec<SplitImpl>> {
151        let core = self.inner.read();
152        core.iter_over_fragments()
153            .flat_map(|(_, fragment)| {
154                fragment
155                    .actors
156                    .iter()
157                    .map(|(actor_id, info)| (*actor_id, info.splits.clone()))
158            })
159            .collect()
160    }
161
162    /// Migrates splits from previous actors to the new actors for a rescheduled fragment.
163    ///
164    /// Very occasionally split removal may happen during scaling, in which case we need to
165    /// use the old splits for reallocation instead of the latest splits (which may be missing),
166    /// so that we can resolve the split removal in the next command.
167    pub fn migrate_splits_for_source_actors(
168        &self,
169        fragment_id: FragmentId,
170        prev_actor_ids: &[ActorId],
171        curr_actor_ids: &[ActorId],
172    ) -> MetaResult<HashMap<ActorId, Vec<SplitImpl>>> {
173        let guard = self.read_guard();
174
175        let prev_splits = prev_actor_ids
176            .iter()
177            .flat_map(|actor_id| {
178                // Note: File Source / Iceberg Source doesn't have splits assigned by meta.
179                guard
180                    .get_fragment(fragment_id)
181                    .and_then(|info| info.actors.get(actor_id))
182                    .map(|actor| actor.splits.clone())
183                    .unwrap_or_default()
184            })
185            .map(|split| (split.id(), split))
186            .collect();
187
188        let empty_actor_splits = curr_actor_ids
189            .iter()
190            .map(|actor_id| (*actor_id, vec![]))
191            .collect();
192
193        let diff = crate::stream::source_manager::reassign_splits(
194            fragment_id,
195            empty_actor_splits,
196            &prev_splits,
197            // pre-allocate splits is the first time getting splits, and it does not have scale-in scene
198            std::default::Default::default(),
199        )
200        .unwrap_or_default();
201
202        Ok(diff)
203    }
204}
205
206impl SharedActorInfos {
207    pub(crate) fn new(notification_manager: NotificationManagerRef) -> Self {
208        Self {
209            inner: Arc::new(Default::default()),
210            notification_manager,
211        }
212    }
213
214    pub(super) fn remove_database(&self, database_id: DatabaseId) {
215        if let Some(database) = self.inner.write().info.remove(&database_id) {
216            let mapping = database
217                .into_values()
218                .map(|fragment| rebuild_fragment_mapping(&fragment))
219                .collect_vec();
220            if !mapping.is_empty() {
221                self.notification_manager
222                    .notify_streaming_fragment_mapping(Operation::Delete, mapping);
223            }
224        }
225    }
226
227    pub(super) fn retain_databases(&self, database_ids: impl IntoIterator<Item = DatabaseId>) {
228        let database_ids: HashSet<_> = database_ids.into_iter().collect();
229
230        let mut mapping = Vec::new();
231        for fragment in self
232            .inner
233            .write()
234            .info
235            .extract_if(|database_id, _| !database_ids.contains(database_id))
236            .flat_map(|(_, fragments)| fragments.into_values())
237        {
238            mapping.push(rebuild_fragment_mapping(&fragment));
239        }
240        if !mapping.is_empty() {
241            self.notification_manager
242                .notify_streaming_fragment_mapping(Operation::Delete, mapping);
243        }
244    }
245
246    pub(super) fn recover_database(
247        &self,
248        database_id: DatabaseId,
249        fragments: impl Iterator<Item = (&InflightFragmentInfo, JobId)>,
250    ) {
251        let mut remaining_fragments: HashMap<_, _> = fragments
252            .map(|info @ (fragment, _)| (fragment.fragment_id, info))
253            .collect();
254        // delete the fragments that exist previously, but not included in the recovered fragments
255        let mut writer = self.start_writer(database_id);
256        let database = writer.write_guard.info.entry(database_id).or_default();
257        for (_, fragment) in database.extract_if(|fragment_id, fragment_infos| {
258            if let Some(info) = remaining_fragments.remove(fragment_id) {
259                let info = info.into();
260                writer
261                    .updated_fragment_mapping
262                    .get_or_insert_default()
263                    .push(rebuild_fragment_mapping(&info));
264                *fragment_infos = info;
265                false
266            } else {
267                true
268            }
269        }) {
270            writer
271                .deleted_fragment_mapping
272                .get_or_insert_default()
273                .push(rebuild_fragment_mapping(&fragment));
274        }
275        for (fragment_id, info) in remaining_fragments {
276            let info = info.into();
277            writer
278                .added_fragment_mapping
279                .get_or_insert_default()
280                .push(rebuild_fragment_mapping(&info));
281            database.insert(fragment_id, info);
282        }
283        writer.finish();
284    }
285
286    pub(super) fn upsert(
287        &self,
288        database_id: DatabaseId,
289        infos: impl IntoIterator<Item = (&InflightFragmentInfo, JobId)>,
290    ) {
291        let mut writer = self.start_writer(database_id);
292        writer.upsert(infos);
293        writer.finish();
294    }
295
296    pub(super) fn start_writer(&self, database_id: DatabaseId) -> SharedActorInfoWriter<'_> {
297        SharedActorInfoWriter {
298            database_id,
299            write_guard: self.inner.write(),
300            notification_manager: &self.notification_manager,
301            added_fragment_mapping: None,
302            updated_fragment_mapping: None,
303            deleted_fragment_mapping: None,
304        }
305    }
306}
307
308pub(super) struct SharedActorInfoWriter<'a> {
309    database_id: DatabaseId,
310    write_guard: parking_lot::RwLockWriteGuard<'a, SharedActorInfosInner>,
311    notification_manager: &'a NotificationManagerRef,
312    added_fragment_mapping: Option<Vec<PbFragmentWorkerSlotMapping>>,
313    updated_fragment_mapping: Option<Vec<PbFragmentWorkerSlotMapping>>,
314    deleted_fragment_mapping: Option<Vec<PbFragmentWorkerSlotMapping>>,
315}
316
317impl SharedActorInfoWriter<'_> {
318    pub(super) fn upsert(
319        &mut self,
320        infos: impl IntoIterator<Item = (&InflightFragmentInfo, JobId)>,
321    ) {
322        let database = self.write_guard.info.entry(self.database_id).or_default();
323        for info @ (fragment, _) in infos {
324            match database.entry(fragment.fragment_id) {
325                Entry::Occupied(mut entry) => {
326                    let info = info.into();
327                    self.updated_fragment_mapping
328                        .get_or_insert_default()
329                        .push(rebuild_fragment_mapping(&info));
330                    entry.insert(info);
331                }
332                Entry::Vacant(entry) => {
333                    let info = info.into();
334                    self.added_fragment_mapping
335                        .get_or_insert_default()
336                        .push(rebuild_fragment_mapping(&info));
337                    entry.insert(info);
338                }
339            }
340        }
341    }
342
343    pub(super) fn remove(&mut self, info: &InflightFragmentInfo) {
344        if let Some(database) = self.write_guard.info.get_mut(&self.database_id)
345            && let Some(fragment) = database.remove(&info.fragment_id)
346        {
347            self.deleted_fragment_mapping
348                .get_or_insert_default()
349                .push(rebuild_fragment_mapping(&fragment));
350        }
351    }
352
353    pub(super) fn finish(self) {
354        if let Some(mapping) = self.added_fragment_mapping {
355            self.notification_manager
356                .notify_streaming_fragment_mapping(Operation::Add, mapping);
357        }
358        if let Some(mapping) = self.updated_fragment_mapping {
359            self.notification_manager
360                .notify_streaming_fragment_mapping(Operation::Update, mapping);
361        }
362        if let Some(mapping) = self.deleted_fragment_mapping {
363            self.notification_manager
364                .notify_streaming_fragment_mapping(Operation::Delete, mapping);
365        }
366    }
367}
368
369#[derive(Debug, Clone)]
370pub(super) struct BarrierInfo {
371    pub prev_epoch: TracedEpoch,
372    pub curr_epoch: TracedEpoch,
373    pub kind: BarrierKind,
374}
375
376impl BarrierInfo {
377    pub(super) fn prev_epoch(&self) -> u64 {
378        self.prev_epoch.value().0
379    }
380
381    pub(super) fn curr_epoch(&self) -> u64 {
382        self.curr_epoch.value().0
383    }
384
385    pub(super) fn epoch(&self) -> EpochPair {
386        EpochPair {
387            curr: self.curr_epoch(),
388            prev: self.prev_epoch(),
389        }
390    }
391}
392
393#[derive(Clone, Debug)]
394pub enum SubscriberType {
395    Subscription(u64),
396    SnapshotBackfill,
397}
398
399#[derive(Debug)]
400pub(super) enum CreateStreamingJobStatus {
401    Init,
402    Creating { tracker: CreateMviewProgressTracker },
403    Created,
404}
405
406#[derive(Debug)]
407pub(super) struct InflightStreamingJobInfo {
408    pub job_id: JobId,
409    pub fragment_infos: HashMap<FragmentId, InflightFragmentInfo>,
410    pub subscribers: HashMap<SubscriberId, SubscriberType>,
411    pub status: CreateStreamingJobStatus,
412    pub cdc_table_backfill_tracker: Option<CdcTableBackfillTracker>,
413}
414
415impl InflightStreamingJobInfo {
416    pub fn fragment_infos(&self) -> impl Iterator<Item = &InflightFragmentInfo> + '_ {
417        self.fragment_infos.values()
418    }
419
420    pub fn snapshot_backfill_actor_ids(
421        fragment_infos: &HashMap<FragmentId, InflightFragmentInfo>,
422    ) -> impl Iterator<Item = ActorId> + '_ {
423        fragment_infos
424            .values()
425            .filter(|fragment| {
426                fragment
427                    .fragment_type_mask
428                    .contains(FragmentTypeFlag::SnapshotBackfillStreamScan)
429            })
430            .flat_map(|fragment| fragment.actors.keys().copied())
431    }
432
433    pub fn tracking_progress_actor_ids(
434        fragment_infos: &HashMap<FragmentId, InflightFragmentInfo>,
435    ) -> Vec<(ActorId, BackfillUpstreamType)> {
436        StreamJobFragments::tracking_progress_actor_ids_impl(fragment_infos.values().map(
437            |fragment| {
438                (
439                    fragment.fragment_type_mask,
440                    &fragment.nodes,
441                    fragment.actors.keys().copied(),
442                )
443            },
444        ))
445    }
446}
447
448impl<'a> IntoIterator for &'a InflightStreamingJobInfo {
449    type Item = &'a InflightFragmentInfo;
450
451    type IntoIter = impl Iterator<Item = &'a InflightFragmentInfo> + 'a;
452
453    fn into_iter(self) -> Self::IntoIter {
454        self.fragment_infos()
455    }
456}
457
458#[derive(Debug)]
459pub struct InflightDatabaseInfo {
460    pub(super) database_id: DatabaseId,
461    jobs: HashMap<JobId, InflightStreamingJobInfo>,
462    fragment_location: HashMap<FragmentId, JobId>,
463    pub(super) shared_actor_infos: SharedActorInfos,
464}
465
466impl InflightDatabaseInfo {
467    pub(super) fn job_ids(&self) -> impl Iterator<Item = JobId> + '_ {
468        self.jobs.keys().copied()
469    }
470
471    pub fn fragment_infos(&self) -> impl Iterator<Item = &InflightFragmentInfo> + '_ {
472        self.jobs.values().flat_map(|job| job.fragment_infos())
473    }
474
475    pub fn contains_job(&self, job_id: JobId) -> bool {
476        self.jobs.contains_key(&job_id)
477    }
478
479    /// Empty if the job is not in this database.
480    pub fn job_fragment_infos(
481        &self,
482        job_id: JobId,
483    ) -> impl Iterator<Item = &InflightFragmentInfo> + '_ {
484        self.jobs
485            .get(&job_id)
486            .into_iter()
487            .flat_map(|job| job.fragment_infos())
488    }
489
490    pub(super) fn job_id_by_fragment(&self, fragment_id: FragmentId) -> Option<JobId> {
491        self.fragment_location.get(&fragment_id).copied()
492    }
493
494    pub fn fragment(&self, fragment_id: FragmentId) -> &InflightFragmentInfo {
495        let job_id = self.fragment_location[&fragment_id];
496        self.jobs
497            .get(&job_id)
498            .expect("should exist")
499            .fragment_infos
500            .get(&fragment_id)
501            .expect("should exist")
502    }
503
504    pub(super) fn backfill_fragment_ids_for_job(
505        &self,
506        job_id: JobId,
507    ) -> MetaResult<HashSet<FragmentId>> {
508        let job = self
509            .jobs
510            .get(&job_id)
511            .ok_or_else(|| MetaError::invalid_parameter(format!("job {} not found", job_id)))?;
512        Ok(job
513            .fragment_infos
514            .iter()
515            .filter_map(|(fragment_id, fragment)| {
516                fragment
517                    .fragment_type_mask
518                    .contains_any([
519                        FragmentTypeFlag::StreamScan,
520                        FragmentTypeFlag::SourceScan,
521                        FragmentTypeFlag::LocalityProvider,
522                    ])
523                    .then_some(*fragment_id)
524            })
525            .collect())
526    }
527
528    pub(super) fn is_backfill_fragment(&self, fragment_id: FragmentId) -> MetaResult<bool> {
529        let job_id = self.fragment_location.get(&fragment_id).ok_or_else(|| {
530            MetaError::invalid_parameter(format!("fragment {} not found", fragment_id))
531        })?;
532        let fragment = self
533            .jobs
534            .get(job_id)
535            .expect("should exist")
536            .fragment_infos
537            .get(&fragment_id)
538            .expect("should exist");
539        Ok(fragment.fragment_type_mask.contains_any([
540            FragmentTypeFlag::StreamScan,
541            FragmentTypeFlag::SourceScan,
542            FragmentTypeFlag::LocalityProvider,
543        ]))
544    }
545
546    pub fn gen_backfill_progress(&self) -> impl Iterator<Item = (JobId, BackfillProgress)> + '_ {
547        self.jobs
548            .iter()
549            .filter_map(|(job_id, job)| match &job.status {
550                CreateStreamingJobStatus::Init => None,
551                CreateStreamingJobStatus::Creating { tracker } => {
552                    let progress = tracker.gen_backfill_progress();
553                    Some((
554                        *job_id,
555                        BackfillProgress {
556                            progress,
557                            backfill_type: PbBackfillType::NormalBackfill,
558                        },
559                    ))
560                }
561                CreateStreamingJobStatus::Created => None,
562            })
563    }
564
565    pub fn gen_cdc_progress(&self) -> impl Iterator<Item = (JobId, CdcProgress)> + '_ {
566        self.jobs.iter().filter_map(|(job_id, job)| {
567            job.cdc_table_backfill_tracker
568                .as_ref()
569                .map(|tracker| (*job_id, tracker.gen_cdc_progress()))
570        })
571    }
572
573    pub fn gen_fragment_backfill_progress(&self) -> Vec<FragmentBackfillProgress> {
574        let mut result = Vec::new();
575        for job in self.jobs.values() {
576            let CreateStreamingJobStatus::Creating { tracker } = &job.status else {
577                continue;
578            };
579            let fragment_progress = tracker.collect_fragment_progress(&job.fragment_infos, true);
580            result.extend(fragment_progress);
581        }
582        result
583    }
584
585    pub(super) fn may_assign_fragment_cdc_backfill_splits(
586        &mut self,
587        fragment_id: FragmentId,
588    ) -> MetaResult<Option<HashMap<ActorId, PbCdcTableSnapshotSplits>>> {
589        let job_id = self.fragment_location[&fragment_id];
590        let job = self.jobs.get_mut(&job_id).expect("should exist");
591        if let Some(tracker) = &mut job.cdc_table_backfill_tracker {
592            let cdc_scan_fragment_id = tracker.cdc_scan_fragment_id();
593            if cdc_scan_fragment_id != fragment_id {
594                return Ok(None);
595            }
596            let actors = job.fragment_infos[&cdc_scan_fragment_id]
597                .actors
598                .keys()
599                .copied()
600                .collect();
601            tracker.reassign_splits(actors).map(Some)
602        } else {
603            Ok(None)
604        }
605    }
606
607    pub(super) fn assign_cdc_backfill_splits(
608        &mut self,
609        job_id: JobId,
610    ) -> MetaResult<Option<HashMap<ActorId, PbCdcTableSnapshotSplits>>> {
611        let job = self.jobs.get_mut(&job_id).expect("should exist");
612        if let Some(tracker) = &mut job.cdc_table_backfill_tracker {
613            let cdc_scan_fragment_id = tracker.cdc_scan_fragment_id();
614            let actors = job.fragment_infos[&cdc_scan_fragment_id]
615                .actors
616                .keys()
617                .copied()
618                .collect();
619            tracker.reassign_splits(actors).map(Some)
620        } else {
621            Ok(None)
622        }
623    }
624
625    pub(super) fn apply_collected_command(
626        &mut self,
627        command: &PostCollectCommand,
628        resps: &HashMap<WorkerId, BarrierCompleteResponse>,
629        version_stats: &HummockVersionStats,
630    ) {
631        if let PostCollectCommand::CreateStreamingJob { info, job_type, .. } = command {
632            match job_type {
633                CreateStreamingJobType::Normal | CreateStreamingJobType::SinkIntoTable(_) => {
634                    let job_id = info.streaming_job.id();
635                    if let Some(job_info) = self.jobs.get_mut(&job_id) {
636                        let CreateStreamingJobStatus::Init = replace(
637                            &mut job_info.status,
638                            CreateStreamingJobStatus::Creating {
639                                tracker: CreateMviewProgressTracker::new(
640                                    info,
641                                    version_stats,
642                                    &job_info.fragment_infos,
643                                ),
644                            },
645                        ) else {
646                            unreachable!("should be init before collect the first barrier")
647                        };
648                    } else {
649                        info!(%job_id, "newly create job get cancelled before first barrier is collected")
650                    }
651                }
652                CreateStreamingJobType::Independent { .. } => {
653                    // The progress of SnapshotBackfill/BatchRefresh won't be tracked here
654                }
655            }
656        }
657        if let PostCollectCommand::Reschedule { reschedules, .. } = command {
658            // During reschedule we expect fragments to be rebuilt with new actors and no vnode bitmap update.
659            debug_assert!(
660                reschedules
661                    .values()
662                    .all(|reschedule| reschedule.vnode_bitmap_updates.is_empty()),
663                "Reschedule should not carry vnode bitmap updates when actors are rebuilt"
664            );
665
666            // Collect jobs that own the rescheduled fragments; de-duplicate via HashSet.
667            let related_job_ids = reschedules
668                .keys()
669                .filter_map(|fragment_id| self.fragment_location.get(fragment_id))
670                .cloned()
671                .collect::<HashSet<_>>();
672            for job_id in related_job_ids {
673                if let Some(job) = self.jobs.get_mut(&job_id)
674                    && let CreateStreamingJobStatus::Creating { tracker, .. } = &mut job.status
675                {
676                    tracker.refresh_after_reschedule(&job.fragment_infos, version_stats);
677                }
678            }
679        }
680        for progress in resps.values().flat_map(|resp| &resp.create_mview_progress) {
681            let Some(job_id) = self.fragment_location.get(&progress.fragment_id) else {
682                warn!(
683                    "update the progress of an non-existent creating streaming job: {progress:?}, which could be cancelled"
684                );
685                continue;
686            };
687            let tracker = match &mut self.jobs.get_mut(job_id).expect("should exist").status {
688                CreateStreamingJobStatus::Init => {
689                    continue;
690                }
691                CreateStreamingJobStatus::Creating { tracker, .. } => tracker,
692                CreateStreamingJobStatus::Created => {
693                    if !progress.done {
694                        warn!("update the progress of an created streaming job: {progress:?}");
695                    }
696                    continue;
697                }
698            };
699            tracker.apply_progress(progress, version_stats);
700        }
701        for progress in resps
702            .values()
703            .flat_map(|resp| &resp.cdc_table_backfill_progress)
704        {
705            let Some(job_id) = self.fragment_location.get(&progress.fragment_id) else {
706                warn!(
707                    "update the cdc progress of an non-existent creating streaming job: {progress:?}, which could be cancelled"
708                );
709                continue;
710            };
711            let Some(tracker) = &mut self
712                .jobs
713                .get_mut(job_id)
714                .expect("should exist")
715                .cdc_table_backfill_tracker
716            else {
717                warn!("update the cdc progress of an created streaming job: {progress:?}");
718                continue;
719            };
720            tracker.update_split_progress(progress);
721        }
722        // Handle CDC source offset updated events
723        for cdc_offset_updated in resps
724            .values()
725            .flat_map(|resp| &resp.cdc_source_offset_updated)
726        {
727            let source_id = cdc_offset_updated.source_id;
728            let job_id = source_id.as_share_source_job_id();
729            if let Some(job) = self.jobs.get_mut(&job_id) {
730                if let CreateStreamingJobStatus::Creating { tracker, .. } = &mut job.status {
731                    tracker.mark_cdc_source_finished();
732                }
733            } else {
734                warn!(
735                    "update cdc source offset for non-existent creating streaming job: source_id={}, job_id={}",
736                    cdc_offset_updated.source_id, job_id
737                );
738            }
739        }
740    }
741
742    fn iter_creating_job_tracker(&self) -> impl Iterator<Item = &CreateMviewProgressTracker> {
743        self.jobs.values().filter_map(|job| match &job.status {
744            CreateStreamingJobStatus::Init => None,
745            CreateStreamingJobStatus::Creating { tracker, .. } => Some(tracker),
746            CreateStreamingJobStatus::Created => None,
747        })
748    }
749
750    fn iter_mut_creating_job_tracker(
751        &mut self,
752    ) -> impl Iterator<Item = &mut CreateMviewProgressTracker> {
753        self.jobs
754            .values_mut()
755            .filter_map(|job| match &mut job.status {
756                CreateStreamingJobStatus::Init => None,
757                CreateStreamingJobStatus::Creating { tracker, .. } => Some(tracker),
758                CreateStreamingJobStatus::Created => None,
759            })
760    }
761
762    pub(super) fn has_pending_finished_jobs(&self) -> bool {
763        self.iter_creating_job_tracker()
764            .any(|tracker| tracker.is_finished())
765    }
766
767    pub(super) fn take_pending_backfill_nodes(&mut self) -> Vec<FragmentId> {
768        self.iter_mut_creating_job_tracker()
769            .flat_map(|tracker| tracker.take_pending_backfill_nodes())
770            .collect()
771    }
772
773    pub(super) fn take_staging_commit_info(&mut self) -> StagingCommitInfo {
774        let mut finished_jobs = vec![];
775        let mut table_ids_to_truncate = vec![];
776        let mut finished_cdc_table_backfill = vec![];
777        for (job_id, job) in &mut self.jobs {
778            if let CreateStreamingJobStatus::Creating { tracker, .. } = &mut job.status {
779                let (is_finished, truncate_table_ids) = tracker.collect_staging_commit_info();
780                table_ids_to_truncate.extend(truncate_table_ids);
781                if is_finished {
782                    let CreateStreamingJobStatus::Creating { tracker, .. } =
783                        replace(&mut job.status, CreateStreamingJobStatus::Created)
784                    else {
785                        unreachable!()
786                    };
787                    finished_jobs.push(tracker.into_tracking_job());
788                }
789            }
790            if let Some(tracker) = &mut job.cdc_table_backfill_tracker
791                && tracker.take_pre_completed()
792            {
793                finished_cdc_table_backfill.push(*job_id);
794            }
795        }
796        StagingCommitInfo {
797            finished_jobs,
798            table_ids_to_truncate,
799            finished_cdc_table_backfill,
800        }
801    }
802
803    pub fn fragment_subscribers(
804        &self,
805        fragment_id: FragmentId,
806    ) -> impl Iterator<Item = SubscriberId> + '_ {
807        let job_id = self.fragment_location[&fragment_id];
808        self.jobs[&job_id].subscribers.keys().copied()
809    }
810
811    pub fn job_subscribers(&self, job_id: JobId) -> impl Iterator<Item = SubscriberId> + '_ {
812        self.jobs[&job_id].subscribers.keys().copied()
813    }
814
815    pub fn subscribed_tables(&self) -> impl Iterator<Item = TableId> + '_ {
816        self.jobs.iter().filter_map(|(job_id, info)| {
817            info.subscribers
818                .values()
819                .any(|subscriber| matches!(subscriber, SubscriberType::Subscription(_)))
820                .then_some(job_id.as_mv_table_id())
821        })
822    }
823
824    pub fn register_subscriber(
825        &mut self,
826        job_id: JobId,
827        subscriber_id: SubscriberId,
828        subscriber: SubscriberType,
829    ) {
830        self.jobs
831            .get_mut(&job_id)
832            .expect("should exist")
833            .subscribers
834            .try_insert(subscriber_id, subscriber)
835            .expect("non duplicate");
836    }
837
838    pub fn unregister_subscriber(
839        &mut self,
840        job_id: JobId,
841        subscriber_id: SubscriberId,
842    ) -> Option<SubscriberType> {
843        self.jobs
844            .get_mut(&job_id)
845            .expect("should exist")
846            .subscribers
847            .remove(&subscriber_id)
848    }
849
850    pub fn update_subscription_retention(
851        &mut self,
852        job_id: JobId,
853        subscriber_id: SubscriberId,
854        retention_second: u64,
855    ) {
856        let job = self.jobs.get_mut(&job_id).expect("should exist");
857        match job.subscribers.get_mut(&subscriber_id) {
858            Some(SubscriberType::Subscription(current_retention)) => {
859                *current_retention = retention_second;
860            }
861            Some(SubscriberType::SnapshotBackfill) => {
862                warn!(
863                    %job_id,
864                    %subscriber_id,
865                    "cannot update retention for snapshot backfill subscriber"
866                );
867            }
868            None => {
869                warn!(%job_id, %subscriber_id, "subscription subscriber not found");
870            }
871        }
872    }
873
874    fn fragment_mut(&mut self, fragment_id: FragmentId) -> (&mut InflightFragmentInfo, JobId) {
875        let job_id = self.fragment_location[&fragment_id];
876        let fragment = self
877            .jobs
878            .get_mut(&job_id)
879            .expect("should exist")
880            .fragment_infos
881            .get_mut(&fragment_id)
882            .expect("should exist");
883        (fragment, job_id)
884    }
885
886    fn empty_inner(database_id: DatabaseId, shared_actor_infos: SharedActorInfos) -> Self {
887        Self {
888            database_id,
889            jobs: Default::default(),
890            fragment_location: Default::default(),
891            shared_actor_infos,
892        }
893    }
894
895    pub fn empty(database_id: DatabaseId, shared_actor_infos: SharedActorInfos) -> Self {
896        // remove the database because it's empty.
897        shared_actor_infos.remove_database(database_id);
898        Self::empty_inner(database_id, shared_actor_infos)
899    }
900
901    pub fn recover(
902        database_id: DatabaseId,
903        jobs: impl Iterator<Item = InflightStreamingJobInfo>,
904        shared_actor_infos: SharedActorInfos,
905    ) -> Self {
906        let mut info = Self::empty_inner(database_id, shared_actor_infos);
907        for job in jobs {
908            info.add_existing(job);
909        }
910        info
911    }
912
913    pub fn is_empty(&self) -> bool {
914        self.jobs.is_empty()
915    }
916
917    pub fn add_existing(&mut self, job: InflightStreamingJobInfo) {
918        let InflightStreamingJobInfo {
919            job_id,
920            fragment_infos,
921            subscribers,
922            status,
923            cdc_table_backfill_tracker,
924        } = job;
925        self.jobs
926            .try_insert(
927                job_id,
928                InflightStreamingJobInfo {
929                    job_id,
930                    subscribers,
931                    fragment_infos: Default::default(), // fill in later in pre_apply_new_fragments
932                    status,
933                    cdc_table_backfill_tracker,
934                },
935            )
936            .expect("non-duplicate");
937        self.pre_apply_new_fragments(
938            fragment_infos
939                .into_iter()
940                .map(|(fragment_id, info)| (fragment_id, job_id, info)),
941        );
942    }
943
944    /// Register a new streaming job entry (with empty `fragment_infos`).
945    pub(crate) fn pre_apply_new_job(
946        &mut self,
947        job_id: JobId,
948        cdc_table_backfill_tracker: Option<CdcTableBackfillTracker>,
949    ) {
950        {
951            self.jobs
952                .try_insert(
953                    job_id,
954                    InflightStreamingJobInfo {
955                        job_id,
956                        fragment_infos: Default::default(),
957                        subscribers: Default::default(), // no subscriber for newly create job
958                        status: CreateStreamingJobStatus::Init,
959                        cdc_table_backfill_tracker,
960                    },
961                )
962                .expect("non-duplicate");
963        }
964    }
965
966    /// Add new fragment infos and update shared actor infos.
967    pub(crate) fn pre_apply_new_fragments(
968        &mut self,
969        fragments: impl IntoIterator<Item = (FragmentId, JobId, InflightFragmentInfo)>,
970    ) {
971        {
972            let shared_infos = self.shared_actor_infos.clone();
973            let mut shared_actor_writer = shared_infos.start_writer(self.database_id);
974            for (fragment_id, job_id, info) in fragments {
975                {
976                    {
977                        let fragment_infos = self.jobs.get_mut(&job_id).expect("should exist");
978                        shared_actor_writer.upsert([(&info, job_id)]);
979                        fragment_infos
980                            .fragment_infos
981                            .try_insert(fragment_id, info)
982                            .expect("non duplicate");
983                        self.fragment_location
984                            .try_insert(fragment_id, job_id)
985                            .expect("non duplicate");
986                    }
987                }
988            }
989            shared_actor_writer.finish();
990        }
991    }
992
993    /// Pre-apply reschedule: update actors, vnode bitmaps, and splits.
994    /// The actual removal of old actors happens in `post_apply_reschedules`.
995    pub(crate) fn pre_apply_reschedule(
996        &mut self,
997        fragment_id: FragmentId,
998        new_actors: HashMap<ActorId, InflightActorInfo>,
999        actor_update_vnode_bitmap: HashMap<ActorId, Bitmap>,
1000        actor_splits: HashMap<ActorId, Vec<SplitImpl>>,
1001    ) {
1002        {
1003            {
1004                {
1005                    {
1006                        let (info, _) = self.fragment_mut(fragment_id);
1007                        let actors = &mut info.actors;
1008                        for (actor_id, new_vnodes) in actor_update_vnode_bitmap {
1009                            actors
1010                                .get_mut(&actor_id)
1011                                .expect("should exist")
1012                                .vnode_bitmap = Some(new_vnodes);
1013                        }
1014                        for (actor_id, actor) in new_actors {
1015                            actors
1016                                .try_insert(actor_id as _, actor)
1017                                .expect("non-duplicate");
1018                        }
1019                        for (actor_id, splits) in actor_splits {
1020                            actors.get_mut(&actor_id).expect("should exist").splits = splits;
1021                        }
1022                        // info will be upserted into shared_actor_infos in post_apply stage
1023                    }
1024                }
1025            }
1026        }
1027    }
1028
1029    /// Replace upstream fragment IDs in merge nodes of a fragment's stream graph.
1030    pub(crate) fn pre_apply_replace_node_upstream(
1031        &mut self,
1032        fragment_id: FragmentId,
1033        replace_map: &HashMap<FragmentId, FragmentId>,
1034    ) {
1035        {
1036            {
1037                {
1038                    {
1039                        let mut remaining_fragment_ids: HashSet<_> =
1040                            replace_map.keys().cloned().collect();
1041                        let (info, _) = self.fragment_mut(fragment_id);
1042                        visit_stream_node_mut(&mut info.nodes, |node| {
1043                            if let NodeBody::Merge(m) = node
1044                                && let Some(new_upstream_fragment_id) =
1045                                    replace_map.get(&m.upstream_fragment_id)
1046                            {
1047                                if !remaining_fragment_ids.remove(&m.upstream_fragment_id) {
1048                                    if cfg!(debug_assertions) {
1049                                        panic!(
1050                                            "duplicate upstream fragment: {:?} {:?}",
1051                                            m, replace_map
1052                                        );
1053                                    } else {
1054                                        warn!(?m, ?replace_map, "duplicate upstream fragment");
1055                                    }
1056                                }
1057                                m.upstream_fragment_id = *new_upstream_fragment_id;
1058                            }
1059                        });
1060                        if cfg!(debug_assertions) {
1061                            assert!(
1062                                remaining_fragment_ids.is_empty(),
1063                                "non-existing fragment to replace: {:?} {:?} {:?}",
1064                                remaining_fragment_ids,
1065                                info.nodes,
1066                                replace_map
1067                            );
1068                        } else {
1069                            warn!(?remaining_fragment_ids, node = ?info.nodes, ?replace_map, "non-existing fragment to replace");
1070                        }
1071                    }
1072                }
1073            }
1074        }
1075    }
1076
1077    /// Add a new upstream sink node to a fragment's `UpstreamSinkUnion`.
1078    pub(crate) fn pre_apply_add_node_upstream(
1079        &mut self,
1080        fragment_id: FragmentId,
1081        new_upstream_info: &PbUpstreamSinkInfo,
1082    ) {
1083        {
1084            {
1085                {
1086                    {
1087                        let (info, _) = self.fragment_mut(fragment_id);
1088                        let mut injected = false;
1089                        visit_stream_node_mut(&mut info.nodes, |node| {
1090                            if let NodeBody::UpstreamSinkUnion(u) = node {
1091                                if cfg!(debug_assertions) {
1092                                    let current_upstream_fragment_ids = u
1093                                        .init_upstreams
1094                                        .iter()
1095                                        .map(|upstream| upstream.upstream_fragment_id)
1096                                        .collect::<HashSet<_>>();
1097                                    if current_upstream_fragment_ids
1098                                        .contains(&new_upstream_info.upstream_fragment_id)
1099                                    {
1100                                        panic!(
1101                                            "duplicate upstream fragment: {:?} {:?}",
1102                                            u, new_upstream_info
1103                                        );
1104                                    }
1105                                }
1106                                u.init_upstreams.push(new_upstream_info.clone());
1107                                injected = true;
1108                            }
1109                        });
1110                        assert!(injected, "should inject upstream into UpstreamSinkUnion");
1111                    }
1112                }
1113            }
1114        }
1115    }
1116
1117    /// Remove upstream sink nodes from a fragment's `UpstreamSinkUnion`.
1118    pub(crate) fn pre_apply_drop_node_upstream(
1119        &mut self,
1120        fragment_id: FragmentId,
1121        drop_upstream_fragment_ids: &[FragmentId],
1122    ) {
1123        if !self.fragment_location.contains_key(&fragment_id) {
1124            warn!(
1125                target_fragment_id = %fragment_id,
1126                drop_upstream_fragment_ids = ?drop_upstream_fragment_ids,
1127                "skip dropping upstream sink fragments for non-existing target fragment"
1128            );
1129            return;
1130        }
1131        {
1132            {
1133                {
1134                    {
1135                        let (info, _) = self.fragment_mut(fragment_id);
1136                        let mut removed = false;
1137                        visit_stream_node_mut(&mut info.nodes, |node| {
1138                            if let NodeBody::UpstreamSinkUnion(u) = node {
1139                                if cfg!(debug_assertions) {
1140                                    let current_upstream_fragment_ids = u
1141                                        .init_upstreams
1142                                        .iter()
1143                                        .map(|upstream| upstream.upstream_fragment_id)
1144                                        .collect::<HashSet<FragmentId>>();
1145                                    for drop_fragment_id in drop_upstream_fragment_ids {
1146                                        if !current_upstream_fragment_ids.contains(drop_fragment_id)
1147                                        {
1148                                            panic!(
1149                                                "non-existing upstream fragment to drop: {:?} {:?} {:?}",
1150                                                u, drop_upstream_fragment_ids, drop_fragment_id
1151                                            );
1152                                        }
1153                                    }
1154                                }
1155                                u.init_upstreams.retain(|upstream| {
1156                                    !drop_upstream_fragment_ids
1157                                        .contains(&upstream.upstream_fragment_id)
1158                                });
1159                                removed = true;
1160                            }
1161                        });
1162                        assert!(removed, "should remove upstream from UpstreamSinkUnion");
1163                    }
1164                }
1165            }
1166        }
1167    }
1168
1169    /// Sync inflight `nodes` so a later reschedule won't materialize new actors from stale data.
1170    pub(crate) fn pre_apply_throttle(
1171        &mut self,
1172        config: &mut ThrottleConfigMap,
1173    ) -> Option<Mutation> {
1174        extract_throttle_config(config, |fragment_id, stream_node| {
1175            if !self.fragment_location.contains_key(&fragment_id) {
1176                return false;
1177            }
1178            self.fragment_mut(fragment_id).0.nodes = stream_node.clone();
1179            true
1180        })
1181    }
1182
1183    /// Update split assignments for actors in fragments.
1184    pub(crate) fn pre_apply_split_assignments(
1185        &mut self,
1186        assignments: impl IntoIterator<Item = (FragmentId, HashMap<ActorId, Vec<SplitImpl>>)>,
1187    ) {
1188        {
1189            let shared_infos = self.shared_actor_infos.clone();
1190            let mut shared_actor_writer = shared_infos.start_writer(self.database_id);
1191            {
1192                {
1193                    for (fragment_id, actor_splits) in assignments {
1194                        let (info, job_id) = self.fragment_mut(fragment_id);
1195                        let actors = &mut info.actors;
1196                        for (actor_id, splits) in actor_splits {
1197                            actors.get_mut(&actor_id).expect("should exist").splits = splits;
1198                        }
1199                        shared_actor_writer.upsert([(&*info, job_id)]);
1200                    }
1201                }
1202            }
1203            shared_actor_writer.finish();
1204        }
1205    }
1206
1207    pub(super) fn build_edge(
1208        &self,
1209        info: Option<(&CreateStreamingJobCommandInfo, bool)>,
1210        replace_job: Option<&ReplaceStreamJobPlan>,
1211        new_upstream_sink: Option<&UpstreamSinkInfo>,
1212        control_stream_manager: &ControlStreamManager,
1213        stream_actors: &HashMap<FragmentId, Vec<StreamActor>>,
1214        actor_location: &HashMap<ActorId, WorkerId>,
1215    ) -> MetaResult<(FragmentEdgeBuildResult, ActorNewNoShuffle)> {
1216        // `existing_fragment_ids` consists of
1217        //  - keys of `info.upstream_fragment_downstreams`, which are the `fragment_id` the upstream fragment of the newly created job
1218        //  - keys of `replace_job.upstream_fragment_downstreams`, which are the `fragment_id` of upstream fragment of replace_job,
1219        // if the upstream fragment previously exists
1220        //  - keys of `replace_upstream`, which are the `fragment_id` of downstream fragments that will update their upstream fragments,
1221        // if creating a new sink-into-table
1222        //  - should contain the `fragment_id` of the downstream table.
1223        let existing_fragment_ids = info
1224            .into_iter()
1225            .flat_map(|(info, _)| info.upstream_fragment_downstreams.keys())
1226            .chain(replace_job.into_iter().flat_map(|replace_job| {
1227                replace_job
1228                    .upstream_fragment_downstreams
1229                    .keys()
1230                    .filter(|fragment_id| {
1231                        info.map(|(info, _)| {
1232                            !info
1233                                .stream_job_fragments
1234                                .fragments
1235                                .contains_key(*fragment_id)
1236                        })
1237                        .unwrap_or(true)
1238                    })
1239                    .chain(replace_job.replace_upstream.keys())
1240            }))
1241            .chain(
1242                new_upstream_sink
1243                    .into_iter()
1244                    .map(|ctx| &ctx.new_sink_downstream.downstream_fragment_id),
1245            )
1246            .cloned();
1247        // Collect new fragments with their partial graph IDs
1248        let new_fragments = info
1249            .into_iter()
1250            .flat_map(|(info, is_snapshot_backfill)| {
1251                let partial_graph_id = to_partial_graph_id(
1252                    self.database_id,
1253                    is_snapshot_backfill.then_some(info.streaming_job.id()),
1254                );
1255                info.stream_job_fragments
1256                    .fragments
1257                    .values()
1258                    .map(move |fragment| (partial_graph_id, fragment))
1259            })
1260            .chain(replace_job.into_iter().flat_map(|replace_job| {
1261                replace_job
1262                    .new_fragments
1263                    .fragments
1264                    .values()
1265                    .chain(
1266                        replace_job
1267                            .auto_refresh_schema_sinks
1268                            .as_ref()
1269                            .into_iter()
1270                            .flat_map(move |sinks| sinks.iter().map(|sink| &sink.new_fragment)),
1271                    )
1272                    .map(|fragment| {
1273                        (
1274                            // we assume that replace job only happens in database partial graph
1275                            to_partial_graph_id(self.database_id, None),
1276                            fragment,
1277                        )
1278                    })
1279            }));
1280
1281        let database_partial_graph_id = to_partial_graph_id(self.database_id, None);
1282        let mut builder = FragmentEdgeBuilder::new()
1283            .add_existing_fragments(
1284                existing_fragment_ids.map(|fragment_id| self.fragment(fragment_id)),
1285                database_partial_graph_id,
1286                control_stream_manager,
1287            )
1288            .add_new_logical_fragments(
1289                new_fragments,
1290                stream_actors,
1291                actor_location,
1292                control_stream_manager,
1293            )
1294            .finish_fragments();
1295        if let Some((info, _)) = info {
1296            builder = builder
1297                .add_relations(&info.upstream_fragment_downstreams)?
1298                .add_relations(&info.stream_job_fragments.downstreams)?;
1299        }
1300        if let Some(replace_job) = replace_job {
1301            builder = builder
1302                .add_relations(&replace_job.upstream_fragment_downstreams)?
1303                .add_relations(&replace_job.new_fragments.downstreams)?;
1304        }
1305        if let Some(new_upstream_sink) = new_upstream_sink {
1306            let sink_fragment_id = new_upstream_sink.sink_fragment_id;
1307            let new_sink_downstream = &new_upstream_sink.new_sink_downstream;
1308            builder = builder.add_edge(sink_fragment_id, new_sink_downstream)?;
1309        }
1310        if let Some(replace_job) = replace_job {
1311            for (fragment_id, fragment_replacement) in &replace_job.replace_upstream {
1312                for (original_upstream_fragment_id, new_upstream_fragment_id) in
1313                    fragment_replacement
1314                {
1315                    builder = builder.replace_upstream(
1316                        *fragment_id,
1317                        *original_upstream_fragment_id,
1318                        *new_upstream_fragment_id,
1319                    );
1320                }
1321            }
1322        }
1323        Ok(builder.build())
1324    }
1325
1326    /// Post-apply reschedule: remove actors that were marked for removal.
1327    pub(crate) fn post_apply_reschedules(
1328        &mut self,
1329        reschedules: impl IntoIterator<Item = (FragmentId, HashSet<ActorId>)>,
1330    ) {
1331        let inner = self.shared_actor_infos.clone();
1332        let mut shared_actor_writer = inner.start_writer(self.database_id);
1333        {
1334            {
1335                {
1336                    for (fragment_id, to_remove) in reschedules {
1337                        let job_id = self.fragment_location[&fragment_id];
1338                        let info = self
1339                            .jobs
1340                            .get_mut(&job_id)
1341                            .expect("should exist")
1342                            .fragment_infos
1343                            .get_mut(&fragment_id)
1344                            .expect("should exist");
1345                        for actor_id in to_remove {
1346                            assert!(info.actors.remove(&actor_id).is_some());
1347                        }
1348                        shared_actor_writer.upsert([(&*info, job_id)]);
1349                    }
1350                }
1351            }
1352        }
1353        shared_actor_writer.finish();
1354    }
1355
1356    /// Post-apply fragment removal: remove fragments and their jobs if empty.
1357    pub(crate) fn post_apply_remove_fragments(
1358        &mut self,
1359        fragment_ids: impl IntoIterator<Item = FragmentId>,
1360    ) {
1361        let inner = self.shared_actor_infos.clone();
1362        let mut shared_actor_writer = inner.start_writer(self.database_id);
1363        {
1364            {
1365                {
1366                    for fragment_id in fragment_ids {
1367                        let job_id = self
1368                            .fragment_location
1369                            .remove(&fragment_id)
1370                            .expect("should exist");
1371                        let job = self.jobs.get_mut(&job_id).expect("should exist");
1372                        let fragment = job
1373                            .fragment_infos
1374                            .remove(&fragment_id)
1375                            .expect("should exist");
1376                        shared_actor_writer.remove(&fragment);
1377                        if job.fragment_infos.is_empty() {
1378                            self.jobs.remove(&job_id).expect("should exist");
1379                        }
1380                    }
1381                }
1382            }
1383        }
1384        shared_actor_writer.finish();
1385    }
1386
1387    pub(crate) fn post_apply_remove_job(
1388        &mut self,
1389        job_id: JobId,
1390    ) -> Option<InflightStreamingJobInfo> {
1391        let job = self.jobs.remove(&job_id)?;
1392        let inner = self.shared_actor_infos.clone();
1393        let mut shared_actor_writer = inner.start_writer(self.database_id);
1394        for (fragment_id, fragment) in &job.fragment_infos {
1395            self.fragment_location
1396                .remove(fragment_id)
1397                .expect("should exist");
1398            shared_actor_writer.remove(fragment);
1399        }
1400        shared_actor_writer.finish();
1401        Some(job)
1402    }
1403}
1404
1405impl InflightFragmentInfo {
1406    /// Returns actor list to collect in the target worker node.
1407    pub(crate) fn actor_ids_to_collect(
1408        infos: impl IntoIterator<Item = &Self>,
1409    ) -> HashMap<WorkerId, HashSet<ActorId>> {
1410        let mut ret: HashMap<_, HashSet<_>> = HashMap::new();
1411        for (actor_id, actor) in infos.into_iter().flat_map(|info| info.actors.iter()) {
1412            assert!(
1413                ret.entry(actor.worker_id)
1414                    .or_default()
1415                    .insert(*actor_id as _)
1416            )
1417        }
1418        ret
1419    }
1420
1421    pub fn existing_table_ids<'a>(
1422        infos: impl IntoIterator<Item = &'a Self> + 'a,
1423    ) -> impl Iterator<Item = TableId> + 'a {
1424        infos
1425            .into_iter()
1426            .flat_map(|info| info.state_table_ids.iter().cloned())
1427    }
1428
1429    pub fn workers<'a>(
1430        infos: impl IntoIterator<Item = &'a Self> + 'a,
1431    ) -> impl Iterator<Item = WorkerId> + 'a {
1432        infos
1433            .into_iter()
1434            .flat_map(|fragment| fragment.actors.values().map(|actor| actor.worker_id))
1435    }
1436
1437    pub fn contains_worker<'a>(
1438        infos: impl IntoIterator<Item = &'a Self> + 'a,
1439        worker_id: WorkerId,
1440    ) -> bool {
1441        Self::workers(infos).any(|existing_worker_id| existing_worker_id == worker_id)
1442    }
1443}
1444
1445impl InflightDatabaseInfo {
1446    pub fn contains_worker(&self, worker_id: WorkerId) -> bool {
1447        InflightFragmentInfo::contains_worker(self.fragment_infos(), worker_id)
1448    }
1449
1450    pub fn existing_table_ids(&self) -> impl Iterator<Item = TableId> + '_ {
1451        InflightFragmentInfo::existing_table_ids(self.fragment_infos())
1452    }
1453}