Skip to main content

risingwave_stream/executor/
mod.rs

1// Copyright 2022 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15mod prelude;
16
17use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
18use std::fmt::Debug;
19use std::future::pending;
20use std::hash::Hash;
21use std::pin::Pin;
22use std::sync::Arc;
23use std::task::Poll;
24use std::vec;
25
26use await_tree::InstrumentAwait;
27use enum_as_inner::EnumAsInner;
28use futures::future::try_join_all;
29use futures::stream::{BoxStream, FusedStream, FuturesUnordered, StreamFuture};
30use futures::{FutureExt, Stream, StreamExt, TryStreamExt};
31use itertools::Itertools;
32use prometheus::core::{AtomicU64, GenericCounter};
33use risingwave_common::array::StreamChunk;
34use risingwave_common::bitmap::Bitmap;
35use risingwave_common::catalog::{Schema, TableId};
36use risingwave_common::config::StreamingConfig;
37use risingwave_common::metrics::LabelGuardedMetric;
38use risingwave_common::row::OwnedRow;
39use risingwave_common::types::{DataType, Datum, DefaultOrd, ScalarImpl};
40use risingwave_common::util::epoch::{Epoch, EpochPair};
41use risingwave_common::util::tracing::TracingContext;
42use risingwave_common::util::value_encoding::{DatumFromProtoExt, DatumToProtoExt};
43use risingwave_common_estimate_size::EstimateSize;
44use risingwave_connector::source::SplitImpl;
45use risingwave_expr::expr::NonStrictExpression;
46use risingwave_pb::common::ThrottleType;
47use risingwave_pb::data::PbEpoch;
48use risingwave_pb::expr::PbInputRef;
49use risingwave_pb::stream_plan::add_mutation::PbNewUpstreamSink;
50use risingwave_pb::stream_plan::barrier::BarrierKind;
51use risingwave_pb::stream_plan::barrier_mutation::Mutation as PbMutation;
52use risingwave_pb::stream_plan::stream_node::PbStreamKind;
53use risingwave_pb::stream_plan::throttle_mutation::ThrottleConfig;
54use risingwave_pb::stream_plan::update_mutation::{DispatcherUpdate, MergeUpdate};
55use risingwave_pb::stream_plan::{
56    IcebergPkIndexCompactionContext, PbBarrier, PbBarrierMutation, PbDispatcher,
57    PbSinkSchemaChange, PbStreamMessageBatch, PbWatermark, SubscriptionUpstreamInfo,
58};
59use smallvec::SmallVec;
60use tokio::sync::mpsc;
61use tokio::time::{Duration, Instant};
62
63use crate::error::StreamResult;
64use crate::executor::exchange::input::{
65    BoxedActorInput, BoxedInput, assert_equal_dispatcher_barrier, new_input,
66};
67use crate::executor::monitor::ActorInputMetrics;
68use crate::executor::prelude::StreamingMetrics;
69use crate::executor::watermark::BufferedWatermarks;
70use crate::task::{ActorId, FragmentId, LocalBarrierManager};
71
72mod actor;
73mod barrier_align;
74pub mod exchange;
75pub mod monitor;
76
77pub mod aggregate;
78pub mod asof_join;
79mod backfill;
80mod barrier_recv;
81mod batch_query;
82mod chain;
83mod changelog;
84mod dedup;
85mod dispatch;
86pub mod dml;
87mod dynamic_filter;
88pub mod eowc;
89pub mod error;
90mod expand;
91mod filter;
92mod gap_fill;
93pub mod hash_join;
94mod hop_window;
95pub(crate) mod iceberg_with_pk_index;
96mod join;
97pub mod locality_provider;
98mod lookup;
99mod lookup_union;
100pub mod match_recognize;
101mod merge;
102mod mview;
103mod nested_loop_temporal_join;
104mod no_op;
105mod now;
106mod over_window;
107pub mod project;
108mod receiver;
109pub mod row_id_gen;
110mod sink;
111pub mod source;
112mod stream_reader;
113pub mod subtask;
114mod temporal_join;
115mod top_n;
116mod troublemaker;
117mod union;
118mod upstream_sink_union;
119mod values;
120mod watermark;
121mod watermark_filter;
122mod wrapper;
123
124mod approx_percentile;
125
126mod row_merge;
127
128#[cfg(test)]
129mod integration_tests;
130mod sync_kv_log_store;
131#[cfg(any(test, feature = "test"))]
132pub mod test_utils;
133mod utils;
134mod vector;
135
136pub use actor::{Actor, ActorContext, ActorContextRef};
137use anyhow::{Context, anyhow};
138pub use approx_percentile::global::GlobalApproxPercentileExecutor;
139pub use approx_percentile::local::LocalApproxPercentileExecutor;
140pub use backfill::arrangement_backfill::*;
141pub use backfill::cdc::{
142    CdcBackfillExecutor, ExternalStorageTable, ParallelizedCdcBackfillExecutor,
143};
144pub use backfill::no_shuffle_backfill::*;
145pub use backfill::snapshot_backfill::*;
146pub use barrier_recv::BarrierRecvExecutor;
147pub use batch_query::BatchQueryExecutor;
148pub use chain::ChainExecutor;
149pub use changelog::{ChangeLogExecutor, ChangeLogMode};
150pub use dedup::AppendOnlyDedupExecutor;
151pub use dispatch::{DispatchExecutor, SyncLogStoreDispatchExecutor};
152pub use dynamic_filter::DynamicFilterExecutor;
153pub use error::{StreamExecutorError, StreamExecutorResult};
154pub use expand::ExpandExecutor;
155pub use filter::{FilterExecutor, UpsertFilterExecutor};
156pub use gap_fill::{GapFillExecutor, GapFillExecutorArgs};
157pub use hash_join::*;
158pub use hop_window::HopWindowExecutor;
159pub use iceberg_with_pk_index::{
160    CompactionResolverExecutor, IcebergWriterImpl, PositionDeleteHandlerImpl,
161    PositionDeleteMergerExecutor, WriterExecutor,
162};
163pub use join::asof_join::{AsOfCpuEncoding, AsOfMemoryEncoding};
164pub use join::row::{CachedJoinRow, CpuEncoding, JoinEncoding, MemoryEncoding};
165pub use join::{AsOfDesc, AsOfJoinType, JoinType};
166pub use lookup::*;
167pub use lookup_union::LookupUnionExecutor;
168pub use merge::MergeExecutor;
169pub(crate) use merge::{MergeExecutorInput, MergeExecutorUpstream};
170pub use mview::{MaterializeExecutor, RefreshableMaterializeArgs};
171pub use nested_loop_temporal_join::NestedLoopTemporalJoinExecutor;
172pub use no_op::NoOpExecutor;
173pub use now::*;
174pub use over_window::*;
175pub use receiver::ReceiverExecutor;
176use risingwave_common::id::SourceId;
177pub use row_merge::RowMergeExecutor;
178pub use sink::SinkExecutor;
179pub use sync_kv_log_store::SyncedKvLogStoreExecutor;
180pub use sync_kv_log_store::metrics::SyncedKvLogStoreMetrics;
181pub use temporal_join::TemporalJoinExecutor;
182pub use top_n::{
183    AppendOnlyGroupTopNExecutor, AppendOnlyTopNExecutor, GroupTopNExecutor, TopNExecutor,
184};
185pub use troublemaker::TroublemakerExecutor;
186pub use union::UnionExecutor;
187pub use upstream_sink_union::{UpstreamFragmentInfo, UpstreamSinkUnionExecutor};
188pub use utils::DummyExecutor;
189pub use values::ValuesExecutor;
190pub use vector::*;
191pub use watermark_filter::{UpsertWatermarkFilterExecutor, WatermarkFilterExecutor};
192pub use wrapper::WrapperExecutor;
193
194use self::barrier_align::AlignedMessageStream;
195
196pub type MessageStreamItemInner<M> = StreamExecutorResult<MessageInner<M>>;
197pub type MessageStreamItem = MessageStreamItemInner<BarrierMutationType>;
198pub type DispatcherMessageStreamItem = StreamExecutorResult<DispatcherMessage>;
199pub type BoxedMessageStream = BoxStream<'static, MessageStreamItem>;
200
201pub use risingwave_common::util::epoch::task_local::{curr_epoch, epoch, prev_epoch};
202use risingwave_connector::sink::catalog::SinkId;
203use risingwave_connector::source::cdc::{
204    CdcTableSnapshotSplitAssignmentWithGeneration,
205    build_actor_cdc_table_snapshot_splits_with_generation,
206};
207use risingwave_pb::id::{ExecutorId, SubscriberId};
208use risingwave_pb::stream_plan::stream_message_batch::{BarrierBatch, StreamMessageBatch};
209
210pub trait MessageStreamInner<M> = Stream<Item = MessageStreamItemInner<M>> + Send;
211pub trait MessageStream = Stream<Item = MessageStreamItem> + Send;
212pub trait DispatcherMessageStream = Stream<Item = DispatcherMessageStreamItem> + Send;
213
214/// Static information of an executor.
215#[derive(Debug, Default, Clone)]
216pub struct ExecutorInfo {
217    /// The schema of the OUTPUT of the executor.
218    pub schema: Schema,
219
220    /// The stream key indices of the OUTPUT of the executor.
221    pub stream_key: StreamKey,
222
223    /// The stream kind of the OUTPUT of the executor.
224    pub stream_kind: PbStreamKind,
225
226    /// Identity of the executor.
227    pub identity: String,
228
229    /// The executor id of the executor.
230    pub id: ExecutorId,
231}
232
233impl ExecutorInfo {
234    pub fn for_test(schema: Schema, stream_key: StreamKey, identity: String, id: u64) -> Self {
235        Self {
236            schema,
237            stream_key,
238            stream_kind: PbStreamKind::Retract, // dummy value for test
239            identity,
240            id: id.into(),
241        }
242    }
243}
244
245/// [`Execute`] describes the methods an executor should implement to handle control messages.
246pub trait Execute: Send + 'static {
247    fn execute(self: Box<Self>) -> BoxedMessageStream;
248
249    fn boxed(self) -> Box<dyn Execute>
250    where
251        Self: Sized + Send + 'static,
252    {
253        Box::new(self)
254    }
255}
256
257/// [`Executor`] combines the static information ([`ExecutorInfo`]) and the executable object to
258/// handle messages ([`Execute`]).
259pub struct Executor {
260    info: ExecutorInfo,
261    execute: Box<dyn Execute>,
262}
263
264impl Executor {
265    pub fn new(info: ExecutorInfo, execute: Box<dyn Execute>) -> Self {
266        Self { info, execute }
267    }
268
269    pub fn info(&self) -> &ExecutorInfo {
270        &self.info
271    }
272
273    pub fn schema(&self) -> &Schema {
274        &self.info.schema
275    }
276
277    pub fn stream_key(&self) -> StreamKeyRef<'_> {
278        &self.info.stream_key
279    }
280
281    pub fn stream_kind(&self) -> PbStreamKind {
282        self.info.stream_kind
283    }
284
285    pub fn identity(&self) -> &str {
286        &self.info.identity
287    }
288
289    pub fn execute(self) -> BoxedMessageStream {
290        self.execute.execute()
291    }
292}
293
294impl std::fmt::Debug for Executor {
295    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
296        f.write_str(self.identity())
297    }
298}
299
300impl From<(ExecutorInfo, Box<dyn Execute>)> for Executor {
301    fn from((info, execute): (ExecutorInfo, Box<dyn Execute>)) -> Self {
302        Self::new(info, execute)
303    }
304}
305
306impl<E> From<(ExecutorInfo, E)> for Executor
307where
308    E: Execute,
309{
310    fn from((info, execute): (ExecutorInfo, E)) -> Self {
311        Self::new(info, execute.boxed())
312    }
313}
314
315pub const INVALID_EPOCH: u64 = 0;
316
317type UpstreamFragmentId = FragmentId;
318type SplitAssignments = HashMap<ActorId, Vec<SplitImpl>>;
319
320#[derive(Debug, Clone, PartialEq)]
321#[cfg_attr(any(test, feature = "test"), derive(Default))]
322pub struct UpdateMutation {
323    pub dispatchers: HashMap<ActorId, Vec<DispatcherUpdate>>,
324    pub merges: HashMap<(ActorId, UpstreamFragmentId), MergeUpdate>,
325    pub vnode_bitmaps: HashMap<ActorId, Arc<Bitmap>>,
326    pub dropped_actors: HashSet<ActorId>,
327    pub actor_splits: SplitAssignments,
328    pub actor_new_dispatchers: HashMap<ActorId, Vec<PbDispatcher>>,
329    pub actor_cdc_table_snapshot_splits: CdcTableSnapshotSplitAssignmentWithGeneration,
330    pub sink_schema_change: HashMap<SinkId, PbSinkSchemaChange>,
331    pub subscriptions_to_drop: Vec<SubscriptionUpstreamInfo>,
332    pub iceberg_pk_index_compaction: Option<IcebergPkIndexCompactionContext>,
333}
334
335#[derive(Debug, Clone, PartialEq)]
336#[cfg_attr(any(test, feature = "test"), derive(Default))]
337pub struct AddMutation {
338    pub adds: HashMap<ActorId, Vec<PbDispatcher>>,
339    pub added_actors: HashSet<ActorId>,
340    pub dropped_actors: HashSet<ActorId>,
341    // TODO: remove this and use `SourceChangesSplit` after we support multiple mutations.
342    pub splits: SplitAssignments,
343    pub pause: bool,
344    /// (`upstream_mv_table_id`,  `subscriber_id`)
345    pub subscriptions_to_add: Vec<(TableId, SubscriberId)>,
346    /// nodes which should start backfill
347    pub backfill_nodes_to_pause: HashSet<FragmentId>,
348    pub actor_cdc_table_snapshot_splits: CdcTableSnapshotSplitAssignmentWithGeneration,
349    pub new_upstream_sinks: HashMap<FragmentId, PbNewUpstreamSink>,
350    pub sink_log_store_flush: HashSet<SinkId>,
351}
352
353#[derive(Debug, Clone, PartialEq)]
354#[cfg_attr(any(test, feature = "test"), derive(Default))]
355pub struct StopMutation {
356    pub dropped_actors: HashSet<ActorId>,
357    pub dropped_sink_fragments: HashSet<FragmentId>,
358}
359
360/// See [`PbMutation`] for the semantics of each mutation.
361#[derive(Debug, Clone, PartialEq)]
362pub enum Mutation {
363    Stop(StopMutation),
364    Update(UpdateMutation),
365    Add(AddMutation),
366    SourceChangeSplit(SplitAssignments),
367    Pause,
368    Resume,
369    Throttle(HashMap<FragmentId, ThrottleConfig>),
370    ConnectorPropsChange(HashMap<u32, HashMap<String, String>>),
371    DropSubscriptions {
372        /// `subscriber` -> `upstream_mv_table_id`
373        subscriptions_to_drop: Vec<SubscriptionUpstreamInfo>,
374    },
375    StartFragmentBackfill {
376        fragment_ids: HashSet<FragmentId>,
377    },
378    RefreshStart {
379        table_id: TableId,
380        associated_source_id: SourceId,
381    },
382    ListFinish {
383        associated_source_id: SourceId,
384    },
385    LoadFinish {
386        associated_source_id: SourceId,
387    },
388    ResetSource {
389        source_id: SourceId,
390    },
391    InjectSourceOffsets {
392        source_id: SourceId,
393        /// Split ID -> offset (JSON-encoded based on connector type)
394        split_offsets: HashMap<String, String>,
395    },
396}
397
398/// The generic type `M` is the mutation type of the barrier.
399///
400/// For barrier of in the dispatcher, `M` is `()`, which means the mutation is erased.
401/// For barrier flowing within the streaming actor, `M` is the normal `BarrierMutationType`.
402#[derive(Debug, Clone)]
403pub struct BarrierInner<M> {
404    pub epoch: EpochPair,
405    pub mutation: M,
406    pub kind: BarrierKind,
407
408    /// Tracing context for the **current** epoch of this barrier.
409    pub tracing_context: TracingContext,
410}
411
412pub type BarrierMutationType = Option<Arc<Mutation>>;
413pub type Barrier = BarrierInner<BarrierMutationType>;
414pub type DispatcherBarrier = BarrierInner<()>;
415
416impl<M: Default> BarrierInner<M> {
417    /// Create a plain barrier.
418    pub fn new_test_barrier(epoch: u64) -> Self {
419        Self {
420            epoch: EpochPair::new_test_epoch(epoch),
421            kind: BarrierKind::Checkpoint,
422            tracing_context: TracingContext::none(),
423            mutation: Default::default(),
424        }
425    }
426
427    pub fn with_prev_epoch_for_test(epoch: u64, prev_epoch: u64) -> Self {
428        Self {
429            epoch: EpochPair::new(epoch, prev_epoch),
430            kind: BarrierKind::Checkpoint,
431            tracing_context: TracingContext::none(),
432            mutation: Default::default(),
433        }
434    }
435}
436
437impl Barrier {
438    pub fn into_dispatcher(self) -> DispatcherBarrier {
439        DispatcherBarrier {
440            epoch: self.epoch,
441            mutation: (),
442            kind: self.kind,
443            tracing_context: self.tracing_context,
444        }
445    }
446
447    #[must_use]
448    pub fn with_mutation(self, mutation: Mutation) -> Self {
449        Self {
450            mutation: Some(Arc::new(mutation)),
451            ..self
452        }
453    }
454
455    #[cfg(any(test, feature = "test"))]
456    #[must_use]
457    pub fn with_iceberg_pk_index_compaction(self, update: IcebergPkIndexCompactionContext) -> Self {
458        self.with_mutation(Mutation::Update(UpdateMutation {
459            iceberg_pk_index_compaction: Some(update),
460            ..Default::default()
461        }))
462    }
463
464    pub fn iceberg_pk_index_compaction(&self) -> Option<&IcebergPkIndexCompactionContext> {
465        match self.mutation.as_deref() {
466            Some(Mutation::Update(update)) => update.iceberg_pk_index_compaction.as_ref(),
467            _ => None,
468        }
469    }
470
471    #[must_use]
472    pub fn with_stop(self) -> Self {
473        self.with_mutation(Mutation::Stop(StopMutation {
474            dropped_actors: Default::default(),
475            dropped_sink_fragments: Default::default(),
476        }))
477    }
478
479    /// Whether this barrier carries stop mutation.
480    pub fn is_with_stop_mutation(&self) -> bool {
481        matches!(self.mutation.as_deref(), Some(Mutation::Stop(_)))
482    }
483
484    /// Whether this barrier is to stop the actor with `actor_id`.
485    pub fn is_stop(&self, actor_id: ActorId) -> bool {
486        self.all_stop_actors()
487            .is_some_and(|actors| actors.contains(&actor_id))
488    }
489
490    pub fn is_checkpoint(&self) -> bool {
491        self.kind == BarrierKind::Checkpoint
492    }
493
494    /// Get the initial split assignments for the actor with `actor_id`.
495    ///
496    /// This should only be called on the initial barrier received by the executor. It must be
497    ///
498    /// - `Add` mutation when it's a new streaming job, or recovery.
499    /// - `Update` mutation when it's created for scaling.
500    ///
501    /// Note that `SourceChangeSplit` is **not** included, because it's only used for changing splits
502    /// of existing executors.
503    pub fn initial_split_assignment(&self, actor_id: ActorId) -> Option<&[SplitImpl]> {
504        match self.mutation.as_deref()? {
505            Mutation::Update(UpdateMutation { actor_splits, .. })
506            | Mutation::Add(AddMutation {
507                splits: actor_splits,
508                ..
509            }) => actor_splits.get(&actor_id),
510
511            _ => {
512                if cfg!(debug_assertions) {
513                    panic!(
514                        "the initial mutation of the barrier should not be {:?}",
515                        self.mutation
516                    );
517                }
518                None
519            }
520        }
521        .map(|s| s.as_slice())
522    }
523
524    /// Get all actors that to be stopped (dropped) by this barrier.
525    pub fn all_stop_actors(&self) -> Option<&HashSet<ActorId>> {
526        self.mutation.as_deref()?.all_stop_actors()
527    }
528
529    /// Whether this barrier is to newly add the actor with `actor_id`. This is used for `Chain` and
530    /// `Values` to decide whether to output the existing (historical) data.
531    ///
532    /// By "newly", we mean the actor belongs to a subgraph of a new streaming job. That is, actors
533    /// added for scaling are not included.
534    pub fn is_newly_added(&self, actor_id: ActorId) -> bool {
535        match self.mutation.as_deref() {
536            Some(Mutation::Add(AddMutation { added_actors, .. })) => {
537                added_actors.contains(&actor_id)
538            }
539            _ => false,
540        }
541    }
542
543    pub fn should_start_fragment_backfill(&self, fragment_id: FragmentId) -> bool {
544        if let Some(Mutation::StartFragmentBackfill { fragment_ids }) = self.mutation.as_deref() {
545            fragment_ids.contains(&fragment_id)
546        } else {
547            false
548        }
549    }
550
551    /// Whether this barrier adds new downstream fragment for the actor with `upstream_actor_id`.
552    ///
553    /// # Use case
554    /// Some optimizations are applied when an actor doesn't have any downstreams ("standalone" actors).
555    /// * Pause a standalone shared `SourceExecutor`.
556    /// * Disable a standalone `MaterializeExecutor`'s conflict check.
557    ///
558    /// This is implemented by checking `actor_context.initial_dispatch_num` on startup, and
559    /// check `has_more_downstream_fragments` on barrier to see whether the optimization
560    /// needs to be turned off.
561    ///
562    /// ## Some special cases not included
563    ///
564    /// Note that this is not `has_new_downstream_actor/fragment`. For our use case, we only
565    /// care about **number of downstream fragments** (more precisely, existence).
566    /// - When scaling, the number of downstream actors is changed, and they are "new", but downstream fragments is not changed.
567    /// - When `ALTER TABLE sink_into_table`, the fragment is replaced with a "new" one, but the number is not changed.
568    pub fn has_more_downstream_fragments(&self, upstream_actor_id: ActorId) -> bool {
569        let Some(mutation) = self.mutation.as_deref() else {
570            return false;
571        };
572        match mutation {
573            // Add is for mv, index and sink creation.
574            Mutation::Add(AddMutation { adds, .. }) => adds.get(&upstream_actor_id).is_some(),
575            Mutation::Update(_)
576            | Mutation::Stop(_)
577            | Mutation::Pause
578            | Mutation::Resume
579            | Mutation::SourceChangeSplit(_)
580            | Mutation::Throttle { .. }
581            | Mutation::DropSubscriptions { .. }
582            | Mutation::ConnectorPropsChange(_)
583            | Mutation::StartFragmentBackfill { .. }
584            | Mutation::RefreshStart { .. }
585            | Mutation::ListFinish { .. }
586            | Mutation::LoadFinish { .. }
587            | Mutation::ResetSource { .. }
588            | Mutation::InjectSourceOffsets { .. } => false,
589        }
590    }
591
592    /// Whether this barrier requires the executor to pause its data stream on startup.
593    pub fn is_pause_on_startup(&self) -> bool {
594        match self.mutation.as_deref() {
595            Some(Mutation::Add(AddMutation { pause, .. })) => *pause,
596            _ => false,
597        }
598    }
599
600    pub fn is_backfill_pause_on_startup(&self, backfill_fragment_id: FragmentId) -> bool {
601        match self.mutation.as_deref() {
602            Some(Mutation::Add(AddMutation {
603                backfill_nodes_to_pause,
604                ..
605            })) => backfill_nodes_to_pause.contains(&backfill_fragment_id),
606            Some(Mutation::Update(_)) => false,
607            _ => {
608                tracing::warn!(
609                    "expected an AddMutation or UpdateMutation on Startup, instead got {:?}",
610                    self
611                );
612                false
613            }
614        }
615    }
616
617    /// Whether this barrier is for resume.
618    pub fn is_resume(&self) -> bool {
619        matches!(self.mutation.as_deref(), Some(Mutation::Resume))
620    }
621
622    /// Returns the [`MergeUpdate`] if this barrier is to update the merge executors for the actor
623    /// with `actor_id`.
624    pub fn as_update_merge(
625        &self,
626        actor_id: ActorId,
627        upstream_fragment_id: UpstreamFragmentId,
628    ) -> Option<&MergeUpdate> {
629        self.mutation
630            .as_deref()
631            .and_then(|mutation| match mutation {
632                Mutation::Update(UpdateMutation { merges, .. }) => {
633                    merges.get(&(actor_id, upstream_fragment_id))
634                }
635                _ => None,
636            })
637    }
638
639    /// Returns the new upstream sink information if this barrier is to add a new upstream sink for
640    /// the specified downstream fragment.
641    pub fn as_new_upstream_sink(&self, fragment_id: FragmentId) -> Option<&PbNewUpstreamSink> {
642        self.mutation
643            .as_deref()
644            .and_then(|mutation| match mutation {
645                Mutation::Add(AddMutation {
646                    new_upstream_sinks, ..
647                }) => new_upstream_sinks.get(&fragment_id),
648                _ => None,
649            })
650    }
651
652    /// Returns the dropped upstream sink-fragment if this barrier is to drop any sink.
653    pub fn as_dropped_upstream_sinks(&self) -> Option<&HashSet<FragmentId>> {
654        self.mutation
655            .as_deref()
656            .and_then(|mutation| match mutation {
657                Mutation::Stop(StopMutation {
658                    dropped_sink_fragments,
659                    ..
660                }) => Some(dropped_sink_fragments),
661                _ => None,
662            })
663    }
664
665    /// Returns the new vnode bitmap if this barrier is to update the vnode bitmap for the actor
666    /// with `actor_id`.
667    ///
668    /// Actually, this vnode bitmap update is only useful for the record accessing validation for
669    /// distributed executors, since the read/write pattern will never be across multiple vnodes.
670    pub fn as_update_vnode_bitmap(&self, actor_id: ActorId) -> Option<Arc<Bitmap>> {
671        self.mutation
672            .as_deref()
673            .and_then(|mutation| match mutation {
674                Mutation::Update(UpdateMutation { vnode_bitmaps, .. }) => {
675                    vnode_bitmaps.get(&actor_id).cloned()
676                }
677                _ => None,
678            })
679    }
680
681    pub fn assume_no_update_vnode_bitmap(&self, actor_id: ActorId) -> StreamExecutorResult<()> {
682        if self.as_update_vnode_bitmap(actor_id).is_some() {
683            return Err(anyhow!("updating vnode bitmap in place is not supported").into());
684        }
685        Ok(())
686    }
687
688    pub fn as_sink_schema_change(&self, sink_id: SinkId) -> Option<PbSinkSchemaChange> {
689        self.mutation
690            .as_deref()
691            .and_then(|mutation| match mutation {
692                Mutation::Update(UpdateMutation {
693                    sink_schema_change, ..
694                }) => sink_schema_change.get(&sink_id).cloned(),
695                _ => None,
696            })
697    }
698
699    pub fn should_flush_sink_log_store(&self, sink_id: SinkId) -> bool {
700        self.mutation
701            .as_deref()
702            .is_some_and(|mutation| match mutation {
703                Mutation::Add(AddMutation {
704                    sink_log_store_flush,
705                    ..
706                }) => sink_log_store_flush.contains(&sink_id),
707                _ => false,
708            })
709    }
710
711    pub fn as_subscriptions_to_drop(&self) -> Option<&[SubscriptionUpstreamInfo]> {
712        match self.mutation.as_deref() {
713            Some(Mutation::DropSubscriptions {
714                subscriptions_to_drop,
715            })
716            | Some(Mutation::Update(UpdateMutation {
717                subscriptions_to_drop,
718                ..
719            })) => Some(subscriptions_to_drop.as_slice()),
720            _ => None,
721        }
722    }
723
724    pub fn get_curr_epoch(&self) -> Epoch {
725        Epoch(self.epoch.curr)
726    }
727
728    /// Retrieve the tracing context for the **current** epoch of this barrier.
729    pub fn tracing_context(&self) -> &TracingContext {
730        &self.tracing_context
731    }
732
733    pub fn added_subscriber_on_mv_table(
734        &self,
735        mv_table_id: TableId,
736    ) -> impl Iterator<Item = SubscriberId> + '_ {
737        if let Some(Mutation::Add(add)) = self.mutation.as_deref() {
738            Some(add)
739        } else {
740            None
741        }
742        .into_iter()
743        .flat_map(move |add| {
744            add.subscriptions_to_add.iter().filter_map(
745                move |(upstream_mv_table_id, subscriber_id)| {
746                    if *upstream_mv_table_id == mv_table_id {
747                        Some(*subscriber_id)
748                    } else {
749                        None
750                    }
751                },
752            )
753        })
754    }
755}
756
757impl<M: PartialEq> PartialEq for BarrierInner<M> {
758    fn eq(&self, other: &Self) -> bool {
759        self.epoch == other.epoch && self.mutation == other.mutation
760    }
761}
762
763impl Mutation {
764    /// Return the backfill throttle configuration for `fragment_id`.
765    pub fn backfill_throttle_config(&self, fragment_id: FragmentId) -> Option<&ThrottleConfig> {
766        let Mutation::Throttle(fragment_throttles) = self else {
767            return None;
768        };
769
770        fragment_throttles
771            .get(&fragment_id)
772            .filter(|config| config.throttle_type() == ThrottleType::Backfill)
773    }
774
775    /// Get all actors to be stopped (dropped) by this mutation.
776    pub fn all_stop_actors(&self) -> Option<&HashSet<ActorId>> {
777        match self {
778            Mutation::Stop(StopMutation { dropped_actors, .. })
779            | Mutation::Update(UpdateMutation { dropped_actors, .. })
780            | Mutation::Add(AddMutation { dropped_actors, .. }) => Some(dropped_actors),
781            _ => None,
782        }
783    }
784
785    /// Return true if the mutation stops the given actor.
786    pub fn is_stop(&self, actor_id: ActorId) -> bool {
787        self.all_stop_actors()
788            .is_some_and(|actors| actors.contains(&actor_id))
789    }
790
791    /// Return true if the mutation is stop.
792    ///
793    /// Note that this does not mean we will stop the current actor.
794    #[cfg(test)]
795    pub fn is_stop_mutation(&self) -> bool {
796        matches!(self, Mutation::Stop(_))
797    }
798
799    #[cfg(test)]
800    fn to_protobuf(&self) -> PbMutation {
801        use risingwave_pb::source::{
802            ConnectorSplit, ConnectorSplits, PbCdcTableSnapshotSplitsWithGeneration,
803        };
804        use risingwave_pb::stream_plan::connector_props_change_mutation::ConnectorPropsInfo;
805        use risingwave_pb::stream_plan::{
806            PbAddMutation, PbConnectorPropsChangeMutation, PbDispatchers,
807            PbDropSubscriptionsMutation, PbPauseMutation, PbResumeMutation,
808            PbSourceChangeSplitMutation, PbStartFragmentBackfillMutation, PbStopMutation,
809            PbThrottleMutation, PbUpdateMutation,
810        };
811        let actor_splits_to_protobuf = |actor_splits: &SplitAssignments| {
812            actor_splits
813                .iter()
814                .map(|(&actor_id, splits)| {
815                    (
816                        actor_id,
817                        ConnectorSplits {
818                            splits: splits.clone().iter().map(ConnectorSplit::from).collect(),
819                        },
820                    )
821                })
822                .collect::<HashMap<_, _>>()
823        };
824
825        match self {
826            Mutation::Stop(StopMutation {
827                dropped_actors,
828                dropped_sink_fragments,
829            }) => PbMutation::Stop(PbStopMutation {
830                actors: dropped_actors.iter().copied().collect(),
831                dropped_sink_fragments: dropped_sink_fragments.iter().copied().collect(),
832            }),
833            Mutation::Update(UpdateMutation {
834                dispatchers,
835                merges,
836                vnode_bitmaps,
837                dropped_actors,
838                actor_splits,
839                actor_new_dispatchers,
840                actor_cdc_table_snapshot_splits,
841                sink_schema_change,
842                subscriptions_to_drop,
843                iceberg_pk_index_compaction,
844            }) => PbMutation::Update(PbUpdateMutation {
845                dispatcher_update: dispatchers.values().flatten().cloned().collect(),
846                merge_update: merges.values().cloned().collect(),
847                actor_vnode_bitmap_update: vnode_bitmaps
848                    .iter()
849                    .map(|(&actor_id, bitmap)| (actor_id, bitmap.to_protobuf()))
850                    .collect(),
851                dropped_actors: dropped_actors.iter().copied().collect(),
852                actor_splits: actor_splits_to_protobuf(actor_splits),
853                actor_new_dispatchers: actor_new_dispatchers
854                    .iter()
855                    .map(|(&actor_id, dispatchers)| {
856                        (
857                            actor_id,
858                            PbDispatchers {
859                                dispatchers: dispatchers.clone(),
860                            },
861                        )
862                    })
863                    .collect(),
864                actor_cdc_table_snapshot_splits: Some(PbCdcTableSnapshotSplitsWithGeneration {
865                    splits:actor_cdc_table_snapshot_splits.splits.iter().map(|(actor_id,(splits, generation))| {
866                        (*actor_id, risingwave_pb::source::PbCdcTableSnapshotSplits {
867                            splits: splits.iter().map(risingwave_connector::source::cdc::build_cdc_table_snapshot_split).collect(),
868                            generation: *generation,
869                        })
870                    }).collect()
871                }),
872                sink_schema_change: sink_schema_change
873                    .iter()
874                    .map(|(sink_id, change)| ((*sink_id).as_raw_id(), change.clone()))
875                    .collect(),
876                subscriptions_to_drop: subscriptions_to_drop.clone(),
877                iceberg_pk_index_compaction: iceberg_pk_index_compaction.clone(),
878            }),
879            Mutation::Add(AddMutation {
880                adds,
881                added_actors,
882                dropped_actors,
883                splits,
884                pause,
885                subscriptions_to_add,
886                backfill_nodes_to_pause,
887                actor_cdc_table_snapshot_splits,
888                new_upstream_sinks,
889                sink_log_store_flush,
890            }) => PbMutation::Add(PbAddMutation {
891                actor_dispatchers: adds
892                    .iter()
893                    .map(|(&actor_id, dispatchers)| {
894                        (
895                            actor_id,
896                            PbDispatchers {
897                                dispatchers: dispatchers.clone(),
898                            },
899                        )
900                    })
901                    .collect(),
902                added_actors: added_actors.iter().copied().collect(),
903                actor_splits: actor_splits_to_protobuf(splits),
904                pause: *pause,
905                subscriptions_to_add: subscriptions_to_add
906                    .iter()
907                    .map(|(table_id, subscriber_id)| SubscriptionUpstreamInfo {
908                        subscriber_id: *subscriber_id,
909                        upstream_mv_table_id: *table_id,
910                    })
911                    .collect(),
912                backfill_nodes_to_pause: backfill_nodes_to_pause.iter().copied().collect(),
913                actor_cdc_table_snapshot_splits:
914                Some(PbCdcTableSnapshotSplitsWithGeneration {
915                    splits:actor_cdc_table_snapshot_splits.splits.iter().map(|(actor_id,(splits, generation))| {
916                        (*actor_id, risingwave_pb::source::PbCdcTableSnapshotSplits {
917                            splits: splits.iter().map(risingwave_connector::source::cdc::build_cdc_table_snapshot_split).collect(),
918                            generation: *generation,
919                        })
920                    }).collect()
921                }),
922                new_upstream_sinks: new_upstream_sinks
923                    .iter()
924                    .map(|(k, v)| (*k, v.clone()))
925                    .collect(),
926                dropped_actors: dropped_actors.iter().copied().collect(),
927                sink_log_store_flush: sink_log_store_flush.iter().copied().collect(),
928            }),
929            Mutation::SourceChangeSplit(changes) => {
930                PbMutation::Splits(PbSourceChangeSplitMutation {
931                    actor_splits: changes
932                        .iter()
933                        .map(|(&actor_id, splits)| {
934                            (
935                                actor_id,
936                                ConnectorSplits {
937                                    splits: splits
938                                        .clone()
939                                        .iter()
940                                        .map(ConnectorSplit::from)
941                                        .collect(),
942                                },
943                            )
944                        })
945                        .collect(),
946                })
947            }
948            Mutation::Pause => PbMutation::Pause(PbPauseMutation {}),
949            Mutation::Resume => PbMutation::Resume(PbResumeMutation {}),
950            Mutation::Throttle (changes) => PbMutation::Throttle(PbThrottleMutation {
951                fragment_throttle: changes.clone(),
952            }),
953            Mutation::DropSubscriptions {
954                subscriptions_to_drop,
955            } => PbMutation::DropSubscriptions(PbDropSubscriptionsMutation {
956                info: subscriptions_to_drop.clone(),
957            }),
958            Mutation::ConnectorPropsChange(map) => {
959                PbMutation::ConnectorPropsChange(PbConnectorPropsChangeMutation {
960                    connector_props_infos: map
961                        .iter()
962                        .map(|(actor_id, options)| {
963                            (
964                                *actor_id,
965                                ConnectorPropsInfo {
966                                    connector_props_info: options
967                                        .iter()
968                                        .map(|(k, v)| (k.clone(), v.clone()))
969                                        .collect(),
970                                },
971                            )
972                        })
973                        .collect(),
974                })
975            }
976            Mutation::StartFragmentBackfill { fragment_ids } => {
977                PbMutation::StartFragmentBackfill(PbStartFragmentBackfillMutation {
978                    fragment_ids: fragment_ids.iter().copied().collect(),
979                })
980            }
981            Mutation::RefreshStart {
982                table_id,
983                associated_source_id,
984            } => PbMutation::RefreshStart(risingwave_pb::stream_plan::RefreshStartMutation {
985                table_id: *table_id,
986                associated_source_id: *associated_source_id,
987            }),
988            Mutation::ListFinish {
989                associated_source_id,
990            } => PbMutation::ListFinish(risingwave_pb::stream_plan::ListFinishMutation {
991                associated_source_id: *associated_source_id,
992            }),
993            Mutation::LoadFinish {
994                associated_source_id,
995            } => PbMutation::LoadFinish(risingwave_pb::stream_plan::LoadFinishMutation {
996                associated_source_id: *associated_source_id,
997            }),
998            Mutation::ResetSource { source_id } => {
999                PbMutation::ResetSource(risingwave_pb::stream_plan::ResetSourceMutation {
1000                    source_id: source_id.as_raw_id(),
1001                })
1002            }
1003            Mutation::InjectSourceOffsets {
1004                source_id,
1005                split_offsets,
1006            } => PbMutation::InjectSourceOffsets(
1007                risingwave_pb::stream_plan::InjectSourceOffsetsMutation {
1008                    source_id: source_id.as_raw_id(),
1009                    split_offsets: split_offsets.clone(),
1010                },
1011            ),
1012        }
1013    }
1014
1015    fn from_protobuf(prost: &PbMutation) -> StreamExecutorResult<Self> {
1016        let mutation = match prost {
1017            PbMutation::Stop(stop) => Mutation::Stop(StopMutation {
1018                dropped_actors: stop.actors.iter().copied().collect(),
1019                dropped_sink_fragments: stop.dropped_sink_fragments.iter().copied().collect(),
1020            }),
1021
1022            PbMutation::Update(update) => Mutation::Update(UpdateMutation {
1023                dispatchers: update
1024                    .dispatcher_update
1025                    .iter()
1026                    .map(|u| (u.actor_id, u.clone()))
1027                    .into_group_map(),
1028                merges: update
1029                    .merge_update
1030                    .iter()
1031                    .map(|u| ((u.actor_id, u.upstream_fragment_id), u.clone()))
1032                    .collect(),
1033                vnode_bitmaps: update
1034                    .actor_vnode_bitmap_update
1035                    .iter()
1036                    .map(|(&actor_id, bitmap)| (actor_id, Arc::new(bitmap.into())))
1037                    .collect(),
1038                dropped_actors: update.dropped_actors.iter().copied().collect(),
1039                actor_splits: update
1040                    .actor_splits
1041                    .iter()
1042                    .map(|(&actor_id, splits)| {
1043                        (
1044                            actor_id,
1045                            splits
1046                                .splits
1047                                .iter()
1048                                .map(|split| split.try_into().unwrap())
1049                                .collect(),
1050                        )
1051                    })
1052                    .collect(),
1053                actor_new_dispatchers: update
1054                    .actor_new_dispatchers
1055                    .iter()
1056                    .map(|(&actor_id, dispatchers)| (actor_id, dispatchers.dispatchers.clone()))
1057                    .collect(),
1058                actor_cdc_table_snapshot_splits:
1059                    build_actor_cdc_table_snapshot_splits_with_generation(
1060                        update
1061                            .actor_cdc_table_snapshot_splits
1062                            .clone()
1063                            .unwrap_or_default(),
1064                    ),
1065                sink_schema_change: update
1066                    .sink_schema_change
1067                    .iter()
1068                    .map(|(sink_id, change)| (SinkId::from(*sink_id), change.clone()))
1069                    .collect(),
1070                subscriptions_to_drop: update.subscriptions_to_drop.clone(),
1071                iceberg_pk_index_compaction: update.iceberg_pk_index_compaction.clone(),
1072            }),
1073
1074            PbMutation::Add(add) => Mutation::Add(AddMutation {
1075                adds: add
1076                    .actor_dispatchers
1077                    .iter()
1078                    .map(|(&actor_id, dispatchers)| (actor_id, dispatchers.dispatchers.clone()))
1079                    .collect(),
1080                added_actors: add.added_actors.iter().copied().collect(),
1081                dropped_actors: add.dropped_actors.iter().copied().collect(),
1082                // TODO: remove this and use `SourceChangesSplit` after we support multiple
1083                // mutations.
1084                splits: add
1085                    .actor_splits
1086                    .iter()
1087                    .map(|(&actor_id, splits)| {
1088                        (
1089                            actor_id,
1090                            splits
1091                                .splits
1092                                .iter()
1093                                .map(|split| split.try_into().unwrap())
1094                                .collect(),
1095                        )
1096                    })
1097                    .collect(),
1098                pause: add.pause,
1099                subscriptions_to_add: add
1100                    .subscriptions_to_add
1101                    .iter()
1102                    .map(
1103                        |SubscriptionUpstreamInfo {
1104                             subscriber_id,
1105                             upstream_mv_table_id,
1106                         }| { (*upstream_mv_table_id, *subscriber_id) },
1107                    )
1108                    .collect(),
1109                backfill_nodes_to_pause: add.backfill_nodes_to_pause.iter().copied().collect(),
1110                actor_cdc_table_snapshot_splits:
1111                    build_actor_cdc_table_snapshot_splits_with_generation(
1112                        add.actor_cdc_table_snapshot_splits
1113                            .clone()
1114                            .unwrap_or_default(),
1115                    ),
1116                new_upstream_sinks: add
1117                    .new_upstream_sinks
1118                    .iter()
1119                    .map(|(k, v)| (*k, v.clone()))
1120                    .collect(),
1121                sink_log_store_flush: add.sink_log_store_flush.iter().copied().collect(),
1122            }),
1123
1124            PbMutation::Splits(s) => {
1125                let mut change_splits: Vec<(ActorId, Vec<SplitImpl>)> =
1126                    Vec::with_capacity(s.actor_splits.len());
1127                for (&actor_id, splits) in &s.actor_splits {
1128                    if !splits.splits.is_empty() {
1129                        change_splits.push((
1130                            actor_id,
1131                            splits
1132                                .splits
1133                                .iter()
1134                                .map(SplitImpl::try_from)
1135                                .try_collect()?,
1136                        ));
1137                    }
1138                }
1139                Mutation::SourceChangeSplit(change_splits.into_iter().collect())
1140            }
1141            PbMutation::Pause(_) => Mutation::Pause,
1142            PbMutation::Resume(_) => Mutation::Resume,
1143            PbMutation::Throttle(changes) => Mutation::Throttle(changes.fragment_throttle.clone()),
1144            PbMutation::DropSubscriptions(drop) => Mutation::DropSubscriptions {
1145                subscriptions_to_drop: drop.info.clone(),
1146            },
1147            PbMutation::ConnectorPropsChange(alter_connector_props) => {
1148                Mutation::ConnectorPropsChange(
1149                    alter_connector_props
1150                        .connector_props_infos
1151                        .iter()
1152                        .map(|(connector_id, options)| {
1153                            (
1154                                *connector_id,
1155                                options
1156                                    .connector_props_info
1157                                    .iter()
1158                                    .map(|(k, v)| (k.clone(), v.clone()))
1159                                    .collect(),
1160                            )
1161                        })
1162                        .collect(),
1163                )
1164            }
1165            PbMutation::StartFragmentBackfill(start_fragment_backfill) => {
1166                Mutation::StartFragmentBackfill {
1167                    fragment_ids: start_fragment_backfill
1168                        .fragment_ids
1169                        .iter()
1170                        .copied()
1171                        .collect(),
1172                }
1173            }
1174            PbMutation::RefreshStart(refresh_start) => Mutation::RefreshStart {
1175                table_id: refresh_start.table_id,
1176                associated_source_id: refresh_start.associated_source_id,
1177            },
1178            PbMutation::ListFinish(list_finish) => Mutation::ListFinish {
1179                associated_source_id: list_finish.associated_source_id,
1180            },
1181            PbMutation::LoadFinish(load_finish) => Mutation::LoadFinish {
1182                associated_source_id: load_finish.associated_source_id,
1183            },
1184            PbMutation::ResetSource(reset_source) => Mutation::ResetSource {
1185                source_id: SourceId::from(reset_source.source_id),
1186            },
1187            PbMutation::InjectSourceOffsets(inject) => Mutation::InjectSourceOffsets {
1188                source_id: SourceId::from(inject.source_id),
1189                split_offsets: inject.split_offsets.clone(),
1190            },
1191        };
1192        Ok(mutation)
1193    }
1194}
1195
1196impl<M> BarrierInner<M> {
1197    fn to_protobuf_inner(&self, barrier_fn: impl FnOnce(&M) -> Option<PbMutation>) -> PbBarrier {
1198        let Self {
1199            epoch,
1200            mutation,
1201            kind,
1202            tracing_context,
1203        } = self;
1204
1205        PbBarrier {
1206            epoch: Some(PbEpoch {
1207                curr: epoch.curr,
1208                prev: epoch.prev,
1209            }),
1210            mutation: barrier_fn(mutation).map(|mutation| PbBarrierMutation {
1211                mutation: Some(mutation),
1212            }),
1213            tracing_context: tracing_context.to_protobuf(),
1214            kind: *kind as _,
1215        }
1216    }
1217
1218    fn from_protobuf_inner(
1219        prost: &PbBarrier,
1220        mutation_from_pb: impl FnOnce(Option<&PbMutation>) -> StreamExecutorResult<M>,
1221    ) -> StreamExecutorResult<Self> {
1222        let epoch = prost.get_epoch()?;
1223
1224        Ok(Self {
1225            kind: prost.kind(),
1226            epoch: EpochPair::new(epoch.curr, epoch.prev),
1227            mutation: mutation_from_pb(
1228                (prost.mutation.as_ref()).and_then(|mutation| mutation.mutation.as_ref()),
1229            )?,
1230            tracing_context: TracingContext::from_protobuf(&prost.tracing_context),
1231        })
1232    }
1233
1234    pub fn map_mutation<M2>(self, f: impl FnOnce(M) -> M2) -> BarrierInner<M2> {
1235        BarrierInner {
1236            epoch: self.epoch,
1237            mutation: f(self.mutation),
1238            kind: self.kind,
1239            tracing_context: self.tracing_context,
1240        }
1241    }
1242}
1243
1244impl DispatcherBarrier {
1245    pub fn to_protobuf(&self) -> PbBarrier {
1246        self.to_protobuf_inner(|_| None)
1247    }
1248}
1249
1250impl Barrier {
1251    #[cfg(test)]
1252    pub fn to_protobuf(&self) -> PbBarrier {
1253        self.to_protobuf_inner(|mutation| mutation.as_ref().map(|mutation| mutation.to_protobuf()))
1254    }
1255
1256    pub fn from_protobuf(prost: &PbBarrier) -> StreamExecutorResult<Self> {
1257        Self::from_protobuf_inner(prost, |mutation| {
1258            mutation
1259                .map(|m| Mutation::from_protobuf(m).map(Arc::new))
1260                .transpose()
1261        })
1262    }
1263}
1264
1265#[derive(Debug, PartialEq, Eq, Clone, EstimateSize)]
1266pub struct Watermark {
1267    pub col_idx: usize,
1268    #[estimate_size(ignore)]
1269    pub data_type: DataType,
1270    pub val: ScalarImpl,
1271}
1272
1273impl PartialOrd for Watermark {
1274    fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
1275        Some(self.cmp(other))
1276    }
1277}
1278
1279impl Ord for Watermark {
1280    fn cmp(&self, other: &Self) -> std::cmp::Ordering {
1281        self.val.default_cmp(&other.val)
1282    }
1283}
1284
1285impl Watermark {
1286    pub fn new(col_idx: usize, data_type: DataType, val: ScalarImpl) -> Self {
1287        Self {
1288            col_idx,
1289            data_type,
1290            val,
1291        }
1292    }
1293
1294    pub async fn transform_with_expr(
1295        self,
1296        expr: &NonStrictExpression,
1297        new_col_idx: usize,
1298    ) -> Option<Self> {
1299        let Self { col_idx, val, .. } = self;
1300        let row = {
1301            let mut row = vec![None; col_idx + 1];
1302            row[col_idx] = Some(val);
1303            OwnedRow::new(row)
1304        };
1305        let val = expr.eval_row_infallible(&row).await?;
1306        Some(Self::new(new_col_idx, expr.inner().return_type(), val))
1307    }
1308
1309    /// Transform the watermark with the given output indices. If this watermark is not in the
1310    /// output, return `None`.
1311    pub fn transform_with_indices(self, output_indices: &[usize]) -> Option<Self> {
1312        output_indices
1313            .iter()
1314            .position(|p| *p == self.col_idx)
1315            .map(|new_col_idx| self.with_idx(new_col_idx))
1316    }
1317
1318    pub fn to_protobuf(&self) -> PbWatermark {
1319        PbWatermark {
1320            column: Some(PbInputRef {
1321                index: self.col_idx as _,
1322                r#type: Some(self.data_type.to_protobuf()),
1323            }),
1324            val: Some(&self.val).to_protobuf().into(),
1325        }
1326    }
1327
1328    pub fn from_protobuf(prost: &PbWatermark) -> StreamExecutorResult<Self> {
1329        let col_ref = prost.get_column()?;
1330        let data_type = DataType::from(col_ref.get_type()?);
1331        let val = Datum::from_protobuf(prost.get_val()?, &data_type)?
1332            .expect("watermark value cannot be null");
1333        Ok(Self::new(col_ref.get_index() as _, data_type, val))
1334    }
1335
1336    pub fn with_idx(self, idx: usize) -> Self {
1337        Self::new(idx, self.data_type, self.val)
1338    }
1339}
1340
1341#[cfg_attr(any(test, feature = "test"), derive(PartialEq))]
1342#[derive(Debug, EnumAsInner, Clone)]
1343pub enum MessageInner<M> {
1344    Chunk(StreamChunk),
1345    Barrier(BarrierInner<M>),
1346    Watermark(Watermark),
1347}
1348
1349impl<M> MessageInner<M> {
1350    pub fn map_mutation<M2>(self, f: impl FnOnce(M) -> M2) -> MessageInner<M2> {
1351        match self {
1352            MessageInner::Chunk(chunk) => MessageInner::Chunk(chunk),
1353            MessageInner::Barrier(barrier) => MessageInner::Barrier(barrier.map_mutation(f)),
1354            MessageInner::Watermark(watermark) => MessageInner::Watermark(watermark),
1355        }
1356    }
1357}
1358
1359pub type Message = MessageInner<BarrierMutationType>;
1360pub type DispatcherMessage = MessageInner<()>;
1361
1362/// `MessageBatchInner` is used exclusively by `Dispatcher` and the `Merger`/`Receiver` for exchanging messages between them.
1363/// It shares the same message type as the fundamental `MessageInner`, but batches multiple barriers into a single message.
1364#[derive(Debug, EnumAsInner, Clone)]
1365pub enum MessageBatchInner<M> {
1366    Chunk(StreamChunk),
1367    BarrierBatch(Vec<BarrierInner<M>>),
1368    Watermark(Watermark),
1369}
1370pub type MessageBatch = MessageBatchInner<BarrierMutationType>;
1371pub type DispatcherBarriers = Vec<DispatcherBarrier>;
1372pub type DispatcherMessageBatch = MessageBatchInner<()>;
1373
1374impl<M> From<MessageInner<M>> for MessageBatchInner<M> {
1375    fn from(m: MessageInner<M>) -> Self {
1376        match m {
1377            MessageInner::Chunk(c) => Self::Chunk(c),
1378            MessageInner::Barrier(b) => Self::BarrierBatch(vec![b]),
1379            MessageInner::Watermark(w) => Self::Watermark(w),
1380        }
1381    }
1382}
1383
1384impl From<StreamChunk> for Message {
1385    fn from(chunk: StreamChunk) -> Self {
1386        Message::Chunk(chunk)
1387    }
1388}
1389
1390impl<'a> TryFrom<&'a Message> for &'a Barrier {
1391    type Error = ();
1392
1393    fn try_from(m: &'a Message) -> std::result::Result<Self, Self::Error> {
1394        match m {
1395            Message::Chunk(_) => Err(()),
1396            Message::Barrier(b) => Ok(b),
1397            Message::Watermark(_) => Err(()),
1398        }
1399    }
1400}
1401
1402impl Message {
1403    /// Return true if the message is a stop barrier, meaning the stream
1404    /// will not continue, false otherwise.
1405    ///
1406    /// Note that this does not mean we will stop the current actor.
1407    #[cfg(test)]
1408    pub fn is_stop(&self) -> bool {
1409        matches!(
1410            self,
1411            Message::Barrier(Barrier {
1412                mutation,
1413                ..
1414            }) if mutation.as_ref().unwrap().is_stop_mutation()
1415        )
1416    }
1417}
1418
1419impl DispatcherMessageBatch {
1420    pub fn to_protobuf(&self) -> PbStreamMessageBatch {
1421        let prost = match self {
1422            Self::Chunk(stream_chunk) => {
1423                let prost_stream_chunk = stream_chunk.to_protobuf();
1424                StreamMessageBatch::StreamChunk(prost_stream_chunk)
1425            }
1426            Self::BarrierBatch(barrier_batch) => StreamMessageBatch::BarrierBatch(BarrierBatch {
1427                barriers: barrier_batch.iter().map(|b| b.to_protobuf()).collect(),
1428            }),
1429            Self::Watermark(watermark) => StreamMessageBatch::Watermark(watermark.to_protobuf()),
1430        };
1431        PbStreamMessageBatch {
1432            stream_message_batch: Some(prost),
1433        }
1434    }
1435
1436    pub fn from_protobuf(prost: &PbStreamMessageBatch) -> StreamExecutorResult<Self> {
1437        let res = match prost.get_stream_message_batch()? {
1438            StreamMessageBatch::StreamChunk(chunk) => {
1439                Self::Chunk(StreamChunk::from_protobuf(chunk)?)
1440            }
1441            StreamMessageBatch::BarrierBatch(barrier_batch) => {
1442                let barriers = barrier_batch
1443                    .barriers
1444                    .iter()
1445                    .map(|barrier| {
1446                        DispatcherBarrier::from_protobuf_inner(barrier, |mutation| {
1447                            if mutation.is_some() {
1448                                if cfg!(debug_assertions) {
1449                                    panic!("should not receive message of barrier with mutation");
1450                                } else {
1451                                    warn!(?barrier, "receive message of barrier with mutation");
1452                                }
1453                            }
1454                            Ok(())
1455                        })
1456                    })
1457                    .try_collect()?;
1458                Self::BarrierBatch(barriers)
1459            }
1460            StreamMessageBatch::Watermark(watermark) => {
1461                Self::Watermark(Watermark::from_protobuf(watermark)?)
1462            }
1463        };
1464        Ok(res)
1465    }
1466
1467    pub fn get_encoded_len(msg: &impl ::prost::Message) -> usize {
1468        ::prost::Message::encoded_len(msg)
1469    }
1470}
1471
1472pub type StreamKey = Vec<usize>;
1473pub type StreamKeyRef<'a> = &'a [usize];
1474pub type StreamKeyDataTypes = SmallVec<[DataType; 1]>;
1475
1476/// Expect the first message of the given `stream` as a barrier.
1477pub async fn expect_first_barrier<M: Debug>(
1478    stream: &mut (impl MessageStreamInner<M> + Unpin),
1479) -> StreamExecutorResult<BarrierInner<M>> {
1480    let message = stream
1481        .next()
1482        .instrument_await("expect_first_barrier")
1483        .await
1484        .context("failed to extract the first message: stream closed unexpectedly")??;
1485    let barrier = message
1486        .into_barrier()
1487        .expect("the first message must be a barrier");
1488    // TODO: Is this check correct?
1489    assert!(matches!(
1490        barrier.kind,
1491        BarrierKind::Checkpoint | BarrierKind::Initial
1492    ));
1493    Ok(barrier)
1494}
1495
1496/// Expect the first message of the given `stream` as a barrier.
1497pub async fn expect_first_barrier_from_aligned_stream(
1498    stream: &mut (impl AlignedMessageStream + Unpin),
1499) -> StreamExecutorResult<Barrier> {
1500    let message = stream
1501        .next()
1502        .instrument_await("expect_first_barrier")
1503        .await
1504        .context("failed to extract the first message: stream closed unexpectedly")??;
1505    let barrier = message
1506        .into_barrier()
1507        .expect("the first message must be a barrier");
1508    Ok(barrier)
1509}
1510
1511/// `StreamConsumer` is the last step in an actor.
1512pub trait StreamConsumer: Send + 'static {
1513    type BarrierStream: Stream<Item = StreamResult<Barrier>> + Send;
1514
1515    fn execute(self: Box<Self>) -> Self::BarrierStream;
1516}
1517
1518type BoxedMessageInput<InputId, M> = BoxedInput<InputId, MessageStreamItemInner<M>>;
1519
1520/// A stream for merging messages from multiple upstreams.
1521/// Can dynamically add and delete upstream streams.
1522/// For the meaning of the generic parameter `M` used, refer to `BarrierInner<M>`.
1523pub struct DynamicReceivers<InputId, M> {
1524    /// The barrier we're aligning to. If this is `None`, then `blocked_upstreams` is empty.
1525    barrier: Option<BarrierInner<M>>,
1526    /// The start timestamp of the current barrier. Used for measuring the alignment duration.
1527    start_ts: Option<Instant>,
1528    /// The upstreams that're blocked by the `barrier`.
1529    blocked: Vec<BoxedMessageInput<InputId, M>>,
1530    /// The upstreams that're not blocked and can be polled.
1531    active: FuturesUnordered<StreamFuture<BoxedMessageInput<InputId, M>>>,
1532    /// watermark column index -> `BufferedWatermarks`
1533    buffered_watermarks: BTreeMap<usize, BufferedWatermarks<InputId>>,
1534    /// Currently only used for union.
1535    barrier_align_duration: Option<LabelGuardedMetric<GenericCounter<AtomicU64>>>,
1536    /// Only for merge. If None, then we don't take `Instant::now()` and `observe` during `poll_next`
1537    merge_barrier_align_duration: Option<LabelGuardedMetric<GenericCounter<AtomicU64>>>,
1538}
1539
1540impl<InputId: Clone + Ord + Hash + std::fmt::Debug + Unpin, M: Clone + Unpin> Stream
1541    for DynamicReceivers<InputId, M>
1542{
1543    type Item = MessageStreamItemInner<M>;
1544
1545    fn poll_next(
1546        mut self: Pin<&mut Self>,
1547        cx: &mut std::task::Context<'_>,
1548    ) -> Poll<Option<Self::Item>> {
1549        if self.is_empty() {
1550            return Poll::Ready(None);
1551        }
1552
1553        loop {
1554            match futures::ready!(self.active.poll_next_unpin(cx)) {
1555                // Directly forward the error.
1556                Some((Some(Err(e)), _)) => {
1557                    return Poll::Ready(Some(Err(e)));
1558                }
1559                // Handle the message from some upstream.
1560                Some((Some(Ok(message)), remaining)) => {
1561                    let input_id = remaining.id();
1562                    match message {
1563                        MessageInner::Chunk(chunk) => {
1564                            // Continue polling this upstream by pushing it back to `active`.
1565                            self.active.push(remaining.into_future());
1566                            return Poll::Ready(Some(Ok(MessageInner::Chunk(chunk))));
1567                        }
1568                        MessageInner::Watermark(watermark) => {
1569                            // Continue polling this upstream by pushing it back to `active`.
1570                            self.active.push(remaining.into_future());
1571                            if let Some(watermark) = self.handle_watermark(input_id, watermark) {
1572                                return Poll::Ready(Some(Ok(MessageInner::Watermark(watermark))));
1573                            }
1574                        }
1575                        MessageInner::Barrier(barrier) => {
1576                            // Block this upstream by pushing it to `blocked`.
1577                            if self.blocked.is_empty() {
1578                                self.start_ts = Some(Instant::now());
1579                            }
1580                            self.blocked.push(remaining);
1581                            if let Some(current_barrier) = self.barrier.as_ref() {
1582                                if current_barrier.epoch != barrier.epoch {
1583                                    return Poll::Ready(Some(Err(
1584                                        StreamExecutorError::align_barrier(
1585                                            current_barrier.clone().map_mutation(|_| None),
1586                                            barrier.map_mutation(|_| None),
1587                                        ),
1588                                    )));
1589                                }
1590                            } else {
1591                                self.barrier = Some(barrier);
1592                            }
1593                        }
1594                    }
1595                }
1596                // We use barrier as the control message of the stream. That is, we always stop the
1597                // actors actively when we receive a `Stop` mutation, instead of relying on the stream
1598                // termination.
1599                //
1600                // Besides, in abnormal cases when the other side of the `Input` closes unexpectedly,
1601                // we also yield an `Err(ExchangeChannelClosed)`, which will hit the `Err` arm above.
1602                // So this branch will never be reached in all cases.
1603                Some((None, remaining)) => {
1604                    return Poll::Ready(Some(Err(StreamExecutorError::channel_closed(format!(
1605                        "upstream input {:?} unexpectedly closed",
1606                        remaining.id()
1607                    )))));
1608                }
1609                // There's no active upstreams. Process the barrier and resume the blocked ones.
1610                None => {
1611                    assert!(!self.blocked.is_empty());
1612
1613                    let start_ts = self
1614                        .start_ts
1615                        .take()
1616                        .expect("should have received at least one barrier");
1617                    if let Some(barrier_align_duration) = &self.barrier_align_duration {
1618                        barrier_align_duration.inc_by(start_ts.elapsed().as_nanos() as u64);
1619                    }
1620                    if let Some(merge_barrier_align_duration) = &self.merge_barrier_align_duration {
1621                        merge_barrier_align_duration.inc_by(start_ts.elapsed().as_nanos() as u64);
1622                    }
1623
1624                    break;
1625                }
1626            }
1627        }
1628
1629        assert!(self.active.is_terminated());
1630
1631        let barrier = self.barrier.take().unwrap();
1632
1633        let upstreams = std::mem::take(&mut self.blocked);
1634        self.extend_active(upstreams);
1635        assert!(!self.active.is_terminated());
1636
1637        Poll::Ready(Some(Ok(MessageInner::Barrier(barrier))))
1638    }
1639}
1640
1641impl<InputId: Clone + Ord + Hash + std::fmt::Debug, M> DynamicReceivers<InputId, M> {
1642    pub fn new(
1643        upstreams: Vec<BoxedMessageInput<InputId, M>>,
1644        barrier_align_duration: Option<LabelGuardedMetric<GenericCounter<AtomicU64>>>,
1645        merge_barrier_align_duration: Option<LabelGuardedMetric<GenericCounter<AtomicU64>>>,
1646    ) -> Self {
1647        let mut this = Self {
1648            barrier: None,
1649            start_ts: None,
1650            blocked: Vec::with_capacity(upstreams.len()),
1651            active: Default::default(),
1652            buffered_watermarks: Default::default(),
1653            merge_barrier_align_duration,
1654            barrier_align_duration,
1655        };
1656        this.extend_active(upstreams);
1657        this
1658    }
1659
1660    /// Extend the active upstreams with the given upstreams. The current stream must be at the
1661    /// clean state right after a barrier.
1662    pub fn extend_active(
1663        &mut self,
1664        upstreams: impl IntoIterator<Item = BoxedMessageInput<InputId, M>>,
1665    ) {
1666        assert!(self.blocked.is_empty() && self.barrier.is_none());
1667
1668        self.active
1669            .extend(upstreams.into_iter().map(|s| s.into_future()));
1670    }
1671
1672    /// Handle a new watermark message. Optionally returns the watermark message to emit.
1673    pub fn handle_watermark(
1674        &mut self,
1675        input_id: InputId,
1676        watermark: Watermark,
1677    ) -> Option<Watermark> {
1678        let col_idx = watermark.col_idx;
1679        // Insert a buffer watermarks when first received from a column.
1680        let upstream_ids: Vec<_> = self.upstream_input_ids().collect();
1681        let watermarks = self
1682            .buffered_watermarks
1683            .entry(col_idx)
1684            .or_insert_with(|| BufferedWatermarks::with_ids(upstream_ids));
1685        watermarks.handle_watermark(input_id, watermark)
1686    }
1687
1688    /// Consume `other` and add its upstreams to `self`. The two streams must be at the clean state
1689    /// right after a barrier.
1690    pub fn add_upstreams_from(
1691        &mut self,
1692        new_inputs: impl IntoIterator<Item = BoxedMessageInput<InputId, M>>,
1693    ) {
1694        assert!(self.blocked.is_empty() && self.barrier.is_none());
1695
1696        let new_inputs: Vec<_> = new_inputs.into_iter().collect();
1697        let input_ids = new_inputs.iter().map(|input| input.id());
1698        self.buffered_watermarks.values_mut().for_each(|buffers| {
1699            // Add buffers to the buffered watermarks for all cols
1700            buffers.add_buffers(input_ids.clone());
1701        });
1702        self.active
1703            .extend(new_inputs.into_iter().map(|s| s.into_future()));
1704    }
1705
1706    /// Remove upstreams from `self` in `upstream_input_ids`. The current stream must be at the
1707    /// clean state right after a barrier.
1708    /// The current container does not necessarily contain all the input ids passed in.
1709    pub fn remove_upstreams(&mut self, upstream_input_ids: &HashSet<InputId>) {
1710        assert!(self.blocked.is_empty() && self.barrier.is_none());
1711
1712        let new_upstreams = std::mem::take(&mut self.active)
1713            .into_iter()
1714            .map(|s| s.into_inner().unwrap())
1715            .filter(|u| !upstream_input_ids.contains(&u.id()));
1716        self.extend_active(new_upstreams);
1717        self.buffered_watermarks.values_mut().for_each(|buffers| {
1718            // Call `check_heap` in case the only upstream(s) that does not have
1719            // watermark in heap is removed
1720            buffers.remove_buffer(upstream_input_ids.clone());
1721        });
1722    }
1723
1724    pub fn merge_barrier_align_duration(
1725        &self,
1726    ) -> Option<LabelGuardedMetric<GenericCounter<AtomicU64>>> {
1727        self.merge_barrier_align_duration.clone()
1728    }
1729
1730    pub fn flush_buffered_watermarks(&mut self) {
1731        self.buffered_watermarks
1732            .values_mut()
1733            .for_each(|buffers| buffers.clear());
1734    }
1735
1736    pub fn upstream_input_ids(&self) -> impl Iterator<Item = InputId> + '_ {
1737        self.blocked
1738            .iter()
1739            .map(|s| s.id())
1740            .chain(self.active.iter().map(|s| s.get_ref().unwrap().id()))
1741    }
1742
1743    pub fn is_empty(&self) -> bool {
1744        self.blocked.is_empty() && self.active.is_empty()
1745    }
1746}
1747
1748// Explanation of why we need `DispatchBarrierBuffer`:
1749//
1750// When we need to create or replace an upstream fragment for the current fragment, the `Merge` operator must
1751// add some new upstream actor inputs. However, the `Merge` operator may still have old upstreams. We must wait
1752// for these old upstreams to completely process their barriers and align before we can safely update the
1753// `upstream-input-set`.
1754//
1755// Meanwhile, the creation of a new upstream actor can only succeed after the channel to the downstream `Merge`
1756// operator has been established. This creates a potential dependency chain: [new_actor_creation ->
1757// downstream_merge_update -> old_actor_processing]
1758//
1759// To address this, we split the application of a barrier's `Mutation` into two steps:
1760// 1. Parse the `Mutation`. If there is an addition on the upstream-set, establish a channel with the upstream
1761//    and cache it.
1762// 2. When the upstream barrier actually arrives, apply the cached upstream changes to the upstream-set
1763//
1764// Additionally, since receiving a barrier from current upstream input and from the `barrier_rx` are
1765// asynchronous, we cannot determine which will arrive first. Therefore, when a barrier is received from an
1766// upstream: if a cached mutation is present, we apply it. Otherwise, we must fetch a new barrier from
1767// `barrier_rx`.
1768pub(crate) struct DispatchBarrierBuffer {
1769    buffer: VecDeque<(Barrier, Option<Vec<BoxedActorInput>>)>,
1770    barrier_rx: mpsc::UnboundedReceiver<Barrier>,
1771    recv_state: BarrierReceiverState,
1772    curr_upstream_fragment_id: FragmentId,
1773    actor_id: ActorId,
1774    // read-only context for building new inputs
1775    build_input_ctx: Arc<BuildInputContext>,
1776}
1777
1778struct BuildInputContext {
1779    pub actor_id: ActorId,
1780    pub local_barrier_manager: LocalBarrierManager,
1781    pub metrics: Arc<StreamingMetrics>,
1782    pub fragment_id: FragmentId,
1783    pub actor_config: Arc<StreamingConfig>,
1784}
1785
1786type BoxedNewInputsFuture =
1787    Pin<Box<dyn Future<Output = StreamExecutorResult<Vec<BoxedActorInput>>> + Send>>;
1788
1789enum BarrierReceiverState {
1790    ReceivingBarrier,
1791    CreatingNewInput(Barrier, BoxedNewInputsFuture),
1792}
1793
1794impl DispatchBarrierBuffer {
1795    pub fn new(
1796        barrier_rx: mpsc::UnboundedReceiver<Barrier>,
1797        actor_id: ActorId,
1798        curr_upstream_fragment_id: FragmentId,
1799        local_barrier_manager: LocalBarrierManager,
1800        metrics: Arc<StreamingMetrics>,
1801        fragment_id: FragmentId,
1802        actor_config: Arc<StreamingConfig>,
1803    ) -> Self {
1804        Self {
1805            buffer: VecDeque::new(),
1806            barrier_rx,
1807            recv_state: BarrierReceiverState::ReceivingBarrier,
1808            curr_upstream_fragment_id,
1809            actor_id,
1810            build_input_ctx: Arc::new(BuildInputContext {
1811                actor_id,
1812                local_barrier_manager,
1813                metrics,
1814                fragment_id,
1815                actor_config,
1816            }),
1817        }
1818    }
1819
1820    pub async fn await_next_message(
1821        &mut self,
1822        stream: &mut (impl Stream<Item = StreamExecutorResult<DispatcherMessage>> + Unpin),
1823        metrics: &ActorInputMetrics,
1824        upstream_is_empty: bool,
1825    ) -> StreamExecutorResult<DispatcherMessage> {
1826        if upstream_is_empty {
1827            while self.buffer.is_empty() {
1828                self.try_fetch_barrier_rx(false).await?;
1829            }
1830            let (barrier, _) = self.buffer.front().unwrap();
1831            return Ok(DispatcherMessage::Barrier(
1832                barrier.clone().into_dispatcher(),
1833            ));
1834        }
1835
1836        let mut start_time = Instant::now();
1837        let interval_duration = Duration::from_secs(15);
1838        let mut interval =
1839            tokio::time::interval_at(start_time + interval_duration, interval_duration);
1840
1841        loop {
1842            tokio::select! {
1843                biased;
1844                msg = stream.try_next() => {
1845                    metrics
1846                        .actor_input_buffer_blocking_duration_ns
1847                        .inc_by(start_time.elapsed().as_nanos() as u64);
1848                    return msg?.ok_or_else(
1849                        || StreamExecutorError::channel_closed("upstream executor closed unexpectedly")
1850                    );
1851                }
1852
1853                e = self.continuously_fetch_barrier_rx() => {
1854                    return Err(e);
1855                }
1856
1857                _ = interval.tick() => {
1858                    start_time = Instant::now();
1859                    metrics.actor_input_buffer_blocking_duration_ns.inc_by(interval_duration.as_nanos() as u64);
1860                    continue;
1861                }
1862            }
1863        }
1864    }
1865
1866    pub async fn pop_barrier_with_inputs(
1867        &mut self,
1868        barrier: DispatcherBarrier,
1869    ) -> StreamExecutorResult<(Barrier, Option<Vec<BoxedActorInput>>)> {
1870        while self.buffer.is_empty() {
1871            self.try_fetch_barrier_rx(false).await?;
1872        }
1873        let (recv_barrier, inputs) = self.buffer.pop_front().unwrap();
1874        assert_equal_dispatcher_barrier(&recv_barrier, &barrier);
1875
1876        Ok((recv_barrier, inputs))
1877    }
1878
1879    async fn continuously_fetch_barrier_rx(&mut self) -> StreamExecutorError {
1880        loop {
1881            if let Err(e) = self.try_fetch_barrier_rx(true).await {
1882                return e;
1883            }
1884        }
1885    }
1886
1887    async fn try_fetch_barrier_rx(&mut self, pending_on_end: bool) -> StreamExecutorResult<()> {
1888        match &mut self.recv_state {
1889            BarrierReceiverState::ReceivingBarrier => {
1890                let Some(barrier) = self.barrier_rx.recv().await else {
1891                    if pending_on_end {
1892                        return pending().await;
1893                    } else {
1894                        return Err(StreamExecutorError::channel_closed(
1895                            "barrier channel closed unexpectedly",
1896                        ));
1897                    }
1898                };
1899                if let Some(fut) = self.pre_apply_barrier(&barrier) {
1900                    self.recv_state = BarrierReceiverState::CreatingNewInput(barrier, fut);
1901                } else {
1902                    self.buffer.push_back((barrier, None));
1903                }
1904            }
1905            BarrierReceiverState::CreatingNewInput(barrier, fut) => {
1906                let new_inputs = fut.await?;
1907                self.buffer.push_back((barrier.clone(), Some(new_inputs)));
1908                self.recv_state = BarrierReceiverState::ReceivingBarrier;
1909            }
1910        }
1911        Ok(())
1912    }
1913
1914    fn pre_apply_barrier(&mut self, barrier: &Barrier) -> Option<BoxedNewInputsFuture> {
1915        let update = barrier.as_update_merge(self.actor_id, self.curr_upstream_fragment_id)?;
1916        let upstream_fragment_id = update
1917            .new_upstream_fragment_id
1918            .unwrap_or(self.curr_upstream_fragment_id);
1919        // Keep the fragment id in sync even when switching to an empty upstream. Otherwise a
1920        // subsequent reattach cannot be recognized as a merge update and its new input may not
1921        // receive the matching first barrier before the barrier is forwarded.
1922        self.curr_upstream_fragment_id = upstream_fragment_id;
1923
1924        if !update.added_upstream_actors.is_empty() {
1925            let ctx = self.build_input_ctx.clone();
1926            let added_upstream_actors = update.added_upstream_actors.clone();
1927            let barrier = barrier.clone();
1928            let fut = async move {
1929                try_join_all(added_upstream_actors.iter().map(|upstream_actor| async {
1930                    let mut new_input = new_input(
1931                        &ctx.local_barrier_manager,
1932                        ctx.metrics.clone(),
1933                        ctx.actor_id,
1934                        ctx.fragment_id,
1935                        upstream_actor,
1936                        upstream_fragment_id,
1937                        ctx.actor_config.clone(),
1938                    )
1939                    .await?;
1940
1941                    // Poll the first barrier from the new upstreams. It must be the same as the one we polled from
1942                    // original upstreams.
1943                    let first_barrier = expect_first_barrier(&mut new_input).await?;
1944                    assert_equal_dispatcher_barrier(&barrier, &first_barrier);
1945
1946                    StreamExecutorResult::Ok(new_input)
1947                }))
1948                .await
1949            }
1950            .boxed();
1951
1952            Some(fut)
1953        } else {
1954            None
1955        }
1956    }
1957}