diff --git a/docs/src/js/classes/Connection.md b/docs/src/js/classes/Connection.md index d46d414d3..75af5102e 100644 --- a/docs/src/js/classes/Connection.md +++ b/docs/src/js/classes/Connection.md @@ -319,6 +319,38 @@ Creates a new Table and initialize it with new data. *** +### createView() + +```ts +abstract createView( + name, + query, + namespacePath?): Promise +``` + +Create a view: a named query the database plans on every read. + +The query is planned once, at creation, so one that cannot be planned is +rejected now rather than at the first read. A view holds no rows, and its +readers see its sources as they are at read time. + +There is no replace: a name already taken is an error, and changing a +view is a drop followed by a create. + +#### Parameters + +* **name**: `string` + +* **query**: `string` + +* **namespacePath?**: `string`[] + +#### Returns + +`Promise`<[`ViewDescription`](../interfaces/ViewDescription.md)> + +*** + ### describeNamespace() ```ts @@ -342,6 +374,27 @@ The namespace's properties *** +### describeView() + +```ts +abstract describeView(name, namespacePath?): Promise +``` + +What this database records about the view named `name`: its defining +query and the schema that query resolved to. + +#### Parameters + +* **name**: `string` + +* **namespacePath?**: `string`[] + +#### Returns + +`Promise`<[`ViewDescription`](../interfaces/ViewDescription.md)> + +*** + ### display() ```ts @@ -498,6 +551,28 @@ on the returned job to know when cleanup has finished. *** +### dropView() + +```ts +abstract dropView(name, namespacePath?): Promise +``` + +Drop the view named `name`. + +The tables it reads are untouched: a view holds no rows of its own. + +#### Parameters + +* **name**: `string` + +* **namespacePath?**: `string`[] + +#### Returns + +`Promise`<`void`> + +*** + ### isOpen() ```ts @@ -636,6 +711,26 @@ A page of table names and an *** +### listViews() + +```ts +abstract listViews(namespacePath?): Promise +``` + +The names of the views in one namespace. + +Names only; a definition comes from [describeView](Connection.md#describeview). + +#### Parameters + +* **namespacePath?**: `string`[] + +#### Returns + +`Promise`<`string`[]> + +*** + ### openJob() ```ts diff --git a/docs/src/js/globals.md b/docs/src/js/globals.md index 6d242a15b..2f49578b1 100644 --- a/docs/src/js/globals.md +++ b/docs/src/js/globals.md @@ -149,6 +149,7 @@ - [UpdateOptions](interfaces/UpdateOptions.md) - [UpdateResult](interfaces/UpdateResult.md) - [Version](interfaces/Version.md) +- [ViewDescription](interfaces/ViewDescription.md) - [WriteExecutionOptions](interfaces/WriteExecutionOptions.md) - [WriteProgress](interfaces/WriteProgress.md) diff --git a/docs/src/js/interfaces/ViewDescription.md b/docs/src/js/interfaces/ViewDescription.md new file mode 100644 index 000000000..d53c884b5 --- /dev/null +++ b/docs/src/js/interfaces/ViewDescription.md @@ -0,0 +1,78 @@ +[**@lancedb/lancedb**](../README.md) • **Docs** + +*** + +[@lancedb/lancedb](../globals.md) / ViewDescription + +# Interface: ViewDescription + +What a database records about one view. + +A view holds no rows: it stores the statement that defines it and the +schema that statement resolved to, so a reader sees its sources as they +are at read time. + +## Properties + +### defaultDatabase + +```ts +defaultDatabase: string; +``` + +The database that unqualified table names in [query](ViewDescription.md#query) resolve against. + +*** + +### defaultNamespacePath + +```ts +defaultNamespacePath: string[]; +``` + +The namespace path those unqualified names resolve against; empty is the +root namespace. + +Recorded with the view because it outlives the session that declared it: +a reader resolving the query against its own default namespace could read +a different table than the view was defined over. + +*** + +### name + +```ts +name: string; +``` + +The view's name within its namespace. + +*** + +### namespacePath + +```ts +namespacePath: string[]; +``` + +The namespace holding the view; empty is the root namespace. + +*** + +### query + +```ts +query: string; +``` + +The defining query, as the database stores it. + +*** + +### schema + +```ts +schema: Schema; +``` + +The schema the defining query resolved to when the view was created. diff --git a/docs/src/python/python.md b/docs/src/python/python.md index d1ba66280..68cd306c4 100644 --- a/docs/src/python/python.md +++ b/docs/src/python/python.md @@ -198,6 +198,10 @@ listing a storage directory. ::: lancedb.materialized_view.MaterializedViewDefinition +## Views + +::: lancedb.view.ViewDescription + ## Expressions Type-safe expression builder for filters and projections. Use these instead diff --git a/nodejs/__test__/remote.test.ts b/nodejs/__test__/remote.test.ts index 40f4c67db..4c1d21604 100644 --- a/nodejs/__test__/remote.test.ts +++ b/nodejs/__test__/remote.test.ts @@ -100,6 +100,67 @@ describe("remote connection", () => { ); }); + it("creates a view and decodes the schema it resolved to", async () => { + await withMockDatabase( + (req, res) => { + expect(req.method).toBe("POST"); + expect(req.url).toBe("/v1/view/analytics$adults/create"); + res.writeHead(200, { "content-type": "application/json" }).end( + JSON.stringify({ + name: "adults", + namespace: ["analytics"], + query: "SELECT name FROM people", + // biome-ignore lint/style/useNamingConvention: the wire field is snake_case + default_database: "db", + // biome-ignore lint/style/useNamingConvention: the wire field is snake_case + default_namespace: ["analytics"], + schema: { + fields: [ + { name: "name", nullable: true, type: { type: "utf8" } }, + ], + }, + }), + ); + }, + async (db) => { + const view = await db.createView("adults", "SELECT name FROM people", [ + "analytics", + ]); + expect(view.name).toBe("adults"); + expect(view.namespacePath).toEqual(["analytics"]); + expect(view.query).toBe("SELECT name FROM people"); + expect(view.defaultDatabase).toBe("db"); + expect(view.defaultNamespacePath).toEqual(["analytics"]); + expect(view.schema.fields.map((f) => f.name)).toEqual(["name"]); + }, + ); + }); + + it("lists and drops views through their own routes", async () => { + await withMockDatabase( + (req, res) => { + expect(req.url).toBe("/v1/namespace/$/view/list"); + res + .writeHead(200, { "content-type": "application/json" }) + .end(JSON.stringify({ views: ["adults"] })); + }, + async (db) => { + expect(await db.listViews()).toEqual(["adults"]); + }, + ); + + await withMockDatabase( + (req, res) => { + expect(req.method).toBe("POST"); + expect(req.url).toBe("/v1/view/adults/drop"); + res.writeHead(200, { "content-type": "application/json" }).end("{}"); + }, + async (db) => { + await db.dropView("adults"); + }, + ); + }); + 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 f5a6679cf..9637d60c2 100644 --- a/nodejs/lancedb/connection.ts +++ b/nodejs/lancedb/connection.ts @@ -31,6 +31,7 @@ import type { ListNamespacesResponse, ListTablesResponse, } from "./native"; +import { ViewDescription, viewDescriptionFromNative } from "./view"; export type { CreateNamespaceResponse, DescribeNamespaceResponse, @@ -377,6 +378,45 @@ export abstract class Connection { namespacePath?: string[], ): Promise; + /** + * Create a view: a named query the database plans on every read. + * + * The query is planned once, at creation, so one that cannot be planned is + * rejected now rather than at the first read. A view holds no rows, and its + * readers see its sources as they are at read time. + * + * There is no replace: a name already taken is an error, and changing a + * view is a drop followed by a create. + */ + abstract createView( + name: string, + query: string, + namespacePath?: string[], + ): Promise; + + /** + * What this database records about the view named `name`: its defining + * query and the schema that query resolved to. + */ + abstract describeView( + name: string, + namespacePath?: string[], + ): Promise; + + /** + * Drop the view named `name`. + * + * The tables it reads are untouched: a view holds no rows of its own. + */ + abstract dropView(name: string, namespacePath?: string[]): Promise; + + /** + * The names of the views in one namespace. + * + * Names only; a definition comes from {@link describeView}. + */ + abstract listViews(namespacePath?: string[]): Promise; + abstract openTable( name: string, namespacePath?: string[], @@ -696,6 +736,33 @@ export class LocalConnection extends Connection { ); } + async createView( + name: string, + query: string, + namespacePath?: string[], + ): Promise { + return viewDescriptionFromNative( + await this.inner.createView(name, query, namespacePath ?? []), + ); + } + + async describeView( + name: string, + namespacePath?: string[], + ): Promise { + return viewDescriptionFromNative( + await this.inner.describeView(name, namespacePath ?? []), + ); + } + + async dropView(name: string, namespacePath?: string[]): Promise { + return this.inner.dropView(name, namespacePath ?? []); + } + + async listViews(namespacePath?: string[]): Promise { + return this.inner.listViews(namespacePath ?? []); + } + async listTables( namespacePathOrOptions?: string[] | Partial, options?: Partial, diff --git a/nodejs/lancedb/index.ts b/nodejs/lancedb/index.ts index c0a4dfcbd..ed6889232 100644 --- a/nodejs/lancedb/index.ts +++ b/nodejs/lancedb/index.ts @@ -26,6 +26,7 @@ export { MaterializedViewDefinition, MaterializedViewSelect, } from "./materialized_view"; +export { ViewDescription } from "./view"; export { JsHeaderProvider as NativeJsHeaderProvider } from "./native.js"; // OpenTelemetry metrics bridge. Only the high-level entry point is public; the diff --git a/nodejs/lancedb/view.ts b/nodejs/lancedb/view.ts new file mode 100644 index 000000000..5ed848648 --- /dev/null +++ b/nodejs/lancedb/view.ts @@ -0,0 +1,48 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The LanceDB Authors + +import { Schema, tableFromIPC } from "apache-arrow"; +import { ViewDescription as NativeViewDescription } from "./native"; + +/** + * What a database records about one view. + * + * A view holds no rows: it stores the statement that defines it and the + * schema that statement resolved to, so a reader sees its sources as they + * are at read time. + */ +export interface ViewDescription { + /** The view's name within its namespace. */ + name: string; + /** The namespace holding the view; empty is the root namespace. */ + namespacePath: string[]; + /** The defining query, as the database stores it. */ + query: string; + /** The database that unqualified table names in {@link query} resolve against. */ + defaultDatabase: string; + /** + * The namespace path those unqualified names resolve against; empty is the + * root namespace. + * + * Recorded with the view because it outlives the session that declared it: + * a reader resolving the query against its own default namespace could read + * a different table than the view was defined over. + */ + defaultNamespacePath: string[]; + /** The schema the defining query resolved to when the view was created. */ + schema: Schema; +} + +/** Decode the schema the binding hands over as an Arrow IPC file. */ +export function viewDescriptionFromNative( + view: NativeViewDescription, +): ViewDescription { + return { + name: view.name, + namespacePath: view.namespacePath, + query: view.query, + defaultDatabase: view.defaultDatabase, + defaultNamespacePath: view.defaultNamespacePath, + schema: tableFromIPC(view.schema).schema, + }; +} diff --git a/nodejs/src/connection.rs b/nodejs/src/connection.rs index 0c07ffd50..ee2a9470b 100644 --- a/nodejs/src/connection.rs +++ b/nodejs/src/connection.rs @@ -55,6 +55,32 @@ pub struct DropNamespaceResponse { pub transaction_id: Option>, } +/// What a database records about one view. +#[napi(object)] +pub struct ViewDescription { + pub name: String, + pub namespace_path: Vec, + pub query: String, + pub default_database: String, + pub default_namespace_path: Vec, + /// The view's schema as an empty Arrow IPC file, the way a table reports + /// its own. + pub schema: Buffer, +} + +impl ViewDescription { + fn from_inner(view: lancedb::view::ViewDescription) -> napi::Result { + Ok(Self { + name: view.name, + namespace_path: view.namespace_path, + query: view.query, + default_database: view.default_database, + default_namespace_path: view.default_namespace_path, + schema: crate::util::schema_to_buffer(&view.schema)?, + }) + } +} + impl Connection { pub(crate) fn inner_new(inner: LanceDBConnection) -> Self { Self { inner: Some(inner) } @@ -368,6 +394,63 @@ impl Connection { .default_error() } + /// Create a view: a named query planned on every read. + #[napi(catch_unwind)] + pub async fn create_view( + &self, + name: String, + query: String, + namespace_path: Option>, + ) -> napi::Result { + let ns = namespace_path.unwrap_or_default(); + let view = self + .get_inner()? + .create_view(&name, &query, &ns) + .await + .default_error()?; + ViewDescription::from_inner(view) + } + + /// What the database records about one view. + #[napi(catch_unwind)] + pub async fn describe_view( + &self, + name: String, + namespace_path: Option>, + ) -> napi::Result { + let ns = namespace_path.unwrap_or_default(); + let view = self + .get_inner()? + .describe_view(&name, &ns) + .await + .default_error()?; + ViewDescription::from_inner(view) + } + + /// Drop a view. The tables it reads are untouched. + #[napi(catch_unwind)] + pub async fn drop_view( + &self, + name: String, + namespace_path: Option>, + ) -> napi::Result<()> { + let ns = namespace_path.unwrap_or_default(); + self.get_inner()? + .drop_view(&name, &ns) + .await + .default_error() + } + + /// The names of the views in one namespace. + #[napi(catch_unwind)] + pub async fn list_views( + &self, + namespace_path: Option>, + ) -> napi::Result> { + let ns = namespace_path.unwrap_or_default(); + self.get_inner()?.list_views(&ns).await.default_error() + } + /// Start dropping a materialized view and return its cleanup job. #[napi(catch_unwind)] pub async fn drop_materialized_view_async( diff --git a/python/python/lancedb/__init__.py b/python/python/lancedb/__init__.py index 362e0cdbf..d34b541a8 100644 --- a/python/python/lancedb/__init__.py +++ b/python/python/lancedb/__init__.py @@ -47,6 +47,7 @@ from .materialized_view import ( MaterializedView, MaterializedViewDefinition, ) +from .view import ViewDescription as ViewDescription from .table import AsyncTable, Table from .types import BaseTokenizerType from ._lancedb import Session @@ -582,6 +583,7 @@ __all__ = [ "AsyncMaterializedView", "MaterializedView", "MaterializedViewDefinition", + "ViewDescription", "connect", "connect_async", "tokenize", diff --git a/python/python/lancedb/_lancedb.pyi b/python/python/lancedb/_lancedb.pyi index 6d873850d..ba4e3a0df 100644 --- a/python/python/lancedb/_lancedb.pyi +++ b/python/python/lancedb/_lancedb.pyi @@ -171,6 +171,18 @@ class Connection(object): async def describe_secret( self, name: str, namespace_path: Optional[List[str]] = None ) -> Tuple[str, int, int]: ... + async def create_view( + self, name: str, query: str, namespace_path: Optional[List[str]] = None + ) -> Tuple[str, List[str], str, str, List[str], pa.Schema]: ... + async def describe_view( + self, name: str, namespace_path: Optional[List[str]] = None + ) -> Tuple[str, List[str], str, str, List[str], pa.Schema]: ... + async def drop_view( + self, name: str, namespace_path: Optional[List[str]] = None + ) -> None: ... + async def list_views( + self, namespace_path: Optional[List[str]] = None + ) -> List[str]: ... async def list_jobs(self) -> List[JobInfo]: ... async def cancel_job(self, job_id: str) -> bool: ... async def execute_query_async( diff --git a/python/python/lancedb/db.py b/python/python/lancedb/db.py index b0db22532..10d52c63e 100644 --- a/python/python/lancedb/db.py +++ b/python/python/lancedb/db.py @@ -75,6 +75,7 @@ from .util import ( get_uri_scheme, validate_table_name, ) +from .view import ViewDescription import deprecation @@ -96,6 +97,28 @@ from .namespace_utils import ( ) +def _view_description( + described: Tuple[str, List[str], str, str, List[str], "pa.Schema"], +) -> ViewDescription: + """Name the fields the binding returns positionally.""" + ( + name, + namespace_path, + query, + default_database, + default_namespace_path, + schema, + ) = described + return ViewDescription( + name=name, + query=query, + default_database=default_database, + schema=schema, + namespace_path=namespace_path, + default_namespace_path=default_namespace_path, + ) + + class DBConnection(EnforceOverrides): """An active LanceDB connection interface.""" @@ -918,6 +941,68 @@ class DBConnection(EnforceOverrides): "Secret operations are not supported for this connection type" ) + def create_view( + self, name: str, query: str, *, namespace_path: Optional[List[str]] = None + ) -> ViewDescription: + """Create a view: a named query the database plans on every read. + + The query is planned once, at creation, so one that cannot be planned + is refused now rather than at the first read. A view holds no rows, and + its readers see its sources as they are at read time. + + There is no replace: a name already taken is an error, and changing a + view is a drop followed by a create. Local connections raise + ``NotImplementedError``. + + >>> import lancedb + >>> db = lancedb.connect("db://my_database") # doctest: +SKIP + >>> view = db.create_view( + ... "adults", "SELECT name FROM people WHERE age >= 18" + ... ) # doctest: +SKIP + >>> view.schema # doctest: +SKIP + name: string + """ + raise NotImplementedError( + "View operations are not supported for this connection type" + ) + + def describe_view( + self, name: str, *, namespace_path: Optional[List[str]] = None + ) -> ViewDescription: + """What this database records about a view: its defining query and the + schema that query resolved to. + + The schema is the one recorded at creation; a source altered since then + shows up when the view is read. Local connections raise + ``NotImplementedError``. + """ + raise NotImplementedError( + "View operations are not supported for this connection type" + ) + + def drop_view( + self, name: str, *, namespace_path: Optional[List[str]] = None + ) -> None: + """Drop a view. + + The tables it reads are untouched: a view holds no rows of its own. + 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. + + Names only; a definition comes from + [describe_view][lancedb.db.DBConnection.describe_view]. Local + connections raise ``NotImplementedError``. + """ + raise NotImplementedError( + "View operations are not supported for this connection type" + ) + def open_job(self, job_id: str) -> Job: """Open a server-side job by id, returning a handle with its record already populated. @@ -1731,6 +1816,30 @@ class LanceDBConnection(DBConnection): ) -> SecretInfo: return LOOP.run(self._conn.describe_secret(name, namespace_path=namespace_path)) + @override + def create_view( + self, name: str, query: str, *, namespace_path: Optional[List[str]] = None + ) -> ViewDescription: + return LOOP.run( + self._conn.create_view(name, query, namespace_path=namespace_path) + ) + + @override + def describe_view( + self, name: str, *, namespace_path: Optional[List[str]] = None + ) -> ViewDescription: + return LOOP.run(self._conn.describe_view(name, namespace_path=namespace_path)) + + @override + def drop_view( + self, name: str, *, namespace_path: Optional[List[str]] = None + ) -> None: + LOOP.run(self._conn.drop_view(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)) + @override def list_jobs(self) -> List[JobInfo]: """List server-side jobs across the database's tables.""" @@ -2689,6 +2798,38 @@ class AsyncConnection(object): updated_at_millis=updated_at_millis, ) + async def create_view( + self, name: str, query: str, *, namespace_path: Optional[List[str]] = None + ) -> ViewDescription: + """Create a view: a named query the database plans on every read. + + See + [DBConnection.create_view][lancedb.DBConnection.create_view]. + """ + return _view_description( + await self._inner.create_view(name, query, list(namespace_path or [])) + ) + + async def describe_view( + self, name: str, *, namespace_path: Optional[List[str]] = None + ) -> ViewDescription: + """What this database records about a view: query and schema.""" + return _view_description( + await self._inner.describe_view(name, list(namespace_path or [])) + ) + + async def drop_view( + self, name: str, *, namespace_path: Optional[List[str]] = None + ) -> None: + """Drop a view. The tables it reads are untouched.""" + await self._inner.drop_view(name, list(namespace_path or [])) + + async def list_views( + self, *, namespace_path: Optional[List[str]] = None + ) -> List[str]: + """The names of the views in one namespace.""" + return await self._inner.list_views(list(namespace_path or [])) + async def list_jobs(self) -> List[JobInfo]: """List server-side jobs across the database's tables.""" return await self._inner.list_jobs() diff --git a/python/python/lancedb/remote/db.py b/python/python/lancedb/remote/db.py index ef40ef007..962ef66e1 100644 --- a/python/python/lancedb/remote/db.py +++ b/python/python/lancedb/remote/db.py @@ -41,6 +41,7 @@ from ..sql import Query as SqlQuery from ..sql import QueryDescription from ..materialized_view import MaterializedView, SelectArg from ..secrets import EnvVarSecret, SecretInfo +from ..view import ViewDescription if TYPE_CHECKING: from .._lancedb import JobInfo @@ -910,6 +911,30 @@ class RemoteDBConnection(DBConnection): ) -> None: LOOP.run(self._conn.drop_secret(name, namespace_path=namespace_path)) + @override + def create_view( + self, name: str, query: str, *, namespace_path: Optional[List[str]] = None + ) -> ViewDescription: + return LOOP.run( + self._conn.create_view(name, query, namespace_path=namespace_path) + ) + + @override + def describe_view( + self, name: str, *, namespace_path: Optional[List[str]] = None + ) -> ViewDescription: + return LOOP.run(self._conn.describe_view(name, namespace_path=namespace_path)) + + @override + def drop_view( + self, name: str, *, namespace_path: Optional[List[str]] = None + ) -> None: + LOOP.run(self._conn.drop_view(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)) + @override def list_jobs(self) -> List["JobInfo"]: """List server-side jobs across the database's tables.""" diff --git a/python/python/lancedb/view.py b/python/python/lancedb/view.py new file mode 100644 index 000000000..715252fc6 --- /dev/null +++ b/python/python/lancedb/view.py @@ -0,0 +1,40 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright The LanceDB Authors + +"""Views: a named query a database stores and plans on every read. + +A view holds no rows. What it stores is the statement that defines it and the +schema that statement resolved to, so a reader sees the sources as they are +now. See ``DBConnection.create_view``. +""" + +from __future__ import annotations + +from dataclasses import dataclass, field +from typing import TYPE_CHECKING, List + +if TYPE_CHECKING: + import pyarrow as pa + + +@dataclass(frozen=True) +class ViewDescription: + """What a database records about one view.""" + + name: str + """The view's name within its namespace.""" + query: str + """The defining query, as the database stores it.""" + default_database: str + """The database that unqualified table names in ``query`` resolve against.""" + schema: "pa.Schema" + """The schema the defining query resolved to when the view was created.""" + namespace_path: List[str] = field(default_factory=list) + """The namespace holding the view; empty is the root namespace.""" + default_namespace_path: List[str] = field(default_factory=list) + """The namespace path those unqualified names resolve against. + + Recorded with the view because it outlives the session that declared it: a + reader resolving the query against its own default namespace could read a + different table than the view was defined over. + """ diff --git a/python/python/tests/test_remote_db.py b/python/python/tests/test_remote_db.py index 1c887b218..c03862a49 100644 --- a/python/python/tests/test_remote_db.py +++ b/python/python/tests/test_remote_db.py @@ -2747,3 +2747,57 @@ def test_remote_job_handle_reports_its_own_detail(): "limit": 500, "filter": "state = 'claim_complete'", } + + +def test_view_crud_addresses_its_own_routes(): + # The view verbs are their own routes, and the schema comes back in the + # namespace spec's JSON encoding, decoded into a pyarrow schema. + paths = [] + + def handler(request): + paths.append((request.command, request.path)) + if request.path.endswith("/view/list"): + body = {"views": ["adults"]} + elif request.path.endswith("/drop"): + body = {} + else: + body = { + "name": "adults", + "namespace": ["analytics"], + "query": "SELECT name FROM people", + "default_database": "dev", + "default_namespace": ["analytics"], + "schema": { + "fields": [ + {"name": "name", "nullable": True, "type": {"type": "utf8"}} + ] + }, + } + request.send_response(200) + request.send_header("Content-Type", "application/json") + request.end_headers() + request.wfile.write(json.dumps(body).encode()) + + with mock_lancedb_connection(handler) as db: + view = db.create_view( + "adults", "SELECT name FROM people", namespace_path=["analytics"] + ) + assert view.name == "adults" + assert view.namespace_path == ["analytics"] + assert view.query == "SELECT name FROM people" + assert view.default_database == "dev" + assert view.default_namespace_path == ["analytics"] + assert view.schema == pa.schema([pa.field("name", pa.utf8(), nullable=True)]) + + described = db.describe_view("adults", namespace_path=["analytics"]) + assert described.schema == view.schema + + assert db.list_views(namespace_path=["analytics"]) == ["adults"] + db.drop_view("adults", namespace_path=["analytics"]) + + assert paths == [ + ("POST", "/v1/view/analytics$adults/create"), + ("POST", "/v1/view/analytics$adults/describe"), + ("GET", "/v1/namespace/analytics/view/list"), + ("POST", "/v1/view/analytics$adults/drop"), + ] diff --git a/python/src/connection.rs b/python/src/connection.rs index 3385e6b2d..69f54c683 100644 --- a/python/src/connection.rs +++ b/python/src/connection.rs @@ -13,7 +13,11 @@ use crate::{ runtime::future_into_py, table::Table, }; -use arrow::{datatypes::Schema, ffi_stream::ArrowArrayStreamReader, pyarrow::FromPyArrow}; +use arrow::{ + datatypes::Schema, + ffi_stream::ArrowArrayStreamReader, + pyarrow::{FromPyArrow, ToPyArrow}, +}; use lancedb::{ connection::Connection as LanceConnection, connection::NamespaceClientPushdownOperation, @@ -100,6 +104,28 @@ fn parse_default_namespace_path(path: Option>) -> PyResult, String, String, Vec, Py); + +/// A view description as a plain tuple, with the schema converted to the +/// pyarrow schema the caller would get from any other lancedb API. The Python +/// layer names the fields; this keeps the binding free of a class that would +/// have to be kept in step with the Rust struct. +fn view_description_to_py(view: lancedb::view::ViewDescription) -> PyResult { + Python::attach(|py| { + let schema = view.schema.to_pyarrow(py)?.unbind(); + Ok(( + view.name, + view.namespace_path, + view.query, + view.default_database, + view.default_namespace_path, + schema, + )) + }) +} + #[pymethods] impl Connection { fn __repr__(&self) -> String { @@ -862,6 +888,66 @@ impl Connection { }) } + #[pyo3(signature = (name, query, namespace_path=None))] + pub fn create_view( + self_: PyRef<'_, Self>, + name: String, + query: 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 { + let view = inner + .create_view(name, query, &namespace_path) + .await + .infer_error()?; + view_description_to_py(view) + }) + } + + #[pyo3(signature = (name, namespace_path=None))] + pub fn describe_view( + 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 { + let view = inner + .describe_view(name, &namespace_path) + .await + .infer_error()?; + view_description_to_py(view) + }) + } + + #[pyo3(signature = (name, namespace_path=None))] + pub fn drop_view( + 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(name, &namespace_path).await.infer_error() + }) + } + + #[pyo3(signature = (namespace_path=None))] + pub fn list_views( + self_: PyRef<'_, Self>, + 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.list_views(&namespace_path).await.infer_error() + }) + } + pub fn list_jobs(self_: PyRef<'_, Self>) -> PyResult> { let inner = self_.get_inner()?.clone(); future_into_py(self_.py(), async move { diff --git a/rust/lancedb/src/connection.rs b/rust/lancedb/src/connection.rs index d43f8f744..962e9e586 100644 --- a/rust/lancedb/src/connection.rs +++ b/rust/lancedb/src/connection.rs @@ -37,7 +37,11 @@ use crate::remote::{ }, }; use crate::secrets::SecretInfo; -use crate::utils::{validate_secret_component, validate_secret_reference}; +use crate::utils::{ + validate_namespace, validate_secret_component, validate_secret_reference, + validate_view_reference, +}; +use crate::view::ViewDescription; use lance::io::ObjectStoreParams; pub use lance_file::version::LanceFileVersion; #[cfg(feature = "remote")] @@ -747,6 +751,94 @@ impl Connection { .await } + /// Create a view: a named query the database plans on every read. + /// + /// The query is planned once, here, so one that cannot be planned is + /// refused now rather than at the first read. A view holds no rows, and + /// its readers see its sources as they are at read time. + /// + /// There is no replace: a name already taken is an error, and changing a + /// view is a drop followed by a create, each authorized against what it + /// actually touches. Local databases return [`Error::NotSupported`]. + /// + /// # Example + /// + /// ```no_run + /// # async fn view_lifecycle( + /// # connection: &lancedb::Connection, + /// # ) -> Result<(), Box> { + /// let namespace = vec!["analytics".to_string()]; + /// + /// let view = connection + /// .create_view( + /// "recent_orders", + /// "SELECT id, total FROM orders WHERE total > 100", + /// &namespace, + /// ) + /// .await?; + /// println!("{} has {} columns", view.name, view.schema.fields().len()); + /// + /// // The query comes back as it was recorded, with the defaults its + /// // unqualified names resolve against. + /// let described = connection.describe_view("recent_orders", &namespace).await?; + /// println!("{} in {:?}", described.query, described.default_namespace_path); + /// + /// let names = connection.list_views(&namespace).await?; + /// assert!(names.iter().any(|name| name == "recent_orders")); + /// + /// // Dropping the view leaves `orders` untouched. + /// connection.drop_view("recent_orders", &namespace).await?; + /// # Ok(()) + /// # } + /// ``` + pub async fn create_view( + &self, + name: impl AsRef, + query: impl AsRef, + namespace_path: &[String], + ) -> Result { + validate_view_reference(name.as_ref(), namespace_path)?; + self.internal + .create_view(name.as_ref(), query.as_ref(), namespace_path) + .await + } + + /// What this database records about one view: its defining query and the + /// schema that query resolved to. + /// + /// The schema is the one recorded at creation. A source altered since then + /// shows up when the view is read, not here. Local databases return + /// [`Error::NotSupported`]. + pub async fn describe_view( + &self, + name: impl AsRef, + namespace_path: &[String], + ) -> Result { + validate_view_reference(name.as_ref(), namespace_path)?; + self.internal + .describe_view(name.as_ref(), namespace_path) + .await + } + + /// Drop a view. + /// + /// The tables it reads are untouched: a view holds no rows of its own. + /// 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 + } + + /// The names of the views in one namespace. + /// + /// Names only; a definition is query metadata and comes from + /// [`Self::describe_view`]. The client walks all server pages before + /// returning. Local databases return [`Error::NotSupported`]. + pub async fn list_views(&self, namespace_path: &[String]) -> Result> { + validate_namespace(namespace_path)?; + self.internal.list_views(namespace_path).await + } + /// Rename a table in the database. /// /// This is only supported in LanceDB Cloud. diff --git a/rust/lancedb/src/database.rs b/rust/lancedb/src/database.rs index 91b200043..086d1df35 100644 --- a/rust/lancedb/src/database.rs +++ b/rust/lancedb/src/database.rs @@ -32,6 +32,7 @@ use crate::job::Job; use crate::materialized_view::CreateMaterializedViewRequest; use crate::secrets::SecretInfo; use crate::table::{BaseTable, WriteOptions}; +use crate::view::ViewDescription; pub mod listing; pub mod namespace; @@ -258,6 +259,12 @@ fn secret_catalog_not_supported() -> Result { }) } +fn view_ops_not_supported() -> Result { + Err(crate::error::Error::NotSupported { + message: "View operations are not supported by this database".to_string(), + }) +} + /// The `Database` trait defines the interface for database implementations. /// /// A database is responsible for managing tables and their metadata. @@ -443,6 +450,35 @@ pub trait Database: async fn describe_secret(&self, _name: &str, _namespace_path: &[String]) -> Result { secret_catalog_not_supported() } + /// Create a view from a defining query. The query is planned once, by the + /// database, so one that cannot be planned is refused now rather than at + /// the first read. + async fn create_view( + &self, + _name: &str, + _query: &str, + _namespace_path: &[String], + ) -> Result { + view_ops_not_supported() + } + /// What the database records about one view: its defining query and the + /// schema that query resolved to. + async fn describe_view( + &self, + _name: &str, + _namespace_path: &[String], + ) -> Result { + view_ops_not_supported() + } + /// Drop a view. 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() + } + /// The names of the views in one namespace. + async fn list_views(&self, _namespace_path: &[String]) -> Result> { + view_ops_not_supported() + } /// Open a job by id, returning a handle with its record already /// populated. Fails with [`crate::Error::JobNotFound`] when the server has /// no such job. diff --git a/rust/lancedb/src/lib.rs b/rust/lancedb/src/lib.rs index 4d9afc7c7..530328cbb 100644 --- a/rust/lancedb/src/lib.rs +++ b/rust/lancedb/src/lib.rs @@ -202,6 +202,7 @@ pub mod table; #[cfg(test)] pub mod test_utils; pub mod utils; +pub mod view; use std::{fmt::Display, str::FromStr}; diff --git a/rust/lancedb/src/remote/db.rs b/rust/lancedb/src/remote/db.rs index 7a59db665..893c584e1 100644 --- a/rust/lancedb/src/remote/db.rs +++ b/rust/lancedb/src/remote/db.rs @@ -14,8 +14,8 @@ use reqwest::header::CONTENT_TYPE; use lance_namespace::models::{ CreateNamespaceRequest, CreateNamespaceResponse, DescribeNamespaceRequest, - DescribeNamespaceResponse, DropNamespaceRequest, DropNamespaceResponse, ListNamespacesRequest, - ListNamespacesResponse, ListTablesRequest, ListTablesResponse, + DescribeNamespaceResponse, DropNamespaceRequest, DropNamespaceResponse, JsonArrowSchema, + ListNamespacesRequest, ListNamespacesResponse, ListTablesRequest, ListTablesResponse, }; use crate::Error; @@ -36,6 +36,7 @@ use crate::secrets::SecretBinding; use crate::secrets::SecretInfo; use crate::table::BaseTable; use crate::utils::{reject_relative_segment, validate_table_name}; +use crate::view::ViewDescription; use super::client::{ ClientConfig, HeaderProvider, HttpSend, ID_DELIMITER, RequestResultExt, RestfulLanceDbClient, @@ -784,6 +785,59 @@ struct RemoteListedSecret { name: String, } +/// Define a view from a query. The name and its namespace are the path +/// identifier, so neither appears here. +#[derive(serde::Serialize)] +struct RemoteCreateViewRequest<'a> { + query: &'a str, +} + +/// What the service reports about one view. The schema arrives as the +/// namespace spec's JSON encoding, which is what `describe_table` uses too. +#[derive(serde::Deserialize)] +struct RemoteViewDescription { + name: String, + #[serde(default)] + namespace: Vec, + query: String, + default_database: String, + /// A path, like `namespace`: the root is the absent field rather than a + /// spelling of its own. + #[serde(default)] + default_namespace: Vec, + schema: JsonArrowSchema, +} + +impl RemoteViewDescription { + fn into_description(self, request_id: String) -> Result { + let schema = + lance_namespace::schema::convert_json_arrow_schema(&self.schema).map_err(|source| { + Error::Http { + source: format!("View '{}' has an undecodable schema: {source}", self.name) + .into(), + request_id, + status_code: None, + } + })?; + Ok(ViewDescription { + name: self.name, + namespace_path: self.namespace, + query: self.query, + default_database: self.default_database, + default_namespace_path: self.default_namespace, + schema: Arc::new(schema), + }) + } +} + +#[derive(serde::Deserialize)] +struct RemoteListViewsResponse { + #[serde(default)] + views: Vec, + #[serde(default)] + page_token: Option, +} + /// Bound on `list_jobs` page walking; a warning is logged when the listing /// is truncated at this many pages. const MAX_LIST_JOBS_PAGES: usize = 100; @@ -1122,6 +1176,79 @@ impl Database for RemoteDatabase { response.json().await.err_to_http(request_id) } + async fn create_view( + &self, + name: &str, + query: &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}/create")) + .json(&RemoteCreateViewRequest { query }); + let (request_id, response) = self.client.send(req).await?; + let response = self.client.check_response(&request_id, response).await?; + let description: RemoteViewDescription = + response.json().await.err_to_http(request_id.clone())?; + description.into_description(request_id) + } + + async fn describe_view( + &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}/describe")); + let (request_id, response) = self.client.send(req).await?; + let response = self.client.check_response(&request_id, response).await?; + let description: RemoteViewDescription = + response.json().await.err_to_http(request_id.clone())?; + description.into_description(request_id) + } + + async fn drop_view(&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(()) + } + + async fn list_views(&self, namespace_path: &[String]) -> Result> { + let namespace_id = build_namespace_identifier(namespace_path)?; + let path = format!("/v1/namespace/{namespace_id}/view/list"); + let mut views = Vec::new(); + let mut page_token: Option = None; + let mut seen_page_tokens = HashSet::new(); + loop { + let mut req = self.client.get(&path); + if let Some(token) = &page_token { + req = req.query(&[("page_token", token)]); + } + let (request_id, response) = self.client.send(req).await?; + let response = self.client.check_response(&request_id, response).await?; + let status = response.status(); + let response: RemoteListViewsResponse = + response.json().await.err_to_http(request_id.clone())?; + views.extend(response.views); + let Some(next_page_token) = response.page_token.filter(|token| !token.is_empty()) + else { + break; + }; + if !seen_page_tokens.insert(next_page_token.clone()) { + return Err(Error::Http { + source: "View listing response repeated a page_token".into(), + request_id, + status_code: Some(status), + }); + } + page_token = Some(next_page_token); + } + Ok(views) + } + async fn open_job(&self, job_id: &str) -> Result { let handle = super::job::RemoteJob::new(self.client.clone(), job_id.to_string()); match crate::job::JobHandle::describe(&handle).await { @@ -3798,6 +3925,166 @@ mod tests { .unwrap(); } + /// A view description carries the schema in the namespace spec's JSON + /// encoding, so the body a test asserts against is the one that encoding + /// produces rather than a hand-written guess at it. + fn view_description_body(name: &str, namespace: &[&str], query: &str) -> String { + let schema = Schema::new(vec![Field::new("id", DataType::Int32, true)]); + serde_json::json!({ + "name": name, + "namespace": namespace, + "query": query, + "default_database": "db", + "default_namespace": namespace, + "schema": lance_namespace::schema::arrow_schema_to_json(&schema).unwrap(), + }) + .to_string() + } + + #[tokio::test] + async fn test_create_view_posts_the_query_and_returns_the_resolved_schema() { + let conn = Connection::new_with_handler(|request| { + assert_eq!(request.method(), &reqwest::Method::POST); + assert_eq!(request.url().path(), "/v1/view/analytics$adults/create"); + let body: serde_json::Value = + serde_json::from_slice(request.body().unwrap().as_bytes().unwrap()).unwrap(); + // The name and its namespace address the view in the path, so the + // body says only what they cannot. + assert_eq!(body, serde_json::json!({"query": "SELECT id FROM people"})); + http::Response::builder() + .status(200) + .body(view_description_body( + "adults", + &["analytics"], + "SELECT id FROM people", + )) + .unwrap() + }); + let view = conn + .create_view("adults", "SELECT id FROM people", &["analytics".into()]) + .await + .unwrap(); + assert_eq!(view.name, "adults"); + assert_eq!(view.namespace_path, vec!["analytics".to_string()]); + assert_eq!(view.query, "SELECT id FROM people"); + assert_eq!(view.default_database, "db"); + assert_eq!(view.default_namespace_path, vec!["analytics".to_string()]); + assert_eq!(view.schema.field(0).name(), "id"); + } + + #[tokio::test] + async fn test_describe_and_drop_address_the_view_in_the_path() { + let conn = Connection::new_with_handler(|request| { + assert_eq!(request.url().path(), "/v1/view/adults/describe"); + assert!(request.body().is_none(), "{:?}", request.body()); + http::Response::builder() + .status(200) + .body(view_description_body( + "adults", + &[], + "SELECT id FROM people", + )) + .unwrap() + }); + let view = conn.describe_view("adults", &[]).await.unwrap(); + assert!(view.namespace_path.is_empty()); + assert_eq!(view.schema.fields().len(), 1); + + let conn = Connection::new_with_handler(|request| { + assert_eq!(request.method(), &reqwest::Method::POST); + assert_eq!(request.url().path(), "/v1/view/analytics$adults/drop"); + assert!(request.body().is_none(), "{:?}", request.body()); + http::Response::builder().status(200).body("{}").unwrap() + }); + conn.drop_view("adults", &["analytics".into()]) + .await + .unwrap(); + } + + /// 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. + #[tokio::test] + async fn test_describe_view_rejects_an_undecodable_schema() { + let conn = Connection::new_with_handler(|_| { + http::Response::builder() + .status(200) + .body( + r#"{"name":"adults","query":"SELECT 1","default_database":"db", + "schema":{"fields":[{"name":"id","type":{"type":"nonesuch"}, + "nullable":true}]}}"#, + ) + .unwrap() + }); + let error = conn.describe_view("adults", &[]).await.unwrap_err(); + assert!(error.to_string().contains("undecodable schema"), "{error}"); + } + + #[tokio::test] + async fn test_list_views_walks_pages() { + let conn = Connection::new_with_handler(|request| { + assert_eq!(request.method(), &reqwest::Method::GET); + assert_eq!(request.url().path(), "/v1/namespace/analytics/view/list"); + let page = request + .url() + .query_pairs() + .find(|(key, _)| key == "page_token") + .map(|(_, value)| value.into_owned()); + let body = match page.as_deref() { + None => r#"{"views":[],"page_token":"p2"}"#, + Some("p2") => r#"{"views":["adults"]}"#, + Some(other) => panic!("unexpected page token: {other}"), + }; + http::Response::builder().status(200).body(body).unwrap() + }); + assert_eq!( + conn.list_views(&["analytics".into()]).await.unwrap(), + vec!["adults".to_string()] + ); + } + + /// A server that keeps handing back the same token would otherwise spin + /// forever. + #[tokio::test] + async fn test_list_views_rejects_a_repeated_page_token() { + let conn = Connection::new_with_handler(|_| { + http::Response::builder() + .status(200) + .body(r#"{"views":["adults"],"page_token":"same"}"#) + .unwrap() + }); + let error = conn.list_views(&[]).await.unwrap_err(); + assert!( + error.to_string().contains("repeated a page_token"), + "{error}" + ); + } + + /// A name carrying the delimiter would split back apart as a different + /// view, so it is refused before it reaches a route. + #[tokio::test] + async fn test_view_names_that_would_resplit_are_refused() { + let conn = Connection::new_with_handler(|_| -> http::Response { + panic!("an invalid identifier must not reach the service") + }); + for (name, namespace) in [ + ("analytics$adults", vec![]), + ("adults", vec!["ana$lytics".to_string()]), + ("", vec![]), + ("..", vec![]), + ] { + let error = match conn.describe_view(name, &namespace).await { + Ok(_) => panic!("accepted {name:?} in {namespace:?}"), + Err(error) => error.to_string(), + }; + // A view is not a table, so the refusal says so. + assert!( + error.contains("view name") || error.contains("view namespace path segment"), + "{error}" + ); + } + } + #[tokio::test] async fn test_create_function_async_sends_canonical_request_and_decodes_typed_job() { const REQUEST: &str = include_str!( diff --git a/rust/lancedb/src/utils/mod.rs b/rust/lancedb/src/utils/mod.rs index f80cc876a..52b7bd5e8 100644 --- a/rust/lancedb/src/utils/mod.rs +++ b/rust/lancedb/src/utils/mod.rs @@ -174,6 +174,27 @@ pub fn validate_secret_reference(name: &str, namespace_path: &[String]) -> Resul validate_secret_component("Secret name", name) } +/// Validate a view name and every segment of the namespace path holding it. +/// +/// Worded for a view rather than deferring to [`validate_table_name`]: a view +/// is not a table, and a caller who mistypes one should not be told their +/// table name is invalid. The rule is the same one every object name follows, +/// which is what lets the `$`-joined identifier split back apart. +pub fn validate_view_reference(name: &str, namespace_path: &[String]) -> Result<()> { + for segment in namespace_path { + validate_view_component("view namespace path segment", segment)?; + } + validate_view_component("view name", name) +} + +/// Validate one component of a view identifier: a view name, or one segment of +/// the namespace path holding it. +pub fn validate_view_component(what: &str, value: &str) -> Result<()> { + check_object_name(value).map_err(|reason| Error::InvalidInput { + message: format!("invalid {what} '{value}': {reason}"), + }) +} + /// Validate all components of a namespace /// /// Iterates through all namespace components and validates each one. diff --git a/rust/lancedb/src/view.rs b/rust/lancedb/src/view.rs new file mode 100644 index 000000000..ca9703209 --- /dev/null +++ b/rust/lancedb/src/view.rs @@ -0,0 +1,44 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The LanceDB Authors + +//! Views: a named query a database stores and plans on every read. +//! +//! A view holds no rows. What it stores is the statement that defines it and +//! the schema that statement resolved to, so a reader sees the sources as they +//! are now rather than as they were when the view was created. That is the +//! whole difference from [`crate::materialized_view`], which holds rows and +//! moves them forward by refresh. +//! +//! The verbs live on [`crate::connection::Connection`]. A view is read through +//! SQL, by name, so there is no `open_view` returning a +//! [`crate::table::Table`]: there would be no rows behind it. + +use arrow_schema::SchemaRef; + +/// What a database records about one view. +/// +/// Returned by [`crate::connection::Connection::describe_view`], and by +/// `create_view` so a caller has the resolved schema without a second call. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ViewDescription { + /// The view's name within its namespace. + pub name: String, + /// The namespace holding the view; empty is the root namespace. + pub namespace_path: Vec, + /// The defining query, as the database stores it. + pub query: String, + /// The database that unqualified table names in `query` resolve against. + /// A view created from a session connected elsewhere can name a database + /// other than its own, so this is part of what the query means rather + /// than a restatement of where the view lives. + pub default_database: String, + /// The namespace path those unqualified names resolve against; empty is + /// the root namespace. + /// + /// Recorded with the view because it outlives the session that declared + /// it: a reader that resolved the query against its own default namespace + /// could read a different table than the view was defined over. + pub default_namespace_path: Vec, + /// The schema the defining query resolved to when the view was created. + pub schema: SchemaRef, +}