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
42 changes: 40 additions & 2 deletions doc/api/sqlite.md
Original file line number Diff line number Diff line change
Expand Up @@ -442,6 +442,11 @@

<!-- YAML
added: v24.10.0
changes:
- version: REPLACEME
pr-url: https://github.com/nodejs/node/pull/65156

Check warning on line 447 in doc/api/sqlite.md

View workflow job for this annotation

GitHub Actions / lint-pr-url

pr-url doesn't match the URL of the current PR.
description: Accessing the invoking database connection from the authorizer
callback now throws.
-->

* `callback` {Function|null} The authorizer function to set, or `null` to
Expand All @@ -467,6 +472,31 @@
* `SQLITE_DENY` - Deny the operation (causes an error).
* `SQLITE_IGNORE` - Ignore the operation (silently skip).

SQLite requires that the authorizer callback not modify the database connection
that invoked it, which includes preparing and stepping statements. Methods that
would do so throw an error with code `ERR_INVALID_STATE` while the callback is
on the stack, including `database.prepare()`, `database.exec()`, the execution
methods of that connection's statements, iterators, and tag stores, and
`database.setAuthorizer()` itself. Other connections remain usable.

The callback can also be invoked from within `statement.run()`,
`statement.get()`, and similar methods, because SQLite may re-prepare a
statement during execution after a schema change.

Separately, a statement that is currently being executed cannot be reentered.
Calling `statement.close()` on it would free the virtual machine that is
running, and re-running it through `statement.run()`, `statement.get()`,
`statement.all()`, `statement.iterate()`, `iterator.next()`,
`iterator.return()`, or the equivalent tag store methods would reset that
virtual machine mid-execution. All of these throw an `ERR_INVALID_STATE` error
instead. This applies to any callback SQLite invokes during execution, such as a
user-defined function. Other statements on the connection remain usable.

Operations that touch no SQLite state stay available from the callback:
`sqlTagStore.clear()`, which only drops cached statements, and `next()` and
`return()` on an already-drained iterator, which keep returning
`{ done: true }`.

