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