This is an automated email from the ASF dual-hosted git repository.

djwang pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/cloudberry.git

commit f9793da16cb4d72baa18a277145181f7fc7e121b
Author: Maxim Smyatkin <[email protected]>
AuthorDate: Wed Sep 6 16:11:06 2023 +0300

    [yagp_hooks_collector] Clean up threading, signal handling, and logging
    
    Mute PG-destined signals in GRPC reconnection thread.  Move debian
    config to CI.  Redirect debug output to log file.  Harden memory
    handling.  Remove thread-unsafe logging and dead code.
---
 debian/compat          |  1 -
 debian/control         | 11 -------
 debian/postinst        |  8 -----
 debian/rules           | 10 ------
 src/EventSender.cpp    | 85 ++++++++++++++++++++++----------------------------
 src/EventSender.h      |  1 -
 src/GrpcConnector.cpp  | 85 +++++++++++++++++++++++++++++++++-----------------
 src/GrpcConnector.h    |  3 +-
 src/SpillInfoWrapper.c | 21 -------------
 9 files changed, 97 insertions(+), 128 deletions(-)

diff --git a/debian/compat b/debian/compat
deleted file mode 100644
index ec635144f60..00000000000
--- a/debian/compat
+++ /dev/null
@@ -1 +0,0 @@
-9
diff --git a/debian/control b/debian/control
deleted file mode 100644
index 07176e94be5..00000000000
--- a/debian/control
+++ /dev/null
@@ -1,11 +0,0 @@
-Source: greenplum-6-yagpcc-hooks
-Section: misc
-Priority: optional
-Maintainer: Maxim Smyatkin <[email protected]>
-Build-Depends: make, gcc, g++, debhelper (>=9), greenplum-db-6 (>=6.19.3), 
ya-grpc (=1.46-57-50820-02384e3918-yandex)
-Standards-Version: 3.9.8
-
-Package: greenplum-6-yagpcc-hooks
-Architecture: any
-Depends: ${misc:Depends}, ${shlibs:Depends}, greenplum-db-6 (>=6.19.3), 
ya-grpc (=1.46-57-50820-02384e3918-yandex)
-Description: Greenplum extension to send query execution metrics to yandex 
command center agent
diff --git a/debian/postinst b/debian/postinst
deleted file mode 100644
index 27ddfc06a7d..00000000000
--- a/debian/postinst
+++ /dev/null
@@ -1,8 +0,0 @@
-#!/bin/bash
-
-set -e
-
-GPADMIN=gpadmin
-GPHOME=/opt/greenplum-db-6
-
-chown -R ${GPADMIN}:${GPADMIN} ${GPHOME}
diff --git a/debian/rules b/debian/rules
deleted file mode 100644
index 6c2c7491067..00000000000
--- a/debian/rules
+++ /dev/null
@@ -1,10 +0,0 @@
-#!/usr/bin/make -f
-# You must remove unused comment lines for the released package.
-export DH_VERBOSE = 1
-
-
-export GPHOME := /opt/greenplum-db-6
-export PATH := $(GPHOME)/bin:$(PATH)
-
-%:
-       dh $@ 
diff --git a/src/EventSender.cpp b/src/EventSender.cpp
index 2810e581313..57fe6f13391 100644
--- a/src/EventSender.cpp
+++ b/src/EventSender.cpp
@@ -1,8 +1,8 @@
 #include "Config.h"
 #include "GrpcConnector.h"
 #include "ProcStats.h"
-#include <ctime>
 #include <chrono>
+#include <ctime>
 
 #define typeid __typeid
 #define operator __operator
@@ -20,14 +20,11 @@ extern "C" {
 
 #include "cdb/cdbdisp.h"
 #include "cdb/cdbexplain.h"
-#include "cdb/cdbvars.h"
 #include "cdb/cdbinterconnect.h"
+#include "cdb/cdbvars.h"
 
 #include "stat_statements_parser/pg_stat_statements_ya_parser.h"
 #include "tcop/utility.h"
-
-void get_spill_info(int ssid, int ccid, int32_t *file_count,
-                    int64_t *total_bytes);
 }
 #undef typeid
 #undef operator
