summaryrefslogtreecommitdiff
path: root/src/curl/websocket.c
blob: ebd8f7975911a35162a2d2b8ad8b4f6a40fb01df (plain)
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;
}