1783 lines
79 KiB
C++
1783 lines
79 KiB
C++
// Copyright 2023 Memgraph Ltd.
|
|
//
|
|
// Use of this software is governed by the Business Source License
|
|
// included in the file licenses/BSL.txt; by using this file, you agree to be bound by the terms of the Business Source
|
|
// License, and you may not use this file except in compliance with the Business Source License.
|
|
//
|
|
// As of the Change Date specified in that file, in accordance with
|
|
// the Business Source License, use of this software will be governed
|
|
// by the Apache License, Version 2.0, included in the file
|
|
// licenses/APL.txt.
|
|
|
|
#include "storage/v2/storage.hpp"
|
|
#include <algorithm>
|
|
#include <atomic>
|
|
#include <boost/exception/exception.hpp>
|
|
#include <cstdint>
|
|
#include <iterator>
|
|
#include <limits>
|
|
#include <memory>
|
|
#include <mutex>
|
|
#include <stop_token>
|
|
#include <string_view>
|
|
#include <variant>
|
|
#include <vector>
|
|
#include "storage/v2/delta.hpp"
|
|
#include "storage/v2/edge.hpp"
|
|
#include "storage/v2/id_types.hpp"
|
|
#include "storage/v2/isolation_level.hpp"
|
|
#include "storage/v2/mvcc.hpp"
|
|
#include "storage/v2/property_store.hpp"
|
|
#include "storage/v2/result.hpp"
|
|
|
|
#include <gflags/gflags.h>
|
|
#include <spdlog/spdlog.h>
|
|
|
|
#include "io/network/endpoint.hpp"
|
|
#include "storage/v2/disk/indices.hpp"
|
|
#include "storage/v2/disk/storage.hpp"
|
|
#include "storage/v2/durability/durability.hpp"
|
|
#include "storage/v2/durability/metadata.hpp"
|
|
#include "storage/v2/durability/paths.hpp"
|
|
#include "storage/v2/durability/snapshot.hpp"
|
|
#include "storage/v2/durability/wal.hpp"
|
|
#include "storage/v2/edge_accessor.hpp"
|
|
#include "storage/v2/mvcc.hpp"
|
|
#include "storage/v2/replication/config.hpp"
|
|
#include "storage/v2/replication/enums.hpp"
|
|
#include "storage/v2/replication/replication_persistence_helper.hpp"
|
|
#include "storage/v2/storage_mode.hpp"
|
|
#include "storage/v2/transaction.hpp"
|
|
#include "storage/v2/vertex_accessor.hpp"
|
|
#include "storage/v2/view.hpp"
|
|
#include "utils/exceptions.hpp"
|
|
#include "utils/file.hpp"
|
|
#include "utils/logging.hpp"
|
|
#include "utils/memory_tracker.hpp"
|
|
#include "utils/message.hpp"
|
|
#include "utils/rocksdb.hpp"
|
|
#include "utils/rw_lock.hpp"
|
|
#include "utils/spin_lock.hpp"
|
|
#include "utils/stat.hpp"
|
|
#include "utils/uuid.hpp"
|
|
|
|
/// REPLICATION ///
|
|
#include "storage/v2/replication/replication_client.hpp"
|
|
#include "storage/v2/replication/replication_server.hpp"
|
|
#include "storage/v2/replication/rpc.hpp"
|
|
#include "storage/v2/storage_error.hpp"
|
|
|
|
#include "storage/v2/disk/rocksdb_storage.hpp"
|
|
|
|
namespace memgraph::storage {
|
|
|
|
using OOMExceptionEnabler = utils::MemoryTracker::OutOfMemoryExceptionEnabler;
|
|
|
|
namespace {
|
|
|
|
constexpr const char *vertexHandle = "vertex";
|
|
constexpr const char *edgeHandle = "edge";
|
|
constexpr const char *main_storage_path = "./rocks_experiment";
|
|
|
|
inline constexpr uint16_t kEpochHistoryRetention = 1000;
|
|
|
|
} // namespace
|
|
|
|
/// TODO: indices should be initialized at some point after
|
|
DiskStorage::DiskStorage(Config config)
|
|
: Storage(config),
|
|
indices_(&constraints_, config.items),
|
|
isolation_level_(IsolationLevel::SNAPSHOT_ISOLATION),
|
|
storage_mode_(StorageMode::IN_MEMORY_TRANSACTIONAL),
|
|
config_(config) {
|
|
if (config_.durability.snapshot_wal_mode == Config::Durability::SnapshotWalMode::DISABLED
|
|
/// TODO(andi): When replication support will be added, uncomment this.
|
|
// && replication_role_ == ReplicationRole::MAIN) {
|
|
) {
|
|
spdlog::warn(
|
|
"The instance has the MAIN replication role, but durability logs and snapshots are disabled. Please consider "
|
|
"enabling durability by using --storage-snapshot-interval-sec and --storage-wal-enabled flags because "
|
|
"without write-ahead logs this instance is not replicating any data.");
|
|
}
|
|
if (config_.durability.snapshot_wal_mode != Config::Durability::SnapshotWalMode::DISABLED ||
|
|
config_.durability.snapshot_on_exit || config_.durability.recover_on_startup) {
|
|
// Create the directory initially to crash the database in case of
|
|
// permission errors. This is done early to crash the database on startup
|
|
// instead of crashing the database for the first time during runtime (which
|
|
// could be an unpleasant surprise).
|
|
utils::EnsureDirOrDie(snapshot_directory_);
|
|
// Same reasoning as above.
|
|
utils::EnsureDirOrDie(wal_directory_);
|
|
|
|
// Verify that the user that started the process is the same user that is
|
|
// the owner of the storage directory.
|
|
durability::VerifyStorageDirectoryOwnerAndProcessUserOrDie(config_.durability.storage_directory);
|
|
|
|
// Create the lock file and open a handle to it. This will crash the
|
|
// database if it can't open the file for writing or if any other process is
|
|
// holding the file opened.
|
|
lock_file_handle_.Open(lock_file_path_, utils::OutputFile::Mode::OVERWRITE_EXISTING);
|
|
MG_ASSERT(lock_file_handle_.AcquireLock(),
|
|
"Couldn't acquire lock on the storage directory {}"
|
|
"!\nAnother Memgraph process is currently running with the same "
|
|
"storage directory, please stop it first before starting this "
|
|
"process!",
|
|
config_.durability.storage_directory);
|
|
}
|
|
/// TODO(andi): Return capabilities for recovering data
|
|
// if (config_.durability.recover_on_startup) {
|
|
// auto info = durability::RecoverData(snapshot_directory_, wal_directory_, &uuid_, &epoch_id_, &epoch_history_,
|
|
// & vertices_, &edges_, &edge_count_, &name_id_mapper_, &indices_,
|
|
// &constraints_, config_.items, &wal_seq_num_);
|
|
// if (info) {
|
|
// vertex_id_ = info->next_vertex_id;
|
|
// edge_id_ = info->next_edge_id;
|
|
// timestamp_ = std::max(timestamp_, info->next_timestamp);
|
|
// if (info->last_commit_timestamp) {
|
|
// last_commit_timestamp_ = *info->last_commit_timestamp;
|
|
// }
|
|
// }
|
|
// } else if (config_.durability.snapshot_wal_mode != Config::Durability::SnapshotWalMode::DISABLED ||
|
|
// config_.durability.snapshot_on_exit) {
|
|
// bool files_moved = false;
|
|
// auto backup_root = config_.durability.storage_directory / durability::kBackupDirectory;
|
|
// for (const auto &[path, dirname, what] :
|
|
// {std::make_tuple(snapshot_directory_, durability::kSnapshotDirectory, "snapshot"),
|
|
// std::make_tuple(wal_directory_, durability::kWalDirectory, "WAL")}) {
|
|
// if (!utils::DirExists(path)) continue;
|
|
// auto backup_curr = backup_root / dirname;
|
|
// std::error_code error_code;
|
|
// for (const auto &item : std::filesystem::directory_iterator(path, error_code)) {
|
|
// utils::EnsureDirOrDie(backup_root);
|
|
// utils::EnsureDirOrDie(backup_curr);
|
|
// std::error_code item_error_code;
|
|
// std::filesystem::rename(item.path(), backup_curr / item.path().filename(), item_error_code);
|
|
// MG_ASSERT(!item_error_code, "Couldn't move {} file {} because of: {}", what, item.path(),
|
|
// item_error_code.message());
|
|
// files_moved = true;
|
|
// }
|
|
// MG_ASSERT(!error_code, "Couldn't backup {} files because of: {}", what, error_code.message());
|
|
// }
|
|
// if (files_moved) {
|
|
// spdlog::warn(
|
|
// "Since Memgraph was not supposed to recover on startup and "
|
|
// "durability is enabled, your current durability files will likely "
|
|
// "be overridden. To prevent important data loss, Memgraph has stored "
|
|
// "those files into a .backup directory inside the storage directory.");
|
|
// }
|
|
// }
|
|
if (config_.durability.snapshot_wal_mode != Config::Durability::SnapshotWalMode::DISABLED) {
|
|
snapshot_runner_.Run("Snapshot", config_.durability.snapshot_interval, [this] {
|
|
if (auto maybe_error = this->CreateSnapshot({true}); maybe_error.HasError()) {
|
|
switch (maybe_error.GetError()) {
|
|
case CreateSnapshotError::DisabledForReplica:
|
|
spdlog::warn(
|
|
utils::MessageWithLink("Snapshots are disabled for replicas.", "https://memgr.ph/replication"));
|
|
break;
|
|
case CreateSnapshotError::DisabledForAnalyticsPeriodicCommit:
|
|
spdlog::warn(utils::MessageWithLink("Periodic snapshots are disabled for analytical mode.",
|
|
"https://memgr.ph/durability"));
|
|
break;
|
|
case storage::Storage::CreateSnapshotError::ReachedMaxNumTries:
|
|
spdlog::warn("Failed to create snapshot. Reached max number of tries. Please contact support");
|
|
break;
|
|
}
|
|
}
|
|
});
|
|
}
|
|
if (config_.gc.type == Config::Gc::Type::PERIODIC) {
|
|
gc_runner_.Run("Storage GC", config_.gc.interval, [this] { this->CollectGarbage<false>(); });
|
|
}
|
|
|
|
if (timestamp_ == kTimestampInitialId) {
|
|
commit_log_.emplace();
|
|
} else {
|
|
commit_log_.emplace(timestamp_);
|
|
}
|
|
|
|
/// TODO(andi): Return capabilities for restoring replicas
|
|
// if (config_.durability.restore_replicas_on_startup) {
|
|
// spdlog::info("Replica's configuration will be stored and will be automatically restored in case of a crash.");
|
|
// utils::EnsureDirOrDie(config_.durability.storage_directory / durability::kReplicationDirectory);
|
|
// storage_ =
|
|
// std::make_unique<kvstore_::kvstore_>(config_.durability.storage_directory /
|
|
// durability::kReplicationDirectory);
|
|
// RestoreReplicas();
|
|
// } else {
|
|
// spdlog::warn("Replicas' configuration will NOT be stored. When the server restarts, replicas will be
|
|
// forgotten.");
|
|
// }
|
|
|
|
std::filesystem::path rocksdb_path = main_storage_path;
|
|
kvstore_ = std::make_unique<RocksDBStorage>();
|
|
utils::EnsureDirOrDie(rocksdb_path);
|
|
kvstore_->options_.create_if_missing = true;
|
|
kvstore_->options_.comparator = new ComparatorWithU64TsImpl();
|
|
logging::AssertRocksDBStatus(rocksdb::DB::Open(kvstore_->options_, rocksdb_path, &kvstore_->db_));
|
|
logging::AssertRocksDBStatus(
|
|
kvstore_->db_->CreateColumnFamily(kvstore_->options_, vertexHandle, &kvstore_->vertex_chandle));
|
|
logging::AssertRocksDBStatus(
|
|
kvstore_->db_->CreateColumnFamily(kvstore_->options_, edgeHandle, &kvstore_->edge_chandle));
|
|
}
|
|
|
|
DiskStorage::~DiskStorage() {
|
|
/// TODO(andi): I think that without destroy column family handle, there are memory leaks
|
|
/// But I also think that DestroyColumnFamilyHandle deletes all data in its handle.
|
|
logging::AssertRocksDBStatus(kvstore_->db_->DropColumnFamily(kvstore_->vertex_chandle));
|
|
logging::AssertRocksDBStatus(kvstore_->db_->DropColumnFamily(kvstore_->edge_chandle));
|
|
logging::AssertRocksDBStatus(kvstore_->db_->DestroyColumnFamilyHandle(kvstore_->vertex_chandle));
|
|
logging::AssertRocksDBStatus(kvstore_->db_->DestroyColumnFamilyHandle(kvstore_->edge_chandle));
|
|
}
|
|
|
|
DiskStorage::DiskAccessor::DiskAccessor(DiskStorage *storage, IsolationLevel isolation_level, StorageMode storage_mode)
|
|
: Accessor(storage, isolation_level, storage_mode), config_(storage->config_.items) {}
|
|
|
|
DiskStorage::DiskAccessor::DiskAccessor(DiskAccessor &&other) noexcept
|
|
: Accessor(std::move(other)), config_(other.config_) {
|
|
// Don't allow the other accessor to abort our transaction in destructor.
|
|
other.is_transaction_active_ = false;
|
|
other.commit_timestamp_.reset();
|
|
}
|
|
|
|
DiskStorage::DiskAccessor::~DiskAccessor() {
|
|
if (is_transaction_active_) {
|
|
Abort();
|
|
}
|
|
|
|
FinalizeTransaction();
|
|
}
|
|
|
|
std::optional<VertexAccessor> DiskStorage::DiskAccessor::DeserializeVertex(const rocksdb::Slice &key,
|
|
const rocksdb::Slice &value) {
|
|
OOMExceptionEnabler oom_exception;
|
|
const std::vector<std::string> vertex_parts = utils::Split(key.ToStringView(), "|");
|
|
auto gid = storage::Gid::FromUint(std::stoull(vertex_parts[1]));
|
|
auto acc = storage_->vertices_.access();
|
|
/// In situations when the vertex wasn't evicted from the cache, it will be found in the cache.
|
|
if (acc.find(gid) != acc.end()) {
|
|
spdlog::debug("Vertex with gid {} already exists in the cache!", gid.AsUint());
|
|
return std::nullopt;
|
|
}
|
|
spdlog::debug("Vertex with gid {} doesn't exist in the cache!", gid.AsUint());
|
|
uint64_t vertex_commit_ts = utils::ExtractTimestampFromDeserializedUserKey(key);
|
|
// Deserialize labels
|
|
std::vector<LabelId> label_ids;
|
|
if (!vertex_parts[0].empty()) {
|
|
auto labels = utils::Split(vertex_parts[0], ",");
|
|
std::transform(labels.begin(), labels.end(), std::back_inserter(label_ids),
|
|
[](const auto &label) { return storage::LabelId::FromUint(std::stoull(label)); });
|
|
}
|
|
auto impl = CreateVertex(gid, vertex_commit_ts, label_ids, value.ToStringView());
|
|
return impl;
|
|
}
|
|
|
|
std::optional<EdgeAccessor> DiskStorage::DiskAccessor::DeserializeEdge(const rocksdb::Slice &key,
|
|
const rocksdb::Slice &value) {
|
|
const auto edge_parts = utils::Split(key.ToStringView(), "|");
|
|
const Gid edge_gid = utils::DeserializeIdType(edge_parts[4]);
|
|
|
|
auto edge_acc = storage_->edges_.access();
|
|
auto res = edge_acc.find(edge_gid);
|
|
if (res != edge_acc.end()) {
|
|
spdlog::debug("Edge with gid {} already exists in the cache!", edge_parts[4]);
|
|
return std::nullopt;
|
|
}
|
|
spdlog::debug("Edge with gid {} doesn't exist in the cache!", edge_parts[4]);
|
|
|
|
auto [from_gid, to_gid] = std::invoke(
|
|
[&](const auto &edge_parts) {
|
|
if (edge_parts[2] == "0") { // out edge
|
|
return std::make_pair(edge_parts[0], edge_parts[1]);
|
|
}
|
|
// in edge
|
|
return std::make_pair(edge_parts[1], edge_parts[0]);
|
|
},
|
|
edge_parts);
|
|
|
|
// load vertex accessors
|
|
auto from_acc = FindVertex(utils::DeserializeIdType(from_gid), View::OLD);
|
|
auto to_acc = FindVertex(utils::DeserializeIdType(to_gid), View::OLD);
|
|
if (!from_acc || !to_acc) {
|
|
throw utils::BasicException("Non-existing vertices found during edge deserialization");
|
|
}
|
|
const auto edge_type_id = storage::EdgeTypeId::FromUint(std::stoull(edge_parts[3]));
|
|
auto maybe_edge =
|
|
CreateEdge(&*from_acc, &*to_acc, edge_type_id, edge_gid, utils::ExtractTimestampFromDeserializedUserKey(key));
|
|
MG_ASSERT(maybe_edge.HasValue());
|
|
// TODO(andi): Initialize deserialized edge
|
|
// static_cast<DiskEdgeAccessor *>(maybe_edge->get())->InitializeDeserializedEdge(edge_type_id, value.ToStringView());
|
|
return *maybe_edge;
|
|
}
|
|
|
|
VerticesIterable DiskStorage::DiskAccessor::Vertices(View view) {
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
rocksdb::ReadOptions ro;
|
|
rocksdb::Slice ts = utils::StringTimestamp(transaction_.start_timestamp);
|
|
ro.timestamp = &ts;
|
|
auto it = std::unique_ptr<rocksdb::Iterator>(
|
|
disk_storage->kvstore_->db_->NewIterator(ro, disk_storage->kvstore_->vertex_chandle));
|
|
for (it->SeekToFirst(); it->Valid(); it->Next()) {
|
|
// When deserializing vertex, key size is set to user key-size
|
|
// To be able to extract timestamp, here a copy can be created
|
|
// with size explicitly added with sizeof(uint64_t)
|
|
DeserializeVertex(it->key(), it->value());
|
|
}
|
|
return VerticesIterable(AllVerticesIterable(storage_->vertices_.access(), &transaction_, view,
|
|
&disk_storage->indices_, &disk_storage->constraints_,
|
|
disk_storage->config_.items));
|
|
}
|
|
|
|
VerticesIterable DiskStorage::DiskAccessor::Vertices(LabelId label, View view) {
|
|
return VerticesIterable(
|
|
static_cast<DiskStorage *>(storage_)->indices_.label_index.Vertices(label, view, &transaction_));
|
|
}
|
|
|
|
VerticesIterable DiskStorage::DiskAccessor::Vertices(LabelId label, PropertyId property, View view) {
|
|
throw utils::NotYetImplemented("DiskStorage::DiskAComparatorWithccessor::Vertices(label, property)");
|
|
}
|
|
|
|
VerticesIterable DiskStorage::DiskAccessor::Vertices(LabelId label, PropertyId property, const PropertyValue &value,
|
|
View view) {
|
|
throw utils::NotYetImplemented("DiskStorage::DiskAccessor::Vertices(label, property, value)");
|
|
}
|
|
|
|
VerticesIterable DiskStorage::DiskAccessor::Vertices(LabelId label, PropertyId property,
|
|
const std::optional<utils::Bound<PropertyValue>> &lower_bound,
|
|
const std::optional<utils::Bound<PropertyValue>> &upper_bound,
|
|
View view) {
|
|
throw utils::NotYetImplemented("DiskStorage::DiskAccessor::Vertices(label, property, lower_bound, upper_bound)");
|
|
}
|
|
|
|
int64_t DiskStorage::DiskAccessor::ApproximateVertexCount() const {
|
|
uint64_t estimate_num_keys = 0;
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
// TODO: This method should probably be organized.
|
|
disk_storage->kvstore_->db_->GetIntProperty(disk_storage->kvstore_->vertex_chandle, "rocksdb.estimate-num-keys",
|
|
&estimate_num_keys);
|
|
return static_cast<int64_t>(estimate_num_keys);
|
|
}
|
|
|
|
VertexAccessor DiskStorage::DiskAccessor::CreateVertex() {
|
|
OOMExceptionEnabler oom_exception;
|
|
auto gid = storage_->vertex_id_.fetch_add(1, std::memory_order_acq_rel);
|
|
auto acc = storage_->vertices_.access();
|
|
|
|
auto *delta = CreateDeleteObjectDelta(&transaction_);
|
|
auto [it, inserted] = acc.insert(Vertex{storage::Gid::FromUint(gid), delta});
|
|
MG_ASSERT(inserted, "The vertex must be inserted here!");
|
|
MG_ASSERT(it != acc.end(), "Invalid Vertex accessor!");
|
|
|
|
if (delta) {
|
|
delta->prev.Set(&*it);
|
|
}
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
|
|
return {&*it, &transaction_, &disk_storage->indices_, &disk_storage->constraints_, config_};
|
|
}
|
|
|
|
VertexAccessor DiskStorage::DiskAccessor::CreateVertex(storage::Gid gid) {
|
|
OOMExceptionEnabler oom_exception;
|
|
// NOTE: When we update the next `vertex_id_` here we perform a RMW
|
|
// (read-modify-write) operation that ISN'T atomic! But, that isn't an issue
|
|
// because this function is only called from the replication delta applier
|
|
// that runs single-threadedly and while this instance is set-up to apply
|
|
// threads (it is the replica), it is guaranteed that no other writes are
|
|
// possible.
|
|
spdlog::debug("Vertex with gid {} doesn't exist in the cache, creating it!", gid.AsUint());
|
|
auto acc = storage_->vertices_.access();
|
|
storage_->vertex_id_.store(std::max(storage_->vertex_id_.load(std::memory_order_acquire), gid.AsUint() + 1),
|
|
std::memory_order_release);
|
|
|
|
auto *delta = CreateDeleteObjectDelta(&transaction_);
|
|
auto [it, inserted] = acc.insert(Vertex{gid, delta});
|
|
/// if the vertex with the given gid doesn't exist on the disk, it must be inserted here.
|
|
MG_ASSERT(inserted, "The vertex must be inserted here!");
|
|
MG_ASSERT(it != acc.end(), "Invalid Vertex accessor!");
|
|
if (delta) {
|
|
delta->prev.Set(&*it);
|
|
}
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
return {&*it, &transaction_, &disk_storage->indices_, &disk_storage->constraints_, config_};
|
|
}
|
|
|
|
/// TODO(andi): This method is the duplicate of CreateVertex(storage::Gid gid), the only thing that is different is
|
|
/// delta creation How to remove this duplication?
|
|
VertexAccessor DiskStorage::DiskAccessor::CreateVertex(storage::Gid gid, uint64_t vertex_commit_ts,
|
|
std::vector<LabelId> label_ids, std::string_view properties) {
|
|
OOMExceptionEnabler oom_exception;
|
|
// NOTE: When we update the next `vertex_id_` here we perform a RMW
|
|
// (read-modify-write) operation that ISN'T atomic! But, that isn't an issue
|
|
// because this function is only called from the replication delta applier
|
|
// that runs single-threadedly and while this instance is set-up to apply
|
|
// threads (it is the replica), it is guaranteed that no other writes are
|
|
// possible.
|
|
auto acc = storage_->vertices_.access();
|
|
storage_->vertex_id_.store(std::max(storage_->vertex_id_.load(std::memory_order_acquire), gid.AsUint() + 1),
|
|
std::memory_order_release);
|
|
auto *delta = CreateDeleteDeserializedObjectDelta(&transaction_, vertex_commit_ts);
|
|
auto [it, inserted] = acc.insert(Vertex{gid, delta});
|
|
// NOTE: If the vertex with the given gid doesn't exist on the disk, it must be inserted here.
|
|
MG_ASSERT(inserted, "The vertex must be inserted here!");
|
|
MG_ASSERT(it != acc.end(), "Invalid Vertex accessor!");
|
|
for (auto label_id : label_ids) {
|
|
it->labels.push_back(label_id);
|
|
}
|
|
it->properties.SetBuffer(properties);
|
|
if (delta) {
|
|
delta->prev.Set(&*it);
|
|
}
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
return {&*it, &transaction_, &disk_storage->indices_, &disk_storage->constraints_, config_};
|
|
}
|
|
|
|
std::optional<VertexAccessor> DiskStorage::DiskAccessor::FindVertex(storage::Gid gid, View view) {
|
|
/// Check if the vertex is in the cache.
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
auto acc = storage_->vertices_.access();
|
|
auto vertex_it = acc.find(gid);
|
|
if (vertex_it != acc.end()) {
|
|
if (vertex_it->delta) {
|
|
spdlog::debug("Vertex with gid {} found in the cache with ts {}!", gid.AsUint(),
|
|
vertex_it->delta->timestamp->load(std::memory_order_acquire));
|
|
} else {
|
|
spdlog::debug("Vertex with gid {} found in the cache! Delta is null.", gid.AsUint());
|
|
}
|
|
return VertexAccessor::Create(&*vertex_it, &transaction_, &disk_storage->indices_, &disk_storage->constraints_,
|
|
config_, view);
|
|
}
|
|
/// If not in the memory, check whether it exists in RocksDB.
|
|
rocksdb::ReadOptions read_opts;
|
|
rocksdb::Slice ts = utils::StringTimestamp(transaction_.start_timestamp);
|
|
read_opts.timestamp = &ts;
|
|
auto it = std::unique_ptr<rocksdb::Iterator>(
|
|
disk_storage->kvstore_->db_->NewIterator(read_opts, disk_storage->kvstore_->vertex_chandle));
|
|
for (it->SeekToFirst(); it->Valid(); it->Next()) {
|
|
const auto &key = it->key();
|
|
// TODO(andi): If we change format of vertex serialization, change vertex_parts[1] to vertex_parts[0].
|
|
if (const auto vertex_parts = utils::Split(key.ToString(), "|"); vertex_parts[1] == utils::SerializeIdType(gid)) {
|
|
return DeserializeVertex(key, it->value());
|
|
}
|
|
}
|
|
return std::nullopt;
|
|
}
|
|
|
|
Result<std::optional<VertexAccessor>> DiskStorage::DiskAccessor::DeleteVertex(VertexAccessor *vertex) {
|
|
MG_ASSERT(vertex->transaction_ == &transaction_,
|
|
"VertexAccessor must be from the same transaction as the storage "
|
|
"accessor when deleting a vertex!");
|
|
auto *vertex_ptr = vertex->vertex_;
|
|
|
|
std::lock_guard<utils::SpinLock> guard(vertex_ptr->lock);
|
|
|
|
if (!PrepareForWrite(&transaction_, vertex_ptr)) return Error::SERIALIZATION_ERROR;
|
|
|
|
if (vertex_ptr->deleted) {
|
|
return std::optional<VertexAccessor>{};
|
|
}
|
|
|
|
if (!vertex_ptr->in_edges.empty() || !vertex_ptr->out_edges.empty()) return Error::VERTEX_HAS_EDGES;
|
|
|
|
CreateAndLinkDelta(&transaction_, vertex_ptr, Delta::RecreateObjectTag());
|
|
vertex_ptr->deleted = true;
|
|
vertices_to_delete_.emplace_back(utils::SerializeVertex(*vertex_ptr));
|
|
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
return std::make_optional<VertexAccessor>(vertex_ptr, &transaction_, &disk_storage->indices_,
|
|
&disk_storage->constraints_, config_, true);
|
|
}
|
|
|
|
Result<std::optional<std::pair<VertexAccessor, std::vector<EdgeAccessor>>>>
|
|
DiskStorage::DiskAccessor::DetachDeleteVertex(VertexAccessor *vertex) {
|
|
using ReturnType = std::pair<VertexAccessor, std::vector<EdgeAccessor>>;
|
|
MG_ASSERT(vertex->transaction_ == &transaction_,
|
|
"VertexAccessor must be from the same transaction as the storage "
|
|
"accessor when deleting a vertex!");
|
|
auto *vertex_ptr = vertex->vertex_;
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
|
|
std::vector<std::tuple<EdgeTypeId, Vertex *, EdgeRef>> in_edges;
|
|
std::vector<std::tuple<EdgeTypeId, Vertex *, EdgeRef>> out_edges;
|
|
|
|
{
|
|
std::lock_guard<utils::SpinLock> guard(vertex_ptr->lock);
|
|
|
|
if (!PrepareForWrite(&transaction_, vertex_ptr)) return Error::SERIALIZATION_ERROR;
|
|
|
|
if (vertex_ptr->deleted) return std::optional<ReturnType>{};
|
|
|
|
in_edges = vertex_ptr->in_edges;
|
|
out_edges = vertex_ptr->out_edges;
|
|
}
|
|
|
|
std::vector<EdgeAccessor> deleted_edges;
|
|
for (const auto &item : in_edges) {
|
|
auto [edge_type, from_vertex, edge] = item;
|
|
EdgeAccessor e(edge, edge_type, from_vertex, vertex_ptr, &transaction_, &disk_storage->indices_,
|
|
&disk_storage->constraints_, config_);
|
|
auto ret = DeleteEdge(&e);
|
|
if (ret.HasError()) {
|
|
MG_ASSERT(ret.GetError() == Error::SERIALIZATION_ERROR, "Invalid database state!");
|
|
return ret.GetError();
|
|
}
|
|
|
|
if (ret.GetValue()) {
|
|
deleted_edges.push_back(*ret.GetValue());
|
|
}
|
|
}
|
|
for (const auto &item : out_edges) {
|
|
auto [edge_type, to_vertex, edge] = item;
|
|
EdgeAccessor e(edge, edge_type, vertex_ptr, to_vertex, &transaction_, &disk_storage->indices_,
|
|
&disk_storage->constraints_, config_);
|
|
auto ret = DeleteEdge(&e);
|
|
if (ret.HasError()) {
|
|
MG_ASSERT(ret.GetError() == Error::SERIALIZATION_ERROR, "Invalid database state!");
|
|
return ret.GetError();
|
|
}
|
|
|
|
if (ret.GetValue()) {
|
|
deleted_edges.push_back(*ret.GetValue());
|
|
}
|
|
}
|
|
|
|
std::lock_guard<utils::SpinLock> guard(vertex_ptr->lock);
|
|
|
|
// We need to check again for serialization errors because we unlocked the
|
|
// vertex. Some other transaction could have modified the vertex in the
|
|
// meantime if we didn't have any edges to delete.
|
|
|
|
if (!PrepareForWrite(&transaction_, vertex_ptr)) return Error::SERIALIZATION_ERROR;
|
|
|
|
MG_ASSERT(!vertex_ptr->deleted, "Invalid database state!");
|
|
|
|
CreateAndLinkDelta(&transaction_, vertex_ptr, Delta::RecreateObjectTag());
|
|
vertex_ptr->deleted = true;
|
|
vertices_to_delete_.emplace_back(utils::SerializeVertex(*vertex_ptr));
|
|
|
|
return std::make_optional<ReturnType>(
|
|
VertexAccessor{vertex_ptr, &transaction_, &disk_storage->indices_, &disk_storage->constraints_, config_, true},
|
|
std::move(deleted_edges));
|
|
}
|
|
|
|
void DiskStorage::DiskAccessor::PrefetchEdges(const auto &prefetch_edge_filter) {
|
|
rocksdb::ReadOptions read_opts;
|
|
rocksdb::Slice ts = utils::StringTimestamp(transaction_.start_timestamp);
|
|
read_opts.timestamp = &ts;
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
auto it = std::unique_ptr<rocksdb::Iterator>(
|
|
disk_storage->kvstore_->db_->NewIterator(read_opts, disk_storage->kvstore_->edge_chandle));
|
|
for (it->SeekToFirst(); it->Valid(); it->Next()) {
|
|
const rocksdb::Slice &key = it->key();
|
|
const auto edge_parts = utils::Split(key.ToStringView(), "|");
|
|
if (prefetch_edge_filter(edge_parts[0], edge_parts[2])) {
|
|
DeserializeEdge(key, it->value());
|
|
}
|
|
}
|
|
}
|
|
|
|
/// TOOD(andi): Add support for fetching in edges only for one vertex. Currently, all in edges are fetched.
|
|
void DiskStorage::DiskAccessor::PrefetchInEdges(const VertexAccessor &vertex_acc) {
|
|
PrefetchEdges([&vertex_acc](const std::string_view disk_edge_gid,
|
|
const std::string_view disk_edge_direction) -> bool {
|
|
return disk_edge_gid == utils::SerializeIdType(vertex_acc.Gid()) && disk_edge_direction == utils::inEdgeDirection;
|
|
});
|
|
}
|
|
|
|
/// TODO(andi): Functionality of this method can probably be merged with the functionality of in-edges
|
|
void DiskStorage::DiskAccessor::PrefetchOutEdges(const VertexAccessor &vertex_acc) {
|
|
PrefetchEdges([&vertex_acc](const std::string_view disk_edge_gid,
|
|
const std::string_view disk_edge_direction) -> bool {
|
|
return disk_edge_gid == utils::SerializeIdType(vertex_acc.Gid()) && disk_edge_direction == utils::outEdgeDirection;
|
|
});
|
|
}
|
|
|
|
Result<EdgeAccessor> DiskStorage::DiskAccessor::CreateEdge(VertexAccessor *from, VertexAccessor *to,
|
|
EdgeTypeId edge_type, storage::Gid gid,
|
|
uint64_t edge_commit_ts) {
|
|
OOMExceptionEnabler oom_exception;
|
|
MG_ASSERT(from->transaction_ == to->transaction_,
|
|
"VertexAccessors must be from the same transaction when creating "
|
|
"an edge!");
|
|
MG_ASSERT(from->transaction_ == &transaction_,
|
|
"VertexAccessors must be from the same transaction in when "
|
|
"creating an edge!");
|
|
|
|
auto from_vertex = from->vertex_;
|
|
auto to_vertex = to->vertex_;
|
|
|
|
// Obtain the locks by `gid` order to avoid lock cycles.
|
|
std::unique_lock<utils::SpinLock> guard_from(from_vertex->lock, std::defer_lock);
|
|
std::unique_lock<utils::SpinLock> guard_to(to_vertex->lock, std::defer_lock);
|
|
if (from_vertex->gid < to_vertex->gid) {
|
|
guard_from.lock();
|
|
guard_to.lock();
|
|
} else if (from_vertex->gid > to_vertex->gid) {
|
|
guard_to.lock();
|
|
guard_from.lock();
|
|
} else {
|
|
// The vertices are the same vertex, only lock one.
|
|
guard_from.lock();
|
|
}
|
|
|
|
if (!PrepareForWrite(&transaction_, from_vertex)) return Error::SERIALIZATION_ERROR;
|
|
if (from_vertex->deleted) return Error::DELETED_OBJECT;
|
|
|
|
if (to_vertex != from_vertex) {
|
|
if (!PrepareForWrite(&transaction_, to_vertex)) return Error::SERIALIZATION_ERROR;
|
|
if (to_vertex->deleted) return Error::DELETED_OBJECT;
|
|
}
|
|
|
|
// NOTE: When we update the next `edge_id_` here we perform a RMW
|
|
// (read-modify-write) operation that ISN'T atomic! But, that isn't an issue
|
|
// because this function is only called from the replication delta applier
|
|
// that runs single-threadedly and while this instance is set-up to apply
|
|
// threads (it is the replica), it is guaranteed that no other writes are
|
|
// possible.
|
|
storage_->edge_id_.store(std::max(storage_->edge_id_.load(std::memory_order_acquire), gid.AsUint() + 1),
|
|
std::memory_order_release);
|
|
|
|
/// TODO(andi): Remove this once we add full support for edges.
|
|
MG_ASSERT(config_.properties_on_edges, "Properties on edges must be enabled currently for Disk version!");
|
|
// if (config_.properties_on_edges) {
|
|
auto acc = storage_->edges_.access();
|
|
auto delta = CreateDeleteDeserializedObjectDelta(&transaction_, edge_commit_ts);
|
|
auto [it, inserted] = acc.insert(Edge(gid, delta));
|
|
MG_ASSERT(inserted, "The edge must be inserted here!");
|
|
MG_ASSERT(it != acc.end(), "Invalid Edge accessor!");
|
|
auto edge = EdgeRef(&*it);
|
|
delta->prev.Set(&*it);
|
|
// }
|
|
from_vertex->out_edges.emplace_back(edge_type, to_vertex, edge);
|
|
to_vertex->in_edges.emplace_back(edge_type, from_vertex, edge);
|
|
|
|
// Increment edge count.
|
|
storage_->edge_count_.fetch_add(1, std::memory_order_acq_rel);
|
|
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
return EdgeAccessor(edge, edge_type, from_vertex, to_vertex, &transaction_, &disk_storage->indices_,
|
|
&disk_storage->constraints_, config_);
|
|
}
|
|
|
|
Result<EdgeAccessor> DiskStorage::DiskAccessor::CreateEdge(VertexAccessor *from, VertexAccessor *to,
|
|
EdgeTypeId edge_type, storage::Gid gid) {
|
|
OOMExceptionEnabler oom_exception;
|
|
MG_ASSERT(from->transaction_ == to->transaction_,
|
|
"VertexAccessors must be from the same transaction when creating "
|
|
"an edge!");
|
|
MG_ASSERT(from->transaction_ == &transaction_,
|
|
"VertexAccessors must be from the same transaction in when "
|
|
"creating an edge!");
|
|
|
|
auto from_vertex = from->vertex_;
|
|
auto to_vertex = to->vertex_;
|
|
|
|
// Obtain the locks by `gid` order to avoid lock cycles.
|
|
std::unique_lock<utils::SpinLock> guard_from(from_vertex->lock, std::defer_lock);
|
|
std::unique_lock<utils::SpinLock> guard_to(to_vertex->lock, std::defer_lock);
|
|
if (from_vertex->gid < to_vertex->gid) {
|
|
guard_from.lock();
|
|
guard_to.lock();
|
|
} else if (from_vertex->gid > to_vertex->gid) {
|
|
guard_to.lock();
|
|
guard_from.lock();
|
|
} else {
|
|
// The vertices are the same vertex, only lock one.
|
|
guard_from.lock();
|
|
}
|
|
|
|
if (!PrepareForWrite(&transaction_, from_vertex)) return Error::SERIALIZATION_ERROR;
|
|
if (from_vertex->deleted) return Error::DELETED_OBJECT;
|
|
|
|
if (to_vertex != from_vertex) {
|
|
if (!PrepareForWrite(&transaction_, to_vertex)) return Error::SERIALIZATION_ERROR;
|
|
if (to_vertex->deleted) return Error::DELETED_OBJECT;
|
|
}
|
|
|
|
// NOTE: When we update the next `edge_id_` here we perform a RMW
|
|
// (read-modify-write) operation that ISN'T atomic! But, that isn't an issue
|
|
// because this function is only called from the replication delta applier
|
|
// that runs single-threadedly and while this instance is set-up to apply
|
|
// threads (it is the replica), it is guaranteed that no other writes are
|
|
// possible.
|
|
storage_->edge_id_.store(std::max(storage_->edge_id_.load(std::memory_order_acquire), gid.AsUint() + 1),
|
|
std::memory_order_release);
|
|
|
|
EdgeRef edge(gid);
|
|
if (config_.properties_on_edges) {
|
|
auto acc = storage_->edges_.access();
|
|
auto delta = CreateDeleteObjectDelta(&transaction_);
|
|
auto [it, inserted] = acc.insert(Edge(gid, delta));
|
|
MG_ASSERT(inserted, "The edge must be inserted here!");
|
|
MG_ASSERT(it != acc.end(), "Invalid Edge accessor!");
|
|
edge = EdgeRef(&*it);
|
|
delta->prev.Set(&*it);
|
|
}
|
|
|
|
CreateAndLinkDelta(&transaction_, from_vertex, Delta::RemoveOutEdgeTag(), edge_type, to_vertex, edge);
|
|
from_vertex->out_edges.emplace_back(edge_type, to_vertex, edge);
|
|
|
|
CreateAndLinkDelta(&transaction_, to_vertex, Delta::RemoveInEdgeTag(), edge_type, from_vertex, edge);
|
|
to_vertex->in_edges.emplace_back(edge_type, from_vertex, edge);
|
|
|
|
// Increment edge count.
|
|
storage_->edge_count_.fetch_add(1, std::memory_order_acq_rel);
|
|
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
return EdgeAccessor(edge, edge_type, from_vertex, to_vertex, &transaction_, &disk_storage->indices_,
|
|
&disk_storage->constraints_, config_);
|
|
}
|
|
|
|
Result<EdgeAccessor> DiskStorage::DiskAccessor::CreateEdge(VertexAccessor *from, VertexAccessor *to,
|
|
EdgeTypeId edge_type) {
|
|
MG_ASSERT(from->transaction_ == &transaction_,
|
|
"VertexAccessor must be from the same transaction as the storage "
|
|
"accessor when deleting a vertex!");
|
|
MG_ASSERT(to->transaction_ == &transaction_,
|
|
"VertexAccessor must be from the same transaction as the storage "
|
|
"accessor when deleting a vertex!");
|
|
|
|
auto from_vertex = from->vertex_;
|
|
auto to_vertex = to->vertex_;
|
|
|
|
// Obtain the locks by `gid` order to avoid lock cycles.
|
|
std::unique_lock<utils::SpinLock> guard_from(from_vertex->lock, std::defer_lock);
|
|
std::unique_lock<utils::SpinLock> guard_to(to_vertex->lock, std::defer_lock);
|
|
if (from_vertex->gid < to_vertex->gid) {
|
|
guard_from.lock();
|
|
guard_to.lock();
|
|
} else if (from_vertex->gid > to_vertex->gid) {
|
|
guard_to.lock();
|
|
guard_from.lock();
|
|
} else {
|
|
// The vertices are the same vertex, only lock one.
|
|
guard_from.lock();
|
|
}
|
|
|
|
if (!PrepareForWrite(&transaction_, from_vertex)) return Error::SERIALIZATION_ERROR;
|
|
if (from_vertex->deleted) return Error::DELETED_OBJECT;
|
|
|
|
if (to_vertex != from_vertex) {
|
|
if (!PrepareForWrite(&transaction_, to_vertex)) return Error::SERIALIZATION_ERROR;
|
|
if (to_vertex->deleted) return Error::DELETED_OBJECT;
|
|
}
|
|
|
|
auto gid = storage::Gid::FromUint(storage_->edge_id_.fetch_add(1, std::memory_order_acq_rel));
|
|
EdgeRef edge(gid);
|
|
if (config_.properties_on_edges) {
|
|
auto acc = storage_->edges_.access();
|
|
auto *delta = CreateDeleteObjectDelta(&transaction_);
|
|
auto [it, inserted] = acc.insert(Edge(gid, delta));
|
|
MG_ASSERT(inserted, "The edge must be inserted here!");
|
|
MG_ASSERT(it != acc.end(), "Invalid Edge accessor!");
|
|
edge = EdgeRef(&*it);
|
|
delta->prev.Set(&*it);
|
|
}
|
|
|
|
CreateAndLinkDelta(&transaction_, from_vertex, Delta::RemoveOutEdgeTag(), edge_type, to_vertex, edge);
|
|
from_vertex->out_edges.emplace_back(edge_type, to_vertex, edge);
|
|
|
|
CreateAndLinkDelta(&transaction_, to_vertex, Delta::RemoveInEdgeTag(), edge_type, from_vertex, edge);
|
|
to_vertex->in_edges.emplace_back(edge_type, from_vertex, edge);
|
|
|
|
// Increment edge count.
|
|
storage_->edge_count_.fetch_add(1, std::memory_order_acq_rel);
|
|
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
return EdgeAccessor(edge, edge_type, from_vertex, to_vertex, &transaction_, &disk_storage->indices_,
|
|
&disk_storage->constraints_, config_);
|
|
}
|
|
|
|
Result<std::optional<EdgeAccessor>> DiskStorage::DiskAccessor::DeleteEdge(EdgeAccessor *edge) {
|
|
MG_ASSERT(edge->transaction_ == &transaction_,
|
|
"EdgeAccessor must be from the same transaction as the storage "
|
|
"accessor when deleting an edge!");
|
|
auto edge_ref = edge->edge_;
|
|
auto edge_type = edge->edge_type_;
|
|
|
|
std::unique_lock<utils::SpinLock> guard;
|
|
if (config_.properties_on_edges) {
|
|
auto *edge_ptr = edge_ref.ptr;
|
|
guard = std::unique_lock<utils::SpinLock>(edge_ptr->lock);
|
|
|
|
if (!PrepareForWrite(&transaction_, edge_ptr)) return Error::SERIALIZATION_ERROR;
|
|
|
|
if (edge_ptr->deleted) return std::optional<EdgeAccessor>{};
|
|
}
|
|
|
|
auto *from_vertex = edge->from_vertex_;
|
|
auto *to_vertex = edge->to_vertex_;
|
|
|
|
// Obtain the locks by `gid` order to avoid lock cycles.
|
|
std::unique_lock<utils::SpinLock> guard_from(from_vertex->lock, std::defer_lock);
|
|
std::unique_lock<utils::SpinLock> guard_to(to_vertex->lock, std::defer_lock);
|
|
if (from_vertex->gid < to_vertex->gid) {
|
|
guard_from.lock();
|
|
guard_to.lock();
|
|
} else if (from_vertex->gid > to_vertex->gid) {
|
|
guard_to.lock();
|
|
guard_from.lock();
|
|
} else {
|
|
// The vertices are the same vertex, only lock one.
|
|
guard_from.lock();
|
|
}
|
|
|
|
if (!PrepareForWrite(&transaction_, from_vertex)) return Error::SERIALIZATION_ERROR;
|
|
MG_ASSERT(!from_vertex->deleted, "Invalid database state!");
|
|
|
|
if (to_vertex != from_vertex) {
|
|
if (!PrepareForWrite(&transaction_, to_vertex)) return Error::SERIALIZATION_ERROR;
|
|
MG_ASSERT(!to_vertex->deleted, "Invalid database state!");
|
|
}
|
|
|
|
auto delete_edge_from_storage = [&edge_type, &edge_ref, this](auto *vertex, auto *edges) {
|
|
std::tuple<EdgeTypeId, Vertex *, EdgeRef> link(edge_type, vertex, edge_ref);
|
|
auto it = std::find(edges->begin(), edges->end(), link);
|
|
if (config_.properties_on_edges) {
|
|
MG_ASSERT(it != edges->end(), "Invalid database state!");
|
|
} else if (it == edges->end()) {
|
|
return false;
|
|
}
|
|
std::swap(*it, *edges->rbegin());
|
|
edges->pop_back();
|
|
return true;
|
|
};
|
|
|
|
auto op1 = delete_edge_from_storage(to_vertex, &from_vertex->out_edges);
|
|
auto op2 = delete_edge_from_storage(from_vertex, &to_vertex->in_edges);
|
|
auto [src_dest_del_key, dest_src_del_key] =
|
|
utils::SerializeEdge(from_vertex->gid, to_vertex->gid, edge_type, edge_ref.ptr);
|
|
edges_to_delete_.emplace_back(src_dest_del_key);
|
|
edges_to_delete_.emplace_back(dest_src_del_key);
|
|
|
|
if (config_.properties_on_edges) {
|
|
MG_ASSERT((op1 && op2), "Invalid database state!");
|
|
} else {
|
|
MG_ASSERT((op1 && op2) || (!op1 && !op2), "Invalid database state!");
|
|
if (!op1 && !op2) {
|
|
// The edge is already deleted.
|
|
return std::optional<EdgeAccessor>{};
|
|
}
|
|
}
|
|
|
|
if (config_.properties_on_edges) {
|
|
auto *edge_ptr = edge_ref.ptr;
|
|
CreateAndLinkDelta(&transaction_, edge_ptr, Delta::RecreateObjectTag());
|
|
edge_ptr->deleted = true;
|
|
}
|
|
// TODO(andi): How to know that the edge is deleted if properties_on_edges are disabled?
|
|
|
|
CreateAndLinkDelta(&transaction_, from_vertex, Delta::AddOutEdgeTag(), edge_type, to_vertex, edge_ref);
|
|
CreateAndLinkDelta(&transaction_, to_vertex, Delta::AddInEdgeTag(), edge_type, from_vertex, edge_ref);
|
|
|
|
// Decrement edge count.
|
|
storage_->edge_count_.fetch_add(-1, std::memory_order_acq_rel);
|
|
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
return std::make_optional<EdgeAccessor>(edge_ref, edge_type, from_vertex, to_vertex, &transaction_,
|
|
&disk_storage->indices_, &disk_storage->constraints_, config_, true);
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
const std::string &DiskStorage::DiskAccessor::LabelToName(LabelId label) const { return storage_->LabelToName(label); }
|
|
|
|
// this should be handled on an above level of abstraction
|
|
const std::string &DiskStorage::DiskAccessor::PropertyToName(PropertyId property) const {
|
|
return storage_->PropertyToName(property);
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
const std::string &DiskStorage::DiskAccessor::EdgeTypeToName(EdgeTypeId edge_type) const {
|
|
return storage_->EdgeTypeToName(edge_type);
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
LabelId DiskStorage::DiskAccessor::NameToLabel(const std::string_view name) { return storage_->NameToLabel(name); }
|
|
|
|
// this should be handled on an above level of abstraction
|
|
PropertyId DiskStorage::DiskAccessor::NameToProperty(const std::string_view name) {
|
|
return storage_->NameToProperty(name);
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
EdgeTypeId DiskStorage::DiskAccessor::NameToEdgeType(const std::string_view name) {
|
|
return storage_->NameToEdgeType(name);
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
void DiskStorage::DiskAccessor::AdvanceCommand() { ++transaction_.command_id; }
|
|
|
|
void DiskStorage::DiskAccessor::FlushCache() {
|
|
/// Flush vertex cache.
|
|
auto vertex_acc = storage_->vertices_.access();
|
|
uint64_t num_ser_edges = 0;
|
|
rocksdb::WriteOptions write_options;
|
|
rocksdb::Slice ts = utils::StringTimestamp(*commit_timestamp_);
|
|
write_options.timestamp = &ts;
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
for (Vertex &vertex : vertex_acc) {
|
|
logging::AssertRocksDBStatus(disk_storage->kvstore_->db_->Put(write_options, disk_storage->kvstore_->vertex_chandle,
|
|
utils::SerializeVertex(vertex),
|
|
utils::SerializeProperties(vertex.properties)));
|
|
spdlog::debug("rocksdb: Saved vertex with key {} and ts {}", utils::SerializeVertex(vertex), *commit_timestamp_);
|
|
spdlog::debug("Vertex {} has {} out edges", vertex.gid.AsUint(), vertex.out_edges.size());
|
|
for (auto &edge_entry : vertex.out_edges) {
|
|
Edge *edge_ptr = std::get<2>(edge_entry).ptr;
|
|
auto [src_dest_key, dest_src_key] =
|
|
utils::SerializeEdge(vertex.gid, std::get<1>(edge_entry)->gid, std::get<0>(edge_entry), edge_ptr);
|
|
logging::AssertRocksDBStatus(disk_storage->kvstore_->db_->Put(write_options, disk_storage->kvstore_->edge_chandle,
|
|
src_dest_key,
|
|
utils::SerializeProperties(edge_ptr->properties)));
|
|
logging::AssertRocksDBStatus(disk_storage->kvstore_->db_->Put(write_options, disk_storage->kvstore_->edge_chandle,
|
|
dest_src_key,
|
|
utils::SerializeProperties(edge_ptr->properties)));
|
|
spdlog::debug("rocksdb: Saved edge with key {} and ts {}", src_dest_key, *commit_timestamp_);
|
|
spdlog::debug("rocksdb: Saved edge with key {} and ts {}", dest_src_key, *commit_timestamp_);
|
|
num_ser_edges++;
|
|
}
|
|
}
|
|
// Delete vertices that were deleted in the current transaction.
|
|
for (const auto &vertex_to_delete : vertices_to_delete_) {
|
|
spdlog::debug("rocksdb: Deleted vertex with key {}", vertex_to_delete);
|
|
logging::AssertRocksDBStatus(
|
|
disk_storage->kvstore_->db_->Delete(write_options, disk_storage->kvstore_->vertex_chandle, vertex_to_delete));
|
|
}
|
|
|
|
// Delete edges that were deleted in the current transaction.
|
|
for (const auto &edge_to_delete : edges_to_delete_) {
|
|
spdlog::debug("rocksdb: Deleted edge with key {}", edge_to_delete);
|
|
logging::AssertRocksDBStatus(
|
|
disk_storage->kvstore_->db_->Delete(write_options, disk_storage->kvstore_->edge_chandle, edge_to_delete));
|
|
}
|
|
}
|
|
|
|
// this will be modified here for a disk-based storage
|
|
utils::BasicResult<StorageDataManipulationError, void> DiskStorage::DiskAccessor::Commit(
|
|
const std::optional<uint64_t> desired_commit_timestamp) {
|
|
MG_ASSERT(is_transaction_active_, "The transaction is already terminated!");
|
|
MG_ASSERT(!transaction_.must_abort, "The transaction can't be committed!");
|
|
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
auto could_replicate_all_sync_replicas = true;
|
|
|
|
if (transaction_.deltas.empty() ||
|
|
std::all_of(transaction_.deltas.begin(), transaction_.deltas.end(),
|
|
[](const Delta &delta) { return delta.action == Delta::Action::DELETE_DESERIALIZED_OBJECT; })) {
|
|
// We don't have to update the commit timestamp here because no one reads
|
|
// it.
|
|
// If there are no deltas, then we don't have to serialize anything on the disk.
|
|
disk_storage->commit_log_->MarkFinished(transaction_.start_timestamp);
|
|
} else {
|
|
// Validate that existence constraints are satisfied for all modified
|
|
// vertices.
|
|
for (const auto &delta : transaction_.deltas) {
|
|
auto prev = delta.prev.Get();
|
|
MG_ASSERT(prev.type != PreviousPtr::Type::NULLPTR, "Invalid pointer!");
|
|
if (prev.type != PreviousPtr::Type::VERTEX) {
|
|
continue;
|
|
}
|
|
// No need to take any locks here because we modified this vertex and no
|
|
// one else can touch it until we commit.
|
|
auto validation_result = ValidateExistenceConstraints(*prev.vertex, disk_storage->constraints_);
|
|
if (validation_result) {
|
|
Abort();
|
|
return StorageDataManipulationError{*validation_result};
|
|
}
|
|
}
|
|
|
|
// Result of validating the vertex against unqiue constraints. It has to be
|
|
// declared outside of the critical section scope because its value is
|
|
// tested for Abort call which has to be done out of the scope.
|
|
std::optional<ConstraintViolation> unique_constraint_violation;
|
|
|
|
// Save these so we can mark them used in the commit log.
|
|
uint64_t start_timestamp = transaction_.start_timestamp;
|
|
|
|
{
|
|
std::unique_lock<utils::SpinLock> engine_guard(storage_->engine_lock_);
|
|
commit_timestamp_.emplace(disk_storage->CommitTimestamp(desired_commit_timestamp));
|
|
|
|
// Before committing and validating vertices against unique constraints,
|
|
// we have to update unique constraints with the vertices that are going
|
|
// to be validated/committed.
|
|
for (const auto &delta : transaction_.deltas) {
|
|
auto prev = delta.prev.Get();
|
|
MG_ASSERT(prev.type != PreviousPtr::Type::NULLPTR, "Invalid pointer!");
|
|
if (prev.type != PreviousPtr::Type::VERTEX) {
|
|
continue;
|
|
}
|
|
disk_storage->constraints_.unique_constraints.UpdateBeforeCommit(prev.vertex, transaction_);
|
|
}
|
|
|
|
// Validate that unique constraints are satisfied for all modified
|
|
// vertices.
|
|
for (const auto &delta : transaction_.deltas) {
|
|
auto prev = delta.prev.Get();
|
|
MG_ASSERT(prev.type != PreviousPtr::Type::NULLPTR, "Invalid pointer!");
|
|
if (prev.type != PreviousPtr::Type::VERTEX) {
|
|
continue;
|
|
}
|
|
|
|
// No need to take any locks here because we modified this vertex and no
|
|
// one else can touch it until we commit.
|
|
unique_constraint_violation =
|
|
disk_storage->constraints_.unique_constraints.Validate(*prev.vertex, transaction_, *commit_timestamp_);
|
|
if (unique_constraint_violation) {
|
|
break;
|
|
}
|
|
}
|
|
|
|
if (!unique_constraint_violation) {
|
|
// Write transaction to WAL while holding the engine lock to make sure
|
|
// that committed transactions are sorted by the commit timestamp in the
|
|
// WAL files. We supply the new commit timestamp to the function so that
|
|
// it knows what will be the final commit timestamp. The WAL must be
|
|
// written before actually committing the transaction (before setting
|
|
// the commit timestamp) so that no other transaction can see the
|
|
// commits before they are written to disk.
|
|
// Replica can log only the write transaction received from Main
|
|
// so the Wal files are consistent
|
|
// if (storage_->replication_role_ == ReplicationRole::MAIN || desired_commit_timestamp.has_value()) {
|
|
if (desired_commit_timestamp.has_value()) {
|
|
could_replicate_all_sync_replicas =
|
|
disk_storage->AppendToWalDataManipulation(transaction_, *commit_timestamp_);
|
|
}
|
|
|
|
// Take committed_transactions lock while holding the engine lock to
|
|
// make sure that committed transactions are sorted by the commit
|
|
// timestamp in the list.
|
|
disk_storage->committed_transactions_.WithLock([&](auto &committed_transactions) {
|
|
// TODO: release lock, and update all deltas to have a local copy
|
|
// of the commit timestamp
|
|
MG_ASSERT(transaction_.commit_timestamp != nullptr, "Invalid database state!");
|
|
transaction_.commit_timestamp->store(*commit_timestamp_, std::memory_order_release);
|
|
// Replica can only update the last commit timestamp with
|
|
// the commits received from main.
|
|
// if (storage_->replication_role_ == ReplicationRole::MAIN || desired_commit_timestamp.has_value()) {
|
|
if (desired_commit_timestamp.has_value()) {
|
|
// Update the last commit timestamp
|
|
storage_->last_commit_timestamp_.store(*commit_timestamp_);
|
|
}
|
|
// Release engine lock because we don't have to hold it anymore
|
|
// and emplace back could take a long time.
|
|
engine_guard.unlock();
|
|
});
|
|
|
|
disk_storage->commit_log_->MarkFinished(start_timestamp);
|
|
}
|
|
}
|
|
|
|
if (unique_constraint_violation) {
|
|
Abort();
|
|
return StorageDataManipulationError{*unique_constraint_violation};
|
|
}
|
|
FlushCache();
|
|
}
|
|
is_transaction_active_ = false;
|
|
|
|
if (!could_replicate_all_sync_replicas) {
|
|
return StorageDataManipulationError{ReplicationError{}};
|
|
}
|
|
|
|
return {};
|
|
}
|
|
|
|
/// TODO(andi): I think we should have one Abort method per storage.
|
|
void DiskStorage::DiskAccessor::Abort() {
|
|
MG_ASSERT(is_transaction_active_, "The transaction is already terminated!");
|
|
|
|
// We collect vertices and edges we've created here and then splice them into
|
|
// `deleted_vertices_` and `deleted_edges_` lists, instead of adding them one
|
|
// by one and acquiring lock every time.
|
|
std::list<Gid> my_deleted_vertices;
|
|
std::list<Gid> my_deleted_edges;
|
|
|
|
for (const auto &delta : transaction_.deltas) {
|
|
auto prev = delta.prev.Get();
|
|
switch (prev.type) {
|
|
case PreviousPtr::Type::VERTEX: {
|
|
auto vertex = prev.vertex;
|
|
std::lock_guard<utils::SpinLock> guard(vertex->lock);
|
|
Delta *current = vertex->delta;
|
|
while (current != nullptr && current->timestamp->load(std::memory_order_acquire) ==
|
|
transaction_.transaction_id.load(std::memory_order_acquire)) {
|
|
switch (current->action) {
|
|
case Delta::Action::REMOVE_LABEL: {
|
|
auto it = std::find(vertex->labels.begin(), vertex->labels.end(), current->label);
|
|
MG_ASSERT(it != vertex->labels.end(), "Invalid database state!");
|
|
std::swap(*it, *vertex->labels.rbegin());
|
|
vertex->labels.pop_back();
|
|
break;
|
|
}
|
|
case Delta::Action::ADD_LABEL: {
|
|
auto it = std::find(vertex->labels.begin(), vertex->labels.end(), current->label);
|
|
MG_ASSERT(it == vertex->labels.end(), "Invalid database state!");
|
|
vertex->labels.push_back(current->label);
|
|
break;
|
|
}
|
|
case Delta::Action::SET_PROPERTY: {
|
|
vertex->properties.SetProperty(current->property.key, current->property.value);
|
|
break;
|
|
}
|
|
case Delta::Action::ADD_IN_EDGE: {
|
|
std::tuple<EdgeTypeId, Vertex *, EdgeRef> link{current->vertex_edge.edge_type,
|
|
current->vertex_edge.vertex, current->vertex_edge.edge};
|
|
auto it = std::find(vertex->in_edges.begin(), vertex->in_edges.end(), link);
|
|
MG_ASSERT(it == vertex->in_edges.end(), "Invalid database state!");
|
|
vertex->in_edges.push_back(link);
|
|
break;
|
|
}
|
|
case Delta::Action::ADD_OUT_EDGE: {
|
|
std::tuple<EdgeTypeId, Vertex *, EdgeRef> link{current->vertex_edge.edge_type,
|
|
current->vertex_edge.vertex, current->vertex_edge.edge};
|
|
auto it = std::find(vertex->out_edges.begin(), vertex->out_edges.end(), link);
|
|
MG_ASSERT(it == vertex->out_edges.end(), "Invalid database state!");
|
|
vertex->out_edges.push_back(link);
|
|
// Increment edge count. We only increment the count here because
|
|
// the information in `ADD_IN_EDGE` and `Edge/RECREATE_OBJECT` is
|
|
// redundant. Also, `Edge/RECREATE_OBJECT` isn't available when
|
|
// edge properties are disabled.
|
|
storage_->edge_count_.fetch_add(1, std::memory_order_acq_rel);
|
|
break;
|
|
}
|
|
case Delta::Action::REMOVE_IN_EDGE: {
|
|
std::tuple<EdgeTypeId, Vertex *, EdgeRef> link{current->vertex_edge.edge_type,
|
|
current->vertex_edge.vertex, current->vertex_edge.edge};
|
|
auto it = std::find(vertex->in_edges.begin(), vertex->in_edges.end(), link);
|
|
MG_ASSERT(it != vertex->in_edges.end(), "Invalid database state!");
|
|
std::swap(*it, *vertex->in_edges.rbegin());
|
|
vertex->in_edges.pop_back();
|
|
break;
|
|
}
|
|
case Delta::Action::REMOVE_OUT_EDGE: {
|
|
std::tuple<EdgeTypeId, Vertex *, EdgeRef> link{current->vertex_edge.edge_type,
|
|
current->vertex_edge.vertex, current->vertex_edge.edge};
|
|
auto it = std::find(vertex->out_edges.begin(), vertex->out_edges.end(), link);
|
|
MG_ASSERT(it != vertex->out_edges.end(), "Invalid database state!");
|
|
std::swap(*it, *vertex->out_edges.rbegin());
|
|
vertex->out_edges.pop_back();
|
|
// Decrement edge count. We only decrement the count here because
|
|
// the information in `REMOVE_IN_EDGE` and `Edge/DELETE_OBJECT` is
|
|
// redundant. Also, `Edge/DELETE_OBJECT` isn't available when edge
|
|
// properties are disabled.
|
|
storage_->edge_count_.fetch_add(-1, std::memory_order_acq_rel);
|
|
break;
|
|
}
|
|
case Delta::Action::DELETE_DESERIALIZED_OBJECT:
|
|
case Delta::Action::DELETE_OBJECT: {
|
|
vertex->deleted = true;
|
|
my_deleted_vertices.push_back(vertex->gid);
|
|
break;
|
|
}
|
|
case Delta::Action::RECREATE_OBJECT: {
|
|
vertex->deleted = false;
|
|
break;
|
|
}
|
|
}
|
|
current = current->next.load(std::memory_order_acquire);
|
|
}
|
|
vertex->delta = current;
|
|
if (current != nullptr) {
|
|
current->prev.Set(vertex);
|
|
}
|
|
|
|
break;
|
|
}
|
|
case PreviousPtr::Type::EDGE: {
|
|
auto edge = prev.edge;
|
|
std::lock_guard<utils::SpinLock> guard(edge->lock);
|
|
Delta *current = edge->delta;
|
|
while (current != nullptr && current->timestamp->load(std::memory_order_acquire) ==
|
|
transaction_.transaction_id.load(std::memory_order_acquire)) {
|
|
switch (current->action) {
|
|
case Delta::Action::SET_PROPERTY: {
|
|
edge->properties.SetProperty(current->property.key, current->property.value);
|
|
break;
|
|
}
|
|
case Delta::Action::DELETE_DESERIALIZED_OBJECT:
|
|
case Delta::Action::DELETE_OBJECT: {
|
|
edge->deleted = true;
|
|
my_deleted_edges.push_back(edge->gid);
|
|
break;
|
|
}
|
|
case Delta::Action::RECREATE_OBJECT: {
|
|
edge->deleted = false;
|
|
break;
|
|
}
|
|
case Delta::Action::REMOVE_LABEL:
|
|
case Delta::Action::ADD_LABEL:
|
|
case Delta::Action::ADD_IN_EDGE:
|
|
case Delta::Action::ADD_OUT_EDGE:
|
|
case Delta::Action::REMOVE_IN_EDGE:
|
|
case Delta::Action::REMOVE_OUT_EDGE: {
|
|
LOG_FATAL("Invalid database state!");
|
|
break;
|
|
}
|
|
}
|
|
current = current->next.load(std::memory_order_acquire);
|
|
}
|
|
edge->delta = current;
|
|
if (current != nullptr) {
|
|
current->prev.Set(edge);
|
|
}
|
|
|
|
break;
|
|
}
|
|
case PreviousPtr::Type::DELTA:
|
|
// pointer probably couldn't be set because allocation failed
|
|
case PreviousPtr::Type::NULLPTR:
|
|
break;
|
|
}
|
|
}
|
|
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
{
|
|
std::unique_lock<utils::SpinLock> engine_guard(storage_->engine_lock_);
|
|
uint64_t mark_timestamp = storage_->timestamp_;
|
|
// Take garbage_undo_buffers lock while holding the engine lock to make
|
|
// sure that entries are sorted by mark timestamp in the list.
|
|
disk_storage->garbage_undo_buffers_.WithLock([&](auto &garbage_undo_buffers) {
|
|
// Release engine lock because we don't have to hold it anymore and
|
|
// emplace back could take a long time.
|
|
engine_guard.unlock();
|
|
garbage_undo_buffers.emplace_back(mark_timestamp, std::move(transaction_.deltas));
|
|
});
|
|
disk_storage->deleted_vertices_.WithLock(
|
|
[&](auto &deleted_vertices) { deleted_vertices.splice(deleted_vertices.begin(), my_deleted_vertices); });
|
|
disk_storage->deleted_edges_.WithLock(
|
|
[&](auto &deleted_edges) { deleted_edges.splice(deleted_edges.begin(), my_deleted_edges); });
|
|
}
|
|
|
|
disk_storage->commit_log_->MarkFinished(transaction_.start_timestamp);
|
|
is_transaction_active_ = false;
|
|
}
|
|
|
|
// maybe will need some usages here
|
|
void DiskStorage::DiskAccessor::FinalizeTransaction() {
|
|
auto *disk_storage = static_cast<DiskStorage *>(storage_);
|
|
if (commit_timestamp_) {
|
|
disk_storage->commit_log_->MarkFinished(*commit_timestamp_);
|
|
disk_storage->committed_transactions_.WithLock(
|
|
[&](auto &committed_transactions) { committed_transactions.emplace_back(std::move(transaction_)); });
|
|
commit_timestamp_.reset();
|
|
}
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
std::optional<uint64_t> DiskStorage::DiskAccessor::GetTransactionId() const {
|
|
if (is_transaction_active_) {
|
|
return transaction_.transaction_id.load(std::memory_order_acquire);
|
|
}
|
|
return {};
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
const std::string &DiskStorage::LabelToName(LabelId label) const { return name_id_mapper_.IdToName(label.AsUint()); }
|
|
|
|
// this should be handled on an above level of abstraction
|
|
const std::string &DiskStorage::PropertyToName(PropertyId property) const {
|
|
return name_id_mapper_.IdToName(property.AsUint());
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
const std::string &DiskStorage::EdgeTypeToName(EdgeTypeId edge_type) const {
|
|
return name_id_mapper_.IdToName(edge_type.AsUint());
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
LabelId DiskStorage::NameToLabel(const std::string_view name) {
|
|
return LabelId::FromUint(name_id_mapper_.NameToId(name));
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
PropertyId DiskStorage::NameToProperty(const std::string_view name) {
|
|
return PropertyId::FromUint(name_id_mapper_.NameToId(name));
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
EdgeTypeId DiskStorage::NameToEdgeType(const std::string_view name) {
|
|
return EdgeTypeId::FromUint(name_id_mapper_.NameToId(name));
|
|
}
|
|
|
|
utils::BasicResult<StorageIndexDefinitionError, void> DiskStorage::CreateIndex(
|
|
LabelId label, const std::optional<uint64_t> desired_commit_timestamp) {
|
|
/// TODO: (andi): Here we will probably use some lock to protect reading from and writing to the RocksDB at the
|
|
// same
|
|
/// time.
|
|
// Since indexes are handled by the storage, we don't have transaction id for indexes methods.
|
|
// However, we need to set timestamp because of RocksDB requirement so we set it to the maximum value.
|
|
// RocksDB will still fetch the most recent value.
|
|
rocksdb::ReadOptions ro;
|
|
rocksdb::Slice ts = utils::StringTimestamp(std::numeric_limits<uint64_t>::max());
|
|
ro.timestamp = &ts;
|
|
auto it = std::unique_ptr<rocksdb::Iterator>(kvstore_->db_->NewIterator(ro, kvstore_->vertex_chandle));
|
|
std::vector<std::tuple<std::string, std::string, uint64_t>> indexed_vertices;
|
|
for (it->SeekToFirst(); it->Valid(); it->Next()) {
|
|
const auto &key = it->key();
|
|
const auto vertex_parts = utils::Split(key.ToStringView(), "|");
|
|
if (const auto labels = utils::Split(vertex_parts[0], ",");
|
|
// TODO: (andi): When you decouple utils::SerializeIdType, modify this to_string call to use
|
|
// utils::SerializeIdType.
|
|
std::find(labels.begin(), labels.end(), std::to_string(label.AsUint())) != labels.end()) {
|
|
spdlog::debug("Found vertex with gid {} for index creation", vertex_parts[1]);
|
|
indexed_vertices.emplace_back(key.ToString(), it->value().ToString(),
|
|
utils::ExtractTimestampFromDeserializedUserKey(key));
|
|
}
|
|
}
|
|
|
|
// std::unique_lock<utils::RWLock> storage_guard(main_lock_);
|
|
// if (!indices_.label_index.CreateIndex(label, indexed_vertices)) {
|
|
// return StorageIndexDefinitionError{IndexDefinitionError{}};
|
|
// }
|
|
// const auto commit_timestamp = CommitTimestamp(desired_commit_timestamp);
|
|
// const auto success =
|
|
// AppendToWalDataDefinition(durability::StorageGlobalOperation::LABEL_INDEX_CREATE, label, {}, commit_timestamp);
|
|
// commit_log_->MarkFinished(commit_timestamp);
|
|
// last_commit_timestamp_ = commit_timestamp;
|
|
|
|
// if (success) {
|
|
// return {};
|
|
// }
|
|
|
|
return StorageIndexDefinitionError{ReplicationError{}};
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
utils::BasicResult<StorageIndexDefinitionError, void> DiskStorage::CreateIndex(
|
|
LabelId label, PropertyId property, const std::optional<uint64_t> desired_commit_timestamp) {
|
|
throw utils::NotYetImplemented("CreateIndex");
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
utils::BasicResult<StorageIndexDefinitionError, void> DiskStorage::DropIndex(
|
|
LabelId label, const std::optional<uint64_t> desired_commit_timestamp) {
|
|
throw utils::NotYetImplemented("DropIndex");
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
utils::BasicResult<StorageIndexDefinitionError, void> DiskStorage::DropIndex(
|
|
LabelId label, PropertyId property, const std::optional<uint64_t> desired_commit_timestamp) {
|
|
throw utils::NotYetImplemented("DropIndex");
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
IndicesInfo DiskStorage::ListAllIndices() const { throw utils::NotYetImplemented("ListAllIndices"); }
|
|
|
|
// this should be handled on an above level of abstraction
|
|
utils::BasicResult<StorageExistenceConstraintDefinitionError, void> DiskStorage::CreateExistenceConstraint(
|
|
LabelId label, PropertyId property, const std::optional<uint64_t> desired_commit_timestamp) {
|
|
throw utils::NotYetImplemented("CreateExistenceConstraint");
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
utils::BasicResult<StorageExistenceConstraintDroppingError, void> DiskStorage::DropExistenceConstraint(
|
|
LabelId label, PropertyId property, const std::optional<uint64_t> desired_commit_timestamp) {
|
|
throw utils::NotYetImplemented("DropExistenceConstraint");
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
utils::BasicResult<StorageUniqueConstraintDefinitionError, UniqueConstraints::CreationStatus>
|
|
DiskStorage::CreateUniqueConstraint(LabelId label, const std::set<PropertyId> &properties,
|
|
const std::optional<uint64_t> desired_commit_timestamp) {
|
|
throw utils::NotYetImplemented("CreateUniqueConstraint");
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
utils::BasicResult<StorageUniqueConstraintDroppingError, UniqueConstraints::DeletionStatus>
|
|
DiskStorage::DropUniqueConstraint(LabelId label, const std::set<PropertyId> &properties,
|
|
const std::optional<uint64_t> desired_commit_timestamp) {
|
|
throw utils::NotYetImplemented("DropUniqueConstraint");
|
|
}
|
|
|
|
// this should be handled on an above level of abstraction
|
|
ConstraintsInfo DiskStorage::ListAllConstraints() const { throw utils::NotYetImplemented("ListAllConstraints"); }
|
|
|
|
// this should be handled on an above level of abstraction
|
|
StorageInfo DiskStorage::GetInfo() const { throw utils::NotYetImplemented("GetInfo"); }
|
|
|
|
Transaction DiskStorage::CreateTransaction(IsolationLevel isolation_level, StorageMode storage_mode) {
|
|
/// We acquire the transaction engine lock here because we access (and
|
|
/// modify) the transaction engine variables (`transaction_id` and
|
|
/// `timestamp`) below.
|
|
uint64_t transaction_id;
|
|
uint64_t start_timestamp;
|
|
{
|
|
std::lock_guard<utils::SpinLock> guard(engine_lock_);
|
|
transaction_id = transaction_id_++;
|
|
/// TODO: when we introduce replication to the disk storage, take care of start_timestamp
|
|
start_timestamp = timestamp_++;
|
|
}
|
|
return {transaction_id, start_timestamp, isolation_level, storage_mode};
|
|
}
|
|
|
|
template <bool force>
|
|
void DiskStorage::CollectGarbage() {
|
|
if constexpr (force) {
|
|
// We take the unique lock on the main storage lock so we can forcefully clean
|
|
// everything we can
|
|
if (!main_lock_.try_lock()) {
|
|
CollectGarbage<false>();
|
|
return;
|
|
}
|
|
} else {
|
|
// Because the garbage collector iterates through the indices and constraints
|
|
// to clean them up, it must take the main lock for reading to make sure that
|
|
// the indices and constraints aren't concurrently being modified.
|
|
main_lock_.lock_shared();
|
|
}
|
|
|
|
utils::OnScopeExit lock_releaser{[&] {
|
|
if constexpr (force) {
|
|
main_lock_.unlock();
|
|
} else {
|
|
main_lock_.unlock_shared();
|
|
}
|
|
}};
|
|
|
|
// Garbage collection must be performed in two phases. In the first phase,
|
|
// deltas that won't be applied by any transaction anymore are unlinked from
|
|
// the version chains. They cannot be deleted immediately, because there
|
|
// might be a transaction that still needs them to terminate the version
|
|
// chain traversal. They are instead marked for deletion and will be deleted
|
|
// in the second GC phase in this GC iteration or some of the following
|
|
// ones.
|
|
std::unique_lock<std::mutex> gc_guard(gc_lock_, std::try_to_lock);
|
|
if (!gc_guard.owns_lock()) {
|
|
return;
|
|
}
|
|
|
|
uint64_t oldest_active_start_timestamp = commit_log_->OldestActive();
|
|
// We don't move undo buffers of unlinked transactions to garbage_undo_buffers
|
|
// list immediately, because we would have to repeatedly take
|
|
// garbage_undo_buffers lock.
|
|
std::list<std::pair<uint64_t, std::list<Delta>>> unlinked_undo_buffers;
|
|
|
|
// We will only free vertices deleted up until now in this GC cycle, and we
|
|
// will do it after cleaning-up the indices. That way we are sure that all
|
|
// vertices that appear in an index also exist in main storage.
|
|
std::list<Gid> current_deleted_edges;
|
|
std::list<Gid> current_deleted_vertices;
|
|
deleted_vertices_->swap(current_deleted_vertices);
|
|
deleted_edges_->swap(current_deleted_edges);
|
|
|
|
// Flag that will be used to determine whether the Index GC should be run. It
|
|
// should be run when there were any items that were cleaned up (there were
|
|
// updates between this run of the GC and the previous run of the GC). This
|
|
// eliminates high CPU usage when the GC doesn't have to clean up anything.
|
|
bool run_index_cleanup = !committed_transactions_->empty() || !garbage_undo_buffers_->empty();
|
|
|
|
while (true) {
|
|
// We don't want to hold the lock on commited transactions for too long,
|
|
// because that prevents other transactions from committing.
|
|
Transaction *transaction;
|
|
{
|
|
auto committed_transactions_ptr = committed_transactions_.Lock();
|
|
if (committed_transactions_ptr->empty()) {
|
|
break;
|
|
}
|
|
transaction = &committed_transactions_ptr->front();
|
|
}
|
|
|
|
auto commit_timestamp = transaction->commit_timestamp->load(std::memory_order_acquire);
|
|
if (commit_timestamp >= oldest_active_start_timestamp) {
|
|
break;
|
|
}
|
|
|
|
// When unlinking a delta which is the first delta in its version chain,
|
|
// special care has to be taken to avoid the following race condition:
|
|
//
|
|
// [Vertex] --> [Delta A]
|
|
//
|
|
// GC thread: Delta A is the first in its chain, it must be unlinked from
|
|
// vertex and marked for deletion
|
|
// TX thread: Update vertex and add Delta B with Delta A as next
|
|
//
|
|
// [Vertex] --> [Delta B] <--> [Delta A]
|
|
//
|
|
// GC thread: Unlink delta from Vertex
|
|
//
|
|
// [Vertex] --> (nullptr)
|
|
//
|
|
// When processing a delta that is the first one in its chain, we
|
|
// obtain the corresponding vertex or edge lock, and then verify that this
|
|
// delta still is the first in its chain.
|
|
// When processing a delta that is in the middle of the chain we only
|
|
// process the final delta of the given transaction in that chain. We
|
|
// determine the owner of the chain (either a vertex or an edge), obtain the
|
|
// corresponding lock, and then verify that this delta is still in the same
|
|
// position as it was before taking the lock.
|
|
//
|
|
// Even though the delta chain is lock-free (both `next` and `prev`) the
|
|
// chain should not be modified without taking the lock from the object that
|
|
// owns the chain (either a vertex or an edge). Modifying the chain without
|
|
// taking the lock will cause subtle race conditions that will leave the
|
|
// chain in a broken state.
|
|
// The chain can be only read without taking any locks.
|
|
|
|
for (Delta &delta : transaction->deltas) {
|
|
while (true) {
|
|
auto prev = delta.prev.Get();
|
|
switch (prev.type) {
|
|
case PreviousPtr::Type::VERTEX: {
|
|
Vertex *vertex = prev.vertex;
|
|
std::lock_guard<utils::SpinLock> vertex_guard(vertex->lock);
|
|
if (vertex->delta != &delta) {
|
|
// Something changed, we're not the first delta in the chain
|
|
// anymore.
|
|
continue;
|
|
}
|
|
vertex->delta = nullptr;
|
|
if (vertex->deleted) {
|
|
current_deleted_vertices.push_back(vertex->gid);
|
|
}
|
|
break;
|
|
}
|
|
case PreviousPtr::Type::EDGE: {
|
|
Edge *edge = prev.edge;
|
|
std::lock_guard<utils::SpinLock> edge_guard(edge->lock);
|
|
if (edge->delta != &delta) {
|
|
// Something changed, we're not the first delta in the chain
|
|
// anymore.
|
|
continue;
|
|
}
|
|
edge->delta = nullptr;
|
|
if (edge->deleted) {
|
|
current_deleted_edges.push_back(edge->gid);
|
|
}
|
|
break;
|
|
}
|
|
case PreviousPtr::Type::DELTA: {
|
|
if (prev.delta->timestamp->load(std::memory_order_acquire) == commit_timestamp) {
|
|
// The delta that is newer than this one is also a delta from this
|
|
// transaction. We skip the current delta and will remove it as a
|
|
// part of the suffix later.
|
|
break;
|
|
}
|
|
std::unique_lock<utils::SpinLock> guard;
|
|
{
|
|
// We need to find the parent object in order to be able to use
|
|
// its lock.
|
|
auto parent = prev;
|
|
while (parent.type == PreviousPtr::Type::DELTA) {
|
|
parent = parent.delta->prev.Get();
|
|
}
|
|
switch (parent.type) {
|
|
case PreviousPtr::Type::VERTEX:
|
|
guard = std::unique_lock<utils::SpinLock>(parent.vertex->lock);
|
|
break;
|
|
case PreviousPtr::Type::EDGE:
|
|
guard = std::unique_lock<utils::SpinLock>(parent.edge->lock);
|
|
break;
|
|
case PreviousPtr::Type::DELTA:
|
|
case PreviousPtr::Type::NULLPTR:
|
|
LOG_FATAL("Invalid database state!");
|
|
}
|
|
}
|
|
if (delta.prev.Get() != prev) {
|
|
// Something changed, we could now be the first delta in the
|
|
// chain.
|
|
continue;
|
|
}
|
|
Delta *prev_delta = prev.delta;
|
|
prev_delta->next.store(nullptr, std::memory_order_release);
|
|
break;
|
|
}
|
|
case PreviousPtr::Type::NULLPTR: {
|
|
LOG_FATAL("Invalid pointer!");
|
|
}
|
|
}
|
|
break;
|
|
}
|
|
}
|
|
|
|
committed_transactions_.WithLock([&](auto &committed_transactions) {
|
|
unlinked_undo_buffers.emplace_back(0, std::move(transaction->deltas));
|
|
committed_transactions.pop_front();
|
|
});
|
|
}
|
|
|
|
// After unlinking deltas from vertices, we refresh the indices. That way
|
|
// we're sure that none of the vertices from `current_deleted_vertices`
|
|
// appears in an index, and we can safely remove the from the main storage
|
|
// after the last currently active transaction is finished.
|
|
if (run_index_cleanup) {
|
|
// This operation is very expensive as it traverses through all of the items
|
|
// in every index every time.
|
|
// RemoveObsoleteEntries(&indices_, oldest_active_start_timestamp);
|
|
constraints_.unique_constraints.RemoveObsoleteEntries(oldest_active_start_timestamp);
|
|
}
|
|
|
|
{
|
|
std::unique_lock<utils::SpinLock> guard(engine_lock_);
|
|
uint64_t mark_timestamp = timestamp_;
|
|
// Take garbage_undo_buffers lock while holding the engine lock to make
|
|
// sure that entries are sorted by mark timestamp in the list.
|
|
garbage_undo_buffers_.WithLock([&](auto &garbage_undo_buffers) {
|
|
// Release engine lock because we don't have to hold it anymore and
|
|
// this could take a long time.
|
|
guard.unlock();
|
|
// TODO(mtomic): holding garbage_undo_buffers_ lock here prevents
|
|
// transactions from aborting until we're done marking, maybe we should
|
|
// add them one-by-one or something
|
|
for (auto &[timestamp, undo_buffer] : unlinked_undo_buffers) {
|
|
timestamp = mark_timestamp;
|
|
}
|
|
garbage_undo_buffers.splice(garbage_undo_buffers.end(), unlinked_undo_buffers);
|
|
});
|
|
for (auto vertex : current_deleted_vertices) {
|
|
garbage_vertices_.emplace_back(mark_timestamp, vertex);
|
|
}
|
|
}
|
|
|
|
garbage_undo_buffers_.WithLock([&](auto &undo_buffers) {
|
|
// if force is set to true we can simply delete all the leftover undos because
|
|
// no transaction is active
|
|
if constexpr (force) {
|
|
undo_buffers.clear();
|
|
} else {
|
|
while (!undo_buffers.empty() && undo_buffers.front().first <= oldest_active_start_timestamp) {
|
|
undo_buffers.pop_front();
|
|
}
|
|
}
|
|
});
|
|
|
|
{
|
|
auto vertex_acc = vertices_.access();
|
|
if constexpr (force) {
|
|
// if force is set to true, then we have unique_lock and no transactions are active
|
|
// so we can clean all of the deleted vertices
|
|
while (!garbage_vertices_.empty()) {
|
|
MG_ASSERT(vertex_acc.remove(garbage_vertices_.front().second), "Invalid database state!");
|
|
garbage_vertices_.pop_front();
|
|
}
|
|
} else {
|
|
while (!garbage_vertices_.empty() && garbage_vertices_.front().first < oldest_active_start_timestamp) {
|
|
MG_ASSERT(vertex_acc.remove(garbage_vertices_.front().second), "Invalid database state!");
|
|
garbage_vertices_.pop_front();
|
|
}
|
|
}
|
|
}
|
|
|
|
{
|
|
auto edge_acc = edges_.access();
|
|
for (auto edge : current_deleted_edges) {
|
|
MG_ASSERT(edge_acc.remove(edge), "Invalid database state!");
|
|
}
|
|
}
|
|
|
|
// TODO(andi): Remove this after we are assured that deserialization works
|
|
edges_.clear();
|
|
vertices_.clear();
|
|
spdlog::debug("Cleared caches");
|
|
}
|
|
|
|
// tell the linker he can find the CollectGarbage definitions here
|
|
template void DiskStorage::CollectGarbage<true>();
|
|
template void DiskStorage::CollectGarbage<false>();
|
|
|
|
bool DiskStorage::InitializeWalFile() {
|
|
if (config_.durability.snapshot_wal_mode != Config::Durability::SnapshotWalMode::PERIODIC_SNAPSHOT_WITH_WAL)
|
|
return false;
|
|
if (!wal_file_) {
|
|
wal_file_.emplace(wal_directory_, uuid_, epoch_id_, config_.items, &name_id_mapper_, wal_seq_num_++,
|
|
&file_retainer_);
|
|
}
|
|
return true;
|
|
}
|
|
|
|
void DiskStorage::FinalizeWalFile() {
|
|
++wal_unsynced_transactions_;
|
|
if (wal_unsynced_transactions_ >= config_.durability.wal_file_flush_every_n_tx) {
|
|
wal_file_->Sync();
|
|
wal_unsynced_transactions_ = 0;
|
|
}
|
|
if (wal_file_->GetSize() / 1024 >= config_.durability.wal_file_size_kibibytes) {
|
|
wal_file_->FinalizeWal();
|
|
wal_file_ = std::nullopt;
|
|
wal_unsynced_transactions_ = 0;
|
|
} else {
|
|
// Try writing the internal buffer if possible, if not
|
|
// the data should be written as soon as it's possible
|
|
// (triggered by the new transaction commit, or some
|
|
// reading thread EnabledFlushing)
|
|
wal_file_->TryFlushing();
|
|
}
|
|
}
|
|
|
|
bool DiskStorage::AppendToWalDataManipulation(const Transaction &transaction, uint64_t final_commit_timestamp) {
|
|
throw utils::NotYetImplemented("AppendToWalDataManipulation");
|
|
}
|
|
|
|
bool DiskStorage::AppendToWalDataDefinition(durability::StorageGlobalOperation operation, LabelId label,
|
|
const std::set<PropertyId> &properties, uint64_t final_commit_timestamp) {
|
|
throw utils::NotYetImplemented("AppendToWalDataDefinition");
|
|
}
|
|
|
|
void DiskStorage::SetStorageMode(StorageMode storage_mode) {
|
|
std::unique_lock main_guard{main_lock_};
|
|
storage_mode_ = storage_mode;
|
|
}
|
|
|
|
StorageMode DiskStorage::GetStorageMode() { return storage_mode_; }
|
|
|
|
utils::BasicResult<DiskStorage::CreateSnapshotError> DiskStorage::CreateSnapshot(std::optional<bool> is_periodic) {
|
|
throw utils::NotYetImplemented("CreateSnapshot");
|
|
}
|
|
|
|
bool DiskStorage::LockPath() { throw utils::NotYetImplemented("LockPath"); }
|
|
|
|
bool DiskStorage::UnlockPath() { throw utils::NotYetImplemented("UnlockPath"); }
|
|
|
|
void DiskStorage::FreeMemory() { throw utils::NotYetImplemented("FreeMemory"); }
|
|
|
|
uint64_t DiskStorage::CommitTimestamp(const std::optional<uint64_t> desired_commit_timestamp) {
|
|
if (!desired_commit_timestamp) {
|
|
return timestamp_++;
|
|
}
|
|
timestamp_ = std::max(timestamp_, *desired_commit_timestamp + 1);
|
|
return *desired_commit_timestamp;
|
|
}
|
|
|
|
bool DiskStorage::SetReplicaRole(io::network::Endpoint endpoint, const replication::ReplicationServerConfig &config) {
|
|
throw utils::NotYetImplemented("SetReplicaRole");
|
|
}
|
|
|
|
bool DiskStorage::SetMainReplicationRole() { throw utils::NotYetImplemented("SetMainReplicationRole"); }
|
|
|
|
utils::BasicResult<DiskStorage::RegisterReplicaError> DiskStorage::RegisterReplica(
|
|
std::string name, io::network::Endpoint endpoint, const replication::ReplicationMode replication_mode,
|
|
const replication::RegistrationMode registration_mode, const replication::ReplicationClientConfig &config) {
|
|
throw utils::NotYetImplemented("RegisterReplica");
|
|
}
|
|
|
|
bool DiskStorage::UnregisterReplica(const std::string &name) { throw utils::NotYetImplemented("UnregisterReplica"); }
|
|
|
|
std::optional<replication::ReplicaState> DiskStorage::GetReplicaState(const std::string_view name) {
|
|
throw utils::NotYetImplemented("GetReplicaState");
|
|
}
|
|
|
|
ReplicationRole DiskStorage::GetReplicationRole() const {
|
|
/// TODO(andi): Since we won't support replication for the beginning don't crash with this, rather just log a
|
|
/// message
|
|
spdlog::debug("GetReplicationRole called, but we don't support replication yet");
|
|
return {};
|
|
}
|
|
|
|
std::vector<DiskStorage::ReplicaInfo> DiskStorage::ReplicasInfo() { throw utils::NotYetImplemented("ReplicasInfo"); }
|
|
|
|
utils::BasicResult<DiskStorage::SetIsolationLevelError> DiskStorage::SetIsolationLevel(IsolationLevel isolation_level) {
|
|
std::unique_lock main_guard{main_lock_};
|
|
if (storage_mode_ == storage::StorageMode::IN_MEMORY_ANALYTICAL) {
|
|
return Storage::SetIsolationLevelError::DisabledForAnalyticalMode;
|
|
}
|
|
|
|
isolation_level_ = isolation_level;
|
|
return {};
|
|
}
|
|
|
|
void DiskStorage::RestoreReplicas() { throw utils::NotYetImplemented("RestoreReplicas"); }
|
|
|
|
bool DiskStorage::ShouldStoreAndRestoreReplicas() const {
|
|
throw utils::NotYetImplemented("ShouldStoreAndRestoreReplicas");
|
|
}
|
|
|
|
} // namespace memgraph::storage
|