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
73 changes: 70 additions & 3 deletions orm_lib/src/DbClientImpl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -144,10 +144,15 @@ void DbClientImpl::execSql(
}
DbConnectionPtr conn;
bool busy = false;
bool unavailable = false;
{
std::lock_guard<std::mutex> guard(connectionsMutex_);

if (readyConnections_.size() == 0)
if (failedSqliteConnections_ == numberOfConnections_)
{
unavailable = true;
}
else if (readyConnections_.size() == 0)
{
if (sqlCmdBuffer_.size() > 200000)
{
Expand Down Expand Up @@ -187,6 +192,12 @@ void DbClientImpl::execSql(
std::move(exceptCallback));
return;
}
if (unavailable)
{
exceptCallback(std::make_exception_ptr(
BrokenConnection("No usable SQLite connections remain")));
return;
}
if (busy)
{
auto exceptPtr =
Expand All @@ -201,9 +212,14 @@ void DbClientImpl::newTransactionAsync(
TransactionType transType)
{
DbConnectionPtr conn;
bool unavailable = false;
{
std::lock_guard<std::mutex> lock(connectionsMutex_);
if (!readyConnections_.empty())
if (failedSqliteConnections_ == numberOfConnections_)
{
unavailable = true;
}
else if (!readyConnections_.empty())
{
auto iter = readyConnections_.begin();
busyConnections_.insert(*iter);
Expand Down Expand Up @@ -255,6 +271,11 @@ void DbClientImpl::newTransactionAsync(
transCallbacks_.push_back({callbackPtr, transType});
}
}
if (unavailable)
{
callback(nullptr);
return;
}
if (conn)
{
makeTrans(conn,
Expand Down Expand Up @@ -337,6 +358,9 @@ std::shared_ptr<Transaction> DbClientImpl::newTransaction(
auto trans = f.get();
if (!trans)
{
std::lock_guard<std::mutex> lock(connectionsMutex_);
if (failedSqliteConnections_ == numberOfConnections_)
throw BrokenConnection("No usable SQLite connections remain");
throw TimeoutError("Timeout, no connection available for transaction");
}
trans->setCommitCallback(commitCallback);
Expand Down Expand Up @@ -432,6 +456,38 @@ DbConnectionPtr DbClientImpl::newConnection(trantor::EventLoop *loop)
auto thisPtr = weakPtr.lock();
if (!thisPtr)
return;
if (thisPtr->type_ == ClientType::Sqlite3)
{
decltype(sqlCmdBuffer_) commands;
decltype(transCallbacks_) transactions;
{
std::lock_guard<std::mutex> guard(thisPtr->connectionsMutex_);
if (thisPtr->connections_.find(closeConnPtr) ==
thisPtr->connections_.end())
return;
thisPtr->readyConnections_.erase(closeConnPtr);
thisPtr->busyConnections_.erase(closeConnPtr);
++thisPtr->failedSqliteConnections_;
if (thisPtr->failedSqliteConnections_ ==
thisPtr->numberOfConnections_)
{
commands.swap(thisPtr->sqlCmdBuffer_);
transactions.swap(thisPtr->transCallbacks_);
}
}
// A replacement would lose connection-local settings and possibly
// an in-memory database. Keep the failed connection out of the pool
// until the client is recreated, but retain ownership of its
// thread.
closeConnPtr->disconnect();
const auto error = std::make_exception_ptr(
BrokenConnection("No usable SQLite connections remain"));
for (const auto &command : commands)
command->exceptionCallback_(error);
for (const auto &transaction : transactions)
(*transaction.first)(nullptr);
return;
}
{
std::lock_guard<std::mutex> guard(thisPtr->connectionsMutex_);
thisPtr->readyConnections_.erase(closeConnPtr);
Expand Down Expand Up @@ -509,6 +565,7 @@ void DbClientImpl::execSqlWithTimeout(
assert(timeout_ > 0.0);
auto cmd = std::make_shared<std::weak_ptr<SqlCmd>>();
bool busy = false;
bool unavailable = false;
auto ecpPtr =
std::make_shared<std::function<void(const std::exception_ptr &)>>(
std::move(ecb));
Expand Down Expand Up @@ -551,7 +608,11 @@ void DbClientImpl::execSqlWithTimeout(
{
std::lock_guard<std::mutex> guard(connectionsMutex_);

if (readyConnections_.size() == 0)
if (failedSqliteConnections_ == numberOfConnections_)
{
unavailable = true;
}
else if (readyConnections_.size() == 0)
{
if (sqlCmdBuffer_.size() > 200000)
{
Expand Down Expand Up @@ -594,6 +655,12 @@ void DbClientImpl::execSqlWithTimeout(
return;
}

if (unavailable)
{
exceptionCallback(std::make_exception_ptr(
BrokenConnection("No usable SQLite connections remain")));
return;
}
if (busy)
{
exceptionCallback(
Expand Down
3 changes: 3 additions & 0 deletions orm_lib/src/DbClientImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,9 @@ class DbClientImpl : public DbClient,
std::unordered_set<DbConnectionPtr> connections_;
std::unordered_set<DbConnectionPtr> readyConnections_;
std::unordered_set<DbConnectionPtr> busyConnections_;
// Retired SQLite connections stay owned until closeAll(), so their threads
// are not destroyed from inside a connection callback.
size_t failedSqliteConnections_{0};

using TransCallbackEntry =
std::pair<std::shared_ptr<std::function<void(
Expand Down
65 changes: 63 additions & 2 deletions orm_lib/src/TransactionImpl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -14,12 +14,65 @@

#include "TransactionImpl.h"
#include "../../lib/src/TaskTimeoutFlag.h"
#if USE_SQLITE3
#include "sqlite3_impl/Sqlite3Connection.h"
#endif
#include <string_view>
#include <trantor/utils/Logger.h>

using namespace drogon::orm;
using namespace drogon;

namespace
{
#if USE_SQLITE3
void rollbackFailedSqliteCommit(const DbConnectionPtr &conn,
const std::function<void()> &usedUpCallback,
const std::function<void(bool)> &commitCallback)
{
// SQLite can leave the transaction open after COMMIT fails. Do not return
// this connection to the pool when the failed command becomes idle.
conn->setIdleCallback([]() {});
conn->loop()->queueInLoop([conn, usedUpCallback, commitCallback]() {
auto sqlite = std::static_pointer_cast<Sqlite3Connection>(conn);
auto finish = [sqlite, usedUpCallback, commitCallback](bool reusable) {
// Let the failed statement and its callbacks finish before closing
// the connection or allowing the pool to dispatch another query.
sqlite->loop()->queueInLoop(
[sqlite, usedUpCallback, commitCallback, reusable]() {
if (!reusable)
sqlite->invalidate();
if (commitCallback)
commitCallback(false);
if (reusable && usedUpCallback)
usedUpCallback();
});
};
// Some COMMIT errors already roll back the entire transaction.
if (!sqlite->hasActiveTransaction())
{
finish(true);
return;
}
conn->execSql(
"rollback",
0,
{},
{},
{},
[sqlite, finish](const Result &) {
LOG_TRACE << "Transaction rolled back after failed commit";
finish(!sqlite->hasActiveTransaction());
},
[sqlite, finish](const std::exception_ptr &) {
LOG_ERROR << "Transaction rollback after failed commit failed";
finish(!sqlite->hasActiveTransaction());
});
});
}
#endif
} // namespace

TransactionImpl::TransactionImpl(ClientType type,
const DbConnectionPtr &connPtr,
std::function<void(bool)> commitCallback,
Expand All @@ -42,9 +95,10 @@ TransactionImpl::~TransactionImpl()
{
auto loop = connectionPtr_->loop();
loop->queueInLoop([conn = connectionPtr_,
type = type_,
ucb = std::move(usedUpCallback_),
commitCb = std::move(commitCallback_)]() {
conn->setIdleCallback([ucb = std::move(ucb)]() {
conn->setIdleCallback([ucb]() {
if (ucb)
ucb();
});
Expand All @@ -61,7 +115,7 @@ TransactionImpl::~TransactionImpl()
commitCb(true);
}
},
[commitCb](const std::exception_ptr &ePtr) {
[conn, type, ucb, commitCb](const std::exception_ptr &ePtr) {
try
{
std::rethrow_exception(ePtr);
Expand All @@ -70,6 +124,13 @@ TransactionImpl::~TransactionImpl()
{
LOG_ERROR << "Transaction submission failed:"
<< e.base().what();
#if USE_SQLITE3
if (type == ClientType::Sqlite3)
{
rollbackFailedSqliteCommit(conn, ucb, commitCb);
return;
}
#endif
if (commitCb)
{
commitCb(false);
Expand Down
19 changes: 19 additions & 0 deletions orm_lib/src/sqlite3_impl/Sqlite3Connection.cc
Original file line number Diff line number Diff line change
Expand Up @@ -375,6 +375,25 @@ int Sqlite3Connection::stmtStep(
return r;
}

bool Sqlite3Connection::hasActiveTransaction() const
{
loop_->assertInLoopThread();
return connectionPtr_ && sqlite3_get_autocommit(connectionPtr_.get()) == 0;
}

void Sqlite3Connection::invalidate()
{
loop_->assertInLoopThread();
if (status_ != ConnectStatus::Ok)
return;
status_ = ConnectStatus::Bad;
// No statement is executing now. Finalize the cache before disconnecting
// so closing also releases the failed transaction and its locks.
stmtsMap_.clear();
stmts_.clear();
closeCallback_(shared_from_this());
}

void Sqlite3Connection::disconnect()
{
std::promise<int> pro;
Expand Down
4 changes: 4 additions & 0 deletions orm_lib/src/sqlite3_impl/Sqlite3Connection.h
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,10 @@ class Sqlite3Connection : public DbConnection,

void disconnect() override;

// Internal failed-COMMIT recovery. Call only on this connection's loop.
bool hasActiveTransaction() const;
void invalidate();

private:
static std::once_flag once_;
void execSqlInQueue(
Expand Down
12 changes: 12 additions & 0 deletions orm_lib/src/sqlite3_impl/test/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,18 @@ set_property(TARGET sqlite3_test1 PROPERTY CXX_STANDARD ${DROGON_CXX_STANDARD})
set_property(TARGET sqlite3_test1 PROPERTY CXX_STANDARD_REQUIRED ON)
set_property(TARGET sqlite3_test1 PROPERTY CXX_EXTENSIONS OFF)

add_executable(sqlite3_failed_commit_test failed_commit_test.cc)
if (TARGET SQLite3_lib)
target_link_libraries(sqlite3_failed_commit_test PRIVATE SQLite3_lib)
else ()
target_link_libraries(sqlite3_failed_commit_test PRIVATE unofficial::sqlite3::sqlite3)
endif ()
set_property(TARGET sqlite3_failed_commit_test PROPERTY CXX_STANDARD ${DROGON_CXX_STANDARD})
set_property(TARGET sqlite3_failed_commit_test PROPERTY CXX_STANDARD_REQUIRED ON)
set_property(TARGET sqlite3_failed_commit_test PROPERTY CXX_EXTENSIONS OFF)
add_test(NAME sqlite3_failed_commit COMMAND sqlite3_failed_commit_test)
set_tests_properties(sqlite3_failed_commit PROPERTIES TIMEOUT 30)

add_executable(sqlite3_disconnect_test Sqlite3DisconnectTest.cc)
set_property(TARGET sqlite3_disconnect_test
PROPERTY CXX_STANDARD ${DROGON_CXX_STANDARD})
Expand Down
Loading
Loading