Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 7 additions & 4 deletions docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,9 @@ You can use the `options` parameter to specify various options. Options supporte
- `authToken`: authentication token for the provider URL (optional).
- `timeout`: number of milliseconds to wait on locked database before returning `SQLITE_BUSY` error
- `defaultQueryTimeout`: default maximum number of milliseconds a query is allowed to run before being interrupted with `SQLITE_INTERRUPT` error
- `defaultBatchSize`: default number of rows that the promise API (`libsql/promise`) reads per native iterator call. It must be an integer from 1 through 10,000 and defaults to 1. The synchronous API ignores it.

With the promise API, use `batchSize` in `queryOptions` to override `defaultBatchSize` for one `all()` or `iterate()` call. When the query has no bind parameters, pass `undefined` before the options: `await statement.all(undefined, { batchSize: 250 })`. The synchronous API does not support `batchSize`.

The function returns a `Database` object.

Expand Down Expand Up @@ -68,7 +71,7 @@ Convenience wrapper that prepares `sql` and executes `Statement.all`. Returns al
| -------------- | ------------------- | -------------------------------------------------------------------- |
| sql | <code>string</code> | The SQL statement string. |
| bindParameters | <code>any</code> | Optional positional or named bind parameters. |
| queryOptions | <code>object</code> | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). |
| queryOptions | <code>object</code> | Optional per-query overrides (for example, `{ queryTimeout: 100, batchSize: 250 }`). |

**Note:** This is an extension in libSQL and not available in `better-sqlite3`.

Expand All @@ -80,7 +83,7 @@ Convenience wrapper that prepares `sql` and executes `Statement.iterate`. Return
| -------------- | ------------------- | -------------------------------------------------------------------- |
| sql | <code>string</code> | The SQL statement string. |
| bindParameters | <code>any</code> | Optional positional or named bind parameters. |
| queryOptions | <code>object</code> | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). |
| queryOptions | <code>object</code> | Optional per-query overrides (for example, `{ queryTimeout: 100, batchSize: 250 }`). |

**Note:** This is an extension in libSQL and not available in `better-sqlite3`.

Expand Down Expand Up @@ -333,7 +336,7 @@ Executes the SQL statement and returns an array of the resulting rows.
| Param | Type | Description |
| -------------- | ----------------------------- | ------------------------------------------------ |
| bindParameters | <code>array of objects</code> | The bind parameters for executing the statement. |
| queryOptions | <code>object</code> | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). |
| queryOptions | <code>object</code> | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). The promise API also accepts `batchSize`. |

### iterate([...bindParameters][, queryOptions]) ⇒ iterator

Expand All @@ -342,7 +345,7 @@ Executes the SQL statement and returns an iterator to the resulting rows.
| Param | Type | Description |
| -------------- | ----------------------------- | ------------------------------------------------ |
| bindParameters | <code>array of objects</code> | The bind parameters for executing the statement. |
| queryOptions | <code>object</code> | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). |
| queryOptions | <code>object</code> | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). The promise API also accepts `batchSize`. |

### pluck([toggleState]) ⇒ this

