diff --git a/include/compat.h b/include/compat.h index c67adbe..9a7cd34 100644 --- a/include/compat.h +++ b/include/compat.h @@ -123,6 +123,8 @@ static inline int compat_unsetenv(const char *name) { #define timegm compat_timegm #define usleep compat_usleep #define sleep compat_sleep +#define strcasecmp _stricmp +#define strncasecmp _strnicmp #define strndup compat_strndup #define memmem compat_memmem #define basename compat_basename @@ -146,6 +148,7 @@ typedef unsigned int gid_t; #include #include #include +#include #endif #endif diff --git a/src/http/websocket.c b/src/http/websocket.c index b759080..27625c9 100644 --- a/src/http/websocket.c +++ b/src/http/websocket.c @@ -1,11 +1,10 @@ -#include -#include +#include // IWYU pragma: keep + #include #include #include #include #include -#include #include #include "base64.h" diff --git a/src/modules/net.c b/src/modules/net.c index f558ed0..c09b1be 100644 --- a/src/modules/net.c +++ b/src/modules/net.c @@ -398,7 +398,7 @@ static void net_socket_sync_state(net_socket_t *socket) { if (socket->destroyed) ready_state = "closed"; else if (socket->connecting) ready_state = "opening"; 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, "pending", js_bool(socket->connecting || (socket->conn == NULL && !socket->destroyed))); js_set(socket->js, obj, "connecting", js_bool(socket->connecting)); 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))); diff --git a/src/modules/stream.c b/src/modules/stream.c index a82452f..fc281a6 100644 --- a/src/modules/stream.c +++ b/src/modules/stream.c @@ -259,6 +259,14 @@ static ant_offset_t stream_readable_buffer_len(ant_t *js, ant_value_t stream_obj return len > head ? len - head : 0; } +static double stream_chunk_size(ant_t *js, ant_value_t chunk, bool object_mode) { + const uint8_t *bytes = NULL; + size_t byte_len = 0; + if (object_mode) return 1; + if (buffer_source_get_bytes(js, chunk, &bytes, &byte_len)) return (double)byte_len; + return 1; +} + static void stream_compact_readable_buffer(ant_t *js, ant_value_t stream_obj) { ant_value_t state = stream_readable_state(js, stream_obj); ant_value_t buffer = stream_readable_buffer(js, stream_obj); @@ -287,9 +295,14 @@ static void stream_compact_readable_buffer(ant_t *js, ant_value_t stream_obj) { static void stream_buffer_push(ant_t *js, ant_value_t stream_obj, ant_value_t value) { ant_value_t state = stream_readable_state(js, stream_obj); ant_value_t buffer = stream_readable_buffer(js, stream_obj); + ant_value_t length = 0; + bool object_mode = false; if (!is_object_type(state) || vtype(buffer) != T_ARR) return; js_arr_push(js, buffer, value); + object_mode = js_truthy(js, js_get(js, state, "objectMode")); + length = js_get(js, state, "length"); + js_set(js, state, "length", js_mknum((vtype(length) == T_NUM ? tod(length) : 0) + stream_chunk_size(js, value, object_mode))); } static ant_value_t stream_buffer_shift(ant_t *js, ant_value_t stream_obj) { @@ -302,7 +315,14 @@ static ant_value_t stream_buffer_shift(ant_t *js, ant_value_t stream_obj) { value = js_arr_get(js, buffer, head); ant_value_t state = stream_readable_state(js, stream_obj); - if (is_object_type(state)) js_set(js, state, "dataEmitted", js_true); + if (is_object_type(state)) { + ant_value_t length = js_get(js, state, "length"); + bool object_mode = js_truthy(js, js_get(js, state, "objectMode")); + double next_length = (vtype(length) == T_NUM ? tod(length) : 0) - stream_chunk_size(js, value, object_mode); + js_set(js, state, "length", js_mknum(next_length > 0 ? next_length : 0)); + js_set(js, state, "dataEmitted", js_true); + } + stream_set_readable_buffer_head(js, stream_obj, head + 1); stream_compact_readable_buffer(js, stream_obj); @@ -390,6 +410,7 @@ static void stream_init_readable(ant_t *js, ant_value_t obj, ant_value_t raw_opt js_set(js, state, "flowingReadScheduled", js_false); js_set(js, state, "reading", js_false); js_set(js, state, "highWaterMark", js_mknum(high_water_mark)); + js_set(js, state, "length", js_mknum(0)); js_set(js, state, "buffer", js_mkarr(js)); js_set(js, state, "bufferHead", js_mknum(0)); js_set(js, obj, "_readableState", state); @@ -807,18 +828,25 @@ ant_value_t stream_readable_flush(ant_t *js, ant_value_t stream_obj) { ant_value_t stream_readable_maybe_read(ant_t *js, ant_value_t stream_obj) { ant_value_t state = stream_readable_state(js, stream_obj); + ant_value_t read_fn = 0; ant_value_t args[1]; + ant_value_t hwm = 0; + ant_value_t length = 0; if (!is_object_type(state)) return js_mkundef(); if (js_truthy(js, js_get(js, stream_obj, "destroyed"))) return js_mkundef(); if (js_truthy(js, js_get(js, state, "reading"))) return js_mkundef(); if (js_truthy(js, js_get(js, state, "ended"))) return js_mkundef(); - if (stream_readable_buffer_len(js, stream_obj) > 0) return js_mkundef(); + + hwm = js_get(js, state, "highWaterMark"); + length = js_get(js, state, "length"); + if (vtype(hwm) == T_NUM && vtype(length) == T_NUM && tod(length) >= tod(hwm)) + return js_mkundef(); read_fn = js_getprop_fallback(js, stream_obj, "_read"); js_set(js, state, "reading", js_true); - args[0] = js_get(js, state, "highWaterMark"); + args[0] = hwm; if (is_callable(read_fn)) stream_call(js, read_fn, stream_obj, args, 1, false); js_set(js, state, "reading", js_false); @@ -874,8 +902,12 @@ static ant_value_t js_readable_read(ant_t *js, ant_value_t *args, int nargs) { state = stream_readable_state(js, stream_obj); if (!is_object_type(state)) return js_mknull(); + if (nargs > 0 && vtype(args[0]) == T_NUM && tod(args[0]) == 0.0) { + stream_readable_maybe_read(js, stream_obj); + return js_mknull(); + } + if (stream_readable_buffer_len(js, stream_obj) == 0) stream_readable_maybe_read(js, stream_obj); - if (nargs > 0 && vtype(args[0]) == T_NUM && tod(args[0]) == 0.0) return js_mknull(); if (stream_readable_buffer_len(js, stream_obj) == 0) return js_mknull(); chunk = stream_buffer_shift(js, stream_obj); diff --git a/src/modules/tls.c b/src/modules/tls.c index 50e42ab..6c5bddc 100644 --- a/src/modules/tls.c +++ b/src/modules/tls.c @@ -22,6 +22,7 @@ #include "modules/events.h" #include "modules/net.h" #include "modules/symbol.h" +#include "modules/timer.h" #include "silver/engine.h" typedef struct @@ -79,6 +80,7 @@ typedef struct ant_tls_socket_s { bool closing; bool had_error; bool ended; + bool read_drain_scheduled; } ant_tls_socket_t; enum { @@ -251,6 +253,48 @@ static void tls_socket_free_read_queue(ant_tls_socket_t *socket) { } } +static ant_value_t tls_socket_drain_read_queue(ant_t *js, ant_value_t *args, int nargs) { + ant_value_t obj = js_get_slot(js_getcurrentfunc(js), SLOT_DATA); + ant_tls_socket_t *socket = tls_socket_data(obj); + + if (!socket) return js_mkundef(); + socket->read_drain_scheduled = false; + + while ( + !socket->destroyed && + socket->read_head && + eventemitter_listener_count(js, obj, "data") > 0 + ) { + tls_read_chunk_t *chunk = socket->read_head; + size_t len = chunk->len - chunk->off; + ant_value_t data = vtype(socket->encoding) == T_STR + ? js_mkstr(js, chunk->data + chunk->off, len) + : tls_make_buffer_chunk(js, chunk->data + chunk->off, len); + + socket->read_head = chunk->next; + if (socket->read_tail == chunk) socket->read_tail = NULL; + socket->read_len -= len; + free(chunk->data); + free(chunk); + + if (is_err(data)) { + socket->had_error = true; + tls_emit(js, obj, "error", &data, 1); + return data; + } + + tls_emit(js, obj, "data", &data, 1); + } + + return js_mkundef(); +} + +static void tls_socket_schedule_read_drain(ant_tls_socket_t *socket) { + if (!socket || socket->read_drain_scheduled || socket->destroyed) return; + socket->read_drain_scheduled = true; + queue_microtask(socket->js, js_heavy_mkfun(socket->js, tls_socket_drain_read_queue, socket->obj)); +} + static bool tls_value_bytes( ant_t *js, ant_value_t value, @@ -469,22 +513,23 @@ static void tls_socket_on_read(uv_stream_t *stream, ssize_t nread, const uv_buf_ if (nread > 0) { ant_value_t chunk = 0; - if (!tls_socket_push_read(socket, buf->base, (size_t)nread, false)) { - ant_value_t err = js_mkerr_typed(js, JS_ERR_TYPE, "Out of memory"); - socket->had_error = true; - tls_emit(js, socket->obj, "error", &err, 1); - tls_socket_close(socket); - goto done; - } - socket->bytes_read += (uint64_t)nread; tls_socket_sync_state(socket); - tls_emit(js, socket->obj, "readable", NULL, 0); - - if (vtype(socket->encoding) == T_STR) - chunk = js_mkstr(js, buf->base, (size_t)nread); - else chunk = tls_make_buffer_chunk(js, buf->base, (size_t)nread); - if (!is_err(chunk)) tls_emit(js, socket->obj, "data", &chunk, 1); + if (eventemitter_listener_count(js, socket->obj, "data") > 0) { + if (vtype(socket->encoding) == T_STR) + chunk = js_mkstr(js, buf->base, (size_t)nread); + else chunk = tls_make_buffer_chunk(js, buf->base, (size_t)nread); + if (!is_err(chunk)) tls_emit(js, socket->obj, "data", &chunk, 1); + } else { + if (!tls_socket_push_read(socket, buf->base, (size_t)nread, false)) { + ant_value_t err = js_mkerr_typed(js, JS_ERR_TYPE, "Out of memory"); + socket->had_error = true; + tls_emit(js, socket->obj, "error", &err, 1); + tls_socket_close(socket); + goto done; + } + tls_emit(js, socket->obj, "readable", NULL, 0); + } } else if (nread == UV_EOF) { socket->ended = true; tls_emit(js, socket->obj, "end", NULL, 0); @@ -579,8 +624,11 @@ static ant_value_t js_tls_socket_unshift(ant_t *js, ant_value_t *args, int nargs if (!socket) return js->thrown_value; if (!tls_parse_write_args(js, args, nargs, &bytes, &len, NULL, &err)) return err; - if (len > 0 && !tls_socket_push_read(socket, (const char *)bytes, len, true)) - return js_mkerr_typed(js, JS_ERR_TYPE, "Out of memory"); + if (len > 0) { + if (!tls_socket_push_read(socket, (const char *)bytes, len, true)) + return js_mkerr_typed(js, JS_ERR_TYPE, "Out of memory"); + tls_socket_schedule_read_drain(socket); + } return js_getthis(js); } diff --git a/tests/test_net_create_connection.cjs b/tests/test_net_create_connection.cjs index 3a819ef..d2c0368 100644 --- a/tests/test_net_create_connection.cjs +++ b/tests/test_net_create_connection.cjs @@ -39,9 +39,14 @@ server.listen(0, '127.0.0.1', () => { const address = server.address(); const client = net.createConnection({ port: address.port, host: '127.0.0.1' }, () => { connected = true; + assert.strictEqual(client.pending, false); + assert.strictEqual(client.connecting, false); client.write('ping'); }); + assert.strictEqual(client.pending, true); + assert.strictEqual(client.connecting, true); + client.setEncoding('utf8'); client.on('error', (err) => { throw err; diff --git a/tests/test_stream_read0_refills_below_hwm.cjs b/tests/test_stream_read0_refills_below_hwm.cjs new file mode 100644 index 0000000..6d5c4aa --- /dev/null +++ b/tests/test_stream_read0_refills_below_hwm.cjs @@ -0,0 +1,44 @@ +const assert = require('node:assert'); +const { Readable } = require('node:stream'); + +let reads = 0; +const readable = new Readable({ + objectMode: true, + highWaterMark: 4, + read() { + reads++; + if (reads === 1) this.push('a'); + else if (reads === 2) this.push('b'); + }, +}); + +readable.read(0); +assert.strictEqual(reads, 1); +assert.strictEqual(readable._readableState.length, 1); + +readable.read(0); +assert.strictEqual(reads, 2); +assert.strictEqual(readable._readableState.length, 2); +assert.strictEqual(readable.read().toString(), 'a'); +assert.strictEqual(readable.read().toString(), 'b'); +assert.strictEqual(readable._readableState.length, 0); + +{ + let byteReads = 0; + const bytes = new Readable({ + highWaterMark: 4, + read() { + byteReads++; + if (byteReads === 1) this.push(Buffer.alloc(4)); + else if (byteReads === 2) this.push(Buffer.alloc(1)); + }, + }); + + bytes.read(0); + assert.strictEqual(byteReads, 1); + assert.strictEqual(bytes._readableState.length, 4); + bytes.read(0); + assert.strictEqual(byteReads, 1); +} + +console.log('stream-read0-refills-below-hwm:ok');