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            assert!(delete_sst_batch_size > 0);
282            let mut sst_ids_to_delete = sst_ids_to_delete.into_iter();
283            while sst_ids_to_delete.len() > 0 {
284                hummock_sstable_info::Entity::delete_many()
285                    .filter(
286                        hummock_sstable_info::Column::SstId
287                            .is_in(sst_ids_to_delete.by_ref().take(delete_sst_batch_size)),
288                    )
289                    .exec(txn)
290                    .await?;
291            }
292            Ok(())
293        }
294        for delta_id_batch in delta_ids_to_delete.chunks(delta_fetch_batch_size) {
295            let mut delta_to_delete_by_id: HashMap<_, _> =
296                hummock_time_travel_delta::Entity::find()
297                    .filter(
298                        hummock_time_travel_delta::Column::VersionId
299                            .is_in(delta_id_batch.iter().copied()),
300                    )
301                    .all(&txn)
302                    .await?
303                    .into_iter()
304                    .map(|delta| (delta.version_id, delta))
305                    .collect();
306            for &delta_id_to_delete in delta_id_batch {
307                let delta_to_delete = delta_to_delete_by_id
308                    .remove(&delta_id_to_delete)
309                    .ok_or_else(|| {
310                        Error::TimeTravel(anyhow!(format!(
311                            "version delta {} not found",
312                            delta_id_to_delete
313                        )))
314                    })?;
315                let delta_to_delete = IncompleteHummockVersionDelta::from_persisted_protobuf_owned(
316                    delta_to_delete.version_delta.to_protobuf(),
317                );
318                let new_sst_ids = delta_to_delete.newly_added_sst_ids();
319                // The SST ids added and then deleted by compaction between the 2 versions.
320                sst_ids_to_delete.extend(new_sst_ids.difference(&retained_snapshot_sst_ids));
321                if sst_ids_to_delete.len() >= delete_sst_batch_size {
322                    delete_sst_in_batch(
323                        &txn,
324                        std::mem::take(&mut sst_ids_to_delete),
325                        delete_sst_batch_size,
326                    )
327                    .await?;
328                }
329                let new_object_ids = delta_to_delete.newly_added_object_ids();
330                object_ids_to_delete
331                    .extend(new_object_ids.difference(&retained_snapshot_object_ids));
332            }
333        }
334        for prev_version_id in version_ids_to_delete {
335            let prev_version = {
336                let prev_version = hummock_time_travel_version::Entity::find_by_id(prev_version_id)
337                    .one(&txn)
338                    .await?
339                    .ok_or_else(|| {
340                        Error::TimeTravel(anyhow!(format!(
341                            "prev_version {} not found",
342                            prev_version_id
343                        )))
344                    })?;
345                IncompleteHummockVersion::from_persisted_protobuf_owned(
346                    prev_version.version.to_protobuf(),
347                )
348            };
349            let sst_ids = prev_version.get_sst_ids();
350            sst_ids_to_delete.extend(sst_ids.difference(&retained_snapshot_sst_ids));
351            if sst_ids_to_delete.len() >= delete_sst_batch_size {
352                delete_sst_in_batch(
353                    &txn,
354                    std::mem::take(&mut sst_ids_to_delete),
355                    delete_sst_batch_size,
356                )
357                .await?;
358            }
359            let new_object_ids: HashSet<_> = prev_version.get_object_ids().collect();
360            object_ids_to_delete.extend(new_object_ids.difference(&retained_snapshot_object_ids));
361        }
362        if !sst_ids_to_delete.is_empty() {
363            delete_sst_in_batch(&txn, sst_ids_to_delete, delete_sst_batch_size).await?;
364        }
365
366        if !object_ids_to_delete.is_empty() {
367            // IMPORTANT: object_ids_to_delete may include objects that are still being used by SSTs not included in time travel metadata.
368            // So it's crucial to filter out those objects before actually deleting them, i.e. when using `try_take_may_delete_object_ids`.
369            self.gc_manager
370                .add_may_delete_object_ids(object_ids_to_delete.into_iter());
371        }
372
373        let res = hummock_time_travel_version::Entity::delete_many()
374            .filter(version_delete_condition)
375            .exec(&txn)
376            .await?;
377        tracing::info!(
378            %watermark_version_id,
379            %latest_valid_version_id,
380            "Deleted {} rows from hummock_time_travel_version.",
381            res.rows_affected
382        );
383
384        let res = hummock_time_travel_delta::Entity::delete_many()
385            .filter(delta_delete_condition)
386            .exec(&txn)
387            .await?;
388        tracing::info!(
389            %watermark_version_id,
390            %latest_valid_version_id,
391            "Deleted {} rows from hummock_time_travel_delta.",
392            res.rows_affected
393        );
394
395        txn.commit().await?;
396        Ok(())
397    }
398
399    pub(crate) async fn filter_out_objects_by_time_travel_v1(
400        &self,
401        objects: impl Iterator<Item = HummockObjectId>,
402    ) -> Result<HashSet<HummockObjectId>> {
403        let batch_size = self
404            .env
405            .opts
406            .hummock_time_travel_filter_out_objects_batch_size;
407        info!("filter out objects by time travel v1, only sst will remain in the result set");
408        // The input object count is much smaller than time travel pinned object count in meta store.
409        // So search input object in meta store.
410        let mut result: HashSet<_> = objects
411            .filter(|object_id| match object_id {
412                HummockObjectId::Sstable(_) => true,
413                HummockObjectId::VectorFile(_) | HummockObjectId::HnswGraphFile(_) => false,
414            })
415            .collect();
416        let mut remain_sst: VecDeque<_> = result.iter().copied().collect();
417        while !remain_sst.is_empty() {
418            let batch = remain_sst
419                .drain(..std::cmp::min(remain_sst.len(), batch_size))
420                .map(|object_id| object_id.as_raw());
421            let reject_object_ids: Vec<risingwave_meta_model::HummockSstableObjectId> =
422                hummock_sstable_info::Entity::find()
423                    .filter(hummock_sstable_info::Column::ObjectId.is_in(batch))
424                    .select_only()
425                    .column(hummock_sstable_info::Column::ObjectId)
426                    .into_tuple()
427                    .all(&self.env.meta_store_ref().conn)
428                    .await?;
429            for reject in reject_object_ids {
430                let object_id = HummockObjectId::Sstable(reject);
431                result.remove(&object_id);
432            }
433        }
434        Ok(result)
435    }
436
437    /// Removes candidates referenced by retained time-travel metadata.
438    ///
439    /// The v2 scan only removes candidates, so it stops once none remain. It does not validate
440    /// the rest of the archive or report errors from skipped reads.
441    pub(crate) async fn filter_out_objects_by_time_travel(
442        &self,
443        objects: impl Iterator<Item = HummockObjectId>,
444    ) -> Result<HashSet<HummockObjectId>> {
445        if self.env.opts.hummock_time_travel_filter_out_objects_v1 {
446            return self.filter_out_objects_by_time_travel_v1(objects).await;
447        }
448        let mut result: HashSet<_> = objects.collect();
449        if result.is_empty() {
450            return Ok(result);
451        }
452
453        // filtered out object id pinned by time travel hummock version
454        {
455            let mut prev_version_id: Option<HummockVersionId> = None;
456            loop {
457                let query = hummock_time_travel_version::Entity::find();
458                let query = if let Some(prev_version_id) = prev_version_id {
459                    query.filter(hummock_time_travel_version::Column::VersionId.gt(prev_version_id))
460                } else {
461                    query
462                };
463                let mut version_stream = query
464                    .order_by_asc(hummock_time_travel_version::Column::VersionId)
465                    .limit(
466                        self.env
467                            .opts
468                            .hummock_time_travel_filter_out_objects_list_version_batch_size
469                            as u64,
470                    )
471                    .stream(&self.env.meta_store_ref().conn)
472                    .await?;
473                let mut next_prev_version_id = None;
474                while let Some(model) = version_stream.try_next().await? {
475                    let version =
476                        HummockVersion::from_persisted_protobuf_owned(model.version.to_protobuf());
477                    for object_id in version.get_object_ids() {
478                        result.remove(&object_id);
479                    }
480                    if result.is_empty() {
481                        return Ok(result);
482                    }
483                    next_prev_version_id = Some(model.version_id);
484                }
485                if let Some(next_prev_version_id) = next_prev_version_id {
486                    prev_version_id = Some(next_prev_version_id);
487                } else {
488                    break;
489                }
490            }
491        }
492
493        // filtered out object ids pinned by time travel hummock version delta
494        {
495            let mut prev_version_id: Option<HummockVersionId> = None;
496            loop {
497                let query = hummock_time_travel_delta::Entity::find();
498                let query = if let Some(prev_version_id) = prev_version_id {
499                    query.filter(hummock_time_travel_delta::Column::VersionId.gt(prev_version_id))
500                } else {
501                    query
502                };
503                let mut version_stream = query
504                    .order_by_asc(hummock_time_travel_delta::Column::VersionId)
505                    .limit(
506                        self.env
507                            .opts
508                            .hummock_time_travel_filter_out_objects_list_delta_batch_size
509                            as u64,
510                    )
511                    .stream(&self.env.meta_store_ref().conn)
512                    .await?;
513                let mut next_prev_version_id = None;
514                while let Some(model) = version_stream.try_next().await? {
515                    let version_delta = HummockVersionDelta::from_persisted_protobuf_owned(
516                        model.version_delta.to_protobuf(),
517                    );
518                    for object_id in version_delta.newly_added_object_ids() {
519                        result.remove(&object_id);
520                    }
521                    if result.is_empty() {
522                        return Ok(result);
523                    }
524                    next_prev_version_id = Some(model.version_id);
525                }
526                if let Some(next_prev_version_id) = next_prev_version_id {
527                    prev_version_id = Some(next_prev_version_id);
528                } else {
529                    break;
530                }
531            }
532        }
533
534        Ok(result)
535    }
536
537    pub(crate) async fn time_travel_pinned_object_count(&self) -> Result<u64> {
538        let count = hummock_sstable_info::Entity::find()
539            .count(&self.env.meta_store_ref().conn)
540            .await?;
541        Ok(count)
542    }
543
544    /// Attempt to locate the version corresponding to `query_epoch`.
545    ///
546    /// The version is retrieved from `hummock_epoch_to_version`, selecting the entry with the largest epoch that's lte `query_epoch`.
547    ///
548    /// The resulted version is complete, i.e. with correct `SstableInfo`.
549    pub async fn epoch_to_version(
550        &self,
551        query_epoch: HummockEpoch,
552        table_id: TableId,
553    ) -> Result<HummockVersion> {
554        let sql_store = self.env.meta_store_ref();
555        let _permit = self.inflight_time_travel_query.try_acquire().map_err(|_| {
556            anyhow!(format!(
557                "too many inflight time travel queries, max_inflight_time_travel_query={}",
558                self.env.opts.max_inflight_time_travel_query
559            ))
560        })?;
561        let epoch_to_version = hummock_epoch_to_version::Entity::find()
562            .filter(
563                Condition::any()
564                    .add(
565                        hummock_epoch_to_version::Column::TableId
566                            .eq(i64::from(table_id.as_raw_id())),
567                    )
568                    // for backward compatibility
569                    .add(hummock_epoch_to_version::Column::TableId.eq(0)),
570            )
571            .filter(
572                hummock_epoch_to_version::Column::Epoch
573                    .lte(risingwave_meta_model::Epoch::try_from(query_epoch).unwrap()),
574            )
575            .order_by_desc(hummock_epoch_to_version::Column::Epoch)
576            .one(&sql_store.conn)
577            .await?
578            .ok_or_else(|| Error::TimeTravelVersionExpired {
579                table_id,
580                epoch: query_epoch,
581            })?;
582        let timer = self
583            .metrics
584            .time_travel_version_replay_latency
585            .start_timer();
586        let actual_version_id = epoch_to_version.version_id;
587        tracing::debug!(
588            query_epoch,
589            query_tz = ?(Epoch(query_epoch).as_timestamptz()),
590            actual_epoch = epoch_to_version.epoch,
591            actual_tz = ?(Epoch(u64::try_from(epoch_to_version.epoch).unwrap()).as_timestamptz()),
592            %actual_version_id,
593            "convert query epoch"
594        );
595
596        let Some(resolved) =
597            resolve_time_travel_version(&sql_store.conn, actual_version_id).await?
598        else {
599            return Err(Error::TimeTravelVersionExpired {
600                table_id,
601                epoch: query_epoch,
602            });
603        };
604        let mut actual_version = replay_archive(
605            resolved.replay_version.to_protobuf(),
606            resolved.deltas.into_iter().map(|delta| delta.to_protobuf()),
607        )?;
608        if actual_version.id != actual_version_id {
609            return Err(Error::TimeTravelVersionExpired {
610                table_id,
611                epoch: query_epoch,
612            });
613        }
614
615        // SstableInfo in actual_version is incomplete before refill_version.
616        let mut sst_ids = actual_version
617            .get_sst_ids()
618            .into_iter()
619            .collect::<VecDeque<_>>();
620        let sst_count = sst_ids.len();
621        let mut sst_id_to_info = HashMap::with_capacity(sst_count);
622        let sst_info_fetch_batch_size = self.env.opts.hummock_time_travel_sst_info_fetch_batch_size;
623        while !sst_ids.is_empty() {
624            let sst_infos = hummock_sstable_info::Entity::find()
625                .filter(hummock_sstable_info::Column::SstId.is_in(
626                    sst_ids.drain(..std::cmp::min(sst_info_fetch_batch_size, sst_ids.len())),
627                ))
628                .all(&sql_store.conn)
629                .await?;
630            for sst_info in sst_infos {
631                let sst_info: SstableInfo = sst_info.sstable_info.to_protobuf().into();
632                sst_id_to_info.insert(sst_info.sst_id, sst_info);
633            }
634        }
635        if sst_count != sst_id_to_info.len() {
636            return Err(Error::TimeTravelVersionExpired {
637                table_id,
638                epoch: query_epoch,
639            });
640        }
641        refill_version(&mut actual_version, &sst_id_to_info, table_id);
642        timer.observe_duration();
643        Ok(actual_version)
644    }
645
646    pub(crate) async fn write_time_travel_metadata(
647        &self,
648        txn: &DatabaseTransaction,
649        version: Option<&HummockVersion>,
650        delta: HummockVersionDelta,
651        time_travel_table_ids: HashSet<StateTableId>,
652        skip_sst_ids: &HashSet<HummockSstableId>,
653        tables_to_commit: impl Iterator<Item = (&TableId, &CompactionGroupId, u64)>,
654    ) -> Result<Option<HashSet<HummockSstableId>>> {
655        let _timer = self
656            .metrics
657            .time_travel_write_metadata_latency
658            .start_timer();
659        if self
660            .env
661            .system_params_reader()
662            .await
663            .time_travel_retention_ms()
664            == 0
665        {
666            return Ok(None);
667        }
668        async fn write_sstable_infos(
669            mut sst_infos: impl Iterator<Item = &SstableInfo>,
670            txn: &DatabaseTransaction,
671            batch_size: usize,
672        ) -> Result<usize> {
673            let mut count = 0;
674            let mut is_finished = false;
675            while !is_finished {
676                let mut remain = batch_size;
677                let mut batch = vec![];
678                while remain > 0 {
679                    let Some(sst_info) = sst_infos.next() else {
680                        is_finished = true;
681                        break;
682                    };
683                    batch.push(hummock_sstable_info::ActiveModel {
684                        sst_id: Set(sst_info.sst_id),
685                        object_id: Set(sst_info.object_id),
686                        sstable_info: Set(SstableInfoV2Backend::from(&sst_info.to_protobuf())),
687                    });
688                    remain -= 1;
689                    count += 1;
690                }
691                if batch.is_empty() {
692                    break;
693                }
694                hummock_sstable_info::Entity::insert_many(batch)
695                    .on_conflict_do_nothing()
696                    .exec(txn)
697                    .await?;
698            }
699            Ok(count)
700        }
701
702        let mut batch = vec![];
703        for (table_id, _cg_id, committed_epoch) in tables_to_commit {
704            if !time_travel_table_ids.contains(table_id) {
705                continue;
706            }
707            let m = hummock_epoch_to_version::ActiveModel {
708                epoch: Set(committed_epoch.try_into().unwrap()),
709                table_id: Set(i64::from(table_id.as_raw_id())),
710                version_id: Set(delta.id),
711            };
712            batch.push(m);
713            if batch.len()
714                >= self
715                    .env
716                    .opts
717                    .hummock_time_travel_epoch_version_insert_batch_size
718            {
719                // There should be no conflict rows.
720                hummock_epoch_to_version::Entity::insert_many(std::mem::take(&mut batch))
721                    .do_nothing()
722                    .exec(txn)
723                    .await?;
724            }
725        }
726        if !batch.is_empty() {
727            // There should be no conflict rows.
728            hummock_epoch_to_version::Entity::insert_many(batch)
729                .do_nothing()
730                .exec(txn)
731                .await?;
732        }
733
734        let mut version_sst_ids = None;
735        if let Some(version) = version {
736            // `version_sst_ids` is used to update `last_time_travel_snapshot_sst_ids`.
737            version_sst_ids = Some(
738                version
739                    .get_sst_infos()
740                    .filter_map(|s| {
741                        if s.table_ids
742                            .iter()
743                            .any(|tid| time_travel_table_ids.contains(tid))
744                        {
745                            return Some(s.sst_id);
746                        }
747                        None
748                    })
749                    .collect(),
750            );
751            write_sstable_infos(
752                version.get_sst_infos().filter(|s| {
753                    !skip_sst_ids.contains(&s.sst_id)
754                        && s.table_ids
755                            .iter()
756                            .any(|tid| time_travel_table_ids.contains(tid))
757                }),
758                txn,
759                self.env.opts.hummock_time_travel_sst_info_insert_batch_size,
760            )
761            .await?;
762            let m = hummock_time_travel_version::ActiveModel {
763                version_id: Set(version.id),
764                version: Set(
765                    (&IncompleteHummockVersion::from((version, &time_travel_table_ids))
766                        .to_protobuf())
767                        .into(),
768                ),
769            };
770            hummock_time_travel_version::Entity::insert(m)
771                .on_conflict_do_nothing()
772                .exec(txn)
773                .await?;
774            // Return early to skip persisting delta.
775            return Ok(version_sst_ids);
776        }
777        let written = write_sstable_infos(
778            delta.newly_added_sst_infos().filter(|s| {
779                !skip_sst_ids.contains(&s.sst_id)
780                    && s.table_ids
781                        .iter()
782                        .any(|tid| time_travel_table_ids.contains(tid))
783            }),
784            txn,
785            self.env.opts.hummock_time_travel_sst_info_insert_batch_size,
786        )
787        .await?;
788        let has_state_table_info_delta = delta
789            .state_table_info_delta
790            .keys()
791            .any(|table_id| time_travel_table_ids.contains(table_id));
792        if written > 0 || has_state_table_info_delta {
793            let m = hummock_time_travel_delta::ActiveModel {
794                version_id: Set(delta.id),
795                version_delta: Set((&IncompleteHummockVersionDelta::from((
796                    &delta,
797                    &time_travel_table_ids,
798                ))
799                .to_protobuf())
800                    .into()),
801            };
802            hummock_time_travel_delta::Entity::insert(m)
803                .on_conflict_do_nothing()
804                .exec(txn)
805                .await?;
806        }
807
808        Ok(version_sst_ids)
809    }
810}
811
812struct ResolvedTimeTravelVersion {
813    replay_version: IncompleteHummockVersion,
814    deltas: Vec<IncompleteHummockVersionDelta>,
815}
816
817async fn resolve_time_travel_version(
818    conn: &impl ConnectionTrait,
819    actual_version_id: HummockVersionId,
820) -> Result<Option<ResolvedTimeTravelVersion>> {
821    let Some(replay_version) = hummock_time_travel_version::Entity::find()
822        .filter(hummock_time_travel_version::Column::VersionId.lte(actual_version_id))
823        .order_by_desc(hummock_time_travel_version::Column::VersionId)
824        .one(conn)
825        .await?
826    else {
827        return Ok(None);
828    };
829    let deltas = hummock_time_travel_delta::Entity::find()
830        .filter(hummock_time_travel_delta::Column::VersionId.gt(replay_version.version_id))
831        .filter(hummock_time_travel_delta::Column::VersionId.lte(actual_version_id))
832        .order_by_asc(hummock_time_travel_delta::Column::VersionId)
833        .all(conn)
834        .await?;
835    Ok(Some(ResolvedTimeTravelVersion {
836        replay_version: IncompleteHummockVersion::from_persisted_protobuf_owned(
837            replay_version.version.to_protobuf(),
838        ),
839        deltas: deltas
840            .into_iter()
841            .map(|delta| {
842                IncompleteHummockVersionDelta::from_persisted_protobuf_owned(
843                    delta.version_delta.to_protobuf(),
844                )
845            })
846            .collect(),
847    }))
848}
849
850/// The `HummockVersion` is actually `InHummockVersion`. It requires `refill_version`.
851fn replay_archive(
852    version: PbHummockVersion,
853    deltas: impl Iterator<Item = PbHummockVersionDelta>,
854) -> Result<HummockVersion> {
855    // The pb version ann pb version delta are actually written by InHummockVersion and InHummockVersionDelta, respectively.
856    // Using HummockVersion make it easier for `refill_version` later.
857    let mut last_version = HummockVersion::from_persisted_protobuf_owned(version);
858    for d in deltas {
859        let d = HummockVersionDelta::from_persisted_protobuf_owned(d);
860        debug_assert!(
861            !should_mark_next_time_travel_version_snapshot(&d),
862            "unexpected time travel delta {:?}",
863            d
864        );
865        if d.prev_id < last_version.id {
866            return Err(Error::TimeTravel(anyhow!(format!(
867                "invalid time travel delta chain: delta {} has prev version {}, but replay has reached {}",
868                d.id, d.prev_id, last_version.id
869            ))));
870        }
871        // Compaction deltas are not included in the time travel archive, so there may be gaps
872        // between the last replayed version and this delta's previous version.
873        last_version.id = d.prev_id;
874        last_version.apply_version_delta(&d);
875    }
876    Ok(last_version)
877}
878
879pub fn require_sql_meta_store_err() -> Error {
880    Error::TimeTravel(anyhow!("require SQL meta store"))
881}
882
883/// Time travel delta replay only expect `NewL0SubLevel`. In all other cases, a new version snapshot should be created.
884pub fn should_mark_next_time_travel_version_snapshot(delta: &HummockVersionDelta) -> bool {
885    delta.group_deltas.iter().any(|(_, deltas)| {
886        deltas
887            .group_deltas
888            .iter()
889            .any(|d| !matches!(d, GroupDeltaCommon::NewL0SubLevel(_)))
890    })
891}
892
893#[cfg(test)]
894mod tests {
895    use super::*;
896
897    fn version(id: u64) -> PbHummockVersion {
898        let mut version = HummockVersion::default();
899        version.id = id.into();
900        version.to_protobuf()
901    }
902
903    fn delta(id: u64, prev_id: u64) -> PbHummockVersionDelta {
904        let mut delta = HummockVersionDelta::default();
905        delta.id = id.into();
906        delta.prev_id = prev_id.into();
907        delta.to_protobuf()
908    }
909
910    #[test]
911    fn test_replay_archive_delta_chain() {
912        let replayed = replay_archive(version(1), [delta(4, 3), delta(5, 4)].into_iter()).unwrap();
913        assert_eq!(replayed.id, HummockVersionId::new(5));
914
915        assert!(replay_archive(version(3), [delta(4, 2)].into_iter()).is_err());
916    }
917}