```cjs
const { DatabaseSync, constants } = require('node:sqlite');
const db = new DatabaseSync(':memory:');
Expand Down Expand Up @@ -674,6 +704,9 @@
<!-- YAML
added: v22.5.0
changes:
- version: REPLACEME
pr-url: https://github.com/nodejs/node/pull/62757

Check warning on line 708 in doc/api/sqlite.md

View workflow job for this annotation

GitHub Actions / lint-pr-url

pr-url doesn't match the URL of the current PR.
description: Add the `persistent` option.
- version: REPLACEME
pr-url: https://github.com/nodejs/node/pull/65157
description: Throw `ERR_INVALID_ARG_VALUE` if `sql` contains no statements.
Expand All @@ -690,10 +723,14 @@
database options or `true`.
* `allowUnknownNamedParameters` {boolean} If `true`, unknown named parameters
are ignored. **Default:** inherited from database options or `false`.
* `persistent` {boolean} If `true`, hints to SQLite that this statement will
be retained for a long time and likely reused many times. SQLite currently
responds to this hint by avoiding lookaside memory. Corresponds to the
[`SQLITE_PREPARE_PERSISTENT`][] flag. **Default:** `false`.
* Returns: {StatementSync} The prepared statement.

Compiles a SQL statement into a [prepared statement][]. This method is a wrapper
around [`sqlite3_prepare_v2()`][].
around [`sqlite3_prepare_v3()`][].

### `database.createTagStore([maxSize])`

Expand Down Expand Up @@ -1860,6 +1897,7 @@
[`SQLITE_DETERMINISTIC`]: https://www.sqlite.org/c3ref/c_deterministic.html
[`SQLITE_DIRECTONLY`]: https://www.sqlite.org/c3ref/c_deterministic.html
[`SQLITE_MAX_FUNCTION_ARG`]: https://www.sqlite.org/limits.html#max_function_arg
[`SQLITE_PREPARE_PERSISTENT`]: https://sqlite.org/c3ref/c_prepare_dont_log.html#sqlitepreparepersistent
[`SQLTagStore`]: #class-sqltagstore
[`database.applyChangeset()`]: #databaseapplychangesetchangeset-options
[`database.createTagStore()`]: #databasecreatetagstoremaxsize
Expand All @@ -1885,7 +1923,7 @@
[`sqlite3_get_autocommit()`]: https://sqlite.org/c3ref/get_autocommit.html
[`sqlite3_last_insert_rowid()`]: https://www.sqlite.org/c3ref/last_insert_rowid.html
[`sqlite3_load_extension()`]: https://www.sqlite.org/c3ref/load_extension.html
[`sqlite3_prepare_v2()`]: https://www.sqlite.org/c3ref/prepare.html
[`sqlite3_prepare_v3()`]: https://www.sqlite.org/c3ref/prepare.html
[`sqlite3_serialize()`]: https://sqlite.org/c3ref/serialize.html
[`sqlite3_set_authorizer()`]: https://sqlite.org/c3ref/set_authorizer.html
[`sqlite3_sql()`]: https://www.sqlite.org/c3ref/expanded_sql.html
Expand Down
243 changes: 215 additions & 28 deletions lib/internal/streams/readable.js
Original file line number Diff line number Diff line change
Expand Up @@ -23,13 +23,18 @@

const {
ArrayPrototypeIndexOf,
AsyncIteratorPrototype,
FunctionPrototypeCall,
NumberIsInteger,
NumberIsNaN,
NumberParseInt,
ObjectDefineProperties,
ObjectKeys,
ObjectSetPrototypeOf,
Promise,
PromisePrototypeThen,
PromiseReject,
PromiseResolve,
ReflectApply,
SafeSet,
Symbol,
Expand Down Expand Up @@ -100,6 +105,7 @@ const FastBuffer = Buffer[SymbolSpecies];

const { StringDecoder } = require('string_decoder');
const from = require('internal/streams/from');
const FixedQueue = require('internal/fixed_queue');

ObjectSetPrototypeOf(Readable.prototype, Stream.prototype);
ObjectSetPrototypeOf(Readable, Stream);
Expand Down Expand Up @@ -1386,10 +1392,22 @@ function streamToAsyncIterator(stream, options) {
return iter;
}

async function* createAsyncIterator(stream, options) {
// Async iterator over a Readable. Requests received while another is
// outstanding are queued and processed in order.
function createAsyncIterator(stream, options) {
let callback = nop;

function next(resolve) {
let error; // undefined: active, null: ended cleanly, else: Error
let started = false;
let completed = false;
let inFlight = false; // An asynchronous request is outstanding
let queue = null; // Requests received while inFlight
let draining = false;
let cleanup;

// Used both as the 'readable' listener (where `this === stream`) and
// as a promise executor storing the resolver that wakes up a pending
// pump().
function wakeup(resolve) {
if (this === stream) {
callback();
callback = nop;
Expand All @@ -1398,32 +1416,23 @@ async function* createAsyncIterator(stream, options) {
}
}

stream.on('readable', next);
function start() {
started = true;

let error;
const cleanup = eos(stream, { writable: false }, (err) => {
error = err ? aggregateTwoErrors(error, err) : null;
callback();
callback = nop;
});
stream.on('readable', wakeup);

cleanup = eos(stream, { writable: false }, (err) => {
error = err ? aggregateTwoErrors(error, err) : null;
callback();
callback = nop;
});
}

// Complete the iterator and either destroy the stream or detach
// from it.
function finalize() {
completed = true;

try {
while (true) {
const chunk = stream.destroyed ? null : stream.read();
if (chunk !== null) {
yield chunk;
} else if (error) {
throw error;
} else if (error === null) {
return;
} else {
await new Promise(next);
}
}
} catch (err) {
error = aggregateTwoErrors(error, err);
throw error;
} finally {
const preserveHalfOpenDuplex =
error === null &&
stream.allowHalfOpen === true &&
Expand All @@ -1437,10 +1446,188 @@ async function* createAsyncIterator(stream, options) {
) {
destroyImpl.destroyer(stream, null);
} else {
stream.off('readable', next);
stream.off('readable', wakeup);
cleanup();
}
}

function settleError(err, reject) {
error = aggregateTwoErrors(error, err);
finalize();
reject(error);
}

function drain() {
// Requests settled synchronously call back into drain(); the guard
// keeps a single loop going instead of recursing once per request.
if (draining) {
return;
}
draining = true;
try {
while (!inFlight && !queue.isEmpty()) {
const req = queue.shift();
if (req.type === 'next') {
processNext(req.resolve, req.reject);
} else if (req.type === 'return') {
processReturn(req.value, req.resolve);
} else {
processThrow(req.value, req.reject);
}
}
} finally {
draining = false;
}
}

// Thenable chunks are unwrapped before delivery; a rejection tears
// down the iterator and the stream.
function onChunkFulfilled(value) {
inFlight = false;
if (queue !== null) drain();
return { done: false, value };
}

function onChunkRejected(err) {
inFlight = false;
error = aggregateTwoErrors(error, err);
finalize();
if (queue !== null) drain();
throw error;
}

// Runs with inFlight === true; settles the request and hands over to
// any requests that queued up behind it.
function pump(resolve, reject) {
const chunk = stream.destroyed ? null : stream.read();
if (chunk !== null) {
// Read `then` only once so that a getter cannot observe (or throw
// on) a second access.
const then = chunk.then;
if (typeof then === 'function') {
FunctionPrototypeCall(then, chunk, (value) => {
inFlight = false;
resolve({ done: false, value });
if (queue !== null) drain();
}, (err) => {
inFlight = false;
settleError(err, reject);
if (queue !== null) drain();
});
return;
}
inFlight = false;
resolve({ done: false, value: chunk });
if (queue !== null) drain();
} else if (error) {
inFlight = false;
settleError(error, reject);
if (queue !== null) drain();
} else if (error === null) {
inFlight = false;
finalize();
resolve({ done: true, value: undefined });
if (queue !== null) drain();
} else {
// No data buffered yet; wait for 'readable' or end-of-stream and
// retry.
PromisePrototypeThen(new Promise(wakeup), () => pump(resolve, reject));
}
}

function processNext(resolve, reject) {
if (completed) {
resolve({ done: true, value: undefined });
return;
}
if (!started) start();
inFlight = true;
pump(resolve, reject);
}

function processReturn(value, resolve) {
if (!completed) {
if (started) {
finalize();
} else {
// Never started: complete without touching the stream.
completed = true;
}
}
resolve({ done: true, value });
}

function processThrow(err, reject) {
if (completed || !started) {
completed = true;
reject(err);
return;
}
settleError(err, reject);
}

return {
__proto__: AsyncIteratorPrototype,
next() {
if (!inFlight && !completed) {
if (!started) start();
// Fast path: a chunk is already buffered.
const chunk = stream.destroyed ? null : stream.read();
if (chunk !== null) {
// Read `then` only once so that a getter cannot observe (or
// throw on) a second access.
const then = chunk.then;
if (typeof then === 'function') {
inFlight = true;
return FunctionPrototypeCall(
then, chunk, onChunkFulfilled, onChunkRejected);
}
return PromiseResolve({ done: false, value: chunk });
}
if (error) {
finalize();
return PromiseReject(error);
}
if (error === null) {
finalize();
return PromiseResolve({ done: true, value: undefined });
}
// No data buffered yet; wait for 'readable' or end-of-stream.
inFlight = true;
return new Promise((resolve, reject) => {
PromisePrototypeThen(new Promise(wakeup), () => pump(resolve, reject));
});
}
return new Promise((resolve, reject) => {
if (inFlight) {
queue ??= new FixedQueue();
queue.push({ __proto__: null, type: 'next', value: undefined, resolve, reject });
} else {
resolve({ done: true, value: undefined });
}
});
},
return(value) {
return new Promise((resolve, reject) => {
if (inFlight) {
queue ??= new FixedQueue();
queue.push({ __proto__: null, type: 'return', value, resolve, reject });
} else {
processReturn(value, resolve);
}
});
},
throw(err) {
return new Promise((resolve, reject) => {
if (inFlight) {
queue ??= new FixedQueue();
queue.push({ __proto__: null, type: 'throw', value: err, resolve, reject });
} else {
processThrow(err, reject);
}
});
},
};
}

let composeImpl;
Expand Down
Loading
Loading