Skip to main content

turbo_tasks_backend/
kv_backing_storage.rs

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/// The fixed keys in the [`KeySpace::Infra`] keyspace.
34#[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
80// We want to invalidate the cache on panic for most users, but this is a band-aid to underlying
81// problems in turbo-tasks.
82//
83// If we invalidate the cache upon panic and it "fixes" the issue upon restart, users typically
84// won't report bugs to us, and we'll never find root-causes for these problems.
85//
86// These overrides let us avoid the cache invalidation / error suppression within Vercel so that we
87// feel these pain points and fix the root causes of bugs.
88fn 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    /// Used when calling [`TurboBackingStorage::invalidate`]. Can be `None` in the
102    /// memory-only/no-op storage case.
103    base_path: Option<PathBuf>,
104    /// Used to skip calling [`invalidate_db`] when the database has already been invalidated.
105    invalidated: Mutex<bool>,
106    /// We configure a panic hook to invalidate the cache. This guard cleans up our panic hook upon
107    /// drop.
108    _panic_hook_guard: Option<PanicHookGuard>,
109}
110
111/// The higher-level backing storage passed to [`TurboTasksBackend::new`], used by
112/// [`crate::turbo_backing_storage`] and [`crate::noop_backing_storage`].
113///
114/// Wraps a low-level [`TurboKeyValueDatabase`] and adapts it into the persistence operations the
115/// backend needs (snapshots, task-candidate lookups, etc.).
116///
117/// [`TurboTasksBackend::new`]: crate::TurboTasksBackend::new
118pub struct TurboBackingStorage {
119    // wrapped so that `register_panic_hook` can hold a weak reference to `inner`.
120    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    /// Handles boilerplate logic for an on-disk persisted database with versioning.
136    ///
137    /// - Creates a directory per version, with a maximum number of old versions and performs
138    ///   automatic cleanup of old versions.
139    /// - Checks for a database invalidation marker file, and cleans up the database as needed.
140    /// - [Registers a dynamic panic hook][turbo_tasks::panic_hooks] to invalidate the database upon
141    ///   a panic. This invalidates the database using [`invalidation_reasons::PANIC`].
142    ///
143    /// Along with returning a [`TurboBackingStorage`], this returns a
144    /// [`StartupCacheState`], which can be used by the application for logging information to the
145    /// user or telemetry about the cache.
146    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                        // If a panic happened that must mean something deep inside of turbopack
166                        // or turbo-tasks failed, and it may be hard to recover. We don't want
167                        // the cache to stick around, as that may persist bugs. Make a
168                        // best-effort attempt to invalidate the database (ignoring failures).
169                        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        // `base_path` is `None` for in-memory backing storage (see `noop_backing_storage`).
189        if let Some(base_path) = &self.base_path {
190            // Invalidation could happen frequently if there's a bunch of panics. We only need to
191            // invalidate once, so grab a lock.
192            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 first, as it's a very fast atomic operation. `prevent_writes` is allowed
200            // to be slower (e.g. wait for a lock) and is allowed to corrupt the database with
201            // partial writes.
202            invalidate_db(base_path, reason_code)?;
203            self.database.prevent_writes();
204            // Avoid redundant invalidations from future panics
205            *invalidated_guard = true;
206        }
207        Ok(())
208    }
209
210    /// Used to read the next free task ID from the database.
211    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    /// Called when the database should be invalidated upon re-initialization.
221    ///
222    /// This typically means that we'll restart the process or `turbo-tasks` soon with a fresh
223    /// database. If this happens, there's no point in writing anything else to disk, or flushing
224    /// during [`TurboTasksBackend::stop`].
225    ///
226    /// [`TurboTasksBackend::stop`]: turbo_tasks::backend::Backend::stop
227    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    /// Reads the persisted GC roots set (see [`InfraKey::GcRoots`]). Empty on a fresh database.
240    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                                // Register the task type only for new tasks.
299                                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                                // TaskCache is MultiValue, delete just this id from the bucket.
318                                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                        // The on-disk byte totals aren't known until the batch is committed
331                        // below; they're filled in from `CommitStats` after `batch.commit()`.
332                        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                    // Safety: `map_collect_owned` has returned, so no concurrent `put` or `delete`
347                    // on these key spaces are in-flight.
348                    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                // Byte totals are the physical on-disk bytes (post-compression, including .sst /
359                // .blob / .meta files) produced and removed by the commit.
360                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            // Checking if the database is empty is a performance optimization
377            // to avoid computing the hash.
378            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    /// Reads the stored `category` for `task_id`.
398    ///
399    /// `None` means the database had no key for this category. That is distinct from `Some`
400    /// of an empty [`TaskStorage`] (a key that decoded to nothing): `MustExist` treats a missing
401    /// requested category as a missing task, regardless of the other category.
402    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    // Safety: save_infra is called after all concurrent writes to Infra are done.
504    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    /// Options used by these tests. `is_short_session` disables background compaction, which
522    /// requires a turbo-tasks context that these tests don't set up.
523    const TEST_STORAGE_OPTIONS: BackingStorageOptions = BackingStorageOptions {
524        is_ci: false,
525        is_short_session: true,
526        skip_compaction: false,
527    };
528
529    /// Helper to write to the database using the concurrent batch API.
530    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    /// Reads the TaskIds stored under `hash` in `TaskCache`, sorted for stable comparison.
546    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    /// Tests that `get_multiple` correctly returns multiple TaskIds when the same hash key
560    /// is used (simulating a hash collision scenario).
561    ///
562    /// This is a lower-level test that verifies the database layer correctly handles
563    /// the case where multiple task IDs are stored under the same hash key.
564    #[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        // Simulate a hash collision by writing multiple TaskIds with the same hash key
572        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 three task IDs under the same hash key (simulating collision)
578        // Each write creates a new SST file, so all three will be returned by get_multiple
579        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        // Now query using get_multiple - should return all three TaskIds
584        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    /// Tests that multiple distinct keys written in a single batch with flush can be read back.
595    /// This mirrors the actual save_snapshot pattern: write many TaskCache entries, flush, commit.
596    // This test is too slow to run under Miri.
597    #[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        // Write all entries in a single batch with flush (like save_snapshot does)
610        {
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            // Flush TaskCache (like the new code does)
622            unsafe { batch.flush(KeySpace::TaskCache) }?;
623            batch.commit()?;
624
625            db.shutdown()?;
626        }
627
628        // Reopen and verify all entries are readable
629        {
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    /// `save_snapshot` delete path: a `Delete` item must erase the task's `TaskMeta` and
652    /// `TaskData` (`SingleValue`) entries and remove *only* that id from its `TaskCache`
653    /// (`MultiValue`) bucket.
654    ///
655    /// The colliding survivor is never read or rewritten — the key-value tombstone names the
656    /// single id it deletes, so anything else in the bucket is untouched whether or not this
657    /// commit knows about it.
658    #[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        // Both ids collide in one TaskCache bucket, purely on disk; the deleted task also has
678        // meta and data entries.
679        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        // Sanity: everything is present before the delete.
695        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        // Snapshot with no task data, just the one deletion.
705        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}