From d8c5695b7b629244cf7461bef3e623e7bbbdd4d1 Mon Sep 17 00:00:00 2001 From: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com> Date: Thu, 6 Aug 2026 04:00:56 +0000 Subject: [PATCH] fix: reject unsafe version cleanup windows --- docs/src/js/classes/Table.md | 3 + docs/src/js/interfaces/OptimizeOptions.md | 9 ++- nodejs/__test__/table.test.ts | 10 ++- nodejs/lancedb/table.ts | 12 +++- nodejs/src/table.rs | 2 + python/python/lancedb/table.py | 50 ++++++++++--- python/python/tests/test_table.py | 12 +++- python/src/table.rs | 2 + rust/lancedb/src/table/optimize.rs | 85 ++++++++++++++++++++++- 9 files changed, 163 insertions(+), 22 deletions(-) diff --git a/docs/src/js/classes/Table.md b/docs/src/js/classes/Table.md index 11fca32d0..c1a12dc8f 100644 --- a/docs/src/js/classes/Table.md +++ b/docs/src/js/classes/Table.md @@ -587,6 +587,9 @@ Modeled after ``VACUUM`` in PostgreSQL. you have added or modified 100,000 or more records or run more than 20 data modification operations. + Cleanup retention must exceed the longest expected write. Retention shorter + than 10 minutes requires `deleteUnverified: true` and exclusive write access. + #### Parameters * **options?**: `Partial`<[`OptimizeOptions`](../interfaces/OptimizeOptions.md)> diff --git a/docs/src/js/interfaces/OptimizeOptions.md b/docs/src/js/interfaces/OptimizeOptions.md index 700632342..6b5f226a3 100644 --- a/docs/src/js/interfaces/OptimizeOptions.md +++ b/docs/src/js/interfaces/OptimizeOptions.md @@ -16,7 +16,10 @@ cleanupOlderThan: Date; If set then all versions older than the given date be removed. The current version will never be removed. -The default is 7 days +The default is 7 days. The resulting retention period must be longer than +the longest expected write. Values shorter than 10 minutes require +`deleteUnverified: true` and are only safe when no other process can write +to the dataset. #### Example @@ -26,8 +29,8 @@ const olderThan = new Date(); olderThan.setDate(olderThan.getDate() - 1)); tbl.optimize({cleanupOlderThan: olderThan}); -// Delete all versions except the current version -tbl.optimize({cleanupOlderThan: new Date()}); +// With exclusive access, delete all versions except the current version +tbl.optimize({cleanupOlderThan: new Date(), deleteUnverified: true}); ``` *** diff --git a/nodejs/__test__/table.test.ts b/nodejs/__test__/table.test.ts index 4cad365af..e561bfa3b 100644 --- a/nodejs/__test__/table.test.ts +++ b/nodejs/__test__/table.test.ts @@ -2196,7 +2196,15 @@ describe("when optimizing a dataset", () => { }); it("cleanups old versions", async () => { - const stats = await table.optimize({ cleanupOlderThan: new Date() }); + await expect( + table.optimize({ cleanupOlderThan: new Date() }), + ).rejects.toThrow("at least 10 minutes"); + expect(await table.version()).toBe(2); + + const stats = await table.optimize({ + cleanupOlderThan: new Date(), + deleteUnverified: true, + }); expect(stats.prune.bytesRemoved).toBeGreaterThan(0); expect(stats.prune.oldVersionsRemoved).toBe(3); }); diff --git a/nodejs/lancedb/table.ts b/nodejs/lancedb/table.ts index 3359a2643..5920391be 100644 --- a/nodejs/lancedb/table.ts +++ b/nodejs/lancedb/table.ts @@ -129,15 +129,18 @@ export interface OptimizeOptions { /** * If set then all versions older than the given date * be removed. The current version will never be removed. - * The default is 7 days + * The default is 7 days. The resulting retention period must be longer than + * the longest expected write. Values shorter than 10 minutes require + * `deleteUnverified: true` and are only safe when no other process can write + * to the dataset. * @example * // Delete all versions older than 1 day * const olderThan = new Date(); * olderThan.setDate(olderThan.getDate() - 1)); * tbl.optimize({cleanupOlderThan: olderThan}); * - * // Delete all versions except the current version - * tbl.optimize({cleanupOlderThan: new Date()}); + * // With exclusive access, delete all versions except the current version + * tbl.optimize({cleanupOlderThan: new Date(), deleteUnverified: true}); */ cleanupOlderThan: Date; /** @@ -744,6 +747,9 @@ export abstract class Table { * optimize should be run frequently. A good rule of thumb is to run optimize if * you have added or modified 100,000 or more records or run more than 20 data * modification operations. + * + * Cleanup retention must exceed the longest expected write. Retention shorter + * than 10 minutes requires `deleteUnverified: true` and exclusive write access. */ abstract optimize(options?: Partial): Promise; /** List all indices that have been created with {@link Table.createIndex} */ diff --git a/nodejs/src/table.rs b/nodejs/src/table.rs index 2ac2fecb2..c175f09d7 100644 --- a/nodejs/src/table.rs +++ b/nodejs/src/table.rs @@ -6,6 +6,7 @@ use std::collections::HashMap; use chrono::{DateTime, Utc}; use lancedb::ipc::{ipc_file_to_batches, ipc_file_to_schema}; +use lancedb::table::optimize::validate_cleanup_options; use lancedb::table::{ AddDataMode, ColumnAlteration as LanceColumnAlteration, Duration, FieldMetadataUpdate as LanceFieldMetadataUpdate, FtsToken as LanceDbFtsToken, @@ -559,6 +560,7 @@ impl Table { } else { None }; + validate_cleanup_options(older_than, delete_unverified).default_error()?; let compaction_stats = inner .optimize(OptimizeAction::Compact { diff --git a/python/python/lancedb/table.py b/python/python/lancedb/table.py index 31a70c298..9ab625d02 100644 --- a/python/python/lancedb/table.py +++ b/python/python/lancedb/table.py @@ -117,6 +117,23 @@ _MODEL_BACKED_TOKENIZER_ERRORS = ( "Failed to initialize default tokenizer", ) +_MIN_SAFE_CLEANUP_AGE = timedelta(minutes=10) + + +def _validate_cleanup_options( + older_than: Optional[timedelta], delete_unverified: bool +) -> None: + if ( + older_than is not None + and older_than < _MIN_SAFE_CLEANUP_AGE + and not delete_unverified + ): + raise ValueError( + "cleanup age must be at least 10 minutes unless delete_unverified is " + "true; short cleanup windows can remove a manifest still needed by an " + "in-progress write" + ) + def _add_unique_note(exception: BaseException, note: str) -> None: existing_notes = getattr(exception, "__notes__", ()) or () @@ -1767,7 +1784,9 @@ class Table(ABC): ---------- older_than: timedelta, default None The minimum age of the version to delete. If None, then this defaults - to two weeks. + to two weeks. This must be longer than the longest expected write. + Values shorter than 10 minutes require `delete_unverified=True` and + are only safe when no other process can write to the dataset. delete_unverified: bool, default False Because they may be part of an in-progress transaction, files newer than 7 days old are not deleted by default. If you are sure that @@ -1835,9 +1854,11 @@ class Table(ABC): Parameters ---------- cleanup_older_than: timedelta, optional default 7 days - All files belonging to versions older than this will be removed. Set - to 0 days to remove all versions except the latest. The latest version - is never removed. + All files belonging to versions older than this will be removed. The + latest version is never removed. This must be longer than the longest + expected write. Values shorter than 10 minutes require + `delete_unverified=True` and are only safe when no other process can + write to the dataset. delete_unverified: bool, default False Files leftover from a failed transaction may appear to be part of an in-progress operation (e.g. appending new data) and these files will not @@ -3782,7 +3803,9 @@ class LanceTable(Table): ---------- older_than: timedelta, default None The minimum age of the version to delete. If None, then this defaults - to two weeks. + to two weeks. This must be longer than the longest expected write. + Values shorter than 10 minutes require `delete_unverified=True` and + are only safe when no other process can write to the dataset. delete_unverified: bool, default False Because they may be part of an in-progress transaction, files newer than 7 days old are not deleted by default. If you are sure that @@ -3795,6 +3818,7 @@ class LanceTable(Table): The stats of the cleanup operation, including how many bytes were freed. """ + _validate_cleanup_options(older_than, delete_unverified) return self.to_lance().cleanup_old_versions( older_than, delete_unverified=delete_unverified ) @@ -3840,9 +3864,11 @@ class LanceTable(Table): Parameters ---------- cleanup_older_than: timedelta, optional default 7 days - All files belonging to versions older than this will be removed. Set - to 0 days to remove all versions except the latest. The latest version - is never removed. + All files belonging to versions older than this will be removed. The + latest version is never removed. This must be longer than the longest + expected write. Values shorter than 10 minutes require + `delete_unverified=True` and are only safe when no other process can + write to the dataset. delete_unverified: bool, default False Files leftover from a failed transaction may appear to be part of an in-progress operation (e.g. appending new data) and these files will not @@ -6034,9 +6060,11 @@ class AsyncTable: Parameters ---------- cleanup_older_than: timedelta, optional default 7 days - All files belonging to versions older than this will be removed. Set - to 0 days to remove all versions except the latest. The latest version - is never removed. + All files belonging to versions older than this will be removed. The + latest version is never removed. This must be longer than the longest + expected write. Values shorter than 10 minutes require + `delete_unverified=True` and are only safe when no other process can + write to the dataset. delete_unverified: bool, default False Files leftover from a failed transaction may appear to be part of an in-progress operation (e.g. appending new data) and these files will not diff --git a/python/python/tests/test_table.py b/python/python/tests/test_table.py index 542c944c8..2ef0b61e8 100644 --- a/python/python/tests/test_table.py +++ b/python/python/tests/test_table.py @@ -2958,6 +2958,9 @@ def test_compact_cleanup(tmp_db: DBConnection): stats = table.cleanup_old_versions() assert stats.bytes_removed == 0 + with pytest.raises(ValueError, match="at least 10 minutes"): + table.cleanup_old_versions(older_than=timedelta(0)) + stats = table.cleanup_old_versions(older_than=timedelta(0), delete_unverified=True) assert stats.bytes_removed > 0 assert table.version == 4 @@ -3374,7 +3377,14 @@ async def test_optimize(mem_db_async: AsyncConnection): assert stats.prune.bytes_removed == 0 assert stats.prune.old_versions_removed == 0 - stats = await table.optimize(cleanup_older_than=timedelta(seconds=0)) + version_before_rejected_cleanup = await table.version() + with pytest.raises(ValueError, match="at least 10 minutes"): + await table.optimize(cleanup_older_than=timedelta(seconds=0)) + assert await table.version() == version_before_rejected_cleanup + + stats = await table.optimize( + cleanup_older_than=timedelta(seconds=0), delete_unverified=True + ) assert stats.prune.bytes_removed > 0 assert stats.prune.old_versions_removed == 3 diff --git a/python/src/table.rs b/python/src/table.rs index f20bcf2bf..61bcac4f2 100644 --- a/python/src/table.rs +++ b/python/src/table.rs @@ -19,6 +19,7 @@ use arrow::{ }; use lancedb::blob::{BlobFile, BlobRangeRequest}; use lancedb::index::scalar::FtsIndexBuilder; +use lancedb::table::optimize::validate_cleanup_options; use lancedb::table::{ AddDataMode, ColumnAlteration, Duration, FieldMetadataUpdate, FtsToken as LanceDbFtsToken, NewColumnTransform, OptimizeAction, OptimizeOptions, Ref, Table as LanceDbTable, @@ -1214,6 +1215,7 @@ impl Table { } else { None }; + validate_cleanup_options(older_than, delete_unverified).infer_error()?; future_into_py(self_.py(), async move { let compaction_stats = inner .optimize(OptimizeAction::Compact { diff --git a/rust/lancedb/src/table/optimize.rs b/rust/lancedb/src/table/optimize.rs index e29445b2d..0ebb18d58 100644 --- a/rust/lancedb/src/table/optimize.rs +++ b/rust/lancedb/src/table/optimize.rs @@ -18,7 +18,33 @@ pub use chrono::Duration; pub use lance::dataset::optimize::CompactionOptions; use super::NativeTable; -use crate::error::Result; +use crate::error::{Error, Result}; + +const MIN_SAFE_CLEANUP_AGE_MINUTES: i64 = 10; + +/// Validate version-cleanup options before starting an optimization operation. +/// +/// This is public for the language bindings, which run compaction and cleanup as +/// separate operations and must reject unsafe cleanup options before compaction +/// changes the table. +#[doc(hidden)] +pub fn validate_cleanup_options( + older_than: Option, + delete_unverified: Option, +) -> Result<()> { + let minimum_age = + Duration::try_minutes(MIN_SAFE_CLEANUP_AGE_MINUTES).expect("minimum cleanup age is valid"); + if older_than.is_some_and(|age| age < minimum_age) && delete_unverified != Some(true) { + return Err(Error::InvalidInput { + message: format!( + "cleanup age must be at least {MIN_SAFE_CLEANUP_AGE_MINUTES} minutes unless \ + delete_unverified is true; short cleanup windows can remove a manifest still \ + needed by an in-progress write" + ), + }); + } + Ok(()) +} /// Optimize the dataset. /// @@ -60,7 +86,9 @@ pub enum OptimizeAction { /// /// Once a version is pruned it can no longer be checked out. Prune { - /// The duration of time to keep versions of the dataset. + /// The duration of time to keep versions of the dataset. This should be longer than the + /// longest expected write. Values shorter than 10 minutes require `delete_unverified` to + /// be true and are only safe when no other process can write to the dataset. older_than: Option, /// Because they may be part of an in-progress transaction, files newer than 7 days old are not deleted by default. /// If you are sure that there are no in-progress transactions, then you can set this to True to delete all files older than `older_than`. @@ -164,6 +192,15 @@ pub(crate) async fn execute_optimize( table: &NativeTable, action: OptimizeAction, ) -> Result { + if let OptimizeAction::Prune { + older_than, + delete_unverified, + .. + } = &action + { + validate_cleanup_options(*older_than, *delete_unverified)?; + } + let mut stats = OptimizeStats { compaction: None, prune: None, @@ -222,7 +259,7 @@ mod tests { use crate::connect; use crate::index::{Index, scalar::BTreeIndexBuilder}; use crate::query::ExecutableQuery; - use crate::table::{CompactionOptions, OptimizeAction, OptimizeStats}; + use crate::table::{CompactionOptions, Duration, OptimizeAction, OptimizeStats}; use futures::TryStreamExt; #[tokio::test] @@ -380,6 +417,48 @@ mod tests { assert_eq!(all_values, expected); } + #[tokio::test] + async fn test_optimize_rejects_unsafe_cleanup_age() { + let conn = connect("memory://").execute().await.unwrap(); + let batch = RecordBatch::try_new( + Arc::new(Schema::new(vec![Field::new("i", DataType::Int32, false)])), + vec![Arc::new(Int32Array::from_iter_values(0..10))], + ) + .unwrap(); + let table = conn + .create_table("test_unsafe_prune", batch.clone()) + .execute() + .await + .unwrap(); + table.add(batch).execute().await.unwrap(); + + let versions_before = table + .list_versions() + .await + .unwrap() + .into_iter() + .map(|version| version.version) + .collect::>(); + let err = table + .optimize(OptimizeAction::Prune { + older_than: Some(Duration::zero()), + delete_unverified: Some(false), + error_if_tagged_old_versions: None, + }) + .await + .unwrap_err(); + + assert!(err.to_string().contains("at least 10 minutes")); + let versions_after = table + .list_versions() + .await + .unwrap() + .into_iter() + .map(|version| version.version) + .collect::>(); + assert_eq!(versions_after, versions_before); + } + #[tokio::test] async fn test_optimize_index() { let conn = connect("memory://").execute().await.unwrap();