Skip to main content

risingwave_stream/executor/
sync_kv_log_store.rs

1// Copyright 2025 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
15//! This contains the synced kv log store implementation.
16//! It's meant to buffer a large number of records emitted from upstream,
17//! to avoid overwhelming the downstream executor.
18//!
19//! The synced kv log store polls two futures:
20//!
21//! 1. Upstream: upstream message source
22//!
23//!   It will write stream messages to the log store buffer. e.g. `Message::Barrier`, `Message::Chunk`, ...
24//!   When writing a stream chunk, if the log store buffer is full, it will:
25//!     a. Flush the buffer to the log store.
26//!     b. Convert the stream chunk into a reference (`LogStoreBufferItem::Flushed`)
27//!       which can read the corresponding chunks in the log store.
28//!       We will compact adjacent references,
29//!       so it can read multiple chunks if there's a build up.
30//!
31//!   On receiving barriers, it will:
32//!     a. Apply truncation to historical data in the logstore.
33//!     b. Flush and checkpoint the logstore data.
34//!
35//! 2. State store + buffer + recently flushed chunks: the storage components of the logstore.
36//!
37//!   It will read all historical data from the logstore first. This can be done just by
38//!   constructing a state store stream, which will read all data until the latest epoch.
39//!   This is a static snapshot of data.
40//!   For any subsequently flushed chunks, we will read them via
41//!   `flushed_chunk_future`. See the next paragraph below.
42//!
43//!   We will next read `flushed_chunk_future` (if there's one pre-existing one), see below for how
44//!   it's constructed, what it is.
45//!
46//!   Finally we will pop the earliest item in the buffer.
47//!   - If it's a chunk yield it.
48//!   - If it's a watermark yield it.
49//!   - If it's a flushed chunk reference (`LogStoreBufferItem::Flushed`),
50//!     we will read the corresponding chunks in the log store.
51//!     This is done by constructing a `flushed_chunk_future` which will read the log store
52//!     using the `seq_id`.
53//!   - Barrier,
54//!     because they are directly propagated from the upstream when polling it.
55//!
56//! TODO(kwannoel):
57//! - [] Add dedicated metrics for sync log store, namespace according to the upstream.
58//! - [] Add tests
59//! - [] Handle paused stream
60
61use std::collections::VecDeque;
62use std::future::pending;
63use std::mem::replace;
64use std::pin::Pin;
65
66use anyhow::anyhow;
67use futures::future::{BoxFuture, Either, select};
68use futures::stream::StreamFuture;
69use futures::{FutureExt, StreamExt, TryStreamExt};
70use futures_async_stream::try_stream;
71use risingwave_common::array::StreamChunk;
72use risingwave_common::bitmap::Bitmap;
73use risingwave_common::catalog::{TableId, TableOption};
74use risingwave_common::must_match;
75use risingwave_common::util::epoch::EpochPair;
76use risingwave_common_estimate_size::EstimateSize;
77use risingwave_connector::sink::log_store::{ChunkId, LogStoreResult};
78use risingwave_pb::id::FragmentId;
79use risingwave_storage::StateStore;
80use risingwave_storage::store::timeout_auto_rebuild::TimeoutAutoRebuildIter;
81use risingwave_storage::store::{
82    LocalStateStore, NewLocalOptions, OpConsistencyLevel, StateStoreRead,
83};
84use rw_futures_util::drop_either_future;
85use tokio::time::{Duration, Instant, Sleep, sleep_until};
86use tokio_stream::adapters::Peekable;
87
88use crate::common::log_store_impl::kv_log_store::buffer::LogStoreBufferItem;
89use crate::common::log_store_impl::kv_log_store::reader::LogStoreReadStateStreamRangeStart;
90use crate::common::log_store_impl::kv_log_store::serde::{
91    KvLogStoreItem, LogStoreItemMergeStream, LogStoreRowSerde,
92};
93use crate::common::log_store_impl::kv_log_store::state::{
94    LogStorePostSealCurrentEpoch, LogStoreReadState, LogStoreStateWriteChunkFuture,
95    LogStoreWriteState, new_log_store_state,
96};
97use crate::common::log_store_impl::kv_log_store::{
98    Epoch, FIRST_SEQ_ID, FlushInfo, LogStoreVnodeProgress, SeqId,
99};
100use crate::executor::prelude::*;
101use crate::executor::sync_kv_log_store::metrics::SyncedKvLogStoreMetrics;
102use crate::executor::{
103    Barrier, BoxedMessageStream, Message, StreamExecutorError, StreamExecutorResult, Watermark,
104};
105
106pub mod metrics {
107    use risingwave_common::id::FragmentId;
108    use risingwave_common::metrics::{LabelGuardedIntCounter, LabelGuardedIntGauge};
109
110    use crate::common::log_store_impl::kv_log_store::KvLogStoreReadMetrics;
111    use crate::executor::monitor::StreamingMetrics;
112    use crate::task::ActorId;
113
114    #[derive(Clone)]
115    pub struct SyncedKvLogStoreMetrics {
116        // state of the log store
117        pub unclean_state: LabelGuardedIntCounter,
118        pub clean_state: LabelGuardedIntCounter,
119        pub wait_next_poll_ns: LabelGuardedIntCounter,
120
121        // Write metrics
122        pub storage_write_count: LabelGuardedIntCounter,
123        pub storage_write_size: LabelGuardedIntCounter,
124        pub pause_duration_ns: LabelGuardedIntCounter,
125
126        // Buffer metrics
127        pub buffer_unconsumed_item_count: LabelGuardedIntGauge,
128        pub buffer_unconsumed_row_count: LabelGuardedIntGauge,
129        pub buffer_unconsumed_epoch_count: LabelGuardedIntGauge,
130        pub buffer_unconsumed_min_epoch: LabelGuardedIntGauge,
131        pub buffer_memory_bytes: LabelGuardedIntGauge,
132        pub buffer_read_count: LabelGuardedIntCounter,
133        pub buffer_read_size: LabelGuardedIntCounter,
134
135        // Read metrics
136        pub total_read_count: LabelGuardedIntCounter,
137        pub total_read_size: LabelGuardedIntCounter,
138        pub persistent_log_read_metrics: KvLogStoreReadMetrics,
139        pub flushed_buffer_read_metrics: KvLogStoreReadMetrics,
140    }
141
142    impl SyncedKvLogStoreMetrics {
143        /// `id`: refers to a unique way to identify the logstore. This can be the sink id,
144        ///       or for joins, it can be the `fragment_id`.
145        /// `name`: refers to the MV / Sink that the log store is associated with.
146        /// `target`: refers to the target of the log store,
147        ///           for instance `MySql` Sink, PG sink, etc...
148        ///           or unaligned join.
149        pub(crate) fn new(
150            metrics: &StreamingMetrics,
151            actor_id: ActorId,
152            fragment_id: FragmentId,
153            name: &str,
154            target: &'static str,
155        ) -> Self {
156            let actor_id_str = actor_id.to_string();
157            let fragment_id_str = fragment_id.to_string();
158            let labels = &[&actor_id_str, target, &fragment_id_str, name];
159
160            let unclean_state = metrics.sync_kv_log_store_state.with_guarded_label_values(&[
161                "dirty",
162                &actor_id_str,
163                target,
164                &fragment_id_str,
165                name,
166            ]);
167            let clean_state = metrics.sync_kv_log_store_state.with_guarded_label_values(&[
168                "clean",
169                &actor_id_str,
170                target,
171                &fragment_id_str,
172                name,
173            ]);
174            let wait_next_poll_ns = metrics
175                .sync_kv_log_store_wait_next_poll_ns
176                .with_guarded_label_values(labels);
177
178            let storage_write_size = metrics
179                .sync_kv_log_store_storage_write_size
180                .with_guarded_label_values(labels);
181            let storage_write_count = metrics
182                .sync_kv_log_store_storage_write_count
183                .with_guarded_label_values(labels);
184            let storage_pause_duration_ns = metrics
185                .sync_kv_log_store_write_pause_duration_ns
186                .with_guarded_label_values(labels);
187
188            let buffer_unconsumed_item_count = metrics
189                .sync_kv_log_store_buffer_unconsumed_item_count
190                .with_guarded_label_values(labels);
191            let buffer_unconsumed_row_count = metrics
192                .sync_kv_log_store_buffer_unconsumed_row_count
193                .with_guarded_label_values(labels);
194            let buffer_unconsumed_epoch_count = metrics
195                .sync_kv_log_store_buffer_unconsumed_epoch_count
196                .with_guarded_label_values(labels);
197            let buffer_unconsumed_min_epoch = metrics
198                .sync_kv_log_store_buffer_unconsumed_min_epoch
199                .with_guarded_label_values(labels);
200            let buffer_memory_bytes = metrics
201                .sync_kv_log_store_buffer_memory_bytes
202                .with_guarded_label_values(labels);
203            let buffer_read_count = metrics
204                .sync_kv_log_store_read_count
205                .with_guarded_label_values(&[
206                    "buffer",
207                    &actor_id_str,
208                    target,
209                    &fragment_id_str,
210                    name,
211                ]);
212
213            let buffer_read_size = metrics
214                .sync_kv_log_store_read_size
215                .with_guarded_label_values(&[
216                    "buffer",
217                    &actor_id_str,
218                    target,
219                    &fragment_id_str,
220                    name,
221                ]);
222
223            let total_read_count = metrics
224                .sync_kv_log_store_read_count
225                .with_guarded_label_values(&[
226                    "total",
227                    &actor_id_str,
228                    target,
229                    &fragment_id_str,
230                    name,
231                ]);
232
233            let total_read_size = metrics
234                .sync_kv_log_store_read_size
235                .with_guarded_label_values(&[
236                    "total",
237                    &actor_id_str,
238                    target,
239                    &fragment_id_str,
240                    name,
241                ]);
242
243            const READ_PERSISTENT_LOG: &str = "persistent_log";
244            const READ_FLUSHED_BUFFER: &str = "flushed_buffer";
245
246            let persistent_log_read_size = metrics
247                .sync_kv_log_store_read_size
248                .with_guarded_label_values(&[
249                    READ_PERSISTENT_LOG,
250                    &actor_id_str,
251                    target,
252                    &fragment_id_str,
253                    name,
254                ]);
255
256            let persistent_log_read_count = metrics
257                .sync_kv_log_store_read_count
258                .with_guarded_label_values(&[
259                    READ_PERSISTENT_LOG,
260                    &actor_id_str,
261                    target,
262                    &fragment_id_str,
263                    name,
264                ]);
265
266            let flushed_buffer_read_size = metrics
267                .sync_kv_log_store_read_size
268                .with_guarded_label_values(&[
269                    READ_FLUSHED_BUFFER,
270                    &actor_id_str,
271                    target,
272                    &fragment_id_str,
273                    name,
274                ]);
275
276            let flushed_buffer_read_count = metrics
277                .sync_kv_log_store_read_count
278                .with_guarded_label_values(&[
279                    READ_FLUSHED_BUFFER,
280                    &actor_id_str,
281                    target,
282                    &fragment_id_str,
283                    name,
284                ]);
285
286            Self {
287                unclean_state,
288                clean_state,
289                wait_next_poll_ns,
290                storage_write_size,
291                storage_write_count,
292                pause_duration_ns: storage_pause_duration_ns,
293                buffer_unconsumed_item_count,
294                buffer_unconsumed_row_count,
295                buffer_unconsumed_epoch_count,
296                buffer_unconsumed_min_epoch,
297                buffer_memory_bytes,
298                buffer_read_count,
299                buffer_read_size,
300                total_read_count,
301                total_read_size,
302                persistent_log_read_metrics: KvLogStoreReadMetrics {
303                    storage_read_size: persistent_log_read_size,
304                    storage_read_count: persistent_log_read_count,
305                },
306                flushed_buffer_read_metrics: KvLogStoreReadMetrics {
307                    storage_read_count: flushed_buffer_read_count,
308                    storage_read_size: flushed_buffer_read_size,
309                },
310            }
311        }
312
313        #[cfg(test)]
314        pub(crate) fn for_test() -> Self {
315            SyncedKvLogStoreMetrics {
316                unclean_state: LabelGuardedIntCounter::test_int_counter::<5>(),
317                clean_state: LabelGuardedIntCounter::test_int_counter::<5>(),
318                wait_next_poll_ns: LabelGuardedIntCounter::test_int_counter::<4>(),
319                storage_write_count: LabelGuardedIntCounter::test_int_counter::<4>(),
320                storage_write_size: LabelGuardedIntCounter::test_int_counter::<4>(),
321                pause_duration_ns: LabelGuardedIntCounter::test_int_counter::<4>(),
322                buffer_unconsumed_item_count: LabelGuardedIntGauge::test_int_gauge::<4>(),
323                buffer_unconsumed_row_count: LabelGuardedIntGauge::test_int_gauge::<4>(),
324                buffer_unconsumed_epoch_count: LabelGuardedIntGauge::test_int_gauge::<4>(),
325                buffer_unconsumed_min_epoch: LabelGuardedIntGauge::test_int_gauge::<4>(),
326                buffer_memory_bytes: LabelGuardedIntGauge::test_int_gauge::<4>(),
327                buffer_read_count: LabelGuardedIntCounter::test_int_counter::<5>(),
328                buffer_read_size: LabelGuardedIntCounter::test_int_counter::<5>(),
329                total_read_count: LabelGuardedIntCounter::test_int_counter::<5>(),
330                total_read_size: LabelGuardedIntCounter::test_int_counter::<5>(),
331                persistent_log_read_metrics: KvLogStoreReadMetrics::for_test(),
332                flushed_buffer_read_metrics: KvLogStoreReadMetrics::for_test(),
333            }
334        }
335    }
336}
337
338pub(crate) type ReadFlushedChunkFuture =
339    BoxFuture<'static, LogStoreResult<(ChunkId, StreamChunk, Epoch)>>;
340
341pub struct SyncKvLogStoreContext<S: StateStore> {
342    pub table_id: TableId,
343    pub fragment_id: FragmentId,
344    pub metrics: SyncedKvLogStoreMetrics,
345    pub serde: LogStoreRowSerde,
346    pub state_store: S,
347    pub max_buffer_size: usize,
348    pub chunk_size: usize,
349    pub pause_duration_ms: Duration,
350    pub aligned: bool,
351}
352
353pub struct SyncedKvLogStoreExecutor<S: StateStore> {
354    actor_context: ActorContextRef,
355    upstream: Executor,
356    logstore_context: SyncKvLogStoreContext<S>,
357}
358
359// Stream interface
360impl<S: StateStore> SyncedKvLogStoreExecutor<S> {
361    #[expect(clippy::too_many_arguments)]
362    pub(crate) fn new(
363        actor_context: ActorContextRef,
364        table_id: TableId,
365        metrics: SyncedKvLogStoreMetrics,
366        serde: LogStoreRowSerde,
367        state_store: S,
368        buffer_size: usize,
369        chunk_size: usize,
370        upstream: Executor,
371        pause_duration_ms: Duration,
372        aligned: bool,
373    ) -> Self {
374        let logstore_context = SyncKvLogStoreContext {
375            table_id,
376            fragment_id: actor_context.fragment_id,
377            metrics,
378            serde,
379            state_store,
380            max_buffer_size: buffer_size,
381            chunk_size,
382            pause_duration_ms,
383            aligned,
384        };
385        Self {
386            actor_context,
387            upstream,
388            logstore_context,
389        }
390    }
391}
392
393pub(crate) struct FlushedChunkInfo {
394    epoch: u64,
395    start_seq_id: SeqId,
396    end_seq_id: SeqId,
397    flush_info: FlushInfo,
398    vnode_bitmap: Bitmap,
399}
400
401pub(crate) enum WriteFuture<S: LocalStateStore> {
402    /// We trigger a brief pause to let the `ReadFuture` be polled in the following scenarios:
403    /// - When seeing an upstream data chunk, when the buffer becomes full, and the state is clean.
404    /// - When seeing a checkpoint barrier, when the buffer is not empty, and the state is clean.
405    ///
406    /// On pausing, we will transition to a dirty state.
407    ///
408    /// We trigger resume to let the `ReadFuture` to be polled in the following scenarios:
409    /// - After the pause duration.
410    /// - After the read future consumes a chunk.
411    Paused {
412        start_instant: Instant,
413        sleep_future: Option<Pin<Box<Sleep>>>,
414        barrier: Barrier,
415        stream: BoxedMessageStream,
416        write_state: LogStoreWriteState<S>, // Just used to hold the state
417    },
418    ReceiveFromUpstream {
419        future: StreamFuture<BoxedMessageStream>,
420        write_state: LogStoreWriteState<S>,
421    },
422    FlushingChunk {
423        epoch: u64,
424        start_seq_id: SeqId,
425        end_seq_id: SeqId,
426        future: Pin<Box<LogStoreStateWriteChunkFuture<S>>>,
427        stream: BoxedMessageStream,
428    },
429    Empty,
430}
431
432pub(crate) enum WriteFutureEvent {
433    UpstreamMessageReceived(Message),
434    ChunkFlushed(FlushedChunkInfo),
435}
436
437impl<S: LocalStateStore> WriteFuture<S> {
438    fn flush_chunk(
439        stream: BoxedMessageStream,
440        write_state: LogStoreWriteState<S>,
441        chunk: StreamChunk,
442        epoch: u64,
443        start_seq_id: SeqId,
444        end_seq_id: SeqId,
445    ) -> Self {
446        tracing::trace!(
447            start_seq_id,
448            end_seq_id,
449            epoch,
450            cardinality = chunk.cardinality(),
451            "write_future: flushing chunk"
452        );
453        Self::FlushingChunk {
454            epoch,
455            start_seq_id,
456            end_seq_id,
457            future: Box::pin(write_state.into_write_chunk_future(
458                chunk,
459                epoch,
460                start_seq_id,
461                end_seq_id,
462            )),
463            stream,
464        }
465    }
466
467    pub(crate) fn receive_from_upstream(
468        stream: BoxedMessageStream,
469        write_state: LogStoreWriteState<S>,
470    ) -> Self {
471        Self::ReceiveFromUpstream {
472            future: stream.into_future(),
473            write_state,
474        }
475    }
476
477    pub(crate) fn paused(
478        duration: Duration,
479        barrier: Barrier,
480        stream: BoxedMessageStream,
481        write_state: LogStoreWriteState<S>,
482    ) -> Self {
483        let now = Instant::now();
484        tracing::trace!(?now, ?duration, "write_future_pause");
485        Self::Paused {
486            start_instant: now,
487            sleep_future: Some(Box::pin(sleep_until(now + duration))),
488            barrier,
489            stream,
490            write_state,
491        }
492    }
493
494    pub(crate) async fn next_event(
495        &mut self,
496        metrics: &SyncedKvLogStoreMetrics,
497    ) -> StreamExecutorResult<(BoxedMessageStream, LogStoreWriteState<S>, WriteFutureEvent)> {
498        match self {
499            WriteFuture::Paused {
500                start_instant,
501                sleep_future,
502                ..
503            } => {
504                if let Some(sleep_future) = sleep_future {
505                    sleep_future.await;
506                    metrics
507                        .pause_duration_ns
508                        .inc_by(start_instant.elapsed().as_nanos() as _);
509                    tracing::trace!("resuming write future");
510                }
511                must_match!(replace(self, WriteFuture::Empty), WriteFuture::Paused { stream, write_state, barrier, .. } => {
512                    Ok((stream, write_state, WriteFutureEvent::UpstreamMessageReceived(Message::Barrier(barrier))))
513                })
514            }
515            WriteFuture::ReceiveFromUpstream { future, .. } => {
516                let (opt, stream) = future.await;
517                must_match!(replace(self, WriteFuture::Empty), WriteFuture::ReceiveFromUpstream { write_state, .. } => {
518                    opt
519                    .ok_or_else(|| anyhow!("end of upstream input").into())
520                    .and_then(|result| result.map(|item| {
521                        (stream, write_state, WriteFutureEvent::UpstreamMessageReceived(item))
522                    }))
523                })
524            }
525            WriteFuture::FlushingChunk { future, .. } => {
526                let (write_state, result) = future.await;
527                let result = must_match!(replace(self, WriteFuture::Empty), WriteFuture::FlushingChunk { epoch, start_seq_id, end_seq_id, stream, ..  } => {
528                    result.map(|(flush_info, vnode_bitmap)| {
529                        (stream, write_state, WriteFutureEvent::ChunkFlushed(FlushedChunkInfo {
530                            epoch,
531                            start_seq_id,
532                            end_seq_id,
533                            flush_info,
534                            vnode_bitmap,
535                        }))
536                    })
537                });
538                result.map_err(Into::into)
539            }
540            WriteFuture::Empty => {
541                unreachable!("should not be polled after ready")
542            }
543        }
544    }
545}
546
547pub(crate) type LocalLogStoreReadState<S> =
548    LogStoreReadState<<<S as StateStore>::Local as LocalStateStore>::FlushedSnapshotReader>;
549pub(crate) type LocalLogStoreWriteState<S> = LogStoreWriteState<<S as StateStore>::Local>;
550
551// Stream interface
552impl<S: StateStore> SyncedKvLogStoreExecutor<S> {
553    pub(crate) async fn init_local_log_store_state(
554        context: &SyncKvLogStoreContext<S>,
555        first_write_epoch: EpochPair,
556    ) -> StreamExecutorResult<(LocalLogStoreReadState<S>, LocalLogStoreWriteState<S>)> {
557        let local_state_store = context
558            .state_store
559            .new_local(NewLocalOptions {
560                table_id: context.table_id,
561                fragment_id: context.fragment_id,
562                op_consistency_level: OpConsistencyLevel::Inconsistent,
563                table_option: TableOption {
564                    retention_seconds: None,
565                },
566                is_replicated: false,
567                vnodes: context.serde.vnodes().clone(),
568                upload_on_flush: false,
569            })
570            .await;
571
572        let (read_state, mut initial_write_state) = new_log_store_state(
573            context.table_id,
574            local_state_store,
575            context.serde.clone(),
576            context.chunk_size,
577        );
578        initial_write_state.init(first_write_epoch).await?;
579        Ok((read_state, initial_write_state))
580    }
581
582    #[try_stream(ok = Message, error = StreamExecutorError)]
583    pub(crate) async fn aligned_message_stream(
584        actor_id: ActorId,
585        input: BoxedMessageStream,
586        read_state: LocalLogStoreReadState<S>,
587        mut initial_write_state: LocalLogStoreWriteState<S>,
588        metrics: SyncedKvLogStoreMetrics,
589        initial_write_epoch: EpochPair,
590    ) {
591        tracing::info!("aligned mode");
592        // We want to realign the buffer and the stream.
593        // We just block the upstream input stream,
594        // and wait until the persisted logstore is empty.
595        // Then after that we can consume the input stream.
596        let log_store_stream = read_state
597            .read_persisted_log_store(
598                metrics.persistent_log_read_metrics.clone(),
599                initial_write_epoch.curr,
600                LogStoreReadStateStreamRangeStart::Unbounded,
601            )
602            .await?;
603
604        #[for_await]
605        for message in log_store_stream {
606            let (_epoch, message) = message?;
607            match message {
608                KvLogStoreItem::Barrier { .. } => {
609                    continue;
610                }
611                KvLogStoreItem::StreamChunk { chunk, .. } => {
612                    yield Message::Chunk(chunk);
613                }
614            }
615        }
616
617        let mut realigned_logstore = false;
618
619        #[for_await]
620        for message in input {
621            match message? {
622                Message::Barrier(barrier) => {
623                    let is_checkpoint = barrier.is_checkpoint();
624                    let mut progress = LogStoreVnodeProgress::None;
625                    progress.apply_aligned(read_state.vnodes().clone(), barrier.epoch.prev, None);
626                    // Truncate the logstore.
627                    let post_seal =
628                        initial_write_state.seal_current_epoch(barrier.epoch.curr, progress.take());
629                    barrier.assume_no_update_vnode_bitmap(actor_id)?;
630                    yield Message::Barrier(barrier);
631                    post_seal.post_yield_barrier(None).await?;
632                    if !realigned_logstore && is_checkpoint {
633                        realigned_logstore = true;
634                        tracing::info!("realigned logstore");
635                    }
636                }
637                Message::Chunk(chunk) => {
638                    yield Message::Chunk(chunk);
639                }
640                Message::Watermark(watermark) => {
641                    yield Message::Watermark(watermark);
642                }
643            }
644        }
645    }
646
647    pub(crate) fn apply_pause_resume_mutation(barrier: &Barrier, pause_stream: &mut bool) {
648        if let Some(mutation) = barrier.mutation.as_deref() {
649            match mutation {
650                Mutation::Pause => {
651                    *pause_stream = true;
652                }
653                Mutation::Resume => {
654                    *pause_stream = false;
655                }
656                _ => {}
657            }
658        }
659    }
660
661    pub(crate) fn process_upstream_chunk(
662        seq_id: SeqId,
663        stream: BoxedMessageStream,
664        write_state: LogStoreWriteState<S::Local>,
665        chunk: StreamChunk,
666        buffer: &mut SyncedLogStoreBuffer,
667    ) -> (SeqId, WriteFuture<S::Local>) {
668        let cardinality = chunk.cardinality();
669        if cardinality == 0 {
670            tracing::warn!(
671                epoch = write_state.epoch().curr,
672                "received empty chunk (cardinality=0), skipping"
673            );
674            return (
675                seq_id,
676                WriteFuture::receive_from_upstream(stream, write_state),
677            );
678        }
679
680        let start_seq_id = seq_id;
681        let new_seq_id = seq_id + cardinality as SeqId;
682        let end_seq_id = new_seq_id - 1;
683        let epoch = write_state.epoch().curr;
684        tracing::trace!(
685            start_seq_id,
686            end_seq_id,
687            new_seq_id,
688            epoch,
689            cardinality,
690            "received chunk"
691        );
692        let next_write_future = if let Some(chunk_to_flush) =
693            buffer.add_or_flush_chunk(start_seq_id, end_seq_id, chunk, epoch)
694        {
695            WriteFuture::flush_chunk(
696                stream,
697                write_state,
698                chunk_to_flush,
699                epoch,
700                start_seq_id,
701                end_seq_id,
702            )
703        } else {
704            WriteFuture::receive_from_upstream(stream, write_state)
705        };
706        (new_seq_id, next_write_future)
707    }
708
709    pub(crate) fn process_flushed_chunk(
710        stream: BoxedMessageStream,
711        write_state: LogStoreWriteState<S::Local>,
712        info: FlushedChunkInfo,
713        buffer: &mut SyncedLogStoreBuffer,
714        metrics: &SyncedKvLogStoreMetrics,
715    ) -> WriteFuture<S::Local> {
716        buffer.add_flushed_item_to_buffer(
717            info.start_seq_id,
718            info.end_seq_id,
719            info.vnode_bitmap,
720            info.epoch,
721        );
722        metrics
723            .storage_write_count
724            .inc_by(info.flush_info.flush_count as _);
725        metrics
726            .storage_write_size
727            .inc_by(info.flush_info.flush_size as _);
728        WriteFuture::receive_from_upstream(stream, write_state)
729    }
730
731    #[try_stream(ok= Message, error = StreamExecutorError)]
732    pub async fn execute_monitored(self) {
733        let wait_next_poll_ns = self.logstore_context.metrics.wait_next_poll_ns.clone();
734        #[for_await]
735        for message in self.execute_inner() {
736            let current_time = Instant::now();
737            yield message?;
738            wait_next_poll_ns.inc_by(current_time.elapsed().as_nanos() as _);
739        }
740    }
741
742    #[try_stream(ok = Message, error = StreamExecutorError)]
743    async fn execute_inner(self) {
744        let mut input = self.upstream.execute();
745
746        // init first epoch + local state store
747        let first_barrier = expect_first_barrier(&mut input).await?;
748        let first_write_epoch = first_barrier.epoch;
749        yield Message::Barrier(first_barrier.clone());
750
751        let (read_state, initial_write_state) =
752            Self::init_local_log_store_state(&self.logstore_context, first_write_epoch).await?;
753
754        let mut pause_stream = first_barrier.is_pause_on_startup();
755        let initial_write_epoch = first_write_epoch;
756
757        if self.logstore_context.aligned {
758            let aligned_stream = Self::aligned_message_stream(
759                self.actor_context.id,
760                input,
761                read_state,
762                initial_write_state,
763                self.logstore_context.metrics.clone(),
764                initial_write_epoch,
765            );
766            #[for_await]
767            for message in aligned_stream {
768                yield message?;
769            }
770            return Ok(());
771        }
772
773        let mut seq_id = FIRST_SEQ_ID;
774        let mut buffer = SyncedLogStoreBuffer::new(
775            self.logstore_context.max_buffer_size,
776            self.logstore_context.chunk_size,
777            &self.logstore_context.metrics,
778        );
779
780        let log_store_stream = read_state
781            .read_persisted_log_store(
782                self.logstore_context
783                    .metrics
784                    .persistent_log_read_metrics
785                    .clone(),
786                initial_write_epoch.curr,
787                LogStoreReadStateStreamRangeStart::Unbounded,
788            )
789            .await?;
790
791        let mut log_store_stream = tokio_stream::StreamExt::peekable(log_store_stream);
792        let mut clean_state = log_store_stream.peek().await.is_none();
793        tracing::trace!(?clean_state);
794
795        let mut read_future_state = ReadFuture::ReadingPersistedStream(log_store_stream);
796
797        let mut write_future_state = WriteFuture::receive_from_upstream(input, initial_write_state);
798
799        let mut progress = LogStoreVnodeProgress::None;
800
801        loop {
802            let select_result = {
803                let read_future = async {
804                    if pause_stream {
805                        pending().await
806                    } else {
807                        read_future_state
808                            .next_message(
809                                &mut progress,
810                                &read_state,
811                                &mut buffer,
812                                &self.logstore_context.metrics,
813                            )
814                            .await
815                    }
816                };
817                pin_mut!(read_future);
818                let write_future = write_future_state.next_event(&self.logstore_context.metrics);
819                pin_mut!(write_future);
820                let output = select(write_future, read_future).await;
821                drop_either_future(output)
822            };
823            match select_result {
824                Either::Left(result) => {
825                    // drop the future to ensure that the future must be reset later
826                    drop(write_future_state);
827                    let (stream, mut write_state, either) = result?;
828                    match either {
829                        WriteFutureEvent::UpstreamMessageReceived(msg) => match msg {
830                            Message::Barrier(barrier) => {
831                                if clean_state && barrier.kind.is_checkpoint() && !buffer.is_empty()
832                                {
833                                    write_future_state = WriteFuture::paused(
834                                        self.logstore_context.pause_duration_ms,
835                                        barrier,
836                                        stream,
837                                        write_state,
838                                    );
839                                    clean_state = false;
840                                    self.logstore_context.metrics.unclean_state.inc();
841                                } else {
842                                    Self::apply_pause_resume_mutation(&barrier, &mut pause_stream);
843                                    let write_state_post_write_barrier = Self::write_barrier(
844                                        self.actor_context.id,
845                                        &mut write_state,
846                                        barrier.clone(),
847                                        &self.logstore_context.metrics,
848                                        progress.take(),
849                                        &mut buffer,
850                                    )
851                                    .await?;
852                                    seq_id = FIRST_SEQ_ID;
853                                    barrier.assume_no_update_vnode_bitmap(self.actor_context.id)?;
854
855                                    yield Message::Barrier(barrier);
856
857                                    write_state_post_write_barrier
858                                        .post_yield_barrier(None)
859                                        .await?;
860
861                                    write_future_state =
862                                        WriteFuture::receive_from_upstream(stream, write_state);
863                                }
864                            }
865                            Message::Chunk(chunk) => {
866                                let (new_seq_id, next_write_future) = Self::process_upstream_chunk(
867                                    seq_id,
868                                    stream,
869                                    write_state,
870                                    chunk,
871                                    &mut buffer,
872                                );
873                                seq_id = new_seq_id;
874                                write_future_state = next_write_future;
875                            }
876                            Message::Watermark(watermark) => {
877                                buffer.add_watermark(write_state.epoch().curr, watermark);
878                                write_future_state =
879                                    WriteFuture::receive_from_upstream(stream, write_state);
880                            }
881                        },
882                        WriteFutureEvent::ChunkFlushed(info) => {
883                            write_future_state = Self::process_flushed_chunk(
884                                stream,
885                                write_state,
886                                info,
887                                &mut buffer,
888                                &self.logstore_context.metrics,
889                            );
890                        }
891                    }
892                }
893                Either::Right(result) => {
894                    let clean_state_reached = read_future_state.mark_clean_state(
895                        &mut clean_state,
896                        &buffer,
897                        &self.logstore_context.metrics,
898                    );
899
900                    if clean_state_reached {
901                        // Let write future resume immediately
902                        if let WriteFuture::Paused { sleep_future, .. } = &mut write_future_state {
903                            tracing::trace!("resuming paused future");
904                            assert!(buffer.has_available_capacity());
905                            *sleep_future = None;
906                        }
907                    }
908                    let message = result?;
909                    if let Message::Chunk(chunk) = &message {
910                        self.logstore_context
911                            .metrics
912                            .total_read_count
913                            .inc_by(chunk.cardinality() as _);
914                    }
915
916                    yield message;
917                }
918            }
919        }
920    }
921}
922
923pub(crate) type PersistedStream<S> =
924    Peekable<Pin<Box<LogStoreItemMergeStream<TimeoutAutoRebuildIter<S>>>>>;
925
926pub(crate) enum ReadFuture<S: StateStoreRead> {
927    ReadingPersistedStream(PersistedStream<S>),
928    ReadingFlushedChunk {
929        future: ReadFlushedChunkFuture,
930        end_seq_id: SeqId,
931    },
932    Idle,
933}
934
935// Read methods
936impl<S: StateStoreRead> ReadFuture<S> {
937    pub(crate) fn mark_clean_state(
938        &self,
939        clean_state: &mut bool,
940        buffer: &SyncedLogStoreBuffer,
941        metrics: &SyncedKvLogStoreMetrics,
942    ) -> bool {
943        if !*clean_state && matches!(self, ReadFuture::Idle) && buffer.is_empty() {
944            *clean_state = true;
945            metrics.clean_state.inc();
946            true
947        } else {
948            false
949        }
950    }
951
952    pub(crate) async fn next_message(
953        &mut self,
954        progress: &mut LogStoreVnodeProgress,
955        read_state: &LogStoreReadState<S>,
956        buffer: &mut SyncedLogStoreBuffer,
957        metrics: &SyncedKvLogStoreMetrics,
958    ) -> StreamExecutorResult<Message> {
959        match self {
960            ReadFuture::ReadingPersistedStream(stream) => {
961                while let Some((epoch, item)) = stream.try_next().await? {
962                    match item {
963                        KvLogStoreItem::Barrier { vnodes, .. } => {
964                            tracing::trace!(epoch, "read logstore barrier");
965                            // update the progress
966                            progress.apply_aligned(vnodes, epoch, None);
967                            continue;
968                        }
969                        KvLogStoreItem::StreamChunk {
970                            chunk,
971                            progress: chunk_progress,
972                        } => {
973                            tracing::trace!("read logstore chunk of size: {}", chunk.cardinality());
974                            progress.apply_per_vnode(epoch, chunk_progress);
975                            return Ok(Message::Chunk(chunk));
976                        }
977                    }
978                }
979                *self = ReadFuture::Idle;
980            }
981            ReadFuture::ReadingFlushedChunk { .. } | ReadFuture::Idle => {}
982        }
983        match self {
984            ReadFuture::ReadingPersistedStream(_) => {
985                unreachable!("must have finished read persisted stream when reaching here")
986            }
987            ReadFuture::ReadingFlushedChunk { .. } => {}
988            ReadFuture::Idle => loop {
989                let Some((item_epoch, item)) = buffer.pop_front() else {
990                    return pending().await;
991                };
992                match item {
993                    SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::StreamChunk {
994                        chunk,
995                        start_seq_id,
996                        end_seq_id,
997                        flushed,
998                        ..
999                    }) => {
1000                        metrics.buffer_read_count.inc_by(chunk.cardinality() as _);
1001                        tracing::trace!(
1002                            start_seq_id,
1003                            end_seq_id,
1004                            flushed,
1005                            cardinality = chunk.cardinality(),
1006                            "read buffered chunk of size"
1007                        );
1008                        progress.apply_aligned(
1009                            read_state.vnodes().clone(),
1010                            item_epoch,
1011                            Some(end_seq_id),
1012                        );
1013                        return Ok(Message::Chunk(chunk));
1014                    }
1015                    SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Flushed {
1016                        vnode_bitmap,
1017                        start_seq_id,
1018                        end_seq_id,
1019                        chunk_id,
1020                    }) => {
1021                        tracing::trace!(start_seq_id, end_seq_id, chunk_id, "read flushed chunk");
1022                        let read_metrics = metrics.flushed_buffer_read_metrics.clone();
1023                        let future = read_state
1024                            .read_flushed_chunk(
1025                                vnode_bitmap,
1026                                chunk_id,
1027                                start_seq_id,
1028                                end_seq_id,
1029                                item_epoch,
1030                                read_metrics,
1031                            )
1032                            .boxed();
1033                        *self = ReadFuture::ReadingFlushedChunk { future, end_seq_id };
1034                        break;
1035                    }
1036                    SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Barrier { .. }) => {
1037                        tracing::trace!(item_epoch, "read buffer barrier");
1038                        progress.apply_aligned(read_state.vnodes().clone(), item_epoch, None);
1039                        continue;
1040                    }
1041                    SyncedLogStoreBufferItem::Watermark(watermark) => {
1042                        return Ok(Message::Watermark(watermark));
1043                    }
1044                }
1045            },
1046        }
1047
1048        let (future, end_seq_id) = match self {
1049            ReadFuture::ReadingPersistedStream(_) | ReadFuture::Idle => {
1050                unreachable!("should be at ReadingFlushedChunk")
1051            }
1052            ReadFuture::ReadingFlushedChunk { future, end_seq_id } => (future, *end_seq_id),
1053        };
1054
1055        let (_, chunk, epoch) = future.await?;
1056        progress.apply_aligned(read_state.vnodes().clone(), epoch, Some(end_seq_id));
1057        tracing::trace!(
1058            end_seq_id,
1059            "read flushed chunk of size: {}",
1060            chunk.cardinality()
1061        );
1062        *self = ReadFuture::Idle;
1063        Ok(Message::Chunk(chunk))
1064    }
1065}
1066
1067// Write methods
1068impl<S: StateStore> SyncedKvLogStoreExecutor<S> {
1069    pub(crate) async fn write_barrier<'a>(
1070        actor_id: ActorId,
1071        write_state: &'a mut LogStoreWriteState<S::Local>,
1072        barrier: Barrier,
1073        metrics: &SyncedKvLogStoreMetrics,
1074        progress: LogStoreVnodeProgress,
1075        buffer: &mut SyncedLogStoreBuffer,
1076    ) -> StreamExecutorResult<LogStorePostSealCurrentEpoch<'a, S::Local>> {
1077        tracing::trace!(%actor_id, ?progress, "applying truncation");
1078        // TODO(kwannoel): As an optimization we can also change flushed chunks to be flushed items
1079        // to reduce memory consumption of logstore.
1080
1081        let epoch = barrier.epoch.prev;
1082        let mut writer = write_state.start_writer(false);
1083        writer.write_barrier(epoch, barrier.is_checkpoint())?;
1084
1085        if barrier.is_checkpoint() {
1086            for (epoch, item) in buffer.buffer.iter_mut().rev() {
1087                match item {
1088                    SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::StreamChunk {
1089                        chunk,
1090                        start_seq_id,
1091                        end_seq_id,
1092                        flushed,
1093                        ..
1094                    }) => {
1095                        if !*flushed {
1096                            writer.write_chunk(chunk, *epoch, *start_seq_id, *end_seq_id)?;
1097                            *flushed = true;
1098                        } else {
1099                            break;
1100                        }
1101                    }
1102                    SyncedLogStoreBufferItem::LogStore(
1103                        LogStoreBufferItem::Flushed { .. } | LogStoreBufferItem::Barrier { .. },
1104                    )
1105                    | SyncedLogStoreBufferItem::Watermark(_) => {}
1106                }
1107            }
1108        }
1109
1110        // Apply truncation
1111        let (flush_info, _) = writer.finish().await?;
1112        metrics
1113            .storage_write_count
1114            .inc_by(flush_info.flush_count as _);
1115        metrics
1116            .storage_write_size
1117            .inc_by(flush_info.flush_size as _);
1118        let post_seal = write_state.seal_current_epoch(barrier.epoch.curr, progress);
1119
1120        // Add to buffer
1121        buffer.buffer.push_back((
1122            epoch,
1123            SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Barrier {
1124                is_checkpoint: barrier.is_checkpoint(),
1125                next_epoch: barrier.epoch.curr,
1126                schema_change: None,
1127                is_stop: false,
1128            }),
1129        ));
1130        buffer.next_chunk_id = 0;
1131        buffer.update_buffer_metrics();
1132
1133        Ok(post_seal)
1134    }
1135}
1136
1137#[derive(EstimateSize)]
1138enum SyncedLogStoreBufferItem {
1139    LogStore(LogStoreBufferItem),
1140    Watermark(Watermark),
1141}
1142
1143pub(crate) struct SyncedLogStoreBuffer {
1144    buffer: VecDeque<(u64, SyncedLogStoreBufferItem)>,
1145    current_size: usize,
1146    max_size: usize,
1147    max_chunk_size: usize,
1148    next_chunk_id: ChunkId,
1149    metrics: SyncedKvLogStoreMetrics,
1150    flushed_count: usize,
1151}
1152
1153impl SyncedLogStoreBuffer {
1154    pub fn new(
1155        max_buffer_size: usize,
1156        chunk_size: usize,
1157        metrics: &SyncedKvLogStoreMetrics,
1158    ) -> Self {
1159        SyncedLogStoreBuffer {
1160            buffer: VecDeque::new(),
1161            current_size: 0,
1162            max_size: max_buffer_size,
1163            max_chunk_size: chunk_size,
1164            next_chunk_id: 0,
1165            metrics: metrics.clone(),
1166            flushed_count: 0,
1167        }
1168    }
1169
1170    pub fn is_empty(&self) -> bool {
1171        self.current_size == 0
1172    }
1173
1174    pub fn has_available_capacity(&self) -> bool {
1175        self.current_size < self.max_size
1176    }
1177
1178    fn add_or_flush_chunk(
1179        &mut self,
1180        start_seq_id: SeqId,
1181        end_seq_id: SeqId,
1182        chunk: StreamChunk,
1183        epoch: u64,
1184    ) -> Option<StreamChunk> {
1185        let current_size = self.current_size;
1186        let chunk_size = chunk.cardinality();
1187
1188        tracing::trace!(
1189            current_size,
1190            chunk_size,
1191            max_size = self.max_size,
1192            "checking chunk size"
1193        );
1194        let should_flush_chunk = current_size + chunk_size > self.max_size;
1195        if should_flush_chunk {
1196            tracing::trace!(start_seq_id, end_seq_id, epoch, "flushing chunk",);
1197            Some(chunk)
1198        } else {
1199            tracing::trace!(start_seq_id, end_seq_id, epoch, "buffering chunk",);
1200            self.add_chunk_to_buffer(chunk, start_seq_id, end_seq_id, epoch);
1201            None
1202        }
1203    }
1204
1205    /// After flushing a chunk, we will preserve a `FlushedItem` inside the buffer.
1206    /// This doesn't contain any data, but it contains the metadata to read the flushed chunk.
1207    fn add_flushed_item_to_buffer(
1208        &mut self,
1209        start_seq_id: SeqId,
1210        end_seq_id: SeqId,
1211        new_vnode_bitmap: Bitmap,
1212        epoch: u64,
1213    ) {
1214        let new_chunk_size = (end_seq_id - start_seq_id + 1) as usize;
1215
1216        if let Some((
1217            item_epoch,
1218            SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Flushed {
1219                start_seq_id: prev_start_seq_id,
1220                end_seq_id: prev_end_seq_id,
1221                vnode_bitmap,
1222                ..
1223            }),
1224        )) = self.buffer.back_mut()
1225            && let flushed_chunk_size = (*prev_end_seq_id - *prev_start_seq_id + 1) as usize
1226            && let projected_flushed_chunk_size = flushed_chunk_size + new_chunk_size
1227            && projected_flushed_chunk_size <= self.max_chunk_size
1228        {
1229            assert!(
1230                *prev_end_seq_id < start_seq_id,
1231                "prev end_seq_id {} should be smaller than current start_seq_id {}",
1232                end_seq_id,
1233                start_seq_id
1234            );
1235            assert_eq!(
1236                epoch, *item_epoch,
1237                "epoch of newly added flushed item must be the same as the last flushed item"
1238            );
1239            *prev_end_seq_id = end_seq_id;
1240            *vnode_bitmap |= new_vnode_bitmap;
1241        } else {
1242            let chunk_id = self.next_chunk_id;
1243            self.next_chunk_id += 1;
1244            self.buffer.push_back((
1245                epoch,
1246                SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Flushed {
1247                    start_seq_id,
1248                    end_seq_id,
1249                    vnode_bitmap: new_vnode_bitmap,
1250                    chunk_id,
1251                }),
1252            ));
1253            self.flushed_count += 1;
1254            tracing::trace!(
1255                "adding flushed item to buffer: start_seq_id: {start_seq_id}, end_seq_id: {end_seq_id}, chunk_id: {chunk_id}"
1256            );
1257        }
1258        // FIXME(kwannoel): Seems these metrics are updated _after_ the flush info is reported.
1259        self.update_buffer_metrics();
1260    }
1261
1262    fn add_chunk_to_buffer(
1263        &mut self,
1264        chunk: StreamChunk,
1265        start_seq_id: SeqId,
1266        end_seq_id: SeqId,
1267        epoch: u64,
1268    ) {
1269        let chunk_id = self.next_chunk_id;
1270        self.next_chunk_id += 1;
1271        self.current_size += chunk.cardinality();
1272        self.buffer.push_back((
1273            epoch,
1274            SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::StreamChunk {
1275                chunk,
1276                start_seq_id,
1277                end_seq_id,
1278                flushed: false,
1279                chunk_id,
1280            }),
1281        ));
1282        self.update_buffer_metrics();
1283    }
1284
1285    pub(crate) fn add_watermark(&mut self, epoch: u64, watermark: Watermark) {
1286        self.buffer
1287            .push_back((epoch, SyncedLogStoreBufferItem::Watermark(watermark)));
1288        self.update_buffer_metrics();
1289    }
1290
1291    fn pop_front(&mut self) -> Option<(u64, SyncedLogStoreBufferItem)> {
1292        let item = self.buffer.pop_front();
1293        match &item {
1294            Some((_, SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Flushed { .. }))) => {
1295                self.flushed_count -= 1;
1296            }
1297            Some((
1298                _,
1299                SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::StreamChunk {
1300                    chunk, ..
1301                }),
1302            )) => {
1303                self.current_size -= chunk.cardinality();
1304            }
1305            _ => {}
1306        }
1307        self.update_buffer_metrics();
1308        item
1309    }
1310
1311    fn update_buffer_metrics(&self) {
1312        let mut epoch_count = 0;
1313        let mut row_count = 0;
1314        let mut memory_bytes = 0;
1315        for (_, item) in &self.buffer {
1316            match item {
1317                SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::StreamChunk {
1318                    chunk,
1319                    ..
1320                }) => {
1321                    row_count += chunk.cardinality();
1322                }
1323                SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Flushed {
1324                    start_seq_id,
1325                    end_seq_id,
1326                    ..
1327                }) => {
1328                    row_count += (end_seq_id - start_seq_id) as usize;
1329                }
1330                SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Barrier { .. }) => {
1331                    epoch_count += 1;
1332                }
1333                SyncedLogStoreBufferItem::Watermark(_) => {}
1334            }
1335            memory_bytes += item.estimated_size();
1336        }
1337        self.metrics.buffer_unconsumed_epoch_count.set(epoch_count);
1338        self.metrics.buffer_unconsumed_row_count.set(row_count as _);
1339        self.metrics
1340            .buffer_unconsumed_item_count
1341            .set(self.buffer.len() as _);
1342        self.metrics.buffer_unconsumed_min_epoch.set(
1343            self.buffer
1344                .front()
1345                .map(|(epoch, _)| *epoch)
1346                .unwrap_or_default() as _,
1347        );
1348        self.metrics.buffer_memory_bytes.set(memory_bytes as _);
1349    }
1350}
1351
1352impl<S> Execute for SyncedKvLogStoreExecutor<S>
1353where
1354    S: StateStore,
1355{
1356    fn execute(self: Box<Self>) -> BoxedMessageStream {
1357        self.execute_monitored().boxed()
1358    }
1359}
1360
1361#[cfg(test)]
1362mod tests {
1363    use itertools::Itertools;
1364    use pretty_assertions::assert_eq;
1365    use risingwave_common::catalog::Field;
1366    use risingwave_common::hash::VirtualNode;
1367    use risingwave_common::test_prelude::*;
1368    use risingwave_common::util::epoch::test_epoch;
1369    use risingwave_storage::memory::MemoryStateStore;
1370
1371    use super::*;
1372    use crate::assert_stream_chunk_eq;
1373    use crate::common::log_store_impl::kv_log_store::KV_LOG_STORE_V2_INFO;
1374    use crate::common::log_store_impl::kv_log_store::test_utils::{
1375        check_stream_chunk_eq, gen_test_log_store_table, test_payload_schema,
1376    };
1377    use crate::executor::sync_kv_log_store::metrics::SyncedKvLogStoreMetrics;
1378    use crate::executor::test_utils::MockSource;
1379
1380    fn init_logger() {
1381        let _ = tracing_subscriber::fmt()
1382            .with_env_filter(tracing_subscriber::EnvFilter::from_default_env())
1383            .with_ansi(false)
1384            .try_init();
1385    }
1386
1387    // test read/write buffer
1388    #[tokio::test]
1389    async fn test_read_write_buffer() {
1390        init_logger();
1391
1392        let pk_info = &KV_LOG_STORE_V2_INFO;
1393        let column_descs = test_payload_schema(pk_info);
1394        let fields = column_descs
1395            .into_iter()
1396            .map(|desc| Field::new(desc.name.clone(), desc.data_type))
1397            .collect_vec();
1398        let schema = Schema { fields };
1399        let stream_key = vec![0];
1400        let (mut tx, source) = MockSource::channel();
1401        let source = source.into_executor(schema.clone(), stream_key.clone());
1402
1403        let vnodes = Some(Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)));
1404
1405        let table = gen_test_log_store_table(pk_info);
1406
1407        let log_store_executor = SyncedKvLogStoreExecutor::new(
1408            ActorContext::for_test(123),
1409            table.id,
1410            SyncedKvLogStoreMetrics::for_test(),
1411            LogStoreRowSerde::new(&table, vnodes, pk_info),
1412            MemoryStateStore::new(),
1413            10,
1414            256,
1415            source,
1416            Duration::from_millis(256),
1417            false,
1418        )
1419        .boxed();
1420
1421        // Init
1422        tx.push_barrier(test_epoch(1), false);
1423
1424        let chunk_1 = StreamChunk::from_pretty(
1425            "  I   T
1426            +  5  10
1427            +  6  10
1428            +  8  10
1429            +  9  10
1430            +  10 11",
1431        );
1432
1433        let chunk_2 = StreamChunk::from_pretty(
1434            "   I   T
1435            -   5  10
1436            -   6  10
1437            -   8  10
1438            U-  9  10
1439            U+ 10  11",
1440        );
1441
1442        tx.push_chunk(chunk_1.clone());
1443        tx.push_int64_watermark(0, 7);
1444        tx.push_chunk(chunk_2.clone());
1445
1446        let mut stream = log_store_executor.execute();
1447
1448        match stream.next().await {
1449            Some(Ok(Message::Barrier(barrier))) => {
1450                assert_eq!(barrier.epoch.curr, test_epoch(1));
1451            }
1452            other => panic!("Expected a barrier message, got {:?}", other),
1453        }
1454
1455        match stream.next().await {
1456            Some(Ok(Message::Chunk(chunk))) => {
1457                assert_stream_chunk_eq!(chunk, chunk_1);
1458            }
1459            other => panic!("Expected a chunk message, got {:?}", other),
1460        }
1461
1462        match stream.next().await {
1463            Some(Ok(Message::Watermark(watermark))) => {
1464                assert_eq!(watermark, Watermark::new(0, DataType::Int64, 7_i64.into()));
1465            }
1466            other => panic!("Expected a watermark message, got {:?}", other),
1467        }
1468
1469        match stream.next().await {
1470            Some(Ok(Message::Chunk(chunk))) => {
1471                assert_stream_chunk_eq!(chunk, chunk_2);
1472            }
1473            other => panic!("Expected a chunk message, got {:?}", other),
1474        }
1475
1476        tx.push_barrier(test_epoch(2), false);
1477
1478        match stream.next().await {
1479            Some(Ok(Message::Barrier(barrier))) => {
1480                assert_eq!(barrier.epoch.curr, test_epoch(2));
1481            }
1482            other => panic!("Expected a barrier message, got {:?}", other),
1483        }
1484    }
1485
1486    // test barrier persisted read
1487    //
1488    // sequence of events (earliest -> latest):
1489    // barrier(1) -> chunk(1) -> watermark -> chunk(2) -> poll(4) items -> barrier(2)
1490    // * poll just means we read from the executor stream.
1491    #[tokio::test]
1492    async fn test_barrier_persisted_read() {
1493        init_logger();
1494
1495        let pk_info = &KV_LOG_STORE_V2_INFO;
1496        let column_descs = test_payload_schema(pk_info);
1497        let fields = column_descs
1498            .into_iter()
1499            .map(|desc| Field::new(desc.name.clone(), desc.data_type))
1500            .collect_vec();
1501        let schema = Schema { fields };
1502        let stream_key = vec![0];
1503        let (mut tx, source) = MockSource::channel();
1504        let source = source.into_executor(schema.clone(), stream_key.clone());
1505
1506        let vnodes = Some(Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)));
1507
1508        let table = gen_test_log_store_table(pk_info);
1509
1510        let log_store_executor = SyncedKvLogStoreExecutor::new(
1511            ActorContext::for_test(123),
1512            table.id,
1513            SyncedKvLogStoreMetrics::for_test(),
1514            LogStoreRowSerde::new(&table, vnodes, pk_info),
1515            MemoryStateStore::new(),
1516            10,
1517            256,
1518            source,
1519            Duration::from_millis(256),
1520            false,
1521        )
1522        .boxed();
1523
1524        // Init
1525        tx.push_barrier(test_epoch(1), false);
1526
1527        let chunk_1 = StreamChunk::from_pretty(
1528            "  I   T
1529            +  5  10
1530            +  6  10
1531            +  8  10
1532            +  9  10
1533            +  10 11",
1534        );
1535
1536        let chunk_2 = StreamChunk::from_pretty(
1537            "   I   T
1538            -   5  10
1539            -   6  10
1540            -   8  10
1541            U- 10  11
1542            U+ 10  10",
1543        );
1544
1545        tx.push_chunk(chunk_1.clone());
1546        tx.push_int64_watermark(0, 7);
1547        tx.push_chunk(chunk_2.clone());
1548
1549        tx.push_barrier(test_epoch(2), false);
1550
1551        let mut stream = log_store_executor.execute();
1552
1553        match stream.next().await {
1554            Some(Ok(Message::Barrier(barrier))) => {
1555                assert_eq!(barrier.epoch.curr, test_epoch(1));
1556            }
1557            other => panic!("Expected a barrier message, got {:?}", other),
1558        }
1559
1560        match stream.next().await {
1561            Some(Ok(Message::Chunk(chunk))) => {
1562                assert_stream_chunk_eq!(chunk, chunk_1);
1563            }
1564            other => panic!("Expected a chunk message, got {:?}", other),
1565        }
1566
1567        match stream.next().await {
1568            Some(Ok(Message::Watermark(watermark))) => {
1569                assert_eq!(watermark, Watermark::new(0, DataType::Int64, 7_i64.into()));
1570            }
1571            other => panic!("Expected a watermark message, got {:?}", other),
1572        }
1573
1574        match stream.next().await {
1575            Some(Ok(Message::Chunk(chunk))) => {
1576                assert_stream_chunk_eq!(chunk, chunk_2);
1577            }
1578            other => panic!("Expected a chunk message, got {:?}", other),
1579        }
1580
1581        match stream.next().await {
1582            Some(Ok(Message::Barrier(barrier))) => {
1583                assert_eq!(barrier.epoch.curr, test_epoch(2));
1584            }
1585            other => panic!("Expected a barrier message, got {:?}", other),
1586        }
1587    }
1588
1589    // When we hit buffer max_chunk, we only store placeholder `FlushedItem`.
1590    // So we just let capacity = 0, and we will always flush incoming chunks to state store.
1591    #[tokio::test]
1592    async fn test_max_chunk_persisted_read() {
1593        init_logger();
1594
1595        let pk_info = &KV_LOG_STORE_V2_INFO;
1596        let column_descs = test_payload_schema(pk_info);
1597        let fields = column_descs
1598            .into_iter()
1599            .map(|desc| Field::new(desc.name.clone(), desc.data_type))
1600            .collect_vec();
1601        let schema = Schema { fields };
1602        let stream_key = vec![0];
1603        let (mut tx, source) = MockSource::channel();
1604        let source = source.into_executor(schema.clone(), stream_key.clone());
1605
1606        let vnodes = Some(Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)));
1607
1608        let table = gen_test_log_store_table(pk_info);
1609
1610        let log_store_executor = SyncedKvLogStoreExecutor::new(
1611            ActorContext::for_test(123),
1612            table.id,
1613            SyncedKvLogStoreMetrics::for_test(),
1614            LogStoreRowSerde::new(&table, vnodes, pk_info),
1615            MemoryStateStore::new(),
1616            0,
1617            256,
1618            source,
1619            Duration::from_millis(256),
1620            false,
1621        )
1622        .boxed();
1623
1624        // Init
1625        tx.push_barrier(test_epoch(1), false);
1626
1627        let chunk_1 = StreamChunk::from_pretty(
1628            "  I   T
1629            +  5  10
1630            +  6  10
1631            +  8  10
1632            +  9  10
1633            +  10 11",
1634        );
1635
1636        let chunk_2 = StreamChunk::from_pretty(
1637            "   I   T
1638            -   5  10
1639            -   6  10
1640            -   8  10
1641            U- 10  11
1642            U+ 10  10",
1643        );
1644
1645        let chunk_3 = StreamChunk::from_pretty(
1646            "   I   T
1647            +  11  12",
1648        );
1649
1650        tx.push_chunk(chunk_1.clone());
1651        tx.push_chunk(chunk_2.clone());
1652        tx.push_int64_watermark(0, 7);
1653        tx.push_chunk(chunk_3.clone());
1654
1655        tx.push_barrier(test_epoch(2), false);
1656
1657        let mut stream = log_store_executor.execute();
1658
1659        for i in 1..=2 {
1660            match stream.next().await {
1661                Some(Ok(Message::Barrier(barrier))) => {
1662                    assert_eq!(barrier.epoch.curr, test_epoch(i));
1663                }
1664                other => panic!("Expected a barrier message, got {:?}", other),
1665            }
1666        }
1667
1668        match stream.next().await {
1669            Some(Ok(Message::Chunk(actual))) => {
1670                let expected = StreamChunk::from_pretty(
1671                    "   I   T
1672                    +   5  10
1673                    +   6  10
1674                    +   8  10
1675                    +   9  10
1676                    +  10  11
1677                    -   5  10
1678                    -   6  10
1679                    -   8  10
1680                    U- 10  11
1681                    U+ 10  10",
1682                );
1683                assert_stream_chunk_eq!(actual, expected);
1684            }
1685            other => panic!("Expected a chunk message, got {:?}", other),
1686        }
1687
1688        match stream.next().await {
1689            Some(Ok(Message::Watermark(watermark))) => {
1690                assert_eq!(watermark, Watermark::new(0, DataType::Int64, 7_i64.into()));
1691            }
1692            other => panic!("Expected a watermark message, got {:?}", other),
1693        }
1694
1695        match stream.next().await {
1696            Some(Ok(Message::Chunk(actual))) => {
1697                assert_stream_chunk_eq!(actual, chunk_3);
1698            }
1699            other => panic!("Expected a chunk message, got {:?}", other),
1700        }
1701    }
1702}