mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-30 00:45:37 +00:00
feat: add view CRUD APIs (#4236)
A view is a named query a database stores and plans on every read. It
holds no rows, which is the whole difference from a materialized view.
## API
| Verb | Route |
| --- | --- |
| `create_view(name, query, namespace_path)` | `POST
/v1/view/{id}/create` |
| `describe_view(name, namespace_path)` | `POST /v1/view/{id}/describe`
|
| `drop_view(name, namespace_path)` | `POST /v1/view/{id}/drop` |
| `list_views(namespace_path)` | `GET /v1/namespace/{id}/view/list` |
On `Connection` and the `Database` trait, with the remote client, Python
(sync and async) and Node bindings. Local databases return
`NotSupported`: the server side is Sophon's, where a view is an object
of the database manifest.
`ViewDescription` carries the defining query, the database *and
namespace path* unqualified names in it resolve against, and the schema
the query resolved to. `create_view` returns one, so a caller has the
schema without a second call.
Both defaults travel with the view because it outlives the session that
declared it: the server re-plans the stored query on every read, so a
reader resolving an unqualified name against its own defaults would read
a different table. `default_namespace_path` crosses the wire as
`default_namespace`, a path like `namespace`, absent for the root.
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.
Querying a view stays SQL's job. There are no rows behind a view, so
there is no `open_view` returning a `Table`.
This commit is contained in:
@@ -319,6 +319,38 @@ Creates a new Table and initialize it with new data.
|
||||
|
||||
***
|
||||
|
||||
### createView()
|
||||
|
||||
```ts
|
||||
abstract createView(
|
||||
name,
|
||||
query,
|
||||
namespacePath?): Promise<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
|
||||
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<ViewDescription>
|
||||
```
|
||||
|
||||
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<void>
|
||||
```
|
||||
|
||||
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<string[]>
|
||||
```
|
||||
|
||||
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
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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<any>;
|
||||
```
|
||||
|
||||
The schema the defining query resolved to when the view was created.
|
||||
@@ -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
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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<Job>;
|
||||
|
||||
/**
|
||||
* 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<ViewDescription>;
|
||||
|
||||
/**
|
||||
* 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<ViewDescription>;
|
||||
|
||||
/**
|
||||
* 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<void>;
|
||||
|
||||
/**
|
||||
* The names of the views in one namespace.
|
||||
*
|
||||
* Names only; a definition comes from {@link describeView}.
|
||||
*/
|
||||
abstract listViews(namespacePath?: string[]): Promise<string[]>;
|
||||
|
||||
abstract openTable(
|
||||
name: string,
|
||||
namespacePath?: string[],
|
||||
@@ -696,6 +736,33 @@ export class LocalConnection extends Connection {
|
||||
);
|
||||
}
|
||||
|
||||
async createView(
|
||||
name: string,
|
||||
query: string,
|
||||
namespacePath?: string[],
|
||||
): Promise<ViewDescription> {
|
||||
return viewDescriptionFromNative(
|
||||
await this.inner.createView(name, query, namespacePath ?? []),
|
||||
);
|
||||
}
|
||||
|
||||
async describeView(
|
||||
name: string,
|
||||
namespacePath?: string[],
|
||||
): Promise<ViewDescription> {
|
||||
return viewDescriptionFromNative(
|
||||
await this.inner.describeView(name, namespacePath ?? []),
|
||||
);
|
||||
}
|
||||
|
||||
async dropView(name: string, namespacePath?: string[]): Promise<void> {
|
||||
return this.inner.dropView(name, namespacePath ?? []);
|
||||
}
|
||||
|
||||
async listViews(namespacePath?: string[]): Promise<string[]> {
|
||||
return this.inner.listViews(namespacePath ?? []);
|
||||
}
|
||||
|
||||
async listTables(
|
||||
namespacePathOrOptions?: string[] | Partial<ListTablesOptions>,
|
||||
options?: Partial<ListTablesOptions>,
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
}
|
||||
@@ -55,6 +55,32 @@ pub struct DropNamespaceResponse {
|
||||
pub transaction_id: Option<Vec<String>>,
|
||||
}
|
||||
|
||||
/// What a database records about one view.
|
||||
#[napi(object)]
|
||||
pub struct ViewDescription {
|
||||
pub name: String,
|
||||
pub namespace_path: Vec<String>,
|
||||
pub query: String,
|
||||
pub default_database: String,
|
||||
pub default_namespace_path: Vec<String>,
|
||||
/// 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<Self> {
|
||||
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<Vec<String>>,
|
||||
) -> napi::Result<ViewDescription> {
|
||||
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<Vec<String>>,
|
||||
) -> napi::Result<ViewDescription> {
|
||||
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<Vec<String>>,
|
||||
) -> 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<Vec<String>>,
|
||||
) -> napi::Result<Vec<String>> {
|
||||
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(
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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."""
|
||||
|
||||
@@ -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.
|
||||
"""
|
||||
@@ -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"),
|
||||
]
|
||||
|
||||
@@ -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<Bound<'_, PyAny>>) -> PyResult<Vec<
|
||||
}
|
||||
}
|
||||
|
||||
/// A view description on its way to Python: name, namespace, query, default
|
||||
/// database, and the schema as pyarrow renders it.
|
||||
type PyViewDescription = (String, Vec<String>, String, String, Vec<String>, Py<PyAny>);
|
||||
|
||||
/// 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<PyViewDescription> {
|
||||
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<Vec<String>>,
|
||||
) -> PyResult<Bound<'_, PyAny>> {
|
||||
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<Vec<String>>,
|
||||
) -> PyResult<Bound<'_, PyAny>> {
|
||||
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<Vec<String>>,
|
||||
) -> PyResult<Bound<'_, PyAny>> {
|
||||
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<Vec<String>>,
|
||||
) -> PyResult<Bound<'_, PyAny>> {
|
||||
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<Bound<'_, PyAny>> {
|
||||
let inner = self_.get_inner()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
|
||||
@@ -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<dyn std::error::Error>> {
|
||||
/// 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<str>,
|
||||
query: impl AsRef<str>,
|
||||
namespace_path: &[String],
|
||||
) -> Result<ViewDescription> {
|
||||
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<str>,
|
||||
namespace_path: &[String],
|
||||
) -> Result<ViewDescription> {
|
||||
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<str>, 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<Vec<String>> {
|
||||
validate_namespace(namespace_path)?;
|
||||
self.internal.list_views(namespace_path).await
|
||||
}
|
||||
|
||||
/// Rename a table in the database.
|
||||
///
|
||||
/// This is only supported in LanceDB Cloud.
|
||||
|
||||
@@ -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<T>() -> Result<T> {
|
||||
})
|
||||
}
|
||||
|
||||
fn view_ops_not_supported<T>() -> Result<T> {
|
||||
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<SecretInfo> {
|
||||
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<ViewDescription> {
|
||||
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<ViewDescription> {
|
||||
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<Vec<String>> {
|
||||
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.
|
||||
|
||||
@@ -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};
|
||||
|
||||
|
||||
@@ -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<String>,
|
||||
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<String>,
|
||||
schema: JsonArrowSchema,
|
||||
}
|
||||
|
||||
impl RemoteViewDescription {
|
||||
fn into_description(self, request_id: String) -> Result<ViewDescription> {
|
||||
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<String>,
|
||||
#[serde(default)]
|
||||
page_token: Option<String>,
|
||||
}
|
||||
|
||||
/// 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<S: HttpSend> Database for RemoteDatabase<S> {
|
||||
response.json().await.err_to_http(request_id)
|
||||
}
|
||||
|
||||
async fn create_view(
|
||||
&self,
|
||||
name: &str,
|
||||
query: &str,
|
||||
namespace_path: &[String],
|
||||
) -> Result<ViewDescription> {
|
||||
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<ViewDescription> {
|
||||
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<Vec<String>> {
|
||||
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<String> = 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<Job> {
|
||||
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<String> {
|
||||
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!(
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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<String>,
|
||||
/// 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<String>,
|
||||
/// The schema the defining query resolved to when the view was created.
|
||||
pub schema: SchemaRef,
|
||||
}
|
||||
Reference in New Issue
Block a user