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            risingwave_pb::stream_plan::IcebergPkIndexCompactionContext {
350                sink_id,
351                task_id: 7.into(),
352                phase: phase as i32,
353                resolver_task_input: (phase == Phase::Begin).then_some(
354                    risingwave_pb::stream_plan::iceberg_pk_index_compaction_context::ResolverTaskInput::default(),
355                ),
356            },
357        )
358    }
359
360    const MERGER_ACTOR_ID: ActorId = ActorId::new(123);
361
362    fn input_schema() -> Schema {
363        Schema::new(vec![
364            Field::unnamed(DataType::Varchar),
365            Field::unnamed(DataType::Int64),
366        ])
367    }
368
369    #[tokio::test]
370    async fn test_position_delete_merger_basic() {
371        let handler = PositionDeleteHandlerMock::new();
372        let written_dvs = handler.written_dvs.clone();
373
374        let (mut tx, source) = MockSource::channel();
375        let source = source.into_executor(input_schema(), vec![]);
376
377        let lbm = LocalBarrierManager::for_test();
378        let mut executor =
379            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
380                .boxed()
381                .execute();
382
383        tx.push_barrier(test_epoch(1), false);
384        assert!(executor.next().await.unwrap().unwrap().is_barrier());
385
386        tx.push_chunk(build_delete_position_chunk(&[
387            ("file1.parquet", 0),
388            ("file1.parquet", 3),
389            ("file2.parquet", 1),
390        ]));
391        tx.push_barrier(test_epoch(2), false);
392
393        assert!(executor.next().await.unwrap().unwrap().is_barrier());
394
395        let dvs = written_dvs.lock().unwrap();
396        assert_eq!(dvs.get("file1.parquet").unwrap(), &BTreeSet::from([0, 3]));
397        assert_eq!(dvs.get("file2.parquet").unwrap(), &BTreeSet::from([1]));
398    }
399
400    #[tokio::test]
401    async fn test_position_delete_merger_merge_with_existing() {
402        let handler = PositionDeleteHandlerMock::new()
403            .with_existing_dv("file1.parquet", BTreeSet::from([0, 5, 10]));
404        let written_dvs = handler.written_dvs.clone();
405
406        let (mut tx, source) = MockSource::channel();
407        let source = source.into_executor(input_schema(), vec![]);
408
409        let lbm = LocalBarrierManager::for_test();
410        let mut executor =
411            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
412                .boxed()
413                .execute();
414
415        tx.push_barrier(test_epoch(1), false);
416        assert!(executor.next().await.unwrap().unwrap().is_barrier());
417
418        tx.push_chunk(build_delete_position_chunk(&[
419            ("file1.parquet", 3),
420            ("file1.parquet", 5),
421            ("file1.parquet", 7),
422        ]));
423        tx.push_barrier(test_epoch(2), false);
424
425        assert!(executor.next().await.unwrap().unwrap().is_barrier());
426
427        let dvs = written_dvs.lock().unwrap();
428        assert_eq!(
429            dvs.get("file1.parquet").unwrap(),
430            &BTreeSet::from([0, 3, 5, 7, 10])
431        );
432    }
433
434    #[tokio::test]
435    async fn test_position_delete_merger_no_deletes() {
436        let handler = PositionDeleteHandlerMock::new();
437        let written_dvs = handler.written_dvs.clone();
438
439        let (mut tx, source) = MockSource::channel();
440        let source = source.into_executor(input_schema(), vec![]);
441
442        let lbm = LocalBarrierManager::for_test();
443        let mut executor =
444            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
445                .boxed()
446                .execute();
447
448        tx.push_barrier(test_epoch(1), false);
449        assert!(executor.next().await.unwrap().unwrap().is_barrier());
450
451        tx.push_barrier(test_epoch(2), false);
452        assert!(executor.next().await.unwrap().unwrap().is_barrier());
453
454        assert!(written_dvs.lock().unwrap().is_empty());
455    }
456
457    #[tokio::test]
458    async fn test_position_delete_merger_multiple_epochs() {
459        let handler = PositionDeleteHandlerMock::new();
460        let written_dvs = handler.written_dvs.clone();
461
462        let (mut tx, source) = MockSource::channel();
463        let source = source.into_executor(input_schema(), vec![]);
464
465        let lbm = LocalBarrierManager::for_test();
466        let mut executor =
467            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
468                .boxed()
469                .execute();
470
471        tx.push_barrier(test_epoch(1), false);
472        assert!(executor.next().await.unwrap().unwrap().is_barrier());
473
474        tx.push_chunk(build_delete_position_chunk(&[("file1.parquet", 0)]));
475        tx.push_barrier(test_epoch(2), false);
476        assert!(executor.next().await.unwrap().unwrap().is_barrier());
477        assert_eq!(
478            written_dvs.lock().unwrap().get("file1.parquet").unwrap(),
479            &BTreeSet::from([0])
480        );
481
482        tx.push_chunk(build_delete_position_chunk(&[("file1.parquet", 2)]));
483        tx.push_barrier(test_epoch(3), false);
484        assert!(executor.next().await.unwrap().unwrap().is_barrier());
485
486        assert_eq!(
487            written_dvs.lock().unwrap().get("file1.parquet").unwrap(),
488            &BTreeSet::from([0, 2])
489        );
490    }
491
492    #[tokio::test]
493    async fn test_compaction_b2_forwards_before_reseed_wait() {
494        let sink_id = SinkId::new(0);
495        let handler = PositionDeleteHandlerMock::new();
496        let seed_epochs = handler.seed_epochs.clone();
497
498        let (mut tx, source) = MockSource::channel();
499        let source = source.into_executor(input_schema(), vec![]);
500        let mut executor = PositionDeleteMergerExecutor::new(
501            MERGER_ACTOR_ID,
502            sink_id,
503            LocalBarrierManager::for_test(),
504            source,
505            handler,
506        )
507        .boxed()
508        .execute();
509
510        tx.push_barrier(test_epoch(1), false);
511        assert!(executor.next().await.unwrap().unwrap().is_barrier());
512        let initial_prev = EpochPair::new_test_epoch(test_epoch(1)).prev;
513
514        tx.push_chunk(build_delete_position_chunk(&[("input.parquet", 0)]));
515        tx.send_barrier(compaction_barrier(2, sink_id, Phase::End));
516        let b2 = executor.next().await.unwrap().unwrap();
517        assert!(b2.is_barrier(), "B2 must be observable before re-seeding");
518        assert_eq!(
519            *seed_epochs.lock().unwrap(),
520            vec![initial_prev],
521            "the generator must not call start_seed until resumed after yielding B2"
522        );
523
524        tx.push_barrier(test_epoch(3), false);
525        assert!(executor.next().await.unwrap().unwrap().is_barrier());
526        assert_eq!(
527            *seed_epochs.lock().unwrap(),
528            vec![initial_prev, EpochPair::new_test_epoch(test_epoch(2)).prev]
529        );
530    }
531
532    #[tokio::test]
533    async fn test_compaction_reseed_buffers_e2_until_next_checkpoint() {
534        let sink_id = SinkId::new(0);
535        let handler = PositionDeleteHandlerMock::new();
536        handler.block_second_seed.store(true, Ordering::SeqCst);
537        let second_seed_started = handler.second_seed_started.clone();
538        let release_second_seed = handler.release_second_seed.clone();
539        let flush_calls = handler.flush_calls.clone();
540        let written_dvs = handler.written_dvs.clone();
541
542        let (mut tx, source) = MockSource::channel();
543        let source = source.into_executor(input_schema(), vec![]);
544        let mut executor = PositionDeleteMergerExecutor::new(
545            MERGER_ACTOR_ID,
546            sink_id,
547            LocalBarrierManager::for_test(),
548            source,
549            handler,
550        )
551        .boxed()
552        .execute();
553
554        tx.push_barrier(test_epoch(1), false);
555        assert!(executor.next().await.unwrap().unwrap().is_barrier());
556        tx.send_barrier(compaction_barrier(2, sink_id, Phase::End));
557        assert!(executor.next().await.unwrap().unwrap().is_barrier());
558        assert!(
559            !second_seed_started.load(Ordering::SeqCst),
560            "B2 must be forwarded before the re-seed starts"
561        );
562
563        tx.push_chunk(build_delete_position_chunk(&[("output.parquet", 9)]));
564        tx.push_barrier(test_epoch(3), false);
565        assert!(
566            tokio::time::timeout(std::time::Duration::from_millis(20), executor.next())
567                .await
568                .is_err(),
569            "the next checkpoint must wait for the post-B2 seed"
570        );
571        assert!(second_seed_started.load(Ordering::SeqCst));
572        assert_eq!(flush_calls.load(Ordering::SeqCst), 2);
573        assert!(
574            written_dvs.lock().unwrap().is_empty(),
575            "E2 deletes must remain buffered until seeding completes"
576        );
577
578        release_second_seed.notify_one();
579        assert!(executor.next().await.unwrap().unwrap().is_barrier());
580        assert_eq!(
581            written_dvs.lock().unwrap().get("output.parquet").unwrap(),
582            &BTreeSet::from([9])
583        );
584    }
585
586    /// Compaction retires the data files the resident delete state is keyed by, and the surviving
587    /// output files carry delete files written by the compaction resolver. The `End` barrier must
588    /// therefore evict the whole cache and re-seed from the post-compaction snapshot, so the next
589    /// delete merges with the resolver's delete instead of clobbering it.
590    #[tokio::test]
591    async fn test_position_delete_merger_compaction_resume_reseeds_from_new_snapshot() {
592        let sink_id = SinkId::new(0);
593        let handler = PositionDeleteHandlerMock::new();
594        let table_state = handler.table_state.clone();
595        let staged = handler.staged.clone();
596        let written_dvs = handler.written_dvs.clone();
597        let seed_epochs = handler.seed_epochs.clone();
598
599        let (mut tx, source) = MockSource::channel();
600        let source = source.into_executor(input_schema(), vec![]);
601
602        let lbm = LocalBarrierManager::for_test();
603        let mut executor =
604            PositionDeleteMergerExecutor::new(MERGER_ACTOR_ID, sink_id, lbm, source, handler)
605                .boxed()
606                .execute();
607
608        tx.push_barrier(test_epoch(1), false);
609        assert!(executor.next().await.unwrap().unwrap().is_barrier());
610
611        // Steady state: delete position 0 of the pre-compaction data file.
612        tx.push_chunk(build_delete_position_chunk(&[("input.parquet", 0)]));
613        tx.push_barrier(test_epoch(2), false);
614        assert!(executor.next().await.unwrap().unwrap().is_barrier());
615        assert_eq!(
616            written_dvs.lock().unwrap().get("input.parquet").unwrap(),
617            &BTreeSet::from([0])
618        );
619
620        // Compaction commits under epoch `test_epoch(2)`: `input.parquet` is retired in favour of
621        // `output.parquet`, whose position 4 is already deleted by a resolver-written delete file.
622        {
623            let mut table_state = table_state.lock().unwrap();
624            table_state.remove("input.parquet");
625            table_state.insert("output.parquet".to_owned(), BTreeSet::from([4]));
626        }
627
628        tx.send_barrier(compaction_barrier(3, sink_id, Phase::End));
629        assert!(executor.next().await.unwrap().unwrap().is_barrier());
630        assert_eq!(
631            *seed_epochs.lock().unwrap(),
632            vec![EpochPair::new_test_epoch(test_epoch(1)).prev],
633            "re-seed starts only after B2 is forwarded"
634        );
635
636        tx.push_chunk(build_delete_position_chunk(&[("output.parquet", 9)]));
637        tx.push_barrier(test_epoch(4), false);
638        assert!(executor.next().await.unwrap().unwrap().is_barrier());
639
640        {
641            let staged = staged.lock().unwrap();
642            assert!(
643                !staged.contains_key("input.parquet"),
644                "the retired data file must be evicted from the resident cache"
645            );
646            assert_eq!(
647                staged.get("output.parquet").unwrap(),
648                &BTreeSet::from([4, 9]),
649                "the resolver's delete and E2 delete must be resident after the re-seed and flush"
650            );
651        }
652
653        assert_eq!(
654            written_dvs.lock().unwrap().get("output.parquet").unwrap(),
655            &BTreeSet::from([4, 9])
656        );
657
658        // Seeded twice: at the first barrier, and at the epoch the compaction committed under.
659        assert_eq!(
660            *seed_epochs.lock().unwrap(),
661            vec![
662                EpochPair::new_test_epoch(test_epoch(1)).prev,
663                EpochPair::new_test_epoch(test_epoch(3)).prev,
664            ]
665        );
666    }
667
668    /// Only the `End` context signals that a compaction has committed. The `Begin` context arrives
669    /// before the resolver has run, so re-seeding there would just reload the same snapshot.
670    #[tokio::test]
671    async fn test_position_delete_merger_compaction_begin_does_not_reseed() {
672        let sink_id = SinkId::new(0);
673        let handler = PositionDeleteHandlerMock::new();
674        let seed_epochs = handler.seed_epochs.clone();
675
676        let (mut tx, source) = MockSource::channel();
677        let source = source.into_executor(input_schema(), vec![]);
678
679        let lbm = LocalBarrierManager::for_test();
680        let mut executor =
681            PositionDeleteMergerExecutor::new(MERGER_ACTOR_ID, sink_id, lbm, source, handler)
682                .boxed()
683                .execute();
684
685        tx.push_barrier(test_epoch(1), false);
686        assert!(executor.next().await.unwrap().unwrap().is_barrier());
687
688        tx.send_barrier(compaction_barrier(2, sink_id, Phase::Begin));
689        assert!(executor.next().await.unwrap().unwrap().is_barrier());
690
691        assert_eq!(seed_epochs.lock().unwrap().len(), 1);
692    }
693
694    /// Compaction is coordinated per sink, and one worker can host mergers of several sinks. A
695    /// `End` for another sink's compaction must not evict this sink's cache.
696    #[tokio::test]
697    async fn test_position_delete_merger_compaction_resume_for_other_sink_is_ignored() {
698        let handler = PositionDeleteHandlerMock::new();
699        let seed_epochs = handler.seed_epochs.clone();
700
701        let (mut tx, source) = MockSource::channel();
702        let source = source.into_executor(input_schema(), vec![]);
703
704        let lbm = LocalBarrierManager::for_test();
705        let mut executor = PositionDeleteMergerExecutor::new(
706            MERGER_ACTOR_ID,
707            SinkId::new(0),
708            lbm,
709            source,
710            handler,
711        )
712        .boxed()
713        .execute();
714
715        tx.push_barrier(test_epoch(1), false);
716        assert!(executor.next().await.unwrap().unwrap().is_barrier());
717
718        tx.send_barrier(compaction_barrier(2, SinkId::new(1), Phase::End));
719        assert!(executor.next().await.unwrap().unwrap().is_barrier());
720
721        assert_eq!(seed_epochs.lock().unwrap().len(), 1);
722    }
723}