Skip to main content

risingwave_meta/barrier/
rpc.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::hash_map::Entry;
16use std::collections::{HashMap, HashSet};
17use std::error::Error;
18use std::fmt::{Debug, Formatter};
19use std::future::poll_fn;
20use std::sync::Arc;
21use std::task::{Context, Poll};
22use std::time::Duration;
23
24use anyhow::anyhow;
25use fail::fail_point;
26use futures::future::{BoxFuture, join_all};
27use futures::{FutureExt, StreamExt};
28use itertools::Itertools;
29use risingwave_common::bail;
30use risingwave_common::catalog::{DatabaseId, FragmentTypeFlag, TableId};
31use risingwave_common::id::JobId;
32use risingwave_common::util::epoch::Epoch;
33use risingwave_common::util::retry::exponential_backoff;
34use risingwave_common::util::stream_graph_visitor::visit_stream_node_cont;
35use risingwave_common::util::tracing::TracingContext;
36use risingwave_connector::source::SplitImpl;
37use risingwave_meta_model::WorkerId;
38use risingwave_pb::common::{HostAddress, WorkerNode};
39use risingwave_pb::hummock::HummockVersionStats;
40use risingwave_pb::id::PartialGraphId;
41use risingwave_pb::source::{PbCdcTableSnapshotSplits, PbCdcTableSnapshotSplitsWithGeneration};
42use risingwave_pb::stream_plan::barrier_mutation::Mutation;
43use risingwave_pb::stream_plan::stream_node::NodeBody;
44use risingwave_pb::stream_plan::{
45    AddMutation, Barrier, BarrierMutation, IcebergPkIndexCompactionContext,
46};
47use risingwave_pb::stream_service::inject_barrier_request::build_actor_info::UpstreamActors;
48use risingwave_pb::stream_service::inject_barrier_request::{
49    BuildActorInfo, FragmentBuildActorInfo,
50};
51use risingwave_pb::stream_service::streaming_control_stream_request::{
52    CreatePartialGraphRequest, PbCreatePartialGraphRequest, PbInitRequest,
53    RemovePartialGraphRequest, ResetPartialGraphsRequest,
54};
55use risingwave_pb::stream_service::{
56    InjectBarrierRequest, StreamingControlStreamRequest, streaming_control_stream_request,
57    streaming_control_stream_response,
58};
59use risingwave_rpc_client::StreamingControlHandle;
60use thiserror_ext::AsReport;
61use tokio::time::{Instant, sleep};
62use tracing::{debug, error, info, warn};
63use uuid::Uuid;
64
65use super::{BarrierKind, TracedEpoch};
66use crate::barrier::BackfillOrderState;
67use crate::barrier::backfill_order_control::get_nodes_with_backfill_dependencies;
68use crate::barrier::cdc_progress::CdcTableBackfillTracker;
69use crate::barrier::checkpoint::{
70    BarrierWorkerState, BatchRefreshJobCheckpointControl, BatchRefreshRenderResult,
71    CreatingStreamingJobControl, DatabaseCheckpointControl, DatabaseCheckpointControlMetrics,
72    IndependentCheckpointJobControl,
73};
74use crate::barrier::context::{GlobalBarrierWorkerContext, GlobalBarrierWorkerContextImpl};
75use crate::barrier::edge_builder::{EdgeBuilderFragmentInfo, FragmentEdgeBuilder};
76use crate::barrier::info::{
77    BarrierInfo, CreateStreamingJobStatus, InflightDatabaseInfo, InflightStreamingJobInfo,
78    SubscriberType,
79};
80use crate::barrier::partial_graph::PartialGraphRecoverer;
81use crate::barrier::progress::CreateMviewProgressTracker;
82use crate::barrier::utils::NodeToCollect;
83use crate::controller::fragment::InflightFragmentInfo;
84use crate::controller::utils::StreamingJobExtraInfo;
85use crate::manager::MetaSrvEnv;
86use crate::model::{
87    ActorId, FragmentDownstreamRelation, FragmentId, StreamActor, StreamJobActorsToCreate,
88    SubscriptionId,
89};
90use crate::stream::cdc::{
91    CdcTableSnapshotSplits, is_parallelized_backfill_enabled_cdc_scan_fragment,
92};
93use crate::stream::{
94    ExtendedFragmentBackfillOrder, StreamFragmentGraph, UserDefinedFragmentBackfillOrder,
95    build_actor_connector_splits,
96};
97use crate::{MetaError, MetaResult};
98
99pub(crate) fn to_partial_graph_id(
100    database_id: DatabaseId,
101    creating_job_id: Option<JobId>,
102) -> PartialGraphId {
103    let raw_job_id = creating_job_id
104        .map(|job_id| {
105            assert_ne!(job_id, u32::MAX);
106            job_id.as_raw_id()
107        })
108        .unwrap_or(u32::MAX);
109    (((database_id.as_raw_id() as u64) << 32) | (raw_job_id as u64)).into()
110}
111
112pub(super) fn from_partial_graph_id(
113    partial_graph_id: PartialGraphId,
114) -> (DatabaseId, Option<JobId>) {
115    let id = partial_graph_id.as_raw_id();
116    let database_id = (id >> 32) as u32;
117    let raw_creating_job_id = (id & ((1 << 32) - 1)) as u32;
118    let creating_job_id = if raw_creating_job_id == u32::MAX {
119        None
120    } else {
121        Some(JobId::new(raw_creating_job_id))
122    };
123    (database_id.into(), creating_job_id)
124}
125
126pub(super) fn build_locality_fragment_state_table_mapping(
127    fragment_infos: &HashMap<FragmentId, InflightFragmentInfo>,
128) -> HashMap<FragmentId, Vec<TableId>> {
129    let mut mapping = HashMap::new();
130
131    for (fragment_id, fragment_info) in fragment_infos {
132        let mut state_table_ids = Vec::new();
133        visit_stream_node_cont(&fragment_info.nodes, |stream_node| {
134            if let Some(NodeBody::LocalityProvider(locality_provider)) =
135                stream_node.node_body.as_ref()
136            {
137                let state_table_id = locality_provider
138                    .state_table
139                    .as_ref()
140                    .expect("must have state table")
141                    .id;
142                state_table_ids.push(state_table_id);
143                false
144            } else {
145                true
146            }
147        });
148        if !state_table_ids.is_empty() {
149            mapping.insert(*fragment_id, state_table_ids);
150        }
151    }
152
153    mapping
154}
155
156pub(super) fn database_partial_graphs<'a>(
157    database_id: DatabaseId,
158    creating_jobs: impl Iterator<Item = JobId> + Sized + 'a,
159) -> impl Iterator<Item = PartialGraphId> + 'a {
160    creating_jobs
161        .map(Some)
162        .chain([None])
163        .map(move |creating_job_id| to_partial_graph_id(database_id, creating_job_id))
164}
165
166struct ControlStreamNode {
167    worker_id: WorkerId,
168    host: HostAddress,
169    handle: StreamingControlHandle,
170}
171
172enum WorkerNodeState {
173    Connected {
174        control_stream: ControlStreamNode,
175        removed: bool,
176    },
177    Reconnecting(BoxFuture<'static, StreamingControlHandle>),
178}
179
180pub(super) struct ControlStreamManager {
181    workers: HashMap<WorkerId, (WorkerNode, WorkerNodeState)>,
182    pub env: MetaSrvEnv,
183}
184
185impl ControlStreamManager {
186    pub(super) fn new(env: MetaSrvEnv) -> Self {
187        Self {
188            workers: Default::default(),
189            env,
190        }
191    }
192
193    pub(super) fn host_addr(&self, worker_id: WorkerId) -> HostAddress {
194        self.workers[&worker_id].0.host.clone().unwrap()
195    }
196
197    pub(super) async fn add_worker(
198        &mut self,
199        node: WorkerNode,
200        partial_graphs: impl Iterator<Item = PartialGraphId>,
201        term_id: &String,
202        context: Arc<impl GlobalBarrierWorkerContext>,
203    ) {
204        let node_id = node.id;
205        if let Entry::Occupied(entry) = self.workers.entry(node_id) {
206            let (existing_node, worker_state) = entry.get();
207            assert_eq!(existing_node.host, node.host);
208            warn!(id = %node.id, host = ?node.host, "node already exists");
209            match worker_state {
210                WorkerNodeState::Connected { .. } => {
211                    warn!(id = %node.id, host = ?node.host, "new node already connected");
212                    return;
213                }
214                WorkerNodeState::Reconnecting(_) => {
215                    warn!(id = %node.id, host = ?node.host, "remove previous pending worker connect request and reconnect");
216                    entry.remove();
217                }
218            }
219        }
220        let node_host = node.host.clone().unwrap();
221        let mut backoff =
222            exponential_backoff(Duration::from_millis(100), 5, Duration::from_secs(3));
223        const MAX_RETRY: usize = 5;
224        for i in 1..=MAX_RETRY {
225            match context
226                .new_control_stream(
227                    &node,
228                    &PbInitRequest {
229                        term_id: term_id.clone(),
230                    },
231                )
232                .await
233            {
234                Ok(mut handle) => {
235                    WorkerNodeConnected {
236                        handle: &mut handle,
237                        node: &node,
238                    }
239                    .initialize(partial_graphs);
240                    info!(?node_host, "add control stream worker");
241                    assert!(
242                        self.workers
243                            .insert(
244                                node_id,
245                                (
246                                    node,
247                                    WorkerNodeState::Connected {
248                                        control_stream: ControlStreamNode {
249                                            worker_id: node_id as _,
250                                            host: node_host,
251                                            handle,
252                                        },
253                                        removed: false
254                                    }
255                                )
256                            )
257                            .is_none()
258                    );
259                    return;
260                }
261                Err(e) => {
262                    // It may happen that the dns information of newly registered worker node
263                    // has not been propagated to the meta node and cause error. Wait for a while and retry
264                    let delay = backoff.next().unwrap();
265                    error!(
266                        attempt = i,
267                        backoff_delay = ?delay,
268                        err = %e.as_report(),
269                        ?node_host,
270                        "failed to resolve the worker node address",
271                    );
272                    sleep(delay).await;
273                }
274            }
275        }
276        error!(?node_host, "failed to create the worker node after retries");
277        assert!(
278            self.workers
279                .insert(
280                    node_id,
281                    (
282                        node.clone(),
283                        WorkerNodeState::Reconnecting(ControlStreamManager::retry_connect(
284                            node,
285                            term_id.to_owned(),
286                            context,
287                        ))
288                    )
289                )
290                .is_none()
291        );
292    }
293
294    pub(super) fn remove_worker(&mut self, node: WorkerNode) {
295        if let Entry::Occupied(mut entry) = self.workers.entry(node.id) {
296            let (_, worker_state) = entry.get_mut();
297            match worker_state {
298                WorkerNodeState::Connected { removed, .. } => {
299                    info!(worker_id = %node.id, "mark connected worker as removed");
300                    *removed = true;
301                }
302                WorkerNodeState::Reconnecting(_) => {
303                    info!(worker_id = %node.id, "remove worker");
304                    entry.remove();
305                }
306            }
307        }
308    }
309
310    fn retry_connect(
311        node: WorkerNode,
312        term_id: String,
313        context: Arc<impl GlobalBarrierWorkerContext>,
314    ) -> BoxFuture<'static, StreamingControlHandle> {
315        async move {
316            let mut attempt = 0;
317            let backoff = exponential_backoff(
318                Duration::from_millis(100),
319                5,
320                Duration::from_mins(1),
321            );
322            let init_request = PbInitRequest { term_id };
323            for delay in backoff {
324                attempt += 1;
325                sleep(delay).await;
326                match context.new_control_stream(&node, &init_request).await {
327                    Ok(handle) => {
328                        return handle;
329                    }
330                    Err(e) => {
331                        warn!(e = %e.as_report(), ?node, attempt, "failed to create the control stream worker");
332                    }
333                }
334            }
335            unreachable!("end of retry backoff")
336        }.boxed()
337    }
338
339    pub(super) async fn recover(
340        env: MetaSrvEnv,
341        nodes: &HashMap<WorkerId, WorkerNode>,
342        term_id: &str,
343        context: Arc<impl GlobalBarrierWorkerContext>,
344    ) -> Self {
345        let reset_start_time = Instant::now();
346        let init_request = PbInitRequest {
347            term_id: term_id.to_owned(),
348        };
349        let init_request = &init_request;
350        let nodes = join_all(nodes.iter().map(|(worker_id, node)| async {
351            let result = context.new_control_stream(node, init_request).await;
352            (*worker_id, node.clone(), result)
353        }))
354        .await;
355        let mut unconnected_workers = HashSet::new();
356        let mut workers = HashMap::new();
357        for (worker_id, node, result) in nodes {
358            match result {
359                Ok(handle) => {
360                    let control_stream = ControlStreamNode {
361                        worker_id: node.id,
362                        host: node.host.clone().unwrap(),
363                        handle,
364                    };
365                    assert!(
366                        workers
367                            .insert(
368                                worker_id,
369                                (
370                                    node,
371                                    WorkerNodeState::Connected {
372                                        control_stream,
373                                        removed: false
374                                    }
375                                )
376                            )
377                            .is_none()
378                    );
379                }
380                Err(e) => {
381                    unconnected_workers.insert(worker_id);
382                    warn!(
383                        e = %e.as_report(),
384                        %worker_id,
385                        ?node,
386                        "failed to connect to node"
387                    );
388                    assert!(
389                        workers
390                            .insert(
391                                worker_id,
392                                (
393                                    node.clone(),
394                                    WorkerNodeState::Reconnecting(Self::retry_connect(
395                                        node,
396                                        term_id.to_owned(),
397                                        context.clone()
398                                    ))
399                                )
400                            )
401                            .is_none()
402                    );
403                }
404            }
405        }
406
407        info!(elapsed=?reset_start_time.elapsed(), ?unconnected_workers, "control stream reset");
408
409        Self { workers, env }
410    }
411
412    /// Clear all nodes and response streams in the manager.
413    pub(super) fn clear(&mut self) {
414        *self = Self::new(self.env.clone());
415    }
416}
417
418pub(super) struct WorkerNodeConnected<'a> {
419    node: &'a WorkerNode,
420    handle: &'a mut StreamingControlHandle,
421}
422
423impl<'a> WorkerNodeConnected<'a> {
424    pub(super) fn initialize(self, partial_graphs: impl Iterator<Item = PartialGraphId>) {
425        for partial_graph_id in partial_graphs {
426            if let Err(e) = self.handle.send_request(StreamingControlStreamRequest {
427                request: Some(
428                    streaming_control_stream_request::Request::CreatePartialGraph(
429                        PbCreatePartialGraphRequest { partial_graph_id },
430                    ),
431                ),
432            }) {
433                warn!(e = %e.as_report(), node = ?self.node, "failed to send initial partial graph request");
434            }
435        }
436    }
437}
438
439pub(super) enum WorkerNodeEvent<'a> {
440    Response(MetaResult<streaming_control_stream_response::Response>),
441    Connected(WorkerNodeConnected<'a>),
442}
443
444impl ControlStreamManager {
445    fn poll_next_event<'a>(
446        this_opt: &mut Option<&'a mut Self>,
447        cx: &mut Context<'_>,
448        term_id: &str,
449        context: &Arc<impl GlobalBarrierWorkerContext>,
450        poll_reconnect: bool,
451    ) -> Poll<(WorkerId, WorkerNodeEvent<'a>)> {
452        let this = this_opt.as_mut().expect("Future polled after completion");
453        if this.workers.is_empty() {
454            return Poll::Pending;
455        }
456        {
457            for (&worker_id, (node, worker_state)) in &mut this.workers {
458                let control_stream = match worker_state {
459                    WorkerNodeState::Connected { control_stream, .. } => control_stream,
460                    WorkerNodeState::Reconnecting(_) if !poll_reconnect => {
461                        continue;
462                    }
463                    WorkerNodeState::Reconnecting(join_handle) => {
464                        match join_handle.poll_unpin(cx) {
465                            Poll::Ready(handle) => {
466                                info!(id=%node.id, host=?node.host, "reconnected to worker");
467                                *worker_state = WorkerNodeState::Connected {
468                                    control_stream: ControlStreamNode {
469                                        worker_id: node.id,
470                                        host: node.host.clone().unwrap(),
471                                        handle,
472                                    },
473                                    removed: false,
474                                };
475                                let this = this_opt.take().expect("should exist");
476                                let (node, worker_state) =
477                                    this.workers.get_mut(&worker_id).expect("should exist");
478                                let WorkerNodeState::Connected { control_stream, .. } =
479                                    worker_state
480                                else {
481                                    unreachable!()
482                                };
483                                return Poll::Ready((
484                                    worker_id,
485                                    WorkerNodeEvent::Connected(WorkerNodeConnected {
486                                        node,
487                                        handle: &mut control_stream.handle,
488                                    }),
489                                ));
490                            }
491                            Poll::Pending => {
492                                continue;
493                            }
494                        }
495                    }
496                };
497                match control_stream.handle.response_stream.poll_next_unpin(cx) {
498                    Poll::Ready(result) => {
499                        {
500                            let result = result
501                                .ok_or_else(|| (false, anyhow!("end of stream").into()))
502                                .and_then(|result| {
503                                    result.map_err(|err| -> (bool, MetaError) { (false, err.into()) }).and_then(|resp| {
504                                        match resp
505                                            .response
506                                            .ok_or_else(|| (false, anyhow!("empty response").into()))?
507                                        {
508                                            streaming_control_stream_response::Response::Shutdown(_) => Err((true, anyhow!(
509                                                "worker node {worker_id} is shutting down"
510                                            )
511                                                .into())),
512                                            streaming_control_stream_response::Response::Init(_) => {
513                                                // This arm should be unreachable.
514                                                Err((false, anyhow!("get unexpected init response").into()))
515                                            }
516                                            resp => {
517                                                if let streaming_control_stream_response::Response::CompleteBarrier(barrier_resp) = &resp {
518                                                    assert_eq!(worker_id, barrier_resp.worker_id);
519                                                }
520                                                Ok(resp)
521                                            }
522                                        }
523                                    })
524                                });
525                            let result = match result {
526                                Ok(resp) => Ok(resp),
527                                Err((shutdown, err)) => {
528                                    warn!(worker_id = %node.id, host = ?node.host, err = %err.as_report(), "get error from response stream");
529                                    let WorkerNodeState::Connected { removed, .. } = worker_state
530                                    else {
531                                        unreachable!("checked connected")
532                                    };
533                                    if *removed || shutdown {
534                                        this.workers.remove(&worker_id);
535                                    } else {
536                                        *worker_state = WorkerNodeState::Reconnecting(
537                                            ControlStreamManager::retry_connect(
538                                                node.clone(),
539                                                term_id.to_owned(),
540                                                context.clone(),
541                                            ),
542                                        );
543                                    }
544                                    Err(err)
545                                }
546                            };
547                            return Poll::Ready((worker_id, WorkerNodeEvent::Response(result)));
548                        }
549                    }
550                    Poll::Pending => {
551                        continue;
552                    }
553                }
554            }
555        };
556
557        Poll::Pending
558    }
559
560    #[await_tree::instrument("control_stream_next_event")]
561    pub(super) async fn next_event<'a>(
562        &'a mut self,
563        term_id: &str,
564        context: &Arc<impl GlobalBarrierWorkerContext>,
565    ) -> (WorkerId, WorkerNodeEvent<'a>) {
566        let mut this = Some(self);
567        poll_fn(|cx| Self::poll_next_event(&mut this, cx, term_id, context, true)).await
568    }
569
570    #[await_tree::instrument("control_stream_next_response")]
571    pub(super) async fn next_response(
572        &mut self,
573        term_id: &str,
574        context: &Arc<impl GlobalBarrierWorkerContext>,
575    ) -> (
576        WorkerId,
577        MetaResult<streaming_control_stream_response::Response>,
578    ) {
579        let mut this = Some(self);
580        let (worker_id, event) =
581            poll_fn(|cx| Self::poll_next_event(&mut this, cx, term_id, context, false)).await;
582        match event {
583            WorkerNodeEvent::Response(result) => (worker_id, result),
584            WorkerNodeEvent::Connected(_) => {
585                unreachable!("set poll_reconnect=false")
586            }
587        }
588    }
589}
590
591pub(super) struct DatabaseInitialBarrierCollector {
592    pub(super) database_id: DatabaseId,
593    pub(super) initializing_partial_graphs: HashSet<PartialGraphId>,
594    pub(super) database: DatabaseCheckpointControl,
595}
596
597impl Debug for DatabaseInitialBarrierCollector {
598    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
599        f.debug_struct("DatabaseInitialBarrierCollector")
600            .field("database_id", &self.database_id)
601            .field("initializing_graphs", &self.initializing_partial_graphs)
602            .finish()
603    }
604}
605
606impl DatabaseInitialBarrierCollector {
607    pub(super) fn is_collected(&self) -> bool {
608        self.initializing_partial_graphs.is_empty()
609    }
610
611    pub(super) fn partial_graph_initialized(&mut self, partial_graph_id: PartialGraphId) {
612        assert!(self.initializing_partial_graphs.remove(&partial_graph_id));
613    }
614
615    pub(super) fn all_partial_graphs(&self) -> impl Iterator<Item = PartialGraphId> + '_ {
616        database_partial_graphs(
617            self.database_id,
618            self.database
619                .independent_checkpoint_job_controls
620                .keys()
621                .copied(),
622        )
623    }
624
625    pub(super) fn finish(self) -> DatabaseCheckpointControl {
626        assert!(self.is_collected());
627        self.database
628    }
629
630    pub(super) fn is_valid_after_worker_err(&self, worker_id: WorkerId) -> bool {
631        self.database.is_valid_after_worker_err(worker_id)
632    }
633}
634
635impl PartialGraphRecoverer<'_> {
636    /// Extract information from the loaded runtime barrier worker snapshot info, and inject the initial barrier.
637    #[expect(clippy::too_many_arguments)]
638    pub(super) fn inject_database_initial_barrier(
639        &mut self,
640        database_id: DatabaseId,
641        jobs: HashMap<JobId, HashMap<FragmentId, InflightFragmentInfo>>,
642        job_extra_info: &HashMap<JobId, StreamingJobExtraInfo>,
643        state_table_committed_epochs: &mut HashMap<TableId, u64>,
644        state_table_log_epochs: &mut HashMap<TableId, Vec<(Vec<u64>, u64)>>,
645        fragment_relations: &FragmentDownstreamRelation,
646        stream_actors: &HashMap<ActorId, StreamActor>,
647        source_splits: &mut HashMap<ActorId, Vec<SplitImpl>>,
648        creating_jobs: &mut HashSet<JobId>,
649        mv_depended_subscriptions: &mut HashMap<TableId, HashMap<SubscriptionId, u64>>,
650        is_paused: bool,
651        hummock_version_stats: &HummockVersionStats,
652        cdc_table_snapshot_splits: &mut HashMap<JobId, CdcTableSnapshotSplits>,
653        batch_refresh: HashMap<JobId, BatchRefreshRenderResult>,
654    ) -> MetaResult<DatabaseCheckpointControl> {
655        fn collect_source_splits(
656            fragment_infos: impl Iterator<Item = &InflightFragmentInfo>,
657            source_splits: &mut HashMap<ActorId, Vec<SplitImpl>>,
658        ) -> HashMap<ActorId, Vec<SplitImpl>> {
659            fragment_infos
660                .flat_map(|info| info.actors.keys())
661                .filter_map(|actor_id| {
662                    let actor_id = *actor_id as ActorId;
663                    source_splits
664                        .remove(&actor_id)
665                        .map(|splits| (actor_id, splits))
666                })
667                .collect()
668        }
669        fn build_mutation(
670            splits: &HashMap<ActorId, Vec<SplitImpl>>,
671            cdc_table_snapshot_split_assignment: HashMap<ActorId, PbCdcTableSnapshotSplits>,
672            backfill_orders: &ExtendedFragmentBackfillOrder,
673            is_paused: bool,
674        ) -> Mutation {
675            let backfill_nodes_to_pause = get_nodes_with_backfill_dependencies(backfill_orders)
676                .into_iter()
677                .collect();
678            Mutation::Add(AddMutation {
679                // Actors built during recovery is not treated as newly added actors.
680                actor_dispatchers: Default::default(),
681                added_actors: Default::default(),
682                actor_splits: build_actor_connector_splits(splits),
683                actor_cdc_table_snapshot_splits: Some(PbCdcTableSnapshotSplitsWithGeneration {
684                    splits: cdc_table_snapshot_split_assignment,
685                }),
686                pause: is_paused,
687                subscriptions_to_add: Default::default(),
688                backfill_nodes_to_pause,
689                new_upstream_sinks: Default::default(),
690                dropped_actors: Default::default(),
691                sink_log_store_flush: Default::default(),
692            })
693        }
694
695        fn resolve_jobs_committed_epoch(
696            state_table_committed_epochs: &mut HashMap<TableId, u64>,
697            table_ids: impl Iterator<Item = TableId>,
698        ) -> u64 {
699            let mut epochs = table_ids.map(|table_id| {
700                (
701                    table_id,
702                    state_table_committed_epochs
703                        .remove(&table_id)
704                        .expect("should exist"),
705                )
706            });
707            let (first_table_id, prev_epoch) = epochs.next().expect("non-empty");
708            for (table_id, epoch) in epochs {
709                assert_eq!(
710                    prev_epoch, epoch,
711                    "{} has different committed epoch to {}",
712                    first_table_id, table_id
713                );
714            }
715            prev_epoch
716        }
717        fn job_backfill_orders(
718            job_extra_info: &HashMap<JobId, StreamingJobExtraInfo>,
719            job_id: JobId,
720        ) -> UserDefinedFragmentBackfillOrder {
721            UserDefinedFragmentBackfillOrder::new(
722                job_extra_info
723                    .get(&job_id)
724                    .and_then(|info| info.backfill_orders.clone())
725                    .map_or_else(HashMap::new, |orders| orders.0),
726            )
727        }
728
729        let mut subscribers: HashMap<_, HashMap<_, _>> = jobs
730            .keys()
731            .filter_map(|job_id| {
732                mv_depended_subscriptions
733                    .remove(&job_id.as_mv_table_id())
734                    .map(|subscriptions| {
735                        (
736                            job_id.as_mv_table_id(),
737                            subscriptions
738                                .into_iter()
739                                .map(|(subscription_id, retention)| {
740                                    (
741                                        subscription_id.as_subscriber_id(),
742                                        SubscriberType::Subscription(retention),
743                                    )
744                                })
745                                .collect(),
746                        )
747                    })
748            })
749            .collect();
750
751        // Batch-refresh jobs are rendered outside `jobs`, but their upstream tables
752        // must still start with log-store-enabled subscribers after recovery.
753        for (job_id, render_result) in &batch_refresh {
754            let snapshot_backfill_info = StreamFragmentGraph::collect_snapshot_backfill_info_impl(
755                render_result
756                    .fragment_infos
757                    .values()
758                    .map(|fragment| (&fragment.nodes, fragment.fragment_type_mask)),
759            )?
760            .0
761            .ok_or_else(|| anyhow!("batch refresh job {} has no snapshot backfill info", job_id))?;
762
763            for upstream_table_id in snapshot_backfill_info
764                .upstream_mv_table_id_to_backfill_epoch
765                .keys()
766            {
767                subscribers
768                    .entry(*upstream_table_id)
769                    .or_default()
770                    .try_insert(job_id.as_subscriber_id(), SubscriberType::SnapshotBackfill)
771                    .expect("non-duplicate");
772            }
773        }
774
775        let mut database_jobs = HashMap::new();
776        let mut snapshot_backfill_jobs = HashMap::new();
777
778        for (job_id, job_fragments) in jobs {
779            if creating_jobs.remove(&job_id) {
780                if job_fragments.values().any(|fragment| {
781                    fragment
782                        .fragment_type_mask
783                        .contains(FragmentTypeFlag::SnapshotBackfillStreamScan)
784                }) {
785                    debug!(%job_id, "recovered snapshot backfill job");
786                    snapshot_backfill_jobs.insert(job_id, job_fragments);
787                } else {
788                    database_jobs.insert(job_id, (job_fragments, true));
789                }
790            } else {
791                database_jobs.insert(job_id, (job_fragments, false));
792            }
793        }
794
795        let database_job_log_epochs: HashMap<_, _> = database_jobs
796            .keys()
797            .filter_map(|job_id| {
798                state_table_log_epochs
799                    .remove(&job_id.as_mv_table_id())
800                    .map(|epochs| (job_id.as_mv_table_id(), epochs))
801            })
802            .collect();
803
804        let prev_epoch = resolve_jobs_committed_epoch(
805            state_table_committed_epochs,
806            InflightFragmentInfo::existing_table_ids(
807                database_jobs.values().flat_map(|(job, _)| job.values()),
808            ),
809        );
810        let prev_epoch = TracedEpoch::new(Epoch(prev_epoch));
811        // Use a different `curr_epoch` for each recovery attempt.
812        let curr_epoch = prev_epoch.next();
813        let barrier_info = BarrierInfo {
814            prev_epoch,
815            curr_epoch,
816            kind: BarrierKind::Initial,
817        };
818
819        let mut ongoing_snapshot_backfill_jobs: HashMap<JobId, _> = HashMap::new();
820        for (job_id, fragment_infos) in snapshot_backfill_jobs {
821            let committed_epoch = resolve_jobs_committed_epoch(
822                state_table_committed_epochs,
823                InflightFragmentInfo::existing_table_ids(fragment_infos.values()),
824            );
825            if committed_epoch == barrier_info.prev_epoch() {
826                info!(
827                    "recovered creating snapshot backfill job {} catch up with upstream already",
828                    job_id
829                );
830                database_jobs
831                    .try_insert(job_id, (fragment_infos, true))
832                    .expect("non-duplicate");
833                continue;
834            }
835            let snapshot_backfill_info = StreamFragmentGraph::collect_snapshot_backfill_info_impl(
836                fragment_infos
837                    .values()
838                    .map(|fragment| (&fragment.nodes, fragment.fragment_type_mask)),
839            )?
840            .0
841            .ok_or_else(|| {
842                anyhow!(
843                    "recovered snapshot backfill job {} has no snapshot backfill info",
844                    job_id
845                )
846            })?;
847            let mut snapshot_epoch = None;
848            let upstream_table_ids: HashSet<_> = snapshot_backfill_info
849                .upstream_mv_table_id_to_backfill_epoch
850                .keys()
851                .cloned()
852                .collect();
853            for (upstream_table_id, epoch) in
854                snapshot_backfill_info.upstream_mv_table_id_to_backfill_epoch
855            {
856                let epoch = epoch.ok_or_else(|| anyhow!("recovered snapshot backfill job {} to upstream {} has not set snapshot epoch", job_id, upstream_table_id))?;
857                let snapshot_epoch = snapshot_epoch.get_or_insert(epoch);
858                if *snapshot_epoch != epoch {
859                    return Err(anyhow!("snapshot epoch {} to upstream {} different to snapshot epoch {} to previous upstream", epoch, upstream_table_id, snapshot_epoch).into());
860                }
861            }
862            let snapshot_epoch = snapshot_epoch.ok_or_else(|| {
863                anyhow!(
864                    "snapshot backfill job {} has not set snapshot epoch",
865                    job_id
866                )
867            })?;
868            for upstream_table_id in &upstream_table_ids {
869                subscribers
870                    .entry(*upstream_table_id)
871                    .or_default()
872                    .try_insert(job_id.as_subscriber_id(), SubscriberType::SnapshotBackfill)
873                    .expect("non-duplicate");
874            }
875            ongoing_snapshot_backfill_jobs
876                .try_insert(
877                    job_id,
878                    (
879                        fragment_infos,
880                        upstream_table_ids,
881                        committed_epoch,
882                        snapshot_epoch,
883                    ),
884                )
885                .expect("non-duplicated");
886        }
887
888        let mut cdc_table_snapshot_split_assignment: HashMap<ActorId, PbCdcTableSnapshotSplits> =
889            HashMap::new();
890
891        let database_jobs: HashMap<JobId, InflightStreamingJobInfo> = {
892            database_jobs
893                .into_iter()
894                .map(|(job_id, (fragment_infos, is_creating))| {
895                    let status = if is_creating {
896                        let backfill_ordering = job_backfill_orders(job_extra_info, job_id);
897                        let backfill_ordering = StreamFragmentGraph::extend_fragment_backfill_ordering_with_locality_backfill(
898                            backfill_ordering,
899                            fragment_relations,
900                            || fragment_infos.iter().map(|(fragment_id, fragment)| {
901                            (*fragment_id, fragment.fragment_type_mask, &fragment.nodes)
902                        }));
903                        let locality_fragment_state_table_mapping =
904                            build_locality_fragment_state_table_mapping(&fragment_infos);
905                        let backfill_order_state = BackfillOrderState::recover_from_fragment_infos(
906                            &backfill_ordering,
907                            &fragment_infos,
908                            locality_fragment_state_table_mapping,
909                        );
910                        CreateStreamingJobStatus::Creating {
911                            tracker: CreateMviewProgressTracker::recover(
912                                job_id,
913                                &fragment_infos,
914                                backfill_order_state,
915                                hummock_version_stats,
916                            ),
917                        }
918                    } else {
919                        CreateStreamingJobStatus::Created
920                    };
921                    let cdc_table_backfill_tracker =
922                        if let Some(splits) = cdc_table_snapshot_splits.remove(&job_id) {
923                            let cdc_fragment = fragment_infos
924                                .values()
925                                .find(|fragment| {
926                                    is_parallelized_backfill_enabled_cdc_scan_fragment(
927                                        fragment.fragment_type_mask,
928                                        &fragment.nodes,
929                                    )
930                                    .is_some()
931                                })
932                                .expect("should have parallel cdc fragment");
933                            let cdc_actors = cdc_fragment.actors.keys().copied().collect();
934                            let mut tracker =
935                                CdcTableBackfillTracker::restore(cdc_fragment.fragment_id, splits);
936                            cdc_table_snapshot_split_assignment
937                                .extend(tracker.reassign_splits(cdc_actors)?);
938                            Some(tracker)
939                        } else {
940                            None
941                        };
942                    Ok((
943                        job_id,
944                        InflightStreamingJobInfo {
945                            job_id,
946                            fragment_infos,
947                            subscribers: subscribers
948                                .remove(&job_id.as_mv_table_id())
949                                .unwrap_or_default(),
950                            status,
951                            cdc_table_backfill_tracker,
952                        },
953                    ))
954                })
955                .try_collect::<_, _, MetaError>()
956        }?;
957
958        let control_stream_manager = self.control_stream_manager();
959        let mut builder = FragmentEdgeBuilder::new(
960            database_jobs
961                .values()
962                .flat_map(|job| {
963                    let partial_graph_id = to_partial_graph_id(database_id, None);
964                    job.fragment_infos().map(move |info| {
965                        (
966                            info.fragment_id,
967                            EdgeBuilderFragmentInfo::from_inflight(
968                                info,
969                                partial_graph_id,
970                                control_stream_manager,
971                            ),
972                        )
973                    })
974                })
975                .chain(ongoing_snapshot_backfill_jobs.iter().flat_map(
976                    |(job_id, (fragments, ..))| {
977                        let partial_graph_id = to_partial_graph_id(database_id, Some(*job_id));
978                        fragments.values().map(move |fragment| {
979                            (
980                                fragment.fragment_id,
981                                EdgeBuilderFragmentInfo::from_inflight(
982                                    fragment,
983                                    partial_graph_id,
984                                    control_stream_manager,
985                                ),
986                            )
987                        })
988                    },
989                )),
990        );
991        builder.add_relations(fragment_relations);
992        let mut edges = builder.build();
993
994        {
995            let new_actors =
996                edges.collect_actors_to_create(database_jobs.values().flat_map(move |job| {
997                    job.fragment_infos.values().map(move |fragment_infos| {
998                        (
999                            fragment_infos.fragment_id,
1000                            &fragment_infos.nodes,
1001                            fragment_infos.actors.iter().map(move |(actor_id, actor)| {
1002                                (
1003                                    stream_actors.get(actor_id).expect("should exist"),
1004                                    actor.worker_id,
1005                                )
1006                            }),
1007                            job.subscribers.keys().copied(),
1008                        )
1009                    })
1010                }));
1011
1012            let nodes_actors =
1013                InflightFragmentInfo::actor_ids_to_collect(database_jobs.values().flatten());
1014            let database_job_source_splits =
1015                collect_source_splits(database_jobs.values().flatten(), source_splits);
1016            let database_backfill_orders =
1017                UserDefinedFragmentBackfillOrder::merge(database_jobs.values().map(|job| {
1018                    if matches!(job.status, CreateStreamingJobStatus::Creating { .. }) {
1019                        job_backfill_orders(job_extra_info, job.job_id)
1020                    } else {
1021                        UserDefinedFragmentBackfillOrder::default()
1022                    }
1023                }));
1024            let database_backfill_orders =
1025                StreamFragmentGraph::extend_fragment_backfill_ordering_with_locality_backfill(
1026                    database_backfill_orders,
1027                    fragment_relations,
1028                    || {
1029                        database_jobs.values().flat_map(|job_fragments| {
1030                            job_fragments
1031                                .fragment_infos
1032                                .iter()
1033                                .map(|(fragment_id, fragment)| {
1034                                    (*fragment_id, fragment.fragment_type_mask, &fragment.nodes)
1035                                })
1036                        })
1037                    },
1038                );
1039            let mutation = build_mutation(
1040                &database_job_source_splits,
1041                cdc_table_snapshot_split_assignment,
1042                &database_backfill_orders,
1043                is_paused,
1044            );
1045
1046            let partial_graph_id = to_partial_graph_id(database_id, None);
1047            self.recover_graph(
1048                partial_graph_id,
1049                mutation,
1050                &barrier_info,
1051                &nodes_actors,
1052                InflightFragmentInfo::existing_table_ids(database_jobs.values().flatten()),
1053                new_actors,
1054                DatabaseCheckpointControlMetrics::new(database_id),
1055            )?;
1056            debug!(
1057                %database_id,
1058                "inject initial barrier"
1059            );
1060        };
1061
1062        let mut independent_checkpoint_job_controls: HashMap<
1063            JobId,
1064            IndependentCheckpointJobControl,
1065        > = HashMap::new();
1066        for (job_id, (info, upstream_table_ids, committed_epoch, snapshot_epoch)) in
1067            ongoing_snapshot_backfill_jobs
1068        {
1069            let node_actors = edges.collect_actors_to_create(info.values().map(|fragment_infos| {
1070                (
1071                    fragment_infos.fragment_id,
1072                    &fragment_infos.nodes,
1073                    fragment_infos.actors.iter().map(move |(actor_id, actor)| {
1074                        (
1075                            stream_actors.get(actor_id).expect("should exist"),
1076                            actor.worker_id,
1077                        )
1078                    }),
1079                    vec![], // no subscribers for backfilling jobs,
1080                )
1081            }));
1082
1083            let database_job_source_splits =
1084                collect_source_splits(database_jobs.values().flatten(), source_splits);
1085            assert!(
1086                !cdc_table_snapshot_splits.contains_key(&job_id),
1087                "snapshot backfill job {job_id} should not have cdc backfill"
1088            );
1089            if is_paused {
1090                bail!("should not pause when having snapshot backfill job {job_id}");
1091            }
1092            let job_backfill_orders = job_backfill_orders(job_extra_info, job_id);
1093            let job_backfill_orders =
1094                StreamFragmentGraph::extend_fragment_backfill_ordering_with_locality_backfill(
1095                    job_backfill_orders,
1096                    fragment_relations,
1097                    || {
1098                        info.iter().map(|(fragment_id, fragment)| {
1099                            (*fragment_id, fragment.fragment_type_mask, &fragment.nodes)
1100                        })
1101                    },
1102                );
1103            let mutation = build_mutation(
1104                &database_job_source_splits,
1105                Default::default(), // no cdc backfill job for
1106                &job_backfill_orders,
1107                false,
1108            );
1109
1110            let job = CreatingStreamingJobControl::recover(
1111                database_id,
1112                job_id,
1113                upstream_table_ids,
1114                &database_job_log_epochs,
1115                snapshot_epoch,
1116                committed_epoch,
1117                &barrier_info,
1118                info,
1119                job_backfill_orders,
1120                fragment_relations,
1121                hummock_version_stats,
1122                node_actors,
1123                mutation.clone(),
1124                self,
1125            )?;
1126            independent_checkpoint_job_controls.insert(
1127                job_id,
1128                IndependentCheckpointJobControl::CreatingStreamingJob(job),
1129            );
1130        }
1131
1132        // Recover batch refresh jobs (both idle and consuming snapshot).
1133        // Actors were already rendered by `render_runtime_info()`.
1134        for (job_id, render_result) in batch_refresh {
1135            creating_jobs.remove(&job_id);
1136            debug!(%job_id, "recovered batch refresh job");
1137
1138            // Resolve committed epoch from state tables.
1139            let committed_epoch = resolve_jobs_committed_epoch(
1140                state_table_committed_epochs,
1141                InflightFragmentInfo::existing_table_ids(render_result.fragment_infos.values()),
1142            );
1143
1144            let snapshot_backfill_info = StreamFragmentGraph::collect_snapshot_backfill_info_impl(
1145                render_result
1146                    .fragment_infos
1147                    .values()
1148                    .map(|fragment| (&fragment.nodes, fragment.fragment_type_mask)),
1149            )?
1150            .0
1151            .ok_or_else(|| anyhow!("batch refresh job {} has no snapshot backfill info", job_id))?;
1152
1153            let upstream_table_ids: HashSet<TableId> = snapshot_backfill_info
1154                .upstream_mv_table_id_to_backfill_epoch
1155                .keys()
1156                .copied()
1157                .collect();
1158            let snapshot_epoch = snapshot_backfill_info
1159                .upstream_mv_table_id_to_backfill_epoch
1160                .values()
1161                .find_map(|e| *e)
1162                .unwrap_or(committed_epoch);
1163
1164            let job_backfill_orders = job_backfill_orders(job_extra_info, job_id);
1165            let job_backfill_orders =
1166                StreamFragmentGraph::extend_fragment_backfill_ordering_with_locality_backfill(
1167                    job_backfill_orders,
1168                    fragment_relations,
1169                    || {
1170                        render_result
1171                            .fragment_infos
1172                            .iter()
1173                            .map(|(fid, f)| (*fid, f.fragment_type_mask, &f.nodes))
1174                    },
1175                );
1176            let mutation = build_mutation(
1177                &Default::default(), // batch refresh has no source splits
1178                Default::default(),
1179                &job_backfill_orders,
1180                false,
1181            );
1182
1183            let refresh_interval_sec = job_extra_info
1184                .get(&job_id)
1185                .and_then(|info| info.refresh_interval_sec)
1186                .expect("batch refresh job should have refresh_interval_sec in job extra info");
1187
1188            let job = BatchRefreshJobCheckpointControl::recover(
1189                database_id,
1190                job_id,
1191                upstream_table_ids,
1192                snapshot_epoch,
1193                committed_epoch,
1194                job_backfill_orders,
1195                hummock_version_stats,
1196                mutation,
1197                render_result,
1198                self,
1199                refresh_interval_sec,
1200            )?;
1201            independent_checkpoint_job_controls
1202                .insert(job_id, IndependentCheckpointJobControl::BatchRefresh(job));
1203        }
1204
1205        self.control_stream_manager()
1206            .env
1207            .shared_actor_infos()
1208            .recover_database(
1209                database_id,
1210                database_jobs
1211                    .values()
1212                    .flat_map(|info| {
1213                        info.fragment_infos()
1214                            .map(move |fragment| (fragment, info.job_id))
1215                    })
1216                    .chain(
1217                        independent_checkpoint_job_controls
1218                            .iter()
1219                            .flat_map(|(job_id, job)| {
1220                                let job_id = *job_id;
1221                                job.fragment_infos()
1222                                    .into_iter()
1223                                    .flat_map(move |infos| infos.values().map(move |f| (f, job_id)))
1224                            }),
1225                    ),
1226            );
1227
1228        let committed_epoch = barrier_info.prev_epoch();
1229        let new_epoch = barrier_info.curr_epoch;
1230        let database_info = InflightDatabaseInfo::recover(
1231            database_id,
1232            database_jobs.into_values(),
1233            self.control_stream_manager()
1234                .env
1235                .shared_actor_infos()
1236                .clone(),
1237        );
1238        let database_state = BarrierWorkerState::recovery(new_epoch, is_paused);
1239        Ok(DatabaseCheckpointControl::recovery(
1240            database_id,
1241            database_state,
1242            committed_epoch,
1243            database_info,
1244            independent_checkpoint_job_controls,
1245        ))
1246    }
1247}
1248
1249impl ControlStreamManager {
1250    fn connected_workers(&self) -> impl Iterator<Item = (WorkerId, &ControlStreamNode)> + '_ {
1251        self.workers
1252            .iter()
1253            .filter_map(|(worker_id, (_, worker_state))| match worker_state {
1254                WorkerNodeState::Connected { control_stream, .. } => {
1255                    Some((*worker_id, control_stream))
1256                }
1257                WorkerNodeState::Reconnecting(_) => None,
1258            })
1259    }
1260
1261    pub(super) fn inject_barrier(
1262        &mut self,
1263        partial_graph_id: PartialGraphId,
1264        mutation: Option<Mutation>,
1265        iceberg_pk_index_compaction: Option<IcebergPkIndexCompactionContext>,
1266        barrier_info: &BarrierInfo,
1267        node_actors: &HashMap<WorkerId, HashSet<ActorId>>,
1268        table_ids_to_sync: impl Iterator<Item = TableId>,
1269        nodes_to_sync_table: impl Iterator<Item = WorkerId>,
1270        mut new_actors: Option<StreamJobActorsToCreate>,
1271    ) -> MetaResult<NodeToCollect> {
1272        fail_point!("inject_barrier_err", |_| risingwave_common::bail!(
1273            "inject_barrier_err"
1274        ));
1275
1276        let nodes_to_sync_table: HashSet<_> = nodes_to_sync_table.collect();
1277
1278        nodes_to_sync_table.iter().for_each(|worker_id| {
1279            assert!(node_actors.contains_key(worker_id), "worker_id {worker_id} in nodes_to_sync_table {nodes_to_sync_table:?} but not in node_actors {node_actors:?}");
1280        });
1281
1282        let mut node_need_collect = NodeToCollect::new();
1283        let table_ids_to_sync = table_ids_to_sync.collect_vec();
1284
1285        node_actors.iter()
1286            .try_for_each(|(worker_id, actor_ids_to_collect)| {
1287                assert!(!actor_ids_to_collect.is_empty(), "empty actor_ids_to_collect on worker {worker_id} in node_actors {node_actors:?}");
1288                let table_ids_to_sync = if nodes_to_sync_table.contains(worker_id) {
1289                    table_ids_to_sync.clone()
1290                } else {
1291                    vec![]
1292                };
1293
1294                let node = if let Some((_, worker_state)) = self.workers.get(worker_id)
1295                    &&
1296                    let WorkerNodeState::Connected { control_stream, .. } = worker_state
1297                {
1298                    control_stream
1299                } else {
1300                    return Err(anyhow!("unconnected worker node {}", worker_id).into());
1301                };
1302
1303                {
1304                    let mutation = mutation.clone();
1305                    let barrier = Barrier {
1306                        epoch: Some(risingwave_pb::data::Epoch {
1307                            curr: barrier_info.curr_epoch(),
1308                            prev: barrier_info.prev_epoch(),
1309                        }),
1310                        mutation: mutation.clone().map(|_| BarrierMutation { mutation }),
1311                        tracing_context: TracingContext::from_span(barrier_info.curr_epoch.span())
1312                            .to_protobuf(),
1313                        kind: barrier_info.kind.to_protobuf() as i32,
1314                        iceberg_pk_index_compaction,
1315                    };
1316
1317                    node.handle
1318                        .request_sender
1319                        .send(StreamingControlStreamRequest {
1320                            request: Some(
1321                                streaming_control_stream_request::Request::InjectBarrier(
1322                                    InjectBarrierRequest {
1323                                        request_id: Uuid::new_v4().to_string(),
1324                                        barrier: Some(barrier),
1325                                        actor_ids_to_collect: actor_ids_to_collect.iter().copied().collect(),
1326                                        table_ids_to_sync,
1327                                        partial_graph_id,
1328                                        actors_to_build: new_actors
1329                                            .as_mut()
1330                                            .map(|new_actors| new_actors.remove(worker_id))
1331                                            .into_iter()
1332                                            .flatten()
1333                                            .flatten()
1334                                            .map(|(fragment_id, (node, actors, initial_subscriber_ids))| {
1335                                                FragmentBuildActorInfo {
1336                                                    fragment_id,
1337                                                    node: Some(node),
1338                                                    actors: actors
1339                                                        .into_iter()
1340                                                        .map(|(actor, upstreams, dispatchers)| {
1341                                                            BuildActorInfo {
1342                                                                actor_id: actor.actor_id,
1343                                                                fragment_upstreams: upstreams
1344                                                                    .into_iter()
1345                                                                    .map(|(fragment_id, upstreams)| {
1346                                                                        (
1347                                                                            fragment_id,
1348                                                                            UpstreamActors {
1349                                                                                actors: upstreams
1350                                                                                    .into_values()
1351                                                                                    .collect(),
1352                                                                            },
1353                                                                        )
1354                                                                    })
1355                                                                    .collect(),
1356                                                                dispatchers,
1357                                                                vnode_bitmap: actor.vnode_bitmap.map(|bitmap| bitmap.to_protobuf()),
1358                                                                mview_definition: actor.mview_definition,
1359                                                                expr_context: actor.expr_context,
1360                                                                config_override: actor.config_override.to_string(),
1361                                                                initial_subscriber_ids: initial_subscriber_ids.iter().copied().collect(),
1362                                                            }
1363                                                        })
1364                                                        .collect(),
1365                                                }
1366                                            })
1367                                            .collect(),
1368                                    },
1369                                ),
1370                            ),
1371                        })
1372                        .map_err(|_| {
1373                            MetaError::from(anyhow!(
1374                                "failed to send request to {} {:?}",
1375                                node.worker_id,
1376                                node.host
1377                            ))
1378                        })?;
1379
1380                    node_need_collect.insert(*worker_id);
1381                    Result::<_, MetaError>::Ok(())
1382                }
1383            })
1384            .inspect_err(|e| {
1385                // Record failure in event log.
1386                use risingwave_pb::meta::event_log;
1387                let event = event_log::EventInjectBarrierFail {
1388                    prev_epoch: barrier_info.prev_epoch(),
1389                    cur_epoch: barrier_info.curr_epoch(),
1390                    error: e.to_report_string(),
1391                };
1392                self.env
1393                    .event_log_manager_ref()
1394                    .add_event_logs(vec![event_log::Event::InjectBarrierFail(event)]);
1395            })?;
1396        Ok(node_need_collect)
1397    }
1398
1399    pub(super) fn add_partial_graph(&mut self, partial_graph_id: PartialGraphId) {
1400        self.connected_workers().for_each(|(_, node)| {
1401            if node
1402                .handle
1403                .request_sender
1404                .send(StreamingControlStreamRequest {
1405                    request: Some(
1406                        streaming_control_stream_request::Request::CreatePartialGraph(
1407                            CreatePartialGraphRequest {
1408                                partial_graph_id,
1409                            },
1410                        ),
1411                    ),
1412                }).is_err() {
1413                let (database_id, creating_job_id) = from_partial_graph_id(partial_graph_id);
1414                warn!(%database_id, ?creating_job_id, worker_id = %node.worker_id, "failed to add the partial graph to the worker")
1415            }
1416        });
1417    }
1418
1419    pub(super) fn remove_partial_graphs(&mut self, partial_graph_ids: Vec<PartialGraphId>) {
1420        self.connected_workers().for_each(|(_, node)| {
1421            if node.handle
1422                .request_sender
1423                .send(StreamingControlStreamRequest {
1424                    request: Some(
1425                        streaming_control_stream_request::Request::RemovePartialGraph(
1426                            RemovePartialGraphRequest {
1427                                partial_graph_ids: partial_graph_ids.clone(),
1428                            },
1429                        ),
1430                    ),
1431                })
1432                .is_err()
1433            {
1434                warn!(worker_id = %node.worker_id,node = ?node.host,"failed to send remove partial graph request");
1435            }
1436        })
1437    }
1438
1439    pub(super) fn reset_partial_graphs(
1440        &mut self,
1441        partial_graph_ids: Vec<PartialGraphId>,
1442    ) -> HashSet<WorkerId> {
1443        self.connected_workers()
1444            .filter_map(|(worker_id, node)| {
1445                if node
1446                    .handle
1447                    .request_sender
1448                    .send(StreamingControlStreamRequest {
1449                        request: Some(
1450                            streaming_control_stream_request::Request::ResetPartialGraphs(
1451                                ResetPartialGraphsRequest {
1452                                    partial_graph_ids: partial_graph_ids.clone(),
1453                                },
1454                            ),
1455                        ),
1456                    })
1457                    .is_err()
1458                {
1459                    warn!(%worker_id, node = ?node.host,"failed to send reset database request");
1460                    None
1461                } else {
1462                    Some(worker_id)
1463                }
1464            })
1465            .collect()
1466    }
1467}
1468
1469impl GlobalBarrierWorkerContextImpl {
1470    pub(super) async fn new_control_stream_impl(
1471        &self,
1472        node: &WorkerNode,
1473        init_request: &PbInitRequest,
1474    ) -> MetaResult<StreamingControlHandle> {
1475        let handle = self
1476            .env
1477            .stream_client_pool()
1478            .get(node)
1479            .await?
1480            .start_streaming_control(init_request.clone())
1481            .await?;
1482        Ok(handle)
1483    }
1484}
1485
1486pub(super) fn merge_node_rpc_errors<E: Error + Send + Sync + 'static>(
1487    message: &str,
1488    errors: impl IntoIterator<Item = (WorkerId, E)>,
1489) -> MetaError {
1490    use std::fmt::Write;
1491
1492    use risingwave_common::error::error_request_copy;
1493    use risingwave_common::error::tonic::extra::Score;
1494
1495    let errors = errors.into_iter().collect_vec();
1496
1497    if errors.is_empty() {
1498        return anyhow!(message.to_owned()).into();
1499    }
1500
1501    // Create the error from the single error.
1502    let single_error = |(worker_id, e)| {
1503        anyhow::Error::from(e)
1504            .context(format!("{message}, in worker node {worker_id}"))
1505            .into()
1506    };
1507
1508    if errors.len() == 1 {
1509        return single_error(errors.into_iter().next().unwrap());
1510    }
1511
1512    // Find the error with the highest score.
1513    let max_score = errors
1514        .iter()
1515        .filter_map(|(_, e)| error_request_copy::<Score>(e))
1516        .max();
1517
1518    if let Some(max_score) = max_score {
1519        let mut errors = errors;
1520        let max_scored = errors
1521            .extract_if(.., |(_, e)| {
1522                error_request_copy::<Score>(e) == Some(max_score)
1523            })
1524            .next()
1525            .unwrap();
1526
1527        return single_error(max_scored);
1528    }
1529
1530    // The errors do not have scores, so simply concatenate them.
1531    let concat: String = errors
1532        .into_iter()
1533        .fold(format!("{message}: "), |mut s, (w, e)| {
1534            write!(&mut s, " in worker node {}, {};", w, e.as_report()).unwrap();
1535            s
1536        });
1537    anyhow!(concat).into()
1538}
1539
1540#[cfg(test)]
1541mod test_partial_graph_id {
1542    use crate::barrier::rpc::{from_partial_graph_id, to_partial_graph_id};
1543
1544    #[test]
1545    fn test_partial_graph_id_conversion() {
1546        let database_id = 233.into();
1547        let job_id = 233.into();
1548        assert_eq!(
1549            (database_id, None),
1550            from_partial_graph_id(to_partial_graph_id(database_id, None))
1551        );
1552        assert_eq!(
1553            (database_id, Some(job_id)),
1554            from_partial_graph_id(to_partial_graph_id(database_id, Some(job_id)))
1555        );
1556    }
1557}