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
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@

## Unreleased

- fixed: statement cleanup no longer waits on a connection that is inside the busy handler. `release/2` defers the same way, and a deferred statement is finalized when that connection unlocks rather than waiting for a later call. See [#360](https://github.com/elixir-sqlite/exqlite/issues/360).

## v0.41.0

- added: `:authorizer` connection option to deny SQL actions on `DBConnection` managed connections via `Exqlite.Sqlite3.set_authorizer/2`.
Expand Down
244 changes: 222 additions & 22 deletions c_src/sqlite3_nif.c
Original file line number Diff line number Diff line change
Expand Up @@ -56,11 +56,30 @@ ErlNifMutex* log_hook_mutex = NULL;
// currently defined SQLite action code is SQLITE_RECURSIVE (33).
#define AUTHORIZER_DENY_SIZE 64

// Allocated at prepare so the statement destructor never allocates.
// statement is set when the Erlang resource is already gone.
// owner is set when release/2 deferred because the connection mutex was held;
// the sqlite3_stmt stays in owner until drain, so an in-flight step can finish.
typedef struct statement statement_t;

typedef struct deferred_finalize
{
sqlite3_stmt* statement;
statement_t* owner;
struct deferred_finalize* next;
} deferred_finalize_t;

typedef struct connection
{
sqlite3* db;
ErlNifMutex* mutex;
ErlNifMutex* interrupt_mutex;

// Guards finalize_head only.
// NOTE: Never hold this across a SQLite call.
ErlNifMutex* finalize_mutex;
deferred_finalize_t* finalize_head;

ErlNifPid update_hook_pid;
int authorizer_deny[AUTHORIZER_DENY_SIZE];

Expand All @@ -72,11 +91,14 @@ typedef struct connection
ErlNifPid caller_pid;
} connection_t;

typedef struct statement
struct statement
{
connection_t* conn;
sqlite3_stmt* statement;
} statement_t;
deferred_finalize_t* slot;
};

static void connection_drain_deferred(connection_t* conn);

static int exqlite_progress_handler(void* arg);

Expand Down Expand Up @@ -282,13 +304,43 @@ connection_acquire_lock(connection_t* conn)
{
assert(conn);
enif_mutex_lock(conn->mutex);

// Dropped statements may have been queued while this connection was
// inside a SQLite call.
//
// NOTE: Finalize them before the next SQLite call.
connection_drain_deferred(conn);
}

static inline void
connection_release_lock(connection_t* conn)
{
assert(conn);
enif_mutex_unlock(conn->mutex);

// Statements may have been queued while this thread was inside SQLite.
// Finalize them before unlocking. Otherwise an idle connection keeps a
// WAL read mark until some later call takes the mutex.
//
// When the queue is empty, unlock conn->mutex before finalize_mutex.
// A destructor can enqueue only while holding finalize_mutex, then
// trylocks conn->mutex. It either observes this thread still holding
// the connection mutex, or it obtains the mutex and drains itself.
if (!conn->finalize_mutex) {
enif_mutex_unlock(conn->mutex);
return;
}

for (;;) {
connection_drain_deferred(conn);

enif_mutex_lock(conn->finalize_mutex);
if (conn->finalize_head == NULL) {
enif_mutex_unlock(conn->mutex);
enif_mutex_unlock(conn->finalize_mutex);
return;
}
enif_mutex_unlock(conn->finalize_mutex);
}
}

static inline void
Expand All @@ -305,6 +357,53 @@ statement_release_lock(statement_t* statement)
connection_release_lock(statement->conn);
}

// Caller must hold conn->mutex, unless the connection is not reachable
// from another thread. sqlite3_finalize runs after finalize_mutex is released
// so a statement destructor can enqueue without waiting for finalize.
static void
connection_drain_deferred(connection_t* conn)
{
deferred_finalize_t* node;
deferred_finalize_t* next;

if (!conn->finalize_mutex) {
return;
}

enif_mutex_lock(conn->finalize_mutex);
node = conn->finalize_head;
conn->finalize_head = NULL;
enif_mutex_unlock(conn->finalize_mutex);

while (node) {
next = node->next;
if (node->owner) {
statement_t* owner = node->owner;
node->owner = NULL;
if (owner->statement) {
sqlite3_finalize(owner->statement);
owner->statement = NULL;
}
// Drops the keep taken when release deferred. The NIF argument
// ref, if any, keeps the destructor from running on this thread.
enif_release_resource(owner);
} else if (node->statement) {
sqlite3_finalize(node->statement);
}
enif_free(node);
node = next;
}
}

static void
connection_defer_finalize(connection_t* conn, deferred_finalize_t* slot)
{
enif_mutex_lock(conn->finalize_mutex);
slot->next = conn->finalize_head;
conn->finalize_head = slot;
enif_mutex_unlock(conn->finalize_mutex);
}

