1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
|
#define AL_LOG_SECTION "curl_ws"
#include <al/log.h>
#include "../socket/net.h"
#include "websocket.h"
static void websocket_handle_events(void *userdata, struct nn_curl *curl)
{
struct nn_websocket *ws = (struct nn_websocket *)userdata;
s32 pending;
CURLMsg *msg;
while ((msg = curl_multi_info_read(curl->multi_handle, &pending))) {
switch (msg->msg) {
case CURLMSG_DONE:
nn_curl_remove_handle(curl);
nn_curl_close(curl);
nn_buffer_free(&ws->frame);
ws->callback(ws->userdata, NULL, NULL, 0);
return;
case CURLMSG_LAST:
break;
case CURLMSG_NONE:
break;
}
}
}
static size_t stream_write_callback(char *buffer, size_t size, size_t nmemb, void *userdata)
{
(void)size;
(void)nmemb;
struct nn_websocket *ws = (struct nn_websocket *)userdata;
const struct curl_ws_frame *m = curl_ws_meta(ws->curl.handle);
size_t frame_size = m->offset + m->len + m->bytesleft;
log_debug("frame_size: %zu, bytesleft: %zu.", frame_size, m->bytesleft);
nn_buffer_ensure_space(&ws->frame, frame_size);
al_memcpy(nn_buffer_get_ptr(&ws->frame, m->offset), buffer, m->len);
if (m->bytesleft == 0) {
ws->callback(ws->userdata, m, nn_buffer_get_ptr(&ws->frame, 0), frame_size);
}
return m->len;
}
bool nn_websocket_init(struct nn_websocket *ws)
{
nn_buffer_init(&ws->frame);
return nn_curl_init(&ws->curl);
}
void nn_websocket_set_url(struct nn_websocket *ws, str *url)
{
nn_curl_set_url(&ws->curl, url);
}
bool nn_websocket_connect(struct nn_websocket *ws, struct nn_event_loop *loop,
void (*callback)(void *, const struct curl_ws_frame *, u8 *, size_t), void *userdata)
{
struct nn_curl *curl = &ws->curl;
al_assert(curl->loop == NULL);
curl->loop = loop;
curl_easy_setopt(curl->handle, CURLOPT_CONNECT_ONLY, 0L);
ws->callback = callback;
ws->userdata = userdata;
curl_easy_setopt(curl->handle, CURLOPT_WRITEFUNCTION, stream_write_callback);
curl_easy_setopt(curl->handle, CURLOPT_WRITEDATA, ws);
curl->handle_events = websocket_handle_events;
curl->userdata = ws;
return nn_curl_add_handle(&ws->curl);
}
void nn_websocket_disconnect(struct nn_websocket *ws, u16 status)
{
size_t sent;
status = nn_htons(status);
curl_ws_send(ws->curl.handle, &status, sizeof(u16), &sent, 0, CURLWS_CLOSE);
al_assert(sent == sizeof(u16));
}
bool nn_websocket_send_frame(struct nn_websocket *ws, u8 *data, size_t size)
{
size_t sent;
CURLcode code = curl_ws_send(ws->curl.handle, data, size, &sent, 0, CURLWS_TEXT);
if (code != CURLE_OK) {
log_error("curl_ws_send() failed (%s).", curl_easy_strerror(code));
return false;
}
al_assert(sent == size);
return true;
}
|