@cryptotaxi247 / netdata-1 / commits / b6d729f96

remove vneg from ACLK-NG (#10980)

removes obsolete version negotiation from ACLK-NG

Timotej S committed Apr 21, 2021 at 15:41 UTC b6d729f96a0b1e17abffb2cf69d1ff29e62004f4
7 files changed +10 -174
aclk/aclk.c
+1 -3
@@ -43,8 +43,6 @@ netdata_mutex_t aclk_shared_state_mutex = NETDATA_MUTEX_INITIALIZER;
43 struct aclk_shared_state aclk_shared_state = {
44 .agent_state = AGENT_INITIALIZING,
45 .last_popcorn_interrupt = 0,
46 - .version_neg = 0,
47 - .version_neg_wait_till = 0,
46 .mqtt_shutdown_msg_id = -1,
47 .mqtt_shutdown_msg_rcvd = 0
48 };
@@ -338,7 +336,7 @@ static inline void mqtt_connected_actions(mqtt_wss_client client)
336 aclk_stats_upd_online(1);
337 aclk_connected = 1;
338 aclk_pubacks_per_conn = 0;
341 - aclk_hello_msg(client);
339 +
340 ACLK_SHARED_STATE_LOCK;
341 if (aclk_shared_state.agent_state != AGENT_INITIALIZING) {
342 error("Sending `connect` payload immediately as popcorning was finished already.");
aclk/aclk.h
+2 -22
@@ -9,23 +9,8 @@ typedef struct aclk_rrdhost_state {
9 #include "../daemon/common.h"
10 #include "aclk_util.h"
11
12 -// minimum and maximum supported version of ACLK
13 -// in this version of agent
14 -#define ACLK_VERSION_MIN 2
15 -#define ACLK_VERSION_MAX 2
16 -
17 -// Version negotiation messages have they own versioning
18 -// this is also used for LWT message as we set that up
19 -// before version negotiation
20 -#define ACLK_VERSION_NEG_VERSION 1
21 -
22 -// Maximum time to wait for version negotiation before aborting
23 -// and defaulting to oldest supported version
24 -#define VERSION_NEG_TIMEOUT 3
25 -
26 -#if ACLK_VERSION_MIN > ACLK_VERSION_MAX
27 -#error "ACLK_VERSION_MAX must be >= than ACLK_VERSION_MIN"
28 -#endif
12 +// version for aclk legacy (old cloud arch)
13 +#define ACLK_VERSION 2
14
15 // Define ACLK Feature Version Boundaries Here
16 #define ACLK_V_COMPRESSION 2
@@ -70,11 +55,6 @@ extern struct aclk_shared_state {
55 ACLK_AGENT_STATE agent_state;
56 time_t last_popcorn_interrupt;
57
73 - // read only while ACLK connected
74 - // protect by lock otherwise
75 - int version_neg;
76 - usec_t version_neg_wait_till;
77 -
58 // To wait for `disconnect` message PUBACK
59 // when shuting down
60 // at the same time if > 0 we know link is
aclk/aclk_query.c
-17
@@ -235,23 +235,6 @@ void *aclk_query_main_thread(void *ptr)
235 {
236 struct aclk_query_thread *info = ptr;
237 while (!netdata_exit) {
238 - ACLK_SHARED_STATE_LOCK;
239 - if (unlikely(!aclk_shared_state.version_neg)) {
240 - if (!aclk_shared_state.version_neg_wait_till || aclk_shared_state.version_neg_wait_till > now_monotonic_usec()) {
241 - ACLK_SHARED_STATE_UNLOCK;
242 - info("Waiting for ACLK Version Negotiation message from Cloud");
243 - sleep(1);
244 - continue;
245 - }
246 - errno = 0;
247 - error("ACLK version negotiation failed. No reply to \"hello\" with \"version\" from cloud in time of %ds."
248 - " Reverting to default ACLK version of %d.", VERSION_NEG_TIMEOUT, ACLK_VERSION_MIN);
249 - aclk_shared_state.version_neg = ACLK_VERSION_MIN;
250 -// When ACLK v3 is implemented you will need this
251 -// aclk_set_rx_handlers(aclk_shared_state.version_neg);
252 - }
253 - ACLK_SHARED_STATE_UNLOCK;
254 -
238 aclk_query_process_msgs(info);
239
240 QUERY_THREAD_LOCK;
aclk/aclk_rx_msgs.c
-88
@@ -166,81 +166,6 @@ error:
166 return 1;
167 }
168
169 -// This handles `version` message from cloud used to negotiate
170 -// protocol version we will use
171 -static int aclk_handle_version_response(struct aclk_request *cloud_to_agent, char *raw_payload)
172 -{
173 - UNUSED(raw_payload);
174 - int version = -1;
175 - errno = 0;
176 -
177 - if (unlikely(cloud_to_agent->version != ACLK_VERSION_NEG_VERSION)) {
178 - error(
179 - "Unsuported version of \"version\" message from cloud. Expected %d, Got %d",
180 - ACLK_VERSION_NEG_VERSION,
181 - cloud_to_agent->version);
182 - return 1;
183 - }
184 - if (unlikely(!cloud_to_agent->min_version)) {
185 - error("Min version missing or 0");
186 - return 1;
187 - }
188 - if (unlikely(!cloud_to_agent->max_version)) {
189 - error("Max version missing or 0");
190 - return 1;
191 - }
192 - if (unlikely(cloud_to_agent->max_version < cloud_to_agent->min_version)) {
193 - error(
194 - "Max version (%d) must be >= than min version (%d)", cloud_to_agent->max_version,
195 - cloud_to_agent->min_version);
196 - return 1;
197 - }
198 -
199 - if (unlikely(cloud_to_agent->min_version > ACLK_VERSION_MAX)) {
200 - error(
201 - "Agent too old for this cloud. Minimum version required by cloud %d."
202 - " Maximum version supported by this agent %d.",
203 - cloud_to_agent->min_version, ACLK_VERSION_MAX);
204 - aclk_kill_link = 1;
205 - aclk_disable_runtime = 1;
206 - return 1;
207 - }
208 - if (unlikely(cloud_to_agent->max_version < ACLK_VERSION_MIN)) {
209 - error(
210 - "Cloud version is too old for this agent. Maximum version supported by cloud %d."
211 - " Minimum (oldest) version supported by this agent %d.",
212 - cloud_to_agent->max_version, ACLK_VERSION_MIN);
213 - aclk_kill_link = 1;
214 - return 1;
215 - }
216 -
217 - version = MIN(cloud_to_agent->max_version, ACLK_VERSION_MAX);
218 -
219 - ACLK_SHARED_STATE_LOCK;
220 - if (unlikely(now_monotonic_usec() > aclk_shared_state.version_neg_wait_till)) {
221 - errno = 0;
222 - error("The \"version\" message came too late ignoring.");
223 - goto err_cleanup;
224 - }
225 - if (unlikely(aclk_shared_state.version_neg)) {
226 - errno = 0;
227 - error("Version has already been set to %d", aclk_shared_state.version_neg);
228 - goto err_cleanup;
229 - }
230 - aclk_shared_state.version_neg = version;
231 - ACLK_SHARED_STATE_UNLOCK;
232 -
233 - info("Choosing version %d of ACLK", version);
234 -
235 - aclk_set_rx_handlers(version);
236 -
237 - return 0;
238 -
239 -err_cleanup:
240 - ACLK_SHARED_STATE_UNLOCK;
241 - return 1;
242 -}
243 -
169 typedef struct aclk_incoming_msg_type{
170 char *name;
171 int(*fnc)(struct aclk_request *, char *);
@@ -248,20 +173,11 @@ typedef struct aclk_incoming_msg_type{
173
174 aclk_incoming_msg_type aclk_incoming_msg_types_compression[] = {
175 { .name = "http", .fnc = aclk_handle_cloud_request_v2 },
251 - { .name = "version", .fnc = aclk_handle_version_response },
176 { .name = NULL, .fnc = NULL }
177 };
178
179 struct aclk_incoming_msg_type *aclk_incoming_msg_types = aclk_incoming_msg_types_compression;
180
257 -void aclk_set_rx_handlers(int version)
258 -{
259 -// ACLK_NG ACLK version support starts at 2
260 -// TODO ACLK v3
261 - UNUSED(version);
262 - aclk_incoming_msg_types = aclk_incoming_msg_types_compression;
263 -}
264 -
181 int aclk_handle_cloud_message(char *payload)
182 {
183 struct aclk_request cloud_to_agent;
@@ -295,10 +211,6 @@ int aclk_handle_cloud_message(char *payload)
211 goto err_cleanup;
212 }
213
298 - if (!aclk_shared_state.version_neg && strcmp(cloud_to_agent.type_id, "version")) {
299 - error("Only \"version\" message is allowed before popcorning and version negotiation is finished. Ignoring");
300 - goto err_cleanup;
301 - }
214
215 for (int i = 0; aclk_incoming_msg_types[i].name; i++) {
216 if (strcmp(cloud_to_agent.type_id, aclk_incoming_msg_types[i].name) == 0) {
aclk/aclk_rx_msgs.h
-1
@@ -9,6 +9,5 @@
9 #include "libnetdata/libnetdata.h"
10
11 int aclk_handle_cloud_message(char *payload);
12 -void aclk_set_rx_handlers(int version);
12
13 #endif /* ACLK_RX_MSGS_H */
aclk/aclk_tx_msgs.c
+7 -41
@@ -211,9 +211,9 @@ void aclk_send_info_metadata(mqtt_wss_client client, int metadata_submitted, RRD
211 // a fake on_connect message then use the real timestamp to indicate it is within the existing
212 // session.
213 if (metadata_submitted)
214 - msg = create_hdr("update", msg_id, 0, 0, aclk_shared_state.version_neg);
214 + msg = create_hdr("update", msg_id, 0, 0, ACLK_VERSION);
215 else
216 - msg = create_hdr("connect", msg_id, aclk_session_sec, aclk_session_us, aclk_shared_state.version_neg);
216 + msg = create_hdr("connect", msg_id, aclk_session_sec, aclk_session_us, ACLK_VERSION);
217
218 payload = json_object_new_object();
219 json_object_object_add(msg, "payload", payload);
@@ -253,9 +253,9 @@ void aclk_send_alarm_metadata(mqtt_wss_client client, int metadata_submitted)
253 // session.
254
255 if (metadata_submitted)
256 - msg = create_hdr("connect_alarms", msg_id, 0, 0, aclk_shared_state.version_neg);
256 + msg = create_hdr("connect_alarms", msg_id, 0, 0, ACLK_VERSION);
257 else
258 - msg = create_hdr("connect_alarms", msg_id, aclk_session_sec, aclk_session_us, aclk_shared_state.version_neg);
258 + msg = create_hdr("connect_alarms", msg_id, aclk_session_sec, aclk_session_us, ACLK_VERSION);
259
260 payload = json_object_new_object();
261 json_object_object_add(msg, "payload", payload);
@@ -277,39 +277,6 @@ void aclk_send_alarm_metadata(mqtt_wss_client client, int metadata_submitted)
277 buffer_free(local_buffer);
278 }
279
280 -void aclk_hello_msg(mqtt_wss_client client)
281 -{
282 - json_object *tmp, *msg;
283 -
284 - char *msg_id = create_uuid();
285 -
286 - ACLK_SHARED_STATE_LOCK;
287 - aclk_shared_state.version_neg = 0;
288 - aclk_shared_state.version_neg_wait_till = now_monotonic_usec() + USEC_PER_SEC * VERSION_NEG_TIMEOUT;
289 - ACLK_SHARED_STATE_UNLOCK;
290 -
291 - //Hello message is versioned separatelly from the rest of the protocol
292 - msg = create_hdr("hello", msg_id, 0, 0, ACLK_VERSION_NEG_VERSION);
293 -
294 - tmp = json_object_new_int(ACLK_VERSION_MIN);
295 - json_object_object_add(msg, "min-version", tmp);
296 -
297 - tmp = json_object_new_int(ACLK_VERSION_MAX);
298 - json_object_object_add(msg, "max-version", tmp);
299 -
300 -#ifdef ACLK_NG
301 - tmp = json_object_new_string("Next Generation");
302 -#else
303 - tmp = json_object_new_string("Legacy");
304 -#endif
305 - json_object_object_add(msg, "aclk-implementation", tmp);
306 -
307 - aclk_send_message_subtopic(client, msg, ACLK_TOPICID_METADATA);
308 -
309 - json_object_put(msg);
310 - freez(msg_id);
311 -}
312 -
280 void aclk_http_msg_v2(mqtt_wss_client client, const char *topic, const char *msg_id, usec_t t_exec, usec_t created, int http_code, const char *payload, size_t payload_len)
281 {
282 json_object *tmp, *msg;
@@ -352,7 +319,7 @@ void aclk_chart_msg(mqtt_wss_client client, RRDHOST *host, const char *chart)
319 return;
320 }
321
355 - msg = create_hdr("chart", NULL, 0, 0, aclk_shared_state.version_neg);
322 + msg = create_hdr("chart", NULL, 0, 0, ACLK_VERSION);
323 json_object_object_add(msg, "payload", payload);
324
325 aclk_send_message_subtopic(client, msg, ACLK_TOPICID_CHART);
@@ -364,11 +331,10 @@ void aclk_chart_msg(mqtt_wss_client client, RRDHOST *host, const char *chart)
331 void aclk_alarm_state_msg(mqtt_wss_client client, json_object *msg)
332 {
333 // we create header here on purpose (and not send message with it already as `msg` param)
367 - // one is version_neg is guaranteed to be done here
368 - // other are timestamps etc. which in ACLK legacy would be wrong (because ACLK legacy
334 + // timestamps etc. which in ACLK legacy would be wrong (because ACLK legacy
335 // send message with timestamps already to Query Queue they would be incorrect at time
336 // when query queue would get to send them)
371 - json_object *obj = create_hdr("status-change", NULL, 0, 0, aclk_shared_state.version_neg);
337 + json_object *obj = create_hdr("status-change", NULL, 0, 0, ACLK_VERSION);
338 json_object_object_add(obj, "payload", msg);
339
340 aclk_send_message_subtopic(client, obj, ACLK_TOPICID_ALARMS);
aclk/aclk_tx_msgs.h
-2
@@ -10,8 +10,6 @@
10 void aclk_send_info_metadata(mqtt_wss_client client, int metadata_submitted, RRDHOST *host);
11 void aclk_send_alarm_metadata(mqtt_wss_client client, int metadata_submitted);
12
13 -void aclk_hello_msg(mqtt_wss_client client);
14 -
13 void aclk_http_msg_v2(mqtt_wss_client client, const char *topic, const char *msg_id, usec_t t_exec, usec_t created, int http_code, const char *payload, size_t payload_len);
14
15 void aclk_chart_msg(mqtt_wss_client client, RRDHOST *host, const char *chart);