static inline void
connection_configure_progress_handler(connection_t* conn)
{
Expand Down Expand Up @@ -484,6 +583,8 @@ exqlite_open(ErlNifEnv* env, int argc, const ERL_NIF_TERM argv[])
conn->db = db;
conn->mutex = mutex;
conn->interrupt_mutex = NULL;
conn->finalize_mutex = NULL;
conn->finalize_head = NULL;
memset(conn->authorizer_deny, 0, sizeof(conn->authorizer_deny));

// Initialize busy handler fields
Expand All @@ -492,6 +593,12 @@ exqlite_open(ErlNifEnv* env, int argc, const ERL_NIF_TERM argv[])
conn->progress_handler_steps = 1000;
conn->callback_env = NULL;

conn->finalize_mutex = enif_mutex_create("exqlite:finalize");
if (conn->finalize_mutex == NULL) {
enif_release_resource(conn);
return make_error_tuple(env, am_failed_to_create_mutex);
}

conn->interrupt_mutex = enif_mutex_create("exqlite:interrupt");
if (conn->interrupt_mutex == NULL) {
enif_release_resource(conn);
Expand Down Expand Up @@ -544,9 +651,10 @@ exqlite_close(ErlNifEnv* env, int argc, const ERL_NIF_TERM argv[])
if (autocommit == 0) {
rc = sqlite3_exec(conn->db, "ROLLBACK;", NULL, NULL, NULL);
if (rc != SQLITE_OK) {
ERL_NIF_TERM error = make_sqlite3_error_tuple(env, rc, conn->db);
connection_clear_caller(conn);
connection_release_lock(conn);
return make_sqlite3_error_tuple(env, rc, conn->db);
return error;
}
}

Expand All @@ -563,9 +671,10 @@ exqlite_close(ErlNifEnv* env, int argc, const ERL_NIF_TERM argv[])
enif_mutex_lock(conn->interrupt_mutex);
rc = sqlite3_close_v2(conn->db);
if (rc != SQLITE_OK) {
ERL_NIF_TERM error = make_sqlite3_error_tuple(env, rc, conn->db);
enif_mutex_unlock(conn->interrupt_mutex);
connection_release_lock(conn);
return make_sqlite3_error_tuple(env, rc, conn->db);
return error;
}
conn->db = NULL;
enif_mutex_unlock(conn->interrupt_mutex);
Expand Down Expand Up @@ -612,9 +721,10 @@ exqlite_execute(ErlNifEnv* env, int argc, const ERL_NIF_TERM argv[])

rc = sqlite3_exec(conn->db, (char*)bin.data, NULL, NULL, NULL);
if (rc != SQLITE_OK) {
ERL_NIF_TERM error = make_sqlite3_error_tuple(env, rc, conn->db);
connection_clear_caller(conn);
connection_release_lock(conn);
return make_sqlite3_error_tuple(env, rc, conn->db);
return error;
}

connection_clear_caller(conn);
Expand Down Expand Up @@ -683,6 +793,15 @@ exqlite_prepare(ErlNifEnv* env, int argc, const ERL_NIF_TERM argv[])
return make_error_tuple(env, am_out_of_memory);
}
statement->statement = NULL;
statement->conn = NULL;
statement->slot = enif_alloc(sizeof(deferred_finalize_t));
if (!statement->slot) {
enif_release_resource(statement);
return make_error_tuple(env, am_out_of_memory);
}
statement->slot->statement = NULL;
statement->slot->owner = NULL;
statement->slot->next = NULL;

enif_keep_resource(conn);
statement->conn = conn;
Expand Down Expand Up @@ -1011,11 +1130,14 @@ exqlite_multi_step(ErlNifEnv* env, int argc, const ERL_NIF_TERM argv[])
rows = enif_make_list_cell(env, row, rows);
break;

default:
default: {
ERL_NIF_TERM error;
sqlite3_reset(statement->statement);
error = make_sqlite3_error_tuple(env, rc, conn->db);
connection_clear_caller(conn);
connection_release_lock(conn);
return make_sqlite3_error_tuple(env, rc, conn->db);
return error;
}
}
}

Expand Down Expand Up @@ -1317,9 +1439,10 @@ exqlite_deserialize(ErlNifEnv* env, int argc, const ERL_NIF_TERM argv[])
memcpy(buffer, serialized.data, size);
rc = sqlite3_deserialize(conn->db, (const char*)database_name.data, buffer, size, size, flags);
if (rc != SQLITE_OK) {
ERL_NIF_TERM error = make_sqlite3_error_tuple(env, rc, conn->db);
sqlite3_free(buffer);
connection_release_lock(conn);
return make_sqlite3_error_tuple(env, rc, conn->db);
return error;
}

