diff --git a/docs/src/js/classes/Connection.md b/docs/src/js/classes/Connection.md index 75af5102e..a51fb986a 100644 --- a/docs/src/js/classes/Connection.md +++ b/docs/src/js/classes/Connection.md @@ -557,9 +557,10 @@ on the returned job to know when cleanup has finished. abstract dropView(name, namespacePath?): Promise ``` -Drop the view named `name`. +Drop the view named `name` and wait for its definition to be deleted. -The tables it reads are untouched: a view holds no rows of its own. +The tables it reads are untouched: a view holds no rows of its own. Use +[dropViewAsync](Connection.md#dropviewasync) to retain the cleanup job instead of waiting on it. #### Parameters @@ -573,6 +574,30 @@ The tables it reads are untouched: a view holds no rows of its own. *** +### dropViewAsync() + +```ts +abstract dropViewAsync(name, namespacePath?): Promise +``` + +Start dropping the view named `name` and return the job deleting its +definition, without waiting for completion. + +The name is free before this resolves. When nothing was bound to it, the +returned job is already finished and has no id. + +#### Parameters + +* **name**: `string` + +* **namespacePath?**: `string`[] + +#### Returns + +`Promise`<[`Job`](Job.md)> + +*** + ### isOpen() ```ts diff --git a/nodejs/__test__/remote.test.ts b/nodejs/__test__/remote.test.ts index 4c1d21604..ffac7aae1 100644 --- a/nodejs/__test__/remote.test.ts +++ b/nodejs/__test__/remote.test.ts @@ -161,6 +161,37 @@ describe("remote connection", () => { ); }); + it("reports the cleanup job when a view drop is accepted", async () => { + await withMockDatabase( + (req, res) => { + expect(req.method).toBe("POST"); + expect(req.url).toBe("/v1/view/adults/drop"); + res + .writeHead(202, { "content-type": "application/json" }) + .end('{"job_id": "j1-do-abc"}'); + }, + async (db) => { + const job = await db.dropViewAsync("adults"); + expect(job.id).toBe("j1-do-abc"); + }, + ); + }); + + it("reports a finished job when a view drop had nothing to delete", async () => { + await withMockDatabase( + (req, res) => { + expect(req.url).toBe("/v1/view/adults/drop"); + res.writeHead(200, { "content-type": "application/json" }).end("{}"); + }, + async (db) => { + // A 200 means the name was not bound, so there is no cleanup to wait on. + const job = await db.dropViewAsync("adults"); + expect(job.id).toBeNull(); + await job.wait(); + }, + ); + }); + it("should accept partial connection options", async () => { await connect("db://test", { apiKey: "fake", diff --git a/nodejs/lancedb/connection.ts b/nodejs/lancedb/connection.ts index 9637d60c2..40b5ce72a 100644 --- a/nodejs/lancedb/connection.ts +++ b/nodejs/lancedb/connection.ts @@ -404,12 +404,22 @@ export abstract class Connection { ): Promise; /** - * Drop the view named `name`. + * Drop the view named `name` and wait for its definition to be deleted. * - * The tables it reads are untouched: a view holds no rows of its own. + * The tables it reads are untouched: a view holds no rows of its own. Use + * {@link dropViewAsync} to retain the cleanup job instead of waiting on it. */ abstract dropView(name: string, namespacePath?: string[]): Promise; + /** + * Start dropping the view named `name` and return the job deleting its + * definition, without waiting for completion. + * + * The name is free before this resolves. When nothing was bound to it, the + * returned job is already finished and has no id. + */ + abstract dropViewAsync(name: string, namespacePath?: string[]): Promise; + /** * The names of the views in one namespace. * @@ -759,6 +769,10 @@ export class LocalConnection extends Connection { return this.inner.dropView(name, namespacePath ?? []); } + async dropViewAsync(name: string, namespacePath?: string[]): Promise { + return new Job(await this.inner.dropViewAsync(name, namespacePath ?? [])); + } + async listViews(namespacePath?: string[]): Promise { return this.inner.listViews(namespacePath ?? []); } diff --git a/nodejs/src/connection.rs b/nodejs/src/connection.rs index ee2a9470b..f43c841ee 100644 --- a/nodejs/src/connection.rs +++ b/nodejs/src/connection.rs @@ -427,7 +427,8 @@ impl Connection { ViewDescription::from_inner(view) } - /// Drop a view. The tables it reads are untouched. + /// Drop a view and wait for its definition to be deleted. The tables it + /// reads are untouched. #[napi(catch_unwind)] pub async fn drop_view( &self, @@ -441,6 +442,22 @@ impl Connection { .default_error() } + /// Start dropping a view and return the job deleting its definition. + #[napi(catch_unwind)] + pub async fn drop_view_async( + &self, + name: String, + namespace_path: Option>, + ) -> napi::Result { + let ns = namespace_path.unwrap_or_default(); + let job = self + .get_inner()? + .drop_view_async(&name, &ns) + .await + .default_error()?; + Ok(crate::job::Job::new(job)) + } + /// The names of the views in one namespace. #[napi(catch_unwind)] pub async fn list_views( diff --git a/python/python/lancedb/_lancedb.pyi b/python/python/lancedb/_lancedb.pyi index ba4e3a0df..672ca5713 100644 --- a/python/python/lancedb/_lancedb.pyi +++ b/python/python/lancedb/_lancedb.pyi @@ -180,6 +180,9 @@ class Connection(object): async def drop_view( self, name: str, namespace_path: Optional[List[str]] = None ) -> None: ... + async def drop_view_async( + self, name: str, namespace_path: Optional[List[str]] = None + ) -> Job: ... async def list_views( self, namespace_path: Optional[List[str]] = None ) -> List[str]: ... diff --git a/python/python/lancedb/db.py b/python/python/lancedb/db.py index 10d52c63e..e53b78896 100644 --- a/python/python/lancedb/db.py +++ b/python/python/lancedb/db.py @@ -983,15 +983,30 @@ class DBConnection(EnforceOverrides): def drop_view( self, name: str, *, namespace_path: Optional[List[str]] = None ) -> None: - """Drop a view. + """Drop a view and wait for its definition to be deleted. - The tables it reads are untouched: a view holds no rows of its own. + The tables it reads are untouched: a view holds no rows of its own. Use + :meth:`drop_view_async` to get the cleanup job instead of waiting on it. Local connections raise ``NotImplementedError``. """ raise NotImplementedError( "View operations are not supported for this connection type" ) + def drop_view_async( + self, name: str, *, namespace_path: Optional[List[str]] = None + ) -> "Job[None]": + """Start dropping a view and return the job deleting its definition. + + The name is free before this returns. Call :meth:`Job.wait` to wait for + the definition dataset to be deleted. When nothing was bound to the + name, the returned job is already finished and has no id. Local + connections raise ``NotImplementedError``. + """ + raise NotImplementedError( + "View operations are not supported for this connection type" + ) + def list_views(self, *, namespace_path: Optional[List[str]] = None) -> List[str]: """The names of the views in one namespace. @@ -1836,6 +1851,14 @@ class LanceDBConnection(DBConnection): ) -> None: LOOP.run(self._conn.drop_view(name, namespace_path=namespace_path)) + @override + def drop_view_async( + self, name: str, *, namespace_path: Optional[List[str]] = None + ) -> "Job[None]": + return Job( + LOOP.run(self._conn.drop_view_async(name, namespace_path=namespace_path)) + ) + @override def list_views(self, *, namespace_path: Optional[List[str]] = None) -> List[str]: return LOOP.run(self._conn.list_views(namespace_path=namespace_path)) @@ -2821,9 +2844,30 @@ class AsyncConnection(object): async def drop_view( self, name: str, *, namespace_path: Optional[List[str]] = None ) -> None: - """Drop a view. The tables it reads are untouched.""" + """Drop a view and wait for its definition to be deleted. + + The tables it reads are untouched. Use :meth:`drop_view_async` to get + the cleanup job instead of waiting on it. + """ await self._inner.drop_view(name, list(namespace_path or [])) + async def drop_view_async( + self, + name: str, + *, + namespace_path: Optional[List[str]] = None, + ) -> AsyncJob[None]: + """Start dropping a view and return the job deleting its definition. + + The name is free before this returns. Await :meth:`AsyncJob.wait` before + assuming the definition dataset is gone. + """ + if namespace_path is None: + namespace_path = [] + return AsyncJob( + await self._inner.drop_view_async(name, namespace_path=namespace_path) + ) + async def list_views( self, *, namespace_path: Optional[List[str]] = None ) -> List[str]: diff --git a/python/python/lancedb/remote/db.py b/python/python/lancedb/remote/db.py index 962ef66e1..cac9a0023 100644 --- a/python/python/lancedb/remote/db.py +++ b/python/python/lancedb/remote/db.py @@ -931,6 +931,13 @@ class RemoteDBConnection(DBConnection): ) -> None: LOOP.run(self._conn.drop_view(name, namespace_path=namespace_path)) + @override + def drop_view_async( + self, name: str, *, namespace_path: Optional[List[str]] = None + ) -> Job[None]: + job = LOOP.run(self._conn.drop_view_async(name, namespace_path=namespace_path)) + return Job(job) + @override def list_views(self, *, namespace_path: Optional[List[str]] = None) -> List[str]: return LOOP.run(self._conn.list_views(namespace_path=namespace_path)) diff --git a/python/src/connection.rs b/python/src/connection.rs index 69f54c683..e5ad4e6df 100644 --- a/python/src/connection.rs +++ b/python/src/connection.rs @@ -936,6 +936,23 @@ impl Connection { }) } + #[pyo3(signature = (name, namespace_path=None))] + pub fn drop_view_async( + self_: PyRef<'_, Self>, + name: String, + namespace_path: Option>, + ) -> PyResult> { + let inner = self_.get_inner()?.clone(); + let namespace_path = namespace_path.unwrap_or_default(); + future_into_py(self_.py(), async move { + inner + .drop_view_async(name, &namespace_path) + .await + .infer_error() + .map(crate::job::Job::new) + }) + } + #[pyo3(signature = (namespace_path=None))] pub fn list_views( self_: PyRef<'_, Self>, diff --git a/python/src/runtime.rs b/python/src/runtime.rs index 33ffd8bf9..8891e2bf2 100644 --- a/python/src/runtime.rs +++ b/python/src/runtime.rs @@ -208,13 +208,12 @@ where let _guard = guard; fut.await; }; + // Detaching is the point: dropping the `JoinHandle` leaves the task running, and the + // `OutstandingGuard` it carries is what keeps it visible to `shutdown`. Written as `drop` + // rather than `let _ =`, which reads as discarding an unpolled future. match runtime::Handle::try_current() { - Ok(handle) => { - let _ = handle.spawn(task); - } - Err(_) => { - let _ = get_runtime().spawn(task); - } + Ok(handle) => drop(handle.spawn(task)), + Err(_) => drop(get_runtime().spawn(task)), } } diff --git a/rust/lancedb/src/connection.rs b/rust/lancedb/src/connection.rs index 962e9e586..0643fa309 100644 --- a/rust/lancedb/src/connection.rs +++ b/rust/lancedb/src/connection.rs @@ -820,15 +820,42 @@ impl Connection { .await } - /// Drop a view. + /// Drop a view and wait for its definition to be deleted. /// /// The tables it reads are untouched: a view holds no rows of its own. - /// Local databases return [`Error::NotSupported`]. + /// Use [`Connection::drop_view_async`] to get the cleanup job instead of + /// waiting on it. Local databases return [`Error::NotSupported`]. pub async fn drop_view(&self, name: impl AsRef, namespace_path: &[String]) -> Result<()> { validate_view_reference(name.as_ref(), namespace_path)?; self.internal.drop_view(name.as_ref(), namespace_path).await } + /// Start dropping a view and return the job deleting its definition. + /// + /// The name is free before this returns; the definition dataset may still + /// be being deleted. Await [`Job::wait`][crate::job::Job::wait] to wait for + /// that. When nothing was bound to the name, the returned job is already + /// finished and has no id. Local databases return [`Error::NotSupported`]. + /// + /// ```no_run + /// # use lancedb::Connection; + /// # async fn drop(conn: &Connection) -> lancedb::Result<()> { + /// let job = conn.drop_view_async("recent_orders", &[]).await?; + /// job.wait().await?; + /// # Ok(()) + /// # } + /// ``` + pub async fn drop_view_async( + &self, + name: impl AsRef, + namespace_path: &[String], + ) -> Result { + validate_view_reference(name.as_ref(), namespace_path)?; + self.internal + .drop_view_async(name.as_ref(), namespace_path) + .await + } + /// The names of the views in one namespace. /// /// Names only; a definition is query metadata and comes from diff --git a/rust/lancedb/src/database.rs b/rust/lancedb/src/database.rs index 086d1df35..aaf6fd31f 100644 --- a/rust/lancedb/src/database.rs +++ b/rust/lancedb/src/database.rs @@ -470,11 +470,16 @@ pub trait Database: ) -> Result { view_ops_not_supported() } - /// Drop a view. Its sources are untouched -- a view holds no rows of its - /// own. + /// Drop a view and wait for its definition to be deleted. Its sources are + /// untouched -- a view holds no rows of its own. async fn drop_view(&self, _name: &str, _namespace_path: &[String]) -> Result<()> { view_ops_not_supported() } + /// Drop a view and return the job deleting its definition, without waiting. + #[doc(hidden)] + async fn drop_view_async(&self, _name: &str, _namespace_path: &[String]) -> Result { + view_ops_not_supported() + } /// The names of the views in one namespace. async fn list_views(&self, _namespace_path: &[String]) -> Result> { view_ops_not_supported() diff --git a/rust/lancedb/src/remote/db.rs b/rust/lancedb/src/remote/db.rs index ad4a1781b..93c7dc82e 100644 --- a/rust/lancedb/src/remote/db.rs +++ b/rust/lancedb/src/remote/db.rs @@ -1210,11 +1210,39 @@ impl Database for RemoteDatabase { } async fn drop_view(&self, name: &str, namespace_path: &[String]) -> Result<()> { + self.drop_view_async(name, namespace_path) + .await? + .wait() + .await + } + + async fn drop_view_async(&self, name: &str, namespace_path: &[String]) -> Result { let view_id = build_object_identifier("View name", name, namespace_path)?; let req = self.client.post(&format!("/v1/view/{view_id}/drop")); let (request_id, response) = self.client.send(req).await?; - self.client.check_response(&request_id, response).await?; - Ok(()) + let response = self.client.check_response(&request_id, response).await?; + let status = response.status(); + let body = response.text().await.err_to_http(request_id.clone())?; + match status { + // Nothing was bound to the name, so nothing is being deleted. + StatusCode::OK => Ok(Job::new_done()), + StatusCode::ACCEPTED => { + let job_id = extract_job_id(&body).ok_or_else(|| Error::Http { + source: "view drop response did not contain a valid job_id".into(), + request_id, + status_code: Some(status), + })?; + Ok(Job::new(Box::new(RemoteJob::new( + self.client.clone(), + job_id, + )))) + } + _ => Err(Error::Http { + source: "view drop must return 200 OK or 202 Accepted".into(), + request_id, + status_code: Some(status), + }), + } } async fn list_views(&self, namespace_path: &[String]) -> Result> { @@ -4046,6 +4074,55 @@ mod tests { .unwrap(); } + /// An accepted drop hands back the job so a caller can wait on the delete, + /// and the waiting `drop_view` does that for them. + #[tokio::test] + async fn test_drop_view_async_reports_the_cleanup_job() { + let db = super::RemoteDatabase::new_mock(|_| { + http::Response::builder() + .status(202) + .body(r#"{"job_id":"j1-do-abc"}"#) + .unwrap() + }); + let job = db.drop_view_async("adults", &[]).await.unwrap(); + assert_eq!(job.id(), Some("j1-do-abc")); + } + + /// Nothing was bound, so nothing is being deleted and the job is already done. + #[tokio::test] + async fn test_drop_view_async_reports_a_finished_job_when_nothing_was_bound() { + let db = super::RemoteDatabase::new_mock(|_| { + http::Response::builder().status(200).body("{}").unwrap() + }); + let job = db.drop_view_async("adults", &[]).await.unwrap(); + assert_eq!(job.id(), None); + assert_eq!(job.status().await.unwrap(), "finished"); + job.wait().await.unwrap(); + } + + #[tokio::test] + async fn test_drop_view_rejects_incomplete_acceptance() { + for body in ["{}", r#"{"job_id":""}"#, r#"{"job_id":null}"#] { + let db = super::RemoteDatabase::new_mock(move |_| { + http::Response::builder().status(202).body(body).unwrap() + }); + let error = db.drop_view_async("adults", &[]).await.err().unwrap(); + assert!(error.to_string().contains("valid job_id"), "{error}"); + } + } + + #[tokio::test] + async fn test_drop_view_rejects_unexpected_success_status() { + let db = super::RemoteDatabase::new_mock(|_| { + http::Response::builder().status(204).body("").unwrap() + }); + let error = db.drop_view_async("adults", &[]).await.err().unwrap(); + assert!( + error.to_string().contains("200 OK or 202 Accepted"), + "{error}" + ); + } + /// A schema the client cannot decode is a broken response, not a view /// with no columns: reporting it as an error keeps a caller from reading /// an empty schema as the truth about the view.