1use std::{
2 borrow::Borrow,
3 env,
4 path::PathBuf,
5 sync::{Arc, LazyLock, Mutex, PoisonError, Weak},
6};
7
8use anyhow::{Context, Result};
9use smallvec::SmallVec;
10use turbo_bincode::{new_turbo_bincode_decoder, turbo_bincode_decode, turbo_bincode_encode};
11use turbo_persistence::CommitStats;
12use turbo_tasks::{
13 DynTaskInputs, RawVc, TaskId,
14 macro_helpers::NativeFunction,
15 panic_hooks::{PanicHookGuard, register_panic_hook},
16 parallel,
17};
18
19use crate::{
20 GitVersionInfo,
21 backend::{SpecificTaskDataCategory, TtlCounter, storage_schema::TaskStorage},
22 backing_storage::{SnapshotItem, SnapshotMeta, compute_task_type_hash_from_components},
23 database::{
24 db_invalidation::{StartupCacheState, check_db_invalidation_and_cleanup, invalidate_db},
25 db_versioning::handle_db_versioning,
26 key_value_database::KeySpace,
27 turbo::{TurboKeyValueDatabase, TurboWriteBatch},
28 write_batch::WriteBuffer,
29 },
30 db_invalidation::invalidation_reasons,
31};
32
33#[derive(Clone, Copy)]
35#[repr(u8)]
36enum InfraKey {
37 NextFreeTaskId = 0,
38 GcRoots = 1,
39}
40
41impl InfraKey {
42 fn key(self) -> ByteKey {
43 ByteKey::new(self as u8)
44 }
45}
46
47struct ByteKey([u8; 1]);
48
49impl ByteKey {
50 fn new(value: u8) -> Self {
51 Self([value])
52 }
53}
54
55impl AsRef<[u8]> for ByteKey {
56 fn as_ref(&self) -> &[u8] {
57 &self.0
58 }
59}
60
61struct IntKey([u8; 4]);
62
63impl IntKey {
64 fn new(value: u32) -> Self {
65 Self(value.to_le_bytes())
66 }
67}
68
69impl AsRef<[u8]> for IntKey {
70 fn as_ref(&self) -> &[u8] {
71 &self.0
72 }
73}
74
75fn as_u32(bytes: impl Borrow<[u8]>) -> Result<u32> {
76 let n = u32::from_le_bytes(bytes.borrow().try_into()?);
77 Ok(n)
78}
79
80fn should_invalidate_on_panic() -> bool {
89 fn env_is_falsy(key: &str) -> bool {
90 env::var_os(key)
91 .is_none_or(|value| ["".as_ref(), "0".as_ref(), "false".as_ref()].contains(&&*value))
92 }
93 static SHOULD_INVALIDATE: LazyLock<bool> = LazyLock::new(|| {
94 env_is_falsy("TURBO_ENGINE_SKIP_INVALIDATE_ON_PANIC") && env_is_falsy("__NEXT_TEST_MODE")
95 });
96 *SHOULD_INVALIDATE
97}
98
99struct TurboBackingStorageInner {
100 database: TurboKeyValueDatabase,
101 base_path: Option<PathBuf>,
104 invalidated: Mutex<bool>,
106 _panic_hook_guard: Option<PanicHookGuard>,
109}
110
111pub struct TurboBackingStorage {
119 inner: Arc<TurboBackingStorageInner>,
121}
122
123impl TurboBackingStorage {
124 pub(crate) fn new_in_memory(database: TurboKeyValueDatabase) -> Self {
125 Self {
126 inner: Arc::new(TurboBackingStorageInner {
127 database,
128 base_path: None,
129 invalidated: Mutex::new(false),
130 _panic_hook_guard: None,
131 }),
132 }
133 }
134
135 pub(crate) fn open_versioned_on_disk(
147 base_path: PathBuf,
148 version_info: &GitVersionInfo,
149 is_ci: bool,
150 database: impl FnOnce(PathBuf) -> Result<TurboKeyValueDatabase>,
151 ) -> Result<(Self, StartupCacheState)> {
152 let startup_cache_state = check_db_invalidation_and_cleanup(&base_path)
153 .context("Failed to check database invalidation and cleanup")?;
154 let versioned_path = handle_db_versioning(&base_path, version_info, is_ci)
155 .context("Failed to handle database versioning")?;
156 let database = (database)(versioned_path).context("Failed to open database")?;
157 let backing_storage = Self {
158 inner: Arc::new_cyclic(move |weak_inner: &Weak<TurboBackingStorageInner>| {
159 let panic_hook_guard = if should_invalidate_on_panic() {
160 let weak_inner = weak_inner.clone();
161 Some(register_panic_hook(Box::new(move |_| {
162 let Some(inner) = weak_inner.upgrade() else {
163 return;
164 };
165 let _ = inner.invalidate(invalidation_reasons::PANIC);
170 })))
171 } else {
172 None
173 };
174 TurboBackingStorageInner {
175 database,
176 base_path: Some(base_path),
177 invalidated: Mutex::new(false),
178 _panic_hook_guard: panic_hook_guard,
179 }
180 }),
181 };
182 Ok((backing_storage, startup_cache_state))
183 }
184}
185
186impl TurboBackingStorageInner {
187 fn invalidate(&self, reason_code: &str) -> Result<()> {
188 if let Some(base_path) = &self.base_path {
190 let mut invalidated_guard = self
193 .invalidated
194 .lock()
195 .unwrap_or_else(PoisonError::into_inner);
196 if *invalidated_guard {
197 return Ok(());
198 }
199 invalidate_db(base_path, reason_code)?;
203 self.database.prevent_writes();
204 *invalidated_guard = true;
206 }
207 Ok(())
208 }
209
210 fn get_infra_u32(&self, key: InfraKey) -> Result<Option<u32>> {
212 self.database
213 .get(KeySpace::Infra, key.key().as_ref())?
214 .map(as_u32)
215 .transpose()
216 }
217}
218
219impl TurboBackingStorage {
220 pub(crate) fn invalidate(&self, reason_code: &str) -> Result<()> {
228 self.inner.invalidate(reason_code)
229 }
230
231 pub(crate) fn next_free_task_id(&self) -> Result<TaskId> {
232 Ok(self
233 .inner
234 .get_infra_u32(InfraKey::NextFreeTaskId)
235 .context("Unable to read next free task id from database")?
236 .map_or(Ok(TaskId::MIN), TaskId::try_from)?)
237 }
238
239 pub(crate) fn roots(&self) -> Result<Vec<(TaskId, TtlCounter)>> {
241 fn get(database: &TurboKeyValueDatabase) -> Result<Vec<(TaskId, TtlCounter)>> {
242 let Some(roots) = database.get(KeySpace::Infra, InfraKey::GcRoots.key().as_ref())?
243 else {
244 return Ok(Vec::new());
245 };
246 let roots = turbo_bincode_decode(roots.borrow())?;
247 Ok(roots)
248 }
249 get(&self.inner.database).context("Unable to read GC roots from database")
250 }
251
252 pub(crate) fn save_snapshot<I>(
253 &self,
254 roots: Option<Vec<(TaskId, TtlCounter)>>,
255 snapshots: Vec<I>,
256 ) -> Result<SnapshotMeta>
257 where
258 I: IntoIterator<Item = SnapshotItem> + Send + Sync,
259 {
260 let _span = tracing::info_span!("save snapshot").entered();
261 let batch = self.inner.database.write_batch()?;
262
263 {
264 let span = tracing::trace_span!("update task data");
265 let mut snapshot_meta =
266 parallel::map_collect_owned::<_, _, Result<Vec<_>>>(snapshots, |shard: I| {
267 let _span = span.clone().entered();
268 let mut max_new_task_id = 0;
269 let mut data_items = 0;
270 let mut meta_items = 0;
271 let mut task_cache_items = 0;
272 for item in shard {
273 match item {
274 SnapshotItem::Put {
275 task_id,
276 meta,
277 data,
278 task_type_hash,
279 } => {
280 let key = IntKey::new(*task_id);
281 let key = key.as_ref();
282 if let Some(meta) = meta {
283 batch.put(
284 KeySpace::TaskMeta,
285 WriteBuffer::Borrowed(key),
286 WriteBuffer::SmallVec(meta),
287 )?;
288 meta_items += 1;
289 }
290 if let Some(data) = data {
291 batch.put(
292 KeySpace::TaskData,
293 WriteBuffer::Borrowed(key),
294 WriteBuffer::SmallVec(data),
295 )?;
296 data_items += 1;
297 }
298 if let Some(task_type_hash) = task_type_hash {
300 batch.put(
301 KeySpace::TaskCache,
302 WriteBuffer::Borrowed(&task_type_hash),
303 WriteBuffer::Borrowed(key),
304 )?;
305 task_cache_items += 1;
306 max_new_task_id = max_new_task_id.max(*task_id);
307 }
308 }
309 SnapshotItem::Delete {
310 task_id,
311 task_type_hash,
312 } => {
313 let key = IntKey::new(*task_id);
314 let key = key.as_ref();
315 batch.delete(KeySpace::TaskMeta, WriteBuffer::Borrowed(key))?;
316 batch.delete(KeySpace::TaskData, WriteBuffer::Borrowed(key))?;
317 batch.delete_value(
319 KeySpace::TaskCache,
320 WriteBuffer::Borrowed(&task_type_hash[..]),
321 WriteBuffer::Borrowed(key),
322 )?;
323 }
324 }
325 }
326 Ok(SnapshotMeta {
327 data_items,
328 meta_items,
329 task_cache_items,
330 bytes_written: 0,
333 bytes_deleted: 0,
334 max_next_task_id: max_new_task_id,
335 })
336 })?
337 .into_iter()
338 .reduce(|t1, t2| t1.merge(t2))
339 .unwrap_or_default();
340
341 let span = tracing::trace_span!("flush task data");
342 parallel::try_for_each(
343 &[KeySpace::TaskMeta, KeySpace::TaskData, KeySpace::TaskCache],
344 |&key_space| {
345 let _span = span.clone().entered();
346 unsafe { batch.flush(key_space) }
349 },
350 )?;
351
352 let mut next_task_id = get_next_free_task_id(&batch)?;
353 next_task_id = next_task_id.max(snapshot_meta.max_next_task_id + 1);
354
355 save_infra(&batch, next_task_id, roots)?;
356 {
357 let _span = tracing::trace_span!("commit").entered();
358 let stats = batch.commit().context("Unable to commit snapshot")?;
361 snapshot_meta.bytes_written = stats.bytes_written;
362 snapshot_meta.bytes_deleted = stats.bytes_deleted;
363 }
364 Ok(snapshot_meta)
365 }
366 }
367
368 pub(crate) fn lookup_task_candidates(
369 &self,
370 native_fn: &'static NativeFunction,
371 this: Option<RawVc>,
372 arg: &dyn DynTaskInputs,
373 ) -> Result<SmallVec<[TaskId; 1]>> {
374 let inner = &*self.inner;
375 if inner.database.is_empty() {
376 return Ok(SmallVec::new());
379 }
380 let hash = compute_task_type_hash_from_components(native_fn, this, arg);
381 let buffers = inner
382 .database
383 .get_multiple(KeySpace::TaskCache, &hash)
384 .with_context(|| {
385 format!("Looking up task id for {native_fn:?}(this={this:?}) from database failed")
386 })?;
387
388 let mut task_ids = SmallVec::with_capacity(buffers.len());
389 for bytes in buffers {
390 let bytes = Borrow::<[u8]>::borrow(&bytes).try_into()?;
391 let id = TaskId::try_from(u32::from_le_bytes(bytes)).unwrap();
392 task_ids.push(id);
393 }
394 Ok(task_ids)
395 }
396
397 pub(crate) fn lookup_data(
403 &self,
404 task_id: TaskId,
405 category: SpecificTaskDataCategory,
406 ) -> Result<Option<TaskStorage>> {
407 let inner = &*self.inner;
408 let Some(bytes) = inner
409 .database
410 .get(category.key_space(), IntKey::new(*task_id).as_ref())
411 .with_context(|| {
412 format!("Looking up task storage for {task_id} from database failed")
413 })?
414 else {
415 return Ok(None);
416 };
417 let mut storage = TaskStorage::default();
418 let mut decoder = new_turbo_bincode_decoder(bytes.borrow());
419 storage
420 .decode(category, &mut decoder)
421 .with_context(|| format!("Failed to decode {category:?}"))?;
422 Ok(Some(storage))
423 }
424
425 pub(crate) fn batch_lookup_data(
426 &self,
427 task_ids: &[TaskId],
428 category: SpecificTaskDataCategory,
429 ) -> Result<Vec<Option<TaskStorage>>> {
430 let inner = &*self.inner;
431 let int_keys: Vec<_> = task_ids.iter().map(|&id| IntKey::new(*id)).collect();
432 let keys = int_keys.iter().map(|k| k.as_ref()).collect::<Vec<_>>();
433 let bytes = inner
434 .database
435 .batch_get(category.key_space(), &keys)
436 .with_context(|| {
437 format!(
438 "Looking up typed data for {} tasks from database failed",
439 task_ids.len()
440 )
441 })?;
442 bytes
443 .into_iter()
444 .map(|opt_bytes| {
445 let Some(bytes) = opt_bytes else {
446 return Ok(None);
447 };
448 let mut storage = TaskStorage::new();
449 let mut decoder = new_turbo_bincode_decoder(bytes.borrow());
450 storage
451 .decode(category, &mut decoder)
452 .map_err(|e| anyhow::anyhow!("Failed to decode {category:?}: {e:?}"))?;
453 Ok(Some(storage))
454 })
455 .collect::<Result<Vec<_>>>()
456 }
457
458 pub(crate) fn compact(&self) -> Result<Option<CommitStats>> {
459 self.inner.database.compact()
460 }
461
462 pub(crate) fn shutdown(&self) -> Result<()> {
463 self.inner.database.shutdown()
464 }
465
466 pub(crate) fn has_unrecoverable_write_error(&self) -> bool {
467 self.inner.database.has_unrecoverable_write_error()
468 }
469}
470
471fn get_next_free_task_id(batch: &TurboWriteBatch<'_>) -> Result<u32, anyhow::Error> {
472 Ok(
473 match batch.get(KeySpace::Infra, InfraKey::NextFreeTaskId.key().as_ref())? {
474 Some(bytes) => u32::from_le_bytes(Borrow::<[u8]>::borrow(&bytes).try_into()?),
475 None => 1,
476 },
477 )
478}
479
480fn save_infra(
481 batch: &TurboWriteBatch<'_>,
482 next_task_id: u32,
483 roots: Option<Vec<(TaskId, TtlCounter)>>,
484) -> Result<(), anyhow::Error> {
485 batch
486 .put(
487 KeySpace::Infra,
488 WriteBuffer::Borrowed(InfraKey::NextFreeTaskId.key().as_ref()),
489 WriteBuffer::Borrowed(&next_task_id.to_le_bytes()),
490 )
491 .context("Unable to write next free task id")?;
492 if let Some(roots) = roots {
493 let _span = tracing::trace_span!("update roots", roots = roots.len()).entered();
494 let roots = turbo_bincode_encode(&roots).context("Unable to serialize GC roots")?;
495 batch
496 .put(
497 KeySpace::Infra,
498 WriteBuffer::Borrowed(InfraKey::GcRoots.key().as_ref()),
499 WriteBuffer::SmallVec(roots),
500 )
501 .context("Unable to write GC roots")?;
502 }
503 unsafe { batch.flush(KeySpace::Infra)? };
505 Ok(())
506}
507
508#[cfg(test)]
509mod tests {
510 use std::borrow::Borrow;
511
512 use turbo_tasks::TaskId;
513
514 use super::*;
515 use crate::{
516 BackingStorageOptions,
517 database::{turbo::TurboKeyValueDatabase, write_batch::WriteBuffer},
518 utils::test_temp_dir::test_temp_dir,
519 };
520
521 const TEST_STORAGE_OPTIONS: BackingStorageOptions = BackingStorageOptions {
524 is_ci: false,
525 is_short_session: true,
526 skip_compaction: false,
527 };
528
529 fn write_task_cache_entry(
531 db: &TurboKeyValueDatabase,
532 hash: u64,
533 task_id: TaskId,
534 ) -> Result<()> {
535 let batch = db.write_batch()?;
536 batch.put(
537 KeySpace::TaskCache,
538 WriteBuffer::Borrowed(&hash.to_le_bytes()),
539 WriteBuffer::Borrowed(&(*task_id).to_le_bytes()),
540 )?;
541 batch.commit()?;
542 Ok(())
543 }
544
545 fn task_cache_ids(db: &TurboKeyValueDatabase, hash: u64) -> Result<Vec<TaskId>> {
547 let mut ids: Vec<TaskId> = db
548 .get_multiple(KeySpace::TaskCache, &hash.to_le_bytes())?
549 .iter()
550 .map(|bytes| {
551 let bytes: [u8; 4] = Borrow::<[u8]>::borrow(bytes).try_into().unwrap();
552 TaskId::try_from(u32::from_le_bytes(bytes)).unwrap()
553 })
554 .collect();
555 ids.sort_by_key(|id| **id);
556 Ok(ids)
557 }
558
559 #[tokio::test(flavor = "multi_thread")]
565 async fn test_hash_collision_returns_multiple_candidates() -> Result<()> {
566 let tempdir = test_temp_dir()?;
567 let path = tempdir.path();
568
569 let db = TurboKeyValueDatabase::new(path.to_path_buf(), TEST_STORAGE_OPTIONS)?;
570
571 let collision_hash: u64 = 0xDEADBEEF;
573 let task_id_1 = TaskId::try_from(100u32).unwrap();
574 let task_id_2 = TaskId::try_from(200u32).unwrap();
575 let task_id_3 = TaskId::try_from(300u32).unwrap();
576
577 write_task_cache_entry(&db, collision_hash, task_id_1)?;
580 write_task_cache_entry(&db, collision_hash, task_id_2)?;
581 write_task_cache_entry(&db, collision_hash, task_id_3)?;
582
583 assert_eq!(
585 task_cache_ids(&db, collision_hash)?,
586 vec![task_id_1, task_id_2, task_id_3],
587 "Should return all 3 task IDs for the colliding hash"
588 );
589
590 db.shutdown()?;
591 Ok(())
592 }
593
594 #[cfg(not(miri))]
598 #[tokio::test(flavor = "multi_thread")]
599 async fn test_batch_write_with_flush_and_reopen() -> Result<()> {
600 let tempdir = test_temp_dir()?;
601 let path = tempdir.path();
602
603 let n = 100_000;
604 let hashes: Vec<u64> = (0..n).map(|i| 0x1000 + i as u64).collect();
605 let task_ids: Vec<TaskId> = (1..=n as u32)
606 .map(|i| TaskId::try_from(i).unwrap())
607 .collect();
608
609 {
611 let db = TurboKeyValueDatabase::new(path.to_path_buf(), TEST_STORAGE_OPTIONS)?;
612 let batch = db.write_batch()?;
613
614 for (hash, task_id) in hashes.iter().zip(task_ids.iter()) {
615 batch.put(
616 KeySpace::TaskCache,
617 WriteBuffer::Borrowed(&hash.to_le_bytes()),
618 WriteBuffer::Borrowed(&(**task_id).to_le_bytes()),
619 )?;
620 }
621 unsafe { batch.flush(KeySpace::TaskCache) }?;
623 batch.commit()?;
624
625 db.shutdown()?;
626 }
627
628 {
630 let db = TurboKeyValueDatabase::new(path.to_path_buf(), TEST_STORAGE_OPTIONS)?;
631 let mut found = 0;
632 let mut missing = 0;
633 for (hash, expected_id) in hashes.iter().zip(task_ids.iter()) {
634 let results = db.get_multiple(KeySpace::TaskCache, &hash.to_le_bytes())?;
635 if results.is_empty() {
636 missing += 1;
637 } else {
638 found += 1;
639 let bytes: [u8; 4] = Borrow::<[u8]>::borrow(&results[0]).try_into().unwrap();
640 let id = TaskId::try_from(u32::from_le_bytes(bytes)).unwrap();
641 assert_eq!(id, *expected_id, "Task ID mismatch for hash {hash:#x}");
642 }
643 }
644 assert_eq!(missing, 0, "Found {found}/{n} entries, missing {missing}");
645 db.shutdown()?;
646 }
647
648 Ok(())
649 }
650
651 #[tokio::test(flavor = "multi_thread")]
659 async fn test_save_snapshot_delete_tombstones_task() -> Result<()> {
660 let tempdir = test_temp_dir()?;
661 let path = tempdir.path();
662
663 let collision_hash: u64 = 0xC0FFEE;
664 let deleted_id = TaskId::try_from(111u32).unwrap();
665 let survivor_id = TaskId::try_from(222u32).unwrap();
666 let deleted_key = (*deleted_id).to_le_bytes();
667
668 let db = TurboKeyValueDatabase::new(
669 path.to_path_buf(),
670 BackingStorageOptions {
671 is_ci: false,
672 is_short_session: true,
673 skip_compaction: false,
674 },
675 )?;
676
677 write_task_cache_entry(&db, collision_hash, deleted_id)?;
680 write_task_cache_entry(&db, collision_hash, survivor_id)?;
681 let batch = db.write_batch()?;
682 batch.put(
683 KeySpace::TaskMeta,
684 WriteBuffer::Borrowed(&deleted_key),
685 WriteBuffer::Borrowed(b"meta-bytes"),
686 )?;
687 batch.put(
688 KeySpace::TaskData,
689 WriteBuffer::Borrowed(&deleted_key),
690 WriteBuffer::Borrowed(b"data-bytes"),
691 )?;
692 batch.commit()?;
693
694 assert!(db.get(KeySpace::TaskMeta, &deleted_key)?.is_some());
696 assert!(db.get(KeySpace::TaskData, &deleted_key)?.is_some());
697 assert_eq!(
698 task_cache_ids(&db, collision_hash)?,
699 vec![deleted_id, survivor_id],
700 );
701
702 let storage = TurboBackingStorage::new_in_memory(db);
703
704 storage.save_snapshot(
706 None,
707 vec![vec![SnapshotItem::Delete {
708 task_id: deleted_id,
709 task_type_hash: collision_hash.to_le_bytes(),
710 }]],
711 )?;
712
713 let db = &storage.inner.database;
714 assert!(
715 db.get(KeySpace::TaskMeta, &deleted_key)?.is_none(),
716 "TaskMeta should be tombstoned"
717 );
718 assert!(
719 db.get(KeySpace::TaskData, &deleted_key)?.is_none(),
720 "TaskData should be tombstoned"
721 );
722 assert_eq!(
723 task_cache_ids(db, collision_hash)?,
724 vec![survivor_id],
725 "save_snapshot should delete only the named id from the bucket"
726 );
727
728 db.shutdown()?;
729 Ok(())
730 }
731}