Skip to main content

risingwave_meta/hummock/manager/
time_travel.rs

1// Copyright 2024 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 std::collections::{HashMap, HashSet, VecDeque};
16
17use anyhow::anyhow;
18use futures::TryStreamExt;
19use risingwave_common::catalog::TableId;
20use risingwave_common::system_param::reader::SystemParamsRead;
21use risingwave_common::util::epoch::Epoch;
22use risingwave_hummock_sdk::compaction_group::StateTableId;
23use risingwave_hummock_sdk::sstable_info::SstableInfo;
24use risingwave_hummock_sdk::time_travel::{
25    IncompleteHummockVersion, IncompleteHummockVersionDelta, refill_version,
26};
27use risingwave_hummock_sdk::version::{GroupDeltaCommon, HummockVersion, HummockVersionDelta};
28use risingwave_hummock_sdk::{CompactionGroupId, HummockEpoch, HummockObjectId, HummockSstableId};
29use risingwave_meta_model::hummock_sstable_info::SstableInfoV2Backend;
30use risingwave_meta_model::{
31    HummockVersionId, hummock_epoch_to_version, hummock_sstable_info, hummock_time_travel_delta,
32    hummock_time_travel_version,
33};
34use risingwave_pb::hummock::{PbHummockVersion, PbHummockVersionDelta};
35use sea_orm::ActiveValue::Set;
36use sea_orm::{
37    ColumnTrait, Condition, ConnectionTrait, DatabaseTransaction, EntityTrait, PaginatorTrait,
38    QueryFilter, QueryOrder, QuerySelect, TransactionTrait,
39};
40use tracing::info;
41
42use crate::hummock::HummockManager;
43use crate::hummock::error::{Error, Result};
44
45/// Time travel.
46impl HummockManager {
47    pub(crate) async fn init_time_travel_state(&self) -> Result<()> {
48        let sql_store = self.env.meta_store_ref();
49        let mut guard = self
50            .versioning
51            .write_with_process_name("init_time_travel_state")
52            .await;
53        guard.mark_next_time_travel_version_snapshot();
54
55        guard.last_time_travel_snapshot_sst_ids = HashSet::new();
56        let Some(version) = hummock_time_travel_version::Entity::find()
57            .order_by_desc(hummock_time_travel_version::Column::VersionId)
58            .one(&sql_store.conn)
59            .await?
60            .map(|v| {
61                IncompleteHummockVersion::from_persisted_protobuf_owned(v.version.to_protobuf())
62            })
63        else {
64            return Ok(());
65        };
66        guard.last_time_travel_snapshot_sst_ids = version.get_sst_ids();
67        Ok(())
68    }
69
70    pub(crate) async fn truncate_time_travel_metadata(
71        &self,
72        epoch_watermark: HummockEpoch,
73        pinned_snapshot_epochs: HashMap<TableId, HashSet<HummockEpoch>>,
74    ) -> Result<()> {
75        let _timer = self
76            .metrics
77            .time_travel_vacuum_metadata_latency
78            .start_timer();
79        let min_pinned_version_id = self
80            .context_info
81            .read_with_process_name("truncate_time_travel_metadata")
82            .await
83            .min_pinned_version_id();
84        let sql_store = self.env.meta_store_ref();
85        let txn = sql_store.conn.begin().await?;
86
87        let epoch_watermark_model =
88            risingwave_meta_model::Epoch::try_from(epoch_watermark).unwrap();
89        let has_expired_epoch = hummock_epoch_to_version::Entity::find()
90            .filter(hummock_epoch_to_version::Column::Epoch.lt(epoch_watermark_model))
91            .select_only()
92            .column(hummock_epoch_to_version::Column::VersionId)
93            .into_tuple::<HummockVersionId>()
94            .one(&txn)
95            .await?
96            .is_some();
97        if !has_expired_epoch {
98            txn.commit().await?;
99            return Ok(());
100        }
101
102        let mut pinned_epoch_rows = HashSet::new();
103        let mut pinned_snapshot_version_ids = HashSet::new();
104        for (table_id, pinned_epochs) in pinned_snapshot_epochs {
105            for pinned_epoch in pinned_epochs {
106                if pinned_epoch >= epoch_watermark {
107                    continue;
108                }
109                let pinned_epoch_model =
110                    risingwave_meta_model::Epoch::try_from(pinned_epoch).unwrap();
111                let epoch_to_version = hummock_epoch_to_version::Entity::find_by_id((
112                    pinned_epoch_model,
113                    i64::from(table_id.as_raw_id()),
114                ))
115                .one(&txn)
116                .await?;
117                let Some(epoch_to_version) = epoch_to_version else {
118                    tracing::warn!(
119                        %table_id,
120                        pinned_epoch,
121                        "pinned snapshot epoch mapping not found, skip pinning"
122                    );
123                    continue;
124                };
125                pinned_epoch_rows.insert((epoch_to_version.epoch, epoch_to_version.table_id));
126                pinned_snapshot_version_ids.insert(epoch_to_version.version_id);
127            }
128        }
129
130        let mut pinned_replay_version_ids = HashSet::new();
131        let mut pinned_delta_ids = HashSet::new();
132        let mut pinned_snapshot_sst_ids = HashSet::new();
133        let mut pinned_snapshot_object_ids = HashSet::new();
134        for pinned_snapshot_version_id in pinned_snapshot_version_ids {
135            let Some(resolved) =
136                resolve_time_travel_version(&txn, pinned_snapshot_version_id).await?
137            else {
138                tracing::warn!(
139                    %pinned_snapshot_version_id,
140                    "time travel version before pinned version not found, skip pinning"
141                );
142                continue;
143            };
144            pinned_replay_version_ids.insert(resolved.replay_version.id);
145            pinned_snapshot_sst_ids.extend(resolved.replay_version.get_sst_ids());
146            pinned_snapshot_object_ids.extend(resolved.replay_version.get_object_ids());
147            for delta in resolved.deltas {
148                pinned_delta_ids.insert(delta.id);
149                pinned_snapshot_sst_ids.extend(delta.newly_added_sst_ids());
150                pinned_snapshot_object_ids.extend(delta.newly_added_object_ids());
151            }
152        }
153
154        // Version ids follow global commit order, while epochs are only monotonic per table.
155        // Use the earliest version needed by any retained mapping as the watermark so its replay
156        // metadata cannot be truncated. If no mapping is retained, only version pins and vacuum
157        // throttling need to constrain the watermark.
158        let version_watermark = hummock_epoch_to_version::Entity::find()
159            .filter(hummock_epoch_to_version::Column::Epoch.gte(epoch_watermark_model))
160            .select_only()
161            .column(hummock_epoch_to_version::Column::VersionId)
162            .order_by_asc(hummock_epoch_to_version::Column::VersionId)
163            .into_tuple::<HummockVersionId>()
164            .one(&txn)
165            .await?;
166        // metadata BELOW watermark_version_id will be truncated.
167        let mut watermark_version_id = version_watermark.map_or(min_pinned_version_id, |id| {
168            std::cmp::min(id, min_pinned_version_id)
169        });
170        if let Some(max_version_count) = self.env.opts.time_travel_vacuum_max_version_count {
171            let mut query = hummock_time_travel_version::Entity::find()
172                .select_only()
173                .column(hummock_time_travel_version::Column::VersionId)
174                .order_by_asc(hummock_time_travel_version::Column::VersionId)
175                .limit(2);
176            // A replay version can be pinned for arbitrarily long. Counting it would keep selecting
177            // the same two earliest versions and prevent incremental vacuum from advancing.
178            if !pinned_replay_version_ids.is_empty() {
179                query = query.filter(
180                    hummock_time_travel_version::Column::VersionId
181                        .is_not_in(pinned_replay_version_ids.iter().copied()),
182                );
183            }
184            let earliest2_version_ids = query.into_tuple::<HummockVersionId>().all(&txn).await?;
185            // Ensure at least 1 version BELOW watermark_version_id if applying time_travel_vacuum_max_version_count.
186            if earliest2_version_ids.len() == 2 {
187                watermark_version_id = std::cmp::min(
188                    watermark_version_id,
189                    HummockVersionId::new(std::cmp::max(
190                        earliest2_version_ids[0]
191                            .as_raw_id()
192                            .saturating_add(max_version_count.into()),
193                        earliest2_version_ids[1].as_raw_id(),
194                    )),
195                );
196            }
197        }
198        let mut delete_epoch_rows = hummock_epoch_to_version::Entity::delete_many()
199            .filter(hummock_epoch_to_version::Column::Epoch.lt(epoch_watermark_model));
200        for (epoch, table_id) in &pinned_epoch_rows {
201            delete_epoch_rows = delete_epoch_rows.filter(
202                Condition::any()
203                    .add(hummock_epoch_to_version::Column::Epoch.ne(*epoch))
204                    .add(hummock_epoch_to_version::Column::TableId.ne(*table_id)),
205            );
206        }
207        let res = delete_epoch_rows.exec(&txn).await?;
208        tracing::info!(
209            epoch_watermark,
210            "Delete {} rows from hummock_epoch_to_version.",
211            res.rows_affected
212        );
213        let latest_valid_version = hummock_time_travel_version::Entity::find()
214            .filter(hummock_time_travel_version::Column::VersionId.lte(watermark_version_id))
215            .order_by_desc(hummock_time_travel_version::Column::VersionId)
216            .one(&txn)
217            .await?
218            .map(|m| {
219                IncompleteHummockVersion::from_persisted_protobuf_owned(m.version.to_protobuf())
220            });
221        let Some(latest_valid_version) = latest_valid_version else {
222            txn.commit().await?;
223            return Ok(());
224        };
225        let latest_valid_version_id = latest_valid_version.id;
226        let mut retained_snapshot_sst_ids = latest_valid_version.get_sst_ids();
227        let mut retained_snapshot_object_ids = latest_valid_version
228            .get_object_ids()
229            .collect::<HashSet<_>>();
230        retained_snapshot_sst_ids.extend(pinned_snapshot_sst_ids);
231        retained_snapshot_object_ids.extend(pinned_snapshot_object_ids);
232        let mut object_ids_to_delete: HashSet<_> = HashSet::default();
233        let mut version_delete_condition = Condition::all()
234            .add(hummock_time_travel_version::Column::VersionId.lt(latest_valid_version_id));
235        if !pinned_replay_version_ids.is_empty() {
236            version_delete_condition = version_delete_condition.add(
237                hummock_time_travel_version::Column::VersionId
238                    .is_not_in(pinned_replay_version_ids.iter().copied()),
239            );
240        }
241        let version_ids_to_delete: Vec<risingwave_meta_model::HummockVersionId> =
242            hummock_time_travel_version::Entity::find()
243                .select_only()
244                .column(hummock_time_travel_version::Column::VersionId)
245                .filter(version_delete_condition.clone())
246                .order_by_desc(hummock_time_travel_version::Column::VersionId)
247                .into_tuple()
248                .all(&txn)
249                .await?;
250        let mut delta_delete_condition = Condition::all()
251            .add(hummock_time_travel_delta::Column::VersionId.lt(latest_valid_version_id));
252        if !pinned_delta_ids.is_empty() {
253            delta_delete_condition = delta_delete_condition.add(
254                hummock_time_travel_delta::Column::VersionId
255                    .is_not_in(pinned_delta_ids.iter().copied()),
256            );
257        }
258        let delta_ids_to_delete: Vec<risingwave_meta_model::HummockVersionId> =
259            hummock_time_travel_delta::Entity::find()
260                .select_only()
261                .column(hummock_time_travel_delta::Column::VersionId)
262                .filter(delta_delete_condition.clone())
263                .into_tuple()
264                .all(&txn)
265                .await?;
266        let delete_sst_batch_size = self
267            .env
268            .opts
269            .hummock_time_travel_epoch_version_insert_batch_size;
270        let delta_fetch_batch_size = self
271            .env
272            .opts
273            .hummock_time_travel_delta_fetch_batch_size
274            .max(1);
275        let mut sst_ids_to_delete: HashSet<_> = HashSet::default();
276        async fn delete_sst_in_batch(
277            txn: &DatabaseTransaction,
278            sst_ids_to_delete: HashSet<HummockSstableId>,
279            delete_sst_batch_size: usize,
280        ) -> Result<()> {
281            for start_idx in 0..=(sst_ids_to_delete.len().saturating_sub(1) / delete_sst_batch_size)
282            {
283                hummock_sstable_info::Entity::delete_many()
284                    .filter(
285                        hummock_sstable_info::Column::SstId.is_in(
286                            sst_ids_to_delete
287                                .iter()
288                                .skip(start_idx * delete_sst_batch_size)
289                                .take(delete_sst_batch_size)
290                                .copied(),
291                        ),
292                    )
293                    .exec(txn)
294                    .await?;
295            }
296            Ok(())
297        }
298        for delta_id_batch in delta_ids_to_delete.chunks(delta_fetch_batch_size) {
299            let mut delta_to_delete_by_id: HashMap<_, _> =
300                hummock_time_travel_delta::Entity::find()
301                    .filter(
302                        hummock_time_travel_delta::Column::VersionId
303                            .is_in(delta_id_batch.iter().copied()),
304                    )
305                    .all(&txn)
306                    .await?
307                    .into_iter()
308                    .map(|delta| (delta.version_id, delta))
309                    .collect();
310            for &delta_id_to_delete in delta_id_batch {
311                let delta_to_delete = delta_to_delete_by_id
312                    .remove(&delta_id_to_delete)
313                    .ok_or_else(|| {
314                        Error::TimeTravel(anyhow!(format!(
315                            "version delta {} not found",
316                            delta_id_to_delete
317                        )))
318                    })?;
319                let delta_to_delete = IncompleteHummockVersionDelta::from_persisted_protobuf_owned(
320                    delta_to_delete.version_delta.to_protobuf(),
321                );
322                let new_sst_ids = delta_to_delete.newly_added_sst_ids();
323                // The SST ids added and then deleted by compaction between the 2 versions.
324                sst_ids_to_delete.extend(&new_sst_ids - &retained_snapshot_sst_ids);
325                if sst_ids_to_delete.len() >= delete_sst_batch_size {
326                    delete_sst_in_batch(
327                        &txn,
328                        std::mem::take(&mut sst_ids_to_delete),
329                        delete_sst_batch_size,
330                    )
331                    .await?;
332                }
333                let new_object_ids = delta_to_delete.newly_added_object_ids();
334                object_ids_to_delete.extend(&new_object_ids - &retained_snapshot_object_ids);
335            }
336        }
337        for prev_version_id in version_ids_to_delete {
338            let prev_version = {
339                let prev_version = hummock_time_travel_version::Entity::find_by_id(prev_version_id)
340                    .one(&txn)
341                    .await?
342                    .ok_or_else(|| {
343                        Error::TimeTravel(anyhow!(format!(
344                            "prev_version {} not found",
345                            prev_version_id
346                        )))
347                    })?;
348                IncompleteHummockVersion::from_persisted_protobuf_owned(
349                    prev_version.version.to_protobuf(),
350                )
351            };
352            let sst_ids = prev_version.get_sst_ids();
353            sst_ids_to_delete.extend(&sst_ids - &retained_snapshot_sst_ids);
354            if sst_ids_to_delete.len() >= delete_sst_batch_size {
355                delete_sst_in_batch(
356                    &txn,
357                    std::mem::take(&mut sst_ids_to_delete),
358                    delete_sst_batch_size,
359                )
360                .await?;
361            }
362            let new_object_ids: HashSet<_> = prev_version.get_object_ids().collect();
363            object_ids_to_delete.extend(&new_object_ids - &retained_snapshot_object_ids);
364        }
365        if !sst_ids_to_delete.is_empty() {
366            delete_sst_in_batch(&txn, sst_ids_to_delete, delete_sst_batch_size).await?;
367        }
368
369        if !object_ids_to_delete.is_empty() {
370            // IMPORTANT: object_ids_to_delete may include objects that are still being used by SSTs not included in time travel metadata.
371            // So it's crucial to filter out those objects before actually deleting them, i.e. when using `try_take_may_delete_object_ids`.
372            self.gc_manager
373                .add_may_delete_object_ids(object_ids_to_delete.into_iter());
374        }
375
376        let res = hummock_time_travel_version::Entity::delete_many()
377            .filter(version_delete_condition)
378            .exec(&txn)
379            .await?;
380        tracing::info!(
381            %watermark_version_id,
382            %latest_valid_version_id,
383            "Deleted {} rows from hummock_time_travel_version.",
384            res.rows_affected
385        );
386
387        let res = hummock_time_travel_delta::Entity::delete_many()
388            .filter(delta_delete_condition)
389            .exec(&txn)
390            .await?;
391        tracing::info!(
392            %watermark_version_id,
393            %latest_valid_version_id,
394            "Deleted {} rows from hummock_time_travel_delta.",
395            res.rows_affected
396        );
397
398        txn.commit().await?;
399        Ok(())
400    }
401
402    pub(crate) async fn filter_out_objects_by_time_travel_v1(
403        &self,
404        objects: impl Iterator<Item = HummockObjectId>,
405    ) -> Result<HashSet<HummockObjectId>> {
406        let batch_size = self
407            .env
408            .opts
409            .hummock_time_travel_filter_out_objects_batch_size;
410        info!("filter out objects by time travel v1, only sst will remain in the result set");
411        // The input object count is much smaller than time travel pinned object count in meta store.
412        // So search input object in meta store.
413        let mut result: HashSet<_> = objects
414            .filter(|object_id| match object_id {
415                HummockObjectId::Sstable(_) => true,
416                HummockObjectId::VectorFile(_) | HummockObjectId::HnswGraphFile(_) => false,
417            })
418            .collect();
419        let mut remain_sst: VecDeque<_> = result.iter().copied().collect();
420        while !remain_sst.is_empty() {
421            let batch = remain_sst
422                .drain(..std::cmp::min(remain_sst.len(), batch_size))
423                .map(|object_id| object_id.as_raw());
424            let reject_object_ids: Vec<risingwave_meta_model::HummockSstableObjectId> =
425                hummock_sstable_info::Entity::find()
426                    .filter(hummock_sstable_info::Column::ObjectId.is_in(batch))
427                    .select_only()
428                    .column(hummock_sstable_info::Column::ObjectId)
429                    .into_tuple()
430                    .all(&self.env.meta_store_ref().conn)
431                    .await?;
432            for reject in reject_object_ids {
433                let object_id = HummockObjectId::Sstable(reject);
434                result.remove(&object_id);
435            }
436        }
437        Ok(result)
438    }
439
440    pub(crate) async fn filter_out_objects_by_time_travel(
441        &self,
442        objects: impl Iterator<Item = HummockObjectId>,
443    ) -> Result<HashSet<HummockObjectId>> {
444        if self.env.opts.hummock_time_travel_filter_out_objects_v1 {
445            return self.filter_out_objects_by_time_travel_v1(objects).await;
446        }
447        let mut result: HashSet<_> = objects.collect();
448
449        // filtered out object id pinned by time travel hummock version
450        {
451            let mut prev_version_id: Option<HummockVersionId> = None;
452            loop {
453                let query = hummock_time_travel_version::Entity::find();
454                let query = if let Some(prev_version_id) = prev_version_id {
455                    query.filter(hummock_time_travel_version::Column::VersionId.gt(prev_version_id))
456                } else {
457                    query
458                };
459                let mut version_stream = query
460                    .order_by_asc(hummock_time_travel_version::Column::VersionId)
461                    .limit(
462                        self.env
463                            .opts
464                            .hummock_time_travel_filter_out_objects_list_version_batch_size
465                            as u64,
466                    )
467                    .stream(&self.env.meta_store_ref().conn)
468                    .await?;
469                let mut next_prev_version_id = None;
470                while let Some(model) = version_stream.try_next().await? {
471                    let version =
472                        HummockVersion::from_persisted_protobuf_owned(model.version.to_protobuf());
473                    for object_id in version.get_object_ids() {
474                        result.remove(&object_id);
475                    }
476                    next_prev_version_id = Some(model.version_id);
477                }
478                if let Some(next_prev_version_id) = next_prev_version_id {
479                    prev_version_id = Some(next_prev_version_id);
480                } else {
481                    break;
482                }
483            }
484        }
485
486        // filtered out object ids pinned by time travel hummock version delta
487        {
488            let mut prev_version_id: Option<HummockVersionId> = None;
489            loop {
490                let query = hummock_time_travel_delta::Entity::find();
491                let query = if let Some(prev_version_id) = prev_version_id {
492                    query.filter(hummock_time_travel_delta::Column::VersionId.gt(prev_version_id))
493                } else {
494                    query
495                };
496                let mut version_stream = query
497                    .order_by_asc(hummock_time_travel_delta::Column::VersionId)
498                    .limit(
499                        self.env
500                            .opts
501                            .hummock_time_travel_filter_out_objects_list_delta_batch_size
502                            as u64,
503                    )
504                    .stream(&self.env.meta_store_ref().conn)
505                    .await?;
506                let mut next_prev_version_id = None;
507                while let Some(model) = version_stream.try_next().await? {
508                    let version_delta = HummockVersionDelta::from_persisted_protobuf_owned(
509                        model.version_delta.to_protobuf(),
510                    );
511                    for object_id in version_delta.newly_added_object_ids() {
512                        result.remove(&object_id);
513                    }
514                    next_prev_version_id = Some(model.version_id);
515                }
516                if let Some(next_prev_version_id) = next_prev_version_id {
517                    prev_version_id = Some(next_prev_version_id);
518                } else {
519                    break;
520                }
521            }
522        }
523
524        Ok(result)
525    }
526
527    pub(crate) async fn time_travel_pinned_object_count(&self) -> Result<u64> {
528        let count = hummock_sstable_info::Entity::find()
529            .count(&self.env.meta_store_ref().conn)
530            .await?;
531        Ok(count)
532    }
533
534    /// Attempt to locate the version corresponding to `query_epoch`.
535    ///
536    /// The version is retrieved from `hummock_epoch_to_version`, selecting the entry with the largest epoch that's lte `query_epoch`.
537    ///
538    /// The resulted version is complete, i.e. with correct `SstableInfo`.
539    pub async fn epoch_to_version(
540        &self,
541        query_epoch: HummockEpoch,
542        table_id: TableId,
543    ) -> Result<HummockVersion> {
544        let sql_store = self.env.meta_store_ref();
545        let _permit = self.inflight_time_travel_query.try_acquire().map_err(|_| {
546            anyhow!(format!(
547                "too many inflight time travel queries, max_inflight_time_travel_query={}",
548                self.env.opts.max_inflight_time_travel_query
549            ))
550        })?;
551        let epoch_to_version = hummock_epoch_to_version::Entity::find()
552            .filter(
553                Condition::any()
554                    .add(
555                        hummock_epoch_to_version::Column::TableId
556                            .eq(i64::from(table_id.as_raw_id())),
557                    )
558                    // for backward compatibility
559                    .add(hummock_epoch_to_version::Column::TableId.eq(0)),
560            )
561            .filter(
562                hummock_epoch_to_version::Column::Epoch
563                    .lte(risingwave_meta_model::Epoch::try_from(query_epoch).unwrap()),
564            )
565            .order_by_desc(hummock_epoch_to_version::Column::Epoch)
566            .one(&sql_store.conn)
567            .await?
568            .ok_or_else(|| Error::TimeTravelVersionExpired {
569                table_id,
570                epoch: query_epoch,
571            })?;
572        let timer = self
573            .metrics
574            .time_travel_version_replay_latency
575            .start_timer();
576        let actual_version_id = epoch_to_version.version_id;
577        tracing::debug!(
578            query_epoch,
579            query_tz = ?(Epoch(query_epoch).as_timestamptz()),
580            actual_epoch = epoch_to_version.epoch,
581            actual_tz = ?(Epoch(u64::try_from(epoch_to_version.epoch).unwrap()).as_timestamptz()),
582            %actual_version_id,
583            "convert query epoch"
584        );
585
586        let Some(resolved) =
587            resolve_time_travel_version(&sql_store.conn, actual_version_id).await?
588        else {
589            return Err(Error::TimeTravelVersionExpired {
590                table_id,
591                epoch: query_epoch,
592            });
593        };
594        let mut actual_version = replay_archive(
595            resolved.replay_version.to_protobuf(),
596            resolved.deltas.into_iter().map(|delta| delta.to_protobuf()),
597        )?;
598        if actual_version.id != actual_version_id {
599            return Err(Error::TimeTravelVersionExpired {
600                table_id,
601                epoch: query_epoch,
602            });
603        }
604
605        // SstableInfo in actual_version is incomplete before refill_version.
606        let mut sst_ids = actual_version
607            .get_sst_ids()
608            .into_iter()
609            .collect::<VecDeque<_>>();
610        let sst_count = sst_ids.len();
611        let mut sst_id_to_info = HashMap::with_capacity(sst_count);
612        let sst_info_fetch_batch_size = self.env.opts.hummock_time_travel_sst_info_fetch_batch_size;
613        while !sst_ids.is_empty() {
614            let sst_infos = hummock_sstable_info::Entity::find()
615                .filter(hummock_sstable_info::Column::SstId.is_in(
616                    sst_ids.drain(..std::cmp::min(sst_info_fetch_batch_size, sst_ids.len())),
617                ))
618                .all(&sql_store.conn)
619                .await?;
620            for sst_info in sst_infos {
621                let sst_info: SstableInfo = sst_info.sstable_info.to_protobuf().into();
622                sst_id_to_info.insert(sst_info.sst_id, sst_info);
623            }
624        }
625        if sst_count != sst_id_to_info.len() {
626            return Err(Error::TimeTravelVersionExpired {
627                table_id,
628                epoch: query_epoch,
629            });
630        }
631        refill_version(&mut actual_version, &sst_id_to_info, table_id);
632        timer.observe_duration();
633        Ok(actual_version)
634    }
635
636    pub(crate) async fn write_time_travel_metadata(
637        &self,
638        txn: &DatabaseTransaction,
639        version: Option<&HummockVersion>,
640        delta: HummockVersionDelta,
641        time_travel_table_ids: HashSet<StateTableId>,
642        skip_sst_ids: &HashSet<HummockSstableId>,
643        tables_to_commit: impl Iterator<Item = (&TableId, &CompactionGroupId, u64)>,
644    ) -> Result<Option<HashSet<HummockSstableId>>> {
645        let _timer = self
646            .metrics
647            .time_travel_write_metadata_latency
648            .start_timer();
649        if self
650            .env
651            .system_params_reader()
652            .await
653            .time_travel_retention_ms()
654            == 0
655        {
656            return Ok(None);
657        }
658        async fn write_sstable_infos(
659            mut sst_infos: impl Iterator<Item = &SstableInfo>,
660            txn: &DatabaseTransaction,
661            batch_size: usize,
662        ) -> Result<usize> {
663            let mut count = 0;
664            let mut is_finished = false;
665            while !is_finished {
666                let mut remain = batch_size;
667                let mut batch = vec![];
668                while remain > 0 {
669                    let Some(sst_info) = sst_infos.next() else {
670                        is_finished = true;
671                        break;
672                    };
673                    batch.push(hummock_sstable_info::ActiveModel {
674                        sst_id: Set(sst_info.sst_id),
675                        object_id: Set(sst_info.object_id),
676                        sstable_info: Set(SstableInfoV2Backend::from(&sst_info.to_protobuf())),
677                    });
678                    remain -= 1;
679                    count += 1;
680                }
681                if batch.is_empty() {
682                    break;
683                }
684                hummock_sstable_info::Entity::insert_many(batch)
685                    .on_conflict_do_nothing()
686                    .exec(txn)
687                    .await?;
688            }
689            Ok(count)
690        }
691
692        let mut batch = vec![];
693        for (table_id, _cg_id, committed_epoch) in tables_to_commit {
694            if !time_travel_table_ids.contains(table_id) {
695                continue;
696            }
697            let m = hummock_epoch_to_version::ActiveModel {
698                epoch: Set(committed_epoch.try_into().unwrap()),
699                table_id: Set(i64::from(table_id.as_raw_id())),
700                version_id: Set(delta.id),
701            };
702            batch.push(m);
703            if batch.len()
704                >= self
705                    .env
706                    .opts
707                    .hummock_time_travel_epoch_version_insert_batch_size
708            {
709                // There should be no conflict rows.
710                hummock_epoch_to_version::Entity::insert_many(std::mem::take(&mut batch))
711                    .do_nothing()
712                    .exec(txn)
713                    .await?;
714            }
715        }
716        if !batch.is_empty() {
717            // There should be no conflict rows.
718            hummock_epoch_to_version::Entity::insert_many(batch)
719                .do_nothing()
720                .exec(txn)
721                .await?;
722        }
723
724        let mut version_sst_ids = None;
725        if let Some(version) = version {
726            // `version_sst_ids` is used to update `last_time_travel_snapshot_sst_ids`.
727            version_sst_ids = Some(
728                version
729                    .get_sst_infos()
730                    .filter_map(|s| {
731                        if s.table_ids
732                            .iter()
733                            .any(|tid| time_travel_table_ids.contains(tid))
734                        {
735                            return Some(s.sst_id);
736                        }
737                        None
738                    })
739                    .collect(),
740            );
741            write_sstable_infos(
742                version.get_sst_infos().filter(|s| {
743                    !skip_sst_ids.contains(&s.sst_id)
744                        && s.table_ids
745                            .iter()
746                            .any(|tid| time_travel_table_ids.contains(tid))
747                }),
748                txn,
749                self.env.opts.hummock_time_travel_sst_info_insert_batch_size,
750            )
751            .await?;
752            let m = hummock_time_travel_version::ActiveModel {
753                version_id: Set(version.id),
754                version: Set(
755                    (&IncompleteHummockVersion::from((version, &time_travel_table_ids))
756                        .to_protobuf())
757                        .into(),
758                ),
759            };
760            hummock_time_travel_version::Entity::insert(m)
761                .on_conflict_do_nothing()
762                .exec(txn)
763                .await?;
764            // Return early to skip persisting delta.
765            return Ok(version_sst_ids);
766        }
767        let written = write_sstable_infos(
768            delta.newly_added_sst_infos().filter(|s| {
769                !skip_sst_ids.contains(&s.sst_id)
770                    && s.table_ids
771                        .iter()
772                        .any(|tid| time_travel_table_ids.contains(tid))
773            }),
774            txn,
775            self.env.opts.hummock_time_travel_sst_info_insert_batch_size,
776        )
777        .await?;
778        let has_state_table_info_delta = delta
779            .state_table_info_delta
780            .keys()
781            .any(|table_id| time_travel_table_ids.contains(table_id));
782        if written > 0 || has_state_table_info_delta {
783            let m = hummock_time_travel_delta::ActiveModel {
784                version_id: Set(delta.id),
785                version_delta: Set((&IncompleteHummockVersionDelta::from((
786                    &delta,
787                    &time_travel_table_ids,
788                ))
789                .to_protobuf())
790                    .into()),
791            };
792            hummock_time_travel_delta::Entity::insert(m)
793                .on_conflict_do_nothing()
794                .exec(txn)
795                .await?;
796        }
797
798        Ok(version_sst_ids)
799    }
800}
801
802struct ResolvedTimeTravelVersion {
803    replay_version: IncompleteHummockVersion,
804    deltas: Vec<IncompleteHummockVersionDelta>,
805}
806
807async fn resolve_time_travel_version(
808    conn: &impl ConnectionTrait,
809    actual_version_id: HummockVersionId,
810) -> Result<Option<ResolvedTimeTravelVersion>> {
811    let Some(replay_version) = hummock_time_travel_version::Entity::find()
812        .filter(hummock_time_travel_version::Column::VersionId.lte(actual_version_id))
813        .order_by_desc(hummock_time_travel_version::Column::VersionId)
814        .one(conn)
815        .await?
816    else {
817        return Ok(None);
818    };
819    let deltas = hummock_time_travel_delta::Entity::find()
820        .filter(hummock_time_travel_delta::Column::VersionId.gt(replay_version.version_id))
821        .filter(hummock_time_travel_delta::Column::VersionId.lte(actual_version_id))
822        .order_by_asc(hummock_time_travel_delta::Column::VersionId)
823        .all(conn)
824        .await?;
825    Ok(Some(ResolvedTimeTravelVersion {
826        replay_version: IncompleteHummockVersion::from_persisted_protobuf_owned(
827            replay_version.version.to_protobuf(),
828        ),
829        deltas: deltas
830            .into_iter()
831            .map(|delta| {
832                IncompleteHummockVersionDelta::from_persisted_protobuf_owned(
833                    delta.version_delta.to_protobuf(),
834                )
835            })
836            .collect(),
837    }))
838}
839
840/// The `HummockVersion` is actually `InHummockVersion`. It requires `refill_version`.
841fn replay_archive(
842    version: PbHummockVersion,
843    deltas: impl Iterator<Item = PbHummockVersionDelta>,
844) -> Result<HummockVersion> {
845    // The pb version ann pb version delta are actually written by InHummockVersion and InHummockVersionDelta, respectively.
846    // Using HummockVersion make it easier for `refill_version` later.
847    let mut last_version = HummockVersion::from_persisted_protobuf_owned(version);
848    for d in deltas {
849        let d = HummockVersionDelta::from_persisted_protobuf_owned(d);
850        debug_assert!(
851            !should_mark_next_time_travel_version_snapshot(&d),
852            "unexpected time travel delta {:?}",
853            d
854        );
855        if d.prev_id < last_version.id {
856            return Err(Error::TimeTravel(anyhow!(format!(
857                "invalid time travel delta chain: delta {} has prev version {}, but replay has reached {}",
858                d.id, d.prev_id, last_version.id
859            ))));
860        }
861        // Compaction deltas are not included in the time travel archive, so there may be gaps
862        // between the last replayed version and this delta's previous version.
863        last_version.id = d.prev_id;
864        last_version.apply_version_delta(&d);
865    }
866    Ok(last_version)
867}
868
869pub fn require_sql_meta_store_err() -> Error {
870    Error::TimeTravel(anyhow!("require SQL meta store"))
871}
872
873/// Time travel delta replay only expect `NewL0SubLevel`. In all other cases, a new version snapshot should be created.
874pub fn should_mark_next_time_travel_version_snapshot(delta: &HummockVersionDelta) -> bool {
875    delta.group_deltas.iter().any(|(_, deltas)| {
876        deltas
877            .group_deltas
878            .iter()
879            .any(|d| !matches!(d, GroupDeltaCommon::NewL0SubLevel(_)))
880    })
881}
882
883#[cfg(test)]
884mod tests {
885    use super::*;
886
887    fn version(id: u64) -> PbHummockVersion {
888        let mut version = HummockVersion::default();
889        version.id = id.into();
890        version.to_protobuf()
891    }
892
893    fn delta(id: u64, prev_id: u64) -> PbHummockVersionDelta {
894        let mut delta = HummockVersionDelta::default();
895        delta.id = id.into();
896        delta.prev_id = prev_id.into();
897        delta.to_protobuf()
898    }
899
900    #[test]
901    fn test_replay_archive_delta_chain() {
902        let replayed = replay_archive(version(1), [delta(4, 3), delta(5, 4)].into_iter()).unwrap();
903        assert_eq!(replayed.id, HummockVersionId::new(5));
904
905        assert!(replay_archive(version(3), [delta(4, 2)].into_iter()).is_err());
906    }
907}