Compare commits

..

1 Commits

Author SHA1 Message Date
Lance Release 5c7e31949d Bump version: 0.37.0-beta.0 → 0.37.1-beta.0 2026-07-29 07:11:47 +00:00
146 changed files with 2257 additions and 13096 deletions
+1 -1
View File
@@ -1,5 +1,5 @@
[tool.bumpversion]
current_version = "0.37.1"
current_version = "0.37.1-beta.0"
parse = """(?x)
(?P<major>0|[1-9]\\d*)\\.
(?P<minor>0|[1-9]\\d*)\\.
-222
View File
@@ -1,222 +0,0 @@
name: Check doc links
# Checking external links is inherently noisy: third-party sites rate-limit
# automated clients, reject non-browser user agents, and go down temporarily.
# Blocking pull requests on that trades a lot of false failures for very little
# signal, so this runs on a schedule and reports findings in a single tracking
# issue instead of failing anyone's build.
on:
schedule:
- cron: "0 7 * * *"
workflow_dispatch:
# The report lives in one repository-global issue, so runs must not overlap: a
# lookup racing a create produces duplicate issues, and a healthy run closing
# the issue while a failing run only rewrites its body would leave a broken
# report closed. The group is deliberately ref-independent so that a manual
# dispatch serializes against the scheduled run.
concurrency:
group: docs-link-check
cancel-in-progress: false
permissions: {}
env:
REPORT_TITLE: "Docs link checker report"
jobs:
scan:
name: Scan links
runs-on: ubuntu-24.04
# lychee-action is pinned by SHA, but its wrapper downloads the lychee
# release tarball at run time without verifying a digest, and hands the
# resulting binary a GitHub token. Release assets remain replaceable, so
# that binary is confined to a job whose token can only read public
# content; everything that writes runs in the report job below.
permissions:
contents: read
outputs:
exit_code: ${{ steps.lychee.outputs.exit_code }}
steps:
- name: Checkout
uses: actions/checkout@v6
with:
# workflow_dispatch can run from any ref, but the report is
# repository-global. Always measure the default branch so a manual
# run from a topic branch cannot close a report that main warrants,
# or overwrite it with branch-only findings.
ref: ${{ github.event.repository.default_branch }}
persist-credentials: false
- name: Check links
id: lychee
uses: lycheeverse/lychee-action@e7477775783ea5526144ba13e8db5eec57747ce8 # v2.9.0
with:
# Restricted to http(s) on purpose. Much of docs/src is generated
# API reference (the js/ tree comes from `npm run docs` in nodejs)
# and the hand-written pages use mkdocstrings cross-references and
# nav-relative paths that only resolve in the site mkdocs builds,
# not in this checkout, so relative links would be reported as
# broken on every run.
args: >-
--scheme https
--scheme http
--no-progress
--max-retries 3
--timeout 20
'docs/src/**/*.md'
format: json
output: ./lychee/out.json
jobSummary: false
# The report, not a red build, is the signal for broken links. The
# validation step below still fails the run if the check itself
# breaks.
fail: false
- name: Validate report
# lychee does not reserve exit code 2 for broken links: its CLI
# parser also exits 2 on an invalid option, before any link was
# checked or any report written. Only a parseable report whose
# counts agree with the exit code counts as a link verdict; anything
# else fails here, and the report job below is skipped entirely, so
# the tracking issue is never touched. Exit 2 covers timeouts as
# well as errors, and a timed-out host is exactly the transient
# unavailability this report exists to surface, so both count as
# findings. Requiring total > 0 also catches a glob that silently
# stopped matching any file.
if: steps.lychee.outputs.exit_code == 0 || steps.lychee.outputs.exit_code == 2
env:
EXIT_CODE: ${{ steps.lychee.outputs.exit_code }}
run: |
jq -e --argjson code "$EXIT_CODE" '
(.total > 0) and
(if $code == 0
then .errors == 0 and .timeouts == 0
and (.error_map | length == 0) and (.timeout_map | length == 0)
else (.errors + .timeouts) > 0
and ((.error_map | length) + (.timeout_map | length)) > 0
end)
' ./lychee/out.json
- name: Upload report
if: steps.lychee.outputs.exit_code == 2
uses: actions/upload-artifact@v7
with:
name: link-report
path: ./lychee/out.json
retention-days: 7
report:
name: Update report issue
needs: scan
runs-on: ubuntu-24.04
# Deliberately no checkout: this job needs the report artifact and the
# issues API, not the repository contents.
permissions:
issues: write
env:
EXIT_CODE: ${{ needs.scan.outputs.exit_code }}
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
steps:
- name: Classify checker result
# lychee exits 0 when every link resolves and 2 when links fail,
# both already cross-checked against the report by the scan job's
# validation step. Anything else (1 runtime, 3 bad config) means the
# check never produced a link verdict, which must surface as a failed
# run rather than be published as "broken documentation links".
run: |
case "$EXIT_CODE" in
0|2)
echo "lychee exit code $EXIT_CODE"
;;
*)
echo "::error::lychee exited with '$EXIT_CODE': the link check did not complete. Leaving the report issue untouched."
exit 1
;;
esac
- name: Find existing report issue
id: report
# Matched on title alone, and through search rather than a listing:
# the issue action applies labels in a separate call after creating the
# issue, so a label filter misses a half-created report, and this
# repository has far more open issues than one listing page holds.
# Closed issues are included because a healthy run closes the report:
# an open-only lookup would forget that identity and the next failing
# run would open a duplicate. The oldest match stays the canonical
# report and is reopened below when links break again.
run: |
match=$(gh issue list --repo "$GITHUB_REPOSITORY" --state all \
--search "in:title \"$REPORT_TITLE\" author:app/github-actions" \
--limit 50 --json number,title,state \
--jq "[.[] | select(.title == \"$REPORT_TITLE\")] | sort_by(.number) | first // empty")
echo "number=$(jq -r '.number // empty' <<<"$match")" >> "$GITHUB_OUTPUT"
echo "state=$(jq -r '.state // empty' <<<"$match")" >> "$GITHUB_OUTPUT"
- name: Download report
if: env.EXIT_CODE == 2
uses: actions/download-artifact@v8
with:
name: link-report
path: ./lychee
- name: Compose report
if: env.EXIT_CODE == 2
run: |
run_url="$GITHUB_SERVER_URL/$GITHUB_REPOSITORY/actions/runs/$GITHUB_RUN_ID"
{
echo "Broken documentation links found by [\`$GITHUB_WORKFLOW\`]($run_url)."
echo
echo "This issue is rewritten by every scheduled run and closed automatically once all links resolve."
echo
echo "Entries can be false positives: some sites rate-limit or block automated clients while working fine in a browser. Confirm before editing the docs, and add persistent offenders to \`--exclude\` in \`.github/workflows/docs-link-check.yml\`."
echo
# Timeouts are reported alongside errors: entries land in
# timeout_map with a status text instead of an HTTP code.
jq -r '
"\(.errors) of \(.total) links failed, \(.timeouts) timed out.",
"",
([(.error_map | to_entries[]), (.timeout_map | to_entries[])]
| group_by(.key)[] |
"### Errors in \(.[0].key)",
"",
(map(.value[])[] | "* [\(.status.code // .status.text // "ERR")] <\(.url)> — \(.status.details // .status.text // "unknown error")"),
"")
' ./lychee/out.json
} > ./lychee/issue.md
- name: Reopen report issue
# A healthy run closes the report, and the issue action below only
# rewrites the body of whatever number it is given. Without an
# explicit reopen, the 2 -> 0 -> 2 sequence would keep rewriting a
# closed issue while links are broken. A CLOSED state implies the
# lookup found a canonical issue, so no separate emptiness check.
if: env.EXIT_CODE == 2 && steps.report.outputs.state == 'CLOSED'
env:
ISSUE_NUMBER: ${{ steps.report.outputs.number }}
run: |
run_url="$GITHUB_SERVER_URL/$GITHUB_REPOSITORY/actions/runs/$GITHUB_RUN_ID"
gh issue reopen "$ISSUE_NUMBER" --repo "$GITHUB_REPOSITORY" \
--comment "Broken documentation links found again in [the latest run]($run_url)."
- name: Report broken links
if: env.EXIT_CODE == 2
uses: peter-evans/create-issue-from-file@fca9117c27cdc29c6c4db3b86c48e4115a786710 # v6.0.0
with:
# Empty on the first failing run, which creates the issue; afterwards
# the same issue is updated in place.
issue-number: ${{ steps.report.outputs.number }}
title: ${{ env.REPORT_TITLE }}
content-filepath: ./lychee/issue.md
labels: documentation
- name: Close report issue once links are healthy
# An OPEN state implies the lookup found a canonical issue; a report
# that is already closed needs nothing.
if: env.EXIT_CODE == 0 && steps.report.outputs.state == 'OPEN'
env:
ISSUE_NUMBER: ${{ steps.report.outputs.number }}
run: |
run_url="$GITHUB_SERVER_URL/$GITHUB_REPOSITORY/actions/runs/$GITHUB_RUN_ID"
gh issue close "$ISSUE_NUMBER" --repo "$GITHUB_REPOSITORY" \
--comment "All documentation links resolved in [the latest run]($run_url)."
+2 -4
View File
@@ -1,8 +1,7 @@
name: Create release commit
# This workflow increments the version, updates lockfiles, tags the final
# commit, and pushes it. All SDKs share a single version, so one tag releases
# all of them.
# This workflow increments the version, tags it, and pushes it. All SDKs share
# a single version, so one tag releases all of them.
# When a tag is pushed, another workflow is triggered that creates a GH release
# and uploads the binaries. This workflow is only for creating the tag.
@@ -64,7 +63,6 @@ jobs:
pip install bump-my-version PyGithub packaging
bash ci/bump_version.sh ${{ inputs.type }} ${{ inputs.bump-minor }}
bash ci/update_lockfiles.sh --amend
bash ci/create_release_tag.sh
- name: Push new version tag
if: ${{ !inputs.dry_run }}
uses: ad-m/github-push-action@881a6320fdb16eb5318c5054f31c218aec2b324c # v1.3.0
+4 -110
View File
@@ -4,16 +4,6 @@ on:
push:
tags:
- 'v*'
workflow_dispatch:
inputs:
release_tag:
description: Stable release tag to recover (for example, v0.37.1)
required: true
type: string
source_run_id:
description: Failed tag workflow run containing the completed non-Windows wheels
required: true
type: string
pull_request:
# This should trigger a dry run (we skip the final publish step)
paths:
@@ -39,7 +29,6 @@ concurrency:
jobs:
linux:
name: Python ${{ matrix.config.package_name }} ${{ matrix.config.platform }} manylinux${{ matrix.config.manylinux }}
if: github.event_name != 'workflow_dispatch'
timeout-minutes: 60
strategy:
matrix:
@@ -95,7 +84,6 @@ jobs:
path: target/wheels/*.whl
if-no-files-found: error
mac:
if: github.event_name != 'workflow_dispatch'
timeout-minutes: 90
runs-on: ${{ matrix.config.runner }}
strategy:
@@ -125,8 +113,7 @@ jobs:
path: target/wheels/lancedb-*.whl
if-no-files-found: error
windows:
if: github.event_name != 'workflow_dispatch'
timeout-minutes: 120
timeout-minutes: 90
runs-on: windows-latest
env:
# link.exe is single-threaded and the long pole on Windows builds. Use
@@ -158,32 +145,6 @@ jobs:
name: wheels-windows
path: target/wheels/lancedb-*.whl
if-no-files-found: error
recover-windows:
name: Recover Windows wheel
if: github.event_name == 'workflow_dispatch'
timeout-minutes: 120
runs-on: windows-latest
env:
CARGO_TARGET_X86_64_PC_WINDOWS_MSVC_LINKER: rust-lld
steps:
- uses: actions/checkout@v6
with:
ref: ${{ inputs.release_tag }}
fetch-depth: 0
lfs: true
- name: Set up Python
uses: actions/setup-python@v6
with:
python-version: "3.13"
- uses: ./.github/workflows/build_windows_wheel
with:
python-minor-version: 10
args: "--release --strip"
- uses: actions/upload-artifact@v7
with:
name: wheels-windows-recovery
path: target/wheels/lancedb-*.whl
if-no-files-found: error
publish:
name: Publish wheels
if: startsWith(github.ref, 'refs/tags/v')
@@ -204,13 +165,11 @@ jobs:
run: ls -la target/wheels
- name: Choose repo
id: choose_repo
env:
RELEASE_REF: ${{ github.ref }}
run: |
if [[ "$RELEASE_REF" == *beta* ]]; then
echo "repo=fury" >> "$GITHUB_OUTPUT"
if [[ ${{ github.ref }} == *beta* ]]; then
echo "repo=fury" >> $GITHUB_OUTPUT
else
echo "repo=pypi" >> "$GITHUB_OUTPUT"
echo "repo=pypi" >> $GITHUB_OUTPUT
fi
- name: Publish to Fury
if: steps.choose_repo.outputs.repo == 'fury'
@@ -237,71 +196,6 @@ jobs:
uses: pypa/gh-action-pypi-publish@release/v1
with:
packages-dir: target/wheels/
recover-publish:
name: Recover PyPI publish
if: github.event_name == 'workflow_dispatch'
needs: [recover-windows]
runs-on: ubuntu-latest
permissions:
actions: read
id-token: write
contents: read
steps:
- uses: actions/checkout@v6
with:
ref: ${{ inputs.release_tag }}
fetch-depth: 0
- name: Verify recovery source
env:
GH_TOKEN: ${{ github.token }}
RELEASE_TAG: ${{ inputs.release_tag }}
SOURCE_RUN_ID: ${{ inputs.source_run_id }}
run: |
if [[ "$RELEASE_TAG" != v* || "$RELEASE_TAG" == *beta* ]]; then
echo "Recovery only supports stable v* release tags" >&2
exit 1
fi
TAG_SHA=$(git rev-parse HEAD)
RUN_SHA=$(gh api "/repos/${{ github.repository }}/actions/runs/$SOURCE_RUN_ID" --jq .head_sha)
if [[ "$TAG_SHA" != "$RUN_SHA" ]]; then
echo "Source run $SOURCE_RUN_ID ($RUN_SHA) does not match $RELEASE_TAG ($TAG_SHA)" >&2
exit 1
fi
- name: Download Linux wheel artifacts
uses: actions/download-artifact@v8
with:
github-token: ${{ github.token }}
repository: ${{ github.repository }}
run-id: ${{ inputs.source_run_id }}
pattern: wheels-linux-*
path: target/wheels
merge-multiple: true
- name: Download macOS wheel artifacts
uses: actions/download-artifact@v8
with:
github-token: ${{ github.token }}
repository: ${{ github.repository }}
run-id: ${{ inputs.source_run_id }}
pattern: wheels-mac-*
path: target/wheels
merge-multiple: true
- name: Download recovered Windows wheel
uses: actions/download-artifact@v8
with:
name: wheels-windows-recovery
path: target/wheels
- name: Validate recovered wheels
run: |
find target/wheels -maxdepth 1 -type f -name '*.whl' -print
WHEEL_COUNT=$(find target/wheels -maxdepth 1 -type f -name '*.whl' | wc -l)
if [[ "$WHEEL_COUNT" -ne 5 ]]; then
echo "Expected 5 wheels, found $WHEEL_COUNT" >&2
exit 1
fi
- name: Publish to PyPI
uses: pypa/gh-action-pypi-publish@release/v1
with:
packages-dir: target/wheels/
report-failure:
name: Report Workflow Failure
runs-on: ubuntu-latest
+6 -8
View File
@@ -296,18 +296,16 @@ jobs:
cargo update -p aws-types --precise 1.3.9
cargo update -p aws-sigv4 --precise 1.3.5
cargo update -p aws-credential-types --precise 1.2.8
# aws-smithy-checksums must stay at or above 0.63.13: OpenDAL's S3
# service needs crc-fast ~1.9, and older releases pin it to ~1.3.
cargo update -p aws-smithy-checksums --precise 0.63.13
cargo update -p aws-smithy-checksums --precise 0.63.9
cargo update -p aws-smithy-runtime --precise 1.9.3
cargo update -p aws-smithy-http --precise 0.62.6
cargo update -p aws-smithy-eventstream --precise 0.60.14
cargo update -p aws-smithy-http --precise 0.62.4
cargo update -p aws-smithy-eventstream --precise 0.60.12
cargo update -p aws-smithy-http-client --precise 1.1.3
cargo update -p aws-smithy-observability --precise 0.1.4
cargo update -p aws-smithy-query --precise 0.60.8
cargo update -p aws-smithy-runtime-api --precise 1.9.3
cargo update -p aws-smithy-async --precise 1.2.7
cargo update -p aws-smithy-types --precise 1.3.6
cargo update -p aws-smithy-runtime-api --precise 1.9.1
cargo update -p aws-smithy-async --precise 1.2.6
cargo update -p aws-smithy-types --precise 1.3.5
cargo update -p aws-smithy-xml --precise 0.60.11
cargo update -p home --precise 0.5.9
- name: cargo +${{ matrix.msrv }} check
-29
View File
@@ -92,8 +92,6 @@ Python bindings changes:
* Should use `LOOP.run()` to call the corresponding `AsyncTable` method.
6. Add concrete sync method to `RemoteTable` class in `python/python/lancedb/remote/table.py`.
7. Add unit test in `python/tests/test_table.py`.
8. If you added a new public class or module-level function (not just a method on an
existing class), expose it in the API reference. See "Python API reference" below.
TypeScript bindings changes:
@@ -105,33 +103,6 @@ TypeScript bindings changes:
5. Add test in `nodejs/__test__/table.test.ts`.
6. Run `npm run docs` to generate TypeScript documentation.
## Python API reference
`docs/src/python/python.md` is the entire Python API reference. It is maintained by
hand, and anything not listed there is not rendered at all, so new public classes and
module-level functions have to be added explicitly. How depends on the module:
* `lancedb.index`, `lancedb.embeddings`, `lancedb.remote`, and `lancedb.rerankers` are
rendered by a single directive each, driven by the module's `__all__`. Add the new
name to `__all__` and it appears; forget, and it is silently omitted.
* Everything else (`lancedb`, `lancedb.table`, `lancedb.query`, `lancedb.db`, ...) is
listed symbol by symbol. Add a `::: lancedb.<module>.<Name>` line to the matching
section, and remember that the page separates synchronous and asynchronous APIs.
Deliberately undocumented: concrete implementations reached through an abstract base
(`LanceTable`, `LanceDBConnection`, `RemoteDBConnection`), query base classes already
covered by `inherited_members`, and internal helpers.
Cross-references in docstrings use mkdocstrings syntax, `[text][lancedb.table.Table]`.
Plain relative links such as `[Table](Table)` do not resolve. To check your work:
```shell
pip install -r docs/requirements.txt
cd docs && PYTHONPATH=. mkdocs build
```
The docs site only builds on pushes to `main`, so this is not covered by PR CI.
## Review Guidelines
Please consider the following when reviewing code contributions.
Generated
+211 -259
View File
File diff suppressed because it is too large Load Diff
+15 -15
View File
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
rust-version = "1.91.0"
[workspace.dependencies]
lance = { "version" = "=10.0.0", default-features = false }
lance-core = "=10.0.0"
lance-datagen = "=10.0.0"
lance-file = "=10.0.0"
lance-io = { "version" = "=10.0.0", default-features = false }
lance-index = "=10.0.0"
lance-linalg = "=10.0.0"
lance-namespace = "=10.0.0"
lance-namespace-impls = { "version" = "=10.0.0", default-features = false }
lance-table = "=10.0.0"
lance-testing = "=10.0.0"
lance-datafusion = "=10.0.0"
lance-encoding = "=10.0.0"
lance-arrow = "=10.0.0"
lance = { "version" = "=10.0.0-beta.5", default-features = false, "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=10.0.0-beta.5", default-features = false, "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=10.0.0-beta.5", default-features = false, "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
ahash = "0.8"
# Note that this one does not include pyarrow
arrow = { version = "58.0.0", optional = false }
@@ -52,7 +52,7 @@ env_logger = "0.11"
half = { "version" = "2.7.1", default-features = false, features = [
"num-traits",
] }
futures = "0.3"
futures = "0"
log = "0.4"
metrics = "0.24"
metrics-util = "0.19"
+9 -4
View File
@@ -6,11 +6,16 @@ HEAD_SHA=$(git rev-parse HEAD)
readonly TAG_PREFIX="v"
readonly SELF_DIR=$(cd "$( dirname "${BASH_SOURCE[0]}" )" && pwd )
readonly BUMP_ARGS="--no-tag"
PREV_TAG=$(git tag --sort='version:refname' | grep ^$TAG_PREFIX | python $SELF_DIR/semver_sort.py $TAG_PREFIX | tail -n 1)
echo "Found previous tag $PREV_TAG"
# Initially, we don't want to tag if we are doing stable, because we will bump
# again later. See comment at end for why.
if [[ "$RELEASE_TYPE" == 'stable' ]]; then
BUMP_ARGS="--no-tag"
fi
# If last is stable and not bumping minor
if [[ $PREV_TAG != *beta* ]]; then
if [[ "$BUMP_MINOR" != "false" ]]; then
@@ -34,12 +39,12 @@ fi
# a stable version, bump the pre-release level ("pre_l") to make it stable.
if [[ $RELEASE_TYPE == 'stable' ]]; then
# X.Y.Z-beta.N -> X.Y.Z
bump-my-version bump -vv $BUMP_ARGS pre_l
bump-my-version bump -vv pre_l
fi
# Validate that we have incremented version appropriately for breaking changes
NEW_VERSION=$(python -c 'import tomllib; print(tomllib.load(open(".bumpversion.toml", "rb"))["tool"]["bumpversion"]["current_version"])')
NEW_TAG="$TAG_PREFIX$NEW_VERSION"
NEW_TAG=$(git describe --tags --exact-match HEAD)
NEW_VERSION=$(echo $NEW_TAG | sed "s/^$TAG_PREFIX//")
LAST_STABLE_RELEASE=$(git tag --sort='version:refname' | grep ^$TAG_PREFIX | grep -v beta | grep -vF "$NEW_TAG" | python $SELF_DIR/semver_sort.py $TAG_PREFIX | tail -n 1)
LAST_STABLE_VERSION=$(echo $LAST_STABLE_RELEASE | sed "s/^$TAG_PREFIX//")
-21
View File
@@ -1,21 +0,0 @@
#!/usr/bin/env bash
set -euo pipefail
RELEASE_VERSION=$(python -c 'import tomllib; print(tomllib.load(open(".bumpversion.toml", "rb"))["tool"]["bumpversion"]["current_version"])')
RELEASE_TAG="v${RELEASE_VERSION}"
if git rev-parse --quiet --verify "refs/tags/${RELEASE_TAG}" >/dev/null; then
echo "Release tag ${RELEASE_TAG} already exists" >&2
exit 1
fi
git tag --annotate "$RELEASE_TAG" --message "Release ${RELEASE_TAG}"
HEAD_SHA=$(git rev-parse HEAD)
TAG_SHA=$(git rev-parse "refs/tags/${RELEASE_TAG}^{}")
if [[ "$TAG_SHA" != "$HEAD_SHA" ]]; then
echo "Release tag ${RELEASE_TAG} points to ${TAG_SHA}, expected ${HEAD_SHA}" >&2
exit 1
fi
echo "Created ${RELEASE_TAG} at ${HEAD_SHA}"
-72
View File
@@ -1,72 +0,0 @@
import tempfile
import unittest
from pathlib import Path
from validate_stable_lance import validate
class ValidateStableLanceTest(unittest.TestCase):
def write_fixture(
self,
root: Path,
*,
rust: str = "=10.0.0",
python: str = "10.0.0",
java: str = "10.0.0",
) -> None:
(root / "python").mkdir()
(root / "java").mkdir()
(root / "Cargo.toml").write_text(
f'[workspace.dependencies]\nlance = "{rust}"\nlance-core = "{rust}"\n'
)
(root / "python" / "pyproject.toml").write_text(
'[project.optional-dependencies]\ntests = ["pylance==' + python + '"]\n'
)
(root / "java" / "pom.xml").write_text(
"<project><properties><lance-core.version>"
+ java
+ "</lance-core.version></properties></project>"
)
def test_accepts_matching_stable_versions(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
root = Path(temp_dir)
self.write_fixture(root)
self.assertEqual(validate(root), "10.0.0")
def test_rejects_prerelease(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
root = Path(temp_dir)
self.write_fixture(root, python="10.0.0rc1")
with self.assertRaisesRegex(ValueError, "not stable"):
validate(root)
def test_rejects_sdk_mismatch(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
root = Path(temp_dir)
self.write_fixture(root, java="9.0.0")
with self.assertRaisesRegex(ValueError, "do not match across SDKs"):
validate(root)
def test_rejects_non_exact_rust_dependency(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
root = Path(temp_dir)
self.write_fixture(root, rust="10.0.0")
with self.assertRaisesRegex(ValueError, "not exact"):
validate(root)
def test_rejects_unpublished_rust_source(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
root = Path(temp_dir)
self.write_fixture(root)
(root / "Cargo.toml").write_text(
"[workspace.dependencies]\n"
'lance = { version = "=10.0.0", git = "https://example.com/lance", '
'tag = "v10.0.0" }\n'
)
with self.assertRaisesRegex(ValueError, "unpublished source fields"):
validate(root)
if __name__ == "__main__":
unittest.main()
Executable → Regular
+26 -104
View File
@@ -1,112 +1,34 @@
#!/usr/bin/env python3
"""Validate that every SDK uses the same published stable Lance release."""
from __future__ import annotations
import re
import xml.etree.ElementTree as ET
from pathlib import Path
import tomllib
STABLE_VERSION = re.compile(r"[0-9]+\.[0-9]+\.[0-9]+")
found_preview_lance = False
with open("Cargo.toml", "rb") as f:
cargo_data = tomllib.load(f)
def _stable_version(raw: str, *, dependency: str, exact: bool = False) -> str:
value = raw.strip()
if exact and not value.startswith("="):
raise ValueError(f"Dependency '{dependency}' is not exact: {raw}")
value = value.removeprefix("=").removeprefix("v")
if STABLE_VERSION.fullmatch(value) is None:
raise ValueError(f"Dependency '{dependency}' is not stable: {raw}")
return value
for name, dep in cargo_data["workspace"]["dependencies"].items():
if name == "lance" or name.startswith("lance-"):
if isinstance(dep, str):
version = dep
elif isinstance(dep, dict):
# Version doesn't have the beta tag in it, so we instead look
# at the git tag.
version = dep.get('tag', dep.get('version'))
else:
raise ValueError("Unexpected type for dependency: " + str(dep))
if "beta" in version:
found_preview_lance = True
print(f"Dependency '{name}' is a preview version: {version}")
def rust_lance_version(repo_root: Path) -> str:
with (repo_root / "Cargo.toml").open("rb") as cargo_file:
dependencies = tomllib.load(cargo_file)["workspace"]["dependencies"]
with open("python/pyproject.toml", "rb") as f:
py_proj_data = tomllib.load(f)
versions: dict[str, str] = {}
for name, dependency in dependencies.items():
if name != "lance" and not name.startswith("lance-"):
continue
for dep in py_proj_data["project"]["dependencies"]:
if dep.startswith("pylance"):
if "b" in dep:
found_preview_lance = True
print(f"Dependency '{dep}' is a preview version")
break # Only one pylance dependency
if isinstance(dependency, str):
raw_version = dependency
elif isinstance(dependency, dict):
forbidden_sources = [
source
for source in ("git", "path", "branch", "rev", "tag")
if source in dependency
]
if forbidden_sources:
joined = ", ".join(forbidden_sources)
raise ValueError(
f"Dependency '{name}' uses unpublished source fields: {joined}"
)
raw_version = dependency.get("version")
if raw_version is None:
raise ValueError(f"Dependency '{name}' has no version")
else:
raise TypeError(f"Dependency '{name}' has an unexpected definition")
versions[name] = _stable_version(raw_version, dependency=name, exact=True)
if not versions:
raise ValueError("No Rust Lance dependencies found")
unique_versions = set(versions.values())
if len(unique_versions) != 1:
details = ", ".join(f"{name}={version}" for name, version in versions.items())
raise ValueError(f"Rust Lance dependency versions do not match: {details}")
return unique_versions.pop()
def python_lance_version(repo_root: Path) -> str:
with (repo_root / "python" / "pyproject.toml").open("rb") as pyproject_file:
pyproject = tomllib.load(pyproject_file)
requirements = pyproject["project"]["optional-dependencies"]["tests"]
pylance_requirements = [
requirement for requirement in requirements if requirement.startswith("pylance")
]
if len(pylance_requirements) != 1:
raise ValueError(
"Expected exactly one pylance requirement in the Python test dependencies"
)
requirement = pylance_requirements[0]
prefix = "pylance=="
if not requirement.startswith(prefix):
raise ValueError(f"Python test dependency is not exact: {requirement}")
return _stable_version(requirement[len(prefix) :], dependency="pylance")
def java_lance_version(repo_root: Path) -> str:
pom_root = ET.parse(repo_root / "java" / "pom.xml").getroot()
versions = [
element.text
for element in pom_root.iter()
if element.tag.rsplit("}", 1)[-1] == "lance-core.version"
]
if len(versions) != 1 or versions[0] is None:
raise ValueError("Expected exactly one Java lance-core.version property")
return _stable_version(versions[0], dependency="Java lance-core")
def validate(repo_root: Path) -> str:
versions = {
"Rust": rust_lance_version(repo_root),
"Python": python_lance_version(repo_root),
"Java": java_lance_version(repo_root),
}
if len(set(versions.values())) != 1:
details = ", ".join(f"{sdk}={version}" for sdk, version in versions.items())
raise ValueError(
f"Lance dependency versions do not match across SDKs: {details}"
)
return versions["Rust"]
if __name__ == "__main__":
version = validate(Path(__file__).resolve().parents[1])
print(f"Validated published stable Lance v{version} across Rust, Python, and Java")
if found_preview_lance:
raise ValueError("Found preview version of Lance in dependencies")
-5
View File
@@ -51,11 +51,6 @@ plugins:
paths: [../python/python]
options:
docstring_style: numpy
docstring_options:
# Attributes documented in a `Parameters` section, and pydantic
# dataclasses whose `__init__` griffe cannot see statically, both
# trip this check. It reports nothing actionable here.
warn_unknown_params: false
heading_level: 3
show_signature_annotations: true
show_root_heading: true
+1 -11
View File
@@ -453,16 +453,6 @@ paths:
The metric type to use for the index. l2, Cosine, Dot are supported.
index_type:
type: string
custom_stop_words:
type: [array, "null"]
items:
type: string
description: |
The custom stop-word list for an FTS index. A non-null
array replaces the language's built-in stop-word list and is only
applied when remove_stop_words is enabled. Null uses the built-in
language list, while an empty array explicitly replaces it with no
stop words.
responses:
"200":
description: Index successfully created
@@ -520,4 +510,4 @@ paths:
"401":
$ref: "#/components/responses/unauthorized"
"404":
$ref: "#/components/responses/not_found"
$ref: "#/components/responses/not_found"
+1 -1
View File
@@ -14,7 +14,7 @@ Add the following dependency to your `pom.xml`:
<dependency>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-core</artifactId>
<version>0.37.1</version>
<version>0.37.1-beta.0</version>
</dependency>
```
+1 -1
View File
@@ -1,7 +1,7 @@
# Contributing to LanceDB Typescript
This document outlines the process for contributing to LanceDB Typescript.
For general contribution guidelines, see [CONTRIBUTING.md](https://github.com/lancedb/lancedb/blob/main/CONTRIBUTING.md).
For general contribution guidelines, see [CONTRIBUTING.md](../CONTRIBUTING.md).
## Project layout
-97
View File
@@ -25,27 +25,6 @@ the underlying connection has been closed.
## Methods
### cancelJob()
```ts
abstract cancelJob(jobId): Promise<boolean>
```
Request cancellation of a server-side job by id.
Resolves to true if the server accepted the cancellation, false if no
such job exists. Cancelling an already-terminal job is a no-op success.
#### Parameters
* **jobId**: `string`
#### Returns
`Promise`&lt;`boolean`&gt;
***
### cloneTable()
```ts
@@ -386,26 +365,6 @@ Drop an existing table.
***
### getJob()
```ts
abstract getJob(jobId): Promise<null | JobDescription>
```
Describe a single server-side job by id.
Resolves to `null` when the server has no such job.
#### Parameters
* **jobId**: `string`
#### Returns
`Promise`&lt;`null` \| [`JobDescription`](../interfaces/JobDescription.md)&gt;
***
### isOpen()
```ts
@@ -420,62 +379,6 @@ Return true if the connection has not been closed
***
### job()
```ts
abstract job(jobId): Job
```
A [Job](Job.md) handle for a server-side job by id.
The handle is constructed without a server round trip; an unknown id
surfaces when the handle is used. Dropping the handle has no effect on
the job itself.
#### Parameters
* **jobId**: `string`
#### Returns
[`Job`](Job.md)
***
### jobHistory()
```ts
abstract jobHistory(jobId?): Promise<Table<any>>
```
The lifecycle event history of a server-side job, as an Arrow table.
Lists history across all jobs when `jobId` is omitted.
#### Parameters
* **jobId?**: `string`
#### Returns
`Promise`&lt;`Table`&lt;`any`&gt;&gt;
***
### listJobs()
```ts
abstract listJobs(): Promise<JobInfo[]>
```
List server-side jobs across the database's tables.
#### Returns
`Promise`&lt;[`JobInfo`](../interfaces/JobInfo.md)[]&gt;
***
### listNamespaces()
```ts
-83
View File
@@ -1,83 +0,0 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / Job
# Class: Job
A handle to an operation that may still be running.
## Constructors
### new Job()
```ts
new Job(): Job
```
#### Returns
[`Job`](Job.md)
## Accessors
### id
```ts
get id(): null | string
```
Identifies the operation on the server that is running it. Operations
that run in this process have no server id. The value is opaque.
#### Returns
`null` \| `string`
## Methods
### cancel()
```ts
cancel(): Promise<void>
```
Request cancellation. Cancelling a finished operation is a no-op.
#### Returns
`Promise`&lt;`void`&gt;
***
### status()
```ts
status(): Promise<string>
```
The operation's current lifecycle state: "running", "finished",
"failed", or "cancelled".
A point snapshot; unlike [Job.wait](Job.md#wait) it does not block or reject
on a terminal failure state. States a newer server reports that this
client version does not know pass through as-is.
#### Returns
`Promise`&lt;`string`&gt;
***
### wait()
```ts
wait(): Promise<void>
```
Wait until the operation reaches a terminal state.
#### Returns
`Promise`&lt;`void`&gt;
-23
View File
@@ -295,29 +295,6 @@ await table.createIndex("my_float_col");
***
### createIndexAsync()
```ts
abstract createIndexAsync(column, options?): Promise<Job>
```
Create an index, returning a handle to the indexing job.
The job may already be complete when returned; callers must not assume
the index exists until [Job.wait](Job.md#wait) resolves.
#### Parameters
* **column**: `string`
* **options?**: `Partial`&lt;[`IndexOptions`](../interfaces/IndexOptions.md)&gt;
#### Returns
`Promise`&lt;[`Job`](Job.md)&gt;
***
### currentBranch()
```ts
-4
View File
@@ -25,7 +25,6 @@
- [Connection](classes/Connection.md)
- [HeaderProvider](classes/HeaderProvider.md)
- [Index](classes/Index.md)
- [Job](classes/Job.md)
- [MakeArrowTableOptions](classes/MakeArrowTableOptions.md)
- [MatchQuery](classes/MatchQuery.md)
- [MergeInsertBuilder](classes/MergeInsertBuilder.md)
@@ -89,9 +88,6 @@
- [IvfFlatOptions](interfaces/IvfFlatOptions.md)
- [IvfPqOptions](interfaces/IvfPqOptions.md)
- [IvfRqOptions](interfaces/IvfRqOptions.md)
- [JobDescription](interfaces/JobDescription.md)
- [JobFailureInfo](interfaces/JobFailureInfo.md)
- [JobInfo](interfaces/JobInfo.md)
- [ListNamespacesOptions](interfaces/ListNamespacesOptions.md)
- [ListNamespacesResponse](interfaces/ListNamespacesResponse.md)
- [LsmWriteSpec](interfaces/LsmWriteSpec.md)
-15
View File
@@ -56,21 +56,6 @@ the experimental FTS V3 format and may introduce breaking changes.
***
### customStopWords?
```ts
optional customStopWords: string[];
```
Custom stop words that replace the built-in list for `language`.
This option only affects tokenization when `removeStopWords` is true.
`undefined` keeps the built-in language list. An empty array explicitly
replaces it with no stop words.
***
### language?
```ts
-66
View File
@@ -1,66 +0,0 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / JobDescription
# Interface: JobDescription
A described job from `Connection.getJob`.
## Properties
### creationMs
```ts
creationMs: number;
```
When the job was created, in milliseconds since the epoch.
***
### failure?
```ts
optional failure: JobFailureInfo;
```
Why the job failed, when the job is failed and the server reports a
reason.
***
### jobId
```ts
jobId: string;
```
***
### jobType
```ts
jobType: string;
```
***
### specJson?
```ts
optional specJson: string;
```
The job-type-specific specification as a JSON string, when present.
***
### state
```ts
state: string;
```
Lifecycle state: "running", "finished", "failed", or "cancelled".
-33
View File
@@ -1,33 +0,0 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / JobFailureInfo
# Interface: JobFailureInfo
The server's account of why a job failed.
## Properties
### message?
```ts
optional message: string;
```
***
### phase?
```ts
optional phase: string;
```
***
### retryable?
```ts
optional retryable: boolean;
```
-58
View File
@@ -1,58 +0,0 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / JobInfo
# Interface: JobInfo
A row from `Connection.listJobs`: one server-side job.
## Properties
### createdAtMillis
```ts
createdAtMillis: number;
```
When the job was created, in milliseconds since the epoch.
***
### jobId
```ts
jobId: string;
```
The job id -- what `Connection.getJob` and `Connection.cancelJob`
accept.
***
### jobType
```ts
jobType: string;
```
***
### state
```ts
state: string;
```
Lifecycle state: "running", "finished", "failed", or "cancelled".
***
### table
```ts
table: string;
```
The table the job runs against, without URI or namespace.
+1 -4
View File
@@ -44,7 +44,4 @@ The number of rows in the table
totalBytes: number;
```
The total size, in bytes, of the table's data files, index files, and
overlay files
Read from the manifest, so this excludes deletion files and manifests.
The total number of bytes in the table
-15
View File
@@ -30,21 +30,6 @@ The tokenizer to use. The default is "simple".
***
### customStopWords?
```ts
optional customStopWords: string[];
```
Custom stop words that replace the built-in list for `language`.
This option only affects tokenization when `removeStopWords` is true.
`undefined` keeps the built-in language list. An empty array explicitly
replaces it with no stop words.
***
### language?
```ts
+52 -141
View File
@@ -26,18 +26,6 @@ is also an [asynchronous API client](#connections-asynchronous).
::: lancedb.db.DBConnection
::: lancedb.Session
## Namespaces (Synchronous)
A namespace-backed connection resolves tables through a
[Lance namespace](https://lance-format.github.io/lance-namespace/) service instead of
listing a storage directory.
::: lancedb.connect_namespace
::: lancedb.namespace.LanceNamespaceDBConnection
## Tables (Synchronous)
::: lancedb.table.Table
@@ -46,12 +34,8 @@ listing a storage directory.
::: lancedb.table.FragmentSummaryStats
::: lancedb.table.TableStatistics
::: lancedb.table.Tags
::: lancedb.table.Branches
## Expressions
Type-safe expression builder for filters and projections. Use these instead
@@ -78,46 +62,29 @@ of raw SQL strings with [where][lancedb.query.LanceQueryBuilder.where] and
::: lancedb.query.LanceHybridQueryBuilder
::: lancedb.query.LanceEmptyQueryBuilder
::: lancedb.query.LanceTakeQueryBuilder
## Full text queries
Structured full text queries can be passed to
[Table.search][lancedb.table.Table.search] or
[AsyncTable.search][lancedb.table.AsyncTable.search] in place of a query string,
and combined with [BooleanQuery][lancedb.query.BooleanQuery].
::: lancedb.query.FullTextQuery
::: lancedb.query.MatchQuery
::: lancedb.query.PhraseQuery
::: lancedb.query.BoostQuery
::: lancedb.query.MultiMatchQuery
::: lancedb.query.BooleanQuery
::: lancedb.query.FullTextOperator
::: lancedb.query.Occur
## Embeddings
::: lancedb.embeddings
options:
show_root_heading: false
show_root_toc_entry: false
::: lancedb.embeddings.registry.EmbeddingFunctionRegistry
::: lancedb.embeddings.base.EmbeddingFunctionConfig
::: lancedb.embeddings.base.EmbeddingFunction
::: lancedb.embeddings.base.TextEmbeddingFunction
::: lancedb.embeddings.sentence_transformers.SentenceTransformerEmbeddings
::: lancedb.embeddings.openai.OpenAIEmbeddings
::: lancedb.embeddings.open_clip.OpenClipEmbeddings
## Remote configuration
::: lancedb.remote
options:
show_root_heading: false
show_root_toc_entry: false
::: lancedb.remote.ClientConfig
::: lancedb.remote.TimeoutConfig
::: lancedb.remote.RetryConfig
## Context
@@ -127,50 +94,11 @@ and combined with [BooleanQuery][lancedb.query.BooleanQuery].
## Full text search
Pass `custom_stop_words` to [lancedb.index.FTS][]:
Use [lancedb.table.Table.create_fts_index][] for the synchronous API or
[lancedb.table.AsyncTable.create_index][] with [lancedb.index.FTS][] for the
asynchronous API.
```python
from lancedb.index import FTS
table.create_index(
"text",
config=FTS(remove_stop_words=True, custom_stop_words=["acme", "internal"]),
)
```
The list replaces the built-in stop words and is used only when
`remove_stop_words=True`:
- `custom_stop_words=None` uses the built-in list for `language`.
- `custom_stop_words=[]` removes no words.
- Values are passed through without trimming, lowercasing, or other rewriting.
The same option is available on `lancedb.tokenize(...)` and the deprecated
[lancedb.table.Table.create_fts_index][] compatibility helper:
```python
import lancedb
tokens = list(lancedb.tokenize("acme makes searchable data",
custom_stop_words=["acme"]))
```
::: lancedb.tokenize
::: lancedb.FtsToken
## Blobs
Blob columns store large binary values out of line so they can be read lazily
instead of being materialized with the rest of the row.
::: lancedb.blob
::: lancedb.BlobType
::: lancedb._blob.BlobFile
options:
show_root_full_path: false
::: lancedb.index.FTS
## Utilities
@@ -178,14 +106,6 @@ instead of being materialized with the rest of the row.
::: lancedb.merge.LanceMergeInsertBuilder
::: lancedb.otel.instrument_lancedb_metrics
## Exceptions
::: lancedb.exceptions.MissingValueError
::: lancedb.exceptions.MissingColumnError
## Integrations
## Pydantic
@@ -194,30 +114,19 @@ instead of being materialized with the rest of the row.
::: lancedb.pydantic.vector
::: lancedb.pydantic.Vector
::: lancedb.pydantic.MultiVector
::: lancedb.pydantic.LanceModel
## PyTorch
::: lancedb.streaming.StreamingDataset
::: lancedb.permutation.permutation_builder
::: lancedb.permutation.PermutationBuilder
::: lancedb.permutation.Permutation
::: lancedb.permutation.Transforms
## Reranking
::: lancedb.rerankers
options:
show_root_heading: false
show_root_toc_entry: false
::: lancedb.rerankers.linear_combination.LinearCombinationReranker
::: lancedb.rerankers.cohere.CohereReranker
::: lancedb.rerankers.colbert.ColbertReranker
::: lancedb.rerankers.cross_encoder.CrossEncoderReranker
::: lancedb.rerankers.openai.OpenaiReranker
## Connections (Asynchronous)
@@ -228,12 +137,6 @@ can be used to create, list, or open tables.
::: lancedb.db.AsyncConnection
## Namespaces (Asynchronous)
::: lancedb.connect_namespace_async
::: lancedb.namespace.AsyncLanceNamespaceDBConnection
## Tables (Asynchronous)
Table hold your actual data as a collection of records / rows.
@@ -242,20 +145,32 @@ Table hold your actual data as a collection of records / rows.
::: lancedb.table.AsyncTags
::: lancedb.table.AsyncBranches
## Indices (Asynchronous)
Indices can be created on a table to speed up queries. This section
lists the indices that LanceDb supports.
::: lancedb.index
options:
show_root_heading: false
show_root_toc_entry: false
# `lang_mapping` is defined in the module rather than imported, so it is
# picked up despite not being in `__all__`. It is an internal lookup table.
filters: ["!^_", "!^lang_mapping$"]
::: lancedb.index.BTree
::: lancedb.index.Bitmap
::: lancedb.index.LabelList
::: lancedb.index.FTS
::: lancedb.index.IvfPq
::: lancedb.index.HnswPq
::: lancedb.index.HnswSq
::: lancedb.index.IvfFlat
::: lancedb.index.IvfSq
::: lancedb.index.IvfRq
::: lancedb.index.HnswFlat
::: lancedb.table.IndexStatistics
@@ -283,7 +198,3 @@ rows nearest to a query vector and can be created with the
::: lancedb.query.AsyncHybridQuery
options:
inherited_members: true
::: lancedb.query.AsyncTakeQuery
options:
inherited_members: true
+1 -1
View File
@@ -8,7 +8,7 @@
<parent>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.37.1-final.0</version>
<version>0.37.1-beta.0</version>
<relativePath>../pom.xml</relativePath>
</parent>
+2 -2
View File
@@ -6,7 +6,7 @@
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.37.1-final.0</version>
<version>0.37.1-beta.0</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>10.0.0</lance-core.version>
<lance-core.version>10.0.0-beta.5</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
View File
@@ -1,7 +1,7 @@
# Contributing to LanceDB Typescript
This document outlines the process for contributing to LanceDB Typescript.
For general contribution guidelines, see [CONTRIBUTING.md](https://github.com/lancedb/lancedb/blob/main/CONTRIBUTING.md).
For general contribution guidelines, see [CONTRIBUTING.md](../CONTRIBUTING.md).
## Project layout
+1 -1
View File
@@ -1,7 +1,7 @@
[package]
name = "lancedb-nodejs"
edition.workspace = true
version = "0.37.1"
version = "0.37.1-beta.0"
publish = false
license.workspace = true
description.workspace = true
-144
View File
@@ -6,9 +6,7 @@ import * as arrow17 from "apache-arrow-17";
import * as arrow18 from "apache-arrow-18";
import {
Vector as CurrentVector,
convertToTable,
tableFromIPC as currentTableFromIPC,
fromBufferToRecordBatch,
fromDataToBuffer,
fromRecordBatchToBuffer,
@@ -21,7 +19,6 @@ import {
FunctionOptions,
} from "../lancedb/embedding/embedding_function";
import { EmbeddingFunctionConfig } from "../lancedb/embedding/registry";
import { sanitizeTable } from "../lancedb/sanitize";
// biome-ignore lint/suspicious/noExplicitAny: skip
function sampleRecords(): Array<Record<string, any>> {
@@ -67,11 +64,7 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
tableFromIPC,
DataType,
Dictionary,
RecordBatch: ArrowRecordBatch,
Table: ArrowTable,
Uint8: ArrowUint8,
makeData: arrowMakeData,
vectorFromArray,
// biome-ignore lint/suspicious/noExplicitAny: <explanation>
} = <any>arrow;
type Schema = ApacheArrow["Schema"];
@@ -204,35 +197,6 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
expect(table.getChild("d")?.toJSON()).toEqual([9n, 10n, null]);
});
it("will use a provided FixedSizeList schema with typed array values", function () {
const schema = new Schema([
new Field("text", new Utf8(), false),
new Field(
"vector",
new FixedSizeList(3, new Field("item", new Float32(), false)),
false,
),
]);
const table = makeArrowTable(
[
{
text: "foo",
vector: new Float32Array([1, 2, 3]),
},
],
{ schema },
);
expect(table.getChild("text")?.toJSON()).toEqual(["foo"]);
expect(
table
.getChild("vector")
?.toJSON()
.map((value) => value.toJSON()),
).toEqual([[1, 2, 3]]);
});
it("will assume the column `vector` is FixedSizeList<Float32> by default", async function () {
const schema = new Schema([
new Field("a", new Float(Precision.DOUBLE), true),
@@ -1061,114 +1025,6 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
});
describe("when using two versions of arrow", function () {
it("preserves a dictionary shared by multiple fields", async function () {
const values = ["alpha", "beta", "alpha"];
const dictionaryVector = vectorFromArray(values);
const batch = new ArrowRecordBatch({
first: dictionaryVector.data[0],
second: dictionaryVector.data[0],
});
const table = new ArrowTable([batch]);
const sanitized = sanitizeTable(table);
expect([...sanitized.getChild("first")!]).toEqual(values);
expect([...sanitized.getChild("second")!]).toEqual(values);
const firstType = sanitized.schema.fields[0].type as {
dictionary: unknown;
};
const secondType = sanitized.schema.fields[1].type as {
dictionary: unknown;
};
expect(secondType.dictionary).toBe(firstType.dictionary);
expect(sanitized.batches[0].data.children[1].dictionary).toBe(
sanitized.batches[0].data.children[0].dictionary,
);
const buf = await fromDataToBuffer(table);
const actual = currentTableFromIPC(buf);
expect([...actual.getChild("first")!]).toEqual(values);
expect([...actual.getChild("second")!]).toEqual(values);
});
it("preserves shared dictionary data from another Arrow version", async function () {
const values = ["alpha", "beta", "alpha"];
const dictionaryVector = vectorFromArray(values);
const firstBatch = new ArrowRecordBatch({
label: dictionaryVector.slice(0, 2).data[0],
});
const secondBatch = new ArrowRecordBatch({
label: dictionaryVector.slice(2).data[0],
});
const table = new ArrowTable([firstBatch, secondBatch]);
const sanitized = sanitizeTable(table);
expect([...sanitized.getChild("label")!]).toEqual(values);
const dictionaries = sanitized.batches.map(
(batch) => batch.data.children[0].dictionary,
);
expect(dictionaries[0]).toBeInstanceOf(CurrentVector);
expect(dictionaries[1]).toBe(dictionaries[0]);
const buf = await fromDataToBuffer(table);
const actual = currentTableFromIPC(buf);
expect([...actual.getChild("label")!]).toEqual(values);
});
it("preserves shared chunks in growing dictionaries", async function () {
const type = new Dictionary(new Utf8(), new Int32(), 42, false);
const firstDictionary = vectorFromArray(["alpha", "beta"], new Utf8());
const secondDictionary = firstDictionary.concat(
vectorFromArray(["gamma"], new Utf8()),
);
const firstData = arrowMakeData({
type,
data: Int32Array.from([0, 1]),
dictionary: firstDictionary,
});
const secondData = arrowMakeData({
type,
data: Int32Array.from([2]),
dictionary: secondDictionary,
});
const table = new ArrowTable([
new ArrowRecordBatch({ label: firstData }),
new ArrowRecordBatch({ label: secondData }),
]);
const sanitized = sanitizeTable(table);
const expected = ["alpha", "beta", "gamma"];
expect([...sanitized.getChild("label")!]).toEqual(expected);
const firstLocalDictionary =
sanitized.batches[0].data.children[0].dictionary!;
const secondLocalDictionary =
sanitized.batches[1].data.children[0].dictionary!;
expect(secondLocalDictionary.data[0]).toBe(
firstLocalDictionary.data[0],
);
const buf = await fromTableToBuffer(sanitized);
const actual = currentTableFromIPC(buf);
expect([...actual.getChild("label")!]).toEqual(expected);
});
it("can serialize list data from another Arrow version", async function () {
const values = [["anime", "action"], [], null];
const vector = vectorFromArray(
values,
new List(new Field("item", new Utf8(), true)),
);
const table = new ArrowTable({ tags: vector });
const buf = await fromDataToBuffer(table);
const actual = currentTableFromIPC(buf);
const actualTags = actual.getChild("tags");
expect(actualTags?.get(0)?.toJSON()).toEqual(values[0]);
expect(actualTags?.get(1)?.toJSON()).toEqual(values[1]);
expect(actualTags?.get(2)).toBeNull();
});
it("can still import data", async function () {
const schema = new arrow15.Schema([
new arrow15.Field("id", new arrow15.Int32()),
-60
View File
@@ -11,11 +11,8 @@ import {
Float16,
Float32,
Float64,
Int32,
Schema,
Utf8,
fromDataToBuffer,
tableFromIPC,
} from "../lancedb/arrow";
import { EmbeddingFunction, LanceSchema } from "../lancedb/embedding";
import { getRegistry, register } from "../lancedb/embedding/registry";
@@ -187,63 +184,6 @@ describe("embedding functions", () => {
const vector0 = JSON.parse(JSON.stringify(arr[0].vector));
expect(vector0).toEqual([1, 2, 3]);
});
it("should append generated vectors to a non-nullable schema", async () => {
@register("non_nullable_schema_test")
class MockEmbeddingFunction extends EmbeddingFunction<string> {
ndims() {
return 3;
}
embeddingDataType(): Float {
return new Float64();
}
async computeSourceEmbeddings(data: string[]) {
return data.map(() => [1, 2, 3]);
}
}
const schema = new Schema([
new Field("id", new Int32()),
new Field("text", new Utf8()),
new Field("type", new Utf8()),
new Field(
"vector",
new FixedSizeList(3, new Field("item", new Float64())),
),
]);
const func = new MockEmbeddingFunction();
const db = await connect(tmpDir.name);
const table = await db.createEmptyTable("test_non_nullable", schema, {
embeddingFunction: {
function: func,
sourceColumn: "text",
},
});
const data = [
{ id: 1, text: "Carrot", type: "vegetable" },
{ id: 2, text: "Apple", type: "fruit" },
];
const buffer = await fromDataToBuffer(
data,
undefined,
await table.schema(),
);
const generatedTable = tableFromIPC(buffer);
const vectorField = generatedTable.schema.fields.find(
(field) => field.name === "vector",
);
expect(vectorField?.nullable).toBe(false);
await table.add(data);
const rows = await table.query().toArray();
expect(rows).toHaveLength(2);
for (const row of rows) {
expect([...row.vector]).toEqual([1, 2, 3]);
}
});
it("should error when appending to a table with an unregistered embedding function", async () => {
@register("mock")
class MockEmbeddingFunction extends EmbeddingFunction<string> {
-14
View File
@@ -1,14 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
import packageJson = require("../package.json");
describe("package metadata", () => {
it("requires Node.js type declarations compatible with the runtime", () => {
expect(packageJson.engines.node).toBe(">= 18");
expect(packageJson.peerDependencies["@types/node"]).toBe(">=18");
expect(packageJson.peerDependenciesMeta["@types/node"]).toEqual({
optional: true,
});
});
});
-75
View File
@@ -110,81 +110,6 @@ describe("Query outputSchema", () => {
});
});
describe("Search pagination", () => {
let tmpDir: tmp.DirResult;
let table: Table;
beforeEach(async () => {
tmpDir = tmp.dirSync({ unsafeCleanup: true });
const db = await connect(tmpDir.name);
const schema = new Schema([
new Field("id", new Int64(), false),
new Field("text", new Utf8(), false),
new Field(
"vector",
new FixedSizeList(2, new Field("item", new Float32())),
false,
),
]);
const data = makeArrowTable(
[
{ id: 1n, text: "common", vector: [0, 0] },
{ id: 2n, text: "common common", vector: [1, 1] },
{ id: 3n, text: "common common common", vector: [2, 2] },
{ id: 4n, text: "common common common common", vector: [3, 3] },
],
{ schema },
);
table = await db.createTable("test", data);
});
afterEach(() => {
tmpDir.removeCallback();
});
it("applies offset after the vector search limit", async () => {
const allResults = await table
.vectorSearch([0, 0])
.select(["id"])
.limit(4)
.toArray();
const secondPage = await table
.vectorSearch([0, 0])
.select(["id"])
.limit(2)
.offset(2)
.toArray();
expect(allResults).toHaveLength(4);
expect(secondPage).toHaveLength(2);
expect(secondPage.map((row) => row.id)).toEqual(
allResults.slice(2, 4).map((row) => row.id),
);
});
it("applies offset after the full-text search limit", async () => {
await table.createIndex("text", { config: Index.fts() });
const allResults = await table
.search("common", "fts")
.select(["id"])
.limit(4)
.toArray();
const secondPage = await table
.search("common", "fts")
.select(["id"])
.limit(2)
.offset(2)
.toArray();
expect(allResults).toHaveLength(4);
expect(secondPage).toHaveLength(2);
expect(secondPage.map((row) => row.id)).toEqual(
allResults.slice(2, 4).map((row) => row.id),
);
});
});
describe("Query orderBy", () => {
let tmpDir: tmp.DirResult;
let table: Table;
+2 -132
View File
@@ -170,38 +170,6 @@ describe("remote connection", () => {
);
});
it("surfaces JSON server errors from remote table operations", async () => {
await withMockDatabase(
(req, res) => {
const path = req.url ?? "";
if (path.endsWith("/describe/")) {
res.writeHead(200, { "Content-Type": "application/json" }).end(
JSON.stringify({
name: "broken_table",
version: 1,
schema: { fields: [] },
}),
);
return;
}
if (path.endsWith("/count_rows/")) {
res
.writeHead(400, { "Content-Type": "application/json" })
.end(JSON.stringify({ error: "count rows failed" }));
return;
}
res.writeHead(404).end();
},
async (db) => {
const table = await db.openTable("broken_table");
await expect(table.countRows()).rejects.toThrow("count rows failed");
},
);
});
it("should pass on requested extra headers", async () => {
await withMockDatabase(
(req, res) => {
@@ -258,7 +226,7 @@ describe("remote connection", () => {
);
});
it("sends FTS options to remote tables", async () => {
it("sends the FTS posting block size to remote tables", async () => {
let createIndexBody: Record<string, unknown> | undefined;
await withMockDatabase(
@@ -296,11 +264,7 @@ describe("remote connection", () => {
async (db) => {
const table = await db.openTable("t");
await table.createIndex("text", {
config: Index.fts({
blockSize: 256,
removeStopWords: true,
customStopWords: ["the"],
}),
config: Index.fts({ blockSize: 256 }),
});
},
);
@@ -308,7 +272,6 @@ describe("remote connection", () => {
expect(createIndexBody?.["column"]).toBe("text");
expect(createIndexBody?.["index_type"]).toBe("FTS");
expect(createIndexBody?.["block_size"]).toBe(256);
expect(createIndexBody?.["custom_stop_words"]).toEqual(["the"]);
});
it("diffs and merges remote branches", async () => {
@@ -909,96 +872,3 @@ describe("remote connection", () => {
});
});
});
describe("remote connection jobs surface", () => {
it("lists, describes, cancels, and reads history", async () => {
const { tableFromArrays, tableToIPC } = await import("apache-arrow");
const eventsTable = tableFromArrays({ state: ["created", "succeeded"] });
const eventsBody = Buffer.from(tableToIPC(eventsTable, "stream"));
await withMockDatabase(
(req, res) => {
let body = "";
req.on("data", (chunk) => {
body += chunk;
});
req.on("end", () => {
const payload = body.length > 0 ? JSON.parse(body) : {};
if (req.url === "/v1/jobs/list") {
if (payload["page_token"] === undefined) {
res
.writeHead(200, { "Content-Type": "application/json" })
.end(
'{"jobs": [{"job_id": "job-1", "table": "t1", ' +
'"job_type": "create_index", "state": "in_progress", ' +
'"created_at_millis": 1000}], "page_token": "next"}',
);
} else {
res
.writeHead(200, { "Content-Type": "application/json" })
.end(
'{"jobs": [{"job_id": "job-2", "table": "t2", ' +
'"job_type": "create_index", "state": "succeeded", ' +
'"created_at_millis": 2000}]}',
);
}
} else if (req.url === "/v1/jobs/describe") {
if (payload["job_id"] !== "job-1") {
res.writeHead(404).end("no such job");
return;
}
res
.writeHead(200, { "Content-Type": "application/json" })
.end(
'{"job_id": "job-1", "job_type": "create_index", ' +
'"job_state": "FAILED", "creation_ms": 1000, ' +
'"spec": {"column": "vec"}, "failure": {"phase": "execute", ' +
'"message": "worker died", "retryable": true}}',
);
} else if (req.url === "/v1/jobs/cancel") {
if (payload["job_id"] !== "job-1") {
res.writeHead(404).end("no such job");
return;
}
res
.writeHead(200, { "Content-Type": "application/json" })
.end('{"job_id": "job-1"}');
} else if (req.url === "/v1/jobs/query_events") {
res
.writeHead(200, {
"Content-Type": "application/vnd.apache.arrow.stream",
})
.end(eventsBody);
} else {
res.writeHead(404).end();
}
});
},
async (db) => {
const jobs = await db.listJobs();
expect(jobs.map((job) => job.jobId)).toEqual(["job-1", "job-2"]);
expect(jobs[0].state).toEqual("running");
expect(jobs[1].state).toEqual("finished");
const description = await db.getJob("job-1");
expect(description?.state).toEqual("failed");
expect(JSON.parse(description?.specJson ?? "")).toEqual({
column: "vec",
});
expect(description?.failure?.message).toEqual("worker died");
expect(await db.getJob("missing")).toBeNull();
expect(await db.cancelJob("job-1")).toBe(true);
expect(await db.cancelJob("missing")).toBe(false);
const history = await db.jobHistory("job-1");
expect(history.numRows).toEqual(2);
const job = db.job("job-1");
expect(job.id).toEqual("job-1");
expect(await job.status()).toEqual("failed");
await expect(job.wait()).rejects.toThrow("worker died");
},
);
});
});
+2 -61
View File
@@ -86,44 +86,6 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
await expect(table.countRows()).resolves.toBe(3);
});
it("should support a foreign Float64 vector schema end to end", async () => {
const conn = await connect(tmpDir.name);
const schema = new arrow.Schema([
new arrow.Field("resource_id", new arrow.Int32(), false),
new arrow.Field(
"vector",
new arrow.FixedSizeList(
3,
new arrow.Field("value", new arrow.Float64(), true),
),
false,
),
]);
const data = [
{
// biome-ignore lint/style/useNamingConvention: matches the reported schema
resource_id: 0,
vector: [0.1, 0.1, 0.1],
},
];
const resources = await conn.createTable("resources", data, { schema });
const existing = await resources
.query()
.where("resource_id = 0")
.limit(1)
.toArray();
expect(existing).toHaveLength(1);
const matched = await resources
.search(Float64Array.from(data[0].vector))
.limit(1)
.toArray();
expect(matched).toHaveLength(1);
expect(matched[0]["resource_id"]).toBe(0);
});
it("should support branches", async () => {
await table.add([{ id: 1 }]);
expect(await table.countRows()).toBe(1);
@@ -277,16 +239,8 @@ 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: 44,
});
// Index files count toward totalBytes too (only deletion files and
// manifests are excluded).
await table.createIndex("id", { config: Index.btree() });
const statsWithIndex = await table.stats();
expect(statsWithIndex.numIndices).toBe(1);
expect(statsWithIndex.totalBytes).toBeGreaterThan(684);
});
it("should overwrite data if asked", async () => {
@@ -897,11 +851,7 @@ describe("When creating an index", () => {
afterEach(() => tmpDir.removeCallback());
it("should create a vector index on vector columns", async () => {
const job = await tbl.createIndexAsync("vec");
expect(job.id).toBeNull();
await job.wait();
// Cancelling a job that already finished succeeds and does nothing.
await job.cancel();
await tbl.createIndex("vec");
// check index directory
const indexDir = path.join(tmpDir.name, "test.lance", "_indices");
@@ -2819,15 +2769,6 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
},
);
test("tokenize supports custom stop words", async () => {
const tokens = await tokenize("the lance data", {
stem: false,
removeStopWords: true,
customStopWords: ["lance"],
});
expect(tokens.map((token) => token.text)).toEqual(["the", "data"]);
});
describe("when calling explainPlan", () => {
let tmpDir: tmp.DirResult;
let table: Table;
+1 -7
View File
@@ -29,14 +29,8 @@ test("full text search", async () => {
const tbl = await db.createTable("myVectors", data, { mode: "overwrite" });
await tbl.createIndex("doc", {
config: lancedb.Index.fts({
stem: false,
removeStopWords: true,
customStopWords: ["banana"],
}),
config: lancedb.Index.fts(),
});
const tokens = await tbl.tokenize("apple banana", { column: "doc" });
expect(tokens.map((token) => token.text)).toEqual(["apple"]);
// --8<-- [start:full_text_search]
const result = await tbl
-62
View File
@@ -1,7 +1,6 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
import { tableFromIPC } from "apache-arrow";
import {
Data,
SchemaLike,
@@ -21,9 +20,6 @@ import type {
CreateNamespaceResponse,
DescribeNamespaceResponse,
DropNamespaceResponse,
Job,
JobDescription,
JobInfo,
ListNamespacesResponse,
} from "./native";
export type {
@@ -440,40 +436,6 @@ export abstract class Connection {
newName: string,
options?: RenameTableOptions,
): Promise<void>;
/**
* A {@link Job} handle for a server-side job by id.
*
* The handle is constructed without a server round trip; an unknown id
* surfaces when the handle is used. Dropping the handle has no effect on
* the job itself.
*/
abstract job(jobId: string): Job;
/** List server-side jobs across the database's tables. */
abstract listJobs(): Promise<JobInfo[]>;
/**
* Describe a single server-side job by id.
*
* Resolves to `null` when the server has no such job.
*/
abstract getJob(jobId: string): Promise<JobDescription | null>;
/**
* Request cancellation of a server-side job by id.
*
* Resolves to true if the server accepted the cancellation, false if no
* such job exists. Cancelling an already-terminal job is a no-op success.
*/
abstract cancelJob(jobId: string): Promise<boolean>;
/**
* The lifecycle event history of a server-side job, as an Arrow table.
*
* Lists history across all jobs when `jobId` is omitted.
*/
abstract jobHistory(jobId?: string): Promise<ArrowTable>;
}
/** @hideconstructor */
@@ -760,30 +722,6 @@ export class LocalConnection extends Connection {
options?.newNamespacePath,
);
}
job(jobId: string): Job {
return this.inner.job(jobId);
}
async listJobs(): Promise<JobInfo[]> {
return this.inner.listJobs();
}
async getJob(jobId: string): Promise<JobDescription | null> {
return this.inner.getJob(jobId);
}
async cancelJob(jobId: string): Promise<boolean> {
return this.inner.cancelJob(jobId);
}
async jobHistory(jobId?: string): Promise<ArrowTable> {
const buf = await this.inner.jobHistory(jobId);
if (buf.length === 0) {
return new ArrowTable();
}
return tableFromIPC(buf);
}
}
/**
+1 -18
View File
@@ -85,13 +85,7 @@ export {
RenameTableOptions,
} from "./connection";
export {
Job,
JobDescription,
JobFailureInfo,
JobInfo,
Session,
} from "./native.js";
export { Session } from "./native.js";
export {
ExecutableQuery,
@@ -200,16 +194,6 @@ export interface TokenizeOptions {
/** Whether to remove stop words. */
removeStopWords?: boolean;
/**
* Custom stop words that replace the built-in list for `language`.
*
* This option only affects tokenization when `removeStopWords` is true.
*
* `undefined` keeps the built-in language list. An empty array explicitly
* replaces it with no stop words.
*/
customStopWords?: string[];
/** Whether to fold ASCII characters. */
asciiFolding?: boolean;
@@ -241,7 +225,6 @@ export async function tokenize(
options?.lowercase,
options?.stem,
options?.removeStopWords,
options?.customStopWords,
options?.asciiFolding,
options?.ngramMinLength,
options?.ngramMaxLength,
-11
View File
@@ -553,16 +553,6 @@ export interface FtsOptions {
*/
removeStopWords?: boolean;
/**
* Custom stop words that replace the built-in list for `language`.
*
* This option only affects tokenization when `removeStopWords` is true.
*
* `undefined` keeps the built-in language list. An empty array explicitly
* replaces it with no stop words.
*/
customStopWords?: string[];
/**
* whether to remove punctuation
*/
@@ -765,7 +755,6 @@ export class Index {
options?.lowercase,
options?.stem,
options?.removeStopWords,
options?.customStopWords,
options?.asciiFolding,
options?.ngramMinLength,
options?.ngramMaxLength,
+29 -174
View File
@@ -9,7 +9,7 @@
// comes from the exact same library instance. This is not always the case
// and so we must sanitize the input to ensure that it is compatible.
import { BufferType, Data, Vector } from "apache-arrow";
import { BufferType, Data } from "apache-arrow";
import type { IntBitWidth, TKeys, TimeBitWidth } from "apache-arrow/type";
import {
Binary,
@@ -74,20 +74,6 @@ import {
Utf8,
} from "./arrow";
type SanitizationContext = {
types: WeakMap<object, DataType>;
vectors: WeakMap<object, Vector>;
data: WeakMap<object, Data<DataType>>;
};
function createSanitizationContext(): SanitizationContext {
return {
types: new WeakMap(),
vectors: new WeakMap(),
data: new WeakMap(),
};
}
export function sanitizeMetadata(
metadataLike?: unknown,
): Map<string, string> | undefined {
@@ -200,13 +186,6 @@ export function sanitizeInterval(typeLike: object) {
}
export function sanitizeList(typeLike: object) {
return sanitizeListWithContext(typeLike, createSanitizationContext());
}
function sanitizeListWithContext(
typeLike: object,
context: SanitizationContext,
) {
if (!("children" in typeLike) || !Array.isArray(typeLike.children)) {
throw Error(
"Expected a List type to have an array-like `children` property",
@@ -215,35 +194,19 @@ function sanitizeListWithContext(
if (typeLike.children.length !== 1) {
throw Error("Expected a List type to have exactly one child");
}
return new List(sanitizeFieldWithContext(typeLike.children[0], context));
return new List(sanitizeField(typeLike.children[0]));
}
export function sanitizeStruct(typeLike: object) {
return sanitizeStructWithContext(typeLike, createSanitizationContext());
}
function sanitizeStructWithContext(
typeLike: object,
context: SanitizationContext,
) {
if (!("children" in typeLike) || !Array.isArray(typeLike.children)) {
throw Error(
"Expected a Struct type to have an array-like `children` property",
);
}
return new Struct(
typeLike.children.map((child) => sanitizeFieldWithContext(child, context)),
);
return new Struct(typeLike.children.map((child) => sanitizeField(child)));
}
export function sanitizeUnion(typeLike: object) {
return sanitizeUnionWithContext(typeLike, createSanitizationContext());
}
function sanitizeUnionWithContext(
typeLike: object,
context: SanitizationContext,
) {
if (
!("typeIds" in typeLike) ||
!("mode" in typeLike) ||
@@ -263,7 +226,7 @@ function sanitizeUnionWithContext(
typeLike.mode,
// biome-ignore lint/suspicious/noExplicitAny: skip
typeLike.typeIds as any,
typeLike.children.map((child) => sanitizeFieldWithContext(child, context)),
typeLike.children.map((child) => sanitizeField(child)),
);
}
@@ -271,19 +234,6 @@ export function sanitizeTypedUnion(
typeLike: object,
// eslint-disable-next-line @typescript-eslint/naming-convention
UnionType: typeof DenseUnion | typeof SparseUnion,
) {
return sanitizeTypedUnionWithContext(
typeLike,
UnionType,
createSanitizationContext(),
);
}
function sanitizeTypedUnionWithContext(
typeLike: object,
// eslint-disable-next-line @typescript-eslint/naming-convention
UnionType: typeof DenseUnion | typeof SparseUnion,
context: SanitizationContext,
) {
if (!("typeIds" in typeLike)) {
throw Error(
@@ -298,7 +248,7 @@ function sanitizeTypedUnionWithContext(
return new UnionType(
typeLike.typeIds as Int32Array | number[],
typeLike.children.map((child) => sanitizeFieldWithContext(child, context)),
typeLike.children.map((child) => sanitizeField(child)),
);
}
@@ -312,16 +262,6 @@ export function sanitizeFixedSizeBinary(typeLike: object) {
}
export function sanitizeFixedSizeList(typeLike: object) {
return sanitizeFixedSizeListWithContext(
typeLike,
createSanitizationContext(),
);
}
function sanitizeFixedSizeListWithContext(
typeLike: object,
context: SanitizationContext,
) {
if (!("listSize" in typeLike) || typeof typeLike.listSize !== "number") {
throw Error("Expected a FixedSizeList type to have a `listSize` property");
}
@@ -335,18 +275,11 @@ function sanitizeFixedSizeListWithContext(
}
return new FixedSizeList(
typeLike.listSize,
sanitizeFieldWithContext(typeLike.children[0], context),
sanitizeField(typeLike.children[0]),
);
}
export function sanitizeMap(typeLike: object) {
return sanitizeMapWithContext(typeLike, createSanitizationContext());
}
function sanitizeMapWithContext(
typeLike: object,
context: SanitizationContext,
) {
if (!("children" in typeLike) || !Array.isArray(typeLike.children)) {
throw Error(
"Expected a Map type to have an array-like `children` property",
@@ -359,10 +292,7 @@ function sanitizeMapWithContext(
throw Error("Expected a Map type to have exactly one child");
}
return new Map_(
sanitizeFieldWithContext(typeLike.children[0], context),
typeLike.keysSorted,
);
return new Map_(sanitizeField(typeLike.children[0]), typeLike.keysSorted);
}
export function sanitizeDuration(typeLike: object) {
@@ -373,13 +303,6 @@ export function sanitizeDuration(typeLike: object) {
}
export function sanitizeDictionary(typeLike: object) {
return sanitizeDictionaryWithContext(typeLike, createSanitizationContext());
}
function sanitizeDictionaryWithContext(
typeLike: object,
context: SanitizationContext,
) {
if (!("id" in typeLike) || typeof typeLike.id !== "number") {
throw Error("Expected a Dictionary type to have an `id` property");
}
@@ -393,8 +316,8 @@ function sanitizeDictionaryWithContext(
throw Error("Expected a Dictionary type to have an `isOrdered` property");
}
return new Dictionary(
sanitizeTypeWithContext(typeLike.dictionary, context),
sanitizeTypeWithContext(typeLike.indices, context) as TKeys,
sanitizeType(typeLike.dictionary),
sanitizeType(typeLike.indices) as TKeys,
typeLike.id,
typeLike.isOrdered,
);
@@ -402,23 +325,12 @@ function sanitizeDictionaryWithContext(
// biome-ignore lint/suspicious/noExplicitAny: skip
export function sanitizeType(typeLike: unknown): DataType<any> {
return sanitizeTypeWithContext(typeLike, createSanitizationContext());
}
function sanitizeTypeWithContext(
typeLike: unknown,
context: SanitizationContext,
): DataType {
if (typeof typeLike === "string") {
return dataTypeFromName(typeLike);
}
if (typeof typeLike !== "object" || typeLike === null) {
throw Error("Expected a Type but object was null/undefined");
}
const cached = context.types.get(typeLike);
if (cached !== undefined) {
return cached;
}
if (
!("typeId" in typeLike) ||
!(
@@ -437,16 +349,6 @@ function sanitizeTypeWithContext(
throw Error("Type's typeId property was not a function or number");
}
const type = sanitizeTypeById(typeLike, typeId, context);
context.types.set(typeLike, type);
return type;
}
function sanitizeTypeById(
typeLike: object,
typeId: Type,
context: SanitizationContext,
): DataType {
switch (typeId) {
case Type.NONE:
throw Error("Received a Type with a typeId of NONE");
@@ -473,21 +375,21 @@ function sanitizeTypeById(
case Type.Interval:
return sanitizeInterval(typeLike);
case Type.List:
return sanitizeListWithContext(typeLike, context);
return sanitizeList(typeLike);
case Type.Struct:
return sanitizeStructWithContext(typeLike, context);
return sanitizeStruct(typeLike);
case Type.Union:
return sanitizeUnionWithContext(typeLike, context);
return sanitizeUnion(typeLike);
case Type.FixedSizeBinary:
return sanitizeFixedSizeBinary(typeLike);
case Type.FixedSizeList:
return sanitizeFixedSizeListWithContext(typeLike, context);
return sanitizeFixedSizeList(typeLike);
case Type.Map:
return sanitizeMapWithContext(typeLike, context);
return sanitizeMap(typeLike);
case Type.Duration:
return sanitizeDuration(typeLike);
case Type.Dictionary:
return sanitizeDictionaryWithContext(typeLike, context);
return sanitizeDictionary(typeLike);
case Type.Int8:
return new Int8();
case Type.Int16:
@@ -531,9 +433,9 @@ function sanitizeTypeById(
case Type.TimestampSecond:
return sanitizeTypedTimestamp(typeLike, TimestampSecond);
case Type.DenseUnion:
return sanitizeTypedUnionWithContext(typeLike, DenseUnion, context);
return sanitizeTypedUnion(typeLike, DenseUnion);
case Type.SparseUnion:
return sanitizeTypedUnionWithContext(typeLike, SparseUnion, context);
return sanitizeTypedUnion(typeLike, SparseUnion);
case Type.IntervalDayTime:
return new IntervalDayTime();
case Type.IntervalYearMonth:
@@ -552,13 +454,6 @@ function sanitizeTypeById(
}
export function sanitizeField(fieldLike: unknown): Field {
return sanitizeFieldWithContext(fieldLike, createSanitizationContext());
}
function sanitizeFieldWithContext(
fieldLike: unknown,
context: SanitizationContext,
): Field {
if (fieldLike instanceof Field) {
return fieldLike;
}
@@ -576,7 +471,7 @@ function sanitizeFieldWithContext(
}
let type: DataType;
try {
type = sanitizeTypeWithContext(fieldLike.type, context);
type = sanitizeType(fieldLike.type);
} catch (error: unknown) {
throw Error(
`Unable to sanitize type for field: ${fieldLike.name} due to error: ${error}`,
@@ -606,13 +501,6 @@ function sanitizeFieldWithContext(
* than lancedb is using.
*/
export function sanitizeSchema(schemaLike: SchemaLike): Schema {
return sanitizeSchemaWithContext(schemaLike, createSanitizationContext());
}
function sanitizeSchemaWithContext(
schemaLike: SchemaLike,
context: SanitizationContext,
): Schema {
if (schemaLike instanceof Schema) {
return schemaLike;
}
@@ -634,7 +522,7 @@ function sanitizeSchemaWithContext(
);
}
const sanitizedFields = schemaLike.fields.map((field) =>
sanitizeFieldWithContext(field, context),
sanitizeField(field),
);
return new Schema(sanitizedFields, metadata);
}
@@ -656,18 +544,13 @@ export function sanitizeTable(tableLike: TableLike): Table {
"The table passed in does not appear to be a table (no 'columns' property)",
);
}
const context = createSanitizationContext();
const schema = sanitizeSchemaWithContext(tableLike.schema, context);
const batches = tableLike.batches.map((batch) =>
sanitizeRecordBatch(batch, context),
);
const schema = sanitizeSchema(tableLike.schema);
const batches = tableLike.batches.map(sanitizeRecordBatch);
return new Table(schema, batches);
}
function sanitizeRecordBatch(
batchLike: RecordBatchLike,
context: SanitizationContext,
): RecordBatch {
function sanitizeRecordBatch(batchLike: RecordBatchLike): RecordBatch {
if (batchLike instanceof RecordBatch) {
return batchLike;
}
@@ -684,43 +567,19 @@ function sanitizeRecordBatch(
"The record batch passed in does not appear to be a record batch (no 'data' property)",
);
}
const schema = sanitizeSchemaWithContext(batchLike.schema, context);
const data = sanitizeData(batchLike.data, context) as Data<Struct>;
const schema = sanitizeSchema(batchLike.schema);
const data = sanitizeData(batchLike.data);
return new RecordBatch(schema, data);
}
type DictionaryVectorLike = {
data: readonly DataLike[];
};
type DictionaryDataLike = DataLike & {
dictionary?: DictionaryVectorLike;
};
function sanitizeData(
dataLike: DataLike,
context: SanitizationContext,
): Data<DataType> {
// biome-ignore lint/suspicious/noExplicitAny: <explanation>
): import("apache-arrow").Data<Struct<any>> {
if (dataLike instanceof Data) {
return dataLike;
}
const cachedData = context.data.get(dataLike);
if (cachedData !== undefined) {
return cachedData;
}
const dictionaryLike = (dataLike as DictionaryDataLike).dictionary;
let dictionary: Vector | undefined;
if (dictionaryLike !== undefined) {
dictionary = context.vectors.get(dictionaryLike);
if (dictionary === undefined) {
dictionary = new Vector(
dictionaryLike.data.map((data) => sanitizeData(data, context)),
);
context.vectors.set(dictionaryLike, dictionary);
}
}
const data = new Data(
sanitizeTypeWithContext(dataLike.type, context),
return new Data(
dataLike.type,
dataLike.offset,
dataLike.length,
dataLike.nullCount,
@@ -730,11 +589,7 @@ function sanitizeData(
[BufferType.VALIDITY]: dataLike.nullBitmap,
[BufferType.TYPE]: dataLike.typeIds,
},
dataLike.children.map((child) => sanitizeData(child, context)),
dictionary,
);
context.data.set(dataLike, data);
return data;
}
const constructorsByTypeName = {
-28
View File
@@ -30,7 +30,6 @@ import {
DropColumnsResult,
IndexConfig,
IndexStatistics,
Job,
Branches as NativeBranches,
OptimizeStats,
TableStatistics,
@@ -359,17 +358,6 @@ export abstract class Table {
options?: Partial<IndexOptions>,
): Promise<void>;
/**
* Create an index, returning a handle to the indexing job.
*
* The job may already be complete when returned; callers must not assume
* the index exists until {@link Job.wait} resolves.
*/
abstract createIndexAsync(
column: string,
options?: Partial<IndexOptions>,
): Promise<Job>;
/**
* Drop an index from the table.
*
@@ -952,22 +940,6 @@ export class LocalTable extends Table {
);
}
async createIndexAsync(
column: string,
options?: Partial<IndexOptions>,
): Promise<Job> {
// biome-ignore lint/suspicious/noExplicitAny: skip
const nativeIndex = (options?.config as any)?.inner;
return await this.inner.createIndexAsync(
nativeIndex,
column,
options?.replace,
options?.waitTimeoutSeconds,
options?.name,
options?.train,
);
}
async dropIndex(name: string): Promise<void> {
await this.inner.dropIndex(name);
}
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-darwin-arm64",
"version": "0.37.1",
"version": "0.37.1-beta.0",
"os": ["darwin"],
"cpu": ["arm64"],
"main": "lancedb.darwin-arm64.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-arm64-gnu",
"version": "0.37.1",
"version": "0.37.1-beta.0",
"os": ["linux"],
"cpu": ["arm64"],
"main": "lancedb.linux-arm64-gnu.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-arm64-musl",
"version": "0.37.1",
"version": "0.37.1-beta.0",
"os": ["linux"],
"cpu": ["arm64"],
"main": "lancedb.linux-arm64-musl.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-x64-gnu",
"version": "0.37.1",
"version": "0.37.1-beta.0",
"os": ["linux"],
"cpu": ["x64"],
"main": "lancedb.linux-x64-gnu.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-x64-musl",
"version": "0.37.1",
"version": "0.37.1-beta.0",
"os": ["linux"],
"cpu": ["x64"],
"main": "lancedb.linux-x64-musl.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-arm64-msvc",
"version": "0.37.1",
"version": "0.37.1-beta.0",
"os": [
"win32"
],
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-x64-msvc",
"version": "0.37.1",
"version": "0.37.1-beta.0",
"os": ["win32"],
"cpu": ["x64"],
"main": "lancedb.win32-x64-msvc.node",
+2 -8
View File
@@ -1,12 +1,12 @@
{
"name": "@lancedb/lancedb",
"version": "0.37.1",
"version": "0.37.0-beta.0",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "@lancedb/lancedb",
"version": "0.37.1",
"version": "0.37.0-beta.0",
"cpu": [
"x64",
"arm64"
@@ -55,13 +55,7 @@
"openai": "4.29.2"
},
"peerDependencies": {
"@types/node": ">=18",
"apache-arrow": ">=15.0.0 <=18.1.0"
},
"peerDependenciesMeta": {
"@types/node": {
"optional": true
}
}
},
"node_modules/@aws-crypto/crc32": {
+1 -7
View File
@@ -11,7 +11,7 @@
"ann"
],
"private": false,
"version": "0.37.1",
"version": "0.37.1-beta.0",
"main": "dist/index.js",
"exports": {
".": "./dist/index.js",
@@ -101,12 +101,6 @@
"openai": "4.29.2"
},
"peerDependencies": {
"@types/node": ">=18",
"apache-arrow": ">=15.0.0 <=18.1.0"
},
"peerDependenciesMeta": {
"@types/node": {
"optional": true
}
}
}
-63
View File
@@ -340,69 +340,6 @@ impl Connection {
self.get_inner()?.drop_all_tables(&ns).await.default_error()
}
/// A `Job` handle for a server-side job by id.
///
/// The handle is constructed without a server round trip; an unknown id
/// surfaces when the handle is used.
#[napi]
pub fn job(&self, job_id: String) -> napi::Result<crate::job::Job> {
let job = self.get_inner()?.job(job_id).default_error()?;
Ok(crate::job::Job::new(job))
}
/// List server-side jobs across the database's tables.
#[napi(catch_unwind)]
pub async fn list_jobs(&self) -> napi::Result<Vec<crate::job::JobInfo>> {
let jobs = self.get_inner()?.list_jobs().await.default_error()?;
Ok(jobs.into_iter().map(Into::into).collect())
}
/// Describe a single server-side job by id. `null` when the server has
/// no such job.
#[napi(catch_unwind)]
pub async fn get_job(
&self,
job_id: String,
) -> napi::Result<Option<crate::job::JobDescription>> {
let description = self.get_inner()?.get_job(&job_id).await.default_error()?;
Ok(description.map(Into::into))
}
/// Request cancellation of a server-side job by id. Returns true if the
/// server accepted the cancellation, false if no such job exists.
#[napi(catch_unwind)]
pub async fn cancel_job(&self, job_id: String) -> napi::Result<bool> {
self.get_inner()?.cancel_job(&job_id).await.default_error()
}
/// The lifecycle event history of a server-side job (all jobs when
/// `job_id` is null), as an Arrow IPC stream buffer. Empty when there is
/// no history.
#[napi(catch_unwind)]
pub async fn job_history(&self, job_id: Option<String>) -> napi::Result<Buffer> {
let batches = self
.get_inner()?
.job_history(job_id.as_deref())
.await
.default_error()?;
let Some(first) = batches.first() else {
return Ok(Buffer::from(Vec::<u8>::new()));
};
let mut out = Vec::new();
let mut writer = arrow_ipc::writer::StreamWriter::try_new(&mut out, &first.schema())
.map_err(|e| napi::Error::from_reason(e.to_string()))?;
for batch in &batches {
writer
.write(batch)
.map_err(|e| napi::Error::from_reason(e.to_string()))?;
}
writer
.finish()
.map_err(|e| napi::Error::from_reason(e.to_string()))?;
drop(writer);
Ok(Buffer::from(out))
}
#[napi(catch_unwind)]
/// Describe a namespace and return its properties.
pub async fn describe_namespace(
-4
View File
@@ -43,7 +43,6 @@ pub fn tokenize(
lower_case: Option<bool>,
stem: Option<bool>,
remove_stop_words: Option<bool>,
custom_stop_words: Option<Vec<String>>,
ascii_folding: Option<bool>,
ngram_min_length: Option<u32>,
ngram_max_length: Option<u32>,
@@ -73,7 +72,6 @@ pub fn tokenize(
if let Some(remove_stop_words) = remove_stop_words {
opts = opts.remove_stop_words(remove_stop_words);
}
opts = opts.custom_stop_words(custom_stop_words);
if let Some(ascii_folding) = ascii_folding {
opts = opts.ascii_folding(ascii_folding);
}
@@ -224,7 +222,6 @@ impl Index {
lower_case: Option<bool>,
stem: Option<bool>,
remove_stop_words: Option<bool>,
custom_stop_words: Option<Vec<String>>,
ascii_folding: Option<bool>,
ngram_min_length: Option<u32>,
ngram_max_length: Option<u32>,
@@ -253,7 +250,6 @@ impl Index {
if let Some(remove_stop_words) = remove_stop_words {
opts = opts.remove_stop_words(remove_stop_words);
}
opts = opts.custom_stop_words(custom_stop_words);
if let Some(ascii_folding) = ascii_folding {
opts = opts.ascii_folding(ascii_folding);
}
-123
View File
@@ -1,123 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
use std::sync::Arc;
use napi_derive::napi;
use crate::error::NapiErrorExt;
/// A handle to an operation that may still be running.
#[napi]
pub struct Job {
inner: Arc<lancedb::Job>,
}
impl Job {
pub(crate) fn new(inner: lancedb::Job) -> Self {
Self {
inner: Arc::new(inner),
}
}
}
#[napi]
impl Job {
/// Identifies the operation on the server that is running it. Operations
/// that run in this process have no server id. The value is opaque.
#[napi(getter)]
pub fn id(&self) -> Option<String> {
self.inner.id().map(str::to_string)
}
/// The operation's current lifecycle state: "running", "finished",
/// "failed", or "cancelled".
///
/// A point snapshot; unlike {@link Job.wait} it does not block or reject
/// on a terminal failure state. States a newer server reports that this
/// client version does not know pass through as-is.
#[napi(catch_unwind)]
pub async fn status(&self) -> napi::Result<String> {
self.inner.status().await.default_error()
}
/// Wait until the operation reaches a terminal state.
#[napi(catch_unwind)]
pub async fn wait(&self) -> napi::Result<()> {
self.inner.wait().await.default_error()
}
/// Request cancellation. Cancelling a finished operation is a no-op.
#[napi(catch_unwind)]
pub async fn cancel(&self) -> napi::Result<()> {
self.inner.cancel().await.default_error()
}
}
/// A row from `Connection.listJobs`: one server-side job.
#[napi(object)]
pub struct JobInfo {
/// The job id -- what `Connection.getJob` and `Connection.cancelJob`
/// accept.
pub job_id: String,
/// The table the job runs against, without URI or namespace.
pub table: String,
pub job_type: String,
/// Lifecycle state: "running", "finished", "failed", or "cancelled".
pub state: String,
/// When the job was created, in milliseconds since the epoch.
pub created_at_millis: i64,
}
impl From<lancedb::database::JobInfo> for JobInfo {
fn from(info: lancedb::database::JobInfo) -> Self {
Self {
job_id: info.job_id,
table: info.table,
job_type: info.job_type,
state: info.state,
created_at_millis: info.created_at_millis,
}
}
}
/// The server's account of why a job failed.
#[napi(object)]
pub struct JobFailureInfo {
pub phase: Option<String>,
pub message: Option<String>,
pub retryable: Option<bool>,
}
/// A described job from `Connection.getJob`.
#[napi(object)]
pub struct JobDescription {
pub job_id: String,
pub job_type: String,
/// Lifecycle state: "running", "finished", "failed", or "cancelled".
pub state: String,
/// When the job was created, in milliseconds since the epoch.
pub creation_ms: i64,
/// The job-type-specific specification as a JSON string, when present.
pub spec_json: Option<String>,
/// Why the job failed, when the job is failed and the server reports a
/// reason.
pub failure: Option<JobFailureInfo>,
}
impl From<lancedb::database::JobDescription> for JobDescription {
fn from(description: lancedb::database::JobDescription) -> Self {
Self {
job_id: description.job_id,
job_type: description.job_type,
state: description.state,
creation_ms: description.creation_ms,
spec_json: (!description.spec.is_null()).then(|| description.spec.to_string()),
failure: description.failure.map(|failure| JobFailureInfo {
phase: failure.phase,
message: failure.message,
retryable: failure.retryable,
}),
}
}
}
-1
View File
@@ -11,7 +11,6 @@ mod error;
mod header;
mod index;
mod iterator;
mod job;
pub mod merge;
pub mod otel;
pub mod permutation;
+3 -43
View File
@@ -168,39 +168,6 @@ impl Table {
builder.execute().await.default_error()
}
#[napi(catch_unwind)]
pub async fn create_index_async(
&self,
index: Option<&Index>,
column: String,
replace: Option<bool>,
wait_timeout_s: Option<i64>,
name: Option<String>,
train: Option<bool>,
) -> napi::Result<crate::job::Job> {
let lancedb_index = if let Some(index) = index {
index.consume()?
} else {
lancedb::index::Index::Auto
};
let mut builder = self.inner_ref()?.create_index(&[column], lancedb_index);
if let Some(replace) = replace {
builder = builder.replace(replace);
}
if let Some(timeout) = wait_timeout_s {
builder =
builder.wait_timeout(std::time::Duration::from_secs(timeout.try_into().unwrap()));
}
if let Some(name) = name {
builder = builder.name(name);
}
if let Some(train) = train {
builder = builder.train(train);
}
let job = builder.execute_async().await.default_error()?;
Ok(crate::job::Job::new(job))
}
#[napi(catch_unwind)]
pub async fn drop_index(&self, index_name: String) -> napi::Result<()> {
self.inner_ref()?
@@ -339,9 +306,7 @@ impl Table {
let transforms = NewColumnTransform::SqlExpressions(transforms);
let res = self
.inner_ref()?
.add_columns()
.transform(transforms)
.execute()
.add_columns(transforms, None)
.await
.default_error()?;
Ok(res.into())
@@ -358,9 +323,7 @@ impl Table {
let transforms = NewColumnTransform::AllNulls(schema);
let res = self
.inner_ref()?
.add_columns()
.transform(transforms)
.execute()
.add_columns(transforms, None)
.await
.default_error()?;
Ok(res.into())
@@ -1043,10 +1006,7 @@ impl From<lancedb::index::IndexStatistics> for IndexStatistics {
#[napi(object)]
pub struct TableStatistics {
/// The total size, in bytes, of the table's data files, index files, and
/// overlay files
///
/// Read from the manifest, so this excludes deletion files and manifests.
/// The total number of bytes in the table
pub total_bytes: i64,
/// The number of rows in the table
+3 -3
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb-python"
version = "0.37.1"
version = "0.37.1-beta.0"
publish = false
edition.workspace = true
description = "Python bindings for LanceDB"
@@ -26,7 +26,7 @@ lance-namespace-impls.workspace = true
lance-io.workspace = true
env_logger.workspace = true
log.workspace = true
pyo3 = { version = "0.28", features = ["extension-module", "abi3-py310", "chrono"] }
pyo3 = { version = "0.28", features = ["extension-module", "abi3-py39", "chrono"] }
chrono = { version = "0.4", default-features = false, features = ["clock"] }
pyo3-async-runtimes = { version = "0.28", features = [
"attributes",
@@ -43,7 +43,7 @@ libc = "0.2"
[build-dependencies]
pyo3-build-config = { version = "0.28", features = [
"extension-module",
"abi3-py310",
"abi3-py39",
] }
[features]
+2 -3
View File
@@ -60,10 +60,10 @@ tests = [
"pytest-asyncio>=0.21",
"duckdb>=0.9.0",
"pytz>=2023.3",
"polars>=0.19, <=1.32.3",
"polars>=0.19, <=1.3.0",
"pyarrow<25",
"pyarrow-stubs>=16.0",
"pylance==10.0.0",
"pylance==9.0.0rc1",
"requests>=2.31.0",
"datafusion>=54,<55",
"opentelemetry-sdk>=1.30.0",
@@ -140,7 +140,6 @@ include = [
"python/lancedb/remote/errors.py",
"python/lancedb/embeddings/__init__.py",
"python/lancedb/_lancedb.pyi",
"python/type_tests/connect.py",
]
exclude = ["python/tests/"]
pythonVersion = "3.13"
+2 -8
View File
@@ -20,7 +20,6 @@ from .remote import ClientConfig
from .remote.db import RemoteDBConnection
from .expr import Expr, col, lit, func
from .schema import blob, vector, BlobType
from .job import AsyncJob, Job
from .table import AsyncTable, Table
from .types import BaseTokenizerType
from ._lancedb import Session
@@ -259,7 +258,6 @@ def tokenize(
lower_case: bool = True,
stem: bool = True,
remove_stop_words: bool = True,
custom_stop_words: Optional[List[str]] = None,
ascii_folding: bool = True,
ngram_min_length: int = 3,
ngram_max_length: int = 3,
@@ -267,10 +265,9 @@ def tokenize(
) -> Iterable[FtsToken]:
"""Tokenize a full-text search query using an explicit tokenizer.
This does not require an FTS index. The tokenizer options match
:class:`lancedb.index.FTS`. ``custom_stop_words`` accepts a list of strings.
This does not require a table or FTS index. The tokenizer options match
:class:`lancedb.index.FTS`.
"""
return _tokenize(
query,
base_tokenizer=base_tokenizer,
@@ -279,7 +276,6 @@ def tokenize(
lower_case=lower_case,
stem=stem,
remove_stop_words=remove_stop_words,
custom_stop_words=custom_stop_words,
ascii_folding=ascii_folding,
ngram_min_length=ngram_min_length,
ngram_max_length=ngram_max_length,
@@ -501,7 +497,6 @@ __all__ = [
"connect_namespace",
"connect_namespace_async",
"AsyncConnection",
"AsyncJob",
"AsyncLanceNamespaceDBConnection",
"AsyncTable",
"FtsToken",
@@ -515,7 +510,6 @@ __all__ = [
"BlobType",
"vector",
"DBConnection",
"Job",
"LanceDBConnection",
"LanceNamespaceDBConnection",
"RemoteDBConnection",
+23 -6
View File
@@ -14,10 +14,14 @@ import pyarrow as pa
from .expr import Expr
from .schema import blob_v2_column_paths
from .types import BlobMode, QueryProjection, QueryProjectionSpec
from .util import get_uri_scheme
if TYPE_CHECKING:
from _typeshed import WriteableBuffer
from .remote.table import RemoteTable
from .table import AsyncTable, Table
BLOB_MODE_TO_HANDLING = {
"lazy": "blobs_descriptions",
"bytes": "all_binary",
@@ -100,6 +104,22 @@ def validate_blob_mode(blob_mode: BlobMode) -> None:
raise ValueError(f"blob_mode must be one of {modes}, got {blob_mode!r}")
def supports_blob_auto_row_id(table: Table | AsyncTable | RemoteTable) -> bool:
"""Blob auto row-id applies to native tables, not LanceDB Cloud."""
from .remote.table import RemoteTable
if isinstance(table, RemoteTable):
return False
inner = getattr(table, "_inner", None)
if inner is not None:
uri = inner.database().uri
if isinstance(uri, str) and get_uri_scheme(uri) == "db":
return False
return True
def projection_includes_blob_column(
projection: QueryProjection,
blob_columns: Iterable[str],
@@ -144,14 +164,16 @@ def v2_projection_needs_row_id(
def blob_auto_row_id_for_scan(
table: Table | AsyncTable | RemoteTable,
schema: pa.Schema,
projection: QueryProjection,
*,
with_row_id: bool | None,
) -> bool:
"""Auto row-id only applies when the caller said nothing about row ids."""
if with_row_id is not None:
return False
if not supports_blob_auto_row_id(table):
return False
return v2_projection_needs_row_id(schema, projection, with_row_id=False)
@@ -164,11 +186,6 @@ def finalize_blob_query_table(
) -> pa.Table:
if user_requested_row_id or not blob_auto_row_id:
return tbl
if "_rowid" not in tbl.column_names:
# A backend that ignores the row-id request leaves nothing to stash. Hand
# back the projection as-is so fetch_blobs raises the error that names the
# ways to supply row ids, rather than failing here about a hidden column.
return tbl
return stash_auto_row_ids(tbl, blob_paths)
-75
View File
@@ -59,7 +59,6 @@ def tokenize(
lower_case: bool = True,
stem: bool = True,
remove_stop_words: bool = True,
custom_stop_words: Optional[List[str]] = None,
ascii_folding: bool = True,
ngram_min_length: int = 3,
ngram_max_length: int = 3,
@@ -146,13 +145,6 @@ class Connection(object):
start_after: Optional[str],
limit: Optional[int],
) -> list[str]: ... # Deprecated: Use list_tables instead
def job(self, job_id: str) -> Job: ...
async def list_jobs(self) -> List[JobInfo]: ...
async def get_job(self, job_id: str) -> Optional[JobDescription]: ...
async def cancel_job(self, job_id: str) -> bool: ...
async def job_history(
self, job_id: Optional[str] = None
) -> List[pa.RecordBatch]: ...
async def create_table(
self,
name: str,
@@ -216,47 +208,6 @@ class BlobFile:
def read_range(self, offset: int, length: int) -> bytes: ...
def read_up_to(self, length: int) -> bytes: ...
class Job:
@property
def id(self) -> Optional[str]: ...
async def status(self) -> str: ...
async def wait(self) -> None: ...
async def cancel(self) -> None: ...
class JobInfo:
@property
def job_id(self) -> str: ...
@property
def table(self) -> str: ...
@property
def job_type(self) -> str: ...
@property
def state(self) -> str: ...
@property
def created_at_millis(self) -> int: ...
class JobFailureInfo:
@property
def phase(self) -> Optional[str]: ...
@property
def message(self) -> Optional[str]: ...
@property
def retryable(self) -> Optional[bool]: ...
class JobDescription:
@property
def job_id(self) -> str: ...
@property
def job_type(self) -> str: ...
@property
def state(self) -> str: ...
@property
def creation_ms(self) -> int: ...
@property
def spec_json(self) -> Optional[str]: ...
@property
def failure(self) -> Optional[JobFailureInfo]: ...
class Table:
def name(self) -> str: ...
def __repr__(self) -> str: ...
@@ -296,28 +247,6 @@ class Table:
name: Optional[str],
train: Optional[bool],
): ...
async def create_index_async(
self,
column: str,
index: Union[
IvfFlat,
IvfSq,
IvfPq,
HnswPq,
HnswSq,
HnswFlat,
BTree,
Bitmap,
LabelList,
Fm,
FTS,
],
replace: Optional[bool],
wait_timeout: Optional[object],
*,
name: Optional[str],
train: Optional[bool],
) -> Job: ...
async def list_versions(self) -> List[Dict[str, Any]]: ...
async def version(self) -> int: ...
async def checkout(self, version: Union[int, str]): ...
@@ -355,10 +284,6 @@ class Table:
async def set_lsm_write_spec(self, spec: LsmWriteSpec) -> None: ...
async def unset_lsm_write_spec(self) -> None: ...
async def get_lsm_write_spec(self) -> Optional[LsmWriteSpec]: ...
async def checkpoint_lsm(self) -> None: ...
async def flush_lsm(self) -> None: ...
async def compact_lsm(self) -> None: ...
async def get_lsm_stats(self, include_generation_rows: bool) -> Optional[dict]: ...
async def close_lsm_writers(self) -> None: ...
@property
def tags(self) -> Tags: ...
+10 -185
View File
@@ -45,7 +45,6 @@ from lance_namespace.errors import NamespaceNotEmptyError, TableNotFoundError
from . import __version__
from ._lancedb import connect as lancedb_connect # type: ignore
from .job import AsyncJob, Job
from .table import (
AsyncTable,
LanceTable,
@@ -64,7 +63,6 @@ if TYPE_CHECKING:
from .pydantic import LanceModel
from ._lancedb import Connection as LanceDbConnection
from ._lancedb import JobDescription, JobInfo
from .common import DATA, URI
from .embeddings import EmbeddingFunctionConfig
from ._lancedb import Session
@@ -180,51 +178,6 @@ class DBConnection(EnforceOverrides):
"Namespace operations are not supported for this connection type"
)
def namespace_exists(self, namespace_id: List[str]) -> bool:
"""Check if a namespace exists.
Parameters
----------
namespace_id: List[str]
The namespace identifier to check.
Returns
-------
bool
True if the namespace exists, False otherwise.
Raises
------
NotImplementedError
If the connection type does not support namespace operations.
"""
raise NotImplementedError(
"Namespace operations are not supported for this connection type"
)
def table_exists(self, table_id: List[str]) -> bool:
"""Check if a table exists.
Parameters
----------
table_id: List[str]
The table identifier to check (full path including namespace
segments and table name).
Returns
-------
bool
True if the table exists, False otherwise.
Raises
------
NotImplementedError
If the connection type does not support namespace operations.
"""
raise NotImplementedError(
"Namespace operations are not supported for this connection type"
)
def list_tables(
self,
namespace_path: Optional[List[str]] = None,
@@ -406,7 +359,7 @@ class DBConnection(EnforceOverrides):
Data is converted to Arrow before being written to disk. For maximum
control over how data is saved, either provide the PyArrow schema to
convert to or else provide a [PyArrow Table][pyarrow.Table] directly.
convert to or else provide a [PyArrow Table](pyarrow.Table) directly.
>>> import pyarrow as pa
>>> custom_schema = pa.schema([
@@ -610,46 +563,6 @@ class DBConnection(EnforceOverrides):
"""
raise NotImplementedError("serialize is not supported for this connection type")
def job(self, job_id: str) -> Job:
"""A [Job][lancedb.job.Job] handle for a server-side job by id.
The handle is constructed without a server round trip; an unknown id
surfaces when the handle is used. Dropping the handle has no effect
on the job itself.
"""
raise NotImplementedError("job is not supported for this connection type")
def list_jobs(self) -> List[JobInfo]:
"""List server-side jobs across the database's tables."""
raise NotImplementedError("list_jobs is not supported for this connection type")
def get_job(self, job_id: str) -> Optional[JobDescription]:
"""Describe a single server-side job by id.
Returns None when the server has no such job.
"""
raise NotImplementedError("get_job is not supported for this connection type")
def cancel_job(self, job_id: str) -> bool:
"""Request cancellation of a server-side job by id.
Returns True if the server accepted the cancellation, False if no
such job exists. Cancelling an already-terminal job is a no-op
success.
"""
raise NotImplementedError(
"cancel_job is not supported for this connection type"
)
def job_history(self, job_id: Optional[str] = None) -> List[pa.RecordBatch]:
"""The lifecycle event history of a server-side job, as Arrow batches.
Lists history across all jobs when `job_id` is None.
"""
raise NotImplementedError(
"job_history is not supported for this connection type"
)
class LanceDBConnection(DBConnection):
"""
@@ -707,9 +620,6 @@ class LanceDBConnection(DBConnection):
self._namespace_client_properties = namespace_client_properties
if _inner is not None:
self._conn = _inner
# Native-derived wrappers resolve this in their async reconstruction
# path so construction never synchronously re-enters LOOP.
self._read_consistency_interval = read_consistency_interval
self._cached_namespace_client = None
return
@@ -759,14 +669,11 @@ class LanceDBConnection(DBConnection):
# storage_options. Also, this class really shouldn't be holding any state
# beyond _conn.
self._conn = AsyncConnection(LOOP.run(do_connect()))
# Keep property access synchronous so debugger introspection cannot wait on
# the background loop while that thread is suspended at a breakpoint.
self._read_consistency_interval = read_consistency_interval
self._cached_namespace_client: Optional[LanceNamespace] = None
@property
def read_consistency_interval(self) -> Optional[timedelta]:
return self._read_consistency_interval
return LOOP.run(self._conn.get_read_consistency_interval())
@property
def session(self) -> Optional[Session]:
@@ -777,19 +684,15 @@ class LanceDBConnection(DBConnection):
return self._conn.uri
@classmethod
def from_inner(
cls,
inner: LanceDbConnection,
read_consistency_interval: Optional[timedelta],
):
return cls(
None,
read_consistency_interval=read_consistency_interval,
_inner=inner,
)
def from_inner(cls, inner: LanceDbConnection):
return cls(None, _inner=inner)
def __repr__(self) -> str:
return f"{self.__class__.__name__}(uri={self._conn.uri!r})"
val = f"{self.__class__.__name__}(uri={self._conn.uri!r}"
if self.read_consistency_interval is not None:
val += f", read_consistency_interval={repr(self.read_consistency_interval)}"
val += ")"
return val
@override
def serialize(self) -> str:
@@ -1226,47 +1129,6 @@ class LanceDBConnection(DBConnection):
)
)
@override
def job(self, job_id: str) -> Job:
"""A [Job][lancedb.job.Job] handle for a server-side job by id.
The handle is constructed without a server round trip; an unknown id
surfaces when the handle is used. Dropping the handle has no effect
on the job itself.
"""
return Job(self._conn.job(job_id))
@override
def list_jobs(self) -> List[JobInfo]:
"""List server-side jobs across the database's tables."""
return LOOP.run(self._conn.list_jobs())
@override
def get_job(self, job_id: str) -> Optional[JobDescription]:
"""Describe a single server-side job by id.
Returns None when the server has no such job.
"""
return LOOP.run(self._conn.get_job(job_id))
@override
def cancel_job(self, job_id: str) -> bool:
"""Request cancellation of a server-side job by id.
Returns True if the server accepted the cancellation, False if no
such job exists. Cancelling an already-terminal job is a no-op
success.
"""
return LOOP.run(self._conn.cancel_job(job_id))
@override
def job_history(self, job_id: Optional[str] = None) -> List[pa.RecordBatch]:
"""The lifecycle event history of a server-side job, as Arrow batches.
Lists history across all jobs when `job_id` is None.
"""
return LOOP.run(self._conn.job_history(job_id))
@override
def namespace_client(self) -> LanceNamespace:
"""Get the equivalent namespace client for this connection.
@@ -1667,7 +1529,7 @@ class AsyncConnection(object):
Data is converted to Arrow before being written to disk. For maximum
control over how data is saved, either provide the PyArrow schema to
convert to or else provide a [PyArrow Table][pyarrow.Table] directly.
convert to or else provide a [PyArrow Table](pyarrow.Table) directly.
>>> import pyarrow as pa
>>> custom_schema = pa.schema([
@@ -1976,43 +1838,6 @@ class AsyncConnection(object):
namespace_path = []
await self._inner.drop_all_tables(namespace_path=namespace_path)
def job(self, job_id: str) -> AsyncJob:
"""An [AsyncJob][lancedb.job.AsyncJob] handle for a server-side job
by id.
The handle is constructed without a server round trip; an unknown id
surfaces when the handle is used. Dropping the handle has no effect
on the job itself.
"""
return AsyncJob(self._inner.job(job_id))
async def list_jobs(self) -> List[JobInfo]:
"""List server-side jobs across the database's tables."""
return await self._inner.list_jobs()
async def get_job(self, job_id: str) -> Optional[JobDescription]:
"""Describe a single server-side job by id.
Returns None when the server has no such job.
"""
return await self._inner.get_job(job_id)
async def cancel_job(self, job_id: str) -> bool:
"""Request cancellation of a server-side job by id.
Returns True if the server accepted the cancellation, False if no
such job exists. Cancelling an already-terminal job is a no-op
success.
"""
return await self._inner.cancel_job(job_id)
async def job_history(self, job_id: Optional[str] = None) -> List[pa.RecordBatch]:
"""The lifecycle event history of a server-side job, as Arrow batches.
Lists history across all jobs when `job_id` is None.
"""
return await self._inner.job_history(job_id)
async def namespace_client(self) -> LanceNamespace:
"""Get the equivalent namespace client for this connection.
@@ -21,32 +21,3 @@ from .watsonx import WatsonxEmbeddings
from .voyageai import VoyageAIEmbeddingFunction
from .colpali import ColPaliEmbeddings
from .siglip import SigLipEmbeddings
# The API reference renders this package with a single mkdocstrings directive,
# which only picks up names listed here. New embedding functions must be added
# to both the imports above and this list, or they will silently go undocumented.
__all__ = [
"EmbeddingFunction",
"EmbeddingFunctionConfig",
"TextEmbeddingFunction",
"EmbeddingFunctionRegistry",
"get_registry",
"register",
"SentenceTransformerEmbeddings",
"OpenAIEmbeddings",
"OpenClipEmbeddings",
"BedRockText",
"CohereEmbeddingFunction",
"GeminiText",
"GteEmbeddings",
"InstructorEmbeddingFunction",
"JinaEmbeddings",
"OllamaEmbeddings",
"TransformersEmbeddingFunction",
"ColbertEmbeddings",
"VoyageAIEmbeddingFunction",
"WatsonxEmbeddings",
"ColPaliEmbeddings",
"ImageBindEmbeddings",
"SigLipEmbeddings",
]
+5 -5
View File
@@ -21,20 +21,20 @@ class BedRockText(TextEmbeddingFunction):
"""
Parameters
----------
name : str, default "amazon.titan-embed-text-v1"
name: str, default "amazon.titan-embed-text-v1"
The model ID of the bedrock model to use. Supported models for are:
- amazon.titan-embed-text-v1
- cohere.embed-english-v3
- cohere.embed-multilingual-v3
region : str, default "us-east-1"
region: str, default "us-east-1"
Optional name of the AWS Region in which the service should be called.
profile_name : str, default None
profile_name: str, default None
Optional name of the AWS profile to use for calling the Bedrock service.
If not specified, the default profile will be used.
assumed_role : str, default None
assumed_role: str, default None
Optional ARN of an AWS IAM role to assume for calling the Bedrock service.
If not specified, the current active credentials will be used.
role_session_name : str, default "lancedb-embeddings"
role_session_name: str, default "lancedb-embeddings"
Optional name of the AWS IAM role session to use for calling the Bedrock
service. If not specified, "lancedb-embeddings" name will be used.
+3 -5
View File
@@ -22,7 +22,7 @@ class CohereEmbeddingFunction(TextEmbeddingFunction):
Parameters
----------
name : str, default "embed-multilingual-v2.0"
name: str, default "embed-multilingual-v2.0"
The name of the model to use. List of acceptable models:
* embed-english-v3.0
@@ -33,14 +33,12 @@ class CohereEmbeddingFunction(TextEmbeddingFunction):
* embed-english-light-v2.0
* embed-multilingual-v2.0
source_input_type : str, default "search_document"
source_input_type: str, default "search_document"
The input type for the source column in the database
query_input_type : str, default "search_query"
query_input_type: str, default "search_query"
The input type for the query column in the database
Notes
-----
Cohere supports following input types:
| Input Type | Description |
+2 -2
View File
@@ -44,7 +44,7 @@ class ColPaliEmbeddings(EmbeddingFunction):
The token pooling strategy to use, by default "hierarchical".
- "hierarchical": Progressively pools tokens to reduce sequence length.
- "lambda": A simpler pooling that uses a custom `pooling_func`.
pooling_func : typing.Callable, optional
pooling_func: typing.Callable, optional
A function to use for pooling when `pooling_strategy` is "lambda".
pool_factor : int
Factor to reduce sequence length if token pooling is enabled (default 2).
@@ -52,7 +52,7 @@ class ColPaliEmbeddings(EmbeddingFunction):
Quantization configuration for the model. (default None, bitsandbytes needed)
batch_size : int
Batch size for processing inputs (default 2).
offload_folder : str, optional
offload_folder: str, optional
Folder to offload model weights if using CPU offloading (default None). This is
useful for large models that do not fit in memory.
"""
@@ -48,16 +48,16 @@ class GeminiText(TextEmbeddingFunction):
Parameters
----------
name : str, default "gemini-embedding-001"
name: str, default "gemini-embedding-001"
The name of the model to use. Supported models include:
- "gemini-embedding-001" (768 dimensions)
Note: The legacy "models/embedding-001" format is also supported but
"gemini-embedding-001" is recommended.
query_task_type : str, default "retrieval_query"
query_task_type: str, default "retrieval_query"
Sets the task type for the queries.
source_task_type : str, default "retrieval_document"
source_task_type: str, default "retrieval_document"
Sets the task type for ingestion.
Examples
+4 -4
View File
@@ -26,13 +26,13 @@ class GteEmbeddings(TextEmbeddingFunction):
Parameters
----------
name : str, default "thenlper/gte-large"
name: str, default "thenlper/gte-large"
The name of the model to use.
device : str, default "cpu"
device: str, default "cpu"
Sets the device type for the model.
normalize : str, default "True"
normalize: str, default "True"
Controls normalize param in encode function for the transformer.
mlx : bool, default False
mlx: bool, default False
Controls which model to use. False for gte-large,True for the mlx version.
Examples
+10 -9
View File
@@ -35,23 +35,23 @@ class InstructorEmbeddingFunction(TextEmbeddingFunction):
Parameters
----------
name : str
name: str
The name of the model to use. Available models are listed at
https://github.com/xlang-ai/instructor-embedding#model-list;
The default model is hkunlp/instructor-base
batch_size : int, default 32
batch_size: int, default 32
The batch size to use when generating embeddings
device : str, default "cpu"
device: str, default "cpu"
The device to use when generating embeddings
show_progress_bar : bool, default True
show_progress_bar: bool, default True
Whether to show a progress bar when generating embeddings
normalize_embeddings : bool, default True
normalize_embeddings: bool, default True
Whether to normalize the embeddings
quantize : bool, default False
quantize: bool, default False
Whether to quantize the model
source_instruction : str, default "represent the document for retrieval"
source_instruction: str, default "represent the document for retrieval"
The instruction for the source column
query_instruction : str, default "represent the document for retrieving the most
query_instruction: str, default "represent the document for retrieving the most
similar documents"
The instruction for the query
@@ -101,7 +101,8 @@ class InstructorEmbeddingFunction(TextEmbeddingFunction):
@weak_lru(maxsize=1)
def ndims(self):
return len(self.generate_embeddings([[self.source_instruction, "foo"]])[0])
model = self.get_model()
return model.encode("foo").shape[0]
def compute_query_embeddings(self, query: str, *args, **kwargs) -> List[np.array]:
return self.generate_embeddings([[self.query_instruction, query]])
+5 -6
View File
@@ -40,10 +40,10 @@ class JinaEmbeddings(EmbeddingFunction):
Parameters
----------
name : str, default "jina-clip-v1". Note that some models support both image
name: str, default "jina-clip-v1". Note that some models support both image
and text embeddings and some just text embedding
api_key : str, default None
api_key: str, default None
The api key to access Jina API. If you pass None, you can set JINA_API_KEY
environment variable
@@ -87,13 +87,12 @@ class JinaEmbeddings(EmbeddingFunction):
if isinstance(image, bytes):
image_dict = {"image": base64.b64encode(image).decode("utf-8")}
elif isinstance(image, (str, Path)):
parsed = urlparse(str(image))
parsed = urlparse.urlparse(image)
# TODO handle drive letter on windows.
PIL_Image = attempt_import_or_raise("PIL.Image", "pillow")
if parsed.scheme == "file":
pil_image = PIL_Image.open(parsed.path)
elif parsed.scheme == "" or (os.name == "nt" and len(parsed.scheme) == 1):
# A Windows drive letter parses as a one-character scheme
# ("C:\\img.png" -> scheme="c"), so treat it as a local path.
elif parsed.scheme == "":
pil_image = PIL_Image.open(image if os.name == "nt" else parsed.path)
elif parsed.scheme.startswith("http"):
pil_image = PIL_Image.open(io.BytesIO(url_retrieve(image)))
@@ -21,13 +21,13 @@ class SentenceTransformerEmbeddings(TextEmbeddingFunction):
Parameters
----------
name : str, default "all-MiniLM-L6-v2"
name: str, default "all-MiniLM-L6-v2"
The name of the model to use.
device : str, default "cpu"
device: str, default "cpu"
The device to use for the model
normalize : bool, default True
normalize: bool, default True
Whether to normalize the embeddings
trust_remote_code : bool, default True
trust_remote_code: bool, default True
Whether to trust the remote code
"""
+2 -2
View File
@@ -167,7 +167,7 @@ class VoyageAIEmbeddingFunction(EmbeddingFunction):
Parameters
----------
name : str
name: str
The name of the model to use. List of acceptable models:
* voyage-4 (1024 dims, general-purpose and multilingual retrieval)
@@ -185,7 +185,7 @@ class VoyageAIEmbeddingFunction(EmbeddingFunction):
* voyage-law-2
* voyage-code-2
output_dimension : int, optional
output_dimension: int, optional
The output dimension for models that support flexible dimensions.
Currently only voyage-multimodal-3.5 supports this feature.
Valid options: 256, 512, 1024 (default), 2048.
-12
View File
@@ -23,15 +23,3 @@ class MissingColumnError(KeyError):
return (
f"Error: Column '{self.column_name}' does not exist in the DataFrame object"
)
class JobFailedError(RuntimeError):
"""Exception raised when an asynchronous job reaches the failed state."""
pass
class JobCancelledError(RuntimeError):
"""Exception raised when an asynchronous job was cancelled."""
pass
+24 -33
View File
@@ -2,7 +2,7 @@
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
from dataclasses import dataclass
from typing import List, Literal, Optional
from typing import Literal, Optional
from ._lancedb import (
IndexConfig,
@@ -151,11 +151,6 @@ class FTS:
remove_stop_words : bool, default True
Whether to remove stop words. Stop words are common words that are often
removed from text before indexing. For example, in English "the" and "and".
custom_stop_words : list of str, optional
Custom words replace the built-in language stop words
and only take effect when ``remove_stop_words`` is True. ``None`` uses
the built-in language list, while an empty list explicitly uses no
stop words.
ascii_folding : bool, default True
Whether to fold ASCII characters. This converts accented characters to
their ASCII equivalent. For example, "café" would be converted to "cafe".
@@ -184,7 +179,6 @@ class FTS:
ngram_max_length: int = 3
prefix_only: bool = False
block_size: int = 128
custom_stop_words: Optional[List[str]] = None
@dataclass
@@ -219,7 +213,7 @@ class HnswPq:
distance has a range of (-, ). If the vectors are normalized (i.e. their
l2 norm is 1), then dot distance is equivalent to the cosine distance.
num_partitions: int, default sqrt(num_rows)
num_partitions, default sqrt(num_rows)
The number of IVF partitions to create.
@@ -228,7 +222,7 @@ class HnswPq:
will require too much memory. Each partition becomes its own HNSW graph, so
setting this value higher reduces the peak memory use of training.
num_sub_vectors: int, default is vector dimension / 16
num_sub_vectors, default is vector dimension / 16
Number of sub-vectors of PQ.
@@ -244,13 +238,13 @@ class HnswPq:
If the dimension is not visible by 8 then we use 1 subvector. This is not
ideal and will likely result in poor performance.
num_bits: int, default 8
num_bits: int, default 8
Number of bits to encode each sub-vector.
This value controls how much the sub-vectors are compressed. The more bits
the more accurate the index but the slower search. Only 4 and 8 are supported.
max_iterations: int, default 50
max_iterations, default 50
Max iterations to train kmeans.
@@ -263,7 +257,7 @@ class HnswPq:
those cases it is unlikely that setting this larger will lead to the index
converging anyways.
sample_rate: int, default 256
sample_rate, default 256
The rate used to calculate the number of training vectors for kmeans.
@@ -279,14 +273,14 @@ class HnswPq:
Increasing this value might improve the quality of the index but in
most cases the default should be sufficient.
m: int, default 20
m, default 20
The number of neighbors to select for each vector in the HNSW graph.
This value controls the tradeoff between search speed and accuracy.
The higher the value the more accurate the search but the slower it will be.
ef_construction: int, default 300
ef_construction, default 300
The number of candidates to evaluate during the construction of the HNSW graph.
@@ -297,7 +291,7 @@ class HnswPq:
This value should be set to a value that is not less than `ef` in the
search phase.
target_partition_size: int, default is 1,048,576
target_partition_size, default is 1,048,576
The target size of each partition.
@@ -351,7 +345,7 @@ class HnswSq:
distance has a range of (-, ). If the vectors are normalized (i.e. their
l2 norm is 1), then dot distance is equivalent to the cosine distance.
num_partitions: int, default sqrt(num_rows)
num_partitions, default sqrt(num_rows)
The number of IVF partitions to create.
@@ -360,7 +354,7 @@ class HnswSq:
will require too much memory. Each partition becomes its own HNSW graph, so
setting this value higher reduces the peak memory use of training.
max_iterations: int, default 50
max_iterations, default 50
Max iterations to train kmeans.
@@ -373,7 +367,7 @@ class HnswSq:
In those cases it is unlikely that setting this larger will lead to
the index converging anyways.
sample_rate: int, default 256
sample_rate, default 256
The rate used to calculate the number of training vectors for kmeans.
@@ -389,14 +383,14 @@ class HnswSq:
Increasing this value might improve the quality of the index but in
most cases the default should be sufficient.
m: int, default 20
m, default 20
The number of neighbors to select for each vector in the HNSW graph.
This value controls the tradeoff between search speed and accuracy.
The higher the value the more accurate the search but the slower it will be.
ef_construction: int, default 300
ef_construction, default 300
The number of candidates to evaluate during the construction of the HNSW graph.
@@ -407,7 +401,7 @@ class HnswSq:
This value should be set to a value that is not less than `ef` in the search
phase.
target_partition_size: int, default is 1,048,576
target_partition_size, default is 1,048,576
The target size of each partition.
@@ -460,7 +454,7 @@ class HnswFlat:
distance has a range of (-, ). If the vectors are normalized (i.e. their
l2 norm is 1), then dot distance is equivalent to the cosine distance.
num_partitions: int, default sqrt(num_rows)
num_partitions, default sqrt(num_rows)
The number of IVF partitions to create.
@@ -470,18 +464,18 @@ class HnswFlat:
graph, so setting this value higher reduces the peak memory use of
training.
max_iterations: int, default 50
max_iterations, default 50
Max iterations to train kmeans.
When training an IVF index we use kmeans to calculate the partitions.
This parameter controls how many iterations of kmeans to run.
sample_rate: int, default 256
sample_rate, default 256
The rate used to calculate the number of training vectors for kmeans.
m: int, default 20
m, default 20
The number of neighbors to select for each vector in the HNSW graph.
@@ -489,7 +483,7 @@ class HnswFlat:
The higher the value the more accurate the search but the slower it
will be.
ef_construction: int, default 300
ef_construction, default 300
The number of candidates to evaluate during the construction of the HNSW
graph.
@@ -501,7 +495,7 @@ class HnswFlat:
than 500. This value should be set to a value that is not less than `ef`
in the search phase.
target_partition_size: int, default is 1,048,576
target_partition_size, default is 1,048,576
The target size of each partition.
"""
@@ -605,7 +599,7 @@ class IvfFlat:
The default value is 256.
target_partition_size: int, default is 8192
target_partition_size, default is 8192
The target size of each partition.
@@ -769,7 +763,7 @@ class IvfPq:
The default value is 256.
target_partition_size: int, default is 8192
target_partition_size, default is 8192
The target size of each partition.
@@ -830,7 +824,7 @@ class IvfRq:
sample_rate: int, default 256
Controls the number of training vectors: sample_rate * num_partitions.
target_partition_size: int, default is 8192
target_partition_size, default is 8192
Target size of each partition.
"""
@@ -845,9 +839,6 @@ class IvfRq:
accelerator: Optional[str] = None
# The API reference renders this module with a single mkdocstrings directive,
# which only picks up names listed here. New public names must be added to this
# list, or they will silently go undocumented.
__all__ = [
"BTree",
"IvfPq",
-105
View File
@@ -1,105 +0,0 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
"""Handles to operations a server may run asynchronously."""
import asyncio
from datetime import timedelta
from typing import Optional
from lancedb.background_loop import LOOP
from . import _lancedb
class AsyncJob:
"""A handle to an operation that may still be running.
The operation may already be complete when the handle is created.
"""
def __init__(self, inner: Optional["_lancedb.Job"]):
self._inner = inner
@property
def id(self) -> Optional[str]:
"""Identifies the operation on the server that is running it.
Returned for correlating with server logs or the jobs API. Operations
that run in this process have no server id and return `None`. The value
is opaque: parsing it or storing it to resume the job later is not
supported.
"""
return self._inner.id if self._inner is not None else None
async def status(self) -> str:
"""The operation's current lifecycle state: "running", "finished",
"failed", or "cancelled".
A point snapshot; unlike `wait` it does not block or raise on a
terminal failure state. States a newer server reports that this
client version does not know pass through as-is.
"""
if self._inner is None:
return "finished"
return await self._inner.status()
async def wait(self, timeout: Optional[timedelta] = None):
"""Wait until the operation reaches a terminal state.
Raises `JobFailedError` if the operation failed, `JobCancelledError`
if it was cancelled, and `TimeoutError` if `timeout` elapses first.
"""
if self._inner is None:
return
if timeout is None:
await self._inner.wait()
else:
await asyncio.wait_for(self._inner.wait(), timeout.total_seconds())
async def cancel(self):
"""Request cancellation. Cancelling a finished operation is a no-op."""
if self._inner is None:
return
await self._inner.cancel()
class Job:
"""Synchronous counterpart of `AsyncJob`."""
def __init__(self, inner: Optional[AsyncJob]):
self._inner = inner
@property
def id(self) -> Optional[str]:
"""Identifies the operation on the server that is running it.
See :attr:`AsyncJob.id`.
"""
return self._inner.id if self._inner is not None else None
def status(self) -> str:
"""The operation's current lifecycle state: "running", "finished",
"failed", or "cancelled".
See :meth:`AsyncJob.status`.
"""
if self._inner is None:
return "finished"
return LOOP.run(self._inner.status())
def wait(self, timeout: Optional[timedelta] = None):
"""Block until the operation reaches a terminal state.
Raises `JobFailedError` if the operation failed, `JobCancelledError`
if it was cancelled, and `TimeoutError` if `timeout` elapses first.
"""
if self._inner is None:
return
LOOP.run(self._inner.wait(timeout))
def cancel(self):
"""Request cancellation. Cancelling a finished operation is a no-op."""
if self._inner is None:
return
LOOP.run(self._inner.cancel())
+1 -3
View File
@@ -92,10 +92,8 @@ class LanceMergeInsertBuilder(object):
self._when_not_matched_by_source_delete = True
if isinstance(condition, Expr):
self._when_not_matched_by_source_condition_expr = condition._inner
self._when_not_matched_by_source_condition = None
else:
elif condition is not None:
self._when_not_matched_by_source_condition = condition
self._when_not_matched_by_source_condition_expr = None
return self
def use_index(self, use_index: bool) -> LanceMergeInsertBuilder:
+1 -95
View File
@@ -38,11 +38,7 @@ from lance_namespace_urllib3_client.models.query_table_request_vector import (
QueryTableRequestVector,
)
from lance_namespace_urllib3_client.models.string_fts_query import StringFtsQuery
from lance_namespace.errors import (
NamespaceNotEmptyError,
NamespaceNotFoundError,
TableNotFoundError,
)
from lance_namespace.errors import NamespaceNotEmptyError, TableNotFoundError
from lancedb._lancedb import (
connect_namespace as _connect_namespace,
connect_namespace_client as _connect_namespace_client,
@@ -57,8 +53,6 @@ from lance_namespace import (
DropNamespaceResponse,
ListNamespacesResponse,
ListTablesResponse,
NamespaceExistsRequest,
TableExistsRequest,
)
from lancedb.table import AsyncTable, LanceTable, Table
from lancedb.util import validate_table_name
@@ -786,51 +780,6 @@ class LanceNamespaceDBConnection(DBConnection):
"""
return LOOP.run(self._inner.describe_namespace(namespace_path))
@override
def namespace_exists(self, namespace_id: List[str]) -> bool:
"""
Check if a namespace exists.
Parameters
----------
namespace_id : List[str]
The namespace identifier to check.
Returns
-------
bool
True if the namespace exists, False otherwise.
"""
request = NamespaceExistsRequest(id=namespace_id)
try:
self._namespace_client.namespace_exists(request)
return True
except NamespaceNotFoundError:
return False
@override
def table_exists(self, table_id: List[str]) -> bool:
"""
Check if a table exists.
Parameters
----------
table_id : List[str]
The table identifier to check (full path including namespace
segments and table name).
Returns
-------
bool
True if the table exists, False otherwise.
"""
request = TableExistsRequest(id=table_id)
try:
self._namespace_client.table_exists(request)
return True
except TableNotFoundError:
return False
@override
def list_tables(
self,
@@ -1284,49 +1233,6 @@ class AsyncLanceNamespaceDBConnection:
"""
return await self._inner.describe_namespace(namespace_path)
async def namespace_exists(self, namespace_id: List[str]) -> bool:
"""
Check if a namespace exists.
Parameters
----------
namespace_id : List[str]
The namespace identifier to check.
Returns
-------
bool
True if the namespace exists, False otherwise.
"""
request = NamespaceExistsRequest(id=namespace_id)
try:
self._namespace_client.namespace_exists(request)
return True
except NamespaceNotFoundError:
return False
async def table_exists(self, table_id: List[str]) -> bool:
"""
Check if a table exists.
Parameters
----------
table_id : List[str]
The table identifier to check (full path including namespace
segments and table name).
Returns
-------
bool
True if the table exists, False otherwise.
"""
request = TableExistsRequest(id=table_id)
try:
self._namespace_client.table_exists(request)
return True
except TableNotFoundError:
return False
async def list_tables(
self,
namespace_path: Optional[List[str]] = None,
+7 -12
View File
@@ -226,7 +226,7 @@ class PermutationBuilder:
async def do_execute():
inner_tbl = await self._async.execute()
return await LanceTable.from_inner(inner_tbl)
return LanceTable.from_inner(inner_tbl)
return LOOP.run(do_execute())
@@ -438,8 +438,7 @@ class Permutation:
_reader: Optional[PermutationReader] = None,
):
"""
Internal constructor. Use
[from_tables][lancedb.permutation.Permutation.from_tables] instead.
Internal constructor. Use [from_tables](#from_tables) instead.
"""
assert base_table is not None, "base_table is required"
assert selection is not None, "selection is required"
@@ -986,9 +985,8 @@ class Permutation:
types. Conversion of strings, lists, and structs will require creating python
objects and this is not zero-copy.
For custom formatting, use
[with_transform][lancedb.permutation.Permutation.with_transform] which
overrides this method.
For custom formatting, use [with_transform](#with_transform) which overrides
this method.
"""
assert format is not None, "format is required"
if format == "python":
@@ -1063,8 +1061,7 @@ class Permutation:
Note: this method returns a new permutation and does not modify `self`
It is provided for compatibility with the huggingface Dataset API.
Use [with_skip][lancedb.permutation.Permutation.with_skip] instead to
avoid confusion.
Use [with_skip](#with_skip) instead to avoid confusion.
"""
return self.with_skip(skip)
@@ -1087,8 +1084,7 @@ class Permutation:
Note: this method returns a new permutation and does not modify `self`
It is provided for compatibility with the huggingface Dataset API.
Use [with_take][lancedb.permutation.Permutation.with_take] instead to
avoid confusion.
Use [with_take](#with_take) instead to avoid confusion.
"""
return self.with_take(limit)
@@ -1111,8 +1107,7 @@ class Permutation:
Note: this method returns a new permutation and does not modify `self`
It is provided for compatibility with the huggingface Dataset API.
Use [with_repeat][lancedb.permutation.Permutation.with_repeat] instead
to avoid confusion.
Use [with_repeat](#with_repeat) instead to avoid confusion.
"""
return self.with_repeat(times)
-1
View File
@@ -1 +0,0 @@
-10
View File
@@ -153,16 +153,6 @@ def Vector(
return FixedSizeList
def _raise_bare_vector_error(*_args):
raise TypeError("Vector must be parameterized with a dimension, e.g. Vector(128).")
# Pydantic v1 and v2 otherwise treat the bare Vector factory as a field validator
# and inspect its signature, which produces misleading errors about internal types.
setattr(Vector, "__get_validators__", _raise_bare_vector_error)
setattr(Vector, "__get_pydantic_core_schema__", _raise_bare_vector_error)
def MultiVector(
dim: int, value_type: pa.DataType = pa.float32(), nullable: bool = True
) -> Type:
+23 -21
View File
@@ -52,6 +52,7 @@ from ._blob import (
finalize_blob_query_table,
replace_v2_blob_columns_with_bytes,
replace_v2_blob_columns_with_bytes_sync,
supports_blob_auto_row_id,
validate_blob_mode,
)
from .types import BlobMode, QueryProjection
@@ -650,8 +651,7 @@ class Query(pydantic.BaseModel):
distance_type : Optional[str]
the distance type to use for vector search
This can be l2 (default), cosine and dot. See
[metric definitions](https://lancedb.com/docs/search/vector-search/) for
This can be l2 (default), cosine and dot. See [metric definitions][search] for
more details.
If this is not a vector search this will be None.
@@ -664,9 +664,8 @@ class Query(pydantic.BaseModel):
- A higher number makes search more accurate but also slower.
- See discussion in
[Querying an ANN Index](https://lancedb.com/docs/indexing/)
for tuning advice.
- See discussion in [Querying an ANN Index][querying-an-ann-index] for
tuning advice.
Will be None if this is not a vector search.
refine_factor : Optional[int]
@@ -674,9 +673,8 @@ class Query(pydantic.BaseModel):
- A higher number makes search more accurate but also slower.
- See discussion in
[Querying an ANN Index](https://lancedb.com/docs/indexing/)
for tuning advice.
- See discussion in [Querying an ANN Index][querying-an-ann-index] for
tuning advice.
Will be None if this is not a vector search.
lower_bound : Optional[float]
@@ -1279,7 +1277,10 @@ class LanceQueryBuilder(ABC):
return self._with_row_id is True
def _blob_auto_row_id_enabled(self) -> bool:
if not supports_blob_auto_row_id(self._table):
return False
return blob_auto_row_id_for_scan(
self._table,
self._table.schema,
self._columns,
with_row_id=self._with_row_id,
@@ -1650,8 +1651,8 @@ class LanceVectorQueryBuilder(LanceQueryBuilder):
Higher values will yield better recall (more likely to find vectors if
they exist) at the expense of latency.
See discussion in [Querying an ANN Index](https://lancedb.com/docs/indexing/)
for tuning advice.
See discussion in [Querying an ANN Index][querying-an-ann-index] for
tuning advice.
This method sets both the minimum and maximum number of probes to the same
value. See `minimum_nprobes` and `maximum_nprobes` for more fine-grained
@@ -1751,8 +1752,8 @@ class LanceVectorQueryBuilder(LanceQueryBuilder):
As an example, a refine factor of 2 will sample 2x as many vectors as
requested, re-ranks them, and returns the top half most relevant results.
See discussion in [Querying an ANN Index](https://lancedb.com/docs/indexing/)
for tuning advice.
See discussion in [Querying an ANN Index][querying-an-ann-index] for
tuning advice.
Parameters
----------
@@ -2697,7 +2698,7 @@ class LanceHybridQueryBuilder(LanceQueryBuilder):
self._fts_query.phrase_query(True)
if self._distance_type:
self._vector_query.metric(self._distance_type)
if self._minimum_nprobes is not None:
if self._minimum_nprobes:
self._vector_query.minimum_nprobes(self._minimum_nprobes)
if self._maximum_nprobes is not None:
self._vector_query.maximum_nprobes(self._maximum_nprobes)
@@ -2770,7 +2771,7 @@ class AsyncQueryBase(object):
)
async def _maybe_add_blob_row_id(self) -> None:
if self._table is None:
if self._table is None or not supports_blob_auto_row_id(self._table):
self._blob_auto_row_id = False
self._blob_paths = ()
return
@@ -2778,6 +2779,7 @@ class AsyncQueryBase(object):
req = self._inner.to_query_request()
schema = await self._table.schema()
self._blob_auto_row_id = blob_auto_row_id_for_scan(
self._table,
schema,
req.select,
with_row_id=self._with_row_id,
@@ -3029,6 +3031,7 @@ class AsyncQueryBase(object):
schema = await self._table.schema()
blob_auto_row_id = blob_auto_row_id_for_scan(
self._table,
schema,
query.columns,
with_row_id=self._with_row_id,
@@ -3376,9 +3379,8 @@ class AsyncQuery(AsyncStandardQuery):
are various ANN search parameters that will let you fine tune your recall
accuracy vs search latency.
Vector searches always have a
[limit][lancedb.query.AsyncVectorQuery.limit]. If `limit` has not been
called then a default `limit` of 10 will be used.
Vector searches always have a [limit][]. If `limit` has not been called then
a default `limit` of 10 will be used.
Typically, a single vector is passed in as the query. However, you can also
pass in multiple vectors. When multiple vectors are passed in, if the vector
@@ -3509,9 +3511,8 @@ class AsyncFTSQuery(AsyncStandardQuery):
are various ANN search parameters that will let you fine tune your recall
accuracy vs search latency.
Hybrid searches always have a
[limit][lancedb.query.AsyncHybridQuery.limit]. If `limit` has not been
called then a default `limit` of 10 will be used.
Hybrid searches always have a [limit][]. If `limit` has not been called then
a default `limit` of 10 will be used.
Typically, a single vector is passed in as the query. However, you can also
pass in multiple vectors. This can be useful if you want to find the nearest
@@ -3874,9 +3875,10 @@ class AsyncHybridQuery(AsyncStandardQuery, AsyncVectorQueryBase):
req = fts_query._inner.to_query_request()
blob_auto_row_id = False
blob_paths: tuple[str, ...] = ()
if self._table is not None:
if self._table is not None and supports_blob_auto_row_id(self._table):
schema = await self._table.schema()
blob_auto_row_id = blob_auto_row_id_for_scan(
self._table,
schema,
req.select,
with_row_id=self._with_row_id,
-3
View File
@@ -11,9 +11,6 @@ from lancedb import __version__
from .header import HeaderProvider
from .oauth import OAuthConfig, OAuthFlowType
# The API reference renders this module with a single mkdocstrings directive,
# which only picks up names listed here. New public names must be added to this
# list, or they will silently go undocumented.
__all__ = [
"TimeoutConfig",
"RetryConfig",
+1 -51
View File
@@ -7,7 +7,7 @@ import json
import logging
from concurrent.futures import ThreadPoolExecutor
import sys
from typing import TYPE_CHECKING, Any, Dict, Iterable, List, Optional, Union
from typing import Any, Dict, Iterable, List, Optional, Union
from urllib.parse import urlparse
import warnings
@@ -23,10 +23,6 @@ import pyarrow as pa
from ..common import DATA
from ..db import DBConnection, LOOP
from ..job import Job
if TYPE_CHECKING:
from .._lancedb import JobDescription, JobInfo
from ..embeddings import EmbeddingFunctionConfig
from lance_namespace import (
LanceNamespace,
@@ -419,11 +415,6 @@ class RemoteDBConnection(DBConnection):
if namespace_path is None:
namespace_path = []
if storage_options is not None:
logging.info(
"storage_options is ignored in LanceDb Cloud"
" (storage is managed; set storage_options on connect() instead)"
)
if index_cache_size is not None:
logging.info(
"index_cache_size is ignored in LanceDb Cloud"
@@ -693,47 +684,6 @@ class RemoteDBConnection(DBConnection):
)
)
@override
def job(self, job_id: str) -> Job:
"""A [Job][lancedb.job.Job] handle for a server-side job by id.
The handle is constructed without a server round trip; an unknown id
surfaces when the handle is used. Dropping the handle has no effect
on the job itself.
"""
return Job(self._conn.job(job_id))
@override
def list_jobs(self) -> List["JobInfo"]:
"""List server-side jobs across the database's tables."""
return LOOP.run(self._conn.list_jobs())
@override
def get_job(self, job_id: str) -> Optional["JobDescription"]:
"""Describe a single server-side job by id.
Returns None when the server has no such job.
"""
return LOOP.run(self._conn.get_job(job_id))
@override
def cancel_job(self, job_id: str) -> bool:
"""Request cancellation of a server-side job by id.
Returns True if the server accepted the cancellation, False if no
such job exists. Cancelling an already-terminal job is a no-op
success.
"""
return LOOP.run(self._conn.cancel_job(job_id))
@override
def job_history(self, job_id: Optional[str] = None) -> List[pa.RecordBatch]:
"""The lifecycle event history of a server-side job, as Arrow batches.
Lists history across all jobs when `job_id` is None.
"""
return LOOP.run(self._conn.job_history(job_id))
@override
def namespace_client(self) -> LanceNamespace:
"""Get the equivalent namespace client for this connection.
+2 -2
View File
@@ -53,9 +53,9 @@ class RetryError(LanceDBClientError):
"""An error that occurs when the client has exceeded the maximum number of retries.
The retry strategy can be adjusted by setting the
[retry_config][lancedb.remote.ClientConfig.retry_config] in the client
[retry_config](lancedb.remote.ClientConfig.retry_config) in the client
configuration. This is passed in the `client_config` argument of
[connect][lancedb.connect] and [connect_async][lancedb.connect_async].
[connect](lancedb.connect) and [connect_async](lancedb.connect_async).
The __cause__ attribute of this exception will be the last exception that
caused the retry to fail. It will be an
+12 -46
View File
@@ -20,7 +20,6 @@ from typing import (
import warnings
from lancedb import __version__
from lancedb._blob import BlobFile
from lancedb._lancedb import (
AddColumnsResult,
@@ -48,7 +47,6 @@ from lancedb.index import (
IvfSq,
LabelList,
)
from lancedb.job import Job
from lancedb.remote.db import LOOP
from lancedb.table import IndexConfigType, KNOWN_METRICS
import pyarrow as pa
@@ -342,7 +340,6 @@ class RemoteTable(Table):
lower_case: bool = True,
stem: bool = True,
remove_stop_words: bool = True,
custom_stop_words: Optional[List[str]] = None,
ascii_folding: bool = True,
ngram_min_length: int = 3,
ngram_max_length: int = 3,
@@ -364,7 +361,6 @@ class RemoteTable(Table):
lower_case=lower_case,
stem=stem,
remove_stop_words=remove_stop_words,
custom_stop_words=custom_stop_words,
ascii_folding=ascii_folding,
ngram_min_length=ngram_min_length,
ngram_max_length=ngram_max_length,
@@ -542,34 +538,6 @@ class RemoteTable(Table):
)
)
def create_index_async(
self,
column: str,
*,
config: IndexConfigType,
replace: Optional[bool] = None,
wait_timeout: Optional[timedelta] = None,
name: Optional[str] = None,
train: bool = True,
) -> Job:
"""Create an index, returning a handle to the indexing job.
The job may already be complete when returned; callers must not assume
the index exists until :meth:`Job.wait` returns.
"""
return Job(
LOOP.run(
self._table.create_index_async(
column,
replace=replace,
config=config,
wait_timeout=wait_timeout,
name=name,
train=train,
)
)
)
def _is_legacy_create_index_call(
self,
first_arg: str,
@@ -610,9 +578,8 @@ class RemoteTable(Table):
progress: Optional[Union[bool, Callable, Any]] = None,
write_parallelism: Optional[int] = None,
) -> AddResult:
"""Add more data to the [Table][lancedb.table.Table].
It has the same API signature as the OSS version.
"""Add more data to the [Table](Table). It has the same API signature as
the OSS version.
Parameters
----------
@@ -672,8 +639,7 @@ class RemoteTable(Table):
fast_search: bool = False,
) -> LanceVectorQueryBuilder:
"""Create a search query to find the nearest neighbors
of the given query vector. We currently support
[vector search](https://lancedb.com/docs/search/vector-search/)
of the given query vector. We currently support [vector search][search]
All query options are defined in
[LanceVectorQueryBuilder][lancedb.query.LanceVectorQueryBuilder].
@@ -1069,22 +1035,22 @@ class RemoteTable(Table):
)
def blob_columns(self) -> list[str]:
return LOOP.run(self._table.blob_columns())
raise NotImplementedError(
"blob_columns() is not yet supported on the LanceDB Cloud"
)
def fetch_blobs(
self, column: str, row_ids: Union[list[int], pa.Table]
) -> pa.LargeBinaryArray:
return LOOP.run(self._table.fetch_blobs(column, row_ids))
def fetch_blobs(self, column: str, row_ids) -> pa.LargeBinaryArray:
raise NotImplementedError("fetch_blobs() is not supported on LanceDB Cloud")
def fetch_blob_ranges(self, column: str, requests) -> pa.LargeBinaryArray:
raise NotImplementedError(
"fetch_blob_ranges() is not supported on LanceDB Cloud"
)
def fetch_blob_files(
self, column: str, row_ids: Union[list[int], pa.Table]
) -> "list[Optional[BlobFile]]":
return LOOP.run(self._table.fetch_blob_files(column, row_ids))
def fetch_blob_files(self, column: str, row_ids):
raise NotImplementedError(
"fetch_blob_files() is not supported on LanceDB Cloud"
)
def head(self, n=5) -> pa.Table:
"""
@@ -14,9 +14,6 @@ from .answerdotai import AnswerdotaiRerankers
from .voyageai import VoyageAIReranker
from .watsonx import WatsonxReranker
# The API reference renders this module with a single mkdocstrings directive,
# which only picks up names listed here. New public names must be added to this
# list, or they will silently go undocumented.
__all__ = [
"Reranker",
"CrossEncoderReranker",
+36 -240
View File
@@ -40,7 +40,6 @@ from ._blob import (
from .types import BlobMode
from lancedb.arrow import peek_reader
from lancedb.background_loop import LOOP, embedding_executor
from lancedb.job import AsyncJob, Job
from .dependencies import (
_check_for_hugging_face,
_check_for_lance,
@@ -108,11 +107,6 @@ def _should_push_down_query_table(
return namespace_client is not None and "QueryTable" in pushdown_operations
def _polars_predicate_pushdown_barrier(frame: Any) -> Any:
"""Return a Polars frame unchanged while blocking predicate pushdown."""
return frame
_MODEL_BACKED_TOKENIZER_PREFIXES = ("jieba", "lindera")
_MODEL_BACKED_TOKENIZER_ERRORS = (
"unknown base tokenizer",
@@ -869,18 +863,12 @@ class Table(ABC):
"""
raise NotImplementedError
def to_polars(self, **kwargs) -> "pl.LazyFrame":
"""Return the table as a Polars LazyFrame.
Note
----
The Polars streaming engine is not supported because it does not currently
implement Python PyArrow dataset scans. Use the default engine when collecting
this LazyFrame.
def to_polars(self, **kwargs) -> "pl.DataFrame":
"""Return the table as a polars.DataFrame.
Returns
-------
polars.LazyFrame
polars.DataFrame
"""
raise NotImplementedError
@@ -989,24 +977,6 @@ class Table(ABC):
"""
raise NotImplementedError
def create_index_async(
self,
column: str,
*,
config: IndexConfigType,
replace: Optional[bool] = None,
wait_timeout: Optional[timedelta] = None,
name: Optional[str] = None,
train: bool = True,
) -> Job:
"""Create an index, returning a handle to the indexing job.
Takes the same arguments as :meth:`create_index`. The job may already
be complete when returned; callers must not assume the index exists
until :meth:`Job.wait` returns.
"""
raise NotImplementedError
def drop_index(self, name: str) -> None:
"""
Drop an index from the table.
@@ -1133,7 +1103,6 @@ class Table(ABC):
lower_case: bool = True,
stem: bool = True,
remove_stop_words: bool = True,
custom_stop_words: Optional[List[str]] = None,
ascii_folding: bool = True,
ngram_min_length: int = 3,
ngram_max_length: int = 3,
@@ -1201,9 +1170,6 @@ class Table(ABC):
remove_stop_words : bool, default True
Whether to remove stop words. Stop words are common words that are often
removed from text before indexing. For example, in English "the" and "and".
custom_stop_words : list of str, optional
Custom words that replace the built-in language stop words. ``None``
uses the built-in list; an empty list explicitly uses no stop words.
ascii_folding : bool, default True
Whether to fold ASCII characters. This converts accented characters to
their ASCII equivalent. For example, "café" would be converted to "cafe".
@@ -1241,7 +1207,7 @@ class Table(ABC):
progress: Optional[Union[bool, Callable, Any]] = None,
write_parallelism: Optional[int] = None,
) -> AddResult:
"""Add more data to the [Table][lancedb.table.Table].
"""Add more data to the [Table](Table).
Parameters
----------
@@ -1373,8 +1339,8 @@ class Table(ABC):
fts_columns: Optional[Union[str, List[str]]] = None,
) -> LanceQueryBuilder:
"""Create a search query to find the nearest neighbors
of the given query vector. We currently support [vector search](https://lancedb.com/docs/search/vector-search/)
and [full-text search](https://lancedb.com/docs/search/full-text-search/).
of the given query vector. We currently support [vector search][search]
and [full-text search][experimental-full-text-search].
All query options are defined in
[LanceQueryBuilder][lancedb.query.LanceQueryBuilder].
@@ -1604,10 +1570,8 @@ class Table(ABC):
"""Open lazy, seekable :class:`~lancedb._blob.BlobFile` handles.
Prefer this over :meth:`fetch_blobs` for large payloads. ``row_ids`` is
a ``list[int]`` or a query ``pyarrow.Table`` carrying row identity via
``_rowid`` or a ``_lance_row_id`` field on the blob descriptor. Null
rows are ``None``. Remote tables require LanceDB Cloud server 0.5.0 or
newer.
a ``list[int]`` or query ``pyarrow.Table`` with ``_rowid`` (or stashed
row-id metadata). Null rows are ``None``. Local tables only.
"""
@abstractmethod
@@ -1810,7 +1774,7 @@ class Table(ABC):
for faster reads.
Arguments are passed onto Lance's
`lance.dataset.DatasetOptimizer.compact_files`.
[compact_files][lance.dataset.DatasetOptimizer.compact_files].
For most cases, the default should be fine.
See Also
@@ -1864,8 +1828,6 @@ class Table(ABC):
retrain: bool, default False
This parameter is no longer used and is deprecated.
Notes
-----
The frequency an application should call optimize is based on the frequency of
data modifications. If data is frequently added, deleted, or updated then
optimize should be run frequently. A good rule of thumb is to run optimize if
@@ -2020,14 +1982,15 @@ class Table(ABC):
change permanent you can use the `[Self::restore]` method.
Any operation that modifies the table will fail while the table is in a checked
out state. To return the table to a normal state use
`[Self::checkout_latest]`.
out state.
Parameters
----------
version: int | str,
The version to check out. A version number (`int`) or a tag
(`str`) can be provided.
To return the table to a normal state use `[Self::checkout_latest]`
"""
@abstractmethod
@@ -2193,15 +2156,11 @@ class LanceTable(Table):
return self.name
@classmethod
async def from_inner(cls, tbl: LanceDBTable):
from .db import AsyncConnection, LanceDBConnection
def from_inner(cls, tbl: LanceDBTable):
from .db import LanceDBConnection
async_tbl = AsyncTable(tbl)
inner_conn = tbl.database()
read_consistency_interval = await AsyncConnection(
inner_conn
).get_read_consistency_interval()
conn = LanceDBConnection.from_inner(inner_conn, read_consistency_interval)
conn = LanceDBConnection.from_inner(tbl.database())
return cls(
conn,
async_tbl.name,
@@ -2505,7 +2464,13 @@ class LanceTable(Table):
return LOOP.run(self._table.count_rows(filter))
def __repr__(self) -> str:
return f"{self.__class__.__name__}(name={self.name!r}, _conn={self._conn!r})"
val = f"{self.__class__.__name__}(name={self.name!r}"
if self._conn.read_consistency_interval is not None:
val += ", read_consistency_interval={!r}".format(
self._conn.read_consistency_interval
)
val += f", _conn={self._conn!r})"
return val
def __str__(self) -> str:
return self.__repr__()
@@ -2580,9 +2545,6 @@ class LanceTable(Table):
2. Currently we've disabled push-down of the filters from polars
because polars pushdown into pyarrow uses pyarrow compute
expressions rather than SQl strings (which LanceDB supports)
3. The Polars streaming engine is not supported because it does not
currently implement Python PyArrow dataset scans. Use the default
engine when collecting this LazyFrame.
Returns
-------
@@ -2591,12 +2553,8 @@ class LanceTable(Table):
from lancedb.integrations.pyarrow import PyarrowDatasetAdapter
dataset = PyarrowDatasetAdapter(self)
# Polars 1.32's non-PyArrow callback path passes batch_size twice. Keep
# the compatible PyArrow path, but block predicates because this adapter
# cannot translate PyArrow expressions into LanceDB filters.
return pl.scan_pyarrow_dataset(dataset, batch_size=batch_size).map_batches(
_polars_predicate_pushdown_barrier,
predicate_pushdown=False,
return pl.scan_pyarrow_dataset(
dataset, allow_pyarrow_filter=False, batch_size=batch_size
)
# New unified API overload
@@ -2821,34 +2779,6 @@ class LanceTable(Table):
)
)
def create_index_async(
self,
column: str,
*,
config: IndexConfigType,
replace: Optional[bool] = None,
wait_timeout: Optional[timedelta] = None,
name: Optional[str] = None,
train: bool = True,
) -> Job:
"""Create an index, returning a handle to the indexing job.
The job may already be complete when returned; callers must not assume
the index exists until :meth:`Job.wait` returns.
"""
return Job(
LOOP.run(
self._table.create_index_async(
column,
replace=replace,
config=config,
wait_timeout=wait_timeout,
name=name,
train=train,
)
)
)
def _is_legacy_create_index_call(
self,
first_arg: str,
@@ -3125,7 +3055,6 @@ class LanceTable(Table):
lower_case: bool = True,
stem: bool = True,
remove_stop_words: bool = True,
custom_stop_words: Optional[List[str]] = None,
ascii_folding: bool = True,
ngram_min_length: int = 3,
ngram_max_length: int = 3,
@@ -3172,7 +3101,6 @@ class LanceTable(Table):
"lower_case": lower_case,
"stem": stem,
"remove_stop_words": remove_stop_words,
"custom_stop_words": custom_stop_words,
"ascii_folding": ascii_folding,
"ngram_min_length": ngram_min_length,
"ngram_max_length": ngram_max_length,
@@ -3180,7 +3108,6 @@ class LanceTable(Table):
}
else:
tokenizer_configs = self.infer_tokenizer_configs(tokenizer_name)
tokenizer_configs["custom_stop_words"] = custom_stop_words
config = FTS(block_size=block_size, **tokenizer_configs)
@@ -3453,8 +3380,8 @@ class LanceTable(Table):
fts_columns: Optional[Union[str, List[str]]] = None,
) -> LanceQueryBuilder:
"""Create a search query to find the nearest neighbors
of the given query vector. We currently support [vector search](https://lancedb.com/docs/search/vector-search/)
and [full-text search](https://lancedb.com/docs/search/full-text-search/).
of the given query vector. We currently support [vector search][search]
and [full-text search][search].
Examples
--------
@@ -3484,9 +3411,8 @@ class LanceTable(Table):
- *default None*.
Acceptable types are: list, np.ndarray, PIL.Image.Image
- If None then the
select/[where][lancedb.query.LanceQueryBuilder.where]/limit clauses
are applied to filter the table
- If None then the select/[where][sql]/limit clauses are applied
to filter the table
vector_column_name: str, optional
The name of the vector column to search.
@@ -3880,8 +3806,6 @@ class LanceTable(Table):
retrain: bool, default False
This parameter is no longer used and is deprecated.
Notes
-----
The frequency an application should call optimize is based on the frequency of
data modifications. If data is frequently added, deleted, or updated then
optimize should be run frequently. A good rule of thumb is to run optimize if
@@ -3976,28 +3900,6 @@ class LanceTable(Table):
[`AsyncTable.get_lsm_write_spec`][lancedb.AsyncTable.get_lsm_write_spec]."""
return LOOP.run(self._table.get_lsm_write_spec())
def checkpoint_lsm(self) -> None:
"""Synchronous version of
[`AsyncTable.checkpoint_lsm`][lancedb.AsyncTable.checkpoint_lsm]."""
return LOOP.run(self._table.checkpoint_lsm())
def flush_lsm(self) -> None:
"""Synchronous version of
[`AsyncTable.flush_lsm`][lancedb.AsyncTable.flush_lsm]."""
return LOOP.run(self._table.flush_lsm())
def compact_lsm(self) -> None:
"""Synchronous version of
[`AsyncTable.compact_lsm`][lancedb.AsyncTable.compact_lsm]."""
return LOOP.run(self._table.compact_lsm())
def get_lsm_stats(self, *, include_generation_rows: bool = False) -> Optional[dict]:
"""Synchronous version of
[`AsyncTable.get_lsm_stats`][lancedb.AsyncTable.get_lsm_stats]."""
return LOOP.run(
self._table.get_lsm_stats(include_generation_rows=include_generation_rows)
)
def close_lsm_writers(self) -> None:
"""Close cached MemWAL shard writers. See
[`AsyncTable.close_lsm_writers`][lancedb.AsyncTable.close_lsm_writers]."""
@@ -4708,67 +4610,6 @@ class AsyncTable:
"""
return await self._inner.get_lsm_write_spec()
async def checkpoint_lsm(self) -> None:
"""Converge this table's LSM write path into its base table.
One flush, sealing every memtable into L0, then compaction triggers
until every generation that existed at that moment has reached base.
The loop runs client-side, reading progress from ``get_lsm_stats``.
Best-effort: generations created *while* it runs are deliberately not
waited on, which is what lets it terminate on a table taking writes.
Idempotent and safe on a cadence.
There is no deadline, and the caller owns that. It returns when the
target generations are gone, raises on a terminal server fault, and
otherwise waits however long the server takes. A slow table and a
stuck one are the same picture from the client: the compactor pool is
shared across every table on the node, so a checkpoint queued behind
unrelated work looks exactly like one that is merging. Wrap this in
``asyncio.wait_for`` for a wall-clock bound; abandoning it partway
costs nothing.
"""
return await self._inner.checkpoint_lsm()
async def flush_lsm(self) -> None:
"""Seal every bucket's active memtable into L0.
Does not touch the base table moving L0 into base is
`compact_lsm`. On a node that has not claimed this table, this claims
it and replays its WAL log first.
"""
return await self._inner.flush_lsm()
async def compact_lsm(self) -> None:
"""Trigger a background L0 to base compaction pass per bucket.
Returns once the passes are dispatched, not once they finish: watch
``get_lsm_stats`` for progress, or use ``checkpoint_lsm`` to loop
until the current L0 has reached base.
"""
return await self._inner.compact_lsm()
async def get_lsm_stats(
self, *, include_generation_rows: bool = False
) -> Optional[dict]:
"""Read live per-bucket LSM state.
Answers "how far behind is my fresh tier", "which bucket is hot", and
"why is my fresh-tier vector search brute-force". Mutates no table
state, though on a node that has not claimed this table it claims it,
exactly as a read would.
Returns ``None`` only when the LSM write path is not enabled.
Parameters
----------
include_generation_rows
Report a row count per L0 generation. Off by default: each count
opens an uncached Lance dataset, and ``checkpoint_lsm`` polls this
needing only generation numbers.
"""
return await self._inner.get_lsm_stats(include_generation_rows)
async def close_lsm_writers(self) -> None:
"""Drain and close any cached MemWAL shard writers for this table.
@@ -4843,7 +4684,7 @@ class AsyncTable:
Parameters
----------
**kwargs
Forwarded to `lance.dataset`.
Forwarded to [`lance.dataset`][lance.dataset].
Returns
-------
@@ -5019,46 +4860,6 @@ class AsyncTable:
)
raise e
async def create_index_async(
self,
column: str,
*,
replace: Optional[bool] = None,
config: Optional[
Union[
IvfFlat,
IvfPq,
IvfRq,
HnswPq,
HnswSq,
HnswFlat,
BTree,
Bitmap,
LabelList,
Fm,
FTS,
]
] = None,
wait_timeout: Optional[timedelta] = None,
name: Optional[str] = None,
train: bool = True,
) -> AsyncJob:
"""Create an index, returning a handle to the indexing job.
Takes the same arguments as :meth:`create_index`. The job may already
be complete when returned; callers must not assume the index exists
until :meth:`AsyncJob.wait` resolves.
"""
job = await self._inner.create_index_async(
column,
index=config,
replace=replace,
wait_timeout=wait_timeout,
name=name,
train=train,
)
return AsyncJob(job)
async def drop_index(self, name: str) -> None:
"""
Drop an index from the table.
@@ -5202,7 +5003,7 @@ class AsyncTable:
progress: Optional[Union[bool, Callable, Any]] = None,
write_parallelism: Optional[int] = None,
) -> AddResult:
"""Add more data to the [AsyncTable][lancedb.table.AsyncTable].
"""Add more data to the [Table](Table).
Parameters
----------
@@ -5404,8 +5205,8 @@ class AsyncTable:
fts_columns: Optional[Union[str, List[str]]] = None,
) -> Union[AsyncHybridQuery, AsyncFTSQuery, AsyncVectorQuery]:
"""Create a search query to find the nearest neighbors
of the given query vector. We currently support [vector search](https://lancedb.com/docs/search/vector-search/)
and [full-text search](https://lancedb.com/docs/search/full-text-search/).
of the given query vector. We currently support [vector search][search]
and [full-text search][experimental-full-text-search].
All query options are defined in [AsyncQuery][lancedb.query.AsyncQuery].
@@ -5966,14 +5767,15 @@ class AsyncTable:
change permanent you can use the `[Self::restore]` method.
Any operation that modifies the table will fail while the table is in a checked
out state. To return the table to a normal state use
`[Self::checkout_latest]`.
out state.
Parameters
----------
version: int | str,
The version to check out. A version number (`int`) or a tag
(`str`) can be provided.
To return the table to a normal state use `[Self::checkout_latest]`
"""
try:
await self._inner.checkout(version)
@@ -6157,8 +5959,6 @@ class AsyncTable:
retrain: bool, default False
This parameter is no longer used and is deprecated.
Notes
-----
The frequency an application should call optimize is based on the frequency of
data modifications. If data is frequently added, deleted, or updated then
optimize should be run frequently. A good rule of thumb is to run optimize if
@@ -6334,9 +6134,7 @@ class TableStatistics:
Attributes
----------
total_bytes: int
The total size, in bytes, of the table's data files, index files, and
overlay files. Read from the manifest, so this excludes deletion files
and manifests.
The total number of bytes in the table.
num_rows: int
The total number of rows in the table.
num_indices: int
@@ -6541,8 +6339,6 @@ class Branches:
dry_run: bool, default False
When True, only preview. When False, attempt the merge.
Notes
-----
A rejected merge returns ``status="rejected"`` instead of raising.
"""
return LOOP.run(self._table.branches.merge(from_branch, dry_run))
-5
View File
@@ -395,11 +395,6 @@ def _(value: dict):
)
@value_to_sql.register(pa.Scalar)
def _(value: pa.Scalar):
return value_to_sql(value.as_py())
@value_to_sql.register(np.ndarray)
def _(value: np.ndarray):
return value_to_sql(value.tolist())
+3 -3
View File
@@ -226,13 +226,13 @@ def test_fetch_blob_ranges_validates_requests():
table = _blob_table("range_validation", [{"id": 1, "image": b"abc"}])
row_id = _row_ids_by_id(table)[1]
with pytest.raises(ValueError, match="exceeds blob size"):
with pytest.raises(RuntimeError, match="exceeds blob size"):
table.fetch_blob_ranges("image", [(row_id, 2, 2)])
with pytest.raises(ValueError, match="offset \\+ length overflowed"):
with pytest.raises(RuntimeError, match="offset \\+ length overflowed"):
table.fetch_blob_ranges("image", [(row_id, 2**64 - 1, 1)])
with pytest.raises(ValueError, match="row IDs"):
with pytest.raises(ValueError, match="row ids"):
table.fetch_blob_ranges("image", [(2**64 - 1, 0, 1)])
-44
View File
@@ -2,11 +2,9 @@
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
import inspect
import re
import sys
from datetime import timedelta
from importlib import resources
import os
from types import SimpleNamespace
@@ -19,10 +17,6 @@ from lance_namespace.errors import NamespaceNotEmptyError, TableNotFoundError
from lancedb.pydantic import LanceModel, Vector
def test_package_includes_pep_561_marker():
assert resources.files(lancedb).joinpath("py.typed").is_file()
def test_basic(tmp_path):
db = lancedb.connect(tmp_path)
@@ -68,44 +62,6 @@ def test_basic(tmp_path):
assert db.open_table("test").name == db["test"].name
def test_sync_debugger_inspection_does_not_use_background_loop(tmp_path, monkeypatch):
from lancedb.background_loop import LOOP
db = lancedb.connect(tmp_path)
table = db.create_table("test", data=[{"id": 1}])
def fail_run(*args, **kwargs):
raise AssertionError("debugger inspection should not use the background loop")
monkeypatch.setattr(LOOP, "run", fail_run)
# Debuggers enumerate and evaluate every exposed attribute when expanding a
# variable. This must remain safe while their breakpoint suspends LOOP's thread.
members = dict(inspect.getmembers(db))
assert members["uri"] == str(tmp_path)
assert members["read_consistency_interval"] is None
assert repr(db) == f"LanceDBConnection(uri={str(tmp_path)!r})"
assert repr(table) == f"LanceTable(name='test', _conn={db!r})"
def test_read_consistency_interval_does_not_use_background_loop(tmp_path, monkeypatch):
from lancedb.background_loop import LOOP
from lancedb.db import LanceDBConnection
consistency_interval = timedelta(seconds=5)
db = lancedb.connect(tmp_path, read_consistency_interval=consistency_interval)
db_from_inner = LanceDBConnection.from_inner(db._inner, consistency_interval)
def fail_run(*args, **kwargs):
raise AssertionError("properties should not use the Python background loop")
monkeypatch.setattr(LOOP, "run", fail_run)
assert db.read_consistency_interval == consistency_interval
assert db_from_inner.read_consistency_interval == consistency_interval
def test_ingest_pd(tmp_path):
db = lancedb.connect(tmp_path)
+27 -51
View File
@@ -64,23 +64,6 @@ def test_embedding_function(tmp_path):
assert np.allclose(actual, expected)
def test_instructor_ndims_uses_instruction():
instructor = get_registry().get("instructor").create()
model = MagicMock()
model.encode.return_value = np.zeros((1, 384))
with patch.object(type(instructor), "get_model", return_value=model):
assert instructor.ndims() == 384
model.encode.assert_called_once_with(
[[instructor.source_instruction, "foo"]],
batch_size=instructor.batch_size,
show_progress_bar=instructor.show_progress_bar,
normalize_embeddings=instructor.normalize_embeddings,
device=instructor.device,
)
def test_embedding_function_variables():
@register("variable-testing")
class VariableTestingFunction(TextEmbeddingFunction):
@@ -132,16 +115,34 @@ def test_embedding_function_variables():
assert func.safe_model_dump()["secret_key"] == "$var:secret"
def test_openai_variables_survive_metadata_round_trip():
def test_parse_functions_with_variables():
@register("variable-parsing-test")
class VariableParsingFunction(TextEmbeddingFunction):
api_key: str
base_url: Optional[str] = None
@staticmethod
def sensitive_keys():
return ["api_key"]
def ndims(self):
return 10
def generate_embeddings(self, texts):
# Mock implementation that just returns random embeddings
# In real usage, this would use the api_key to call an API
return [np.random.rand(self.ndims()).tolist() for _ in texts]
registry = EmbeddingFunctionRegistry.get_instance()
registry.set_var("test_api_key", "sk-test-key-12345")
registry.set_var("test_base_url", "https://api.example.com")
conf = EmbeddingFunctionConfig(
source_column="text",
vector_column="vector",
function=registry.get("openai").create(
api_key="$var:test_api_key", base_url="https://api.example.com"
function=registry.get("variable-parsing-test").create(
api_key="$var:test_api_key", base_url="$var:test_base_url"
),
)
@@ -149,10 +150,7 @@ def test_openai_variables_survive_metadata_round_trip():
# Create a mock arrow table with the metadata
schema = pa.schema(
[
pa.field("text", pa.string()),
pa.field("vector", pa.list_(pa.float32(), 1536)),
]
[pa.field("text", pa.string()), pa.field("vector", pa.list_(pa.float32(), 10))]
)
table = pa.table({"text": [], "vector": []}, schema=schema)
table = table.replace_schema_metadata(metadata)
@@ -166,15 +164,13 @@ def test_openai_variables_survive_metadata_round_trip():
assert parsed_func.api_key == "sk-test-key-12345"
assert parsed_func.base_url == "https://api.example.com"
embeddings = parsed_func.generate_embeddings(["test text"])
assert len(embeddings) == 1
assert len(embeddings[0]) == 10
assert parsed_func.safe_model_dump()["api_key"] == "$var:test_api_key"
with patch("lancedb.embeddings.openai.attempt_import_or_raise") as import_openai:
parsed_func._openai_client
import_openai.return_value.OpenAI.assert_called_once_with(
api_key="sk-test-key-12345", base_url="https://api.example.com"
)
def test_embedding_with_bad_results(tmp_path):
@register("null-embedding")
@@ -631,23 +627,3 @@ def test_url_retrieve_downloads_image():
image_bytes = url_retrieve(image_url)
img = Image.open(io.BytesIO(image_bytes))
assert img.size[0] > 0 and img.size[1] > 0
def test_jina_generate_image_input_dict_local_path(tmp_path):
"""
JinaEmbeddings._generate_image_input_dict must accept a local image path
(str or Path), not just bytes. Previously it crashed with
`AttributeError: 'function' object has no attribute 'urlparse'` on any
str/Path input because it called `urlparse.urlparse(image)` instead of
`urlparse(image)` (urlparse was imported as a function, not a module).
"""
Image = pytest.importorskip("PIL.Image")
from lancedb.embeddings.jinaai import JinaEmbeddings
image_path = tmp_path / "test.png"
Image.new("RGB", (4, 4), color="red").save(image_path, format="PNG")
for image in (str(image_path), image_path):
image_dict = JinaEmbeddings._generate_image_input_dict(image)
assert "image" in image_dict
assert isinstance(image_dict["image"], str) and len(image_dict["image"]) > 0
-20
View File
@@ -219,13 +219,11 @@ def test_create_inverted_index(table, with_position):
table.create_fts_index(
"text",
with_position=with_position,
custom_stop_words=["puppy"],
name="custom_fts_index",
)
indices = table.list_indices()
fts_indices = [i for i in indices if i.index_type == "FTS"]
assert any(i.name == "custom_fts_index" for i in fts_indices)
assert fts_indices[0].index_details["custom_stop_words"] == ["puppy"]
@pytest.mark.parametrize("block_size", [128, 256])
@@ -245,24 +243,6 @@ def test_create_inverted_index_rejects_invalid_block_size(table):
table.create_index("text", config=FTS(block_size=129))
def test_custom_stop_words_list(table):
table.create_index(
"text",
config=FTS(stem=False, custom_stop_words=["lance"]),
)
assert table.list_indices()[0].index_details["custom_stop_words"] == ["lance"]
tokens = table.tokenize("the lance data", column="text")
assert [token.text for token in tokens] == ["the", "data"]
empty_tokens = ldb.tokenize("the lance data", stem=False, custom_stop_words=[])
assert [token.text for token in empty_tokens] == ["the", "lance", "data"]
with pytest.raises(TypeError, match=r"custom_stop_words.*int"):
ldb.tokenize(
"the lance data",
custom_stop_words=["lance", 42],
)
def test_search_fts(table):
table.create_fts_index("text")
results = table.search("puppy").select(["id", "text"]).limit(5).to_list()
+1 -94
View File
@@ -12,7 +12,7 @@ import pyarrow.compute as pc
import pytest
import pytest_asyncio
from lancedb.index import BTree, FTS, IvfPq
from lancedb.index import FTS
from lancedb.table import AsyncTable, Table
@@ -99,86 +99,6 @@ async def test_async_hybrid_query_filters(table: AsyncTable):
assert result["text"].to_pylist() == ["cat", "b"]
@pytest.mark.asyncio
async def test_hybrid_query_with_stale_fixed_size_binary_prefilter(
tmpdir_factory,
):
tmp_path = str(tmpdir_factory.mktemp("stale_scalar_prefilter"))
db = await lancedb.connect_async(tmp_path)
def fixed_size_binary(value: int) -> bytes:
return value.to_bytes(16, byteorder="big")
num_rows = 1000
data = pa.table(
{
"space_id": pa.array(
[fixed_size_binary(i) for i in range(num_rows)],
type=pa.binary(16),
),
"text": ["book"] * num_rows,
"vector": pa.array(
[[float(i), float(i)] for i in range(num_rows)],
type=pa.list_(pa.float32(), 2),
),
}
)
table = await db.create_table("test", data)
await table.create_index(
"vector", config=IvfPq(num_partitions=4, num_sub_vectors=2)
)
await table.create_index("space_id", config=BTree())
await table.create_index("text", config=FTS(with_position=False))
# Advance the search indices without advancing the scalar index. This is the
# state that previously let hybrid search use an incomplete scalar prefilter.
await table.add(data)
lance_dataset = await table.to_lance()
lance_dataset.optimize.optimize_indices(index_names=["vector_idx", "text_idx"])
await table.checkout_latest()
scalar_stats = await table.index_stats("space_id_idx")
assert scalar_stats is not None
assert scalar_stats.num_indexed_rows == num_rows
assert scalar_stats.num_unindexed_rows == num_rows
for index_name in ["vector_idx", "text_idx"]:
search_stats = await table.index_stats(index_name)
assert search_stats is not None
assert search_stats.num_indexed_rows == num_rows * 2
assert search_stats.num_unindexed_rows == 0
matching_ids = [5, 10, 15, 20, 25, 30]
literals = [
f"arrow_cast(0x{fixed_size_binary(i).hex()}, 'FixedSizeBinary(16)')"
for i in matching_ids
]
predicate = f"space_id IN ({', '.join(literals)})"
expected_ids = sorted(fixed_size_binary(i) for i in matching_ids for _ in range(2))
vector_query = (
table.query().where(predicate).nearest_to([5.0, 5.0]).limit(num_rows * 2)
)
vector_results = await vector_query.to_arrow()
assert sorted(vector_results["space_id"].to_pylist()) == expected_ids
fts_query = (
table.query().where(predicate).nearest_to_text("book").limit(num_rows * 2)
)
fts_results = await fts_query.to_arrow()
assert sorted(fts_results["space_id"].to_pylist()) == expected_ids
hybrid_results = await (
table.query()
.where(predicate)
.nearest_to([5.0, 5.0])
.nearest_to_text("book")
.limit(num_rows * 2)
.to_arrow()
)
assert sorted(hybrid_results["space_id"].to_pylist()) == expected_ids
@pytest.mark.asyncio
async def test_async_hybrid_query_default_limit(table: AsyncTable):
# add 10 new rows
@@ -203,19 +123,6 @@ async def test_async_hybrid_query_default_limit(table: AsyncTable):
assert texts.count("a") == 1
def test_hybrid_query_minimum_nprobes_zero_raises(sync_table: Table):
# minimum_nprobes(0) must raise the same validation error a plain vector
# query raises, not silently no-op because 0 is falsy.
with pytest.raises(ValueError, match="minimum_nprobes must be greater than 0"):
(
sync_table.search(query_type="hybrid")
.vector([0.0, 0.4])
.text("dog")
.minimum_nprobes(0)
.to_arrow()
)
def test_hybrid_query_distance_range(sync_table: Table):
reranker = RRFReranker(return_score="all")
result = (
-33
View File
@@ -1,33 +0,0 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
import re
import shutil
import subprocess
import sys
import lancedb._lancedb as _lancedb
import pytest
@pytest.mark.skipif(sys.platform != "linux", reason="ldd is Linux-specific")
def test_native_extension_does_not_link_openssl():
"""OpenSSL-linked wheels abort when imported on RHEL hosts in FIPS mode."""
ldd = shutil.which("ldd")
if ldd is None:
pytest.skip("ldd is not installed")
result = subprocess.run(
[ldd, _lancedb.__file__],
check=True,
capture_output=True,
text=True,
)
openssl_libraries = re.findall(
r"^\s*(lib(?:crypto|ssl)\S*)\s+=>", result.stdout, flags=re.MULTILINE
)
assert not openssl_libraries, (
"the LanceDB native extension must use rustls instead of linking OpenSSL: "
f"{openssl_libraries}"
)
-34
View File
@@ -84,15 +84,6 @@ async def binary_table(db_async):
)
@pytest.mark.asyncio
async def test_create_index_async_returns_done_job(some_table: AsyncTable):
job = await some_table.create_index_async("id", config=BTree())
assert job.id is None
await job.wait()
assert len(await some_table.list_indices()) == 1
await job.cancel()
@pytest.mark.asyncio
async def test_create_scalar_index(some_table: AsyncTable):
# Can create
@@ -372,31 +363,6 @@ async def test_create_vector_index(some_table: AsyncTable):
assert stats.num_indices == 1
@pytest.mark.asyncio
async def test_create_ivf_index_reports_unsplittable_partitions(db_async):
dim = 8
num_partitions = 300 # More than 256 selects hierarchical k-means.
base_vectors = [[float(row == column) for column in range(dim)] for row in range(5)]
vectors = pa.array(base_vectors * 200, pa.list_(pa.float32(), dim))
table = await db_async.create_table(
"unsplittable_partitions",
pa.table({"vector": vectors}),
)
error_pattern = (
rf"Cannot create {num_partitions} IVF partitions: k-means could only form"
)
with pytest.raises(RuntimeError, match=error_pattern):
await table.create_index(
"vector",
config=IvfFlat(
distance_type="dot",
num_partitions=num_partitions,
max_iterations=10,
),
)
@pytest.mark.asyncio
async def test_create_4bit_ivfpq_index(some_table: AsyncTable):
# Can create
@@ -18,7 +18,6 @@ Tests verify:
"""
import copy
import os
import shutil
import sys
import tempfile
@@ -240,7 +239,7 @@ def create_tracking_namespace(
dir_props = {f"storage.{k}": v for k, v in storage_options_with_refresh.items()}
if os.path.isabs(bucket_name) or bucket_name.startswith("file://"):
if bucket_name.startswith("/") or bucket_name.startswith("file://"):
dir_props["root"] = f"{bucket_name}/namespace_root"
else:
dir_props["root"] = f"s3://{bucket_name}/namespace_root"
@@ -768,70 +767,3 @@ def test_namespace_with_schema_only(s3_bucket: str, use_custom: bool):
# Verify data was added
assert table.count_rows() == 2
@pytest.mark.parametrize("use_custom", [False, True], ids=["DirectoryNS", "CustomNS"])
def test_namespace_exists(use_custom: bool):
"""
Test namespace_exists returns True for existing and False for non-existent.
"""
temp_dir = tempfile.mkdtemp()
try:
ns_client, _ = create_tracking_namespace(
bucket_name=temp_dir,
storage_options={},
credential_expires_in_seconds=3600,
use_custom=use_custom,
)
db = LanceNamespaceDBConnection(ns_client)
namespace_name = f"test_ns_{uuid.uuid4().hex[:8]}"
db.create_namespace([namespace_name])
# Existing namespace should return True
assert db.namespace_exists(namespace_id=[namespace_name]) is True
# Non-existent namespace should return False
assert db.namespace_exists(namespace_id=["nonexistent_ns"]) is False
finally:
shutil.rmtree(temp_dir, ignore_errors=True)
@pytest.mark.parametrize("use_custom", [False, True], ids=["DirectoryNS", "CustomNS"])
def test_table_exists(use_custom: bool):
"""
Test table_exists returns True for existing table and False for non-existent.
"""
temp_dir = tempfile.mkdtemp()
try:
ns_client, _ = create_tracking_namespace(
bucket_name=temp_dir,
storage_options={},
credential_expires_in_seconds=3600,
use_custom=use_custom,
)
db = LanceNamespaceDBConnection(ns_client)
namespace_name = f"test_ns_{uuid.uuid4().hex[:8]}"
db.create_namespace([namespace_name])
table_name = f"test_table_{uuid.uuid4().hex}"
namespace_path = [namespace_name]
schema = pa.schema(
[
pa.field("id", pa.int64()),
pa.field("vector", pa.list_(pa.float32(), 2)),
pa.field("text", pa.string()),
]
)
db.create_table(table_name, schema=schema, namespace_path=namespace_path)
# Existing table should return True
table_id = namespace_path + [table_name]
assert db.table_exists(table_id=table_id) is True
# Non-existent table should return False
assert db.table_exists(table_id=namespace_path + ["nonexistent_table"]) is False
finally:
shutil.rmtree(temp_dir, ignore_errors=True)
@@ -1,42 +0,0 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
import importlib
import re
import sys
from pathlib import Path
import pytest
def test_pyo3_abi_matches_minimum_supported_python():
project_dir = Path(__file__).parents[2]
pyproject = (project_dir / "pyproject.toml").read_text()
cargo_manifest = (project_dir / "Cargo.toml").read_text()
minimum_python = re.search(
r'^requires-python\s*=\s*">=(\d+)\.(\d+)"$', pyproject, re.MULTILINE
)
assert minimum_python is not None
major, minor = minimum_python.groups()
expected_abi = f"abi3-py{major}{minor}"
configured_abis = re.findall(r'"(abi3-py\d+)"', cargo_manifest)
assert configured_abis == [expected_abi, expected_abi], (
"the pyo3 runtime and build ABI features must both match requires-python"
)
@pytest.mark.skipif(sys.platform != "win32", reason="Windows wheel regression test")
def test_windows_wheel_tag_and_native_import():
project_dir = Path(__file__).parents[2]
wheels = list((project_dir.parent / "target" / "wheels").glob("lancedb-*.whl"))
if not wheels:
pytest.skip("no wheel artifact is available in this development environment")
assert len(wheels) == 1
assert wheels[0].name.endswith("-cp310-abi3-win_amd64.whl")
native_module = importlib.import_module("lancedb._lancedb")
assert Path(native_module.__file__).suffix == ".pyd"
-20
View File
@@ -6,7 +6,6 @@ import math
import pytest
from lancedb import DBConnection, Table, connect
from lancedb.background_loop import LOOP
from lancedb.permutation import Permutation, Permutations, permutation_builder
@@ -32,25 +31,6 @@ def test_split_random_ratios(mem_db):
assert 65 <= split_1_count <= 75 # ~70% ± tolerance
def test_execute_does_not_reenter_background_loop(tmp_path, monkeypatch):
import threading
db = connect(tmp_path)
tbl = db.create_table("test_table", pa.table({"x": range(10)}))
original_run = LOOP.run
def fail_on_reentry(future):
assert threading.current_thread() is not LOOP.thread
return original_run(future)
monkeypatch.setattr(LOOP, "run", fail_on_reentry)
permutation_tbl = permutation_builder(tbl).execute()
assert permutation_tbl.count_rows() == 10
assert permutation_tbl._conn.read_consistency_interval is None
def test_split_random_counts(mem_db):
"""Test random splitting with absolute counts."""
tbl = mem_db.create_table(
-11
View File
@@ -415,17 +415,6 @@ def test_nullable_vector():
assert schema == pa.schema([pa.field("vec", pa.list_(pa.float32(), 16), True)])
def test_bare_vector_raises_clear_error():
namespace = {
"__name__": "test_model_without_pyarrow",
"LanceModel": LanceModel,
"Vector": Vector,
}
with pytest.raises(TypeError, match=r"Vector must be parameterized.*Vector\(128\)"):
exec("class TestModel(LanceModel):\n vector: Vector", namespace)
def test_fixed_size_list_field():
class TestModel(pydantic.BaseModel):
vec: Vector(16)

Some files were not shown because too many files have changed in this diff Show More