connection_release_lock(conn);
Expand Down Expand Up @@ -1350,14 +1473,59 @@ exqlite_release(ErlNifEnv* env, int argc, const ERL_NIF_TERM argv[])
return make_error_tuple(env, am_invalid_statement);
}

statement_acquire_lock(statement);
// Lock the statement's connection, not argv[0]. A caller can pass another
// connection while this one is inside the busy handler. Blocking here stalls
// the dirty scheduler that holds the write lock the handler is waiting for.
connection_t* owner = statement->conn;
if (!owner || !owner->mutex) {
return am_ok;
}

if (statement->statement) {
sqlite3_finalize(statement->statement);
statement->statement = NULL;
if (enif_mutex_trylock(owner->mutex) == 0) {
connection_drain_deferred(owner);
if (statement->statement) {
sqlite3_finalize(statement->statement);
statement->statement = NULL;
}
connection_release_lock(owner);
return am_ok;
}

statement_release_lock(statement);
// Contended. Blocking here would stall the dirty scheduler that holds
// the write lock. Without a finalize queue, waiting is the only option.
if (!owner->finalize_mutex) {
statement_acquire_lock(statement);
if (statement->statement) {
sqlite3_finalize(statement->statement);
statement->statement = NULL;
}
statement_release_lock(statement);
return am_ok;
}

// Do not read statement->statement. Another release writes it under
// conn->mutex, and the holder may be stepping it. slot is stolen only
// under finalize_mutex. Drain finalizes after those calls return.
enif_keep_resource(statement);
enif_mutex_lock(owner->finalize_mutex);
if (!statement->slot) {
enif_mutex_unlock(owner->finalize_mutex);
enif_release_resource(statement);
return am_ok;
}

deferred_finalize_t* slot = statement->slot;
statement->slot = NULL;
slot->statement = NULL;
slot->owner = statement;
slot->next = owner->finalize_head;
owner->finalize_head = slot;
enif_mutex_unlock(owner->finalize_mutex);

// The holder may have unlocked after our first trylock. Drain if so.
if (enif_mutex_trylock(owner->mutex) == 0) {
connection_release_lock(owner);
}

return am_ok;
}
Expand Down Expand Up @@ -1399,6 +1567,11 @@ connection_type_destructor(ErlNifEnv* env, void* arg)
conn->mutex = NULL;
}

if (conn->finalize_mutex) {
enif_mutex_destroy(conn->finalize_mutex);
conn->finalize_mutex = NULL;
}

if (conn->interrupt_mutex) {
enif_mutex_destroy(conn->interrupt_mutex);
conn->interrupt_mutex = NULL;
Expand All @@ -1411,17 +1584,44 @@ statement_type_destructor(ErlNifEnv* env, void* arg)
assert(env);
assert(arg);

statement_t* statement = (statement_t*)arg;
statement_acquire_lock(statement);
// Runs on whichever normal scheduler drops the last reference, not as a
// dirty NIF. A connection inside the busy handler holds conn->mutex and
// SQLite's db mutex until busy_timeout expires. Taking either mutex here
// stalls this scheduler and the writer the handler is waiting for.
statement_t* statement = (statement_t*)arg;
connection_t* conn = statement->conn;
deferred_finalize_t* slot = statement->slot;
sqlite3_stmt* sqlite_stmt = statement->statement;

statement->conn = NULL;
statement->slot = NULL;
statement->statement = NULL;

if (statement->statement) {
sqlite3_finalize(statement->statement);
statement->statement = NULL;
if (!sqlite_stmt) {
if (slot) {
enif_free(slot);
}
} else if (conn && conn->mutex && enif_mutex_trylock(conn->mutex) == 0) {
sqlite3_finalize(sqlite_stmt);
if (slot) {
enif_free(slot);
}
connection_release_lock(conn);
} else if (conn && conn->finalize_mutex && slot) {
slot->statement = sqlite_stmt;
slot->owner = NULL;
connection_defer_finalize(conn, slot);
// If the holder unlocked after our trylock failed, finish the queue now.
if (conn->mutex && enif_mutex_trylock(conn->mutex) == 0) {
connection_release_lock(conn);
}
} else if (slot) {
enif_free(slot);
}

statement_release_lock(statement);
enif_release_resource(statement->conn);
statement->conn = NULL;
if (conn) {
enif_release_resource(conn);
}
}

int
Expand Down
2 changes: 1 addition & 1 deletion mix.exs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
defmodule Exqlite.MixProject do
use Mix.Project

@version "0.41.0"
@version "0.42.0-dev"

def project do
[
Expand Down
Loading
Loading