1use 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
45impl 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 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 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 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 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 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 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 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 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 {
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 {
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 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 .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 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 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 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 = 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 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
850fn replay_archive(
852 version: PbHummockVersion,
853 deltas: impl Iterator<Item = PbHummockVersionDelta>,
854) -> Result<HummockVersion> {
855 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 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
883pub 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}