diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index 0de998382..df7584a9e 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -46,7 +46,7 @@ set(mg_single_node_sources storage/common/property_value_store.cpp storage/single_node/record_accessor.cpp storage/single_node/vertex_accessor.cpp - transactions/single_node/engine_single_node.cpp + transactions/single_node/engine.cpp memgraph_init.cpp ) diff --git a/src/database/counters.hpp b/src/database/distributed/counters.hpp similarity index 100% rename from src/database/counters.hpp rename to src/database/distributed/counters.hpp diff --git a/src/database/distributed/distributed_counters.hpp b/src/database/distributed/distributed_counters.hpp index c83a83c5e..7ac12a6be 100644 --- a/src/database/distributed/distributed_counters.hpp +++ b/src/database/distributed/distributed_counters.hpp @@ -6,7 +6,7 @@ #include #include "data_structures/concurrent/concurrent_map.hpp" -#include "database/counters.hpp" +#include "database/distributed/counters.hpp" namespace communication::rpc { class Server; diff --git a/src/database/distributed/distributed_graph_db.cpp b/src/database/distributed/distributed_graph_db.cpp index f964c0b98..b33ba2d2a 100644 --- a/src/database/distributed/distributed_graph_db.cpp +++ b/src/database/distributed/distributed_graph_db.cpp @@ -25,7 +25,7 @@ #include "durability/distributed/snapshooter.hpp" // TODO: Why do we depend on query here? #include "query/exceptions.hpp" -#include "storage/common/concurrent_id_mapper.hpp" +#include "storage/distributed/concurrent_id_mapper.hpp" #include "storage/distributed/concurrent_id_mapper_master.hpp" #include "storage/distributed/concurrent_id_mapper_worker.hpp" #include "storage/distributed/storage_gc_master.hpp" diff --git a/src/database/distributed/graph_db.hpp b/src/database/distributed/graph_db.hpp index a9d54ba01..6c3e8b575 100644 --- a/src/database/distributed/graph_db.hpp +++ b/src/database/distributed/graph_db.hpp @@ -5,16 +5,16 @@ #include #include -#include "database/counters.hpp" +#include "database/distributed/counters.hpp" #include "durability/distributed/recovery.hpp" #include "durability/distributed/wal.hpp" #include "io/network/endpoint.hpp" -#include "storage/common/concurrent_id_mapper.hpp" #include "storage/common/types.hpp" +#include "storage/distributed/concurrent_id_mapper.hpp" #include "storage/distributed/storage.hpp" #include "storage/distributed/storage_gc.hpp" #include "storage/distributed/vertex_accessor.hpp" -#include "transactions/engine.hpp" +#include "transactions/distributed/engine.hpp" #include "utils/scheduler.hpp" namespace database { diff --git a/src/database/single_node/single_node_counters.hpp b/src/database/single_node/counters.hpp similarity index 77% rename from src/database/single_node/single_node_counters.hpp rename to src/database/single_node/counters.hpp index 4e801868e..5cd681f2b 100644 --- a/src/database/single_node/single_node_counters.hpp +++ b/src/database/single_node/counters.hpp @@ -6,20 +6,19 @@ #include #include "data_structures/concurrent/concurrent_map.hpp" -#include "database/counters.hpp" namespace database { /// Implementation for the single-node memgraph -class SingleNodeCounters : public Counters { +class Counters { public: - int64_t Get(const std::string &name) override { + int64_t Get(const std::string &name) { return counters_.access() .emplace(name, std::make_tuple(name), std::make_tuple(0)) .first->second.fetch_add(1); } - void Set(const std::string &name, int64_t value) override { + void Set(const std::string &name, int64_t value) { auto name_counter_pair = counters_.access().emplace( name, std::make_tuple(name), std::make_tuple(value)); if (!name_counter_pair.second) name_counter_pair.first->second.store(value); diff --git a/src/database/single_node/graph_db.cpp b/src/database/single_node/graph_db.cpp index bda854ab3..1fc079656 100644 --- a/src/database/single_node/graph_db.cpp +++ b/src/database/single_node/graph_db.cpp @@ -4,220 +4,26 @@ #include +#include "database/single_node/counters.hpp" #include "database/single_node/graph_db_accessor.hpp" -#include "database/single_node/single_node_counters.hpp" #include "durability/paths.hpp" #include "durability/single_node/recovery.hpp" #include "durability/single_node/snapshooter.hpp" -#include "storage/single_node/concurrent_id_mapper_single_node.hpp" -#include "storage/single_node/storage_gc_single_node.hpp" -#include "transactions/single_node/engine_single_node.hpp" +#include "storage/single_node/concurrent_id_mapper.hpp" +#include "storage/single_node/storage_gc.hpp" +#include "transactions/single_node/engine.hpp" #include "utils/file.hpp" namespace database { -namespace { - -////////////////////////////////////////////////////////////////////// -// RecordAccessor and GraphDbAccessor implementations -////////////////////////////////////////////////////////////////////// - -template -class SingleNodeRecordAccessor final { - public: - typename RecordAccessor::AddressT GlobalAddress( - const RecordAccessor &record_accessor) { - // TODO: This is still coupled to distributed storage, albeit loosely. - int worker_id = 0; - CHECK(record_accessor.is_local()); - return storage::Address>(record_accessor.gid(), - worker_id); - } - - void SetOldNew(const RecordAccessor &record_accessor, TRecord **old, - TRecord **newr) { - auto &dba = record_accessor.db_accessor(); - const auto &address = record_accessor.address(); - CHECK(record_accessor.is_local()); - address.local()->find_set_old_new(dba.transaction(), old, newr); - } - - TRecord *FindNew(const RecordAccessor &record_accessor) { - const auto &address = record_accessor.address(); - auto &dba = record_accessor.db_accessor(); - CHECK(address.is_local()); - return address.local()->update(dba.transaction()); - } - - void ProcessDelta(const RecordAccessor &record_accessor, - const database::StateDelta &delta) { - CHECK(record_accessor.is_local()); - record_accessor.db_accessor().wal().Emplace(delta); - } - - int64_t CypherId(const RecordAccessor &record_accessor) { - return record_accessor.address().local()->cypher_id(); - } -}; - -class VertexAccessorImpl final : public ::VertexAccessor::Impl { - SingleNodeRecordAccessor accessor_; - - public: - typename RecordAccessor::AddressT GlobalAddress( - const RecordAccessor &ra) override { - return accessor_.GlobalAddress(ra); - } - - void SetOldNew(const RecordAccessor &ra, Vertex **old_record, - Vertex **new_record) override { - return accessor_.SetOldNew(ra, old_record, new_record); - } - - Vertex *FindNew(const RecordAccessor &ra) override { - return accessor_.FindNew(ra); - } - - void ProcessDelta(const RecordAccessor &ra, - const database::StateDelta &delta) override { - return accessor_.ProcessDelta(ra, delta); - } - - void AddLabel(const VertexAccessor &va, - const storage::Label &label) override { - CHECK(va.is_local()); - auto &dba = va.db_accessor(); - auto delta = StateDelta::AddLabel(dba.transaction_id(), va.gid(), label, - dba.LabelName(label)); - Vertex &vertex = va.update(); - // not a duplicate label, add it - if (!utils::Contains(vertex.labels_, label)) { - vertex.labels_.emplace_back(label); - dba.wal().Emplace(delta); - dba.UpdateLabelIndices(label, va, &vertex); - } - } - - void RemoveLabel(const VertexAccessor &va, - const storage::Label &label) override { - CHECK(va.is_local()); - auto &dba = va.db_accessor(); - auto delta = StateDelta::RemoveLabel(dba.transaction_id(), va.gid(), label, - dba.LabelName(label)); - Vertex &vertex = va.update(); - if (utils::Contains(vertex.labels_, label)) { - auto &labels = vertex.labels_; - auto found = std::find(labels.begin(), labels.end(), delta.label); - std::swap(*found, labels.back()); - labels.pop_back(); - dba.wal().Emplace(delta); - } - } - - int64_t CypherId(const RecordAccessor &ra) override { - return accessor_.CypherId(ra); - } -}; - -class EdgeAccessorImpl final : public ::RecordAccessor::Impl { - SingleNodeRecordAccessor accessor_; - - public: - typename RecordAccessor::AddressT GlobalAddress( - const RecordAccessor &ra) override { - return accessor_.GlobalAddress(ra); - } - - void SetOldNew(const RecordAccessor &ra, Edge **old_record, - Edge **new_record) override { - return accessor_.SetOldNew(ra, old_record, new_record); - } - - Edge *FindNew(const RecordAccessor &ra) override { - return accessor_.FindNew(ra); - } - - void ProcessDelta(const RecordAccessor &ra, - const database::StateDelta &delta) override { - return accessor_.ProcessDelta(ra, delta); - } - - int64_t CypherId(const RecordAccessor &ra) override { - return accessor_.CypherId(ra); - } -}; - -class SingleNodeAccessor : public GraphDbAccessor { - // Shared implementations of record accessors. - static VertexAccessorImpl vertex_accessor_; - static EdgeAccessorImpl edge_accessor_; - - public: - explicit SingleNodeAccessor(GraphDb &db) : GraphDbAccessor(db) {} - SingleNodeAccessor(GraphDb &db, tx::TransactionId tx_id) - : GraphDbAccessor(db, tx_id) {} - - ::VertexAccessor::Impl *GetVertexImpl() override { return &vertex_accessor_; } - - ::RecordAccessor::Impl *GetEdgeImpl() override { - return &edge_accessor_; - } -}; - -VertexAccessorImpl SingleNodeAccessor::vertex_accessor_; -EdgeAccessorImpl SingleNodeAccessor::edge_accessor_; - -} // namespace - -////////////////////////////////////////////////////////////////////// -// SingleNode GraphDb implementation -////////////////////////////////////////////////////////////////////// - -template