Skip to main content

risingwave_meta/barrier/context/
context_impl.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::{HashMap, HashSet};
16use std::sync::Arc;
17
18use anyhow::Context;
19use futures::StreamExt;
20use futures::stream::FuturesUnordered;
21use risingwave_common::catalog::{DatabaseId, FragmentTypeFlag, TableId};
22use risingwave_common::id::{JobId, PartialGraphId, SinkId};
23use risingwave_hummock_sdk::change_log::TableChangeLogs;
24use risingwave_meta_model::ActorId;
25use risingwave_meta_model::streaming_job::BackfillOrders;
26use risingwave_pb::common::WorkerNode;
27use risingwave_pb::hummock::HummockVersionStats;
28use risingwave_pb::id::SourceId;
29use risingwave_pb::meta::PbTableRefillRuntimeConfig;
30use risingwave_pb::meta::subscribe_response::{Info, Operation};
31use risingwave_pb::stream_service::barrier_complete_response::{
32    PbIcebergPkIndexSinkMetadata, PbListFinishedSource, PbLoadFinishedSource,
33};
34use risingwave_pb::stream_service::streaming_control_stream_request::PbInitRequest;
35use risingwave_rpc_client::StreamingControlHandle;
36use thiserror_ext::AsReport;
37
38use crate::barrier::cdc_progress::CdcTableBackfillTracker;
39use crate::barrier::checkpoint::independent_job::BatchRefreshJobTriggerContext;
40use crate::barrier::command::{
41    PostCollectCommand, ResumeBackfillTarget, SinceTimestampResolvedEpoch,
42};
43use crate::barrier::context::{GlobalBarrierWorkerContext, GlobalBarrierWorkerContextImpl};
44use crate::barrier::progress::TrackingJob;
45use crate::barrier::schedule::MarkReadyOptions;
46use crate::barrier::{
47    BarrierManagerStatus, BarrierWorkerRuntimeInfoSnapshot, BatchRefreshInfo, Command,
48    CreateStreamingJobCommandInfo, CreateStreamingJobType, DatabaseRuntimeInfoSnapshot,
49    RecoveryReason, ReplaceStreamJobPlan, Scheduled,
50};
51use crate::hummock::CommitEpochInfo;
52use crate::manager::LocalNotification;
53use crate::model::FragmentDownstreamRelation;
54use crate::serving::{fetch_serving_infos, sync_serving_table_vnode_mappings_to_hummock};
55use crate::stream::{SourceChange, cleanup_dropped_streaming_jobs};
56use crate::{MetaError, MetaResult};
57
58fn resolve_since_timestamp_log_store_epoch(
59    table_id: TableId,
60    since_epoch: u64,
61    upstream_committed_epoch: u64,
62    table_change_log: &TableChangeLogs,
63) -> MetaResult<SinceTimestampResolvedEpoch> {
64    let change_log = table_change_log.get(&table_id).ok_or_else(|| {
65        anyhow::anyhow!(
66            "no table changelog found for upstream table {} when resolving since_timestamp",
67            table_id
68        )
69    })?;
70    let Some(first_log) = change_log.first() else {
71        return Err(anyhow::anyhow!(
72            "empty table changelog found for upstream table {} when resolving since_timestamp",
73            table_id
74        )
75        .into());
76    };
77    let first_checkpoint_epoch = first_log.checkpoint_epoch;
78    if since_epoch < first_checkpoint_epoch {
79        return Err(anyhow::anyhow!(
80            "since_timestamp is earlier than the retained changelog of upstream table {}: requested epoch {}, first retained checkpoint epoch {}",
81            table_id,
82            since_epoch,
83            first_checkpoint_epoch,
84        )
85        .into());
86    }
87    if since_epoch >= upstream_committed_epoch {
88        return Err(anyhow::anyhow!(
89            "since_timestamp is not before the committed epoch of upstream table {}: requested epoch {}, committed epoch {}",
90            table_id,
91            since_epoch,
92            upstream_committed_epoch,
93        )
94        .into());
95    }
96    let latest_log = change_log.last().expect("checked non-empty");
97    if upstream_committed_epoch != latest_log.checkpoint_epoch {
98        return Err(anyhow::anyhow!(
99            "upstream committed epoch {} does not match latest changelog epoch {} for upstream table {}",
100            upstream_committed_epoch,
101            latest_log.checkpoint_epoch,
102            table_id,
103        )
104        .into());
105    }
106
107    // `binary_search_by_checkpoint_epoch` searches only by the checkpoint epoch
108    // of each changelog entry. An entry may still cover earlier non-checkpoint
109    // epochs, e.g. `{ non_checkpoint_epochs: [30, 35], checkpoint_epoch: 40 }`
110    // covers `(20, 40]` if the previous checkpoint is 20. For `since_epoch =
111    // 35`, the search returns `Err(index_of_40)`, and 40 is the snapshot epoch.
112    // For `since_epoch = 40`, it returns `Ok(index_of_40)`. In both cases, the
113    // resolved snapshot epoch is the least checkpoint epoch >= `since_epoch`.
114    let snapshot_epoch_index = match change_log.binary_search_by_checkpoint_epoch(since_epoch) {
115        Ok(index) => index,
116        Err(index) => index,
117    };
118    let snapshot_epoch = change_log
119        .get(snapshot_epoch_index)
120        .ok_or_else(|| {
121            anyhow::anyhow!(
122                "since_timestamp is later than the latest changelog of upstream table {}: requested epoch {}, latest changelog epoch {}",
123                table_id,
124                since_epoch,
125                upstream_committed_epoch,
126            )
127        })?
128        .checkpoint_epoch;
129    // A request inside the latest changelog entry may be smaller than
130    // `upstream_committed_epoch`, but still resolve to the latest checkpoint.
131    // That end-of-log case is not a usable historical snapshot epoch.
132    if snapshot_epoch >= upstream_committed_epoch {
133        return Err(anyhow::anyhow!(
134            "since_timestamp is too new for upstream table {}: requested epoch {}, resolved snapshot epoch {}, latest changelog epoch {}",
135            table_id,
136            since_epoch,
137            snapshot_epoch,
138            upstream_committed_epoch,
139        )
140        .into());
141    }
142    let epochs = change_log
143        .range((snapshot_epoch_index + 1)..)
144        .map(|epoch_log| {
145            (
146                epoch_log.non_checkpoint_epochs.clone(),
147                epoch_log.checkpoint_epoch,
148            )
149        })
150        .collect::<Vec<_>>();
151
152    Ok((snapshot_epoch, epochs))
153}
154
155impl GlobalBarrierWorkerContext for GlobalBarrierWorkerContextImpl {
156    #[await_tree::instrument]
157    async fn commit_epoch(&self, commit_info: CommitEpochInfo) -> MetaResult<HummockVersionStats> {
158        self.hummock_manager.commit_epoch(commit_info).await?;
159        Ok(self.hummock_manager.get_version_stats().await)
160    }
161
162    #[await_tree::instrument("next_scheduled_barrier")]
163    async fn next_scheduled(&self) -> Scheduled {
164        self.scheduled_barriers.next_scheduled().await
165    }
166
167    fn abort_and_mark_blocked(
168        &self,
169        database_id: Option<DatabaseId>,
170        recovery_reason: RecoveryReason,
171    ) {
172        if database_id.is_none() {
173            self.set_status(BarrierManagerStatus::Recovering(recovery_reason));
174        }
175
176        // Mark blocked and abort buffered schedules, they might be dirty already.
177        self.scheduled_barriers
178            .abort_and_mark_blocked(database_id, "cluster is under recovering");
179    }
180
181    fn mark_ready(&self, options: MarkReadyOptions) {
182        let is_global = matches!(&options, MarkReadyOptions::Global { .. });
183        self.scheduled_barriers.mark_ready(options);
184        if is_global {
185            self.set_status(BarrierManagerStatus::Running);
186        }
187    }
188
189    async fn resolve_log_store_epoch<'a>(
190        &'a self,
191        upstream_table_ids: impl Iterator<Item = TableId> + Send + 'a,
192        since_epoch: u64,
193    ) -> MetaResult<SinceTimestampResolvedEpoch> {
194        let upstream_table_ids = upstream_table_ids.collect::<Vec<_>>();
195        if upstream_table_ids.is_empty() {
196            return Err(
197                anyhow::anyhow!("since_timestamp requires at least one upstream table").into(),
198            );
199        }
200
201        self.hummock_manager
202            .on_current_version_and_table_change_log(|version, table_change_log| {
203                let mut unified_log_epochs = None;
204                for &upstream_table_id in &upstream_table_ids {
205                    let upstream_committed_epoch = version
206                        .state_table_info
207                        .info()
208                        .get(&upstream_table_id)
209                        .map(|info| info.committed_epoch)
210                        .ok_or_else(|| {
211                            anyhow::anyhow!(
212                                "cannot get committed epoch for upstream table {}",
213                                upstream_table_id
214                            )
215                        })?;
216                    let table_log_epochs = resolve_since_timestamp_log_store_epoch(
217                        upstream_table_id,
218                        since_epoch,
219                        upstream_committed_epoch,
220                        table_change_log,
221                    )?;
222                    if let Some(unified_log_epochs) = &unified_log_epochs {
223                        if unified_log_epochs != &table_log_epochs {
224                            return Err(anyhow::anyhow!(
225                                "resolved since_timestamp log epochs for upstream table {} do not match previous upstream table log epochs: {:?} vs {:?}",
226                                upstream_table_id,
227                                table_log_epochs,
228                                unified_log_epochs,
229                            )
230                            .into());
231                        }
232                    } else {
233                        unified_log_epochs = Some(table_log_epochs);
234                    }
235                }
236                unified_log_epochs.ok_or_else(|| {
237                    anyhow::anyhow!("since_timestamp requires at least one upstream table").into()
238                })
239            })
240            .await
241    }
242
243    async fn refresh_table_refill_runtime_state_after_recovery(&self) -> MetaResult<()> {
244        let policies = self
245            .metadata_manager
246            .catalog_controller
247            .table_cache_refill_policies_snapshot()
248            .await?;
249        let (serving_workers, fragment_serving_infos) =
250            fetch_serving_infos(&self.metadata_manager).await?;
251
252        self.env
253            .notification_manager()
254            .notify_hummock(
255                Operation::Update,
256                Info::TableRefillRuntimeConfig(PbTableRefillRuntimeConfig {
257                    table_cache_refill_policies: Some(policies),
258                    ..Default::default()
259                }),
260            )
261            .await;
262        sync_serving_table_vnode_mappings_to_hummock(
263            &self.env.notification_manager_ref(),
264            &self.serving_vnode_mapping,
265            &serving_workers,
266            &fragment_serving_infos,
267        )
268        .await;
269        Ok(())
270    }
271
272    #[await_tree::instrument("post_collect_command({command})")]
273    async fn post_collect_command(&self, command: PostCollectCommand) -> MetaResult<()> {
274        Box::pin(command.post_collect(self)).await
275    }
276
277    async fn notify_creating_job_failed(&self, database_id: Option<DatabaseId>, err: String) {
278        self.metadata_manager
279            .notify_finish_failed(database_id, err)
280            .await
281    }
282
283    #[await_tree::instrument("finish_creating_job({job})")]
284    async fn finish_creating_job(&self, job: TrackingJob) -> MetaResult<()> {
285        let job_id = job.job_id();
286        job.finish(&self.metadata_manager, &self.source_manager)
287            .await?;
288        self.env
289            .notification_manager()
290            .notify_local_subscribers(LocalNotification::StreamingJobBackfillFinished(job_id));
291        Ok(())
292    }
293
294    #[await_tree::instrument("finish_cdc_table_backfill({job})")]
295    async fn finish_cdc_table_backfill(&self, job: JobId) -> MetaResult<()> {
296        CdcTableBackfillTracker::mark_complete_job(&self.env.meta_store().conn, job).await
297    }
298
299    #[await_tree::instrument("new_control_stream({})", node.id)]
300    async fn new_control_stream(
301        &self,
302        node: &WorkerNode,
303        init_request: &PbInitRequest,
304    ) -> MetaResult<StreamingControlHandle> {
305        self.new_control_stream_impl(node, init_request).await
306    }
307
308    async fn reload_runtime_info(&self) -> MetaResult<BarrierWorkerRuntimeInfoSnapshot> {
309        self.reload_runtime_info_impl().await
310    }
311
312    async fn reload_database_runtime_info(
313        &self,
314        database_id: DatabaseId,
315    ) -> MetaResult<DatabaseRuntimeInfoSnapshot> {
316        self.reload_database_runtime_info_impl(database_id).await
317    }
318
319    async fn handle_list_finished_source_ids(
320        &self,
321        list_finished: Vec<PbListFinishedSource>,
322    ) -> MetaResult<()> {
323        let mut list_finished_info: HashMap<(TableId, SourceId), HashSet<ActorId>> = HashMap::new();
324
325        for list_finished in list_finished {
326            let table_id = list_finished.table_id;
327            let associated_source_id = list_finished.associated_source_id;
328            list_finished_info
329                .entry((table_id, associated_source_id))
330                .or_default()
331                .insert(list_finished.reporter_actor_id);
332        }
333
334        for ((table_id, associated_source_id), actors) in list_finished_info {
335            let allow_yield = self
336                .refresh_manager
337                .mark_list_stage_finished(table_id, &actors)?;
338
339            if !allow_yield {
340                continue;
341            }
342
343            let Some(database_id) = self
344                .get_source_database_id_for_refresh_stage(table_id, associated_source_id, "list")
345                .await?
346            else {
347                continue;
348            };
349
350            // Create ListFinish command
351            let list_finish_command = Command::ListFinish {
352                table_id,
353                associated_source_id,
354            };
355
356            // Schedule the command through the barrier system without waiting
357            self.barrier_scheduler
358                .run_command_no_wait(database_id, list_finish_command)
359                .context("Failed to schedule ListFinish command")?;
360
361            tracing::info!(
362                %table_id,
363                %associated_source_id,
364                "ListFinish command scheduled successfully"
365            );
366        }
367        Ok(())
368    }
369
370    async fn handle_load_finished_source_ids(
371        &self,
372        load_finished: Vec<PbLoadFinishedSource>,
373    ) -> MetaResult<()> {
374        let mut load_finished_info: HashMap<(TableId, SourceId), HashSet<ActorId>> = HashMap::new();
375
376        for load_finished in load_finished {
377            let table_id = load_finished.table_id;
378            let associated_source_id = load_finished.associated_source_id;
379            load_finished_info
380                .entry((table_id, associated_source_id))
381                .or_default()
382                .insert(load_finished.reporter_actor_id);
383        }
384
385        for ((table_id, associated_source_id), actors) in load_finished_info {
386            let allow_yield = self
387                .refresh_manager
388                .mark_load_stage_finished(table_id, &actors)?;
389
390            if !allow_yield {
391                continue;
392            }
393
394            let Some(database_id) = self
395                .get_source_database_id_for_refresh_stage(table_id, associated_source_id, "load")
396                .await?
397            else {
398                continue;
399            };
400
401            // Create LoadFinish command
402            let load_finish_command = Command::LoadFinish {
403                table_id,
404                associated_source_id,
405            };
406
407            // Schedule the command through the barrier system without waiting
408            self.barrier_scheduler
409                .run_command_no_wait(database_id, load_finish_command)
410                .context("Failed to schedule LoadFinish command")?;
411
412            tracing::info!(
413                %table_id,
414                %associated_source_id,
415                "LoadFinish command scheduled successfully"
416            );
417        }
418
419        Ok(())
420    }
421
422    async fn handle_refresh_finished_table_ids(
423        &self,
424        refresh_finished_table_job_ids: Vec<JobId>,
425    ) -> MetaResult<()> {
426        for job_id in refresh_finished_table_job_ids {
427            let table_id = job_id.as_mv_table_id();
428
429            self.refresh_manager.mark_refresh_complete(table_id).await?;
430        }
431
432        Ok(())
433    }
434
435    async fn load_batch_refresh_trigger_context(
436        &self,
437        job_id: JobId,
438        database_id: DatabaseId,
439        last_committed_epoch: u64,
440    ) -> MetaResult<BatchRefreshJobTriggerContext> {
441        self.load_batch_refresh_trigger_context_impl(job_id, database_id, last_committed_epoch)
442            .await
443    }
444
445    #[await_tree::instrument]
446    async fn pre_commit_iceberg_pk_index_sink_metadata(
447        &self,
448        reports: Vec<PbIcebergPkIndexSinkMetadata>,
449    ) -> MetaResult<Vec<SinkId>> {
450        let grouped = group_reports_by_sink(reports)?;
451        let success_ids: Vec<SinkId> = grouped.keys().cloned().collect();
452        let futs = FuturesUnordered::new();
453        for (sink_id, (prev_epoch, reports)) in grouped {
454            if reports.is_empty() {
455                continue;
456            }
457            let manager = &self.iceberg_pk_index_sink_manager;
458            futs.push(async move {
459                (
460                    sink_id,
461                    manager.pre_commit_epoch(sink_id, prev_epoch, reports).await,
462                )
463            });
464        }
465
466        // Drain all futures regardless of individual failures, so that no coordinator is left with
467        // state inconsistent vs. the caller's view.
468        let results: Vec<(SinkId, anyhow::Result<()>)> = futs.collect().await;
469        let errs: Vec<(SinkId, anyhow::Error)> = results
470            .into_iter()
471            .filter_map(|(id, r)| r.err().map(|e| (id, e)))
472            .collect();
473        if errs.is_empty() {
474            Ok(success_ids)
475        } else {
476            Err(aggregate_sink_errors("pre-commit", errs).into())
477        }
478    }
479
480    #[await_tree::instrument]
481    async fn commit_iceberg_pk_index_sink_metadata(&self, sink_ids: Vec<SinkId>) -> MetaResult<()> {
482        let futs = FuturesUnordered::new();
483        for sink_id in sink_ids {
484            let manager = &self.iceberg_pk_index_sink_manager;
485            futs.push(async move { (sink_id, manager.commit_epoch(sink_id).await) });
486        }
487
488        let results: Vec<(SinkId, anyhow::Result<()>)> = futs.collect().await;
489        let errs: Vec<(SinkId, anyhow::Error)> = results
490            .into_iter()
491            .filter_map(|(id, r)| r.err().map(|e| (id, e)))
492            .collect();
493        if errs.is_empty() {
494            Ok(())
495        } else {
496            Err(aggregate_sink_errors("commit", errs).into())
497        }
498    }
499
500    fn advance_iceberg_pk_index_sink_committed_epochs(
501        &self,
502        epochs: impl IntoIterator<Item = (PartialGraphId, u64)>,
503    ) {
504        self.iceberg_pk_index_sink_manager
505            .advance_committed_epochs(epochs);
506    }
507}
508
509/// Combine per-sink errors from a fan-out into a single `anyhow::Error`. The first failing
510/// sink's error is used as the source so the original chain is preserved; the message lists
511/// every failing `sink_id` and its error stringified.
512fn aggregate_sink_errors(
513    phase: &'static str,
514    mut errs: Vec<(SinkId, anyhow::Error)>,
515) -> anyhow::Error {
516    debug_assert!(!errs.is_empty());
517    let sink_ids: Vec<String> = errs.iter().map(|(id, _)| id.to_string()).collect();
518    let details: Vec<String> = errs
519        .iter()
520        .map(|(id, e)| format!("sink {}: {}", id, e.as_report()))
521        .collect();
522    // Preserve the first error's chain as the cause.
523    let (_, first_err) = errs.remove(0);
524    first_err.context(format!(
525        "iceberg v3 sink {} failed for sink_id(s) [{}]: {}",
526        phase,
527        sink_ids.join(", "),
528        details.join("; ")
529    ))
530}
531
532fn group_reports_by_sink(
533    reports: Vec<PbIcebergPkIndexSinkMetadata>,
534) -> MetaResult<HashMap<SinkId, (u64, Vec<PbIcebergPkIndexSinkMetadata>)>> {
535    let mut grouped: HashMap<SinkId, (u64, Vec<PbIcebergPkIndexSinkMetadata>)> = HashMap::new();
536    for r in reports {
537        let sink_id = r.sink_id;
538        let prev_epoch = r.prev_epoch;
539        let entry = grouped.entry(sink_id).or_insert((prev_epoch, Vec::new()));
540        if entry.0 != prev_epoch {
541            return Err(anyhow::anyhow!(
542                "iceberg v3 sink {} reports disagree on prev_epoch: {} vs {}",
543                sink_id,
544                entry.0,
545                prev_epoch
546            )
547            .into());
548        }
549        entry.1.push(r);
550    }
551    Ok(grouped)
552}
553
554impl GlobalBarrierWorkerContextImpl {
555    async fn get_source_database_id_for_refresh_stage(
556        &self,
557        table_id: TableId,
558        associated_source_id: SourceId,
559        stage: &'static str,
560    ) -> MetaResult<Option<DatabaseId>> {
561        match self
562            .metadata_manager
563            .catalog_controller
564            .get_object_database_id(associated_source_id)
565            .await
566        {
567            Ok(database_id) => Ok(Some(database_id)),
568            Err(err) if err.is_catalog_id_not_found("object") => {
569                tracing::warn!(
570                    %table_id,
571                    %associated_source_id,
572                    stage,
573                    "skip refresh finish command because associated source is already dropped"
574                );
575                Ok(None)
576            }
577            Err(err) => Err(err)
578                .with_context(|| {
579                    format!(
580                        "failed to get database id for refresh stage: table_id={}, associated_source_id={}, stage={stage}",
581                        table_id, associated_source_id
582                    )
583                })
584                .map_err(Into::into),
585        }
586    }
587
588    fn set_status(&self, new_status: BarrierManagerStatus) {
589        self.status.store(Arc::new(new_status));
590    }
591
592    /// Load the context metadata and resolve upstream log epochs for a batch refresh trigger.
593    async fn load_batch_refresh_trigger_context_impl(
594        &self,
595        job_id: JobId,
596        database_id: DatabaseId,
597        last_committed_epoch: u64,
598    ) -> MetaResult<BatchRefreshJobTriggerContext> {
599        use itertools::Itertools;
600        use sea_orm::TransactionTrait;
601
602        use crate::controller::scale::load_fragment_context_for_jobs;
603
604        // Load metadata from the catalog under a single transaction.
605        let inner = self
606            .metadata_manager
607            .catalog_controller
608            .get_inner_read_guard()
609            .await;
610        let txn = inner.db.begin().await?;
611
612        // 1. Load fragment context (job model, database model).
613        let fragment_context =
614            load_fragment_context_for_jobs(&txn, HashSet::from([job_id])).await?;
615
616        let streaming_job_model = fragment_context
617            .job_map
618            .get(&job_id)
619            .ok_or_else(|| anyhow::anyhow!("streaming job model not found for job {}", job_id))?
620            .clone();
621
622        let database_model = fragment_context
623            .database_map
624            .get(&database_id)
625            .ok_or_else(|| {
626                anyhow::anyhow!("database model not found for database {}", database_id)
627            })?;
628        let database_resource_group = database_model.resource_group.clone();
629
630        // 2. Load job definition.
631        let mut job_extra_info = self
632            .metadata_manager
633            .catalog_controller
634            .get_streaming_job_extra_info_in_txn(&txn, vec![job_id])
635            .await?;
636        let definition = job_extra_info
637            .remove(&job_id)
638            .ok_or_else(|| anyhow::anyhow!("extra info not found for job {}", job_id))?
639            .job_definition;
640
641        // 3. Get fragments from fragment_context and load downstream relations.
642        let fragments = fragment_context
643            .job_fragments
644            .get(&job_id)
645            .ok_or_else(|| anyhow::anyhow!("fragments not found for job {}", job_id))?
646            .clone();
647
648        // Derive upstream table IDs from the snapshot backfill scan nodes in the fragments.
649        let upstream_table_ids: HashSet<TableId> = {
650            use crate::stream::StreamFragmentGraph;
651            let snapshot_backfill_info = StreamFragmentGraph::collect_snapshot_backfill_info_impl(
652                fragments.values().map(|f| (&f.nodes, f.fragment_type_mask)),
653            )?
654            .0
655            .ok_or_else(|| {
656                anyhow::anyhow!("batch refresh job {} has no snapshot backfill info", job_id)
657            })?;
658            snapshot_backfill_info
659                .upstream_mv_table_id_to_backfill_epoch
660                .into_keys()
661                .collect()
662        };
663
664        let fragment_ids: Vec<_> = fragments.keys().copied().collect();
665        let downstreams = self
666            .metadata_manager
667            .catalog_controller
668            .get_fragment_downstream_relations_in_txn(&txn, fragment_ids)
669            .await?;
670
671        txn.commit().await?;
672        drop(inner);
673
674        // Resolve upstream log epochs from the hummock changelog.
675        let (upstream_table_log_epochs, target_upstream_epoch) = self
676            .hummock_manager
677            .on_current_version_and_table_change_log(|version, table_change_log| {
678                let mut target_upstream_epoch = last_committed_epoch;
679                let mut log_epochs: HashMap<TableId, Vec<(Vec<u64>, u64)>> = HashMap::new();
680
681                for &upstream_table_id in &upstream_table_ids {
682                    let upstream_committed_epoch = version
683                        .state_table_info
684                        .info()
685                        .get(&upstream_table_id)
686                        .map(|info| info.committed_epoch)
687                        .ok_or_else(|| {
688                            anyhow::anyhow!(
689                                "cannot get committed epoch for upstream table {}",
690                                upstream_table_id
691                            )
692                        })?;
693
694                    target_upstream_epoch =
695                        std::cmp::max(target_upstream_epoch, upstream_committed_epoch);
696
697                    if upstream_committed_epoch <= last_committed_epoch {
698                        continue;
699                    }
700
701                    if let Some(change_log) = table_change_log.get(&upstream_table_id) {
702                        let epochs = change_log
703                            .filter_epoch((last_committed_epoch, upstream_committed_epoch))
704                            .map(|epoch_log| {
705                                (
706                                    epoch_log.non_checkpoint_epochs.clone(),
707                                    epoch_log.checkpoint_epoch,
708                                )
709                            })
710                            .collect_vec();
711                        if !epochs.is_empty() {
712                            log_epochs.insert(upstream_table_id, epochs);
713                        }
714                    } else {
715                        anyhow::bail!(
716                            "upstream table {} has lagged downstream on epoch {} but no table change log (upstream committed: {})",
717                            upstream_table_id,
718                            last_committed_epoch,
719                            upstream_committed_epoch,
720                        );
721                    }
722                }
723
724                Ok((log_epochs, target_upstream_epoch))
725            })
726            .await?;
727
728        Ok(BatchRefreshJobTriggerContext {
729            fragments,
730            downstreams,
731            streaming_job_model,
732            definition,
733            database_resource_group,
734            upstream_table_log_epochs,
735            target_upstream_epoch,
736        })
737    }
738}
739
740impl PostCollectCommand {
741    /// Do some stuffs after barriers are collected and the new storage version is committed, for
742    /// the given command.
743    pub async fn post_collect(
744        self,
745        barrier_manager_context: &GlobalBarrierWorkerContextImpl,
746    ) -> MetaResult<()> {
747        match self {
748            PostCollectCommand::Command(_) => {}
749            PostCollectCommand::SourceChangeSplit {
750                split_assignment: assignment,
751            } => {
752                barrier_manager_context
753                    .metadata_manager
754                    .update_fragment_splits(&assignment)
755                    .await?;
756            }
757
758            PostCollectCommand::DropStreamingJobs => {}
759            PostCollectCommand::ConnectorPropsChange(obj_id_map_props) => {
760                // todo: we dont know the type of the object id, it can be a source or a sink. Should carry more info in the barrier command.
761                barrier_manager_context
762                    .source_manager
763                    .apply_source_change(SourceChange::UpdateSourceProps {
764                        // Only sources are managed in source manager. Convert object IDs to source IDs and let
765                        // source manager ignore unknown/unregistered sources.
766                        source_id_map_new_props: obj_id_map_props
767                            .iter()
768                            .map(|(object_id, props)| (object_id.as_source_id(), props.clone()))
769                            .collect(),
770                    })
771                    .await;
772            }
773            PostCollectCommand::ResumeBackfill { target } => match target {
774                ResumeBackfillTarget::Job(job_id) => {
775                    barrier_manager_context
776                        .metadata_manager
777                        .catalog_controller
778                        .update_backfill_orders_by_job_id(job_id, None)
779                        .await?;
780                }
781                ResumeBackfillTarget::Fragment(fragment_id) => {
782                    let mut job_ids = barrier_manager_context
783                        .metadata_manager
784                        .catalog_controller
785                        .get_fragment_job_id(vec![fragment_id])
786                        .await?;
787                    let job_id = job_ids.pop().ok_or_else(|| {
788                        MetaError::invalid_parameter("fragment not found".to_owned())
789                    })?;
790                    let job_id = JobId::new(job_id.as_raw_id());
791
792                    let extra_info = barrier_manager_context
793                        .metadata_manager
794                        .catalog_controller
795                        .get_streaming_job_extra_info(vec![job_id])
796                        .await?;
797                    let mut backfill_orders: BackfillOrders = extra_info
798                        .get(&job_id)
799                        .cloned()
800                        .ok_or_else(|| MetaError::invalid_parameter("job not found".to_owned()))?
801                        .backfill_orders
802                        .unwrap_or_default();
803
804                    let resumed_fragment_id = fragment_id.as_raw_id();
805                    for children in backfill_orders.0.values_mut() {
806                        children.retain(|child| *child != resumed_fragment_id);
807                    }
808                    backfill_orders.0.retain(|_, children| !children.is_empty());
809
810                    barrier_manager_context
811                        .metadata_manager
812                        .catalog_controller
813                        .update_backfill_orders_by_job_id(job_id, Some(backfill_orders))
814                        .await?;
815                }
816            },
817            PostCollectCommand::CreateStreamingJob {
818                info,
819                job_type,
820                cross_db_snapshot_backfill_info,
821                resolved_split_assignment,
822            } => {
823                match &job_type {
824                    CreateStreamingJobType::SinkIntoTable(_) | CreateStreamingJobType::Normal => {
825                        barrier_manager_context
826                            .metadata_manager
827                            .catalog_controller
828                            .fill_snapshot_backfill_epoch(
829                                info.stream_job_fragments.fragments.iter().filter_map(
830                                    |(fragment_id, fragment)| {
831                                        if fragment.fragment_type_mask.contains(
832                                            FragmentTypeFlag::CrossDbSnapshotBackfillStreamScan,
833                                        ) {
834                                            Some(*fragment_id as _)
835                                        } else {
836                                            None
837                                        }
838                                    },
839                                ),
840                                None,
841                                &cross_db_snapshot_backfill_info,
842                            )
843                            .await?
844                    }
845                    CreateStreamingJobType::SnapshotBackfill {
846                        snapshot_backfill_info,
847                        ..
848                    }
849                    | CreateStreamingJobType::BatchRefresh(BatchRefreshInfo {
850                        snapshot_backfill_info,
851                        ..
852                    }) => {
853                        barrier_manager_context
854                            .metadata_manager
855                            .catalog_controller
856                            .fill_snapshot_backfill_epoch(
857                                info.stream_job_fragments.fragments.iter().filter_map(
858                                    |(fragment_id, fragment)| {
859                                        if fragment.fragment_type_mask.contains_any([
860                                            FragmentTypeFlag::SnapshotBackfillStreamScan,
861                                            FragmentTypeFlag::CrossDbSnapshotBackfillStreamScan,
862                                        ]) {
863                                            Some(*fragment_id as _)
864                                        } else {
865                                            None
866                                        }
867                                    },
868                                ),
869                                Some(snapshot_backfill_info),
870                                &cross_db_snapshot_backfill_info,
871                            )
872                            .await?
873                    }
874                }
875
876                // Do `post_collect_job_fragments` of the original streaming job in the end, so that in any previous failure,
877                // we won't mark the job as `Creating`, and then the job will be later clean by the recovery triggered by the returned error.
878                let CreateStreamingJobCommandInfo {
879                    stream_job_fragments,
880                    upstream_fragment_downstreams,
881                    replace_sink,
882                    ..
883                } = info;
884                let new_job_id = stream_job_fragments.stream_job_id();
885                let new_sink_downstream =
886                    if let CreateStreamingJobType::SinkIntoTable(ctx) = job_type {
887                        let new_downstreams = ctx.new_sink_downstream.clone();
888                        let new_downstreams = FragmentDownstreamRelation::from([(
889                            ctx.sink_fragment_id,
890                            vec![new_downstreams],
891                        )]);
892                        Some(new_downstreams)
893                    } else {
894                        None
895                    };
896
897                let old_state_table_ids = barrier_manager_context
898                    .metadata_manager
899                    .catalog_controller
900                    .post_collect_job_fragments(
901                        new_job_id,
902                        &upstream_fragment_downstreams,
903                        new_sink_downstream,
904                        Some(&resolved_split_assignment),
905                        replace_sink.as_ref(),
906                        replace_sink.is_none(),
907                    )
908                    .await?;
909
910                let added_source_fragments = stream_job_fragments.stream_source_fragments();
911                let added_backfill_fragments = stream_job_fragments.source_backfill_fragments();
912                // Skip the source-manager notification (and its `core` lock) when the job has no source/backfill fragments.
913                if !added_source_fragments.is_empty() || !added_backfill_fragments.is_empty() {
914                    barrier_manager_context
915                        .source_manager
916                        .apply_source_change(SourceChange::CreateJob {
917                            added_source_fragments,
918                            added_backfill_fragments,
919                        })
920                        .await;
921                }
922
923                if let Some(old_sink_id) = replace_sink {
924                    barrier_manager_context
925                        .sink_manager
926                        .stop_sink_coordinator(vec![old_sink_id])
927                        .await;
928                    cleanup_dropped_streaming_jobs(
929                        &barrier_manager_context.refresh_manager,
930                        &barrier_manager_context.hummock_manager,
931                        &barrier_manager_context.metadata_manager,
932                        [old_sink_id.as_job_id()],
933                        old_state_table_ids.expect("replace sink should return old state tables"),
934                        "replace_sink",
935                    )
936                    .await?;
937                }
938            }
939            PostCollectCommand::Reschedule { reschedules, .. } => {
940                let fragment_splits = reschedules
941                    .iter()
942                    .map(|(fragment_id, reschedule)| {
943                        (*fragment_id, reschedule.actor_splits.clone())
944                    })
945                    .collect();
946
947                barrier_manager_context
948                    .metadata_manager
949                    .update_fragment_splits(&fragment_splits)
950                    .await?;
951            }
952
953            PostCollectCommand::ReplaceStreamJob {
954                plan: replace_plan,
955                resolved_split_assignment,
956            } => {
957                let ReplaceStreamJobPlan {
958                    old_fragments,
959                    new_fragments,
960                    upstream_fragment_downstreams,
961                    to_drop_state_table_ids,
962                    auto_refresh_schema_sinks,
963                    ..
964                } = &replace_plan;
965                // Update actors and actor_dispatchers for new table fragments.
966                barrier_manager_context
967                    .metadata_manager
968                    .catalog_controller
969                    .post_collect_job_fragments(
970                        new_fragments.stream_job_id,
971                        upstream_fragment_downstreams,
972                        None,
973                        Some(&resolved_split_assignment),
974                        None,
975                        false,
976                    )
977                    .await?;
978
979                if let Some(sinks) = auto_refresh_schema_sinks {
980                    for sink in sinks {
981                        barrier_manager_context
982                            .metadata_manager
983                            .catalog_controller
984                            .post_collect_job_fragments(
985                                sink.tmp_sink_id.as_job_id(),
986                                &Default::default(), // upstream_fragment_downstreams is already inserted in the job of upstream table
987                                None, // no replace plan
988                                None, // no init split assignment
989                                None,
990                                false,
991                            )
992                            .await?;
993                    }
994                }
995
996                // Apply the split changes in source manager.
997                barrier_manager_context
998                    .source_manager
999                    .handle_replace_job(
1000                        old_fragments,
1001                        new_fragments.stream_source_fragments(),
1002                        &replace_plan,
1003                    )
1004                    .await;
1005                cleanup_dropped_streaming_jobs(
1006                    &barrier_manager_context.refresh_manager,
1007                    &barrier_manager_context.hummock_manager,
1008                    &barrier_manager_context.metadata_manager,
1009                    [],
1010                    to_drop_state_table_ids.clone(),
1011                    "replace_streaming_job",
1012                )
1013                .await?;
1014            }
1015
1016            PostCollectCommand::CreateSubscription { subscription_id } => {
1017                barrier_manager_context
1018                    .metadata_manager
1019                    .catalog_controller
1020                    .finish_create_subscription_catalog(subscription_id)
1021                    .await?
1022            }
1023        }
1024
1025        Ok(())
1026    }
1027}
1028
1029#[cfg(test)]
1030mod tests {
1031    use risingwave_hummock_sdk::change_log::{EpochNewChangeLog, TableChangeLog};
1032
1033    use super::*;
1034
1035    fn test_change_logs(table_id: TableId) -> TableChangeLogs {
1036        TableChangeLogs::from_iter([(
1037            table_id,
1038            TableChangeLog::new([
1039                EpochNewChangeLog {
1040                    new_value: vec![],
1041                    old_value: vec![],
1042                    non_checkpoint_epochs: vec![10],
1043                    checkpoint_epoch: 20,
1044                },
1045                EpochNewChangeLog {
1046                    new_value: vec![],
1047                    old_value: vec![],
1048                    non_checkpoint_epochs: vec![30],
1049                    checkpoint_epoch: 40,
1050                },
1051                EpochNewChangeLog {
1052                    new_value: vec![],
1053                    old_value: vec![],
1054                    non_checkpoint_epochs: vec![50],
1055                    checkpoint_epoch: 60,
1056                },
1057                EpochNewChangeLog {
1058                    new_value: vec![],
1059                    old_value: vec![],
1060                    non_checkpoint_epochs: vec![],
1061                    checkpoint_epoch: 80,
1062                },
1063            ]),
1064        )])
1065    }
1066
1067    #[test]
1068    fn test_resolve_since_timestamp_log_store_epoch() {
1069        let table_id = TableId::new(233);
1070        let change_logs = test_change_logs(table_id);
1071
1072        let resolved =
1073            resolve_since_timestamp_log_store_epoch(table_id, 45, 80, &change_logs).unwrap();
1074
1075        assert_eq!(resolved, (60, vec![(vec![], 80)]));
1076    }
1077
1078    #[test]
1079    fn test_resolve_since_timestamp_log_store_epoch_rejects_invalid_range() {
1080        let table_id = TableId::new(233);
1081        let change_logs = test_change_logs(table_id);
1082
1083        let err = resolve_since_timestamp_log_store_epoch(table_id, 9, 80, &change_logs)
1084            .unwrap_err()
1085            .to_string();
1086        assert!(err.contains("earlier than the retained changelog"));
1087
1088        let err = resolve_since_timestamp_log_store_epoch(table_id, 80, 80, &change_logs)
1089            .unwrap_err()
1090            .to_string();
1091        assert!(err.contains("is not before the committed epoch"));
1092
1093        let err = resolve_since_timestamp_log_store_epoch(table_id, 45, 60, &change_logs)
1094            .unwrap_err()
1095            .to_string();
1096        assert!(err.contains("does not match latest changelog epoch"));
1097    }
1098
1099    #[test]
1100    fn test_resolve_since_timestamp_log_store_epoch_rejects_empty_or_missing_log() {
1101        let table_id = TableId::new(233);
1102        let empty_change_logs = TableChangeLogs::from_iter([(table_id, TableChangeLog::new([]))]);
1103
1104        let err = resolve_since_timestamp_log_store_epoch(table_id, 30, 60, &empty_change_logs)
1105            .unwrap_err()
1106            .to_string();
1107        assert!(err.contains("empty table changelog"));
1108
1109        let err = resolve_since_timestamp_log_store_epoch(table_id, 30, 60, &Default::default())
1110            .unwrap_err()
1111            .to_string();
1112        assert!(err.contains("no table changelog"));
1113    }
1114
1115    #[test]
1116    fn test_resolve_since_timestamp_log_store_epoch_keeps_latest_non_checkpoint_epochs() {
1117        let table_id = TableId::new(233);
1118        let change_logs = TableChangeLogs::from_iter([(
1119            table_id,
1120            TableChangeLog::new([
1121                EpochNewChangeLog {
1122                    new_value: vec![],
1123                    old_value: vec![],
1124                    non_checkpoint_epochs: vec![],
1125                    checkpoint_epoch: 20,
1126                },
1127                EpochNewChangeLog {
1128                    new_value: vec![],
1129                    old_value: vec![],
1130                    non_checkpoint_epochs: vec![30],
1131                    checkpoint_epoch: 40,
1132                },
1133            ]),
1134        )]);
1135
1136        let resolved =
1137            resolve_since_timestamp_log_store_epoch(table_id, 20, 40, &change_logs).unwrap();
1138        assert_eq!(resolved, (20, vec![(vec![30], 40)]));
1139
1140        let err = resolve_since_timestamp_log_store_epoch(table_id, 30, 40, &change_logs)
1141            .unwrap_err()
1142            .to_string();
1143        assert!(err.contains("since_timestamp is too new"));
1144    }
1145
1146    #[test]
1147    fn test_skip_refresh_finish_when_associated_source_missing() {
1148        let err = MetaError::catalog_id_not_found("object", 42);
1149        assert!(err.is_catalog_id_not_found("object"));
1150    }
1151
1152    #[test]
1153    fn test_do_not_skip_refresh_finish_for_other_not_found_types() {
1154        let err = MetaError::catalog_id_not_found("table", 42);
1155        assert!(!err.is_catalog_id_not_found("object"));
1156    }
1157}