Expand Down
5 changes: 5 additions & 0 deletions index.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,10 +14,12 @@ export interface Options {
encryptionKey?: string
remoteEncryptionKey?: string
defaultQueryTimeout?: number
defaultBatchSize?: number
}
/** Per-query execution options. */
export interface QueryOptions {
queryTimeout?: number
batchSize?: number
Comment thread
xoxohorses marked this conversation as resolved.
}
export declare function connect(path: string, opts?: Options | undefined | null): Promise<Database>
/** Result of a database sync operation. */
Expand Down Expand Up @@ -164,6 +166,7 @@ export declare class Database {
}
/** SQLite statement object. */
export declare class Statement {
get defaultBatchSize(): number
/**
* Executes a SQL statement.
*
Expand Down Expand Up @@ -200,6 +203,8 @@ export declare class Statement {
/** A raw iterator over rows. The JavaScript layer wraps this in a iterable. */
export declare class RowsIterator {
next(): Promise<Record>
/** Reads one batch of rows. The batch size must be an integer between 1 and 10,000. */
nextBatch(maxRows: number): Promise<unknown[]>
close(): void
}
export declare class Record {
Expand Down
152 changes: 147 additions & 5 deletions integration-tests/tests/async.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,61 @@ test.serial("Statement.iterate()", async (t) => {
}
});

test.serial("Statement.iterate() preserves concurrent next() order", async (t) => {
const db = t.context.db;
const stmt = await db.prepare("SELECT 1 AS value UNION ALL SELECT 2 UNION ALL SELECT 3");
const iterator = await stmt.iterate(undefined, { batchSize: 2 });
Comment thread
xoxohorses marked this conversation as resolved.

t.deepEqual(await Promise.all([
iterator.next(),
iterator.next(),
iterator.next(),
iterator.next(),
]), [
{ done: false, value: { value: 1 } },
{ done: false, value: { value: 2 } },
{ done: false, value: { value: 3 } },
{ done: true, value: null },
]);
});

test.serial("Statement.iterate() discards buffered rows after return()", async (t) => {
const db = t.context.db;
const stmt = await db.prepare("SELECT * FROM users ORDER BY id");
const iterator = await stmt.iterate(undefined, { batchSize: 2 });

await iterator.next();
iterator.return();

t.deepEqual(await iterator.next(), { done: true, value: null });
});

test.serial("Statement.iterate() return() during an in-flight batch releases the statement", async (t) => {
const path = genDatabaseFilename();
const [conn1] = await connect(path);
await conn1.exec("CREATE TABLE t(x)");
await conn1.exec(`
WITH RECURSIVE numbers(x) AS (SELECT 1 UNION ALL SELECT x + 1 FROM numbers WHERE x < 20000)
INSERT INTO t SELECT x FROM numbers
`);
const stmt = await conn1.prepare("SELECT x FROM t");
const iterator = await stmt.iterate(undefined, { batchSize: 10_000 });

const pending = iterator.next();
// Let the wrapper issue nextBatch() before closing the iterator.
await null;
iterator.return();
await pending;

// An active reader would hold a SHARED lock and make this write fail with SQLITE_BUSY.
const [conn2] = await connect(path);
await t.notThrowsAsync(() => conn2.exec("INSERT INTO t VALUES (0)"));

conn1.close();
conn2.close();
fs.unlinkSync(path);
});

test.serial("Statement.iterate() with invalid bind parameter", async (t) => {
const db = t.context.db;

Expand All @@ -158,6 +213,65 @@ test.serial("Statement.all()", async (t) => {
{ id: 2, name: "Bob", email: "bob@example.com" },
];
t.deepEqual(await stmt.all(), expected);

const namedStmt = await db.prepare("SELECT :batchSize AS value");
t.deepEqual(await namedStmt.all({ batchSize: 2 }), [{ value: 2 }]);
});

test.serial("Statement.all() rejects invalid batch sizes", async (t) => {
const db = t.context.db;
const stmt = await db.prepare("SELECT * FROM users");

const iterator = await stmt.iterate(undefined, { batchSize: 0 });
await t.throwsAsync(() => iterator.next(), {
message: "maxRows must be an integer between 1 and 10000",
});
iterator.return();
t.deepEqual(await iterator.next(), { done: true, value: null });

for (const batchSize of [-1, 1.5, NaN, Infinity, 10_001]) {
await t.throwsAsync(() => stmt.all(undefined, { batchSize }), {
message: "maxRows must be an integer between 1 and 10000",
});
}

t.is((await stmt.all(undefined, { batchSize: 10_000 })).length, 2);
});

test.serial("Invalid batch size in for await releases the statement", async (t) => {
const path = genDatabaseFilename();
const [conn1] = await connect(path);
await conn1.exec("CREATE TABLE t(x)");
await conn1.exec("INSERT INTO t VALUES (1), (2)");
const stmt = await conn1.prepare("SELECT x FROM t");

// A rejected next() does not call return(), so the statement must be released natively.
await t.throwsAsync(async () => {
for await (const _ of await stmt.iterate(undefined, { batchSize: 0 })) {
}
}, {
message: "maxRows must be an integer between 1 and 10000",
});

// An active reader would hold a SHARED lock and make this write fail with SQLITE_BUSY.
const [conn2] = await connect(path);
await t.notThrowsAsync(() => conn2.exec("INSERT INTO t VALUES (0)"));

conn1.close();
conn2.close();
fs.unlinkSync(path);
});

test.serial("defaultBatchSize applies and batchSize overrides it", async (t) => {
const [db] = await connect(":memory:", { defaultBatchSize: 0 });
const stmt = await db.prepare("SELECT 1 AS value");

await t.throwsAsync(() => stmt.all(), {
message: "maxRows must be an integer between 1 and 10000",
});
t.deepEqual(await stmt.all(undefined, { batchSize: 1 }), [{ value: 1 }]);

db.close();
});

test.serial("Statement.all() [raw]", async (t) => {
Expand All @@ -168,7 +282,7 @@ test.serial("Statement.all() [raw]", async (t) => {
[1, "Alice", "alice@example.org"],
[2, "Bob", "bob@example.com"],
];
t.deepEqual(await stmt.raw().all(), expected);
t.deepEqual(await stmt.raw().all(undefined, { batchSize: 250 }), expected);
});

test.serial("Statement.all() [pluck]", async (t) => {
Expand All @@ -179,7 +293,7 @@ test.serial("Statement.all() [pluck]", async (t) => {
1,
2,
];
t.deepEqual(await stmt.pluck().all(), expected);
t.deepEqual(await stmt.pluck().all(undefined, { batchSize: 250 }), expected);
});

test.serial("Statement.all() [default safe integers]", async (t) => {
Expand All @@ -190,7 +304,7 @@ test.serial("Statement.all() [default safe integers]", async (t) => {
[1n, "Alice", "alice@example.org"],
[2n, "Bob", "bob@example.com"],
];
t.deepEqual(await stmt.raw().all(), expected);
t.deepEqual(await stmt.raw().all(undefined, { batchSize: 250 }), expected);
});

test.serial("Statement.all() [statement safe integers]", async (t) => {
Expand All @@ -201,7 +315,7 @@ test.serial("Statement.all() [statement safe integers]", async (t) => {
[1n, "Alice", "alice@example.org"],
[2n, "Bob", "bob@example.com"],
];
t.deepEqual(await stmt.raw().all(), expected);
t.deepEqual(await stmt.raw().all(undefined, { batchSize: 250 }), expected);
});

test.serial("Statement.raw() [failure]", async (t) => {
Expand Down Expand Up @@ -453,6 +567,34 @@ test.serial("Query timeout option interrupts long-running query", async (t) => {
db.close();
});

test.serial("Query timeout resets batched Statement.all() for reuse", async (t) => {
const [db, errorType] = await connect(":memory:");
const stmt = await db.prepare(`
WITH RECURSIVE numbers(value) AS (
SELECT 1
UNION ALL
SELECT value + 1 FROM numbers WHERE value < ?
)
SELECT value FROM numbers
`);

await t.throwsAsync(async () => {
await stmt.all(1_000_000_000, { queryTimeout: 100, batchSize: 100 });
}, {
instanceOf: errorType,
message: "interrupted",
code: "SQLITE_INTERRUPT",
});

t.deepEqual(await stmt.all(3, { batchSize: 2 }), [
{ value: 1 },
{ value: 2 },
{ value: 3 },
]);

db.close();
});

test.serial("Query timeout option interrupts long-running Statement.get()", async (t) => {
const [db, errorType] = await connect(":memory:", { defaultQueryTimeout: 100 });
const stmt = await db.prepare(`
Expand Down Expand Up @@ -500,7 +642,7 @@ test.serial("Stale timeout guard from exhausted iterator does not interrupt late
// interrupt unrelated later queries.
const stmt = await db.prepare("SELECT * FROM t ORDER BY x ASC");
for (let i = 0; i < 150; i++) {
const rows = await stmt.all();
const rows = await stmt.all(undefined, { batchSize: 250 });
t.is(rows.length, 2_000);
}

Expand Down
63 changes: 63 additions & 0 deletions perf/perf-libsql-batched-rows.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
// Run `npm run build`, then from this directory run:
// `npm install && node --expose-gc perf-libsql-batched-rows.js`
import { baseline, bench, group, run } from 'mitata';

// Import the parent checkout so this benchmark uses its locally built native module.
import libsql from '../promise.js';

const { connect } = libsql;

const BATCH_SIZE = 250;
const ROW_COUNTS = [1_000, 10_000, 100_000, 1_000_000];
const MAX_ROW_COUNT = ROW_COUNTS[ROW_COUNTS.length - 1];

const db = await connect(':memory:', {});
await db.exec(`
CREATE TABLE benchmark_rows (value INTEGER PRIMARY KEY);
WITH RECURSIVE numbers(value) AS (
SELECT 1
UNION ALL
SELECT value + 1 FROM numbers WHERE value < ${MAX_ROW_COUNT}
)
INSERT INTO benchmark_rows SELECT value FROM numbers;
`);

const stmt = await db.prepare(`
SELECT value
FROM benchmark_rows
WHERE value <= ?
ORDER BY value
`);

for (const rowCount of ROW_COUNTS) {
group(`${rowCount.toLocaleString('en-US')} rows`, () => {
baseline('all()', async () => {
validateRows(await stmt.all(rowCount), rowCount);
});
bench(`all() with batchSize ${BATCH_SIZE}`, async () => {
validateRows(await stmt.all(rowCount, { batchSize: BATCH_SIZE }), rowCount);
});
});
}

await run({
units: false,
silent: false,
avg: true,
json: false,
colors: process.stdout.isTTY,
min_max: true,
percentiles: true,
});

db.close();

function validateRows(rows, expectedCount) {
if (
rows.length !== expectedCount ||
rows[0]?.value !== 1 ||
rows[rows.length - 1]?.value !== expectedCount
) {
throw new Error(`Expected rows 1 through ${expectedCount}`);
}
}
Loading