From 0267c5c3a0af1824c37ca85f320a6ee459ee54c3 Mon Sep 17 00:00:00 2001 From: Julie Yu Date: Sat, 12 Sep 2026 00:56:53 +0000 Subject: [PATCH 01/14] Batch async row iteration --- index.d.ts | 1 + integration-tests/tests/async.test.js | 35 ++++++ promise.js | 29 +++++ src/lib.rs | 163 +++++++++++++++++++++++--- 4 files changed, 209 insertions(+), 19 deletions(-) diff --git a/index.d.ts b/index.d.ts index 91e6140..a7926c8 100644 --- a/index.d.ts +++ b/index.d.ts @@ -200,6 +200,7 @@ export declare class Statement { /** A raw iterator over rows. The JavaScript layer wraps this in a iterable. */ export declare class RowsIterator { next(): Promise + nextBatch(maxRows: number): Promise<{ records: unknown[]; done: boolean }> close(): void } export declare class Record { diff --git a/integration-tests/tests/async.test.js b/integration-tests/tests/async.test.js index e41b768..7815e5d 100644 --- a/integration-tests/tests/async.test.js +++ b/integration-tests/tests/async.test.js @@ -160,6 +160,41 @@ test.serial("Statement.all()", async (t) => { t.deepEqual(await stmt.all(), expected); }); +test.serial("Statement.allBatched() returns rows from multiple native batches", async (t) => { + const db = t.context.db; + const stmt = await db.prepare(` + WITH RECURSIVE numbers(value) AS ( + SELECT 1 + UNION ALL + SELECT value + 1 FROM numbers WHERE value < 501 + ) + SELECT value FROM numbers ORDER BY value + `); + + const rows = await stmt.allBatched(100); + + t.is(rows.length, 501); + t.is(rows[0].value, 1); + t.is(rows[500].value, 501); +}); + +test.serial("Statement.allBatched() [raw]", async (t) => { + const db = t.context.db; + const stmt = await db.prepare("SELECT * FROM users ORDER BY id"); + + t.deepEqual(await stmt.raw().allBatched(1), [ + [1, "Alice", "alice@example.org"], + [2, "Bob", "bob@example.com"], + ]); +}); + +test.serial("Statement.allBatched() [pluck and safe integers]", async (t) => { + const db = t.context.db; + const stmt = await db.prepare("SELECT id FROM users ORDER BY id"); + + t.deepEqual(await stmt.pluck().safeIntegers().allBatched(1), [1n, 2n]); +}); + test.serial("Statement.all() [raw]", async (t) => { const db = t.context.db; diff --git a/promise.js b/promise.js index 0b8baf5..8e52ef0 100644 --- a/promise.js +++ b/promise.js @@ -503,6 +503,35 @@ class Statement { } } + /** + * Executes the SQL statement and returns all resulting rows in native batches. + * + * @param {number} batchSize - The maximum number of rows to read per native call. + * @param bindParameters - The bind parameters for executing the statement. + */ + async allBatched(batchSize, ...bindParameters) { + try { + const { params, queryOptions } = splitBindParameters(bindParameters); + const result = []; + const iterator = await this.stmt.iterate(params, queryOptions); + try { + while (true) { + const batch = await iterator.nextBatch(batchSize); + result.push(...batch.records); + if (batch.done) { + return result; + } + } + } finally { + if (typeof iterator.close === "function") { + iterator.close(); + } + } + } catch (err) { + throw convertError(err); + } + } + /** * Interrupts the statement. */ diff --git a/src/lib.rs b/src/lib.rs index 0ff748b..5893a5b 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -906,6 +906,14 @@ fn throw_database_closed_error(env: &Env) -> napi::Error { err } +fn database_not_open_error() -> napi::Error { + throw_sqlite_error( + "The database connection is not open".to_string(), + "SQLITE_NOTOPEN".to_string(), + 0, + ) +} + fn query_timeout_duration(timeout_ms: f64) -> Option { if timeout_ms.is_finite() && timeout_ms > 0.0 { Some(Duration::from_millis(timeout_ms as u64)) @@ -1272,11 +1280,7 @@ impl Statement { fn handle(&self) -> Result { match &*self.slot.0.lock().unwrap() { Some(handle) => Ok(handle.clone()), - None => Err(throw_sqlite_error( - "The database connection is not open".to_string(), - "SQLITE_NOTOPEN".to_string(), - 0, - )), + None => Err(database_not_open_error()), } } @@ -1578,6 +1582,18 @@ impl RowsIteratorState { } } +impl RowsIteratorState { + /// Drops the rows if the database was closed while they were locked. + /// Called after unlocking so that either the caller or release() drops them. + fn drop_rows_if_released(&self) { + if self.released.load(Ordering::SeqCst) { + if let Ok(mut rows) = self.rows.try_lock() { + rows.take(); + } + } + } +} + impl ConnectionResource for RowsIteratorState { fn release(&self) { self.released.store(true, Ordering::SeqCst); @@ -1629,22 +1645,10 @@ impl RowsIterator { } match rows_slot.as_mut() { Some(rows) => rows.next().await, - None => { - return Err(throw_sqlite_error( - "The database connection is not open".to_string(), - "SQLITE_NOTOPEN".to_string(), - 0, - )); - } + None => return Err(database_not_open_error()), } }; - // The database may have been closed while we held the rows. Checking - // after unlocking guarantees either we or release() drops them. - if self.state.released.load(Ordering::SeqCst) { - if let Ok(mut rows_slot) = self.state.rows.try_lock() { - rows_slot.take(); - } - } + self.state.drop_rows_if_released(); let row = match result { Ok(row) => row, Err(err) => { @@ -1670,6 +1674,77 @@ impl RowsIterator { }) } + #[napi(ts_return_type = "Promise<{ records: unknown[]; done: boolean }>")] + pub fn next_batch(&self, env: Env, max_rows: u32) -> Result { + if max_rows == 0 { + return Err(napi::Error::from_reason( + "maxRows must be greater than zero", + )); + } + + let state = self.state.clone(); + let column_count = self.column_names.len(); + let future = async move { + let result = async { + let mut rows_slot = state.rows.lock().await; + if state.released.load(Ordering::SeqCst) { + rows_slot.take(); + } + let Some(rows) = rows_slot.as_mut() else { + return Err(database_not_open_error()); + }; + let mut records = Vec::with_capacity(max_rows as usize); + let mut done = false; + + for _ in 0..max_rows { + let row = match rows.next().await { + Ok(row) => row, + Err(err) => { + state.release_operation_resources(); + return Err(Error::from(err).into()); + } + }; + let Some(row) = row else { + state.release_operation_resources(); + done = true; + break; + }; + let values = match read_row_values(&row, column_count) { + Ok(values) => values, + Err(err) => { + state.release_operation_resources(); + return Err(err); + } + }; + records.push(values); + } + + Ok::<_, napi::Error>((records, done)) + } + .await; + state.drop_rows_if_released(); + result + }; + let column_names = self.column_names.clone(); + let safe_ints = self.safe_ints; + let raw = self.raw; + let pluck = self.pluck; + env.execute_tokio_future(future, move |&mut env, (records, done)| { + let mut js_records = env.create_array(records.len() as u32)?; + for (index, values) in records.iter().enumerate() { + js_records.set( + index as u32, + map_values(&env, &column_names, values, safe_ints, raw, pluck)?, + )?; + } + + let mut result = env.create_object()?; + result.set_named_property("records", js_records)?; + result.set_named_property("done", env.get_boolean(done)?)?; + Ok(result) + }) + } + #[napi] pub fn close(&self) { self.state.release_operation_resources(); @@ -1781,6 +1856,15 @@ pub(crate) fn pin_module_in_memory() { }); } +fn read_row_values(row: &libsql::Row, column_count: usize) -> Result> { + (0..column_count) + .map(|index| { + row.get_value(index as i32) + .map_err(|error| napi::Error::from_reason(error.to_string())) + }) + .collect() +} + fn map_row( env: &Env, column_names: &[std::ffi::CString], @@ -1886,6 +1970,47 @@ fn map_row_raw( Ok(arr.coerce_to_object()?.into_unknown()) } +fn map_values( + env: &Env, + column_names: &[std::ffi::CString], + values: &[libsql::Value], + safe_ints: bool, + raw: bool, + pluck: bool, +) -> Result { + if pluck { + return values + .first() + .map(|value| convert_value_to_js(env, value, safe_ints)) + .transpose()? + .map_or_else(|| Ok(env.get_null()?.into_unknown()), Ok); + } + + if raw { + let mut result = env.create_array(values.len() as u32)?; + for (index, value) in values.iter().enumerate() { + result.set(index as u32, convert_value_to_js(env, value, safe_ints)?)?; + } + return Ok(result.coerce_to_object()?.into_unknown()); + } + + let result = env.create_object()?; + let result = unsafe { napi::JsObject::to_napi_value(env.raw(), result)? }; + for (column_name, value) in column_names.iter().zip(values) { + let js_value = convert_value_to_js(env, value, safe_ints)?; + unsafe { + napi::sys::napi_set_named_property( + env.raw(), + result, + column_name.as_ptr(), + napi::JsUnknown::to_napi_value(env.raw(), js_value)?, + ); + } + } + let result: napi::JsObject = unsafe { napi::JsObject::from_napi_value(env.raw(), result)? }; + Ok(result.into_unknown()) +} + static LOGGER_INIT: OnceCell<()> = OnceCell::new(); fn ensure_logger() { From 67e2b8c8d54fe0f98bb11183fa1b198c59e5e260 Mon Sep 17 00:00:00 2001 From: Julie Yu Date: Sat, 12 Sep 2026 02:04:52 +0000 Subject: [PATCH 02/14] Harden batched row iteration --- index.d.ts | 1 + integration-tests/tests/async.test.js | 62 ++++++++++++++++++++++++++- promise.js | 5 ++- src/lib.rs | 28 ++++++++---- 4 files changed, 86 insertions(+), 10 deletions(-) diff --git a/index.d.ts b/index.d.ts index a7926c8..ad72aa6 100644 --- a/index.d.ts +++ b/index.d.ts @@ -200,6 +200,7 @@ export declare class Statement { /** A raw iterator over rows. The JavaScript layer wraps this in a iterable. */ export declare class RowsIterator { next(): Promise + /** Reads one batch of rows. The batch size must be an integer between 1 and 10,000. */ nextBatch(maxRows: number): Promise<{ records: unknown[]; done: boolean }> close(): void } diff --git a/integration-tests/tests/async.test.js b/integration-tests/tests/async.test.js index 7815e5d..0cfd326 100644 --- a/integration-tests/tests/async.test.js +++ b/integration-tests/tests/async.test.js @@ -178,6 +178,19 @@ test.serial("Statement.allBatched() returns rows from multiple native batches", t.is(rows[500].value, 501); }); +test.serial("Statement.allBatched() rejects invalid batch sizes", async (t) => { + const db = t.context.db; + const stmt = await db.prepare("SELECT * FROM users"); + + for (const batchSize of [0, -1, 1.5, NaN, Infinity, 10_001]) { + await t.throwsAsync(() => stmt.allBatched(batchSize), { + message: "maxRows must be an integer between 1 and 10000", + }); + } + + t.is((await stmt.allBatched(10_000)).length, 2); +}); + test.serial("Statement.allBatched() [raw]", async (t) => { const db = t.context.db; const stmt = await db.prepare("SELECT * FROM users ORDER BY id"); @@ -190,7 +203,7 @@ test.serial("Statement.allBatched() [raw]", async (t) => { test.serial("Statement.allBatched() [pluck and safe integers]", async (t) => { const db = t.context.db; - const stmt = await db.prepare("SELECT id FROM users ORDER BY id"); + const stmt = await db.prepare("SELECT id, email FROM users ORDER BY id"); t.deepEqual(await stmt.pluck().safeIntegers().allBatched(1), [1n, 2n]); }); @@ -488,6 +501,34 @@ test.serial("Query timeout option interrupts long-running query", async (t) => { db.close(); }); +test.serial("Query timeout resets Statement.allBatched() 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.allBatched(100, 1_000_000_000, { queryTimeout: 100 }); + }, { + instanceOf: errorType, + message: "interrupted", + code: "SQLITE_INTERRUPT", + }); + + t.deepEqual(await stmt.allBatched(2, 3), [ + { 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(` @@ -542,6 +583,25 @@ test.serial("Stale timeout guard from exhausted iterator does not interrupt late db.close(); }); +test.serial("Stale timeout guard from exhausted batched iterator does not interrupt later queries", async (t) => { + t.timeout(30_000); + const [db] = await connect(":memory:", { defaultQueryTimeout: 500 }); + + await db.exec("CREATE TABLE t(x INTEGER)"); + const insert = await db.prepare("INSERT INTO t VALUES (?)"); + for (let i = 0; i < 2_000; i++) { + await insert.run(i); + } + + const stmt = await db.prepare("SELECT * FROM t ORDER BY x ASC"); + for (let i = 0; i < 500; i++) { + const rows = await stmt.allBatched(250); + t.is(rows.length, 2_000); + } + + db.close(); +}); + test.serial("Per-query timeout option interrupts long-running Statement.all()", async (t) => { const [db, errorType] = await connect(":memory:"); const stmt = await db.prepare( diff --git a/promise.js b/promise.js index 8e52ef0..3a9b063 100644 --- a/promise.js +++ b/promise.js @@ -507,6 +507,7 @@ class Statement { * Executes the SQL statement and returns all resulting rows in native batches. * * @param {number} batchSize - The maximum number of rows to read per native call. + * Must be an integer between 1 and 10,000. * @param bindParameters - The bind parameters for executing the statement. */ async allBatched(batchSize, ...bindParameters) { @@ -517,7 +518,9 @@ class Statement { try { while (true) { const batch = await iterator.nextBatch(batchSize); - result.push(...batch.records); + for (const record of batch.records) { + result.push(record); + } if (batch.done) { return result; } diff --git a/src/lib.rs b/src/lib.rs index 5893a5b..8e9579e 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1554,6 +1554,8 @@ fn map_value(value: JsUnknown) -> Result { } } +const MAX_ROW_BATCH_SIZE: usize = 10_000; + /// A raw iterator over rows. The JavaScript layer wraps this in a iterable. #[napi] pub struct RowsIterator { @@ -1674,16 +1676,26 @@ impl RowsIterator { }) } + /// Reads one batch of rows. The batch size must be an integer between 1 and 10,000. #[napi(ts_return_type = "Promise<{ records: unknown[]; done: boolean }>")] - pub fn next_batch(&self, env: Env, max_rows: u32) -> Result { - if max_rows == 0 { - return Err(napi::Error::from_reason( - "maxRows must be greater than zero", - )); + pub fn next_batch(&self, env: Env, max_rows: f64) -> Result { + if !max_rows.is_finite() + || max_rows.fract() != 0.0 + || max_rows < 1.0 + || max_rows > MAX_ROW_BATCH_SIZE as f64 + { + return Err(napi::Error::from_reason(format!( + "maxRows must be an integer between 1 and {MAX_ROW_BATCH_SIZE}" + ))); } + let max_rows = max_rows as usize; let state = self.state.clone(); - let column_count = self.column_names.len(); + let value_count = if self.pluck { + self.column_names.len().min(1) + } else { + self.column_names.len() + }; let future = async move { let result = async { let mut rows_slot = state.rows.lock().await; @@ -1693,7 +1705,7 @@ impl RowsIterator { let Some(rows) = rows_slot.as_mut() else { return Err(database_not_open_error()); }; - let mut records = Vec::with_capacity(max_rows as usize); + let mut records = Vec::with_capacity(max_rows); let mut done = false; for _ in 0..max_rows { @@ -1709,7 +1721,7 @@ impl RowsIterator { done = true; break; }; - let values = match read_row_values(&row, column_count) { + let values = match read_row_values(&row, value_count) { Ok(values) => values, Err(err) => { state.release_operation_resources(); From 2e5e35a4f31d92e8ad96e862ce493098cb22042d Mon Sep 17 00:00:00 2001 From: Julie Yu Date: Sat, 12 Sep 2026 02:27:27 +0000 Subject: [PATCH 03/14] Add batched row benchmark --- perf/perf-libsql-batched-rows.js | 63 ++++++++++++++++++++++++++++++++ 1 file changed, 63 insertions(+) create mode 100644 perf/perf-libsql-batched-rows.js diff --git a/perf/perf-libsql-batched-rows.js b/perf/perf-libsql-batched-rows.js new file mode 100644 index 0000000..08d04ab --- /dev/null +++ b/perf/perf-libsql-batched-rows.js @@ -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(`allBatched(${BATCH_SIZE})`, async () => { + validateRows(await stmt.allBatched(BATCH_SIZE, rowCount), 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}`); + } +} From 395417813e330d0396e4872bc1af393e5e68f3bb Mon Sep 17 00:00:00 2001 From: Julie Yu Date: Sat, 12 Sep 2026 02:35:37 +0000 Subject: [PATCH 04/14] Benchmark multiple row batch sizes --- perf/perf-libsql-batched-rows.js | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/perf/perf-libsql-batched-rows.js b/perf/perf-libsql-batched-rows.js index 08d04ab..4f215b2 100644 --- a/perf/perf-libsql-batched-rows.js +++ b/perf/perf-libsql-batched-rows.js @@ -8,6 +8,7 @@ import libsql from '../promise.js'; const { connect } = libsql; const BATCH_SIZE = 250; +const LARGE_RESULT_BATCH_SIZES = [100, BATCH_SIZE, 1_000]; const ROW_COUNTS = [1_000, 10_000, 100_000, 1_000_000]; const MAX_ROW_COUNT = ROW_COUNTS[ROW_COUNTS.length - 1]; @@ -30,13 +31,18 @@ const stmt = await db.prepare(` `); for (const rowCount of ROW_COUNTS) { + const batchSizes = rowCount === MAX_ROW_COUNT + ? LARGE_RESULT_BATCH_SIZES + : [BATCH_SIZE]; group(`${rowCount.toLocaleString('en-US')} rows`, () => { baseline('all()', async () => { validateRows(await stmt.all(rowCount), rowCount); }); - bench(`allBatched(${BATCH_SIZE})`, async () => { - validateRows(await stmt.allBatched(BATCH_SIZE, rowCount), rowCount); - }); + for (const batchSize of batchSizes) { + bench(`allBatched(${batchSize})`, async () => { + validateRows(await stmt.allBatched(batchSize, rowCount), rowCount); + }); + } }); } From d8caa75c236f848199da65faf0e72aad56da327d Mon Sep 17 00:00:00 2001 From: Julie Yu Date: Sat, 12 Sep 2026 02:39:19 +0000 Subject: [PATCH 05/14] Keep row count benchmark --- perf/perf-libsql-batched-rows.js | 12 +++--------- 1 file changed, 3 insertions(+), 9 deletions(-) diff --git a/perf/perf-libsql-batched-rows.js b/perf/perf-libsql-batched-rows.js index 4f215b2..08d04ab 100644 --- a/perf/perf-libsql-batched-rows.js +++ b/perf/perf-libsql-batched-rows.js @@ -8,7 +8,6 @@ import libsql from '../promise.js'; const { connect } = libsql; const BATCH_SIZE = 250; -const LARGE_RESULT_BATCH_SIZES = [100, BATCH_SIZE, 1_000]; const ROW_COUNTS = [1_000, 10_000, 100_000, 1_000_000]; const MAX_ROW_COUNT = ROW_COUNTS[ROW_COUNTS.length - 1]; @@ -31,18 +30,13 @@ const stmt = await db.prepare(` `); for (const rowCount of ROW_COUNTS) { - const batchSizes = rowCount === MAX_ROW_COUNT - ? LARGE_RESULT_BATCH_SIZES - : [BATCH_SIZE]; group(`${rowCount.toLocaleString('en-US')} rows`, () => { baseline('all()', async () => { validateRows(await stmt.all(rowCount), rowCount); }); - for (const batchSize of batchSizes) { - bench(`allBatched(${batchSize})`, async () => { - validateRows(await stmt.allBatched(batchSize, rowCount), rowCount); - }); - } + bench(`allBatched(${BATCH_SIZE})`, async () => { + validateRows(await stmt.allBatched(BATCH_SIZE, rowCount), rowCount); + }); }); } From 116d3c06f8d27714b55014423eb2ffc779659c59 Mon Sep 17 00:00:00 2001 From: Julie Yu Date: Sat, 12 Sep 2026 20:01:46 +0000 Subject: [PATCH 06/14] Simplify batched row iteration --- index.d.ts | 2 +- integration-tests/tests/async.test.js | 19 +++------- promise.js | 4 +- src/lib.rs | 53 ++++++++------------------- 4 files changed, 25 insertions(+), 53 deletions(-) diff --git a/index.d.ts b/index.d.ts index ad72aa6..cc13739 100644 --- a/index.d.ts +++ b/index.d.ts @@ -201,7 +201,7 @@ export declare class Statement { export declare class RowsIterator { next(): Promise /** Reads one batch of rows. The batch size must be an integer between 1 and 10,000. */ - nextBatch(maxRows: number): Promise<{ records: unknown[]; done: boolean }> + nextBatch(maxRows: number): Promise close(): void } export declare class Record { diff --git a/integration-tests/tests/async.test.js b/integration-tests/tests/async.test.js index 0cfd326..d25bafd 100644 --- a/integration-tests/tests/async.test.js +++ b/integration-tests/tests/async.test.js @@ -162,20 +162,13 @@ test.serial("Statement.all()", async (t) => { test.serial("Statement.allBatched() returns rows from multiple native batches", async (t) => { const db = t.context.db; - const stmt = await db.prepare(` - WITH RECURSIVE numbers(value) AS ( - SELECT 1 - UNION ALL - SELECT value + 1 FROM numbers WHERE value < 501 - ) - SELECT value FROM numbers ORDER BY value - `); + const stmt = await db.prepare("SELECT 1 AS value UNION ALL SELECT 2 UNION ALL SELECT 3"); - const rows = await stmt.allBatched(100); - - t.is(rows.length, 501); - t.is(rows[0].value, 1); - t.is(rows[500].value, 501); + t.deepEqual(await stmt.allBatched(2), [ + { value: 1 }, + { value: 2 }, + { value: 3 }, + ]); }); test.serial("Statement.allBatched() rejects invalid batch sizes", async (t) => { diff --git a/promise.js b/promise.js index 3a9b063..e476556 100644 --- a/promise.js +++ b/promise.js @@ -518,10 +518,10 @@ class Statement { try { while (true) { const batch = await iterator.nextBatch(batchSize); - for (const record of batch.records) { + for (const record of batch) { result.push(record); } - if (batch.done) { + if (batch.length < batchSize) { return result; } } diff --git a/src/lib.rs b/src/lib.rs index 8e9579e..621016b 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1677,12 +1677,11 @@ impl RowsIterator { } /// Reads one batch of rows. The batch size must be an integer between 1 and 10,000. - #[napi(ts_return_type = "Promise<{ records: unknown[]; done: boolean }>")] + #[napi(ts_return_type = "Promise")] pub fn next_batch(&self, env: Env, max_rows: f64) -> Result { if !max_rows.is_finite() || max_rows.fract() != 0.0 - || max_rows < 1.0 - || max_rows > MAX_ROW_BATCH_SIZE as f64 + || !(1.0..=MAX_ROW_BATCH_SIZE as f64).contains(&max_rows) { return Err(napi::Error::from_reason(format!( "maxRows must be an integer between 1 and {MAX_ROW_BATCH_SIZE}" @@ -1697,7 +1696,7 @@ impl RowsIterator { self.column_names.len() }; let future = async move { - let result = async { + let result: Result<_> = async { let mut rows_slot = state.rows.lock().await; if state.released.load(Ordering::SeqCst) { rows_slot.take(); @@ -1706,42 +1705,29 @@ impl RowsIterator { return Err(database_not_open_error()); }; let mut records = Vec::with_capacity(max_rows); - let mut done = false; - for _ in 0..max_rows { - let row = match rows.next().await { - Ok(row) => row, - Err(err) => { - state.release_operation_resources(); - return Err(Error::from(err).into()); - } - }; - let Some(row) = row else { - state.release_operation_resources(); - done = true; + let Some(row) = rows.next().await.map_err(Error::from)? else { break; }; - let values = match read_row_values(&row, value_count) { - Ok(values) => values, - Err(err) => { - state.release_operation_resources(); - return Err(err); - } - }; - records.push(values); + records.push(read_row_values(&row, value_count).map_err(Error::from)?); } - - Ok::<_, napi::Error>((records, done)) + Ok(records) } .await; state.drop_rows_if_released(); + if result + .as_ref() + .map_or(true, |records| records.len() < max_rows) + { + state.release_operation_resources(); + } result }; let column_names = self.column_names.clone(); let safe_ints = self.safe_ints; let raw = self.raw; let pluck = self.pluck; - env.execute_tokio_future(future, move |&mut env, (records, done)| { + env.execute_tokio_future(future, move |&mut env, records| { let mut js_records = env.create_array(records.len() as u32)?; for (index, values) in records.iter().enumerate() { js_records.set( @@ -1749,11 +1735,7 @@ impl RowsIterator { map_values(&env, &column_names, values, safe_ints, raw, pluck)?, )?; } - - let mut result = env.create_object()?; - result.set_named_property("records", js_records)?; - result.set_named_property("done", env.get_boolean(done)?)?; - Ok(result) + Ok(js_records) }) } @@ -1868,12 +1850,9 @@ pub(crate) fn pin_module_in_memory() { }); } -fn read_row_values(row: &libsql::Row, column_count: usize) -> Result> { +fn read_row_values(row: &libsql::Row, column_count: usize) -> libsql::Result> { (0..column_count) - .map(|index| { - row.get_value(index as i32) - .map_err(|error| napi::Error::from_reason(error.to_string())) - }) + .map(|index| row.get_value(index as i32)) .collect() } From e2a5d9dc2c6ff8737692140cdfa0a47711aa9ec2 Mon Sep 17 00:00:00 2001 From: Julie Yu Date: Mon, 14 Sep 2026 17:25:51 +0000 Subject: [PATCH 07/14] Batch Promise iterator row reads --- index.d.ts | 1 + integration-tests/tests/async.test.js | 109 ++++++++++++-------------- perf/perf-libsql-batched-rows.js | 4 +- promise.js | 89 +++++++++++---------- src/lib.rs | 2 + 5 files changed, 104 insertions(+), 101 deletions(-) diff --git a/index.d.ts b/index.d.ts index cc13739..c93774a 100644 --- a/index.d.ts +++ b/index.d.ts @@ -18,6 +18,7 @@ export interface Options { /** Per-query execution options. */ export interface QueryOptions { queryTimeout?: number + batchSize?: number } export declare function connect(path: string, opts?: Options | undefined | null): Promise /** Result of a database sync operation. */ diff --git a/integration-tests/tests/async.test.js b/integration-tests/tests/async.test.js index d25bafd..3387654 100644 --- a/integration-tests/tests/async.test.js +++ b/integration-tests/tests/async.test.js @@ -110,6 +110,7 @@ test.serial("Statement.get() [named]", async (t) => { t.is((await stmt.get({ id: 0 })), undefined); t.is((await stmt.get({ id: 1 })).name, "Alice"); t.is((await stmt.get({ id: 2 })).name, "Bob"); + }); @@ -139,6 +140,35 @@ 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 }); + + 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() with invalid bind parameter", async (t) => { const db = t.context.db; @@ -158,47 +188,29 @@ test.serial("Statement.all()", async (t) => { { id: 2, name: "Bob", email: "bob@example.com" }, ]; t.deepEqual(await stmt.all(), expected); -}); -test.serial("Statement.allBatched() returns rows from multiple native batches", async (t) => { - const db = t.context.db; - const stmt = await db.prepare("SELECT 1 AS value UNION ALL SELECT 2 UNION ALL SELECT 3"); - - t.deepEqual(await stmt.allBatched(2), [ - { value: 1 }, - { value: 2 }, - { value: 3 }, - ]); + const namedStmt = await db.prepare("SELECT :batchSize AS value"); + t.deepEqual(await namedStmt.all({ batchSize: 2 }), [{ value: 2 }]); }); -test.serial("Statement.allBatched() rejects invalid batch sizes", async (t) => { +test.serial("Statement.all() rejects invalid batch sizes", async (t) => { const db = t.context.db; const stmt = await db.prepare("SELECT * FROM users"); - for (const batchSize of [0, -1, 1.5, NaN, Infinity, 10_001]) { - await t.throwsAsync(() => stmt.allBatched(batchSize), { + 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.allBatched(10_000)).length, 2); -}); - -test.serial("Statement.allBatched() [raw]", async (t) => { - const db = t.context.db; - const stmt = await db.prepare("SELECT * FROM users ORDER BY id"); - - t.deepEqual(await stmt.raw().allBatched(1), [ - [1, "Alice", "alice@example.org"], - [2, "Bob", "bob@example.com"], - ]); -}); - -test.serial("Statement.allBatched() [pluck and safe integers]", async (t) => { - const db = t.context.db; - const stmt = await db.prepare("SELECT id, email FROM users ORDER BY id"); - - t.deepEqual(await stmt.pluck().safeIntegers().allBatched(1), [1n, 2n]); + t.is((await stmt.all(undefined, { batchSize: 10_000 })).length, 2); }); test.serial("Statement.all() [raw]", async (t) => { @@ -209,7 +221,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) => { @@ -220,7 +232,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) => { @@ -231,7 +243,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) => { @@ -242,7 +254,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) => { @@ -494,7 +506,7 @@ test.serial("Query timeout option interrupts long-running query", async (t) => { db.close(); }); -test.serial("Query timeout resets Statement.allBatched() for reuse", async (t) => { +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 ( @@ -506,14 +518,14 @@ test.serial("Query timeout resets Statement.allBatched() for reuse", async (t) = `); await t.throwsAsync(async () => { - await stmt.allBatched(100, 1_000_000_000, { queryTimeout: 100 }); + await stmt.all(1_000_000_000, { queryTimeout: 100, batchSize: 100 }); }, { instanceOf: errorType, message: "interrupted", code: "SQLITE_INTERRUPT", }); - t.deepEqual(await stmt.allBatched(2, 3), [ + t.deepEqual(await stmt.all(3, { batchSize: 2 }), [ { value: 1 }, { value: 2 }, { value: 3 }, @@ -569,26 +581,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(); - t.is(rows.length, 2_000); - } - - db.close(); -}); - -test.serial("Stale timeout guard from exhausted batched iterator does not interrupt later queries", async (t) => { - t.timeout(30_000); - const [db] = await connect(":memory:", { defaultQueryTimeout: 500 }); - - await db.exec("CREATE TABLE t(x INTEGER)"); - const insert = await db.prepare("INSERT INTO t VALUES (?)"); - for (let i = 0; i < 2_000; i++) { - await insert.run(i); - } - - const stmt = await db.prepare("SELECT * FROM t ORDER BY x ASC"); - for (let i = 0; i < 500; i++) { - const rows = await stmt.allBatched(250); + const rows = await stmt.all(undefined, { batchSize: 250 }); t.is(rows.length, 2_000); } diff --git a/perf/perf-libsql-batched-rows.js b/perf/perf-libsql-batched-rows.js index 08d04ab..ebb9622 100644 --- a/perf/perf-libsql-batched-rows.js +++ b/perf/perf-libsql-batched-rows.js @@ -34,8 +34,8 @@ for (const rowCount of ROW_COUNTS) { baseline('all()', async () => { validateRows(await stmt.all(rowCount), rowCount); }); - bench(`allBatched(${BATCH_SIZE})`, async () => { - validateRows(await stmt.allBatched(BATCH_SIZE, rowCount), rowCount); + bench(`all() with batchSize ${BATCH_SIZE}`, async () => { + validateRows(await stmt.all(rowCount, { batchSize: BATCH_SIZE }), rowCount); }); }); } diff --git a/promise.js b/promise.js index e476556..c1188d8 100644 --- a/promise.js +++ b/promise.js @@ -4,6 +4,8 @@ const { Database: NativeDb, connect: nativeConnect } = require("./index.js"); const SqliteError = require("./sqlite-error.js"); const { Authorization, Action } = require("./auth"); +const DEFAULT_ROW_BATCH_SIZE = 1; + /** * @import {Options as NativeOptions, Statement as NativeStatement} from './index.js' */ @@ -41,18 +43,22 @@ function convertError(err) { return err; } -function isQueryOptions(value) { +function isQueryOptions(value, allowBatchSize) { return value != null && typeof value === "object" && !Array.isArray(value) - && Object.prototype.hasOwnProperty.call(value, "queryTimeout"); + && (Object.prototype.hasOwnProperty.call(value, "queryTimeout") + || (allowBatchSize && Object.prototype.hasOwnProperty.call(value, "batchSize"))); } -function splitBindParameters(bindParameters) { +function splitBindParameters(bindParameters, allowBatchSize = false) { if (bindParameters.length === 0) { return { params: undefined, queryOptions: undefined }; } - if (isQueryOptions(bindParameters[bindParameters.length - 1])) { + if (isQueryOptions( + bindParameters[bindParameters.length - 1], + allowBatchSize && bindParameters.length > 1, + )) { if (bindParameters.length === 1) { return { params: undefined, queryOptions: bindParameters[0] }; } @@ -470,9 +476,9 @@ class Statement { */ async iterate(...bindParameters) { try { - const { params, queryOptions } = splitBindParameters(bindParameters); + const { params, queryOptions } = splitBindParameters(bindParameters, true); const it = await this.stmt.iterate(params, queryOptions); - return wrappedIter(it); + return wrappedIter(it, queryOptions?.batchSize); } catch (err) { throw convertError(err); } @@ -503,38 +509,6 @@ class Statement { } } - /** - * Executes the SQL statement and returns all resulting rows in native batches. - * - * @param {number} batchSize - The maximum number of rows to read per native call. - * Must be an integer between 1 and 10,000. - * @param bindParameters - The bind parameters for executing the statement. - */ - async allBatched(batchSize, ...bindParameters) { - try { - const { params, queryOptions } = splitBindParameters(bindParameters); - const result = []; - const iterator = await this.stmt.iterate(params, queryOptions); - try { - while (true) { - const batch = await iterator.nextBatch(batchSize); - for (const record of batch) { - result.push(record); - } - if (batch.length < batchSize) { - return result; - } - } - } finally { - if (typeof iterator.close === "function") { - iterator.close(); - } - } - } catch (err) { - throw convertError(err); - } - } - /** * Interrupts the statement. */ @@ -559,14 +533,47 @@ class Statement { } } -function wrappedIter(it) { +function wrappedIter(it, batchSize = DEFAULT_ROW_BATCH_SIZE) { + let batch = []; + let done = false; + let index = 0; + let pending = Promise.resolve(); + return { next() { - return it.next().catch((err) => { - throw convertError(err); + pending = pending.then(async () => { + if (index === batch.length) { + if (done) { + return { done: true, value: null }; + } + let nextBatch; + try { + nextBatch = await it.nextBatch(batchSize); + } catch (error) { + if (done) { + return { done: true, value: null }; + } + throw convertError(error); + } + if (done) { + return { done: true, value: null }; + } + batch = nextBatch; + done = batch.length < batchSize; + index = 0; + } + if (batch.length === 0) { + return { done: true, value: null }; + } + return { done: false, value: batch[index++] }; }); + return pending; }, return(value) { + done = true; + batch = []; + index = 0; + pending = pending.catch(() => {}); if (typeof it.close === "function") { it.close(); } diff --git a/src/lib.rs b/src/lib.rs index 621016b..4acfb83 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -212,6 +212,8 @@ pub struct Options { pub struct QueryOptions { // Maximum time in milliseconds that this query is allowed to run. pub queryTimeout: Option, + // Maximum number of rows to read per native iterator call. + pub batchSize: Option, } /// Access mode. From e2f0ba075851f2093f2757b869a30fb3ef42bd53 Mon Sep 17 00:00:00 2001 From: Julie Yu Date: Mon, 14 Sep 2026 17:30:21 +0000 Subject: [PATCH 08/14] Inline default row batch size --- promise.js | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/promise.js b/promise.js index c1188d8..28da60c 100644 --- a/promise.js +++ b/promise.js @@ -4,8 +4,6 @@ const { Database: NativeDb, connect: nativeConnect } = require("./index.js"); const SqliteError = require("./sqlite-error.js"); const { Authorization, Action } = require("./auth"); -const DEFAULT_ROW_BATCH_SIZE = 1; - /** * @import {Options as NativeOptions, Statement as NativeStatement} from './index.js' */ @@ -533,7 +531,7 @@ class Statement { } } -function wrappedIter(it, batchSize = DEFAULT_ROW_BATCH_SIZE) { +function wrappedIter(it, batchSize = 1) { let batch = []; let done = false; let index = 0; From cd0e069372f7cf9cb4bbd6fdd12404c1ef9de822 Mon Sep 17 00:00:00 2001 From: Julie Yu Date: Mon, 14 Sep 2026 17:44:52 +0000 Subject: [PATCH 09/14] Inline row value collection --- src/lib.rs | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/src/lib.rs b/src/lib.rs index 4acfb83..5026eb4 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1711,7 +1711,12 @@ impl RowsIterator { let Some(row) = rows.next().await.map_err(Error::from)? else { break; }; - records.push(read_row_values(&row, value_count).map_err(Error::from)?); + records.push( + (0..value_count) + .map(|index| row.get_value(index as i32)) + .collect::>>() + .map_err(Error::from)?, + ); } Ok(records) } @@ -1852,12 +1857,6 @@ pub(crate) fn pin_module_in_memory() { }); } -fn read_row_values(row: &libsql::Row, column_count: usize) -> libsql::Result> { - (0..column_count) - .map(|index| row.get_value(index as i32)) - .collect() -} - fn map_row( env: &Env, column_names: &[std::ffi::CString], From f1cb6b80b82170fbcee502eaa9f9be929d79e403 Mon Sep 17 00:00:00 2001 From: Julie Yu Date: Mon, 14 Sep 2026 18:01:49 +0000 Subject: [PATCH 10/14] Remove stray test whitespace --- integration-tests/tests/async.test.js | 1 - 1 file changed, 1 deletion(-) diff --git a/integration-tests/tests/async.test.js b/integration-tests/tests/async.test.js index 3387654..888855e 100644 --- a/integration-tests/tests/async.test.js +++ b/integration-tests/tests/async.test.js @@ -110,7 +110,6 @@ test.serial("Statement.get() [named]", async (t) => { t.is((await stmt.get({ id: 0 })), undefined); t.is((await stmt.get({ id: 1 })).name, "Alice"); t.is((await stmt.get({ id: 2 })).name, "Bob"); - }); From 1f76a11d66455b2e115923e0c7b78552e7801d5c Mon Sep 17 00:00:00 2001 From: Julie Yu Date: Mon, 14 Sep 2026 18:28:17 +0000 Subject: [PATCH 11/14] Add default row batch size --- docs/api.md | 11 +++++++---- index.d.ts | 2 ++ integration-tests/tests/async.test.js | 12 ++++++++++++ promise.js | 2 +- src/lib.rs | 19 +++++++++++++++++++ 5 files changed, 41 insertions(+), 5 deletions(-) diff --git a/docs/api.md b/docs/api.md index a449290..7865bf6 100644 --- a/docs/api.md +++ b/docs/api.md @@ -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 reads per native iterator call. It must be an integer from 1 through 10,000 and defaults to 1. + +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: `statement.all(undefined, { batchSize: 250 })`. The function returns a `Database` object. @@ -68,7 +71,7 @@ Convenience wrapper that prepares `sql` and executes `Statement.all`. Returns al | -------------- | ------------------- | -------------------------------------------------------------------- | | sql | string | The SQL statement string. | | bindParameters | any | Optional positional or named bind parameters. | -| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). | +| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100, batchSize: 250 }`). | **Note:** This is an extension in libSQL and not available in `better-sqlite3`. @@ -80,7 +83,7 @@ Convenience wrapper that prepares `sql` and executes `Statement.iterate`. Return | -------------- | ------------------- | -------------------------------------------------------------------- | | sql | string | The SQL statement string. | | bindParameters | any | Optional positional or named bind parameters. | -| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). | +| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100, batchSize: 250 }`). | **Note:** This is an extension in libSQL and not available in `better-sqlite3`. @@ -333,7 +336,7 @@ Executes the SQL statement and returns an array of the resulting rows. | Param | Type | Description | | -------------- | ----------------------------- | ------------------------------------------------ | | bindParameters | array of objects | The bind parameters for executing the statement. | -| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). | +| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100, batchSize: 250 }`). | ### iterate([...bindParameters][, queryOptions]) ⇒ iterator @@ -342,7 +345,7 @@ Executes the SQL statement and returns an iterator to the resulting rows. | Param | Type | Description | | -------------- | ----------------------------- | ------------------------------------------------ | | bindParameters | array of objects | The bind parameters for executing the statement. | -| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). | +| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100, batchSize: 250 }`). | ### pluck([toggleState]) ⇒ this diff --git a/index.d.ts b/index.d.ts index c93774a..f78f970 100644 --- a/index.d.ts +++ b/index.d.ts @@ -14,6 +14,7 @@ export interface Options { encryptionKey?: string remoteEncryptionKey?: string defaultQueryTimeout?: number + defaultBatchSize?: number } /** Per-query execution options. */ export interface QueryOptions { @@ -165,6 +166,7 @@ export declare class Database { } /** SQLite statement object. */ export declare class Statement { + get defaultBatchSize(): number /** * Executes a SQL statement. * diff --git a/integration-tests/tests/async.test.js b/integration-tests/tests/async.test.js index 888855e..4368112 100644 --- a/integration-tests/tests/async.test.js +++ b/integration-tests/tests/async.test.js @@ -212,6 +212,18 @@ test.serial("Statement.all() rejects invalid batch sizes", async (t) => { t.is((await stmt.all(undefined, { batchSize: 10_000 })).length, 2); }); +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) => { const db = t.context.db; diff --git a/promise.js b/promise.js index 28da60c..f507a52 100644 --- a/promise.js +++ b/promise.js @@ -476,7 +476,7 @@ class Statement { try { const { params, queryOptions } = splitBindParameters(bindParameters, true); const it = await this.stmt.iterate(params, queryOptions); - return wrappedIter(it, queryOptions?.batchSize); + return wrappedIter(it, queryOptions?.batchSize ?? this.stmt.defaultBatchSize); } catch (err) { throw convertError(err); } diff --git a/src/lib.rs b/src/lib.rs index 5026eb4..693ae79 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -205,6 +205,8 @@ pub struct Options { pub remoteEncryptionKey: Option, // Default maximum time in milliseconds that a query is allowed to run. pub defaultQueryTimeout: Option, + // Default maximum number of rows to read per native iterator call. + pub defaultBatchSize: Option, } /// Per-query execution options. @@ -240,6 +242,8 @@ pub struct Database { memory: bool, // Maximum time in milliseconds that a query is allowed to run. query_timeout: Option, + // Default maximum number of rows to read per native iterator call. + default_batch_size: f64, // Statements and iterators that hold references to the connection. resources: Arc, } @@ -405,12 +409,17 @@ pub async fn connect(path: String, opts: Option) -> Result { .as_ref() .and_then(|o| o.defaultQueryTimeout) .and_then(query_timeout_duration); + let default_batch_size = opts + .as_ref() + .and_then(|o| o.defaultBatchSize) + .unwrap_or(1.0); Ok(Database { db: Some(db), conn: Some(conn), default_safe_integers, memory, query_timeout, + default_batch_size, resources: Arc::new(OpenResources::default()), }) } @@ -491,6 +500,7 @@ impl Database { stmt, mode, self.query_timeout, + self.default_batch_size, self.resources.clone(), )) } @@ -987,6 +997,8 @@ pub struct Statement { mode: AccessMode, // Maximum time in milliseconds that a query is allowed to run. query_timeout: Option, + // Default maximum number of rows to read per native iterator call. + default_batch_size: f64, } #[napi] @@ -1003,6 +1015,7 @@ impl Statement { stmt: libsql::Statement, mode: AccessMode, query_timeout: Option, + default_batch_size: f64, resources: Arc, ) -> Self { let column_names: Vec = stmt @@ -1023,9 +1036,15 @@ impl Statement { column_names, mode, query_timeout, + default_batch_size, } } + #[napi(getter)] + pub fn default_batch_size(&self) -> f64 { + self.default_batch_size + } + /// Executes a SQL statement. /// /// # Arguments From 19254454b9565788e65e1736375425b35d91152b Mon Sep 17 00:00:00 2001 From: Pekka Enberg Date: Tue, 15 Sep 2026 15:38:54 +0300 Subject: [PATCH 12/14] Stop in-flight row batch when iterator is closed Closing the iterator resets the statement, but a nextBatch() already running kept stepping it, which restarts the query from the first row. If the batch filled up, the statement was left active with its timeout guard dropped, holding a read lock until the statement was reused. Track closure with a flag that the batch loop checks before each step, and reset the statement after the loop whenever the iterator was closed. --- integration-tests/tests/async.test.js | 26 ++++++++++++++++++++++++++ src/lib.rs | 21 ++++++++++++++++++--- 2 files changed, 44 insertions(+), 3 deletions(-) diff --git a/integration-tests/tests/async.test.js b/integration-tests/tests/async.test.js index 4368112..a95ab59 100644 --- a/integration-tests/tests/async.test.js +++ b/integration-tests/tests/async.test.js @@ -168,6 +168,32 @@ test.serial("Statement.iterate() discards buffered rows after return()", async ( 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; diff --git a/src/lib.rs b/src/lib.rs index 693ae79..6505fb5 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1594,6 +1594,8 @@ struct RowsIteratorState { timeout_guard: Mutex>, // Set when the database is closed while next() holds the rows. released: AtomicBool, + // Set by close() so that an in-flight batch stops stepping the reset statement. + closed: AtomicBool, } impl RowsIteratorState { @@ -1606,6 +1608,12 @@ impl RowsIteratorState { } impl RowsIteratorState { + /// Whether the iterator or its database was closed. Stepping the reset + /// statement afterwards would restart the query from the first row. + fn is_closed(&self) -> bool { + self.closed.load(Ordering::SeqCst) || self.released.load(Ordering::SeqCst) + } + /// Drops the rows if the database was closed while they were locked. /// Called after unlocking so that either the caller or release() drops them. fn drop_rows_if_released(&self) { @@ -1647,6 +1655,7 @@ impl RowsIterator { stmt: Mutex::new(Some(stmt)), timeout_guard: Mutex::new(timeout_guard), released: AtomicBool::new(false), + closed: AtomicBool::new(false), }); let weak_state: Weak = Arc::downgrade(&state) as _; resources.register(weak_state); @@ -1727,6 +1736,10 @@ impl RowsIterator { }; let mut records = Vec::with_capacity(max_rows); for _ in 0..max_rows { + // Stepping after close() would restart the query from the first row. + if state.is_closed() { + break; + } let Some(row) = rows.next().await.map_err(Error::from)? else { break; }; @@ -1741,9 +1754,10 @@ impl RowsIterator { } .await; state.drop_rows_if_released(); - if result - .as_ref() - .map_or(true, |records| records.len() < max_rows) + if state.is_closed() + || result + .as_ref() + .map_or(true, |records| records.len() < max_rows) { state.release_operation_resources(); } @@ -1767,6 +1781,7 @@ impl RowsIterator { #[napi] pub fn close(&self) { + self.state.closed.store(true, Ordering::SeqCst); self.state.release_operation_resources(); } } From dc6f8bc97a577a36e92ecd3b3c84f06f159d6a27 Mon Sep 17 00:00:00 2001 From: Pekka Enberg Date: Tue, 15 Sep 2026 15:39:30 +0300 Subject: [PATCH 13/14] Release statement when nextBatch() rejects the batch size The statement has already been stepped by query() when nextBatch() validates maxRows, so returning early left it active with its timeout guard registered. A rejected next() does not trigger return() in `for await`, so nothing else released it. --- integration-tests/tests/async.test.js | 24 ++++++++++++++++++++++++ src/lib.rs | 2 ++ 2 files changed, 26 insertions(+) diff --git a/integration-tests/tests/async.test.js b/integration-tests/tests/async.test.js index a95ab59..f6536a7 100644 --- a/integration-tests/tests/async.test.js +++ b/integration-tests/tests/async.test.js @@ -238,6 +238,30 @@ test.serial("Statement.all() rejects invalid batch sizes", async (t) => { 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"); diff --git a/src/lib.rs b/src/lib.rs index 6505fb5..d4c0270 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1713,6 +1713,8 @@ impl RowsIterator { || max_rows.fract() != 0.0 || !(1.0..=MAX_ROW_BATCH_SIZE as f64).contains(&max_rows) { + // A rejected next() does not trigger return() in `for await`. + self.state.release_operation_resources(); return Err(napi::Error::from_reason(format!( "maxRows must be an integer between 1 and {MAX_ROW_BATCH_SIZE}" ))); From 11c8ea07eb218eb4ab64766128f045140e31259e Mon Sep 17 00:00:00 2001 From: Pekka Enberg Date: Tue, 15 Sep 2026 15:40:05 +0300 Subject: [PATCH 14/14] Document batchSize as promise API only The synchronous API reads rows one at a time and its option parsing only recognizes queryTimeout, so `statement.all(undefined, { batchSize: 250 })` binds the options object and throws, and defaultBatchSize is ignored. --- docs/api.md | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/docs/api.md b/docs/api.md index 7865bf6..195439e 100644 --- a/docs/api.md +++ b/docs/api.md @@ -23,9 +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 reads per native iterator call. It must be an integer from 1 through 10,000 and defaults to 1. +- `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. -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: `statement.all(undefined, { batchSize: 250 })`. +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. @@ -336,7 +336,7 @@ Executes the SQL statement and returns an array of the resulting rows. | Param | Type | Description | | -------------- | ----------------------------- | ------------------------------------------------ | | bindParameters | array of objects | The bind parameters for executing the statement. | -| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100, batchSize: 250 }`). | +| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). The promise API also accepts `batchSize`. | ### iterate([...bindParameters][, queryOptions]) ⇒ iterator @@ -345,7 +345,7 @@ Executes the SQL statement and returns an iterator to the resulting rows. | Param | Type | Description | | -------------- | ----------------------------- | ------------------------------------------------ | | bindParameters | array of objects | The bind parameters for executing the statement. | -| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100, batchSize: 250 }`). | +| queryOptions | object | Optional per-query overrides (for example, `{ queryTimeout: 100 }`). The promise API also accepts `batchSize`. | ### pluck([toggleState]) ⇒ this