diff --git a/docs/src/js/classes/BlobFile.md b/docs/src/js/classes/BlobFile.md new file mode 100644 index 000000000..84a596d1a --- /dev/null +++ b/docs/src/js/classes/BlobFile.md @@ -0,0 +1,62 @@ +[**@lancedb/lancedb**](../README.md) • **Docs** + +*** + +[@lancedb/lancedb](../globals.md) / BlobFile + +# Class: BlobFile + +A lazy handle to blob bytes. Create one with [Table.fetchBlobFiles](Table.md#fetchblobfiles). + +## Methods + +### read() + +```ts +read(): Promise +``` + +Reads from the cursor to the end and advances the cursor. + +A second call returns an empty buffer. [BlobFile.readRange](BlobFile.md#readrange) does +not move the cursor. + +#### Returns + +`Promise`<`Buffer`> + +*** + +### readRange() + +```ts +readRange(start, end): Promise +``` + +Reads the half-open byte range `[start, end)`. + +Fails when `end` is past the blob size. Does not move the cursor. + +#### Parameters + +* **start**: `bigint` + +* **end**: `bigint` + +#### Returns + +`Promise`<`Buffer`> + +*** + +### size() + +```ts +size(): bigint +``` + +Returns the blob size in bytes. + +#### Returns + +`bigint` diff --git a/docs/src/js/classes/Table.md b/docs/src/js/classes/Table.md index 3847f3a39..ef6e9535a 100644 --- a/docs/src/js/classes/Table.md +++ b/docs/src/js/classes/Table.md @@ -137,6 +137,20 @@ containing the new version number of the table after altering the columns. *** +### blobColumns() + +```ts +abstract blobColumns(): Promise +``` + +Blob v2 columns, including nested dotted paths. + +#### Returns + +`Promise`<`string`[]> + +*** + ### branches() ```ts @@ -499,6 +513,54 @@ Drop an index from the table. *** +### fetchBlobFiles() + +```ts +abstract fetchBlobFiles(column, rowIds): Promise<(null | BlobFile)[]> +``` + +Opens lazy blob handles for `column` at the given row IDs using the +table's current checkout. + +Preserves input order, duplicates, and nulls. Use this for large payloads. +See [Table.fetchBlobs](Table.md#fetchblobs) for row-ID validity across versions. + +#### Parameters + +* **column**: `string` + +* **rowIds**: readonly (`number` \| `bigint`)[] + +#### Returns + +`Promise`<(`null` \| [`BlobFile`](BlobFile.md))[]> + +*** + +### fetchBlobs() + +```ts +abstract fetchBlobs(column, rowIds): Promise<(null | Buffer)[]> +``` + +Bytes for `column` at row IDs from [Query.withRowId](Query.md#withrowid). + +Reads the table's current checkout. IDs from another version can fail after +compaction unless stable row ids are enabled. Results keep input order and +duplicates. Null blobs are `null`. Empty blobs are empty buffers. + +#### Parameters + +* **column**: `string` + +* **rowIds**: readonly (`number` \| `bigint`)[] + +#### Returns + +`Promise`<(`null` \| `Buffer`)[]> + +*** + ### flushLsm() ```ts diff --git a/docs/src/js/functions/blob.md b/docs/src/js/functions/blob.md new file mode 100644 index 000000000..20a734cd4 --- /dev/null +++ b/docs/src/js/functions/blob.md @@ -0,0 +1,55 @@ +[**@lancedb/lancedb**](../README.md) • **Docs** + +*** + +[@lancedb/lancedb](../globals.md) / blob + +# Function: blob() + +```ts +function blob(name, options): Field +``` + +Declares a `lance.blob.v2` column. + +Query results are descriptors, not payload bytes. Use [Table.fetchBlobs](../classes/Table.md#fetchblobs) +or [Table.fetchBlobFiles](../classes/Table.md#fetchblobfiles) to read bytes. + +## Parameters + +* **name**: `string` + +* **options**: [`BlobOptions`](../type-aliases/BlobOptions.md) = `{}` + +## Returns + +`Field` + +## Example + +```ts +import { readFile } from "node:fs/promises"; +import { Field, Int64, Schema } from "apache-arrow"; +import { blob, connect } from "@lancedb/lancedb"; + +const db = await connect("./data"); +const video = await readFile("clip.mp4"); +const table = await db.createTable( + "videos", + [{ id: 1n, video }], + { + schema: new Schema([ + new Field("id", new Int64()), + blob("video"), + ]), + }, +); + +const rows = await table.query().select(["id"]).withRowId().toArray(); +const rowIds = rows.map((row) => row._rowid as bigint); +const bytes = await table.fetchBlobs("video", rowIds); + +const [handle] = await table.fetchBlobFiles("video", rowIds); +const size = handle!.size(); +const header = await handle!.readRange(0n, size < 65536n ? size : 65536n); +``` diff --git a/docs/src/js/functions/isBlobField.md b/docs/src/js/functions/isBlobField.md new file mode 100644 index 000000000..944309f90 --- /dev/null +++ b/docs/src/js/functions/isBlobField.md @@ -0,0 +1,22 @@ +[**@lancedb/lancedb**](../README.md) • **Docs** + +*** + +[@lancedb/lancedb](../globals.md) / isBlobField + +# Function: isBlobField() + +```ts +function isBlobField(field): boolean +``` + +Checks for the `lance.blob.v2` extension marker. Does not validate the +field's storage type. + +## Parameters + +* **field**: `Field`<`any`> + +## Returns + +`boolean` diff --git a/docs/src/js/globals.md b/docs/src/js/globals.md index eb0fc7d5a..4a5effaae 100644 --- a/docs/src/js/globals.md +++ b/docs/src/js/globals.md @@ -19,6 +19,7 @@ ## Classes - [AutoQuery](classes/AutoQuery.md) +- [BlobFile](classes/BlobFile.md) - [BooleanQuery](classes/BooleanQuery.md) - [BoostQuery](classes/BoostQuery.md) - [BranchContents](classes/BranchContents.md) @@ -143,6 +144,7 @@ - [AnalyzePlanDistributedMetrics](type-aliases/AnalyzePlanDistributedMetrics.md) - [BaseTokenizer](type-aliases/BaseTokenizer.md) +- [BlobOptions](type-aliases/BlobOptions.md) - [Data](type-aliases/Data.md) - [DataLike](type-aliases/DataLike.md) - [FieldLike](type-aliases/FieldLike.md) @@ -158,9 +160,11 @@ ## Functions - [RecordBatchIterator](functions/RecordBatchIterator.md) +- [blob](functions/blob.md) - [connect](functions/connect.md) - [connectNamespace](functions/connectNamespace.md) - [instrumentLanceDbMetrics](functions/instrumentLanceDbMetrics.md) +- [isBlobField](functions/isBlobField.md) - [makeArrowTable](functions/makeArrowTable.md) - [packBits](functions/packBits.md) - [permutationBuilder](functions/permutationBuilder.md) diff --git a/docs/src/js/type-aliases/BlobOptions.md b/docs/src/js/type-aliases/BlobOptions.md new file mode 100644 index 000000000..41dbe92c2 --- /dev/null +++ b/docs/src/js/type-aliases/BlobOptions.md @@ -0,0 +1,48 @@ +[**@lancedb/lancedb**](../README.md) • **Docs** + +*** + +[@lancedb/lancedb](../globals.md) / BlobOptions + +# Type Alias: BlobOptions + +```ts +type BlobOptions: object; +``` + +## Type declaration + +### dedicatedSizeThreshold? + +```ts +optional dedicatedSizeThreshold: number; +``` + +Max payload bytes stored in a packed sidecar before a dedicated file. Must +be a positive safe integer. + +### inlineSizeThreshold? + +```ts +optional inlineSizeThreshold: number; +``` + +Max payload bytes kept inline in the data file. Zero is allowed. Must be a +safe integer. + +### nullable? + +```ts +optional nullable: boolean; +``` + +Defaults to true. + +### packFileSizeThreshold? + +```ts +optional packFileSizeThreshold: number; +``` + +Max bytes in one packed sidecar before starting another. Must be a positive +safe integer. diff --git a/nodejs/__test__/blob.test.ts b/nodejs/__test__/blob.test.ts new file mode 100644 index 000000000..e8e9357e3 --- /dev/null +++ b/nodejs/__test__/blob.test.ts @@ -0,0 +1,185 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The LanceDB Authors + +import { Field, Int64, List, Schema, Struct, Utf8 } from "apache-arrow"; +import { makeArrowTable } from "../lancedb/arrow"; +import { BlobFile, blob, coerceBlobValue, isBlobField } from "../lancedb/blob"; + +describe("blob()", () => { + it("marks the field as lance.blob.v2", () => { + const field = blob("image", { nullable: false }); + expect(field.nullable).toBe(false); + expect(isBlobField(field)).toBe(true); + expect(field.metadata.get("ARROW:extension:name")).toBe("lance.blob.v2"); + }); + + it("writes encoding thresholds as field metadata", () => { + const field = blob("video", { + inlineSizeThreshold: 1024, + dedicatedSizeThreshold: 2 * 1024 * 1024, + packFileSizeThreshold: 64 * 1024 * 1024, + }); + expect( + field.metadata.get("lance-encoding:blob-inline-size-threshold"), + ).toBe("1024"); + expect( + field.metadata.get("lance-encoding:blob-dedicated-size-threshold"), + ).toBe(String(2 * 1024 * 1024)); + expect( + field.metadata.get("lance-encoding:blob-pack-file-size-threshold"), + ).toBe(String(64 * 1024 * 1024)); + }); + + it("rejects invalid thresholds", () => { + expect(() => blob("image", { inlineSizeThreshold: -1 })).toThrow( + /inlineSizeThreshold must be non-negative/, + ); + expect(() => blob("image", { dedicatedSizeThreshold: 0 })).toThrow( + /dedicatedSizeThreshold must be positive/, + ); + expect(() => blob("image", { packFileSizeThreshold: 1.5 })).toThrow( + /packFileSizeThreshold must be a safe integer/, + ); + expect(() => + blob("image", { dedicatedSizeThreshold: Number.MAX_SAFE_INTEGER + 1 }), + ).toThrow(/dedicatedSizeThreshold must be a safe integer/); + }); +}); + +describe("coerceBlobValue", () => { + it.each([ + ["Buffer", Buffer.from("x"), { data: Buffer.from("x"), uri: null }], + [ + "Uint8Array", + new Uint8Array([120]), + { data: new Uint8Array([120]), uri: null }, + ], + ["URI string", "s3://bucket/key", { data: null, uri: "s3://bucket/key" }], + [ + "data struct", + { data: Buffer.from("y") }, + { data: Buffer.from("y"), uri: null }, + ], + [ + "uri struct", + { uri: "s3://bucket/key" }, + { data: null, uri: "s3://bucket/key" }, + ], + ["null", null, null], + ])("accepts %s", (_name, input, expected) => { + expect(coerceBlobValue(input)).toEqual(expected); + }); + + it.each([ + ["empty URI", "", /uri cannot be empty/], + ["object without data or uri", { position: 0 }, /data' or 'uri/], + [ + "Int16Array", + new Int16Array([1]), + /Blob data must be Buffer or Uint8Array/, + ], + [ + "both data and uri", + { data: Buffer.from("y"), uri: "s3://bucket/key" }, + /exactly one of 'data' or 'uri'/, + ], + [ + "neither data nor uri", + { data: null, uri: null }, + /exactly one of 'data' or 'uri'/, + ], + ])("rejects %s", (_name, input, message) => { + expect(() => coerceBlobValue(input)).toThrow(message); + }); +}); + +describe("BlobFile", () => { + it("rejects constructing BlobFile without a native handle", () => { + expect(() => new (BlobFile as unknown as { new (): BlobFile })()).toThrow( + /fetchBlobFiles/, + ); + }); +}); + +describe("makeArrowTable blob columns", () => { + it("coerces Buffer input onto a blob field", () => { + const schema = new Schema([ + new Field("id", new Int64(), true), + blob("image"), + ]); + const table = makeArrowTable([{ id: 1n, image: Buffer.from("hello") }], { + schema, + }); + expect(isBlobField(table.schema.fields[1])).toBe(true); + const image = table.getChild("image")!; + expect(image.nullCount).toBe(0); + expect(image.getChild("uri")!.get(0)).toBeNull(); + expect(image.getChild("data")!.nullCount).toBe(0); + expect(Buffer.from(image.getChild("data")!.get(0)!).toString()).toBe( + "hello", + ); + }); + + it("coerces Buffer elements inside a list and keeps null slots", () => { + const schema = new Schema([ + new Field("id", new Int64(), true), + new Field("images", new List(blob("image")), true), + ]); + const table = makeArrowTable( + [ + { id: 1n, images: [Buffer.from("a"), Buffer.from("bb")] }, + { id: 2n, images: null }, + { id: 3n, images: [Buffer.from("c"), null] }, + { id: 4n, images: [] }, + ], + { schema }, + ); + const images = table.getChild("images")!; + expect(images.nullCount).toBe(1); + const rows = images.toArray(); + expect(rows[1]).toBeNull(); + expect(Array.from(rows[3] as Iterable)).toHaveLength(0); + const first = Array.from(rows[0] as Iterable<{ data: Uint8Array | null }>); + expect(Buffer.from(first[0].data!).toString()).toBe("a"); + expect(Buffer.from(first[1].data!).toString()).toBe("bb"); + const third = Array.from( + rows[2] as Iterable<{ data: Uint8Array | null } | null>, + ); + expect(Buffer.from(third[0]!.data!).toString()).toBe("c"); + expect(third[1]).toBeNull(); + }); + + it("coerces Buffer fields inside list structs", () => { + const schema = new Schema([ + new Field("id", new Int64(), true), + new Field( + "items", + new List( + new Field( + "item", + new Struct([new Field("name", new Utf8(), true), blob("image")]), + true, + ), + ), + true, + ), + ]); + const table = makeArrowTable( + [ + { + id: 1n, + items: [{ name: "one", image: Buffer.from("alpha") }], + }, + ], + { schema }, + ); + const items = Array.from( + table.getChild("items")!.toArray()[0] as Iterable<{ + name: string; + image: { data: Uint8Array | null }; + }>, + ); + expect(items[0].name).toBe("one"); + expect(Buffer.from(items[0].image.data!).toString()).toBe("alpha"); + }); +}); diff --git a/nodejs/__test__/table.test.ts b/nodejs/__test__/table.test.ts index e8b97bf77..cae01d9d5 100644 --- a/nodejs/__test__/table.test.ts +++ b/nodejs/__test__/table.test.ts @@ -18,6 +18,7 @@ import { Query, Table, VectorQuery, + blob, connect, tokenize, } from "../lancedb"; @@ -2401,6 +2402,276 @@ describe("when dealing with versioning", () => { }); }); +describe("when dealing with blob columns", () => { + let tmpDir: tmp.DirResult; + beforeEach(() => { + tmpDir = tmp.dirSync({ unsafeCleanup: true }); + }); + afterEach(() => { + tmpDir.removeCallback(); + }); + + it("discovers blob columns", async () => { + const { table } = await openBlobTable(); + expect(await table.blobColumns()).toEqual(["image"]); + }); + + it("preserves order, duplicates, and nulls", async () => { + const { table, rowIds } = await openBlobTable(); + const [alphaId, betaId, nullId] = rowIds; + const bytes = await table.fetchBlobs("image", [ + betaId, + alphaId, + betaId, + nullId, + ]); + expect(bytes.map((b) => (b == null ? null : b.toString()))).toEqual([ + "beta", + "alpha", + "beta", + null, + ]); + const files = await table.fetchBlobFiles("image", [ + betaId, + nullId, + alphaId, + ]); + expect(files.map((f) => f == null)).toEqual([false, true, false]); + }); + + it("reads full blob contents", async () => { + const { table, rowIds, alpha, beta } = await openBlobTable(); + const bytes = await table.fetchBlobs("image", rowIds); + expect(bytes[0]!.equals(alpha)).toBe(true); + expect(bytes[1]!.equals(beta)).toBe(true); + const files = await table.fetchBlobFiles("image", rowIds); + expect(files[0]!.size()).toBe(BigInt(alpha.length)); + expect(Buffer.from(await files[0]!.read()).toString()).toBe("alpha"); + expect(Buffer.from(await files[1]!.read()).toString()).toBe("beta"); + }); + + it("reads a half-open range", async () => { + const { table, rowIds } = await openBlobTable(); + const files = await table.fetchBlobFiles("image", rowIds); + expect(Buffer.from(await files[0]!.readRange(0n, 2n)).toString()).toBe( + "al", + ); + }); + + it("readRange does not move the cursor", async () => { + const { table, rowIds, alpha } = await openBlobTable(); + const [handle] = await table.fetchBlobFiles("image", rowIds); + expect((await handle!.readRange(1n, 3n)).toString()).toBe("lp"); + expect(await handle!.read()).toEqual(alpha); + expect(await handle!.read()).toEqual(Buffer.alloc(0)); + }); + + it("fails when readRange end is past the blob size", async () => { + const { table, rowIds, alpha } = await openBlobTable(); + const files = await table.fetchBlobFiles("image", rowIds); + await expect( + files[0]!.readRange(0n, BigInt(alpha.length + 1)), + ).rejects.toThrow(/exceeds blob size/); + }); + + it("rejects fetchBlobs on a non-blob column", async () => { + const { table, rowIds } = await openBlobTable(); + await expect(table.fetchBlobs("id", rowIds)).rejects.toThrow(/blob/i); + }); + + it("discovers and fetches nested blob columns", async () => { + const db = await connect(tmpDir.name); + const schema = new Schema([ + new Field("id", new Int64(), true), + new Field("info", new Struct([blob("image")]), true), + ]); + const payload = Buffer.from("nested"); + const table = await db.createTable( + "nested_blobs", + [{ id: 1n, info: { image: payload } }], + { schema }, + ); + expect(await table.blobColumns()).toEqual(["info.image"]); + const rows = await table.query().withRowId().toArray(); + const bytes = await table.fetchBlobs("info.image", [ + rows[0]._rowid as bigint, + ]); + expect(bytes[0]!.equals(payload)).toBe(true); + }); + + it("creates and adds list blob columns", async () => { + const db = await connect(tmpDir.name); + const schema = new Schema([ + new Field("id", new Int64(), true), + new Field("images", new List(blob("image")), true), + ]); + const alpha = Buffer.from("alpha"); + const beta = Buffer.from("beta"); + const gamma = Buffer.from("gamma"); + const table = await db.createTable( + "list_blobs", + [{ id: 1n, images: [alpha, beta] }], + { schema }, + ); + await table.add([ + { id: 2n, images: null }, + { id: 3n, images: [gamma, null] }, + { id: 4n, images: [] }, + ]); + expect(await table.blobColumns()).toEqual(["images.image"]); + const rows = await table.query().toArray(); + const byId = new Map(rows.map((row) => [Number(row.id), row])); + expect(descriptorSizes(byId.get(1)!.images)).toEqual([ + alpha.length, + beta.length, + ]); + expect(byId.get(2)!.images).toBeNull(); + expect(descriptorSizes(byId.get(3)!.images)).toEqual([gamma.length, null]); + expect(Array.from(byId.get(4)!.images as Iterable)).toHaveLength( + 0, + ); + }); + + it("creates and adds list struct blob columns", async () => { + const db = await connect(tmpDir.name); + const schema = new Schema([ + new Field("id", new Int64(), true), + new Field( + "items", + new List( + new Field( + "item", + new Struct([new Field("name", new Utf8(), true), blob("image")]), + true, + ), + ), + true, + ), + ]); + const alpha = Buffer.from("nested-alpha"); + const beta = Buffer.from("nested-beta"); + const table = await db.createTable( + "list_struct_blobs", + [{ id: 1n, items: [{ name: "one", image: alpha }] }], + { schema }, + ); + await table.add([ + { + id: 2n, + items: [ + { name: "two", image: beta }, + { name: "three", image: null }, + ], + }, + ]); + const rows = await table.query().toArray(); + const byId = new Map(rows.map((row) => [Number(row.id), row])); + expect( + descriptorSizes( + Array.from(byId.get(1)!.items as Iterable<{ image: unknown }>).map( + (item) => item.image, + ), + ), + ).toEqual([alpha.length]); + expect( + descriptorSizes( + Array.from(byId.get(2)!.items as Iterable<{ image: unknown }>).map( + (item) => item.image, + ), + ), + ).toEqual([beta.length, null]); + }); + + it("rejects blob fields inside a fixed-size list", async () => { + const db = await connect(tmpDir.name); + const schema = new Schema([ + new Field("id", new Int64(), true), + new Field("frames", new FixedSizeList(2, blob("frame")), true), + ]); + await expect( + db.createTable( + "fsl_blobs", + [{ id: 1n, frames: [Buffer.from("a"), Buffer.from("b")] }], + { schema }, + ), + ).rejects.toThrow( + "Blob fields inside FixedSizeList are not supported. Use List instead.", + ); + }); + + it("rejects blob fields inside a nested fixed-size list", async () => { + const db = await connect(tmpDir.name); + const schema = new Schema([ + new Field("id", new Int64(), true), + new Field( + "clip", + new Struct([ + new Field("frames", new FixedSizeList(2, blob("frame")), true), + ]), + true, + ), + ]); + await expect( + db.createTable( + "nested_fsl_blobs", + [ + { + id: 1n, + clip: { frames: [Buffer.from("a"), Buffer.from("b")] }, + }, + ], + { schema }, + ), + ).rejects.toThrow( + "Blob fields inside FixedSizeList are not supported. Use List instead.", + ); + }); + + it("rejects an Arrow table with blob fields inside a fixed-size list", async () => { + const db = await connect(tmpDir.name); + const schema = new Schema([ + new Field("id", new Int64(), true), + new Field("frames", new FixedSizeList(2, blob("frame")), true), + ]); + await expect( + db.createTable("fsl_blobs_ipc", new ArrowTable(schema)), + ).rejects.toThrow( + "Blob fields inside FixedSizeList are not supported. Use List instead.", + ); + }); + + function descriptorSizes(values: unknown): (number | null)[] { + return Array.from( + values as Iterable<{ size?: bigint | number } | null>, + ).map((value) => (value == null ? null : Number(value.size))); + } + + async function openBlobTable() { + const db = await connect(tmpDir.name); + const schema = new Schema([ + new Field("id", new Int64(), true), + blob("image"), + ]); + const alpha = Buffer.from("alpha"); + const beta = Buffer.from("beta"); + const table = await db.createTable( + "blobs", + [ + { id: 1n, image: alpha }, + { id: 2n, image: beta }, + { id: 3n, image: null }, + ], + { schema }, + ); + const rows = await table.query().withRowId().toArray(); + const rowIdById = new Map( + rows.map((r) => [Number(r.id), r._rowid as bigint]), + ); + const rowIds = [1, 2, 3].map((id) => rowIdById.get(id)!); + return { table, rowIds, alpha, beta }; + } +}); + describe("when dealing with tags", () => { let tmpDir: tmp.DirResult; beforeEach(() => { diff --git a/nodejs/lancedb/arrow.ts b/nodejs/lancedb/arrow.ts index 119887704..81b140da7 100644 --- a/nodejs/lancedb/arrow.ts +++ b/nodejs/lancedb/arrow.ts @@ -40,6 +40,7 @@ import { } from "apache-arrow"; import { Buffers } from "apache-arrow/data"; import { typedArrayToArrowType } from "./arrow_type"; +import { coerceBlobValue, isBlobField } from "./blob"; import { type EmbeddingFunction } from "./embedding/embedding_function"; import { EmbeddingFunctionConfig, @@ -430,12 +431,14 @@ export function makeArrowTable( throw new Error("A schema must be provided if data is empty"); } else { schema = new Schema(schema.fields, schemaMetadata); + validateBlobSchema(schema); return new ArrowTable(schema); } } let inferredSchema = inferSchema(data, schema, opt); inferredSchema = new Schema(inferredSchema.fields, schemaMetadata); + validateBlobSchema(inferredSchema); const finalColumns: Record = {}; for (const field of inferredSchema.fields) { @@ -445,6 +448,35 @@ export function makeArrowTable( return new ArrowTable(inferredSchema, finalColumns); } +function validateBlobSchema(schema: Schema): void { + for (const field of schema.fields) { + validateBlobField(field); + } +} + +function validateBlobField(field: Field): void { + if ( + isFixedSizeList(field.type) && + containsBlobField(field.type.children[0]) + ) { + throw new Error( + "Blob fields inside FixedSizeList are not supported. Use List instead.", + ); + } + for (const child of field.type.children ?? []) { + validateBlobField(child); + } +} + +function containsBlobField(field: Field): boolean { + if (isBlobField(field)) { + return true; + } + return (field.type.children ?? []).some((child: Field) => + containsBlobField(child), + ); +} + function isObject(value: unknown): value is Record { return ( typeof value === "object" && @@ -480,6 +512,32 @@ function transposeData( path: string[] = [], ): Vector { const valuesPath = [...path, field.name]; + if (isBlobField(field) && field.type instanceof Struct) { + const blobRows = data.map((datum) => + coerceBlobValue(valueAtPath(datum, valuesPath)), + ); + const childVectors = field.type.children.map((child) => { + const values = blobRows.map((row) => + row == null ? null : (row[child.name as "data" | "uri"] ?? null), + ); + return makeVector(values, child.type, undefined, child.nullable); + }); + const nullCount = blobRows.filter((row) => row === null).length; + const structData = makeData({ + type: field.type, + length: blobRows.length, + nullCount, + nullBitmap: + nullCount > 0 + ? arrowUtil.packBools(blobRows.map((row) => row !== null)) + : undefined, + children: childVectors.map((v) => v.data[0]), + }); + return arrowMakeVector(structData); + } + if (isList(field.type) && containsBlobField(field.type.children[0])) { + return transposeListData(data, field, valuesPath); + } const values = data.map((datum) => valueAtPath(datum, valuesPath)); if (field.type instanceof Struct) { const childFields = field.type.children; @@ -495,7 +553,7 @@ function transposeData( nullCount > 0 ? arrowUtil.packBools(values.map((value) => value !== null)) : undefined, - children: childVectors as unknown as ArrowData[], + children: childVectors.map((v) => v.data[0]), }); return arrowMakeVector(structData); } else { @@ -503,6 +561,48 @@ function transposeData( } } +function transposeListData( + data: Record[], + field: Field, + valuesPath: string[], +): Vector { + const listType = field.type as List; + const childField = listType.children[0]; + const lists = data.map((datum) => valueAtPath(datum, valuesPath)); + const flattened: Record[] = []; + const validity: boolean[] = []; + const offsets: number[] = [0]; + + for (const list of lists) { + if (list == null) { + validity.push(false); + offsets.push(flattened.length); + continue; + } + if (!Array.isArray(list)) { + throw new Error(`expected an array for list field '${field.name}'`); + } + validity.push(true); + for (const element of list) { + flattened.push({ [childField.name]: element }); + } + offsets.push(flattened.length); + } + + const childVector = transposeData(flattened, childField, []); + const nullCount = validity.filter((valid) => !valid).length; + return arrowMakeVector( + makeData({ + type: listType, + length: lists.length, + nullCount, + nullBitmap: nullCount > 0 ? arrowUtil.packBools(validity) : undefined, + valueOffsets: Int32Array.from(offsets), + child: childVector.data[0], + }), + ); +} + /** * Create an empty Arrow table with the provided schema */ @@ -952,6 +1052,7 @@ export async function fromTableToBuffer( schema = sanitizeSchema(schema); } const tableWithEmbeddings = await applyEmbeddings(table, embeddings, schema); + validateBlobSchema(tableWithEmbeddings.schema); const writer = RecordBatchFileWriter.writeAll(tableWithEmbeddings); return Buffer.from(await writer.toUint8Array()); } diff --git a/nodejs/lancedb/blob.ts b/nodejs/lancedb/blob.ts new file mode 100644 index 000000000..244faac09 --- /dev/null +++ b/nodejs/lancedb/blob.ts @@ -0,0 +1,236 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The LanceDB Authors + +import { Field, LargeBinary, Struct, Utf8 } from "apache-arrow"; +import { BlobFile as NativeBlobFile } from "./native"; + +const BLOB_V2_EXTENSION_NAME = "lance.blob.v2"; + +const INLINE_SIZE_THRESHOLD_KEY = "lance-encoding:blob-inline-size-threshold"; +const DEDICATED_SIZE_THRESHOLD_KEY = + "lance-encoding:blob-dedicated-size-threshold"; +const PACK_FILE_SIZE_THRESHOLD_KEY = + "lance-encoding:blob-pack-file-size-threshold"; + +export type BlobInput = { + data: Buffer | Uint8Array | null; + uri: string | null; +}; + +export type BlobOptions = { + /** Defaults to true. */ + nullable?: boolean; + /** + * Max payload bytes kept inline in the data file. Zero is allowed. Must be a + * safe integer. + */ + inlineSizeThreshold?: number; + /** + * Max payload bytes stored in a packed sidecar before a dedicated file. Must + * be a positive safe integer. + */ + dedicatedSizeThreshold?: number; + /** + * Max bytes in one packed sidecar before starting another. Must be a positive + * safe integer. + */ + packFileSizeThreshold?: number; +}; + +/** + * Declares a `lance.blob.v2` column. + * + * Query results are descriptors, not payload bytes. Use {@link Table.fetchBlobs} + * or {@link Table.fetchBlobFiles} to read bytes. + * + * @example + * ```ts + * import { readFile } from "node:fs/promises"; + * import { Field, Int64, Schema } from "apache-arrow"; + * import { blob, connect } from "@lancedb/lancedb"; + * + * const db = await connect("./data"); + * const video = await readFile("clip.mp4"); + * const table = await db.createTable( + * "videos", + * [{ id: 1n, video }], + * { + * schema: new Schema([ + * new Field("id", new Int64()), + * blob("video"), + * ]), + * }, + * ); + * + * const rows = await table.query().select(["id"]).withRowId().toArray(); + * const rowIds = rows.map((row) => row._rowid as bigint); + * const bytes = await table.fetchBlobs("video", rowIds); + * + * const [handle] = await table.fetchBlobFiles("video", rowIds); + * const size = handle!.size(); + * const header = await handle!.readRange(0n, size < 65536n ? size : 65536n); + * ``` + */ +export function blob(name: string, options: BlobOptions = {}): Field { + const metadata = new Map([ + ["ARROW:extension:name", BLOB_V2_EXTENSION_NAME], + ]); + setThreshold( + metadata, + INLINE_SIZE_THRESHOLD_KEY, + "inlineSizeThreshold", + options.inlineSizeThreshold, + 0, + ); + setThreshold( + metadata, + DEDICATED_SIZE_THRESHOLD_KEY, + "dedicatedSizeThreshold", + options.dedicatedSizeThreshold, + 1, + ); + setThreshold( + metadata, + PACK_FILE_SIZE_THRESHOLD_KEY, + "packFileSizeThreshold", + options.packFileSizeThreshold, + 1, + ); + return new Field( + name, + new Struct([ + new Field("data", new LargeBinary(), true), + new Field("uri", new Utf8(), true), + ]), + options.nullable ?? true, + metadata, + ); +} + +/** + * Checks for the `lance.blob.v2` extension marker. Does not validate the + * field's storage type. + */ +export function isBlobField(field: Field): boolean { + return field.metadata?.get("ARROW:extension:name") === BLOB_V2_EXTENSION_NAME; +} + +/** + * A lazy handle to blob bytes. Create one with {@link Table.fetchBlobFiles}. + * + * @hideconstructor + */ +export class BlobFile { + private readonly inner: NativeBlobFile; + + private constructor(inner: NativeBlobFile) { + if (!(inner instanceof NativeBlobFile)) { + throw new Error("BlobFile handles come from Table.fetchBlobFiles"); + } + this.inner = inner; + } + + /** @ignore */ + static fromNative(inner: NativeBlobFile): BlobFile { + return new BlobFile(inner); + } + + /** Returns the blob size in bytes. */ + size(): bigint { + return this.inner.size(); + } + + /** + * Reads from the cursor to the end and advances the cursor. + * + * A second call returns an empty buffer. {@link BlobFile.readRange} does + * not move the cursor. + */ + read(): Promise { + return this.inner.read(); + } + + /** + * Reads the half-open byte range `[start, end)`. + * + * Fails when `end` is past the blob size. Does not move the cursor. + */ + readRange(start: bigint, end: bigint): Promise { + return this.inner.readRange(start, end); + } +} + +export function coerceBlobValue(value: unknown): BlobInput | null { + if (value == null) { + return null; + } + if (isBlobBytes(value)) { + return { data: value, uri: null }; + } + if (ArrayBuffer.isView(value)) { + throw new Error("Blob data must be Buffer or Uint8Array"); + } + if (typeof value === "string") { + if (value === "") { + throw new Error("Blob uri cannot be empty"); + } + return { data: null, uri: value }; + } + if (typeof value === "object") { + const record = value as Record; + if (!("data" in record) && !("uri" in record)) { + throw new Error( + "Blob struct values must include a 'data' or 'uri' field", + ); + } + const uri = record.uri; + if (uri === "") { + throw new Error("Blob uri cannot be empty"); + } + if (uri != null && typeof uri !== "string") { + throw new Error(`Blob uri must be a string or null, got ${typeof uri}`); + } + const data = record.data; + if (data != null && !isBlobBytes(data)) { + throw new Error("Blob data must be Buffer, Uint8Array, or null"); + } + const bytes = (data as Buffer | Uint8Array | null | undefined) ?? null; + const uriValue = uri ?? null; + if ((bytes == null) === (uriValue == null)) { + throw new Error( + "Blob struct values must set exactly one of 'data' or 'uri'", + ); + } + return { data: bytes, uri: uriValue }; + } + throw new Error( + "Blob column values must be Buffer, Uint8Array, a URI string, null, or { data?, uri? }", + ); +} + +function isBlobBytes(value: unknown): value is Buffer | Uint8Array { + return Buffer.isBuffer(value) || value instanceof Uint8Array; +} + +function setThreshold( + metadata: Map, + key: string, + optionName: string, + value: number | undefined, + minimum: number, +): void { + if (value === undefined) { + return; + } + if (!Number.isSafeInteger(value)) { + throw new Error(`${optionName} must be a safe integer`); + } + if (value < minimum) { + throw new Error( + minimum <= 0 + ? `${optionName} must be non-negative` + : `${optionName} must be positive`, + ); + } + metadata.set(key, String(value)); +} diff --git a/nodejs/lancedb/index.ts b/nodejs/lancedb/index.ts index 4f8ff77e5..d94007a11 100644 --- a/nodejs/lancedb/index.ts +++ b/nodejs/lancedb/index.ts @@ -77,6 +77,9 @@ export { VectorColumnOptions, } from "./arrow"; +export { blob, isBlobField, BlobFile } from "./blob"; +export type { BlobOptions } from "./blob"; + export { Connection, CreateTableOptions, diff --git a/nodejs/lancedb/table.ts b/nodejs/lancedb/table.ts index d1fb8acd8..eac9f490d 100644 --- a/nodejs/lancedb/table.ts +++ b/nodejs/lancedb/table.ts @@ -17,6 +17,7 @@ import { tableFromIPC, } from "./arrow"; +import { BlobFile } from "./blob"; import { EmbeddingFunctionConfig, getRegistry } from "./embedding/registry"; import { IndexOptions } from "./indices"; import { Job } from "./job"; @@ -510,6 +511,35 @@ export abstract class Table { */ abstract takeRowIds(rowIds: readonly (bigint | number)[]): TakeQuery; + /** + * Blob v2 columns, including nested dotted paths. + */ + abstract blobColumns(): Promise; + + /** + * Bytes for `column` at row IDs from {@link Query.withRowId}. + * + * Reads the table's current checkout. IDs from another version can fail after + * compaction unless stable row ids are enabled. Results keep input order and + * duplicates. Null blobs are `null`. Empty blobs are empty buffers. + */ + abstract fetchBlobs( + column: string, + rowIds: readonly (bigint | number)[], + ): Promise<(Buffer | null)[]>; + + /** + * Opens lazy blob handles for `column` at the given row IDs using the + * table's current checkout. + * + * Preserves input order, duplicates, and nulls. Use this for large payloads. + * See {@link Table.fetchBlobs} for row-ID validity across versions. + */ + abstract fetchBlobFiles( + column: string, + rowIds: readonly (bigint | number)[], + ): Promise<(BlobFile | null)[]>; + /** * Create a search query to find the nearest neighbors * of the given query @@ -1160,23 +1190,34 @@ export class LocalTable extends Table { } takeRowIds(rowIds: readonly (bigint | number)[]): TakeQuery { - const ids = rowIds.map((id) => { - if (typeof id === "bigint") { - return id; - } - if (!Number.isInteger(id)) { - throw new Error("Row id must be an integer (or bigint)"); - } - if (id < 0) { - throw new Error("Row id cannot be negative"); - } - if (!Number.isSafeInteger(id)) { - throw new Error("Row id is too large for number; use bigint instead"); - } - return BigInt(id); - }); + return new TakeQuery(this.inner.takeRowIds(rowIdsToBigInts(rowIds))); + } - return new TakeQuery(this.inner.takeRowIds(ids)); + blobColumns(): Promise { + return this.inner.blobColumns(); + } + + async fetchBlobs( + column: string, + rowIds: readonly (bigint | number)[], + ): Promise<(Buffer | null)[]> { + const values = await this.inner.fetchBlobs(column, rowIdsToBigInts(rowIds)); + // N-API Option maps missing values to undefined. Collapse those to null. + return values.map((value) => value ?? null); + } + + async fetchBlobFiles( + column: string, + rowIds: readonly (bigint | number)[], + ): Promise<(BlobFile | null)[]> { + const files = await this.inner.fetchBlobFiles( + column, + rowIdsToBigInts(rowIds), + ); + // N-API Option maps missing values to undefined. Collapse those to null. + return files.map((file) => + file == null ? null : BlobFile.fromNative(file), + ); } query(): Query { @@ -1733,3 +1774,21 @@ export class Branches { )) as unknown as CherryPickResult; } } + +function rowIdsToBigInts(rowIds: readonly (bigint | number)[]): bigint[] { + return rowIds.map((id) => { + if (typeof id === "bigint") { + return id; + } + if (!Number.isInteger(id)) { + throw new Error("Row id must be an integer (or bigint)"); + } + if (id < 0) { + throw new Error("Row id cannot be negative"); + } + if (!Number.isSafeInteger(id)) { + throw new Error("Row id is too large for number; use bigint instead"); + } + return BigInt(id); + }); +} diff --git a/nodejs/src/blob.rs b/nodejs/src/blob.rs new file mode 100644 index 000000000..0e19a9a8a --- /dev/null +++ b/nodejs/src/blob.rs @@ -0,0 +1,95 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The LanceDB Authors + +use std::ops::Range; +use std::sync::Arc; + +use arrow_array::{Array, LargeBinaryArray}; +use lancedb::blob::BlobFile as LanceBlobFile; +use napi::bindgen_prelude::*; +use napi_derive::napi; + +use crate::error::convert_error; + +#[napi] +pub struct BlobFile { + inner: Arc, +} + +impl BlobFile { + pub(crate) fn new(inner: LanceBlobFile) -> Self { + Self { + inner: Arc::new(inner), + } + } +} + +#[napi] +impl BlobFile { + #[napi] + pub fn size(&self) -> BigInt { + BigInt::from(self.inner.size()) + } + + #[napi] + pub async fn read(&self) -> napi::Result { + let bytes = self.inner.read().await.map_err(|err| convert_error(&err))?; + Ok(Buffer::from(bytes.as_ref())) + } + + #[napi] + pub async fn read_range(&self, start: BigInt, end: BigInt) -> napi::Result { + let range = bigint_range(start, end)?; + let bytes = self + .inner + .read_range(range) + .await + .map_err(|err| convert_error(&err))?; + Ok(Buffer::from(bytes.as_ref())) + } +} + +fn bigint_range(start: BigInt, end: BigInt) -> napi::Result> { + let start = parse_u64(start, "start")?; + let end = parse_u64(end, "end")?; + if start > end { + return Err(napi::Error::from_reason(format!( + "invalid blob range: start ({start}) > end ({end})" + ))); + } + Ok(start..end) +} + +pub(crate) fn parse_u64(value: BigInt, name: &str) -> napi::Result { + let (negative, value, lossless) = value.get_u64(); + if negative { + return Err(napi::Error::from_reason(format!( + "{name} cannot be negative" + ))); + } + if !lossless { + return Err(napi::Error::from_reason(format!( + "{name} is too large to fit in u64" + ))); + } + Ok(value) +} + +pub(crate) fn parse_row_ids(row_ids: Vec) -> napi::Result> { + row_ids + .into_iter() + .map(|id| parse_u64(id, "row id")) + .collect() +} + +pub(crate) fn copy_blob_buffers(array: LargeBinaryArray) -> Vec> { + (0..array.len()) + .map(|i| { + if array.is_null(i) { + None + } else { + Some(Buffer::from(array.value(i).to_vec())) + } + }) + .collect() +} diff --git a/nodejs/src/lib.rs b/nodejs/src/lib.rs index 1110f6203..288a2b925 100644 --- a/nodejs/src/lib.rs +++ b/nodejs/src/lib.rs @@ -10,6 +10,7 @@ use std::collections::HashMap; use env_logger::Env; use napi_derive::*; +mod blob; mod connection; mod error; mod header; diff --git a/nodejs/src/table.rs b/nodejs/src/table.rs index db74d38fa..7ae4402ab 100644 --- a/nodejs/src/table.rs +++ b/nodejs/src/table.rs @@ -15,6 +15,7 @@ use napi::bindgen_prelude::*; use napi::threadsafe_function::{ThreadsafeFunction, ThreadsafeFunctionCallMode}; use napi_derive::napi; +use crate::blob::{BlobFile, copy_blob_buffers, parse_row_ids}; use crate::error::NapiErrorExt; use crate::index::Index; use crate::merge::NativeMergeInsertBuilder; @@ -329,6 +330,44 @@ impl Table { )) } + #[napi(catch_unwind)] + pub async fn blob_columns(&self) -> napi::Result> { + self.inner_ref()?.blob_columns().await.default_error() + } + + #[napi(catch_unwind)] + pub async fn fetch_blobs( + &self, + column: String, + row_ids: Vec, + ) -> napi::Result>> { + let row_ids = parse_row_ids(row_ids)?; + let array = self + .inner_ref()? + .fetch_blobs(column.as_str(), &row_ids) + .await + .default_error()?; + Ok(copy_blob_buffers(array)) + } + + #[napi(catch_unwind)] + pub async fn fetch_blob_files( + &self, + column: String, + row_ids: Vec, + ) -> napi::Result>> { + let row_ids = parse_row_ids(row_ids)?; + let files = self + .inner_ref()? + .fetch_blob_files(column.as_str(), &row_ids) + .await + .default_error()?; + Ok(files + .into_iter() + .map(|file| file.map(BlobFile::new)) + .collect()) + } + #[napi(catch_unwind)] pub fn vector_search(&self, vector: Float32Array) -> napi::Result { self.query()?.nearest_to(vector)