integrated object cache into session and added message bus to object_cache and update tests
This commit is contained in:
@@ -5,6 +5,7 @@
|
||||
#include "matador/object/object_resolver.hpp"
|
||||
|
||||
#include "matador/utils/identifier.hpp"
|
||||
#include "matador/utils/message_bus.hpp"
|
||||
|
||||
#include <unordered_map>
|
||||
#include <memory>
|
||||
@@ -35,6 +36,23 @@ struct cache_entry : cache_entry_base {
|
||||
}
|
||||
};
|
||||
|
||||
struct object_cache_event {
|
||||
std::type_index type{typeid(void)};
|
||||
utils::identifier id{};
|
||||
std::chrono::steady_clock::time_point timestamp{};
|
||||
};
|
||||
|
||||
struct object_cache_accessed_event : object_cache_event {};
|
||||
struct object_cache_proxy_added_event : object_cache_event {};
|
||||
struct object_cache_entity_attached_event : object_cache_event {};
|
||||
struct object_cache_imported_event : object_cache_event {};
|
||||
struct object_cache_erased_event : object_cache_event {};
|
||||
|
||||
struct object_cache_sweep_event {
|
||||
std::size_t removed{};
|
||||
std::chrono::steady_clock::time_point timestamp{};
|
||||
};
|
||||
|
||||
/**
|
||||
* @brief Thread-sicherer Cache für Objekt-Proxies und (optional) geladene Entities.
|
||||
*
|
||||
@@ -54,7 +72,8 @@ public:
|
||||
/**
|
||||
* @brief Erzeugt einen leeren Cache.
|
||||
*/
|
||||
object_cache() = default;
|
||||
explicit object_cache(utils::message_bus &bus)
|
||||
: bus_(bus) {}
|
||||
|
||||
/**
|
||||
* @brief Liefert einen Proxy für (T, id) und erstellt ihn bei Bedarf.
|
||||
@@ -83,6 +102,7 @@ public:
|
||||
// found entry, return std::shared_ptr of proxy
|
||||
auto *entry = entry_cast_<Type>(it->second.get());
|
||||
if (auto proxy_ptr = entry->proxy.lock()) {
|
||||
bus_.publish<object_cache_accessed_event>({k.type, id, std::chrono::steady_clock::now()});
|
||||
return proxy_ptr;
|
||||
}
|
||||
|
||||
@@ -96,6 +116,8 @@ public:
|
||||
|
||||
entry->proxy = proxy_ptr;
|
||||
|
||||
bus_.publish<object_cache_proxy_added_event>({k.type, id, std::chrono::steady_clock::now()});
|
||||
|
||||
// return the shared_ptr of the proxy
|
||||
return proxy_ptr;
|
||||
}
|
||||
@@ -117,10 +139,65 @@ public:
|
||||
entry->proxy = proxy_ptr;
|
||||
map_.emplace(k, std::move(entry_ptr));
|
||||
|
||||
bus_.publish<object_cache_proxy_added_event>({k.type, id, std::chrono::steady_clock::now()});
|
||||
|
||||
// return the shared_ptr of the proxy
|
||||
return proxy_ptr;
|
||||
}
|
||||
|
||||
template<typename Type, typename ResolverPointerType>
|
||||
bool import(const utils::identifier &id, const std::shared_ptr<object_proxy<Type>> &proxy, ResolverPointerType &&resolver_ptr) {
|
||||
if (!proxy) {
|
||||
return false;
|
||||
}
|
||||
|
||||
auto obj = proxy->object();
|
||||
auto weak_resolver = to_weak<Type>(std::forward<ResolverPointerType>(resolver_ptr));
|
||||
const auto k = make_key<Type>(id);
|
||||
|
||||
std::unique_lock lock(mutex_);
|
||||
|
||||
auto it = map_.find(k);
|
||||
if (it != map_.end()) {
|
||||
auto *entry = entry_cast_<Type>(it->second.get());
|
||||
|
||||
if (auto existing_proxy = entry->proxy.lock()) {
|
||||
if (existing_proxy != proxy) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
proxy->primary_key(id);
|
||||
proxy->resolver(std::move(weak_resolver));
|
||||
|
||||
if (obj) {
|
||||
entry->entity = obj;
|
||||
proxy->attach(std::move(obj));
|
||||
}
|
||||
|
||||
entry->proxy = proxy;
|
||||
bus_.publish<object_cache_imported_event>({k.type, id, std::chrono::steady_clock::now()});
|
||||
return true;
|
||||
}
|
||||
|
||||
auto entry_ptr = std::make_unique<cache_entry<Type>>();
|
||||
auto *entry = entry_ptr.get();
|
||||
|
||||
proxy->primary_key(id);
|
||||
proxy->resolver(std::move(weak_resolver));
|
||||
|
||||
if (obj) {
|
||||
entry->entity = obj;
|
||||
proxy->attach(std::move(obj));
|
||||
}
|
||||
|
||||
entry->proxy = proxy;
|
||||
map_.emplace(k, std::move(entry_ptr));
|
||||
|
||||
bus_.publish<object_cache_imported_event>({k.type, id, std::chrono::steady_clock::now()});
|
||||
|
||||
return true;
|
||||
}
|
||||
/**
|
||||
* @brief Verknüpft eine geladene Entity mit einem Cache-Eintrag (Type, id).
|
||||
*
|
||||
@@ -145,6 +222,7 @@ public:
|
||||
auto *entry = entry_ptr.get();
|
||||
entry->entity = obj;
|
||||
map_.emplace(k, std::move(entry_ptr));
|
||||
bus_.publish<object_cache_entity_attached_event>({k.type, id, std::chrono::steady_clock::now()});
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -155,6 +233,8 @@ public:
|
||||
proxy->attach(std::move(obj));
|
||||
}
|
||||
|
||||
bus_.publish<object_cache_entity_attached_event>({k.type, id, std::chrono::steady_clock::now()});
|
||||
|
||||
prune_if_dead(it);
|
||||
}
|
||||
|
||||
@@ -179,6 +259,9 @@ public:
|
||||
|
||||
auto *entry = entry_cast_<Type>(it->second.get());
|
||||
auto entity_ptr = entry->entity.lock();
|
||||
if (entity_ptr) {
|
||||
bus_.publish<object_cache_accessed_event>({k.type, id, std::chrono::steady_clock::now()});
|
||||
}
|
||||
prune_if_dead(it);
|
||||
return entity_ptr;
|
||||
}
|
||||
@@ -204,6 +287,9 @@ public:
|
||||
|
||||
auto *entry = entry_cast_<Type>(it->second.get());
|
||||
const bool loaded = !entry->entity.expired();
|
||||
if (loaded) {
|
||||
bus_.publish<object_cache_accessed_event>({k.type, id, std::chrono::steady_clock::now()});
|
||||
}
|
||||
prune_if_dead(it);
|
||||
return loaded;
|
||||
}
|
||||
@@ -224,6 +310,7 @@ public:
|
||||
}
|
||||
|
||||
map_.erase(it);
|
||||
bus_.publish<object_cache_erased_event>({k.type, id, std::chrono::steady_clock::now()});
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -246,6 +333,10 @@ public:
|
||||
}
|
||||
}
|
||||
|
||||
if (removed > 0) {
|
||||
bus_.publish<object_cache_sweep_event>({removed, std::chrono::steady_clock::now()});
|
||||
}
|
||||
|
||||
return removed;
|
||||
}
|
||||
|
||||
@@ -325,6 +416,7 @@ private:
|
||||
|
||||
private:
|
||||
mutable std::shared_mutex mutex_{};
|
||||
utils::message_bus &bus_;
|
||||
std::unordered_map<key, std::unique_ptr<cache_entry_base>, key_hash> map_;
|
||||
};
|
||||
|
||||
|
||||
@@ -53,6 +53,18 @@ public:
|
||||
}
|
||||
}
|
||||
|
||||
void resolver(std::weak_ptr<object_resolver<Type>> resolver) {
|
||||
std::lock_guard lock(mutex_);
|
||||
resolver_ = std::move(resolver);
|
||||
}
|
||||
|
||||
[[nodiscard]] std::shared_ptr<Type> object() const {
|
||||
if (!obj_) {
|
||||
std::ignore = resolve();
|
||||
}
|
||||
return obj_;
|
||||
}
|
||||
|
||||
void invalidate() {
|
||||
std::lock_guard lock(mutex_);
|
||||
obj_.reset();
|
||||
|
||||
@@ -57,6 +57,8 @@ public:
|
||||
void reset() { proxy_.reset(); }
|
||||
void reset(const std::shared_ptr<object_proxy<Type>>& proxy) { proxy_ = proxy; }
|
||||
|
||||
[[nodiscard]] std::shared_ptr<object_proxy<Type>> proxy() const { return proxy_; }
|
||||
|
||||
operator bool() const { return valid(); }
|
||||
[[nodiscard]] bool valid() const { return proxy_ != nullptr && !proxy_->empty(); }
|
||||
|
||||
|
||||
Reference in New Issue
Block a user