1use std::collections::HashMap;
16use std::sync::Arc;
17
18use futures::future::{Either as FutureEither, pending, select};
19use futures::{StreamExt, TryStreamExt, pin_mut};
20use futures_async_stream::try_stream;
21use itertools::Itertools;
22use risingwave_common::array::{DataChunk, Op, StreamChunk};
23use risingwave_common::catalog::Schema;
24use risingwave_common::hash::{VirtualNode, VnodeBitmapExt};
25use risingwave_common::row::{OwnedRow, Row, RowExt};
26use risingwave_common::types::{Datum, ToOwnedDatum};
27use risingwave_common::util::chunk_coalesce::DataChunkBuilder;
28use risingwave_common::util::sort_util::cmp_datum_iter;
29use risingwave_common_rate_limit::{MonitoredRateLimiter, RateLimit, RateLimiter};
30use risingwave_pb::common::ThrottleType;
31use risingwave_storage::StateStore;
32use risingwave_storage::store::PrefetchOptions;
33
34use crate::common::table::state_table::{FlushedStateTableReader, StateTable};
35use crate::executor::backfill::utils::create_builder;
36use crate::executor::prelude::*;
37use crate::task::{CreateMviewProgressReporter, FragmentId};
38
39type Builders = HashMap<VirtualNode, DataChunkBuilder>;
40
41#[derive(Clone, Debug, PartialEq, Eq)]
43enum LocalityBackfillProgress {
44 NotStarted,
46 InProgress {
48 current_pos: OwnedRow,
50 processed_rows: u64,
52 },
53 Completed {
55 final_pos: OwnedRow,
57 total_rows: u64,
59 },
60}
61
62#[derive(Clone, Debug)]
64struct LocalityBackfillState {
65 per_vnode: HashMap<VirtualNode, LocalityBackfillProgress>,
67 total_snapshot_rows: u64,
69}
70
71impl LocalityBackfillState {
72 fn new(vnodes: impl Iterator<Item = VirtualNode>) -> Self {
73 let per_vnode = vnodes
74 .map(|vnode| (vnode, LocalityBackfillProgress::NotStarted))
75 .collect();
76 Self {
77 per_vnode,
78 total_snapshot_rows: 0,
79 }
80 }
81
82 fn is_completed(&self) -> bool {
83 self.per_vnode
84 .values()
85 .all(|progress| matches!(progress, LocalityBackfillProgress::Completed { .. }))
86 }
87
88 fn vnodes(&self) -> impl Iterator<Item = (VirtualNode, &LocalityBackfillProgress)> {
89 self.per_vnode
90 .iter()
91 .map(|(&vnode, progress)| (vnode, progress))
92 }
93
94 fn has_progress(&self) -> bool {
95 self.per_vnode
96 .values()
97 .any(|progress| matches!(progress, LocalityBackfillProgress::InProgress { .. }))
98 }
99
100 fn update_progress(&mut self, vnode: VirtualNode, new_pos: OwnedRow, row_count_delta: u64) {
101 let progress = self.per_vnode.get_mut(&vnode).unwrap();
102 match progress {
103 LocalityBackfillProgress::NotStarted => {
104 *progress = LocalityBackfillProgress::InProgress {
105 current_pos: new_pos,
106 processed_rows: row_count_delta,
107 };
108 }
109 LocalityBackfillProgress::InProgress { processed_rows, .. } => {
110 *progress = LocalityBackfillProgress::InProgress {
111 current_pos: new_pos,
112 processed_rows: *processed_rows + row_count_delta,
113 };
114 }
115 LocalityBackfillProgress::Completed { .. } => {
116 }
118 }
119 self.total_snapshot_rows += row_count_delta;
120 }
121
122 fn finish_vnode(&mut self, vnode: VirtualNode, pk_len: usize) {
123 let progress = self.per_vnode.get_mut(&vnode).unwrap();
124 match progress {
125 LocalityBackfillProgress::NotStarted => {
126 let final_pos = OwnedRow::new(vec![None; pk_len]);
128 *progress = LocalityBackfillProgress::Completed {
129 final_pos,
130 total_rows: 0,
131 };
132 }
133 LocalityBackfillProgress::InProgress {
134 current_pos,
135 processed_rows,
136 } => {
137 *progress = LocalityBackfillProgress::Completed {
138 final_pos: current_pos.clone(),
139 total_rows: *processed_rows,
140 };
141 }
142 LocalityBackfillProgress::Completed { .. } => {
143 }
145 }
146 }
147
148 fn get_progress(&self, vnode: &VirtualNode) -> &LocalityBackfillProgress {
149 self.per_vnode.get(vnode).unwrap()
150 }
151}
152
153pub struct LocalityProviderExecutor<S: StateStore> {
165 upstream: Executor,
167
168 #[expect(dead_code)]
170 locality_columns: Vec<usize>,
171
172 state_table: StateTable<S>,
174
175 progress_table: StateTable<S>,
177
178 input_schema: Schema,
179
180 progress: CreateMviewProgressReporter,
182
183 fragment_id: FragmentId,
184
185 actor_id: ActorId,
186
187 metrics: Arc<StreamingMetrics>,
189
190 chunk_size: usize,
192
193 rate_limiter: MonitoredRateLimiter,
194}
195
196impl<S: StateStore> LocalityProviderExecutor<S> {
197 #[expect(clippy::too_many_arguments)]
198 pub fn new(
199 upstream: Executor,
200 locality_columns: Vec<usize>,
201 state_table: StateTable<S>,
202 progress_table: StateTable<S>,
203 input_schema: Schema,
204 progress: CreateMviewProgressReporter,
205 metrics: Arc<StreamingMetrics>,
206 chunk_size: usize,
207 fragment_id: FragmentId,
208 rate_limit: RateLimit,
209 ) -> Self {
210 let rate_limiter = RateLimiter::new(rate_limit).monitored(state_table.table_id());
211 Self {
212 upstream,
213 locality_columns,
214 state_table,
215 progress_table,
216 input_schema,
217 actor_id: progress.actor_id(),
218 progress,
219 metrics,
220 chunk_size,
221 fragment_id,
222 rate_limiter,
223 }
224 }
225
226 fn apply_throttle(
228 rate_limiter: &MonitoredRateLimiter,
229 fragment_id: FragmentId,
230 barrier: &Barrier,
231 ) -> Option<RateLimit> {
232 let Some(Mutation::Throttle(fragment_to_apply)) = barrier.mutation.as_deref() else {
233 return None;
234 };
235 let entry = fragment_to_apply.get(&fragment_id)?;
236 if entry.throttle_type() != ThrottleType::Backfill {
237 return None;
238 }
239 let new_rate_limit = entry.rate_limit.into();
240 let old_rate_limit = rate_limiter.update(new_rate_limit);
241 (old_rate_limit != new_rate_limit).then(|| {
242 tracing::info!(
243 ?old_rate_limit,
244 ?new_rate_limit,
245 %fragment_id,
246 "locality backfill rate limit changed"
247 );
248 new_rate_limit
249 })
250 }
251
252 #[try_stream(ok = (VirtualNode, OwnedRow), error = StreamExecutorError)]
254 async fn make_snapshot_stream<'a>(
255 reader: FlushedStateTableReader<S>,
256 backfill_state: LocalityBackfillState,
257 rate_limiter: &'a MonitoredRateLimiter,
258 ) {
259 for vnode in reader.vnodes().iter_vnodes() {
261 let progress = backfill_state.get_progress(&vnode);
262
263 let current_pos = match progress {
264 LocalityBackfillProgress::NotStarted => None,
265 LocalityBackfillProgress::Completed { .. } => {
266 continue;
268 }
269 LocalityBackfillProgress::InProgress { current_pos, .. } => {
270 Some(current_pos.clone())
271 }
272 };
273
274 let range_bounds = if let Some(ref pos) = current_pos {
276 let start_bound = std::ops::Bound::Excluded(pos.as_inner());
277 (start_bound, std::ops::Bound::<&[Datum]>::Unbounded)
278 } else {
279 (
280 std::ops::Bound::<&[Datum]>::Unbounded,
281 std::ops::Bound::<&[Datum]>::Unbounded,
282 )
283 };
284
285 let iter = reader
287 .iter_with_vnode(
288 vnode,
289 &range_bounds,
290 PrefetchOptions::prefetch_for_small_range_scan(),
291 )
292 .await?;
293 pin_mut!(iter);
294
295 while let Some(row) = iter.try_next().await? {
296 rate_limiter.wait(1).await;
297 yield (vnode, row);
298 }
299 }
300 }
301
302 async fn persist_backfill_state(
304 progress_table: &mut StateTable<S>,
305 backfill_state: &LocalityBackfillState,
306 ) -> StreamExecutorResult<()> {
307 for (vnode, progress) in &backfill_state.per_vnode {
308 let (is_finished, current_pos, row_count) = match progress {
309 LocalityBackfillProgress::NotStarted => continue, LocalityBackfillProgress::InProgress {
311 current_pos,
312 processed_rows,
313 } => (false, current_pos.clone(), *processed_rows),
314 LocalityBackfillProgress::Completed {
315 final_pos,
316 total_rows,
317 } => (true, final_pos.clone(), *total_rows),
318 };
319
320 let mut row_data = vec![Some(vnode.to_scalar().into())];
322 row_data.extend(current_pos);
323 row_data.push(Some(risingwave_common::types::ScalarImpl::Bool(
324 is_finished,
325 )));
326 row_data.push(Some(risingwave_common::types::ScalarImpl::Int64(
327 row_count as i64,
328 )));
329
330 let new_row = OwnedRow::new(row_data);
331
332 let key_data = vec![Some(vnode.to_scalar().into())];
335 let key = OwnedRow::new(key_data);
336
337 if let Some(existing_row) = progress_table.get_row(&key).await? {
338 progress_table.update(existing_row, new_row);
340 } else {
341 progress_table.insert(new_row);
343 }
344 }
345 Ok(())
346 }
347
348 async fn load_backfill_state(
350 progress_table: &StateTable<S>,
351 ) -> StreamExecutorResult<LocalityBackfillState> {
352 let mut backfill_state = LocalityBackfillState::new(progress_table.vnodes().iter_vnodes());
353 let mut total_snapshot_rows = 0;
354
355 for vnode in progress_table.vnodes().iter_vnodes() {
357 let key_data = vec![Some(vnode.to_scalar().into())];
359
360 let key = OwnedRow::new(key_data);
361
362 if let Some(row) = progress_table.get_row(&key).await? {
363 let finished_col_idx = row.len() - 2;
365 let is_finished = row
366 .datum_at(finished_col_idx)
367 .map(|d| d.into_bool())
368 .unwrap_or(false);
369
370 let row_count = row
372 .datum_at(row.len() - 1)
373 .map(|d| d.into_int64() as u64)
374 .unwrap_or(0);
375
376 let current_pos_data: Vec<Datum> = (1..finished_col_idx)
377 .map(|i| row.datum_at(i).to_owned_datum())
378 .collect();
379 let current_pos = OwnedRow::new(current_pos_data);
380
381 let progress = if is_finished {
383 LocalityBackfillProgress::Completed {
384 final_pos: current_pos,
385 total_rows: row_count,
386 }
387 } else {
388 LocalityBackfillProgress::InProgress {
389 current_pos,
390 processed_rows: row_count,
391 }
392 };
393
394 backfill_state.per_vnode.insert(vnode, progress);
395 total_snapshot_rows += row_count;
396 }
397 }
399
400 backfill_state.total_snapshot_rows = total_snapshot_rows;
401 Ok(backfill_state)
402 }
403
404 fn mark_chunk(
406 chunk: StreamChunk,
407 backfill_state: &LocalityBackfillState,
408 state_table: &StateTable<S>,
409 ) -> StreamExecutorResult<StreamChunk> {
410 let chunk = chunk.compact_vis();
411 let (data, ops) = chunk.into_parts();
412 let mut new_visibility = risingwave_common::bitmap::BitmapBuilder::with_capacity(ops.len());
413
414 let pk_indices = state_table.pk_indices();
415 let pk_order = state_table.pk_serde().get_order_types();
416
417 for row in data.rows() {
418 let pk = row.project(pk_indices);
420 let vnode = state_table.compute_vnode_by_pk(pk);
421
422 let visible = match backfill_state.get_progress(&vnode) {
423 LocalityBackfillProgress::Completed { .. } => true,
424 LocalityBackfillProgress::NotStarted => false,
425 LocalityBackfillProgress::InProgress { current_pos, .. } => {
426 cmp_datum_iter(pk.iter(), current_pos.iter(), pk_order.iter().copied()).is_le()
428 }
429 };
430
431 new_visibility.append(visible);
432 }
433
434 let (columns, _) = data.into_parts();
435 let chunk = StreamChunk::with_visibility(ops, columns, new_visibility.finish());
436 Ok(chunk)
437 }
438
439 fn handle_snapshot_chunk(
440 data_chunk: DataChunk,
441 vnode: VirtualNode,
442 pk_indices: &[usize],
443 backfill_state: &mut LocalityBackfillState,
444 cur_barrier_snapshot_processed_rows: &mut u64,
445 ) -> StreamExecutorResult<StreamChunk> {
446 let chunk = StreamChunk::from_parts(vec![Op::Insert; data_chunk.cardinality()], data_chunk);
447 let chunk_cardinality = chunk.cardinality() as u64;
448
449 if let Some(last_row) = chunk.rows().last() {
452 let pk = last_row.1.project(pk_indices);
453 let pk_owned = pk.into_owned_row();
454 backfill_state.update_progress(vnode, pk_owned, chunk_cardinality);
455 }
456
457 *cur_barrier_snapshot_processed_rows += chunk_cardinality;
458 Ok(chunk)
459 }
460}
461
462impl<S: StateStore> Execute for LocalityProviderExecutor<S> {
463 fn execute(self: Box<Self>) -> BoxedMessageStream {
464 self.execute_inner().boxed()
465 }
466}
467
468impl<S: StateStore> LocalityProviderExecutor<S> {
469 #[try_stream(ok = Message, error = StreamExecutorError)]
470 async fn execute_inner(mut self) {
471 let mut upstream = self.upstream.execute();
472
473 let first_barrier = expect_first_barrier(&mut upstream).await?;
475 let first_epoch = first_barrier.epoch;
476
477 yield Message::Barrier(first_barrier);
479
480 let mut state_table = self.state_table;
481 let mut progress_table = self.progress_table;
482 let rate_limiter = self.rate_limiter;
483
484 state_table.init_epoch(first_epoch).await?;
486 progress_table.init_epoch(first_epoch).await?;
487
488 let mut backfill_state = Self::load_backfill_state(&progress_table).await?;
490
491 let pk_indices = state_table.pk_indices().iter().cloned().collect_vec();
493
494 let need_backfill = !backfill_state.is_completed();
495 let mut report_finished_on_first_barrier = !need_backfill;
496
497 let need_buffering = backfill_state
498 .per_vnode
499 .values()
500 .all(|progress| matches!(progress, LocalityBackfillProgress::NotStarted));
501 if need_buffering {
503 let mut start_backfill = false;
505
506 #[for_await]
507 for msg in upstream.by_ref() {
508 let msg = msg?;
509
510 match msg {
511 Message::Watermark(_) => {
512 }
514 Message::Chunk(chunk) => {
515 state_table.write_chunk(chunk);
516 state_table.try_flush().await?;
517 }
518 Message::Barrier(barrier) => {
519 let epoch = barrier.epoch;
520 Self::apply_throttle(&rate_limiter, self.fragment_id, &barrier);
521
522 if let Some(mutation) = barrier.mutation.as_deref() {
524 use crate::executor::Mutation;
525 if let Mutation::StartFragmentBackfill { fragment_ids } = mutation
526 && fragment_ids.contains(&self.fragment_id)
527 {
528 tracing::info!(
529 "Start backfill of locality provider with fragment id: {:?}",
530 &self.fragment_id
531 );
532 start_backfill = true;
533 }
534 }
535
536 barrier.assume_no_update_vnode_bitmap(self.actor_id)?;
538 state_table
539 .commit_assert_no_update_vnode_bitmap(epoch)
540 .await?;
541 progress_table
542 .commit_assert_no_update_vnode_bitmap(epoch)
543 .await?;
544
545 yield Message::Barrier(barrier);
546
547 if start_backfill {
549 break;
550 }
551 }
552 }
553 }
554 }
555
556 if need_backfill {
583 let mut upstream_chunk_buffer: Vec<StreamChunk> = vec![];
584
585 let metrics = self
586 .metrics
587 .new_backfill_metrics(state_table.table_id(), self.actor_id);
588
589 let snapshot_data_types = self.input_schema.data_types();
591 let vnodes = state_table.vnodes().clone();
592 let new_builders = |rate_limit| -> Builders {
593 vnodes
594 .iter_vnodes()
595 .map(|vnode| {
596 let builder = create_builder(
597 rate_limit,
598 self.chunk_size,
599 snapshot_data_types.clone(),
600 );
601 (vnode, builder)
602 })
603 .collect()
604 };
605 let mut builders = new_builders(rate_limiter.rate_limit());
606
607 let snapshot_reader = state_table.flushed_snapshot_reader();
608 let snapshot_stream = Self::make_snapshot_stream(
609 snapshot_reader.clone(),
610 backfill_state.clone(),
611 &rate_limiter,
612 );
613 pin_mut!(snapshot_stream);
614
615 'backfill_loop: loop {
616 let mut cur_barrier_snapshot_processed_rows: u64 = 0;
617 let mut cur_barrier_upstream_processed_rows: u64 = 0;
618
619 let barrier = loop {
622 let upstream_next = upstream.next();
623 let mut snapshot_stream_ref = snapshot_stream.as_mut();
624 let snapshot_paused = rate_limiter.rate_limit().is_paused();
625 let snapshot_next = async move {
626 if snapshot_paused {
627 pending().await
628 } else {
629 snapshot_stream_ref.next().await
630 }
631 };
632 pin_mut!(upstream_next);
633 pin_mut!(snapshot_next);
634
635 match select(upstream_next, snapshot_next).await {
636 FutureEither::Left((msg, _)) => match msg.transpose()? {
637 Some(Message::Barrier(barrier)) => {
638 break barrier;
640 }
641 Some(Message::Chunk(chunk)) => {
642 upstream_chunk_buffer.push(chunk.compact_vis());
644 }
645 Some(Message::Watermark(_)) => {
646 }
648 None => {
649 return Err(anyhow::anyhow!(
650 "locality provider upstream ended unexpectedly during backfill"
651 )
652 .into());
653 }
654 },
655 FutureEither::Right((msg, _)) => match msg.transpose()? {
656 Some((vnode, row)) => {
657 let builder = builders.get_mut(&vnode).unwrap();
659 if let Some(data_chunk) = builder.append_one_row(row) {
660 let chunk = Self::handle_snapshot_chunk(
662 data_chunk,
663 vnode,
664 &pk_indices,
665 &mut backfill_state,
666 &mut cur_barrier_snapshot_processed_rows,
667 )?;
668 yield Message::Chunk(chunk);
669 }
670 }
673 None => {
674 for (vnode, builder) in &mut builders {
677 if let Some(data_chunk) = builder.consume_all() {
678 let chunk = Self::handle_snapshot_chunk(
679 data_chunk,
680 *vnode,
681 &pk_indices,
682 &mut backfill_state,
683 &mut cur_barrier_snapshot_processed_rows,
684 )?;
685 yield Message::Chunk(chunk);
686 }
687 }
688
689 for chunk in upstream_chunk_buffer.drain(..) {
691 let chunk_cardinality = chunk.cardinality() as u64;
692 cur_barrier_upstream_processed_rows += chunk_cardinality;
693 yield Message::Chunk(chunk);
694 }
695 metrics
696 .backfill_snapshot_read_row_count
697 .inc_by(cur_barrier_snapshot_processed_rows);
698 metrics
699 .backfill_upstream_output_row_count
700 .inc_by(cur_barrier_upstream_processed_rows);
701 break 'backfill_loop;
702 }
703 },
704 }
705 };
706
707 for (vnode, builder) in &mut builders {
709 if let Some(data_chunk) = builder.consume_all() {
710 let chunk = Self::handle_snapshot_chunk(
711 data_chunk,
712 *vnode,
713 &pk_indices,
714 &mut backfill_state,
715 &mut cur_barrier_snapshot_processed_rows,
716 )?;
717 yield Message::Chunk(chunk);
718 }
719 }
720
721 if let Some(new_rate_limit) =
722 Self::apply_throttle(&rate_limiter, self.fragment_id, &barrier)
723 {
724 builders = new_builders(new_rate_limit);
725 }
726
727 let should_refresh_snapshot = !upstream_chunk_buffer.is_empty();
729 for chunk in upstream_chunk_buffer.drain(..) {
730 cur_barrier_upstream_processed_rows += chunk.cardinality() as u64;
731
732 if backfill_state.has_progress() {
734 let marked_chunk =
735 Self::mark_chunk(chunk.clone(), &backfill_state, &state_table)?;
736 yield Message::Chunk(marked_chunk);
737 }
738
739 state_table.write_chunk(chunk);
742 }
743
744 let barrier_epoch = barrier.epoch;
745 barrier.assume_no_update_vnode_bitmap(self.actor_id)?;
746 state_table
747 .commit_assert_no_update_vnode_bitmap(barrier_epoch)
748 .await?;
749
750 let total_snapshot_processed_rows: u64 = backfill_state
753 .vnodes()
754 .map(|(_, progress)| match *progress {
755 LocalityBackfillProgress::InProgress { processed_rows, .. } => {
756 processed_rows
757 }
758 LocalityBackfillProgress::Completed { total_rows, .. } => total_rows,
759 LocalityBackfillProgress::NotStarted => 0,
760 })
761 .sum();
762
763 self.progress.update_with_buffered_rows(
764 barrier.epoch,
765 barrier.epoch.curr, total_snapshot_processed_rows,
767 0,
768 );
769
770 Self::persist_backfill_state(&mut progress_table, &backfill_state).await?;
772 progress_table
773 .commit_assert_no_update_vnode_bitmap(barrier_epoch)
774 .await?;
775
776 metrics
777 .backfill_snapshot_read_row_count
778 .inc_by(cur_barrier_snapshot_processed_rows);
779 metrics
780 .backfill_upstream_output_row_count
781 .inc_by(cur_barrier_upstream_processed_rows);
782
783 yield Message::Barrier(barrier);
784
785 if should_refresh_snapshot {
786 snapshot_stream.set(Self::make_snapshot_stream(
787 snapshot_reader.clone(),
788 backfill_state.clone(),
789 &rate_limiter,
790 ));
791 }
792 }
793 }
794
795 tracing::debug!("Locality provider backfill finished, forwarding upstream directly");
796
797 if need_backfill && !backfill_state.is_completed() {
799 while let Some(Ok(msg)) = upstream.next().await {
800 match msg {
801 Message::Barrier(barrier) => {
802 barrier.assume_no_update_vnode_bitmap(self.actor_id)?;
803
804 state_table
806 .commit_assert_no_update_vnode_bitmap(barrier.epoch)
807 .await?;
808
809 for vnode in state_table.vnodes().iter_vnodes() {
811 backfill_state.finish_vnode(vnode, pk_indices.len());
812 }
813
814 let total_snapshot_processed_rows: u64 = backfill_state
816 .vnodes()
817 .map(|(_, progress)| match *progress {
818 LocalityBackfillProgress::Completed { total_rows, .. } => {
819 total_rows
820 }
821 LocalityBackfillProgress::InProgress { processed_rows, .. } => {
822 processed_rows
823 }
824 LocalityBackfillProgress::NotStarted => 0,
825 })
826 .sum();
827
828 self.progress.finish_with_buffered_rows(
831 barrier.epoch,
832 total_snapshot_processed_rows,
833 total_snapshot_processed_rows,
834 );
835
836 Self::persist_backfill_state(&mut progress_table, &backfill_state).await?;
838 progress_table
839 .commit_assert_no_update_vnode_bitmap(barrier.epoch)
840 .await?;
841
842 yield Message::Barrier(barrier);
843 break; }
845 Message::Chunk(chunk) => {
846 yield Message::Chunk(chunk);
848 }
849 Message::Watermark(watermark) => {
850 yield Message::Watermark(watermark);
852 }
853 }
854 }
855 }
856
857 #[for_await]
859 for msg in upstream {
860 let msg = msg?;
861
862 match msg {
863 Message::Barrier(barrier) => {
864 barrier.assume_no_update_vnode_bitmap(self.actor_id)?;
865
866 state_table
868 .commit_assert_no_update_vnode_bitmap(barrier.epoch)
869 .await?;
870 progress_table
871 .commit_assert_no_update_vnode_bitmap(barrier.epoch)
872 .await?;
873 if report_finished_on_first_barrier {
874 self.progress.finish_with_buffered_rows(
876 barrier.epoch,
877 backfill_state.total_snapshot_rows,
878 backfill_state.total_snapshot_rows,
879 );
880 report_finished_on_first_barrier = false;
881 }
882 yield Message::Barrier(barrier);
883 }
884 _ => {
885 yield msg;
887 }
888 }
889 }
890 }
891}