diff --git a/docs/src/js/classes/Table.md b/docs/src/js/classes/Table.md index 712c15ad0..4479bf4e4 100644 --- a/docs/src/js/classes/Table.md +++ b/docs/src/js/classes/Table.md @@ -79,8 +79,9 @@ input leaves the value computed at fill time; recomputing means dropping the column and declaring it again. While a declaration reads a column, that column cannot be renamed, retyped or dropped. -Computed columns are local-only: LanceDB Cloud and Enterprise reject a -declaration. +On LanceDB Cloud and Enterprise the expression is planned by the +server, and the refresh runs as a server job -- see +[Table#refreshColumnAsync](Table.md#refreshcolumnasync). #### Parameters @@ -754,7 +755,8 @@ Fill the rows of a computed column that hold no value yet. Rows appended since the last refresh are filled by the next one; rows already filled are left as they are, so the call is idempotent and does -not observe a mutated input. Local tables only. +not observe a mutated input. Local tables only: a remote refresh runs +as a server job, through [Table#refreshColumnAsync](Table.md#refreshcolumnasync). #### Parameters @@ -782,7 +784,8 @@ job instead of blocking until it completes. The job may already be complete when returned; callers must not assume the column is filled until [Job.wait](Job.md#wait) resolves. Invalid input -- an unknown column, or one that is not computed -- rejects here rather -than failing the job. Local tables only. +than failing the job. On local tables the job runs in-process; on +LanceDB Cloud and Enterprise it is the server's backfill job. #### Parameters diff --git a/nodejs/lancedb/table.ts b/nodejs/lancedb/table.ts index 4469e41a0..a7dc8def1 100644 --- a/nodejs/lancedb/table.ts +++ b/nodejs/lancedb/table.ts @@ -537,8 +537,9 @@ export abstract class Table { * the column and declaring it again. While a declaration reads a column, * that column cannot be renamed, retyped or dropped. * - * Computed columns are local-only: LanceDB Cloud and Enterprise reject a - * declaration. + * On LanceDB Cloud and Enterprise the expression is planned by the + * server, and the refresh runs as a server job -- see + * {@link Table#refreshColumnAsync}. * @param {AddColumnsSql[] | Field | Field[] | Schema} newColumnTransforms Either: * - An array of objects with column names and SQL expressions to calculate values * - A single Arrow Field defining one column with its data type (column will be initialized with null values) @@ -567,7 +568,8 @@ export abstract class Table { * * Rows appended since the last refresh are filled by the next one; rows * already filled are left as they are, so the call is idempotent and does - * not observe a mutated input. Local tables only. + * not observe a mutated input. Local tables only: a remote refresh runs + * as a server job, through {@link Table#refreshColumnAsync}. * @param {string} column The name of the computed column to fill. * @returns {Promise} A promise that resolves to the * number of rows filled and the new version number of the table. @@ -581,7 +583,8 @@ export abstract class Table { * The job may already be complete when returned; callers must not assume * the column is filled until {@link Job.wait} resolves. Invalid input -- * an unknown column, or one that is not computed -- rejects here rather - * than failing the job. Local tables only. + * than failing the job. On local tables the job runs in-process; on + * LanceDB Cloud and Enterprise it is the server's backfill job. * @param {string} column The name of the computed column to fill. * @example * ```ts diff --git a/python/python/lancedb/remote/table.py b/python/python/lancedb/remote/table.py index b1bc5bded..aa822b913 100644 --- a/python/python/lancedb/remote/table.py +++ b/python/python/lancedb/remote/table.py @@ -964,17 +964,13 @@ class RemoteTable(Table): *, computed: Dict[str, str] | None = None, ) -> AddColumnsResult: - if computed: - raise NotImplementedError( - "computed columns are supported only on local tables" - ) - return LOOP.run(self._table.add_columns(transforms)) + return LOOP.run(self._table.add_columns(transforms, computed=computed)) def refresh_column(self, column: str): - raise NotImplementedError("computed columns are supported only on local tables") + return LOOP.run(self._table.refresh_column(column)) def refresh_column_async(self, column: str) -> Job: - raise NotImplementedError("computed columns are supported only on local tables") + return Job(LOOP.run(self._table.refresh_column_async(column))) def alter_columns( self, *alterations: Iterable[Dict[str, str]] diff --git a/python/python/lancedb/table.py b/python/python/lancedb/table.py index 9c5925cb7..4ecf6e836 100644 --- a/python/python/lancedb/table.py +++ b/python/python/lancedb/table.py @@ -1954,8 +1954,10 @@ class Table(ABC): dropping the column and declaring it again. While a declaration reads a column, that column cannot be renamed, retyped or dropped. - Local tables only; LanceDB Cloud and Enterprise raise - ``NotImplementedError``. Cannot be combined with ``transforms``. + On LanceDB Cloud and Enterprise the expression is planned by the + server, and the refresh runs as a server job -- see + [`refresh_column_async`][lancedb.table.Table.refresh_column_async]. + Cannot be combined with ``transforms``. Returns ------- @@ -1987,8 +1989,8 @@ class Table(ABC): by the next one; rows already filled are left as they are, so the call is idempotent and does not observe a mutated input. - Local tables only; LanceDB Cloud and Enterprise raise - ``NotImplementedError``. + Local tables only: a remote refresh runs as a server job, through + [`refresh_column_async`][lancedb.table.Table.refresh_column_async]. Parameters ---------- @@ -2011,8 +2013,8 @@ class Table(ABC): The job may already be complete when returned; callers must not assume the column is filled until :meth:`Job.wait` returns. Invalid input -- an unknown column, or one that is not computed -- raises here rather - than failing the job. Local tables only; LanceDB Cloud and Enterprise - raise ``NotImplementedError``. + than failing the job. On local tables the job runs in-process; on + LanceDB Cloud and Enterprise it is the server's backfill job. Examples -------- @@ -5999,7 +6001,8 @@ class AsyncTable: declaration reads a column, that column cannot be renamed, retyped or dropped. - Local tables only. Cannot be combined with ``transforms``. + On LanceDB Cloud and Enterprise the expression is planned by + the server. Cannot be combined with ``transforms``. Returns ------- @@ -6035,8 +6038,8 @@ class AsyncTable: by the next one; rows already filled are left as they are, so the call is idempotent and does not observe a mutated input. - Local tables only; LanceDB Cloud and Enterprise raise - ``NotImplementedError``. + Local tables only: a remote refresh runs as a server job, through + [`refresh_column_async`][lancedb.table.Table.refresh_column_async]. Parameters ---------- @@ -6058,8 +6061,9 @@ class AsyncTable: The job may already be complete when returned; callers must not assume the column is filled until :meth:`AsyncJob.wait` resolves. Invalid input -- an unknown column, or one that is not computed -- raises here - rather than failing the job. Local tables only; LanceDB Cloud and - Enterprise raise ``NotImplementedError``. + rather than failing the job. On local tables the job runs + in-process; on LanceDB Cloud and Enterprise it is the server's + backfill job. Examples -------- diff --git a/rust/lancedb/src/remote/table.rs b/rust/lancedb/src/remote/table.rs index 8ca84a520..a0a4cebc2 100644 --- a/rust/lancedb/src/remote/table.rs +++ b/rust/lancedb/src/remote/table.rs @@ -33,7 +33,9 @@ use crate::table::lsm_stats::GetLsmStatsResponse; use crate::table::merge::MergeFilter; use crate::table::query::create_multi_vector_plan; use crate::table::write_progress::FinishOnDrop; -use crate::table::{AlterColumnsResult, FieldMetadataUpdate, UpdateFieldMetadataResult}; +use crate::table::{ + AlterColumnsResult, FieldMetadataUpdate, RefreshColumnResult, UpdateFieldMetadataResult, +}; use crate::table::{AnyQuery, Filter, Predicate, PreprocessingOutput, TableStatistics}; use crate::utils::background_cache::BackgroundCache; use crate::utils::{ @@ -140,6 +142,40 @@ impl FreshnessHeaders { } } +/// A backfill job whose successful wait establishes a read-freshness +/// baseline on the submitting handle, so a later read cannot be served +/// from a cache older than the completed fill. A handle pinned by checkout +/// at completion keeps its time-travel view instead. +struct FreshnessJob { + inner: RemoteJob, + freshness: Arc>, + version: Arc>>, +} + +#[async_trait] +impl crate::job::JobHandle for FreshnessJob { + fn id(&self) -> Option<&str> { + crate::job::JobHandle::id(&self.inner) + } + + async fn status(&self) -> Result { + crate::job::JobHandle::status(&self.inner).await + } + + async fn wait(&self) -> Result<()> { + crate::job::JobHandle::wait(&self.inner).await?; + let version = self.version.read().await; + if version.is_none() { + self.freshness.lock().unwrap().checkout_baseline = Some(SystemTime::now()); + } + Ok(()) + } + + async fn cancel(&self) -> Result<()> { + crate::job::JobHandle::cancel(&self.inner).await + } +} + fn compute_min_timestamp( state: &FreshnessState, interval: Option, @@ -274,10 +310,10 @@ pub struct RemoteTable { identifier: String, server_version: ServerVersion, - version: RwLock>, + version: Arc>>, location: RwLock>, schema_cache: BackgroundCache, - freshness: Mutex, + freshness: Arc>, /// The branch this handle is scoped to, or `None` for the main branch. /// Stamped onto every branch-accepting request so reads and writes resolve /// on the branch's own version chain rather than main's. @@ -415,10 +451,10 @@ impl RemoteTable { namespace, identifier, server_version, - version: RwLock::new(None), + version: Arc::new(RwLock::new(None)), location: RwLock::new(None), schema_cache: BackgroundCache::new(SCHEMA_CACHE_TTL, SCHEMA_CACHE_REFRESH_WINDOW), - freshness: Mutex::new(FreshnessState::default()), + freshness: Arc::new(Mutex::new(FreshnessState::default())), branch: None, } } @@ -447,10 +483,10 @@ impl RemoteTable { namespace: self.namespace.clone(), identifier: self.identifier.clone(), server_version: self.server_version.clone(), - version: RwLock::new(None), + version: Arc::new(RwLock::new(None)), location: RwLock::new(None), schema_cache: BackgroundCache::new(SCHEMA_CACHE_TTL, SCHEMA_CACHE_REFRESH_WINDOW), - freshness: Mutex::new(FreshnessState::default()), + freshness: Arc::new(Mutex::new(FreshnessState::default())), branch, } } @@ -1268,10 +1304,10 @@ mod test_utils { namespace: vec![], identifier: name, server_version: version.map(ServerVersion).unwrap_or_default(), - version: RwLock::new(None), + version: Arc::new(RwLock::new(None)), location: RwLock::new(None), schema_cache: BackgroundCache::new(SCHEMA_CACHE_TTL, SCHEMA_CACHE_REFRESH_WINDOW), - freshness: Mutex::new(FreshnessState::default()), + freshness: Arc::new(Mutex::new(FreshnessState::default())), branch: None, } } @@ -1292,10 +1328,10 @@ mod test_utils { namespace: vec![], identifier: name, server_version: ServerVersion::default(), - version: RwLock::new(None), + version: Arc::new(RwLock::new(None)), location: RwLock::new(None), schema_cache: BackgroundCache::new(SCHEMA_CACHE_TTL, SCHEMA_CACHE_REFRESH_WINDOW), - freshness: Mutex::new(FreshnessState::default()), + freshness: Arc::new(Mutex::new(FreshnessState::default())), branch: None, } } @@ -1325,10 +1361,10 @@ mod test_utils { namespace: vec![], identifier: name, server_version: version.map(ServerVersion).unwrap_or_default(), - version: RwLock::new(None), + version: Arc::new(RwLock::new(None)), location: RwLock::new(None), schema_cache: BackgroundCache::new(SCHEMA_CACHE_TTL, SCHEMA_CACHE_REFRESH_WINDOW), - freshness: Mutex::new(FreshnessState::default()), + freshness: Arc::new(Mutex::new(FreshnessState::default())), branch: None, } } @@ -2700,13 +2736,6 @@ impl BaseTable for RemoteTable { Ok(result) } - // A declaration reaches here as AllNulls, which the remote protocol - // has no representation for. - NewColumnTransform::AllNulls(_) => { - return Err(Error::NotSupported { - message: "computed columns are supported only on local tables".into(), - }); - } _ => { return Err(Error::NotSupported { message: "Only SQL expressions are supported for adding columns".into(), @@ -2715,6 +2744,86 @@ impl BaseTable for RemoteTable { } } + async fn add_computed_columns(&self, columns: &[(String, String)]) -> Result { + self.check_mutable().await?; + // The server plans the declaration: expression validation, type + // inference and the persisted binding all happen there. + let entries = columns + .iter() + .map( + |(name, expression)| lance_namespace::models::AddColumnsEntry { + name: name.clone(), + computed: Some(Some(expression.clone())), + ..Default::default() + }, + ) + .collect::>(); + let mut body = serde_json::json!({ "new_columns": entries }); + self.apply_branch_body(&mut body); + let request = self + .client + .post(&format!("/v1/table/{}/add_columns/", self.identifier)) + .json(&body); + let (request_id, response) = self.send(request, true).await?; + let response = self.check_table_response(&request_id, response).await?; + let body = response.text().await.err_to_http(request_id.clone())?; + + if body.trim().is_empty() { + // Backward compatible with old servers + return Ok(AddColumnsResult { version: 0 }); + } + + let result: AddColumnsResult = serde_json::from_str(&body).map_err(|e| Error::Http { + source: format!("Failed to parse add_columns response: {}", e).into(), + request_id, + status_code: None, + })?; + + self.invalidate_schema_cache(); + self.track_write_version(result.version); + + Ok(result) + } + + async fn refresh_column(&self, _column: &str) -> Result { + // The server runs a refresh as a job and does not report a fill + // count, so the blocking form has no honest result to return. + Err(Error::NotSupported { + message: "a remote refresh runs as a server job; use refresh_column_async and \ + wait on the returned handle" + .into(), + }) + } + + async fn refresh_column_async(&self, column: &str) -> Result { + self.check_mutable().await?; + let mut body = serde_json::json!({ "column": column }); + self.apply_branch_body(&mut body); + let request = self + .client + .post(&format!("/v1/table/{}/backfill_column", self.identifier)) + .json(&body); + let (request_id, response) = self.send(request, true).await?; + let response = self.check_table_response(&request_id, response).await?; + let body = response.text().await.err_to_http(request_id.clone())?; + + #[derive(serde::Deserialize)] + struct BackfillResponse { + job_id: String, + } + let response: BackfillResponse = serde_json::from_str(&body).map_err(|e| Error::Http { + source: format!("Failed to parse backfill_column response: {}", e).into(), + request_id, + status_code: None, + })?; + + Ok(Job::new(Box::new(FreshnessJob { + inner: RemoteJob::new(self.client.clone(), response.job_id), + freshness: self.freshness.clone(), + version: self.version.clone(), + }))) + } + async fn alter_columns(&self, alterations: &[ColumnAlteration]) -> Result { self.check_mutable().await?; let body = alterations @@ -6456,37 +6565,346 @@ mod tests { assert_eq!(result.version, if old_server { 0 } else { 43 }); } - /// Computed columns are local-only. Both halves say so here rather than - /// reaching the wire and failing somewhere less legible. + /// A declaration is sent as `{name, computed}` entries for the server to + /// plan; the client never types the expression itself. #[tokio::test] - async fn test_computed_columns_are_refused() { - let table = Table::new_with_handler("my_table", |request| -> http::Response { - panic!("unexpected request: {}", request.url().path()) + async fn test_add_computed_columns_sends_the_expression() { + let table = Table::new_with_handler("my_table", |request| { + assert_eq!(request.method(), "POST"); + assert_eq!(request.url().path(), "/v1/table/my_table/add_columns/"); + let body = request.body().unwrap().as_bytes().unwrap(); + let value: serde_json::Value = serde_json::from_slice(body).unwrap(); + assert_eq!( + value["new_columns"], + serde_json::json!([{"name": "doubled", "computed": "x * 2"}]) + ); + http::Response::builder() + .status(200) + .body(r#"{"version": 7}"#) + .unwrap() }); - let declared = Arc::new(Schema::new(vec![Field::new( - "doubled", - DataType::Int32, - true, - )])); - let err = table + let result = table .add_columns() - .transform(NewColumnTransform::AllNulls(declared)) + .computed("doubled", "x * 2") .execute() .await - .unwrap_err(); - assert!( - matches!(&err, Error::NotSupported { message } if message.contains("local tables")), - "{err:?}" - ); + .unwrap(); + assert_eq!(result.version, 7); + } + + /// A remote refresh is a server job: the async form returns its handle, + /// and the blocking form refuses rather than invent a fill count. + #[tokio::test] + async fn test_refresh_column_async_submits_a_backfill_job() { + let table = Table::new_with_handler("my_table", |request| { + assert_eq!(request.method(), "POST"); + assert_eq!(request.url().path(), "/v1/table/my_table/backfill_column"); + let body = request.body().unwrap().as_bytes().unwrap(); + let value: serde_json::Value = serde_json::from_slice(body).unwrap(); + assert_eq!(value["column"], "doubled"); + http::Response::builder() + .status(202) + .body(r#"{"job_id": "j-42"}"#) + .unwrap() + }); + + let job = table.refresh_column_async("doubled").await.unwrap(); + assert_eq!(job.id(), Some("j-42")); let err = table.refresh_column("doubled").await.unwrap_err(); assert!( - matches!(&err, Error::NotSupported { message } if message.contains("local tables")), + matches!(&err, Error::NotSupported { message } + if message.contains("refresh_column_async")), "{err:?}" ); } + /// The gate's reproducer: after a successful wait, a same-handle read + /// must carry a freshness baseline so a stale server cache cannot serve + /// the pre-backfill snapshot. + #[tokio::test] + async fn test_backfill_wait_establishes_read_freshness() { + let saw_min_timestamp = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let saw = saw_min_timestamp.clone(); + let table = + Table::new_with_handler("my_table", move |request| match request.url().path() { + "/v1/table/my_table/backfill_column" => http::Response::builder() + .status(202) + .body(r#"{"job_id": "j-7"}"#.to_string()) + .unwrap(), + "/v1/jobs/describe" => http::Response::builder() + .status(200) + .body(r#"{"job_id": "j-7", "job_state": "DONE"}"#.to_string()) + .unwrap(), + "/v1/table/my_table/count_rows/" => { + saw.store( + request.headers().contains_key("x-lancedb-min-timestamp"), + std::sync::atomic::Ordering::SeqCst, + ); + http::Response::builder() + .status(200) + .body("1".to_string()) + .unwrap() + } + path => panic!("unexpected request: {path}"), + }); + + let job = table.refresh_column_async("doubled").await.unwrap(); + job.wait().await.unwrap(); + table.count_rows(None).await.unwrap(); + assert!( + saw_min_timestamp.load(std::sync::atomic::Ordering::SeqCst), + "read after wait carried no freshness baseline" + ); + } + + /// A checkout after submission wins over the completion fence: the + /// pinned view must not regain a timestamp floor from the job. + #[tokio::test] + async fn test_checkout_after_submit_beats_the_completion_fence() { + let saw_min_timestamp = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let saw = saw_min_timestamp.clone(); + let table = + Table::new_with_handler("my_table", move |request| match request.url().path() { + "/v1/table/my_table/backfill_column" => http::Response::builder() + .status(202) + .body(r#"{"job_id": "j-8"}"#.to_string()) + .unwrap(), + "/v1/jobs/describe" => http::Response::builder() + .status(200) + .body(r#"{"job_id": "j-8", "job_state": "DONE"}"#.to_string()) + .unwrap(), + "/v1/table/my_table/describe/" => { + let schema = Schema::new(vec![Field::new("x", DataType::Int32, true)]); + http::Response::builder() + .status(200) + .body(describe_response(&schema)) + .unwrap() + } + "/v1/table/my_table/count_rows/" => { + saw.store( + request.headers().contains_key("x-lancedb-min-timestamp"), + std::sync::atomic::Ordering::SeqCst, + ); + http::Response::builder() + .status(200) + .body("1".to_string()) + .unwrap() + } + path => panic!("unexpected request: {path}"), + }); + + let job = table.refresh_column_async("doubled").await.unwrap(); + table.checkout(3).await.unwrap(); + job.wait().await.unwrap(); + table.count_rows(None).await.unwrap(); + assert!( + !saw_min_timestamp.load(std::sync::atomic::Ordering::SeqCst), + "completion fence overrode an explicit checkout" + ); + } + + /// Tag checkout resets freshness state wholesale; the fence must not + /// survive it. + #[tokio::test] + async fn test_tag_checkout_after_submit_beats_the_completion_fence() { + let saw_min_timestamp = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let saw = saw_min_timestamp.clone(); + let table = + Table::new_with_handler("my_table", move |request| match request.url().path() { + "/v1/table/my_table/backfill_column" => http::Response::builder() + .status(202) + .body(r#"{"job_id": "j-9"}"#.to_string()) + .unwrap(), + "/v1/jobs/describe" => http::Response::builder() + .status(200) + .body(r#"{"job_id": "j-9", "job_state": "DONE"}"#.to_string()) + .unwrap(), + "/v1/table/my_table/tags/version/" => http::Response::builder() + .status(200) + .body(r#"{"version": 5}"#.to_string()) + .unwrap(), + "/v1/table/my_table/describe/" => { + let schema = Schema::new(vec![Field::new("x", DataType::Int32, true)]); + http::Response::builder() + .status(200) + .body(describe_response(&schema)) + .unwrap() + } + "/v1/table/my_table/count_rows/" => { + saw.store( + request.headers().contains_key("x-lancedb-min-timestamp"), + std::sync::atomic::Ordering::SeqCst, + ); + http::Response::builder() + .status(200) + .body("1".to_string()) + .unwrap() + } + path => panic!("unexpected request: {path}"), + }); + + let job = table.refresh_column_async("doubled").await.unwrap(); + table.checkout_tag("v1").await.unwrap(); + job.wait().await.unwrap(); + table.count_rows(None).await.unwrap(); + assert!( + !saw_min_timestamp.load(std::sync::atomic::Ordering::SeqCst), + "completion fence overrode a tag checkout" + ); + } + + /// A checkout landing while the submission request is in flight advances + /// the epoch past the token captured at submit. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_checkout_during_submission_beats_the_completion_fence() { + let saw_min_timestamp = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let saw = saw_min_timestamp.clone(); + let (release_tx, release_rx) = std::sync::mpsc::channel::<()>(); + let release_rx = Arc::new(std::sync::Mutex::new(release_rx)); + let (arrived_tx, arrived_rx) = std::sync::mpsc::channel::<()>(); + let arrived_tx = Arc::new(std::sync::Mutex::new(arrived_tx)); + let table = Table::new_with_handler("my_table", move |request| { + match request.url().path() { + "/v1/table/my_table/backfill_column" => { + // Signal arrival, then hold the response until the + // test's checkout completes. + arrived_tx.lock().unwrap().send(()).unwrap(); + release_rx + .lock() + .unwrap() + .recv_timeout(std::time::Duration::from_secs(10)) + .unwrap(); + http::Response::builder() + .status(202) + .body(r#"{"job_id": "j-10"}"#.to_string()) + .unwrap() + } + "/v1/jobs/describe" => http::Response::builder() + .status(200) + .body(r#"{"job_id": "j-10", "job_state": "DONE"}"#.to_string()) + .unwrap(), + "/v1/table/my_table/describe/" => { + let schema = Schema::new(vec![Field::new("x", DataType::Int32, true)]); + http::Response::builder() + .status(200) + .body(describe_response(&schema)) + .unwrap() + } + "/v1/table/my_table/count_rows/" => { + saw.store( + request.headers().contains_key("x-lancedb-min-timestamp"), + std::sync::atomic::Ordering::SeqCst, + ); + http::Response::builder() + .status(200) + .body("1".to_string()) + .unwrap() + } + path => panic!("unexpected request: {path}"), + } + }); + + let submit = tokio::spawn({ + let table = table.clone(); + async move { table.refresh_column_async("doubled").await } + }); + tokio::task::spawn_blocking(move || { + arrived_rx + .recv_timeout(std::time::Duration::from_secs(10)) + .unwrap() + }) + .await + .unwrap(); + table.checkout(7).await.unwrap(); + release_tx.send(()).unwrap(); + + let job = submit.await.unwrap().unwrap(); + job.wait().await.unwrap(); + table.count_rows(None).await.unwrap(); + assert!( + !saw_min_timestamp.load(std::sync::atomic::Ordering::SeqCst), + "completion fence overrode a checkout that landed mid-submission" + ); + } + + /// checkout_latest keeps the handle on latest, so a completed backfill + /// must still establish its post-fill baseline -- strictly later than the + /// checkout's own, or a pre-fill cache could still serve. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_checkout_latest_during_submission_keeps_the_fence() { + let seen_min_timestamp = Arc::new(std::sync::Mutex::new(None::)); + let saw = seen_min_timestamp.clone(); + let (release_tx, release_rx) = std::sync::mpsc::channel::<()>(); + let release_rx = Arc::new(std::sync::Mutex::new(release_rx)); + let (arrived_tx, arrived_rx) = std::sync::mpsc::channel::<()>(); + let arrived_tx = Arc::new(std::sync::Mutex::new(arrived_tx)); + let table = + Table::new_with_handler("my_table", move |request| match request.url().path() { + "/v1/table/my_table/backfill_column" => { + arrived_tx.lock().unwrap().send(()).unwrap(); + release_rx + .lock() + .unwrap() + .recv_timeout(std::time::Duration::from_secs(10)) + .unwrap(); + http::Response::builder() + .status(202) + .body(r#"{"job_id": "j-11"}"#.to_string()) + .unwrap() + } + "/v1/jobs/describe" => http::Response::builder() + .status(200) + .body(r#"{"job_id": "j-11", "job_state": "DONE"}"#.to_string()) + .unwrap(), + "/v1/table/my_table/count_rows/" => { + *saw.lock().unwrap() = request + .headers() + .get("x-lancedb-min-timestamp") + .map(|v| v.to_str().unwrap().to_string()); + http::Response::builder() + .status(200) + .body("1".to_string()) + .unwrap() + } + path => panic!("unexpected request: {path}"), + }); + + let submit = tokio::spawn({ + let table = table.clone(); + async move { table.refresh_column_async("doubled").await } + }); + tokio::task::spawn_blocking(move || { + arrived_rx + .recv_timeout(std::time::Duration::from_secs(10)) + .unwrap() + }) + .await + .unwrap(); + table.checkout_latest().await.unwrap(); + let after_checkout = SystemTime::now(); + // Real separation between the checkout baseline and completion. + tokio::time::sleep(std::time::Duration::from_millis(50)).await; + release_tx.send(()).unwrap(); + + let job = submit.await.unwrap().unwrap(); + job.wait().await.unwrap(); + table.count_rows(None).await.unwrap(); + let header = seen_min_timestamp + .lock() + .unwrap() + .clone() + .expect("no baseline"); + let sent: SystemTime = chrono::DateTime::parse_from_rfc3339(&header) + .unwrap() + .into(); + assert!( + sent > after_checkout, + "baseline {header} did not advance past the checkout" + ); + } + #[tokio::test] async fn test_prewarm_index() { let table = Table::new_with_handler("my_table", |request| { diff --git a/rust/lancedb/src/table.rs b/rust/lancedb/src/table.rs index 093d63438..2e16b0940 100644 --- a/rust/lancedb/src/table.rs +++ b/rust/lancedb/src/table.rs @@ -748,6 +748,10 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync { read_columns: Option>, ) -> Result; /// Declare computed columns, each defined by a SQL expression. + /// + /// Where the declaration is planned depends on the backend: a local table + /// validates and types the expression itself, a remote one sends the text + /// for the server to plan. async fn add_computed_columns( &self, _columns: &[(String, String)], @@ -1672,7 +1676,8 @@ impl Table { /// filled are left as they are, so the call is idempotent and does not /// observe a mutated input. /// - /// Local tables only. + /// Local tables only: a remote refresh runs as a server job, through + /// [`Table::refresh_column_async`]. /// /// ``` /// # use lancedb::Table; @@ -1692,8 +1697,9 @@ impl Table { /// The job may already be complete when returned, and callers must not /// assume the column is filled until [`Job::wait`] returns. Invalid input /// -- an unknown column, or one that is not computed -- is reported by - /// this call rather than by the job. Local tables only: LanceDB Cloud and - /// Enterprise reject with `NotSupported`. + /// this call rather than by the job. On local tables the job runs as an + /// in-process task; on LanceDB Cloud and Enterprise it is the server's + /// backfill job. /// /// ``` /// # use lancedb::Table; diff --git a/rust/lancedb/src/table/add_columns.rs b/rust/lancedb/src/table/add_columns.rs index 6aa2ce86a..67764c346 100644 --- a/rust/lancedb/src/table/add_columns.rs +++ b/rust/lancedb/src/table/add_columns.rs @@ -61,8 +61,9 @@ impl AddColumnsBuilder { /// column and declaring it again. An input cannot be renamed, retyped or /// dropped while a declaration reads it, since the expression names it. /// - /// Local tables only: LanceDB Cloud and Enterprise reject a declaration - /// with `NotSupported`. + /// On LanceDB Cloud and Enterprise the expression is planned by the + /// server, and the refresh runs as a server job -- see + /// [`Table::refresh_column_async`](super::Table::refresh_column_async). /// /// ``` /// # use lancedb::Table;