From 32c88dbbd2425983a5a7248d254fd316d6681f5e Mon Sep 17 00:00:00 2001 From: theMackabu Date: Thu, 26 Mar 2026 22:37:32 -0700 Subject: [PATCH] add node:net module --- examples/demo/net.cjs | 26 ++ examples/demo/server.cjs | 3 + examples/demo/welcome.js | 18 +- include/gc/modules.h | 1 + include/net/connection.h | 14 +- include/net/listener.h | 27 +- src/gc/objects.c | 1 + src/modules/net.c | 884 ++++++++++++++++++++++++++++++++++++++- src/modules/server.c | 11 +- src/modules/v8.c | 2 - src/net/connection.c | 202 ++++++++- src/net/listener.c | 6 +- 12 files changed, 1129 insertions(+), 66 deletions(-) create mode 100644 examples/demo/net.cjs create mode 100644 examples/demo/server.cjs diff --git a/examples/demo/net.cjs b/examples/demo/net.cjs new file mode 100644 index 0000000..ee75051 --- /dev/null +++ b/examples/demo/net.cjs @@ -0,0 +1,26 @@ +const net = require('node:net'); + +const httpResponse = `HTTP/1.1 200 OK\r +Content-Type: text/plain\r +\r +Hello, World from net module! +`; + +const server = net.createServer(socket => { + console.log('Client connected'); + + socket.on('data', () => { + socket.write(httpResponse); + socket.end(); + }); + + socket.on('end', () => { + console.log('Client disconnected'); + }); + + socket.on('error', err => { + console.error(`Socket error: ${err.message}`); + }); +}); + +server.listen(3000, () => console.log(`server on http://localhost:3000`)); diff --git a/examples/demo/server.cjs b/examples/demo/server.cjs new file mode 100644 index 0000000..2188615 --- /dev/null +++ b/examples/demo/server.cjs @@ -0,0 +1,3 @@ +require('node:http') + .createServer((_req, res) => res.end('ant!')) + .listen(3000); diff --git a/examples/demo/welcome.js b/examples/demo/welcome.js index 4432db5..35f7750 100644 --- a/examples/demo/welcome.js +++ b/examples/demo/welcome.js @@ -2,14 +2,16 @@ console.log(`Hello from ${navigator.userAgent}! Did you know:\n`); const rows = [ ['Feature', 'Description'], - ['Async/Await', 'Full coroutine support with minicoro'], - ['HTTP Server', 'Built-in Ant.serve() with TLS support'], - ['Fetch API', 'HTTP client with TLS via tlsuv'], - ['File System', 'Async/sync fs with ant:fs module'], - ['FFI', 'Native library integration'], - ['Web Locks', 'Navigator.locks API'], - ['TypeScript', 'Built-in type stripping via oxc'], - ['Garbage Collection', 'Mark-copy compacting + Boehm-Demers'], + ['Async/Await', 'Native promises, microtasks, and async functions'], + ['HTTP Server', 'Module-first server runtime with export default { fetch }'], + ['WinterTC APIs', 'Request, Response, Headers, fetch, streams, Blob, FormData'], + ['Fetch API', 'Request/Response/Headers/FormData with streamed bodies'], + ['File System', 'Async and sync filesystem APIs for scripts and tooling'], + ['FFI', 'Call native libraries directly with libffi-backed bindings'], + ['Web Locks', 'Web-standard Navigator.locks coordination primitives'], + ['TypeScript', 'Built-in .ts execution with automatic type stripping'], + ['Compression', 'CompressionStream, DecompressionStream, gzip, deflate, brotli'], + ['Garbage Collection', 'Adaptive generational GC with nursery and major cycles'], null, [null, 'And more...'] ]; diff --git a/include/gc/modules.h b/include/gc/modules.h index 3d5e105..88ffc91 100644 --- a/include/gc/modules.h +++ b/include/gc/modules.h @@ -13,6 +13,7 @@ void gc_mark_child_process(ant_t *js, gc_mark_fn mark); void gc_mark_readline(ant_t *js, gc_mark_fn mark); void gc_mark_process(ant_t *js, gc_mark_fn mark); void gc_mark_navigator(ant_t *js, gc_mark_fn mark); +void gc_mark_net(ant_t *js, gc_mark_fn mark); void gc_mark_server(ant_t *js, gc_mark_fn mark); void gc_mark_events(ant_t *js, gc_mark_fn mark); void gc_mark_lmdb(ant_t *js, gc_mark_fn mark); diff --git a/include/net/connection.h b/include/net/connection.h index f3de59d..a67e0bc 100644 --- a/include/net/connection.h +++ b/include/net/connection.h @@ -3,6 +3,7 @@ #include #include +#include #include typedef struct ant_listener_s ant_listener_t; @@ -16,17 +17,26 @@ struct ant_conn_s { char *buffer; size_t buffer_len; size_t buffer_cap; - int timeout_secs; + uint64_t timeout_ms; + uint64_t bytes_read; + uint64_t bytes_written; int close_handles; bool closing; + bool read_paused; + bool read_eof; + bool has_local_addr; bool has_remote_addr; + char local_addr[64]; char remote_addr[64]; + char local_family[8]; + char remote_family[8]; + int local_port; int remote_port; struct ant_conn_s *next; }; void ant_conn_start(ant_conn_t *conn); int ant_conn_accept(ant_conn_t *conn, uv_stream_t *server_stream); -ant_conn_t *ant_conn_create_tcp(ant_listener_t *listener, int timeout_secs); +ant_conn_t *ant_conn_create_tcp(ant_listener_t *listener, uint64_t timeout_ms); #endif diff --git a/include/net/listener.h b/include/net/listener.h index aa83da7..a2ee5af 100644 --- a/include/net/listener.h +++ b/include/net/listener.h @@ -3,6 +3,7 @@ #include #include +#include #include typedef struct ant_listener_s ant_listener_t; @@ -17,6 +18,9 @@ typedef void (*ant_conn_write_cb)( typedef struct { void (*on_accept)(ant_listener_t *listener, ant_conn_t *conn, void *user_data); void (*on_read)(ant_conn_t *conn, ssize_t nread, void *user_data); + void (*on_end)(ant_conn_t *conn, void *user_data); + void (*on_error)(ant_conn_t *conn, int status, void *user_data); + void (*on_timeout)(ant_conn_t *conn, void *user_data); void (*on_conn_close)(ant_conn_t *conn, void *user_data); void (*on_listener_close)(ant_listener_t *listener, void *user_data); } ant_listener_callbacks_t; @@ -27,7 +31,7 @@ struct ant_listener_s { ant_conn_t *connections; ant_listener_callbacks_t callbacks; void *user_data; - int idle_timeout_secs; + uint64_t idle_timeout_ms; int port; int backlog; bool started; @@ -41,7 +45,7 @@ int ant_listener_listen_tcp( const char *hostname, int port, int backlog, - int idle_timeout_secs, + uint64_t idle_timeout_ms, const ant_listener_callbacks_t *callbacks, void *user_data ); @@ -49,26 +53,41 @@ int ant_listener_listen_tcp( void ant_listener_stop(ant_listener_t *listener, bool force); void *ant_listener_get_user_data(const ant_listener_t *listener); +int ant_conn_local_port(const ant_conn_t *conn); int ant_listener_port(const ant_listener_t *listener); -int ant_conn_timeout(const ant_conn_t *conn); int ant_conn_remote_port(const ant_conn_t *conn); +int ant_conn_set_no_delay(ant_conn_t *conn, bool enable); +int ant_conn_set_keep_alive(ant_conn_t *conn, bool enable, unsigned int delay_secs); int ant_conn_write(ant_conn_t *conn, char *data, size_t len, ant_conn_write_cb cb, void *user_data); bool ant_listener_has_connections(const ant_listener_t *listener); bool ant_listener_is_closed(const ant_listener_t *listener); bool ant_conn_is_closing(const ant_conn_t *conn); +bool ant_conn_has_local_addr(const ant_conn_t *conn); bool ant_conn_has_remote_addr(const ant_conn_t *conn); void ant_conn_set_user_data(ant_conn_t *conn, void *user_data); void *ant_conn_get_user_data(const ant_conn_t *conn); void ant_conn_pause_read(ant_conn_t *conn); +void ant_conn_resume_read(ant_conn_t *conn); void ant_conn_close(ant_conn_t *conn); -void ant_conn_set_timeout(ant_conn_t *conn, int timeout_secs); +void ant_conn_shutdown(ant_conn_t *conn); +void ant_conn_ref(ant_conn_t *conn); +void ant_conn_unref(ant_conn_t *conn); +void ant_conn_set_timeout_ms(ant_conn_t *conn, uint64_t timeout_ms); +void ant_conn_consume(ant_conn_t *conn, size_t len); const char *ant_conn_buffer(const ant_conn_t *conn); +const char *ant_conn_local_addr(const ant_conn_t *conn); const char *ant_conn_remote_addr(const ant_conn_t *conn); +const char *ant_conn_local_family(const ant_conn_t *conn); +const char *ant_conn_remote_family(const ant_conn_t *conn); size_t ant_conn_buffer_len(const ant_conn_t *conn); ant_listener_t *ant_conn_listener(const ant_conn_t *conn); +uint64_t ant_conn_timeout_ms(const ant_conn_t *conn); +uint64_t ant_conn_bytes_read(const ant_conn_t *conn); +uint64_t ant_conn_bytes_written(const ant_conn_t *conn); + #endif diff --git a/src/gc/objects.c b/src/gc/objects.c index 95af9b6..a15652d 100644 --- a/src/gc/objects.c +++ b/src/gc/objects.c @@ -446,6 +446,7 @@ static void gc_mark_roots(ant_t *js) { gc_mark_readline(js, gc_mark_value); gc_mark_process(js, gc_mark_value); gc_mark_navigator(js, gc_mark_value); + gc_mark_net(js, gc_mark_value); gc_mark_server(js, gc_mark_value); gc_mark_events(js, gc_mark_value); gc_mark_lmdb(js, gc_mark_value); diff --git a/src/modules/net.c b/src/modules/net.c index 0d542cf..af00417 100644 --- a/src/modules/net.c +++ b/src/modules/net.c @@ -1,6 +1,8 @@ -// stub: minimal node:net implementation (isIP, isIPv4, isIPv6) -// just enough for vite to resolve the module and validate hostnames +#include // IWYU pragma: keep +#include +#include +#include #include #ifdef _WIN32 @@ -11,48 +13,886 @@ #endif #include "ant.h" -#include "internal.h" // IWYU pragma: keep +#include "ptr.h" +#include "errors.h" +#include "internal.h" -static ant_value_t net_isIP(ant_t *js, ant_value_t *args, int nargs) { - if (nargs < 1) return js_mknum(0); - size_t len; - const char *host = js_getstr(js, args[0], &len); - if (!host) return js_mknum(0); +#include "gc/modules.h" +#include "net/listener.h" +#include "silver/engine.h" + +#include "modules/buffer.h" +#include "modules/events.h" +#include "modules/net.h" +#include "modules/symbol.h" + +typedef struct net_server_s net_server_t; +typedef struct net_socket_s net_socket_t; + +struct net_socket_s { + ant_t *js; + ant_value_t obj; + ant_conn_t *conn; + net_server_t *server; + ant_value_t encoding; + struct net_socket_s *next_active; + struct net_socket_s *next_in_server; + bool allow_half_open; + bool destroyed; + bool had_error; +}; + +struct net_server_s { + ant_t *js; + ant_value_t obj; + ant_listener_t listener; + net_socket_t *sockets; + struct net_server_s *next_active; + char *host; + int port; + int backlog; + int max_connections; + bool listening; + bool closing; + bool allow_half_open; + bool pause_on_connect; + bool no_delay; + bool keep_alive; + unsigned int keep_alive_initial_delay_secs; +}; + +static ant_value_t g_net_server_proto = 0; +static ant_value_t g_net_socket_proto = 0; +static ant_value_t g_net_server_ctor = 0; +static ant_value_t g_net_socket_ctor = 0; + +static net_server_t *g_active_servers = NULL; +static net_socket_t *g_active_sockets = NULL; + +static bool g_default_auto_select_family = true; +static double g_default_auto_select_family_attempt_timeout = 250.0; + +enum { + NET_SERVER_NATIVE_TAG = 0x4e455453u, + NET_SOCKET_NATIVE_TAG = 0x4e45544bu, +}; + +static net_server_t *net_server_data(ant_value_t value) { + if (!js_check_native_tag(value, NET_SERVER_NATIVE_TAG)) return NULL; + return (net_server_t *)js_get_native_ptr(value); +} + +static net_socket_t *net_socket_data(ant_value_t value) { + if (!js_check_native_tag(value, NET_SOCKET_NATIVE_TAG)) return NULL; + return (net_socket_t *)js_get_native_ptr(value); +} + +static void net_add_active_server(net_server_t *server) { + server->next_active = g_active_servers; + g_active_servers = server; +} + +static void net_remove_active_server(net_server_t *server) { + net_server_t **it = NULL; + for (it = &g_active_servers; *it; it = &(*it)->next_active) { + if (*it == server) { + *it = server->next_active; + server->next_active = NULL; + return; + }} +} + +static void net_add_active_socket(net_socket_t *socket) { + socket->next_active = g_active_sockets; + g_active_sockets = socket; +} + +static void net_remove_active_socket(net_socket_t *socket) { + net_socket_t **it = NULL; + for (it = &g_active_sockets; *it; it = &(*it)->next_active) { + if (*it == socket) { + *it = socket->next_active; + socket->next_active = NULL; + return; + }} +} + +static ant_value_t net_call_value( + ant_t *js, + ant_value_t fn, + ant_value_t this_val, + ant_value_t *args, + int nargs +) { + ant_value_t saved_this = js->this_val; + ant_value_t result = js_mkundef(); + + js->this_val = this_val; + if (vtype(fn) == T_CFUNC) + result = ((ant_value_t (*)(ant_t *, ant_value_t *, int))vdata(fn))(js, args, nargs); + else result = sv_vm_call(js->vm, js, fn, this_val, args, nargs, NULL, false); + js->this_val = saved_this; + return result; +} + +static ant_value_t net_call_method( + ant_t *js, + ant_value_t target, + const char *name, + ant_value_t *args, + int nargs +) { + ant_value_t fn = js_get(js, target, name); + if (!is_callable(fn)) return js_mkundef(); + return net_call_value(js, fn, target, args, nargs); +} + +static void net_emit(ant_t *js, ant_value_t target, const char *event, ant_value_t *args, int nargs) { + ant_value_t emit_args[8] = {0}; + int i = 0; + + if (nargs > 7) nargs = 7; + emit_args[0] = js_mkstr(js, event, strlen(event)); + for (i = 0; i < nargs; i++) emit_args[i + 1] = args[i]; + net_call_method(js, target, "emit", emit_args, nargs + 1); +} + +static ant_value_t net_make_buffer_chunk(ant_t *js, const char *data, size_t len) { + ArrayBufferData *ab = create_array_buffer_data(len); + if (!ab) return js_mkerr_typed(js, JS_ERR_TYPE, "Out of memory"); + if (len > 0 && data) memcpy(ab->data, data, len); + return create_typed_array(js, TYPED_ARRAY_UINT8, ab, 0, len, "Buffer"); +} + +static void net_socket_sync_state(net_socket_t *socket) { + ant_value_t obj = 0; + const char *ready_state = "open"; + + if (!socket) return; + obj = socket->obj; + if (!is_object_type(obj)) return; + + if (socket->destroyed) ready_state = "closed"; + js_set(socket->js, obj, "destroyed", js_bool(socket->destroyed)); + js_set(socket->js, obj, "pending", js_bool(socket->conn == NULL && !socket->destroyed)); + js_set(socket->js, obj, "connecting", js_false); + js_set(socket->js, obj, "readyState", js_mkstr(socket->js, ready_state, strlen(ready_state))); + js_set(socket->js, obj, "bytesRead", js_mknum((double)(socket->conn ? ant_conn_bytes_read(socket->conn) : 0))); + js_set(socket->js, obj, "bytesWritten", js_mknum((double)(socket->conn ? ant_conn_bytes_written(socket->conn) : 0))); + js_set(socket->js, obj, "timeout", js_mknum((double)(socket->conn ? ant_conn_timeout_ms(socket->conn) : 0))); +} + +static void net_server_sync_state(net_server_t *server) { + if (!server || !is_object_type(server->obj)) return; + js_set(server->js, server->obj, "listening", js_bool(server->listening)); + js_set(server->js, server->obj, "maxConnections", js_mknum((double)server->max_connections)); + js_set(server->js, server->obj, "dropMaxConnection", js_false); +} + +static int net_server_socket_count(net_server_t *server) { + int count = 0; + net_socket_t *socket = NULL; + + for (socket = server ? server->sockets : NULL; socket; socket = socket->next_in_server) + count++; + return count; +} + +static void net_server_maybe_finish_close(net_server_t *server) { + if (!server || !server->closing) return; + if (!ant_listener_is_closed(&server->listener)) return; + if (server->sockets) return; + + server->closing = false; + server->listening = false; + net_server_sync_state(server); + net_emit(server->js, server->obj, "close", NULL, 0); + net_remove_active_server(server); +} + +static void net_socket_detach(net_socket_t *socket) { + if (!socket) return; + + if (socket->server) { + net_socket_t **it = NULL; + for (it = &socket->server->sockets; *it; it = &(*it)->next_in_server) { + if (*it == socket) { + *it = socket->next_in_server; + break; + } + } + socket->next_in_server = NULL; + } + + net_remove_active_socket(socket); + if (is_object_type(socket->obj)) { + js_set_native_ptr(socket->obj, NULL); + js_set_native_tag(socket->obj, 0); + } + net_server_maybe_finish_close(socket->server); + free(socket); +} + +static void net_socket_emit_error(net_socket_t *socket, const char *message) { + ant_value_t arg = js_mkerr_typed(socket->js, JS_ERR_TYPE, "%s", message); + socket->had_error = true; + net_emit(socket->js, socket->obj, "error", &arg, 1); +} + +static net_socket_t *net_socket_create(ant_t *js, bool allow_half_open) { + ant_value_t obj = js_mkobj(js); + ant_value_t proto = js_instance_proto_from_new_target(js, g_net_socket_proto); + net_socket_t *socket = calloc(1, sizeof(*socket)); + if (!socket) return NULL; + if (is_object_type(proto)) js_set_proto_init(obj, proto); + + socket->js = js; + socket->obj = obj; + socket->encoding = js_mkundef(); + socket->allow_half_open = allow_half_open; + + js_set_native_ptr(obj, socket); + js_set_native_tag(obj, NET_SOCKET_NATIVE_TAG); + js_set(js, obj, "remoteAddress", js_mkundef()); + js_set(js, obj, "remotePort", js_mkundef()); + js_set(js, obj, "remoteFamily", js_mkundef()); + js_set(js, obj, "localAddress", js_mkundef()); + js_set(js, obj, "localPort", js_mkundef()); + js_set(js, obj, "localFamily", js_mkundef()); + net_socket_sync_state(socket); + + return socket; +} + +static void net_socket_attach_conn(net_socket_t *socket, ant_conn_t *conn) { + if (!socket || !conn) return; + + socket->conn = conn; + ant_conn_set_user_data(conn, socket); + if (ant_conn_has_remote_addr(conn)) { + js_set(socket->js, socket->obj, "remoteAddress", js_mkstr(socket->js, ant_conn_remote_addr(conn), strlen(ant_conn_remote_addr(conn)))); + js_set(socket->js, socket->obj, "remotePort", js_mknum(ant_conn_remote_port(conn))); + js_set(socket->js, socket->obj, "remoteFamily", js_mkstr(socket->js, ant_conn_remote_family(conn), strlen(ant_conn_remote_family(conn)))); + } + if (ant_conn_has_local_addr(conn)) { + js_set(socket->js, socket->obj, "localAddress", js_mkstr(socket->js, ant_conn_local_addr(conn), strlen(ant_conn_local_addr(conn)))); + js_set(socket->js, socket->obj, "localPort", js_mknum(ant_conn_local_port(conn))); + js_set(socket->js, socket->obj, "localFamily", js_mkstr(socket->js, ant_conn_local_family(conn), strlen(ant_conn_local_family(conn)))); + } + net_socket_sync_state(socket); +} + +static net_server_t *net_server_create(ant_t *js) { + ant_value_t obj = js_mkobj(js); + ant_value_t proto = js_instance_proto_from_new_target(js, g_net_server_proto); + net_server_t *server = calloc(1, sizeof(*server)); + + if (!server) return NULL; + if (is_object_type(proto)) js_set_proto_init(obj, proto); + + server->js = js; + server->obj = obj; + server->host = strdup("0.0.0.0"); + server->backlog = 511; + if (!server->host) { + free(server); + return NULL; + } + + js_set_native_ptr(obj, server); + js_set_native_tag(obj, NET_SERVER_NATIVE_TAG); + net_server_sync_state(server); + return server; +} + +static ant_value_t net_not_implemented(ant_t *js, const char *what) { + return js_mkerr_typed(js, JS_ERR_TYPE, "%s is not implemented yet", what); +} + +static ant_value_t net_isIP(ant_t *js, ant_value_t *args, int nargs) { + size_t len = 0; + const char *host = NULL; struct in_addr addr4; struct in6_addr addr6; + if (nargs < 1) return js_mknum(0); + host = js_getstr(js, args[0], &len); + if (!host) return js_mknum(0); + if (inet_pton(AF_INET, host, &addr4) == 1) return js_mknum(4); if (inet_pton(AF_INET6, host, &addr6) == 1) return js_mknum(6); return js_mknum(0); } static ant_value_t net_isIPv4(ant_t *js, ant_value_t *args, int nargs) { - if (nargs < 1) return js_false; - size_t len; - const char *host = js_getstr(js, args[0], &len); - if (!host) return js_false; - - struct in_addr addr; - return inet_pton(AF_INET, host, &addr) == 1 ? js_true : js_false; + if (js_getnum(net_isIP(js, args, nargs)) == 4.0) return js_true; + return js_false; } static ant_value_t net_isIPv6(ant_t *js, ant_value_t *args, int nargs) { - if (nargs < 1) return js_false; - size_t len; - const char *host = js_getstr(js, args[0], &len); - if (!host) return js_false; + if (js_getnum(net_isIP(js, args, nargs)) == 6.0) return js_true; + return js_false; +} + +static void net_server_apply_options(ant_t *js, net_server_t *server, ant_value_t options) { + ant_value_t value = 0; + + if (!server || vtype(options) != T_OBJ) return; + + value = js_get(js, options, "allowHalfOpen"); + if (vtype(value) != T_UNDEF) server->allow_half_open = js_truthy(js, value); + + value = js_get(js, options, "pauseOnConnect"); + if (vtype(value) != T_UNDEF) server->pause_on_connect = js_truthy(js, value); + + value = js_get(js, options, "noDelay"); + if (vtype(value) != T_UNDEF) server->no_delay = js_truthy(js, value); + + value = js_get(js, options, "keepAlive"); + if (vtype(value) != T_UNDEF) server->keep_alive = js_truthy(js, value); - struct in6_addr addr; - return inet_pton(AF_INET6, host, &addr) == 1 ? js_true : js_false; + value = js_get(js, options, "keepAliveInitialDelay"); + if (vtype(value) == T_NUM && js_getnum(value) > 0) + server->keep_alive_initial_delay_secs = (unsigned int)(js_getnum(value) / 1000.0); + + value = js_get(js, options, "backlog"); + if (vtype(value) == T_NUM && js_getnum(value) > 0) + server->backlog = (int)js_getnum(value); +} + +static void net_socket_on_read(ant_conn_t *conn, ssize_t nread, void *user_data) { + net_socket_t *socket = (net_socket_t *)ant_conn_get_user_data(conn); + const char *buffer = NULL; + ant_value_t chunk = 0; + size_t total = 0; + size_t offset = 0; + + if (!socket || !conn || nread <= 0) return; + + total = ant_conn_buffer_len(conn); + if ((size_t)nread > total) return; + offset = total - (size_t)nread; + buffer = ant_conn_buffer(conn) + offset; + + if (vtype(socket->encoding) == T_STR) + chunk = js_mkstr(socket->js, buffer, (size_t)nread); + else chunk = net_make_buffer_chunk(socket->js, buffer, (size_t)nread); + + ant_conn_consume(conn, total); + net_socket_sync_state(socket); + net_emit(socket->js, socket->obj, "data", &chunk, 1); +} + +static void net_socket_on_end(ant_conn_t *conn, void *user_data) { + net_socket_t *socket = (net_socket_t *)ant_conn_get_user_data(conn); + + if (!socket) return; + net_emit(socket->js, socket->obj, "end", NULL, 0); + if (!socket->allow_half_open && socket->conn) ant_conn_shutdown(socket->conn); +} + +static void net_socket_on_error(ant_conn_t *conn, int status, void *user_data) { + net_socket_t *socket = (net_socket_t *)ant_conn_get_user_data(conn); + + if (!socket) return; + socket->had_error = true; { + ant_value_t err = js_mkerr_typed(socket->js, JS_ERR_TYPE, "%s", uv_strerror(status)); + net_emit(socket->js, socket->obj, "error", &err, 1); + } +} + +static void net_socket_on_timeout(ant_conn_t *conn, void *user_data) { + net_socket_t *socket = (net_socket_t *)ant_conn_get_user_data(conn); + if (!socket) return; + net_emit(socket->js, socket->obj, "timeout", NULL, 0); +} + +static void net_server_on_conn_close(ant_conn_t *conn, void *user_data) { + net_socket_t *socket = (net_socket_t *)ant_conn_get_user_data(conn); + ant_value_t arg = 0; + + if (!socket) return; + arg = js_bool(socket->had_error); + socket->conn = NULL; + socket->destroyed = true; + net_socket_sync_state(socket); + net_emit(socket->js, socket->obj, "close", &arg, 1); + ant_conn_set_user_data(conn, NULL); + net_socket_detach(socket); +} + +static void net_server_on_listener_close(ant_listener_t *listener, void *user_data) { + net_server_maybe_finish_close((net_server_t *)user_data); +} + +static void net_server_on_accept(ant_listener_t *listener, ant_conn_t *conn, void *user_data) { + net_server_t *server = (net_server_t *)user_data; + net_socket_t *socket = NULL; + ant_value_t arg = 0; + + if (!server || !conn) return; + + if (server->max_connections > 0 && net_server_socket_count(server) >= server->max_connections) { + ant_conn_close(conn); + return; + } + + socket = net_socket_create(server->js, server->allow_half_open); + if (!socket) { + ant_conn_close(conn); + return; + } + + socket->server = server; + socket->next_in_server = server->sockets; + server->sockets = socket; + net_add_active_socket(socket); + net_socket_attach_conn(socket, conn); + + if (server->no_delay) ant_conn_set_no_delay(conn, true); + if (server->keep_alive) + ant_conn_set_keep_alive(conn, true, server->keep_alive_initial_delay_secs); + if (server->pause_on_connect) ant_conn_pause_read(conn); + + arg = socket->obj; + net_emit(server->js, server->obj, "connection", &arg, 1); +} + +static bool net_server_parse_host(const char *input, const char **out) { + if (!input || !*input) { + *out = "0.0.0.0"; + return true; + } + if (strcmp(input, "localhost") == 0) { + *out = "127.0.0.1"; + return true; + } + *out = input; + return true; +} + +static ant_value_t js_net_socket_ctor(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = NULL; + bool allow_half_open = false; + + if (nargs > 0 && vtype(args[0]) == T_OBJ) { + ant_value_t value = js_get(js, args[0], "allowHalfOpen"); + allow_half_open = js_truthy(js, value); + } + + socket = net_socket_create(js, allow_half_open); + if (!socket) return js_mkerr_typed(js, JS_ERR_TYPE, "Out of memory"); + return socket->obj; +} + +static ant_value_t js_net_socket_address(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + ant_value_t out = js_mkobj(js); + + if (!socket || !socket->conn || !ant_conn_has_local_addr(socket->conn)) return out; + js_set(js, out, "address", js_mkstr(js, ant_conn_local_addr(socket->conn), strlen(ant_conn_local_addr(socket->conn)))); + js_set(js, out, "port", js_mknum(ant_conn_local_port(socket->conn))); + js_set(js, out, "family", js_mkstr(js, ant_conn_local_family(socket->conn), strlen(ant_conn_local_family(socket->conn)))); + return out; +} + +static ant_value_t js_net_socket_pause(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + if (socket && socket->conn) ant_conn_pause_read(socket->conn); + return js_getthis(js); +} + +static ant_value_t js_net_socket_resume(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + if (socket && socket->conn) ant_conn_resume_read(socket->conn); + return js_getthis(js); +} + +static ant_value_t js_net_socket_setEncoding(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + if (!socket) return js_getthis(js); + socket->encoding = nargs > 0 && vtype(args[0]) != T_UNDEF ? js_tostring_val(js, args[0]) : js_mkundef(); + return js_getthis(js); +} + +static ant_value_t js_net_socket_setTimeout(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + double timeout = 0; + ant_value_t once_args[2] = {0}; + + if (!socket) return js_getthis(js); + if (nargs > 0 && vtype(args[0]) == T_NUM) timeout = js_getnum(args[0]); + if (socket->conn) ant_conn_set_timeout_ms(socket->conn, timeout > 0 ? (uint64_t)timeout : 0); + net_socket_sync_state(socket); + + if (nargs > 1 && is_callable(args[1])) { + once_args[0] = js_mkstr(js, "timeout", 7); + once_args[1] = args[1]; + net_call_method(js, socket->obj, "once", once_args, 2); + } + + return js_getthis(js); +} + +static ant_value_t js_net_socket_setNoDelay(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + bool enable = nargs == 0 || js_truthy(js, args[0]); + if (socket && socket->conn) ant_conn_set_no_delay(socket->conn, enable); + return js_getthis(js); +} + +static ant_value_t js_net_socket_setKeepAlive(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + bool enable = nargs > 0 && js_truthy(js, args[0]); + unsigned int delay = nargs > 1 && vtype(args[1]) == T_NUM ? (unsigned int)(js_getnum(args[1]) / 1000.0) : 0; + if (socket && socket->conn) ant_conn_set_keep_alive(socket->conn, enable, delay); + return js_getthis(js); +} + +static ant_value_t js_net_socket_ref(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + if (socket && socket->conn) ant_conn_ref(socket->conn); + return js_getthis(js); +} + +static ant_value_t js_net_socket_unref(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + if (socket && socket->conn) ant_conn_unref(socket->conn); + return js_getthis(js); +} + +static ant_value_t js_net_socket_write(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + const uint8_t *bytes = NULL; + size_t len = 0; + char *copy = NULL; + + if (!socket || !socket->conn) return js_false; + if (nargs < 1) return js_true; + + if (!buffer_source_get_bytes(js, args[0], &bytes, &len)) { + size_t slen = 0; + const char *str = js_getstr(js, js_tostring_val(js, args[0]), &slen); + if (!str) return js_false; + bytes = (const uint8_t *)str; + len = slen; + } + + copy = malloc(len); + if (!copy) return js_mkerr_typed(js, JS_ERR_TYPE, "Out of memory"); + if (len > 0) memcpy(copy, bytes, len); + + if (ant_conn_write(socket->conn, copy, len, NULL, NULL) != 0) { + free(copy); + return js_false; + } + + net_socket_sync_state(socket); + if (nargs > 2 && is_callable(args[2])) net_call_value(js, args[2], js_mkundef(), NULL, 0); + return js_true; +} + +static ant_value_t js_net_socket_end(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + ant_value_t result = js_getthis(js); + + if (!socket || !socket->conn) return result; + if (nargs > 0 && vtype(args[0]) != T_UNDEF && vtype(args[0]) != T_NULL) + js_net_socket_write(js, args, nargs); + ant_conn_shutdown(socket->conn); + if (nargs > 2 && is_callable(args[2])) net_call_value(js, args[2], js_mkundef(), NULL, 0); + return result; +} + +static ant_value_t js_net_socket_destroy(ant_t *js, ant_value_t *args, int nargs) { + net_socket_t *socket = net_socket_data(js_getthis(js)); + if (!socket) return js_getthis(js); + + if (nargs > 0 && vtype(args[0]) != T_UNDEF && vtype(args[0]) != T_NULL) { + ant_value_t err = args[0]; + socket->had_error = true; + net_emit(js, socket->obj, "error", &err, 1); + } + + if (socket->conn) ant_conn_close(socket->conn); + return js_getthis(js); +} + +static ant_value_t js_net_socket_connect(ant_t *js, ant_value_t *args, int nargs) { + return net_not_implemented(js, "net.Socket.connect"); +} + +static ant_value_t js_net_server_ctor(ant_t *js, ant_value_t *args, int nargs) { + net_server_t *server = net_server_create(js); + ant_value_t on_args[2] = {0}; + + if (!server) return js_mkerr_typed(js, JS_ERR_TYPE, "Out of memory"); + if (nargs > 0 && vtype(args[0]) == T_OBJ) net_server_apply_options(js, server, args[0]); + if (nargs > 0 && is_callable(args[0])) { + on_args[0] = js_mkstr(js, "connection", 10); + on_args[1] = args[0]; + net_call_method(js, server->obj, "on", on_args, 2); + } else if (nargs > 1 && is_callable(args[1])) { + on_args[0] = js_mkstr(js, "connection", 10); + on_args[1] = args[1]; + net_call_method(js, server->obj, "on", on_args, 2); + } + + return server->obj; +} + +static ant_value_t js_net_server_listen(ant_t *js, ant_value_t *args, int nargs) { + net_server_t *server = net_server_data(js_getthis(js)); + ant_listener_callbacks_t callbacks = {0}; + const char *host = "0.0.0.0"; + ant_value_t cb = js_mkundef(); + int port = 0; + int backlog = 511; + int rc = 0; + + if (!server) return js_mkerr_typed(js, JS_ERR_TYPE, "Invalid net.Server"); + if (server->listening) return js_mkerr_typed(js, JS_ERR_TYPE, "Server is already listening"); + + if (nargs > 0 && vtype(args[0]) == T_NUM) { + port = (int)js_getnum(args[0]); + if (nargs > 1 && vtype(args[1]) == T_STR) { + size_t len = 0; + host = js_getstr(js, args[1], &len); + } + if (nargs > 2 && vtype(args[2]) == T_NUM) backlog = (int)js_getnum(args[2]); + if (nargs > 3 && is_callable(args[3])) cb = args[3]; + else if (nargs > 2 && is_callable(args[2])) cb = args[2]; + else if (nargs > 1 && is_callable(args[1])) cb = args[1]; + } else if (nargs > 0 && vtype(args[0]) == T_OBJ) { + ant_value_t value = js_get(js, args[0], "path"); + if (vtype(value) != T_UNDEF) return net_not_implemented(js, "IPC server listen"); + + value = js_get(js, args[0], "port"); + if (vtype(value) == T_NUM) port = (int)js_getnum(value); + value = js_get(js, args[0], "host"); + if (vtype(value) == T_STR) { + size_t len = 0; + host = js_getstr(js, value, &len); + } + value = js_get(js, args[0], "backlog"); + if (vtype(value) == T_NUM) backlog = (int)js_getnum(value); + if (nargs > 1 && is_callable(args[1])) cb = args[1]; + } else if (nargs > 0 && vtype(args[0]) == T_STR) { + return net_not_implemented(js, "IPC server listen"); + } else if (nargs > 0 && is_callable(args[0])) { + cb = args[0]; + } + + net_server_parse_host(host, &host); + free(server->host); + server->host = strdup(host ? host : "0.0.0.0"); + server->port = port; + server->backlog = backlog > 0 ? backlog : server->backlog; + + callbacks.on_accept = net_server_on_accept; + callbacks.on_read = net_socket_on_read; + callbacks.on_end = net_socket_on_end; + callbacks.on_error = net_socket_on_error; + callbacks.on_timeout = net_socket_on_timeout; + callbacks.on_conn_close = net_server_on_conn_close; + callbacks.on_listener_close = net_server_on_listener_close; + + rc = ant_listener_listen_tcp( + &server->listener, + uv_default_loop(), + server->host, + port, + server->backlog, + 0, + &callbacks, + server + ); + if (rc != 0) return js_mkerr_typed(js, JS_ERR_TYPE, "%s", uv_strerror(rc)); + + server->port = ant_listener_port(&server->listener); + server->listening = true; + server->closing = false; + net_server_sync_state(server); + net_add_active_server(server); + + if (is_callable(cb)) net_call_value(js, cb, js_mkundef(), NULL, 0); + net_emit(js, server->obj, "listening", NULL, 0); + return js_getthis(js); +} + +static ant_value_t js_net_server_close(ant_t *js, ant_value_t *args, int nargs) { + net_server_t *server = net_server_data(js_getthis(js)); + ant_value_t once_args[2] = {0}; + + if (!server) return js_getthis(js); + + if (nargs > 0 && is_callable(args[0])) { + once_args[0] = js_mkstr(js, "close", 5); + once_args[1] = args[0]; + net_call_method(js, server->obj, "once", once_args, 2); + } + + if (!server->listening && !server->closing) { + if (nargs > 0 && is_callable(args[0])) { + ant_value_t err = js_mkerr_typed(js, JS_ERR_TYPE, "Server is not running"); + net_call_value(js, args[0], js_mkundef(), &err, 1); + } + return js_getthis(js); + } + + server->closing = true; + ant_listener_stop(&server->listener, false); + net_server_maybe_finish_close(server); + return js_getthis(js); +} + +static ant_value_t js_net_server_address(ant_t *js, ant_value_t *args, int nargs) { + net_server_t *server = net_server_data(js_getthis(js)); + ant_value_t out = js_mknull(); + if (!server || !server->listening) return out; + + out = js_mkobj(js); + js_set(js, out, "port", js_mknum(server->port)); + js_set(js, out, "family", js_mkstr(js, "IPv4", 4)); + js_set(js, out, "address", js_mkstr(js, server->host, strlen(server->host))); + return out; +} + +static ant_value_t js_net_server_getConnections(ant_t *js, ant_value_t *args, int nargs) { + net_server_t *server = net_server_data(js_getthis(js)); + ant_value_t cb = nargs > 0 ? args[0] : js_mkundef(); + ant_value_t argv[2] = { js_mknull(), js_mknum((double)net_server_socket_count(server)) }; + + if (is_callable(cb)) net_call_value(js, cb, js_mkundef(), argv, 2); + return js_getthis(js); +} + +static ant_value_t js_net_server_ref(ant_t *js, ant_value_t *args, int nargs) { + net_server_t *server = net_server_data(js_getthis(js)); + if (server && !uv_is_closing((uv_handle_t *)&server->listener.handle)) + uv_ref((uv_handle_t *)&server->listener.handle); + return js_getthis(js); +} + +static ant_value_t js_net_server_unref(ant_t *js, ant_value_t *args, int nargs) { + net_server_t *server = net_server_data(js_getthis(js)); + if (server && !uv_is_closing((uv_handle_t *)&server->listener.handle)) + uv_unref((uv_handle_t *)&server->listener.handle); + return js_getthis(js); +} + +static ant_value_t js_net_createServer(ant_t *js, ant_value_t *args, int nargs) { + return js_net_server_ctor(js, args, nargs); +} + +static ant_value_t js_net_createConnection(ant_t *js, ant_value_t *args, int nargs) { + return net_not_implemented(js, "net.createConnection"); +} + +static ant_value_t js_net_connect(ant_t *js, ant_value_t *args, int nargs) { + return js_net_createConnection(js, args, nargs); +} + +static ant_value_t js_net_getDefaultAutoSelectFamily(ant_t *js, ant_value_t *args, int nargs) { + return js_bool(g_default_auto_select_family); +} + +static ant_value_t js_net_setDefaultAutoSelectFamily(ant_t *js, ant_value_t *args, int nargs) { + if (nargs > 0) g_default_auto_select_family = js_truthy(js, args[0]); + return js_mkundef(); +} + +static ant_value_t js_net_getDefaultAutoSelectFamilyAttemptTimeout(ant_t *js, ant_value_t *args, int nargs) { + return js_mknum(g_default_auto_select_family_attempt_timeout); +} + +static ant_value_t js_net_setDefaultAutoSelectFamilyAttemptTimeout(ant_t *js, ant_value_t *args, int nargs) { + if (nargs > 0 && vtype(args[0]) == T_NUM) { + double value = js_getnum(args[0]); + if (value > 0 && value < 10) value = 10; + if (value > 0) g_default_auto_select_family_attempt_timeout = value; + } + return js_mkundef(); +} + +static void net_init_constructors(ant_t *js) { + ant_value_t events = 0; + ant_value_t ee_ctor = 0; + ant_value_t ee_proto = 0; + + if (g_net_server_ctor && g_net_socket_ctor) return; + + events = events_library(js); + ee_ctor = js_get(js, events, "EventEmitter"); + ee_proto = js_get(js, ee_ctor, "prototype"); + + g_net_socket_proto = js_mkobj(js); + js_set_proto_init(g_net_socket_proto, ee_proto); + js_set(js, g_net_socket_proto, "address", js_mkfun(js_net_socket_address)); + js_set(js, g_net_socket_proto, "pause", js_mkfun(js_net_socket_pause)); + js_set(js, g_net_socket_proto, "resume", js_mkfun(js_net_socket_resume)); + js_set(js, g_net_socket_proto, "setEncoding", js_mkfun(js_net_socket_setEncoding)); + js_set(js, g_net_socket_proto, "setTimeout", js_mkfun(js_net_socket_setTimeout)); + js_set(js, g_net_socket_proto, "setNoDelay", js_mkfun(js_net_socket_setNoDelay)); + js_set(js, g_net_socket_proto, "setKeepAlive", js_mkfun(js_net_socket_setKeepAlive)); + js_set(js, g_net_socket_proto, "write", js_mkfun(js_net_socket_write)); + js_set(js, g_net_socket_proto, "end", js_mkfun(js_net_socket_end)); + js_set(js, g_net_socket_proto, "destroy", js_mkfun(js_net_socket_destroy)); + js_set(js, g_net_socket_proto, "connect", js_mkfun(js_net_socket_connect)); + js_set(js, g_net_socket_proto, "ref", js_mkfun(js_net_socket_ref)); + js_set(js, g_net_socket_proto, "unref", js_mkfun(js_net_socket_unref)); + js_set_sym(js, g_net_socket_proto, get_toStringTag_sym(), js_mkstr(js, "Socket", 6)); + g_net_socket_ctor = js_make_ctor(js, js_net_socket_ctor, g_net_socket_proto, "Socket", 6); + + g_net_server_proto = js_mkobj(js); + js_set_proto_init(g_net_server_proto, ee_proto); + js_set(js, g_net_server_proto, "listen", js_mkfun(js_net_server_listen)); + js_set(js, g_net_server_proto, "close", js_mkfun(js_net_server_close)); + js_set(js, g_net_server_proto, "address", js_mkfun(js_net_server_address)); + js_set(js, g_net_server_proto, "getConnections", js_mkfun(js_net_server_getConnections)); + js_set(js, g_net_server_proto, "ref", js_mkfun(js_net_server_ref)); + js_set(js, g_net_server_proto, "unref", js_mkfun(js_net_server_unref)); + js_set_sym(js, g_net_server_proto, get_toStringTag_sym(), js_mkstr(js, "Server", 6)); + g_net_server_ctor = js_make_ctor(js, js_net_server_ctor, g_net_server_proto, "Server", 6); } ant_value_t net_library(ant_t *js) { ant_value_t lib = js_mkobj(js); + net_init_constructors(js); + + js_set(js, lib, "Server", g_net_server_ctor); + js_set(js, lib, "Socket", g_net_socket_ctor); + js_set(js, lib, "createServer", js_mkfun(js_net_createServer)); + js_set(js, lib, "createConnection", js_mkfun(js_net_createConnection)); + js_set(js, lib, "connect", js_mkfun(js_net_connect)); js_set(js, lib, "isIP", js_mkfun(net_isIP)); js_set(js, lib, "isIPv4", js_mkfun(net_isIPv4)); js_set(js, lib, "isIPv6", js_mkfun(net_isIPv6)); - + js_set(js, lib, "getDefaultAutoSelectFamily", js_mkfun(js_net_getDefaultAutoSelectFamily)); + js_set(js, lib, "setDefaultAutoSelectFamily", js_mkfun(js_net_setDefaultAutoSelectFamily)); + js_set(js, lib, "getDefaultAutoSelectFamilyAttemptTimeout", js_mkfun(js_net_getDefaultAutoSelectFamilyAttemptTimeout)); + js_set(js, lib, "setDefaultAutoSelectFamilyAttemptTimeout", js_mkfun(js_net_setDefaultAutoSelectFamilyAttemptTimeout)); + js_set(js, lib, "default", lib); + js_set_sym(js, lib, get_toStringTag_sym(), js_mkstr(js, "net", 3)); return lib; } + +void gc_mark_net(ant_t *js, gc_mark_fn mark) { + net_server_t *server = NULL; + net_socket_t *socket = NULL; + + if (g_net_server_proto) mark(js, g_net_server_proto); + if (g_net_socket_proto) mark(js, g_net_socket_proto); + if (g_net_server_ctor) mark(js, g_net_server_ctor); + if (g_net_socket_ctor) mark(js, g_net_socket_ctor); + + for (server = g_active_servers; server; server = server->next_active) + mark(js, server->obj); + + for (socket = g_active_sockets; socket; socket = socket->next_active) { + mark(js, socket->obj); + if (vtype(socket->encoding) != T_UNDEF) mark(js, socket->encoding); + } +} diff --git a/src/modules/server.c b/src/modules/server.c index 4a92e97..8d86ea4 100644 --- a/src/modules/server.c +++ b/src/modules/server.c @@ -625,7 +625,8 @@ static ant_value_t server_timeout(ant_t *js, ant_value_t *args, int nargs) { if (!req || !req->conn) return js_mkundef(); timeout = (int)js_getnum(args[1]); - ant_conn_set_timeout(req->conn, timeout); + ant_conn_set_timeout_ms(req->conn, (uint64_t)timeout * 1000ULL); + return js_mkundef(); } @@ -728,6 +729,10 @@ static void server_on_read(ant_conn_t *conn, ssize_t nread, void *user_data) { } } +static void server_on_end(ant_conn_t *conn, void *user_data) { + if (conn) ant_conn_close(conn); +} + static void server_on_conn_close(ant_conn_t *conn, void *user_data) { server_runtime_t *server = (server_runtime_t *)user_data; server_request_t *req = (server_request_t *)ant_conn_get_user_data(conn); @@ -899,13 +904,14 @@ ant_value_t server_start_from_export(ant_t *js, ant_value_t default_export) { server->sigterm_handle.data = server; callbacks.on_read = server_on_read; + callbacks.on_end = server_on_end; callbacks.on_conn_close = server_on_conn_close; callbacks.on_listener_close = server_on_listener_close; rc = ant_listener_listen_tcp( &server->listener, server->loop, server->hostname, server->port, - 128, server->idle_timeout_secs, &callbacks, server + 128, (uint64_t)server->idle_timeout_secs * 1000ULL, &callbacks, server ); if (rc != 0) { @@ -915,7 +921,6 @@ ant_value_t server_start_from_export(ant_t *js, ant_value_t default_export) { } server->port = ant_listener_port(&server->listener); - uv_signal_start(&server->sigint_handle, server_signal_cb, SIGINT); uv_signal_start(&server->sigterm_handle, server_signal_cb, SIGTERM); diff --git a/src/modules/v8.c b/src/modules/v8.c index c9dc615..d41eecc 100644 --- a/src/modules/v8.c +++ b/src/modules/v8.c @@ -1,5 +1,3 @@ -// stub: functional where possible - #include #include #include diff --git a/src/net/connection.c b/src/net/connection.c index 86133ad..3a499e2 100644 --- a/src/net/connection.c +++ b/src/net/connection.c @@ -18,6 +18,10 @@ typedef struct { void *user_data; } ant_conn_write_req_t; +typedef struct { + uv_shutdown_t req; +} ant_conn_shutdown_req_t; + static void ant_conn_restart_timer(ant_conn_t *conn); static void ant_conn_close_cb(uv_handle_t *handle); @@ -32,42 +36,85 @@ static void ant_listener_remove_conn(ant_listener_t *listener, ant_conn_t *conn) }} } -static bool ant_conn_store_peer_addr(ant_conn_t *conn) { +static bool ant_conn_store_addr( + uv_tcp_t *handle, + int (*get_name)(const uv_tcp_t *, struct sockaddr *, int *), + char *out_addr, + size_t out_addr_len, + char *out_family, + size_t out_family_len, + int *out_port, + bool *out_has_addr +) { struct sockaddr_storage addr; int len = sizeof(addr); - if (uv_tcp_getpeername(&conn->handle, (struct sockaddr *)&addr, &len) != 0) return false; + if (get_name(handle, (struct sockaddr *)&addr, &len) != 0) return false; if (addr.ss_family == AF_INET) { struct sockaddr_in *in = (struct sockaddr_in *)&addr; - uv_ip4_name(in, conn->remote_addr, sizeof(conn->remote_addr)); - conn->remote_port = ntohs(in->sin_port); - conn->has_remote_addr = true; + uv_ip4_name(in, out_addr, out_addr_len); + memcpy(out_family, "IPv4", out_family_len < 5 ? out_family_len : 5); + *out_port = ntohs(in->sin_port); + *out_has_addr = true; return true; } + if (addr.ss_family == AF_INET6) { struct sockaddr_in6 *in6 = (struct sockaddr_in6 *)&addr; - uv_ip6_name(in6, conn->remote_addr, sizeof(conn->remote_addr)); - conn->remote_port = ntohs(in6->sin6_port); - conn->has_remote_addr = true; + uv_ip6_name(in6, out_addr, out_addr_len); + memcpy(out_family, "IPv6", out_family_len < 5 ? out_family_len : 5); + *out_port = ntohs(in6->sin6_port); + *out_has_addr = true; return true; } + return false; } +static bool ant_conn_store_peer_addr(ant_conn_t *conn) { + return ant_conn_store_addr( + &conn->handle, + (int (*)(const uv_tcp_t *, struct sockaddr *, int *))uv_tcp_getpeername, + conn->remote_addr, + sizeof(conn->remote_addr), + conn->remote_family, + sizeof(conn->remote_family), + &conn->remote_port, + &conn->has_remote_addr + ); +} + +static bool ant_conn_store_local_addr(ant_conn_t *conn) { + return ant_conn_store_addr( + &conn->handle, + (int (*)(const uv_tcp_t *, struct sockaddr *, int *))uv_tcp_getsockname, + conn->local_addr, + sizeof(conn->local_addr), + conn->local_family, + sizeof(conn->local_family), + &conn->local_port, + &conn->has_local_addr + ); +} + static void ant_conn_timeout_cb(uv_timer_t *handle) { ant_conn_t *conn = (ant_conn_t *)handle->data; + ant_listener_t *listener = conn ? conn->listener : NULL; + + if (!conn) return; + if (listener && listener->callbacks.on_timeout) { + listener->callbacks.on_timeout(conn, listener->user_data); + return; + } + ant_conn_close(conn); } static void ant_conn_restart_timer(ant_conn_t *conn) { - uint64_t timeout_ms = 0; - if (!conn || conn->closing) return; uv_timer_stop(&conn->timer); - if (conn->timeout_secs <= 0) return; - - timeout_ms = (uint64_t)conn->timeout_secs * 1000ULL; - uv_timer_start(&conn->timer, ant_conn_timeout_cb, timeout_ms, 0); + if (conn->timeout_ms == 0) return; + uv_timer_start(&conn->timer, ant_conn_timeout_cb, conn->timeout_ms, 0); } static void ant_conn_alloc_cb(uv_handle_t *handle, size_t suggested_size, uv_buf_t *buf) { @@ -103,15 +150,28 @@ static void ant_conn_alloc_cb(uv_handle_t *handle, size_t suggested_size, uv_buf static void ant_conn_read_cb(uv_stream_t *stream, ssize_t nread, const uv_buf_t *buf) { ant_conn_t *conn = (ant_conn_t *)stream->data; ant_listener_t *listener = conn ? conn->listener : NULL; - (void)buf; - if (!conn || !listener) return; + if (nread < 0) { + if (nread == UV_EOF) { + conn->read_eof = true; + uv_read_stop((uv_stream_t *)&conn->handle); + if (listener->callbacks.on_end) + listener->callbacks.on_end(conn, listener->user_data); + return; + } + + if (listener->callbacks.on_error) listener->callbacks.on_error( + conn, (int)nread, listener->user_data + ); + ant_conn_close(conn); return; } + if (nread == 0) return; + conn->bytes_read += (uint64_t)nread; conn->buffer_len += (size_t)nread; ant_conn_restart_timer(conn); @@ -122,11 +182,19 @@ static void ant_conn_read_cb(uv_stream_t *stream, ssize_t nread, const uv_buf_t static void ant_conn_write_cb_impl(uv_write_t *req, int status) { ant_conn_write_req_t *wr = (ant_conn_write_req_t *)req; + if (status >= 0 && wr->conn) + wr->conn->bytes_written += (uint64_t)wr->buf.len; if (wr->cb) wr->cb(wr->conn, status, wr->user_data); + free(wr->buf.base); free(wr); } +static void ant_conn_shutdown_cb(uv_shutdown_t *req, int status) { + ant_conn_shutdown_req_t *shutdown_req = (ant_conn_shutdown_req_t *)req; + free(shutdown_req); +} + static void ant_conn_close_cb(uv_handle_t *handle) { ant_conn_t *conn = (ant_conn_t *)handle->data; ant_listener_t *listener = conn ? conn->listener : NULL; @@ -142,7 +210,7 @@ static void ant_conn_close_cb(uv_handle_t *handle) { free(conn); } -ant_conn_t *ant_conn_create_tcp(ant_listener_t *listener, int timeout_secs) { +ant_conn_t *ant_conn_create_tcp(ant_listener_t *listener, uint64_t timeout_ms) { ant_conn_t *conn = NULL; if (!listener || !listener->loop) return NULL; @@ -151,7 +219,7 @@ ant_conn_t *ant_conn_create_tcp(ant_listener_t *listener, int timeout_secs) { if (!conn) return NULL; conn->listener = listener; - conn->timeout_secs = timeout_secs; + conn->timeout_ms = timeout_ms; conn->buffer_cap = ANT_CONN_READ_BUFFER_SIZE; conn->buffer = malloc(conn->buffer_cap); if (!conn->buffer) { @@ -173,6 +241,7 @@ int ant_conn_accept(ant_conn_t *conn, uv_stream_t *server_stream) { return UV_ECONNABORTED; ant_conn_store_peer_addr(conn); + ant_conn_store_local_addr(conn); conn->next = conn->listener->connections; conn->listener->connections = conn; return 0; @@ -180,6 +249,7 @@ int ant_conn_accept(ant_conn_t *conn, uv_stream_t *server_stream) { void ant_conn_start(ant_conn_t *conn) { if (!conn || conn->closing) return; + conn->read_paused = false; ant_conn_restart_timer(conn); uv_read_start((uv_stream_t *)&conn->handle, ant_conn_alloc_cb, ant_conn_read_cb); } @@ -204,37 +274,117 @@ size_t ant_conn_buffer_len(const ant_conn_t *conn) { return conn ? conn->buffer_len : 0; } -void ant_conn_set_timeout(ant_conn_t *conn, int timeout_secs) { +void ant_conn_set_timeout_ms(ant_conn_t *conn, uint64_t timeout_ms) { if (!conn) return; - conn->timeout_secs = timeout_secs < 0 ? 0 : timeout_secs; + conn->timeout_ms = timeout_ms; ant_conn_restart_timer(conn); } -int ant_conn_timeout(const ant_conn_t *conn) { - return conn ? conn->timeout_secs : 0; +uint64_t ant_conn_timeout_ms(const ant_conn_t *conn) { + return conn ? conn->timeout_ms : 0; } bool ant_conn_is_closing(const ant_conn_t *conn) { return !conn || conn->closing; } +bool ant_conn_has_local_addr(const ant_conn_t *conn) { + return conn && conn->has_local_addr; +} + bool ant_conn_has_remote_addr(const ant_conn_t *conn) { return conn && conn->has_remote_addr; } +const char *ant_conn_local_addr(const ant_conn_t *conn) { + return conn ? conn->local_addr : NULL; +} + const char *ant_conn_remote_addr(const ant_conn_t *conn) { return conn ? conn->remote_addr : NULL; } +const char *ant_conn_local_family(const ant_conn_t *conn) { + return conn ? conn->local_family : NULL; +} + +const char *ant_conn_remote_family(const ant_conn_t *conn) { + return conn ? conn->remote_family : NULL; +} + +int ant_conn_local_port(const ant_conn_t *conn) { + return conn ? conn->local_port : 0; +} + int ant_conn_remote_port(const ant_conn_t *conn) { return conn ? conn->remote_port : 0; } void ant_conn_pause_read(ant_conn_t *conn) { if (!conn || conn->closing) return; + conn->read_paused = true; uv_read_stop((uv_stream_t *)&conn->handle); } +void ant_conn_resume_read(ant_conn_t *conn) { + if (!conn || conn->closing || conn->read_eof || !conn->read_paused) return; + conn->read_paused = false; + ant_conn_restart_timer(conn); + uv_read_start((uv_stream_t *)&conn->handle, ant_conn_alloc_cb, ant_conn_read_cb); +} + +void ant_conn_shutdown(ant_conn_t *conn) { + ant_conn_shutdown_req_t *req = NULL; + int rc = 0; + + if (!conn || conn->closing) return; + + req = calloc(1, sizeof(*req)); + if (!req) { + ant_conn_close(conn); + return; + } + + rc = uv_shutdown(&req->req, (uv_stream_t *)&conn->handle, ant_conn_shutdown_cb); + if (rc != 0) { + free(req); + ant_conn_close(conn); + } +} + +void ant_conn_ref(ant_conn_t *conn) { + if (!conn) return; + if (!uv_is_closing((uv_handle_t *)&conn->handle)) uv_ref((uv_handle_t *)&conn->handle); + if (!uv_is_closing((uv_handle_t *)&conn->timer)) uv_ref((uv_handle_t *)&conn->timer); +} + +void ant_conn_unref(ant_conn_t *conn) { + if (!conn) return; + if (!uv_is_closing((uv_handle_t *)&conn->handle)) uv_unref((uv_handle_t *)&conn->handle); + if (!uv_is_closing((uv_handle_t *)&conn->timer)) uv_unref((uv_handle_t *)&conn->timer); +} + +int ant_conn_set_no_delay(ant_conn_t *conn, bool enable) { + if (!conn || conn->closing) return UV_EINVAL; + return uv_tcp_nodelay(&conn->handle, enable ? 1 : 0); +} + +int ant_conn_set_keep_alive(ant_conn_t *conn, bool enable, unsigned int delay_secs) { + if (!conn || conn->closing) return UV_EINVAL; + return uv_tcp_keepalive(&conn->handle, enable ? 1 : 0, delay_secs); +} + +void ant_conn_consume(ant_conn_t *conn, size_t len) { + if (!conn || conn->buffer_len == 0 || len == 0) return; + if (len >= conn->buffer_len) { + conn->buffer_len = 0; + return; + } + + memmove(conn->buffer, conn->buffer + len, conn->buffer_len - len); + conn->buffer_len -= len; +} + void ant_conn_close(ant_conn_t *conn) { if (!conn || conn->closing) return; conn->closing = true; @@ -284,3 +434,11 @@ int ant_conn_write(ant_conn_t *conn, char *data, size_t len, ant_conn_write_cb c return 0; } + +uint64_t ant_conn_bytes_read(const ant_conn_t *conn) { + return conn ? conn->bytes_read : 0; +} + +uint64_t ant_conn_bytes_written(const ant_conn_t *conn) { + return conn ? conn->bytes_written : 0; +} diff --git a/src/net/listener.c b/src/net/listener.c index 83a0385..24879ef 100644 --- a/src/net/listener.c +++ b/src/net/listener.c @@ -22,7 +22,7 @@ static void ant_listener_accept_cb(uv_stream_t *server_stream, int status) { if (status < 0 || !listener || listener->closing) return; - conn = ant_conn_create_tcp(listener, listener->idle_timeout_secs); + conn = ant_conn_create_tcp(listener, listener->idle_timeout_ms); if (!conn) return; if (ant_conn_accept(conn, server_stream) != 0) { @@ -42,7 +42,7 @@ int ant_listener_listen_tcp( const char *hostname, int port, int backlog, - int idle_timeout_secs, + uint64_t idle_timeout_ms, const ant_listener_callbacks_t *callbacks, void *user_data ) { @@ -55,7 +55,7 @@ int ant_listener_listen_tcp( listener->loop = loop; listener->user_data = user_data; - listener->idle_timeout_secs = idle_timeout_secs; + listener->idle_timeout_ms = idle_timeout_ms; listener->port = port; listener->backlog = backlog > 0 ? backlog : 128; if (callbacks) listener->callbacks = *callbacks; -- 2.51.2