Skip to main content

risingwave_stream/executor/iceberg_with_pk_index/
position_delete_merger.rs

1// Copyright 2026 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use anyhow::Context;
16use risingwave_common::bail;
17use risingwave_common::id::SinkId;
18use risingwave_connector::sink::Result as SinkResult;
19use risingwave_pb::connector_service::SinkMetadata;
20use risingwave_pb::id::IcebergCompactionTaskId;
21use risingwave_pb::stream_plan::iceberg_pk_index_compaction_context::Phase;
22use risingwave_pb::stream_service::PbIcebergPkIndexSinkRole;
23
24use crate::executor::prelude::*;
25use crate::task::LocalBarrierManager;
26
27/// Trait abstracting position-delete file operations for testability.
28///
29/// Implementations are responsible for reading existing position deletes
30/// (V3 Puffin deletion vectors or V2 Parquet position-delete files),
31/// merging new delete positions, writing the resulting delete file, and
32/// returning the commit metadata for the current barrier.
33#[async_trait::async_trait]
34pub trait PositionDeleteHandler: Send + 'static {
35    /// Start (or restart) seeding the resident delete state from an iceberg snapshot that includes
36    /// everything committed through `wait_epoch`.
37    ///
38    /// Called once after the executor's first barrier, and again after every compaction commit (see
39    /// [`PositionDeleteMergerExecutor`]). A restart MUST discard all previously seeded and written
40    /// state.
41    fn start_seed(&mut self, wait_epoch: u64);
42
43    fn write(&mut self, path: &str, pos: i64) -> SinkResult<()>;
44
45    async fn flush(&mut self) -> SinkResult<Option<SinkMetadata>>;
46}
47
48/// Position-delete merger executor for iceberg pk-index sink without Equality Delete.
49///
50/// This stateless executor receives [`file_path`, `position`] messages from the Writer Executor,
51/// merges them with existing position deletes, and reports the merged delete-file metadata to
52/// meta on each barrier. Depending on the table format version, the written delete file is a
53/// V3 Puffin deletion vector or a V2 file-scoped Parquet position-delete file.
54///
55/// The upstream plan shards messages by `file_path`, so each actor only merges delete
56/// positions for the files assigned to its shard.
57///
58/// # Compaction
59///
60/// Compaction rewrites the very data files this executor's resident delete state is keyed by. On
61/// the `End` half of a coordinated compaction barrier pair, the merger therefore discards that
62/// state and re-seeds from the post-compaction snapshot.
63///
64/// Input schema: [`file_path`: Varchar, `position`: int64]
65/// Output: Barriers and watermarks only; no data chunks (terminal executor in the stream graph).
66pub struct PositionDeleteMergerExecutor<H>
67where
68    H: PositionDeleteHandler,
69{
70    actor_id: ActorId,
71    sink_id: SinkId,
72    local_barrier_manager: LocalBarrierManager,
73    input: Option<Executor>,
74    handler: H,
75}
76
77impl<H> PositionDeleteMergerExecutor<H>
78where
79    H: PositionDeleteHandler,
80{
81    pub fn new(
82        actor_id: ActorId,
83        sink_id: SinkId,
84        local_barrier_manager: LocalBarrierManager,
85        input: Executor,
86        handler: H,
87    ) -> Self {
88        Self {
89            actor_id,
90            sink_id,
91            local_barrier_manager,
92            input: Some(input),
93            handler,
94        }
95    }
96
97    #[try_stream(ok = Message, error = StreamExecutorError)]
98    async fn execute_inner(mut self) {
99        let mut input = self.input.take().unwrap().execute();
100
101        // Consume the first barrier. Its `prev` epoch is the point the previous actors (if this is a
102        // scale) committed up to; seed only from a snapshot that includes it.
103        let barrier = expect_first_barrier(&mut input).await?;
104        self.handler.start_seed(barrier.epoch.prev);
105        yield Message::Barrier(barrier);
106
107        #[for_await]
108        for msg in input {
109            match msg? {
110                Message::Chunk(chunk) => {
111                    for (op, row) in chunk.rows() {
112                        debug_assert_eq!(op, risingwave_common::array::Op::Insert);
113                        let file_path = row
114                            .datum_at(0)
115                            .map(|d| d.into_utf8())
116                            .context("file_path should not be null")?;
117                        let position = row
118                            .datum_at(1)
119                            .context("position should not be null")?
120                            .into_int64();
121                        self.handler
122                            .write(file_path, position)
123                            .map_err(|e| StreamExecutorError::sink_error(e, self.sink_id))?;
124                    }
125                }
126                Message::Barrier(barrier) => {
127                    barrier.assume_no_update_vnode_bitmap(self.actor_id)?;
128                    let compaction_resumed = self.compaction_resume_task_id(&barrier);
129                    if compaction_resumed.is_some() && !barrier.is_checkpoint() {
130                        bail!(
131                            "iceberg pk-index merger {} received a non-checkpoint compaction resume barrier {:?}",
132                            self.sink_id,
133                            barrier.epoch
134                        );
135                    }
136
137                    let mut metadata = None;
138                    if barrier.is_checkpoint() {
139                        metadata = self
140                            .handler
141                            .flush()
142                            .await
143                            .map_err(|e| StreamExecutorError::sink_error(e, self.sink_id))?;
144                    }
145
146                    if let Some(metadata) = metadata
147                        && metadata.metadata.is_some()
148                    {
149                        self.local_barrier_manager
150                            .report_iceberg_pk_index_sink_metadata(
151                                barrier.epoch,
152                                self.sink_id,
153                                self.actor_id,
154                                PbIcebergPkIndexSinkRole::PositionDeleteMerger,
155                                Some(metadata),
156                            );
157                    }
158
159                    let reseed = compaction_resumed.map(|_| barrier.epoch.prev);
160                    yield Message::Barrier(barrier);
161
162                    if let Some(epoch) = reseed {
163                        self.handler.start_seed(epoch);
164                    }
165                }
166                Message::Watermark(w) => {
167                    yield Message::Watermark(w);
168                }
169            }
170        }
171    }
172
173    /// The compaction task id if `barrier` carries the `End` half of a coordinated compaction
174    /// barrier pair for this sink, meaning the compaction commits under `barrier.epoch.prev`.
175    fn compaction_resume_task_id(&self, barrier: &Barrier) -> Option<IcebergCompactionTaskId> {
176        match barrier.iceberg_pk_index_compaction() {
177            Some(context)
178                if context.sink_id == self.sink_id && context.phase == Phase::End as i32 =>
179            {
180                Some(context.task_id)
181            }
182            _ => None,
183        }
184    }
185}
186
187impl<H> Execute for PositionDeleteMergerExecutor<H>
188where
189    H: PositionDeleteHandler,
190{
191    fn execute(self: Box<Self>) -> BoxedMessageStream {
192        self.execute_inner().boxed()
193    }
194}
195
196#[cfg(test)]
197mod tests {
198    use std::collections::BTreeSet;
199    use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
200    use std::sync::{Arc, Mutex};
201
202    use hashbrown::HashMap;
203    use risingwave_common::array::{Array, ArrayBuilder, I64ArrayBuilder, Op, Utf8ArrayBuilder};
204    use risingwave_common::catalog::{Field, Schema};
205    use risingwave_common::id::SinkId;
206    use risingwave_common::types::DataType;
207    use risingwave_common::util::epoch::{EpochPair, test_epoch};
208
209    use super::*;
210    use crate::executor::test_utils::MockSource;
211    use crate::task::LocalBarrierManager;
212
213    fn build_delete_position_chunk(positions: &[(&str, i64)]) -> StreamChunk {
214        let len = positions.len();
215        let mut file_path_builder = Utf8ArrayBuilder::new(len);
216        let mut position_builder = I64ArrayBuilder::new(len);
217
218        for (path, offset) in positions {
219            file_path_builder.append(Some(*path));
220            position_builder.append(Some(*offset));
221        }
222
223        StreamChunk::from_parts(
224            vec![Op::Insert; len],
225            risingwave_common::array::DataChunk::new(
226                vec![
227                    file_path_builder.finish().into_ref(),
228                    position_builder.finish().into_ref(),
229                ],
230                len,
231            ),
232        )
233    }
234
235    type Dvs = HashMap<String, BTreeSet<i64>>;
236
237    /// In-memory stand-in for [`PositionDeleteHandlerImpl`](super::PositionDeleteHandlerImpl),
238    /// modelling the three layers that matter for the re-seed contract:
239    ///
240    /// - `table_state`: the committed iceberg snapshot. Tests mutate it directly to simulate an
241    ///   external commit (e.g. compaction replacing data files and adding resolver delete files).
242    /// - `staged`: the resident per-shard cache. Populated only by `start_seed` (from a *clone* of
243    ///   `table_state`, so later external commits are invisible until the next seed) and appended
244    ///   to by `flush`.
245    /// - `pending`: positions buffered since the last flush.
246    #[derive(Clone)]
247    struct PositionDeleteHandlerMock {
248        table_state: Arc<Mutex<Dvs>>,
249        staged: Arc<Mutex<Dvs>>,
250        pending: Arc<Mutex<Dvs>>,
251        written_dvs: Arc<Mutex<Dvs>>,
252        /// Every `wait_epoch` passed to `start_seed`, in call order.
253        seed_epochs: Arc<Mutex<Vec<u64>>>,
254        block_second_seed: Arc<AtomicBool>,
255        second_seed_started: Arc<AtomicBool>,
256        release_second_seed: Arc<tokio::sync::Notify>,
257        flush_calls: Arc<AtomicUsize>,
258    }
259
260    impl PositionDeleteHandlerMock {
261        fn new() -> Self {
262            Self {
263                table_state: Arc::new(Mutex::new(HashMap::new())),
264                staged: Arc::new(Mutex::new(HashMap::new())),
265                pending: Arc::new(Mutex::new(HashMap::new())),
266                written_dvs: Arc::new(Mutex::new(HashMap::new())),
267                seed_epochs: Arc::new(Mutex::new(Vec::new())),
268                block_second_seed: Arc::new(AtomicBool::new(false)),
269                second_seed_started: Arc::new(AtomicBool::new(false)),
270                release_second_seed: Arc::new(tokio::sync::Notify::new()),
271                flush_calls: Arc::new(AtomicUsize::new(0)),
272            }
273        }
274
275        /// Pre-commit a delete vector into the table, as if a previous incarnation of this actor
276        /// (or the compaction resolver) had written it. Picked up by the next `start_seed`.
277        fn with_existing_dv(self, file_path: &str, positions: BTreeSet<i64>) -> Self {
278            self.table_state
279                .lock()
280                .unwrap()
281                .insert(file_path.to_owned(), positions);
282            self
283        }
284    }
285
286    #[async_trait::async_trait]
287    impl PositionDeleteHandler for PositionDeleteHandlerMock {
288        fn start_seed(&mut self, wait_epoch: u64) {
289            let mut seed_epochs = self.seed_epochs.lock().unwrap();
290            seed_epochs.push(wait_epoch);
291            let is_second_seed = seed_epochs.len() == 2;
292            drop(seed_epochs);
293            if is_second_seed && self.block_second_seed.load(Ordering::SeqCst) {
294                self.second_seed_started.store(true, Ordering::SeqCst);
295                return;
296            }
297            // Discard the resident cache wholesale and re-derive it from the (possibly advanced)
298            // committed snapshot. This is the eviction + resync the real handler gets by dropping
299            // its `SeededState`.
300            *self.staged.lock().unwrap() = self.table_state.lock().unwrap().clone();
301        }
302
303        fn write(&mut self, path: &str, pos: i64) -> SinkResult<()> {
304            self.pending
305                .lock()
306                .unwrap()
307                .entry_ref(path)
308                .or_default()
309                .insert(pos);
310            Ok(())
311        }
312
313        async fn flush(&mut self) -> SinkResult<Option<SinkMetadata>> {
314            self.flush_calls.fetch_add(1, Ordering::SeqCst);
315            if self.second_seed_started.load(Ordering::SeqCst)
316                && self.block_second_seed.load(Ordering::SeqCst)
317            {
318                self.release_second_seed.notified().await;
319                *self.staged.lock().unwrap() = self.table_state.lock().unwrap().clone();
320            }
321            let pending = std::mem::take(&mut *self.pending.lock().unwrap());
322            if pending.is_empty() {
323                return Ok(None);
324            }
325
326            let mut staged = self.staged.lock().unwrap();
327            let mut table_state = self.table_state.lock().unwrap();
328            let mut written_dvs = self.written_dvs.lock().unwrap();
329            for (file_path, positions) in pending {
330                // Merge against the resident state only, exactly like the real handler: anything
331                // committed to the table but not seeded is invisible and would be clobbered.
332                let mut merged = staged.get(&file_path).cloned().unwrap_or_default();
333                merged.extend(positions);
334                staged.insert(file_path.clone(), merged.clone());
335                table_state.insert(file_path.clone(), merged.clone());
336                written_dvs.insert(file_path, merged);
337            }
338            // The mock asserts on side effects (`written_dvs`) rather than emitted
339            // metadata, so we still return Ok(None). The report-on-barrier path
340            // is exercised by the SLT integration tests.
341            // TODO: add unit test for report-on-barrier path once test infra is
342            // available to capture `LocalBarrierEvent`s on the receiver side.
343            Ok(None)
344        }
345    }
346
347    fn compaction_barrier(epoch: u64, sink_id: SinkId, phase: Phase) -> Barrier {
348        Barrier::new_test_barrier(test_epoch(epoch)).with_iceberg_pk_index_compaction(
349            crate::executor::IcebergPkIndexCompactionContext {
350                sink_id,
351                task_id: 7.into(),
352                phase: phase as i32,
353            },
354        )
355    }
356
357    const MERGER_ACTOR_ID: ActorId = ActorId::new(123);
358
359    fn input_schema() -> Schema {
360        Schema::new(vec![
361            Field::unnamed(DataType::Varchar),
362            Field::unnamed(DataType::Int64),
363        ])
364    }
365
366    #[tokio::test]
367    async fn test_position_delete_merger_basic() {
368        let handler = PositionDeleteHandlerMock::new();
369        let written_dvs = handler.written_dvs.clone();
370
371        let (mut tx, source) = MockSource::channel();
372        let source = source.into_executor(input_schema(), vec![]);
373
374        let lbm = LocalBarrierManager::for_test();
375        let mut executor =
376            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
377                .boxed()
378                .execute();
379
380        tx.push_barrier(test_epoch(1), false);
381        assert!(executor.next().await.unwrap().unwrap().is_barrier());
382
383        tx.push_chunk(build_delete_position_chunk(&[
384            ("file1.parquet", 0),
385            ("file1.parquet", 3),
386            ("file2.parquet", 1),
387        ]));
388        tx.push_barrier(test_epoch(2), false);
389
390        assert!(executor.next().await.unwrap().unwrap().is_barrier());
391
392        let dvs = written_dvs.lock().unwrap();
393        assert_eq!(dvs.get("file1.parquet").unwrap(), &BTreeSet::from([0, 3]));
394        assert_eq!(dvs.get("file2.parquet").unwrap(), &BTreeSet::from([1]));
395    }
396
397    #[tokio::test]
398    async fn test_position_delete_merger_merge_with_existing() {
399        let handler = PositionDeleteHandlerMock::new()
400            .with_existing_dv("file1.parquet", BTreeSet::from([0, 5, 10]));
401        let written_dvs = handler.written_dvs.clone();
402
403        let (mut tx, source) = MockSource::channel();
404        let source = source.into_executor(input_schema(), vec![]);
405
406        let lbm = LocalBarrierManager::for_test();
407        let mut executor =
408            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
409                .boxed()
410                .execute();
411
412        tx.push_barrier(test_epoch(1), false);
413        assert!(executor.next().await.unwrap().unwrap().is_barrier());
414
415        tx.push_chunk(build_delete_position_chunk(&[
416            ("file1.parquet", 3),
417            ("file1.parquet", 5),
418            ("file1.parquet", 7),
419        ]));
420        tx.push_barrier(test_epoch(2), false);
421
422        assert!(executor.next().await.unwrap().unwrap().is_barrier());
423
424        let dvs = written_dvs.lock().unwrap();
425        assert_eq!(
426            dvs.get("file1.parquet").unwrap(),
427            &BTreeSet::from([0, 3, 5, 7, 10])
428        );
429    }
430
431    #[tokio::test]
432    async fn test_position_delete_merger_no_deletes() {
433        let handler = PositionDeleteHandlerMock::new();
434        let written_dvs = handler.written_dvs.clone();
435
436        let (mut tx, source) = MockSource::channel();
437        let source = source.into_executor(input_schema(), vec![]);
438
439        let lbm = LocalBarrierManager::for_test();
440        let mut executor =
441            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
442                .boxed()
443                .execute();
444
445        tx.push_barrier(test_epoch(1), false);
446        assert!(executor.next().await.unwrap().unwrap().is_barrier());
447
448        tx.push_barrier(test_epoch(2), false);
449        assert!(executor.next().await.unwrap().unwrap().is_barrier());
450
451        assert!(written_dvs.lock().unwrap().is_empty());
452    }
453
454    #[tokio::test]
455    async fn test_position_delete_merger_multiple_epochs() {
456        let handler = PositionDeleteHandlerMock::new();
457        let written_dvs = handler.written_dvs.clone();
458
459        let (mut tx, source) = MockSource::channel();
460        let source = source.into_executor(input_schema(), vec![]);
461
462        let lbm = LocalBarrierManager::for_test();
463        let mut executor =
464            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
465                .boxed()
466                .execute();
467
468        tx.push_barrier(test_epoch(1), false);
469        assert!(executor.next().await.unwrap().unwrap().is_barrier());
470
471        tx.push_chunk(build_delete_position_chunk(&[("file1.parquet", 0)]));
472        tx.push_barrier(test_epoch(2), false);
473        assert!(executor.next().await.unwrap().unwrap().is_barrier());
474        assert_eq!(
475            written_dvs.lock().unwrap().get("file1.parquet").unwrap(),
476            &BTreeSet::from([0])
477        );
478
479        tx.push_chunk(build_delete_position_chunk(&[("file1.parquet", 2)]));
480        tx.push_barrier(test_epoch(3), false);
481        assert!(executor.next().await.unwrap().unwrap().is_barrier());
482
483        assert_eq!(
484            written_dvs.lock().unwrap().get("file1.parquet").unwrap(),
485            &BTreeSet::from([0, 2])
486        );
487    }
488
489    #[tokio::test]
490    async fn test_compaction_b2_forwards_before_reseed_wait() {
491        let sink_id = SinkId::new(0);
492        let handler = PositionDeleteHandlerMock::new();
493        let seed_epochs = handler.seed_epochs.clone();
494
495        let (mut tx, source) = MockSource::channel();
496        let source = source.into_executor(input_schema(), vec![]);
497        let mut executor = PositionDeleteMergerExecutor::new(
498            MERGER_ACTOR_ID,
499            sink_id,
500            LocalBarrierManager::for_test(),
501            source,
502            handler,
503        )
504        .boxed()
505        .execute();
506
507        tx.push_barrier(test_epoch(1), false);
508        assert!(executor.next().await.unwrap().unwrap().is_barrier());
509        let initial_prev = EpochPair::new_test_epoch(test_epoch(1)).prev;
510
511        tx.push_chunk(build_delete_position_chunk(&[("input.parquet", 0)]));
512        tx.send_barrier(compaction_barrier(2, sink_id, Phase::End));
513        let b2 = executor.next().await.unwrap().unwrap();
514        assert!(b2.is_barrier(), "B2 must be observable before re-seeding");
515        assert_eq!(
516            *seed_epochs.lock().unwrap(),
517            vec![initial_prev],
518            "the generator must not call start_seed until resumed after yielding B2"
519        );
520
521        tx.push_barrier(test_epoch(3), false);
522        assert!(executor.next().await.unwrap().unwrap().is_barrier());
523        assert_eq!(
524            *seed_epochs.lock().unwrap(),
525            vec![initial_prev, EpochPair::new_test_epoch(test_epoch(2)).prev]
526        );
527    }
528
529    #[tokio::test]
530    async fn test_compaction_reseed_buffers_e2_until_next_checkpoint() {
531        let sink_id = SinkId::new(0);
532        let handler = PositionDeleteHandlerMock::new();
533        handler.block_second_seed.store(true, Ordering::SeqCst);
534        let second_seed_started = handler.second_seed_started.clone();
535        let release_second_seed = handler.release_second_seed.clone();
536        let flush_calls = handler.flush_calls.clone();
537        let written_dvs = handler.written_dvs.clone();
538
539        let (mut tx, source) = MockSource::channel();
540        let source = source.into_executor(input_schema(), vec![]);
541        let mut executor = PositionDeleteMergerExecutor::new(
542            MERGER_ACTOR_ID,
543            sink_id,
544            LocalBarrierManager::for_test(),
545            source,
546            handler,
547        )
548        .boxed()
549        .execute();
550
551        tx.push_barrier(test_epoch(1), false);
552        assert!(executor.next().await.unwrap().unwrap().is_barrier());
553        tx.send_barrier(compaction_barrier(2, sink_id, Phase::End));
554        assert!(executor.next().await.unwrap().unwrap().is_barrier());
555        assert!(
556            !second_seed_started.load(Ordering::SeqCst),
557            "B2 must be forwarded before the re-seed starts"
558        );
559
560        tx.push_chunk(build_delete_position_chunk(&[("output.parquet", 9)]));
561        tx.push_barrier(test_epoch(3), false);
562        assert!(
563            tokio::time::timeout(std::time::Duration::from_millis(20), executor.next())
564                .await
565                .is_err(),
566            "the next checkpoint must wait for the post-B2 seed"
567        );
568        assert!(second_seed_started.load(Ordering::SeqCst));
569        assert_eq!(flush_calls.load(Ordering::SeqCst), 2);
570        assert!(
571            written_dvs.lock().unwrap().is_empty(),
572            "E2 deletes must remain buffered until seeding completes"
573        );
574
575        release_second_seed.notify_one();
576        assert!(executor.next().await.unwrap().unwrap().is_barrier());
577        assert_eq!(
578            written_dvs.lock().unwrap().get("output.parquet").unwrap(),
579            &BTreeSet::from([9])
580        );
581    }
582
583    /// Compaction retires the data files the resident delete state is keyed by, and the surviving
584    /// output files carry delete files written by the compaction resolver. The `End` barrier must
585    /// therefore evict the whole cache and re-seed from the post-compaction snapshot, so the next
586    /// delete merges with the resolver's delete instead of clobbering it.
587    #[tokio::test]
588    async fn test_position_delete_merger_compaction_resume_reseeds_from_new_snapshot() {
589        let sink_id = SinkId::new(0);
590        let handler = PositionDeleteHandlerMock::new();
591        let table_state = handler.table_state.clone();
592        let staged = handler.staged.clone();
593        let written_dvs = handler.written_dvs.clone();
594        let seed_epochs = handler.seed_epochs.clone();
595
596        let (mut tx, source) = MockSource::channel();
597        let source = source.into_executor(input_schema(), vec![]);
598
599        let lbm = LocalBarrierManager::for_test();
600        let mut executor =
601            PositionDeleteMergerExecutor::new(MERGER_ACTOR_ID, sink_id, lbm, source, handler)
602                .boxed()
603                .execute();
604
605        tx.push_barrier(test_epoch(1), false);
606        assert!(executor.next().await.unwrap().unwrap().is_barrier());
607
608        // Steady state: delete position 0 of the pre-compaction data file.
609        tx.push_chunk(build_delete_position_chunk(&[("input.parquet", 0)]));
610        tx.push_barrier(test_epoch(2), false);
611        assert!(executor.next().await.unwrap().unwrap().is_barrier());
612        assert_eq!(
613            written_dvs.lock().unwrap().get("input.parquet").unwrap(),
614            &BTreeSet::from([0])
615        );
616
617        // Compaction commits under epoch `test_epoch(2)`: `input.parquet` is retired in favour of
618        // `output.parquet`, whose position 4 is already deleted by a resolver-written delete file.
619        {
620            let mut table_state = table_state.lock().unwrap();
621            table_state.remove("input.parquet");
622            table_state.insert("output.parquet".to_owned(), BTreeSet::from([4]));
623        }
624
625        tx.send_barrier(compaction_barrier(3, sink_id, Phase::End));
626        assert!(executor.next().await.unwrap().unwrap().is_barrier());
627        assert_eq!(
628            *seed_epochs.lock().unwrap(),
629            vec![EpochPair::new_test_epoch(test_epoch(1)).prev],
630            "re-seed starts only after B2 is forwarded"
631        );
632
633        tx.push_chunk(build_delete_position_chunk(&[("output.parquet", 9)]));
634        tx.push_barrier(test_epoch(4), false);
635        assert!(executor.next().await.unwrap().unwrap().is_barrier());
636
637        {
638            let staged = staged.lock().unwrap();
639            assert!(
640                !staged.contains_key("input.parquet"),
641                "the retired data file must be evicted from the resident cache"
642            );
643            assert_eq!(
644                staged.get("output.parquet").unwrap(),
645                &BTreeSet::from([4, 9]),
646                "the resolver's delete and E2 delete must be resident after the re-seed and flush"
647            );
648        }
649
650        assert_eq!(
651            written_dvs.lock().unwrap().get("output.parquet").unwrap(),
652            &BTreeSet::from([4, 9])
653        );
654
655        // Seeded twice: at the first barrier, and at the epoch the compaction committed under.
656        assert_eq!(
657            *seed_epochs.lock().unwrap(),
658            vec![
659                EpochPair::new_test_epoch(test_epoch(1)).prev,
660                EpochPair::new_test_epoch(test_epoch(3)).prev,
661            ]
662        );
663    }
664
665    /// Only the `End` context signals that a compaction has committed. The `Begin` context arrives
666    /// before the resolver has run, so re-seeding there would just reload the same snapshot.
667    #[tokio::test]
668    async fn test_position_delete_merger_compaction_begin_does_not_reseed() {
669        let sink_id = SinkId::new(0);
670        let handler = PositionDeleteHandlerMock::new();
671        let seed_epochs = handler.seed_epochs.clone();
672
673        let (mut tx, source) = MockSource::channel();
674        let source = source.into_executor(input_schema(), vec![]);
675
676        let lbm = LocalBarrierManager::for_test();
677        let mut executor =
678            PositionDeleteMergerExecutor::new(MERGER_ACTOR_ID, sink_id, lbm, source, handler)
679                .boxed()
680                .execute();
681
682        tx.push_barrier(test_epoch(1), false);
683        assert!(executor.next().await.unwrap().unwrap().is_barrier());
684
685        tx.send_barrier(compaction_barrier(2, sink_id, Phase::Begin));
686        assert!(executor.next().await.unwrap().unwrap().is_barrier());
687
688        assert_eq!(seed_epochs.lock().unwrap().len(), 1);
689    }
690
691    /// Compaction is coordinated per sink, and one worker can host mergers of several sinks. A
692    /// `End` for another sink's compaction must not evict this sink's cache.
693    #[tokio::test]
694    async fn test_position_delete_merger_compaction_resume_for_other_sink_is_ignored() {
695        let handler = PositionDeleteHandlerMock::new();
696        let seed_epochs = handler.seed_epochs.clone();
697
698        let (mut tx, source) = MockSource::channel();
699        let source = source.into_executor(input_schema(), vec![]);
700
701        let lbm = LocalBarrierManager::for_test();
702        let mut executor = PositionDeleteMergerExecutor::new(
703            MERGER_ACTOR_ID,
704            SinkId::new(0),
705            lbm,
706            source,
707            handler,
708        )
709        .boxed()
710        .execute();
711
712        tx.push_barrier(test_epoch(1), false);
713        assert!(executor.next().await.unwrap().unwrap().is_barrier());
714
715        tx.send_barrier(compaction_barrier(2, SinkId::new(1), Phase::End));
716        assert!(executor.next().await.unwrap().unwrap().is_barrier());
717
718        assert_eq!(seed_epochs.lock().unwrap().len(), 1);
719    }
720}