mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-30 00:45:37 +00:00
feat: return a cleanup job from drop_function (#4237)
Dropping a Function can leave its content to a server-side cleanup job, so the name drop and the content deletion become separate events a caller may want to wait on. `drop_function_async` returns the unbind result alongside a `Job` for the cleanup, the same shape `drop_table_async` and `drop_materialized_view_async` already use: a `202` carries the job id, and a `200` — an inline deletion, or a name that was not bound — yields an already-finished job with no id. A `202` without a usable job id is rejected rather than silently reported as finished. `drop_function` keeps its `bool` result and now delegates, so nothing changes for callers that do not care when the content goes. Available on `Connection` in Rust and on both the sync and asyncio Python connections.
This commit is contained in:
@@ -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: ...
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -768,6 +768,21 @@ impl Connection {
|
||||
})
|
||||
}
|
||||
|
||||
pub fn drop_function_async(
|
||||
self_: PyRef<'_, Self>,
|
||||
name: String,
|
||||
version: String,
|
||||
) -> PyResult<Bound<'_, PyAny>> {
|
||||
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>,
|
||||
|
||||
@@ -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<str>,
|
||||
version: impl AsRef<str>,
|
||||
) -> 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
|
||||
|
||||
@@ -393,6 +393,18 @@ pub trait Database:
|
||||
async fn drop_function(&self, _name: &str, _version: &str) -> Result<bool> {
|
||||
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(
|
||||
|
||||
@@ -1005,6 +1005,10 @@ impl<S: HttpSend> Database for RemoteDatabase<S> {
|
||||
}
|
||||
|
||||
async fn drop_function(&self, name: &str, version: &str) -> Result<bool> {
|
||||
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<S: HttpSend> Database for RemoteDatabase<S> {
|
||||
}));
|
||||
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| {
|
||||
|
||||
Reference in New Issue
Block a user