mirror of
https://github.com/lancedb/lancedb.git
synced 2026-08-29 09:28:27 +00:00
Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 333b0c7c7c | |||
| fd98b845ea | |||
| be48ada352 | |||
| 9ad2dfe601 | |||
| f909df3e87 |
+1
-1
@@ -1,5 +1,5 @@
|
|||||||
[tool.bumpversion]
|
[tool.bumpversion]
|
||||||
current_version = "0.28.0-beta.7"
|
current_version = "0.28.0-beta.8"
|
||||||
parse = """(?x)
|
parse = """(?x)
|
||||||
(?P<major>0|[1-9]\\d*)\\.
|
(?P<major>0|[1-9]\\d*)\\.
|
||||||
(?P<minor>0|[1-9]\\d*)\\.
|
(?P<minor>0|[1-9]\\d*)\\.
|
||||||
|
|||||||
Generated
+3
-3
@@ -4576,7 +4576,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "lancedb"
|
name = "lancedb"
|
||||||
version = "0.28.0-beta.7"
|
version = "0.28.0-beta.8"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"ahash",
|
"ahash",
|
||||||
"anyhow",
|
"anyhow",
|
||||||
@@ -4658,7 +4658,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "lancedb-nodejs"
|
name = "lancedb-nodejs"
|
||||||
version = "0.28.0-beta.7"
|
version = "0.28.0-beta.8"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"arrow-array",
|
"arrow-array",
|
||||||
"arrow-buffer",
|
"arrow-buffer",
|
||||||
@@ -4680,7 +4680,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "lancedb-python"
|
name = "lancedb-python"
|
||||||
version = "0.31.0-beta.7"
|
version = "0.31.0-beta.8"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"arrow",
|
"arrow",
|
||||||
"async-trait",
|
"async-trait",
|
||||||
|
|||||||
@@ -14,7 +14,7 @@ Add the following dependency to your `pom.xml`:
|
|||||||
<dependency>
|
<dependency>
|
||||||
<groupId>com.lancedb</groupId>
|
<groupId>com.lancedb</groupId>
|
||||||
<artifactId>lancedb-core</artifactId>
|
<artifactId>lancedb-core</artifactId>
|
||||||
<version>0.28.0-beta.7</version>
|
<version>0.28.0-beta.8</version>
|
||||||
</dependency>
|
</dependency>
|
||||||
```
|
```
|
||||||
|
|
||||||
|
|||||||
@@ -8,7 +8,7 @@
|
|||||||
<parent>
|
<parent>
|
||||||
<groupId>com.lancedb</groupId>
|
<groupId>com.lancedb</groupId>
|
||||||
<artifactId>lancedb-parent</artifactId>
|
<artifactId>lancedb-parent</artifactId>
|
||||||
<version>0.28.0-beta.7</version>
|
<version>0.28.0-beta.8</version>
|
||||||
<relativePath>../pom.xml</relativePath>
|
<relativePath>../pom.xml</relativePath>
|
||||||
</parent>
|
</parent>
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -6,7 +6,7 @@
|
|||||||
|
|
||||||
<groupId>com.lancedb</groupId>
|
<groupId>com.lancedb</groupId>
|
||||||
<artifactId>lancedb-parent</artifactId>
|
<artifactId>lancedb-parent</artifactId>
|
||||||
<version>0.28.0-beta.7</version>
|
<version>0.28.0-beta.8</version>
|
||||||
<packaging>pom</packaging>
|
<packaging>pom</packaging>
|
||||||
<name>${project.artifactId}</name>
|
<name>${project.artifactId}</name>
|
||||||
<description>LanceDB Java SDK Parent POM</description>
|
<description>LanceDB Java SDK Parent POM</description>
|
||||||
|
|||||||
+1
-1
@@ -1,7 +1,7 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "lancedb-nodejs"
|
name = "lancedb-nodejs"
|
||||||
edition.workspace = true
|
edition.workspace = true
|
||||||
version = "0.28.0-beta.7"
|
version = "0.28.0-beta.8"
|
||||||
license.workspace = true
|
license.workspace = true
|
||||||
description.workspace = true
|
description.workspace = true
|
||||||
repository.workspace = true
|
repository.workspace = true
|
||||||
|
|||||||
@@ -1,6 +1,8 @@
|
|||||||
// SPDX-License-Identifier: Apache-2.0
|
// SPDX-License-Identifier: Apache-2.0
|
||||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||||
|
|
||||||
|
import { spawn } from "node:child_process";
|
||||||
|
import * as path from "node:path";
|
||||||
import { RecordBatch } from "apache-arrow";
|
import { RecordBatch } from "apache-arrow";
|
||||||
import * as tmp from "tmp";
|
import * as tmp from "tmp";
|
||||||
import { Connection, Index, Table, connect, makeArrowTable } from "../lancedb";
|
import { Connection, Index, Table, connect, makeArrowTable } from "../lancedb";
|
||||||
@@ -76,4 +78,91 @@ describe("rerankers", function () {
|
|||||||
|
|
||||||
expect(result).toHaveLength(2);
|
expect(result).toHaveLength(2);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("does not keep process alive after rerank query", async function () {
|
||||||
|
const script = `
|
||||||
|
import * as lancedb from "./dist/index.js";
|
||||||
|
import * as os from "node:os";
|
||||||
|
import * as path from "node:path";
|
||||||
|
import * as fs from "node:fs/promises";
|
||||||
|
|
||||||
|
const dir = await fs.mkdtemp(path.join(os.tmpdir(), "lancedb-rerank-exit-"));
|
||||||
|
const db = await lancedb.connect(dir);
|
||||||
|
const table = await db.createTable("test", [{ text: "hello", vector: [1, 2, 3] }], {
|
||||||
|
mode: "overwrite",
|
||||||
|
});
|
||||||
|
await table.createIndex("text", { config: lancedb.Index.fts() });
|
||||||
|
await table.waitForIndex(["text_idx"], 30);
|
||||||
|
|
||||||
|
const reranker = await lancedb.rerankers.RRFReranker.create();
|
||||||
|
await table
|
||||||
|
.query()
|
||||||
|
.nearestTo([1, 2, 3])
|
||||||
|
.fullTextSearch("hello")
|
||||||
|
.rerank(reranker)
|
||||||
|
.toArray();
|
||||||
|
|
||||||
|
table.close();
|
||||||
|
db.close();
|
||||||
|
`;
|
||||||
|
|
||||||
|
await new Promise<void>((resolve, reject) => {
|
||||||
|
const child = spawn(
|
||||||
|
process.execPath,
|
||||||
|
["--input-type=module", "-e", script],
|
||||||
|
{
|
||||||
|
cwd: path.resolve(__dirname, ".."),
|
||||||
|
stdio: ["ignore", "pipe", "pipe"],
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
let stdout = "";
|
||||||
|
let stderr = "";
|
||||||
|
|
||||||
|
child.stdout.on("data", (chunk) => {
|
||||||
|
stdout += chunk.toString();
|
||||||
|
});
|
||||||
|
|
||||||
|
child.stderr.on("data", (chunk) => {
|
||||||
|
stderr += chunk.toString();
|
||||||
|
});
|
||||||
|
|
||||||
|
const timeout = setTimeout(() => {
|
||||||
|
child.kill();
|
||||||
|
reject(
|
||||||
|
new Error(
|
||||||
|
`child process did not exit in time\nstdout:\n${stdout}\nstderr:\n${stderr}`,
|
||||||
|
),
|
||||||
|
);
|
||||||
|
}, 20_000);
|
||||||
|
|
||||||
|
child.on("error", (err) => {
|
||||||
|
clearTimeout(timeout);
|
||||||
|
reject(err);
|
||||||
|
});
|
||||||
|
|
||||||
|
child.on("exit", (code, signal) => {
|
||||||
|
clearTimeout(timeout);
|
||||||
|
if (signal !== null) {
|
||||||
|
reject(
|
||||||
|
new Error(
|
||||||
|
`child process exited with signal ${signal}\nstdout:\n${stdout}\nstderr:\n${stderr}`,
|
||||||
|
),
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (code !== 0) {
|
||||||
|
reject(
|
||||||
|
new Error(
|
||||||
|
`child process exited with code ${code}\nstdout:\n${stdout}\nstderr:\n${stderr}`,
|
||||||
|
),
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
resolve();
|
||||||
|
});
|
||||||
|
});
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@lancedb/lancedb-darwin-arm64",
|
"name": "@lancedb/lancedb-darwin-arm64",
|
||||||
"version": "0.28.0-beta.7",
|
"version": "0.28.0-beta.8",
|
||||||
"os": ["darwin"],
|
"os": ["darwin"],
|
||||||
"cpu": ["arm64"],
|
"cpu": ["arm64"],
|
||||||
"main": "lancedb.darwin-arm64.node",
|
"main": "lancedb.darwin-arm64.node",
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@lancedb/lancedb-linux-arm64-gnu",
|
"name": "@lancedb/lancedb-linux-arm64-gnu",
|
||||||
"version": "0.28.0-beta.7",
|
"version": "0.28.0-beta.8",
|
||||||
"os": ["linux"],
|
"os": ["linux"],
|
||||||
"cpu": ["arm64"],
|
"cpu": ["arm64"],
|
||||||
"main": "lancedb.linux-arm64-gnu.node",
|
"main": "lancedb.linux-arm64-gnu.node",
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@lancedb/lancedb-linux-arm64-musl",
|
"name": "@lancedb/lancedb-linux-arm64-musl",
|
||||||
"version": "0.28.0-beta.7",
|
"version": "0.28.0-beta.8",
|
||||||
"os": ["linux"],
|
"os": ["linux"],
|
||||||
"cpu": ["arm64"],
|
"cpu": ["arm64"],
|
||||||
"main": "lancedb.linux-arm64-musl.node",
|
"main": "lancedb.linux-arm64-musl.node",
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@lancedb/lancedb-linux-x64-gnu",
|
"name": "@lancedb/lancedb-linux-x64-gnu",
|
||||||
"version": "0.28.0-beta.7",
|
"version": "0.28.0-beta.8",
|
||||||
"os": ["linux"],
|
"os": ["linux"],
|
||||||
"cpu": ["x64"],
|
"cpu": ["x64"],
|
||||||
"main": "lancedb.linux-x64-gnu.node",
|
"main": "lancedb.linux-x64-gnu.node",
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@lancedb/lancedb-linux-x64-musl",
|
"name": "@lancedb/lancedb-linux-x64-musl",
|
||||||
"version": "0.28.0-beta.7",
|
"version": "0.28.0-beta.8",
|
||||||
"os": ["linux"],
|
"os": ["linux"],
|
||||||
"cpu": ["x64"],
|
"cpu": ["x64"],
|
||||||
"main": "lancedb.linux-x64-musl.node",
|
"main": "lancedb.linux-x64-musl.node",
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@lancedb/lancedb-win32-arm64-msvc",
|
"name": "@lancedb/lancedb-win32-arm64-msvc",
|
||||||
"version": "0.28.0-beta.7",
|
"version": "0.28.0-beta.8",
|
||||||
"os": [
|
"os": [
|
||||||
"win32"
|
"win32"
|
||||||
],
|
],
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@lancedb/lancedb-win32-x64-msvc",
|
"name": "@lancedb/lancedb-win32-x64-msvc",
|
||||||
"version": "0.28.0-beta.7",
|
"version": "0.28.0-beta.8",
|
||||||
"os": ["win32"],
|
"os": ["win32"],
|
||||||
"cpu": ["x64"],
|
"cpu": ["x64"],
|
||||||
"main": "lancedb.win32-x64-msvc.node",
|
"main": "lancedb.win32-x64-msvc.node",
|
||||||
|
|||||||
Generated
+2
-2
@@ -1,12 +1,12 @@
|
|||||||
{
|
{
|
||||||
"name": "@lancedb/lancedb",
|
"name": "@lancedb/lancedb",
|
||||||
"version": "0.28.0-beta.7",
|
"version": "0.28.0-beta.8",
|
||||||
"lockfileVersion": 3,
|
"lockfileVersion": 3,
|
||||||
"requires": true,
|
"requires": true,
|
||||||
"packages": {
|
"packages": {
|
||||||
"": {
|
"": {
|
||||||
"name": "@lancedb/lancedb",
|
"name": "@lancedb/lancedb",
|
||||||
"version": "0.28.0-beta.7",
|
"version": "0.28.0-beta.8",
|
||||||
"cpu": [
|
"cpu": [
|
||||||
"x64",
|
"x64",
|
||||||
"arm64"
|
"arm64"
|
||||||
|
|||||||
+1
-1
@@ -11,7 +11,7 @@
|
|||||||
"ann"
|
"ann"
|
||||||
],
|
],
|
||||||
"private": false,
|
"private": false,
|
||||||
"version": "0.28.0-beta.7",
|
"version": "0.28.0-beta.8",
|
||||||
"main": "dist/index.js",
|
"main": "dist/index.js",
|
||||||
"exports": {
|
"exports": {
|
||||||
".": "./dist/index.js",
|
".": "./dist/index.js",
|
||||||
|
|||||||
@@ -18,6 +18,7 @@ type RerankHybridFn = ThreadsafeFunction<
|
|||||||
RerankHybridCallbackArgs,
|
RerankHybridCallbackArgs,
|
||||||
Status,
|
Status,
|
||||||
false,
|
false,
|
||||||
|
true,
|
||||||
>;
|
>;
|
||||||
|
|
||||||
/// Reranker implementation that "wraps" a NodeJS Reranker implementation.
|
/// Reranker implementation that "wraps" a NodeJS Reranker implementation.
|
||||||
@@ -32,7 +33,10 @@ impl Reranker {
|
|||||||
pub fn new(
|
pub fn new(
|
||||||
rerank_hybrid: Function<RerankHybridCallbackArgs, Promise<Buffer>>,
|
rerank_hybrid: Function<RerankHybridCallbackArgs, Promise<Buffer>>,
|
||||||
) -> napi::Result<Self> {
|
) -> napi::Result<Self> {
|
||||||
let rerank_hybrid = rerank_hybrid.build_threadsafe_function().build()?;
|
let rerank_hybrid = rerank_hybrid
|
||||||
|
.build_threadsafe_function()
|
||||||
|
.weak::<true>()
|
||||||
|
.build()?;
|
||||||
Ok(Self { rerank_hybrid })
|
Ok(Self { rerank_hybrid })
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
[tool.bumpversion]
|
[tool.bumpversion]
|
||||||
current_version = "0.31.0-beta.7"
|
current_version = "0.31.0-beta.8"
|
||||||
parse = """(?x)
|
parse = """(?x)
|
||||||
(?P<major>0|[1-9]\\d*)\\.
|
(?P<major>0|[1-9]\\d*)\\.
|
||||||
(?P<minor>0|[1-9]\\d*)\\.
|
(?P<minor>0|[1-9]\\d*)\\.
|
||||||
|
|||||||
+1
-1
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "lancedb-python"
|
name = "lancedb-python"
|
||||||
version = "0.31.0-beta.7"
|
version = "0.31.0-beta.8"
|
||||||
edition.workspace = true
|
edition.workspace = true
|
||||||
description = "Python bindings for LanceDB"
|
description = "Python bindings for LanceDB"
|
||||||
license.workspace = true
|
license.workspace = true
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "lancedb"
|
name = "lancedb"
|
||||||
version = "0.28.0-beta.7"
|
version = "0.28.0-beta.8"
|
||||||
edition.workspace = true
|
edition.workspace = true
|
||||||
description = "LanceDB: A serverless, low-latency vector database for AI applications"
|
description = "LanceDB: A serverless, low-latency vector database for AI applications"
|
||||||
license.workspace = true
|
license.workspace = true
|
||||||
|
|||||||
@@ -26,6 +26,7 @@ use crate::connection::NamespaceClientPushdownOperation;
|
|||||||
use crate::database::ReadConsistency;
|
use crate::database::ReadConsistency;
|
||||||
use crate::error::{Error, Result};
|
use crate::error::{Error, Result};
|
||||||
use crate::table::NativeTable;
|
use crate::table::NativeTable;
|
||||||
|
use lance::dataset::WriteMode;
|
||||||
|
|
||||||
use super::{
|
use super::{
|
||||||
BaseTable, CloneTableRequest, CreateTableMode, CreateTableRequest as DbCreateTableRequest,
|
BaseTable, CloneTableRequest, CreateTableMode, CreateTableRequest as DbCreateTableRequest,
|
||||||
@@ -184,6 +185,7 @@ impl Database for LanceNamespaceDatabase {
|
|||||||
async fn create_table(&self, request: DbCreateTableRequest) -> Result<Arc<dyn BaseTable>> {
|
async fn create_table(&self, request: DbCreateTableRequest) -> Result<Arc<dyn BaseTable>> {
|
||||||
let mut table_id = request.namespace_path.clone();
|
let mut table_id = request.namespace_path.clone();
|
||||||
table_id.push(request.name.clone());
|
table_id.push(request.name.clone());
|
||||||
|
let mut existing_table = None;
|
||||||
|
|
||||||
match request.mode {
|
match request.mode {
|
||||||
CreateTableMode::Create => {}
|
CreateTableMode::Create => {}
|
||||||
@@ -192,20 +194,7 @@ impl Database for LanceNamespaceDatabase {
|
|||||||
id: Some(table_id.clone()),
|
id: Some(table_id.clone()),
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
let describe_result = self.namespace.describe_table(describe_request).await;
|
existing_table = self.namespace.describe_table(describe_request).await.ok();
|
||||||
if describe_result.is_ok() {
|
|
||||||
// Drop the existing table - must succeed
|
|
||||||
let drop_request = DropTableRequest {
|
|
||||||
id: Some(table_id.clone()),
|
|
||||||
..Default::default()
|
|
||||||
};
|
|
||||||
self.namespace
|
|
||||||
.drop_table(drop_request)
|
|
||||||
.await
|
|
||||||
.map_err(|e| Error::Runtime {
|
|
||||||
message: format!("Failed to drop existing table for overwrite: {}", e),
|
|
||||||
})?;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
CreateTableMode::ExistOk(_) => {
|
CreateTableMode::ExistOk(_) => {
|
||||||
let describe_request = DescribeTableRequest {
|
let describe_request = DescribeTableRequest {
|
||||||
@@ -240,40 +229,82 @@ impl Database for LanceNamespaceDatabase {
|
|||||||
};
|
};
|
||||||
|
|
||||||
let (location, initial_storage_options, managed_versioning) = {
|
let (location, initial_storage_options, managed_versioning) = {
|
||||||
let response = self
|
if let Some(response) = existing_table {
|
||||||
.namespace
|
let loc = response.location.ok_or_else(|| Error::Runtime {
|
||||||
.declare_table(declare_request)
|
message: "Table location is missing from describe_table response".to_string(),
|
||||||
.await
|
|
||||||
.map_err(|e| {
|
|
||||||
let err_str = e.to_string();
|
|
||||||
if matches!(request.mode, CreateTableMode::Create)
|
|
||||||
&& (err_str.contains("already exists")
|
|
||||||
|| err_str.contains("TableAlreadyExists")
|
|
||||||
|| err_str.contains("table already exists"))
|
|
||||||
{
|
|
||||||
Error::TableAlreadyExists {
|
|
||||||
name: request.name.clone(),
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
Error::Runtime {
|
|
||||||
message: format!("Failed to declare table: {}", e),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})?;
|
})?;
|
||||||
let loc = response.location.ok_or_else(|| Error::Runtime {
|
let opts = response
|
||||||
message: "Table location is missing from declare_table response".to_string(),
|
.storage_options
|
||||||
})?;
|
.or_else(|| Some(self.storage_options.clone()))
|
||||||
// Use storage options from response, fall back to self.storage_options
|
.filter(|o| !o.is_empty());
|
||||||
let opts = response
|
(loc, opts, response.managed_versioning)
|
||||||
.storage_options
|
} else {
|
||||||
.or_else(|| Some(self.storage_options.clone()))
|
match self.namespace.declare_table(declare_request).await {
|
||||||
.filter(|o| !o.is_empty());
|
Ok(response) => {
|
||||||
(loc, opts, response.managed_versioning)
|
let loc = response.location.ok_or_else(|| Error::Runtime {
|
||||||
|
message: "Table location is missing from declare_table response"
|
||||||
|
.to_string(),
|
||||||
|
})?;
|
||||||
|
let opts = response
|
||||||
|
.storage_options
|
||||||
|
.or_else(|| Some(self.storage_options.clone()))
|
||||||
|
.filter(|o: &HashMap<String, String>| !o.is_empty());
|
||||||
|
(loc, opts, response.managed_versioning)
|
||||||
|
}
|
||||||
|
Err(e)
|
||||||
|
if matches!(request.mode, CreateTableMode::Create) && {
|
||||||
|
let err_str = e.to_string();
|
||||||
|
err_str.contains("already exists")
|
||||||
|
|| err_str.contains("TableAlreadyExists")
|
||||||
|
|| err_str.contains("table already exists")
|
||||||
|
} =>
|
||||||
|
{
|
||||||
|
let response = self
|
||||||
|
.namespace
|
||||||
|
.describe_table(DescribeTableRequest {
|
||||||
|
id: Some(table_id.clone()),
|
||||||
|
..Default::default()
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.map_err(|describe_err| Error::Runtime {
|
||||||
|
message: format!(
|
||||||
|
"Failed to describe existing declared table after declare conflict: {}",
|
||||||
|
describe_err
|
||||||
|
),
|
||||||
|
})?;
|
||||||
|
|
||||||
|
if response.version.is_some() && response.schema.is_some() {
|
||||||
|
return Err(Error::TableAlreadyExists {
|
||||||
|
name: request.name.clone(),
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
let loc = response.location.ok_or_else(|| Error::Runtime {
|
||||||
|
message: "Table location is missing from describe_table response"
|
||||||
|
.to_string(),
|
||||||
|
})?;
|
||||||
|
let opts = response
|
||||||
|
.storage_options
|
||||||
|
.or_else(|| Some(self.storage_options.clone()))
|
||||||
|
.filter(|o: &HashMap<String, String>| !o.is_empty());
|
||||||
|
(loc, opts, response.managed_versioning)
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
return Err(Error::Runtime {
|
||||||
|
message: format!("Failed to declare table: {}", e),
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
// Build write params with storage options and commit handler
|
// Build write params with storage options and commit handler
|
||||||
let mut params = request.write_options.lance_write_params.unwrap_or_default();
|
let mut params = request.write_options.lance_write_params.unwrap_or_default();
|
||||||
|
|
||||||
|
if matches!(request.mode, CreateTableMode::Overwrite) {
|
||||||
|
params.mode = WriteMode::Overwrite;
|
||||||
|
}
|
||||||
|
|
||||||
// Set up storage options if provided
|
// Set up storage options if provided
|
||||||
if let Some(storage_opts) = initial_storage_options {
|
if let Some(storage_opts) = initial_storage_options {
|
||||||
let store_params = params
|
let store_params = params
|
||||||
@@ -730,6 +761,58 @@ mod tests {
|
|||||||
assert_eq!(id_col.value(2), 30);
|
assert_eq!(id_col.value(2), 30);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn test_namespace_create_table_after_declare_conflict() {
|
||||||
|
let tmp_dir = tempdir().unwrap();
|
||||||
|
let root_path = tmp_dir.path().to_str().unwrap().to_string();
|
||||||
|
|
||||||
|
let mut properties = HashMap::new();
|
||||||
|
properties.insert("root".to_string(), root_path);
|
||||||
|
|
||||||
|
let conn = connect_namespace("dir", properties)
|
||||||
|
.execute()
|
||||||
|
.await
|
||||||
|
.expect("Failed to connect to namespace");
|
||||||
|
|
||||||
|
conn.create_namespace(CreateNamespaceRequest {
|
||||||
|
id: Some(vec!["test_ns".into()]),
|
||||||
|
..Default::default()
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.expect("Failed to create namespace");
|
||||||
|
|
||||||
|
let namespace_client = conn.namespace_client().await.unwrap();
|
||||||
|
namespace_client
|
||||||
|
.declare_table(DeclareTableRequest {
|
||||||
|
id: Some(vec!["test_ns".into(), "declared_test".into()]),
|
||||||
|
..Default::default()
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.expect("Failed to declare table");
|
||||||
|
|
||||||
|
let test_data = create_test_data();
|
||||||
|
let table = conn
|
||||||
|
.create_table("declared_test", test_data)
|
||||||
|
.namespace(vec!["test_ns".into()])
|
||||||
|
.execute()
|
||||||
|
.await
|
||||||
|
.expect("Failed to create table after declare conflict");
|
||||||
|
|
||||||
|
let results = table
|
||||||
|
.query()
|
||||||
|
.execute()
|
||||||
|
.await
|
||||||
|
.expect("Failed to query table")
|
||||||
|
.try_collect::<Vec<_>>()
|
||||||
|
.await
|
||||||
|
.expect("Failed to collect results");
|
||||||
|
|
||||||
|
assert_eq!(results.len(), 1);
|
||||||
|
assert_eq!(results[0].num_rows(), 5);
|
||||||
|
assert_eq!(table.namespace(), &["test_ns"]);
|
||||||
|
assert_eq!(table.id(), "test_ns$declared_test");
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_namespace_create_table_exist_ok_mode() {
|
async fn test_namespace_create_table_exist_ok_mode() {
|
||||||
// Setup: Create a temporary directory for the namespace
|
// Setup: Create a temporary directory for the namespace
|
||||||
|
|||||||
Reference in New Issue
Block a user