removed namespace interface

This commit is contained in:
Sascha Kühl
2025-08-11 16:59:16 +02:00
parent 2e2fcd01b6
commit 789be1174b
12 changed files with 247 additions and 101 deletions
+5 -2
View File
@@ -14,7 +14,7 @@ std::type_index message::type() const {
return type_;
}
const void * message::msg_ptr() const {
const void * message::raw_ptr() const {
return owner_ ? owner_.get() : ptr_;
}
@@ -49,6 +49,9 @@ bool subscription::valid() const {
return bus_ != nullptr;
}
message_bus::message_bus()
: next_id_{1} {}
void message_bus::unsubscribe(const std::type_index type, HandlerId id) {
std::unique_lock writeLock(mutex_);
const auto it = handlers_.find(type);
@@ -89,7 +92,7 @@ void message_bus::publish(const message &msg) const {
}
snapshot = it->second;
}
const void* p = msg.msg_ptr();
const void* p = msg.raw_ptr();
for (const auto &e : snapshot) {
if (!e.filter || e.filter(p)) e.handler(p);
}
+2 -2
View File
@@ -19,10 +19,10 @@ public:
explicit connection_statement_proxy(std::unique_ptr<statement_impl>&& stmt)
: statement_proxy(std::move(stmt)) {}
utils::result<size_t, utils::error> execute(interface::parameter_binder& bindings) override {
utils::result<size_t, utils::error> execute(parameter_binder& bindings) override {
return statement_->execute(bindings);
}
utils::result<std::unique_ptr<query_result_impl>, utils::error> fetch(interface::parameter_binder& bindings) override {
utils::result<std::unique_ptr<query_result_impl>, utils::error> fetch(parameter_binder& bindings) override {
return statement_->fetch(bindings);
}
};
+2 -2
View File
@@ -6,11 +6,11 @@ statement_impl::statement_impl(query_context query)
: query_(std::move(query))
{}
void statement_impl::bind(const size_t pos, const char *value, const size_t size, interface::parameter_binder& bindings) const {
void statement_impl::bind(const size_t pos, const char *value, const size_t size, parameter_binder& bindings) const {
utils::data_type_traits<const char*>::bind_value(bindings, adjust_index(pos), value, size);
}
void statement_impl::bind(const size_t pos, std::string &val, const size_t size, interface::parameter_binder& bindings) const {
void statement_impl::bind(const size_t pos, std::string &val, const size_t size, parameter_binder& bindings) const {
utils::data_type_traits<std::string>::bind_value(bindings, adjust_index(pos), val, size);
}
+2 -2
View File
@@ -4,10 +4,10 @@ namespace matador::sql {
statement_proxy::statement_proxy(std::unique_ptr<statement_impl>&& stmt)
: statement_(std::move(stmt)){}
void statement_proxy::bind(const size_t pos, const char* value, const size_t size, interface::parameter_binder& bindings) const {
void statement_proxy::bind(const size_t pos, const char* value, const size_t size, parameter_binder& bindings) const {
statement_->bind(pos, value, size, bindings);
}
void statement_proxy::bind(const size_t pos, std::string& val, const size_t size, interface::parameter_binder& bindings) const {
void statement_proxy::bind(const size_t pos, std::string& val, const size_t size, parameter_binder& bindings) const {
statement_->bind(pos, val, size, bindings);
}
+12 -8
View File
@@ -24,24 +24,26 @@ public:
size_t lock_attempts{0};
};
explicit statement_cache_proxy(std::unique_ptr<statement_impl>&& stmt, connection_pool &pool, const size_t connection_id)
statement_cache_proxy(utils::message_bus &bus, std::unique_ptr<statement_impl>&& stmt, connection_pool &pool, const size_t connection_id)
: statement_proxy(std::move(stmt))
, pool_(pool)
, connection_id_(connection_id) {}
, connection_id_(connection_id)
, bus_(bus) {}
utils::result<size_t, utils::error> execute(interface::parameter_binder& bindings) override {
utils::result<size_t, utils::error> execute(parameter_binder& bindings) override {
execution_metrics metrics{std::chrono::steady_clock::now()};
auto result = try_with_retry([this, &bindings, &metrics]() -> utils::result<size_t, utils::error> {
if (!try_lock()) {
++metrics.lock_attempts;
bus_.publish<statement_lock_failed_event>({sql(), std::chrono::steady_clock::now(), metrics.lock_attempt_start});
return utils::failure(utils::error{
error_code::STATEMENT_LOCKED,
"Failed to execute statement because it is already in use"
});
}
metrics.lock_acquired = std::chrono::steady_clock::now();
// publish
bus_.publish<statement_lock_acquired_event>({sql(), std::chrono::steady_clock::now(), metrics.lock_attempt_start, metrics.lock_acquired});
auto guard = statement_guard(*this);
if (const auto conn = pool_.acquire(connection_id_); !conn.valid()) {
@@ -55,14 +57,15 @@ public:
auto execution_result = statement_->execute(bindings);
metrics.execution_end = std::chrono::steady_clock::now();
// publish
bus_.publish<statement_execution_event>({sql(), std::chrono::steady_clock::now(), metrics.execution_start, metrics.execution_end});
return execution_result;
});
return result;
}
utils::result<std::unique_ptr<query_result_impl>, utils::error> fetch(interface::parameter_binder& bindings) override {
utils::result<std::unique_ptr<query_result_impl>, utils::error> fetch(parameter_binder& bindings) override {
execution_metrics metrics{std::chrono::steady_clock::now()};
auto result = try_with_retry([this, &bindings, &metrics]() -> utils::result<std::unique_ptr<query_result_impl>, utils::error> {
@@ -104,7 +107,7 @@ protected:
if (attempt + 1 < config_.max_attempts) {
std::this_thread::sleep_for(current_wait);
current_wait = std::min(current_wait * 2, config_.max_wait);
current_wait = (std::min)(current_wait * 2, config_.max_wait);
}
}
@@ -132,6 +135,7 @@ private:
connection_pool &pool_;
size_t connection_id_{};
retry_config config_{};
utils::message_bus &bus_;
};
}
@@ -178,7 +182,7 @@ utils::result<statement, utils::error> statement_cache::acquire(const query_cont
const auto it = cache_map_.insert({
key,
{statement{
std::make_shared<internal::statement_cache_proxy>(std::move(stmt), pool_, id)},
std::make_shared<internal::statement_cache_proxy>(bus_, std::move(stmt), pool_, id)},
std::chrono::steady_clock::now(),
usage_list_.begin()}
}).first;