diff options
Diffstat (limited to 'src/curl')
| -rw-r--r-- | src/curl/curl.c | 134 | ||||
| -rw-r--r-- | src/curl/curl.h | 23 | ||||
| -rw-r--r-- | src/curl/http.c | 226 | ||||
| -rw-r--r-- | src/curl/http.h | 56 | ||||
| -rw-r--r-- | src/curl/websocket.c | 81 | ||||
| -rw-r--r-- | src/curl/websocket.h | 17 |
6 files changed, 537 insertions, 0 deletions
diff --git a/src/curl/curl.c b/src/curl/curl.c new file mode 100644 index 0000000..3f4cb9e --- /dev/null +++ b/src/curl/curl.c @@ -0,0 +1,134 @@ +#include <al/log.h> + +#include "curl.h" + +static void curl_socket_action_callback(struct ev_loop *loop, ev_io *w, s32 revents) +{ + (void)loop; + struct aki_curl *curl = (struct aki_curl *)w->data; + s32 action = ((revents & EV_READ) ? CURL_CSELECT_IN : 0) | ((revents & EV_WRITE) ? CURL_CSELECT_OUT : 0); + s32 running; + CURLMcode mc = curl_multi_socket_action(curl->multi_handle, curl->sock, action, &running); + if (mc != CURLM_OK) { + al_log_error("curl", "curl_multi_socket_action() failed (%s).", curl_multi_strerror(mc)); + return; + } + curl->handle_events(curl); +} + +static void set_sock(struct aki_curl *curl, s32 what) +{ + s32 action = ((what & CURL_POLL_IN) ? EV_READ : 0) | ((what & CURL_POLL_OUT) ? EV_WRITE : 0); + if (curl->event.active) ev_io_stop(curl->loop->ev, &curl->event); + ev_io_init(&curl->event, curl_socket_action_callback, curl->sock, action); + ev_io_start(curl->loop->ev, &curl->event); +} + +static s32 sock_callback(CURL *e, curl_socket_t s, s32 what, void *cbp, void *sockp) +{ + (void)e; + struct aki_curl *curl = (struct aki_curl *)cbp; + if (what == CURL_POLL_REMOVE) { + ev_io_stop(curl->loop->ev, &curl->event); + curl_multi_assign(curl->multi_handle, curl->sock, NULL); + curl->sock = 0; + } else { + if (!sockp) { + curl->sock = s; + curl_multi_assign(curl->multi_handle, curl->sock, curl); + set_sock(curl, what); + } else { + set_sock(curl, what); + } + } + return 0; +} + +static void timeout_callback(struct ev_loop *loop, ev_timer *w, s32 revents) +{ + (void)loop; + struct aki_curl *curl = (struct aki_curl *)w->data; + (void)revents; + ev_timer_stop(curl->loop->ev, w); + curl->timer_started = false; + s32 running; + curl_multi_socket_action(curl->multi_handle, CURL_SOCKET_TIMEOUT, 0, &running); + curl->handle_events(curl); +} + +static s32 timer_callback(CURLM *multi, s64 timeout_ms, void *userp) +{ + (void)multi; + struct aki_curl *curl = (struct aki_curl *)userp; + if (timeout_ms == -1) { + ev_timer_stop(curl->loop->ev, &curl->timer); + return 0; + } + if (curl->timer_started) { + ev_timer_stop(curl->loop->ev, &curl->timer); + } + curl->timer.data = curl; + ev_timer_init(&curl->timer, timeout_callback, timeout_ms / 1000.0, 0.0); + curl->timer_started = true; + ev_timer_start(curl->loop->ev, &curl->timer); + return 0; +} + +bool aki_curl_init(struct aki_curl *curl) +{ + curl->handle = curl_easy_init(); + if (!curl->handle) goto err; +#ifdef _DEBUG_ + //curl_easy_setopt(curl->handle, CURLOPT_VERBOSE, 1L); +#endif + curl_easy_setopt(curl->handle, CURLOPT_NOPROGRESS, 1L); + curl->multi_handle = curl_multi_init(); + if (!curl->multi_handle) goto err; + curl_multi_setopt(curl->multi_handle, CURLMOPT_SOCKETFUNCTION, sock_callback); + curl_multi_setopt(curl->multi_handle, CURLMOPT_SOCKETDATA, curl); + curl_multi_setopt(curl->multi_handle, CURLMOPT_TIMERFUNCTION, timer_callback); + curl_multi_setopt(curl->multi_handle, CURLMOPT_TIMERDATA, curl); + curl->timer_started = false; + curl->loop = NULL; + curl->event.data = curl; + curl->event.active = 0; + return true; +err: + aki_curl_close(curl); + return false; +} + +void aki_curl_set_url(struct aki_curl *curl, str *url) +{ + char *c_str = al_str_to_c_str(url); + curl_easy_setopt(curl->handle, CURLOPT_URL, c_str); + al_free(c_str); +} + +bool aki_curl_add_handle(struct aki_curl *curl) +{ + // CURLM_RECURSIVE_API_CALL (8) also means we can't even call this + // while a callback is running on another thread. + CURLMcode mc = curl_multi_add_handle(curl->multi_handle, curl->handle); + if (mc != CURLM_OK) { + al_log_error("curl", "curl_multi_add_handle() failed (%s).", curl_multi_strerror(mc)); + return false; + } + return true; +} + +bool aki_curl_remove_handle(struct aki_curl *curl) +{ + CURLMcode mc = curl_multi_remove_handle(curl->multi_handle, curl->handle); + if (mc != CURLM_OK) { + al_log_error("curl", "curl_multi_remove_handle() failed (%s).", curl_multi_strerror(mc)); + return false; + } + return true; +} + +void aki_curl_close(struct aki_curl *curl) +{ + if (curl->handle) curl_easy_cleanup(curl->handle); + if (curl->multi_handle) curl_multi_cleanup(curl->multi_handle); +} diff --git a/src/curl/curl.h b/src/curl/curl.h new file mode 100644 index 0000000..1850dcf --- /dev/null +++ b/src/curl/curl.h @@ -0,0 +1,23 @@ +#pragma once + +#include <curl/curl.h> +#include <al/str.h> + +#include "../loop.h" + +struct aki_curl { + ev_io event; + ev_timer timer; + bool timer_started; + struct aki_event_loop *loop; + CURL *handle; + CURLM *multi_handle; + curl_socket_t sock; + void (*handle_events)(struct aki_curl *); +}; + +bool aki_curl_init(struct aki_curl *curl); +void aki_curl_set_url(struct aki_curl *curl, str *url); +bool aki_curl_add_handle(struct aki_curl *curl); +bool aki_curl_remove_handle(struct aki_curl *curl); +void aki_curl_close(struct aki_curl *curl); diff --git a/src/curl/http.c b/src/curl/http.c new file mode 100644 index 0000000..ecdec2d --- /dev/null +++ b/src/curl/http.c @@ -0,0 +1,226 @@ +#include <al/log.h> + +#include "http.h" + +static void check_status_codes(struct aki_curl *curl) +{ + struct aki_http *http = (struct aki_http *)curl; + if (http->response_code <= 0) { + curl_easy_getinfo(curl->handle, CURLINFO_RESPONSE_CODE, &http->response_code); + if (http->response_code > 0) { + http->callback(http->userdata, AKI_HTTP_RESPONSE_CODE, NULL, http->response_code); + } + } + if (http->response_code > 0 && http->content_length < 0) { + curl_easy_getinfo(curl->handle, CURLINFO_CONTENT_LENGTH_DOWNLOAD_T, &http->content_length); + if (http->content_length >= 0) { + http->callback(http->userdata, AKI_HTTP_CONTENT_LENGTH, NULL, (s64)http->content_length); + } + } +} + +static void http_handle_events(struct aki_curl *curl) +{ + s32 pending; + CURLMsg *msg; + struct aki_http *http = (struct aki_http *)curl; + while ((msg = curl_multi_info_read(curl->multi_handle, &pending))) { + switch (msg->msg) { + case CURLMSG_DONE: + al_assert(msg->easy_handle == curl->handle); + CURLcode code = msg->data.result; + if (code != CURLE_OK) { + al_log_error("http", "Curl error: %s.", curl_easy_strerror(code)); + return; + } + if (!aki_curl_remove_handle(curl)) return; + check_status_codes(curl); + if (http->response_code == 200) { + http->callback(http->userdata, AKI_HTTP_FINISHED, NULL, http->response_code); + } else if (http->response_code / 100 == 3) { + char *redirect_url = NULL; + curl_easy_getinfo(curl->handle, CURLINFO_REDIRECT_URL, &redirect_url); + if (!redirect_url) { + http->callback(http->userdata, AKI_HTTP_ERROR, NULL, http->response_code); + return; + } + aki_curl_set_url(curl, al_str_c(redirect_url)); + http->response_code = -1; + http->content_length = -1; + if (!aki_curl_add_handle(curl)) return; + http->callback(http->userdata, AKI_HTTP_REDIRECT, NULL, http->response_code); + } else { + http->callback(http->userdata, AKI_HTTP_ERROR, NULL, http->response_code); + } + return; + case CURLMSG_LAST: + al_log_debug("http", "Unhandled CURLMSG_LAST."); + break; + case CURLMSG_NONE: + al_log_debug("http", "Unhandled CURLMSG_NONE."); + break; + } + } +} + +bool aki_http_init(struct aki_http *http) +{ + http->headers = NULL; + http->response_code = -1; + http->content_length = -1; + return aki_curl_init(&http->curl); +} + +void aki_http_set_url(struct aki_http *http, str *url) +{ + aki_curl_set_url(&http->curl, url); +} + +void aki_http_set_user_agent(struct aki_http *http, str *user_agent) +{ + char *c_str = al_str_to_c_str(user_agent); + curl_easy_setopt(http->curl.handle, CURLOPT_USERAGENT, c_str); + al_free(c_str); +} + +void aki_http_add_header(struct aki_http *http, str *header) +{ + char *c_str = al_str_to_c_str(header); + http->headers = curl_slist_append(http->headers, c_str); + al_free(c_str); +} + +void aki_http_set_range(struct aki_http *http, ssize_t start, ssize_t end) +{ + char range[64]; + if (start < 0) al_sprintf(range, "-%zd", end); + else if (end < 0) al_sprintf(range, "%zd-", start); + else al_sprintf(range, "%zd-%zd", start, end); + curl_easy_setopt(http->curl.handle, CURLOPT_RANGE, range); +} + +void aki_http_close(struct aki_http *http) +{ + aki_curl_close(&http->curl); +} + +static size_t stream_write_callback(char *buffer, size_t size, size_t nmemb, void *userdata) +{ + struct aki_http *http = (struct aki_http *)userdata; + check_status_codes(&http->curl); + return http->callback(http->userdata, AKI_HTTP_WRITE, (u8 *)buffer, size * nmemb); +} + +static size_t stream_read_callback(char *buffer, size_t size, size_t nmemb, void *userdata) +{ + struct aki_http *http = (struct aki_http *)userdata; + return http->callback(http->userdata, AKI_HTTP_READ, (u8 *)buffer, size * nmemb); +} + +static size_t request_callback(void *userdata, u8 op, u8 *buf, s64 int0) +{ + struct aki_http_request *request = (struct aki_http_request *)userdata; + switch (op) { + case AKI_HTTP_FINISHED: { + request->callback(request->userdata, request, true); + break; + } + case AKI_HTTP_ERROR: { + request->callback(request->userdata, request, false); + break; + } + case AKI_HTTP_RESPONSE_CODE: { + break; + } + case AKI_HTTP_CONTENT_LENGTH: { + aki_buffer_ensure_space(&request->response, int0); + break; + } + case AKI_HTTP_WRITE: { + aki_buffer_append(&request->response, buf, int0); + break; + } + case AKI_HTTP_READ: { + al_assert(request->pointer >= 0); + size_t size = aki_buffer_get_size(&request->payload); + if (request->pointer + int0 > (off_t)size) { + int0 = size - request->pointer; + } + aki_buffer_read(&request->payload, buf, request->pointer, int0); + request->pointer += int0; + break; + } + } + return int0; +} + +static bool init_request_internal(struct aki_http *http, u8 method) +{ + struct aki_curl *curl = &http->curl; + switch (method) { + case AKI_HTTP_HEAD: + curl_easy_setopt(curl->handle, CURLOPT_NOBODY, 1L); + break; + case AKI_HTTP_GET: + curl_easy_setopt(curl->handle, CURLOPT_HTTPGET, 1L); + break; + case AKI_HTTP_POST: + curl_easy_setopt(curl->handle, CURLOPT_POST, 1L); + break; + } + curl_easy_setopt(curl->handle, CURLOPT_WRITEFUNCTION, stream_write_callback); + curl_easy_setopt(curl->handle, CURLOPT_WRITEDATA, http); + curl_easy_setopt(curl->handle, CURLOPT_READFUNCTION, stream_read_callback); + curl_easy_setopt(curl->handle, CURLOPT_READDATA, http); + if (http->headers) { + curl_easy_setopt(curl->handle, CURLOPT_HTTPHEADER, http->headers); + } + curl->handle_events = http_handle_events; + return aki_curl_add_handle(curl); +} + +bool aki_http_request_stream(struct aki_http *http, u8 method, struct aki_event_loop *loop, + size_t (*callback)(void *, u8, u8 *, s64), void *userdata) +{ + struct aki_curl *curl = &http->curl; + al_assert(curl->loop == NULL); + al_assert(method != AKI_HTTP_POST); // Not implemented. + curl->loop = loop; + http->callback = callback; + http->userdata = userdata; + return init_request_internal(http, method); +} + +bool aki_http_request_init(struct aki_http_request *request) +{ + request->pointer = 0; + aki_buffer_init(&request->payload); + aki_buffer_init(&request->response); + return aki_http_init(&request->http); +} + +bool aki_http_request(struct aki_http_request *request, u8 method, struct aki_event_loop *loop, + void (*callback)(void *, struct aki_http_request *, bool), void *userdata) +{ + struct aki_curl *curl = &request->http.curl; + al_assert(curl->loop == NULL); + curl->loop = loop; + request->http.callback = request_callback; + request->http.userdata = request; + size_t payload = aki_buffer_get_size(&request->payload); + if (payload > 0) { + curl_easy_setopt(curl->handle, CURLOPT_POSTFIELDSIZE, (s64)payload); + } else { + al_assert(method != AKI_HTTP_POST); + } + request->callback = callback; + request->userdata = userdata; + return init_request_internal(&request->http, method); +} + +void aki_http_request_close(struct aki_http_request *request) +{ + aki_buffer_free(&request->payload); + aki_buffer_free(&request->response); + aki_http_close(&request->http); +} diff --git a/src/curl/http.h b/src/curl/http.h new file mode 100644 index 0000000..28942b9 --- /dev/null +++ b/src/curl/http.h @@ -0,0 +1,56 @@ +#pragma once + +#include <al/str.h> + +#include "../util/buffer.h" + +#include "../loop.h" + +#include "curl.h" + +enum { + AKI_HTTP_HEAD = 0, + AKI_HTTP_GET, + AKI_HTTP_POST +}; + +enum { + AKI_HTTP_WRITE = 0, + AKI_HTTP_READ, + AKI_HTTP_RESPONSE_CODE, + AKI_HTTP_CONTENT_LENGTH, + AKI_HTTP_FINISHED, + AKI_HTTP_REDIRECT, + AKI_HTTP_ERROR +}; + +struct aki_http { + struct aki_curl curl; + s64 response_code; + curl_off_t content_length; + struct curl_slist *headers; + size_t (*callback)(void *, u8, u8 *, s64); + void *userdata; +}; + +struct aki_http_request { + struct aki_http http; + off_t pointer; + struct aki_buffer payload; + struct aki_buffer response; + void (*callback)(void *, struct aki_http_request *, bool); + void *userdata; +}; + +bool aki_http_init(struct aki_http *http); +void aki_http_set_url(struct aki_http *http, str *url); +void aki_http_set_user_agent(struct aki_http *http, str *user_agent); +void aki_http_add_header(struct aki_http *http, str *header); +void aki_http_set_range(struct aki_http *http, ssize_t start, ssize_t end); +void aki_http_close(struct aki_http *http); +bool aki_http_request_stream(struct aki_http *http, u8 method, struct aki_event_loop *loop, + size_t (*callback)(void *, u8, u8 *, s64), void *userdata); +bool aki_http_request_init(struct aki_http_request *request); +bool aki_http_request(struct aki_http_request *request, u8 method, struct aki_event_loop *loop, + void (*callback)(void *, struct aki_http_request *, bool), void *userdata); +void aki_http_request_close(struct aki_http_request *request); diff --git a/src/curl/websocket.c b/src/curl/websocket.c new file mode 100644 index 0000000..7e1d9b9 --- /dev/null +++ b/src/curl/websocket.c @@ -0,0 +1,81 @@ +#include <al/log.h> +#include <aki/socket.h> + +#include "websocket.h" + +static void websocket_handle_events(struct aki_curl *curl) +{ + s32 pending; + CURLMsg *msg; + while ((msg = curl_multi_info_read(curl->multi_handle, &pending))) { + switch (msg->msg) { + case CURLMSG_DONE: + aki_curl_remove_handle(curl); + aki_curl_close(curl); + 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 aki_websocket *ws = (struct aki_websocket *)userdata; + const struct curl_ws_frame *m = curl_ws_meta(ws->curl.handle); + size_t frame_size = m->offset + m->len + m->bytesleft; + al_log_debug("ws", "frame_size: %zu, bytesleft: %zu", frame_size, m->bytesleft); + aki_buffer_ensure_space(&ws->frame, frame_size); + al_memcpy(aki_buffer_get_ptr(&ws->frame, m->offset), buffer, m->len); + if (m->bytesleft == 0) { + ws->callback(ws->userdata, m, aki_buffer_get_ptr(&ws->frame, 0), frame_size); + } + return m->len; +} + +bool aki_websocket_init(struct aki_websocket *ws) +{ + aki_buffer_init(&ws->frame); + return aki_curl_init(&ws->curl); +} + +void aki_websocket_set_url(struct aki_websocket *ws, str *url) +{ + aki_curl_set_url(&ws->curl, url); +} + +bool aki_websocket_connect(struct aki_websocket *ws, struct aki_event_loop *loop, + void (*callback)(void *, const struct curl_ws_frame *, u8 *, size_t), void *userdata) +{ + al_assert(ws->curl.loop == NULL); + ws->curl.loop = loop; + ws->callback = callback; + ws->userdata = userdata; + curl_easy_setopt(ws->curl.handle, CURLOPT_WRITEFUNCTION, stream_write_callback); + curl_easy_setopt(ws->curl.handle, CURLOPT_WRITEDATA, ws); + ws->curl.handle_events = websocket_handle_events; + return aki_curl_add_handle(&ws->curl); +} + +void aki_websocket_disconnect(struct aki_websocket *ws, u16 status) +{ + size_t sent; + status = aki_htons(status); + curl_ws_send(ws->curl.handle, &status, sizeof(u16), &sent, 0, CURLWS_CLOSE); + al_assert(sent == sizeof(u16)); +} + +//aki_websocket_close(); +//aki_buffer_free(&ws->frame); + +bool aki_websocket_send_frame(struct aki_websocket *ws, u8 *data, size_t size) +{ + size_t sent; + curl_ws_send(ws->curl.handle, data, size, &sent, 0, CURLWS_TEXT); + al_assert(sent == size); + return true; +} diff --git a/src/curl/websocket.h b/src/curl/websocket.h new file mode 100644 index 0000000..3dc740a --- /dev/null +++ b/src/curl/websocket.h @@ -0,0 +1,17 @@ +#pragma once + +#include "http.h" + +struct aki_websocket { + struct aki_curl curl; + struct aki_buffer frame; + void (*callback)(void *, const struct curl_ws_frame *, u8 *, size_t); + void *userdata; +}; + +bool aki_websocket_init(struct aki_websocket *ws); +void aki_websocket_set_url(struct aki_websocket *ws, str *url); +bool aki_websocket_connect(struct aki_websocket *ws, struct aki_event_loop *loop, + void (*callback)(void *, const struct curl_ws_frame *, u8 *, size_t), void *userdata); +void aki_websocket_disconnect(struct aki_websocket *ws, u16 status); +bool aki_websocket_send_frame(struct aki_websocket *ws, u8 *data, size_t size); |