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
36 changes: 21 additions & 15 deletions cpp/OPBridge.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<bool> { 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<int>(commandCount),
};
return BatchResult{
.affectedRows = affectedRows,
.commands = static_cast<int>(commandCount),
};
});
}

} // namespace opsqlite
5 changes: 5 additions & 0 deletions cpp/OPBridge.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
9 changes: 9 additions & 0 deletions cpp/OPDatabase.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
42 changes: 15 additions & 27 deletions cpp/OPUtils.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<bool> { 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

Expand Down
43 changes: 43 additions & 0 deletions cpp/OPUtils.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@
#include <sqlite3.h>
#endif
#include <ReactCommon/CallInvoker.h>
#include <optional>
#include <stdexcept>
#include <string>
#include <vector>
#include "OPThreadPool.hpp"
Expand Down Expand Up @@ -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 <typename Execute, typename InTransaction, typename Work>
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<bool> 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<bool> 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);
Expand Down
79 changes: 53 additions & 26 deletions cpp/libsql/OPLibsqlBridge.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
#include "OPUtils.hpp"
#include <filesystem>
#include <iostream>
#include <optional>
#include <unordered_map>
#include <variant>

Expand Down Expand Up @@ -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<std::string> 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<JSVariant> *values) {
const char *err;
Expand Down Expand Up @@ -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);

Expand All @@ -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<int>(changes),
.insertId = static_cast<double>(insert_row_id)};
}
Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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);

Expand All @@ -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<int>(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<bool> { 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<int>(commandCount),
};
});
}

} // namespace opsqlite
28 changes: 19 additions & 9 deletions cpp/turso/OPTursoBridge.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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<bool> { 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<int>(command_count)};
return BatchResult{.affectedRows = affected_rows,
.commands = static_cast<int>(command_count)};
});
}

} // namespace opsqlite
8 changes: 1 addition & 7 deletions docs/docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
Loading
Loading