Dynamic memory cleanup for Pub/Sub exporting connector (#9112)
* Update dynamic memory cleanup * Fix unit tests * Rename a function * Delete GRPC objects * Unlock a mutex * Delete an odd file
Vladimir Kobal committed
May 21, 2020 at 23:29 UTC
e7428a1d31da1c695cd869e0bcc8dd535383b10e
5 files changed
+70
-3
exporting/pubsub/pubsub.c
+38
-2
@@ -56,6 +56,35 @@ int init_pubsub_instance(struct instance *instance)
56
return 0;
57
}
58
59
+/**
60
+ * Clean a PubSub connector instance
61
+ *
62
+ * @param instance an instance data structure.
63
+ */
64
+void clean_pubsub_instance(struct instance *instance)
65
+{
66
+ info("EXPORTING: cleaning up instance %s ...", instance->config.name);
67
+
68
+ struct pubsub_specific_data *connector_specific_data =
69
+ (struct pubsub_specific_data *)instance->connector_specific_data;
70
+ pubsub_cleanup(connector_specific_data);
71
+ freez(connector_specific_data);
72
+
73
+ buffer_free(instance->buffer);
74
+
75
+ struct pubsub_specific_config *connector_specific_config =
76
+ (struct pubsub_specific_config *)instance->config.connector_specific_config;
77
+ freez(connector_specific_config->credentials_file);
78
+ freez(connector_specific_config->project_id);
79
+ freez(connector_specific_config->topic_id);
80
+ freez(connector_specific_config);
81
+
82
+ info("EXPORTING: instance %s exited", instance->config.name);
83
+ instance->exited = 1;
84
+
85
+ return;
86
+}
87
+
88
/**
89
* Pub/Sub connector worker
90
*
@@ -69,13 +98,18 @@ void pubsub_connector_worker(void *instance_p)
98
struct pubsub_specific_config *connector_specific_config = instance->config.connector_specific_config;
99
struct pubsub_specific_data *connector_specific_data = instance->connector_specific_data;
100
72
- while (!netdata_exit) {
101
+ while (!instance->engine->exit) {
102
struct stats *stats = &instance->stats;
103
char error_message[ERROR_LINE_MAX + 1] = "";
104
105
uv_mutex_lock(&instance->mutex);
106
uv_cond_wait(&instance->cond_var, &instance->mutex);
107
108
+ if (unlikely(instance->engine->exit)) {
109
+ uv_mutex_unlock(&instance->mutex);
110
+ break;
111
+ }
112
+
113
// reset the monitoring chart counters
114
stats->received_bytes =
115
stats->sent_bytes =
@@ -149,7 +183,9 @@ void pubsub_connector_worker(void *instance_p)
183
uv_mutex_unlock(&instance->mutex);
184
185
#ifdef UNIT_TESTING
152
- break;
186
+ return;
187
#endif
188
}
189
+
190
+ clean_pubsub_instance(instance);
191
}
exporting/pubsub/pubsub.h
+1
@@ -8,6 +8,7 @@
8
#include "pubsub_publish.h"
9
10
int init_pubsub_instance(struct instance *instance);
11
+void clean_pubsub_instance(struct instance *instance);
12
void pubsub_connector_worker(void *instance_p);
13
14
#endif //NETDATA_EXPORTING_PUBSUB_H
exporting/pubsub/pubsub_publish.cc
+29
@@ -79,6 +79,35 @@ int pubsub_init(
79
return 0;
80
}
81
82
+/**
83
+ * Clean the PubSub connector instance specific data
84
+ */
85
+void pubsub_cleanup(void *pubsub_specific_data_p)
86
+{
87
+ struct pubsub_specific_data *connector_specific_data = (struct pubsub_specific_data *)pubsub_specific_data_p;
88
+
89
+ std::list<struct response> *responses = (std::list<struct response> *)connector_specific_data->responses;
90
+ std::list<struct response>::iterator response;
91
+ for (response = responses->begin(); response != responses->end(); ++response) {
92
+ // TODO: If we do this, there are a huge amount of possibly lost records. We need to find a right way of
93
+ // cleaning up contexts
94
+ // delete response->context;
95
+ delete response->publish_response;
96
+ delete response->status;
97
+ }
98
+ delete responses;
99
+
100
+ ((grpc::CompletionQueue *)connector_specific_data->completion_queue)->Shutdown();
101
+ delete (grpc::CompletionQueue *)connector_specific_data->completion_queue;
102
+ delete (google::pubsub::v1::PublishRequest *)connector_specific_data->request;
103
+ delete (google::pubsub::v1::Publisher::Stub *)connector_specific_data->stub;
104
+
105
+ // TODO: Find how to shutdown grpc gracefully. grpc_shutdown() doesn't seem to work.
106
+ // grpc_shutdown();
107
+
108
+ return;
109
+}
110
+
111
/**
112
* Add data to a Pub/Sub request message.
113
*
exporting/pubsub/pubsub_publish.h
+1
@@ -21,6 +21,7 @@ struct pubsub_specific_data {
21
int pubsub_init(
22
void *pubsub_specific_data_p, char *error_message, const char *destination, const char *credentials_file,
23
const char *project_id, const char *topic_id);
24
+void pubsub_cleanup(void *pubsub_specific_data_p);
25
26
int pubsub_add_message(void *pubsub_specific_data_p, char *data);
27
exporting/tests/test_exporting_engine.c
+1
-1
@@ -1117,7 +1117,7 @@ static void rrd_stats_api_v1_charts_allmetrics_prometheus(void **state)
1117
"netdata_info{instance=\"test_hostname\",application=\"(null)\",version=\"(null)\"} 1\n"
1118
"netdata_host_tags_info{key1=\"value1\",key2=\"value2\"} 1\n"
1119
"netdata_host_tags{key1=\"value1\",key2=\"value2\"} 1\n"
1120
- "# COMMENT TYPE test_prefix_test_context gauge\n"
1120
+ "# TYPE test_prefix_test_context gauge\n"
1121
"test_prefix_test_context{chart=\"chart_name\",family=\"test_family\",dimension=\"dimension_name\"} 690565856.0000000\n");
1122
1123
buffer_flush(buffer);