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 } => {
281 let key = IntKey::new(*task_id);
282 let key = key.as_ref();
283 if let Some(meta) = meta {
284 batch.put(
285 KeySpace::TaskMeta,
286 WriteBuffer::Borrowed(key),
287 WriteBuffer::SmallVec(meta),
288 )?;
289 meta_items += 1;
290 }
291 if let Some(data) = data {
292 batch.put(
293 KeySpace::TaskData,
294 WriteBuffer::Borrowed(key),
295 WriteBuffer::SmallVec(data),
296 )?;
297 data_items += 1;
298 }
299 if let Some(task_type_hash) = task_type_hash {
301 batch.put(
302 KeySpace::TaskCache,
303 WriteBuffer::Borrowed(&task_type_hash),
304 WriteBuffer::Borrowed(key),
305 )?;
306 task_cache_items += 1;
307 max_new_task_id = max_new_task_id.max(*task_id);
308 }
309 }
310 SnapshotItem::Delete {
311 task_id,
312 task_type_hash,
313 } => {
314 let key = IntKey::new(*task_id);
315 let key = key.as_ref();
316 batch.delete(KeySpace::TaskMeta, WriteBuffer::Borrowed(key))?;
317 batch.delete(KeySpace::TaskData, WriteBuffer::Borrowed(key))?;
318 batch.delete_value(
320 KeySpace::TaskCache,
321 WriteBuffer::Borrowed(&task_type_hash[..]),
322 WriteBuffer::Borrowed(key),
323 )?;
324 }
325 }
326 }
327 Ok(SnapshotMeta {
328 data_items,
329 meta_items,
330 task_cache_items,
331 bytes_written: 0,
334 bytes_deleted: 0,
335 max_next_task_id: max_new_task_id,
336 })
337 })?
338 .into_iter()
339 .reduce(|t1, t2| t1.merge(t2))
340 .unwrap_or_default();
341
342 let span = tracing::trace_span!("flush task data");
343 parallel::try_for_each(
344 &[KeySpace::TaskMeta, KeySpace::TaskData, KeySpace::TaskCache],
345 |&key_space| {
346 let _span = span.clone().entered();
347 unsafe { batch.flush(key_space) }
350 },
351 )?;
352
353 let mut next_task_id = get_next_free_task_id(&batch)?;
354 next_task_id = next_task_id.max(snapshot_meta.max_next_task_id + 1);
355
356 save_infra(&batch, next_task_id, roots)?;
357 {
358 let _span = tracing::trace_span!("commit").entered();
359 let stats = batch.commit().context("Unable to commit snapshot")?;
362 snapshot_meta.bytes_written = stats.bytes_written;
363 snapshot_meta.bytes_deleted = stats.bytes_deleted;
364 }
365 Ok(snapshot_meta)
366 }
367 }
368
369 pub(crate) fn lookup_task_candidates(
370 &self,
371 native_fn: &'static NativeFunction,
372 this: Option<RawVc>,
373 arg: &dyn DynTaskInputs,
374 ) -> Result<SmallVec<[TaskId; 1]>> {
375 let inner = &*self.inner;
376 if inner.database.is_empty() {
377 return Ok(SmallVec::new());
380 }
381 let hash = compute_task_type_hash_from_components(native_fn, this, arg);
382 let buffers = inner
383 .database
384 .get_multiple(KeySpace::TaskCache, &hash)
385 .with_context(|| {
386 format!("Looking up task id for {native_fn:?}(this={this:?}) from database failed")
387 })?;
388
389 let mut task_ids = SmallVec::with_capacity(buffers.len());
390 for bytes in buffers {
391 let bytes = Borrow::<[u8]>::borrow(&bytes).try_into()?;
392 let id = TaskId::try_from(u32::from_le_bytes(bytes)).unwrap();
393 task_ids.push(id);
394 }
395 Ok(task_ids)
396 }
397
398 pub(crate) fn lookup_data(
404 &self,
405 task_id: TaskId,
406 category: SpecificTaskDataCategory,
407 ) -> Result<Option<TaskStorage>> {
408 let inner = &*self.inner;
409 let Some(bytes) = inner
410 .database
411 .get(category.key_space(), IntKey::new(*task_id).as_ref())
412 .with_context(|| {
413 format!("Looking up task storage for {task_id} from database failed")
414 })?
415 else {
416 return Ok(None);
417 };
418 let mut storage = TaskStorage::default();
419 let mut decoder = new_turbo_bincode_decoder(bytes.borrow());
420 storage
421 .decode(category, &mut decoder)
422 .with_context(|| format!("Failed to decode {category:?}"))?;
423 Ok(Some(storage))
424 }
425
426 pub(crate) fn batch_lookup_data(
427 &self,
428 task_ids: &[TaskId],
429 category: SpecificTaskDataCategory,
430 ) -> Result<Vec<Option<TaskStorage>>> {
431 let inner = &*self.inner;
432 let int_keys: Vec<_> = task_ids.iter().map(|&id| IntKey::new(*id)).collect();
433 let keys = int_keys.iter().map(|k| k.as_ref()).collect::<Vec<_>>();
434 let bytes = inner
435 .database
436 .batch_get(category.key_space(), &keys)
437 .with_context(|| {
438 format!(
439 "Looking up typed data for {} tasks from database failed",
440 task_ids.len()
441 )
442 })?;
443 bytes
444 .into_iter()
445 .map(|opt_bytes| {
446 let Some(bytes) = opt_bytes else {
447 return Ok(None);
448 };
449 let mut storage = TaskStorage::new();
450 let mut decoder = new_turbo_bincode_decoder(bytes.borrow());
451 storage
452 .decode(category, &mut decoder)
453 .map_err(|e| anyhow::anyhow!("Failed to decode {category:?}: {e:?}"))?;
454 Ok(Some(storage))
455 })
456 .collect::<Result<Vec<_>>>()
457 }
458
459 pub(crate) fn compact(&self) -> Result<Option<CommitStats>> {
460 self.inner.database.compact()
461 }
462
463 pub(crate) fn shutdown(&self) {
464 self.inner.database.shutdown()
465 }
466
467 pub(crate) fn has_unrecoverable_write_error(&self) -> bool {
468 self.inner.database.has_unrecoverable_write_error()
469 }
470}
471
472fn get_next_free_task_id(batch: &TurboWriteBatch<'_>) -> Result<u32, anyhow::Error> {
473 Ok(
474 match batch.get(KeySpace::Infra, InfraKey::NextFreeTaskId.key().as_ref())? {
475 Some(bytes) => u32::from_le_bytes(Borrow::<[u8]>::borrow(&bytes).try_into()?),
476 None => 1,
477 },
478 )
479}
480
481fn save_infra(
482 batch: &TurboWriteBatch<'_>,
483 next_task_id: u32,
484 roots: Option<Vec<(TaskId, TtlCounter)>>,
485) -> Result<(), anyhow::Error> {
486 batch
487 .put(
488 KeySpace::Infra,
489 WriteBuffer::Borrowed(InfraKey::NextFreeTaskId.key().as_ref()),
490 WriteBuffer::Borrowed(&next_task_id.to_le_bytes()),
491 )
492 .context("Unable to write next free task id")?;
493 if let Some(roots) = roots {
494 let _span = tracing::trace_span!("update roots", roots = roots.len()).entered();
495 let roots = turbo_bincode_encode(&roots).context("Unable to serialize GC roots")?;
496 batch
497 .put(
498 KeySpace::Infra,
499 WriteBuffer::Borrowed(InfraKey::GcRoots.key().as_ref()),
500 WriteBuffer::SmallVec(roots),
501 )
502 .context("Unable to write GC roots")?;
503 }
504 unsafe { batch.flush(KeySpace::Infra)? };
506 Ok(())
507}
508
509#[cfg(test)]
510mod tests {
511 use std::borrow::Borrow;
512
513 use turbo_tasks::TaskId;
514
515 use super::*;
516 use crate::{
517 BackingStorageOptions,
518 database::{turbo::TurboKeyValueDatabase, write_batch::WriteBuffer},
519 utils::test_temp_dir::test_temp_dir,
520 };
521
522 const TEST_STORAGE_OPTIONS: BackingStorageOptions = BackingStorageOptions {
525 is_ci: false,
526 is_short_session: true,
527 skip_compaction: false,
528 };
529
530 fn write_task_cache_entry(
532 db: &TurboKeyValueDatabase,
533 hash: u64,
534 task_id: TaskId,
535 ) -> Result<()> {
536 let batch = db.write_batch()?;
537 batch.put(
538 KeySpace::TaskCache,
539 WriteBuffer::Borrowed(&hash.to_le_bytes()),
540 WriteBuffer::Borrowed(&(*task_id).to_le_bytes()),
541 )?;
542 batch.commit()?;
543 Ok(())
544 }
545
546 fn task_cache_ids(db: &TurboKeyValueDatabase, hash: u64) -> Result<Vec<TaskId>> {
548 let mut ids: Vec<TaskId> = db
549 .get_multiple(KeySpace::TaskCache, &hash.to_le_bytes())?
550 .iter()
551 .map(|bytes| {
552 let bytes: [u8; 4] = Borrow::<[u8]>::borrow(bytes).try_into().unwrap();
553 TaskId::try_from(u32::from_le_bytes(bytes)).unwrap()
554 })
555 .collect();
556 ids.sort_by_key(|id| **id);
557 Ok(ids)
558 }
559
560 #[tokio::test(flavor = "multi_thread")]
566 async fn test_hash_collision_returns_multiple_candidates() -> Result<()> {
567 let tempdir = test_temp_dir()?;
568 let path = tempdir.path();
569
570 let db = TurboKeyValueDatabase::new(path.to_path_buf(), TEST_STORAGE_OPTIONS)?;
571
572 let collision_hash: u64 = 0xDEADBEEF;
574 let task_id_1 = TaskId::try_from(100u32).unwrap();
575 let task_id_2 = TaskId::try_from(200u32).unwrap();
576 let task_id_3 = TaskId::try_from(300u32).unwrap();
577
578 write_task_cache_entry(&db, collision_hash, task_id_1)?;
581 write_task_cache_entry(&db, collision_hash, task_id_2)?;
582 write_task_cache_entry(&db, collision_hash, task_id_3)?;
583
584 assert_eq!(
586 task_cache_ids(&db, collision_hash)?,
587 vec![task_id_1, task_id_2, task_id_3],
588 "Should return all 3 task IDs for the colliding hash"
589 );
590
591 db.shutdown();
592 Ok(())
593 }
594
595 #[cfg(not(miri))]
599 #[tokio::test(flavor = "multi_thread")]
600 async fn test_batch_write_with_flush_and_reopen() -> Result<()> {
601 let tempdir = test_temp_dir()?;
602 let path = tempdir.path();
603
604 let n = 100_000;
605 let hashes: Vec<u64> = (0..n).map(|i| 0x1000 + i as u64).collect();
606 let task_ids: Vec<TaskId> = (1..=n as u32)
607 .map(|i| TaskId::try_from(i).unwrap())
608 .collect();
609
610 {
612 let db = TurboKeyValueDatabase::new(path.to_path_buf(), TEST_STORAGE_OPTIONS)?;
613 let batch = db.write_batch()?;
614
615 for (hash, task_id) in hashes.iter().zip(task_ids.iter()) {
616 batch.put(
617 KeySpace::TaskCache,
618 WriteBuffer::Borrowed(&hash.to_le_bytes()),
619 WriteBuffer::Borrowed(&(**task_id).to_le_bytes()),
620 )?;
621 }
622 unsafe { batch.flush(KeySpace::TaskCache) }?;
624 batch.commit()?;
625
626 db.shutdown();
627 }
628
629 {
631 let db = TurboKeyValueDatabase::new(path.to_path_buf(), TEST_STORAGE_OPTIONS)?;
632 let mut found = 0;
633 let mut missing = 0;
634 for (hash, expected_id) in hashes.iter().zip(task_ids.iter()) {
635 let results = db.get_multiple(KeySpace::TaskCache, &hash.to_le_bytes())?;
636 if results.is_empty() {
637 missing += 1;
638 } else {
639 found += 1;
640 let bytes: [u8; 4] = Borrow::<[u8]>::borrow(&results[0]).try_into().unwrap();
641 let id = TaskId::try_from(u32::from_le_bytes(bytes)).unwrap();
642 assert_eq!(id, *expected_id, "Task ID mismatch for hash {hash:#x}");
643 }
644 }
645 assert_eq!(missing, 0, "Found {found}/{n} entries, missing {missing}");
646 db.shutdown();
647 }
648
649 Ok(())
650 }
651
652 #[tokio::test(flavor = "multi_thread")]
660 async fn test_save_snapshot_delete_tombstones_task() -> Result<()> {
661 let tempdir = test_temp_dir()?;
662 let path = tempdir.path();
663
664 let collision_hash: u64 = 0xC0FFEE;
665 let deleted_id = TaskId::try_from(111u32).unwrap();
666 let survivor_id = TaskId::try_from(222u32).unwrap();
667 let deleted_key = (*deleted_id).to_le_bytes();
668
669 let db = TurboKeyValueDatabase::new(
670 path.to_path_buf(),
671 BackingStorageOptions {
672 is_ci: false,
673 is_short_session: true,
674 skip_compaction: false,
675 },
676 )?;
677
678 write_task_cache_entry(&db, collision_hash, deleted_id)?;
681 write_task_cache_entry(&db, collision_hash, survivor_id)?;
682 let batch = db.write_batch()?;
683 batch.put(
684 KeySpace::TaskMeta,
685 WriteBuffer::Borrowed(&deleted_key),
686 WriteBuffer::Borrowed(b"meta-bytes"),
687 )?;
688 batch.put(
689 KeySpace::TaskData,
690 WriteBuffer::Borrowed(&deleted_key),
691 WriteBuffer::Borrowed(b"data-bytes"),
692 )?;
693 batch.commit()?;
694
695 assert!(db.get(KeySpace::TaskMeta, &deleted_key)?.is_some());
697 assert!(db.get(KeySpace::TaskData, &deleted_key)?.is_some());
698 assert_eq!(
699 task_cache_ids(&db, collision_hash)?,
700 vec![deleted_id, survivor_id],
701 );
702
703 let storage = TurboBackingStorage::new_in_memory(db);
704
705 storage.save_snapshot(
707 None,
708 vec![vec![SnapshotItem::Delete {
709 task_id: deleted_id,
710 task_type_hash: collision_hash.to_le_bytes(),
711 }]],
712 )?;
713
714 let db = &storage.inner.database;
715 assert!(
716 db.get(KeySpace::TaskMeta, &deleted_key)?.is_none(),
717 "TaskMeta should be tombstoned"
718 );
719 assert!(
720 db.get(KeySpace::TaskData, &deleted_key)?.is_none(),
721 "TaskData should be tombstoned"
722 );
723 assert_eq!(
724 task_cache_ids(db, collision_hash)?,
725 vec![survivor_id],
726 "save_snapshot should delete only the named id from the bucket"
727 );
728
729 db.shutdown();
730 Ok(())
731 }
732}