/src/gdal/curl/lib/mqtt.c
Line | Count | Source |
1 | | /*************************************************************************** |
2 | | * _ _ ____ _ |
3 | | * Project ___| | | | _ \| | |
4 | | * / __| | | | |_) | | |
5 | | * | (__| |_| | _ <| |___ |
6 | | * \___|\___/|_| \_\_____| |
7 | | * |
8 | | * Copyright (C) Daniel Stenberg, <daniel@haxx.se>, et al. |
9 | | * Copyright (C) Björn Stenberg, <bjorn@haxx.se> |
10 | | * |
11 | | * This software is licensed as described in the file COPYING, which |
12 | | * you should have received as part of this distribution. The terms |
13 | | * are also available at https://curl.se/docs/copyright.html. |
14 | | * |
15 | | * You may opt to use, copy, modify, merge, publish, distribute and/or sell |
16 | | * copies of the Software, and permit persons to whom the Software is |
17 | | * furnished to do so, under the terms of the COPYING file. |
18 | | * |
19 | | * This software is distributed on an "AS IS" basis, WITHOUT WARRANTY OF ANY |
20 | | * KIND, either express or implied. |
21 | | * |
22 | | * SPDX-License-Identifier: curl |
23 | | * |
24 | | ***************************************************************************/ |
25 | | #include "curl_setup.h" |
26 | | #include "urldata.h" |
27 | | |
28 | | #ifndef CURL_DISABLE_MQTT |
29 | | |
30 | | #include "transfer.h" |
31 | | #include "sendf.h" |
32 | | #include "curl_trc.h" |
33 | | #include "progress.h" |
34 | | #include "mqtt.h" |
35 | | #include "select.h" |
36 | | #include "url.h" |
37 | | #include "escape.h" |
38 | | #include "rand.h" |
39 | | #include "cfilters.h" |
40 | | #include "connect.h" |
41 | | |
42 | | /* first byte is command. |
43 | | second byte is for flags. */ |
44 | 0 | #define MQTT_MSG_CONNECT 0x10 |
45 | | /* #define MQTT_MSG_CONNACK 0x20 */ |
46 | 0 | #define MQTT_MSG_PUBLISH 0x30 |
47 | 0 | #define MQTT_MSG_SUBSCRIBE 0x82 |
48 | 0 | #define MQTT_MSG_SUBACK 0x90 |
49 | 0 | #define MQTT_MSG_DISCONNECT 0xe0 |
50 | | /* #define MQTT_MSG_PINGREQ 0xC0 */ |
51 | 0 | #define MQTT_MSG_PINGRESP 0xD0 |
52 | | |
53 | 0 | #define MQTT_CONNACK_LEN 2 |
54 | 0 | #define MQTT_SUBACK_LEN 3 |
55 | 0 | #define MQTT_CLIENTID_LEN 12 /* "curl0123abcd" */ |
56 | | |
57 | | /* meta key for storing protocol meta at easy handle */ |
58 | 0 | #define CURL_META_MQTT_EASY "meta:proto:mqtt:easy" |
59 | | /* meta key for storing protocol meta at connection */ |
60 | 0 | #define CURL_META_MQTT_CONN "meta:proto:mqtt:conn" |
61 | | |
62 | | enum mqttstate { |
63 | | MQTT_FIRST, /* 0 */ |
64 | | MQTT_REMAINING_LENGTH, /* 1 */ |
65 | | MQTT_CONNACK, /* 2 */ |
66 | | MQTT_SUBACK, /* 3 */ |
67 | | MQTT_SUBACK_COMING, /* 4 - the SUBACK remainder */ |
68 | | MQTT_PUBWAIT, /* 5 - wait for publish */ |
69 | | MQTT_PUB_REMAIN, /* 6 - wait for the remainder of the publish */ |
70 | | MQTT_POST_DRAIN, /* 7 - finish sending PUBLISH */ |
71 | | MQTT_DISCONNECT_DRAIN, /* 8 - finish sending DISCONNECT */ |
72 | | |
73 | | MQTT_NOSTATE /* 9 - never used an actual state */ |
74 | | }; |
75 | | |
76 | | struct mqtt_conn { |
77 | | enum mqttstate state; |
78 | | enum mqttstate nextstate; /* switch to this after remaining length is |
79 | | done */ |
80 | | unsigned int packetid; |
81 | | }; |
82 | | |
83 | | /* protocol-specific transfer-related data */ |
84 | | struct MQTT { |
85 | | struct dynbuf sendbuf; |
86 | | /* when receiving */ |
87 | | struct dynbuf recvbuf; |
88 | | size_t npacket; /* byte counter */ |
89 | | size_t remaining_length; |
90 | | unsigned char pkt_hd[4]; /* for decoding the arriving packet length */ |
91 | | struct curltime lastTime; /* last time we sent or received data */ |
92 | | unsigned char firstbyte; |
93 | | BIT(pingsent); /* 1 while we wait for ping response */ |
94 | | }; |
95 | | |
96 | | static void mqtt_easy_dtor(const void *key, size_t klen, void *entry) |
97 | 0 | { |
98 | 0 | struct MQTT *mq = entry; |
99 | 0 | (void)key; |
100 | 0 | (void)klen; |
101 | 0 | curlx_dyn_free(&mq->sendbuf); |
102 | 0 | curlx_dyn_free(&mq->recvbuf); |
103 | 0 | curlx_free(mq); |
104 | 0 | } |
105 | | |
106 | | static void mqtt_conn_dtor(const void *key, size_t klen, void *entry) |
107 | 0 | { |
108 | 0 | (void)key; |
109 | 0 | (void)klen; |
110 | 0 | curlx_free(entry); |
111 | 0 | } |
112 | | |
113 | | static CURLcode mqtt_setup_conn(struct Curl_easy *data, |
114 | | struct connectdata *conn) |
115 | 0 | { |
116 | | /* setup MQTT specific meta data at easy handle and connection */ |
117 | 0 | struct mqtt_conn *mqtt; |
118 | 0 | struct MQTT *mq; |
119 | |
|
120 | 0 | mqtt = curlx_calloc(1, sizeof(*mqtt)); |
121 | 0 | if(!mqtt || |
122 | 0 | Curl_conn_meta_set(conn, CURL_META_MQTT_CONN, mqtt, mqtt_conn_dtor)) |
123 | 0 | return CURLE_OUT_OF_MEMORY; |
124 | | |
125 | 0 | mq = curlx_calloc(1, sizeof(struct MQTT)); |
126 | 0 | if(!mq) |
127 | 0 | return CURLE_OUT_OF_MEMORY; |
128 | 0 | curlx_dyn_init(&mq->recvbuf, DYN_MQTT_RECV); |
129 | 0 | curlx_dyn_init(&mq->sendbuf, DYN_MQTT_SEND); |
130 | 0 | if(Curl_meta_set(data, CURL_META_MQTT_EASY, mq, mqtt_easy_dtor)) |
131 | 0 | return CURLE_OUT_OF_MEMORY; |
132 | 0 | return CURLE_OK; |
133 | 0 | } |
134 | | |
135 | | static CURLcode mqtt_send(struct Curl_easy *data, |
136 | | const char *buf, size_t len) |
137 | 0 | { |
138 | 0 | size_t n; |
139 | 0 | CURLcode result; |
140 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
141 | |
|
142 | 0 | if(!mq) |
143 | 0 | return CURLE_FAILED_INIT; |
144 | | |
145 | | /* A new packet must not replace a previously queued packet tail. */ |
146 | 0 | DEBUGASSERT(!curlx_dyn_len(&mq->sendbuf)); |
147 | 0 | result = Curl_xfer_send(data, buf, len, FALSE, &n); |
148 | 0 | if(result) |
149 | 0 | return result; |
150 | 0 | mq->lastTime = *Curl_pgrs_now(data); |
151 | 0 | Curl_debug(data, CURLINFO_HEADER_OUT, buf, n); |
152 | 0 | if(len != n) |
153 | 0 | result = curlx_dyn_addn(&mq->sendbuf, &buf[n], len - n); |
154 | 0 | return result; |
155 | 0 | } |
156 | | |
157 | | /* Send the queued tail, preserving any bytes that still cannot be sent. */ |
158 | | static CURLcode mqtt_flush(struct Curl_easy *data) |
159 | 0 | { |
160 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
161 | 0 | size_t len, n; |
162 | 0 | CURLcode result; |
163 | |
|
164 | 0 | if(!mq) |
165 | 0 | return CURLE_FAILED_INIT; |
166 | 0 | len = curlx_dyn_len(&mq->sendbuf); |
167 | 0 | if(!len) |
168 | 0 | return CURLE_OK; |
169 | | |
170 | 0 | result = Curl_xfer_send(data, curlx_dyn_ptr(&mq->sendbuf), len, FALSE, &n); |
171 | 0 | if(result) |
172 | 0 | return result; |
173 | 0 | mq->lastTime = *Curl_pgrs_now(data); |
174 | 0 | Curl_debug(data, CURLINFO_HEADER_OUT, curlx_dyn_ptr(&mq->sendbuf), n); |
175 | 0 | return curlx_dyn_tail(&mq->sendbuf, len - n); |
176 | 0 | } |
177 | | |
178 | | /* Generic function called by the multi interface to figure out what socket(s) |
179 | | to wait for and for what actions during the DOING and PROTOCONNECT |
180 | | states */ |
181 | | static CURLcode mqtt_pollset(struct Curl_easy *data, |
182 | | struct easy_pollset *ps) |
183 | 0 | { |
184 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
185 | 0 | if(mq && curlx_dyn_len(&mq->sendbuf)) |
186 | 0 | return Curl_pollset_add_out(data, ps, data->conn->sock[FIRSTSOCKET]); |
187 | 0 | return Curl_pollset_add_in(data, ps, data->conn->sock[FIRSTSOCKET]); |
188 | 0 | } |
189 | | |
190 | | static int mqtt_encode_len(char *buf, size_t len) |
191 | 0 | { |
192 | 0 | int i; |
193 | |
|
194 | 0 | for(i = 0; (len > 0) && (i < 4); i++) { |
195 | 0 | unsigned char encoded; |
196 | 0 | encoded = len % 0x80; |
197 | 0 | len /= 0x80; |
198 | 0 | if(len) |
199 | 0 | encoded |= 0x80; |
200 | 0 | buf[i] = (char)encoded; |
201 | 0 | } |
202 | |
|
203 | 0 | return i; |
204 | 0 | } |
205 | | |
206 | | /* add the passwd to the CONNECT packet */ |
207 | | static int add_passwd(const char *passwd, const size_t plen, |
208 | | char *pkt, const size_t start, int remain_pos) |
209 | 0 | { |
210 | | /* magic number that need to be set properly */ |
211 | 0 | const size_t conn_flags_pos = remain_pos + 8; |
212 | 0 | if(plen > 0xffff) |
213 | 0 | return 1; |
214 | | |
215 | | /* set password flag */ |
216 | 0 | pkt[conn_flags_pos] |= 0x40; |
217 | | |
218 | | /* length of password provided */ |
219 | 0 | pkt[start] = (char)((plen >> 8) & 0xFF); |
220 | 0 | pkt[start + 1] = (char)(plen & 0xFF); |
221 | 0 | memcpy(&pkt[start + 2], passwd, plen); |
222 | 0 | return 0; |
223 | 0 | } |
224 | | |
225 | | /* add user to the CONNECT packet */ |
226 | | static int add_user(const char *username, const size_t ulen, |
227 | | unsigned char *pkt, const size_t start, int remain_pos) |
228 | 0 | { |
229 | | /* magic number that need to be set properly */ |
230 | 0 | const size_t conn_flags_pos = remain_pos + 8; |
231 | 0 | if(ulen > 0xffff) |
232 | 0 | return 1; |
233 | | |
234 | | /* set username flag */ |
235 | 0 | pkt[conn_flags_pos] |= 0x80; |
236 | | /* length of username provided */ |
237 | 0 | pkt[start] = (unsigned char)((ulen >> 8) & 0xFF); |
238 | 0 | pkt[start + 1] = (unsigned char)(ulen & 0xFF); |
239 | 0 | memcpy(&pkt[start + 2], username, ulen); |
240 | 0 | return 0; |
241 | 0 | } |
242 | | |
243 | | /* add client ID to the CONNECT packet */ |
244 | | static int add_client_id(const char *client_id, const size_t client_id_len, |
245 | | char *pkt, const size_t start) |
246 | 0 | { |
247 | 0 | if(client_id_len != MQTT_CLIENTID_LEN) |
248 | 0 | return 1; |
249 | 0 | pkt[start] = 0x00; |
250 | 0 | pkt[start + 1] = MQTT_CLIENTID_LEN; |
251 | 0 | memcpy(&pkt[start + 2], client_id, MQTT_CLIENTID_LEN); |
252 | 0 | return 0; |
253 | 0 | } |
254 | | |
255 | | /* Set initial values of CONNECT packet */ |
256 | | static int init_connpack(char *packet, char *remain, int remain_pos) |
257 | 0 | { |
258 | | /* Fixed header starts */ |
259 | | /* packet type */ |
260 | 0 | packet[0] = MQTT_MSG_CONNECT; |
261 | | /* remaining length field */ |
262 | 0 | memcpy(&packet[1], remain, remain_pos); |
263 | | /* Fixed header ends */ |
264 | | |
265 | | /* Variable header starts */ |
266 | | /* protocol length */ |
267 | 0 | packet[remain_pos + 1] = 0x00; |
268 | 0 | packet[remain_pos + 2] = 0x04; |
269 | | /* protocol name */ |
270 | 0 | packet[remain_pos + 3] = 'M'; |
271 | 0 | packet[remain_pos + 4] = 'Q'; |
272 | 0 | packet[remain_pos + 5] = 'T'; |
273 | 0 | packet[remain_pos + 6] = 'T'; |
274 | | /* protocol level */ |
275 | 0 | packet[remain_pos + 7] = 0x04; |
276 | | /* CONNECT flag: CleanSession */ |
277 | 0 | packet[remain_pos + 8] = 0x02; |
278 | | /* keep-alive 0 = disabled */ |
279 | 0 | packet[remain_pos + 9] = 0x00; |
280 | 0 | packet[remain_pos + 10] = 0x3c; |
281 | | /* end of variable header */ |
282 | 0 | return remain_pos + 10; |
283 | 0 | } |
284 | | |
285 | | static CURLcode mqtt_connect(struct Curl_easy *data) |
286 | 0 | { |
287 | 0 | CURLcode result = CURLE_OK; |
288 | 0 | int pos = 0; |
289 | 0 | int rc = 0; |
290 | | /* remain length */ |
291 | 0 | int remain_pos = 0; |
292 | 0 | char remain[4] = { 0 }; |
293 | 0 | size_t packetlen = 0; |
294 | 0 | size_t start_user = 0; |
295 | 0 | size_t start_pwd = 0; |
296 | 0 | char client_id[MQTT_CLIENTID_LEN + 1] = "curl"; |
297 | 0 | const size_t clen = CURL_CSTRLEN("curl"); |
298 | 0 | char *packet = NULL; |
299 | | |
300 | | /* extracting username from request */ |
301 | 0 | struct Curl_creds *creds = data->state.creds; |
302 | 0 | const size_t ulen = creds ? strlen(creds->user) : 0; |
303 | 0 | const size_t plen = creds ? strlen(creds->passwd) : 0; |
304 | 0 | const size_t payloadlen = ulen + plen + MQTT_CLIENTID_LEN + 2 + |
305 | | /* The plus 2s below are for the MSB and LSB describing the length of the |
306 | | string to be added on the payload. Refer to spec 1.5.2 and 1.5.4 */ |
307 | 0 | (ulen ? 2 : 0) + |
308 | 0 | (plen ? 2 : 0); |
309 | | |
310 | | /* getting how much occupy the remain length */ |
311 | 0 | remain_pos = mqtt_encode_len(remain, payloadlen + 10); |
312 | | |
313 | | /* 10 length of variable header and 1 the first byte of the fixed header */ |
314 | 0 | packetlen = payloadlen + 10 + remain_pos + 1; |
315 | | |
316 | | /* allocating packet */ |
317 | 0 | if(packetlen > 0xFFFFFFF) |
318 | 0 | return CURLE_WEIRD_SERVER_REPLY; |
319 | 0 | packet = curlx_calloc(1, packetlen); |
320 | 0 | if(!packet) |
321 | 0 | return CURLE_OUT_OF_MEMORY; |
322 | | |
323 | | /* set initial values for the CONNECT packet */ |
324 | 0 | pos = init_connpack(packet, remain, remain_pos); |
325 | |
|
326 | 0 | result = Curl_rand_alnum(data, (unsigned char *)&client_id[clen], |
327 | 0 | MQTT_CLIENTID_LEN - clen + 1); |
328 | | /* add client id */ |
329 | 0 | rc = add_client_id(client_id, strlen(client_id), packet, pos + 1); |
330 | 0 | if(rc) { |
331 | 0 | failf(data, "Client ID length mismatched: [%zu]", strlen(client_id)); |
332 | 0 | result = CURLE_WEIRD_SERVER_REPLY; |
333 | 0 | goto end; |
334 | 0 | } |
335 | 0 | infof(data, "Using client id '%s'", client_id); |
336 | | |
337 | | /* position where the user payload starts */ |
338 | 0 | start_user = pos + 3 + MQTT_CLIENTID_LEN; |
339 | | /* position where the password payload starts */ |
340 | 0 | start_pwd = start_user + ulen; |
341 | | /* if username was provided, add it to the packet */ |
342 | 0 | if(ulen) { |
343 | 0 | start_pwd += 2; |
344 | |
|
345 | 0 | rc = add_user(creds->user, ulen, |
346 | 0 | (unsigned char *)packet, start_user, remain_pos); |
347 | 0 | if(rc) { |
348 | 0 | failf(data, "Username too long: [%zu]", ulen); |
349 | 0 | result = CURLE_WEIRD_SERVER_REPLY; |
350 | 0 | goto end; |
351 | 0 | } |
352 | 0 | } |
353 | | |
354 | | /* if passwd was provided, add it to the packet */ |
355 | 0 | if(plen) { |
356 | 0 | rc = add_passwd(creds->passwd, plen, packet, start_pwd, remain_pos); |
357 | 0 | if(rc) { |
358 | 0 | failf(data, "Password too long: [%zu]", plen); |
359 | 0 | result = CURLE_WEIRD_SERVER_REPLY; |
360 | 0 | goto end; |
361 | 0 | } |
362 | 0 | } |
363 | | |
364 | 0 | if(!result) |
365 | 0 | result = mqtt_send(data, packet, packetlen); |
366 | |
|
367 | 0 | end: |
368 | 0 | if(packet) { |
369 | 0 | curlx_memzero(packet, packetlen); |
370 | 0 | curlx_free(packet); |
371 | 0 | } |
372 | 0 | Curl_creds_unlink(&data->state.creds); |
373 | 0 | return result; |
374 | 0 | } |
375 | | |
376 | | static CURLcode mqtt_disconnect(struct Curl_easy *data) |
377 | 0 | { |
378 | 0 | return mqtt_send(data, "\xe0\x00", 2); |
379 | 0 | } |
380 | | |
381 | | static CURLcode mqtt_recv_atleast(struct Curl_easy *data, size_t nbytes) |
382 | 0 | { |
383 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
384 | 0 | size_t rlen; |
385 | 0 | CURLcode result; |
386 | |
|
387 | 0 | if(!mq) |
388 | 0 | return CURLE_FAILED_INIT; |
389 | 0 | rlen = curlx_dyn_len(&mq->recvbuf); |
390 | |
|
391 | 0 | if(rlen < nbytes) { |
392 | 0 | unsigned char readbuf[1024]; |
393 | 0 | size_t nread; |
394 | |
|
395 | 0 | DEBUGASSERT(nbytes - rlen < sizeof(readbuf)); |
396 | 0 | result = Curl_xfer_recv(data, (char *)readbuf, nbytes - rlen, &nread); |
397 | 0 | if(result) |
398 | 0 | return result; |
399 | 0 | if(!nread) /* EOF */ |
400 | 0 | return CURLE_RECV_ERROR; |
401 | 0 | if(curlx_dyn_addn(&mq->recvbuf, readbuf, nread)) |
402 | 0 | return CURLE_OUT_OF_MEMORY; |
403 | 0 | rlen = curlx_dyn_len(&mq->recvbuf); |
404 | 0 | } |
405 | 0 | return (rlen >= nbytes) ? CURLE_OK : CURLE_AGAIN; |
406 | 0 | } |
407 | | |
408 | | static void mqtt_recv_consume(struct Curl_easy *data, size_t nbytes) |
409 | 0 | { |
410 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
411 | 0 | DEBUGASSERT(mq); |
412 | 0 | if(mq) { |
413 | 0 | size_t rlen = curlx_dyn_len(&mq->recvbuf); |
414 | 0 | if(rlen <= nbytes) |
415 | 0 | curlx_dyn_reset(&mq->recvbuf); |
416 | 0 | else |
417 | 0 | curlx_dyn_tail(&mq->recvbuf, rlen - nbytes); |
418 | 0 | } |
419 | 0 | } |
420 | | |
421 | | static CURLcode mqtt_verify_connack(struct Curl_easy *data) |
422 | 0 | { |
423 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
424 | 0 | CURLcode result; |
425 | 0 | const char *ptr; |
426 | |
|
427 | 0 | DEBUGASSERT(mq); |
428 | 0 | if(!mq) |
429 | 0 | return CURLE_FAILED_INIT; |
430 | 0 | if(mq->remaining_length != 2) { |
431 | 0 | failf(data, "CONNACK expected Remaining Length 2, got %zu", |
432 | 0 | mq->remaining_length); |
433 | 0 | return CURLE_WEIRD_SERVER_REPLY; |
434 | 0 | } |
435 | | |
436 | 0 | result = mqtt_recv_atleast(data, MQTT_CONNACK_LEN); |
437 | 0 | if(result) |
438 | 0 | return result; |
439 | | |
440 | | /* verify CONNACK */ |
441 | 0 | DEBUGASSERT(curlx_dyn_len(&mq->recvbuf) >= MQTT_CONNACK_LEN); |
442 | 0 | ptr = curlx_dyn_ptr(&mq->recvbuf); |
443 | 0 | Curl_debug(data, CURLINFO_HEADER_IN, ptr, MQTT_CONNACK_LEN); |
444 | |
|
445 | 0 | if(ptr[0] != 0x00 || ptr[1] != 0x00) { |
446 | 0 | failf(data, "Expected %02x%02x but got %02x%02x", |
447 | 0 | 0x00U, 0x00U, (unsigned char)ptr[0], (unsigned char)ptr[1]); |
448 | 0 | curlx_dyn_reset(&mq->recvbuf); |
449 | 0 | return CURLE_WEIRD_SERVER_REPLY; |
450 | 0 | } |
451 | 0 | mqtt_recv_consume(data, MQTT_CONNACK_LEN); |
452 | 0 | return CURLE_OK; |
453 | 0 | } |
454 | | |
455 | | static CURLcode mqtt_get_topic(struct Curl_easy *data, |
456 | | char **topic, size_t *topiclen) |
457 | 0 | { |
458 | 0 | const char *path = data->state.up.path; |
459 | 0 | CURLcode result = CURLE_URL_MALFORMAT; |
460 | 0 | if(strlen(path) > 1) { |
461 | 0 | result = Curl_urldecode(path + 1, 0, topic, topiclen, REJECT_CTRL); |
462 | 0 | if(!result && (*topiclen > 0xffff)) { |
463 | 0 | failf(data, "Too long MQTT topic"); |
464 | 0 | result = CURLE_URL_MALFORMAT; |
465 | 0 | } |
466 | 0 | } |
467 | 0 | else |
468 | 0 | failf(data, "No MQTT topic found. Forgot to URL encode it?"); |
469 | |
|
470 | 0 | return result; |
471 | 0 | } |
472 | | |
473 | | static CURLcode mqtt_subscribe(struct Curl_easy *data) |
474 | 0 | { |
475 | 0 | CURLcode result = CURLE_OK; |
476 | 0 | char *topic = NULL; |
477 | 0 | size_t topiclen; |
478 | 0 | unsigned char *packet = NULL; |
479 | 0 | size_t packetlen; |
480 | 0 | char encodedsize[4]; |
481 | 0 | size_t n; |
482 | 0 | struct connectdata *conn = data->conn; |
483 | 0 | struct mqtt_conn *mqtt = Curl_conn_meta_get(conn, CURL_META_MQTT_CONN); |
484 | |
|
485 | 0 | if(!mqtt) |
486 | 0 | return CURLE_FAILED_INIT; |
487 | | |
488 | 0 | result = mqtt_get_topic(data, &topic, &topiclen); |
489 | 0 | if(result) |
490 | 0 | goto fail; |
491 | | |
492 | 0 | mqtt->packetid++; |
493 | |
|
494 | 0 | packetlen = topiclen + 5; /* packetid + topic (has a two byte length field) |
495 | | + 2 bytes topic length + QoS byte */ |
496 | 0 | n = mqtt_encode_len((char *)encodedsize, packetlen); |
497 | 0 | packetlen += n + 1; /* add one for the control packet type byte */ |
498 | |
|
499 | 0 | packet = curlx_malloc(packetlen); |
500 | 0 | if(!packet) { |
501 | 0 | result = CURLE_OUT_OF_MEMORY; |
502 | 0 | goto fail; |
503 | 0 | } |
504 | | |
505 | 0 | packet[0] = MQTT_MSG_SUBSCRIBE; |
506 | 0 | memcpy(&packet[1], encodedsize, n); |
507 | 0 | packet[1 + n] = (mqtt->packetid >> 8) & 0xff; |
508 | 0 | packet[2 + n] = mqtt->packetid & 0xff; |
509 | 0 | packet[3 + n] = (topiclen >> 8) & 0xff; |
510 | 0 | packet[4 + n] = topiclen & 0xff; |
511 | 0 | memcpy(&packet[5 + n], topic, topiclen); |
512 | 0 | packet[5 + n + topiclen] = 0; /* QoS zero */ |
513 | |
|
514 | 0 | result = mqtt_send(data, (const char *)packet, packetlen); |
515 | |
|
516 | 0 | fail: |
517 | 0 | curlx_free(topic); |
518 | 0 | curlx_free(packet); |
519 | 0 | return result; |
520 | 0 | } |
521 | | |
522 | | /* |
523 | | * Called when the first byte was already read. |
524 | | */ |
525 | | static CURLcode mqtt_verify_suback(struct Curl_easy *data) |
526 | 0 | { |
527 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
528 | 0 | struct connectdata *conn = data->conn; |
529 | 0 | struct mqtt_conn *mqtt = Curl_conn_meta_get(conn, CURL_META_MQTT_CONN); |
530 | 0 | CURLcode result; |
531 | 0 | const char *ptr; |
532 | |
|
533 | 0 | if(!mqtt || !mq) |
534 | 0 | return CURLE_FAILED_INIT; |
535 | | |
536 | 0 | if(mq->remaining_length != 3) { |
537 | 0 | failf(data, "SUBACK expected Remaining Length 3, got %zu", |
538 | 0 | mq->remaining_length); |
539 | 0 | return CURLE_WEIRD_SERVER_REPLY; |
540 | 0 | } |
541 | | |
542 | 0 | result = mqtt_recv_atleast(data, MQTT_SUBACK_LEN); |
543 | 0 | if(result) |
544 | 0 | goto fail; |
545 | | |
546 | | /* verify SUBACK */ |
547 | 0 | DEBUGASSERT(curlx_dyn_len(&mq->recvbuf) >= MQTT_SUBACK_LEN); |
548 | 0 | ptr = curlx_dyn_ptr(&mq->recvbuf); |
549 | 0 | Curl_debug(data, CURLINFO_HEADER_IN, ptr, MQTT_SUBACK_LEN); |
550 | |
|
551 | 0 | if(((unsigned char)ptr[0]) != ((mqtt->packetid >> 8) & 0xff) || |
552 | 0 | ((unsigned char)ptr[1]) != (mqtt->packetid & 0xff) || |
553 | 0 | ptr[2] != 0x00) { |
554 | 0 | curlx_dyn_reset(&mq->recvbuf); |
555 | 0 | result = CURLE_WEIRD_SERVER_REPLY; |
556 | 0 | goto fail; |
557 | 0 | } |
558 | 0 | mqtt_recv_consume(data, MQTT_SUBACK_LEN); |
559 | 0 | fail: |
560 | 0 | return result; |
561 | 0 | } |
562 | | |
563 | 0 | #define MAX_MQTT_MESSAGE_SIZE 0xFFFFFFF |
564 | | |
565 | | static CURLcode mqtt_publish(struct Curl_easy *data) |
566 | 0 | { |
567 | 0 | CURLcode result; |
568 | 0 | char *payload = data->set.postfields; |
569 | 0 | size_t payloadlen; |
570 | 0 | char *topic = NULL; |
571 | 0 | size_t topiclen; |
572 | 0 | unsigned char *pkt = NULL; |
573 | 0 | size_t i = 0; |
574 | 0 | size_t remaininglength; |
575 | 0 | size_t encodelen; |
576 | 0 | char encodedbytes[4]; |
577 | 0 | curl_off_t postfieldsize = data->set.postfieldsize; |
578 | |
|
579 | 0 | if(!payload) { |
580 | 0 | DEBUGF(infof(data, "mqtt_publish without payload, return bad arg")); |
581 | 0 | return CURLE_BAD_FUNCTION_ARGUMENT; |
582 | 0 | } |
583 | 0 | if(!curlx_sotouz_fits(postfieldsize, &payloadlen)) { |
584 | 0 | if(postfieldsize > 0) /* off_t does not fit into size_t */ |
585 | 0 | return CURLE_BAD_FUNCTION_ARGUMENT; |
586 | 0 | payloadlen = strlen(payload); |
587 | 0 | } |
588 | | |
589 | 0 | result = mqtt_get_topic(data, &topic, &topiclen); |
590 | 0 | if(result) |
591 | 0 | goto fail; |
592 | | |
593 | | /* silly check for silly analyzers. On 32-bit systems a payload length close |
594 | | to 4GB can wrap, but it is impossible to have such a large buffer in |
595 | | memory on those systems */ |
596 | 0 | if((SIZE_MAX - topiclen - 2) < payloadlen) { |
597 | 0 | result = CURLE_TOO_LARGE; |
598 | 0 | goto fail; |
599 | 0 | } |
600 | | |
601 | 0 | remaininglength = payloadlen + 2 + topiclen; |
602 | 0 | encodelen = mqtt_encode_len(encodedbytes, remaininglength); |
603 | 0 | if(remaininglength > (MAX_MQTT_MESSAGE_SIZE - encodelen - 1)) { |
604 | 0 | result = CURLE_TOO_LARGE; |
605 | 0 | goto fail; |
606 | 0 | } |
607 | | |
608 | | /* add the control byte and the encoded remaining length */ |
609 | 0 | pkt = curlx_malloc(remaininglength + 1 + encodelen); |
610 | 0 | if(!pkt) { |
611 | 0 | result = CURLE_OUT_OF_MEMORY; |
612 | 0 | goto fail; |
613 | 0 | } |
614 | | |
615 | | /* assemble packet */ |
616 | 0 | pkt[i++] = MQTT_MSG_PUBLISH; |
617 | 0 | memcpy(&pkt[i], encodedbytes, encodelen); |
618 | 0 | i += encodelen; |
619 | 0 | pkt[i++] = (topiclen >> 8) & 0xff; |
620 | 0 | pkt[i++] = (topiclen & 0xff); |
621 | 0 | memcpy(&pkt[i], topic, topiclen); |
622 | 0 | i += topiclen; |
623 | 0 | memcpy(&pkt[i], payload, payloadlen); |
624 | 0 | i += payloadlen; |
625 | 0 | result = mqtt_send(data, (const char *)pkt, i); |
626 | |
|
627 | 0 | fail: |
628 | 0 | curlx_free(pkt); |
629 | 0 | curlx_free(topic); |
630 | 0 | return result; |
631 | 0 | } |
632 | | |
633 | | /* return FALSE on success, TRUE on error */ |
634 | | static bool mqtt_decode_len(size_t *lenp, const unsigned char *buf, |
635 | | size_t buflen) |
636 | 0 | { |
637 | 0 | size_t len = 0; |
638 | 0 | size_t mult = 1; |
639 | 0 | size_t i; |
640 | 0 | unsigned char encoded = 128; |
641 | |
|
642 | 0 | for(i = 0; (i < buflen) && (encoded & 128); i++) { |
643 | 0 | if(i == 4) |
644 | 0 | return TRUE; /* bad size */ |
645 | 0 | encoded = buf[i]; |
646 | 0 | len += (encoded & 127) * mult; |
647 | 0 | mult *= 128; |
648 | 0 | } |
649 | 0 | if(encoded & 128) |
650 | | /* truncated size */ |
651 | 0 | return TRUE; |
652 | | |
653 | 0 | *lenp = len; |
654 | 0 | return FALSE; |
655 | 0 | } |
656 | | |
657 | | #if defined(DEBUGBUILD) && defined(CURLVERBOSE) |
658 | | static const char * const statenames[] = { |
659 | | "MQTT_FIRST", |
660 | | "MQTT_REMAINING_LENGTH", |
661 | | "MQTT_CONNACK", |
662 | | "MQTT_SUBACK", |
663 | | "MQTT_SUBACK_COMING", |
664 | | "MQTT_PUBWAIT", |
665 | | "MQTT_PUB_REMAIN", |
666 | | "MQTT_POST_DRAIN", |
667 | | "MQTT_DISCONNECT_DRAIN", |
668 | | |
669 | | "NOT A STATE" |
670 | | }; |
671 | | #endif |
672 | | |
673 | | /* The only way to change state */ |
674 | | static void mqstate(struct Curl_easy *data, |
675 | | enum mqttstate state, |
676 | | enum mqttstate nextstate) /* used if state == FIRST */ |
677 | 0 | { |
678 | 0 | struct connectdata *conn = data->conn; |
679 | 0 | struct mqtt_conn *mqtt = Curl_conn_meta_get(conn, CURL_META_MQTT_CONN); |
680 | 0 | DEBUGASSERT(mqtt); |
681 | 0 | if(!mqtt) |
682 | 0 | return; |
683 | | #ifdef DEBUGBUILD |
684 | | infof(data, "%s (from %s) (next is %s)", |
685 | | statenames[state], |
686 | | statenames[mqtt->state], |
687 | | (state == MQTT_FIRST) ? statenames[nextstate] : ""); |
688 | | #endif |
689 | 0 | mqtt->state = state; |
690 | 0 | if(state == MQTT_FIRST) |
691 | 0 | mqtt->nextstate = nextstate; |
692 | 0 | } |
693 | | |
694 | | static CURLcode mqtt_read_publish(struct Curl_easy *data, bool *done) |
695 | 0 | { |
696 | 0 | CURLcode result = CURLE_OK; |
697 | 0 | struct connectdata *conn = data->conn; |
698 | 0 | size_t nread; |
699 | 0 | size_t remlen; |
700 | 0 | struct mqtt_conn *mqtt = Curl_conn_meta_get(conn, CURL_META_MQTT_CONN); |
701 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
702 | 0 | unsigned char packet; |
703 | |
|
704 | 0 | DEBUGASSERT(mqtt); |
705 | 0 | if(!mqtt || !mq) |
706 | 0 | return CURLE_FAILED_INIT; |
707 | | |
708 | 0 | switch(mqtt->state) { |
709 | 0 | MQTT_SUBACK_COMING: |
710 | 0 | case MQTT_SUBACK_COMING: |
711 | 0 | result = mqtt_verify_suback(data); |
712 | 0 | if(result) |
713 | 0 | break; |
714 | | |
715 | 0 | mqstate(data, MQTT_FIRST, MQTT_PUBWAIT); |
716 | 0 | break; |
717 | | |
718 | 0 | case MQTT_SUBACK: |
719 | 0 | case MQTT_PUBWAIT: |
720 | | /* we are expecting PUBLISH or SUBACK */ |
721 | 0 | packet = mq->firstbyte & 0xf0; |
722 | 0 | if(packet == MQTT_MSG_PUBLISH) |
723 | 0 | mqstate(data, MQTT_PUB_REMAIN, MQTT_NOSTATE); |
724 | 0 | else if(packet == MQTT_MSG_SUBACK) { |
725 | 0 | mqstate(data, MQTT_SUBACK_COMING, MQTT_NOSTATE); |
726 | 0 | goto MQTT_SUBACK_COMING; |
727 | 0 | } |
728 | 0 | else if(packet == MQTT_MSG_DISCONNECT) { |
729 | 0 | infof(data, "Got DISCONNECT"); |
730 | 0 | *done = TRUE; |
731 | 0 | goto end; |
732 | 0 | } |
733 | 0 | else { |
734 | 0 | result = CURLE_WEIRD_SERVER_REPLY; |
735 | 0 | goto end; |
736 | 0 | } |
737 | | |
738 | | /* -- switched state -- */ |
739 | 0 | remlen = mq->remaining_length; |
740 | 0 | infof(data, "Remaining length: %zu bytes", remlen); |
741 | 0 | if(data->set.max_filesize && |
742 | 0 | (curl_off_t)remlen > data->set.max_filesize) { |
743 | 0 | failf(data, "Maximum file size exceeded"); |
744 | 0 | result = CURLE_FILESIZE_EXCEEDED; |
745 | 0 | goto end; |
746 | 0 | } |
747 | 0 | Curl_pgrsSetDownloadSize(data, remlen); |
748 | 0 | data->req.bytecount = 0; |
749 | 0 | data->req.size = remlen; |
750 | 0 | mq->npacket = remlen; /* get this many bytes */ |
751 | 0 | FALLTHROUGH(); |
752 | 0 | case MQTT_PUB_REMAIN: { |
753 | | /* read rest of packet, but no more. Cap to buffer size */ |
754 | 0 | char buffer[4 * 1024]; |
755 | 0 | size_t rest = mq->npacket; |
756 | 0 | if(rest > sizeof(buffer)) |
757 | 0 | rest = sizeof(buffer); |
758 | 0 | result = Curl_xfer_recv(data, buffer, rest, &nread); |
759 | 0 | if(result) { |
760 | 0 | if(result == CURLE_AGAIN) { |
761 | 0 | infof(data, "EEEE AAAAGAIN"); |
762 | 0 | } |
763 | 0 | goto end; |
764 | 0 | } |
765 | 0 | if(!nread) { |
766 | 0 | infof(data, "server disconnected"); |
767 | 0 | result = CURLE_PARTIAL_FILE; |
768 | 0 | goto end; |
769 | 0 | } |
770 | | |
771 | | /* we received something */ |
772 | 0 | mq->lastTime = *Curl_pgrs_now(data); |
773 | | |
774 | | /* if QoS is set, message contains packet id */ |
775 | 0 | result = Curl_client_write(data, CLIENTWRITE_BODY, buffer, nread); |
776 | 0 | if(result) |
777 | 0 | goto end; |
778 | | |
779 | 0 | mq->npacket -= nread; |
780 | 0 | if(!mq->npacket) |
781 | | /* no more PUBLISH payload, back to subscribe wait state */ |
782 | 0 | mqstate(data, MQTT_FIRST, MQTT_PUBWAIT); |
783 | 0 | break; |
784 | 0 | } |
785 | 0 | default: |
786 | 0 | DEBUGASSERT(NULL); /* illegal state */ |
787 | 0 | result = CURLE_WEIRD_SERVER_REPLY; |
788 | 0 | goto end; |
789 | 0 | } |
790 | 0 | end: |
791 | 0 | return result; |
792 | 0 | } |
793 | | |
794 | | static CURLcode mqtt_do(struct Curl_easy *data, bool *done) |
795 | 0 | { |
796 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
797 | 0 | CURLcode result = CURLE_OK; |
798 | 0 | *done = FALSE; /* unconditionally */ |
799 | |
|
800 | 0 | if(!mq) |
801 | 0 | return CURLE_FAILED_INIT; |
802 | 0 | mq->lastTime = *Curl_pgrs_now(data); |
803 | 0 | mq->pingsent = FALSE; |
804 | |
|
805 | 0 | result = mqtt_connect(data); |
806 | 0 | if(result) { |
807 | 0 | failf(data, "Error %d sending MQTT CONNECT request", (int)result); |
808 | 0 | return result; |
809 | 0 | } |
810 | 0 | mqstate(data, MQTT_FIRST, MQTT_CONNACK); |
811 | 0 | return CURLE_OK; |
812 | 0 | } |
813 | | |
814 | | static CURLcode mqtt_done(struct Curl_easy *data, |
815 | | CURLcode status, bool premature) |
816 | 0 | { |
817 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
818 | 0 | (void)status; |
819 | 0 | (void)premature; |
820 | 0 | if(mq) { |
821 | 0 | curlx_dyn_free(&mq->sendbuf); |
822 | 0 | curlx_dyn_free(&mq->recvbuf); |
823 | 0 | } |
824 | 0 | return CURLE_OK; |
825 | 0 | } |
826 | | |
827 | | /* we ping regularly to avoid being disconnected by the server */ |
828 | | static CURLcode mqtt_ping(struct Curl_easy *data) |
829 | 0 | { |
830 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
831 | 0 | CURLcode result = CURLE_OK; |
832 | 0 | struct connectdata *conn = data->conn; |
833 | 0 | struct mqtt_conn *mqtt = Curl_conn_meta_get(conn, CURL_META_MQTT_CONN); |
834 | |
|
835 | 0 | if(!mqtt || !mq) |
836 | 0 | return CURLE_FAILED_INIT; |
837 | | |
838 | 0 | if(mqtt->state == MQTT_FIRST && |
839 | 0 | !mq->pingsent && |
840 | 0 | data->set.upkeep_interval_ms > 0) { |
841 | 0 | struct curltime t = *Curl_pgrs_now(data); |
842 | 0 | timediff_t diff = curlx_ptimediff_ms(&t, &mq->lastTime); |
843 | |
|
844 | 0 | if(diff > data->set.upkeep_interval_ms) { |
845 | | /* 0xC0 is PINGREQ, and 0x00 is remaining length */ |
846 | 0 | unsigned char packet[2] = { 0xC0, 0x00 }; |
847 | 0 | size_t packetlen = sizeof(packet); |
848 | |
|
849 | 0 | result = mqtt_send(data, (char *)packet, packetlen); |
850 | 0 | if(!result) { |
851 | 0 | mq->pingsent = TRUE; |
852 | 0 | } |
853 | 0 | infof(data, "mqtt_ping: sent ping request."); |
854 | 0 | } |
855 | 0 | } |
856 | 0 | return result; |
857 | 0 | } |
858 | | |
859 | | static CURLcode mqtt_doing(struct Curl_easy *data, bool *done) |
860 | 0 | { |
861 | 0 | struct MQTT *mq = Curl_meta_get(data, CURL_META_MQTT_EASY); |
862 | 0 | CURLcode result = CURLE_OK; |
863 | 0 | size_t nread; |
864 | 0 | unsigned char recvbyte; |
865 | 0 | struct mqtt_conn *mqtt = Curl_conn_meta_get(data->conn, CURL_META_MQTT_CONN); |
866 | |
|
867 | 0 | if(!mqtt || !mq) |
868 | 0 | return CURLE_FAILED_INIT; |
869 | | |
870 | 0 | *done = FALSE; |
871 | |
|
872 | 0 | if(curlx_dyn_len(&mq->sendbuf)) { |
873 | | /* send the remainder of an outgoing packet */ |
874 | 0 | result = mqtt_flush(data); |
875 | | /* CURLE_OK can still mean a short write. Wait for writable progress |
876 | | before sending another packet or processing a reply. */ |
877 | 0 | if(result || curlx_dyn_len(&mq->sendbuf)) |
878 | 0 | return result; |
879 | 0 | } |
880 | | |
881 | 0 | result = mqtt_ping(data); |
882 | | /* Packet processing may send more output, so finish PINGREQ first. */ |
883 | 0 | if(result || curlx_dyn_len(&mq->sendbuf)) |
884 | 0 | return result; |
885 | | |
886 | 0 | infof(data, "mqtt_doing: state [%d]", (int)mqtt->state); |
887 | 0 | switch(mqtt->state) { |
888 | 0 | case MQTT_FIRST: |
889 | | /* Read the initial byte only */ |
890 | 0 | result = Curl_xfer_recv(data, (char *)&mq->firstbyte, 1, &nread); |
891 | 0 | if(result) |
892 | 0 | break; |
893 | 0 | else if(!nread) { |
894 | 0 | failf(data, "Connection disconnected"); |
895 | 0 | *done = TRUE; |
896 | 0 | result = CURLE_RECV_ERROR; |
897 | 0 | break; |
898 | 0 | } |
899 | 0 | Curl_debug(data, CURLINFO_HEADER_IN, (const char *)&mq->firstbyte, 1); |
900 | | |
901 | | /* we received something */ |
902 | 0 | mq->lastTime = *Curl_pgrs_now(data); |
903 | | |
904 | | /* remember the first byte */ |
905 | 0 | mq->npacket = 0; |
906 | 0 | mqstate(data, MQTT_REMAINING_LENGTH, MQTT_NOSTATE); |
907 | 0 | FALLTHROUGH(); |
908 | 0 | case MQTT_REMAINING_LENGTH: |
909 | 0 | do { |
910 | 0 | result = Curl_xfer_recv(data, (char *)&recvbyte, 1, &nread); |
911 | 0 | if(result || !nread) |
912 | 0 | break; |
913 | 0 | Curl_debug(data, CURLINFO_HEADER_IN, (const char *)&recvbyte, 1); |
914 | 0 | mq->pkt_hd[mq->npacket++] = recvbyte; |
915 | 0 | } while((recvbyte & 0x80) && (mq->npacket < 4)); |
916 | 0 | if(!result && nread && (recvbyte & 0x80)) |
917 | | /* MQTT supports up to 127 * 128^0 + 127 * 128^1 + 127 * 128^2 + |
918 | | 127 * 128^3 bytes. server tried to send more */ |
919 | 0 | result = CURLE_WEIRD_SERVER_REPLY; |
920 | 0 | if(result) |
921 | 0 | break; |
922 | 0 | if(mqtt_decode_len(&mq->remaining_length, mq->pkt_hd, mq->npacket)) { |
923 | 0 | result = CURLE_WEIRD_SERVER_REPLY; |
924 | 0 | break; |
925 | 0 | } |
926 | 0 | mq->npacket = 0; |
927 | | /* PINGRESP and DISCONNECT must have remaining_length == 0 and |
928 | | * reserved bits (low nibble) must be zero per MQTT 3.1.1 |
929 | | * sections 2.2.2, 3.13.1 and 3.14.1. Reject before state |
930 | | * dispatch to prevent nextstate confusion. */ |
931 | 0 | { |
932 | 0 | const unsigned char type = mq->firstbyte & 0xF0; |
933 | 0 | const unsigned char reserved = mq->firstbyte & 0x0F; |
934 | 0 | if((type == MQTT_MSG_DISCONNECT || type == MQTT_MSG_PINGRESP) && |
935 | 0 | (mq->remaining_length || reserved)) { |
936 | 0 | failf(data, |
937 | 0 | "Broker sent malformed %s " |
938 | 0 | "(remaining_length=%zu, header byte=0x%02x)", |
939 | 0 | type == MQTT_MSG_DISCONNECT ? "DISCONNECT" : "PINGRESP", |
940 | 0 | mq->remaining_length, mq->firstbyte); |
941 | 0 | result = CURLE_WEIRD_SERVER_REPLY; |
942 | 0 | break; |
943 | 0 | } |
944 | 0 | } |
945 | 0 | if(mq->remaining_length) { |
946 | 0 | mqstate(data, mqtt->nextstate, MQTT_NOSTATE); |
947 | 0 | break; |
948 | 0 | } |
949 | 0 | mqstate(data, MQTT_FIRST, MQTT_FIRST); |
950 | |
|
951 | 0 | if(mq->firstbyte == MQTT_MSG_DISCONNECT) { |
952 | 0 | infof(data, "Got DISCONNECT"); |
953 | 0 | *done = TRUE; |
954 | 0 | } |
955 | | |
956 | | /* ping response */ |
957 | 0 | if(mq->firstbyte == MQTT_MSG_PINGRESP) { |
958 | 0 | infof(data, "Received ping response."); |
959 | 0 | mq->pingsent = FALSE; |
960 | 0 | mqstate(data, MQTT_FIRST, MQTT_PUBWAIT); |
961 | 0 | } |
962 | 0 | break; |
963 | 0 | case MQTT_CONNACK: |
964 | 0 | result = mqtt_verify_connack(data); |
965 | 0 | if(result) |
966 | 0 | break; |
967 | | |
968 | 0 | if(data->state.httpreq != HTTPREQ_POST) { |
969 | 0 | result = mqtt_subscribe(data); |
970 | 0 | if(!result) |
971 | 0 | mqstate(data, MQTT_FIRST, MQTT_SUBACK); |
972 | 0 | break; |
973 | 0 | } |
974 | | |
975 | 0 | result = mqtt_publish(data); |
976 | 0 | if(result) |
977 | 0 | break; |
978 | 0 | mqstate(data, MQTT_POST_DRAIN, MQTT_NOSTATE); |
979 | 0 | FALLTHROUGH(); |
980 | 0 | case MQTT_POST_DRAIN: |
981 | | /* PUBLISH must be complete before DISCONNECT can use the send buffer. */ |
982 | 0 | if(curlx_dyn_len(&mq->sendbuf)) |
983 | 0 | break; |
984 | 0 | result = mqtt_disconnect(data); |
985 | 0 | if(result) |
986 | 0 | break; |
987 | 0 | mqstate(data, MQTT_DISCONNECT_DRAIN, MQTT_NOSTATE); |
988 | 0 | FALLTHROUGH(); |
989 | 0 | case MQTT_DISCONNECT_DRAIN: |
990 | | /* Completing now would discard any DISCONNECT bytes still queued. */ |
991 | 0 | if(!curlx_dyn_len(&mq->sendbuf)) |
992 | 0 | *done = TRUE; |
993 | 0 | break; |
994 | | |
995 | 0 | case MQTT_SUBACK: |
996 | 0 | case MQTT_PUBWAIT: |
997 | 0 | case MQTT_PUB_REMAIN: |
998 | 0 | result = mqtt_read_publish(data, done); |
999 | 0 | break; |
1000 | | |
1001 | 0 | default: |
1002 | 0 | failf(data, "State not handled yet"); |
1003 | 0 | *done = TRUE; |
1004 | 0 | break; |
1005 | 0 | } |
1006 | | |
1007 | 0 | if(result == CURLE_AGAIN) |
1008 | 0 | result = CURLE_OK; |
1009 | 0 | return result; |
1010 | 0 | } |
1011 | | |
1012 | | #ifdef USE_SSL |
1013 | | |
1014 | | static CURLcode mqtts_connecting(struct Curl_easy *data, bool *done) |
1015 | 0 | { |
1016 | 0 | struct connectdata *conn = data->conn; |
1017 | 0 | CURLcode result; |
1018 | |
|
1019 | 0 | result = Curl_conn_connect(data, FIRSTSOCKET, TRUE, done); |
1020 | 0 | if(result) |
1021 | 0 | connclose(conn); |
1022 | 0 | return result; |
1023 | 0 | } |
1024 | | |
1025 | | /* |
1026 | | * MQTTS protocol. |
1027 | | */ |
1028 | | const struct Curl_protocol Curl_protocol_mqtts = { |
1029 | | mqtt_setup_conn, /* setup_connection */ |
1030 | | mqtt_do, /* do_it */ |
1031 | | mqtt_done, /* done */ |
1032 | | ZERO_NULL, /* do_more */ |
1033 | | ZERO_NULL, /* connect_it */ |
1034 | | mqtts_connecting, /* connecting */ |
1035 | | mqtt_doing, /* doing */ |
1036 | | ZERO_NULL, /* proto_pollset */ |
1037 | | mqtt_pollset, /* doing_pollset */ |
1038 | | ZERO_NULL, /* domore_pollset */ |
1039 | | ZERO_NULL, /* perform_pollset */ |
1040 | | ZERO_NULL, /* disconnect */ |
1041 | | ZERO_NULL, /* write_resp */ |
1042 | | ZERO_NULL, /* write_resp_hd */ |
1043 | | ZERO_NULL, /* connection_is_dead */ |
1044 | | ZERO_NULL, /* attach connection */ |
1045 | | ZERO_NULL, /* follow */ |
1046 | | }; |
1047 | | |
1048 | | #endif |
1049 | | |
1050 | | /* |
1051 | | * MQTT protocol. |
1052 | | */ |
1053 | | const struct Curl_protocol Curl_protocol_mqtt = { |
1054 | | mqtt_setup_conn, /* setup_connection */ |
1055 | | mqtt_do, /* do_it */ |
1056 | | mqtt_done, /* done */ |
1057 | | ZERO_NULL, /* do_more */ |
1058 | | ZERO_NULL, /* connect_it */ |
1059 | | ZERO_NULL, /* connecting */ |
1060 | | mqtt_doing, /* doing */ |
1061 | | ZERO_NULL, /* proto_pollset */ |
1062 | | mqtt_pollset, /* doing_pollset */ |
1063 | | ZERO_NULL, /* domore_pollset */ |
1064 | | ZERO_NULL, /* perform_pollset */ |
1065 | | ZERO_NULL, /* disconnect */ |
1066 | | ZERO_NULL, /* write_resp */ |
1067 | | ZERO_NULL, /* write_resp_hd */ |
1068 | | ZERO_NULL, /* connection_is_dead */ |
1069 | | ZERO_NULL, /* attach connection */ |
1070 | | ZERO_NULL, /* follow */ |
1071 | | }; |
1072 | | |
1073 | | #endif /* CURL_DISABLE_MQTT */ |