mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-11 07:42:26 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7b8ceb355c | ||
|
|
d3077b7641 | ||
|
|
3e3878b223 | ||
|
|
19fb665c76 | ||
|
|
1f95398c34 |
+1
-1
@@ -1,5 +1,5 @@
|
||||
[tool.bumpversion]
|
||||
current_version = "0.39.0-beta.4"
|
||||
current_version = "0.39.0-beta.6"
|
||||
parse = """(?x)
|
||||
(?P<major>0|[1-9]\\d*)\\.
|
||||
(?P<minor>0|[1-9]\\d*)\\.
|
||||
|
||||
Generated
+46
-45
@@ -3526,8 +3526,8 @@ checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c"
|
||||
|
||||
[[package]]
|
||||
name = "fsst"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"rand 0.9.5",
|
||||
@@ -4886,8 +4886,8 @@ checksum = "e037a2e1d8d5fdbd49b16a4ea09d5d6401c1f29eca5ff29d03d3824dba16256a"
|
||||
|
||||
[[package]]
|
||||
name = "lance"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"arrow",
|
||||
@@ -4959,8 +4959,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-arrow"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -4982,7 +4982,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-arrow-scalar"
|
||||
version = "58.0.0"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -4996,7 +4996,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-arrow-stats"
|
||||
version = "58.0.0"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5005,8 +5005,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-bitpacking"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrayref",
|
||||
"crunchy",
|
||||
@@ -5016,8 +5016,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-core"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -5054,8 +5054,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-datafusion"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5085,8 +5085,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-datagen"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5103,8 +5103,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-derive"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
@@ -5113,8 +5113,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-encoding"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-arith",
|
||||
"arrow-array",
|
||||
@@ -5147,8 +5147,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-file"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-arith",
|
||||
"arrow-array",
|
||||
@@ -5179,8 +5179,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-index"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"arrow",
|
||||
@@ -5244,8 +5244,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-index-core"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5267,8 +5267,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-io"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5308,8 +5308,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-linalg"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5323,8 +5323,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-namespace"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
@@ -5338,8 +5338,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-namespace-impls"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-ipc",
|
||||
@@ -5392,8 +5392,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-select"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -5407,8 +5407,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-table"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5448,8 +5448,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-testing"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5462,8 +5462,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-tokenizer"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"frostem",
|
||||
"icu_segmenter",
|
||||
@@ -5476,7 +5476,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lancedb"
|
||||
version = "0.39.0-beta.4"
|
||||
version = "0.39.0-beta.6"
|
||||
dependencies = [
|
||||
"ahash",
|
||||
"anyhow",
|
||||
@@ -5553,6 +5553,7 @@ dependencies = [
|
||||
"serde_json",
|
||||
"serde_with",
|
||||
"serial_test",
|
||||
"sha2 0.10.9",
|
||||
"snafu 0.8.9",
|
||||
"tempfile",
|
||||
"test-log",
|
||||
@@ -5567,7 +5568,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lancedb-nodejs"
|
||||
version = "0.39.0-beta.4"
|
||||
version = "0.39.0-beta.6"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -5592,7 +5593,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lancedb-python"
|
||||
version = "0.39.0-beta.4"
|
||||
version = "0.39.0-beta.6"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
|
||||
+14
-14
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
|
||||
rust-version = "1.91.0"
|
||||
|
||||
[workspace.dependencies]
|
||||
lance = { "version" = "=12.0.0-beta.14", default-features = false, "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=12.0.0-beta.14", default-features = false, "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=12.0.0-beta.14", default-features = false, "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance = { "version" = "=12.0.0-beta.15", default-features = false, "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=12.0.0-beta.15", default-features = false, "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=12.0.0-beta.15", default-features = false, "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lancedb = { path = "rust/lancedb", default-features = false }
|
||||
ahash = "0.8"
|
||||
# Note that this one does not include pyarrow
|
||||
|
||||
@@ -14,7 +14,7 @@ Add the following dependency to your `pom.xml`:
|
||||
<dependency>
|
||||
<groupId>com.lancedb</groupId>
|
||||
<artifactId>lancedb-core</artifactId>
|
||||
<version>0.39.0-beta.4</version>
|
||||
<version>0.39.0-beta.6</version>
|
||||
</dependency>
|
||||
```
|
||||
|
||||
|
||||
@@ -74,10 +74,10 @@ now: the column is committed with no values, and rows get them from
|
||||
[Table#refreshColumn](Table.md#refreshcolumn). Declaring one therefore costs the same on a
|
||||
large table as on an empty one.
|
||||
|
||||
A refresh does not revisit rows it has already filled, so mutating an
|
||||
input leaves the value computed at fill time; recomputing means dropping
|
||||
the column and declaring it again. While a declaration reads a column,
|
||||
that column cannot be renamed, retyped or dropped.
|
||||
A refresh also recomputes the rows whose inputs changed since they were
|
||||
computed, so a mutated input is reflected by the next refresh. While a
|
||||
declaration reads a column, that column cannot be renamed, retyped or
|
||||
dropped.
|
||||
|
||||
On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
server, and the refresh runs as a server job -- see
|
||||
@@ -854,10 +854,10 @@ abstract refreshColumn(column): Promise<RefreshColumnResult>
|
||||
|
||||
Fill the rows of a computed column that hold no value yet.
|
||||
|
||||
Rows appended since the last refresh are filled by the next one; rows
|
||||
already filled are left as they are, so the call is idempotent and does
|
||||
not observe a mutated input. Local tables only: a remote refresh runs
|
||||
as a server job, through [Table#refreshColumnAsync](Table.md#refreshcolumnasync).
|
||||
Rows appended since the last refresh are filled by the next one, and
|
||||
rows whose inputs changed since they were computed are recomputed;
|
||||
everything else is left as it is. Local tables only: a remote refresh
|
||||
runs as a server job, through [Table#refreshColumnAsync](Table.md#refreshcolumnasync).
|
||||
|
||||
#### Parameters
|
||||
|
||||
|
||||
@@ -125,10 +125,6 @@ listing a storage directory.
|
||||
|
||||
::: lancedb.functions.UdfDefinition
|
||||
|
||||
::: lancedb.secrets.EnvVarSecret
|
||||
|
||||
::: lancedb.secrets.SecretInfo
|
||||
|
||||
::: lancedb.functions.FunctionRegistrationRequest
|
||||
|
||||
::: lancedb.functions.FunctionArtifactRequest
|
||||
|
||||
@@ -8,7 +8,7 @@
|
||||
<parent>
|
||||
<groupId>com.lancedb</groupId>
|
||||
<artifactId>lancedb-parent</artifactId>
|
||||
<version>0.39.0-beta.4</version>
|
||||
<version>0.39.0-beta.6</version>
|
||||
<relativePath>../pom.xml</relativePath>
|
||||
</parent>
|
||||
|
||||
|
||||
+2
-2
@@ -6,7 +6,7 @@
|
||||
|
||||
<groupId>com.lancedb</groupId>
|
||||
<artifactId>lancedb-parent</artifactId>
|
||||
<version>0.39.0-beta.4</version>
|
||||
<version>0.39.0-beta.6</version>
|
||||
<packaging>pom</packaging>
|
||||
<name>${project.artifactId}</name>
|
||||
<description>LanceDB Java SDK Parent POM</description>
|
||||
@@ -28,7 +28,7 @@
|
||||
<properties>
|
||||
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
|
||||
<arrow.version>15.0.0</arrow.version>
|
||||
<lance-core.version>12.0.0-beta.14</lance-core.version>
|
||||
<lance-core.version>12.0.0-beta.15</lance-core.version>
|
||||
<spotless.skip>false</spotless.skip>
|
||||
<spotless.version>2.30.0</spotless.version>
|
||||
<spotless.java.googlejavaformat.version>1.7</spotless.java.googlejavaformat.version>
|
||||
|
||||
+1
-1
@@ -1,7 +1,7 @@
|
||||
[package]
|
||||
name = "lancedb-nodejs"
|
||||
edition.workspace = true
|
||||
version = "0.39.0-beta.4"
|
||||
version = "0.39.0-beta.6"
|
||||
publish = false
|
||||
license.workspace = true
|
||||
description.workspace = true
|
||||
|
||||
@@ -281,7 +281,7 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
|
||||
numIndices: 0,
|
||||
numRows: 3,
|
||||
// Full on-disk size of the two data files, footers and metadata included.
|
||||
totalBytes: 684,
|
||||
totalBytes: 550,
|
||||
});
|
||||
|
||||
// Index files count toward totalBytes too (only deletion files and
|
||||
@@ -289,7 +289,7 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
|
||||
await table.createIndex("id", { config: Index.btree() });
|
||||
const statsWithIndex = await table.stats();
|
||||
expect(statsWithIndex.numIndices).toBe(1);
|
||||
expect(statsWithIndex.totalBytes).toBeGreaterThan(684);
|
||||
expect(statsWithIndex.totalBytes).toBeGreaterThan(550);
|
||||
});
|
||||
|
||||
it("should overwrite data if asked", async () => {
|
||||
|
||||
@@ -542,10 +542,10 @@ export abstract class Table {
|
||||
* {@link Table#refreshColumn}. Declaring one therefore costs the same on a
|
||||
* large table as on an empty one.
|
||||
*
|
||||
* A refresh does not revisit rows it has already filled, so mutating an
|
||||
* input leaves the value computed at fill time; recomputing means dropping
|
||||
* the column and declaring it again. While a declaration reads a column,
|
||||
* that column cannot be renamed, retyped or dropped.
|
||||
* A refresh also recomputes the rows whose inputs changed since they were
|
||||
* computed, so a mutated input is reflected by the next refresh. While a
|
||||
* declaration reads a column, that column cannot be renamed, retyped or
|
||||
* dropped.
|
||||
*
|
||||
* On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
* server, and the refresh runs as a server job -- see
|
||||
@@ -576,10 +576,10 @@ export abstract class Table {
|
||||
/**
|
||||
* Fill the rows of a computed column that hold no value yet.
|
||||
*
|
||||
* Rows appended since the last refresh are filled by the next one; rows
|
||||
* already filled are left as they are, so the call is idempotent and does
|
||||
* not observe a mutated input. Local tables only: a remote refresh runs
|
||||
* as a server job, through {@link Table#refreshColumnAsync}.
|
||||
* Rows appended since the last refresh are filled by the next one, and
|
||||
* rows whose inputs changed since they were computed are recomputed;
|
||||
* everything else is left as it is. Local tables only: a remote refresh
|
||||
* runs as a server job, through {@link Table#refreshColumnAsync}.
|
||||
* @param {string} column The name of the computed column to fill.
|
||||
* @returns {Promise<RefreshColumnResult>} A promise that resolves to the
|
||||
* number of rows filled and the new version number of the table.
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-darwin-arm64",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["darwin"],
|
||||
"cpu": ["arm64"],
|
||||
"main": "lancedb.darwin-arm64.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-arm64-gnu",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["linux"],
|
||||
"cpu": ["arm64"],
|
||||
"main": "lancedb.linux-arm64-gnu.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-arm64-musl",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["linux"],
|
||||
"cpu": ["arm64"],
|
||||
"main": "lancedb.linux-arm64-musl.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-x64-gnu",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["linux"],
|
||||
"cpu": ["x64"],
|
||||
"main": "lancedb.linux-x64-gnu.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-x64-musl",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["linux"],
|
||||
"cpu": ["x64"],
|
||||
"main": "lancedb.linux-x64-musl.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-win32-arm64-msvc",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": [
|
||||
"win32"
|
||||
],
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-win32-x64-msvc",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["win32"],
|
||||
"cpu": ["x64"],
|
||||
"main": "lancedb.win32-x64-msvc.node",
|
||||
|
||||
+1
-1
@@ -11,7 +11,7 @@
|
||||
"ann"
|
||||
],
|
||||
"private": false,
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"main": "dist/index.js",
|
||||
"exports": {
|
||||
".": "./dist/index.js",
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "lancedb-python"
|
||||
version = "0.39.0-beta.4"
|
||||
version = "0.39.0-beta.6"
|
||||
publish = false
|
||||
edition.workspace = true
|
||||
description = "Python bindings for LanceDB"
|
||||
|
||||
@@ -37,8 +37,6 @@ from .functions import (
|
||||
UdfDefinition as UdfDefinition,
|
||||
udf as udf,
|
||||
)
|
||||
from .secrets import EnvVarSecret as EnvVarSecret
|
||||
from .secrets import SecretInfo as SecretInfo
|
||||
from .materialized_view import (
|
||||
AsyncMaterializedView,
|
||||
MaterializedView,
|
||||
|
||||
@@ -153,17 +153,6 @@ class Connection(object):
|
||||
async def get_function(self, name: str, version: str) -> str: ...
|
||||
async def list_functions(self) -> List[str]: ...
|
||||
async def drop_function(self, name: str, version: str) -> bool: ...
|
||||
async def create_secret(
|
||||
self, name: str, value: str, namespace_path: List[str]
|
||||
) -> None: ...
|
||||
async def alter_secret(
|
||||
self, name: str, value: str, namespace_path: List[str]
|
||||
) -> None: ...
|
||||
async def list_secrets(self, namespace_path: List[str]) -> List[str]: ...
|
||||
async def drop_secret(self, name: str, namespace_path: List[str]) -> None: ...
|
||||
async def describe_secret(
|
||||
self, name: str, namespace_path: List[str]
|
||||
) -> Dict[str, str]: ...
|
||||
async def list_jobs(self) -> List[JobInfo]: ...
|
||||
async def cancel_job(self, job_id: str) -> bool: ...
|
||||
async def execute_query_async(
|
||||
|
||||
+10
-206
@@ -17,7 +17,6 @@ from typing import (
|
||||
List,
|
||||
Literal,
|
||||
Optional,
|
||||
Sequence,
|
||||
Union,
|
||||
)
|
||||
from uuid import UUID
|
||||
@@ -58,12 +57,6 @@ from .materialized_view import (
|
||||
SelectArg,
|
||||
normalize_select,
|
||||
)
|
||||
from .secrets import (
|
||||
EnvVarSecret,
|
||||
SecretInfo,
|
||||
validate_namespace_path,
|
||||
validate_secret_name,
|
||||
)
|
||||
from .table import (
|
||||
AsyncTable,
|
||||
LanceTable,
|
||||
@@ -699,47 +692,15 @@ class DBConnection(EnforceOverrides):
|
||||
"""
|
||||
raise NotImplementedError("serialize is not supported for this connection type")
|
||||
|
||||
def create_function(
|
||||
self,
|
||||
definition: UdfDefinition,
|
||||
*,
|
||||
secrets: Optional[Sequence[EnvVarSecret]] = None,
|
||||
) -> FunctionVersion:
|
||||
def create_function(self, definition: UdfDefinition) -> FunctionVersion:
|
||||
"""Register a scalar Python UDF and wait for its immutable version.
|
||||
|
||||
This is the blocking counterpart of :meth:`create_function_async`.
|
||||
Local connections raise ``NotImplementedError``.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
definition : UdfDefinition
|
||||
A callable decorated with [udf][lancedb.udf].
|
||||
secrets : sequence of EnvVarSecret, optional
|
||||
One [EnvVarSecret][lancedb.secrets.EnvVarSecret] per credential the
|
||||
Function needs, each naming a Secret and the environment variable
|
||||
its value arrives in. The Function's source is unchanged by this;
|
||||
it reads the variable the way it already did.
|
||||
|
||||
Examples
|
||||
--------
|
||||
```python
|
||||
db.create_secret("openai-prod", os.environ["OPENAI_API_KEY"])
|
||||
db.create_function(
|
||||
analyze_caption,
|
||||
secrets=[
|
||||
EnvVarSecret(secret="openai-prod", env_variable="OPENAI_API_KEY")
|
||||
],
|
||||
)
|
||||
```
|
||||
"""
|
||||
return self.create_function_async(definition, secrets=secrets).wait()
|
||||
return self.create_function_async(definition).wait()
|
||||
|
||||
def create_function_async(
|
||||
self,
|
||||
definition: UdfDefinition,
|
||||
*,
|
||||
secrets: Optional[Sequence[EnvVarSecret]] = None,
|
||||
) -> Job[FunctionVersion]:
|
||||
def create_function_async(self, definition: UdfDefinition) -> Job[FunctionVersion]:
|
||||
"""Register a scalar Python UDF through the remote Function catalog.
|
||||
|
||||
Submission returns a typed job. The immutable Function version becomes
|
||||
@@ -784,70 +745,6 @@ class DBConnection(EnforceOverrides):
|
||||
"Function catalog operations are not supported for this connection type"
|
||||
)
|
||||
|
||||
def create_secret(
|
||||
self, name: str, value: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
"""Create a named Secret in this database.
|
||||
|
||||
Fails if the name is taken, so a create never silently becomes a
|
||||
rotation. Nothing reads the value back: it is bound to a Function by
|
||||
name and resolved by the service when that Function runs. Local
|
||||
connections raise ``NotImplementedError``.
|
||||
"""
|
||||
raise NotImplementedError(
|
||||
"Secret operations are not supported for this connection type"
|
||||
)
|
||||
|
||||
def alter_secret(
|
||||
self, name: str, value: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
"""Replace the credential behind an existing Secret.
|
||||
|
||||
Fails if it does not exist. Every Function bound to the Secret uses the
|
||||
new value from its next job, and no new Function version is created --
|
||||
which is how a rotation reaches columns pinned to a version registered
|
||||
before it. Local connections raise ``NotImplementedError``.
|
||||
"""
|
||||
raise NotImplementedError(
|
||||
"Secret operations are not supported for this connection type"
|
||||
)
|
||||
|
||||
def list_secrets(self, *, namespace_path: Optional[List[str]] = None) -> List[str]:
|
||||
"""The names of every Secret in this database.
|
||||
|
||||
Names only. No method returns a stored credential, by construction
|
||||
rather than by policy. Local connections raise ``NotImplementedError``.
|
||||
"""
|
||||
raise NotImplementedError(
|
||||
"Secret operations are not supported for this connection type"
|
||||
)
|
||||
|
||||
def drop_secret(
|
||||
self, name: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
"""Drop a Secret.
|
||||
|
||||
Functions bound to it fail at their next job, naming the Secret; that
|
||||
is the revocation path. The name becomes free to reuse, and a new
|
||||
Secret under it is picked up by everything still bound to that name.
|
||||
Local connections raise ``NotImplementedError``.
|
||||
"""
|
||||
raise NotImplementedError(
|
||||
"Secret operations are not supported for this connection type"
|
||||
)
|
||||
|
||||
def describe_secret(
|
||||
self, name: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> SecretInfo:
|
||||
"""What this database records about a Secret: name and timestamps.
|
||||
|
||||
Never the value -- there is no code path that could return one. Local
|
||||
connections raise ``NotImplementedError``.
|
||||
"""
|
||||
raise NotImplementedError(
|
||||
"Secret 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.
|
||||
@@ -1560,13 +1457,8 @@ class LanceDBConnection(DBConnection):
|
||||
return Job(LOOP.run(self._conn.open_job(job_id)))
|
||||
|
||||
@override
|
||||
def create_function_async(
|
||||
self,
|
||||
definition: UdfDefinition,
|
||||
*,
|
||||
secrets: Optional[Sequence[EnvVarSecret]] = None,
|
||||
) -> Job[FunctionVersion]:
|
||||
job = LOOP.run(self._conn.create_function_async(definition, secrets=secrets))
|
||||
def create_function_async(self, definition: UdfDefinition) -> Job[FunctionVersion]:
|
||||
job = LOOP.run(self._conn.create_function_async(definition))
|
||||
return Job(job)
|
||||
|
||||
@override
|
||||
@@ -1581,34 +1473,6 @@ class LanceDBConnection(DBConnection):
|
||||
def drop_function(self, name: str, *, version: str) -> bool:
|
||||
return LOOP.run(self._conn.drop_function(name, version=version))
|
||||
|
||||
@override
|
||||
def create_secret(
|
||||
self, name: str, value: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
LOOP.run(self._conn.create_secret(name, value, namespace_path=namespace_path))
|
||||
|
||||
@override
|
||||
def alter_secret(
|
||||
self, name: str, value: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
LOOP.run(self._conn.alter_secret(name, value, namespace_path=namespace_path))
|
||||
|
||||
@override
|
||||
def list_secrets(self, *, namespace_path: Optional[List[str]] = None) -> List[str]:
|
||||
return LOOP.run(self._conn.list_secrets(namespace_path=namespace_path))
|
||||
|
||||
@override
|
||||
def drop_secret(
|
||||
self, name: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
LOOP.run(self._conn.drop_secret(name, namespace_path=namespace_path))
|
||||
|
||||
@override
|
||||
def describe_secret(
|
||||
self, name: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> SecretInfo:
|
||||
return LOOP.run(self._conn.describe_secret(name, namespace_path=namespace_path))
|
||||
|
||||
@override
|
||||
def list_jobs(self) -> List[JobInfo]:
|
||||
"""List server-side jobs across the database's tables."""
|
||||
@@ -2401,23 +2265,18 @@ class AsyncConnection(object):
|
||||
return AsyncJob(await self._inner.open_job(job_id))
|
||||
|
||||
async def create_function_async(
|
||||
self,
|
||||
definition: UdfDefinition,
|
||||
*,
|
||||
secrets: Optional[Sequence[EnvVarSecret]] = None,
|
||||
self, definition: UdfDefinition
|
||||
) -> AsyncJob[FunctionVersion]:
|
||||
"""Register a scalar Python UDF through the remote Function catalog.
|
||||
|
||||
The returned typed job resolves to the immutable Function version.
|
||||
``secrets`` is a sequence of
|
||||
[EnvVarSecret][lancedb.secrets.EnvVarSecret], each naming a Secret and
|
||||
the environment variable its value arrives in. Local connections raise
|
||||
``NotImplementedError``.
|
||||
Local connections raise ``NotImplementedError``.
|
||||
"""
|
||||
if not isinstance(definition, UdfDefinition):
|
||||
raise TypeError("create_function_async requires a @udf definition")
|
||||
request = definition.bind_secrets(secrets)
|
||||
inner = await self._inner.create_function_async(request.to_canonical_json())
|
||||
inner = await self._inner.create_function_async(
|
||||
definition.registration_request.to_canonical_json()
|
||||
)
|
||||
return _typed_job(inner, FunctionVersion.from_json)
|
||||
|
||||
async def get_function(self, name: str, *, version: str) -> FunctionVersion:
|
||||
@@ -2439,61 +2298,6 @@ class AsyncConnection(object):
|
||||
"""Drop one exact immutable Function version from the remote catalog."""
|
||||
return await self._inner.drop_function(name, version)
|
||||
|
||||
async def create_secret(
|
||||
self, name: str, value: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
"""Create a named Secret in this database.
|
||||
|
||||
Fails if the name is taken, so a create never silently becomes a
|
||||
rotation. Nothing reads the value back.
|
||||
"""
|
||||
await self._inner.create_secret(
|
||||
validate_secret_name(name),
|
||||
value,
|
||||
list(validate_namespace_path(namespace_path)),
|
||||
)
|
||||
|
||||
async def alter_secret(
|
||||
self, name: str, value: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
"""Replace the credential behind an existing Secret.
|
||||
|
||||
Fails if it does not exist. Bound Functions use the new value from
|
||||
their next job, with no new Function version.
|
||||
"""
|
||||
await self._inner.alter_secret(
|
||||
validate_secret_name(name),
|
||||
value,
|
||||
list(validate_namespace_path(namespace_path)),
|
||||
)
|
||||
|
||||
async def list_secrets(
|
||||
self, *, namespace_path: Optional[List[str]] = None
|
||||
) -> List[str]:
|
||||
"""The names of every Secret in this database. Names only."""
|
||||
return await self._inner.list_secrets(
|
||||
list(validate_namespace_path(namespace_path))
|
||||
)
|
||||
|
||||
async def drop_secret(
|
||||
self, name: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
"""Drop a Secret. Bound Functions fail at their next job."""
|
||||
await self._inner.drop_secret(
|
||||
validate_secret_name(name), list(validate_namespace_path(namespace_path))
|
||||
)
|
||||
|
||||
async def describe_secret(
|
||||
self, name: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> SecretInfo:
|
||||
"""What this database records about a Secret. Never the value."""
|
||||
return SecretInfo.from_json(
|
||||
await self._inner.describe_secret(
|
||||
validate_secret_name(name),
|
||||
list(validate_namespace_path(namespace_path)),
|
||||
)
|
||||
)
|
||||
|
||||
async def list_jobs(self) -> List[JobInfo]:
|
||||
"""List server-side jobs across the database's tables."""
|
||||
return await self._inner.list_jobs()
|
||||
|
||||
@@ -4,7 +4,7 @@
|
||||
"""Canonical Function values exchanged with LanceDB Enterprise services.
|
||||
|
||||
These immutable models contain client/wire state only. Catalog persistence,
|
||||
environment bake, secret resolution, and execution are owned by Sophon.
|
||||
environment bake, and execution are owned by Sophon.
|
||||
``RefreshColumnResult`` is also the backend-neutral result of a local
|
||||
expression-backed refresh job.
|
||||
"""
|
||||
@@ -25,7 +25,7 @@ import re
|
||||
import sys
|
||||
import textwrap
|
||||
import types
|
||||
from collections.abc import Mapping, Sequence
|
||||
from collections.abc import Mapping
|
||||
from datetime import date, datetime
|
||||
from typing import (
|
||||
Annotated,
|
||||
@@ -50,7 +50,6 @@ from pydantic import (
|
||||
)
|
||||
|
||||
from .schema import is_blob_v2_field as _is_blob_v2_field
|
||||
from .secrets import EnvVarSecret
|
||||
|
||||
_Int32 = conint(strict=True, ge=-(2**31), le=2**31 - 1)
|
||||
_UInt32 = conint(strict=True, ge=0, le=2**32 - 1)
|
||||
@@ -227,19 +226,6 @@ class FunctionOutput(_OpenRemoteValue):
|
||||
fields: tuple[FunctionResultField, ...] = ()
|
||||
|
||||
|
||||
class SecretReference(_RemoteValue):
|
||||
"""Where a Secret lives, carried as its parts rather than as one string.
|
||||
|
||||
A joined id would need a delimiter, and a delimiter has to be excluded from
|
||||
every name and segment forever, agreed on by both sides, and re-agreed each
|
||||
time either grows a new way to be configured. Naming the parts settles all
|
||||
of that: nothing here is parsed, so nothing can parse two ways.
|
||||
"""
|
||||
|
||||
name: str
|
||||
namespace_path: tuple[str, ...] = ()
|
||||
|
||||
|
||||
class FunctionSignature(_RemoteValue):
|
||||
inputs: tuple[FunctionParameter, ...]
|
||||
output: FunctionOutput
|
||||
@@ -323,7 +309,6 @@ class FunctionVersion(_RemoteValue):
|
||||
runtime: PythonRuntimeSpec
|
||||
runtime_digest: str
|
||||
environment_digest: str
|
||||
secret_env_bindings: Mapping[str, SecretReference] = {}
|
||||
created_at: str
|
||||
|
||||
def __call__(self, **inputs: Any) -> FunctionApplication:
|
||||
@@ -385,18 +370,12 @@ class FunctionVersion(_RemoteValue):
|
||||
|
||||
|
||||
class FunctionRegistrationRequest(_RemoteValue):
|
||||
"""Stable remote registration envelope produced by :func:`udf`.
|
||||
|
||||
Credential values deliberately have no field here. The only secret-shaped
|
||||
thing a client sends is ``secret_env_bindings``: the name of a Secret the
|
||||
database already holds, which the remote service resolves at execution.
|
||||
"""
|
||||
"""Stable remote registration envelope produced by :func:`udf`."""
|
||||
|
||||
name: str
|
||||
artifact: FunctionArtifactRequest
|
||||
signature: FunctionSignature
|
||||
runtime: PythonRuntimeSpec
|
||||
secret_env_bindings: Mapping[str, SecretReference] = {}
|
||||
|
||||
|
||||
class FunctionVersionRef(_OpenRemoteValue):
|
||||
@@ -545,8 +524,6 @@ class RefreshColumnResult(_RemoteValue):
|
||||
|
||||
|
||||
_FUNCTION_NAME = re.compile(r"^[A-Za-z_][A-Za-z0-9_.-]*$")
|
||||
_DECLARED_SECRET = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$")
|
||||
|
||||
_FUNCTION_BLOB_V2_TYPE = "blob_v2"
|
||||
_ARROW_EXTENSION_NAME_KEY = "ARROW:extension:name"
|
||||
_BLOB_V2_EXTENSION_NAME = "lance.blob.v2"
|
||||
@@ -1288,63 +1265,9 @@ class UdfDefinition:
|
||||
|
||||
@property
|
||||
def registration_request(self) -> FunctionRegistrationRequest:
|
||||
"""The immutable request sent by ``create_function_async``.
|
||||
|
||||
Carries no secret bindings. Binding is a registration-time decision,
|
||||
so a Function bound to Secrets is registered through :meth:`bind_secrets`,
|
||||
which is what ``create_function`` calls.
|
||||
"""
|
||||
"""The immutable request sent by ``create_function_async``."""
|
||||
return self._request
|
||||
|
||||
def bind_secrets(
|
||||
self, secrets: Optional[Sequence[EnvVarSecret]]
|
||||
) -> FunctionRegistrationRequest:
|
||||
"""The registration request for this definition bound to ``secrets``.
|
||||
|
||||
Binding does not change the Function's source: each
|
||||
[EnvVarSecret][lancedb.secrets.EnvVarSecret] names a Secret and the
|
||||
environment variable its value should arrive in, and the Function reads
|
||||
that variable the way it already did. Whether the named Secrets exist is
|
||||
the server's answer, not this one.
|
||||
"""
|
||||
bindings = () if secrets is None else tuple(secrets)
|
||||
wrong_type = [
|
||||
binding for binding in bindings if not isinstance(binding, EnvVarSecret)
|
||||
]
|
||||
if wrong_type:
|
||||
kinds = sorted({type(binding).__name__ for binding in wrong_type})
|
||||
raise TypeError(
|
||||
f"Function secrets must be EnvVarSecret values, not {kinds!r}; a "
|
||||
"credential value is never sent to this API"
|
||||
)
|
||||
variables = [binding.env_variable for binding in bindings]
|
||||
duplicates = sorted({name for name in variables if variables.count(name) > 1})
|
||||
if duplicates:
|
||||
raise ValueError(
|
||||
"a Function binds each environment variable once; duplicated: "
|
||||
f"{duplicates!r}"
|
||||
)
|
||||
# `env` is ordinary configuration carried in the definition, so a name in
|
||||
# both would have a value visible in the Function's record and a value
|
||||
# that is not. Refuse rather than pick.
|
||||
environment = self._request.runtime.env or {}
|
||||
overlap = sorted(set(environment) & set(variables))
|
||||
if overlap:
|
||||
raise ValueError(
|
||||
f"Function env and secret bindings must be disjoint: {overlap!r}"
|
||||
)
|
||||
if not bindings:
|
||||
return self._request
|
||||
# The binding records the full id -- path plus name -- because that is
|
||||
# what the service resolves. At the root it is the bare name.
|
||||
resolved = {
|
||||
binding.env_variable: SecretReference(
|
||||
name=binding.secret, namespace_path=tuple(binding.namespace_path)
|
||||
)
|
||||
for binding in bindings
|
||||
}
|
||||
return self._request._copy(update={"secret_env_bindings": resolved})
|
||||
|
||||
def __call__(self, *args, **kwargs):
|
||||
return self._function(*args, **kwargs)
|
||||
|
||||
@@ -1409,9 +1332,7 @@ def udf(
|
||||
conda_channels : sequence of str, optional
|
||||
Conda channels in priority order; requires ``conda``.
|
||||
env : mapping of str to str, optional
|
||||
Environment variables included in the Function definition. Not for
|
||||
credentials -- these are ordinary configuration, stored with the
|
||||
Function and visible wherever it is.
|
||||
Environment variables included in the Function definition.
|
||||
python_version : str, optional
|
||||
Remote Python major/minor version. Defaults to the client version.
|
||||
gpu : bool, default False
|
||||
|
||||
@@ -7,16 +7,7 @@ import json
|
||||
import logging
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
import sys
|
||||
from typing import (
|
||||
TYPE_CHECKING,
|
||||
Any,
|
||||
Dict,
|
||||
Iterable,
|
||||
List,
|
||||
Optional,
|
||||
Sequence,
|
||||
Union,
|
||||
)
|
||||
from typing import TYPE_CHECKING, Any, Dict, Iterable, List, Optional, Union
|
||||
from urllib.parse import urlparse
|
||||
from uuid import UUID
|
||||
import warnings
|
||||
@@ -38,7 +29,6 @@ from ..job import AsyncJob, Job
|
||||
from ..sql import Query as SqlQuery
|
||||
from ..sql import QueryDescription
|
||||
from ..materialized_view import MaterializedView, SelectArg
|
||||
from ..secrets import EnvVarSecret, SecretInfo
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from .._lancedb import JobInfo
|
||||
@@ -756,14 +746,8 @@ class RemoteDBConnection(DBConnection):
|
||||
return Job(LOOP.run(self._conn.open_job(job_id)))
|
||||
|
||||
@override
|
||||
def create_function_async(
|
||||
self,
|
||||
definition: UdfDefinition,
|
||||
*,
|
||||
secrets: Optional[Sequence[EnvVarSecret]] = None,
|
||||
) -> Job[FunctionVersion]:
|
||||
job = LOOP.run(self._conn.create_function_async(definition, secrets=secrets))
|
||||
return Job(job)
|
||||
def create_function_async(self, definition: UdfDefinition) -> Job[FunctionVersion]:
|
||||
return Job(LOOP.run(self._conn.create_function_async(definition)))
|
||||
|
||||
@override
|
||||
def get_function(self, name: str, *, version: str) -> FunctionVersion:
|
||||
@@ -777,34 +761,6 @@ class RemoteDBConnection(DBConnection):
|
||||
def drop_function(self, name: str, *, version: str) -> bool:
|
||||
return LOOP.run(self._conn.drop_function(name, version=version))
|
||||
|
||||
@override
|
||||
def create_secret(
|
||||
self, name: str, value: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
LOOP.run(self._conn.create_secret(name, value, namespace_path=namespace_path))
|
||||
|
||||
@override
|
||||
def alter_secret(
|
||||
self, name: str, value: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
LOOP.run(self._conn.alter_secret(name, value, namespace_path=namespace_path))
|
||||
|
||||
@override
|
||||
def describe_secret(
|
||||
self, name: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> SecretInfo:
|
||||
return LOOP.run(self._conn.describe_secret(name, namespace_path=namespace_path))
|
||||
|
||||
@override
|
||||
def list_secrets(self, *, namespace_path: Optional[List[str]] = None) -> List[str]:
|
||||
return LOOP.run(self._conn.list_secrets(namespace_path=namespace_path))
|
||||
|
||||
@override
|
||||
def drop_secret(
|
||||
self, name: str, *, namespace_path: Optional[List[str]] = None
|
||||
) -> None:
|
||||
LOOP.run(self._conn.drop_secret(name, namespace_path=namespace_path))
|
||||
|
||||
@override
|
||||
def list_jobs(self) -> List["JobInfo"]:
|
||||
"""List server-side jobs across the database's tables."""
|
||||
|
||||
@@ -1,209 +0,0 @@
|
||||
# SPDX-License-Identifier: Apache-2.0
|
||||
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
"""Named Secrets, and the bindings that deliver them to Functions.
|
||||
|
||||
A Secret is a database-scoped named credential. Nothing in this module holds a
|
||||
value: :class:`EnvVarSecret` names one and says which environment variable it
|
||||
should arrive in, and the value is resolved by the remote service when a
|
||||
Function bound to it runs. No API returns a stored credential, by construction
|
||||
rather than by policy -- there is no code path that could.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import re
|
||||
|
||||
# Periods are legal inside a name: excluding them would put Secrets out of reach
|
||||
# of any namespace whose name carries one, which LanceDB namespaces allow. Not
|
||||
# at either end, which is RFC 1123's shape and Kubernetes' rule for object
|
||||
# names: it rules out `.` and `..` and anything that reads as a hidden file or a
|
||||
# path fragment in a listing. Matches the service, which rejects the same shapes.
|
||||
_SECRET_NAME = re.compile(r"^[A-Za-z0-9]([A-Za-z0-9_.-]{0,253}[A-Za-z0-9])?$")
|
||||
_ENV_VARIABLE = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$")
|
||||
|
||||
|
||||
def validate_secret_name(name: str) -> str:
|
||||
"""Check a Secret name locally and return it unchanged."""
|
||||
if not isinstance(name, str):
|
||||
raise TypeError(f"Secret name must be a string, not {type(name).__name__}")
|
||||
if not _SECRET_NAME.fullmatch(name):
|
||||
raise ValueError(f"invalid Secret name: {name!r}")
|
||||
return name
|
||||
|
||||
|
||||
def validate_namespace_path(namespace_path=None):
|
||||
"""Check a namespace path locally and return it as a tuple.
|
||||
|
||||
``None`` and ``[]`` both mean the root namespace. Segments follow the same
|
||||
rule as Secret names: a binding carries the path and the name as separate
|
||||
fields, so neither is ever parsed out of the other.
|
||||
"""
|
||||
if namespace_path is None:
|
||||
return ()
|
||||
if isinstance(namespace_path, str):
|
||||
raise TypeError(
|
||||
"namespace_path must be a list of segments, not a string; "
|
||||
f"did you mean [{namespace_path!r}]?"
|
||||
)
|
||||
segments = tuple(namespace_path)
|
||||
for segment in segments:
|
||||
if not isinstance(segment, str):
|
||||
raise TypeError(
|
||||
f"namespace path segment must be a string, not {type(segment).__name__}"
|
||||
)
|
||||
if not _SECRET_NAME.fullmatch(segment):
|
||||
raise ValueError(f"invalid namespace path segment: {segment!r}")
|
||||
return segments
|
||||
|
||||
|
||||
def validate_env_variable(name: str) -> str:
|
||||
"""Check an environment variable name locally and return it unchanged."""
|
||||
if not isinstance(name, str):
|
||||
raise TypeError(
|
||||
f"environment variable name must be a string, not {type(name).__name__}"
|
||||
)
|
||||
if not _ENV_VARIABLE.fullmatch(name):
|
||||
raise ValueError(f"invalid environment variable name: {name!r}")
|
||||
return name
|
||||
|
||||
|
||||
class EnvVarSecret:
|
||||
"""A Secret bound to the environment variable a Function's library reads.
|
||||
|
||||
Pass these in the ``secrets`` sequence of
|
||||
[DBConnection.create_function][lancedb.db.DBConnection.create_function]. The
|
||||
Function's source is unchanged by binding: it reads ``OPENAI_API_KEY`` the
|
||||
way it always did, and the binding is what puts a value there.
|
||||
|
||||
This is a local value. Constructing it contacts no server, so it always
|
||||
succeeds and says nothing about whether the Secret exists; that is checked
|
||||
at registration, where a mistyped Secret name surfaces as a clear "does not
|
||||
exist" naming both the Secret and the variable bound to it. A mistyped
|
||||
*variable* name cannot be caught anywhere -- nothing knows which variables a
|
||||
Function reads -- so it surfaces on the first rows instead.
|
||||
|
||||
The type exists so a credential cannot be passed by accident. A bare string
|
||||
in the same position is a plausible-looking mistake with the opposite
|
||||
meaning, and it reads identically in a diff.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
secret : str
|
||||
The Secret's database-scoped name.
|
||||
env_variable : str
|
||||
The environment variable the Function reads it from.
|
||||
|
||||
Examples
|
||||
--------
|
||||
>>> from lancedb import EnvVarSecret
|
||||
>>> binding = EnvVarSecret(secret="openai-prod", env_variable="OPENAI_API_KEY")
|
||||
>>> binding.secret, binding.env_variable
|
||||
('openai-prod', 'OPENAI_API_KEY')
|
||||
"""
|
||||
|
||||
__slots__ = ("_secret", "_env_variable", "_namespace_path")
|
||||
|
||||
def __init__(self, secret: str, env_variable: str, *, namespace_path=None):
|
||||
self._secret = validate_secret_name(secret)
|
||||
self._env_variable = validate_env_variable(env_variable)
|
||||
self._namespace_path = validate_namespace_path(namespace_path)
|
||||
|
||||
@property
|
||||
def secret(self) -> str:
|
||||
"""The Secret's database-scoped name."""
|
||||
return self._secret
|
||||
|
||||
@property
|
||||
def env_variable(self) -> str:
|
||||
"""The environment variable the value is delivered in."""
|
||||
return self._env_variable
|
||||
|
||||
@property
|
||||
def namespace_path(self):
|
||||
"""The namespace path the Secret is addressed within, root when empty."""
|
||||
return list(self._namespace_path)
|
||||
|
||||
def __repr__(self) -> str:
|
||||
path = (
|
||||
f", namespace_path={list(self._namespace_path)!r}"
|
||||
if self._namespace_path
|
||||
else ""
|
||||
)
|
||||
return (
|
||||
f"EnvVarSecret(secret={self._secret!r}, "
|
||||
f"env_variable={self._env_variable!r}{path})"
|
||||
)
|
||||
|
||||
def __eq__(self, other: object) -> bool:
|
||||
return (
|
||||
isinstance(other, EnvVarSecret)
|
||||
and other._secret == self._secret
|
||||
and other._env_variable == self._env_variable
|
||||
and other._namespace_path == self._namespace_path
|
||||
)
|
||||
|
||||
def __hash__(self) -> int:
|
||||
return hash(
|
||||
(EnvVarSecret, self._secret, self._env_variable, self._namespace_path)
|
||||
)
|
||||
|
||||
|
||||
class SecretInfo:
|
||||
"""What a database records about a Secret. Never its value.
|
||||
|
||||
Returned by
|
||||
[DBConnection.describe_secret][lancedb.db.DBConnection.describe_secret].
|
||||
"""
|
||||
|
||||
__slots__ = ("_name", "_created_at", "_updated_at")
|
||||
|
||||
def __init__(self, name: str, created_at: str, updated_at: str):
|
||||
self._name = name
|
||||
self._created_at = created_at
|
||||
self._updated_at = updated_at
|
||||
|
||||
@property
|
||||
def name(self) -> str:
|
||||
"""The Secret's database-scoped name."""
|
||||
return self._name
|
||||
|
||||
@property
|
||||
def created_at(self) -> str:
|
||||
"""When the Secret was created, as an RFC 3339 timestamp."""
|
||||
return self._created_at
|
||||
|
||||
@property
|
||||
def updated_at(self) -> str:
|
||||
"""When the Secret's value was last rotated, as an RFC 3339 timestamp."""
|
||||
return self._updated_at
|
||||
|
||||
@classmethod
|
||||
def from_json(cls, value: dict) -> "SecretInfo":
|
||||
return cls(
|
||||
name=value["name"],
|
||||
created_at=value["created_at"],
|
||||
updated_at=value["updated_at"],
|
||||
)
|
||||
|
||||
def __repr__(self) -> str:
|
||||
return (
|
||||
f"SecretInfo(name={self._name!r}, created_at={self._created_at!r}, "
|
||||
f"updated_at={self._updated_at!r})"
|
||||
)
|
||||
|
||||
def __eq__(self, other: object) -> bool:
|
||||
return (
|
||||
isinstance(other, SecretInfo)
|
||||
and other._name == self._name
|
||||
and other._created_at == self._created_at
|
||||
and other._updated_at == self._updated_at
|
||||
)
|
||||
|
||||
|
||||
__all__ = [
|
||||
"EnvVarSecret",
|
||||
"SecretInfo",
|
||||
"validate_env_variable",
|
||||
"validate_secret_name",
|
||||
]
|
||||
@@ -2188,10 +2188,10 @@ class Table(ABC):
|
||||
Declaring one therefore costs the same on a large table as on an
|
||||
empty one.
|
||||
|
||||
A refresh does not revisit rows it has already filled, so mutating
|
||||
an input leaves the value computed at fill time; recomputing means
|
||||
dropping the column and declaring it again. While a declaration
|
||||
reads a column, that column cannot be renamed, retyped or dropped.
|
||||
A refresh also recomputes the rows whose inputs changed since they
|
||||
were computed, so a mutated input is reflected by the next refresh.
|
||||
While a declaration reads a column, that column cannot be renamed,
|
||||
retyped or dropped.
|
||||
|
||||
On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
server, and the refresh runs as a server job -- see
|
||||
@@ -2211,7 +2211,7 @@ class Table(ABC):
|
||||
>>> table.add_columns(computed={"doubled": "x * 2"})
|
||||
AddColumnsResult(version=2)
|
||||
>>> table.refresh_column("doubled")
|
||||
RefreshColumnResult(rows_filled=2, version=3)
|
||||
RefreshColumnResult(rows_filled=2, version=4)
|
||||
>>> table.to_arrow().sort_by("x").to_pandas()
|
||||
x doubled
|
||||
0 1 2
|
||||
@@ -2225,8 +2225,8 @@ class Table(ABC):
|
||||
|
||||
Declared with ``add_columns(computed=...)``, a column starts empty and
|
||||
gets its values here. Rows appended since the last refresh are filled
|
||||
by the next one; rows already filled are left as they are, so the call
|
||||
is idempotent and does not observe a mutated input.
|
||||
by the next one, and rows whose inputs changed since they were computed
|
||||
are recomputed; everything else is left as it is.
|
||||
|
||||
Local tables only: a remote refresh runs as a server job, through
|
||||
[`refresh_column_async`][lancedb.table.Table.refresh_column_async].
|
||||
@@ -4318,13 +4318,14 @@ class LanceTable(Table):
|
||||
return LOOP.run(self._table.add_columns(transforms, computed=computed))
|
||||
|
||||
def refresh_column(self, column: str) -> "RefreshColumnResult":
|
||||
"""Fill a computed column's unfilled rows. See
|
||||
"""Fill a computed column's unfilled rows and recompute those whose
|
||||
inputs changed. See
|
||||
[`AsyncTable.refresh_column`][lancedb.AsyncTable.refresh_column]."""
|
||||
return LOOP.run(self._table.refresh_column(column))
|
||||
|
||||
def refresh_column_async(self, column: str) -> Job[RefreshColumnJobResult]:
|
||||
"""Fill a computed column's unfilled rows, returning a handle to the
|
||||
refresh job. See
|
||||
"""Fill a computed column's unfilled rows and recompute those whose
|
||||
inputs changed, returning a handle to the refresh job. See
|
||||
[`Table.refresh_column_async`][lancedb.table.Table.refresh_column_async].
|
||||
"""
|
||||
return Job(LOOP.run(self._table.refresh_column_async(column)))
|
||||
@@ -6312,10 +6313,10 @@ class AsyncTable:
|
||||
them from
|
||||
[`refresh_column`][lancedb.table.AsyncTable.refresh_column].
|
||||
|
||||
A refresh does not revisit rows it has already filled, so mutating
|
||||
an input leaves the value computed at fill time. While a
|
||||
declaration reads a column, that column cannot be renamed, retyped
|
||||
or dropped.
|
||||
A refresh also recomputes the rows whose inputs changed since they
|
||||
were computed, so a mutated input is reflected by the next refresh.
|
||||
While a declaration reads a column, that column cannot be renamed,
|
||||
retyped or dropped.
|
||||
|
||||
On LanceDB Cloud and Enterprise the expression is planned by
|
||||
the server. Cannot be combined with ``transforms``.
|
||||
@@ -6377,8 +6378,8 @@ class AsyncTable:
|
||||
|
||||
Declared with ``add_columns(computed=...)``, a column starts empty and
|
||||
gets its values here. Rows appended since the last refresh are filled
|
||||
by the next one; rows already filled are left as they are, so the call
|
||||
is idempotent and does not observe a mutated input.
|
||||
by the next one, and rows whose inputs changed since they were computed
|
||||
are recomputed; everything else is left as it is.
|
||||
|
||||
Local tables only: a remote refresh runs as a server job, through
|
||||
[`refresh_column_async`][lancedb.table.Table.refresh_column_async].
|
||||
|
||||
@@ -297,7 +297,10 @@ def test_blob_v2_projection_sources_use_typed_column_name():
|
||||
|
||||
|
||||
def _legacy_v1_table(name):
|
||||
db = lancedb.connect("memory:///")
|
||||
# Legacy v1 blob columns are only writable at file version <= 2.1.
|
||||
db = lancedb.connect(
|
||||
"memory:///", storage_options={"new_table_data_storage_version": "2.1"}
|
||||
)
|
||||
schema = pa.schema(
|
||||
[
|
||||
pa.field("id", pa.int64()),
|
||||
|
||||
@@ -14,7 +14,6 @@ from lancedb.functions import (
|
||||
FunctionVersion,
|
||||
PythonRuntimeSpec,
|
||||
RefreshColumnResult,
|
||||
SecretReference,
|
||||
)
|
||||
from lancedb.table import AsyncTable
|
||||
|
||||
@@ -38,22 +37,6 @@ def job_result(name: str) -> dict:
|
||||
return json.loads(fixture(name))["result"]
|
||||
|
||||
|
||||
def assert_no_secret_values(value):
|
||||
"""No client value models a resolved credential, at any nesting depth."""
|
||||
if isinstance(value, dict):
|
||||
for key, child in value.items():
|
||||
assert key not in {
|
||||
"secret_value",
|
||||
"secret_values",
|
||||
"resolved_secret",
|
||||
"resolved_secrets",
|
||||
}
|
||||
assert_no_secret_values(child)
|
||||
elif isinstance(value, list):
|
||||
for child in value:
|
||||
assert_no_secret_values(child)
|
||||
|
||||
|
||||
def test_public_function_values_are_in_api_reference():
|
||||
docs = Path(__file__).parents[3] / "docs" / "src" / "python" / "python.md"
|
||||
rendered = docs.read_text()
|
||||
@@ -111,9 +94,6 @@ def test_function_version_identity_is_immutable_and_exact():
|
||||
version = FunctionVersion.from_json(json.dumps(value))
|
||||
assert version.name == "embed"
|
||||
assert version.version == "fv_01K3EXACT"
|
||||
assert dict(version.secret_env_bindings) == {
|
||||
"HF_TOKEN": SecretReference(name="hf-prod")
|
||||
}
|
||||
|
||||
with pytest.raises((TypeError, ValueError)):
|
||||
version.version = "fv_changed"
|
||||
@@ -296,25 +276,6 @@ def test_refresh_result_rejects_non_u64_values(field):
|
||||
RefreshColumnResult.from_json(json.dumps(value))
|
||||
|
||||
|
||||
def test_canonical_client_values_carry_bindings_and_no_credentials():
|
||||
"""A binding names a Secret; the credential behind it has no client field."""
|
||||
version = FunctionVersion.from_json(
|
||||
json.dumps(job_result("remote_function_job.json"))
|
||||
)
|
||||
canonical = json.loads(version.to_canonical_json())
|
||||
assert canonical["secret_env_bindings"] == {"HF_TOKEN": {"name": "hf-prod"}}
|
||||
assert_no_secret_values(canonical)
|
||||
|
||||
|
||||
def test_a_version_without_bindings_keeps_the_original_wire_shape():
|
||||
"""Every Function registered before Secrets existed serializes unchanged."""
|
||||
value = job_result("remote_function_job.json")
|
||||
del value["secret_env_bindings"]
|
||||
version = FunctionVersion.from_json(json.dumps(value))
|
||||
assert dict(version.secret_env_bindings) == {}
|
||||
assert "secret_env_bindings" not in json.loads(version.to_canonical_json())
|
||||
|
||||
|
||||
class _FunctionDeclarationInner:
|
||||
def __init__(self):
|
||||
self.calls = []
|
||||
|
||||
@@ -11,7 +11,6 @@ import types
|
||||
from datetime import date
|
||||
import http.server
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import subprocess
|
||||
import sys
|
||||
@@ -22,16 +21,13 @@ import pyarrow as pa
|
||||
import pytest
|
||||
|
||||
import lancedb
|
||||
from lancedb.background_loop import LOOP
|
||||
from lancedb.functions import (
|
||||
PythonRuntimeSpec,
|
||||
SecretReference,
|
||||
UdfDefinition,
|
||||
_canonical_arrow_type,
|
||||
_GRAMMAR_PRIMITIVES,
|
||||
udf,
|
||||
)
|
||||
from lancedb.secrets import EnvVarSecret
|
||||
|
||||
THRESHOLD = 20
|
||||
_CACHE = None
|
||||
@@ -57,15 +53,6 @@ def normalize_score(value: float) -> float:
|
||||
return value / 100.0
|
||||
|
||||
|
||||
@udf(
|
||||
pip=["openai==3.7.0"],
|
||||
env={"MODE": "test"},
|
||||
python_version="3.12",
|
||||
)
|
||||
def analyze_caption(caption: str) -> str:
|
||||
return caption.strip()
|
||||
|
||||
|
||||
def test_scalar_udf_matches_shared_registration_golden_and_remains_callable():
|
||||
assert isinstance(normalize_score, UdfDefinition)
|
||||
assert normalize_score(25.0) == 0.25
|
||||
@@ -82,289 +69,6 @@ def test_scalar_udf_matches_shared_registration_golden_and_remains_callable():
|
||||
}
|
||||
|
||||
|
||||
def test_secret_bound_udf_matches_its_shared_registration_golden():
|
||||
assert analyze_caption(" hello ") == "hello"
|
||||
bound = analyze_caption.bind_secrets(
|
||||
[EnvVarSecret(secret="openai-prod", env_variable="OPENAI_API_KEY")]
|
||||
)
|
||||
assert (
|
||||
bound.to_canonical_json()
|
||||
== (FIXTURES / "remote_function_secret_registration_request.canonical.json")
|
||||
.read_text()
|
||||
.strip()
|
||||
)
|
||||
|
||||
|
||||
def test_a_namespaced_binding_records_the_path_and_the_name():
|
||||
"""A binding names the parts, so nothing has to be parsed back out.
|
||||
|
||||
A root binding carries no path at all, which is what keeps its wire shape
|
||||
identical to one written before namespaces existed.
|
||||
"""
|
||||
root = EnvVarSecret(secret="openai-prod", env_variable="OPENAI_API_KEY")
|
||||
assert root.namespace_path == []
|
||||
|
||||
nested = EnvVarSecret(
|
||||
secret="openai-prod",
|
||||
env_variable="OPENAI_API_KEY",
|
||||
namespace_path=["prod", "vision"],
|
||||
)
|
||||
assert nested.namespace_path == ["prod", "vision"]
|
||||
assert nested != root
|
||||
|
||||
bound = analyze_caption.bind_secrets([nested])
|
||||
assert bound.secret_env_bindings == {
|
||||
"OPENAI_API_KEY": SecretReference(
|
||||
name="openai-prod", namespace_path=("prod", "vision")
|
||||
)
|
||||
}
|
||||
|
||||
at_root = analyze_caption.bind_secrets([root])
|
||||
assert at_root.secret_env_bindings == {
|
||||
"OPENAI_API_KEY": SecretReference(name="openai-prod")
|
||||
}
|
||||
canonical = json.loads(at_root.to_canonical_json())
|
||||
assert canonical["secret_env_bindings"] == {
|
||||
"OPENAI_API_KEY": {"name": "openai-prod"}
|
||||
}
|
||||
|
||||
|
||||
def test_a_namespace_path_is_validated_locally():
|
||||
# The charset is the service's, not a delimiter's: a reference is never
|
||||
# joined, so a segment cannot make anything parse two ways.
|
||||
with pytest.raises(ValueError):
|
||||
EnvVarSecret(
|
||||
secret="openai-prod", env_variable="K", namespace_path=["with$delim"]
|
||||
)
|
||||
with pytest.raises(ValueError):
|
||||
EnvVarSecret(secret="openai-prod", env_variable="K", namespace_path=["a/b"])
|
||||
# A bare string is a plausible mistake with the wrong meaning.
|
||||
with pytest.raises(TypeError):
|
||||
EnvVarSecret(secret="openai-prod", env_variable="K", namespace_path="prod")
|
||||
|
||||
|
||||
def test_an_unbound_request_carries_no_binding_at_all():
|
||||
"""Binding is a registration-time decision, so the definition holds none.
|
||||
|
||||
The decorator declares nothing about secrets, which is what makes the PRD's
|
||||
claim true: a Function's source and its registration request are identical
|
||||
whether or not a credential is later bound to it.
|
||||
"""
|
||||
unbound = json.loads(analyze_caption.registration_request.to_canonical_json())
|
||||
assert "secret_env_bindings" not in unbound
|
||||
assert "OPENAI_API_KEY" not in json.dumps(unbound)
|
||||
|
||||
|
||||
def test_binding_a_secret_leaves_the_packaged_artifact_untouched():
|
||||
"""The artifact is source bytes and nothing else, with or without secrets."""
|
||||
bound = analyze_caption.bind_secrets(
|
||||
[EnvVarSecret(secret="openai-prod", env_variable="OPENAI_API_KEY")]
|
||||
)
|
||||
assert bound.artifact == analyze_caption.registration_request.artifact
|
||||
assert bound.artifact.digest == analyze_caption.registration_request.artifact.digest
|
||||
|
||||
|
||||
def test_a_function_declaring_no_secret_is_registered_exactly_as_before():
|
||||
"""The compatibility claim: nothing about the no-secret path moves."""
|
||||
assert (
|
||||
normalize_score.bind_secrets(None).to_canonical_json()
|
||||
== normalize_score.registration_request.to_canonical_json()
|
||||
)
|
||||
assert (
|
||||
"secret_env_bindings"
|
||||
not in normalize_score.registration_request.to_canonical_json()
|
||||
)
|
||||
|
||||
|
||||
def test_a_function_binds_each_variable_once():
|
||||
with pytest.raises(ValueError, match="binds each environment variable once"):
|
||||
analyze_caption.bind_secrets(
|
||||
[
|
||||
EnvVarSecret(secret="openai-prod", env_variable="OPENAI_API_KEY"),
|
||||
EnvVarSecret(secret="openai-staging", env_variable="OPENAI_API_KEY"),
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
def test_bindings_may_not_collide_with_plain_configuration():
|
||||
"""`env` is stored with the Function; a Secret is not. Refuse, do not pick."""
|
||||
with pytest.raises(ValueError, match="must be disjoint"):
|
||||
analyze_caption.bind_secrets(
|
||||
[EnvVarSecret(secret="mode-prod", env_variable="MODE")]
|
||||
)
|
||||
|
||||
|
||||
def test_a_function_binds_at_most_sixteen_secrets():
|
||||
"""The cap lives in Rust, so no language surface can be talked past it.
|
||||
|
||||
Registering through the typed API and hand-rolling the request envelope
|
||||
reach the same boundary, and neither reaches the wire.
|
||||
"""
|
||||
bindings = [
|
||||
EnvVarSecret(secret=f"secret-{index}", env_variable=f"TOKEN_{index}")
|
||||
for index in range(17)
|
||||
]
|
||||
with _mock_remote_function_catalog() as (host, state):
|
||||
db = lancedb.connect(
|
||||
"db://dev",
|
||||
api_key="fake",
|
||||
host_override=host,
|
||||
client_config={"retry_config": {"retries": 0}},
|
||||
)
|
||||
with pytest.raises(ValueError, match="at most 16 secrets"):
|
||||
db.create_function(normalize_score, secrets=bindings)
|
||||
|
||||
envelope = json.loads(normalize_score.registration_request.to_canonical_json())
|
||||
envelope["secret_env_bindings"] = {
|
||||
f"TOKEN_{index}": {"name": f"secret-{index}"} for index in range(17)
|
||||
}
|
||||
|
||||
async def submit_envelope():
|
||||
return await db._conn._inner.create_function_async(json.dumps(envelope))
|
||||
|
||||
with pytest.raises(ValueError, match="at most 16 secrets"):
|
||||
LOOP.run(submit_envelope())
|
||||
|
||||
assert state["requests"] == []
|
||||
|
||||
|
||||
def test_binding_names_are_validated_below_the_python_api():
|
||||
"""The low-level entry point reaches the same validator the typed API does.
|
||||
|
||||
Registration envelopes can be hand-rolled past ``bind_secrets``, so the
|
||||
grammar and the disjointness rule live in Rust, above the backend.
|
||||
"""
|
||||
with _mock_remote_function_catalog() as (host, state):
|
||||
db = lancedb.connect(
|
||||
"db://dev",
|
||||
api_key="fake",
|
||||
host_override=host,
|
||||
client_config={"retry_config": {"retries": 0}},
|
||||
)
|
||||
envelope = json.loads(analyze_caption.registration_request.to_canonical_json())
|
||||
envelope["secret_env_bindings"] = {
|
||||
"BAD=NAME": {"name": "openai-prod"},
|
||||
"TOKEN_0": {"name": "secret-0"},
|
||||
}
|
||||
envelope["runtime"]["env"]["TOKEN_0"] = "public"
|
||||
|
||||
async def submit_envelope():
|
||||
return await db._conn._inner.create_function_async(json.dumps(envelope))
|
||||
|
||||
with pytest.raises(ValueError, match="portable"):
|
||||
LOOP.run(submit_envelope())
|
||||
|
||||
assert state["requests"] == []
|
||||
|
||||
|
||||
_SECRET_DEBUG_LOG_SOURCE = """
|
||||
import http.server
|
||||
import json
|
||||
import threading
|
||||
|
||||
import lancedb
|
||||
|
||||
|
||||
class Handler(http.server.BaseHTTPRequestHandler):
|
||||
def log_message(self, *args):
|
||||
pass
|
||||
|
||||
def do_POST(self):
|
||||
self.rfile.read(int(self.headers.get("Content-Length", "0")))
|
||||
payload = json.dumps({}).encode()
|
||||
self.send_response(200)
|
||||
self.send_header("Content-Type", "application/json")
|
||||
self.send_header("Content-Length", str(len(payload)))
|
||||
self.end_headers()
|
||||
self.wfile.write(payload)
|
||||
|
||||
|
||||
server = http.server.ThreadingHTTPServer(("localhost", 0), Handler)
|
||||
threading.Thread(target=server.serve_forever, daemon=True).start()
|
||||
try:
|
||||
db = lancedb.connect(
|
||||
"db://dev",
|
||||
api_key="API_KEY_SENTINEL",
|
||||
host_override="http://localhost:%d" % server.server_address[1],
|
||||
client_config={"retry_config": {"retries": 0}},
|
||||
)
|
||||
db.create_secret("openai-prod", "SECRET_VALUE_SENTINEL")
|
||||
finally:
|
||||
server.shutdown()
|
||||
"""
|
||||
|
||||
|
||||
def test_a_credential_never_reaches_a_debug_log(tmp_path):
|
||||
"""The logger sees the serialized body, so no value-side redaction reaches it.
|
||||
|
||||
Runs in a subprocess because the Rust logger reads ``LANCEDB_LOG`` once, at
|
||||
import.
|
||||
"""
|
||||
script = tmp_path / "write_secret.py"
|
||||
script.write_text(_SECRET_DEBUG_LOG_SOURCE)
|
||||
|
||||
result = subprocess.run(
|
||||
[sys.executable, str(script)],
|
||||
check=True,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
env={**os.environ, "LANCEDB_LOG": "debug"},
|
||||
)
|
||||
output = result.stdout + result.stderr
|
||||
|
||||
# Without this the test passes when debug logging is simply off.
|
||||
assert "Sending request_id=" in output, output
|
||||
assert "SECRET_VALUE_SENTINEL" not in output
|
||||
assert "API_KEY_SENTINEL" not in output
|
||||
|
||||
|
||||
def test_a_credential_value_is_rejected_in_the_binding_position():
|
||||
"""The one mistake the typed binding exists to stop."""
|
||||
with pytest.raises(TypeError, match="EnvVarSecret"):
|
||||
analyze_caption.bind_secrets(["sk-live-0001"])
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("secret", "variable", "message"),
|
||||
[
|
||||
("openai-prod", "not-a-var", "invalid environment variable name"),
|
||||
("openai-prod", "API-TOKEN", "invalid environment variable name"),
|
||||
("not a name", "API_TOKEN", "invalid Secret name"),
|
||||
("openai$prod", "API_TOKEN", "invalid Secret name"),
|
||||
],
|
||||
)
|
||||
def test_a_binding_validates_both_names_locally(secret, variable, message):
|
||||
with pytest.raises(ValueError, match=message):
|
||||
EnvVarSecret(secret=secret, env_variable=variable)
|
||||
|
||||
|
||||
def test_a_period_is_legal_inside_a_name_and_not_at_its_edges():
|
||||
"""LanceDB namespaces already permit a period, so a Secret must be nameable
|
||||
alongside them -- inside the name.
|
||||
|
||||
At either end it is ruled out instead, which is RFC 1123's shape and
|
||||
Kubernetes' rule for object names: it stops `.` and `..` and anything that
|
||||
reads as a hidden file or a path fragment in a listing. The service rejects
|
||||
the same shapes, so this is a local answer to the same rule rather than a
|
||||
second one.
|
||||
"""
|
||||
binding = EnvVarSecret(secret="openai.prod.v1", env_variable="OPENAI_API_KEY")
|
||||
assert binding.secret == "openai.prod.v1"
|
||||
|
||||
for name in [".", "..", ".hidden", "trailing.", "-lead", "trail-", "_x"]:
|
||||
with pytest.raises(ValueError, match="invalid Secret name"):
|
||||
EnvVarSecret(secret=name, env_variable="OPENAI_API_KEY")
|
||||
|
||||
# A namespace segment follows the same rule, for the same reason.
|
||||
for segment in [".", "..", ".hidden", "trailing."]:
|
||||
with pytest.raises(ValueError, match="invalid namespace path segment"):
|
||||
EnvVarSecret(
|
||||
secret="openai-prod",
|
||||
env_variable="OPENAI_API_KEY",
|
||||
namespace_path=[segment],
|
||||
)
|
||||
|
||||
|
||||
def _main_udf_source(
|
||||
*, threshold: int = 20, input_annotation: str = "int", comparison: str = ">="
|
||||
) -> str:
|
||||
@@ -1526,7 +1230,6 @@ def _mock_remote_function_catalog():
|
||||
"runtime": body["runtime"],
|
||||
"runtime_digest": "sha256:runtime",
|
||||
"environment_digest": "sha256:environment",
|
||||
"secret_env_bindings": body.get("secret_env_bindings", {}),
|
||||
"created_at": "2026-08-21T00:00:00Z",
|
||||
}
|
||||
response = {"job_id": "job-register"}
|
||||
@@ -1567,21 +1270,6 @@ def _mock_remote_function_catalog():
|
||||
"version": "fv_exact",
|
||||
}
|
||||
response = {"dropped": True}
|
||||
elif self.path in ("/v1/secrets/create", "/v1/secrets/alter"):
|
||||
assert set(body) == {"name", "value"}
|
||||
response = {}
|
||||
elif self.path == "/v1/secrets/list":
|
||||
if "page_token" not in body:
|
||||
response = {
|
||||
"secrets": [{"name": "openai-prod"}],
|
||||
"page_token": "next",
|
||||
}
|
||||
else:
|
||||
assert body["page_token"] == "next"
|
||||
response = {"secrets": [{"name": "hf-prod"}]}
|
||||
elif self.path == "/v1/secrets/drop":
|
||||
assert body == {"name": "openai-prod"}
|
||||
response = {}
|
||||
else:
|
||||
status = 404
|
||||
response = {"error": "not found"}
|
||||
@@ -1624,75 +1312,6 @@ def test_remote_registration_job_and_exact_version_reopen_round_trip():
|
||||
)
|
||||
|
||||
|
||||
def test_remote_registration_sends_bindings_and_never_a_credential():
|
||||
with _mock_remote_function_catalog() as (host, state):
|
||||
db = lancedb.connect(
|
||||
"db://dev",
|
||||
api_key="fake",
|
||||
host_override=host,
|
||||
client_config={"retry_config": {"retries": 0}},
|
||||
)
|
||||
created = db.create_function(
|
||||
analyze_caption,
|
||||
secrets=[EnvVarSecret(secret="openai-prod", env_variable="OPENAI_API_KEY")],
|
||||
)
|
||||
|
||||
assert dict(created.secret_env_bindings) == {
|
||||
"OPENAI_API_KEY": SecretReference(name="openai-prod")
|
||||
}
|
||||
path, create_request = state["requests"][0]
|
||||
assert path == "/v1/functions/create"
|
||||
assert create_request["secret_env_bindings"] == {
|
||||
"OPENAI_API_KEY": {"name": "openai-prod"}
|
||||
}
|
||||
# The request names a Secret and carries nothing that could be one.
|
||||
assert create_request == json.loads(
|
||||
analyze_caption.bind_secrets(
|
||||
[EnvVarSecret(secret="openai-prod", env_variable="OPENAI_API_KEY")]
|
||||
).to_canonical_json()
|
||||
)
|
||||
|
||||
|
||||
def test_remote_secret_verbs_round_trip():
|
||||
with _mock_remote_function_catalog() as (host, state):
|
||||
db = lancedb.connect(
|
||||
"db://dev",
|
||||
api_key="fake",
|
||||
host_override=host,
|
||||
client_config={"retry_config": {"retries": 0}},
|
||||
)
|
||||
assert db.create_secret("openai-prod", "sk-live-0001") is None
|
||||
assert db.alter_secret("openai-prod", "sk-live-0002") is None
|
||||
assert db.list_secrets() == ["openai-prod", "hf-prod"]
|
||||
assert db.drop_secret("openai-prod") is None
|
||||
|
||||
routes = [path for path, _ in state["requests"]]
|
||||
assert routes == [
|
||||
"/v1/secrets/create",
|
||||
"/v1/secrets/alter",
|
||||
"/v1/secrets/list",
|
||||
"/v1/secrets/list",
|
||||
"/v1/secrets/drop",
|
||||
]
|
||||
assert state["requests"][0][1] == {"name": "openai-prod", "value": "sk-live-0001"}
|
||||
# The listing returns names, and the client has no way to ask for more.
|
||||
assert state["requests"][2][1] == {}
|
||||
|
||||
|
||||
def test_building_a_binding_contacts_no_server():
|
||||
"""A binding is a local value: it says nothing about whether the Secret exists.
|
||||
|
||||
Existence is the server's answer at registration, where a mistyped name is a
|
||||
clear error rather than a client-side check that was already stale.
|
||||
"""
|
||||
with _mock_remote_function_catalog() as (_host, state):
|
||||
binding = EnvVarSecret(secret="openai-prod", env_variable="OPENAI_API_KEY")
|
||||
assert binding.secret == "openai-prod"
|
||||
assert binding.env_variable == "OPENAI_API_KEY"
|
||||
|
||||
assert state["requests"] == []
|
||||
|
||||
|
||||
def test_blocking_remote_registration_returns_function_version():
|
||||
with _mock_remote_function_catalog() as (host, state):
|
||||
db = lancedb.connect(
|
||||
|
||||
@@ -193,7 +193,13 @@ class TestNamespaceConnection:
|
||||
),
|
||||
)
|
||||
|
||||
table = db.create_table("blob_table", data, namespace_path=["test_ns"])
|
||||
# Legacy v1 blob columns are only writable at file version <= 2.1.
|
||||
table = db.create_table(
|
||||
"blob_table",
|
||||
data,
|
||||
namespace_path=["test_ns"],
|
||||
storage_options={"new_table_data_storage_version": "2.1"},
|
||||
)
|
||||
df = table.to_pandas(blob_mode="lazy").sort_values("id")
|
||||
|
||||
blob = df["blob"].iloc[0]
|
||||
|
||||
@@ -40,6 +40,10 @@ from utils import exception_output
|
||||
from importlib.util import find_spec
|
||||
|
||||
|
||||
# Legacy v1 blob columns are only writable at file version <= 2.1.
|
||||
LEGACY_BLOB_STORAGE_OPTIONS = {"new_table_data_storage_version": "2.1"}
|
||||
|
||||
|
||||
def _blob_query_data():
|
||||
return pa.table(
|
||||
{
|
||||
@@ -119,13 +123,17 @@ def _assert_blob_bytes_projection(df):
|
||||
|
||||
def _blob_query_table(db, name, blob_schema):
|
||||
if blob_schema == "v1":
|
||||
return db.create_table(name, _blob_query_data())
|
||||
return db.create_table(
|
||||
name, _blob_query_data(), storage_options=LEGACY_BLOB_STORAGE_OPTIONS
|
||||
)
|
||||
return _create_blob_v2_query_table(db, name)
|
||||
|
||||
|
||||
async def _blob_query_table_async(db, name, blob_schema):
|
||||
if blob_schema == "v1":
|
||||
return await db.create_table(name, _blob_query_data())
|
||||
return await db.create_table(
|
||||
name, _blob_query_data(), storage_options=LEGACY_BLOB_STORAGE_OPTIONS
|
||||
)
|
||||
return await _create_blob_v2_query_table_async(db, name)
|
||||
|
||||
|
||||
@@ -275,7 +283,9 @@ async def test_query_to_pandas_kwargs(table, table_async):
|
||||
def test_plain_scan_query_to_pandas_blob_modes(tmp_db, blob_mode):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table(
|
||||
f"test_query_to_pandas_blob_{blob_mode}", _blob_query_data()
|
||||
f"test_query_to_pandas_blob_{blob_mode}",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
|
||||
df = (
|
||||
@@ -322,7 +332,9 @@ def test_plain_scan_query_to_pandas_blob_mode_does_not_collect_arrow(
|
||||
):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table(
|
||||
"test_query_to_pandas_blob_no_arrow_collect", _blob_query_data()
|
||||
"test_query_to_pandas_blob_no_arrow_collect",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
query = table.search().where("id = 1").select(["id", "blob"])
|
||||
|
||||
@@ -347,7 +359,9 @@ def test_plain_scan_query_to_pandas_blob_descriptions_flatten_uses_scanner(
|
||||
):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table(
|
||||
"test_query_to_pandas_blob_desc_flatten", _blob_query_data()
|
||||
"test_query_to_pandas_blob_desc_flatten",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
query = table.search().where("id = 1").select(["id", "blob"])
|
||||
|
||||
@@ -365,7 +379,11 @@ def test_plain_scan_query_to_pandas_blob_descriptions_flatten_uses_scanner(
|
||||
def test_plain_scan_query_to_pandas_scanner_state(tmp_db):
|
||||
pytest.importorskip("lance")
|
||||
data = _blob_query_data()
|
||||
table = tmp_db.create_table("test_query_to_pandas_scanner_state", data.slice(0, 2))
|
||||
table = tmp_db.create_table(
|
||||
"test_query_to_pandas_scanner_state",
|
||||
data.slice(0, 2),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
table.add(data.slice(2, 2))
|
||||
|
||||
fragments = table.to_lance().get_fragments()
|
||||
@@ -400,7 +418,9 @@ def test_plain_scan_query_to_pandas_scanner_state(tmp_db):
|
||||
async def test_async_plain_scan_query_to_pandas_blob_projection(tmp_db_async):
|
||||
pytest.importorskip("lance")
|
||||
table = await tmp_db_async.create_table(
|
||||
"test_async_query_to_pandas_blob_projection", _blob_query_data()
|
||||
"test_async_query_to_pandas_blob_projection",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
|
||||
lazy_df = await (
|
||||
@@ -452,7 +472,9 @@ async def test_async_plain_scan_query_to_pandas_blob_mode_does_not_collect_arrow
|
||||
):
|
||||
pytest.importorskip("lance")
|
||||
table = await tmp_db_async.create_table(
|
||||
"test_async_query_to_pandas_blob_no_arrow_collect", _blob_query_data()
|
||||
"test_async_query_to_pandas_blob_no_arrow_collect",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
query = table.query().where("id = 1").select(["id", "blob"])
|
||||
|
||||
@@ -474,7 +496,11 @@ async def test_async_plain_scan_query_to_pandas_blob_mode_does_not_collect_arrow
|
||||
|
||||
def test_vector_query_to_pandas_blob_mode_requires_native_path(tmp_db):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table("test_vector_query_blob_mode", _blob_query_data())
|
||||
table = tmp_db.create_table(
|
||||
"test_vector_query_blob_mode",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
|
||||
with pytest.raises(RuntimeError, match="Lance native pandas conversion"):
|
||||
table.search([1.0, 0.0]).select(["blob", "vector"]).limit(1).to_pandas(
|
||||
@@ -485,7 +511,9 @@ def test_vector_query_to_pandas_blob_mode_requires_native_path(tmp_db):
|
||||
def test_vector_query_to_pandas_blob_descriptions_requires_plain_scan(tmp_db):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table(
|
||||
"test_vector_query_blob_descriptions", _blob_query_data()
|
||||
"test_vector_query_blob_descriptions",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
|
||||
with pytest.raises(RuntimeError, match="plain scan query"):
|
||||
|
||||
@@ -64,15 +64,23 @@ async def _blob_v2_table_async(db: AsyncConnection, name: str):
|
||||
return table
|
||||
|
||||
|
||||
# Legacy v1 blob columns are only writable at file version <= 2.1.
|
||||
LEGACY_BLOB_STORAGE_OPTIONS = {"new_table_data_storage_version": "2.1"}
|
||||
|
||||
|
||||
def _blob_table(db: DBConnection, name: str, blob_schema: str):
|
||||
if blob_schema == "v1":
|
||||
return db.create_table(name, data=_blob_test_data())
|
||||
return db.create_table(
|
||||
name, data=_blob_test_data(), storage_options=LEGACY_BLOB_STORAGE_OPTIONS
|
||||
)
|
||||
return _blob_v2_table(db, name)
|
||||
|
||||
|
||||
async def _blob_table_async(db: AsyncConnection, name: str, blob_schema: str):
|
||||
if blob_schema == "v1":
|
||||
return await db.create_table(name, data=_blob_test_data())
|
||||
return await db.create_table(
|
||||
name, data=_blob_test_data(), storage_options=LEGACY_BLOB_STORAGE_OPTIONS
|
||||
)
|
||||
return await _blob_v2_table_async(db, name)
|
||||
|
||||
|
||||
@@ -147,7 +155,11 @@ def test_table_to_pandas_invalid_blob_mode_non_blob_table(tmp_db: DBConnection):
|
||||
@pytest.mark.parametrize("blob_mode", ["lazy", "bytes", "descriptions"])
|
||||
def test_table_to_pandas_blob_modes(tmp_db: DBConnection, blob_mode):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table(f"test_to_pandas_blob_{blob_mode}", _blob_test_data())
|
||||
table = tmp_db.create_table(
|
||||
f"test_to_pandas_blob_{blob_mode}",
|
||||
_blob_test_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
|
||||
df = table.to_pandas(blob_mode=blob_mode)
|
||||
|
||||
@@ -3959,7 +3971,7 @@ def test_stats(mem_db: DBConnection):
|
||||
print(f"{stats=}")
|
||||
assert stats == {
|
||||
# Full on-disk size of the data file, footer and metadata included.
|
||||
"total_bytes": 633,
|
||||
"total_bytes": 637,
|
||||
"num_rows": 2,
|
||||
"num_indices": 0,
|
||||
"fragment_stats": {
|
||||
@@ -4171,13 +4183,14 @@ def test_refresh_column_async_returns_job(tmp_path):
|
||||
assert result.rows_failed == 0
|
||||
assert result.rows_remaining == 0
|
||||
assert result.source_version == 2
|
||||
assert result.published_version == 3
|
||||
# The fill lands at 3; the stamp recording its inputs is published at 4.
|
||||
assert result.published_version == 4
|
||||
assert job.status() == "finished"
|
||||
assert sorted(table.to_arrow()["doubled"].to_pylist()) == [2, 4]
|
||||
|
||||
no_op = table.refresh_column_async("doubled").wait()
|
||||
assert no_op.rows_assigned == 0
|
||||
assert no_op.source_version == 3
|
||||
assert no_op.source_version == 4
|
||||
assert no_op.published_version is None
|
||||
|
||||
# Bad input raises at the call, not through the job.
|
||||
@@ -4196,6 +4209,6 @@ async def test_refresh_column_async_job_async_table(tmp_path):
|
||||
assert isinstance(result, lancedb.RefreshColumnResult)
|
||||
assert result.rows_assigned == 1
|
||||
assert result.source_version == 2
|
||||
assert result.published_version == 3
|
||||
assert result.published_version == 4
|
||||
assert await job.status() == "finished"
|
||||
assert (await table.to_arrow())["tripled"].to_pylist() == [9]
|
||||
|
||||
@@ -704,78 +704,6 @@ impl Connection {
|
||||
})
|
||||
}
|
||||
|
||||
pub fn create_secret(
|
||||
self_: PyRef<'_, Self>,
|
||||
name: String,
|
||||
value: String,
|
||||
namespace_path: Vec<String>,
|
||||
) -> PyResult<Bound<'_, PyAny>> {
|
||||
let inner = self_.get_inner()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
inner
|
||||
.create_secret(name, value, &namespace_path)
|
||||
.await
|
||||
.infer_error()
|
||||
})
|
||||
}
|
||||
|
||||
pub fn alter_secret(
|
||||
self_: PyRef<'_, Self>,
|
||||
name: String,
|
||||
value: String,
|
||||
namespace_path: Vec<String>,
|
||||
) -> PyResult<Bound<'_, PyAny>> {
|
||||
let inner = self_.get_inner()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
inner
|
||||
.alter_secret(name, value, &namespace_path)
|
||||
.await
|
||||
.infer_error()
|
||||
})
|
||||
}
|
||||
|
||||
pub fn list_secrets(
|
||||
self_: PyRef<'_, Self>,
|
||||
namespace_path: Vec<String>,
|
||||
) -> PyResult<Bound<'_, PyAny>> {
|
||||
let inner = self_.get_inner()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
inner.list_secrets(&namespace_path).await.infer_error()
|
||||
})
|
||||
}
|
||||
|
||||
pub fn drop_secret(
|
||||
self_: PyRef<'_, Self>,
|
||||
name: String,
|
||||
namespace_path: Vec<String>,
|
||||
) -> PyResult<Bound<'_, PyAny>> {
|
||||
let inner = self_.get_inner()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
inner.drop_secret(name, &namespace_path).await.infer_error()
|
||||
})
|
||||
}
|
||||
|
||||
/// Name and timestamps as a plain mapping. `SecretInfo` carries no value,
|
||||
/// so there is none to filter out here.
|
||||
pub fn describe_secret(
|
||||
self_: PyRef<'_, Self>,
|
||||
name: String,
|
||||
namespace_path: Vec<String>,
|
||||
) -> PyResult<Bound<'_, PyAny>> {
|
||||
let inner = self_.get_inner()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
let info = inner
|
||||
.describe_secret(name, &namespace_path)
|
||||
.await
|
||||
.infer_error()?;
|
||||
Ok(HashMap::from([
|
||||
("name".to_string(), info.name),
|
||||
("created_at".to_string(), info.created_at),
|
||||
("updated_at".to_string(), info.updated_at),
|
||||
]))
|
||||
})
|
||||
}
|
||||
|
||||
pub fn list_jobs(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
|
||||
let inner = self_.get_inner()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
|
||||
Generated
+95
-98
@@ -10,9 +10,6 @@ resolution-markers = [
|
||||
"python_full_version < '3.11'",
|
||||
]
|
||||
|
||||
[options]
|
||||
prerelease-mode = "allow"
|
||||
|
||||
[[package]]
|
||||
name = "accelerate"
|
||||
version = "1.14.0"
|
||||
@@ -802,7 +799,7 @@ name = "cuda-bindings"
|
||||
version = "13.3.1"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "cuda-pathfinder" },
|
||||
{ name = "cuda-pathfinder", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
|
||||
]
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/a9/21/8464d133752951c154feafb3b65c297e7d80f301183d220bec4c830f1441/cuda_bindings-13.3.1-cp310-cp310-manylinux_2_24_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:120fcc53d57903df529c3486962c56528cba5b7d6c57c99537320ed9922c8b86", size = 6073403, upload-time = "2026-05-29T23:11:36.22Z" },
|
||||
@@ -837,37 +834,37 @@ wheels = [
|
||||
|
||||
[package.optional-dependencies]
|
||||
cublas = [
|
||||
{ name = "nvidia-cublas" },
|
||||
{ name = "nvidia-cublas", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
|
||||
]
|
||||
cudart = [
|
||||
{ name = "nvidia-cuda-runtime" },
|
||||
{ name = "nvidia-cuda-runtime", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
|
||||
]
|
||||
cufft = [
|
||||
{ name = "nvidia-cufft" },
|
||||
{ name = "nvidia-cufft", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
|
||||
]
|
||||
cufile = [
|
||||
{ name = "nvidia-cufile" },
|
||||
{ name = "nvidia-cufile", marker = "sys_platform == 'linux'" },
|
||||
]
|
||||
cupti = [
|
||||
{ name = "nvidia-cuda-cupti" },
|
||||
{ name = "nvidia-cuda-cupti", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
|
||||
]
|
||||
curand = [
|
||||
{ name = "nvidia-curand" },
|
||||
{ name = "nvidia-curand", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
|
||||
]
|
||||
cusolver = [
|
||||
{ name = "nvidia-cusolver" },
|
||||
{ name = "nvidia-cusolver", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
|
||||
]
|
||||
cusparse = [
|
||||
{ name = "nvidia-cusparse" },
|
||||
{ name = "nvidia-cusparse", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
|
||||
]
|
||||
nvjitlink = [
|
||||
{ name = "nvidia-nvjitlink" },
|
||||
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
|
||||
]
|
||||
nvrtc = [
|
||||
{ name = "nvidia-cuda-nvrtc" },
|
||||
{ name = "nvidia-cuda-nvrtc", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
|
||||
]
|
||||
nvtx = [
|
||||
{ name = "nvidia-nvtx" },
|
||||
{ name = "nvidia-nvtx", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -1026,7 +1023,7 @@ name = "exceptiongroup"
|
||||
version = "1.3.1"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "typing-extensions" },
|
||||
{ name = "typing-extensions", marker = "python_full_version < '3.11'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/50/79/66800aadf48771f6b62f7eb014e352e5d06856655206165d775e675a02c9/exceptiongroup-1.3.1.tar.gz", hash = "sha256:8b412432c6055b0b7d14c310000ae93352ed6754f70fa8f7c34141f91c4e3219", size = 30371, upload-time = "2025-11-21T23:01:54.787Z" }
|
||||
wheels = [
|
||||
@@ -1443,16 +1440,16 @@ resolution-markers = [
|
||||
"python_full_version < '3.11'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "cachetools" },
|
||||
{ name = "certifi" },
|
||||
{ name = "httpx" },
|
||||
{ name = "ibm-cos-sdk" },
|
||||
{ name = "lomond" },
|
||||
{ name = "packaging" },
|
||||
{ name = "pandas", version = "2.2.3", source = { registry = "https://pypi.org/simple" } },
|
||||
{ name = "requests" },
|
||||
{ name = "tabulate" },
|
||||
{ name = "urllib3" },
|
||||
{ name = "cachetools", marker = "python_full_version < '3.11'" },
|
||||
{ name = "certifi", marker = "python_full_version < '3.11'" },
|
||||
{ name = "httpx", marker = "python_full_version < '3.11'" },
|
||||
{ name = "ibm-cos-sdk", marker = "python_full_version < '3.11'" },
|
||||
{ name = "lomond", marker = "python_full_version < '3.11'" },
|
||||
{ name = "packaging", marker = "python_full_version < '3.11'" },
|
||||
{ name = "pandas", version = "2.2.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
|
||||
{ name = "requests", marker = "python_full_version < '3.11'" },
|
||||
{ name = "tabulate", marker = "python_full_version < '3.11'" },
|
||||
{ name = "urllib3", marker = "python_full_version < '3.11'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/c7/56/2e3df38a1f13062095d7bde23c87a92f3898982993a15186b1bfecbd206f/ibm_watsonx_ai-1.3.42.tar.gz", hash = "sha256:ee5be59009004245d957ce97d1227355516df95a2640189749487614fef674ff", size = 688651, upload-time = "2025-10-01T13:35:41.527Z" }
|
||||
wheels = [
|
||||
@@ -1471,17 +1468,17 @@ resolution-markers = [
|
||||
"python_full_version == '3.11.*'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "cachetools" },
|
||||
{ name = "certifi" },
|
||||
{ name = "httpx" },
|
||||
{ name = "ibm-cos-sdk" },
|
||||
{ name = "lomond" },
|
||||
{ name = "packaging" },
|
||||
{ name = "pandas", version = "2.3.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.14'" },
|
||||
{ name = "cachetools", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "certifi", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "httpx", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "ibm-cos-sdk", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "lomond", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "packaging", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "pandas", version = "2.3.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
|
||||
{ name = "pandas", version = "3.0.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.14'" },
|
||||
{ name = "requests" },
|
||||
{ name = "tabulate" },
|
||||
{ name = "urllib3" },
|
||||
{ name = "requests", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "tabulate", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "urllib3", marker = "python_full_version >= '3.11'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/29/a3/c756b534696ab2f3f29882fdb7ca7198b7a5c94e10c0a3a327853d6d6b79/ibm_watsonx_ai-1.5.14.tar.gz", hash = "sha256:a756488bd57e87c0fc51be42dcba871143cfe0ac1e805c497c5047e1e4f13e9d", size = 735804, upload-time = "2026-06-22T12:32:43.85Z" }
|
||||
wheels = [
|
||||
@@ -1557,17 +1554,17 @@ resolution-markers = [
|
||||
"python_full_version < '3.11'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "colorama", marker = "sys_platform == 'win32'" },
|
||||
{ name = "decorator" },
|
||||
{ name = "exceptiongroup" },
|
||||
{ name = "jedi" },
|
||||
{ name = "matplotlib-inline" },
|
||||
{ name = "pexpect", marker = "sys_platform != 'emscripten' and sys_platform != 'win32'" },
|
||||
{ name = "prompt-toolkit" },
|
||||
{ name = "pygments" },
|
||||
{ name = "stack-data" },
|
||||
{ name = "traitlets" },
|
||||
{ name = "typing-extensions" },
|
||||
{ name = "colorama", marker = "python_full_version < '3.11' and sys_platform == 'win32'" },
|
||||
{ name = "decorator", marker = "python_full_version < '3.11'" },
|
||||
{ name = "exceptiongroup", marker = "python_full_version < '3.11'" },
|
||||
{ name = "jedi", marker = "python_full_version < '3.11'" },
|
||||
{ name = "matplotlib-inline", marker = "python_full_version < '3.11'" },
|
||||
{ name = "pexpect", marker = "python_full_version < '3.11' and sys_platform != 'emscripten' and sys_platform != 'win32'" },
|
||||
{ name = "prompt-toolkit", marker = "python_full_version < '3.11'" },
|
||||
{ name = "pygments", marker = "python_full_version < '3.11'" },
|
||||
{ name = "stack-data", marker = "python_full_version < '3.11'" },
|
||||
{ name = "traitlets", marker = "python_full_version < '3.11'" },
|
||||
{ name = "typing-extensions", marker = "python_full_version < '3.11'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/40/18/f8598d287006885e7136451fdea0755af4ebcbfe342836f24deefaed1164/ipython-8.39.0.tar.gz", hash = "sha256:4110ae96012c379b8b6db898a07e186c40a2a1ef5d57a7fa83166047d9da7624", size = 5513971, upload-time = "2026-03-27T10:02:13.94Z" }
|
||||
wheels = [
|
||||
@@ -1586,18 +1583,18 @@ resolution-markers = [
|
||||
"python_full_version == '3.11.*'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "colorama", marker = "sys_platform == 'win32'" },
|
||||
{ name = "decorator" },
|
||||
{ name = "ipython-pygments-lexers" },
|
||||
{ name = "jedi" },
|
||||
{ name = "matplotlib-inline" },
|
||||
{ name = "pexpect", marker = "sys_platform != 'emscripten' and sys_platform != 'win32'" },
|
||||
{ name = "prompt-toolkit" },
|
||||
{ name = "psutil", marker = "sys_platform != 'cygwin' and sys_platform != 'emscripten'" },
|
||||
{ name = "pygments" },
|
||||
{ name = "stack-data" },
|
||||
{ name = "traitlets" },
|
||||
{ name = "typing-extensions", marker = "python_full_version < '3.12'" },
|
||||
{ name = "colorama", marker = "python_full_version >= '3.11' and sys_platform == 'win32'" },
|
||||
{ name = "decorator", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "ipython-pygments-lexers", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "jedi", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "matplotlib-inline", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "pexpect", marker = "python_full_version >= '3.11' and sys_platform != 'emscripten' and sys_platform != 'win32'" },
|
||||
{ name = "prompt-toolkit", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "psutil", marker = "python_full_version >= '3.11' and sys_platform != 'cygwin' and sys_platform != 'emscripten'" },
|
||||
{ name = "pygments", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "stack-data", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "traitlets", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "typing-extensions", marker = "python_full_version == '3.11.*'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/53/59/165d3b4d75cc34add3122c4417ecb229085140ac573103c223cd01dde96f/ipython-9.15.0.tar.gz", hash = "sha256:da2819ce2aa83135257df830660b1176d986c3d2876db24df01974fa955b2756", size = 4442580, upload-time = "2026-06-26T11:03:35.913Z" }
|
||||
wheels = [
|
||||
@@ -1609,7 +1606,7 @@ name = "ipython-pygments-lexers"
|
||||
version = "1.1.1"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "pygments" },
|
||||
{ name = "pygments", marker = "python_full_version >= '3.11'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/ef/4c/5dd1d8af08107f88c7f741ead7a40854b8ac24ddf9ae850afbcf698aa552/ipython_pygments_lexers-1.1.1.tar.gz", hash = "sha256:09c0138009e56b6854f9535736f4171d855c8c08a563a0dcd8022f78355c7e81", size = 8393, upload-time = "2025-01-17T11:24:34.505Z" }
|
||||
wheels = [
|
||||
@@ -2861,7 +2858,7 @@ name = "nvidia-cudnn-cu13"
|
||||
version = "9.19.0.56"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "nvidia-cublas" },
|
||||
{ name = "nvidia-cublas", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
|
||||
]
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/f1/84/26025437c1e6b61a707442184fa0c03d083b661adf3a3eecfd6d21677740/nvidia_cudnn_cu13-9.19.0.56-py3-none-manylinux_2_27_aarch64.whl", hash = "sha256:6ed29ffaee1176c612daf442e4dd6cfeb6a0caa43ddcbeb59da94953030b1be4", size = 433781201, upload-time = "2026-02-03T20:40:53.805Z" },
|
||||
@@ -2873,7 +2870,7 @@ name = "nvidia-cufft"
|
||||
version = "12.0.0.61"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "nvidia-nvjitlink" },
|
||||
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
|
||||
]
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/8b/ae/f417a75c0259e85c1d2f83ca4e960289a5f814ed0cea74d18c353d3e989d/nvidia_cufft-12.0.0.61-py3-none-manylinux2014_aarch64.manylinux_2_17_aarch64.whl", hash = "sha256:2708c852ef8cd89d1d2068bdbece0aa188813a0c934db3779b9b1faa8442e5f5", size = 214053554, upload-time = "2025-09-04T08:31:38.196Z" },
|
||||
@@ -2903,9 +2900,9 @@ name = "nvidia-cusolver"
|
||||
version = "12.0.4.66"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "nvidia-cublas" },
|
||||
{ name = "nvidia-cusparse" },
|
||||
{ name = "nvidia-nvjitlink" },
|
||||
{ name = "nvidia-cublas", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
|
||||
{ name = "nvidia-cusparse", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
|
||||
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
|
||||
]
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/c8/c3/b30c9e935fc01e3da443ec0116ed1b2a009bb867f5324d3f2d7e533e776b/nvidia_cusolver-12.0.4.66-py3-none-manylinux_2_27_aarch64.whl", hash = "sha256:02c2457eaa9e39de20f880f4bd8820e6a1cfb9f9a34f820eb12a155aa5bc92d2", size = 223467760, upload-time = "2025-09-04T08:33:04.222Z" },
|
||||
@@ -2917,7 +2914,7 @@ name = "nvidia-cusparse"
|
||||
version = "12.6.3.3"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "nvidia-nvjitlink" },
|
||||
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
|
||||
]
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/f8/94/5c26f33738ae35276672f12615a64bd008ed5be6d1ebcb23579285d960a9/nvidia_cusparse-12.6.3.3-py3-none-manylinux2014_aarch64.manylinux_2_17_aarch64.whl", hash = "sha256:80bcc4662f23f1054ee334a15c72b8940402975e0eab63178fc7e670aa59472c", size = 162155568, upload-time = "2025-09-04T08:33:42.864Z" },
|
||||
@@ -3094,10 +3091,10 @@ resolution-markers = [
|
||||
"python_full_version < '3.11'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } },
|
||||
{ name = "python-dateutil" },
|
||||
{ name = "pytz" },
|
||||
{ name = "tzdata" },
|
||||
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
|
||||
{ name = "python-dateutil", marker = "python_full_version < '3.11'" },
|
||||
{ name = "pytz", marker = "python_full_version < '3.11'" },
|
||||
{ name = "tzdata", marker = "python_full_version < '3.11'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/9c/d6/9f8431bacc2e19dca897724cd097b1bb224a6ad5433784a44b587c7c13af/pandas-2.2.3.tar.gz", hash = "sha256:4f18ba62b61d7e192368b84517265a99b4d7ee8912f8708660fb4a366cc82667", size = 4399213, upload-time = "2024-09-20T13:10:04.827Z" }
|
||||
wheels = [
|
||||
@@ -3146,11 +3143,11 @@ resolution-markers = [
|
||||
"python_full_version == '3.11.*'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12' or python_full_version >= '3.14'" },
|
||||
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
|
||||
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12' and python_full_version < '3.14'" },
|
||||
{ name = "python-dateutil" },
|
||||
{ name = "pytz" },
|
||||
{ name = "tzdata" },
|
||||
{ name = "python-dateutil", marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
|
||||
{ name = "pytz", marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
|
||||
{ name = "tzdata", marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/33/01/d40b85317f86cf08d853a4f495195c73815fdf205eef3993821720274518/pandas-2.3.3.tar.gz", hash = "sha256:e05e1af93b977f7eafa636d043f9f94c7ee3ac81af99c13508215942e64c993b", size = 4495223, upload-time = "2025-09-29T23:34:51.853Z" }
|
||||
wheels = [
|
||||
@@ -3213,9 +3210,9 @@ resolution-markers = [
|
||||
"python_full_version >= '3.14' and sys_platform != 'emscripten' and sys_platform != 'win32'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" } },
|
||||
{ name = "python-dateutil" },
|
||||
{ name = "tzdata", marker = "sys_platform == 'emscripten' or sys_platform == 'win32'" },
|
||||
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.14'" },
|
||||
{ name = "python-dateutil", marker = "python_full_version >= '3.14'" },
|
||||
{ name = "tzdata", marker = "(python_full_version >= '3.14' and sys_platform == 'emscripten') or (python_full_version >= '3.14' and sys_platform == 'win32')" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/f8/87/4341c6252d1c47b08768c3d25ac487362bf403f0313ddae4a2a26c9b1b4c/pandas-3.0.3.tar.gz", hash = "sha256:696a4a00a2a2a35d4e5deb3fc946641b96c944f02230e4f76137fe35d806c4fc", size = 4651414, upload-time = "2026-05-11T18:54:29.21Z" }
|
||||
wheels = [
|
||||
@@ -3323,7 +3320,7 @@ name = "pexpect"
|
||||
version = "4.9.0"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "ptyprocess" },
|
||||
{ name = "ptyprocess", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/42/92/cc564bf6381ff43ce1f4d06852fc19a2f11d180f23dc32d9588bee2f149d/pexpect-4.9.0.tar.gz", hash = "sha256:ee7d41123f3c9911050ea2c2dac107568dc43b2d3b0c7557a33212c398ead30f", size = 166450, upload-time = "2023-11-25T09:07:26.339Z" }
|
||||
wheels = [
|
||||
@@ -3915,8 +3912,8 @@ crypto = [
|
||||
|
||||
[[package]]
|
||||
name = "pylance"
|
||||
version = "9.0.0rc1"
|
||||
source = { registry = "https://pypi.fury.io/lance-format/" }
|
||||
version = "7.0.0"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "lance-namespace" },
|
||||
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
|
||||
@@ -3925,12 +3922,12 @@ dependencies = [
|
||||
{ name = "pyarrow" },
|
||||
]
|
||||
wheels = [
|
||||
{ url = "https://pypi.fury.io/lance-format/-/ver_vEHBE/pylance-9.0.0rc1-cp310-abi3-macosx_11_0_arm64.whl", hash = "sha256:f0b6b02a1808bb3072ee7fe4e36614cae6f86302513e73ec7f55b2234a963b24" },
|
||||
{ url = "https://pypi.fury.io/lance-format/-/ver_1Jipm4/pylance-9.0.0rc1-cp310-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:30f0ebf0d88034301819eb964f9236ce555aaa58e7ab89c5975a3e2250bbb405" },
|
||||
{ url = "https://pypi.fury.io/lance-format/-/ver_IvKxo/pylance-9.0.0rc1-cp310-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:44609ea2615ea6e684b85478d1694af2026458f61cf7895ecc75e238bfd17aa8" },
|
||||
{ url = "https://pypi.fury.io/lance-format/-/ver_2hidj1/pylance-9.0.0rc1-cp310-abi3-manylinux_2_28_aarch64.whl", hash = "sha256:182167a8dba9eeabffbffd53bd5b8548613d4d459b7cd7b34a840dd00cbb806f" },
|
||||
{ url = "https://pypi.fury.io/lance-format/-/ver_1dFx3r/pylance-9.0.0rc1-cp310-abi3-manylinux_2_28_x86_64.whl", hash = "sha256:8a63b11e814b7eab758bcaf0d6f97eb05ea86203d9fb0af718c462c24c7d6c9c" },
|
||||
{ url = "https://pypi.fury.io/lance-format/-/ver_2a8dSh/pylance-9.0.0rc1-cp310-abi3-win_amd64.whl", hash = "sha256:2ff8b953ae2b0550490c1a7efd210aa91bc223d200ffac28849056cfd7436d97" },
|
||||
{ url = "https://files.pythonhosted.org/packages/ac/ad/2f64921bf346e7075aef24a72595db44821724a3d89a9a92dd24e79632aa/pylance-7.0.0-cp39-abi3-macosx_11_0_arm64.whl", hash = "sha256:98422021975be76e72b1572f41b8c9abb3bee5bdc9bfa5e9ce731110a65ed4d1", size = 62134146, upload-time = "2026-05-27T21:59:37.459Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/73/1c/c5a01bee0160b55d9a98895cbd33091d038f0a0995b121ab72e629008d02/pylance-7.0.0-cp39-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:4bec86ee5b6fbd8bfc493e653f0a1fba0303cfe5492b9b46fc25ab908edc7183", size = 65373684, upload-time = "2026-05-27T22:04:01.584Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/eb/da/1fe8b8f7dbfe734d76af76acc994fc360a0d0c79a4874ef69f5a72a58fe3/pylance-7.0.0-cp39-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:881491432c53184e52f8d1db8d5f872f39a03f36fb104bec77b33d379519d8b5", size = 69458555, upload-time = "2026-05-27T22:16:50.567Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/76/f0/dd505cf3fd0226ab9d94759acd713125af1d3bfacfd80bbd52e3b9f89509/pylance-7.0.0-cp39-abi3-manylinux_2_28_aarch64.whl", hash = "sha256:18453999e7fff4f76b16d6b7882c9df0628bd142ff95e2461bd7dd5ee3fe0af3", size = 65394430, upload-time = "2026-05-27T22:05:30.923Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/17/ba/2357b81034f28eb00790e258ed140289a6a887a7468ca9df6349fd186b27/pylance-7.0.0-cp39-abi3-manylinux_2_28_x86_64.whl", hash = "sha256:04a58051d408c60fe76d41a220dcaf8fea8fb6d1aa0ca78a709b60bc3cc8d19a", size = 69473470, upload-time = "2026-05-27T22:17:18.935Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/1f/ec/5c00b6303a67d787f9475141832cbdc513d674ac3dcaeef8a7b169905e65/pylance-7.0.0-cp39-abi3-win_amd64.whl", hash = "sha256:467d4864af047eaab4e1370e2f1e88e2c6f507c079874421116cb41d78bc3629", size = 74792863, upload-time = "2026-05-27T22:19:23.875Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -4686,10 +4683,10 @@ resolution-markers = [
|
||||
"python_full_version < '3.11'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "joblib" },
|
||||
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } },
|
||||
{ name = "scipy", version = "1.15.3", source = { registry = "https://pypi.org/simple" } },
|
||||
{ name = "threadpoolctl" },
|
||||
{ name = "joblib", marker = "python_full_version < '3.11'" },
|
||||
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
|
||||
{ name = "scipy", version = "1.15.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
|
||||
{ name = "threadpoolctl", marker = "python_full_version < '3.11'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/98/c2/a7855e41c9d285dfe86dc50b250978105dce513d6e459ea66a6aeb0e1e0c/scikit_learn-1.7.2.tar.gz", hash = "sha256:20e9e49ecd130598f1ca38a1d85090e1a600147b9c02fa6f15d69cb53d968fda", size = 7193136, upload-time = "2025-09-09T08:21:29.075Z" }
|
||||
wheels = [
|
||||
@@ -4737,13 +4734,13 @@ resolution-markers = [
|
||||
"python_full_version == '3.11.*'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "joblib" },
|
||||
{ name = "narwhals" },
|
||||
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" },
|
||||
{ name = "joblib", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "narwhals", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
|
||||
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" },
|
||||
{ name = "scipy", version = "1.17.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" },
|
||||
{ name = "scipy", version = "1.17.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
|
||||
{ name = "scipy", version = "1.18.0", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" },
|
||||
{ name = "threadpoolctl" },
|
||||
{ name = "threadpoolctl", marker = "python_full_version >= '3.11'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/fa/6f/37092bdb25f712817231799fc5674d8e704066a8a70c1d2d40517e18b4ab/scikit_learn-1.9.0.tar.gz", hash = "sha256:8833266989d3a5110178a9fae30783675460724d0e1efb13b14901d2c660c557", size = 7750767, upload-time = "2026-06-02T11:54:32.706Z" }
|
||||
wheels = [
|
||||
@@ -4787,7 +4784,7 @@ resolution-markers = [
|
||||
"python_full_version < '3.11'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } },
|
||||
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/0f/37/6964b830433e654ec7485e45a00fc9a27cf868d622838f6b6d9c5ec0d532/scipy-1.15.3.tar.gz", hash = "sha256:eae3cf522bc7df64b42cad3925c876e1b0b6c35c1337c93e12c0f366f55b0eaf", size = 59419214, upload-time = "2025-05-08T16:13:05.955Z" }
|
||||
wheels = [
|
||||
@@ -4846,7 +4843,7 @@ resolution-markers = [
|
||||
"python_full_version == '3.11.*'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" } },
|
||||
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/7a/97/5a3609c4f8d58b039179648e62dd220f89864f56f7357f5d4f45c29eb2cc/scipy-1.17.1.tar.gz", hash = "sha256:95d8e012d8cb8816c226aef832200b1d45109ed4464303e997c5b13122b297c0", size = 30573822, upload-time = "2026-02-23T00:26:24.851Z" }
|
||||
wheels = [
|
||||
@@ -4923,7 +4920,7 @@ resolution-markers = [
|
||||
"python_full_version >= '3.12' and python_full_version < '3.14'",
|
||||
]
|
||||
dependencies = [
|
||||
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" } },
|
||||
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/a7/25/c2700dfaf6442b4effaa91af24ebce5dc9d31bb4a69706313aae70d72cd0/scipy-1.18.0.tar.gz", hash = "sha256:67b2ad2ad54c72ca6d04975a9b2df8c3638c34ddd5b28738e94fc2b57929d378", size = 30774447, upload-time = "2026-06-19T15:01:43.456Z" }
|
||||
wheels = [
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "lancedb"
|
||||
version = "0.39.0-beta.4"
|
||||
version = "0.39.0-beta.6"
|
||||
edition.workspace = true
|
||||
description = "LanceDB: A serverless, low-latency vector database for AI applications"
|
||||
license.workspace = true
|
||||
@@ -95,13 +95,14 @@ candle-transformers = { version = "0.9.1", optional = true }
|
||||
candle-nn = { version = "0.9.1", optional = true }
|
||||
tokenizers = { version = "0.19.1", optional = true }
|
||||
semver = { workspace = true }
|
||||
roaring = "0.11.4"
|
||||
sha2 = "0.10"
|
||||
|
||||
[dev-dependencies]
|
||||
anyhow = "1"
|
||||
lance-testing = { workspace = true }
|
||||
tempfile = { workspace = true }
|
||||
random_word = { version = "0.4.3", features = ["en"] }
|
||||
roaring = "0.11.4"
|
||||
tokio = { workspace = true, features = ["io-util", "macros", "net", "test-util"] }
|
||||
uuid = { workspace = true }
|
||||
walkdir = "2"
|
||||
|
||||
@@ -532,10 +532,11 @@ mod tests {
|
||||
fn storage_version_bumps_to_v2_2() {
|
||||
let mut params = WriteParams::default();
|
||||
ensure_blob_storage_version(&blob_schema(), &mut params);
|
||||
assert_eq!(
|
||||
params.data_storage_version.unwrap().resolve(),
|
||||
ConcreteFileVersion::V2_2
|
||||
);
|
||||
let resolved = params
|
||||
.data_storage_version
|
||||
.unwrap_or(LanceFileVersion::Stable)
|
||||
.resolve();
|
||||
assert_eq!(resolved, ConcreteFileVersion::V2_2);
|
||||
assert!(!params.enable_stable_row_ids);
|
||||
}
|
||||
|
||||
@@ -547,10 +548,11 @@ mod tests {
|
||||
};
|
||||
ensure_blob_storage_version(&blob_schema(), &mut params);
|
||||
assert!(params.enable_stable_row_ids);
|
||||
assert_eq!(
|
||||
params.data_storage_version.unwrap().resolve(),
|
||||
ConcreteFileVersion::V2_2
|
||||
);
|
||||
let resolved = params
|
||||
.data_storage_version
|
||||
.unwrap_or(LanceFileVersion::Stable)
|
||||
.resolve();
|
||||
assert_eq!(resolved, ConcreteFileVersion::V2_2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -24,7 +24,7 @@ use crate::data::scannable::Scannable;
|
||||
use crate::database::listing::ListingDatabase;
|
||||
use crate::database::{
|
||||
CloneTableRequest, Database, DatabaseOptions, JobInfo, OpenTableRequest, ReadConsistency,
|
||||
SecretInfo, TableNamesRequest,
|
||||
TableNamesRequest,
|
||||
};
|
||||
use crate::embeddings::{EmbeddingRegistry, MemoryRegistry};
|
||||
use crate::error::{Error, Result};
|
||||
@@ -586,15 +586,10 @@ impl Connection {
|
||||
/// Registration is remote-only and always asynchronous. Waiting on the
|
||||
/// returned typed job yields the durable [`crate::function::FunctionVersion`].
|
||||
/// Local databases return [`Error::NotSupported`].
|
||||
///
|
||||
/// The request's binding shape is validated here rather than in any one
|
||||
/// language binding, so every client surface rejects the same envelopes
|
||||
/// before one reaches the wire.
|
||||
pub async fn create_function_async(
|
||||
&self,
|
||||
request: crate::function::FunctionRegistrationRequest,
|
||||
) -> Result<crate::job::Job<crate::function::FunctionVersion>> {
|
||||
request.validate()?;
|
||||
self.internal.create_function_async(request).await
|
||||
}
|
||||
|
||||
@@ -650,85 +645,6 @@ impl Connection {
|
||||
.await
|
||||
}
|
||||
|
||||
/// Create a named Secret in this database.
|
||||
///
|
||||
/// Fails if the name is taken, so a create can never silently become a
|
||||
/// rotation. There is no API that reads a stored credential back; the only
|
||||
/// consumer is a Function that binds the Secret by name. Local databases
|
||||
/// return [`Error::NotSupported`].
|
||||
pub async fn create_secret(
|
||||
&self,
|
||||
name: impl AsRef<str>,
|
||||
value: impl AsRef<str>,
|
||||
namespace_path: &[String],
|
||||
) -> Result<()> {
|
||||
let value = value.as_ref();
|
||||
crate::function::validate_secret_value(value)?;
|
||||
self.internal
|
||||
.create_secret(name.as_ref(), value, namespace_path)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Replace the credential behind an existing Secret.
|
||||
///
|
||||
/// Fails if it does not exist. Every Function bound to the Secret resolves
|
||||
/// the new value from its next execution, and no new Function version is
|
||||
/// minted -- which is what lets a rotation reach columns pinned to a
|
||||
/// version registered before it. Local databases return
|
||||
/// [`Error::NotSupported`].
|
||||
pub async fn alter_secret(
|
||||
&self,
|
||||
name: impl AsRef<str>,
|
||||
value: impl AsRef<str>,
|
||||
namespace_path: &[String],
|
||||
) -> Result<()> {
|
||||
let value = value.as_ref();
|
||||
crate::function::validate_secret_value(value)?;
|
||||
self.internal
|
||||
.alter_secret(name.as_ref(), value, namespace_path)
|
||||
.await
|
||||
}
|
||||
|
||||
/// The names of every Secret in this database.
|
||||
///
|
||||
/// Names only. No path in this API returns a stored credential, by
|
||||
/// construction rather than by policy. Local databases return
|
||||
/// [`Error::NotSupported`].
|
||||
pub async fn list_secrets(&self, namespace_path: &[String]) -> Result<Vec<String>> {
|
||||
self.internal.list_secrets(namespace_path).await
|
||||
}
|
||||
|
||||
/// Drop a Secret.
|
||||
///
|
||||
/// Functions bound to it fail at their next job, naming the Secret; that
|
||||
/// is the revocation path. The name becomes free to reuse, and a new
|
||||
/// Secret under it is picked up by everything still bound to that name.
|
||||
/// Local databases return [`Error::NotSupported`].
|
||||
pub async fn drop_secret(
|
||||
&self,
|
||||
name: impl AsRef<str>,
|
||||
namespace_path: &[String],
|
||||
) -> Result<()> {
|
||||
self.internal
|
||||
.drop_secret(name.as_ref(), namespace_path)
|
||||
.await
|
||||
}
|
||||
|
||||
/// What this database records about one Secret: its name and timestamps.
|
||||
///
|
||||
/// Never the value. The type it returns has no field for one, so this is a
|
||||
/// property of the API rather than of what the caller chooses to read.
|
||||
/// Local databases return [`Error::NotSupported`].
|
||||
pub async fn describe_secret(
|
||||
&self,
|
||||
name: impl AsRef<str>,
|
||||
namespace_path: &[String],
|
||||
) -> Result<SecretInfo> {
|
||||
self.internal
|
||||
.describe_secret(name.as_ref(), namespace_path)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Rename a table in the database.
|
||||
///
|
||||
/// This is only supported in LanceDB Cloud.
|
||||
|
||||
@@ -249,29 +249,9 @@ fn function_catalog_not_supported<T>() -> Result<T> {
|
||||
})
|
||||
}
|
||||
|
||||
fn secret_catalog_not_supported<T>() -> Result<T> {
|
||||
Err(crate::error::Error::NotSupported {
|
||||
message: "Secret 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.
|
||||
/// What a database records about a Secret. Never its value.
|
||||
///
|
||||
/// Returned by [`crate::connection::Connection::describe_secret`]. There is no
|
||||
/// field for the credential and no method that could produce one.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, serde::Deserialize)]
|
||||
pub struct SecretInfo {
|
||||
/// The Secret's database-scoped name.
|
||||
pub name: String,
|
||||
/// When the Secret was created, as an RFC 3339 timestamp.
|
||||
pub created_at: String,
|
||||
/// When the Secret's value was last rotated, as an RFC 3339 timestamp.
|
||||
pub updated_at: String,
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
pub trait Database:
|
||||
Send + Sync + std::any::Any + std::fmt::Debug + std::fmt::Display + 'static
|
||||
@@ -337,44 +317,6 @@ pub trait Database:
|
||||
async fn drop_function(&self, _name: &str, _version: &str) -> Result<bool> {
|
||||
function_catalog_not_supported()
|
||||
}
|
||||
/// Create a named Secret in this database. Fails if the name is taken, so
|
||||
/// a create can never silently become a rotation.
|
||||
async fn create_secret(
|
||||
&self,
|
||||
_name: &str,
|
||||
_value: &str,
|
||||
_namespace_path: &[String],
|
||||
) -> Result<()> {
|
||||
secret_catalog_not_supported()
|
||||
}
|
||||
/// Replace the credential behind an existing Secret. Fails if it does not
|
||||
/// exist. Every Function bound to it resolves the new value from its next
|
||||
/// execution, with no new Function version.
|
||||
async fn alter_secret(
|
||||
&self,
|
||||
_name: &str,
|
||||
_value: &str,
|
||||
_namespace_path: &[String],
|
||||
) -> Result<()> {
|
||||
secret_catalog_not_supported()
|
||||
}
|
||||
/// The names of every Secret in this database.
|
||||
///
|
||||
/// Names only. No API path returns a stored credential, by construction
|
||||
/// rather than by policy.
|
||||
async fn list_secrets(&self, _namespace_path: &[String]) -> Result<Vec<String>> {
|
||||
secret_catalog_not_supported()
|
||||
}
|
||||
/// Drop a Secret. Functions bound to it fail at their next job, which is
|
||||
/// the revocation path.
|
||||
async fn drop_secret(&self, _name: &str, _namespace_path: &[String]) -> Result<()> {
|
||||
secret_catalog_not_supported()
|
||||
}
|
||||
/// What the database records about one Secret: its name and timestamps,
|
||||
/// never its value.
|
||||
async fn describe_secret(&self, _name: &str, _namespace_path: &[String]) -> Result<SecretInfo> {
|
||||
secret_catalog_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.
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
|
||||
//! Namespace-based database implementation that delegates table management to lance-namespace
|
||||
|
||||
use lance_datafusion::utils::StreamingWriteSource;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
@@ -304,6 +305,10 @@ impl Database for LanceNamespaceDatabase {
|
||||
}
|
||||
|
||||
async fn create_table(&self, request: DbCreateTableRequest) -> Result<Arc<dyn BaseTable>> {
|
||||
// Refuse a bad declaration before the namespace records a table.
|
||||
crate::table::computed_columns::ensure_declarations_are_planned(
|
||||
&request.data.arrow_schema(),
|
||||
)?;
|
||||
let mut table_id = request.namespace_path.clone();
|
||||
table_id.push(request.name.clone());
|
||||
let mut existing_table = None;
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
//! backend-neutral terminal result of a computed-column refresh.
|
||||
//!
|
||||
//! This module contains client/wire values only. Catalog persistence,
|
||||
//! environment bake, secret resolution, and execution are owned by Sophon.
|
||||
//! environment bake, and execution are owned by Sophon.
|
||||
|
||||
use std::collections::BTreeMap;
|
||||
|
||||
@@ -409,8 +409,6 @@ pub struct FunctionVersion {
|
||||
runtime: PythonRuntimeSpec,
|
||||
runtime_digest: String,
|
||||
environment_digest: String,
|
||||
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
|
||||
secret_env_bindings: BTreeMap<String, SecretReference>,
|
||||
created_at: String,
|
||||
}
|
||||
|
||||
@@ -443,16 +441,6 @@ impl FunctionVersion {
|
||||
&self.environment_digest
|
||||
}
|
||||
|
||||
/// Declared environment variable name to the Secret each one resolves.
|
||||
///
|
||||
/// Bindings are part of this version's identity; the credentials behind
|
||||
/// them are not, and resolve at execution. Rotating a bound Secret
|
||||
/// therefore changes what the same version runs with, and no value has a
|
||||
/// field in this model.
|
||||
pub fn secret_env_bindings(&self) -> &BTreeMap<String, SecretReference> {
|
||||
&self.secret_env_bindings
|
||||
}
|
||||
|
||||
pub fn created_at(&self) -> &str {
|
||||
&self.created_at
|
||||
}
|
||||
@@ -493,203 +481,13 @@ pub struct FunctionArtifactRequest {
|
||||
pub adapter: PythonAdapterSpec,
|
||||
}
|
||||
|
||||
/// A Function binds at most this many Secrets to environment variables.
|
||||
///
|
||||
/// Each bound Secret is one extra read on the launch path of every fragment, so
|
||||
/// the count needs a bound for the same reason a credential needs a size limit.
|
||||
pub const MAX_FUNCTION_SECRET_ENV_BINDINGS: usize = 16;
|
||||
|
||||
/// Where a Secret lives, carried as its parts rather than as one string.
|
||||
///
|
||||
/// A joined id would need a delimiter, and a delimiter has to be excluded from
|
||||
/// every name and segment forever, agreed on by both sides, and re-agreed each
|
||||
/// time either grows a new way to be configured. Naming the parts costs one
|
||||
/// object and settles all of that: nothing here is parsed, so nothing can parse
|
||||
/// two ways.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct SecretReference {
|
||||
pub name: String,
|
||||
/// The namespace holding the Secret. Empty is the root, and is omitted from
|
||||
/// the wire so a root binding carries no trace of a feature it does not use.
|
||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||
pub namespace_path: Vec<String>,
|
||||
}
|
||||
|
||||
impl SecretReference {
|
||||
/// A Secret in the root namespace.
|
||||
pub fn new(name: impl Into<String>) -> Self {
|
||||
Self {
|
||||
name: name.into(),
|
||||
namespace_path: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
/// A Secret in `namespace_path`.
|
||||
pub fn in_namespace(name: impl Into<String>, namespace_path: Vec<String>) -> Self {
|
||||
Self {
|
||||
name: name.into(),
|
||||
namespace_path,
|
||||
}
|
||||
}
|
||||
|
||||
fn validate(&self) -> Result<()> {
|
||||
validate_secret_component("Secret name", &self.name)?;
|
||||
for segment in &self.namespace_path {
|
||||
validate_secret_component("Secret namespace path segment", segment)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
/// A Secret name or one namespace path segment.
|
||||
///
|
||||
/// Periods are legal here and delimiters are not a concern: a reference is
|
||||
/// never joined into one string, so the only rule left is the character set the
|
||||
/// service stores.
|
||||
fn validate_secret_component(what: &str, value: &str) -> Result<()> {
|
||||
if value.is_empty() || value.len() > MAX_SECRET_NAME_BYTES {
|
||||
return Err(Error::InvalidInput {
|
||||
message: format!(
|
||||
"{what} must be 1..={MAX_SECRET_NAME_BYTES} bytes, got {}",
|
||||
value.len()
|
||||
),
|
||||
});
|
||||
}
|
||||
if let Some(bad) = value
|
||||
.chars()
|
||||
.find(|c| !c.is_ascii_alphanumeric() && *c != '_' && *c != '-' && *c != '.')
|
||||
{
|
||||
return Err(Error::InvalidInput {
|
||||
message: format!("{what} must match [A-Za-z0-9_.-], and {bad:?} does not"),
|
||||
});
|
||||
}
|
||||
// RFC 1123's shape, which is Kubernetes' rule for object names: it rules out
|
||||
// `.` and `..` and anything reading as a hidden file or a path fragment.
|
||||
let edges_are_alphanumeric = value
|
||||
.chars()
|
||||
.next()
|
||||
.is_some_and(|c| c.is_ascii_alphanumeric())
|
||||
&& value
|
||||
.chars()
|
||||
.next_back()
|
||||
.is_some_and(|c| c.is_ascii_alphanumeric());
|
||||
if !edges_are_alphanumeric {
|
||||
return Err(Error::InvalidInput {
|
||||
message: format!(
|
||||
"{what} must start and end with a letter or digit, and '{value}' does not"
|
||||
),
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Longest Secret name or namespace path segment, matching the service.
|
||||
pub const MAX_SECRET_NAME_BYTES: usize = 255;
|
||||
|
||||
/// Largest credential a Secret may hold, matching the limit the service
|
||||
/// enforces. Bounded because the value is destined for a process environment.
|
||||
pub const MAX_SECRET_VALUE_BYTES: usize = 64 * 1024;
|
||||
|
||||
/// Whether `name` is a portable POSIX environment variable name.
|
||||
///
|
||||
/// Leading letter or underscore, then letters, digits, or underscores. Names
|
||||
/// reserved by the execution sandbox are deliberately not checked here: that
|
||||
/// list belongs to the runtime that owns it, and a copy in the client would
|
||||
/// drift from it silently.
|
||||
fn is_portable_env_name(name: &str) -> bool {
|
||||
let mut bytes = name.bytes();
|
||||
bytes
|
||||
.next()
|
||||
.is_some_and(|byte| byte == b'_' || byte.is_ascii_alphabetic())
|
||||
&& bytes.all(|byte| byte == b'_' || byte.is_ascii_alphanumeric())
|
||||
}
|
||||
|
||||
/// Reject a credential the service would refuse on size alone.
|
||||
///
|
||||
/// Checked before the request body is built, so an oversized value is never
|
||||
/// serialized or uploaded.
|
||||
pub(crate) fn validate_secret_value(value: &str) -> Result<()> {
|
||||
if value.is_empty() {
|
||||
return Err(Error::InvalidInput {
|
||||
message: "a Secret value must not be empty".to_string(),
|
||||
});
|
||||
}
|
||||
if value.contains('\0') {
|
||||
return Err(Error::InvalidInput {
|
||||
message: "a Secret value must not contain NUL".to_string(),
|
||||
});
|
||||
}
|
||||
if value.len() > MAX_SECRET_VALUE_BYTES {
|
||||
return Err(Error::InvalidInput {
|
||||
message: format!(
|
||||
"a Secret value is at most {MAX_SECRET_VALUE_BYTES} bytes, not {}",
|
||||
value.len()
|
||||
),
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Stable request envelope for remote immutable Function registration.
|
||||
///
|
||||
/// Credential values deliberately have no field here. The only secret-shaped
|
||||
/// thing a client sends is `secret_env_bindings`: the name of a Secret the
|
||||
/// database already holds, which Sophon resolves inside the remote runtime.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct FunctionRegistrationRequest {
|
||||
pub name: String,
|
||||
pub artifact: FunctionArtifactRequest,
|
||||
pub signature: FunctionSignature,
|
||||
pub runtime: PythonRuntimeSpec,
|
||||
/// Declared environment variable name to the Secret it binds. A binding is
|
||||
/// a reference: whether the Secret exists is answered when a column is
|
||||
/// declared against this version, not here.
|
||||
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
|
||||
pub secret_env_bindings: BTreeMap<String, SecretReference>,
|
||||
}
|
||||
|
||||
impl FunctionRegistrationRequest {
|
||||
/// Reject a registration whose bindings exceed what a launch can deliver.
|
||||
///
|
||||
/// Shape only, and deliberately not a check that each bound Secret exists:
|
||||
/// that is the service's answer, and it is asked for the first time when a
|
||||
/// column is declared against the registered version.
|
||||
pub fn validate(&self) -> Result<()> {
|
||||
if self.secret_env_bindings.len() > MAX_FUNCTION_SECRET_ENV_BINDINGS {
|
||||
return Err(Error::InvalidInput {
|
||||
message: format!(
|
||||
"a Function binds at most {MAX_FUNCTION_SECRET_ENV_BINDINGS} secrets, not {}",
|
||||
self.secret_env_bindings.len()
|
||||
),
|
||||
});
|
||||
}
|
||||
for (variable, secret) in &self.secret_env_bindings {
|
||||
secret.validate()?;
|
||||
if !is_portable_env_name(variable) {
|
||||
return Err(Error::InvalidInput {
|
||||
message: format!(
|
||||
"secret_env_bindings key '{variable}' is not a portable \
|
||||
environment variable name"
|
||||
),
|
||||
});
|
||||
}
|
||||
// `env` travels with the Function and is readable wherever its
|
||||
// record is; a bound Secret is not. One name carrying both would
|
||||
// resolve by delivery order, so refuse rather than pick.
|
||||
if self
|
||||
.runtime
|
||||
.env()
|
||||
.is_some_and(|env| env.contains_key(variable))
|
||||
{
|
||||
return Err(Error::InvalidInput {
|
||||
message: format!(
|
||||
"secret_env_bindings key '{variable}' is already set by runtime.env"
|
||||
),
|
||||
});
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
impl_json!(FunctionRegistrationRequest);
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -24,8 +24,8 @@ use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
use arrow_array::cast::AsArray;
|
||||
use arrow_array::types::UInt64Type;
|
||||
use arrow_array::{RecordBatch, UInt64Array};
|
||||
use arrow_schema::{Schema as ArrowSchema, SchemaRef};
|
||||
use arrow_array::{RecordBatch, UInt64Array, new_null_array};
|
||||
use arrow_schema::{FieldRef, Schema as ArrowSchema, SchemaRef};
|
||||
use datafusion::common::ScalarValue;
|
||||
use datafusion::error::DataFusionError;
|
||||
use datafusion::physical_plan::SendableRecordBatchStream;
|
||||
@@ -34,7 +34,7 @@ use datafusion::prelude::{col, lit};
|
||||
use futures::{StreamExt, TryStreamExt};
|
||||
use lance::Dataset;
|
||||
use lance::dataset::mem_wal::DatasetMemWalExt;
|
||||
use lance::dataset::transaction::{Operation, Transaction};
|
||||
use lance::dataset::transaction::{Operation, Transaction, UpdateMode};
|
||||
use lance::dataset::write::delete::DeleteBuilder;
|
||||
use lance::dataset::write::merge_insert::inserted_rows::{
|
||||
KeyExistenceFilter, KeyExistenceFilterBuilder, KeyValue,
|
||||
@@ -51,6 +51,9 @@ use super::{
|
||||
definition_to_metadata,
|
||||
};
|
||||
use crate::database::OpenTableRequest;
|
||||
use crate::table::computed_columns::{
|
||||
computed_column_from_field, computed_columns, ensure_declarations_are_planned,
|
||||
};
|
||||
use crate::table::{NativeTable, NativeTableExt, Table};
|
||||
use crate::{Error, Result};
|
||||
|
||||
@@ -167,30 +170,52 @@ pub(crate) async fn execute_refresh(
|
||||
.map(|p| (p.output.clone(), p.expression.clone()))
|
||||
.collect();
|
||||
validate_inputs(&source_ds, definition)?;
|
||||
let (replanned, mut planned_fields, _renames) = super::plan(
|
||||
let (replanned, planned_fields, _renames) = super::plan(
|
||||
source_schema,
|
||||
&definition.source_table,
|
||||
&definition.source_namespace,
|
||||
&projections,
|
||||
Some(&projections),
|
||||
definition.filter.as_deref(),
|
||||
definition.limit,
|
||||
)?;
|
||||
let mut planned_fields = planned_fields;
|
||||
planned_fields.push(arrow_schema::Field::new(
|
||||
SOURCE_ROW_ID_COLUMN,
|
||||
arrow_schema::DataType::UInt64,
|
||||
false,
|
||||
));
|
||||
// A computed column is not planned from the source: refresh writes it
|
||||
// NULL and its declaration's owner fills it. Its declaration must still
|
||||
// be complete, and it must be able to hold NULL.
|
||||
let physical = ArrowSchema::from(view_ds.schema());
|
||||
let planned_shape: Vec<_> = planned_fields
|
||||
.iter()
|
||||
.map(|f| (f.name().clone(), f.data_type().clone(), f.is_nullable()))
|
||||
.collect();
|
||||
let physical_shape: Vec<_> = physical
|
||||
let mut computed = computed_columns(&physical).into_iter().map(|c| c.name);
|
||||
if let Some(name) = computed.by_ref().find(|name| {
|
||||
physical
|
||||
.field_with_name(name)
|
||||
.is_ok_and(|f| !f.is_nullable())
|
||||
}) {
|
||||
return Err(Error::Schema {
|
||||
message: format!(
|
||||
"computed column '{name}' of view '{}' cannot hold NULL; recreate the view",
|
||||
view.name()
|
||||
),
|
||||
});
|
||||
}
|
||||
ensure_declarations_are_planned(&physical)?;
|
||||
let physical_planned: Vec<&FieldRef> = physical
|
||||
.fields()
|
||||
.iter()
|
||||
.map(|f| (f.name().clone(), f.data_type().clone(), f.is_nullable()))
|
||||
.filter(|f| computed_column_from_field(f).is_none())
|
||||
.collect();
|
||||
if planned_shape != physical_shape {
|
||||
// A projected column that became nullable at the source still fits the
|
||||
// view's nullable field; the reverse would not.
|
||||
let matches = planned_fields.len() == physical_planned.len()
|
||||
&& planned_fields.iter().zip(&physical_planned).all(|(e, p)| {
|
||||
e.name() == p.name()
|
||||
&& e.data_type() == p.data_type()
|
||||
&& (p.is_nullable() || !e.is_nullable())
|
||||
});
|
||||
if !matches {
|
||||
return Err(Error::Schema {
|
||||
message: format!(
|
||||
"the stored definition of view '{}' does not produce this \
|
||||
@@ -229,11 +254,18 @@ pub(crate) async fn execute_refresh(
|
||||
.get(SOURCE_VERSION_TS_META_KEY)
|
||||
.and_then(|raw| raw.parse().ok());
|
||||
// The watermark speaks only for the view state its refresh left behind;
|
||||
// any other commit on the view since then is drift.
|
||||
let view_intact = metadata
|
||||
// any other commit on the view since then is drift, except a fill of its
|
||||
// computed columns, which rewrites nothing refresh certifies.
|
||||
let recorded_view_version = metadata
|
||||
.get(VIEW_VERSION_META_KEY)
|
||||
.and_then(|raw| raw.parse::<u64>().ok())
|
||||
== Some(view_ds.version().version);
|
||||
.and_then(|raw| raw.parse::<u64>().ok());
|
||||
let view_intact = match recorded_view_version {
|
||||
Some(recorded) if recorded == view_ds.version().version => true,
|
||||
Some(recorded) if recorded < view_ds.version().version => {
|
||||
only_computed_rewrites_since(&view_ds, recorded).await?
|
||||
}
|
||||
_ => false,
|
||||
};
|
||||
|
||||
if !full && watermark == Some(source_version) && view_intact && recorded_ts == Some(source_ts) {
|
||||
return Ok(RefreshMaterializedViewResult {
|
||||
@@ -1090,6 +1122,83 @@ struct RowScope {
|
||||
limit: Option<u64>,
|
||||
}
|
||||
|
||||
/// Whether every commit on the view after `recorded` is a fill of its
|
||||
/// computed columns: a column rewrite or data replacement touching only
|
||||
/// those fields and neither adding nor removing rows, or the freshness
|
||||
/// stamp a fill leaves on them. A version whose transaction cannot be read
|
||||
/// is not proven, so it counts as drift.
|
||||
async fn only_computed_rewrites_since(view_ds: &Dataset, recorded: u64) -> Result<bool> {
|
||||
// A fill may write any field under a computed column, so the whole
|
||||
// subtree counts, not only the root.
|
||||
let physical = ArrowSchema::from(view_ds.schema());
|
||||
fn subtree(field: &lance_core::datatypes::Field, ids: &mut Vec<u32>) {
|
||||
ids.push(field.id as u32);
|
||||
for child in &field.children {
|
||||
subtree(child, ids);
|
||||
}
|
||||
}
|
||||
let mut computed_fields = Vec::new();
|
||||
for column in computed_columns(&physical) {
|
||||
if let Some(field) = view_ds.schema().field(&column.name) {
|
||||
subtree(field, &mut computed_fields);
|
||||
}
|
||||
}
|
||||
if computed_fields.is_empty() {
|
||||
return Ok(false);
|
||||
}
|
||||
for version in recorded + 1..=view_ds.version().version {
|
||||
let Some(transaction) = view_ds.read_transaction_by_version(version).await? else {
|
||||
return Ok(false);
|
||||
};
|
||||
let fill = match &transaction.operation {
|
||||
Operation::Update {
|
||||
removed_fragment_ids,
|
||||
new_fragments,
|
||||
fields_modified,
|
||||
update_mode: Some(UpdateMode::RewriteColumns),
|
||||
..
|
||||
} => {
|
||||
removed_fragment_ids.is_empty()
|
||||
&& new_fragments.is_empty()
|
||||
&& !fields_modified.is_empty()
|
||||
&& fields_modified
|
||||
.iter()
|
||||
.all(|field| computed_fields.contains(field))
|
||||
}
|
||||
// What `refresh_column` commits for a SQL declaration.
|
||||
Operation::DataReplacement { replacements } => {
|
||||
!replacements.is_empty()
|
||||
&& replacements.iter().all(|group| {
|
||||
!group.1.fields.is_empty()
|
||||
&& group
|
||||
.1
|
||||
.fields
|
||||
.iter()
|
||||
.all(|field| computed_fields.contains(&(*field as u32)))
|
||||
})
|
||||
}
|
||||
// The stamp `refresh_column` writes after its fill (see
|
||||
// `table::freshness`): field metadata on computed columns, no data.
|
||||
Operation::UpdateConfig {
|
||||
config_updates: None,
|
||||
table_metadata_updates: None,
|
||||
schema_metadata_updates: None,
|
||||
field_metadata_updates,
|
||||
} => {
|
||||
!field_metadata_updates.is_empty()
|
||||
&& field_metadata_updates
|
||||
.keys()
|
||||
.all(|field| computed_fields.contains(&(*field as u32)))
|
||||
}
|
||||
_ => false,
|
||||
};
|
||||
if !fill {
|
||||
return Ok(false);
|
||||
}
|
||||
}
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
async fn compute_stream(
|
||||
source: &Dataset,
|
||||
definition: &MaterializedViewDefinition,
|
||||
@@ -1158,6 +1267,10 @@ async fn compute_stream(
|
||||
let batch = batch.map_err(|e| DataFusionError::External(Box::new(e)))?;
|
||||
let mut columns = Vec::with_capacity(out_schema.fields().len());
|
||||
for field in out_schema.fields() {
|
||||
if computed_column_from_field(field).is_some() {
|
||||
columns.push(new_null_array(field.data_type(), batch.num_rows()));
|
||||
continue;
|
||||
}
|
||||
let name = if field.name() == SOURCE_ROW_ID_COLUMN {
|
||||
ROW_ID
|
||||
} else {
|
||||
@@ -2768,7 +2881,7 @@ mod tests {
|
||||
let (conn, source) = db_with_source(vec![1]).await;
|
||||
let prepared = crate::materialized_view::prepare_declaration(
|
||||
&source,
|
||||
&[("x".into(), "x".into()), ("twice".into(), "x * 2".into())],
|
||||
Some(&[("x".into(), "x".into()), ("twice".into(), "x * 2".into())]),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
@@ -3132,4 +3245,464 @@ mod tests {
|
||||
let err = view.refresh().execute().await.unwrap_err();
|
||||
assert!(err.to_string().contains("source table 'src'"), "{err}");
|
||||
}
|
||||
|
||||
/// A view with a computed column, declared over `people` and refreshed.
|
||||
async fn refreshed_computed_view(conn: &Connection) -> MaterializedView {
|
||||
use crate::materialized_view::tests::{computed_field, people, test_binding};
|
||||
let source = people(conn).await;
|
||||
let view = crate::materialized_view::prepare_declaration(
|
||||
&source,
|
||||
Some(&[
|
||||
("id".to_string(), "id".to_string()),
|
||||
("name".to_string(), "name".to_string()),
|
||||
]),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.with_computed_columns(
|
||||
vec![(2, computed_field("emb", "fb_1", "name"))],
|
||||
&[test_binding("fb_1", "name", "emb")],
|
||||
)
|
||||
.unwrap()
|
||||
.create("v")
|
||||
.await
|
||||
.unwrap();
|
||||
let result = view.refresh().execute().await.unwrap();
|
||||
assert_eq!(result.mode, RefreshMode::Rebuild);
|
||||
view
|
||||
}
|
||||
|
||||
async fn unfilled(view: &MaterializedView) -> usize {
|
||||
view.table()
|
||||
.count_rows(Some("emb IS NULL".to_string()))
|
||||
.await
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
async fn append_people(conn: &Connection, ids: Vec<i32>, names: Vec<&str>) {
|
||||
let batch = record_batch!(("id", Int32, ids), ("name", Utf8, names)).unwrap();
|
||||
conn.open_table("people")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap()
|
||||
.add(batch)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// Commit the fill job's shape on the view: a column rewrite of
|
||||
/// `fields`, touching no rows. The data is left as it is; what matters
|
||||
/// here is how the next refresh classifies the commit.
|
||||
async fn commit_column_rewrite(view: &MaterializedView, fields: &[&str]) {
|
||||
let native = view.table().as_native().unwrap();
|
||||
native.dataset.reload().await.unwrap();
|
||||
let dataset = native.dataset.get().await.unwrap().as_ref().clone();
|
||||
let fields_modified = fields
|
||||
.iter()
|
||||
.map(|name| dataset.schema().field(name).unwrap().id as u32)
|
||||
.collect();
|
||||
let updated_fragments = dataset
|
||||
.get_fragments()
|
||||
.iter()
|
||||
.map(|fragment| fragment.metadata().clone())
|
||||
.collect();
|
||||
let operation = Operation::Update {
|
||||
removed_fragment_ids: Vec::new(),
|
||||
updated_fragments,
|
||||
new_fragments: Vec::new(),
|
||||
fields_modified,
|
||||
compacted_sstables: Vec::new(),
|
||||
fields_for_preserving_frag_bitmap: Vec::new(),
|
||||
update_mode: Some(UpdateMode::RewriteColumns),
|
||||
inserted_rows_filter: None,
|
||||
updated_fragment_offsets: None,
|
||||
};
|
||||
let read_version = dataset.version().version;
|
||||
CommitBuilder::new(WriteDestination::Dataset(Arc::new(dataset)))
|
||||
.execute(Transaction::new(read_version, operation, None))
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// Refresh never computes a computed column: every row it writes, on a
|
||||
/// rebuild, an append and a rewrite, carries NULL there, and the
|
||||
/// declaration survives all three.
|
||||
#[tokio::test]
|
||||
async fn test_computed_columns_are_written_null_and_kept() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
assert_eq!(unfilled(&view).await, 3);
|
||||
|
||||
append_people(&conn, vec![4, 5], vec!["d", "e"]).await;
|
||||
let result = view.refresh().execute().await.unwrap();
|
||||
assert_eq!(result.mode, RefreshMode::Incremental);
|
||||
assert_eq!(unfilled(&view).await, 5);
|
||||
|
||||
conn.open_table("people")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap()
|
||||
.update()
|
||||
.column("name", "'z'")
|
||||
.only_if("id = 1")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
view.refresh().execute().await.unwrap();
|
||||
assert_eq!(unfilled(&view).await, 5);
|
||||
assert_eq!(read(view.table(), "id").await, vec![1, 2, 3, 4, 5]);
|
||||
|
||||
let schema = view.table().schema().await.unwrap();
|
||||
assert!(
|
||||
crate::table::computed_columns::function_bindings(&schema)
|
||||
.unwrap()
|
||||
.iter()
|
||||
.any(|b| b.binding_id() == "fb_1"),
|
||||
"the binding envelope was lost"
|
||||
);
|
||||
assert!(
|
||||
computed_column_from_field(schema.field_with_name("emb").unwrap()).is_some(),
|
||||
"the declaration was lost"
|
||||
);
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
}
|
||||
|
||||
/// Field metadata on `field` only, the commit shape of the freshness
|
||||
/// stamp `refresh_column` leaves after its fill.
|
||||
async fn commit_field_metadata(view: &MaterializedView, field: &str, key: &str) {
|
||||
let native = view.table().as_native().unwrap();
|
||||
native.dataset.reload().await.unwrap();
|
||||
let mut dataset = native.dataset.get().await.unwrap().as_ref().clone();
|
||||
dataset
|
||||
.update_field_metadata()
|
||||
.update(field, [(key.to_string(), "{}".to_string())])
|
||||
.unwrap()
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// The stamp is metadata on the computed column and rewrites nothing
|
||||
/// refresh certifies, so it is not drift; the same commit shape on a
|
||||
/// projected column is, like any other write to it.
|
||||
#[tokio::test]
|
||||
async fn test_a_freshness_stamp_is_not_drift() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
|
||||
commit_field_metadata(
|
||||
&view,
|
||||
"emb",
|
||||
crate::table::computed_columns::SOURCE_SIGNATURE_META_KEY,
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
|
||||
commit_field_metadata(&view, "id", "probe").await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Rebuild
|
||||
);
|
||||
}
|
||||
|
||||
/// The fill job's commit rewrites only computed columns. It is the one
|
||||
/// commit on a view that is not drift: the next refresh carries on from
|
||||
/// its watermark instead of rebuilding, which would null what the fill
|
||||
/// just wrote.
|
||||
#[tokio::test]
|
||||
async fn test_a_computed_column_fill_is_not_drift() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
|
||||
commit_column_rewrite(&view, &["emb"]).await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
|
||||
commit_column_rewrite(&view, &["emb"]).await;
|
||||
append_people(&conn, vec![4], vec!["d"]).await;
|
||||
let result = view.refresh().execute().await.unwrap();
|
||||
assert_eq!(result.mode, RefreshMode::Incremental);
|
||||
assert_eq!(result.rows_written, 1);
|
||||
assert_eq!(read(view.table(), "id").await, vec![1, 2, 3, 4]);
|
||||
}
|
||||
|
||||
/// A column rewrite that reaches a projected column is drift like any
|
||||
/// other write: refresh certifies those columns and must recompute them.
|
||||
#[tokio::test]
|
||||
async fn test_a_rewrite_of_a_projected_column_is_drift() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
|
||||
commit_column_rewrite(&view, &["emb", "name"]).await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Rebuild
|
||||
);
|
||||
}
|
||||
|
||||
/// The declaration contract is checked before any refresh mutation: a
|
||||
/// missing binding envelope and a column that lost its declaration both
|
||||
/// fail closed.
|
||||
#[tokio::test]
|
||||
async fn test_a_broken_declaration_is_refused_before_refresh() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
let native = view.table().as_native().unwrap();
|
||||
let mut dataset = native.dataset.get().await.unwrap().as_ref().clone();
|
||||
dataset
|
||||
.update_schema_metadata(vec![(
|
||||
crate::table::computed_columns::FUNCTION_BINDINGS_META_KEY.to_string(),
|
||||
None,
|
||||
)])
|
||||
.await
|
||||
.unwrap();
|
||||
let err = view.refresh().execute().await.unwrap_err().to_string();
|
||||
assert!(err.contains("references missing binding 'fb_1'"), "{err}");
|
||||
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
let native = view.table().as_native().unwrap();
|
||||
let mut dataset = native.dataset.get().await.unwrap().as_ref().clone();
|
||||
dataset
|
||||
.replace_field_metadata(vec![(
|
||||
dataset.schema().field("emb").unwrap().id as u32,
|
||||
HashMap::new(),
|
||||
)])
|
||||
.await
|
||||
.unwrap();
|
||||
let err = view.refresh().execute().await.unwrap_err().to_string();
|
||||
assert!(err.contains("does not match binding 'fb_1'"), "{err}");
|
||||
}
|
||||
|
||||
/// An input the view does not project is materialized on every refresh
|
||||
/// path, before the provenance column, with the source's values.
|
||||
#[tokio::test]
|
||||
async fn test_internal_inputs_are_materialized_and_refreshed() {
|
||||
use crate::materialized_view::tests::{computed_field, strict_people, test_binding};
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let source = strict_people(&conn).await;
|
||||
let mut prepared = crate::materialized_view::prepare_declaration(
|
||||
&source,
|
||||
Some(&[("id".to_string(), "id".to_string())]),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let input = prepared.input_column("name").unwrap();
|
||||
let view = prepared
|
||||
.with_computed_columns(
|
||||
vec![(1, computed_field("emb", "fb_1", &input))],
|
||||
&[test_binding("fb_1", &input, "emb")],
|
||||
)
|
||||
.unwrap()
|
||||
.create("v")
|
||||
.await
|
||||
.unwrap();
|
||||
let names: Vec<String> = view
|
||||
.table()
|
||||
.schema()
|
||||
.await
|
||||
.unwrap()
|
||||
.fields()
|
||||
.iter()
|
||||
.map(|f| f.name().clone())
|
||||
.collect();
|
||||
assert_eq!(names, ["id", "emb", "__input_name", SOURCE_ROW_ID_COLUMN]);
|
||||
|
||||
let unfilled_inputs = || async {
|
||||
view.table()
|
||||
.count_rows(Some("__input_name IS NULL".to_string()))
|
||||
.await
|
||||
.unwrap()
|
||||
};
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Rebuild
|
||||
);
|
||||
assert_eq!(view.table().count_rows(None).await.unwrap(), 3);
|
||||
assert_eq!(unfilled_inputs().await, 0);
|
||||
|
||||
let more = arrow_array::RecordBatch::try_new(
|
||||
source.schema().await.unwrap(),
|
||||
vec![
|
||||
Arc::new(Int32Array::from(vec![4])),
|
||||
Arc::new(arrow_array::StringArray::from(vec!["d"])),
|
||||
],
|
||||
)
|
||||
.unwrap();
|
||||
source.add(more).execute().await.unwrap();
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Incremental
|
||||
);
|
||||
assert_eq!(unfilled_inputs().await, 0);
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.count_rows(Some("__input_name = 'd'".to_string()))
|
||||
.await
|
||||
.unwrap(),
|
||||
1
|
||||
);
|
||||
|
||||
source
|
||||
.update()
|
||||
.column("name", "'z'")
|
||||
.only_if("id = 1")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
view.refresh().execute().await.unwrap();
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.count_rows(Some("__input_name = 'z'".to_string()))
|
||||
.await
|
||||
.unwrap(),
|
||||
1
|
||||
);
|
||||
assert_eq!(
|
||||
unfilled(&view).await,
|
||||
4,
|
||||
"rewritten and new rows are unfilled"
|
||||
);
|
||||
}
|
||||
|
||||
/// A SQL declaration is filled by `refresh_column` on the view, which
|
||||
/// commits a data replacement and then its freshness stamp; the next
|
||||
/// refresh continues from its watermark and keeps what the fill wrote,
|
||||
/// and only rows the view added since come back unfilled.
|
||||
#[tokio::test]
|
||||
async fn test_a_sql_fill_is_not_drift() {
|
||||
use crate::materialized_view::tests::{people, sql_field};
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let source = people(&conn).await;
|
||||
let view = crate::materialized_view::prepare_declaration(
|
||||
&source,
|
||||
Some(&[("id".to_string(), "id".to_string())]),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.with_computed_columns(
|
||||
vec![(
|
||||
1,
|
||||
sql_field("next", arrow_schema::DataType::Int32, "id + 1", r#"["id"]"#),
|
||||
)],
|
||||
&[],
|
||||
)
|
||||
.unwrap()
|
||||
.create("v")
|
||||
.await
|
||||
.unwrap();
|
||||
let filled = || async {
|
||||
view.table()
|
||||
.count_rows(Some("next = id + 1".to_string()))
|
||||
.await
|
||||
.unwrap()
|
||||
};
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Rebuild
|
||||
);
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.refresh_column("next")
|
||||
.await
|
||||
.unwrap()
|
||||
.rows_filled,
|
||||
3
|
||||
);
|
||||
assert_eq!(filled().await, 3);
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
assert_eq!(filled().await, 3);
|
||||
|
||||
append_people(&conn, vec![4], vec!["d"]).await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Incremental
|
||||
);
|
||||
assert_eq!(filled().await, 3);
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.refresh_column("next")
|
||||
.await
|
||||
.unwrap()
|
||||
.rows_filled,
|
||||
1
|
||||
);
|
||||
assert_eq!(filled().await, 4);
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
}
|
||||
|
||||
/// A fill of a nested computed column writes its child fields; that is
|
||||
/// still a fill, not drift.
|
||||
#[tokio::test]
|
||||
async fn test_a_nested_sql_fill_is_not_drift() {
|
||||
use crate::materialized_view::tests::{people, sql_field};
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let source = people(&conn).await;
|
||||
let payload = sql_field(
|
||||
"payload",
|
||||
arrow_schema::DataType::Struct(
|
||||
vec![arrow_schema::Field::new(
|
||||
"value",
|
||||
arrow_schema::DataType::Utf8,
|
||||
true,
|
||||
)]
|
||||
.into(),
|
||||
),
|
||||
"named_struct('value', name)",
|
||||
r#"["name"]"#,
|
||||
);
|
||||
let view = crate::materialized_view::prepare_declaration(
|
||||
&source,
|
||||
Some(&[("name".to_string(), "name".to_string())]),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.with_computed_columns(vec![(1, payload)], &[])
|
||||
.unwrap()
|
||||
.create("v")
|
||||
.await
|
||||
.unwrap();
|
||||
view.refresh().execute().await.unwrap();
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.refresh_column("payload")
|
||||
.await
|
||||
.unwrap()
|
||||
.rows_filled,
|
||||
3
|
||||
);
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.count_rows(Some("payload.value = name".to_string()))
|
||||
.await
|
||||
.unwrap(),
|
||||
3
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -404,18 +404,6 @@ fn validate_dns_hostname(hostname: &str) -> Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Whether a route's request body is a credential rather than a description of
|
||||
/// one.
|
||||
///
|
||||
/// Matched on the path segment rather than a versioned prefix, so a `/v2/` bump
|
||||
/// or a route added under the namespace later is covered without anyone
|
||||
/// remembering to extend this. Every secrets route is denied, not only the two
|
||||
/// that carry a value: their bodies hold names and page tokens, which are worth
|
||||
/// nothing in a debug log next to the risk of a new verb landing here unnoticed.
|
||||
fn route_carries_credential(path: &str) -> bool {
|
||||
path.split('/').any(|segment| segment == "secrets")
|
||||
}
|
||||
|
||||
impl RestfulLanceDbClient<Sender> {
|
||||
fn get_timeout(passed: Option<Duration>, env_var: &str) -> Result<Option<Duration>> {
|
||||
if let Some(passed) = passed {
|
||||
@@ -622,14 +610,12 @@ impl<S: HttpSend> RestfulLanceDbClient<S> {
|
||||
) -> Result<HeaderMap> {
|
||||
let mut headers = HeaderMap::new();
|
||||
if !api_key.is_empty() {
|
||||
// `log_request` prints the request's Debug, which prints headers.
|
||||
// Marking the value sensitive is what makes that print `Sensitive`
|
||||
// instead of the key itself.
|
||||
let mut key = HeaderValue::from_str(api_key).map_err(|_| Error::InvalidInput {
|
||||
message: "non-ascii api key provided".to_string(),
|
||||
})?;
|
||||
key.set_sensitive(true);
|
||||
headers.insert(HeaderName::from_static("x-api-key"), key);
|
||||
headers.insert(
|
||||
HeaderName::from_static("x-api-key"),
|
||||
HeaderValue::from_str(api_key).map_err(|_| Error::InvalidInput {
|
||||
message: "non-ascii api key provided".to_string(),
|
||||
})?,
|
||||
);
|
||||
}
|
||||
if region == "local" {
|
||||
let host = format!("{}.local.api.lancedb.com", db_name);
|
||||
@@ -859,12 +845,7 @@ impl<S: HttpSend> RestfulLanceDbClient<S> {
|
||||
.headers()
|
||||
.get("content-type")
|
||||
.map(|v| v.to_str().unwrap());
|
||||
if route_carries_credential(request.url().path()) {
|
||||
debug!(
|
||||
"Sending request_id={}: {:?} with body suppressed",
|
||||
request_id, request
|
||||
);
|
||||
} else if content_type == Some("application/json") {
|
||||
if content_type == Some("application/json") {
|
||||
let body = request.body().as_ref().unwrap().as_bytes().unwrap();
|
||||
let body = String::from_utf8_lossy(body);
|
||||
debug!(
|
||||
@@ -1211,56 +1192,6 @@ mod tests {
|
||||
assert_eq!(headers.get("x-api-key").unwrap(), "api-key");
|
||||
}
|
||||
|
||||
/// `log_request` prints the request's Debug, and Debug for a request prints
|
||||
/// its headers. Marking the value sensitive is the only thing standing
|
||||
/// between the API key and every debug line; assert on the header map's own
|
||||
/// Debug, which is what that printing reduces to.
|
||||
#[test]
|
||||
fn test_api_key_is_redacted_in_debug_output() {
|
||||
let headers = RestfulLanceDbClient::<Sender>::default_headers(
|
||||
"sk-live-sentinel",
|
||||
"us-east-1",
|
||||
"db-name",
|
||||
false,
|
||||
&RemoteOptions::default(),
|
||||
None,
|
||||
&ClientConfig::default(),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(headers.get("x-api-key").unwrap(), "sk-live-sentinel");
|
||||
assert!(
|
||||
!format!("{:?}", headers).contains("sk-live-sentinel"),
|
||||
"the API key must not survive Debug formatting"
|
||||
);
|
||||
}
|
||||
|
||||
/// Denial follows the path segment, so a verb that does not exist yet and a
|
||||
/// future API version are both covered without an edit here.
|
||||
#[test]
|
||||
fn test_secrets_routes_never_log_a_body() {
|
||||
for route in [
|
||||
"/v1/secrets/create",
|
||||
"/v1/secrets/alter",
|
||||
"/v1/secrets/list",
|
||||
"/v1/secrets/drop",
|
||||
"/v1/secrets/describe",
|
||||
"/v2/secrets/rotate",
|
||||
] {
|
||||
assert!(
|
||||
route_carries_credential(route),
|
||||
"{route} must never log a body"
|
||||
);
|
||||
}
|
||||
for route in [
|
||||
"/v1/functions/create",
|
||||
"/v1/table/foo/query",
|
||||
"/v1/jobs/list",
|
||||
] {
|
||||
assert!(!route_carries_credential(route), "{route} is not a secret");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_rejects_invalid_cloud_dns_hostname() {
|
||||
let invalid_database_names = ["a".repeat(64), "invalid..database".to_string()];
|
||||
|
||||
@@ -21,7 +21,7 @@ use lance_namespace::models::{
|
||||
use crate::Error;
|
||||
use crate::database::{
|
||||
CloneTableRequest, CreateTableMode, CreateTableRequest, Database, DatabaseOptions, JobInfo,
|
||||
OpenTableRequest, ReadConsistency, SecretInfo, TableNamesRequest,
|
||||
OpenTableRequest, ReadConsistency, TableNamesRequest,
|
||||
};
|
||||
use crate::error::Result;
|
||||
use crate::function::{FunctionRegistrationRequest, FunctionVersion};
|
||||
@@ -277,22 +277,6 @@ pub struct RemoteHostOverrides {
|
||||
pub sql: Option<String>,
|
||||
}
|
||||
|
||||
/// Attach a namespace path to a Secret request body.
|
||||
///
|
||||
/// A root path is omitted rather than sent empty, so a root request is byte
|
||||
/// identical to one from a client that predates namespace addressing.
|
||||
fn add_namespace_path(body: &mut serde_json::Value, namespace_path: &[String]) {
|
||||
if namespace_path.is_empty() {
|
||||
return;
|
||||
}
|
||||
body["namespace_path"] = serde_json::Value::Array(
|
||||
namespace_path
|
||||
.iter()
|
||||
.map(|segment| serde_json::Value::String(segment.clone()))
|
||||
.collect(),
|
||||
);
|
||||
}
|
||||
|
||||
impl RemoteDatabase {
|
||||
pub(crate) fn try_new(
|
||||
uri: &str,
|
||||
@@ -368,28 +352,6 @@ impl RemoteDatabase {
|
||||
}
|
||||
|
||||
impl<S: HttpSend> RemoteDatabase<S> {
|
||||
/// `create` and `alter` differ only in which name state the server
|
||||
/// requires, so they share one request shape. The value is a request field
|
||||
/// and never a path segment or query parameter, which keeps it out of
|
||||
/// access logs and proxy traces.
|
||||
async fn write_secret(
|
||||
&self,
|
||||
route: &str,
|
||||
name: &str,
|
||||
value: &str,
|
||||
namespace_path: &[String],
|
||||
) -> Result<()> {
|
||||
let mut body = serde_json::json!({
|
||||
"name": name,
|
||||
"value": value,
|
||||
});
|
||||
add_namespace_path(&mut body, namespace_path);
|
||||
let req = self.client.post(route).json(&body);
|
||||
let (request_id, response) = self.client.send(req).await?;
|
||||
self.client.check_response(&request_id, response).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn submit_drop_table(
|
||||
&self,
|
||||
name: &str,
|
||||
@@ -608,21 +570,6 @@ struct RemoteDropFunctionResponse {
|
||||
dropped: bool,
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
struct RemoteListSecretsResponse {
|
||||
#[serde(default)]
|
||||
secrets: Vec<RemoteListedSecret>,
|
||||
#[serde(default)]
|
||||
page_token: Option<String>,
|
||||
}
|
||||
|
||||
/// An object rather than a bare name so a later listing can carry a Secret's
|
||||
/// type or last-updated time without breaking this one.
|
||||
#[derive(serde::Deserialize)]
|
||||
struct RemoteListedSecret {
|
||||
name: 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;
|
||||
@@ -724,72 +671,6 @@ impl<S: HttpSend> Database for RemoteDatabase<S> {
|
||||
Ok(response.dropped)
|
||||
}
|
||||
|
||||
async fn create_secret(
|
||||
&self,
|
||||
name: &str,
|
||||
value: &str,
|
||||
namespace_path: &[String],
|
||||
) -> Result<()> {
|
||||
self.write_secret("/v1/secrets/create", name, value, namespace_path)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn alter_secret(&self, name: &str, value: &str, namespace_path: &[String]) -> Result<()> {
|
||||
self.write_secret("/v1/secrets/alter", name, value, namespace_path)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn list_secrets(&self, namespace_path: &[String]) -> Result<Vec<String>> {
|
||||
let mut names = Vec::new();
|
||||
let mut page_token: Option<String> = None;
|
||||
let mut seen_page_tokens = HashSet::new();
|
||||
loop {
|
||||
let mut body = serde_json::json!({});
|
||||
if let Some(token) = &page_token {
|
||||
body["page_token"] = serde_json::Value::String(token.clone());
|
||||
}
|
||||
add_namespace_path(&mut body, namespace_path);
|
||||
let req = self.client.post("/v1/secrets/list").json(&body);
|
||||
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: RemoteListSecretsResponse =
|
||||
response.json().await.err_to_http(request_id.clone())?;
|
||||
names.extend(response.secrets.into_iter().map(|secret| secret.name));
|
||||
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: "Secret listing response repeated a page_token".into(),
|
||||
request_id,
|
||||
status_code: Some(status),
|
||||
});
|
||||
}
|
||||
page_token = Some(next_page_token);
|
||||
}
|
||||
Ok(names)
|
||||
}
|
||||
|
||||
async fn drop_secret(&self, name: &str, namespace_path: &[String]) -> Result<()> {
|
||||
let mut body = serde_json::json!({ "name": name });
|
||||
add_namespace_path(&mut body, namespace_path);
|
||||
let req = self.client.post("/v1/secrets/drop").json(&body);
|
||||
let (request_id, response) = self.client.send(req).await?;
|
||||
self.client.check_response(&request_id, response).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn describe_secret(&self, name: &str, namespace_path: &[String]) -> Result<SecretInfo> {
|
||||
let mut body = serde_json::json!({ "name": name });
|
||||
add_namespace_path(&mut body, namespace_path);
|
||||
let req = self.client.post("/v1/secrets/describe").json(&body);
|
||||
let (request_id, response) = self.client.send(req).await?;
|
||||
let response = self.client.check_response(&request_id, response).await?;
|
||||
response.json().await.err_to_http(request_id)
|
||||
}
|
||||
|
||||
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 {
|
||||
@@ -2900,113 +2781,6 @@ mod tests {
|
||||
assert_eq!(batches[0].schema(), schema);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_create_and_alter_secret_send_the_value_in_the_request_body() {
|
||||
for (route, call) in [("/v1/secrets/create", true), ("/v1/secrets/alter", false)] {
|
||||
let conn = Connection::new_with_handler(move |request| {
|
||||
assert_eq!(request.method(), &reqwest::Method::POST);
|
||||
assert_eq!(request.url().path(), route);
|
||||
// Never a path segment or query parameter, which is what keeps
|
||||
// it out of access logs and proxy traces.
|
||||
assert!(request.url().query().is_none(), "{:?}", request.url());
|
||||
let body: serde_json::Value =
|
||||
serde_json::from_slice(request.body().unwrap().as_bytes().unwrap()).unwrap();
|
||||
assert_eq!(body["name"], "openai-prod");
|
||||
assert_eq!(body["value"], "sk-live-0001");
|
||||
http::Response::builder().status(200).body("{}").unwrap()
|
||||
});
|
||||
if call {
|
||||
conn.create_secret("openai-prod", "sk-live-0001", &[])
|
||||
.await
|
||||
.unwrap();
|
||||
} else {
|
||||
conn.alter_secret("openai-prod", "sk-live-0001", &[])
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_list_secrets_walks_pages_and_returns_names_only() {
|
||||
let conn = Connection::new_with_handler(|request| {
|
||||
assert_eq!(request.url().path(), "/v1/secrets/list");
|
||||
let body: serde_json::Value =
|
||||
serde_json::from_slice(request.body().unwrap().as_bytes().unwrap()).unwrap();
|
||||
let page = body.get("page_token").and_then(|token| token.as_str());
|
||||
let body = match page {
|
||||
None => r#"{"secrets":[{"name":"openai-prod"}],"page_token":"p2"}"#,
|
||||
Some("p2") => r#"{"secrets":[{"name":"hf-prod"}]}"#,
|
||||
Some(other) => panic!("unexpected page token: {other}"),
|
||||
};
|
||||
http::Response::builder().status(200).body(body).unwrap()
|
||||
});
|
||||
assert_eq!(
|
||||
conn.list_secrets(&[]).await.unwrap(),
|
||||
vec!["openai-prod".to_string(), "hf-prod".to_string()]
|
||||
);
|
||||
}
|
||||
|
||||
/// A server that keeps handing back the same token would otherwise spin
|
||||
/// forever.
|
||||
#[tokio::test]
|
||||
async fn test_list_secrets_rejects_a_repeated_page_token() {
|
||||
let conn = Connection::new_with_handler(|_| {
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(r#"{"secrets":[{"name":"openai-prod"}],"page_token":"same"}"#)
|
||||
.unwrap()
|
||||
});
|
||||
let error = conn.list_secrets(&[]).await.unwrap_err();
|
||||
assert!(
|
||||
error.to_string().contains("repeated a page_token"),
|
||||
"{error}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_drop_secret_posts_the_name_alone() {
|
||||
let conn = Connection::new_with_handler(|request| {
|
||||
assert_eq!(request.url().path(), "/v1/secrets/drop");
|
||||
let body: serde_json::Value =
|
||||
serde_json::from_slice(request.body().unwrap().as_bytes().unwrap()).unwrap();
|
||||
assert_eq!(body, serde_json::json!({"name": "openai-prod"}));
|
||||
http::Response::builder().status(200).body("{}").unwrap()
|
||||
});
|
||||
conn.drop_secret("openai-prod", &[]).await.unwrap();
|
||||
}
|
||||
|
||||
/// A namespace path is sent when there is one and omitted when there is
|
||||
/// not, so a root request stays byte identical to one from a client that
|
||||
/// predates namespace addressing -- which is what lets the parameter ship
|
||||
/// before every server implements it.
|
||||
#[tokio::test]
|
||||
async fn test_a_namespace_path_is_sent_only_when_it_is_not_root() {
|
||||
let conn = Connection::new_with_handler(|request| {
|
||||
let body: serde_json::Value =
|
||||
serde_json::from_slice(request.body().unwrap().as_bytes().unwrap()).unwrap();
|
||||
assert_eq!(
|
||||
body,
|
||||
serde_json::json!({
|
||||
"name": "openai-prod",
|
||||
"namespace_path": ["prod", "vision"],
|
||||
})
|
||||
);
|
||||
http::Response::builder().status(200).body("{}").unwrap()
|
||||
});
|
||||
conn.drop_secret("openai-prod", &["prod".to_string(), "vision".to_string()])
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let conn = Connection::new_with_handler(|request| {
|
||||
let body: serde_json::Value =
|
||||
serde_json::from_slice(request.body().unwrap().as_bytes().unwrap()).unwrap();
|
||||
assert!(body.get("namespace_path").is_none(), "{body}");
|
||||
http::Response::builder().status(200).body("{}").unwrap()
|
||||
});
|
||||
conn.drop_secret("openai-prod", &[]).await.unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_create_function_async_sends_canonical_request_and_decodes_typed_job() {
|
||||
const REQUEST: &str = include_str!(
|
||||
|
||||
@@ -75,6 +75,7 @@ mod create_index;
|
||||
pub mod datafusion;
|
||||
pub(crate) mod dataset;
|
||||
pub mod delete;
|
||||
pub mod freshness;
|
||||
pub mod lsm_stats;
|
||||
pub mod merge;
|
||||
pub mod optimize;
|
||||
@@ -778,7 +779,8 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
|
||||
message: "Function columns are supported only on LanceDB Cloud and Enterprise".into(),
|
||||
})
|
||||
}
|
||||
/// Fill a computed column's unfilled rows.
|
||||
/// Fill a computed column's unfilled rows and recompute those whose
|
||||
/// inputs changed.
|
||||
///
|
||||
/// The default returns `NotSupported`; Lance-backed tables override it.
|
||||
async fn refresh_column(&self, _column: &str) -> Result<RefreshColumnResult> {
|
||||
@@ -786,8 +788,8 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
|
||||
message: "computed columns are supported only on local tables".into(),
|
||||
})
|
||||
}
|
||||
/// Fill a computed column's unfilled rows, returning a [`Job`] tracking
|
||||
/// the operation.
|
||||
/// Fill a computed column's unfilled rows and recompute those whose
|
||||
/// inputs changed, returning a [`Job`] tracking the operation.
|
||||
async fn refresh_column_async(
|
||||
&self,
|
||||
_column: &str,
|
||||
@@ -1749,9 +1751,10 @@ impl Table {
|
||||
/// Declared with
|
||||
/// [`AddColumnsBuilder::computed`](add_columns::AddColumnsBuilder::computed),
|
||||
/// a column starts empty and gets its values here. Fragments appended
|
||||
/// since the last refresh are filled by the next one; fragments already
|
||||
/// filled are left as they are, so the call is idempotent and does not
|
||||
/// observe a mutated input.
|
||||
/// since the last refresh are filled by the next one, and fragments whose
|
||||
/// inputs changed since they were computed are recomputed (see
|
||||
/// [`freshness`](crate::table::freshness)); everything else is left as
|
||||
/// it is.
|
||||
///
|
||||
/// Local tables only: a remote refresh runs as a server job, through
|
||||
/// [`Table::refresh_column_async`].
|
||||
@@ -2804,7 +2807,7 @@ impl NativeTable {
|
||||
namespace_client: Option<Arc<dyn LanceNamespace>>,
|
||||
pushdown_operations: HashSet<NamespaceClientPushdownOperation>,
|
||||
) -> Result<Self> {
|
||||
computed_columns::ensure_no_foreign_declarations(batches.arrow_schema().fields())?;
|
||||
let batches = computed_columns::admit_create_source(batches)?;
|
||||
// Default params uses format v1.
|
||||
let params = params.unwrap_or(WriteParams {
|
||||
..Default::default()
|
||||
@@ -2904,6 +2907,7 @@ impl NativeTable {
|
||||
pushdown_operations: HashSet<NamespaceClientPushdownOperation>,
|
||||
session: Option<Arc<lance::session::Session>>,
|
||||
) -> Result<Self> {
|
||||
let batches = computed_columns::admit_create_source(batches)?;
|
||||
// Build table_id from namespace + name for the storage options provider
|
||||
let mut table_id = namespace.clone();
|
||||
table_id.push(name.to_string());
|
||||
@@ -5677,7 +5681,7 @@ mod tests {
|
||||
TableStatistics {
|
||||
num_rows: 250,
|
||||
num_indices: 0,
|
||||
total_bytes: 8925,
|
||||
total_bytes: 8969,
|
||||
fragment_stats: FragmentStatistics {
|
||||
num_fragments: 11,
|
||||
num_small_fragments: 11,
|
||||
|
||||
@@ -60,10 +60,11 @@ impl AddColumnsBuilder {
|
||||
/// every fragment that has none -- including fragments appended since the
|
||||
/// last refresh.
|
||||
///
|
||||
/// Refresh does not revisit a fragment it has filled, so mutating an input
|
||||
/// leaves the value computed at fill time; recomputing means dropping the
|
||||
/// column and declaring it again. An input cannot be renamed, retyped or
|
||||
/// dropped while a declaration reads it, since the expression names it.
|
||||
/// A refresh also recomputes the rows of a fragment whose inputs changed
|
||||
/// since it was computed (see [`freshness`](super::freshness)), so a
|
||||
/// mutated input is reflected by the next refresh. An input cannot be
|
||||
/// renamed, retyped or dropped while a declaration reads it, since the
|
||||
/// expression names it.
|
||||
///
|
||||
/// On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
/// server, and the refresh runs as a server job -- see
|
||||
|
||||
@@ -21,6 +21,7 @@
|
||||
//! [`computed_columns`] and [`computed_column_from_field`] read declarations
|
||||
//! back off a schema.
|
||||
|
||||
use futures::StreamExt;
|
||||
use std::collections::{BTreeSet, HashMap, HashSet};
|
||||
use std::sync::Arc;
|
||||
|
||||
@@ -70,6 +71,22 @@ pub const FUNCTION_BINDINGS_META_KEY: &str = "lancedb::function_bindings";
|
||||
/// Version of the schema-level Function binding envelope.
|
||||
pub const FUNCTION_BINDINGS_VERSION: u32 = 1;
|
||||
|
||||
/// Field metadata key holding `{fragment id -> input signature}` as JSON,
|
||||
/// recorded by the refresh that last computed each fragment. Outside the
|
||||
/// declaration namespace on purpose: a declaration is immutable through
|
||||
/// metadata edits, this is rewritten by every refresh. Seeded empty at
|
||||
/// declaration, so a column is tracked from birth; a column without it was
|
||||
/// declared before signatures existed.
|
||||
pub const SOURCE_SIGNATURE_META_KEY: &str = "computed_refresh.source_signature";
|
||||
|
||||
/// Field metadata key holding the definition digest a column was last
|
||||
/// computed under. A change to it makes every row stale.
|
||||
pub const DEFINITION_VERSION_META_KEY: &str = "computed_refresh.definition_version";
|
||||
|
||||
/// Field metadata key holding the table version the signature map describes:
|
||||
/// where a refresh starts following compactions to carry freshness forward.
|
||||
pub const RECORDED_AT_VERSION_META_KEY: &str = "computed_refresh.recorded_at_version";
|
||||
|
||||
/// Value of [`KIND_META_KEY`] for a column defined by a SQL expression.
|
||||
pub const SQL_KIND: &str = "sql";
|
||||
|
||||
@@ -138,6 +155,7 @@ fn computed_column_metadata(expression: &str, inputs: &[String]) -> HashMap<Stri
|
||||
INPUTS_META_KEY.to_string(),
|
||||
serde_json::to_string(inputs).unwrap_or_else(|_| "[]".to_string()),
|
||||
),
|
||||
(SOURCE_SIGNATURE_META_KEY.to_string(), "{}".to_string()),
|
||||
])
|
||||
}
|
||||
|
||||
@@ -162,6 +180,7 @@ pub fn function_computed_column_metadata(
|
||||
INPUTS_META_KEY.to_string(),
|
||||
serde_json::to_string(inputs).unwrap_or_else(|_| "[]".to_string()),
|
||||
),
|
||||
(SOURCE_SIGNATURE_META_KEY.to_string(), "{}".to_string()),
|
||||
])
|
||||
}
|
||||
|
||||
@@ -1294,7 +1313,8 @@ pub(crate) fn ensure_not_an_input(schema: &SchemaRef, paths: &[&str]) -> Result<
|
||||
}
|
||||
|
||||
/// Reject a write that supplies values for a computed column directly:
|
||||
/// only refresh materializes one, and refresh never revisits a filled row.
|
||||
/// only refresh materializes one, and only refresh decides what it
|
||||
/// recomputes.
|
||||
pub(crate) fn ensure_not_written<'a>(
|
||||
schema: &ArrowSchema,
|
||||
written: impl IntoIterator<Item = &'a str>,
|
||||
@@ -1338,6 +1358,106 @@ pub(crate) fn ensure_batch_writes_no_computed_values(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Validate every computed-column declaration `schema` carries against the
|
||||
/// schema itself: every field with declaration metadata is a complete
|
||||
/// declaration, a SQL declaration re-plans to the field it declares, a
|
||||
/// Function declaration satisfies the binding contract, and no declaration
|
||||
/// reads another computed column. What passes here is what `refresh_column`
|
||||
/// can execute.
|
||||
pub(crate) fn ensure_declarations_are_planned(schema: &ArrowSchema) -> Result<()> {
|
||||
let invalid = |message: String| Error::InvalidInput { message };
|
||||
// A field with any declaration key is a declaration; a partial one is
|
||||
// not "no declaration", it is a broken one.
|
||||
for field in schema.fields() {
|
||||
if field.metadata().keys().any(|k| is_declaration_key(k))
|
||||
&& computed_column_from_field(field).is_none()
|
||||
{
|
||||
return Err(invalid(format!(
|
||||
"field '{}' carries an incomplete computed-column declaration",
|
||||
field.name()
|
||||
)));
|
||||
}
|
||||
}
|
||||
let declared: HashSet<String> = computed_columns(schema)
|
||||
.into_iter()
|
||||
.map(|c| c.name)
|
||||
.collect();
|
||||
for column in computed_columns(schema) {
|
||||
let field = schema.field_with_name(&column.name)?;
|
||||
if !field.is_nullable() {
|
||||
return Err(invalid(format!(
|
||||
"computed column '{}' must be nullable until a refresh fills it",
|
||||
column.name
|
||||
)));
|
||||
}
|
||||
match &column.kind {
|
||||
ComputedColumnKind::Sql { expression } => {
|
||||
let others: Vec<ArrowField> = schema
|
||||
.fields()
|
||||
.iter()
|
||||
.filter(|f| f.name() != &column.name)
|
||||
.map(|f| f.as_ref().clone())
|
||||
.collect();
|
||||
let bound = bind(Arc::new(ArrowSchema::new(others)), &column.name, expression)?;
|
||||
if let Some(input) = bound.roots.iter().find(|r| declared.contains(*r)) {
|
||||
return Err(invalid(format!(
|
||||
"computed column '{}' reads computed column '{input}'",
|
||||
column.name
|
||||
)));
|
||||
}
|
||||
if &bound.data_type != field.data_type() {
|
||||
return Err(invalid(format!(
|
||||
"computed column '{}' is declared as {} but its expression yields {}",
|
||||
column.name,
|
||||
field.data_type(),
|
||||
bound.data_type
|
||||
)));
|
||||
}
|
||||
let mut declared_inputs = column.inputs.clone();
|
||||
declared_inputs.sort();
|
||||
if declared_inputs != bound.inputs {
|
||||
return Err(invalid(format!(
|
||||
"computed column '{}' declares inputs {:?} but its expression reads {:?}",
|
||||
column.name, declared_inputs, bound.inputs
|
||||
)));
|
||||
}
|
||||
}
|
||||
ComputedColumnKind::Function { binding_id, .. } => {
|
||||
// The binding validator resolves each input's leaf; the
|
||||
// no-computed-input rule is about the root it hangs from.
|
||||
let bindings = function_bindings(schema)?;
|
||||
let Some(binding) = bindings.iter().find(|b| b.binding_id() == binding_id) else {
|
||||
continue; // reported by the binding validator below
|
||||
};
|
||||
// Roots come from the canonical path parser: a quoted
|
||||
// top-level name may itself contain a dot.
|
||||
if let Some(input) = binding
|
||||
.inputs()
|
||||
.iter()
|
||||
.filter_map(|input| resolve_field_path(schema, &input.field_path).ok())
|
||||
.map(|resolved| resolved.root.name().as_str())
|
||||
.find(|r| declared.contains(*r))
|
||||
{
|
||||
return Err(invalid(format!(
|
||||
"computed column '{}' reads computed column '{input}'",
|
||||
column.name
|
||||
)));
|
||||
}
|
||||
}
|
||||
ComputedColumnKind::Unrecognized { kind } => {
|
||||
return Err(Error::NotSupported {
|
||||
message: format!(
|
||||
"computed column '{}' is defined by '{kind}', which this version \
|
||||
of lancedb cannot fill",
|
||||
column.name
|
||||
),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
ensure_supported_function_metadata(schema)
|
||||
}
|
||||
|
||||
/// Reject fields carrying declaration metadata that did not come through
|
||||
/// [`plan`]. One authority for creation, overwrite and raw transforms.
|
||||
pub(crate) fn ensure_no_foreign_declarations<'a>(
|
||||
@@ -1369,7 +1489,9 @@ fn ensure_no_foreign_declaration(field: &ArrowField) -> Result<()> {
|
||||
/// kind, the expression, the inputs -- would bypass that validation or move
|
||||
/// a binding out from under a refresh. Drop the column and declare it again.
|
||||
pub(crate) fn is_declaration_key(key: &str) -> bool {
|
||||
key == COMPUTED_COLUMN_META_KEY || key.starts_with("computed_column.")
|
||||
key == COMPUTED_COLUMN_META_KEY
|
||||
|| key.starts_with("computed_column.")
|
||||
|| key.starts_with("computed_refresh.")
|
||||
}
|
||||
|
||||
/// Reject retyping a computed column itself.
|
||||
@@ -1796,6 +1918,54 @@ pub(super) async fn add_foreign_kind(table: &crate::Table, name: &str, kind: &st
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// Admit a table's initial data: every declaration it carries is validated,
|
||||
/// and the stream refuses any batch with values in a computed column, whose
|
||||
/// values come from refresh alone. One boundary for every way a table is
|
||||
/// created.
|
||||
pub(crate) fn admit_create_source<S: lance_datafusion::utils::StreamingWriteSource>(
|
||||
batches: S,
|
||||
) -> Result<UnfilledDeclarations<S>> {
|
||||
let schema = batches.arrow_schema();
|
||||
ensure_declarations_are_planned(&schema)?;
|
||||
let declared = computed_columns(&schema)
|
||||
.into_iter()
|
||||
.map(|c| c.name)
|
||||
.collect();
|
||||
Ok(UnfilledDeclarations {
|
||||
inner: batches,
|
||||
declared,
|
||||
})
|
||||
}
|
||||
|
||||
/// A write source whose computed columns must arrive unfilled.
|
||||
pub(crate) struct UnfilledDeclarations<S> {
|
||||
inner: S,
|
||||
declared: Vec<String>,
|
||||
}
|
||||
|
||||
impl<S: lance_datafusion::utils::StreamingWriteSource> lance_datafusion::utils::StreamingWriteSource
|
||||
for UnfilledDeclarations<S>
|
||||
{
|
||||
fn arrow_schema(&self) -> SchemaRef {
|
||||
self.inner.arrow_schema()
|
||||
}
|
||||
|
||||
fn into_stream(self) -> datafusion_physical_plan::SendableRecordBatchStream {
|
||||
if self.declared.is_empty() {
|
||||
return self.inner.into_stream();
|
||||
}
|
||||
let schema = self.inner.arrow_schema();
|
||||
let declared = self.declared;
|
||||
let stream = self.inner.into_stream().map(move |batch| {
|
||||
let batch = batch?;
|
||||
ensure_batch_writes_no_computed_values(&declared, &batch)
|
||||
.map_err(|e| datafusion_common::DataFusionError::External(Box::new(e)))?;
|
||||
Ok(batch)
|
||||
});
|
||||
Box::pin(datafusion_physical_plan::stream::RecordBatchStreamAdapter::new(schema, stream))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
/// The gate's reproducer: the validator applies the same schema-level
|
||||
@@ -2646,6 +2816,8 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// A create carries a declaration only if it re-plans completely; this
|
||||
/// one lacks its inputs and is refused before its forged value matters.
|
||||
#[tokio::test]
|
||||
async fn test_create_table_cannot_inject_a_declaration() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
@@ -2673,7 +2845,7 @@ mod tests {
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(
|
||||
matches!(&err, Error::InvalidInput { message } if message.contains("computed()")),
|
||||
matches!(&err, Error::InvalidInput { message } if message.contains("computed column 'doubled'")),
|
||||
"{err:?}"
|
||||
);
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -3,9 +3,9 @@
|
||||
|
||||
//! Filling computed columns.
|
||||
//!
|
||||
//! A row without a value gets one; a row that has one keeps it. Refresh is
|
||||
//! therefore idempotent and does not observe input mutation -- once a row is
|
||||
//! filled, changing what the expression reads leaves the stored result alone.
|
||||
//! A row without a value gets one; a row that has one keeps it unless its
|
||||
//! fragment's inputs moved since it was computed, which `freshness` decides
|
||||
//! from the manifest and stamps after every fill.
|
||||
//!
|
||||
//! A column's computed inputs are filled first -- the dependency graph is
|
||||
//! walked once, each reachable column filled once in dependency order, each
|
||||
@@ -31,6 +31,7 @@
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
|
||||
use arrow_array::{
|
||||
Array, ArrayRef, BooleanArray, LargeBinaryArray, RecordBatch, RecordBatchOptions, StructArray,
|
||||
@@ -48,6 +49,7 @@ use lance_core::datatypes::{BlobHandling, Schema as LanceSchema};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use super::computed_columns::{BoundExpression, ComputedColumnKind, computed_column_from_field};
|
||||
use super::freshness::{self, SignatureMap, StalenessPlan};
|
||||
use super::{BaseTable, NativeTable};
|
||||
use crate::job::Job;
|
||||
use crate::{Error, Result};
|
||||
@@ -110,28 +112,81 @@ async fn execute_refresh_column_with_source(
|
||||
};
|
||||
let output_is_blob = field.is_blob_v2();
|
||||
|
||||
// Which fragments the null filter cannot speak for: their inputs moved
|
||||
// since they were computed, or the definition did. Decided once, from
|
||||
// the manifest the values are read from.
|
||||
let inputs = freshness::fields_for_paths(dataset.schema(), &bound.inputs)?;
|
||||
let definition = freshness::definition_version(&expression);
|
||||
let staleness = freshness::staleness_against(&dataset, column, &definition, &inputs).await?;
|
||||
|
||||
let mut rows_filled = 0u64;
|
||||
let mut replacements = Vec::new();
|
||||
// Fragments this refresh computed in full, signed at the version read.
|
||||
let mut computed = SignatureMap::new();
|
||||
for fragment in dataset.get_fragments() {
|
||||
let gained = count_fragment_gains(&dataset, &fragment, &bound, column).await?;
|
||||
if gained == 0 {
|
||||
continue;
|
||||
let fragment_id = u32::try_from(fragment.id()).map_err(|_| Error::Runtime {
|
||||
message: format!("fragment id {} does not fit a signature map", fragment.id()),
|
||||
})?;
|
||||
// A recompute rewrites every live row, so it is staged without the
|
||||
// probe and counted as it fills; a null fill probes first, since a
|
||||
// fragment with nothing to gain is not worth a write.
|
||||
let recompute = staleness.is_dirty(fragment_id);
|
||||
let whole = recompute || {
|
||||
let (gained, unfilled) =
|
||||
count_fragment_gains(&dataset, &fragment, &bound, column).await?;
|
||||
if gained == 0 {
|
||||
continue;
|
||||
}
|
||||
rows_filled += gained;
|
||||
unfilled == u64::try_from(fragment.count_rows(None).await?).unwrap_or(u64::MAX)
|
||||
};
|
||||
if whole {
|
||||
computed.insert(
|
||||
fragment_id,
|
||||
freshness::fragment_input_signature(fragment.metadata(), &inputs)?,
|
||||
);
|
||||
}
|
||||
rows_filled += gained;
|
||||
let values =
|
||||
fill_stream(&dataset, &fragment, bound.clone(), column, output_is_blob).await?;
|
||||
let gained = Arc::new(AtomicU64::new(0));
|
||||
let values = fill_stream(
|
||||
&dataset,
|
||||
&fragment,
|
||||
bound.clone(),
|
||||
column,
|
||||
output_is_blob,
|
||||
recompute,
|
||||
gained.clone(),
|
||||
)
|
||||
.await?;
|
||||
replacements.push(fragment.write_columns(values, &column_schema).await?);
|
||||
if recompute {
|
||||
rows_filled += gained.load(Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
|
||||
let source_version = dataset.version().version;
|
||||
if replacements.is_empty() {
|
||||
// Nothing to fill; the stamp may still have something to record -- a
|
||||
// column not yet enrolled, or fragments a compaction carried.
|
||||
let mut latest = (*dataset).clone();
|
||||
let stamped = record(
|
||||
&mut latest,
|
||||
(&dataset, &staleness),
|
||||
column,
|
||||
&definition,
|
||||
&inputs,
|
||||
computed,
|
||||
)
|
||||
.await;
|
||||
if stamped.is_some() {
|
||||
table.dataset.update(latest);
|
||||
}
|
||||
return Ok(RefreshExecution {
|
||||
result: RefreshColumnResult {
|
||||
rows_filled: 0,
|
||||
version: source_version,
|
||||
version: stamped.unwrap_or(source_version),
|
||||
},
|
||||
source_version,
|
||||
published_version: None,
|
||||
published_version: stamped,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -149,7 +204,17 @@ async fn execute_refresh_column_with_source(
|
||||
)
|
||||
.await?;
|
||||
|
||||
let version = new_dataset.version().version;
|
||||
let mut new_dataset = new_dataset;
|
||||
let version = record(
|
||||
&mut new_dataset,
|
||||
(&dataset, &staleness),
|
||||
column,
|
||||
&definition,
|
||||
&inputs,
|
||||
computed,
|
||||
)
|
||||
.await
|
||||
.unwrap_or(new_dataset.version().version);
|
||||
table.dataset.update(new_dataset);
|
||||
Ok(RefreshExecution {
|
||||
result: RefreshColumnResult {
|
||||
@@ -161,6 +226,31 @@ async fn execute_refresh_column_with_source(
|
||||
})
|
||||
}
|
||||
|
||||
/// Stamp the input state the refresh computed from (see
|
||||
/// [`freshness::record_freshness`]); the version the stamp landed at, which
|
||||
/// is the last one the refresh wrote. Never fails the refresh: the values
|
||||
/// are committed, and a missing stamp only costs a recompute next time.
|
||||
async fn record(
|
||||
latest: &mut Dataset,
|
||||
pinned: (&Dataset, &StalenessPlan),
|
||||
column: &str,
|
||||
definition: &str,
|
||||
inputs: &freshness::InputFields,
|
||||
computed: SignatureMap,
|
||||
) -> Option<u64> {
|
||||
match freshness::record_freshness(latest, Some(pinned), column, definition, inputs, computed)
|
||||
.await
|
||||
{
|
||||
Ok(record) => record.version,
|
||||
Err(error) => {
|
||||
log::warn!(
|
||||
"could not record the input state computed column '{column}' was refreshed from ({error}); its fragments will recompute on the next refresh"
|
||||
);
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Refuse while a computed input still has rows a refresh of it would fill:
|
||||
/// read now, its placeholder null would be evaluated as a value and kept.
|
||||
async fn ensure_inputs_filled(
|
||||
@@ -188,7 +278,9 @@ async fn ensure_inputs_filled(
|
||||
let input_bound = super::computed_columns::bind(schema.clone(), input, expression)?;
|
||||
let mut unfilled = 0u64;
|
||||
for fragment in dataset.get_fragments() {
|
||||
unfilled += count_fragment_gains(dataset, &fragment, &input_bound, input).await?;
|
||||
unfilled += count_fragment_gains(dataset, &fragment, &input_bound, input)
|
||||
.await?
|
||||
.0;
|
||||
}
|
||||
if unfilled > 0 {
|
||||
return Err(Error::InvalidInput {
|
||||
@@ -436,32 +528,35 @@ fn blob_array_from_binary(
|
||||
/// Scans only the unfilled live rows -- deleted rows never reach the
|
||||
/// expression here, the filter having already excluded them -- and counts the
|
||||
/// non-null results. Exact, so it is both the staging decision and the
|
||||
/// fragment's contribution to `rows_filled`.
|
||||
/// fragment's contribution to `rows_filled`. Returns the gains and the rows
|
||||
/// scanned.
|
||||
async fn count_fragment_gains(
|
||||
dataset: &Dataset,
|
||||
fragment: &FileFragment,
|
||||
bound: &BoundExpression,
|
||||
column: &str,
|
||||
) -> Result<u64> {
|
||||
) -> Result<(u64, u64)> {
|
||||
let mut scanner = dataset.scan();
|
||||
scanner
|
||||
.with_fragments(vec![fragment.metadata().clone()])
|
||||
.with_row_id()
|
||||
.filter(&format!("{} IS NULL", quote_identifier(column)))?
|
||||
.project(&bound.roots)?;
|
||||
.project(&bound.roots)?
|
||||
.filter(&format!("{} IS NULL", quote_identifier(column)))?;
|
||||
configure_blob_inputs(&mut scanner, dataset.schema(), bound, None)?;
|
||||
|
||||
let mut gained = 0u64;
|
||||
let mut considered = 0u64;
|
||||
let mut batches = scanner.try_into_stream().await?;
|
||||
while let Some(batch) = batches.try_next().await? {
|
||||
let evaluated = evaluate(bound, &evaluation_batch(&batch, bound, None)?)?;
|
||||
gained += (batch.num_rows() - evaluated.null_count()) as u64;
|
||||
considered += batch.num_rows() as u64;
|
||||
}
|
||||
Ok(gained)
|
||||
Ok((gained, considered))
|
||||
}
|
||||
|
||||
/// Stream one fragment's column in physical order, filling the unfilled live
|
||||
/// rows and keeping every other value.
|
||||
/// rows -- every live row, for a recompute -- and keeping every other value.
|
||||
///
|
||||
/// Deleted rows are carried through so the values line up positionally with
|
||||
/// the fragment's data files; they are never read back, but the column file
|
||||
@@ -472,6 +567,8 @@ async fn fill_stream(
|
||||
bound: Arc<BoundExpression>,
|
||||
column: &str,
|
||||
output_is_blob: bool,
|
||||
recompute: bool,
|
||||
gained: Arc<AtomicU64>,
|
||||
) -> Result<impl Stream<Item = lance_core::Result<RecordBatch>> + Send + use<>> {
|
||||
let mut projection: Vec<String> = bound.roots.clone();
|
||||
projection.push(column.to_string());
|
||||
@@ -521,14 +618,20 @@ async fn fill_stream(
|
||||
.column_by_name(ROW_ID)
|
||||
.ok_or_else(|| missing(ROW_ID))?;
|
||||
|
||||
// Only an unfilled live row gains a value; a deleted row has a null
|
||||
// row id and keeps its (null) slot.
|
||||
let unfilled = arrow::compute::is_null(existing.as_ref())?;
|
||||
// Only an unfilled live row gains a value, or every live row under a
|
||||
// recompute; a deleted row has a null row id and keeps its (null) slot.
|
||||
let live = arrow::compute::is_not_null(row_ids.as_ref())?;
|
||||
let fill = arrow::compute::and(&unfilled, &live)?;
|
||||
let fill = if recompute {
|
||||
live
|
||||
} else {
|
||||
let unfilled = arrow::compute::is_null(existing.as_ref())?;
|
||||
arrow::compute::and(&unfilled, &live)?
|
||||
};
|
||||
let keep = arrow::compute::not(&fill)?;
|
||||
|
||||
let computed = evaluate(&bound, &evaluation_batch(&batch, &bound, Some(&keep))?)?;
|
||||
let values = arrow::compute::and(&fill, &arrow::compute::is_not_null(&computed)?)?;
|
||||
gained.fetch_add(values.true_count() as u64, Ordering::Relaxed);
|
||||
let merged = arrow_select::zip::zip(&fill, &computed, existing)?;
|
||||
let merged = if output_is_blob {
|
||||
blob_array_from_binary(&merged, projected.field(0))?
|
||||
@@ -737,7 +840,8 @@ mod tests {
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(no_op.rows_assigned, 0);
|
||||
assert_eq!(no_op.source_version, 3);
|
||||
// The fill, then the stamp recording what it computed from.
|
||||
assert_eq!(no_op.source_version, 4);
|
||||
assert_eq!(no_op.published_version, None);
|
||||
}
|
||||
|
||||
@@ -777,8 +881,8 @@ mod tests {
|
||||
}
|
||||
|
||||
/// A row is filled only by gaining a value, so an expression yielding null
|
||||
/// settles at once instead of re-selecting the same rows forever. Nothing
|
||||
/// is staged, so the version does not move either.
|
||||
/// settles at once instead of re-selecting the same rows forever: the
|
||||
/// second refresh finds the fragment signed and moves nothing.
|
||||
#[tokio::test]
|
||||
async fn test_refresh_converges_on_a_null_result() {
|
||||
let table = table_with("refresh_null_result", vec![1, 2, 3]).await;
|
||||
@@ -792,28 +896,145 @@ mod tests {
|
||||
|
||||
let first = table.refresh_column("maybe").await.unwrap();
|
||||
assert_eq!(first.rows_filled, 0);
|
||||
assert_eq!(first.version, declared);
|
||||
assert!(first.version > declared);
|
||||
assert_eq!(read(&table, "maybe").await, vec![None, None, None]);
|
||||
|
||||
let again = table.refresh_column("maybe").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 0);
|
||||
assert_eq!(again.version, declared);
|
||||
assert_eq!(again.version, first.version);
|
||||
}
|
||||
|
||||
/// The contract's boundary: a filled fragment is not revisited, so
|
||||
/// mutating an input leaves the value computed at fill time.
|
||||
/// A filled row whose input moved is recomputed: the update rewrites
|
||||
/// the row into a fragment the stamp never signed, and only that one.
|
||||
#[tokio::test]
|
||||
async fn test_refresh_does_not_observe_input_mutation() {
|
||||
let table = table_with("refresh_mutation", vec![1]).await;
|
||||
async fn test_refresh_recomputes_a_row_whose_input_moved() {
|
||||
let table = table_with("refresh_mutation", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
append(&table, vec![5]).await;
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2)]);
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(10)]
|
||||
);
|
||||
|
||||
table.update().column("x", "3").execute().await.unwrap();
|
||||
table
|
||||
.update()
|
||||
.column("x", "7")
|
||||
.only_if("x = 5")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let again = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 1);
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(14)]
|
||||
);
|
||||
|
||||
let settled = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(settled.rows_filled, 0);
|
||||
assert_eq!(settled.version, again.version);
|
||||
}
|
||||
|
||||
/// A deleted row is never computed and the rows that stay keep their
|
||||
/// values: a delete recomputes nothing and stamps nothing.
|
||||
#[tokio::test]
|
||||
async fn test_a_delete_recomputes_nothing() {
|
||||
let table = table_with("refresh_delete", vec![1, 2, 3]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
let filled = table.refresh_column("doubled").await.unwrap();
|
||||
|
||||
table.delete("x = 2").await.unwrap();
|
||||
let deleted = table.version().await.unwrap();
|
||||
|
||||
let again = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 0);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2)]);
|
||||
assert_eq!(again.version, deleted);
|
||||
assert!(deleted > filled.version);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2), Some(6)]);
|
||||
}
|
||||
|
||||
/// Compaction copies inputs unchanged, so a fragment it builds from
|
||||
/// signed ones is fresh: the refresh recomputes nothing and only records
|
||||
/// the new fragment.
|
||||
#[tokio::test]
|
||||
async fn test_a_compaction_of_signed_fragments_recomputes_nothing() {
|
||||
let table = table_with("refresh_compact_signed", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
append(&table, vec![5]).await;
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
|
||||
table
|
||||
.optimize(crate::table::OptimizeAction::Compact {
|
||||
options: crate::table::CompactionOptions::default(),
|
||||
remap_options: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
let compacted = table.version().await.unwrap();
|
||||
|
||||
let carried = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(carried.rows_filled, 0);
|
||||
assert_eq!(carried.version, compacted + 1);
|
||||
let settled = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(settled.version, carried.version);
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(10)]
|
||||
);
|
||||
}
|
||||
|
||||
/// A column declared before signatures existed has no map. Its first
|
||||
/// refresh keeps the null-fill contract and enrolls what it read from;
|
||||
/// from then on a moved input is recomputed like any other.
|
||||
#[tokio::test]
|
||||
async fn test_an_unsigned_column_is_enrolled_by_its_first_refresh() {
|
||||
let table = table_with("refresh_legacy", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
table
|
||||
.update()
|
||||
.column("x", "3")
|
||||
.only_if("x = 1")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let native = table.as_native().unwrap();
|
||||
let mut dataset = (*native.dataset.get().await.unwrap()).clone();
|
||||
let declaration = dataset
|
||||
.schema()
|
||||
.field("doubled")
|
||||
.unwrap()
|
||||
.metadata
|
||||
.iter()
|
||||
.filter(|(key, _)| !key.starts_with("computed_refresh."))
|
||||
.map(|(key, value)| (key.clone(), value.clone()))
|
||||
.collect::<Vec<_>>();
|
||||
dataset
|
||||
.update_field_metadata()
|
||||
.replace("doubled", declaration)
|
||||
.unwrap()
|
||||
.await
|
||||
.unwrap();
|
||||
table.checkout_latest().await.unwrap();
|
||||
|
||||
// Null-fill only: the moved row keeps the value it was filled with.
|
||||
let enrolled = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(enrolled.rows_filled, 0);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2), Some(4)]);
|
||||
|
||||
table
|
||||
.update()
|
||||
.column("x", "5")
|
||||
.only_if("x = 3")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let again = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 1);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(4), Some(10)]);
|
||||
}
|
||||
|
||||
/// A row rewrite before the first refresh materializes the declared
|
||||
@@ -831,11 +1052,11 @@ mod tests {
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(6)]);
|
||||
}
|
||||
|
||||
/// The contract holds row by row, not fragment by fragment: revisiting a
|
||||
/// fragment to fill one row must not recompute a filled row sitting beside
|
||||
/// it, even where the input behind it has since changed.
|
||||
/// A fragment compacted out of one the stamp never signed cannot vouch
|
||||
/// for any of its rows: every live row is recomputed, the moved one
|
||||
/// included.
|
||||
#[tokio::test]
|
||||
async fn test_refresh_does_not_recompute_a_filled_row_beside_an_unfilled_one() {
|
||||
async fn test_a_compaction_of_an_unsigned_fragment_recomputes_it() {
|
||||
let table = table_with("refresh_mixed", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
@@ -857,18 +1078,70 @@ mod tests {
|
||||
.unwrap();
|
||||
|
||||
let result = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(result.rows_filled, 1);
|
||||
// 2 is the mutated row keeping the value it was filled with, not 200.
|
||||
assert_eq!(result.rows_filled, 3);
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(4), Some(10), Some(200)]
|
||||
);
|
||||
}
|
||||
|
||||
/// The gate's reproducer: a raw lance append may carry a value for the
|
||||
/// computed column. Compaction cannot certify it, so the product is
|
||||
/// recomputed and the supplied value replaced.
|
||||
#[tokio::test]
|
||||
async fn test_raw_append_values_are_not_trusted_after_compaction() {
|
||||
use arrow_array::RecordBatchIterator;
|
||||
use lance::Dataset;
|
||||
use lance::dataset::{WriteMode, WriteParams};
|
||||
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let conn = connect(dir.path().to_str().unwrap())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let batch = record_batch!(("x", Int32, [1, 2])).unwrap();
|
||||
let table = conn
|
||||
.create_table("raw_append", batch)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
|
||||
let batch = record_batch!(("x", Int32, [5]), ("doubled", Int32, [Some(999_i32)])).unwrap();
|
||||
let schema = batch.schema();
|
||||
let uri = table.uri().await.unwrap();
|
||||
Dataset::write(
|
||||
RecordBatchIterator::new(vec![Ok(batch)], schema),
|
||||
&uri,
|
||||
Some(WriteParams {
|
||||
mode: WriteMode::Append,
|
||||
..Default::default()
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
table.checkout_latest().await.unwrap();
|
||||
table
|
||||
.optimize(crate::table::OptimizeAction::Compact {
|
||||
options: crate::table::CompactionOptions::default(),
|
||||
remap_options: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let result = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(result.rows_filled, 3);
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(10)]
|
||||
);
|
||||
}
|
||||
|
||||
/// Filling a fragment must not disturb the values it already holds, which
|
||||
/// is what makes a compaction-mixed fragment safe to revisit.
|
||||
/// An appended fragment holds no values, so compacting it into a signed
|
||||
/// one leaves the product fresh: only the appended rows are filled.
|
||||
#[tokio::test]
|
||||
async fn test_refresh_preserves_already_filled_rows() {
|
||||
async fn test_a_compaction_with_an_appended_fragment_fills_only_its_rows() {
|
||||
let table = table_with("refresh_preserves", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
@@ -995,7 +1268,8 @@ mod tests {
|
||||
assert_eq!(result.rows_failed, 0);
|
||||
assert_eq!(result.rows_remaining, 0);
|
||||
assert_eq!(result.source_version, 2);
|
||||
assert_eq!(result.published_version, Some(3));
|
||||
// The fill lands at 3; the stamp recording its inputs is published at 4.
|
||||
assert_eq!(result.published_version, Some(4));
|
||||
assert_eq!(job.status().await.unwrap(), "finished");
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
@@ -1059,31 +1333,32 @@ mod tests {
|
||||
assert_eq!(read(&table, "quotient").await, vec![Some(10)]);
|
||||
}
|
||||
|
||||
/// The gate's reproducer: an already-filled row's value must not be
|
||||
/// re-evaluated either -- its input may have mutated into one the
|
||||
/// expression chokes on.
|
||||
/// A filled row whose input moved is re-evaluated, and a row whose
|
||||
/// input did not move is not: the untouched fragment is never read, so
|
||||
/// its poison input is never reached.
|
||||
#[tokio::test]
|
||||
async fn test_a_filled_rows_value_is_never_evaluated() {
|
||||
let table = table_with("refresh_filled_poison", vec![1, 2]).await;
|
||||
async fn test_only_a_moved_rows_value_is_re_evaluated() {
|
||||
let table = table_with("refresh_filled_poison", vec![1, 0]).await;
|
||||
table
|
||||
.add_columns()
|
||||
.computed("quotient", "10 / x")
|
||||
.computed("quotient", "10 / coalesce(nullif(x, 0), 1)")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
table.refresh_column("quotient").await.unwrap();
|
||||
assert_eq!(read(&table, "quotient").await, vec![Some(10), Some(10)]);
|
||||
|
||||
append(&table, vec![5]).await;
|
||||
table
|
||||
.update()
|
||||
.column("x", "0")
|
||||
.column("x", "2")
|
||||
.only_if("x = 1")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
append(&table, vec![5]).await;
|
||||
|
||||
let result = table.refresh_column("quotient").await.unwrap();
|
||||
assert_eq!(result.rows_filled, 1);
|
||||
assert_eq!(result.rows_filled, 2);
|
||||
assert_eq!(
|
||||
read(&table, "quotient").await,
|
||||
vec![Some(2), Some(5), Some(10)]
|
||||
|
||||
@@ -19,6 +19,7 @@ use lancedb::{
|
||||
connect, connect_namespace,
|
||||
database::listing::{
|
||||
ListingDatabaseOptions, NewTableConfig, OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS,
|
||||
OPT_NEW_TABLE_STORAGE_VERSION,
|
||||
},
|
||||
query::{ExecutableQuery, QueryBase},
|
||||
table::{AddDataMode, CompactionOptions, OptimizeAction, OptimizeStats, WriteOptions},
|
||||
@@ -146,7 +147,10 @@ async fn non_blob_table_keeps_default_format_and_row_id_setting() -> Result<()>
|
||||
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int64, false)]));
|
||||
let table = db.create_empty_table("t", schema).execute().await?;
|
||||
|
||||
assert!(!supports_blob_v2(storage_format_version(&table).await));
|
||||
assert_eq!(
|
||||
storage_format_version(&table).await,
|
||||
LanceFileVersion::Stable.resolve()
|
||||
);
|
||||
assert!(!uses_stable_row_ids(&table).await);
|
||||
Ok(())
|
||||
}
|
||||
@@ -809,7 +813,11 @@ async fn fetch_blobs_rejects_unknown_column() -> Result<()> {
|
||||
#[tokio::test]
|
||||
async fn fetch_blobs_rejects_legacy_v1_blob_column() -> Result<()> {
|
||||
let tmp = tempdir().unwrap();
|
||||
let db = connect(tmp.path().to_str().unwrap()).execute().await?;
|
||||
// Legacy v1 blob columns are only writable at file version <= 2.1.
|
||||
let db = connect(tmp.path().to_str().unwrap())
|
||||
.storage_options([(OPT_NEW_TABLE_STORAGE_VERSION, "2.1")])
|
||||
.execute()
|
||||
.await?;
|
||||
let legacy = Field::new("image", DataType::LargeBinary, true).with_metadata(
|
||||
std::collections::HashMap::from([("lance-encoding:blob".to_string(), "true".to_string())]),
|
||||
);
|
||||
|
||||
@@ -1,12 +1,11 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
use std::collections::BTreeMap;
|
||||
use std::fs;
|
||||
use std::path::PathBuf;
|
||||
|
||||
use lancedb::function::{
|
||||
FunctionApplication, FunctionBinding, FunctionVersion, RefreshColumnResult, SecretReference,
|
||||
FunctionApplication, FunctionBinding, FunctionVersion, RefreshColumnResult,
|
||||
};
|
||||
use serde_json::Value;
|
||||
|
||||
@@ -21,26 +20,6 @@ fn job_result(name: &str) -> Value {
|
||||
serde_json::from_str::<Value>(&fixture(name)).expect("remote Job fixture")["result"].clone()
|
||||
}
|
||||
|
||||
/// No client value models a resolved credential, at any nesting depth.
|
||||
fn assert_no_secret_values(value: &Value) {
|
||||
match value {
|
||||
Value::Object(values) => {
|
||||
for (key, value) in values {
|
||||
assert!(
|
||||
!matches!(
|
||||
key.as_str(),
|
||||
"secret_value" | "secret_values" | "resolved_secret" | "resolved_secrets"
|
||||
),
|
||||
"client canonical value must not model resolved secret material"
|
||||
);
|
||||
assert_no_secret_values(value);
|
||||
}
|
||||
}
|
||||
Value::Array(values) => values.iter().for_each(assert_no_secret_values),
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn function_version_job_result_matches_shared_canonical_golden() {
|
||||
let result = job_result("remote_function_job.json");
|
||||
@@ -49,10 +28,6 @@ fn function_version_job_result_matches_shared_canonical_golden() {
|
||||
assert_eq!(version.name(), "embed");
|
||||
assert_eq!(version.version(), "fv_01K3EXACT");
|
||||
assert_eq!(version.runtime_digest(), "sha256:runtime");
|
||||
assert_eq!(
|
||||
version.secret_env_bindings(),
|
||||
&BTreeMap::from([("HF_TOKEN".to_string(), SecretReference::new("hf-prod"))])
|
||||
);
|
||||
assert_eq!(
|
||||
version.to_canonical_json().expect("canonical JSON"),
|
||||
fixture("remote_function_version.canonical.json").trim()
|
||||
@@ -167,40 +142,3 @@ fn floating_point_application_literals_are_rejected_consistently() {
|
||||
.contains("floating-point Function literals")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn canonical_client_values_carry_bindings_and_no_credentials() {
|
||||
let result = job_result("remote_function_job.json");
|
||||
let version = FunctionVersion::from_json(&result.to_string()).expect("FunctionVersion result");
|
||||
let canonical: Value = serde_json::from_str(
|
||||
&version
|
||||
.to_canonical_json()
|
||||
.expect("canonical FunctionVersion"),
|
||||
)
|
||||
.expect("canonical JSON");
|
||||
|
||||
assert_eq!(
|
||||
canonical["secret_env_bindings"],
|
||||
serde_json::json!({"HF_TOKEN": {"name": "hf-prod"}})
|
||||
);
|
||||
assert_no_secret_values(&canonical);
|
||||
}
|
||||
|
||||
/// Every Function registered before Secrets existed serializes unchanged.
|
||||
#[test]
|
||||
fn a_version_without_bindings_keeps_the_original_wire_shape() {
|
||||
let mut result = job_result("remote_function_job.json");
|
||||
result
|
||||
.as_object_mut()
|
||||
.expect("Function version object")
|
||||
.remove("secret_env_bindings");
|
||||
let version = FunctionVersion::from_json(&result.to_string()).expect("FunctionVersion result");
|
||||
|
||||
assert!(version.secret_env_bindings().is_empty());
|
||||
assert!(
|
||||
!version
|
||||
.to_canonical_json()
|
||||
.expect("canonical FunctionVersion")
|
||||
.contains("secret_env_bindings")
|
||||
);
|
||||
}
|
||||
|
||||
@@ -5,11 +5,7 @@ use std::fs;
|
||||
use std::path::PathBuf;
|
||||
|
||||
use lancedb::Error;
|
||||
use lancedb::function::{
|
||||
FunctionRegistrationRequest, MAX_FUNCTION_SECRET_ENV_BINDINGS, MAX_SECRET_VALUE_BYTES,
|
||||
SecretReference,
|
||||
};
|
||||
use serde_json::Value;
|
||||
use lancedb::function::FunctionRegistrationRequest;
|
||||
|
||||
fn fixture(name: &str) -> String {
|
||||
let path = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
|
||||
@@ -18,26 +14,6 @@ fn fixture(name: &str) -> String {
|
||||
fs::read_to_string(path).expect("fixture must be readable")
|
||||
}
|
||||
|
||||
/// A registration request never models a resolved credential, at any depth.
|
||||
fn assert_no_secret_values(value: &Value) {
|
||||
match value {
|
||||
Value::Object(values) => {
|
||||
for (key, value) in values {
|
||||
assert!(
|
||||
!matches!(
|
||||
key.as_str(),
|
||||
"secret_value" | "secret_values" | "resolved_secret" | "resolved_secrets"
|
||||
),
|
||||
"registration requests must not model resolved secret material"
|
||||
);
|
||||
assert_no_secret_values(value);
|
||||
}
|
||||
}
|
||||
Value::Array(values) => values.iter().for_each(assert_no_secret_values),
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn registration_request_matches_shared_canonical_golden() {
|
||||
let request = FunctionRegistrationRequest::from_json(&fixture(
|
||||
@@ -46,45 +22,10 @@ fn registration_request_matches_shared_canonical_golden() {
|
||||
.expect("registration request");
|
||||
assert_eq!(request.name, "normalize_score");
|
||||
assert_eq!(request.artifact.adapter.kind, "scalar_to_arrow_batch");
|
||||
// The unchanged path: a Function that binds nothing serializes today's
|
||||
// bytes, with no `secret_env_bindings` key at all.
|
||||
assert!(request.secret_env_bindings.is_empty());
|
||||
assert_eq!(
|
||||
request.to_canonical_json().expect("canonical request"),
|
||||
fixture("remote_function_registration_request.canonical.json").trim()
|
||||
);
|
||||
|
||||
let value: Value =
|
||||
serde_json::from_str(&request.to_canonical_json().expect("canonical request"))
|
||||
.expect("request JSON");
|
||||
assert_no_secret_values(&value);
|
||||
}
|
||||
|
||||
/// The same shared golden as the Python suite builds from `@udf(secrets=...)`
|
||||
/// plus `bind_secrets`, so both clients agree byte for byte on a bound request.
|
||||
#[test]
|
||||
fn secret_bound_registration_request_matches_shared_canonical_golden() {
|
||||
let request = FunctionRegistrationRequest::from_json(&fixture(
|
||||
"remote_function_secret_registration_request.json",
|
||||
))
|
||||
.expect("registration request");
|
||||
assert_eq!(request.name, "analyze_caption");
|
||||
assert_eq!(
|
||||
request.secret_env_bindings,
|
||||
std::collections::BTreeMap::from([(
|
||||
"OPENAI_API_KEY".to_string(),
|
||||
SecretReference::new("openai-prod")
|
||||
)])
|
||||
);
|
||||
assert_eq!(
|
||||
request.to_canonical_json().expect("canonical request"),
|
||||
fixture("remote_function_secret_registration_request.canonical.json").trim()
|
||||
);
|
||||
|
||||
let value: Value =
|
||||
serde_json::from_str(&request.to_canonical_json().expect("canonical request"))
|
||||
.expect("request JSON");
|
||||
assert_no_secret_values(&value);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -116,111 +57,3 @@ async fn local_function_catalog_operations_return_stable_not_supported() {
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
/// The cap is enforced above the backend, so every database and every language
|
||||
/// surface rejects the same envelope. A local connection would otherwise answer
|
||||
/// `NotSupported` first, which is what makes it the honest probe here.
|
||||
#[tokio::test]
|
||||
async fn a_function_binds_at_most_sixteen_secrets() {
|
||||
let directory = tempfile::tempdir().unwrap();
|
||||
let connection = lancedb::connect(directory.path().to_str().unwrap())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let mut request = FunctionRegistrationRequest::from_json(&fixture(
|
||||
"remote_function_registration_request.json",
|
||||
))
|
||||
.unwrap();
|
||||
request.secret_env_bindings = (0..=MAX_FUNCTION_SECRET_ENV_BINDINGS)
|
||||
.map(|index| {
|
||||
(
|
||||
format!("TOKEN_{index}"),
|
||||
SecretReference::new(format!("secret-{index}")),
|
||||
)
|
||||
})
|
||||
.collect();
|
||||
|
||||
let error = connection.create_function_async(request).await.unwrap_err();
|
||||
assert!(matches!(
|
||||
error,
|
||||
Error::InvalidInput { message } if message.contains("at most 16 secrets")
|
||||
));
|
||||
}
|
||||
|
||||
/// The binding contract is enforced above the backend in full, not just its
|
||||
/// count: a caller that skips a language binding still cannot register a name
|
||||
/// the runtime could not deliver.
|
||||
#[tokio::test]
|
||||
async fn binding_names_are_validated_before_dispatch() {
|
||||
let directory = tempfile::tempdir().unwrap();
|
||||
let connection = lancedb::connect(directory.path().to_str().unwrap())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let mut invalid_name = FunctionRegistrationRequest::from_json(&fixture(
|
||||
"remote_function_registration_request.json",
|
||||
))
|
||||
.unwrap();
|
||||
invalid_name.secret_env_bindings =
|
||||
[("BAD=NAME".to_string(), SecretReference::new("openai-prod"))].into();
|
||||
let error = connection
|
||||
.create_function_async(invalid_name)
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(matches!(
|
||||
error,
|
||||
Error::InvalidInput { message } if message.contains("portable")
|
||||
));
|
||||
|
||||
// `env` is readable wherever the Function's record is; a bound Secret is
|
||||
// not. The same name cannot mean both.
|
||||
let mut overlapping = FunctionRegistrationRequest::from_json(&fixture(
|
||||
"remote_function_registration_request.json",
|
||||
))
|
||||
.unwrap();
|
||||
let bound = overlapping
|
||||
.runtime
|
||||
.env()
|
||||
.and_then(|env| env.keys().next().cloned())
|
||||
.expect("fixture runtime declares env");
|
||||
overlapping.secret_env_bindings = [(bound.clone(), SecretReference::new("openai-prod"))].into();
|
||||
let error = connection
|
||||
.create_function_async(overlapping)
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(matches!(
|
||||
error,
|
||||
Error::InvalidInput { message } if message.contains("already set by runtime.env")
|
||||
));
|
||||
}
|
||||
|
||||
/// An oversized credential is refused before a body is built, so it is never
|
||||
/// serialized or uploaded to be refused by the service instead.
|
||||
#[tokio::test]
|
||||
async fn an_oversized_secret_value_is_refused_before_the_wire() {
|
||||
let directory = tempfile::tempdir().unwrap();
|
||||
let connection = lancedb::connect(directory.path().to_str().unwrap())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
for value in ["", &"x".repeat(MAX_SECRET_VALUE_BYTES + 1)] {
|
||||
let error = connection
|
||||
.create_secret("openai-prod", value, &[])
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(
|
||||
matches!(error, Error::InvalidInput { .. }),
|
||||
"expected InvalidInput, got {error:?}"
|
||||
);
|
||||
}
|
||||
|
||||
// A local database refuses the verb outright, which is what proves the
|
||||
// size check ran ahead of the backend rather than instead of it.
|
||||
let error = connection
|
||||
.create_secret("openai-prod", "x".repeat(MAX_SECRET_VALUE_BYTES), &[])
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(matches!(error, Error::NotSupported { .. }));
|
||||
}
|
||||
|
||||
+6
-32
@@ -3,9 +3,7 @@
|
||||
"job_type": "create_function",
|
||||
"job_state": "DONE",
|
||||
"creation_ms": 1787270400000,
|
||||
"spec": {
|
||||
"name": "embed"
|
||||
},
|
||||
"spec": {"name": "embed"},
|
||||
"result": {
|
||||
"name": "embed",
|
||||
"version": "fv_01K3EXACT",
|
||||
@@ -15,42 +13,18 @@
|
||||
"entrypoint": "embed"
|
||||
},
|
||||
"signature": {
|
||||
"inputs": [
|
||||
{
|
||||
"name": "text",
|
||||
"arrow_type": "utf8",
|
||||
"nullable": true
|
||||
}
|
||||
],
|
||||
"output": {
|
||||
"kind": "scalar",
|
||||
"arrow_type": "list<float32>",
|
||||
"nullable": false
|
||||
}
|
||||
"inputs": [{"name": "text", "arrow_type": "utf8", "nullable": true}],
|
||||
"output": {"kind": "scalar", "arrow_type": "list<float32>", "nullable": false}
|
||||
},
|
||||
"runtime": {
|
||||
"kind": "python",
|
||||
"python_version": "3.12",
|
||||
"environment": {
|
||||
"kind": "pip",
|
||||
"packages": [
|
||||
"sentence-transformers>=3"
|
||||
]
|
||||
},
|
||||
"env": {
|
||||
"TOKENIZERS_PARALLELISM": "false"
|
||||
}
|
||||
"environment": {"kind": "pip", "packages": ["sentence-transformers>=3"]},
|
||||
"env": {"TOKENIZERS_PARALLELISM": "false"}
|
||||
},
|
||||
"runtime_digest": "sha256:runtime",
|
||||
"environment_digest": "sha256:environment",
|
||||
"secret_env_bindings": {
|
||||
"HF_TOKEN": {
|
||||
"name": "hf-prod"
|
||||
}
|
||||
},
|
||||
"created_at": "2026-08-21T00:00:00Z"
|
||||
},
|
||||
"future_job": {
|
||||
"trace_id": "trace-1"
|
||||
}
|
||||
"future_job": {"trace_id": "trace-1"}
|
||||
}
|
||||
|
||||
-1
@@ -1 +0,0 @@
|
||||
{"artifact":{"adapter":{"kind":"scalar_to_arrow_batch","version":1},"content":{"data":"ZnJvbSBfX2Z1dHVyZV9fIGltcG9ydCBhbm5vdGF0aW9ucwoKZGVmIGFuYWx5emVfY2FwdGlvbihjYXB0aW9uOiBzdHIpIC0+IHN0cjoKICAgIHJldHVybiBjYXB0aW9uLnN0cmlwKCkK","encoding":"base64"},"digest":"sha256:800462c9ad15151a80f83f85b8912ff149300c1563e07f58448f099afcd0d077","entrypoint":"analyze_caption","kind":"python_callable"},"name":"analyze_caption","runtime":{"env":{"MODE":"test"},"environment":{"kind":"pip","packages":["openai==3.7.0"]},"kind":"python","python_version":"3.12"},"secret_env_bindings":{"OPENAI_API_KEY":{"name":"openai-prod"}},"signature":{"inputs":[{"arrow_type":"utf8","name":"caption","nullable":false}],"output":{"arrow_type":"utf8","kind":"scalar","nullable":false}}}
|
||||
-48
@@ -1,48 +0,0 @@
|
||||
{
|
||||
"artifact": {
|
||||
"adapter": {
|
||||
"kind": "scalar_to_arrow_batch",
|
||||
"version": 1
|
||||
},
|
||||
"content": {
|
||||
"data": "ZnJvbSBfX2Z1dHVyZV9fIGltcG9ydCBhbm5vdGF0aW9ucwoKZGVmIGFuYWx5emVfY2FwdGlvbihjYXB0aW9uOiBzdHIpIC0+IHN0cjoKICAgIHJldHVybiBjYXB0aW9uLnN0cmlwKCkK",
|
||||
"encoding": "base64"
|
||||
},
|
||||
"digest": "sha256:800462c9ad15151a80f83f85b8912ff149300c1563e07f58448f099afcd0d077",
|
||||
"entrypoint": "analyze_caption",
|
||||
"kind": "python_callable"
|
||||
},
|
||||
"name": "analyze_caption",
|
||||
"runtime": {
|
||||
"env": {
|
||||
"MODE": "test"
|
||||
},
|
||||
"environment": {
|
||||
"kind": "pip",
|
||||
"packages": [
|
||||
"openai==3.7.0"
|
||||
]
|
||||
},
|
||||
"kind": "python",
|
||||
"python_version": "3.12"
|
||||
},
|
||||
"secret_env_bindings": {
|
||||
"OPENAI_API_KEY": {
|
||||
"name": "openai-prod"
|
||||
}
|
||||
},
|
||||
"signature": {
|
||||
"inputs": [
|
||||
{
|
||||
"arrow_type": "utf8",
|
||||
"name": "caption",
|
||||
"nullable": false
|
||||
}
|
||||
],
|
||||
"output": {
|
||||
"arrow_type": "utf8",
|
||||
"kind": "scalar",
|
||||
"nullable": false
|
||||
}
|
||||
}
|
||||
}
|
||||
Vendored
+1
-1
@@ -1 +1 @@
|
||||
{"artifact":{"digest":"sha256:code","entrypoint":"embed","kind":"python_callable"},"created_at":"2026-08-21T00:00:00Z","environment_digest":"sha256:environment","name":"embed","runtime":{"env":{"TOKENIZERS_PARALLELISM":"false"},"environment":{"kind":"pip","packages":["sentence-transformers>=3"]},"kind":"python","python_version":"3.12"},"runtime_digest":"sha256:runtime","secret_env_bindings":{"HF_TOKEN":{"name":"hf-prod"}},"signature":{"inputs":[{"arrow_type":"utf8","name":"text","nullable":true}],"output":{"arrow_type":"list<float32>","kind":"scalar","nullable":false}},"version":"fv_01K3EXACT"}
|
||||
{"artifact":{"digest":"sha256:code","entrypoint":"embed","kind":"python_callable"},"created_at":"2026-08-21T00:00:00Z","environment_digest":"sha256:environment","name":"embed","runtime":{"env":{"TOKENIZERS_PARALLELISM":"false"},"environment":{"kind":"pip","packages":["sentence-transformers>=3"]},"kind":"python","python_version":"3.12"},"runtime_digest":"sha256:runtime","signature":{"inputs":[{"arrow_type":"utf8","name":"text","nullable":true}],"output":{"arrow_type":"list<float32>","kind":"scalar","nullable":false}},"version":"fv_01K3EXACT"}
|
||||
|
||||
Reference in New Issue
Block a user