summaryrefslogtreecommitdiff
path: root/src/curl
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2023-11-06 10:33:45 -0500
committerAndrew Opalach <andrew@akon.city> 2023-11-06 10:33:45 -0500
commit1517f2874ab44fcf210c14e9fdd3f9c1bfd6ba76 (patch)
tree641dd55f6677bea75c691b6010c0d7dd2a5f37e8 /src/curl
parenta723efcc67b6c14ff0dd1efd1ba88596e959daaa (diff)
downloadlibnaunet-1517f2874ab44fcf210c14e9fdd3f9c1bfd6ba76.tar.gz
libnaunet-1517f2874ab44fcf210c14e9fdd3f9c1bfd6ba76.tar.bz2
libnaunet-1517f2874ab44fcf210c14e9fdd3f9c1bfd6ba76.zip
Update
- Refactor curl support, add curl websockets - Remove aki_socket config structure - Add the ability to change aki_timer repeat - Various fixes Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/curl')
-rw-r--r--src/curl/curl.c134
-rw-r--r--src/curl/curl.h23
-rw-r--r--src/curl/http.c226
-rw-r--r--src/curl/http.h56
-rw-r--r--src/curl/websocket.c81
-rw-r--r--src/curl/websocket.h17
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);