Compare commits

..

4 Commits

Author SHA1 Message Date
XYZhan 6c4269bbb0 refactor(lsm): remove the index-catchup activation surface
Follows the lance change that removes the feature bit. With one set of
semantics there is no mode to switch into, so `require_mem_wal_index_catchup`
goes from the trait, `NativeTable`, and the LSM merge module.

`exclusion_watermarks` loses its `catchup_required` argument and keeps the
conservative branch: an index with no `index_catchup` entry is not known to
hold the compacted rows, so its generations stay readable from their
SSTables. That is what makes an existing table safe to read the moment the
new binary starts -- nothing is excluded until an index records that it
covers it.

Two tests changed because they encoded the branch that is gone, and they
had conflated two different things: an index that is *caught up* and an
index that is *untracked* both fell back to the compaction watermark.
Untracked now retains everything, so the tests assert that and cover the
genuinely-caught-up case separately.

484 lancedb lib tests pass.
2026-08-20 18:29:56 -04:00
XYZhan 41e0161067 feat(lsm): move catch-up activation to an explicit API 2026-08-10 09:10:05 -04:00
XYZhan f75279b69f feat(lsm): strict catch-up semantics and activation, on the 10.1 line
lancedb#3780 plus the two pieces it deferred: a missing index_catchup entry
means not-caught-up once the feature bit is set, and turning WAL on requires
catch-up while the table is provably clean. Local only, for end-to-end work
before the lance bump.
2026-08-09 23:19:21 -04:00
lancedb automation 667cf32e78 chore: update lance dependency to v10.1.0-beta.2 2026-08-02 00:12:47 +00:00
86 changed files with 827 additions and 6526 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
Generated
+215 -229
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.1.0-beta.2", default-features = false, "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=10.1.0-beta.2", "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=10.1.0-beta.2", "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=10.1.0-beta.2", "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=10.1.0-beta.2", default-features = false, "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=10.1.0-beta.2", "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=10.1.0-beta.2", "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=10.1.0-beta.2", "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=10.1.0-beta.2", default-features = false, "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=10.1.0-beta.2", "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=10.1.0-beta.2", "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=10.1.0-beta.2", "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=10.1.0-beta.2", "tag" = "v10.1.0-beta.2", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=10.1.0-beta.2", "tag" = "v10.1.0-beta.2", "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")
+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 -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
+1 -1
View File
@@ -31,7 +31,7 @@ is also an [asynchronous API client](#connections-asynchronous).
## Namespaces (Synchronous)
A namespace-backed connection resolves tables through a
[Lance namespace](https://lance-format.github.io/lance-namespace/) service instead of
[Lance namespace](https://lancedb.github.io/lance-namespace/) service instead of
listing a storage directory.
::: lancedb.connect_namespace
+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.1.0-beta.2</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 @@
[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;
-32
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) => {
+1 -47
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 () => {
+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 = {
+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.1-beta.0",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "@lancedb/lancedb",
"version": "0.37.1",
"version": "0.37.1-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
}
}
}
+3 -10
View File
@@ -339,9 +339,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 +356,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 +1039,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"
-4
View File
@@ -355,10 +355,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: ...
+3 -17
View File
@@ -707,9 +707,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 +756,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,16 +771,8 @@ 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})"
@@ -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]])
+3 -4
View File
@@ -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)))
+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 -1
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())
-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:
+1 -1
View File
@@ -2697,7 +2697,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)
+11 -118
View File
@@ -108,11 +108,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 +864,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
@@ -1606,8 +1595,8 @@ class Table(ABC):
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.
rows are ``None``. Unsupported on LanceDB Cloud, where
:meth:`fetch_blobs` returns full bytes instead.
"""
@abstractmethod
@@ -2193,15 +2182,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,
@@ -2580,9 +2565,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 +2573,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
@@ -3976,28 +3954,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 +4664,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.
@@ -6334,9 +6229,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
-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())
+2 -31
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,21 @@ 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):
def test_sync_repr_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")
raise AssertionError("repr should not use the Python 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
+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}"
)
-25
View File
@@ -372,31 +372,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
@@ -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)
-9
View File
@@ -570,15 +570,6 @@ def test_query_builder(table):
assert all(np.array(rs[0]["vector"]) == [1, 2])
def test_query_multiple_vectors(table):
results = table.search([np.array([1, 2]), np.array([4, 5])]).limit(1).to_list()
assert len(results) == 2
results_by_query = {result["query_index"]: result for result in results}
assert results_by_query[0]["id"] == 1
assert results_by_query[1]["id"] == 2
def test_with_row_id(table: lancedb.table.Table):
rs = table.search().with_row_id(True).to_arrow()
assert "_rowid" in rs.column_names
+2 -39
View File
@@ -35,12 +35,6 @@ def make_mock_http_handler(handler):
return MockLanceDBHandler
@pytest.mark.parametrize("db_name", ["a" * 64, "invalid..database"])
def test_connect_rejects_invalid_cloud_dns_hostname(db_name):
with pytest.raises(ValueError, match="DNS labels must contain 1 to 63 bytes"):
lancedb.connect(f"db://{db_name}", api_key="fake")
@contextlib.contextmanager
def mock_lancedb_connection(handler):
with http.server.HTTPServer(
@@ -2061,24 +2055,6 @@ def blob_remote_table(*, server_version=Version("0.5.0")):
request.send_header("phalanx-version", str(server_version))
request.end_headers()
request.wfile.write(json.dumps(BLOB_DESCRIBE_RESPONSE).encode())
elif request.path.startswith("/v1/table/test/blob/image/"):
path = request.path.partition("?")[0]
row_id = int(path.split("/")[-2])
payload = {10: b"alpha", 20: None, 30: b"gamma"}[row_id]
if payload is None:
request.send_response(204)
request.end_headers()
return
byte_range = request.headers["Range"].removeprefix("bytes=")
start_text, end_text = byte_range.split("-", maxsplit=1)
start = int(start_text)
end = int(end_text) if end_text else len(payload) - 1
chunk = payload[start : end + 1]
request.send_response(206)
request.send_header("Content-Range", f"bytes {start}-{end}/{len(payload)}")
request.send_header("Content-Length", str(len(chunk)))
request.end_headers()
request.wfile.write(chunk)
elif request.path == "/v1/table/test/query/":
content_len = int(request.headers.get("Content-Length", 0))
body = json.loads(request.rfile.read(content_len))
@@ -2116,21 +2092,8 @@ def test_remote_blob_columns_and_fetch():
assert table.blob_columns() == ["image"]
blobs = table.fetch_blobs("image", [10, 20, 30])
assert blobs.to_pylist() == [b"alpha", None, b"gamma"]
def test_remote_blob_files_are_lazy_seekable_handles():
with blob_remote_table() as table:
files = table.fetch_blob_files("image", [10, 20, 30])
assert len(files) == 3
alpha, null_row, gamma = files
assert null_row is None
assert alpha is not None
assert gamma is not None
assert alpha.size() == 5
assert alpha.read_range(1, 3) == b"lph"
gamma.seek(2)
assert gamma.read() == b"mma"
with pytest.raises(NotImplementedError, match="Use fetch_blobs for full bytes"):
table.fetch_blob_files("image", [10, 20, 30])
def test_remote_blob_fetch_accepts_query_table():
+4 -325
View File
@@ -2,14 +2,10 @@
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
import ctypes
import gc
import os
import sys
import threading
import warnings
import weakref
from concurrent.futures import ThreadPoolExecutor
from datetime import date, datetime, timedelta
from time import sleep
from typing import List
@@ -102,30 +98,6 @@ def test_basic(mem_db: DBConnection):
assert table.to_arrow() == expected_data
def test_search_preserves_nulls_from_sliced_arrow_table(mem_db: DBConnection):
data = pa.table(
{
"id": [0, 1, 2, 3, 4],
"score_cn": [None, 22, None, 5, 8],
"score_mt": [None, 42, None, 5, 8],
"vector": [
[20, 19, -1, -1],
[41, 38, 22, 42],
[10, 10, -1, -1],
[5, 5, 5, 5],
[8, 8, 8, 8],
],
}
).slice(1)
table = mem_db.create_table("sliced_nullable", data=data)
result = table.search([41, 38, 22, 42]).limit(1).to_arrow()
assert result["id"].to_pylist() == [1]
assert result["score_cn"].to_pylist() == [22]
assert result["score_mt"].to_pylist() == [42]
def test_table_to_pandas_default_matches_arrow(tmp_db: DBConnection):
pd = pytest.importorskip("pandas")
data = pa.table({"id": [1, 2], "text": ["one", "two"]})
@@ -462,38 +434,6 @@ def test_add(mem_db: DBConnection):
_add(table, schema)
def test_add_releases_arrow_buffers_without_gc(mem_db: DBConnection):
"""Regression test for https://github.com/lancedb/lancedb/issues/2512."""
schema = pa.schema([pa.field("x", pa.int64())])
table = mem_db.create_table("test_add_releases_arrow_buffers", schema=schema)
class BufferOwner:
def __init__(self, size: int):
self.memory = ctypes.create_string_buffer(size)
owner_refs = []
gc_was_enabled = gc.isenabled()
gc.disable()
try:
for _ in range(3):
size = 8 * 1024
owner = BufferOwner(size)
arrow_buffer = pa.foreign_buffer(
ctypes.addressof(owner.memory), size, owner
)
array = pa.Array.from_buffers(pa.int64(), 1024, [None, arrow_buffer])
batch = pa.RecordBatch.from_arrays([array], schema=schema)
owner_refs.append(weakref.ref(owner))
table.add(batch)
del batch, array, arrow_buffer, owner
assert all(owner_ref() is None for owner_ref in owner_refs)
finally:
if gc_was_enabled:
gc.enable()
def test_add_write_parallelism(mem_db: DBConnection):
schema = pa.schema([pa.field("id", pa.int64())])
table = mem_db.create_table("test", schema=schema)
@@ -929,7 +869,6 @@ def test_polars(mem_db: DBConnection):
# enter table to polars dataframe
result = table.to_polars()
assert isinstance(result, pl.LazyFrame)
assert np.allclose(result.collect()["vector"].to_list(), data["vector"])
# make sure filtering isn't broken
@@ -1846,27 +1785,6 @@ def test_add_with_empty_fixed_size_list_drops_bad_rows(mem_db: DBConnection):
assert np.allclose(data["embedding"].to_pylist()[0], np.array([0.1] * 16))
def test_add_nullable_fixed_size_list_with_none(mem_db: DBConnection):
"""Regression test for issue #2340."""
table = mem_db.create_table(
"test_nullable_fixed_size_list",
schema=pa.schema(
[
pa.field("id", pa.string()),
pa.field("feature", pa.list_(pa.float32(), 256)),
pa.field("tags", pa.list_(pa.string())),
]
),
)
table.add([{"id": "1", "feature": None, "tags": ["tag1", "tag2"]}])
result = table.to_arrow()
assert result.to_pylist() == [
{"id": "1", "feature": None, "tags": ["tag1", "tag2"]}
]
def test_add_nullable_struct_with_none(mem_db: DBConnection):
"""Regression test for issue #2654: a nullable struct column whose
first batch contains only None values must not crash in
@@ -1906,33 +1824,6 @@ def test_add_nullable_struct_with_none(mem_db: DBConnection):
assert result.column("data").to_pylist() == [{"x": 1.0}, None]
def test_read_mostly_null_list_v2_2_page_boundary(tmp_path):
# Regression test for #3194. This row/value count crosses a v2.2 structural
# encoding page boundary where Lance 3.0.0 sliced repetition/definition
# levels by row offset and decoded child arrays at different lengths.
num_rows = 64_885
num_values = 217
list_type = pa.list_(pa.float32())
source = pa.table(
{
"id": np.arange(num_rows, dtype=np.int64),
"coords": pa.array(
[[1.0, 2.0, 3.0, 4.0]] * num_values + [None] * (num_rows - num_values),
type=list_type,
),
}
)
db = lancedb.connect(
tmp_path,
storage_options={"new_table_data_storage_version": "2.2"},
)
table = db.create_table("test_sparse_nullable_list", data=source)
result = table.search().select(["id", "coords"]).limit(num_rows).to_arrow()
assert result.equals(source)
def test_add_with_integer_embeddings_preserves_casting(mem_db: DBConnection):
class Schema(LanceModel):
text: str
@@ -2218,45 +2109,6 @@ def test_merge(tmp_db: DBConnection, tmp_path):
table.merge(other_dataset, left_on="id")
@pytest.mark.parametrize("storage_version", ["legacy", "stable"])
def test_search_after_merge(tmp_path, storage_version):
pytest.importorskip("lance")
pd = pytest.importorskip("pandas")
db = lancedb.connect(
tmp_path,
storage_options={"new_table_data_storage_version": storage_version},
)
rng = np.random.default_rng(42)
row_count = 512
vectors = rng.standard_normal((row_count, 8)).astype(np.float32)
table = db.create_table(
"search_after_merge",
data=pd.DataFrame(
{
"id": [str(i) for i in range(row_count)],
"vector": list(vectors),
}
),
)
table.create_index("vector", config=IvfPq(num_partitions=1, num_sub_vectors=2))
links = pd.DataFrame(
{
"id": [str(i) for i in range(row_count // 2)],
"link": [f"https://example.com/{i}" for i in range(row_count // 2)],
}
)
table.merge(links, left_on="id")
query = table.search(vectors[-1]).refine_factor(50).limit(10)
assert "ANN" in query.explain_plan(verbose=True)
result = query.to_arrow()
links_by_id = dict(zip(result["id"].to_pylist(), result["link"].to_pylist()))
assert links_by_id[str(row_count - 1)] is None
def test_delete(mem_db: DBConnection):
table = mem_db.create_table(
"my_table",
@@ -2272,27 +2124,6 @@ def test_delete(mem_db: DBConnection):
assert table.to_arrow()["id"].to_pylist() == [1]
def test_concurrent_deletes_are_thread_safe(mem_db: DBConnection):
num_workers = 8
table = mem_db.create_table(
"my_table", data=[{"id": row_id} for row_id in range(num_workers)]
)
barrier = threading.Barrier(num_workers)
def delete(row_id: int):
barrier.wait()
return table.delete(f"id = {row_id}")
with ThreadPoolExecutor(max_workers=num_workers) as pool:
results = list(pool.map(delete, range(num_workers)))
assert all(result.num_deleted_rows == 1 for result in results)
assert sorted(result.version for result in results) == list(
range(2, num_workers + 2)
)
assert table.count_rows() == 0
def test_delete_expr(mem_db: DBConnection):
table = mem_db.create_table(
"my_table",
@@ -2343,20 +2174,6 @@ def test_update(mem_db: DBConnection):
assert np.allclose(v, np.array([[1.2, 1.9], [1.1, 1.1]]))
def test_update_with_arrow_scalar(mem_db: DBConnection):
schema = pa.schema({"id": pa.int64(), "vector": pa.list_(pa.float32(), 4)})
table = mem_db.create_table("my_table", schema=schema)
table.add([{"id": 1, "vector": [1.0, 2.0, 3.0, 4.0]}])
value = table.search().select(["vector"]).limit(1).to_arrow()["vector"][0]
assert isinstance(value, pa.FixedSizeListScalar)
result = table.update(where="id == 1", values={"vector": value})
assert result.rows_updated == 1
assert table.to_arrow()["vector"].to_pylist() == [[1.0, 2.0, 3.0, 4.0]]
def test_update_types(mem_db: DBConnection):
table = mem_db.create_table(
"my_table",
@@ -2524,55 +2341,6 @@ def test_merge_insert(mem_db: DBConnection):
)
def test_merge_insert_nullable_pandas_into_pydantic_schema(mem_db: DBConnection):
# Regression test for https://github.com/lancedb/lancedb/issues/2366
pd = pytest.importorskip("pandas")
class Document(LanceModel):
id: int
title: str
content: str
table = mem_db.create_table("documents", schema=Document)
table.add(
pd.DataFrame(
{
"title": ["Old title", "Unchanged"],
"id": [2, 3],
"content": ["Old content", "Keep this"],
}
)
)
# Pandas produces nullable Arrow fields, in an order that differs from the
# non-nullable Pydantic schema. This is valid as long as the data has no nulls.
new_data = pd.DataFrame(
{
"title": ["Inserted", "Updated"],
"id": [1, 2],
"content": ["New row", "New content"],
}
)
result = (
table.merge_insert("id")
.when_matched_update_all()
.when_not_matched_insert_all()
.execute(new_data)
)
assert result.num_inserted_rows == 1
assert result.num_updated_rows == 1
expected = pa.Table.from_pylist(
[
{"id": 1, "title": "Inserted", "content": "New row"},
{"id": 2, "title": "Updated", "content": "New content"},
{"id": 3, "title": "Unchanged", "content": "Keep this"},
],
schema=Document.to_arrow_schema(),
)
assert table.to_arrow().sort_by("id") == expected
def test_merge_insert_by_source_delete_expr(mem_db: DBConnection):
table = mem_db.create_table(
"my_table",
@@ -2596,29 +2364,6 @@ def test_merge_insert_by_source_delete_expr(mem_db: DBConnection):
assert table.to_arrow().sort_by("a") == expected
def test_merge_insert_by_source_delete_reconfigure(mem_db: DBConnection):
# Calling when_not_matched_by_source_delete() again with no condition must
# widen the delete to unconditional, not keep the earlier condition around.
table = mem_db.create_table(
"my_table",
data=pa.table({"a": [1, 2, 3], "b": ["a", "b", "c"]}),
)
new_data = pa.table({"a": [2, 4], "b": ["x", "z"]})
merge_insert_res = (
table.merge_insert("a")
.when_matched_update_all()
.when_not_matched_insert_all()
.when_not_matched_by_source_delete("a > 2")
.when_not_matched_by_source_delete()
.execute(new_data)
)
assert merge_insert_res.num_deleted_rows == 2
expected = pa.table({"a": [2, 4], "b": ["x", "z"]})
assert table.to_arrow().sort_by("a") == expected
@pytest.mark.asyncio
async def test_merge_insert_by_source_delete_expr_async(
mem_db_async: AsyncConnection,
@@ -2673,36 +2418,6 @@ def test_merge_insert_subschema(mem_db: DBConnection, data_format):
assert table.to_arrow().sort_by("id") == expected
def test_repeated_partial_merge_insert_with_scalar_index(mem_db: DBConnection):
def make_batch(start: int) -> pa.Table:
return pa.table(
{
"id": [f"id-{i:04}" for i in range(start, start + 100)],
"category": ["A"] * 100,
"value_a": [float(i) for i in range(start, start + 100)],
"value_b": [float(i) / 10 for i in range(100)],
}
)
table = mem_db.create_table("my_table", data=make_batch(0))
table.add(make_batch(100))
table.add(make_batch(200))
table.create_index("id", config=BTree())
ids = [f"id-{i:04}" for i in range(100, 200)]
for value in (999.0, 888.0):
result = (
table.merge_insert("id")
.when_matched_update_all()
.execute(pa.table({"id": ids, "value_a": [value] * 100}))
)
assert result.num_updated_rows == 100
actual = table.to_arrow().sort_by("id")
assert actual.num_rows == 300
assert actual["value_a"].to_pylist()[100:200] == [888.0] * 100
@pytest.mark.asyncio
async def test_merge_insert_async(mem_db_async: AsyncConnection):
data = pa.table({"a": [1, 2, 3], "b": ["a", "b", "c"]})
@@ -2799,40 +2514,15 @@ def test_create_with_embedding_function(mem_db: DBConnection):
assert actual == expected
def test_create_f16_table_from_arrow_data(mem_db: DBConnection):
dimension = 32
num_rows = 512
values = pa.array(
np.random.default_rng(42)
.standard_normal(num_rows * dimension)
.astype(np.float16)
)
df = pa.table(
{
"text": [f"s-{i}" for i in range(num_rows)],
"vector": pa.FixedSizeListArray.from_arrays(values, dimension),
}
)
table = mem_db.create_table("f16_tbl", data=df)
assert table.schema.field("vector").type == pa.list_(pa.float16(), dimension)
table.create_index(num_partitions=2, num_sub_vectors=2)
query = df["vector"][2].as_py()
expected = table.search(query).limit(2).to_arrow()
assert "s-2" in expected["text"].to_pylist()
def test_create_f16_table(mem_db: DBConnection):
class MyTable(LanceModel):
text: str
vector: Vector(32, value_type=pa.float16())
rng = np.random.default_rng(42)
df = pa.table(
{
"text": [f"s-{i}" for i in range(512)],
"vector": [rng.standard_normal(32).astype(np.float16) for _ in range(512)],
"vector": [np.random.randn(32).astype(np.float16) for _ in range(512)],
}
)
table = mem_db.create_table(
@@ -3713,8 +3403,7 @@ def test_stats(mem_db: DBConnection):
stats = table.stats()
print(f"{stats=}")
assert stats == {
# Full on-disk size of the data file, footer and metadata included.
"total_bytes": 633,
"total_bytes": 60,
"num_rows": 2,
"num_indices": 0,
"fragment_stats": {
@@ -3732,13 +3421,6 @@ def test_stats(mem_db: DBConnection):
},
}
# Index files count toward total_bytes too (only deletion files and
# manifests are excluded).
table.create_index("id", config=BTree())
stats_with_index = table.stats()
assert stats_with_index["num_indices"] == 1
assert stats_with_index["total_bytes"] > stats["total_bytes"]
def test_create_table_empty_list_with_schema(mem_db: DBConnection):
"""Test creating table with empty list data and schema
@@ -3762,8 +3444,8 @@ def test_create_table_empty_list_no_schema_error(mem_db: DBConnection):
mem_db.create_table("test_empty_no_schema", data=[])
def test_create_table_without_data_with_vector_schema(tmp_path):
"""Test exact scenario from issue #1968.
def test_add_table_with_empty_embeddings(tmp_path):
"""Test exact scenario from issue #1968
Regression test for issue #1968:
https://github.com/lancedb/lancedb/issues/1968
@@ -3775,9 +3457,6 @@ def test_create_table_without_data_with_vector_schema(tmp_path):
embedding: Vector(16)
table = db.create_table("test", schema=MySchema)
assert table.count_rows() == 0
assert table.schema == MySchema.to_arrow_schema()
table.add(
[{"text": "bar", "embedding": [0.1] * 16}],
on_bad_vectors="drop",
@@ -75,22 +75,6 @@ class TestVoyageAIModelRegistration:
with pytest.raises(ValueError, match="not supported"):
func.ndims()
def test_voyage3_source_embeddings_use_text_api(self, mock_voyageai_client):
"""Regression test for text table data being sent to the multimodal API."""
mock_voyageai_client.tokenize.return_value = [["hello", "world"]]
mock_voyageai_client.embed.return_value.embeddings = [[0.1] * 1024]
registry = get_registry()
func = registry.get("voyageai").create(name="voyage-3")
embeddings = func.compute_source_embeddings("hello world")
assert embeddings == [[0.1] * 1024]
mock_voyageai_client.embed.assert_called_once_with(
texts=["hello world"], model="voyage-3", input_type="document"
)
mock_voyageai_client.multimodal_embed.assert_not_called()
@pytest.mark.parametrize(
"model_name",
[
-15
View File
@@ -1,15 +0,0 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
from typing import assert_type
import lancedb
from lancedb import AsyncConnection, DBConnection
def check_connect_type() -> None:
assert_type(lancedb.connect("memory://"), DBConnection)
async def check_connect_async_type() -> None:
assert_type(await lancedb.connect_async("memory://"), AsyncConnection)
+14 -147
View File
@@ -28,72 +28,11 @@ use pyo3::{
Bound, FromPyObject, Py, PyAny, PyRef, PyResult, Python,
exceptions::{PyRuntimeError, PyValueError},
pyclass, pyfunction, pymethods,
types::{IntoPyDict, PyAnyMethods, PyBytes, PyDict, PyDictMethods, PyList, PyListMethods},
types::{IntoPyDict, PyAnyMethods, PyBytes, PyDict, PyDictMethods},
};
mod scannable;
/// Convert `LsmStats` to a Python dict, preserving the per-bucket list.
///
/// Deliberately not flattened to a table-level summary: a table is N
/// buckets on one node, and the per-bucket detail is the reason the
/// endpoint exists — flattening hides the single hot bucket someone opened
/// it to find.
fn lsm_stats_to_py(py: Python<'_>, stats: &lancedb::table::LsmStats) -> PyResult<Py<PyDict>> {
let out = PyDict::new(py);
let buckets = PyList::empty(py);
for b in &stats.buckets {
let e = PyDict::new(py);
e.set_item("shard_id", &b.shard_id)?;
e.set_item("status", &b.status)?;
e.set_item("writer_epoch", b.writer_epoch)?;
e.set_item("manifest_version", b.manifest_version)?;
e.set_item("current_generation", b.current_generation)?;
e.set_item(
"replay_after_wal_entry_position",
b.replay_after_wal_entry_position,
)?;
e.set_item(
"wal_entry_position_last_seen",
b.wal_entry_position_last_seen,
)?;
let generations = PyList::empty(py);
for g in &b.generations {
let ge = PyDict::new(py);
ge.set_item("generation", g.generation)?;
ge.set_item("bytes", g.bytes)?;
ge.set_item("rows", g.rows)?;
generations.append(ge)?;
}
e.set_item("generations", generations)?;
e.set_item("compacting", b.compacting)?;
e.set_item(
"memtables",
b.memtables
.as_ref()
.map(|ms| {
let l = PyList::empty(py);
for m in ms {
let d = PyDict::new(py);
d.set_item("generation", m.generation)?;
d.set_item("rows", m.rows)?;
d.set_item("bytes", m.bytes)?;
d.set_item("batches", m.batches)?;
d.set_item("indexes", m.indexes.clone())?;
l.append(d)?;
}
PyResult::Ok(l.unbind())
})
.transpose()?,
)?;
buckets.append(e)?;
}
out.set_item("buckets", buckets)?;
Ok(out.unbind())
}
#[derive(FromPyObject)]
enum PredicateArg {
Expr(PyExpr),
@@ -487,11 +426,9 @@ pub struct PyBlobFile {
impl PyBlobFile {
fn read_bytes(self_: PyRef<'_, Self>) -> PyResult<Py<PyBytes>> {
let inner = self_.inner.clone();
let py = self_.py();
let bytes = py
.detach(move || block_on(async move { inner.read().await }))
let bytes = block_on(async move { inner.read().await })
.map_err(|e| PyRuntimeError::new_err(format!("blob read failed: {e}")))?;
Ok(PyBytes::new(py, bytes.as_ref()).unbind())
Ok(PyBytes::new(self_.py(), bytes.as_ref()).unbind())
}
pub fn read(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
@@ -507,32 +444,24 @@ impl PyBlobFile {
fn close(self_: PyRef<'_, Self>) -> PyResult<()> {
let inner = self_.inner.clone();
self_
.py()
.detach(move || block_on(async move { inner.close().await }))
block_on(async move { inner.close().await })
.map_err(|e| PyRuntimeError::new_err(format!("blob close failed: {e}")))
}
fn is_closed(self_: PyRef<'_, Self>) -> bool {
let inner = self_.inner.clone();
self_
.py()
.detach(move || block_on(async move { inner.is_closed().await }))
block_on(async move { inner.is_closed().await })
}
fn seek(self_: PyRef<'_, Self>, position: u64) -> PyResult<()> {
let inner = self_.inner.clone();
self_
.py()
.detach(move || block_on(async move { inner.seek(position).await }))
block_on(async move { inner.seek(position).await })
.map_err(|e| PyRuntimeError::new_err(format!("blob seek failed: {e}")))
}
fn tell(self_: PyRef<'_, Self>) -> PyResult<u64> {
let inner = self_.inner.clone();
self_
.py()
.detach(move || block_on(async move { inner.tell().await }))
block_on(async move { inner.tell().await })
.map_err(|e| PyRuntimeError::new_err(format!("blob tell failed: {e}")))
}
@@ -546,20 +475,16 @@ impl PyBlobFile {
.checked_add(length as u64)
.ok_or_else(|| PyValueError::new_err("offset + length overflowed"))?;
let inner = self_.inner.clone();
let py = self_.py();
let bytes = py
.detach(move || block_on(async move { inner.read_range(offset..end).await }))
let bytes = block_on(async move { inner.read_range(offset..end).await })
.map_err(|e| PyRuntimeError::new_err(format!("blob read_range failed: {e}")))?;
Ok(PyBytes::new(py, bytes.as_ref()).unbind())
Ok(PyBytes::new(self_.py(), bytes.as_ref()).unbind())
}
fn read_up_to(self_: PyRef<'_, Self>, length: usize) -> PyResult<Py<PyBytes>> {
let inner = self_.inner.clone();
let py = self_.py();
let bytes = py
.detach(move || block_on(async move { inner.read_up_to(length).await }))
.map_err(|e| PyRuntimeError::new_err(format!("blob read_up_to failed: {e}")))?;
Ok(PyBytes::new(py, bytes.as_ref()).unbind())
let bytes = block_on(async move { inner.read_up_to(length).await })
.map_err(|e| PyRuntimeError::new_err(format!("blob read failed: {e}")))?;
Ok(PyBytes::new(self_.py(), bytes.as_ref()).unbind())
}
}
@@ -806,9 +731,6 @@ impl Table {
#[allow(private_interfaces)]
pub fn delete(self_: PyRef<'_, Self>, condition: PredicateArg) -> PyResult<Bound<'_, PyAny>> {
// Do not hold the Python borrow across the await. The cloned Rust table
// handle is thread-safe and allows deletes on the same Python table to
// run concurrently without PyO3 reporting "Already borrowed".
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
let result = match &condition {
@@ -1400,51 +1322,6 @@ impl Table {
})
}
/// Converge the table's LSM write path into its base table.
///
/// Best-effort: with writes flowing, new rows may land after the last
/// pass. Errors if the table stops making progress.
pub fn checkpoint_lsm(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
inner.checkpoint_lsm().await.infer_error()
})
}
/// Seal every bucket's active memtable into L0.
pub fn flush_lsm(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.inner_ref()?.clone();
future_into_py(
self_.py(),
async move { inner.flush_lsm().await.infer_error() },
)
}
/// Trigger a background L0 → base pass per bucket. Returns once the
/// passes are dispatched, not once they finish — watch `get_lsm_stats`.
pub fn compact_lsm(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
inner.compact_lsm().await.infer_error()
})
}
/// Live LSM state, or `None` when the LSM write path is not enabled.
#[pyo3(signature = (include_generation_rows=false))]
pub fn get_lsm_stats(
self_: PyRef<'_, Self>,
include_generation_rows: bool,
) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
let stats = inner
.get_lsm_stats(include_generation_rows)
.await
.infer_error()?;
Python::attach(|py| stats.map(|s| lsm_stats_to_py(py, &s)).transpose())
})
}
pub fn close_lsm_writers(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
@@ -1484,12 +1361,7 @@ impl Table {
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
let result = inner
.add_columns()
.transform(definitions)
.execute()
.await
.infer_error()?;
let result = inner.add_columns(definitions, None).await.infer_error()?;
Ok(AddColumnsResult::from(result))
})
}
@@ -1503,12 +1375,7 @@ impl Table {
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
let result = inner
.add_columns()
.transform(transform)
.execute()
.await
.infer_error()?;
let result = inner.add_columns(transform, None).await.infer_error()?;
Ok(AddColumnsResult::from(result))
})
}
+96 -96
View File
@@ -799,7 +799,7 @@ name = "cuda-bindings"
version = "13.3.1"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "cuda-pathfinder" },
{ name = "cuda-pathfinder", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/a9/21/8464d133752951c154feafb3b65c297e7d80f301183d220bec4c830f1441/cuda_bindings-13.3.1-cp310-cp310-manylinux_2_24_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:120fcc53d57903df529c3486962c56528cba5b7d6c57c99537320ed9922c8b86", size = 6073403, upload-time = "2026-05-29T23:11:36.22Z" },
@@ -834,37 +834,37 @@ wheels = [
[package.optional-dependencies]
cublas = [
{ name = "nvidia-cublas" },
{ name = "nvidia-cublas", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
]
cudart = [
{ name = "nvidia-cuda-runtime" },
{ name = "nvidia-cuda-runtime", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
]
cufft = [
{ name = "nvidia-cufft" },
{ name = "nvidia-cufft", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
]
cufile = [
{ name = "nvidia-cufile" },
{ name = "nvidia-cufile", marker = "sys_platform == 'linux'" },
]
cupti = [
{ name = "nvidia-cuda-cupti" },
{ name = "nvidia-cuda-cupti", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
]
curand = [
{ name = "nvidia-curand" },
{ name = "nvidia-curand", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
]
cusolver = [
{ name = "nvidia-cusolver" },
{ name = "nvidia-cusolver", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
]
cusparse = [
{ name = "nvidia-cusparse" },
{ name = "nvidia-cusparse", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
]
nvjitlink = [
{ name = "nvidia-nvjitlink" },
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
]
nvrtc = [
{ name = "nvidia-cuda-nvrtc" },
{ name = "nvidia-cuda-nvrtc", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
]
nvtx = [
{ name = "nvidia-nvtx" },
{ name = "nvidia-nvtx", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
]
[[package]]
@@ -1023,7 +1023,7 @@ name = "exceptiongroup"
version = "1.3.1"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "typing-extensions" },
{ name = "typing-extensions", marker = "python_full_version < '3.11'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/50/79/66800aadf48771f6b62f7eb014e352e5d06856655206165d775e675a02c9/exceptiongroup-1.3.1.tar.gz", hash = "sha256:8b412432c6055b0b7d14c310000ae93352ed6754f70fa8f7c34141f91c4e3219", size = 30371, upload-time = "2025-11-21T23:01:54.787Z" }
wheels = [
@@ -1440,16 +1440,16 @@ resolution-markers = [
"python_full_version < '3.11'",
]
dependencies = [
{ name = "cachetools" },
{ name = "certifi" },
{ name = "httpx" },
{ name = "ibm-cos-sdk" },
{ name = "lomond" },
{ name = "packaging" },
{ name = "pandas", version = "2.2.3", source = { registry = "https://pypi.org/simple" } },
{ name = "requests" },
{ name = "tabulate" },
{ name = "urllib3" },
{ name = "cachetools", marker = "python_full_version < '3.11'" },
{ name = "certifi", marker = "python_full_version < '3.11'" },
{ name = "httpx", marker = "python_full_version < '3.11'" },
{ name = "ibm-cos-sdk", marker = "python_full_version < '3.11'" },
{ name = "lomond", marker = "python_full_version < '3.11'" },
{ name = "packaging", marker = "python_full_version < '3.11'" },
{ name = "pandas", version = "2.2.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
{ name = "requests", marker = "python_full_version < '3.11'" },
{ name = "tabulate", marker = "python_full_version < '3.11'" },
{ name = "urllib3", marker = "python_full_version < '3.11'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/c7/56/2e3df38a1f13062095d7bde23c87a92f3898982993a15186b1bfecbd206f/ibm_watsonx_ai-1.3.42.tar.gz", hash = "sha256:ee5be59009004245d957ce97d1227355516df95a2640189749487614fef674ff", size = 688651, upload-time = "2025-10-01T13:35:41.527Z" }
wheels = [
@@ -1468,17 +1468,17 @@ resolution-markers = [
"python_full_version == '3.11.*'",
]
dependencies = [
{ name = "cachetools" },
{ name = "certifi" },
{ name = "httpx" },
{ name = "ibm-cos-sdk" },
{ name = "lomond" },
{ name = "packaging" },
{ name = "pandas", version = "2.3.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.14'" },
{ name = "cachetools", marker = "python_full_version >= '3.11'" },
{ name = "certifi", marker = "python_full_version >= '3.11'" },
{ name = "httpx", marker = "python_full_version >= '3.11'" },
{ name = "ibm-cos-sdk", marker = "python_full_version >= '3.11'" },
{ name = "lomond", marker = "python_full_version >= '3.11'" },
{ name = "packaging", marker = "python_full_version >= '3.11'" },
{ name = "pandas", version = "2.3.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
{ name = "pandas", version = "3.0.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.14'" },
{ name = "requests" },
{ name = "tabulate" },
{ name = "urllib3" },
{ name = "requests", marker = "python_full_version >= '3.11'" },
{ name = "tabulate", marker = "python_full_version >= '3.11'" },
{ name = "urllib3", marker = "python_full_version >= '3.11'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/29/a3/c756b534696ab2f3f29882fdb7ca7198b7a5c94e10c0a3a327853d6d6b79/ibm_watsonx_ai-1.5.14.tar.gz", hash = "sha256:a756488bd57e87c0fc51be42dcba871143cfe0ac1e805c497c5047e1e4f13e9d", size = 735804, upload-time = "2026-06-22T12:32:43.85Z" }
wheels = [
@@ -1554,17 +1554,17 @@ resolution-markers = [
"python_full_version < '3.11'",
]
dependencies = [
{ name = "colorama", marker = "sys_platform == 'win32'" },
{ name = "decorator" },
{ name = "exceptiongroup" },
{ name = "jedi" },
{ name = "matplotlib-inline" },
{ name = "pexpect", marker = "sys_platform != 'emscripten' and sys_platform != 'win32'" },
{ name = "prompt-toolkit" },
{ name = "pygments" },
{ name = "stack-data" },
{ name = "traitlets" },
{ name = "typing-extensions" },
{ name = "colorama", marker = "python_full_version < '3.11' and sys_platform == 'win32'" },
{ name = "decorator", marker = "python_full_version < '3.11'" },
{ name = "exceptiongroup", marker = "python_full_version < '3.11'" },
{ name = "jedi", marker = "python_full_version < '3.11'" },
{ name = "matplotlib-inline", marker = "python_full_version < '3.11'" },
{ name = "pexpect", marker = "python_full_version < '3.11' and sys_platform != 'emscripten' and sys_platform != 'win32'" },
{ name = "prompt-toolkit", marker = "python_full_version < '3.11'" },
{ name = "pygments", marker = "python_full_version < '3.11'" },
{ name = "stack-data", marker = "python_full_version < '3.11'" },
{ name = "traitlets", marker = "python_full_version < '3.11'" },
{ name = "typing-extensions", marker = "python_full_version < '3.11'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/40/18/f8598d287006885e7136451fdea0755af4ebcbfe342836f24deefaed1164/ipython-8.39.0.tar.gz", hash = "sha256:4110ae96012c379b8b6db898a07e186c40a2a1ef5d57a7fa83166047d9da7624", size = 5513971, upload-time = "2026-03-27T10:02:13.94Z" }
wheels = [
@@ -1583,18 +1583,18 @@ resolution-markers = [
"python_full_version == '3.11.*'",
]
dependencies = [
{ name = "colorama", marker = "sys_platform == 'win32'" },
{ name = "decorator" },
{ name = "ipython-pygments-lexers" },
{ name = "jedi" },
{ name = "matplotlib-inline" },
{ name = "pexpect", marker = "sys_platform != 'emscripten' and sys_platform != 'win32'" },
{ name = "prompt-toolkit" },
{ name = "psutil", marker = "sys_platform != 'cygwin' and sys_platform != 'emscripten'" },
{ name = "pygments" },
{ name = "stack-data" },
{ name = "traitlets" },
{ name = "typing-extensions", marker = "python_full_version < '3.12'" },
{ name = "colorama", marker = "python_full_version >= '3.11' and sys_platform == 'win32'" },
{ name = "decorator", marker = "python_full_version >= '3.11'" },
{ name = "ipython-pygments-lexers", marker = "python_full_version >= '3.11'" },
{ name = "jedi", marker = "python_full_version >= '3.11'" },
{ name = "matplotlib-inline", marker = "python_full_version >= '3.11'" },
{ name = "pexpect", marker = "python_full_version >= '3.11' and sys_platform != 'emscripten' and sys_platform != 'win32'" },
{ name = "prompt-toolkit", marker = "python_full_version >= '3.11'" },
{ name = "psutil", marker = "python_full_version >= '3.11' and sys_platform != 'cygwin' and sys_platform != 'emscripten'" },
{ name = "pygments", marker = "python_full_version >= '3.11'" },
{ name = "stack-data", marker = "python_full_version >= '3.11'" },
{ name = "traitlets", marker = "python_full_version >= '3.11'" },
{ name = "typing-extensions", marker = "python_full_version == '3.11.*'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/53/59/165d3b4d75cc34add3122c4417ecb229085140ac573103c223cd01dde96f/ipython-9.15.0.tar.gz", hash = "sha256:da2819ce2aa83135257df830660b1176d986c3d2876db24df01974fa955b2756", size = 4442580, upload-time = "2026-06-26T11:03:35.913Z" }
wheels = [
@@ -1606,7 +1606,7 @@ name = "ipython-pygments-lexers"
version = "1.1.1"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "pygments" },
{ name = "pygments", marker = "python_full_version >= '3.11'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/ef/4c/5dd1d8af08107f88c7f741ead7a40854b8ac24ddf9ae850afbcf698aa552/ipython_pygments_lexers-1.1.1.tar.gz", hash = "sha256:09c0138009e56b6854f9535736f4171d855c8c08a563a0dcd8022f78355c7e81", size = 8393, upload-time = "2025-01-17T11:24:34.505Z" }
wheels = [
@@ -1998,14 +1998,14 @@ requires-dist = [
{ name = "pillow", marker = "extra == 'clip'", specifier = ">=12.1.1" },
{ name = "pillow", marker = "extra == 'embeddings'", specifier = ">=12.1.1" },
{ name = "pillow", marker = "extra == 'siglip'", specifier = ">=12.1.1" },
{ name = "polars", marker = "extra == 'tests'", specifier = ">=0.19,<=1.32.3" },
{ name = "polars", marker = "extra == 'tests'", specifier = ">=0.19,<=1.3.0" },
{ name = "pre-commit", marker = "extra == 'dev'", specifier = ">=3.5.0" },
{ name = "pyarrow", specifier = ">=16" },
{ name = "pyarrow", marker = "extra == 'tests'", specifier = "<25" },
{ name = "pyarrow-stubs", marker = "extra == 'tests'", specifier = ">=16.0" },
{ name = "pydantic", specifier = ">=1.10" },
{ name = "pylance", marker = "extra == 'pylance'", specifier = ">=5.0.0b5" },
{ name = "pylance", marker = "extra == 'tests'", specifier = "==10.0.0" },
{ name = "pylance", marker = "extra == 'tests'", specifier = "==9.0.0rc1" },
{ name = "pyright", marker = "extra == 'dev'", specifier = ">=1.1.350" },
{ name = "pytest", marker = "extra == 'tests'", specifier = ">=7.0" },
{ name = "pytest-asyncio", marker = "extra == 'tests'", specifier = ">=0.21" },
@@ -2858,7 +2858,7 @@ name = "nvidia-cudnn-cu13"
version = "9.19.0.56"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "nvidia-cublas" },
{ name = "nvidia-cublas", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/f1/84/26025437c1e6b61a707442184fa0c03d083b661adf3a3eecfd6d21677740/nvidia_cudnn_cu13-9.19.0.56-py3-none-manylinux_2_27_aarch64.whl", hash = "sha256:6ed29ffaee1176c612daf442e4dd6cfeb6a0caa43ddcbeb59da94953030b1be4", size = 433781201, upload-time = "2026-02-03T20:40:53.805Z" },
@@ -2870,7 +2870,7 @@ name = "nvidia-cufft"
version = "12.0.0.61"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "nvidia-nvjitlink" },
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/8b/ae/f417a75c0259e85c1d2f83ca4e960289a5f814ed0cea74d18c353d3e989d/nvidia_cufft-12.0.0.61-py3-none-manylinux2014_aarch64.manylinux_2_17_aarch64.whl", hash = "sha256:2708c852ef8cd89d1d2068bdbece0aa188813a0c934db3779b9b1faa8442e5f5", size = 214053554, upload-time = "2025-09-04T08:31:38.196Z" },
@@ -2900,9 +2900,9 @@ name = "nvidia-cusolver"
version = "12.0.4.66"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "nvidia-cublas" },
{ name = "nvidia-cusparse" },
{ name = "nvidia-nvjitlink" },
{ name = "nvidia-cublas", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
{ name = "nvidia-cusparse", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/c8/c3/b30c9e935fc01e3da443ec0116ed1b2a009bb867f5324d3f2d7e533e776b/nvidia_cusolver-12.0.4.66-py3-none-manylinux_2_27_aarch64.whl", hash = "sha256:02c2457eaa9e39de20f880f4bd8820e6a1cfb9f9a34f820eb12a155aa5bc92d2", size = 223467760, upload-time = "2025-09-04T08:33:04.222Z" },
@@ -2914,7 +2914,7 @@ name = "nvidia-cusparse"
version = "12.6.3.3"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "nvidia-nvjitlink" },
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/f8/94/5c26f33738ae35276672f12615a64bd008ed5be6d1ebcb23579285d960a9/nvidia_cusparse-12.6.3.3-py3-none-manylinux2014_aarch64.manylinux_2_17_aarch64.whl", hash = "sha256:80bcc4662f23f1054ee334a15c72b8940402975e0eab63178fc7e670aa59472c", size = 162155568, upload-time = "2025-09-04T08:33:42.864Z" },
@@ -3091,10 +3091,10 @@ resolution-markers = [
"python_full_version < '3.11'",
]
dependencies = [
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } },
{ name = "python-dateutil" },
{ name = "pytz" },
{ name = "tzdata" },
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
{ name = "python-dateutil", marker = "python_full_version < '3.11'" },
{ name = "pytz", marker = "python_full_version < '3.11'" },
{ name = "tzdata", marker = "python_full_version < '3.11'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/9c/d6/9f8431bacc2e19dca897724cd097b1bb224a6ad5433784a44b587c7c13af/pandas-2.2.3.tar.gz", hash = "sha256:4f18ba62b61d7e192368b84517265a99b4d7ee8912f8708660fb4a366cc82667", size = 4399213, upload-time = "2024-09-20T13:10:04.827Z" }
wheels = [
@@ -3143,11 +3143,11 @@ resolution-markers = [
"python_full_version == '3.11.*'",
]
dependencies = [
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12' or python_full_version >= '3.14'" },
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12' and python_full_version < '3.14'" },
{ name = "python-dateutil" },
{ name = "pytz" },
{ name = "tzdata" },
{ name = "python-dateutil", marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
{ name = "pytz", marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
{ name = "tzdata", marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/33/01/d40b85317f86cf08d853a4f495195c73815fdf205eef3993821720274518/pandas-2.3.3.tar.gz", hash = "sha256:e05e1af93b977f7eafa636d043f9f94c7ee3ac81af99c13508215942e64c993b", size = 4495223, upload-time = "2025-09-29T23:34:51.853Z" }
wheels = [
@@ -3210,9 +3210,9 @@ resolution-markers = [
"python_full_version >= '3.14' and sys_platform != 'emscripten' and sys_platform != 'win32'",
]
dependencies = [
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" } },
{ name = "python-dateutil" },
{ name = "tzdata", marker = "sys_platform == 'emscripten' or sys_platform == 'win32'" },
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.14'" },
{ name = "python-dateutil", marker = "python_full_version >= '3.14'" },
{ name = "tzdata", marker = "(python_full_version >= '3.14' and sys_platform == 'emscripten') or (python_full_version >= '3.14' and sys_platform == 'win32')" },
]
sdist = { url = "https://files.pythonhosted.org/packages/f8/87/4341c6252d1c47b08768c3d25ac487362bf403f0313ddae4a2a26c9b1b4c/pandas-3.0.3.tar.gz", hash = "sha256:696a4a00a2a2a35d4e5deb3fc946641b96c944f02230e4f76137fe35d806c4fc", size = 4651414, upload-time = "2026-05-11T18:54:29.21Z" }
wheels = [
@@ -3320,7 +3320,7 @@ name = "pexpect"
version = "4.9.0"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "ptyprocess" },
{ name = "ptyprocess", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
]
sdist = { url = "https://files.pythonhosted.org/packages/42/92/cc564bf6381ff43ce1f4d06852fc19a2f11d180f23dc32d9588bee2f149d/pexpect-4.9.0.tar.gz", hash = "sha256:ee7d41123f3c9911050ea2c2dac107568dc43b2d3b0c7557a33212c398ead30f", size = 166450, upload-time = "2023-11-25T09:07:26.339Z" }
wheels = [
@@ -3912,7 +3912,7 @@ crypto = [
[[package]]
name = "pylance"
version = "10.0.0"
version = "7.0.0"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "lance-namespace" },
@@ -3922,12 +3922,12 @@ dependencies = [
{ name = "pyarrow" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/82/b2/c81de196076c4c8d768f485324a0043e113ce0950a339977d83fee9783ef/pylance-10.0.0-cp310-abi3-macosx_11_0_arm64.whl", hash = "sha256:d4bba56ae829202b7e9cdc82c92c00a4b52f03e679dc3edfad1db09d1285be2e", size = 69279797, upload-time = "2026-08-07T18:25:24.812Z" },
{ url = "https://files.pythonhosted.org/packages/d2/e9/af671a6225740bd70c1e70bd83fa08091d628b960397cbecc477f16aaaf1/pylance-10.0.0-cp310-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:489b944827c0271e16a62b4006f8c75acc3bd7fbc381f35388a2afaaae38d438", size = 72775384, upload-time = "2026-08-07T18:31:29.686Z" },
{ url = "https://files.pythonhosted.org/packages/5d/22/e07194195bb3bbdf062b0c31690fc92fccb2686d5e768efa1dc379a93350/pylance-10.0.0-cp310-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:018efe7d437d326b9049c1223458bd955e85e48509b4b0bbf8b2d3d94075dbff", size = 76632194, upload-time = "2026-08-07T18:44:29.001Z" },
{ url = "https://files.pythonhosted.org/packages/f2/71/ed9956cf657e86a5fee5d5ddcfa7bfed0d7c2968bee39ed85d12629d8c93/pylance-10.0.0-cp310-abi3-manylinux_2_28_aarch64.whl", hash = "sha256:254bad5d765c14db6c4eddd5323133dee3c76ff6d40c4777afe0b0e91d4c04d2", size = 72804222, upload-time = "2026-08-07T18:31:24.264Z" },
{ url = "https://files.pythonhosted.org/packages/94/96/de449c246b2892df9d5a0af775b79a061e9897e6942bcee5c4ab6d6071f8/pylance-10.0.0-cp310-abi3-manylinux_2_28_x86_64.whl", hash = "sha256:9f0c089e30389d9b7765a4c16f3d9325b05ed4fbf6cc2bbc9418292983472e54", size = 76602015, upload-time = "2026-08-07T18:46:27.724Z" },
{ url = "https://files.pythonhosted.org/packages/53/81/4a5a9072b6d68c4dbb8c7b5530a381c7a6bea4ba92f9eca58ac885142722/pylance-10.0.0-cp310-abi3-win_amd64.whl", hash = "sha256:9fecf46e6835dc64b2d71767d5eab4f73787543a852432e9ba8ecf8527b3cdab", size = 82799247, upload-time = "2026-08-07T18:47:46.855Z" },
{ url = "https://files.pythonhosted.org/packages/ac/ad/2f64921bf346e7075aef24a72595db44821724a3d89a9a92dd24e79632aa/pylance-7.0.0-cp39-abi3-macosx_11_0_arm64.whl", hash = "sha256:98422021975be76e72b1572f41b8c9abb3bee5bdc9bfa5e9ce731110a65ed4d1", size = 62134146, upload-time = "2026-05-27T21:59:37.459Z" },
{ url = "https://files.pythonhosted.org/packages/73/1c/c5a01bee0160b55d9a98895cbd33091d038f0a0995b121ab72e629008d02/pylance-7.0.0-cp39-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:4bec86ee5b6fbd8bfc493e653f0a1fba0303cfe5492b9b46fc25ab908edc7183", size = 65373684, upload-time = "2026-05-27T22:04:01.584Z" },
{ url = "https://files.pythonhosted.org/packages/eb/da/1fe8b8f7dbfe734d76af76acc994fc360a0d0c79a4874ef69f5a72a58fe3/pylance-7.0.0-cp39-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:881491432c53184e52f8d1db8d5f872f39a03f36fb104bec77b33d379519d8b5", size = 69458555, upload-time = "2026-05-27T22:16:50.567Z" },
{ url = "https://files.pythonhosted.org/packages/76/f0/dd505cf3fd0226ab9d94759acd713125af1d3bfacfd80bbd52e3b9f89509/pylance-7.0.0-cp39-abi3-manylinux_2_28_aarch64.whl", hash = "sha256:18453999e7fff4f76b16d6b7882c9df0628bd142ff95e2461bd7dd5ee3fe0af3", size = 65394430, upload-time = "2026-05-27T22:05:30.923Z" },
{ url = "https://files.pythonhosted.org/packages/17/ba/2357b81034f28eb00790e258ed140289a6a887a7468ca9df6349fd186b27/pylance-7.0.0-cp39-abi3-manylinux_2_28_x86_64.whl", hash = "sha256:04a58051d408c60fe76d41a220dcaf8fea8fb6d1aa0ca78a709b60bc3cc8d19a", size = 69473470, upload-time = "2026-05-27T22:17:18.935Z" },
{ url = "https://files.pythonhosted.org/packages/1f/ec/5c00b6303a67d787f9475141832cbdc513d674ac3dcaeef8a7b169905e65/pylance-7.0.0-cp39-abi3-win_amd64.whl", hash = "sha256:467d4864af047eaab4e1370e2f1e88e2c6f507c079874421116cb41d78bc3629", size = 74792863, upload-time = "2026-05-27T22:19:23.875Z" },
]
[[package]]
@@ -4683,10 +4683,10 @@ resolution-markers = [
"python_full_version < '3.11'",
]
dependencies = [
{ name = "joblib" },
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } },
{ name = "scipy", version = "1.15.3", source = { registry = "https://pypi.org/simple" } },
{ name = "threadpoolctl" },
{ name = "joblib", marker = "python_full_version < '3.11'" },
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
{ name = "scipy", version = "1.15.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
{ name = "threadpoolctl", marker = "python_full_version < '3.11'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/98/c2/a7855e41c9d285dfe86dc50b250978105dce513d6e459ea66a6aeb0e1e0c/scikit_learn-1.7.2.tar.gz", hash = "sha256:20e9e49ecd130598f1ca38a1d85090e1a600147b9c02fa6f15d69cb53d968fda", size = 7193136, upload-time = "2025-09-09T08:21:29.075Z" }
wheels = [
@@ -4734,13 +4734,13 @@ resolution-markers = [
"python_full_version == '3.11.*'",
]
dependencies = [
{ name = "joblib" },
{ name = "narwhals" },
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" },
{ name = "joblib", marker = "python_full_version >= '3.11'" },
{ name = "narwhals", marker = "python_full_version >= '3.11'" },
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" },
{ name = "scipy", version = "1.17.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" },
{ name = "scipy", version = "1.17.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
{ name = "scipy", version = "1.18.0", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" },
{ name = "threadpoolctl" },
{ name = "threadpoolctl", marker = "python_full_version >= '3.11'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/fa/6f/37092bdb25f712817231799fc5674d8e704066a8a70c1d2d40517e18b4ab/scikit_learn-1.9.0.tar.gz", hash = "sha256:8833266989d3a5110178a9fae30783675460724d0e1efb13b14901d2c660c557", size = 7750767, upload-time = "2026-06-02T11:54:32.706Z" }
wheels = [
@@ -4784,7 +4784,7 @@ resolution-markers = [
"python_full_version < '3.11'",
]
dependencies = [
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } },
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/0f/37/6964b830433e654ec7485e45a00fc9a27cf868d622838f6b6d9c5ec0d532/scipy-1.15.3.tar.gz", hash = "sha256:eae3cf522bc7df64b42cad3925c876e1b0b6c35c1337c93e12c0f366f55b0eaf", size = 59419214, upload-time = "2025-05-08T16:13:05.955Z" }
wheels = [
@@ -4843,7 +4843,7 @@ resolution-markers = [
"python_full_version == '3.11.*'",
]
dependencies = [
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" } },
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/7a/97/5a3609c4f8d58b039179648e62dd220f89864f56f7357f5d4f45c29eb2cc/scipy-1.17.1.tar.gz", hash = "sha256:95d8e012d8cb8816c226aef832200b1d45109ed4464303e997c5b13122b297c0", size = 30573822, upload-time = "2026-02-23T00:26:24.851Z" }
wheels = [
@@ -4920,7 +4920,7 @@ resolution-markers = [
"python_full_version >= '3.12' and python_full_version < '3.14'",
]
dependencies = [
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" } },
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/a7/25/c2700dfaf6442b4effaa91af24ebce5dc9d31bb4a69706313aae70d72cd0/scipy-1.18.0.tar.gz", hash = "sha256:67b2ad2ad54c72ca6d04975a9b2df8c3638c34ddd5b28738e94fc2b57929d378", size = 30774447, upload-time = "2026-06-19T15:01:43.456Z" }
wheels = [
+5 -5
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb"
version = "0.37.1"
version = "0.37.1-beta.0"
edition.workspace = true
description = "LanceDB: A serverless, low-latency vector database for AI applications"
license.workspace = true
@@ -49,6 +49,8 @@ lance-namespace = { workspace = true }
lance-namespace-impls = { workspace = true }
metrics = { workspace = true, optional = true }
metrics-util = { workspace = true, optional = true }
# Keep the direct dependency aligned with the version required by OpenDAL.
goosefs-sdk = { version = "=0.1.8", optional = true }
moka = { workspace = true }
pin-project = { workspace = true }
tokio = { version = "1.23", features = ["rt-multi-thread", "sync"] }
@@ -73,8 +75,6 @@ reqwest = { version = "0.12.0", default-features = false, features = [
"http2",
"json",
"macos-system-configuration",
# Avoid linking OpenSSL into Python wheels, which breaks on FIPS hosts.
"rustls-tls-native-roots",
"stream",
], optional = true }
http = { version = "1", optional = true } # Matching what is in reqwest
@@ -98,8 +98,7 @@ anyhow = "1"
lance-testing = { workspace = true }
tempfile = "3.5.0"
random_word = { version = "0.4.3", features = ["en"] }
roaring = "0.11.4"
tokio = { version = "1.23", features = ["io-util", "macros", "net", "rt-multi-thread", "sync", "test-util"] }
tokio = { version = "1.23", features = ["io-util", "macros", "net", "rt-multi-thread", "sync"] }
uuid = { version = "1.7.0", features = ["v4"] }
walkdir = "2"
aws-sdk-dynamodb = { version = "1.55.0" }
@@ -134,6 +133,7 @@ azure = [
]
cos = ["lance/tencent", "lance-io/tencent"]
goosefs = [
"dep:goosefs-sdk",
"lance/goosefs",
"lance-io/goosefs",
"lance-namespace-impls/dir-goosefs",
+3 -199
View File
@@ -9,7 +9,6 @@
//!
//! Blob tables require Lance file format >= 2.2 and stable row ids at create.
use std::ops::Range;
use std::sync::Arc;
use arrow_array::LargeBinaryArray;
@@ -17,203 +16,11 @@ use arrow_array::builder::LargeBinaryBuilder;
use arrow_schema::{DataType, Field, Schema};
use lance::dataset::{BlobRangeRequest as LanceBlobRangeRequest, Dataset, WriteParams};
use lance_arrow::FieldExt;
use lance_file::version::LanceFileVersion;
use lance_io::object_store::ObjectStore;
use object_store::path::Path;
use lance_encoding::version::LanceFileVersion;
use crate::error::{Error, Result};
/// Seekable handle for one blob value, backed by local storage or a remote
/// HTTP byte-range endpoint.
#[derive(Debug)]
pub struct BlobFile {
inner: BlobFileInner,
}
#[derive(Debug)]
enum BlobFileInner {
Native(lance::dataset::BlobFile),
#[cfg(feature = "remote")]
Remote(Box<crate::remote::table::blobs::RemoteBlobFile>),
}
impl From<lance::dataset::BlobFile> for BlobFile {
fn from(value: lance::dataset::BlobFile) -> Self {
Self {
inner: BlobFileInner::Native(value),
}
}
}
#[cfg(feature = "remote")]
impl From<crate::remote::table::blobs::RemoteBlobFile> for BlobFile {
fn from(value: crate::remote::table::blobs::RemoteBlobFile) -> Self {
Self {
inner: BlobFileInner::Remote(Box::new(value)),
}
}
}
impl BlobFile {
/// Inline reader over a data-file slice.
pub fn new_inline(
object_store: Arc<ObjectStore>,
path: Path,
position: u64,
size: u64,
) -> Self {
lance::dataset::BlobFile::new_inline(object_store, path, position, size).into()
}
/// Dedicated sidecar-file reader.
pub fn new_dedicated(object_store: Arc<ObjectStore>, path: Path, size: u64) -> Self {
lance::dataset::BlobFile::new_dedicated(object_store, path, size).into()
}
/// Packed reader for a slice in a shared sidecar.
pub fn new_packed(
object_store: Arc<ObjectStore>,
path: Path,
position: u64,
size: u64,
) -> Self {
lance::dataset::BlobFile::new_packed(object_store, path, position, size).into()
}
/// External reader at a resolved object location.
pub fn new_external(
object_store: Arc<ObjectStore>,
path: Path,
uri: String,
position: u64,
size: u64,
) -> Self {
lance::dataset::BlobFile::new_external(object_store, path, uri, position, size).into()
}
/// Close the handle.
pub async fn close(&self) -> lance_core::Result<()> {
match &self.inner {
BlobFileInner::Native(file) => file.close().await,
#[cfg(feature = "remote")]
BlobFileInner::Remote(file) => file.close().await,
}
}
/// Whether the handle is closed.
pub async fn is_closed(&self) -> bool {
match &self.inner {
BlobFileInner::Native(file) => file.is_closed().await,
#[cfg(feature = "remote")]
BlobFileInner::Remote(file) => file.is_closed(),
}
}
/// Read a range without moving the cursor.
pub async fn read_range(&self, range: Range<u64>) -> lance_core::Result<bytes::Bytes> {
match &self.inner {
BlobFileInner::Native(file) => file.read_range(range).await,
#[cfg(feature = "remote")]
BlobFileInner::Remote(file) => file.read_range(range).await,
}
}
/// Read ranges without moving the cursor.
pub async fn read_ranges(
&self,
ranges: &[Range<u64>],
) -> lance_core::Result<Vec<bytes::Bytes>> {
match &self.inner {
BlobFileInner::Native(file) => file.read_ranges(ranges).await,
#[cfg(feature = "remote")]
BlobFileInner::Remote(file) => file.read_ranges(ranges).await,
}
}
/// Read from the cursor to the end.
pub async fn read(&self) -> lance_core::Result<bytes::Bytes> {
match &self.inner {
BlobFileInner::Native(file) => file.read().await,
#[cfg(feature = "remote")]
BlobFileInner::Remote(file) => file.read().await,
}
}
/// Read up to `len` bytes and advance the cursor.
pub async fn read_up_to(&self, len: usize) -> lance_core::Result<bytes::Bytes> {
match &self.inner {
BlobFileInner::Native(file) => file.read_up_to(len).await,
#[cfg(feature = "remote")]
BlobFileInner::Remote(file) => file.read_up_to(len).await,
}
}
/// Move the cursor to `new_cursor`.
pub async fn seek(&self, new_cursor: u64) -> lance_core::Result<()> {
match &self.inner {
BlobFileInner::Native(file) => file.seek(new_cursor).await,
#[cfg(feature = "remote")]
BlobFileInner::Remote(file) => file.seek(new_cursor).await,
}
}
/// Current cursor position.
pub async fn tell(&self) -> lance_core::Result<u64> {
match &self.inner {
BlobFileInner::Native(file) => file.tell().await,
#[cfg(feature = "remote")]
BlobFileInner::Remote(file) => file.tell().await,
}
}
/// Blob length in bytes.
pub fn size(&self) -> u64 {
match &self.inner {
BlobFileInner::Native(file) => file.size(),
#[cfg(feature = "remote")]
BlobFileInner::Remote(file) => file.size(),
}
}
/// Physical byte offset in the data file. `None` on remote handles. The
/// Cloud byte-range route does not expose storage layout.
pub fn position(&self) -> Option<u64> {
match &self.inner {
BlobFileInner::Native(file) => Some(file.position()),
#[cfg(feature = "remote")]
BlobFileInner::Remote(_) => None,
}
}
/// Path of the data file holding the blob. `None` on remote handles. The
/// Cloud byte-range route does not expose storage layout.
pub fn data_path(&self) -> Option<&Path> {
match &self.inner {
BlobFileInner::Native(file) => Some(file.data_path()),
#[cfg(feature = "remote")]
BlobFileInner::Remote(_) => None,
}
}
/// Native storage layout. `None` on remote handles. The Cloud byte-range
/// route does not expose layout.
pub fn kind(&self) -> Option<lance_core::datatypes::BlobKind> {
match &self.inner {
BlobFileInner::Native(file) => Some(file.kind()),
#[cfg(feature = "remote")]
BlobFileInner::Remote(_) => None,
}
}
/// External URI for native handles. Remote handles do not expose storage URIs.
pub fn uri(&self) -> Option<&str> {
match &self.inner {
BlobFileInner::Native(file) => file.uri(),
#[cfg(feature = "remote")]
BlobFileInner::Remote(_) => None,
}
}
}
pub use lance::dataset::BlobFile;
/// One row-specific blob range read request.
///
@@ -457,10 +264,7 @@ pub(crate) async fn take_blob_files_aligned(
let handles = dataset.take_blobs(row_ids, column).await?;
ensure_all_row_ids_resolved(column, row_ids.len(), handles.len())?;
Ok(handles
.into_iter()
.map(|handle| handle.map(Into::into))
.collect())
Ok(handles)
}
#[cfg(test)]
+1 -1
View File
@@ -34,7 +34,7 @@ use crate::remote::{
db::{OPT_REMOTE_API_KEY, OPT_REMOTE_HOST_OVERRIDE, OPT_REMOTE_REGION},
};
use lance::io::ObjectStoreParams;
pub use lance_file::version::LanceFileVersion;
pub use lance_encoding::version::LanceFileVersion;
#[cfg(feature = "remote")]
use lance_io::object_store::StorageOptions;
use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider};
@@ -202,17 +202,6 @@ mod tests {
assert_eq!(table.count_rows(None).await.unwrap(), 0);
}
#[tokio::test]
async fn create_table_in_named_memory_database() {
let db = connect("memory://foo").execute().await.unwrap();
let batch = record_batch!(("id", Int64, [1, 2, 3])).unwrap();
let table = db.create_table("my_table", batch).execute().await.unwrap();
assert_eq!(table.uri().await.unwrap(), "memory://foo/my_table.lance");
assert_eq!(table.count_rows(None).await.unwrap(), 3);
}
async fn test_create_table_with_data<T>(data: T)
where
T: Scannable + 'static,
+3 -167
View File
@@ -12,7 +12,7 @@ use lance::dataset::refs::Ref;
use lance::dataset::{ReadParams, WriteMode, builder::DatasetBuilder};
use lance::io::{ObjectStore, ObjectStoreParams, WrappingObjectStore};
use lance_datafusion::utils::StreamingWriteSource;
use lance_file::version::LanceFileVersion;
use lance_encoding::version::LanceFileVersion;
use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider};
use lance_table::io::commit::commit_handler_from_url;
use object_store::local::LocalFileSystem;
@@ -1294,11 +1294,9 @@ mod tests {
use crate::connection::ConnectRequest;
use crate::data::scannable::Scannable;
use crate::database::{CreateTableMode, CreateTableRequest};
use crate::query::QueryRequest;
use crate::table::{AnyQuery, WriteOptions};
use crate::table::WriteOptions;
use arrow_array::{Int32Array, RecordBatch, StringArray};
use arrow_schema::{DataType, Field, Schema};
use futures::TryStreamExt;
use std::path::PathBuf;
use tempfile::tempdir;
@@ -1378,156 +1376,6 @@ mod tests {
assert!(!tempdir.path().join("__manifest").exists());
}
/// Regression test for https://github.com/lancedb/lancedb/issues/1600.
///
/// Opening a table used to create a separate object-store client instead of
/// reusing the one that successfully connected to the database. Repeating
/// credential discovery made S3 table opens intermittent, especially in AWS
/// Lambda, and the failed open was reported as `TableNotFound`.
#[tokio::test]
async fn test_open_table_reuses_connection_object_store() {
let tempdir = tempdir().unwrap();
let uri = tempdir.path().to_str().unwrap();
let registry = Arc::new(lance_io::object_store::ObjectStoreRegistry::default());
let session = Arc::new(lance::session::Session::new(16, 16, registry.clone()));
let request = ConnectRequest {
uri: uri.to_string(),
#[cfg(feature = "remote")]
client_config: Default::default(),
options: Default::default(),
namespace_client_properties: Default::default(),
manifest_enabled: false,
read_consistency_interval: None,
session: Some(session),
};
let db = ListingDatabase::connect_with_options(&request)
.await
.unwrap();
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
db.create_table(CreateTableRequest {
name: "test".to_string(),
namespace_path: vec![],
data: Box::new(RecordBatch::new_empty(schema)) as Box<dyn Scannable>,
mode: CreateTableMode::Create,
write_options: Default::default(),
location: None,
namespace_client: None,
})
.await
.unwrap();
let before_open = registry.stats();
for _ in 0..3 {
let table = db
.open_table(OpenTableRequest {
name: "test".to_string(),
namespace_path: vec![],
index_cache_size: None,
lance_read_params: None,
location: None,
namespace_client: None,
managed_versioning: None,
})
.await
.unwrap();
assert_eq!(table.count_rows(None).await.unwrap(), 0);
}
let after_open = registry.stats();
assert_eq!(after_open.misses, before_open.misses);
assert!(after_open.hits >= before_open.hits + 3);
}
/// Regression test for https://github.com/lancedb/lancedb/issues/3197.
#[cfg(unix)]
#[tokio::test]
async fn test_open_table_follows_hugging_face_symlinks() {
let (tempdir, db) = setup_database().await;
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
db.create_table(CreateTableRequest {
name: "test".to_string(),
namespace_path: vec![],
data: Box::new(
RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(vec![1, 2, 3]))])
.unwrap(),
) as Box<dyn Scannable>,
mode: CreateTableMode::Create,
write_options: Default::default(),
location: None,
namespace_client: None,
})
.await
.unwrap();
let table_dir = tempdir.path().join("test.lance");
let versions_dir = table_dir.join("_versions");
let manifest_path = std::fs::read_dir(&versions_dir)
.unwrap()
.map(|entry| entry.unwrap().path())
.find(|path| path.extension().is_some_and(|ext| ext == "manifest"))
.unwrap();
let data_path = std::fs::read_dir(table_dir.join("data"))
.unwrap()
.map(|entry| entry.unwrap().path())
.find(|path| path.extension().is_some_and(|ext| ext == "lance"))
.unwrap();
// Hugging Face snapshots keep dataset objects in a separate blob directory and
// expose them through relative symlinks.
let blobs_dir = tempdir.path().join("blobs");
std::fs::create_dir(&blobs_dir).unwrap();
let manifest_blob = "9b603c63d0e692e05d58be25605f2f2064cc781e5ff94fe983a405059547b816";
let data_blob = "be64f20e5723bd0a27cfdbdb41cf7d6fad94cd572a71973b717fb8340f4310c5";
std::fs::rename(&manifest_path, blobs_dir.join(manifest_blob)).unwrap();
std::fs::rename(&data_path, blobs_dir.join(data_blob)).unwrap();
std::os::unix::fs::symlink(Path::new("../../blobs").join(manifest_blob), &manifest_path)
.unwrap();
std::os::unix::fs::symlink(Path::new("../../blobs").join(data_blob), &data_path).unwrap();
let symlink_len = std::fs::symlink_metadata(&manifest_path).unwrap().len();
let target_len = std::fs::metadata(&manifest_path).unwrap().len();
assert_ne!(symlink_len, target_len);
drop(db);
let db = ListingDatabase::connect_with_options(&ConnectRequest {
uri: tempdir.path().to_str().unwrap().to_string(),
#[cfg(feature = "remote")]
client_config: Default::default(),
options: Default::default(),
namespace_client_properties: Default::default(),
manifest_enabled: false,
read_consistency_interval: None,
session: None,
})
.await
.unwrap();
let table = db
.open_table(OpenTableRequest {
name: "test".to_string(),
namespace_path: vec![],
index_cache_size: None,
lance_read_params: None,
location: None,
namespace_client: None,
managed_versioning: None,
})
.await
.unwrap();
let batches = table
.query(
&AnyQuery::Query(QueryRequest::default()),
Default::default(),
)
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap();
assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 3);
}
#[tokio::test]
async fn test_clone_table_basic() {
let (_tempdir, db) = setup_database().await;
@@ -2432,7 +2280,7 @@ mod tests {
#[tokio::test]
async fn test_table_uri() {
let (_tempdir, mut db) = setup_database().await;
let (_tempdir, db) = setup_database().await;
let mut pb = PathBuf::new();
pb.push(db.uri.clone());
@@ -2441,18 +2289,6 @@ mod tests {
let expected = pb.to_str().unwrap();
let uri = db.table_uri("test").ok().unwrap();
assert_eq!(uri, expected);
// URI paths always use forward slashes, even on Windows. Using
// `Path::join` here used to produce `az://container/prefix\\test.lance`,
// which Azure treated as a different object from the table returned by
// `table_names` (https://github.com/lancedb/lancedb/issues/1072).
for base_uri in ["az://container/prefix", "az://container/prefix/"] {
db.uri = base_uri.to_string();
assert_eq!(
db.table_uri("test").unwrap(),
"az://container/prefix/test.lance"
);
}
}
/// Regression: connecting via a URL-style URI (which goes through
+2 -2
View File
@@ -201,7 +201,7 @@ impl LanceNamespaceDatabase {
&self,
request: &DbCreateTableRequest,
) -> Result<(
Option<lance_file::version::LanceFileVersion>,
Option<lance_encoding::version::LanceFileVersion>,
Option<bool>,
Option<bool>,
)> {
@@ -214,7 +214,7 @@ impl LanceNamespaceDatabase {
let storage_version_override = storage_options
.and_then(|opts| opts.get(OPT_NEW_TABLE_STORAGE_VERSION))
.map(|s| s.parse::<lance_file::version::LanceFileVersion>())
.map(|s| s.parse::<lance_encoding::version::LanceFileVersion>())
.transpose()?;
let v2_manifest_override = storage_options
@@ -12,7 +12,8 @@ use lance_encoding::decoder::{DecoderPlugins, FilterExpression};
use lance_file::{
reader::{FileReader, FileReaderOptions},
version::ConcreteFileVersion,
writer::{FileWriter, FileWriterOptions},
versions,
writer::FileWriterOptions,
};
use lance_io::{
ReadBatchParams,
@@ -153,13 +154,11 @@ impl Shuffler {
source: None,
})?;
let object_writer = object_store.create(&path).await?;
let writer = FileWriter::try_new(
let writer = versions::create_writer(
ConcreteFileVersion::V2_1,
object_writer,
schema.clone(),
FileWriterOptions {
format_version: Some(ConcreteFileVersion::V2_1.into()),
..Default::default()
},
FileWriterOptions::default(),
)?;
file_writers.push(writer);
}
-70
View File
@@ -169,12 +169,6 @@ impl From<DataFusionError> for Error {
impl From<lance::Error> for Error {
fn from(source: lance::Error) -> Self {
if has_unsupported_local_filesystem_source(&source) {
return Self::NotSupported {
message: "the filesystem does not support an operation required for safe Lance commits (such as atomic rename). Object-storage mounts such as Mountpoint for Amazon S3 are not supported; use the native object-store URI (for example, s3://bucket/path) instead".to_string(),
};
}
// Try to unwrap external errors that were wrapped by lance
match source {
lance::Error::Wrapped { error, .. } => Self::from_box_error(error),
@@ -187,27 +181,6 @@ impl From<lance::Error> for Error {
}
}
fn has_unsupported_local_filesystem_source(error: &(dyn std::error::Error + 'static)) -> bool {
let mut current = Some(error);
let mut is_local_filesystem = false;
let mut is_unsupported = false;
while let Some(error) = current {
is_local_filesystem |= error
.downcast_ref::<object_store::Error>()
.is_some_and(|error| {
matches!(error, object_store::Error::Generic { store, .. } if *store == "LocalFileSystem")
});
is_unsupported |= error
.downcast_ref::<std::io::Error>()
.is_some_and(|error| error.kind() == std::io::ErrorKind::Unsupported);
if is_local_filesystem && is_unsupported {
return true;
}
current = error.source();
}
false
}
impl Error {
fn from_box_error(mut source: Box<dyn std::error::Error + Send + Sync>) -> Self {
source = match source.downcast::<Self>() {
@@ -297,46 +270,3 @@ impl From<candle_core::Error> for Error {
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn unsupported_filesystem_operations_have_actionable_error() {
let object_store_error = object_store::Error::Generic {
store: "LocalFileSystem",
source: Box::new(std::io::Error::from(std::io::ErrorKind::Unsupported)),
};
let lance_error = lance::Error::io_source(Box::new(object_store_error));
let error = Error::from(lance_error);
assert!(matches!(
error,
Error::NotSupported { message }
if message.contains("Mountpoint for Amazon S3")
&& message.contains("s3://bucket/path")
));
}
#[test]
fn other_io_errors_remain_lance_errors() {
let object_store_error = object_store::Error::Generic {
store: "LocalFileSystem",
source: Box::new(std::io::Error::from(std::io::ErrorKind::PermissionDenied)),
};
let lance_error = lance::Error::io_source(Box::new(object_store_error));
assert!(matches!(Error::from(lance_error), Error::Lance { .. }));
}
#[test]
fn unsupported_non_filesystem_errors_remain_lance_errors() {
let lance_error = lance::Error::io_source(Box::new(std::io::Error::from(
std::io::ErrorKind::Unsupported,
)));
assert!(matches!(Error::from(lance_error), Error::Lance { .. }));
}
}
+4 -143
View File
@@ -132,14 +132,9 @@ impl ObjectStore for MirroringObjectStore {
if to.primary_only() {
self.primary.copy_opts(from, to, options).await
} else {
// The secondary store can be process-local and less durable than the
// primary, so a source written by another process may not exist here
// or may be evicted before the copy begins.
match self.secondary.copy_opts(from, to, options.clone()).await {
Ok(()) | Err(Error::NotFound { .. }) => {}
Err(err) => return Err(err),
}
self.primary.copy_opts(from, to, options).await
self.secondary.copy_opts(from, to, options.clone()).await?;
self.primary.copy_opts(from, to, options).await?;
Ok(())
}
}
}
@@ -197,8 +192,7 @@ mod test {
use futures::TryStreamExt;
use lance::{dataset::WriteParams, io::ObjectStoreParams};
use lance_testing::datagen::{BatchGenerator, IncrementingInt32, RandomVector};
use object_store::{local::LocalFileSystem, memory::InMemory};
use std::time::Duration;
use object_store::local::LocalFileSystem;
use tempfile;
use crate::{
@@ -207,139 +201,6 @@ mod test {
table::WriteOptions,
};
#[derive(Debug)]
struct EvictBeforeCopyStore {
inner: Arc<dyn ObjectStore>,
}
impl std::fmt::Display for EvictBeforeCopyStore {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
write!(f, "EvictBeforeCopyStore")
}
}
#[async_trait]
impl ObjectStore for EvictBeforeCopyStore {
async fn put_opts(
&self,
location: &Path,
payload: PutPayload,
options: PutOptions,
) -> Result<PutResult> {
self.inner.put_opts(location, payload, options).await
}
async fn put_multipart_opts(
&self,
location: &Path,
options: PutMultipartOptions,
) -> Result<Box<dyn MultipartUpload>> {
self.inner.put_multipart_opts(location, options).await
}
async fn get_opts(&self, location: &Path, options: GetOptions) -> Result<GetResult> {
self.inner.get_opts(location, options).await
}
fn delete_stream(
&self,
locations: BoxStream<'static, Result<Path>>,
) -> BoxStream<'static, Result<Path>> {
self.inner.delete_stream(locations)
}
fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>> {
self.inner.list(prefix)
}
async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result<ListResult> {
self.inner.list_with_delimiter(prefix).await
}
async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> Result<()> {
self.inner.delete(from).await?;
self.inner.copy_opts(from, to, options).await
}
}
#[tokio::test]
async fn test_copy_when_source_is_missing_from_secondary() {
let primary_dir = tempfile::tempdir().unwrap();
let secondary_dir = tempfile::tempdir().unwrap();
let primary: Arc<dyn ObjectStore> =
Arc::new(LocalFileSystem::new_with_prefix(primary_dir.path()).unwrap());
let secondary: Arc<dyn ObjectStore> =
Arc::new(LocalFileSystem::new_with_prefix(secondary_dir.path()).unwrap());
let store = MirroringObjectStore {
primary: primary.clone(),
secondary: secondary.clone(),
};
let staging = Path::from("_versions/1.manifest-staging");
let finalized = Path::from("_versions/1.manifest");
primary
.put(&staging, "manifest contents".into())
.await
.unwrap();
tokio::time::timeout(Duration::from_secs(5), store.copy(&staging, &finalized))
.await
.expect("copy should not hang when the secondary source is missing")
.unwrap();
let copied = primary
.get(&finalized)
.await
.unwrap()
.bytes()
.await
.unwrap();
assert_eq!(copied, "manifest contents");
assert!(matches!(
secondary.head(&finalized).await,
Err(Error::NotFound { .. })
));
}
#[tokio::test]
async fn test_copy_when_secondary_source_disappears_after_head() {
let primary: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let secondary_inner: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let secondary: Arc<dyn ObjectStore> = Arc::new(EvictBeforeCopyStore {
inner: secondary_inner.clone(),
});
let store = MirroringObjectStore {
primary: primary.clone(),
secondary,
};
let staging = Path::from("_versions/1.manifest-staging");
let finalized = Path::from("_versions/1.manifest");
primary
.put(&staging, "manifest contents".into())
.await
.unwrap();
secondary_inner
.put(&staging, "manifest contents".into())
.await
.unwrap();
store.copy(&staging, &finalized).await.unwrap();
let copied = primary
.get(&finalized)
.await
.unwrap()
.bytes()
.await
.unwrap();
assert_eq!(copied, "manifest contents");
assert!(matches!(
secondary_inner.head(&finalized).await,
Err(Error::NotFound { .. })
));
}
// This test is ignored because lance 3.0 introduced LocalWriter optimization
// that bypasses the object store wrapper for local writes. The mirroring feature
// still works for remote/cloud storage, but can't be tested with local storage.
+28 -4
View File
@@ -1661,8 +1661,14 @@ mod tests {
#[tokio::test]
async fn test_setters_getters() {
// TODO: Switch back to memory://foo after https://github.com/lancedb/lancedb/issues/1051
// is fixed
let tmp_dir = tempdir().unwrap();
let dataset_path = tmp_dir.path().join("test.lance");
let uri = dataset_path.to_str().unwrap();
let batches = make_test_batches();
let conn = connect("memory://foo").execute().await.unwrap();
let conn = connect(uri).execute().await.unwrap();
let table = conn
.create_table("my_table", batches)
.execute()
@@ -1757,8 +1763,14 @@ mod tests {
#[tokio::test]
async fn test_execute() {
// TODO: Switch back to memory://foo after https://github.com/lancedb/lancedb/issues/1051
// is fixed
let tmp_dir = tempdir().unwrap();
let dataset_path = tmp_dir.path().join("test.lance");
let uri = dataset_path.to_str().unwrap();
let batches = make_non_empty_batches();
let conn = connect("memory://foo").execute().await.unwrap();
let conn = connect(uri).execute().await.unwrap();
let table = conn
.create_table("my_table", batches)
.execute()
@@ -1877,8 +1889,14 @@ mod tests {
#[tokio::test]
async fn test_select_with_transform() {
// TODO: Switch back to memory://foo after https://github.com/lancedb/lancedb/issues/1051
// is fixed
let tmp_dir = tempdir().unwrap();
let dataset_path = tmp_dir.path().join("test.lance");
let uri = dataset_path.to_str().unwrap();
let batches = make_non_empty_batches();
let conn = connect("memory://foo").execute().await.unwrap();
let conn = connect(uri).execute().await.unwrap();
let table = conn
.create_table("my_table", batches)
.execute()
@@ -1975,9 +1993,15 @@ mod tests {
#[tokio::test]
async fn test_execute_no_vector() {
// TODO: Switch back to memory://foo after https://github.com/lancedb/lancedb/issues/1051
// is fixed
let tmp_dir = tempdir().unwrap();
let dataset_path = tmp_dir.path().join("test.lance");
let uri = dataset_path.to_str().unwrap();
// test that it's ok to not specify a query vector (just filter / limit)
let batches = make_non_empty_batches();
let conn = connect("memory://foo").execute().await.unwrap();
let conn = connect(uri).execute().await.unwrap();
let table = conn
.create_table("my_table", batches)
.execute()
+1 -59
View File
@@ -373,37 +373,6 @@ pub fn parse_db_url(db_url: &str) -> Result<ParsedDbUrl> {
Ok(ParsedDbUrl { db_name, db_prefix })
}
fn validate_dns_hostname(hostname: &str) -> Result<()> {
let ascii_hostname = match url::Host::parse(hostname) {
Ok(url::Host::Domain(hostname)) => hostname,
Ok(_) => {
return Err(Error::InvalidInput {
message: "LanceDB Cloud database URI or region produced a non-DNS hostname"
.to_string(),
});
}
Err(err) => {
return Err(Error::InvalidInput {
message: format!(
"LanceDB Cloud database URI or region produced an invalid hostname: {err}"
),
});
}
};
if ascii_hostname.len() > 253
|| ascii_hostname
.split('.')
.any(|label| label.is_empty() || label.len() > 63)
{
return Err(Error::InvalidInput {
message: "LanceDB Cloud database URI or region produced an invalid hostname: DNS labels must contain 1 to 63 bytes and the full hostname must not exceed 253 bytes".to_string(),
});
}
Ok(())
}
impl RestfulLanceDbClient<Sender> {
fn get_timeout(passed: Option<Duration>, env_var: &str) -> Result<Option<Duration>> {
if let Some(passed) = passed {
@@ -511,11 +480,7 @@ impl RestfulLanceDbClient<Sender> {
let host = match host_override {
Some(host_override) => host_override,
None => {
let hostname = format!("{}.{}.api.lancedb.com", parsed_url.db_name, region);
validate_dns_hostname(&hostname)?;
format!("https://{hostname}")
}
None => format!("https://{}.{}.api.lancedb.com", parsed_url.db_name, region),
};
debug!("Created client for host: {}", host);
let retry_config = client_config.retry_config.clone().try_into()?;
@@ -1192,29 +1157,6 @@ mod tests {
assert_eq!(headers.get("x-api-key").unwrap(), "api-key");
}
#[test]
fn test_rejects_invalid_cloud_dns_hostname() {
let invalid_database_names = ["a".repeat(64), "invalid..database".to_string()];
for db_name in invalid_database_names {
let parsed_url = parse_db_url(&format!("db://{db_name}")).unwrap();
let error = RestfulLanceDbClient::<Sender>::try_new(
&parsed_url,
"us-east-1",
None,
HeaderMap::new(),
ClientConfig::default(),
None,
)
.unwrap_err();
assert!(
matches!(error, Error::InvalidInput { ref message } if message.contains("DNS labels must contain 1 to 63 bytes")),
"unexpected error: {error}"
);
}
}
// Test implementation of HeaderProvider
#[derive(Debug, Clone)]
struct TestHeaderProvider {
+50 -612
View File
@@ -1,7 +1,7 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
pub mod blobs;
mod blobs;
pub mod insert;
use self::insert::{RemoteWriteExec, WriteOp};
@@ -23,13 +23,11 @@ use crate::table::AddResult;
use crate::table::BranchDiff;
use crate::table::DeleteResult;
use crate::table::DropColumnsResult;
use crate::table::LsmStats;
use crate::table::LsmWriteSpec;
use crate::table::MergeBranchResult;
use crate::table::MergeResult;
use crate::table::Tags;
use crate::table::UpdateResult;
use crate::table::lsm_stats::GetLsmStatsResponse;
use crate::table::merge::MergeFilter;
use crate::table::query::create_multi_vector_plan;
use crate::table::write_progress::FinishOnDrop;
@@ -993,18 +991,6 @@ impl<S: HttpSend> RemoteTable<S> {
}
}
/// Send an LSM operator request with the transport retry layer **off**.
///
/// Retry policy on these routes belongs to the checkpoint loop, which
/// reads the status and can tell contention from a lost claim. Leaving the
/// transport layer on would re-ask on its own schedule first, and surface
/// an `Error::Retry` whose status the loop would then have to unwrap.
async fn send_lsm_route(&self, request: RequestBuilder) -> Result<(String, reqwest::Response)> {
let (request_id, response) = self.send(request, false).await?;
let response = self.check_table_response(&request_id, response).await?;
Ok((request_id, response))
}
/// Build a POST request and attach the read-freshness headers
/// (`x-lancedb-min-version`, `x-lancedb-min-timestamp`).
fn post_read(&self, uri: &str) -> RequestBuilder {
@@ -2482,40 +2468,6 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
})
}
async fn flush_lsm(&self) -> Result<()> {
let request = self
.client
.post(&format!("/v1/table/{}/flush_lsm/", self.identifier));
self.send_lsm_route(request).await?;
Ok(())
}
async fn compact_lsm(&self) -> Result<()> {
let request = self
.client
.post(&format!("/v1/table/{}/compact_lsm/", self.identifier));
self.send_lsm_route(request).await?;
Ok(())
}
async fn get_lsm_stats(&self, include_generation_rows: bool) -> Result<Option<LsmStats>> {
// Read-semantics POST, like `get_lsm_write_spec`.
let request = self
.post_read(&format!("/v1/table/{}/get_lsm_stats/", self.identifier))
.json(&serde_json::json!({
"include_generation_rows": include_generation_rows,
}));
let (request_id, response) = self.send_lsm_route(request).await?;
let body = response.text().await.err_to_http(request_id.clone())?;
let parsed: GetLsmStatsResponse = serde_json::from_str(&body).map_err(|e| Error::Http {
source: format!("Failed to parse get_lsm_stats response: {e}").into(),
request_id,
status_code: None,
})?;
// `null` — and only — when the table has no LSM write path.
Ok(parsed.lsm_stats)
}
async fn set_lsm_write_spec(&self, spec: LsmWriteSpec) -> Result<()> {
self.check_mutable().await?;
@@ -2839,10 +2791,9 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
}
async fn index_stats(&self, index_name: &str) -> Result<Option<IndexStatistics>> {
let encoded_name = urlencoding::encode(index_name);
let mut request = self.post_read(&format!(
"/v1/table/{}/index/{encoded_name}/stats/",
self.identifier
"/v1/table/{}/index/{}/stats/",
self.identifier, index_name
));
let version = self.current_version().await;
let mut body = serde_json::json!({ "version": version });
@@ -2869,10 +2820,9 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
}
async fn drop_index(&self, index_name: &str) -> Result<()> {
let encoded_name = urlencoding::encode(index_name);
let request = self.apply_branch_query(self.client.post(&format!(
"/v1/table/{}/index/{encoded_name}/drop/",
self.identifier
"/v1/table/{}/index/{}/drop/",
self.identifier, index_name
)));
let (request_id, response) = self.send(request, true).await?;
if response.status() == StatusCode::NOT_FOUND {
@@ -2885,10 +2835,9 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
}
async fn prewarm_index(&self, index_name: &str) -> Result<()> {
let encoded_name = urlencoding::encode(index_name);
let request = self.client.post(&format!(
"/v1/table/{}/index/{encoded_name}/prewarm/",
self.identifier
"/v1/table/{}/index/{}/prewarm/",
self.identifier, index_name
));
let (request_id, response) = self.send(request, true).await?;
if response.status() == StatusCode::NOT_FOUND {
@@ -3140,12 +3089,10 @@ mod tests {
Box::pin(table.delete("false").map_ok(|_| ())),
Box::pin(
table
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![(
"x".into(),
"y".into(),
)]))
.execute()
.add_columns(
NewColumnTransform::SqlExpressions(vec![("x".into(), "y".into())]),
None,
)
.map_ok(|_| ()),
),
Box::pin(async {
@@ -4353,9 +4300,32 @@ mod tests {
"fetch_blobs",
);
let message = table
.fetch_blob_files("image", &[1])
.await
.unwrap_err()
.to_string();
assert!(
message.contains("fetch_blob_files is not supported on LanceDB Cloud"),
"got: {message}"
);
assert!(
!message.contains("Use fetch_blobs"),
"old server must not be told to use fetch_blobs, got: {message}"
);
}
#[tokio::test]
async fn test_blob_files_point_at_fetch_blobs_on_a_blob_capable_server() {
let table = Table::new_with_handler_version(
"my_table",
semver::Version::new(0, 5, 0),
|_| -> http::Response<String> { panic!("fetch_blob_files must not reach the server") },
);
assert_not_supported_error(
table.fetch_blob_files("image", &[1]).await.unwrap_err(),
"requires LanceDB Cloud server 0.5.0 or newer",
"Use fetch_blobs for full bytes",
);
}
@@ -5955,7 +5925,6 @@ mod tests {
.await
.unwrap();
// Lance 10 retains original token positions after stop-word removal.
assert_eq!(
tokens,
vec![
@@ -6442,12 +6411,13 @@ mod tests {
});
let result = table
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![
("b".into(), "a + 1".into()),
("x".into(), "cast(NULL as int32)".into()),
]))
.execute()
.add_columns(
NewColumnTransform::SqlExpressions(vec![
("b".into(), "a + 1".into()),
("x".into(), "cast(NULL as int32)".into()),
]),
None,
)
.await
.unwrap();
@@ -6541,41 +6511,6 @@ mod tests {
assert!(matches!(e, Error::IndexNotFound { .. }));
}
/// Index names are unvalidated, so reserved characters must be
/// percent-encoded or they restructure the request path.
#[tokio::test]
async fn test_per_index_paths_encode_reserved_characters() {
const NAME: &str = "my/index?a#b c";
const PREFIX: &str = "/v1/table/my_table/index/my%2Findex%3Fa%23b%20c";
let table = Table::new_with_handler("my_table", |request| {
assert_eq!(request.url().path(), format!("{PREFIX}/stats/"));
let body = serde_json::json!({
"num_indexed_rows": 1,
"num_unindexed_rows": 0,
"index_type": "IVF_PQ",
"distance_type": "l2"
});
http::Response::builder()
.status(200)
.body(serde_json::to_string(&body).unwrap())
.unwrap()
});
assert!(table.index_stats(NAME).await.unwrap().is_some());
let table = Table::new_with_handler("my_table", |request| {
assert_eq!(request.url().path(), format!("{PREFIX}/drop/"));
http::Response::builder().status(200).body("{}").unwrap()
});
table.drop_index(NAME).await.unwrap();
let table = Table::new_with_handler("my_table", |request| {
assert_eq!(request.url().path(), format!("{PREFIX}/prewarm/"));
http::Response::builder().status(200).body("{}").unwrap()
});
table.prewarm_index(NAME).await.unwrap();
}
#[tokio::test]
async fn test_set_lsm_write_spec_unsharded() {
let table = Table::new_with_handler("my_table", |request| {
@@ -6729,499 +6664,6 @@ mod tests {
assert!(table.get_lsm_write_spec().await.unwrap().is_none());
}
/// Build a `get_lsm_stats` body for one bucket holding `generations`.
fn stats_body(generations: &[u64], compacting: bool) -> String {
serde_json::json!({
"lsm_stats": {
"buckets": [{
"shard_id": "b0",
"status": "Active",
"writer_epoch": 1,
"manifest_version": 1,
"current_generation": generations.iter().max().copied().unwrap_or(0) + 1,
"replay_after_wal_entry_position": 0,
"wal_entry_position_last_seen": 0,
"generations": generations.iter()
.map(|g| serde_json::json!({ "generation": g, "bytes": 1 }))
.collect::<Vec<_>>(),
"compacting": compacting,
"memtables": [],
}],
}
})
.to_string()
}
/// `flush_lsm` / `compact_lsm` answer 202 with no body at all.
fn accepted() -> http::Response<String> {
http::Response::builder()
.status(202)
.body(String::new())
.unwrap()
}
fn ok_json(body: String) -> http::Response<String> {
http::Response::builder().status(200).body(body).unwrap()
}
/// A flush landing in an empty L0 finishes on the opening stats read
/// alone. Asserting zero compacts is the point: "it returned Ok" is also
/// true of a loop that ran a pointless pass.
#[tokio::test(start_paused = true)]
async fn test_checkpoint_short_circuits_on_empty_l0() {
let compacts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = compacts.clone();
let table = Table::new_with_handler("my_table", move |request| {
let path = request.url().path().to_string();
if path.contains("compact_lsm") {
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
panic!("an already-converged table must issue no compact calls");
}
if path.contains("flush_lsm") {
return accepted();
}
assert_eq!(path, "/v1/table/my_table/get_lsm_stats/");
ok_json(stats_body(&[], false))
});
table.checkpoint_lsm().await.unwrap();
assert_eq!(compacts.load(std::sync::atomic::Ordering::SeqCst), 0);
}
/// The loop triggers compaction until every generation that existed at
/// the start is gone, one bounded prefix per pass.
#[tokio::test(start_paused = true)]
async fn test_checkpoint_triggers_until_targets_are_drained() {
let compacts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = compacts.clone();
let table = Table::new_with_handler("my_table", move |request| {
let path = request.url().path().to_string();
if path.contains("flush_lsm") {
return accepted();
}
if path.contains("compact_lsm") {
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
return accepted();
}
// Each pass drains the oldest generation.
let drained = seen.load(std::sync::atomic::Ordering::SeqCst);
let left: Vec<u64> = [1u64, 2, 3].into_iter().skip(drained).collect();
ok_json(stats_body(&left, false))
});
table.checkpoint_lsm().await.unwrap();
assert_eq!(
compacts.load(std::sync::atomic::Ordering::SeqCst),
3,
"one trigger per generation prefix, then stop"
);
}
/// Generations created *during* the checkpoint are not waited on, which
/// is what lets the loop terminate on a table taking writes where "L0 is
/// empty" never becomes true.
#[tokio::test(start_paused = true)]
async fn test_checkpoint_ignores_generations_created_while_it_runs() {
let compacts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = compacts.clone();
let table = Table::new_with_handler("my_table", move |request| {
let path = request.url().path().to_string();
if path.contains("flush_lsm") {
return accepted();
}
if path.contains("compact_lsm") {
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
return accepted();
}
// Target is 5. One pass drains it; a writer keeps adding above.
let n = seen.load(std::sync::atomic::Ordering::SeqCst);
let body = if n == 0 {
stats_body(&[5], false)
} else {
stats_body(&[6, 7], false)
};
ok_json(body)
});
table.checkpoint_lsm().await.unwrap();
assert_eq!(
compacts.load(std::sync::atomic::Ordering::SeqCst),
1,
"the loop must not chase generations written after it started"
);
}
/// Contention is a 429 and must be retried. The server keeps it off 503
/// precisely so the client can act on the status alone — reading it as
/// terminal stops the checkpoint early on a healthy node.
#[tokio::test(start_paused = true)]
async fn test_checkpoint_retries_contention() {
let compacts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = compacts.clone();
let table = Table::new_with_handler("my_table", move |request| {
let path = request.url().path().to_string();
if path.contains("flush_lsm") {
return accepted();
}
if path.contains("compact_lsm") {
// First two triggers: every bucket already latched.
if seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2 {
return http::Response::builder()
.status(429)
.body(r#"{"code":21,"error":"Too many concurrent writes"}"#.to_string())
.unwrap();
}
return accepted();
}
let accepted_triggers = seen
.load(std::sync::atomic::Ordering::SeqCst)
.saturating_sub(2);
let left: Vec<u64> = if accepted_triggers == 0 {
vec![1]
} else {
vec![]
};
ok_json(stats_body(&left, false))
});
table
.checkpoint_lsm()
.await
.expect("contention must not abort the checkpoint");
assert_eq!(
compacts.load(std::sync::atomic::Ordering::SeqCst),
3,
"assert the retry count, not just the outcome"
);
}
/// A transient fault on the poll must not abort the checkpoint. This route
/// meets the most contention — it runs every `POLL_INTERVAL` for the
/// checkpoint's whole life, with the transport retry layer disabled — yet
/// was the one call reached with a bare `?`.
#[tokio::test(start_paused = true)]
async fn test_checkpoint_retries_a_contended_stats_poll() {
let polls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = polls.clone();
let table = Table::new_with_handler("my_table", move |request| {
let path = request.url().path().to_string();
if path.contains("flush_lsm") || path.contains("compact_lsm") {
return accepted();
}
// The opening read lands; the next two polls are latched out.
let n = seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if (1..3).contains(&n) {
return http::Response::builder()
.status(429)
.body(r#"{"code":21,"error":"Too many concurrent writes"}"#.to_string())
.unwrap();
}
ok_json(stats_body(if n < 4 { &[1] } else { &[] }, false))
});
table
.checkpoint_lsm()
.await
.expect("a contended poll must be retried, not surfaced");
assert_eq!(
polls.load(std::sync::atomic::Ordering::SeqCst),
5,
"the two rejected polls must be re-issued, not skipped"
);
}
/// Contention and a lost claim draw on separate budgets: five straight
/// 429s on `flush`, more than `MAX_REISSUES`, must still converge. On one
/// shared counter this spent the re-issue cap and then reported a lost
/// claim nothing had ever reported.
#[tokio::test(start_paused = true)]
async fn test_contention_does_not_exhaust_the_reissue_budget() {
let flushes = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = flushes.clone();
let table = Table::new_with_handler("my_table", move |request| {
let path = request.url().path().to_string();
if path.contains("flush_lsm") {
if seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 5 {
return http::Response::builder()
.status(429)
.body(r#"{"code":21,"error":"Too many concurrent writes"}"#.to_string())
.unwrap();
}
return accepted();
}
if path.contains("compact_lsm") {
return accepted();
}
ok_json(stats_body(&[], false))
});
table
.checkpoint_lsm()
.await
.expect("contention must not be reported as a lost claim");
assert_eq!(
flushes.load(std::sync::atomic::Ordering::SeqCst),
6,
"five retries against one seal, then it lands"
);
}
/// An exhausted retry budget surfaces the fault that consumed it, not a
/// message the loop invented: "429, nine times" points an operator at a
/// saturated pool, a generic runtime error points them nowhere.
#[tokio::test(start_paused = true)]
async fn test_exhausted_retries_surface_the_underlying_fault() {
let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = calls.clone();
let table = Table::new_with_handler("my_table", move |_request| {
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
http::Response::builder()
.status(429)
.body(r#"{"code":21,"error":"Too many concurrent writes"}"#.to_string())
.unwrap()
});
let err = table.checkpoint_lsm().await.unwrap_err();
assert!(
matches!(&err, Error::Http { status_code: Some(s), .. } if s.as_u16() == 429),
"the fault that spent the budget must be the one reported: {err:?}"
);
assert_eq!(
calls.load(std::sync::atomic::Ordering::SeqCst),
9,
"one call plus MAX_RETRIES — the re-issue budget is not spent on top"
);
}
/// A draining node is terminal, but the client does not know that from the
/// status: draining and a proxy blip are both 503, and telling them apart
/// takes parsing the body for a namespace code. So it spends the retry
/// budget and then reports what the server said — the drain gate never
/// releases, so the answer does not change, and the operator still reads
/// "WAL node draining" in the error.
#[tokio::test(start_paused = true)]
async fn test_draining_surfaces_after_the_retry_budget() {
let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = calls.clone();
let table = Table::new_with_handler("my_table", move |_request| {
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
http::Response::builder()
.status(503)
.body(r#"{"code":19,"error":"WAL node draining"}"#.to_string())
.unwrap()
});
let err = table.checkpoint_lsm().await.unwrap_err();
let message = err.to_string();
assert!(
matches!(&err, Error::Http { status_code: Some(s), .. } if s.as_u16() == 503),
"the 503 must surface as itself: {err:?}"
);
assert!(
message.contains("WAL node draining"),
"the server's own diagnosis must survive to the caller: {message}"
);
assert_eq!(
calls.load(std::sync::atomic::Ordering::SeqCst),
9,
"one call plus MAX_RETRIES, then it reports rather than spinning"
);
}
/// A long stall with nothing compacting must keep waiting, not fail. The
/// client cannot judge this: a checkpoint queued behind unrelated tables
/// on the pod-wide compactor pool reports exactly these numbers — flat
/// generations, an idle latch — as one whose merges are failing. The
/// deadline is the caller's.
#[tokio::test(start_paused = true)]
async fn test_checkpoint_waits_out_a_long_stall_rather_than_failing() {
let polls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = polls.clone();
let table = Table::new_with_handler("my_table", move |request| {
let path = request.url().path().to_string();
if path.contains("flush_lsm") || path.contains("compact_lsm") {
return accepted();
}
// Flat for far longer than any bound this loop ever had, with
// `compacting: false` throughout — then it drains.
let n = seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
ok_json(stats_body(if n < 40 { &[1, 2] } else { &[] }, false))
});
table
.checkpoint_lsm()
.await
.expect("a stall is the server being slow, not the client's call to make");
assert!(
polls.load(std::sync::atomic::Ordering::SeqCst) > 40,
"the loop must have kept polling well past the old ten-poll bound"
);
}
/// A pass already owns the latch on every outstanding bucket, so the loop
/// waits rather than piling on triggers it would only refuse. This is the
/// sole thing `compacting` is read for.
#[tokio::test(start_paused = true)]
async fn test_checkpoint_waits_while_a_pass_is_running() {
let polls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let compacts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen_polls = polls.clone();
let seen_compacts = compacts.clone();
let table = Table::new_with_handler("my_table", move |request| {
let path = request.url().path().to_string();
if path.contains("flush_lsm") {
return accepted();
}
if path.contains("compact_lsm") {
seen_compacts.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
return accepted();
}
// Latched for many polls, then done.
let n = seen_polls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
ok_json(if n > 15 {
stats_body(&[], false)
} else {
stats_body(&[1], true)
})
});
table
.checkpoint_lsm()
.await
.expect("a running pass is progress, not a stall");
assert_eq!(
compacts.load(std::sync::atomic::Ordering::SeqCst),
0,
"never trigger against a bucket already compacting"
);
}
/// WAL off ⇒ `None`; WAL on ⇒ a fully populated `Some` with no field
/// defaulting to a zero it did not measure. `include_generation_rows`
/// rides in the body and is off unless asked for.
#[tokio::test]
async fn test_get_lsm_stats_round_trip() {
let table = Table::new_with_handler("my_table", |request| {
assert_eq!(request.url().path(), "/v1/table/my_table/get_lsm_stats/");
let body = request.body().unwrap().as_bytes().unwrap();
let body: serde_json::Value = serde_json::from_slice(body).unwrap();
assert_eq!(
body["include_generation_rows"], true,
"the flag must reach the server, not be silently dropped"
);
let response = serde_json::json!({
"lsm_stats": {
"buckets": [{
"shard_id": "b0",
"status": "Active",
"writer_epoch": 3,
"manifest_version": 11,
"current_generation": 9,
"replay_after_wal_entry_position": 100,
"wal_entry_position_last_seen": 140,
"generations": [{ "generation": 8, "bytes": 4096, "rows": 30 }],
"compacting": false,
"memtables": [
{ "generation": 9, "rows": 12, "bytes": 900, "batches": 2,
"indexes": ["vec_idx"] }
],
}],
}
});
http::Response::builder()
.status(200)
.body(response.to_string())
.unwrap()
});
let stats = table
.get_lsm_stats(true)
.await
.unwrap()
.expect("a WAL-backed table reports Some");
let bucket = &stats.buckets[0];
assert_eq!(bucket.replay_after_wal_entry_position, 100);
assert_eq!(bucket.wal_entry_position_last_seen, 140);
assert!(!bucket.compacting);
assert_eq!(bucket.generations[0].generation, 8);
assert_eq!(bucket.generations[0].rows, Some(30));
// The line that answers "why is my fresh-tier vector search
// brute-force" — an absent index name is the whole explanation.
let memtables = bucket.memtables.as_ref().unwrap();
assert_eq!(memtables[0].indexes, vec!["vec_idx".to_string()]);
}
/// A 404 arrives as `TableNotFound`, not as a lost claim the loop
/// re-issues from flush until its cap. The two are distinguished by
/// status: 404 is "no such table", 421 is "this node holds no claim".
/// They shared 404 once, and the loop chased a name that never existed.
#[tokio::test(start_paused = true)]
async fn test_missing_table_is_not_read_as_a_lost_claim() {
let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = calls.clone();
let table = Table::new_with_handler("my_table", move |_request| {
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
http::Response::builder()
.status(404)
.body(r#"{"code":4,"error":"Not found: Table not found: my_table"}"#.to_string())
.unwrap()
});
let err = table.checkpoint_lsm().await.unwrap_err();
assert!(
matches!(err, Error::TableNotFound { .. }),
"a missing table must say so: {err:?}"
);
assert_eq!(
calls.load(std::sync::atomic::Ordering::SeqCst),
1,
"no point re-claiming a table that does not exist"
);
}
/// A lost claim — 421, not 404 — does re-issue from flush, the call that
/// re-claims and replays.
#[tokio::test(start_paused = true)]
async fn test_registry_miss_reissues_from_flush() {
let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = calls.clone();
let table = Table::new_with_handler("my_table", move |request| {
let path = request.url().path().to_string();
let n = seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if path.contains("flush_lsm") {
// First flush lands; the claim is then lost, and the
// re-issued flush succeeds.
return accepted();
}
if path.contains("compact_lsm") {
if n < 4 {
return http::Response::builder()
.status(421)
.body(r#"{"code":19,"error":"table not claimed"}"#.to_string())
.unwrap();
}
return accepted();
}
ok_json(stats_body(if n < 6 { &[1] } else { &[] }, false))
});
table
.checkpoint_lsm()
.await
.expect("a lost claim must be recovered by re-flushing, not surfaced");
}
#[tokio::test]
async fn test_get_lsm_stats_absent_when_wal_off() {
let table = Table::new_with_handler("my_table", |_request| {
http::Response::builder()
.status(200)
.body(serde_json::json!({ "lsm_stats": null }).to_string())
.unwrap()
});
assert!(table.get_lsm_stats(false).await.unwrap().is_none());
}
#[tokio::test]
async fn test_wait_for_index() {
let table = _make_table_with_indices(0);
@@ -7700,12 +7142,10 @@ mod tests {
}
"add_columns" => {
let _ = table
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![(
"c".into(),
"a + 1".into(),
)]))
.execute()
.add_columns(
NewColumnTransform::SqlExpressions(vec![("c".into(), "a + 1".into())]),
None,
)
.await;
}
"drop_columns" => {
@@ -10463,12 +9903,10 @@ mod tests {
.await
.unwrap();
branch
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![(
"b".into(),
"a + 1".into(),
)]))
.execute()
.add_columns(
NewColumnTransform::SqlExpressions(vec![("b".into(), "a + 1".into())]),
None,
)
.await
.unwrap();
branch
File diff suppressed because it is too large Load Diff
+13 -336
View File
@@ -3,7 +3,6 @@
//! LanceDB Table APIs
use crate::blob::BlobFile;
use arrow_array::{LargeBinaryArray, RecordBatch, RecordBatchReader};
use arrow_schema::{Schema, SchemaRef};
use async_trait::async_trait;
@@ -13,6 +12,7 @@ use datafusion_physical_plan::ExecutionPlan;
use datafusion_physical_plan::display::DisplayableExecutionPlan;
use futures::StreamExt;
use futures::stream::FuturesUnordered;
use lance::dataset::BlobFile;
pub use lance::dataset::ColumnAlteration;
pub use lance::dataset::NewColumnTransform;
pub use lance::dataset::ReadParams;
@@ -65,15 +65,12 @@ use crate::utils::{PatchReadParam, PatchWriteParam, resolve_arrow_field_path};
use self::dataset::DatasetConsistencyWrapper;
use self::merge::MergeInsertBuilder;
pub mod add_columns;
mod add_data;
pub mod branch_merge;
pub mod checkpoint;
mod create_index;
pub mod datafusion;
pub(crate) mod dataset;
pub mod delete;
pub mod lsm_stats;
pub mod merge;
pub mod optimize;
mod primary_key;
@@ -82,7 +79,6 @@ pub mod schema_evolution;
pub mod update;
pub mod write_progress;
use crate::index::waiter::wait_for_index;
pub use add_columns::AddColumnsBuilder;
#[cfg(feature = "remote")]
pub(crate) use add_data::PreprocessingOutput;
pub use add_data::{AddDataBuilder, AddDataMode, AddResult, NaNVectorBehavior};
@@ -95,8 +91,8 @@ pub use delete::DeleteResult;
use futures::future::join_all;
pub use lance::dataset::refs::{BranchContents, Ref, TagContents, Tags as LanceTags};
pub use lance::dataset::scanner::DatasetRecordBatchStream;
use lance::dataset::statistics::DatasetStatisticsExt;
pub use lance_index::optimize::OptimizeOptions;
pub use lsm_stats::{BucketStats, GenerationStats, LsmStats, MemtableStats};
pub use optimize::{CompactionOptions, OptimizeAction, OptimizeStats};
pub use schema_evolution::{
AddColumnsResult, AlterColumnsResult, DropColumnsResult, FieldMetadataUpdate,
@@ -687,31 +683,6 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
message: "get_lsm_write_spec is not supported on this table type".into(),
})
}
/// Seal every bucket's active memtable into L0.
///
/// The default implementation returns `NotSupported`.
async fn flush_lsm(&self) -> Result<()> {
Err(Error::NotSupported {
message: "flush_lsm is not supported on this table type".into(),
})
}
/// Trigger a background L0 → base compaction pass per bucket.
///
/// The default implementation returns `NotSupported`.
async fn compact_lsm(&self) -> Result<()> {
Err(Error::NotSupported {
message: "compact_lsm is not supported on this table type".into(),
})
}
/// Read live LSM state, or `None` when the LSM write path is not
/// enabled for this table.
///
/// The default implementation returns `NotSupported`.
async fn get_lsm_stats(&self, _include_generation_rows: bool) -> Result<Option<LsmStats>> {
Err(Error::NotSupported {
message: "get_lsm_stats is not supported on this table type".into(),
})
}
/// Drain and close any cached MemWAL shard writers for this table.
///
/// The default implementation is a no-op; table types that maintain
@@ -1649,8 +1620,12 @@ impl Table {
}
/// Add new columns to the table, providing values to fill in.
pub fn add_columns(&self) -> AddColumnsBuilder {
AddColumnsBuilder::new(self.inner.clone())
pub async fn add_columns(
&self,
transforms: NewColumnTransform,
read_columns: Option<Vec<String>>,
) -> Result<AddColumnsResult> {
self.inner.add_columns(transforms, read_columns).await
}
/// Change a column's name or nullability.
@@ -1753,85 +1728,6 @@ impl Table {
self.inner.get_lsm_write_spec().await
}
/// Converge this table's LSM write path into its base table.
///
/// One `flush` to seal 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`, so
/// there is no held socket and nothing to reconcile if you drop this
/// future partway through.
///
/// **Best-effort.** Generations created *after* the opening flush are
/// deliberately not waited on — that is what lets this terminate on a
/// table taking writes. Idempotent and safe on a cadence: an
/// already-converged table costs two round trips and triggers nothing.
///
/// **No deadline, and the caller owns that.** It returns when the target
/// generations are gone, propagates a terminal server fault, and
/// otherwise waits however long the server takes. A slow table and a
/// stuck one are the same picture from here: the compactor pool is shared
/// across every table on the node, so a checkpoint queued behind
/// unrelated work is indistinguishable from one that is merging. Wrap
/// this in `tokio::time::timeout` for a wall-clock bound; abandoning it
/// partway costs nothing.
///
/// # Example
///
/// ```no_run
/// # use lancedb::Table;
/// # async fn example(table: &Table) -> Result<(), Box<dyn std::error::Error>> {
/// let before = table.get_lsm_stats(false).await?;
/// table.checkpoint_lsm().await?;
/// let after = table.get_lsm_stats(false).await?;
/// # Ok(())
/// # }
/// ```
pub async fn checkpoint_lsm(&self) -> Result<()> {
checkpoint::checkpoint_lsm(self).await
}
/// Seal every bucket's active memtable into L0 without touching the
/// base table.
///
/// Independently useful: flushing makes memtable rows readable from L0 at
/// a lower per-query cost. On a node that has not claimed this table it
/// claims it and replays the WAL log first — reporting "nothing to flush"
/// without replaying would lie about durable data.
pub async fn flush_lsm(&self) -> Result<()> {
self.inner.flush_lsm().await
}
/// Run one bounded L0 → base compaction pass per bucket, reporting what
/// it merged and what is left.
///
/// One pass, not convergence: that bounds each request's cost and gives a
/// caller driving its own cadence a progress signal per round trip.
pub async fn compact_lsm(&self) -> Result<()> {
self.inner.compact_lsm().await
}
/// 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.
///
/// `include_generation_rows` reports 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.
///
/// `Ok(None)` only when the LSM write path is not enabled, matching
/// [`Table::get_lsm_write_spec`]. Stats is fresh-tier only, so with the
/// WAL off there is no manifest to report and a struct of zeros would
/// read as measurements.
///
/// Do not build a checkpoint's termination on this: the completion
/// predicate lives in the `flush` and `compact` responses.
pub async fn get_lsm_stats(&self, include_generation_rows: bool) -> Result<Option<LsmStats>> {
self.inner.get_lsm_stats(include_generation_rows).await
}
/// Drain and close any cached MemWAL shard writers held for this table.
///
/// When an [`LsmWriteSpec`] is installed, `merge_insert` opens MemWAL shard
@@ -3547,24 +3443,9 @@ impl BaseTable for NativeTable {
let num_rows = self.count_rows(None).await?;
let num_indices = self.list_indices().await?.len();
let ds = self.dataset.get().await?;
// Sizes come from the manifest. Summing per-field `bytes_on_disk` instead
// would open every data file to read its column metadata, which costs one
// IO per fragment and reports 0 for legacy v1 storage.
//
// The manifest summary covers only the fragments' base data files, so
// overlay files (recorded on each fragment) and index files (recorded in
// the manifest's index section) are added separately.
let mut total_bytes = ds.manifest().summary().total_files_size as usize;
for frag in ds.manifest().fragments.iter() {
for overlay in &frag.overlays {
if let Some(size) = overlay.data_file.file_size_bytes.get() {
total_bytes += size.get() as usize;
}
}
}
for index in ds.load_indices().await?.iter() {
total_bytes += index.total_size_bytes().unwrap_or(0) as usize;
}
let ds_clone = (*ds).clone();
let ds_stats = Arc::new(ds_clone).calculate_data_stats().await?;
let total_bytes = ds_stats.fields.iter().map(|f| f.bytes_on_disk).sum::<u64>() as usize;
let frags = ds.get_fragments();
let mut sorted_sizes = join_all(
@@ -3636,12 +3517,7 @@ impl BaseTable for NativeTable {
#[skip_serializing_none]
#[derive(Debug, Deserialize, PartialEq)]
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,
/// and it excludes any file whose size the manifest does not record
/// (tables and indices written before writers persisted file sizes).
/// The total number of bytes in the table
pub total_bytes: usize,
/// The number of rows in the table
@@ -3702,7 +3578,6 @@ mod tests {
use super::*;
use crate::connect;
use crate::connection::ConnectBuilder;
use crate::io::object_store::io_tracking::IoTrackingStore;
use crate::query::Select;
use crate::query::{ExecutableQuery, QueryBase};
use crate::test_utils::connection::new_test_connection;
@@ -5151,16 +5026,12 @@ mod tests {
let res = table.stats().await.unwrap();
println!("{:#?}", res);
// `total_bytes` is the full on-disk size of the 11 data files (this table
// has no index or overlay files), so it is well above the 2000 bytes of
// column data these 250 int32 pairs hold: each file carries its own footer
// and metadata.
assert_eq!(
res,
TableStatistics {
num_rows: 250,
num_indices: 0,
total_bytes: 8925,
total_bytes: 2300,
fragment_stats: FragmentStatistics {
num_fragments: 11,
num_small_fragments: 11,
@@ -5200,198 +5071,4 @@ mod tests {
}
)
}
/// `total_bytes` counts more than the base data files: index files and
/// overlay files recorded in the manifest are included too.
#[tokio::test]
pub async fn test_stats_includes_index_and_overlay_files() {
use lance::dataset::WriteDestination;
use lance::dataset::transaction::{DataOverlayGroup, Operation};
use lance_file::version::{ConcreteFileVersion, LanceFileVersion};
use lance_file::writer::{FileWriter, FileWriterOptions};
use lance_io::utils::CachedFileSize;
use lance_table::format::DataFile;
use lance_table::format::overlay::{DataOverlayFile, OverlayCoverage};
use roaring::RoaringBitmap;
let tmp_dir = tempdir().unwrap();
let uri = tmp_dir.path().to_str().unwrap();
let conn = ConnectBuilder::new(uri)
.read_consistency_interval(Duration::from_secs(0))
.execute()
.await
.unwrap();
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int32, false),
Field::new("foo", DataType::Int32, true),
]));
let batch = RecordBatch::try_new(
schema.clone(),
vec![
Arc::new(Int32Array::from_iter_values(0..100)),
Arc::new(Int32Array::from_iter_values(0..100)),
],
)
.unwrap();
let table = conn
.create_table("test_stats_extra_files", batch)
.execute()
.await
.unwrap();
let data_only = table.stats().await.unwrap().total_bytes;
assert!(data_only > 0);
// A scalar index adds index files whose sizes are recorded in the
// manifest's index section.
table
.create_index(&["id"], Index::Auto)
.execute()
.await
.unwrap();
let with_index = table.stats().await.unwrap().total_bytes;
let dataset = {
let native = table.as_native().unwrap();
(*native.dataset.get().await.unwrap()).clone()
};
let index_bytes: usize = dataset
.load_indices()
.await
.unwrap()
.iter()
.map(|idx| idx.total_size_bytes().unwrap_or(0) as usize)
.sum();
assert!(index_bytes > 0);
assert_eq!(with_index, data_only + index_bytes);
// Commit an overlay file supplying new `foo` values for the first three
// rows of fragment 0. There is no high-level API that writes overlays
// yet, so write the overlay's data file and commit the `DataOverlay`
// operation by hand.
let read_version = dataset.version().version;
let fragment_id = dataset.get_fragments()[0].id() as u64;
let foo_field_id = dataset.schema().field("foo").unwrap().id;
let overlay_schema = dataset.schema().project_by_ids(&[foo_field_id], true);
let file_version = ConcreteFileVersion::from(LanceFileVersion::Stable);
let filename = "overlay.lance".to_string();
let store = dataset.object_store(None).await.unwrap();
let path = dataset.data_dir().child(filename.clone());
let obj_writer = store.create(&path).await.unwrap();
let mut writer = FileWriter::try_new(
obj_writer,
overlay_schema,
FileWriterOptions {
format_version: Some(file_version.into()),
..Default::default()
},
)
.unwrap();
writer
.write_column(0, Arc::new(Int32Array::from(vec![1000, 1001, 1002])) as _)
.await
.unwrap();
let summary = writer.finish().await.unwrap();
let overlay_bytes = summary.size_bytes as usize;
assert!(overlay_bytes > 0);
let mut data_file = DataFile::new_unstarted(filename, file_version);
data_file.fields = writer
.field_id_to_column_indices()
.iter()
.map(|(field_id, _)| *field_id as i32)
.collect::<Vec<_>>()
.into();
data_file.column_indices = writer
.field_id_to_column_indices()
.iter()
.map(|(_, column_index)| *column_index as i32)
.collect::<Vec<_>>()
.into();
data_file.file_size_bytes = CachedFileSize::new(summary.size_bytes);
let overlay = DataOverlayFile {
data_file,
coverage: OverlayCoverage::dense(RoaringBitmap::from_iter(0..3)),
committed_version: 0,
};
Dataset::commit(
WriteDestination::Dataset(Arc::new(dataset)),
Operation::DataOverlay {
groups: vec![DataOverlayGroup {
fragment_id,
overlays: vec![overlay],
}],
},
Some(read_version),
None,
None,
Arc::new(Default::default()),
false,
)
.await
.unwrap();
table.checkout_latest().await.unwrap();
let with_overlay = table.stats().await.unwrap().total_bytes;
assert_eq!(with_overlay, with_index + overlay_bytes);
}
/// `stats()` must stay manifest-only. Summing per-field `bytes_on_disk`
/// instead opens every data file, so cost would grow with fragment count.
#[tokio::test]
pub async fn test_stats_does_not_read_data_files() {
let tmp_dir = tempdir().unwrap();
let uri = tmp_dir.path().to_str().unwrap();
let conn = ConnectBuilder::new(uri).execute().await.unwrap();
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
let batch = RecordBatch::try_new(
schema.clone(),
vec![Arc::new(Int32Array::from_iter_values(0..10))],
)
.unwrap();
conn.create_table("test_stats_io", batch.clone())
.execute()
.await
.unwrap();
let table = conn.open_table("test_stats_io").execute().await.unwrap();
const NUM_APPENDS: usize = 20;
for _ in 0..NUM_APPENDS {
table.add(batch.clone()).execute().await.unwrap();
}
// Reopen through a tracking store so the counters cover `stats()` alone and
// not the writes above.
let (wrapper, io_stats) = IoTrackingStore::new_wrapper();
let table = conn
.open_table("test_stats_io")
.lance_read_params(ReadParams {
store_options: Some(ObjectStoreParams {
object_store_wrapper: Some(wrapper),
..Default::default()
}),
..Default::default()
})
.execute()
.await
.unwrap();
io_stats.lock().unwrap().read_iops = 0;
let stats = table.stats().await.unwrap();
let read_iops = io_stats.lock().unwrap().read_iops;
assert_eq!(stats.fragment_stats.num_fragments, NUM_APPENDS + 1);
assert!(stats.total_bytes > 0);
// Reading the fragments' data files would take at least one IOP each.
assert!(
read_iops < stats.fragment_stats.num_fragments as u64,
"stats() issued {} read IOPs across {} fragments",
read_iops,
stats.fragment_stats.num_fragments
);
}
}
-161
View File
@@ -1,161 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
//! Builder for adding columns to a table.
use std::sync::Arc;
use lance::dataset::NewColumnTransform;
use super::BaseTable;
use super::schema_evolution::AddColumnsResult;
use crate::{Error, Result};
/// Adds columns to a table. See [`Table::add_columns`](super::Table::add_columns).
pub struct AddColumnsBuilder {
parent: Arc<dyn BaseTable>,
transform: Option<NewColumnTransform>,
read_columns: Option<Vec<String>>,
}
impl std::fmt::Debug for AddColumnsBuilder {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AddColumnsBuilder")
.field("parent", &self.parent)
.field("has_transform", &self.transform.is_some())
.field("read_columns", &self.read_columns)
.finish()
}
}
impl AddColumnsBuilder {
pub(crate) fn new(parent: Arc<dyn BaseTable>) -> Self {
Self {
parent,
transform: None,
read_columns: None,
}
}
/// Set how the new columns' values are produced. Required.
pub fn transform(mut self, transform: NewColumnTransform) -> Self {
self.transform = Some(transform);
self
}
/// Limit which existing columns a [`NewColumnTransform::BatchUDF`] mapper
/// receives. Every other transform determines what it reads, so setting
/// this alongside one is an error rather than a silent no-op.
pub fn read_columns(mut self, columns: impl IntoIterator<Item = impl Into<String>>) -> Self {
self.read_columns = Some(columns.into_iter().map(Into::into).collect());
self
}
/// Add the columns.
pub async fn execute(self) -> Result<AddColumnsResult> {
let Self {
parent,
transform,
read_columns,
} = self;
let Some(transform) = transform else {
return Err(Error::InvalidInput {
message: "add_columns requires a transform".into(),
});
};
if read_columns.is_some() && !matches!(transform, NewColumnTransform::BatchUDF(_)) {
return Err(Error::InvalidInput {
message: "read_columns applies only to a BatchUDF transform; \
every other transform determines what it reads"
.into(),
});
}
parent.add_columns(transform, read_columns).await
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use arrow_array::{Int32Array, RecordBatch, record_batch};
use arrow_schema::{DataType, Field, Schema};
use lance::dataset::{BatchUDF, NewColumnTransform};
use crate::Table;
use crate::connect;
async fn table_with_two_columns(name: &str) -> Table {
let conn = connect("memory://").execute().await.unwrap();
let batch = record_batch!(("x", Int32, [1, 2, 3]), ("y", Int32, [10, 20, 30])).unwrap();
conn.create_table(name, batch).execute().await.unwrap()
}
#[tokio::test]
async fn test_requires_a_transform() {
let table = table_with_two_columns("no_transform").await;
let err = table.add_columns().execute().await.unwrap_err();
assert!(
err.to_string().contains("requires a transform"),
"got: {err}"
);
}
#[tokio::test]
async fn test_read_columns_with_sql_expressions_is_rejected() {
let table = table_with_two_columns("read_cols_sql").await;
let err = table
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![(
"doubled".into(),
"x * 2".into(),
)]))
.read_columns(["x"])
.execute()
.await
.unwrap_err();
assert!(err.to_string().contains("BatchUDF"), "got: {err}");
let schema = table.schema().await.unwrap();
assert!(
schema.field_with_name("doubled").is_err(),
"a rejected call must not commit"
);
}
#[tokio::test]
async fn test_read_columns_limits_what_a_batch_udf_sees() {
let table = table_with_two_columns("read_cols_udf").await;
let output_schema = Arc::new(Schema::new(vec![Field::new("sum", DataType::Int32, true)]));
let mapper_schema = output_schema.clone();
let udf = BatchUDF {
mapper: Box::new(move |batch: &RecordBatch| {
assert!(batch.column_by_name("x").is_some());
assert!(batch.column_by_name("y").is_none(), "y was not requested");
let x = batch["x"].as_any().downcast_ref::<Int32Array>().unwrap();
let doubled: Int32Array = x.iter().map(|v| v.map(|v| v * 2)).collect();
Ok(RecordBatch::try_new(
mapper_schema.clone(),
vec![Arc::new(doubled)],
)?)
}),
output_schema,
result_checkpoint: None,
};
table
.add_columns()
.transform(NewColumnTransform::BatchUDF(udf))
.read_columns(["x"])
.execute()
.await
.unwrap();
let schema = table.schema().await.unwrap();
assert!(schema.field_with_name("sum").is_ok());
}
}
+5 -9
View File
@@ -576,12 +576,10 @@ mod tests {
// Add a new physical column AFTER the embedding column.
table
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![(
"score".into(),
"42.0".into(),
)]))
.execute()
.add_columns(
NewColumnTransform::SqlExpressions(vec![("score".into(), "42.0".into())]),
None,
)
.await
.unwrap();
@@ -685,9 +683,7 @@ mod tests {
true,
)]));
table
.add_columns()
.transform(NewColumnTransform::AllNulls(nested_schema))
.execute()
.add_columns(NewColumnTransform::AllNulls(nested_schema), None)
.await
.unwrap();
-315
View File
@@ -1,315 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
//! Converging a table's LSM write path into its base table.
//!
//! `checkpoint_lsm` seals once, then triggers compaction and watches
//! generation numbers until the L0 that existed at the start is gone.
//!
//! The loop runs in the client, not the server: `compact_lsm` dispatches a
//! pass and returns, so nothing holds a socket and a client can vanish
//! mid-operation with nothing to reconcile. Completion is read from
//! generation numbers in the shard manifest — durable state, unlike a count
//! in a compact response, which a concurrent write invalidates.
//!
//! The target set is fixed at the start, so generations created *during* the
//! checkpoint are ignored. That is what lets it terminate under write load,
//! and what makes it best-effort: it converges the fresh tier as of some
//! instant. Idempotent, abandonable at any point, safe on a cadence.
//!
//! No liveness bound — the caller owns the deadline. The compactor pool is
//! shared pod-wide, so a checkpoint queued behind unrelated tables looks
//! exactly like one that is merging.
use std::collections::HashMap;
use std::future::Future;
use std::time::Duration;
use crate::{Error, Result, Table};
/// The HTTP status a failed request carried, if it carried one.
///
/// `None` for anything with no retry story: a `TableNotFound` that
/// `check_table_response` already translated, or a connection failure that
/// never reached the server. Both are terminal.
fn status_of(e: &Error) -> Option<u16> {
#[cfg(feature = "remote")]
{
match e {
Error::Http {
status_code: Some(status),
..
} => Some(status.as_u16()),
_ => None,
}
}
#[cfg(not(feature = "remote"))]
{
let _ = e;
None
}
}
/// 429 (latch held, pool saturated, or the pod replaying its WAL) and 503 (a
/// draining node, or a proxy between here and it).
///
/// The status is the whole signal: the server deliberately keeps contention
/// off 503, so a latch collision is a 429. A draining node *is* terminal, but
/// it is also a 503 that stays a 503, so retrying spends one budget and then
/// reports the server's own message — cheaper than parsing the body for the
/// namespace code it would take to tell the two apart.
fn is_retryable(e: &Error) -> bool {
matches!(status_of(e), Some(429 | 503))
}
/// 421: the owning node holds no claim. Only `flush` re-claims and replays,
/// so this cannot be retried in place — the caller has to start over.
fn is_lost_claim(e: &Error) -> bool {
status_of(e) == Some(421)
}
/// Interval between `get_lsm_stats` polls. One interval is roughly one
/// compaction pass, the granularity at which the answer can change.
///
/// Fixed rather than configurable, matching `wait_for_index`. It costs
/// nothing on an already-converged table and at most one interval of tail
/// latency after the final pass lands.
const POLL_INTERVAL: Duration = Duration::from_secs(5);
/// Cap on re-issues from `flush` after a 421, so a crash-looping node cannot
/// turn flush → compact → 421 → flush into a spin.
///
/// Deliberately not shared with [`MAX_RETRIES`]: a claim that keeps
/// evaporating is a broken node, while contention is routine and wants a real
/// budget. One shared counter let a merely contended table exhaust this cap
/// and then blame a claim it never lost.
const MAX_REISSUES: usize = 3;
/// Retryable faults tolerated on a *single* request, reset on every success —
/// scattered contention across a long checkpoint must not accumulate toward a
/// cap. Roughly 16s of retrying against the backoff below.
const MAX_RETRIES: usize = 8;
/// Backoff between retries, doubling up to [`RETRY_BACKOFF_MAX`]. Latch
/// contention clears in about the time one pass takes, so start small; a
/// saturated pool wants the ceiling.
const RETRY_BACKOFF_BASE: Duration = Duration::from_millis(100);
const RETRY_BACKOFF_MAX: Duration = Duration::from_secs(5);
/// Sleep before re-issuing a retryable request.
async fn backoff(attempt: usize) {
let delay = RETRY_BACKOFF_BASE
.saturating_mul(1u32 << attempt.min(8) as u32)
.min(RETRY_BACKOFF_MAX);
tokio::time::sleep(delay).await;
}
/// Whether the drain loop finished or needs the table re-claimed first.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CheckpointOutcome {
Done,
ReissueFromFlush,
}
/// What one LSM request produced: its value, or word that the owning node
/// holds no claim and only `flush` can get it back.
enum Attempt<T> {
Ok(T),
ReissueFromFlush,
}
/// Issue one LSM request, retrying in place while the fault is retryable.
///
/// The two recoverable faults have separate budgets: contention clears on its
/// own and retries here against [`MAX_RETRIES`], while a 421 needs `flush` to
/// re-claim, which only the caller can drive.
///
/// An exhausted budget propagates the last error *as itself* rather than a
/// synthesized one — "429 after nine tries" beats "checkpoint failed", and a
/// draining node arrives carrying the server's own message.
async fn issue<T, F, Fut>(mut call: F) -> Result<Attempt<T>>
where
F: FnMut() -> Fut,
Fut: Future<Output = Result<T>>,
{
let mut retries = 0;
loop {
let e = match call().await {
Ok(value) => return Ok(Attempt::Ok(value)),
Err(e) => e,
};
if is_lost_claim(&e) {
return Ok(Attempt::ReissueFromFlush);
}
if !is_retryable(&e) || retries >= MAX_RETRIES {
return Err(e);
}
backoff(retries).await;
retries += 1;
}
}
/// Drive [`Table::checkpoint_lsm`]: seal once, fix the target watermark
/// from the resulting L0, then trigger and poll until it drains.
pub(crate) async fn checkpoint_lsm(table: &Table) -> Result<()> {
for reissue in 0..=MAX_REISSUES {
// The seal turns everything written before this call into a
// generation, so the watermark has to be read after it. Idempotent:
// sealing an empty memtable is a no-op, so a re-issue does not churn
// empty generations.
match issue(|| table.flush_lsm()).await? {
Attempt::Ok(()) => {}
Attempt::ReissueFromFlush => {
backoff(reissue).await;
continue;
}
}
let stats = match issue(|| table.get_lsm_stats(false)).await? {
Attempt::Ok(stats) => stats,
Attempt::ReissueFromFlush => {
backoff(reissue).await;
continue;
}
};
let Some(stats) = stats else {
// Not WAL-backed; `flush_lsm` would have errored first but for a race.
return Ok(());
};
let targets: HashMap<String, u64> = stats
.buckets
.iter()
.filter_map(|b| Some((b.shard_id.clone(), b.newest_generation()?)))
.collect();
if targets.is_empty() {
return Ok(());
}
match drain_to_targets(table, &targets).await? {
CheckpointOutcome::Done => return Ok(()),
CheckpointOutcome::ReissueFromFlush => {
backoff(reissue).await;
continue;
}
}
}
Err(Error::Runtime {
message: "checkpoint_lsm: the owning node kept losing its claim; \
re-issued from flush the maximum number of times"
.into(),
})
}
/// Trigger and poll until no bucket holds a generation at or below its
/// target.
///
/// No liveness bound, deliberately. The pod-wide compactor pool (a semaphore
/// of 2 by default, shared across every table on the node) is taken *inside*
/// the pass, after the bucket latch, so a checkpoint queued behind unrelated
/// tables is indistinguishable from one that is merging. An idle-poll counter
/// here could only ever have fired on a table that would have finished.
async fn drain_to_targets(
table: &Table,
targets: &HashMap<String, u64>,
) -> Result<CheckpointOutcome> {
loop {
let stats = match issue(|| table.get_lsm_stats(false)).await? {
Attempt::Ok(stats) => stats,
Attempt::ReissueFromFlush => return Ok(CheckpointOutcome::ReissueFromFlush),
};
let Some(stats) = stats else {
return Ok(CheckpointOutcome::Done);
};
// `compacting` is the bucket's compaction latch, held from dispatch
// until the pass ends — including while it waits on the pod-wide
// permit. So it answers one question only: do not pile on. Buckets
// with nothing outstanding are skipped, not counted as idle.
let mut outstanding = 0;
let mut all_compacting = true;
for b in &stats.buckets {
let Some(target) = targets.get(&b.shard_id) else {
continue;
};
let n = b.outstanding_generations(*target);
if n > 0 {
outstanding += n;
all_compacting &= b.compacting;
}
}
if outstanding == 0 {
return Ok(CheckpointOutcome::Done);
}
if !all_compacting {
match table.compact_lsm().await {
Ok(()) => {}
Err(e) if is_lost_claim(&e) => return Ok(CheckpointOutcome::ReissueFromFlush),
Err(e) if !is_retryable(&e) => return Err(e),
// A 429 here means the server could latch no bucket at all,
// which the poll above already handles. Not retried in place:
// the latch it would contend for is the one doing the work, so
// fall through and re-read — `POLL_INTERVAL` is the backoff.
Err(_) => {}
}
}
tokio::time::sleep(POLL_INTERVAL).await;
}
}
#[cfg(all(test, feature = "remote"))]
mod tests {
use super::*;
fn http(status: u16) -> Error {
Error::Http {
source: "server said no".into(),
request_id: "rid".into(),
status_code: reqwest::StatusCode::from_u16(status).ok(),
}
}
/// Every status the loop acts on. The two predicates are checked together
/// because their overlap is what would be wrong: a status must never be
/// both, and 421 in particular must not read as retryable — retrying it in
/// place re-issues the call that just said the node holds no claim.
#[test]
fn taxonomy_round_trips() {
for status in [429, 503] {
assert!(is_retryable(&http(status)), "{status} must retry");
assert!(
!is_lost_claim(&http(status)),
"{status} is not a lost claim"
);
}
assert!(is_lost_claim(&http(421)), "a lost claim must re-claim");
assert!(
!is_retryable(&http(421)),
"retrying a lost claim in place only asks the same node again"
);
for status in [400, 404, 409, 500] {
assert!(!is_retryable(&http(status)), "{status} is terminal");
assert!(!is_lost_claim(&http(status)), "{status} is terminal");
}
}
/// An error carrying no status has no retry story and must be terminal —
/// a connection that never reached the server, or a `TableNotFound` that
/// `check_table_response` translated before the loop saw it.
#[test]
fn errors_without_a_status_are_terminal() {
let no_status = Error::Http {
source: "connection reset".into(),
request_id: "rid".into(),
status_code: None,
};
assert!(!is_retryable(&no_status));
assert!(!is_lost_claim(&no_status));
let translated = Error::TableNotFound {
name: "t".into(),
source: "gone".into(),
};
assert!(!is_retryable(&translated));
assert!(!is_lost_claim(&translated));
}
}
-162
View File
@@ -1,162 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
//! Live per-bucket LSM state — the shape [`crate::Table::get_lsm_stats`]
//! returns and [`super::checkpoint`] polls.
//!
//! Nothing here is derived: sums and differences (total L0 bytes, WAL lag)
//! are the caller's to compute. There is no "WAL is off" shape — that case is
//! `None`, because a struct of zeros would read as measurements.
use serde::Deserialize;
/// One flushed L0 generation.
#[derive(Debug, Clone, Deserialize)]
pub struct GenerationStats {
pub generation: u64,
pub bytes: u64,
/// Present only when `include_generation_rows` was requested. Off by
/// default because each count opens an uncached Lance dataset, and the
/// checkpoint loop polls this route needing only generation numbers.
#[serde(default)]
pub rows: Option<u64>,
}
/// One in-memory memtable.
#[derive(Debug, Clone, Deserialize)]
pub struct MemtableStats {
pub generation: u64,
pub rows: u64,
pub bytes: u64,
pub batches: u64,
/// Names of the indexes this memtable carries. An absent name is the whole
/// answer to "why is my fresh-tier search on that column brute-force".
pub indexes: Vec<String>,
}
/// Live state of one bucket. A table is N buckets on one node; flattening to
/// a single number hides the one hot bucket that is usually why someone
/// opened this endpoint.
#[derive(Debug, Clone, Deserialize)]
pub struct BucketStats {
pub shard_id: String,
/// `Active` | `Sealed` (drop-table 2PC in flight).
pub status: String,
pub writer_epoch: u64,
pub manifest_version: u64,
pub current_generation: u64,
pub replay_after_wal_entry_position: u64,
pub wal_entry_position_last_seen: u64,
pub generations: Vec<GenerationStats>,
/// Whether a pass owns this bucket's compaction latch right now. Says *a*
/// driver is running, not *whose*, and the latch is held from dispatch —
/// including while the pass queues for a pod-wide compactor permit. Read
/// it as "do not pile on", never as "mine is progressing".
pub compacting: bool,
/// Oldest first, active last. Absent for a `Sealed` bucket, whose
/// in-memory state is torn down.
#[serde(default)]
pub memtables: Option<Vec<MemtableStats>>,
}
impl BucketStats {
/// The newest flushed generation, or `None` when L0 is empty.
pub(crate) fn newest_generation(&self) -> Option<u64> {
self.generations.iter().map(|g| g.generation).max()
}
/// How many generations at or below `target` are still in L0.
///
/// A count, not a boolean: one pass drains a bounded prefix rather than
/// the whole target set, so a boolean would read as "no progress" for
/// every pass but the last. Compaction drains oldest-first, so this
/// decreases monotonically.
pub(crate) fn outstanding_generations(&self, target: u64) -> usize {
self.generations
.iter()
.filter(|g| g.generation <= target)
.count()
}
}
/// Live LSM state, one entry per bucket.
#[derive(Debug, Clone, Deserialize)]
pub struct LsmStats {
pub buckets: Vec<BucketStats>,
}
/// Server-side JSON envelope for `get_lsm_stats`. `lsm_stats` is null when
/// the table has no LSM write path.
#[derive(Debug, Deserialize)]
pub(crate) struct GetLsmStatsResponse {
#[serde(default)]
pub lsm_stats: Option<LsmStats>,
}
#[cfg(test)]
mod tests {
use super::*;
fn bucket(shard: &str, generations: &[u64], compacting: bool) -> BucketStats {
BucketStats {
shard_id: shard.into(),
status: "Active".into(),
writer_epoch: 1,
manifest_version: 1,
current_generation: generations.iter().max().copied().unwrap_or(0) + 1,
replay_after_wal_entry_position: 0,
wal_entry_position_last_seen: 0,
generations: generations
.iter()
.map(|g| GenerationStats {
generation: *g,
bytes: 1,
rows: None,
})
.collect(),
compacting,
memtables: None,
}
}
/// The target watermark is the newest generation at the start, and a
/// generation created after it must not hold the loop open — that is why
/// the predicate terminates under write load.
#[test]
fn newer_generations_do_not_extend_the_target() {
let start = bucket("b0", &[7, 8], false);
let target = start.newest_generation().expect("L0 is non-empty");
assert_eq!(target, 8);
// Compaction drained 7 and 8; 9 and 10 arrived while it ran.
let later = bucket("b0", &[9, 10], false);
assert_eq!(
later.outstanding_generations(target),
0,
"generations above the target are somebody else's problem"
);
// Still holding 8 means still outstanding.
assert_eq!(
bucket("b0", &[8, 9], false).outstanding_generations(target),
1
);
}
/// The metric counts generations, not buckets: a pass drains a bounded
/// prefix, so one bucket going 3 → 2 → 1 → 0 is three steps.
#[test]
fn progress_is_measured_in_generations() {
let target = 3;
let counts: Vec<usize> = [&[1u64, 2, 3][..], &[2, 3][..], &[3][..], &[][..]]
.iter()
.map(|gens| bucket("b0", gens, false).outstanding_generations(target))
.collect();
assert_eq!(counts, vec![3, 2, 1, 0]);
}
#[test]
fn empty_l0_has_no_target() {
assert!(bucket("b0", &[], false).newest_generation().is_none());
}
}
+1 -70
View File
@@ -315,10 +315,7 @@ pub(crate) async fn execute_merge_insert(
#[cfg(test)]
mod tests {
use arrow_array::builder::FixedSizeBinaryBuilder;
use arrow_array::{
Int32Array, RecordBatch, RecordBatchIterator, RecordBatchReader, StringArray, UInt64Array,
};
use arrow_array::{Int32Array, RecordBatch, RecordBatchIterator, RecordBatchReader};
use arrow_schema::{DataType, Field, Schema};
use std::sync::Arc;
@@ -340,42 +337,6 @@ mod tests {
Box::new(RecordBatchIterator::new(vec![Ok(batch)], schema))
}
fn fixed_size_binary_merge_batch(
id_range: std::ops::Range<u64>,
price: u64,
) -> Box<dyn RecordBatchReader + Send> {
let ids = id_range.collect::<Vec<_>>();
let mut id_builder = FixedSizeBinaryBuilder::new(16);
for id in &ids {
let mut bytes = [0; 16];
bytes[..8].copy_from_slice(&id.to_le_bytes());
id_builder.append_value(bytes).unwrap();
}
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::FixedSizeBinary(16), false),
Field::new("id_as_int", DataType::UInt64, false),
Field::new("name", DataType::Utf8, false),
Field::new("market", DataType::Utf8, false),
]));
let batch = RecordBatch::try_new(
schema.clone(),
vec![
Arc::new(id_builder.finish()),
Arc::new(UInt64Array::from_iter_values(ids.iter().copied())),
Arc::new(StringArray::from_iter_values(
ids.iter().map(|id| format!("name{id}")),
)),
Arc::new(StringArray::from_iter_values(std::iter::repeat_n(
format!("market_{price}"),
ids.len(),
))),
],
)
.unwrap();
Box::new(RecordBatchIterator::new(vec![Ok(batch)], schema))
}
#[tokio::test]
async fn test_merge_insert() {
let conn = connect("memory://").execute().await.unwrap();
@@ -427,36 +388,6 @@ mod tests {
);
}
#[tokio::test]
async fn test_merge_insert_fixed_size_binary_non_nullable() {
// Regression test for #2869: an unrelated FixedSizeBinary column used to corrupt the
// outer join that implements when_not_matched_by_source_delete.
let conn = connect("memory://").execute().await.unwrap();
let table = conn
.create_table(
"fixed_size_binary_merge",
fixed_size_binary_merge_batch(0..256, 100),
)
.execute()
.await
.unwrap();
let mut merge_insert = table.merge_insert(&["id_as_int"]);
merge_insert
.when_matched_update_all(None)
.when_not_matched_insert_all()
.when_not_matched_by_source_delete(None);
let result = merge_insert
.execute(fixed_size_binary_merge_batch(100..356, 200))
.await
.unwrap();
assert_eq!(result.num_updated_rows, 156);
assert_eq!(result.num_inserted_rows, 100);
assert_eq!(result.num_deleted_rows, 100);
assert_eq!(table.count_rows(None).await.unwrap(), 256);
}
#[tokio::test]
async fn test_merge_insert_use_index() {
let conn = connect("memory://").execute().await.unwrap();
+1
View File
@@ -18,6 +18,7 @@
use std::collections::HashMap;
use std::sync::Arc;
use std::time::Duration;
use arrow_array::cast::AsArray;
+1 -148
View File
@@ -214,17 +214,12 @@ pub(crate) async fn execute_optimize(
#[cfg(test)]
mod tests {
use arrow_array::{
Array, FixedSizeListArray, Float32Array, Int32Array, RecordBatch, StringArray,
};
use arrow_array::{Int32Array, RecordBatch, StringArray};
use arrow_schema::{DataType, Field, Schema};
use lance_arrow::FixedSizeListArrayExt;
use rstest::rstest;
use std::sync::Arc;
use crate::connect;
use crate::database::listing::OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS;
use crate::index::vector::IvfRqIndexBuilder;
use crate::index::{Index, scalar::BTreeIndexBuilder};
use crate::query::ExecutableQuery;
use crate::table::{CompactionOptions, OptimizeAction, OptimizeStats};
@@ -309,96 +304,6 @@ mod tests {
assert_eq!(all_values, expected);
}
#[tokio::test]
async fn test_compact_with_concurrent_add() {
const NUM_FRAGMENTS: usize = 5;
const ROWS_PER_FRAGMENT: i32 = 300;
let tmpdir = tempfile::tempdir().unwrap();
let conn = connect(tmpdir.path().to_str().unwrap())
.execute()
.await
.unwrap();
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
let batch = RecordBatch::try_new(
schema,
vec![Arc::new(Int32Array::from_iter_values(0..ROWS_PER_FRAGMENT))],
)
.unwrap();
let table = conn
.create_table("test_concurrent_compact", batch.clone())
.execute()
.await
.unwrap();
table
.create_index(&["id"], Index::BTree(BTreeIndexBuilder::default()))
.execute()
.await
.unwrap();
for _ in 0..NUM_FRAGMENTS {
table.add(batch.clone()).execute().await.unwrap();
}
// Use separate handles so the two writes actually overlap, as they can
// when different Node connections operate on the same S3 table.
let compact_table = conn
.open_table("test_concurrent_compact")
.execute()
.await
.unwrap();
let append_table = conn
.open_table("test_concurrent_compact")
.execute()
.await
.unwrap();
let compact_task = tokio::spawn(async move {
compact_table
.optimize(OptimizeAction::Compact {
options: CompactionOptions {
target_rows_per_fragment: 1_000,
..Default::default()
},
remap_options: None,
})
.await
});
tokio::task::yield_now().await;
for _ in 0..NUM_FRAGMENTS {
append_table.add(batch.clone()).execute().await.unwrap();
}
compact_task.await.unwrap().unwrap();
let table = conn
.open_table("test_concurrent_compact")
.execute()
.await
.unwrap();
let dataset = table.dataset().unwrap().get().await.unwrap();
let fragment_ids = dataset
.get_fragments()
.iter()
.map(|fragment| fragment.id())
.collect::<Vec<_>>();
assert!(fragment_ids.windows(2).all(|ids| ids[0] < ids[1]));
// A second compaction exposed the original out-of-order row-id bug.
table
.optimize(OptimizeAction::Compact {
options: CompactionOptions {
target_rows_per_fragment: 1_000,
..Default::default()
},
remap_options: None,
})
.await
.unwrap();
assert_eq!(
table.count_rows(None).await.unwrap(),
ROWS_PER_FRAGMENT as usize * (NUM_FRAGMENTS * 2 + 1)
);
}
#[tokio::test]
async fn test_optimize_prune_versions() {
let conn = connect("memory://").execute().await.unwrap();
@@ -537,58 +442,6 @@ mod tests {
assert_eq!(final_row_count, 200);
}
#[tokio::test]
async fn test_optimize_vector_index_after_delete_with_stable_row_ids() {
const NUM_ROWS: i32 = 400;
const DIMENSION: i32 = 32;
let conn = connect("memory://").execute().await.unwrap();
let vectors = FixedSizeListArray::try_new_from_values(
Float32Array::from_iter_values((0..NUM_ROWS).flat_map(|id| {
(0..DIMENSION).map(move |offset| ((id as f32 * 0.1) + (offset as f32 * 0.3)).sin())
})),
DIMENSION,
)
.unwrap();
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int32, false),
Field::new("vector", vectors.data_type().clone(), false),
]));
let batch = RecordBatch::try_new(
schema,
vec![
Arc::new(Int32Array::from_iter_values(0..NUM_ROWS)),
Arc::new(vectors),
],
)
.unwrap();
let table = conn
.create_table("test_vector_index_optimize_after_delete", batch)
.storage_option(OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS, "true")
.execute()
.await
.unwrap();
table
.create_index(
&["vector"],
Index::IvfRq(IvfRqIndexBuilder::default().num_partitions(4)),
)
.execute()
.await
.unwrap();
table.delete("id % 3 = 0").await.unwrap();
// Regression test for #3330: deleted stable row IDs used to become
// misaligned with row addresses while joining small IVF partitions.
table
.optimize(OptimizeAction::Index(Default::default()))
.await
.unwrap();
assert_eq!(table.count_rows(None).await.unwrap(), 266);
}
#[tokio::test]
async fn test_optimize_all() {
let conn = connect("memory://").execute().await.unwrap();
+166 -54
View File
@@ -84,9 +84,8 @@ pub(super) async fn create_lsm_plan(
let pk_columns = pk_columns(&ds_ref)?;
// The base index an indexed arm relies on may lag compaction; resolve it so the
// snapshot retains SSTables the index has not yet caught up to.
let arm_index = arm_maintained_index_name(&ds_ref, &query, &details).await?;
let (snapshots, in_memory) =
build_read_context(table, &ds_ref, &details, arm_index.as_deref()).await?;
let arm_indexes = arm_maintained_index_names(&ds_ref, &query, &details).await?;
let (snapshots, in_memory) = build_read_context(table, &ds_ref, &details, &arm_indexes).await?;
let limit = query.base.limit;
let offset = query.base.offset;
@@ -232,28 +231,40 @@ fn pk_columns(dataset: &Dataset) -> Result<Vec<String>> {
Ok(pk)
}
/// Per-shard SSTable exclusion watermark: the generation at or below which SSTables
/// are safe to drop for this arm. A generation is droppable only once it is
/// compacted into the base table AND covered by `index_name`'s catch-up (for an
/// indexed arm); a plain scan (`index_name == None`) uses the compaction watermark
/// alone. Capping at the index catch-up keeps rows the base index has not yet
/// indexed visible through their SSTable. First occurrence per shard mirrors Lance's
/// `compacted_generation_for_shard`.
/// Per-shard SSTable exclusion watermark: the generation at or below which
/// SSTables are safe to drop for this query.
///
/// A generation is droppable only once it is compacted into the base table AND
/// covered by the catch-up of every index the query relies on, so the watermark
/// is the minimum across `index_names`. Gating on fewer than all of them would
/// drop SSTables holding rows an uncounted index has not yet indexed, and that
/// arm would silently return fewer rows.
///
/// See [`arm_maintained_index_names`] for which indexes are collected today: a
/// vector search with a scalar prefilter is not yet among them.
///
/// An empty `index_names` (a plain scan) uses the compaction watermark alone.
/// First occurrence per shard mirrors Lance's `compacted_generation_for_shard`.
fn exclusion_watermarks(
details: &MemWalIndexDetails,
index_name: Option<&str>,
index_names: &[String],
) -> HashMap<Uuid, u64> {
let mut exclude: HashMap<Uuid, u64> = HashMap::new();
for entry in &details.compacted_sstables {
let mut watermark = entry.generation;
if let Some(name) = index_name
&& let Some(caught_up) = details
for name in index_names {
match details
.index_catchup
.iter()
.find(|icp| icp.index_name == name)
.find(|icp| icp.index_name == *name)
.and_then(|icp| icp.caught_up_generation_for_shard(&entry.shard_id))
{
watermark = watermark.min(caught_up);
{
Some(caught_up) => watermark = watermark.min(caught_up),
// No entry means the index is *not* known to hold these rows,
// and the base arm is index-only -- so every generation stays
// readable from its SSTable.
None => watermark = 0,
}
}
exclude.entry(entry.shard_id).or_insert(watermark);
}
@@ -271,9 +282,9 @@ async fn build_read_context(
table: &NativeTable,
dataset: &Dataset,
details: &MemWalIndexDetails,
index_name: Option<&str>,
index_names: &[String],
) -> Result<(Vec<ShardSnapshot>, HashMap<Uuid, InMemoryMemTables>)> {
let exclude = exclusion_watermarks(details, index_name);
let exclude = exclusion_watermarks(details, index_names);
let shard_ids = dataset.list_mem_wal_latest_shard_ids().await?;
// Use the dataset's own object store (not `ObjectStore::from_uri`, which
@@ -487,19 +498,33 @@ async fn index_maintained(
}))
}
/// The maintained base index the query's arm relies on (vector index for ANN, FTS
/// index for full-text), used to gate SSTable compaction exclusion by index catch-up.
/// `None` for a plain scan or when no maintained index covers the searched column.
async fn arm_maintained_index_name(
/// Every maintained base index this query relies on, used to gate SSTable
/// exclusion by index catch-up.
///
/// Returns a list because the watermark must be the lowest across every index a
/// query relies on. Today it never holds more than one: `reject_unsupported`
/// refuses hybrid search, so the vector and full-text arms are mutually
/// exclusive.
///
/// The case that is genuinely multi-index -- a vector search with a scalar or
/// bitmap prefilter -- is **not collected yet**. Identifying those needs the
/// planner's chosen indexes, not the columns the filter names, and no Lance API
/// exposes them. Until it does, such a query is gated on its vector index alone.
///
/// Empty for a plain scan, or when no maintained index covers the searched
/// column.
async fn arm_maintained_index_names(
dataset: &Dataset,
query: &VectorQueryRequest,
details: &MemWalIndexDetails,
) -> Result<Option<String>> {
) -> Result<Vec<String>> {
use lance::index::DatasetIndexExt;
// Resolve the arm's searched column, the index-detail type it relies on, and a
// Each arm's searched column, the index-detail type it relies on, and a
// label for diagnostics — catch-up is taken from the vector/FTS index
// specifically, not a BTree on the same column.
let (column, type_url_suffix, arm) = if !query.query_vector.is_empty() {
let mut arms: Vec<(String, &str, &str)> = Vec::new();
if !query.query_vector.is_empty() {
let arrow_schema = ArrowSchema::from(dataset.schema());
let column = match &query.column {
Some(column) => column.clone(),
@@ -508,31 +533,43 @@ async fn arm_maintained_index_name(
default_vector_column(&arrow_schema, dim)?
}
};
(column, "VectorIndexDetails", "vector")
} else if let Some(fts) = &query.base.full_text_search {
match fts.columns().into_iter().next() {
Some(column) => (column, "InvertedIndexDetails", "full-text"),
None => return Ok(None),
}
} else {
return Ok(None);
};
let Some(field) = dataset.schema().field(&column) else {
return Ok(None);
};
arms.push((column, "VectorIndexDetails", "vector"));
}
if let Some(fts) = &query.base.full_text_search
&& let Some(column) = fts.columns().into_iter().next()
{
arms.push((column, "InvertedIndexDetails", "full-text"));
}
if arms.is_empty() {
return Ok(Vec::new());
}
let indices = dataset.load_indices().await?;
let segment_names: Vec<String> = indices
.iter()
.filter(|idx| {
idx.fields.contains(&field.id)
&& idx
.index_details
.as_ref()
.is_some_and(|d| d.type_url.ends_with(type_url_suffix))
})
.map(|idx| idx.name.clone())
.collect();
resolve_single_index(segment_names, &details.maintained_indexes, arm, &column)
let mut names = Vec::with_capacity(arms.len());
for (column, type_url_suffix, arm) in arms {
let Some(field) = dataset.schema().field(&column) else {
continue;
};
let segment_names: Vec<String> = indices
.iter()
.filter(|idx| {
idx.fields.contains(&field.id)
&& idx
.index_details
.as_ref()
.is_some_and(|d| d.type_url.ends_with(type_url_suffix))
})
.map(|idx| idx.name.clone())
.collect();
if let Some(name) =
resolve_single_index(segment_names, &details.maintained_indexes, arm, &column)?
{
names.push(name);
}
}
names.sort();
names.dedup();
Ok(names)
}
/// Resolve the single logical index from the names of its matching physical
@@ -734,24 +771,99 @@ mod tests {
};
// Plain scan: drop every compacted generation (through 5).
assert_eq!(exclusion_watermarks(&details, None).get(&shard), Some(&5));
assert_eq!(exclusion_watermarks(&details, &[]).get(&shard), Some(&5));
// FTS arm with a lagging index: exclusion is capped at the index catch-up
// (2), so SSTable generations 3..=5 are retained until the index covers
// them — otherwise those documents would silently vanish from FTS results.
assert_eq!(
exclusion_watermarks(&details, Some("fts_idx")).get(&shard),
exclusion_watermarks(&details, &["fts_idx".to_string()]).get(&shard),
Some(&2)
);
// A caught-up index — or one untracked in index_catchup — falls back to the
// compaction watermark.
// An index with no entry has not recorded that it holds these rows, so
// nothing is excluded. This is the case a table written before catch-up
// was maintained lands in, and it errs toward reading the SSTables.
assert_eq!(
exclusion_watermarks(&details, Some("caught_up_idx")).get(&shard),
exclusion_watermarks(&details, &["untracked_idx".to_string()]).get(&shard),
Some(&0)
);
// An index recorded as covering the compaction watermark excludes up to it.
let caught_up = MemWalIndexDetails {
index_catchup: vec![IndexCatchupProgress::new(
"caught_up_idx".to_string(),
vec![CompactedSsTable::new(shard, 5)],
)],
..details.clone()
};
assert_eq!(
exclusion_watermarks(&caught_up, &["caught_up_idx".to_string()]).get(&shard),
Some(&5)
);
}
/// A hybrid search reads a vector and a full-text index, and either may lag.
/// Retaining to the lower of the two is what keeps both arms complete;
/// gating on one alone would drop SSTables the other has not indexed.
#[test]
fn exclusion_watermark_takes_the_minimum_across_every_index_used() {
let shard = Uuid::from_u128(1);
let details = MemWalIndexDetails {
compacted_sstables: vec![CompactedSsTable::new(shard, 9)],
index_catchup: vec![
IndexCatchupProgress::new(
"vec_idx".to_string(),
vec![CompactedSsTable::new(shard, 7)],
),
IndexCatchupProgress::new(
"fts_idx".to_string(),
vec![CompactedSsTable::new(shard, 4)],
),
],
maintained_indexes: vec!["vec_idx".to_string(), "fts_idx".to_string()],
..Default::default()
};
// Each index alone stops at its own catch-up.
assert_eq!(
exclusion_watermarks(&details, &["vec_idx".to_string()]).get(&shard),
Some(&7)
);
assert_eq!(
exclusion_watermarks(&details, &["fts_idx".to_string()]).get(&shard),
Some(&4)
);
// Used together, the lower one governs regardless of order.
let both = ["vec_idx".to_string(), "fts_idx".to_string()];
assert_eq!(exclusion_watermarks(&details, &both).get(&shard), Some(&4));
let reversed = ["fts_idx".to_string(), "vec_idx".to_string()];
assert_eq!(
exclusion_watermarks(&details, &reversed).get(&shard),
Some(&4)
);
}
/// An index with no catch-up entry is not known to hold anything, so it
/// governs over a lagging sibling rather than the other way round.
#[test]
fn an_untracked_index_retains_everything() {
let shard = Uuid::from_u128(1);
let details = MemWalIndexDetails {
compacted_sstables: vec![CompactedSsTable::new(shard, 9)],
index_catchup: vec![IndexCatchupProgress::new(
"fts_idx".to_string(),
vec![CompactedSsTable::new(shard, 4)],
)],
maintained_indexes: vec!["fts_idx".to_string(), "untracked_idx".to_string()],
..Default::default()
};
let both = ["fts_idx".to_string(), "untracked_idx".to_string()];
assert_eq!(exclusion_watermarks(&details, &both).get(&shard), Some(&0));
}
#[test]
fn resolve_single_index_dedupes_segments() {
let maintained = vec!["fts_idx".to_string()];
+19 -24
View File
@@ -193,12 +193,10 @@ mod tests {
// Add a computed column
let result = table
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![(
"doubled".into(),
"id * 2".into(),
)]))
.execute()
.add_columns(
NewColumnTransform::SqlExpressions(vec![("doubled".into(), "id * 2".into())]),
None,
)
.await
.unwrap();
@@ -253,12 +251,13 @@ mod tests {
// Add multiple columns at once
table
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![
("y".into(), "x + 1".into()),
("z".into(), "x * x".into()),
]))
.execute()
.add_columns(
NewColumnTransform::SqlExpressions(vec![
("y".into(), "x + 1".into()),
("z".into(), "x * x".into()),
]),
None,
)
.await
.unwrap();
@@ -284,12 +283,10 @@ mod tests {
// Add a column with a constant value
table
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![(
"constant".into(),
"42".into(),
)]))
.execute()
.add_columns(
NewColumnTransform::SqlExpressions(vec![("constant".into(), "42".into())]),
None,
)
.await
.unwrap();
@@ -662,12 +659,10 @@ mod tests {
// Add column increments version
let add_result = table
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![(
"c".into(),
"a + b".into(),
)]))
.execute()
.add_columns(
NewColumnTransform::SqlExpressions(vec![("c".into(), "a + b".into())]),
None,
)
.await
.unwrap();
assert!(add_result.version > v1);
+3 -255
View File
@@ -9,17 +9,14 @@ use arrow_array::{
};
use arrow_schema::{DataType, Field, Fields, Schema};
use futures::TryStreamExt;
use lance::Dataset;
use lance_file::version::LanceFileVersion;
use lance_encoding::version::LanceFileVersion;
use lancedb::{
Connection, Error, Result, Table,
blob::{BlobRangeRequest, blob},
connect, connect_namespace,
database::listing::{
ListingDatabaseOptions, NewTableConfig, OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS,
},
database::listing::OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS,
query::{ExecutableQuery, QueryBase},
table::{AddDataMode, CompactionOptions, OptimizeAction, OptimizeStats},
table::{AddDataMode, CompactionOptions, OptimizeAction},
};
use tempfile::tempdir;
@@ -1078,252 +1075,3 @@ async fn fetch_blob_files_aligns_across_fragments_with_nulls_and_dups() -> Resul
}
Ok(())
}
/// Rows exercising the null/empty interleavings from
/// <https://github.com/lancedb/lancedb/issues/3744>: a payload, a null, a valid
/// empty value, then payloads whose descriptors a fragment rewrite used to zero.
fn null_empty_input_batch() -> RecordBatch {
let owned = [
Some(dedicated_blob_bytes(1)),
None,
Some(Vec::new()),
Some(dedicated_blob_bytes(4)),
Some(dedicated_blob_bytes(5)),
Some(dedicated_blob_bytes(6)),
];
let payloads: Vec<Option<&[u8]>> = owned.iter().map(|payload| payload.as_deref()).collect();
binary_input_batch(&[1, 2, 3, 4, 5, 6], &payloads)
}
/// One `(id, Some((payload length, first byte)))` per live row, or `(id, None)`
/// for a null blob. Comparing lengths and first bytes keeps failure output
/// readable where comparing whole payloads would not.
type BlobSummary = Vec<(i64, Option<(usize, Option<u8>)>)>;
/// The rows [`null_empty_input_batch`] leaves behind after `id IN (1, 4)` is
/// deleted: a null, a valid empty value, and the two payloads that follow them.
fn expected_null_empty_survivors() -> BlobSummary {
vec![
(2, None),
(3, Some((0, None))),
(5, Some((DEDICATED_BLOB_LEN, Some(5)))),
(6, Some((DEDICATED_BLOB_LEN, Some(6)))),
]
}
/// `optimize()` only rewrites a fragment when lance's compaction planner selects
/// it — here because the delete pushes the fragment past
/// `materialize_deletions_threshold` (0.1 by default; these tests delete 2 of 6
/// rows). Without this check, a planner or threshold change upstream would leave
/// both regression tests green while no rewrite happened at all.
fn assert_compacted(stats: &OptimizeStats) {
let metrics = stats
.compaction
.as_ref()
.expect("OptimizeAction::All runs compaction");
assert!(
metrics.fragments_removed >= 1,
"optimize() rewrote no fragment, so this test proves nothing: {metrics:?}"
);
}
fn summarize(rows: &[(i64, Option<Vec<u8>>)]) -> BlobSummary {
rows.iter()
.map(|(id, payload)| {
(
*id,
payload
.as_ref()
.map(|bytes| (bytes.len(), bytes.first().copied())),
)
})
.collect()
}
async fn sorted_id_rowid(table: &Table) -> Result<Vec<(i64, u64)>> {
let mut pairs = collect_id_rowid(table).await?;
pairs.sort_by_key(|(id, _)| *id);
Ok(pairs)
}
/// `{position, size}` descriptors of a legacy v1 blob column, keyed by `id`.
async fn v1_blob_descriptors(table: &Table) -> Result<Vec<(i64, Option<(u64, u64)>)>> {
let batches = table
.query()
.execute()
.await?
.try_collect::<Vec<_>>()
.await?;
let batch = arrow_select::concat::concat_batches(&batches[0].schema(), &batches).unwrap();
let ids = batch
.column_by_name("id")
.unwrap()
.as_any()
.downcast_ref::<Int64Array>()
.unwrap();
let descriptors = batch
.column_by_name("image")
.unwrap()
.as_any()
.downcast_ref::<StructArray>()
.expect("v1 blob column reads back as a descriptor struct");
let position = descriptors
.column_by_name("position")
.unwrap()
.as_any()
.downcast_ref::<UInt64Array>()
.unwrap();
let size = descriptors
.column_by_name("size")
.unwrap()
.as_any()
.downcast_ref::<UInt64Array>()
.unwrap();
let mut rows: Vec<(i64, Option<(u64, u64)>)> = (0..batch.num_rows())
.map(|row| {
let descriptor =
(!descriptors.is_null(row)).then(|| (position.value(row), size.value(row)));
(ids.value(row), descriptor)
})
.collect();
rows.sort_by_key(|(id, _)| *id);
Ok(rows)
}
/// Payload bytes of every live row of a legacy v1 blob column, keyed by `id`.
/// [`Table::fetch_blobs`] rejects v1 columns, so read them through lance.
async fn v1_blob_payloads(dataset_uri: &str, table: &Table) -> Result<Vec<(i64, Option<Vec<u8>>)>> {
let pairs = sorted_id_rowid(table).await?;
let row_ids: Vec<u64> = pairs.iter().map(|(_, row_id)| *row_id).collect();
let dataset = Arc::new(Dataset::open(dataset_uri).await?);
let files = dataset.take_blobs(&row_ids, "image").await?;
assert_eq!(
files.len(),
pairs.len(),
"take_blobs returned {} handles for {} live rows",
files.len(),
pairs.len()
);
let mut rows = Vec::with_capacity(pairs.len());
for ((id, _), file) in pairs.iter().zip(files) {
let payload = match file {
Some(file) => Some(file.read().await?.to_vec()),
None => None,
};
rows.push((*id, payload));
}
Ok(rows)
}
/// Length and first byte of every live blob v2 value, keyed by `id`.
async fn blob_v2_values(table: &Table) -> Result<BlobSummary> {
let pairs = sorted_id_rowid(table).await?;
let row_ids: Vec<u64> = pairs.iter().map(|(_, row_id)| *row_id).collect();
let bytes = table.fetch_blobs("image", &row_ids).await?;
Ok(pairs
.iter()
.enumerate()
.map(|(slot, (id, _))| {
let value = (!bytes.is_null(slot))
.then(|| (bytes.value(slot).len(), bytes.value(slot).first().copied()));
(*id, value)
})
.collect())
}
/// Regression test for [#3744]: on storage 2.0 (legacy v1 descriptors),
/// compaction rewrote every payload following a null or empty value in the same
/// fragment as `{position: 0, size: 0}`, so the payload bytes read back as `b""`
/// and the new fragment no longer referenced them at all.
///
/// [#3744]: https://github.com/lancedb/lancedb/issues/3744
#[tokio::test]
async fn optimize_preserves_v1_blob_payloads_with_null_and_empty() -> Result<()> {
let tmp = tempdir().unwrap();
let db_uri = tmp.path().to_str().unwrap().to_string();
let db = connect(&db_uri)
.database_options(&ListingDatabaseOptions {
new_table_config: NewTableConfig {
data_storage_version: Some(LanceFileVersion::V2_0),
..Default::default()
},
..Default::default()
})
.execute()
.await?;
let legacy = Field::new("image", DataType::LargeBinary, true).with_metadata(
std::collections::HashMap::from([("lance-encoding:blob".to_string(), "true".to_string())]),
);
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int64, false),
legacy,
]));
let table = db.create_empty_table("t", schema).execute().await?;
table.add(null_empty_input_batch()).execute().await?;
assert_eq!(
storage_format_version(&table).await,
LanceFileVersion::V2_0.resolve(),
"v1 blob descriptors only exist below storage 2.2"
);
let dataset_uri = table.uri().await?;
// Any rewrite triggers it; deleting rows is the shape from the issue.
table.delete("id IN (1, 4)").await?;
let descriptors_before = v1_blob_descriptors(&table).await?;
let before = v1_blob_payloads(&dataset_uri, &table).await?;
assert_eq!(
summarize(&before),
expected_null_empty_survivors(),
"test setup no longer produces the null/empty/payload mix"
);
let stats = table.optimize(OptimizeAction::All).await?;
assert_compacted(&stats);
let descriptors_after = v1_blob_descriptors(&table).await?;
let after = v1_blob_payloads(&dataset_uri, &table).await?;
assert_eq!(
summarize(&after),
summarize(&before),
"optimize() lost blob payloads; descriptors before={descriptors_before:?} after={descriptors_after:?}"
);
assert!(after == before, "optimize() changed blob payload bytes");
Ok(())
}
/// Regression test for the blob v2 half of [#3744]: compaction rewrote a valid
/// empty value as null, destroying the null-vs-empty distinction.
///
/// [#3744]: https://github.com/lancedb/lancedb/issues/3744
#[tokio::test]
async fn optimize_preserves_blob_v2_null_and_empty_distinction() -> Result<()> {
let tmp = tempdir().unwrap();
let db = connect(tmp.path().to_str().unwrap()).execute().await?;
let table = db
.create_empty_table("t", blob_table_schema())
.execute()
.await?;
table.add(null_empty_input_batch()).execute().await?;
assert!(
storage_format_version(&table).await >= LanceFileVersion::V2_2,
"blob v2 columns require storage >= 2.2"
);
table.delete("id IN (1, 4)").await?;
let before = blob_v2_values(&table).await?;
assert_eq!(
before,
expected_null_empty_survivors(),
"test setup no longer produces the null/empty/payload mix"
);
let stats = table.optimize(OptimizeAction::All).await?;
assert_compacted(&stats);
assert_eq!(
blob_v2_values(&table).await?,
before,
"optimize() changed blob v2 values"
);
Ok(())
}