fix: reject unsafe version cleanup windows

This commit is contained in:
Gatefixer
2026-08-06 04:00:56 +00:00
parent c7ea91f3ea
commit d8c5695b7b
9 changed files with 163 additions and 22 deletions
+3
View File
@@ -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)>
+6 -3
View File
@@ -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});
```
***
+9 -1
View File
@@ -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);
});
+9 -3
View File
@@ -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<OptimizeOptions>): Promise<OptimizeStats>;
/** List all indices that have been created with {@link Table.createIndex} */
+2
View File
@@ -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 {
+39 -11
View File
@@ -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
+11 -1
View File
@@ -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
+2
View File
@@ -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 {
+82 -3
View File
@@ -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<Duration>,
delete_unverified: Option<bool>,
) -> 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<Duration>,
/// 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<OptimizeStats> {
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::<Vec<_>>();
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::<Vec<_>>();
assert_eq!(versions_after, versions_before);
}
#[tokio::test]
async fn test_optimize_index() {
let conn = connect("memory://").execute().await.unwrap();