Compare commits

...

61 Commits

Author SHA1 Message Date
gvolfing
451da803dc Add another test to ExpandOne 2022-09-20 18:40:23 +02:00
gvolfing
d1f51aa3f9 Add checkf for edge types in the ExpandOne handle 2022-09-20 18:04:32 +02:00
gvolfing
aec3bc3d05 Fix ExpandOne behavior 2022-09-20 16:39:52 +02:00
gvolfing
7fe5eb2c49 Clean-up tests 2022-09-20 09:52:57 +02:00
gvolfing
26e57d1988 Add the rest of the implementation for ExpandOne handle, fix segfault issue 2022-09-19 19:58:20 +02:00
gvolfing
f066d4f0b8 Save for tomorrow 2022-09-18 13:29:31 +02:00
gvolfing
359ac507a3 Merge branch 'T0919-MG-implement-message-based-actions' into T0919-MG-save-branch-for-outcommented-messages 2022-09-15 19:22:34 +02:00
gvolfing
4530e8fe50 Make ConvertPropertyMap take its argument as rvalue& 2022-09-15 17:18:23 +02:00
gvolfing
ef205dce4d Fix clang-tidy warnings 2022-09-15 13:36:32 +02:00
gvolfing
fd1668403a Restructure EdgeId message struct 2022-09-15 11:57:57 +02:00
gvolfing
b8bca320c8 Make enum types CamelCase, remove dead comments 2022-09-15 11:34:40 +02:00
gvolfing
3140f81d61 Apply suggestions from code review
Co-authored-by: Jure Bajic <jure.bajic@memgraph.com>
2022-09-15 11:07:46 +02:00
gvolfing
59bb54c26e React to PR comments 3 2022-09-14 17:12:43 +02:00
gvolfing
c6bc9d39ab Merge branch 'E118-MG-lexicographically-ordered-storage' into T0919-MG-implement-message-based-actions 2022-09-14 16:51:28 +02:00
gvolfing
d0930f66c8 React to PR comments 2 2022-09-14 16:48:07 +02:00
gvolfing
7db0f9ee23 React to PR comments 1 2022-09-14 14:19:44 +02:00
János Benjamin Antal
3a956d2266 Fix random hangs in simulator 2022-09-13 13:16:26 +02:00
János Benjamin Antal
9284ba0805 Small fixes 2022-09-13 11:09:32 +02:00
gvolfing
4c53713217 Remove outcommented code 2022-09-13 11:08:54 +02:00
gvolfing
ae04e717f9 Enable back the all of the simulator tests 2022-09-13 10:56:14 +02:00
gvolfing
1706be3341 Enable back the rest of the simulator tests 2022-09-13 10:53:46 +02:00
gvolfing
4d53d53859 Clean up used messages and their handles 2022-09-13 10:34:13 +02:00
gvolfing
f5d2bcd979 Cleanup 2022-09-12 18:09:28 +02:00
gvolfing
16028f50e8 Finish impl of ScanAll handle and its test, uncomment non-priority cases, general cleanup 2022-09-12 17:13:54 +02:00
János Benjamin Antal
bfdc69fdbd Fix Value special member functions 2022-09-12 14:06:30 +02:00
János Benjamin Antal
f49c962c5e Fix promise timeout handling 2022-09-12 14:06:02 +02:00
gvolfing
c71f16d9aa Save progress 2022-09-12 13:35:36 +02:00
gvolfing
2751d1fc94 Add test for DeleteEdges message handle 2022-09-11 21:04:29 +02:00
gvolfing
0d728d8a0a Test AddEdges message handle 2022-09-11 20:41:45 +02:00
gvolfing
0cfd0703ab Test UpdateVertices message handle 2022-09-11 19:56:02 +02:00
gvolfing
8ca0c0df53 Add test for DeleteVertices message handle 2022-09-11 15:54:46 +02:00
gvolfing
3feaa79688 Make Test and CreateVertices message handle work 2022-09-11 09:59:58 +02:00
gvolfing
a119fa3f73 Merge branch 'E118-MG-lexicographically-ordered-storage' into T0919-MG-implement-message-based-actions 2022-09-09 13:26:23 +02:00
gvolfing
f50d0f3a75 Add additional constructors to Value and modify test for message handlers 2022-09-09 13:20:13 +02:00
gvolfing
4d650aa595 Add test file 2022-09-08 17:09:53 +02:00
gvolfing
1316a90f8a Add additional dependencies and delete unused headers and type aliases 2022-09-08 16:46:21 +02:00
gvolfing
8bb6e00317 Add files for tests 2022-09-08 14:01:34 +02:00
gvolfing
c690b349f7 Add message routing capabilites to ShardRsm 2022-09-08 13:09:21 +02:00
gvolfing
9c2ecbd0a1 Add definition for ScanVerticesHandle 2022-09-08 12:47:57 +02:00
gvolfing
5b31ddc4fa Add move ctors where needed 2022-09-08 11:09:02 +02:00
gvolfing
a308bf16a2 Uncomment parts that require dinamically allocated resources from the union 2022-09-08 07:53:03 +02:00
gvolfing
3b5a736a97 Add partial implementation of ScanAll 2022-09-07 21:57:00 +02:00
gvolfing
8ebd85b6ea Update messages 2022-09-07 14:03:10 +02:00
gvolfing
db190baeef Add handle declaration for ScanVertices request 2022-09-06 20:27:57 +02:00
gvolfing
a1ec8d19d1 TODO 2022-09-06 15:31:54 +02:00
gvolfing
c576d02974 Add definition for UpdateEdge handle, modify its message 2022-09-06 15:26:14 +02:00
gvolfing
941ef01b88 Add definition for UpdateVertexHandle 2022-09-06 15:26:14 +02:00
gvolfing
6a7b2e586c Add implementation for DeleteEdges handle 2022-09-06 15:26:14 +02:00
gvolfing
10aa6b977c Modify existing messages, add new one for updates, more cleanup 2022-09-06 15:26:14 +02:00
gvolfing
c8ca234fec General cleanup 2022-09-06 15:26:14 +02:00
gvolfing
3e05b690a9 Make handle implementations exit the loops on the first encountered error 2022-09-06 15:26:14 +02:00
gvolfing
316fbfda52 Add implementation to CreateEdge handle, and modify messages 2022-09-06 15:26:14 +02:00
gvolfing
81b1173336 Add message types related to edges and the declarations for their handles 2022-09-06 15:26:14 +02:00
gvolfing
88240a39d9 Modify DeleteVertices message type, add implementation for its handle 2022-09-06 15:26:14 +02:00
gvolfing
31581ce4aa Refactor bridging helper functions 2022-09-06 15:26:14 +02:00
gvolfing
df60eae124 Change the underlying datastructure of property map in the defined messages, and bridge it to the VertexAccessor 2022-09-06 15:26:14 +02:00
gvolfing
950338833e Add debug messages and refactor the CreateVertices handle 2022-09-06 15:26:14 +02:00
gvolfing
921dc8b740 Add other message types relating to vertices 2022-09-06 15:26:14 +02:00
gvolfing
02f63fd68d Rename class ShardMessageHandler to ShardRsm 2022-09-06 15:26:14 +02:00
gvolfing
b1feee5b08 Rename file shard_message_handler to shard_rsm 2022-09-06 15:26:14 +02:00
gvolfing
d7c89d9cd5 Add CreateVerticies skeleton 2022-09-06 15:26:14 +02:00
10 changed files with 2243 additions and 51 deletions

View File