@@ -48,7 +45,6 @@ std::string *get_user_name() {
 std::string *get_db_name() {
   char *dbname = get_database_name(MyDatabaseId);
   std::string *result = dbname ? new std::string(dbname) : nullptr;
-  pfree(dbname);
   return result;
 }
 
@@ -63,7 +59,6 @@ std::string *get_rg_name() {
   if (rgname == nullptr)
     return nullptr;
   auto result = new std::string(rgname);
-  pfree(rgname);
   return result;
 }
 
@@ -114,14 +109,12 @@ void set_query_plan(yagpcc::QueryInfo *qi, QueryDesc 
*query_desc) {
   StringInfo norm_plan = gen_normplan(qi->plan_text().c_str());
   *qi->mutable_template_plan_text() = std::string(norm_plan->data);
   qi->set_plan_id(hash_any((unsigned char *)norm_plan->data, norm_plan->len));
-  // TODO: free stringinfo?
 }
 
 void set_query_text(yagpcc::QueryInfo *qi, QueryDesc *query_desc) {
   *qi->mutable_query_text() = query_desc->sourceText;
   char *norm_query = gen_normquery(query_desc->sourceText);
   *qi->mutable_template_query_text() = std::string(norm_query);
-  pfree(norm_query);
 }
 
 void set_query_info(yagpcc::SetQueryReq *req, QueryDesc *query_desc,
@@ -182,16 +175,7 @@ void 
set_metric_instrumentation(yagpcc::MetricInstrumentation *metrics,
 
 decltype(std::chrono::high_resolution_clock::now()) query_start_time;
 
-void set_gp_metrics(yagpcc::GPMetrics *metrics, QueryDesc *query_desc,
-                    bool need_spillinfo) {
-  if (need_spillinfo) {
-    int32_t n_spill_files = 0;
-    int64_t n_spill_bytes = 0;
-    get_spill_info(gp_session_id, gp_command_count, &n_spill_files,
-                   &n_spill_bytes);
-    metrics->mutable_spill()->set_filecount(n_spill_files);
-    metrics->mutable_spill()->set_totalbytes(n_spill_bytes);
-  }
+void set_gp_metrics(yagpcc::GPMetrics *metrics, QueryDesc *query_desc) {
   if (query_desc->planstate && query_desc->planstate->instrument) {
     set_metric_instrumentation(metrics->mutable_instrumentation(), query_desc);
   }
@@ -254,6 +238,9 @@ void EventSender::query_metrics_collect(QueryMetricsStatus 
status, void *arg) {
 
 void EventSender::executor_before_start(QueryDesc *query_desc,
                                         int /* eflags*/) {
+  if (!connector) {
+    return;
+  }
   if (!need_collect()) {
     return;
   }
@@ -275,71 +262,75 @@ void EventSender::executor_before_start(QueryDesc 
*query_desc,
 }
 
 void EventSender::executor_after_start(QueryDesc *query_desc, int /* eflags*/) 
{
+  if (!connector) {
+    return;
+  }
   if ((Gp_role == GP_ROLE_DISPATCH || Gp_role == GP_ROLE_EXECUTE) &&
       need_collect()) {
     auto req =
         create_query_req(query_desc, yagpcc::QueryStatus::QUERY_STATUS_START);
     set_query_info(&req, query_desc, false, true);
-    send_query_info(&req, "started");
+    connector->set_metric_query(req, "started");
   }
 }
 
 void EventSender::executor_end(QueryDesc *query_desc) {
+  if (!connector) {
+    return;
+  }
   if (!need_collect() ||
       (Gp_role != GP_ROLE_DISPATCH && Gp_role != GP_ROLE_EXECUTE)) {
     return;
   }
-  if (query_desc->totaltime && Config::enable_analyze() &&
-      Config::enable_cdbstats()) {
-    if (query_desc->estate->dispatcherState &&
-        query_desc->estate->dispatcherState->primaryResults) {
-      cdbdisp_checkDispatchResult(query_desc->estate->dispatcherState,
-                                  DISPATCH_WAIT_NONE);
-    }
-    InstrEndLoop(query_desc->totaltime);
-  }
+  /* TODO: when querying via CURSOR this call freezes. Need to investigate.
+     To reproduce - uncomment it and run installchecks. It will freeze around 
join test.
+     Needs investigation
+    
+    if (Gp_role == GP_ROLE_DISPATCH && Config::enable_analyze() &&
+      Config::enable_cdbstats() && query_desc->estate->dispatcherState &&
+      query_desc->estate->dispatcherState->primaryResults) {
+    cdbdisp_checkDispatchResult(query_desc->estate->dispatcherState,
+                                DISPATCH_WAIT_NONE);
+  }*/
   auto req =
       create_query_req(query_desc, yagpcc::QueryStatus::QUERY_STATUS_END);
   // NOTE: there are no cummulative spillinfo stats AFAIU, so no need to
   // gather it here. It only makes sense when doing regular stat checks.
-  set_gp_metrics(req.mutable_query_metrics(), query_desc,
-                 /*need_spillinfo*/ false);
-  send_query_info(&req, "ended");
+  set_gp_metrics(req.mutable_query_metrics(), query_desc);
+  connector->set_metric_query(req, "ended");
 }
 
 void EventSender::collect_query_submit(QueryDesc *query_desc) {
+  if (!connector) {
+    return;
+  }
   if (need_collect()) {
     auto req =
         create_query_req(query_desc, yagpcc::QueryStatus::QUERY_STATUS_SUBMIT);
     set_query_info(&req, query_desc, true, false);
-    send_query_info(&req, "submit");
+    connector->set_metric_query(req, "submit");
   }
 }
 
 void EventSender::collect_query_done(QueryDesc *query_desc,
                                      const std::string &status) {
+  if (!connector) {
+    return;
+  }
   if (need_collect()) {
     auto req =
         create_query_req(query_desc, yagpcc::QueryStatus::QUERY_STATUS_DONE);
-    send_query_info(&req, status);
-  }
-}
-
-void EventSender::send_query_info(yagpcc::SetQueryReq *req,
-                                  const std::string &event) {
-  auto result = connector->set_metric_query(*req);
-  if (result.error_code() == yagpcc::METRIC_RESPONSE_STATUS_CODE_ERROR) {
-    ereport(WARNING,
-            (errmsg("Query {%d-%d-%d} %s reporting failed with an error %s",
-                    req->query_key().tmid(), req->query_key().ssid(),
-                    req->query_key().ccnt(), event.c_str(),
-                    result.error_text().c_str())));
+    connector->set_metric_query(req, status);
   }
 }
 
 EventSender::EventSender() {
   if (Config::enable_collector()) {
-    connector = new GrpcConnector();
+    try {
+      connector = new GrpcConnector();
+    } catch (const std::exception &e) {
+      ereport(INFO, (errmsg("Unable to start query tracing %s", e.what())));
+    }
   }
 }
 
diff --git a/src/EventSender.h b/src/EventSender.h
index f53648bed36..ee0db2f0938 100644
--- a/src/EventSender.h
+++ b/src/EventSender.h
@@ -23,7 +23,6 @@ public:
 private:
   void collect_query_submit(QueryDesc *query_desc);
   void collect_query_done(QueryDesc *query_desc, const std::string &status);
-  void send_query_info(yagpcc::SetQueryReq *req, const std::string &event);
   GrpcConnector *connector;
   int nesting_level = 0;
 };
\ No newline at end of file
diff --git a/src/GrpcConnector.cpp b/src/GrpcConnector.cpp
index 966bfb4a780..73c1944fa04 100644
--- a/src/GrpcConnector.cpp
+++ b/src/GrpcConnector.cpp
@@ -7,45 +7,72 @@
 #include <grpc++/channel.h>
 #include <grpc++/grpc++.h>
 #include <mutex>
+#include <pthread.h>
+#include <signal.h>
 #include <string>
 #include <thread>
 
-extern "C"
-{
+extern "C" {
 #include "postgres.h"
 #include "cdb/cdbvars.h"
 }
 
-class GrpcConnector::Impl
-{
+/*
+ * Set up the thread signal mask, we don't want to run our signal handlers
+ * in downloading and uploading threads.
+ */
+static void MaskThreadSignals() {
+  sigset_t sigs;
+
+  if (pthread_equal(main_tid, pthread_self())) {
+    ereport(ERROR, (errmsg("thread_mask is called from main thread!")));
+    return;
+  }
+
+  sigemptyset(&sigs);
+
+  /* make our thread to ignore these signals (which should allow that they be
+   * delivered to the main thread) */
+  sigaddset(&sigs, SIGHUP);
+  sigaddset(&sigs, SIGINT);
+  sigaddset(&sigs, SIGTERM);
+  sigaddset(&sigs, SIGALRM);
+  sigaddset(&sigs, SIGUSR1);
+  sigaddset(&sigs, SIGUSR2);
+
+  pthread_sigmask(SIG_BLOCK, &sigs, NULL);
+}
+
+class GrpcConnector::Impl {
 public:
-  Impl() : SOCKET_FILE("unix://" + Config::uds_path())
-  {
+  Impl() : SOCKET_FILE("unix://" + Config::uds_path()) {
     GOOGLE_PROTOBUF_VERIFY_VERSION;
     channel =
         grpc::CreateChannel(SOCKET_FILE, grpc::InsecureChannelCredentials());
     stub = yagpcc::SetQueryInfo::NewStub(channel);
     connected = true;
+    reconnected = false;
     done = false;
     reconnect_thread = std::thread(&Impl::reconnect, this);
   }
 
-  ~Impl()
-  {
+  ~Impl() {
     done = true;
     cv.notify_one();
     reconnect_thread.join();
   }
 
-  yagpcc::MetricResponse set_metric_query(yagpcc::SetQueryReq req)
-  {
+  yagpcc::MetricResponse set_metric_query(const yagpcc::SetQueryReq &req,
+                                          const std::string &event) {
     yagpcc::MetricResponse response;
-    if (!connected)
-    {
+    if (!connected) {
       response.set_error_code(yagpcc::METRIC_RESPONSE_STATUS_CODE_ERROR);
       response.set_error_text(
-          "Not tracing this query connection to agent has been lost");
+          "Not tracing this query because grpc connection has been lost");
       return response;
+    } else if (reconnected) {
+      reconnected = false;
+      ereport(LOG, (errmsg("GRPC connection is restored")));
     }
     grpc::ClientContext context;
     int timeout = Gp_role == GP_ROLE_DISPATCH ? 500 : 250;
@@ -53,12 +80,16 @@ public:
         std::chrono::system_clock::now() + std::chrono::milliseconds(timeout);
     context.set_deadline(deadline);
     grpc::Status status = (stub->SetMetricQuery)(&context, req, &response);
-    if (!status.ok())
-    {
-      response.set_error_text("Connection lost: " + status.error_message() +
-                              "; " + status.error_details());
+    if (!status.ok()) {
+      response.set_error_text("GRPC error: " + status.error_message() + "; " +
+                              status.error_details());
       response.set_error_code(yagpcc::METRIC_RESPONSE_STATUS_CODE_ERROR);
+      ereport(LOG, (errmsg("Query {%d-%d-%d} %s tracing failed with error %s",
+                           req.query_key().tmid(), req.query_key().ssid(),
+                           req.query_key().ccnt(), event.c_str(),
+                           response.error_text().c_str())));
       connected = false;
+      reconnected = false;
       cv.notify_one();
     }
 
@@ -69,25 +100,23 @@ private:
   const std::string SOCKET_FILE;
   std::unique_ptr<yagpcc::SetQueryInfo::Stub> stub;
   std::shared_ptr<grpc::Channel> channel;
-  std::atomic_bool connected;
+  std::atomic_bool connected, reconnected, done;
   std::thread reconnect_thread;
   std::condition_variable cv;
   std::mutex mtx;
-  bool done;
 
-  void reconnect()
-  {
-    while (!done)
-    {
+  void reconnect() {
+    MaskThreadSignals();
+    while (!done) {
       {
         std::unique_lock<std::mutex> lock(mtx);
         cv.wait(lock);
       }
-      while (!connected && !done)
-      {
+      while (!connected && !done) {
         auto deadline =
             std::chrono::system_clock::now() + std::chrono::milliseconds(100);
         connected = channel->WaitForConnected(deadline);
+        reconnected = connected.load();
       }
     }
   }
@@ -98,7 +127,7 @@ GrpcConnector::GrpcConnector() { impl = new Impl(); }
 GrpcConnector::~GrpcConnector() { delete impl; }
 
 yagpcc::MetricResponse
-GrpcConnector::set_metric_query(yagpcc::SetQueryReq req)
-{
-  return impl->set_metric_query(req);
+GrpcConnector::set_metric_query(const yagpcc::SetQueryReq &req,
+                                const std::string &event) {
+  return impl->set_metric_query(req, event);
 }
\ No newline at end of file
diff --git a/src/GrpcConnector.h b/src/GrpcConnector.h
index 4fca6960a4e..6571c626dfd 100644
--- a/src/GrpcConnector.h
+++ b/src/GrpcConnector.h
@@ -6,7 +6,8 @@ class GrpcConnector {
 public:
   GrpcConnector();
   ~GrpcConnector();
-  yagpcc::MetricResponse set_metric_query(yagpcc::SetQueryReq req);
+  yagpcc::MetricResponse set_metric_query(const yagpcc::SetQueryReq &req,
+                                          const std::string &event);
 
 private:
   class Impl;
diff --git a/src/SpillInfoWrapper.c b/src/SpillInfoWrapper.c
deleted file mode 100644
index c6ace0a693f..00000000000
--- a/src/SpillInfoWrapper.c
+++ /dev/null
@@ -1,21 +0,0 @@
-#include "postgres.h"
-#include "utils/workfile_mgr.h"
-
-void get_spill_info(int ssid, int ccid, int32_t* file_count, int64_t* 
total_bytes);
-
-void get_spill_info(int ssid, int ccid, int32_t* file_count, int64_t* 
total_bytes)
-{
-    int count = 0;
-    int i = 0;
-    workfile_set *workfiles = workfile_mgr_cache_entries_get_copy(&count);
-    workfile_set *wf_iter = workfiles;
-    for (i = 0; i < count; ++i, ++wf_iter)
-    {
-        if (wf_iter->active && wf_iter->session_id == ssid && 
wf_iter->command_count == ccid)
-        {
-            *file_count += wf_iter->num_files;
-            *total_bytes += wf_iter->total_bytes;
-        }
-    }
-    pfree(workfiles);
-}
\ No newline at end of file


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to