diff --git a/cpp/OPBridge.cpp b/cpp/OPBridge.cpp index b6ac8084..80b559b5 100644 --- a/cpp/OPBridge.cpp +++ b/cpp/OPBridge.cpp @@ -889,6 +889,10 @@ void opsqlite_deregister_rollback_hook(sqlite3 *db) { sqlite3_rollback_hook(db, nullptr, nullptr); } +bool opsqlite_in_transaction(sqlite3 *db) { + return sqlite3_get_autocommit(db) == 0; +} + void opsqlite_load_extension(sqlite3 *db, std::string &path, std::string &entry_point) { #ifdef OP_SQLITE_USE_PHONE_VERSION @@ -932,22 +936,24 @@ opsqlite_execute_batch(sqlite3 *db, throw std::runtime_error("No SQL commands provided"); } - int affectedRows = 0; - // opsqlite_execute(db, "BEGIN EXCLUSIVE TRANSACTION", nullptr); - for (int i = 0; i < commandCount; i++) { - const auto &command = commands->at(i); - // We do not provide a datastructure to receive query data because we - // don't need/want to handle this results in a batch execution - // There is also no need to commit/catch this transaction, this is done - // in the JS code - auto result = opsqlite_execute(db, command.sql, &command.params); - affectedRows += result.affectedRows; - } + return run_in_transaction( + [db](const char *sql) { opsqlite_execute(db, sql, nullptr); }, + [db]() -> std::optional { return opsqlite_in_transaction(db); }, + "BEGIN TRANSACTION", [&]() { + int affectedRows = 0; + for (int i = 0; i < commandCount; i++) { + const auto &command = commands->at(i); + // We do not provide a datastructure to receive query data because + // we don't need/want to handle this results in a batch execution + auto result = opsqlite_execute(db, command.sql, &command.params); + affectedRows += result.affectedRows; + } - return BatchResult{ - .affectedRows = affectedRows, - .commands = static_cast(commandCount), - }; + return BatchResult{ + .affectedRows = affectedRows, + .commands = static_cast(commandCount), + }; + }); } } // namespace opsqlite diff --git a/cpp/OPBridge.hpp b/cpp/OPBridge.hpp index 26cd1439..921965b2 100644 --- a/cpp/OPBridge.hpp +++ b/cpp/OPBridge.hpp @@ -85,6 +85,11 @@ void opsqlite_deregister_commit_hook(sqlite3 *db); void opsqlite_register_rollback_hook(sqlite3 *db, void *opsqlite_db_ptr); void opsqlite_deregister_rollback_hook(sqlite3 *db); +/// Whether the connection has an open transaction, i.e. it is not in +/// autocommit mode. SQLite leaves it on its own after RAISE(ROLLBACK) or a +/// failed COMMIT, so callers check this before issuing a ROLLBACK. +bool opsqlite_in_transaction(sqlite3 *db); + sqlite3_stmt *opsqlite_prepare_statement(sqlite3 *db, std::string const &query); void opsqlite_finalize_statement(sqlite3_stmt *statement); diff --git a/cpp/OPDatabase.cpp b/cpp/OPDatabase.cpp index 94609d43..c4a21bb0 100644 --- a/cpp/OPDatabase.cpp +++ b/cpp/OPDatabase.cpp @@ -396,6 +396,15 @@ void OPDatabase::create_jsi_functions(jsi::Runtime &rt, #endif })); + // libsql exposes no way to read the autocommit state, the JS side falls + // back to always attempting the ROLLBACK when this is missing +#ifndef OP_SQLITE_USE_LIBSQL + js_object.setProperty(rt, "inTransaction", HFN(this) { + throw_if_closed("inTransaction"); + return jsi::Value(opsqlite_in_transaction(db)); + })); +#endif + js_object.setProperty(rt, "delete", HFN(this) { if (count != 0) { throw std::runtime_error("[op-sqlite] Delete no longer takes arguments"); diff --git a/cpp/OPUtils.cpp b/cpp/OPUtils.cpp index 88160e88..7f3ccb2d 100644 --- a/cpp/OPUtils.cpp +++ b/cpp/OPUtils.cpp @@ -327,34 +327,22 @@ BatchResult import_sql_file(sqlite3 *db, std::string path) { throw std::runtime_error("Could not open file: " + path); } - try { - int affectedRows = 0; - int commands = 0; - opsqlite_execute(db, "BEGIN EXCLUSIVE TRANSACTION", nullptr); - while (std::getline(sqFile, line, '\n')) { - if (!line.empty()) { - try { - auto result = opsqlite_execute(db, line, nullptr); - affectedRows += result.affectedRows; - commands++; - } catch (std::exception &) { - opsqlite_execute(db, "ROLLBACK", nullptr); - sqFile.close(); - // Rethrow the original exception object: `throw exc` would copy it - // into a plain std::exception, dropping both the message and the - // SQLite result codes. - throw; + // sqFile is closed by its destructor, on success and on failure alike + return run_in_transaction( + [db](const char *sql) { opsqlite_execute(db, sql, nullptr); }, + [db]() -> std::optional { return opsqlite_in_transaction(db); }, + "BEGIN EXCLUSIVE TRANSACTION", [&]() { + int affectedRows = 0; + int commands = 0; + while (std::getline(sqFile, line, '\n')) { + if (!line.empty()) { + auto result = opsqlite_execute(db, line, nullptr); + affectedRows += result.affectedRows; + commands++; + } } - } - } - sqFile.close(); - opsqlite_execute(db, "COMMIT", nullptr); - return {"", affectedRows, commands}; - } catch (std::exception &) { - sqFile.close(); - opsqlite_execute(db, "ROLLBACK", nullptr); - throw; - } + return BatchResult{"", affectedRows, commands}; + }); } #endif diff --git a/cpp/OPUtils.hpp b/cpp/OPUtils.hpp index b9827ca2..46c2e9a8 100644 --- a/cpp/OPUtils.hpp +++ b/cpp/OPUtils.hpp @@ -10,6 +10,8 @@ #include #endif #include +#include +#include #include #include #include "OPThreadPool.hpp" @@ -45,6 +47,47 @@ void to_batch_arguments(jsi::Runtime &rt, jsi::Array const &batch_params, BatchResult import_sql_file(sqlite3 *db, std::string path); +/// Runs `work` between `begin` and COMMIT, rolling back if either throws. +/// +/// SQLite sometimes rolls the transaction back on its own before we get to it +/// (RAISE(ROLLBACK) in a trigger, a COMMIT failing with SQLITE_FULL or +/// SQLITE_IOERR*). A ROLLBACK then fails with "cannot rollback - no +/// transaction is active", which must never replace the original exception. +/// +/// `in_transaction` returns std::nullopt on backends that cannot report the +/// transaction state (libsql): the ROLLBACK is always attempted there and its +/// failure dropped. +template +auto run_in_transaction(Execute &&execute, InTransaction &&in_transaction, + const char *begin, Work &&work) -> decltype(work()) { + execute(begin); + try { + auto result = work(); + execute("COMMIT"); + return result; + } catch (std::exception &error) { + std::optional open = in_transaction(); + if (!open.has_value() || *open) { + try { + execute("ROLLBACK"); + } catch (std::exception &rollback_error) { + // Every following statement would silently run inside the failed + // transaction, which matters more than the original error + std::optional still_open = in_transaction(); + if (still_open.has_value() && *still_open) { + throw std::runtime_error( + std::string("[op-sqlite] ROLLBACK failed and the connection is " + "still inside a transaction: ") + + rollback_error.what() + ". Original error: " + error.what()); + } + } + } + // Rethrow the original exception object: `throw error` would copy it + // into a plain std::exception, dropping the SQLite result codes. + throw; + } +} + bool folder_exists(const std::string &name); bool file_exists(const std::string &path); diff --git a/cpp/libsql/OPLibsqlBridge.cpp b/cpp/libsql/OPLibsqlBridge.cpp index 422b0938..2b966365 100644 --- a/cpp/libsql/OPLibsqlBridge.cpp +++ b/cpp/libsql/OPLibsqlBridge.cpp @@ -6,6 +6,7 @@ #include "OPUtils.hpp" #include #include +#include #include #include @@ -204,6 +205,17 @@ void opsqlite_libsql_remove(DB &db, std::string const &name, remove(full_path.c_str()); } +/// libsql_next_row is where the statement actually runs, so errors raised +/// while running it (constraint violations, RAISE() in a trigger...) are +/// reported there and not by libsql_query_stmt. Copied before the rows and +/// statement are released, then thrown once they are. +std::optional step_error(int status, const char *err) { + if (status == 0) { + return std::nullopt; + } + return std::string(err != nullptr ? err : "libsql error"); +} + void opsqlite_libsql_bind_statement(libsql_stmt_t statement, const std::vector *values) { const char *err; @@ -344,9 +356,7 @@ BridgeResult opsqlite_libsql_execute_prepared_statement( err = nullptr; } - if (status != 0) { - fprintf(stderr, "%s\n", err); - } + auto error = step_error(status, err); libsql_free_rows(rows); @@ -355,6 +365,10 @@ BridgeResult opsqlite_libsql_execute_prepared_statement( libsql_reset_stmt(stmt, &err); + if (error) { + throw std::runtime_error(*error); + } + return {.affectedRows = static_cast(changes), .insertId = static_cast(insert_row_id)}; } @@ -478,9 +492,15 @@ BridgeResult opsqlite_libsql_execute(DB const &db, std::string const &query, status = libsql_next_row(rows, &row, &err); } + auto error = step_error(status, err); + libsql_free_rows(rows); libsql_free_stmt(stmt); + if (error) { + throw std::runtime_error(*error); + } + unsigned long long changes = libsql_changes(db.c); long long insert_row_id = libsql_last_insert_rowid(db.c); @@ -603,13 +623,15 @@ BridgeResult opsqlite_libsql_execute_with_host_objects( err = nullptr; } - if (status != 0) { - fprintf(stderr, "%s\n", err); - } + auto error = step_error(status, err); libsql_free_rows(rows); libsql_free_stmt(stmt); + if (error) { + throw std::runtime_error(*error); + } + unsigned long long changes = libsql_changes(db.c); long long insert_row_id = libsql_last_insert_rowid(db.c); @@ -727,13 +749,15 @@ opsqlite_libsql_execute_raw(DB const &db, std::string const &query, err = nullptr; } - if (status != 0) { - fprintf(stderr, "%s\n", err); - } + auto error = step_error(status, err); libsql_free_rows(rows); libsql_free_stmt(stmt); + if (error) { + throw std::runtime_error(*error); + } + unsigned long long changes = libsql_changes(db.c); long long insert_row_id = libsql_last_insert_rowid(db.c); @@ -750,23 +774,26 @@ opsqlite_libsql_execute_batch(DB const &db, throw std::runtime_error("No SQL commands provided"); } - int affectedRows = 0; - // Transaction control (BEGIN/COMMIT/ROLLBACK) is left to the JS side, so - // any exception here must propagate to reject the JS promise instead of - // being swallowed - otherwise the wrapping COMMIT would persist a partial - // batch instead of the ROLLBACK the caller expects. - for (int i = 0; i < commandCount; i++) { - auto command = commands->at(i); - // We do not provide a datastructure to receive query data because - // we don't need/want to handle this results in a batch execution - auto result = opsqlite_libsql_execute(db, command.sql, &command.params); - affectedRows += result.affectedRows; - } - - return BatchResult{ - .affectedRows = affectedRows, - .commands = static_cast(commandCount), - }; + return run_in_transaction( + [&db](const char *sql) { opsqlite_libsql_execute(db, sql, nullptr); }, + // libsql exposes no way to read the autocommit state + []() -> std::optional { return std::nullopt; }, + "BEGIN TRANSACTION", [&]() { + int affectedRows = 0; + for (int i = 0; i < commandCount; i++) { + const auto &command = commands->at(i); + // We do not provide a datastructure to receive query data because + // we don't need/want to handle this results in a batch execution + auto result = + opsqlite_libsql_execute(db, command.sql, &command.params); + affectedRows += result.affectedRows; + } + + return BatchResult{ + .affectedRows = affectedRows, + .commands = static_cast(commandCount), + }; + }); } } // namespace opsqlite diff --git a/cpp/turso/OPTursoBridge.cpp b/cpp/turso/OPTursoBridge.cpp index fd69a1d8..6fb0d0d5 100644 --- a/cpp/turso/OPTursoBridge.cpp +++ b/cpp/turso/OPTursoBridge.cpp @@ -884,6 +884,12 @@ void opsqlite_register_rollback_hook( void opsqlite_deregister_rollback_hook([[maybe_unused]] sqlite3 *db) {} +bool opsqlite_in_transaction(sqlite3 *db) { + auto connection = + require_turso_connection(to_turso_db(db), "in_transaction"); + return !turso_connection_get_autocommit(connection); +} + void opsqlite_load_extension([[maybe_unused]] sqlite3 *db, [[maybe_unused]] std::string &path, [[maybe_unused]] std::string &entry_point) { @@ -899,16 +905,20 @@ opsqlite_execute_batch(sqlite3 *db, throw std::runtime_error("No SQL commands provided"); } - int affected_rows = 0; - - for (size_t i = 0; i < command_count; i++) { - const auto &command = commands->at(i); - auto result = opsqlite_execute(db, command.sql, &command.params); - affected_rows += result.affectedRows; - } + return run_in_transaction( + [db](const char *sql) { opsqlite_execute(db, sql, nullptr); }, + [db]() -> std::optional { return opsqlite_in_transaction(db); }, + "BEGIN TRANSACTION", [&]() { + int affected_rows = 0; + for (size_t i = 0; i < command_count; i++) { + const auto &command = commands->at(i); + auto result = opsqlite_execute(db, command.sql, &command.params); + affected_rows += result.affectedRows; + } - return BatchResult{.affectedRows = affected_rows, - .commands = static_cast(command_count)}; + return BatchResult{.affectedRows = affected_rows, + .commands = static_cast(command_count)}; + }); } } // namespace opsqlite diff --git a/docs/docs/api.md b/docs/docs/api.md index f11cce16..834c1e7c 100644 --- a/docs/docs/api.md +++ b/docs/docs/api.md @@ -314,13 +314,7 @@ const res = await db.executeBatch(commands); console.log(`Batch affected ${result.rowsAffected} rows`); ``` -`executeBatch` runs the `BEGIN`/`COMMIT`/`ROLLBACK` statements that wrap the batch asynchronously, off the JS thread. For very large batches the `COMMIT` is where SQLite actually writes the WAL frames/fsyncs, so this keeps the JS thread free while that happens. - -If you need those transaction boundaries to run synchronously on the JS thread instead, use `executeBatchSync`, which has the same signature and behavior otherwise: - -```tsx -const res = await db.executeBatchSync(commands); -``` +The `BEGIN`/`COMMIT`/`ROLLBACK` statements that wrap the batch run natively, in the same background call as the batch itself, so the JS thread stays free while SQLite commits. In some scenarios, dynamic applications may need to get some metadata information about the returned result set. diff --git a/docs/docs/changelog.md b/docs/docs/changelog.md index c82d886e..8b6882c3 100644 --- a/docs/docs/changelog.md +++ b/docs/docs/changelog.md @@ -4,9 +4,14 @@ sidebar_position: 11 # API Changes -## 18.1.0 +## 19.0.0 - Errors coming from SQLite now carry their result codes: rejected/thrown `Error`s from `execute`, `executeSync`, `executeRaw`, `executeRawSync`, `executeBatch`, `prepareStatement`, `attach`, `detach`, `loadExtension` and `open` expose `code` (primary) and `extendedCode` (extended), and both are repeated in the message. Only the plain SQLite3 and SQLCipher backends report them; libsql, Turso, web and node only expose a message. See [Error codes](./api.md#error-codes). +- `executeBatch`, `transaction` and `loadFile` now always reject with the error that made them fail. Before, when SQLite had already rolled the transaction back on its own (`RAISE(ROLLBACK)` in a trigger, or a `COMMIT` failing with `SQLITE_FULL`/`SQLITE_IOERR`), the wrapper's `ROLLBACK` failed with `cannot rollback - no transaction is active` and that error replaced the real one. If the connection is left inside the transaction after a failed `ROLLBACK`, a new error saying so is thrown instead. In `transaction`, a failed `ROLLBACK` is attached to the original error as `rollbackError`, and in the stuck case the original error is kept as `cause`. +- The `BEGIN`/`COMMIT`/`ROLLBACK` around `executeBatch` now run natively, in the same background call as the batch, instead of as separate calls from JS. +- **Breaking:** Removed `executeBatchSync`. Its only difference was running `BEGIN`/`COMMIT` synchronously on the JS thread, and those now run natively for `executeBatch`. Use `executeBatch` instead. +- `tx.rollback()` no longer throws when SQLite has already rolled the transaction back (not on libsql, which cannot report the transaction state). +- libsql: errors raised while a statement runs (constraint violations, `RAISE()` in a trigger, a failing `COMMIT`...) are now thrown. They used to be printed to stderr and dropped, so the call resolved as if it had succeeded. Only errors raised while preparing the statement, such as a syntax error or a missing table, were reported. ## 18.0.0 diff --git a/example/src/tests/index.ts b/example/src/tests/index.ts index 383f1e52..6a8df070 100644 --- a/example/src/tests/index.ts +++ b/example/src/tests/index.ts @@ -6,6 +6,7 @@ import "./hooks"; import "./preparedStatements"; import "./queries"; import "./reactive"; +import "./rollback"; import "./storage"; import "./tokenizer"; import "./transactions"; diff --git a/example/src/tests/queries.ts b/example/src/tests/queries.ts index b21a1515..6bae7100 100644 --- a/example/src/tests/queries.ts +++ b/example/src/tests/queries.ts @@ -483,46 +483,6 @@ describe("Queries tests", () => { ]); }); - it("executeBatchSync", async () => { - const id1 = chance.integer(); - const name1 = chance.name(); - const age1 = chance.integer(); - const networth1 = chance.floating(); - - const id2 = chance.integer(); - const name2 = chance.name(); - const age2 = chance.integer(); - const networth2 = chance.floating(); - - const commands: SQLBatchTuple[] = [ - ['SELECT * FROM "User"', []], - ['SELECT * FROM "User"'], - [ - 'INSERT INTO "User" (id, name, age, networth) VALUES(?, ?, ?, ?)', - [id1, name1, age1, networth1], - ], - [ - 'INSERT INTO "User" (id, name, age, networth) VALUES(?, ?, ?, ?)', - [[id2, name2, age2, networth2]], - ], - ]; - - await db.executeBatchSync(commands); - - const res = await db.execute("SELECT * FROM User"); - - expect(res.rows).toDeepEqual([ - { id: id1, name: name1, age: age1, networth: networth1, nickname: null }, - { - id: id2, - name: name2, - age: age2, - networth: networth2, - nickname: null, - }, - ]); - }); - it("executeBatch rolls back on error", async () => { const id1 = chance.integer(); const name1 = chance.name(); @@ -548,31 +508,6 @@ describe("Queries tests", () => { expect(res.rows).toDeepEqual([]); }); - it("executeBatchSync rolls back on error", async () => { - const id1 = chance.integer(); - const name1 = chance.name(); - const age1 = chance.integer(); - const networth1 = chance.floating(); - - const commands: SQLBatchTuple[] = [ - [ - 'INSERT INTO "User" (id, name, age, networth) VALUES(?, ?, ?, ?)', - [id1, name1, age1, networth1], - ], - ["INSERT INTO [tableThatDoesNotExist] (id) VALUES(1)"], - ]; - - try { - await db.executeBatchSync(commands); - throw new Error("Should not resolve"); - } catch (e) { - expect(((e as Error)?.message?.length ?? 0) > 0).toBe(true); - } - - const res = await db.execute("SELECT * FROM User"); - expect(res.rows).toDeepEqual([]); - }); - it("Batch execute with BLOB", async () => { const db = open({ name: "queries.sqlite", diff --git a/example/src/tests/rollback.ts b/example/src/tests/rollback.ts new file mode 100644 index 00000000..94225792 --- /dev/null +++ b/example/src/tests/rollback.ts @@ -0,0 +1,125 @@ +import { type DB, isLibsql, isTurso, open } from "@op-engineering/op-sqlite"; +import { afterEach, beforeEach, describe, expect, it } from "@op-engineering/op-test"; + +async function captureError(fn: () => unknown): Promise { + try { + await fn(); + } catch (e) { + return e as Error; + } + + throw new Error("Expected the call to fail, it did not"); +} + +// A trigger doing RAISE(ROLLBACK) makes SQLite roll the transaction back on its +// own, so the wrapper's ROLLBACK must not replace the "vetoed" error with +// "cannot rollback - no transaction is active". +describe("Rollback after SQLite already rolled back", () => { + let db: DB; + + beforeEach(async () => { + db = open({ + name: "rollback.sqlite", + encryptionKey: "test", + }); + + await db.execute("DROP TABLE IF EXISTS t;"); + await db.execute("CREATE TABLE t (x INTEGER);"); + // Turso's engine does not support triggers + if (!isTurso()) { + await db.execute( + "CREATE TRIGGER veto BEFORE INSERT ON t WHEN NEW.x = 13 BEGIN SELECT RAISE(ROLLBACK, 'vetoed'); END;", + ); + } + }); + + afterEach(() => { + if (db) { + db.delete(); + // @ts-expect-error + db = null; + } + }); + + it("executeBatch rejects with the original error", async () => { + if (isTurso()) { + return; + } + + const error = await captureError(() => + db.executeBatch([ + ["INSERT INTO t VALUES (?)", [1]], + ["INSERT INTO t VALUES (?)", [13]], + ]), + ); + + expect(error.message).toContain("vetoed"); + + const res = await db.execute("SELECT COUNT(*) AS count FROM t"); + expect(res.rows[0]!.count).toEqual(0); + }); + + it("transaction rejects with the original error", async () => { + if (isTurso()) { + return; + } + + const error = await captureError(() => + db.transaction(async (tx) => { + await tx.execute("INSERT INTO t VALUES (?)", [1]); + await tx.execute("INSERT INTO t VALUES (?)", [13]); + }), + ); + + expect(error.message).toContain("vetoed"); + + const res = await db.execute("SELECT COUNT(*) AS count FROM t"); + expect(res.rows[0]!.count).toEqual(0); + }); + + it("tx.rollback() after SQLite rolled back does not throw", async () => { + // libsql cannot report the transaction state, so the ROLLBACK is issued + // and fails with "no transaction is active" + if (isTurso() || isLibsql()) { + return; + } + + await db.transaction(async (tx) => { + try { + await tx.execute("INSERT INTO t VALUES (?)", [13]); + } catch { + tx.rollback(); + } + }); + + const res = await db.execute("SELECT COUNT(*) AS count FROM t"); + expect(res.rows[0]!.count).toEqual(0); + }); + + it("the connection is usable after the veto", async () => { + if (isTurso()) { + return; + } + + await captureError(() => db.executeBatch([["INSERT INTO t VALUES (?)", [13]]])); + + await db.executeBatch([["INSERT INTO t VALUES (?)", [1]]]); + + const res = await db.execute("SELECT COUNT(*) AS count FROM t"); + expect(res.rows[0]!.count).toEqual(1); + }); + + it("a failed statement still rolls back the batch", async () => { + const error = await captureError(() => + db.executeBatch([ + ["INSERT INTO t VALUES (?)", [1]], + ["INSERT INTO tableThatDoesNotExist VALUES (?)", [2]], + ]), + ); + + expect(error.message).toContain("no such table"); + + const res = await db.execute("SELECT COUNT(*) AS count FROM t"); + expect(res.rows[0]!.count).toEqual(0); + }); +}); diff --git a/src/functions.ts b/src/functions.ts index ab6dd440..fdf1185e 100644 --- a/src/functions.ts +++ b/src/functions.ts @@ -12,6 +12,7 @@ import type { SQLBatchTuple, Transaction, } from "./types"; +import { rollbackAndRethrow } from "./rollback"; declare global { var __OPSQLiteProxy: object | undefined; @@ -67,6 +68,9 @@ function enhanceDB(db: _InternalDB, options: DBParams): DB { } }; + // Not available on libsql, which has no way to read the autocommit state + const inTransaction = db.inTransaction; + // spreading the object does not work with HostObjects (db) // We need to manually assign the fields const enhancedDb = { @@ -91,54 +95,15 @@ function enhanceDB(db: _InternalDB, options: DBParams): DB { }, flushPendingReactiveQueries: db.flushPendingReactiveQueries, executeBatch: async (commands: SQLBatchTuple[]): Promise => { + // BEGIN/COMMIT/ROLLBACK run natively around the batch. The lock is still + // needed so the batch never starts inside an open transaction() async function run() { try { - await enhancedDb.execute("BEGIN TRANSACTION;"); - const res = await db.executeBatch(commands as any[]); - await enhancedDb.execute("COMMIT;"); - await db.flushPendingReactiveQueries(); return res; - } catch (executionError) { - await enhancedDb.execute("ROLLBACK;"); - - throw executionError; - } finally { - lock.inProgress = false; - startNextTransaction(); - } - } - - return await new Promise((resolve, reject) => { - const tx: _PendingTransaction = { - start: () => { - run().then(resolve).catch(reject); - }, - }; - - lock.queue.push(tx); - startNextTransaction(); - }); - }, - executeBatchSync: async (commands: SQLBatchTuple[]): Promise => { - async function run() { - try { - enhancedDb.executeSync("BEGIN TRANSACTION;"); - - const res = await db.executeBatch(commands as any[]); - - enhancedDb.executeSync("COMMIT;"); - - await db.flushPendingReactiveQueries(); - - return res; - } catch (executionError) { - enhancedDb.executeSync("ROLLBACK;"); - - throw executionError; } finally { lock.inProgress = false; startNextTransaction(); @@ -257,6 +222,12 @@ function enhanceDB(db: _InternalDB, options: DBParams): DB { }. Cannot execute query on finalized transaction`, ); } + // SQLite may already have rolled back on its own (e.g. RAISE(ROLLBACK) + // caught inside fn), in which case ROLLBACK would throw + if (inTransaction?.() === false) { + isFinalized = true; + return { rowsAffected: 0, rows: [] }; + } const result = enhancedDb.executeSync("ROLLBACK;"); isFinalized = true; return result; @@ -277,7 +248,12 @@ function enhanceDB(db: _InternalDB, options: DBParams): DB { } } catch (executionError) { if (!isFinalized) { - rollback(); + isFinalized = true; + await rollbackAndRethrow( + executionError, + () => enhancedDb.executeSync("ROLLBACK;"), + inTransaction, + ); } throw executionError; diff --git a/src/functions.web.ts b/src/functions.web.ts index d464840b..ebaadf14 100644 --- a/src/functions.web.ts +++ b/src/functions.web.ts @@ -14,6 +14,7 @@ import type { SQLBatchTuple, Transaction, } from "./types"; +import { rollbackAndRethrow } from "./rollback"; type WorkerPromiser = (type: string, args?: Record) => Promise; @@ -257,7 +258,10 @@ function enhanceWebDb(db: _InternalDB, options: { name?: string; location?: stri } catch (error) { if (!finalized) { - await rollback(); + finalized = true; + // Not the user-facing rollback() above, which only reports that + // the sync API is unsupported on web + await rollbackAndRethrow(error, () => db.execute("ROLLBACK;")); } throw error; @@ -291,8 +295,7 @@ function enhanceWebDb(db: _InternalDB, options: { name?: string; location?: stri await db.execute("COMMIT;"); } catch (error) { - await db.execute("ROLLBACK;"); - throw error; + await rollbackAndRethrow(error, () => db.execute("ROLLBACK;")); } }); @@ -300,10 +303,6 @@ function enhanceWebDb(db: _InternalDB, options: { name?: string; location?: stri rowsAffected: 0, }; }, - // Web has no synchronous native APIs, so there is no distinct blocking behavior to offer. - executeBatchSync: async (commands: SQLBatchTuple[]): Promise => { - return enhancedDb.executeBatch(commands); - }, loadFile: async (_location: string): Promise => { throw new Error("[op-sqlite] loadFile() is not supported on web."); }, diff --git a/src/rollback.ts b/src/rollback.ts new file mode 100644 index 00000000..3a006879 --- /dev/null +++ b/src/rollback.ts @@ -0,0 +1,62 @@ +/** + * Rolls back the open transaction after `error` and rethrows `error`. + * + * SQLite sometimes rolls a transaction back on its own before we get to it: + * RAISE(ROLLBACK) in a trigger, or a COMMIT failing with SQLITE_FULL, + * SQLITE_IOERR* or SQLITE_CORRUPT_VTAB. A ROLLBACK issued after that fails with + * "cannot rollback - no transaction is active", and that failure must never + * replace the error that caused it. + * + * `inTransaction` is undefined on backends that cannot report the transaction + * state (libsql, web). There the ROLLBACK is always attempted and a failure is + * attached to the original error as `rollbackError`. + */ +export async function rollbackAndRethrow( + error: unknown, + rollback: () => unknown, + inTransaction?: () => boolean, +): Promise { + if (inTransaction?.() === false) { + throw error; + } + + try { + await rollback(); + } catch (rollbackError) { + throw withRollbackError(error, rollbackError, inTransaction); + } + + throw error; +} + +function withRollbackError( + error: unknown, + rollbackError: unknown, + inTransaction?: () => boolean, +): unknown { + // The connection is stuck inside the failed transaction: every following + // statement would silently run in it and be lost. That matters more than the + // original error, which is kept as `cause`. + if (inTransaction?.() === true) { + const stuck = new Error( + `[op-sqlite] ROLLBACK failed and the connection is still inside a transaction: ${messageOf(rollbackError)}`, + ) as Error & { cause?: unknown; rollbackError?: unknown }; + stuck.cause = error; + stuck.rollbackError = rollbackError; + return stuck; + } + + if (error !== null && typeof error === "object") { + try { + (error as { rollbackError?: unknown }).rollbackError = rollbackError; + } catch { + // Frozen or otherwise non-extensible error, rethrow it untouched + } + } + + return error; +} + +function messageOf(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} diff --git a/src/types.ts b/src/types.ts index 84a50987..93e25091 100644 --- a/src/types.ts +++ b/src/types.ts @@ -208,6 +208,11 @@ export type _InternalDB = { setReservedBytes: (reservedBytes: number) => void; getReservedBytes: () => number; flushPendingReactiveQueries: () => Promise; + /** + * Whether the connection has an open transaction (sqlite3_get_autocommit == 0). + * Undefined on backends that cannot report it (libsql, web). + */ + inTransaction?: () => boolean; }; export type DB = { @@ -287,23 +292,12 @@ export type DB = { * * It's faster than executing single queries as data is sent to the native side only once * - * The BEGIN/COMMIT/ROLLBACK statements that wrap the batch are executed asynchronously, - * off the JS thread. Use this over `executeBatchSync` unless you specifically need the - * transaction boundaries to block the JS thread. + * The BEGIN/COMMIT/ROLLBACK statements that wrap the batch run natively, off the JS + * thread, in the same call as the batch itself. * @param commands * @returns Promise */ executeBatch: (commands: SQLBatchTuple[]) => Promise; - /** - * Same as `executeBatch` but the BEGIN/COMMIT/ROLLBACK statements that wrap the batch - * are executed synchronously on the JS thread. For large batches this can block the JS - * thread for a noticeable amount of time (the COMMIT is where SQLite writes the WAL - * frames/fsyncs), so prefer `executeBatch` unless you have a specific reason to need - * synchronous transaction boundaries. - * @param commands - * @returns Promise - */ - executeBatchSync: (commands: SQLBatchTuple[]) => Promise; /** * Loads a SQLite Dump from disk. It will be the fastest way to execute a large set of queries as no JS is involved */