@@ -55,15 +55,17 @@ class SimulatorHandle {
void TimeoutPromisesPastDeadline() {
const Time now = cluster_wide_time_microseconds_;
for (auto &[promise_key, dop] : promises_) {
for (auto it = promises_.begin(); it != promises_.end();) {
auto &[promise_key, dop] = *it;
if (dop.deadline < now) {
spdlog::debug("timing out request from requester {} to replier {}.", promise_key.requester_address.ToString(),
promise_key.replier_address.ToString());
spdlog::info("timing out request from requester {} to replier {}.", promise_key.requester_address.ToString(),
promise_key.replier_address.ToString());
std::move(dop).promise.TimeOut();
promises_.erase(promise_key);
it = promises_.erase(it);
stats_.timed_out_requests++;
} else {
++it;
}
}
}

View File

@@ -67,7 +67,7 @@ class Io {
I implementation_;
Address address_;
RequestId request_id_counter_ = 0;
Duration default_timeout_ = std::chrono::microseconds{50000};
Duration default_timeout_ = std::chrono::microseconds{100000};
public:
Io(I io, Address address) : implementation_(io), address_(address) {}

View File

@@ -16,52 +16,66 @@
#include <map>
#include <optional>
#include <unordered_map>
#include <utility>
#include <variant>
#include <vector>
#include "coordinator/hybrid_logical_clock.hpp"
#include "storage/v3/id_types.hpp"
#include "storage/v3/property_value.hpp"
/// Hybrid-logical clock
struct Hlc {
uint64_t logical_id;
using Duration = std::chrono::microseconds;
using Time = std::chrono::time_point<std::chrono::system_clock, Duration>;
Time coordinator_wall_clock;
namespace memgraph::msgs {
bool operator==(const Hlc &other) const = default;
};
// TODO(gvolfing) introduce structure for property updates
// TODO(gvolfing) remove tuples from ExpandOne and add specific structs instead
// TODO(gvolfing) discuss with kostas:
// - EdgeId
// - Everything about ExpandOne
// - Edge properties
// - Ctors
using coordinator::Hlc;
using storage::v3::LabelId;
struct Value;
struct Label {
size_t id;
LabelId id;
};
// TODO(kostasrim) update this with CompoundKey, same for the rest of the file.
using PrimaryKey = std::vector<memgraph::storage::v3::PropertyValue>;
using PrimaryKey = std::vector<Value>;
using VertexId = std::pair<Label, PrimaryKey>;
using Gid = size_t;
using PropertyId = memgraph::storage::v3::PropertyId;
// struct PropertyUpdate{
// PropertyId id;
// Value value;
// };
struct EdgeType {
std::string name;
uint64_t id;
};
struct EdgeId {
VertexId id;
Gid gid;
};
struct Vertex {
VertexId id;
std::vector<Label> labels;
};
struct Edge {
VertexId src;
VertexId dst;
std::optional<std::vector<std::pair<PropertyId, Value>>> properties;
EdgeId id;
EdgeType type;
};
struct Vertex {
VertexId id;
std::vector<Label> labels;
};
struct PathPart {
Vertex dst;
Gid edge;
@@ -75,11 +89,224 @@ struct Path {
struct Null {};
struct Value {
enum Type { NILL, BOOL, INT64, DOUBLE, STRING, LIST, MAP, VERTEX, EDGE, PATH };
Value() : null_v{} {}
explicit Value(const bool val) : type(Type::Bool), bool_v(val) {}
explicit Value(const int64_t val) : type(Type::Int64), int_v(val) {}
explicit Value(const double val) : type(Type::Double), double_v(val) {}
explicit Value(const Vertex val) : type(Type::Vertex), vertex_v(val) {}
explicit Value(const std::string &val) : type(Type::String) { new (&string_v) std::string(val); }
explicit Value(const char *val) : type(Type::String) { new (&string_v) std::string(val); }
explicit Value(const std::vector<Value> &val) : type(Type::List) { new (&list_v) std::vector<Value>(val); }
explicit Value(const std::map<std::string, Value> &val) : type(Type::Map) {
new (&map_v) std::map<std::string, Value>(val);
}
explicit Value(std::string &&val) noexcept : type(Type::String) { new (&string_v) std::string(std::move(val)); }
explicit Value(std::vector<Value> &&val) noexcept : type(Type::List) {
new (&list_v) std::vector<Value>(std::move(val));
}
explicit Value(std::map<std::string, Value> &&val) noexcept : type(Type::Map) {
new (&map_v) std::map<std::string, Value>(std::move(val));
}
~Value() { DestroyValue(); }
void DestroyValue() noexcept {
switch (type) {
case Type::Null:
case Type::Bool:
case Type::Int64:
case Type::Double:
return;
case Type::String:
std::destroy_at(&string_v);
return;
case Type::List:
std::destroy_at(&list_v);
return;
case Type::Map:
std::destroy_at(&map_v);
return;
case Type::Vertex:
std::destroy_at(&vertex_v);
return;
case Type::Path:
std::destroy_at(&path_v);
return;
case Type::Edge:
std::destroy_at(&edge_v);
}
}
Value(const Value &other) : type(other.type) {
switch (other.type) {
case Type::Null:
return;
case Type::Bool:
this->bool_v = other.bool_v;
return;
case Type::Int64:
this->int_v = other.int_v;
return;
case Type::Double:
this->double_v = other.double_v;
return;
case Type::String:
new (&string_v) std::string(other.string_v);
return;
case Type::List:
new (&list_v) std::vector<Value>(other.list_v);
return;
case Type::Map:
new (&map_v) std::map<std::string, Value>(other.map_v);
return;
case Type::Vertex:
new (&vertex_v) Vertex(other.vertex_v);
return;
case Type::Edge:
new (&edge_v) Edge(other.edge_v);
return;
case Type::Path:
new (&path_v) Path(other.path_v);
return;
}
}
Value(Value &&other) noexcept : type(other.type) {
switch (other.type) {
case Type::Null:
break;
case Type::Bool:
this->bool_v = other.bool_v;
break;
case Type::Int64:
this->int_v = other.int_v;
break;
case Type::Double:
this->double_v = other.double_v;
break;
case Type::String:
new (&string_v) std::string(std::move(other.string_v));
break;
case Type::List:
new (&list_v) std::vector<Value>(std::move(other.list_v));
break;
case Type::Map:
new (&map_v) std::map<std::string, Value>(std::move(other.map_v));
break;
case Type::Vertex:
new (&vertex_v) Vertex(std::move(other.vertex_v));
break;
case Type::Edge:
new (&edge_v) Edge(std::move(other.edge_v));
break;
case Type::Path:
new (&path_v) Path(std::move(other.path_v));
break;
}
other.DestroyValue();
other.type = Type::Null;
}
Value &operator=(const Value &other) {
if (this == &other) return *this;
DestroyValue();
type = other.type;
switch (other.type) {
case Type::Null:
break;
case Type::Bool:
this->bool_v = other.bool_v;
break;
case Type::Int64:
this->int_v = other.int_v;
break;
case Type::Double:
this->double_v = other.double_v;
break;
case Type::String:
new (&string_v) std::string(other.string_v);
break;
case Type::List:
new (&list_v) std::vector<Value>(other.list_v);
break;
case Type::Map:
new (&map_v) std::map<std::string, Value>(other.map_v);
break;
case Type::Vertex:
new (&vertex_v) Vertex(other.vertex_v);
break;
case Type::Edge:
new (&edge_v) Edge(other.edge_v);
break;
case Type::Path:
new (&path_v) Path(other.path_v);
break;
}
return *this;
}
Value &operator=(Value &&other) noexcept {
if (this == &other) return *this;
DestroyValue();
type = other.type;
switch (other.type) {
case Type::Null:
break;
case Type::Bool:
this->bool_v = other.bool_v;
break;
case Type::Int64:
this->int_v = other.int_v;
break;
case Type::Double:
this->double_v = other.double_v;
break;
case Type::String:
new (&string_v) std::string(std::move(other.string_v));
break;
case Type::List:
new (&list_v) std::vector<Value>(std::move(other.list_v));
break;
case Type::Map:
new (&map_v) std::map<std::string, Value>(std::move(other.map_v));
break;
case Type::Vertex:
new (&vertex_v) Vertex(std::move(other.vertex_v));
break;
case Type::Edge:
new (&edge_v) Edge(std::move(other.edge_v));
break;
case Type::Path:
new (&path_v) Path(std::move(other.path_v));
break;
}
other.DestroyValue();
other.type = Type::Null;
return *this;
}
enum class Type : uint8_t { Null, Bool, Int64, Double, String, List, Map, Vertex, Edge, Path };
Type type{Type::Null};
union {
Null null_v;
bool bool_v;
uint64_t int_v;
int64_t int_v;
double double_v;
std::string string_v;
std::vector<Value> list_v;
@@ -88,8 +315,6 @@ struct Value {
Edge edge_v;
Path path_v;
};
Type type;
};
struct ValuesMap {
@@ -125,17 +350,23 @@ enum class StorageView { OLD = 0, NEW = 1 };
struct ScanVerticesRequest {
Hlc transaction_id;
size_t start_id;
std::optional<std::vector<std::string>> props_to_return;
VertexId start_id;
std::optional<std::vector<PropertyId>> props_to_return;
std::optional<std::vector<std::string>> filter_expressions;
std::optional<size_t> batch_limit;
StorageView storage_view;
};
struct ScanResultRow {
Value vertex;
// empty() is no properties returned
std::map<PropertyId, Value> props;
};
struct ScanVerticesResponse {
bool success;
Values values;
std::optional<VertexId> next_start_id;
std::vector<ScanResultRow> results;
};
using VertexOrEdgeIds = std::variant<VertexId, EdgeId>;
@@ -194,17 +425,40 @@ struct ExpandOneResultRow {
// The drawback of this is currently the key of the map is always interpreted as a string in Value, not as an
// integer, which should be in case of mapped properties.
Vertex src_vertex;
std::optional<Values> src_vertex_properties;
Values edges;
std::optional<std::map<PropertyId, Value>> src_vertex_properties;
// NOTE: If the desired edges are specified in the request,
// edges_with_specific_properties will have a value and it will
// return the properties as a vector of property values. The order
// of the values returned should be the same as the PropertyIds
// were defined in the request.
std::optional<std::vector<std::tuple<VertexId, Gid, std::map<PropertyId, Value>>>> edges_with_all_properties;
std::optional<std::vector<std::tuple<VertexId, Gid, std::vector<Value>>>> edges_with_specific_properties;
};
struct ExpandOneResponse {
std::vector<ExpandOneResultRow> result;
};
struct UpdateVertexProp {
PrimaryKey primary_key;
std::vector<std::pair<PropertyId, Value>> property_updates;
};
struct UpdateEdgeProp {
EdgeId edge_id;
VertexId src;
VertexId dst;
std::vector<std::pair<PropertyId, Value>> property_updates;
};
/*
* Vertices
*/
struct NewVertex {
std::vector<Label> label_ids;
std::map<PropertyId, Value> properties;
PrimaryKey primary_key;
std::vector<std::pair<PropertyId, Value>> properties;
};
struct CreateVerticesRequest {
@@ -216,8 +470,62 @@ struct CreateVerticesResponse {
bool success;
};
struct DeleteVerticesRequest {
enum class DeletionType { DELETE, DETACH_DELETE };
Hlc transaction_id;
std::vector<std::vector<Value>> primary_keys;
DeletionType deletion_type;
};
struct DeleteVerticesResponse {
bool success;
};
struct UpdateVerticesRequest {
Hlc transaction_id;
std::vector<UpdateVertexProp> new_properties;
};
struct UpdateVerticesResponse {
bool success;
};
/*
* Edges
*/
struct CreateEdgesRequest {
Hlc transaction_id;
std::vector<Edge> edges;
};
struct CreateEdgesResponse {
bool success;
};
struct DeleteEdgesRequest {
Hlc transaction_id;
std::vector<Edge> edges;
};
struct DeleteEdgesResponse {
bool success;
};
struct UpdateEdgesRequest {
Hlc transaction_id;
std::vector<UpdateEdgeProp> new_properties;
};
struct UpdateEdgesResponse {
bool success;
};
using ReadRequests = std::variant<ExpandOneRequest, GetPropertiesRequest, ScanVerticesRequest>;
using ReadResponses = std::variant<ExpandOneResponse, GetPropertiesResponse, ScanVerticesResponse>;
using WriteRequests = CreateVerticesRequest;
using WriteResponses = CreateVerticesResponse;
using WriteRequests = std::variant<CreateVerticesRequest, DeleteVerticesRequest, UpdateVerticesRequest,
CreateEdgesRequest, DeleteEdgesRequest, UpdateEdgesRequest>;
using WriteResponses = std::variant<CreateVerticesResponse, DeleteVerticesResponse, UpdateVerticesResponse,
CreateEdgesResponse, DeleteEdgesResponse, UpdateEdgesResponse>;
} // namespace memgraph::msgs

View File

@@ -15,7 +15,8 @@ set(storage_v3_src_files
schemas.cpp
schema_validator.cpp
shard.cpp
storage.cpp)
storage.cpp
shard_rsm.cpp)
# #### Replication #####
define_add_lcp(add_lcp_storage lcp_storage_cpp_files generated_lcp_storage_files)

View File

@@ -521,6 +521,44 @@ ResultSchema<VertexAccessor> Shard::Accessor::CreateVertexAndValidate(
return vertex_acc;
}
ResultSchema<VertexAccessor> Shard::Accessor::CreateVertexAndValidate(
LabelId primary_label, const std::vector<LabelId> &labels, const std::vector<PropertyValue> &primary_properties,
const std::vector<std::pair<PropertyId, PropertyValue>> &properties) {
if (primary_label != shard_->primary_label_) {
throw utils::BasicException("Cannot add vertex to shard which does not hold the given primary label!");
}
auto maybe_schema_violation = GetSchemaValidator().ValidateVertexCreate(primary_label, labels, properties);
if (maybe_schema_violation) {
return {std::move(*maybe_schema_violation)};
}
OOMExceptionEnabler oom_exception;
auto acc = shard_->vertices_.access();
auto *delta = CreateDeleteObjectDelta(&transaction_);
auto [it, inserted] = acc.insert({Vertex{delta, primary_properties}});
VertexAccessor vertex_acc{&it->vertex, &transaction_, &shard_->indices_,
&shard_->constraints_, config_, shard_->vertex_validator_};
MG_ASSERT(inserted, "The vertex must be inserted here!");
MG_ASSERT(it != acc.end(), "Invalid Vertex accessor!");
// TODO(jbajic) Improve, maybe delay index update
for (const auto &[property_id, property_value] : properties) {
if (!shard_->schemas_.IsPropertyKey(primary_label, property_id)) {
if (const auto err = vertex_acc.SetProperty(property_id, property_value); err.HasError()) {
return {err.GetError()};
}
}
}
// Set secondary labels
for (auto label : labels) {
if (const auto err = vertex_acc.AddLabel(label); err.HasError()) {
return {err.GetError()};
}
}
delta->prev.Set(&it->vertex);
return vertex_acc;
}
std::optional<VertexAccessor> Shard::Accessor::FindVertex(std::vector<PropertyValue> primary_key, View view) {
auto acc = shard_->vertices_.access();
// Later on use label space

View File

@@ -237,11 +237,19 @@ class Shard final {
~Accessor();
// TODO(gvolfing) this is just a workaround for stitching remove this later.
LabelId GetPrimaryLabel() const noexcept { return shard_->primary_label_; }
/// @throw std::bad_alloc
ResultSchema<VertexAccessor> CreateVertexAndValidate(
LabelId primary_label, const std::vector<LabelId> &labels,
const std::vector<std::pair<PropertyId, PropertyValue>> &properties);
/// @throw std::bad_alloc
ResultSchema<VertexAccessor> CreateVertexAndValidate(
LabelId primary_label, const std::vector<LabelId> &labels, const std::vector<PropertyValue> &primary_properties,
const std::vector<std::pair<PropertyId, PropertyValue>> &properties);
std::optional<VertexAccessor> FindVertex(std::vector<PropertyValue> primary_key, View view);
VerticesIterable Vertices(View view) {

View File

@@ -0,0 +1,934 @@
// Copyright 2022 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 <functional>
#include <iterator>
#include <utility>
#include "query/v2/requests.hpp"
#include "storage/v3/shard_rsm.hpp"
#include "storage/v3/vertex_accessor.hpp"
using memgraph::msgs::Label;
using memgraph::msgs::PropertyId;
using memgraph::msgs::Value;
using memgraph::msgs::VertexId;
namespace {
// TODO(gvolfing) use come algorithm instead of explicit for loops
// TODO(gvolfing) make this use rvalue& again instead of copy!!!!
memgraph::storage::v3::PropertyValue ToPropertyValue(Value &&value) {
using PV = memgraph::storage::v3::PropertyValue;
PV ret;
switch (value.type) {
case Value::Type::Null:
return PV{};
case Value::Type::Bool:
return PV(value.bool_v);
case Value::Type::Int64:
return PV(static_cast<int64_t>(value.int_v));
case Value::Type::Double:
return PV(value.double_v);
case Value::Type::String:
return PV(value.string_v);
case Value::Type::List: {
std::vector<PV> list;
for (auto &elem : value.list_v) {
list.emplace_back(ToPropertyValue(std::move(elem)));
}
return PV(list);
}
case Value::Type::Map: {
std::map<std::string, PV> map;
for (auto &[key, value] : value.map_v) {
map.emplace(std::make_pair(key, ToPropertyValue(std::move(value))));
}
return PV(map);
}
// These are not PropertyValues
case Value::Type::Vertex:
case Value::Type::Edge:
case Value::Type::Path:
MG_ASSERT(false, "Not PropertyValue");
}
return ret;
}
Value ToValue(const memgraph::storage::v3::PropertyValue &pv) {
using memgraph::storage::v3::PropertyValue;
switch (pv.type()) {
case PropertyValue::Type::Bool:
return Value(pv.ValueBool());
case PropertyValue::Type::Double:
return Value(pv.ValueDouble());
case PropertyValue::Type::Int:
return Value(pv.ValueInt());
case PropertyValue::Type::List: {
std::vector<Value> list(pv.ValueList().size());
for (const auto &elem : pv.ValueList()) {
list.emplace_back(ToValue(elem));
}
return Value(list);
}
case PropertyValue::Type::Map: {
std::map<std::string, Value> map;
for (const auto &[key, val] : pv.ValueMap()) {
// maybe use std::make_pair once the && issue is resolved.
map.emplace(std::make_pair(key, ToValue(val)));
}
return Value(map);
}
case PropertyValue::Type::Null:
return Value{};
case PropertyValue::Type::String:
return Value(pv.ValueString());
case PropertyValue::Type::TemporalData: {
// TBD -> we need to specify this in the messages, not a priority.
MG_ASSERT(false, "Temporal datatypes are not yet implemented on Value!");
return Value{};
}
}
}
std::vector<std::pair<memgraph::storage::v3::PropertyId, memgraph::storage::v3::PropertyValue>> ConvertPropertyMap(
std::vector<std::pair<PropertyId, Value>> &&properties) {
std::vector<std::pair<memgraph::storage::v3::PropertyId, memgraph::storage::v3::PropertyValue>> ret;
ret.reserve(properties.size());
for (auto &[key, value] : properties) {
ret.emplace_back(std::make_pair(key, ToPropertyValue(std::move(value))));
}
return ret;
}
std::vector<memgraph::storage::v3::PropertyValue> ConvertPropertyVector(std::vector<Value> &&vec) {
std::vector<memgraph::storage::v3::PropertyValue> ret;
ret.reserve(vec.size());
for (auto &elem : vec) {
ret.push_back(ToPropertyValue(std::move(elem)));
}
return ret;
}
std::vector<Value> ConvertValueVector(const std::vector<memgraph::storage::v3::PropertyValue> &vec) {
std::vector<Value> ret;
ret.reserve(vec.size());
for (const auto &elem : vec) {
ret.push_back(ToValue(elem));
}
return ret;
}
std::optional<std::map<PropertyId, Value>> CollectPropertiesFromAccessor(
const memgraph::storage::v3::VertexAccessor &acc, const std::vector<memgraph::storage::v3::PropertyId> &props,
memgraph::storage::v3::View view) {
std::map<PropertyId, Value> ret;
for (const auto &prop : props) {
auto result = acc.GetProperty(prop, view);
if (result.HasError()) {
spdlog::debug("Encountered an Error while trying to get a vertex property.");
continue;
}
auto &value = result.GetValue();
if (value.IsNull()) {
spdlog::debug("The specified property does not exist but it should");
continue;
}
ret.emplace(std::make_pair(prop, ToValue(value)));
}
return ret;
}
std::optional<std::map<PropertyId, Value>> CollectAllPropertiesFromAccessor(
const memgraph::storage::v3::VertexAccessor &acc, memgraph::storage::v3::View view) {
std::map<PropertyId, Value> ret;
auto iter = acc.Properties(view);
if (iter.HasError()) {
spdlog::debug("Encountered an error while trying to get vertex properties.");
}
for (const auto &[prop_key, prop_val] : iter.GetValue()) {
ret.emplace(prop_key, ToValue(prop_val));
}
return ret;
}
Value ConstructValueVertex(const memgraph::storage::v3::VertexAccessor &acc, memgraph::storage::v3::View view) {
// Get the vertex id
auto prim_label = acc.PrimaryLabel(view).GetValue();
Label value_label{.id = prim_label};
auto prim_key = ConvertValueVector(acc.PrimaryKey(view).GetValue());
VertexId vertex_id = std::make_pair(value_label, prim_key);
// Get the labels
auto vertex_labels = acc.Labels(view).GetValue();
std::vector<Label> value_labels;
for (const auto &label : vertex_labels) {
Label l = {.id = label};
value_labels.push_back(l);
}
return Value({.id = vertex_id, .labels = value_labels});
}
bool DoesEdgeTypeMatch(const memgraph::msgs::ExpandOneRequest &req, const memgraph::storage::v3::EdgeAccessor &edge) {
for (const auto &edge_type : req.edge_types) {
if (memgraph::storage::v3::EdgeTypeId::FromUint(edge_type.id) == edge.EdgeType()) {
return true;
}
}
return false;
}
} // namespace
namespace memgraph::storage::v3 {
msgs::WriteResponses ShardRsm::ApplyWrite(msgs::CreateVerticesRequest &&req) {
auto acc = shard_->Access();
// Workaround untill we have access to CreateVertexAndValidate()
// with the new signature that does not require the primary label.
const auto prim_label = acc.GetPrimaryLabel();
bool action_successful = true;
for (auto &new_vertex : req.new_vertices) {
/// TODO(gvolfing) Remove this. In the new implementation each shard
/// should have a predetermined primary label, so there is no point in
/// specifying it in the accessor functions. Their signature will
/// change.
/// TODO(gvolfing) Consider other methods than converting. Change either
/// the way that the property map is stored in the messages, or the
/// signature of CreateVertexAndValidate.
auto converted_property_map = ConvertPropertyMap(std::move(new_vertex.properties));
// TODO(gvolfing) make sure if this conversion is actually needed.
std::vector<memgraph::storage::v3::LabelId> converted_label_ids;
converted_label_ids.reserve(new_vertex.label_ids.size());
for (const auto &label_id : new_vertex.label_ids) {
converted_label_ids.emplace_back(label_id.id);
}
auto result_schema =
acc.CreateVertexAndValidate(prim_label, converted_label_ids,
ConvertPropertyVector(std::move(new_vertex.primary_key)), converted_property_map);
if (result_schema.HasError()) {
auto &error = result_schema.GetError();
std::visit(
[]<typename T>(T &&) {
using ErrorType = std::remove_cvref_t<T>;
if constexpr (std::is_same_v<ErrorType, SchemaViolation>) {
spdlog::debug("Creating vertex failed with error: SchemaViolation");
} else if constexpr (std::is_same_v<ErrorType, Error>) {
spdlog::debug("Creating vertex failed with error: Error");
} else {
static_assert(kAlwaysFalse<T>, "Missing type from variant visitor");
}
},
error);
action_successful = false;
break;
}
}
msgs::CreateVerticesResponse resp{};
resp.success = action_successful;
if (action_successful) {
auto result = acc.Commit(req.transaction_id.logical_id);
if (result.HasError()) {
resp.success = false;
spdlog::debug(&"ConstraintViolation, commiting vertices was unsuccesfull with transaction id: "[req.transaction_id
.logical_id]);
}
}
return resp;
}
msgs::WriteResponses ShardRsm::ApplyWrite(msgs::UpdateVerticesRequest &&req) {
auto acc = shard_->Access();
bool action_successful = true;
for (auto &vertex : req.new_properties) {
if (!action_successful) {
break;
}
auto vertex_to_update = acc.FindVertex(ConvertPropertyVector(std::move(vertex.primary_key)), View::OLD);
if (!vertex_to_update) {
action_successful = false;
spdlog::debug(
&"Vertex could not be found while trying to update its properties. Transaction id: "[req.transaction_id
.logical_id]);
continue;
}
for (auto &update_prop : vertex.property_updates) {
// TODO(gvolfing) Maybe check if the setting is valid if SetPropertyAndValidate()
// does not do that alreaedy.
auto result_schema =
vertex_to_update->SetPropertyAndValidate(update_prop.first, ToPropertyValue(std::move(update_prop.second)));
if (result_schema.HasError()) {
auto &error = result_schema.GetError();
std::visit(
[&action_successful]<typename T>(T &&) {
using ErrorType = std::remove_cvref_t<T>;
if constexpr (std::is_same_v<ErrorType, SchemaViolation>) {
action_successful = false;
spdlog::debug("Updating vertex failed with error: SchemaViolation");
} else if constexpr (std::is_same_v<ErrorType, Error>) {
action_successful = false;
spdlog::debug("Updating vertex failed with error: Error");
} else {
static_assert(kAlwaysFalse<T>, "Missing type from variant visitor");
}
},
error);
break;
}
}
}
msgs::UpdateVerticesResponse resp{};
resp.success = action_successful;
if (action_successful) {
auto result = acc.Commit(req.transaction_id.logical_id);
if (result.HasError()) {
resp.success = false;
spdlog::debug(&"ConstraintViolation, commiting vertices was unsuccesfull with transaction id:"[req.transaction_id
.logical_id]);
}
}
return resp;
}
msgs::WriteResponses ShardRsm::ApplyWrite(msgs::DeleteVerticesRequest &&req) {
bool action_successful = true;
auto acc = shard_->Access();
for (auto &propval : req.primary_keys) {
if (!action_successful) {
break;
}
auto vertex_acc = acc.FindVertex(ConvertPropertyVector(std::move(propval)), View::OLD);
if (!vertex_acc) {
spdlog::debug(
&"Error while trying to delete vertex. Vertex to delete does not exist. Transaction id: "[req.transaction_id
.logical_id]);
action_successful = false;
} else {
// TODO(gvolfing)
// Since we will not have different kinds of deletion types in one transaction,
// we dont have to enter the switch statement on every iteration. Optimize this.
switch (req.deletion_type) {
case msgs::DeleteVerticesRequest::DeletionType::DELETE: {
auto result = acc.DeleteVertex(&vertex_acc.value());
if (result.HasError() || !(result.GetValue().has_value())) {
action_successful = false;
spdlog::debug(&"Error while trying to delete vertex. Transaction id: "[req.transaction_id.logical_id]);
}
break;
}
case msgs::DeleteVerticesRequest::DeletionType::DETACH_DELETE: {
auto result = acc.DetachDeleteVertex(&vertex_acc.value());
if (result.HasError() || !(result.GetValue().has_value())) {
action_successful = false;
spdlog::debug(
&"Error while trying to detach and delete vertex. Transaction id: "[req.transaction_id.logical_id]);
}
break;
}
}
}
}
msgs::DeleteVerticesResponse resp{};
resp.success = action_successful;
if (action_successful) {
auto result = acc.Commit(req.transaction_id.logical_id);
if (result.HasError()) {
resp.success = false;
spdlog::debug(&"ConstraintViolation, commiting vertices was unsuccesfull with transaction id: "[req.transaction_id
.logical_id]);
}
}
return resp;
}
msgs::WriteResponses ShardRsm::ApplyWrite(msgs::CreateEdgesRequest &&req) {
auto acc = shard_->Access();
bool action_successful = true;
for (auto &edge : req.edges) {
auto vertex_acc_from_primary_key = edge.src.second;
auto vertex_from_acc = acc.FindVertex(ConvertPropertyVector(std::move(vertex_acc_from_primary_key)), View::OLD);
auto vertex_acc_to_primary_key = edge.dst.second;
auto vertex_to_acc = acc.FindVertex(ConvertPropertyVector(std::move(vertex_acc_to_primary_key)), View::OLD);
if (!vertex_from_acc || !vertex_to_acc) {
action_successful = false;
spdlog::debug(
&"Error while trying to insert edge, vertex does not exist. Transaction id: "[req.transaction_id.logical_id]);
break;
}
auto from_vertex_id = VertexId(edge.src.first.id, ConvertPropertyVector(std::move(edge.src.second)));
auto to_vertex_id = VertexId(edge.dst.first.id, ConvertPropertyVector(std::move(edge.dst.second)));
auto edge_acc =
acc.CreateEdge(from_vertex_id, to_vertex_id, EdgeTypeId::FromUint(edge.type.id), Gid::FromUint(edge.id.gid));
if (edge_acc.HasError()) {
action_successful = false;
spdlog::debug(&"Creating edge was not successful. Transaction id: "[req.transaction_id.logical_id]);
break;
}
// Add properties to the edge if there is any
if (edge.properties) {
for (auto &[edge_prop_key, edge_prop_val] : edge.properties.value()) {
auto set_result = edge_acc->SetProperty(edge_prop_key, ToPropertyValue(std::move(edge_prop_val)));
if (set_result.HasError()) {
action_successful = false;
spdlog::debug(&"Adding property to edge was not successful. Transaction id: "[req.transaction_id.logical_id]);
break;
}
}
}
}
msgs::CreateEdgesResponse resp{};
resp.success = action_successful;
if (action_successful) {
auto result = acc.Commit(req.transaction_id.logical_id);
if (result.HasError()) {
resp.success = false;
spdlog::debug(
&"ConstraintViolation, commiting edge creation was unsuccesfull with transaction id: "[req.transaction_id
.logical_id]);
}
}
return resp;
}
msgs::WriteResponses ShardRsm::ApplyWrite(msgs::DeleteEdgesRequest &&req) {
bool action_successful = true;
auto acc = shard_->Access();
for (auto &edge : req.edges) {
if (!action_successful) {
break;
}
auto edge_acc = acc.DeleteEdge(VertexId(edge.src.first.id, ConvertPropertyVector(std::move(edge.src.second))),
VertexId(edge.dst.first.id, ConvertPropertyVector(std::move(edge.dst.second))),
Gid::FromUint(edge.id.gid));
if (edge_acc.HasError() || !edge_acc.HasValue()) {
spdlog::debug(&"Error while trying to delete edge. Transaction id: "[req.transaction_id.logical_id]);
action_successful = false;
continue;
}
}
msgs::DeleteEdgesResponse resp{};
resp.success = action_successful;
if (action_successful) {
auto result = acc.Commit(req.transaction_id.logical_id);
if (result.HasError()) {
resp.success = false;
spdlog::debug(
&"ConstraintViolation, commiting edge creation was unsuccesfull with transaction id: "[req.transaction_id
.logical_id]);
}
}
return resp;
}
msgs::WriteResponses ShardRsm::ApplyWrite(msgs::UpdateEdgesRequest &&req) {
auto acc = shard_->Access();
bool action_successful = true;
for (auto &edge : req.new_properties) {
if (!action_successful) {
break;
}
auto vertex_acc = acc.FindVertex(ConvertPropertyVector(std::move(edge.src.second)), View::OLD);
if (!vertex_acc) {
action_successful = false;
spdlog::debug(
&"Encountered an error while trying to acquire VertexAccessor with transaction id: "[req.transaction_id
.logical_id]);
continue;
}
// Since we are using the source vertex of the edge we are only intrested
// in the vertex's out-going edges
auto edges_res = vertex_acc->OutEdges(View::OLD);
if (edges_res.HasError()) {
action_successful = false;
spdlog::debug(
&"Encountered an error while trying to acquire EdgeAccessor with transaction id: "[req.transaction_id
.logical_id]);
continue;
}
auto &edge_accessors = edges_res.GetValue();
// Look for the appropriate edge accessor
bool edge_accessor_did_match = false;
for (auto &edge_accessor : edge_accessors) {
if (edge_accessor.Gid().AsUint() == edge.edge_id.gid) { // Found the appropriate accessor
edge_accessor_did_match = true;
for (auto &[key, value] : edge.property_updates) {
// TODO(gvolfing)
// Check if the property was set if SetProperty does not do that itself.
auto res = edge_accessor.SetProperty(key, ToPropertyValue(std::move(value)));
if (res.HasError()) {
spdlog::debug(&"Encountered an error while trying to set the property of an Edge with transaction id: "
[req.transaction_id.logical_id]);
}
}
}
}
if (!edge_accessor_did_match) {
action_successful = false;
spdlog::debug(&"Could not find the Edge with the specified Gid. Transaction id: "[req.transaction_id.logical_id]);
continue;
}
}
msgs::UpdateEdgesResponse resp{};
resp.success = action_successful;
if (action_successful) {
auto result = acc.Commit(req.transaction_id.logical_id);
if (result.HasError()) {
resp.success = false;
spdlog::debug(
&"ConstraintViolation, commiting edge update was unsuccesfull with transaction id: "[req.transaction_id
.logical_id]);
}
}
return resp;
}
msgs::ReadResponses ShardRsm::HandleRead(msgs::ScanVerticesRequest &&req) {
auto acc = shard_->Access();
bool action_successful = true;
std::vector<msgs::ScanResultRow> results;
std::optional<msgs::VertexId> next_start_id;
const auto view = View(req.storage_view);
auto vertex_iterable = acc.Vertices(view);
bool did_reach_starting_point = false;
uint64_t sample_counter = 0;
for (auto it = vertex_iterable.begin(); it != vertex_iterable.end(); ++it) {
const auto &vertex = *it;
if (ConvertPropertyVector(std::move(req.start_id.second)) == vertex.PrimaryKey(View(req.storage_view)).GetValue()) {
did_reach_starting_point = true;
}
if (did_reach_starting_point) {
std::optional<std::map<PropertyId, Value>> found_props;
if (req.props_to_return) {
found_props = CollectPropertiesFromAccessor(vertex, req.props_to_return.value(), view);
} else {
found_props = CollectAllPropertiesFromAccessor(vertex, view);
}
if (!found_props) {
continue;
}
results.emplace_back(
msgs::ScanResultRow{.vertex = ConstructValueVertex(vertex, view), .props = found_props.value()});
++sample_counter;
if (sample_counter == req.batch_limit) {
// Reached the maximum specified batch size.
// Get the next element before exiting.
const auto &next_vertex = *(++it);
next_start_id = ConstructValueVertex(next_vertex, view).vertex_v.id;
break;
}
}
}
msgs::ScanVerticesResponse resp{};
resp.success = action_successful;
if (action_successful) {
resp.next_start_id = next_start_id;
resp.results = std::move(results);
}
return resp;
}
msgs::ReadResponses ShardRsm::HandleRead(msgs::ExpandOneRequest &&req) {
auto acc = shard_->Access();
bool action_successful = true;
using EdgeProperties = std::variant<std::map<PropertyId, msgs::Value>, std::vector<msgs::Value>>;
std::function<EdgeProperties(const EdgeAccessor &)> get_edge_properties;
if (!req.edge_properties) {
get_edge_properties = [&req, &action_successful](const EdgeAccessor &edge) {
std::map<PropertyId, msgs::Value> ret;
auto property_results = edge.Properties(View::OLD);
if (property_results.HasError()) {
spdlog::debug(
&"Encountered an error while trying to get out-going EdgeAccessors. Transaction id: "[req.transaction_id
.logical_id]);
action_successful = false;
return ret;
}
for (const auto &[prop_key, prop_val] : property_results.GetValue()) {
ret.insert(std::make_pair(prop_key, ToValue(prop_val)));
}
return ret;
};
} else {
// TODO(gvolfing) - do we want to set the action_successful here?
get_edge_properties = [&req, &action_successful](const EdgeAccessor &edge) {
std::vector<msgs::Value> ret;
ret.reserve(req.edge_properties.value().size());
for (const auto &edge_prop : req.edge_properties.value()) {
// TODO(gvolfing) maybe check for the absence of certain properties
ret.push_back(ToValue(edge.GetProperty(edge_prop, View::OLD).GetValue()));
}
return ret;
};
}
std::vector<msgs::ExpandOneResultRow> results;
for (auto &src_vertex : req.src_vertices) {
msgs::ExpandOneResultRow current_row;
msgs::Vertex source_vertex;
/// The empty optional means return all of the properties, while an empty
/// list means do not return any properties.
std::optional<std::map<PropertyId, Value>> src_vertex_properties_opt;
std::map<PropertyId, Value> src_vertex_properties;
auto v_acc = acc.FindVertex(ConvertPropertyVector(std::move(src_vertex.second)), View::OLD);
/// Fill up source vertex
auto secondary_labels = v_acc->Labels(View::OLD);
if (secondary_labels.HasError()) {
spdlog::debug(&"Encountered an error while trying to get the secondary labels of a vertex. Transaction id: "
[req.transaction_id.logical_id]);
action_successful = false;
break;
}
source_vertex.id = src_vertex;
source_vertex.labels.reserve(secondary_labels.GetValue().size());
for (auto label_id : secondary_labels.GetValue()) {
source_vertex.labels.push_back({.id = label_id});
}
/// Fill up source vertex properties
if (!req.src_vertex_properties) {
auto props = v_acc->Properties(View::OLD);
if (props.HasError()) {
spdlog::debug(
&"Encountered an error while trying to access vertex properties. Transaction id: "[req.transaction_id
.logical_id]);
action_successful = false;
break;
}
for (auto &[key, val] : props.GetValue()) {
src_vertex_properties.insert(std::make_pair(key, ToValue(val)));
}
src_vertex_properties_opt = src_vertex_properties;
} else if (req.src_vertex_properties.value().empty()) {
src_vertex_properties_opt = {};
} else {
for (const auto &prop : req.src_vertex_properties.value()) {
const auto &prop_val = v_acc->GetProperty(prop, View::OLD);
src_vertex_properties.insert(std::make_pair(prop, ToValue(prop_val.GetValue())));
}
src_vertex_properties_opt = src_vertex_properties;
}
/// Fill up connecting edges
std::vector<EdgeAccessor> in_edges;
std::vector<EdgeAccessor> out_edges;
switch (req.direction) {
case msgs::EdgeDirection::OUT: {
auto out_edges_result = v_acc->OutEdges(View::OLD);
if (out_edges_result.HasError()) {
spdlog::debug(
&"Encountered an error while trying to get out-going EdgeAccessors. Transaction id: "[req.transaction_id
.logical_id]);
action_successful = false;
break;
}
out_edges = out_edges_result.GetValue();
break;
}
case msgs::EdgeDirection::IN: {
auto in_edges_result = v_acc->InEdges(View::OLD);
if (in_edges_result.HasError()) {
spdlog::debug(
&"Encountered an error while trying to get in-going EdgeAccessors. Transaction id: "[req.transaction_id
.logical_id]);
action_successful = false;
break;
}
in_edges = in_edges_result.GetValue();
break;
}
case msgs::EdgeDirection::BOTH: {
auto in_edges_result = v_acc->InEdges(View::OLD);
if (in_edges_result.HasError()) {
spdlog::debug(
&"Encountered an error while trying to get in-going EdgeAccessors. Transaction id: "[req.transaction_id
.logical_id]);
action_successful = false;
break;
}
in_edges = in_edges_result.GetValue();
auto out_edges_result = v_acc->InEdges(View::OLD);
if (out_edges_result.HasError()) {
spdlog::debug(
&"Encountered an error while trying to get out-going EdgeAccessors. Transaction id: "[req.transaction_id
.logical_id]);
action_successful = false;
break;
}
out_edges = out_edges_result.GetValue();
break;
}
}
// Check for stoppage here because of the switch case
if (!action_successful) {
break;
}
/// Assemble the edge properties
std::optional<std::vector<std::tuple<msgs::VertexId, msgs::Gid, std::map<PropertyId, msgs::Value>>>>
edges_with_all_properties;
std::optional<std::vector<std::tuple<msgs::VertexId, msgs::Gid, std::vector<msgs::Value>>>>
edges_with_specific_properties;
if (!req.edge_properties) {
std::vector<std::tuple<msgs::VertexId, msgs::Gid, std::map<PropertyId, msgs::Value>>> ret_in;
// ret_in.reserve(in_edges.size());
std::vector<std::tuple<msgs::VertexId, msgs::Gid, std::map<PropertyId, msgs::Value>>> ret_out;
// ret_out.reserve(out_edges.size());
for (const auto &edge : in_edges) {
if (!DoesEdgeTypeMatch(req, edge)) {
continue;
}
std::tuple<msgs::VertexId, msgs::Gid, std::map<PropertyId, msgs::Value>> ret_tuple;
msgs::Label label;
label.id = edge.FromVertex().primary_label;
msgs::VertexId other_vertex = std::make_pair(label, ConvertValueVector(edge.FromVertex().primary_key));
const auto &edge_props_var = get_edge_properties(edge);
auto edge_props = std::get<std::map<PropertyId, msgs::Value>>(edge_props_var);
msgs::Gid gid = edge.Gid().AsUint();
ret_tuple = {other_vertex, gid, edge_props};
ret_in.push_back(ret_tuple);
}
for (const auto &edge : out_edges) {
if (!DoesEdgeTypeMatch(req, edge)) {
continue;
}
std::tuple<msgs::VertexId, msgs::Gid, std::map<PropertyId, msgs::Value>> ret_tuple;
msgs::Label label;
label.id = edge.FromVertex().primary_label;
msgs::VertexId other_vertex = std::make_pair(label, ConvertValueVector(edge.ToVertex().primary_key));
const auto &edge_props_var = get_edge_properties(edge);
auto edge_props = std::get<std::map<PropertyId, msgs::Value>>(edge_props_var);
msgs::Gid gid = edge.Gid().AsUint();
ret_tuple = {other_vertex, gid, edge_props};
ret_out.push_back(ret_tuple);
}
// Set one of the options to the actual datastructure and the otherone to nullopt
switch (req.direction) {
case msgs::EdgeDirection::OUT: {
edges_with_all_properties = ret_out;
break;
}
case msgs::EdgeDirection::IN: {
edges_with_all_properties = ret_in;
break;
}
case msgs::EdgeDirection::BOTH: {
std::vector<std::tuple<msgs::VertexId, msgs::Gid, std::map<PropertyId, msgs::Value>>> ret;
ret.resize(ret_out.size() + ret_in.size());
ret.insert(ret.end(), ret_in.begin(), ret_in.end());
ret.insert(ret.end(), ret_out.begin(), ret_out.end());
edges_with_all_properties = ret;
break;
}
}
edges_with_specific_properties = {};
} else {
// when user specifies specific properties, its enough to return just a vector
std::vector<std::tuple<msgs::VertexId, msgs::Gid, std::vector<msgs::Value>>> ret_in;
// ret_in.reserve(in_edges.size());
std::vector<std::tuple<msgs::VertexId, msgs::Gid, std::vector<msgs::Value>>> ret_out;
// ret_out.reserve(out_edges.size());
for (const auto &edge : in_edges) {
if (!DoesEdgeTypeMatch(req, edge)) {
continue;
}
std::tuple<msgs::VertexId, msgs::Gid, std::vector<msgs::Value>> ret_tuple;
msgs::Label label;
label.id = edge.FromVertex().primary_label;
msgs::VertexId other_vertex = std::make_pair(label, ConvertValueVector(edge.FromVertex().primary_key));
const auto &edge_props_var = get_edge_properties(edge);
auto edge_props = std::get<std::vector<msgs::Value>>(edge_props_var);
msgs::Gid gid = edge.Gid().AsUint();
ret_tuple = {other_vertex, gid, edge_props};
ret_in.push_back(ret_tuple);
}
for (const auto &edge : out_edges) {
if (!DoesEdgeTypeMatch(req, edge)) {
continue;
}
std::tuple<msgs::VertexId, msgs::Gid, std::vector<msgs::Value>> ret_tuple;
msgs::Label label;
label.id = edge.FromVertex().primary_label;
msgs::VertexId other_vertex = std::make_pair(label, ConvertValueVector(edge.ToVertex().primary_key));
const auto &edge_props_var = get_edge_properties(edge);
auto edge_props = std::get<std::vector<msgs::Value>>(edge_props_var);
msgs::Gid gid = edge.Gid().AsUint();
ret_tuple = {other_vertex, gid, edge_props};
ret_out.push_back(ret_tuple);
}
// Set one of the options to the actual datastructure and the otherone to nullopt
switch (req.direction) {
case msgs::EdgeDirection::OUT: {
edges_with_specific_properties = ret_out;
break;
}
case msgs::EdgeDirection::IN: {
edges_with_specific_properties = ret_in;
break;
}
case msgs::EdgeDirection::BOTH: {
std::vector<std::tuple<msgs::VertexId, msgs::Gid, std::vector<msgs::Value>>> ret;
ret.resize(ret_out.size() + ret_in.size());
ret.insert(ret.end(), ret_in.begin(), ret_in.end());
ret.insert(ret.end(), ret_out.begin(), ret_out.end());
edges_with_specific_properties = ret;
break;
}
}
edges_with_all_properties = {};
}
results.emplace_back(msgs::ExpandOneResultRow{.src_vertex = src_vertex,
.src_vertex_properties = src_vertex_properties,
.edges_with_all_properties = edges_with_all_properties,
.edges_with_specific_properties = edges_with_specific_properties});
}
msgs::ExpandOneResponse resp{};
if (action_successful) {
resp.result = std::move(results);
}
return resp;
}
// msgs::WriteResponses ShardRsm::ApplyWrite(msgs::UpdateVerticesRequest && /*req*/) {
// return msgs::UpdateVerticesResponse{};
// }
// msgs::WriteResponses ShardRsm::ApplyWrite(msgs::DeleteEdgesRequest && /*req*/) { return msgs::DeleteEdgesResponse{};
// } msgs::WriteResponses ShardRsm::ApplyWrite(msgs::UpdateEdgesRequest && /*req*/) { return
// msgs::UpdateEdgesResponse{}; } msgs::ReadResponses ShardRsm::HandleRead(msgs::ExpandOneRequest && /*req*/) { return
// msgs::ExpandOneResponse{}; } NOLINTNEXTLINE(readability-convert-member-functions-to-static)
msgs::ReadResponses ShardRsm::HandleRead(msgs::GetPropertiesRequest && /*req*/) {
return msgs::GetPropertiesResponse{};
}
} // namespace memgraph::storage::v3

View File

@@ -0,0 +1,58 @@
// Copyright 2022 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.
#pragma once
#include <memory>
#include <variant>
#include <openssl/ec.h>
#include "query/v2/requests.hpp"
#include "storage/v3/shard.hpp"
#include "storage/v3/vertex_accessor.hpp"
namespace memgraph::storage::v3 {
template <typename>
constexpr auto kAlwaysFalse = false;
class ShardRsm {
std::unique_ptr<Shard> shard_;
msgs::ReadResponses HandleRead(msgs::ExpandOneRequest &&req);
msgs::ReadResponses HandleRead(msgs::GetPropertiesRequest &&req);
msgs::ReadResponses HandleRead(msgs::ScanVerticesRequest &&req);
msgs::WriteResponses ApplyWrite(msgs::CreateVerticesRequest &&req);
msgs::WriteResponses ApplyWrite(msgs::DeleteVerticesRequest &&req);
msgs::WriteResponses ApplyWrite(msgs::UpdateVerticesRequest &&req);
msgs::WriteResponses ApplyWrite(msgs::CreateEdgesRequest &&req);
msgs::WriteResponses ApplyWrite(msgs::DeleteEdgesRequest &&req);
msgs::WriteResponses ApplyWrite(msgs::UpdateEdgesRequest &&req);
public:
explicit ShardRsm(std::unique_ptr<Shard> &&shard) : shard_(std::move(shard)){};
// NOLINTNEXTLINE(readability-convert-member-functions-to-static)
msgs::ReadResponses Read(msgs::ReadRequests requests) {
return std::visit([&](auto &&request) mutable { return HandleRead(std::forward<decltype(request)>(request)); },
std::move(requests));
}
// NOLINTNEXTLINE(readability-convert-member-functions-to-static)
msgs::WriteResponses Apply(msgs::WriteRequests requests) {
return std::visit([&](auto &&request) mutable { return ApplyWrite(std::forward<decltype(request)>(request)); },
std::move(requests));
}
};
} // namespace memgraph::storage::v3

View File

@@ -1,34 +1,28 @@
set(test_prefix memgraph__simulation__)
find_package(gflags)
find_package(Boost REQUIRED)
find_package(OpenSSL REQUIRED)
add_custom_target(memgraph__simulation)
function(add_simulation_test test_cpp san)
function(add_simulation_test test_cpp)
# get exec name (remove extension from the abs path)
get_filename_component(exec_name ${test_cpp} NAME_WE)
set(target_name ${test_prefix}${exec_name})
add_executable(${target_name} ${test_cpp})
# OUTPUT_NAME sets the real name of a target when it is built and can be
# used to help create two targets of the same name even though CMake
# requires unique logical target names
set_target_properties(${target_name} PROPERTIES OUTPUT_NAME ${exec_name})
target_link_libraries(${target_name} gtest gmock mg-utils mg-io mg-io-simulator)
# sanitize
target_compile_options(${target_name} PRIVATE -fsanitize=${san})
target_link_options(${target_name} PRIVATE -fsanitize=${san})
target_link_libraries(${target_name} mg-storage-v3 mg-communication gtest gmock mg-utils mg-io mg-io-simulator Boost::headers)
# register test
add_test(${target_name} ${exec_name})
add_dependencies(memgraph__simulation ${target_name})
endfunction(add_simulation_test)
add_simulation_test(basic_request.cpp address)
add_simulation_test(raft.cpp address)
add_simulation_test(trial_query_storage/query_storage_test.cpp address)
add_simulation_test(sharded_map.cpp address)
add_simulation_test(basic_request.cpp)
add_simulation_test(raft.cpp)
add_simulation_test(trial_query_storage/query_storage_test.cpp)
add_simulation_test(sharded_map.cpp)
add_simulation_test(shard_rsm.cpp)

View File

@@ -0,0 +1,849 @@
// Copyright 2022 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 <cstdint>
#include <iostream>
#include <optional>
#include <thread>
#include <utility>
#include <vector>
#include "io/address.hpp"
#include "io/errors.hpp"
#include "io/rsm/raft.hpp"
#include "io/rsm/rsm_client.hpp"
#include "io/simulator/simulator.hpp"
#include "io/simulator/simulator_transport.hpp"
#include "query/v2/requests.hpp"
#include "storage/v3/id_types.hpp"
#include "storage/v3/property_value.hpp"
#include "storage/v3/shard.hpp"
#include "storage/v3/shard_rsm.hpp"
#include "storage/v3/view.hpp"
#include "utils/result.hpp"
namespace memgraph::storage::v3::tests {
using io::Address;
using io::Io;
using io::ResponseEnvelope;
using io::ResponseFuture;
using io::Time;
using io::TimedOut;
using io::rsm::Raft;
using io::rsm::ReadRequest;
using io::rsm::ReadResponse;
using io::rsm::RsmClient;
using io::rsm::WriteRequest;
using io::rsm::WriteResponse;
using io::simulator::Simulator;
using io::simulator::SimulatorConfig;
using io::simulator::SimulatorStats;
using io::simulator::SimulatorTransport;
using utils::BasicResult;
using msgs::ReadRequests;
using msgs::ReadResponses;
using msgs::WriteRequests;
using msgs::WriteResponses;
using ShardClient = RsmClient<SimulatorTransport, WriteRequests, WriteResponses, ReadRequests, ReadResponses>;
using ConcreteShardRsm = Raft<SimulatorTransport, ShardRsm, WriteRequests, WriteResponses, ReadRequests, ReadResponses>;
// TODO(gvolfing) test vertex deletion with DETACH_DELETE as well
template <typename IoImpl>
void RunShardRaft(Raft<IoImpl, ShardRsm, WriteRequests, WriteResponses, ReadRequests, ReadResponses> server) {
server.Run();
}
namespace {
uint64_t GetTransactionId() {
static uint64_t transaction_id = 0;
return transaction_id++;
}
uint64_t GetUniqueInteger() {
static uint64_t prop_val_val = 1001;
return prop_val_val++;
}
LabelId get_primary_label() { return LabelId::FromUint(0); }
SchemaProperty get_schema_property() {
return {.property_id = PropertyId::FromUint(0), .type = common::SchemaType::INT};
}
msgs::PrimaryKey GetPrimaryKey(int64_t value) {
msgs::Value prop_val(static_cast<int64_t>(value));
msgs::PrimaryKey primary_key = {prop_val};
return primary_key;
}
msgs::NewVertex GetNewVertex(int64_t value) {
// Specify Labels.
msgs::Label label1 = {.id = LabelId::FromUint(1)};
std::vector<msgs::Label> label_ids = {label1};
// Specify primary key.
msgs::PrimaryKey primary_key = GetPrimaryKey(value);
// Specify properties
auto val1 = msgs::Value(static_cast<int64_t>(value));
auto prop1 = std::make_pair(PropertyId::FromUint(1), val1);
auto val3 = msgs::Value(static_cast<int64_t>(value));
auto prop3 = std::make_pair(PropertyId::FromUint(2), val3);
//(VERIFY) does the schema has to be specified with the properties or the primarykey?
auto val2 = msgs::Value(static_cast<int64_t>(value));
auto prop2 = std::make_pair(PropertyId::FromUint(0), val2);
std::vector<std::pair<PropertyId, msgs::Value>> properties{prop1, prop2, prop3};
// NewVertex
return {.label_ids = label_ids, .primary_key = primary_key, .properties = properties};
}
// TODO(gvolfing) maybe rename that something that makes sense.
std::vector<std::vector<msgs::Value>> GetValuePrimaryKeysWithValue(int64_t value) {
msgs::Value val(static_cast<int64_t>(value));
return {{val}};
}
} // namespace
// attempts to sending different requests
namespace {
bool AttemptToCreateVertex(ShardClient &client, int64_t value) {
msgs::NewVertex vertex = GetNewVertex(value);
auto create_req = msgs::CreateVerticesRequest{};
create_req.new_vertices = {vertex};
create_req.transaction_id.logical_id = GetTransactionId();
while (true) {
auto write_res = client.SendWriteRequest(create_req);
if (write_res.HasError()) {
continue;
}
auto write_response_result = write_res.GetValue();
auto write_response = std::get<msgs::CreateVerticesResponse>(write_response_result);
return write_response.success;
}
}
bool AttemptToDeleteVertex(ShardClient &client, int64_t value) {
auto delete_req = msgs::DeleteVerticesRequest{};
delete_req.deletion_type = msgs::DeleteVerticesRequest::DeletionType::DELETE;
delete_req.primary_keys = GetValuePrimaryKeysWithValue(value);
delete_req.transaction_id.logical_id = GetTransactionId();
while (true) {
auto write_res = client.SendWriteRequest(delete_req);
if (write_res.HasError()) {
continue;
}
auto write_response_result = write_res.GetValue();
auto write_response = std::get<msgs::DeleteVerticesResponse>(write_response_result);
return write_response.success;
}
}
bool AttemptToUpdateVertex(ShardClient &client, int64_t value) {
auto vertex_id = GetValuePrimaryKeysWithValue(value)[0];
std::vector<std::pair<PropertyId, msgs::Value>> property_updates;
auto property_update = std::make_pair(PropertyId::FromUint(2), msgs::Value(static_cast<int64_t>(10000)));
auto vertex_prop = msgs::UpdateVertexProp{};
vertex_prop.primary_key = vertex_id;
vertex_prop.property_updates = {property_update};
auto update_req = msgs::UpdateVerticesRequest{};
update_req.transaction_id.logical_id = GetTransactionId();
update_req.new_properties = {vertex_prop};
while (true) {
auto write_res = client.SendWriteRequest(update_req);
if (write_res.HasError()) {
continue;
}
auto write_response_result = write_res.GetValue();
auto write_response = std::get<msgs::UpdateVerticesResponse>(write_response_result);
return write_response.success;
}
}
bool AttemptToAddEdge(ShardClient &client, int64_t value_of_vertex_1, int64_t value_of_vertex_2, int64_t edge_gid,
int64_t edge_type_id) {
auto id = msgs::EdgeId{};
msgs::Label label = {.id = get_primary_label()};
auto src = std::make_pair(label, GetPrimaryKey(value_of_vertex_1));
auto dst = std::make_pair(label, GetPrimaryKey(value_of_vertex_2));
id.gid = edge_gid;
auto type = msgs::EdgeType{};
type.id = edge_type_id;
auto edge = msgs::Edge{};
edge.id = id;
edge.type = type;
edge.src = src;
edge.dst = dst;
edge.properties = std::nullopt;
msgs::CreateEdgesRequest create_req{};
create_req.edges = {edge};
create_req.transaction_id.logical_id = GetTransactionId();
while (true) {
auto write_res = client.SendWriteRequest(create_req);
if (write_res.HasError()) {
continue;
}
auto write_response_result = write_res.GetValue();
auto write_response = std::get<msgs::CreateEdgesResponse>(write_response_result);
return write_response.success;
}
}
bool AttemptToAddEdgeWithProperties(ShardClient &client, int64_t value_of_vertex_1, int64_t value_of_vertex_2,
int64_t edge_gid, uint64_t edge_prop_id, int64_t edge_prop_val,
const std::vector<uint64_t> &edge_type_id) {
auto id1 = msgs::EdgeId{};
msgs::Label label = {.id = get_primary_label()};
auto src = std::make_pair(label, GetPrimaryKey(value_of_vertex_1));
auto dst = std::make_pair(label, GetPrimaryKey(value_of_vertex_2));
id1.gid = edge_gid;
auto type1 = msgs::EdgeType{};
type1.id = edge_type_id[0];
auto edge_prop = std::make_pair(PropertyId::FromUint(edge_prop_id), msgs::Value(edge_prop_val));
auto edge = msgs::Edge{};
edge.id = id1;
edge.type = type1;
edge.src = src;
edge.dst = dst;
edge.properties = {edge_prop};
msgs::CreateEdgesRequest create_req{};
create_req.edges = {edge};
create_req.transaction_id.logical_id = GetTransactionId();
while (true) {
auto write_res = client.SendWriteRequest(create_req);
if (write_res.HasError()) {
continue;
}
auto write_response_result = write_res.GetValue();
auto write_response = std::get<msgs::CreateEdgesResponse>(write_response_result);
return write_response.success;
}
}
bool AttemptToDeleteEdge(ShardClient &client, int64_t value_of_vertex_1, int64_t value_of_vertex_2, int64_t edge_gid,
int64_t edge_type_id) {
auto id = msgs::EdgeId{};
msgs::Label label = {.id = get_primary_label()};
auto src = std::make_pair(label, GetPrimaryKey(value_of_vertex_1));
auto dst = std::make_pair(label, GetPrimaryKey(value_of_vertex_2));
id.gid = edge_gid;
auto type = msgs::EdgeType{};
type.id = edge_type_id;
auto edge = msgs::Edge{};
edge.id = id;
edge.type = type;
edge.src = {src};
edge.dst = {dst};
msgs::DeleteEdgesRequest delete_req{};
delete_req.edges = {edge};
delete_req.transaction_id.logical_id = GetTransactionId();
while (true) {
auto write_res = client.SendWriteRequest(delete_req);
if (write_res.HasError()) {
continue;
}
auto write_response_result = write_res.GetValue();
auto write_response = std::get<msgs::DeleteEdgesResponse>(write_response_result);
return write_response.success;
}
}
bool AttemptToUpdateEdge(ShardClient &client, int64_t value_of_vertex_1, int64_t value_of_vertex_2, int64_t edge_gid,
int64_t edge_type_id, uint64_t edge_prop_id, int64_t edge_prop_val) {
auto id = msgs::EdgeId{};
msgs::Label label = {.id = get_primary_label()};
auto src = std::make_pair(label, GetPrimaryKey(value_of_vertex_1));
auto dst = std::make_pair(label, GetPrimaryKey(value_of_vertex_2));
id.gid = edge_gid;
auto type = msgs::EdgeType{};
type.id = edge_type_id;
auto edge = msgs::Edge{};
edge.id = id;
edge.type = type;
auto edge_prop = std::vector<std::pair<PropertyId, msgs::Value>>{
std::make_pair(PropertyId::FromUint(edge_prop_id), msgs::Value(edge_prop_val))};
msgs::UpdateEdgeProp update_props{.src = src, .dst = dst, .edge_id = id, .property_updates = edge_prop};
msgs::UpdateEdgesRequest update_req{};
update_req.transaction_id.logical_id = GetTransactionId();
update_req.new_properties = {update_props};
while (true) {
auto write_res = client.SendWriteRequest(update_req);
if (write_res.HasError()) {
continue;
}
auto write_response_result = write_res.GetValue();
auto write_response = std::get<msgs::UpdateEdgesResponse>(write_response_result);
return write_response.success;
}
}
std::tuple<size_t, std::optional<msgs::VertexId>> AttemptToScanAllWithBatchLimit(ShardClient &client,
msgs::VertexId start_id,
uint64_t batch_limit) {
msgs::ScanVerticesRequest scan_req{};
scan_req.batch_limit = batch_limit;
scan_req.filter_expressions = std::nullopt;
scan_req.props_to_return = std::nullopt;
scan_req.start_id = start_id;
scan_req.storage_view = msgs::StorageView::OLD;
scan_req.transaction_id.logical_id = GetTransactionId();
while (true) {
auto read_res = client.SendReadRequest(scan_req);
if (read_res.HasError()) {
continue;
}
auto write_response_result = read_res.GetValue();
auto write_response = std::get<msgs::ScanVerticesResponse>(write_response_result);
MG_ASSERT(write_response.success);
return {write_response.results.size(), write_response.next_start_id};
}
}
void AttemptToExpandOneWithWrongEdgeType(ShardClient &client, uint64_t src_vertex_val, uint64_t edge_type_id) {
// Source vertex
msgs::Label label = {.id = get_primary_label()};
auto src_vertex = std::make_pair(label, GetPrimaryKey(src_vertex_val));
// Edge type
auto edge_type = msgs::EdgeType{};
edge_type.id = edge_type_id + 1;
// Edge direction
auto edge_direction = msgs::EdgeDirection::OUT;
// Source Vertex properties to look for
std::optional<std::vector<PropertyId>> src_vertex_properties = {};
// Edge properties to look for
std::optional<std::vector<PropertyId>> edge_properties = {};
std::vector<msgs::Expression> expressions;
std::optional<std::vector<msgs::OrderBy>> order_by = {};
std::optional<size_t> limit = {};
std::optional<msgs::Filter> filter = {};
msgs::ExpandOneRequest expand_one_req{};
expand_one_req.direction = edge_direction;
expand_one_req.edge_properties = edge_properties;
expand_one_req.edge_types = {edge_type};
expand_one_req.expressions = expressions;
expand_one_req.filter = filter;
expand_one_req.limit = limit;
expand_one_req.order_by = order_by;
expand_one_req.src_vertex_properties = src_vertex_properties;
expand_one_req.src_vertices = {src_vertex};
expand_one_req.transaction_id.logical_id = GetTransactionId();
while (true) {
auto read_res = client.SendReadRequest(expand_one_req);
if (read_res.HasError()) {
continue;
}
auto write_response_result = read_res.GetValue();
auto write_response = std::get<msgs::ExpandOneResponse>(write_response_result);
MG_ASSERT(write_response.result.size() == 1);
MG_ASSERT(write_response.result[0].edges_with_all_properties);
MG_ASSERT(write_response.result[0].edges_with_all_properties->size() == 0);
MG_ASSERT(!write_response.result[0].edges_with_specific_properties);
break;
}
}
void AttemptToExpandOneSimple(ShardClient &client, uint64_t src_vertex_val, uint64_t edge_type_id) {
// Source vertex
msgs::Label label = {.id = get_primary_label()};
auto src_vertex = std::make_pair(label, GetPrimaryKey(src_vertex_val));
// Edge type
auto edge_type = msgs::EdgeType{};
edge_type.id = edge_type_id;
// Edge direction
auto edge_direction = msgs::EdgeDirection::OUT;
// Source Vertex properties to look for
std::optional<std::vector<PropertyId>> src_vertex_properties = {};
// Edge properties to look for
std::optional<std::vector<PropertyId>> edge_properties = {};
std::vector<msgs::Expression> expressions;
std::optional<std::vector<msgs::OrderBy>> order_by = {};
std::optional<size_t> limit = {};
std::optional<msgs::Filter> filter = {};
msgs::ExpandOneRequest expand_one_req{};
expand_one_req.direction = edge_direction;
expand_one_req.edge_properties = edge_properties;
expand_one_req.edge_types = {edge_type};
expand_one_req.expressions = expressions;
expand_one_req.filter = filter;
expand_one_req.limit = limit;
expand_one_req.order_by = order_by;
expand_one_req.src_vertex_properties = src_vertex_properties;
expand_one_req.src_vertices = {src_vertex};
expand_one_req.transaction_id.logical_id = GetTransactionId();
while (true) {
auto read_res = client.SendReadRequest(expand_one_req);
if (read_res.HasError()) {
continue;
}
auto write_response_result = read_res.GetValue();
auto write_response = std::get<msgs::ExpandOneResponse>(write_response_result);
MG_ASSERT(write_response.result.size() == 1);
MG_ASSERT(write_response.result[0].edges_with_all_properties->size() == 2);
auto number_of_properties_on_edge =
(std::get<std::map<PropertyId, msgs::Value>>(write_response.result[0].edges_with_all_properties.value()[0]))
.size();
MG_ASSERT(number_of_properties_on_edge == 1);
break;
}
}
void AttemptToExpandOneWithSpecifiedSrcVertexProperties(ShardClient &client, uint64_t src_vertex_val,
uint64_t edge_type_id) {
// Source vertex
msgs::Label label = {.id = get_primary_label()};
auto src_vertex = std::make_pair(label, GetPrimaryKey(src_vertex_val));
// Edge type
auto edge_type = msgs::EdgeType{};
edge_type.id = edge_type_id;
// Edge direction
auto edge_direction = msgs::EdgeDirection::OUT;
// Source Vertex properties to look for
std::vector<PropertyId> desired_src_vertex_props{PropertyId::FromUint(2)};
std::optional<std::vector<PropertyId>> src_vertex_properties = desired_src_vertex_props;
// Edge properties to look for
std::optional<std::vector<PropertyId>> edge_properties = {};
std::vector<msgs::Expression> expressions;
std::optional<std::vector<msgs::OrderBy>> order_by = {};
std::optional<size_t> limit = {};
std::optional<msgs::Filter> filter = {};
msgs::ExpandOneRequest expand_one_req{};
expand_one_req.direction = edge_direction;
expand_one_req.edge_properties = edge_properties;
expand_one_req.edge_types = {edge_type};
expand_one_req.expressions = expressions;
expand_one_req.filter = filter;
expand_one_req.limit = limit;
expand_one_req.order_by = order_by;
expand_one_req.src_vertex_properties = src_vertex_properties;
expand_one_req.src_vertices = {src_vertex};
expand_one_req.transaction_id.logical_id = GetTransactionId();
while (true) {
auto read_res = client.SendReadRequest(expand_one_req);
if (read_res.HasError()) {
continue;
}
auto write_response_result = read_res.GetValue();
auto write_response = std::get<msgs::ExpandOneResponse>(write_response_result);
MG_ASSERT(write_response.result.size() == 1);
auto src_vertex_props_size = write_response.result[0].src_vertex_properties->size();
MG_ASSERT(src_vertex_props_size == 1);
MG_ASSERT(write_response.result[0].edges_with_all_properties->size() == 2);
auto number_of_properties_on_edge =
(std::get<std::map<PropertyId, msgs::Value>>(write_response.result[0].edges_with_all_properties.value()[0]))
.size();
MG_ASSERT(number_of_properties_on_edge == 1);
break;
}
}
void AttemptToExpandOneWithSpecifiedEdgeProperties(ShardClient &client, uint64_t src_vertex_val, uint64_t edge_type_id,
uint64_t edge_prop_id) {
// Source vertex
msgs::Label label = {.id = get_primary_label()};
auto src_vertex = std::make_pair(label, GetPrimaryKey(src_vertex_val));
// Edge type
auto edge_type = msgs::EdgeType{};
edge_type.id = edge_type_id;
// Edge direction
auto edge_direction = msgs::EdgeDirection::OUT;
// Source Vertex properties to look for
std::optional<std::vector<PropertyId>> src_vertex_properties = {};
// Edge properties to look for
std::vector<PropertyId> specified_edge_prop{PropertyId::FromUint(edge_prop_id)};
std::optional<std::vector<PropertyId>> edge_properties = {specified_edge_prop};
std::vector<msgs::Expression> expressions;
std::optional<std::vector<msgs::OrderBy>> order_by = {};
std::optional<size_t> limit = {};
std::optional<msgs::Filter> filter = {};
msgs::ExpandOneRequest expand_one_req{};
expand_one_req.direction = edge_direction;
expand_one_req.edge_properties = edge_properties;
expand_one_req.edge_types = {edge_type};
expand_one_req.expressions = expressions;
expand_one_req.filter = filter;
expand_one_req.limit = limit;
expand_one_req.order_by = order_by;
expand_one_req.src_vertex_properties = src_vertex_properties;
expand_one_req.src_vertices = {src_vertex};
expand_one_req.transaction_id.logical_id = GetTransactionId();
while (true) {
auto read_res = client.SendReadRequest(expand_one_req);
if (read_res.HasError()) {
continue;
}
auto write_response_result = read_res.GetValue();
auto write_response = std::get<msgs::ExpandOneResponse>(write_response_result);
MG_ASSERT(write_response.result.size() == 1);
auto specific_properties_size =
(std::get<std::vector<msgs::Value>>(write_response.result[0].edges_with_specific_properties.value()[0]));
MG_ASSERT(specific_properties_size.size() == 1);
break;
}
}
} // namespace
// tests
namespace {
void TestCreateVertices(ShardClient &client) { MG_ASSERT(AttemptToCreateVertex(client, GetUniqueInteger())); }
void TestCreateAndDeleteVertices(ShardClient &client) {
auto unique_prop_val = GetUniqueInteger();
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val));
MG_ASSERT(AttemptToDeleteVertex(client, unique_prop_val));
}
void TestCreateAndUpdateVertices(ShardClient &client) {
auto unique_prop_val = GetUniqueInteger();
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val));
MG_ASSERT(AttemptToUpdateVertex(client, unique_prop_val));
}
void TestCreateEdge(ShardClient &client) {
auto unique_prop_val_1 = GetUniqueInteger();
auto unique_prop_val_2 = GetUniqueInteger();
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_1));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_2));
auto edge_gid = GetUniqueInteger();
auto edge_type_id = GetUniqueInteger();
MG_ASSERT(AttemptToAddEdge(client, unique_prop_val_1, unique_prop_val_2, edge_gid, edge_type_id));
}
void TestCreateAndDeleteEdge(ShardClient &client) {
// Add the Edge
auto unique_prop_val_1 = GetUniqueInteger();
auto unique_prop_val_2 = GetUniqueInteger();
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_1));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_2));
auto edge_gid = GetUniqueInteger();
auto edge_type_id = GetUniqueInteger();
MG_ASSERT(AttemptToAddEdge(client, unique_prop_val_1, unique_prop_val_2, edge_gid, edge_type_id));
// Delete the Edge
MG_ASSERT(AttemptToDeleteEdge(client, unique_prop_val_1, unique_prop_val_2, edge_gid, edge_type_id));
}
void TestUpdateEdge(ShardClient &client) {
// Add the Edge
auto unique_prop_val_1 = GetUniqueInteger();
auto unique_prop_val_2 = GetUniqueInteger();
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_1));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_2));
auto edge_gid = GetUniqueInteger();
auto edge_type_id = GetUniqueInteger();
auto edge_prop_id = GetUniqueInteger();
auto edge_prop_val_old = GetUniqueInteger();
auto edge_prop_val_new = GetUniqueInteger();
MG_ASSERT(AttemptToAddEdgeWithProperties(client, unique_prop_val_1, unique_prop_val_2, edge_gid, edge_prop_id,
edge_prop_val_old, {edge_type_id}));
// Update the Edge
MG_ASSERT(AttemptToUpdateEdge(client, unique_prop_val_1, unique_prop_val_2, edge_gid, edge_type_id, edge_prop_id,
edge_prop_val_new));
}
void TestScanAllOneGo(ShardClient &client) {
auto unique_prop_val_1 = GetUniqueInteger();
auto unique_prop_val_2 = GetUniqueInteger();
auto unique_prop_val_3 = GetUniqueInteger();
auto unique_prop_val_4 = GetUniqueInteger();
auto unique_prop_val_5 = GetUniqueInteger();
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_1));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_2));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_3));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_4));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_5));
msgs::Label prim_label = {.id = get_primary_label()};
msgs::PrimaryKey prim_key = {msgs::Value(static_cast<int64_t>(unique_prop_val_1))};
msgs::VertexId v_id = {prim_label, prim_key};
auto [result_size, next_id] = AttemptToScanAllWithBatchLimit(client, v_id, 5);
MG_ASSERT(result_size == 5);
}
void TestScanAllWithSmallBatchSize(ShardClient &client) {
auto unique_prop_val_1 = GetUniqueInteger();
auto unique_prop_val_2 = GetUniqueInteger();
auto unique_prop_val_3 = GetUniqueInteger();
auto unique_prop_val_4 = GetUniqueInteger();
auto unique_prop_val_5 = GetUniqueInteger();
auto unique_prop_val_6 = GetUniqueInteger();
auto unique_prop_val_7 = GetUniqueInteger();
auto unique_prop_val_8 = GetUniqueInteger();
auto unique_prop_val_9 = GetUniqueInteger();
auto unique_prop_val_10 = GetUniqueInteger();
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_1));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_2));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_3));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_4));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_5));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_6));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_7));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_8));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_9));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_10));
msgs::Label prim_label = {.id = get_primary_label()};
msgs::PrimaryKey prim_key1 = {msgs::Value(static_cast<int64_t>(unique_prop_val_1))};
msgs::VertexId v_id_1 = {prim_label, prim_key1};
auto [result_size1, next_id1] = AttemptToScanAllWithBatchLimit(client, v_id_1, 3);
MG_ASSERT(result_size1 == 3);
auto [result_size2, next_id2] = AttemptToScanAllWithBatchLimit(client, next_id1.value(), 3);
MG_ASSERT(result_size2 == 3);
auto [result_size3, next_id3] = AttemptToScanAllWithBatchLimit(client, next_id2.value(), 3);
MG_ASSERT(result_size3 == 3);
auto [result_size4, next_id4] = AttemptToScanAllWithBatchLimit(client, next_id3.value(), 3);
MG_ASSERT(result_size4 == 1);
MG_ASSERT(!next_id4);
}
void TestExpandOne(ShardClient &client) {
// ExpandOneSimple
auto unique_prop_val_1 = GetUniqueInteger();
auto unique_prop_val_2 = GetUniqueInteger();
auto unique_prop_val_3 = GetUniqueInteger();
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_1));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_2));
MG_ASSERT(AttemptToCreateVertex(client, unique_prop_val_3));
auto edge_type_id = GetUniqueInteger();
auto edge_gid_1 = GetUniqueInteger();
auto edge_gid_2 = GetUniqueInteger();
auto edge_prop_id = GetUniqueInteger();
auto edge_prop_val = GetUniqueInteger();
// (V1)-[edge_type_id]->(V2)
MG_ASSERT(AttemptToAddEdgeWithProperties(client, unique_prop_val_1, unique_prop_val_2, edge_gid_1, edge_prop_id,
edge_prop_val, {edge_type_id}));
// (V1)-[edge_type_id]->(V3)
MG_ASSERT(AttemptToAddEdgeWithProperties(client, unique_prop_val_1, unique_prop_val_3, edge_gid_2, edge_prop_id,
edge_prop_val, {edge_type_id}));
AttemptToExpandOneSimple(client, unique_prop_val_1, edge_type_id);
AttemptToExpandOneWithWrongEdgeType(client, unique_prop_val_1, edge_type_id);
AttemptToExpandOneWithSpecifiedSrcVertexProperties(client, unique_prop_val_1, edge_type_id);
AttemptToExpandOneWithSpecifiedEdgeProperties(client, unique_prop_val_1, edge_type_id, edge_prop_id);
}
} // namespace
int TestMessages() {
SimulatorConfig config{
.drop_percent = 0,
.perform_timeouts = false,
.scramble_messages = false,
.rng_seed = 0,
.start_time = Time::min() + std::chrono::microseconds{256 * 1024},
.abort_time = Time::min() + std::chrono::microseconds{4 * 8 * 1024 * 1024},
};
auto simulator = Simulator(config);
Io<SimulatorTransport> shard_server_io_1 = simulator.RegisterNew();
const auto shard_server_1_address = shard_server_io_1.GetAddress();
Io<SimulatorTransport> shard_server_io_2 = simulator.RegisterNew();
const auto shard_server_2_address = shard_server_io_2.GetAddress();
Io<SimulatorTransport> shard_server_io_3 = simulator.RegisterNew();
const auto shard_server_3_address = shard_server_io_3.GetAddress();
Io<SimulatorTransport> shard_client_io = simulator.RegisterNew();
PropertyValue min_pk(static_cast<int64_t>(0));
std::vector<PropertyValue> min_prim_key = {min_pk};
PropertyValue max_pk(static_cast<int64_t>(10000000));
std::vector<PropertyValue> max_prim_key = {max_pk};
auto shard_ptr1 = std::make_unique<Shard>(get_primary_label(), min_prim_key, max_prim_key);
auto shard_ptr2 = std::make_unique<Shard>(get_primary_label(), min_prim_key, max_prim_key);
auto shard_ptr3 = std::make_unique<Shard>(get_primary_label(), min_prim_key, max_prim_key);
shard_ptr1->CreateSchema(get_primary_label(), {get_schema_property()});
shard_ptr2->CreateSchema(get_primary_label(), {get_schema_property()});
shard_ptr3->CreateSchema(get_primary_label(), {get_schema_property()});
std::vector<Address> address_for_1{shard_server_2_address, shard_server_3_address};
std::vector<Address> address_for_2{shard_server_1_address, shard_server_3_address};
std::vector<Address> address_for_3{shard_server_1_address, shard_server_2_address};
ConcreteShardRsm shard_server1(std::move(shard_server_io_1), address_for_1, ShardRsm(std::move(shard_ptr1)));
ConcreteShardRsm shard_server2(std::move(shard_server_io_2), address_for_2, ShardRsm(std::move(shard_ptr2)));
ConcreteShardRsm shard_server3(std::move(shard_server_io_3), address_for_3, ShardRsm(std::move(shard_ptr3)));
auto server_thread1 = std::jthread([&shard_server1]() { shard_server1.Run(); });
auto server_thread2 = std::jthread([&shard_server2]() { shard_server2.Run(); });
auto server_thread3 = std::jthread([&shard_server3]() { shard_server3.Run(); });
simulator.IncrementServerCountAndWaitForQuiescentState(shard_server_1_address);
simulator.IncrementServerCountAndWaitForQuiescentState(shard_server_2_address);
simulator.IncrementServerCountAndWaitForQuiescentState(shard_server_3_address);
std::cout << "Beginning test after servers have become quiescent." << std::endl;
std::vector server_addrs = {shard_server_1_address, shard_server_2_address, shard_server_3_address};
ShardClient client(shard_client_io, shard_server_1_address, server_addrs);
// Vertex tests
TestCreateVertices(client);
TestCreateAndDeleteVertices(client);
TestCreateAndUpdateVertices(client);
// Edge tests
TestCreateEdge(client);
TestCreateAndDeleteEdge(client);
TestUpdateEdge(client);
// ScanAll tests
TestScanAllOneGo(client);
TestScanAllWithSmallBatchSize(client);
// ExpandOne tests
TestExpandOne(client);
simulator.ShutDown();
SimulatorStats stats = simulator.Stats();
std::cout << "total messages: " << stats.total_messages << std::endl;
std::cout << "dropped messages: " << stats.dropped_messages << std::endl;
std::cout << "timed out requests: " << stats.timed_out_requests << std::endl;
std::cout << "total requests: " << stats.total_requests << std::endl;
std::cout << "total responses: " << stats.total_responses << std::endl;
std::cout << "simulator ticks: " << stats.simulator_ticks << std::endl;
std::cout << "========================== SUCCESS :) ==========================" << std::endl;
return 0;
}
} // namespace memgraph::storage::v3::tests
int main() { return memgraph::storage::v3::tests::TestMessages(); }