diff --git a/python/python/lancedb/_lancedb.pyi b/python/python/lancedb/_lancedb.pyi index 7c6611608..6d873850d 100644 --- a/python/python/lancedb/_lancedb.pyi +++ b/python/python/lancedb/_lancedb.pyi @@ -153,6 +153,9 @@ class Connection(object): async def get_function(self, name: str, version: str) -> str: ... async def list_functions(self) -> List[str]: ... async def drop_function(self, name: str, version: str) -> bool: ... + async def drop_function_async( + self, name: str, version: str + ) -> Tuple[bool, Job]: ... async def create_secret( self, name: str, value: str, namespace_path: Optional[List[str]] = None ) -> None: ... diff --git a/python/python/lancedb/db.py b/python/python/lancedb/db.py index d17e897cf..b0db22532 100644 --- a/python/python/lancedb/db.py +++ b/python/python/lancedb/db.py @@ -18,6 +18,7 @@ from typing import ( Literal, Optional, Sequence, + Tuple, Union, ) from uuid import UUID @@ -828,12 +829,26 @@ class DBConnection(EnforceOverrides): ) def drop_function(self, name: str, *, version: str) -> bool: - """Remove the current Function name binding from the remote catalog. + """Drop a Function name and the object it was bound to. - The requested version must exist in the currently named object. - Object history and existing computed-column references are retained. - Returns True when the name was removed and False when it was absent. - Local connections raise NotImplementedError. + The requested version must exist in the currently named object. Returns + True when the name was removed and False when it was absent. The + object's content is deleted with it, which may finish after this + returns; use :meth:`drop_function_async` to wait for that. Local + connections raise NotImplementedError. + """ + raise NotImplementedError( + "Function catalog operations are not supported for this connection type" + ) + + def drop_function_async(self, name: str, *, version: str) -> Tuple[bool, Job]: + """Drop a Function name and return its cleanup job. + + The name is unbound before this returns; the object's content may still + be being deleted. Call :meth:`Job.wait` to wait for that to finish. When + the server deletes inline, or when nothing was bound, the returned job + is already finished and has no id. Local connections raise + NotImplementedError. """ raise NotImplementedError( "Function catalog operations are not supported for this connection type" @@ -1683,6 +1698,11 @@ class LanceDBConnection(DBConnection): def drop_function(self, name: str, *, version: str) -> bool: return LOOP.run(self._conn.drop_function(name, version=version)) + @override + def drop_function_async(self, name: str, *, version: str) -> Tuple[bool, Job]: + dropped, job = LOOP.run(self._conn.drop_function_async(name, version=version)) + return dropped, Job(job) + @override def create_secret( self, name: str, value: str, *, namespace_path: Optional[List[str]] = None @@ -2596,9 +2616,21 @@ class AsyncConnection(object): ] async def drop_function(self, name: str, *, version: str) -> bool: - """Remove the current name binding, retaining the object and its history.""" + """Drop a Function name and the object it was bound to.""" return await self._inner.drop_function(name, version) + async def drop_function_async( + self, name: str, *, version: str + ) -> Tuple[bool, AsyncJob]: + """Drop a Function name and return its cleanup job. + + The name is unbound before this returns; the object's content may still + be being deleted. Await :meth:`AsyncJob.wait` to wait for that to + finish. + """ + dropped, job = await self._inner.drop_function_async(name, version) + return dropped, AsyncJob(job) + async def create_secret( self, name: str, value: str, *, namespace_path: Optional[List[str]] = None ) -> None: diff --git a/python/python/lancedb/remote/db.py b/python/python/lancedb/remote/db.py index c9e857495..ef40ef007 100644 --- a/python/python/lancedb/remote/db.py +++ b/python/python/lancedb/remote/db.py @@ -16,6 +16,7 @@ from typing import ( List, Optional, Sequence, + Tuple, Union, ) from urllib.parse import urlparse @@ -876,6 +877,11 @@ class RemoteDBConnection(DBConnection): def drop_function(self, name: str, *, version: str) -> bool: return LOOP.run(self._conn.drop_function(name, version=version)) + @override + def drop_function_async(self, name: str, *, version: str) -> Tuple[bool, Job]: + dropped, job = LOOP.run(self._conn.drop_function_async(name, version=version)) + return dropped, Job(job) + @override def create_secret( self, name: str, value: str, *, namespace_path: Optional[List[str]] = None diff --git a/python/src/connection.rs b/python/src/connection.rs index a972f4351..3385e6b2d 100644 --- a/python/src/connection.rs +++ b/python/src/connection.rs @@ -768,6 +768,21 @@ impl Connection { }) } + pub fn drop_function_async( + self_: PyRef<'_, Self>, + name: String, + version: String, + ) -> PyResult> { + let inner = self_.get_inner()?.clone(); + future_into_py(self_.py(), async move { + inner + .drop_function_async(name, version) + .await + .infer_error() + .map(|(dropped, job)| (dropped, crate::job::Job::new(job))) + }) + } + #[pyo3(signature = (name, value, namespace_path=None))] pub fn create_secret( self_: PyRef<'_, Self>, diff --git a/rust/lancedb/src/connection.rs b/rust/lancedb/src/connection.rs index 7234a911f..d43f8f744 100644 --- a/rust/lancedb/src/connection.rs +++ b/rust/lancedb/src/connection.rs @@ -649,6 +649,22 @@ impl Connection { .await } + /// Start dropping a Function and return its cleanup job. + /// + /// The name is unbound before this returns; the object's content may still be being + /// deleted. Await [`Job::wait`][crate::job::Job::wait] to wait for that to finish. When + /// the server deletes inline, or when nothing was bound, the returned job is already + /// finished and has no id. Local databases return [`Error::NotSupported`]. + pub async fn drop_function_async( + &self, + name: impl AsRef, + version: impl AsRef, + ) -> Result<(bool, crate::job::Job)> { + self.internal + .drop_function_async(name.as_ref(), version.as_ref()) + .await + } + /// Create a named Secret in this database. /// /// Fails if the name is taken, so a create can never silently become a diff --git a/rust/lancedb/src/database.rs b/rust/lancedb/src/database.rs index 0adac778c..91b200043 100644 --- a/rust/lancedb/src/database.rs +++ b/rust/lancedb/src/database.rs @@ -393,6 +393,18 @@ pub trait Database: async fn drop_function(&self, _name: &str, _version: &str) -> Result { function_catalog_not_supported() } + /// Start dropping a Function and return a handle to the cleanup job. + /// + /// Backends without asynchronous cleanup complete the drop before returning an + /// already-finished job. + async fn drop_function_async( + &self, + name: &str, + version: &str, + ) -> Result<(bool, crate::job::Job)> { + let dropped = self.drop_function(name, version).await?; + Ok((dropped, crate::job::Job::new_done())) + } /// Create a named Secret in this database. Fails if the name is taken, so /// a create can never silently become a rotation. async fn create_secret( diff --git a/rust/lancedb/src/remote/db.rs b/rust/lancedb/src/remote/db.rs index 7061ac979..7a59db665 100644 --- a/rust/lancedb/src/remote/db.rs +++ b/rust/lancedb/src/remote/db.rs @@ -1005,6 +1005,10 @@ impl Database for RemoteDatabase { } async fn drop_function(&self, name: &str, version: &str) -> Result { + Ok(self.drop_function_async(name, version).await?.0) + } + + async fn drop_function_async(&self, name: &str, version: &str) -> Result<(bool, Job)> { let function_id = build_object_identifier("Function name", name, &[])?; let req = self .client @@ -1014,8 +1018,34 @@ impl Database for RemoteDatabase { })); let (request_id, response) = self.client.send(req).await?; let response = self.client.check_response(&request_id, response).await?; - let response: RemoteDropFunctionResponse = response.json().await.err_to_http(request_id)?; - Ok(response.dropped) + let status = response.status(); + let body = response.text().await.err_to_http(request_id.clone())?; + let dropped: RemoteDropFunctionResponse = + serde_json::from_str(&body).map_err(|source| Error::Http { + source: Box::new(source), + request_id: request_id.clone(), + status_code: Some(status), + })?; + let job = match status { + StatusCode::OK => Job::new_done(), + StatusCode::ACCEPTED => { + let job_id = extract_job_id(&body).ok_or_else(|| Error::Http { + source: "asynchronous Function drop response did not contain a valid job_id" + .into(), + request_id, + status_code: Some(status), + })?; + Job::new(Box::new(RemoteJob::new(self.client.clone(), job_id))) + } + _ => { + return Err(Error::Http { + source: "Function drop must return 200 OK or 202 Accepted".into(), + request_id, + status_code: Some(status), + }); + } + }; + Ok((dropped.dropped, job)) } async fn create_secret( @@ -1783,6 +1813,53 @@ mod tests { assert_eq!(job.id(), Some("j1-mv-drop")); } + #[tokio::test] + async fn test_drop_function_async_returns_job() { + let db = super::RemoteDatabase::new_mock(|request| { + assert_eq!(request.method(), "POST"); + assert_eq!(request.url().path(), "/v1/function/embed/drop"); + http::Response::builder() + .status(202) + .body(serde_json::json!({"dropped": true, "job_id": "j1-fn-drop"}).to_string()) + .unwrap() + }); + let (dropped, job) = db.drop_function_async("embed", "1").await.unwrap(); + assert!(dropped); + assert_eq!(job.id(), Some("j1-fn-drop")); + } + + /// An unbound name and an inline deletion both answer `200`: the name is gone and + /// nothing is left to wait for, so the job is already finished. + #[tokio::test] + async fn test_drop_function_async_completed_inline() { + let db = super::RemoteDatabase::new_mock(|_| { + http::Response::builder() + .status(200) + .body(serde_json::json!({"dropped": false}).to_string()) + .unwrap() + }); + let (dropped, job) = db.drop_function_async("embed", "1").await.unwrap(); + assert!(!dropped); + assert_eq!(job.id(), None); + assert_eq!(job.status().await.unwrap(), "finished"); + job.wait().await.unwrap(); + } + + #[tokio::test] + async fn test_drop_function_async_rejects_incomplete_acceptance() { + for body in [ + r#"{"dropped":true}"#, + r#"{"dropped":true,"job_id":""}"#, + r#"{"dropped":true,"job_id":null}"#, + ] { + let db = super::RemoteDatabase::new_mock(move |_| { + http::Response::builder().status(202).body(body).unwrap() + }); + let error = db.drop_function_async("embed", "1").await.err().unwrap(); + assert!(error.to_string().contains("valid job_id")); + } + } + #[tokio::test] async fn test_drop_materialized_view_completed_inline() { let db = super::RemoteDatabase::new_mock(|request| {