From 2e1a717dcbcd18ce8252ddccf5ea4628e6a1b4cc Mon Sep 17 00:00:00 2001 From: Kostas Kyrimis Date: Tue, 6 Jul 2021 18:44:47 +0300 Subject: [PATCH] Add telemetry for streams and triggers (#188) * Add event counters to streams and triggers --- src/query/interpreter.cpp | 5 +++++ src/query/streams.cpp | 6 ++++++ src/query/trigger.cpp | 6 ++++++ src/utils/event_counter.cpp | 6 +++++- 4 files changed, 22 insertions(+), 1 deletion(-) diff --git a/src/query/interpreter.cpp b/src/query/interpreter.cpp index 3cbd8d7aa..2fb049626 100644 --- a/src/query/interpreter.cpp +++ b/src/query/interpreter.cpp @@ -43,6 +43,9 @@ extern Event ReadWriteQuery; extern const Event LabelIndexCreated; extern const Event LabelPropertyIndexCreated; + +extern const Event StreamsCreated; +extern const Event TriggersCreated; } // namespace EventCounter namespace query { @@ -476,6 +479,7 @@ Callback HandleStreamQuery(StreamQuery *stream_query, const Parameters ¶mete Callback callback; switch (stream_query->action_) { case StreamQuery::Action::CREATE_STREAM: { + EventCounter::IncrementCounter(EventCounter::StreamsCreated); constexpr std::string_view kDefaultConsumerGroup = "mg_consumer"; std::string consumer_group{stream_query->consumer_group_.empty() ? kDefaultConsumerGroup : stream_query->consumer_group_}; @@ -1282,6 +1286,7 @@ PreparedQuery PrepareTriggerQuery(ParsedQuery parsed_query, const bool in_explic auto callback = [trigger_query, interpreter_context, dba, &user_parameters] { switch (trigger_query->action_) { case TriggerQuery::Action::CREATE_TRIGGER: + EventCounter::IncrementCounter(EventCounter::TriggersCreated); return CreateTrigger(trigger_query, user_parameters, interpreter_context, dba); case TriggerQuery::Action::DROP_TRIGGER: return DropTrigger(trigger_query, interpreter_context); diff --git a/src/query/streams.cpp b/src/query/streams.cpp index 6f2250b34..44b3cce1e 100644 --- a/src/query/streams.cpp +++ b/src/query/streams.cpp @@ -12,10 +12,15 @@ #include "query/procedure/mg_procedure_impl.hpp" #include "query/procedure/module.hpp" #include "query/typed_value.hpp" +#include "utils/event_counter.hpp" #include "utils/memory.hpp" #include "utils/on_scope_exit.hpp" #include "utils/pmr/string.hpp" +namespace EventCounter { +extern const Event MessagesConsumed; +} // namespace EventCounter + namespace query { using Consumer = integrations::kafka::Consumer; @@ -335,6 +340,7 @@ Streams::StreamsMap::iterator Streams::CreateConsumer(StreamsMap &map, const std result = mgp_result{nullptr, memory_resource}]( const std::vector &messages) mutable { auto accessor = interpreter_context->db->Access(); + EventCounter::IncrementCounter(EventCounter::MessagesConsumed, messages.size()); CallCustomTransformation(transformation_name, messages, result, accessor, *memory_resource, stream_name); DiscardValueResultStream stream; diff --git a/src/query/trigger.cpp b/src/query/trigger.cpp index 1c5bd39ff..edd4a08e9 100644 --- a/src/query/trigger.cpp +++ b/src/query/trigger.cpp @@ -11,8 +11,13 @@ #include "query/serialization/property_value.hpp" #include "query/typed_value.hpp" #include "storage/v2/property_value.hpp" +#include "utils/event_counter.hpp" #include "utils/memory.hpp" +namespace EventCounter { +extern const Event TriggersExecuted; +} // namespace EventCounter + namespace query { namespace { auto IdentifierString(const TriggerIdentifierTag tag) noexcept { @@ -228,6 +233,7 @@ void Trigger::Execute(DbAccessor *dba, utils::MonotonicBufferResource *execution ; cursor->Shutdown(); + EventCounter::IncrementCounter(EventCounter::TriggersExecuted); } namespace { diff --git a/src/utils/event_counter.cpp b/src/utils/event_counter.cpp index 9ff4be604..f399c626b 100644 --- a/src/utils/event_counter.cpp +++ b/src/utils/event_counter.cpp @@ -41,7 +41,11 @@ \ M(FailedQuery, "Number of times executing a query failed.") \ M(LabelIndexCreated, "Number of times a label index was created.") \ - M(LabelPropertyIndexCreated, "Number of times a label property index was created.") + M(LabelPropertyIndexCreated, "Number of times a label property index was created.") \ + M(StreamsCreated, "Number of Streams created.") \ + M(MessagesConsumed, "Number of consumed streamed messages.") \ + M(TriggersCreated, "Number of Triggers created.") \ + M(TriggersExecuted, "Number of Triggers executed.") namespace EventCounter {