From be8e870a97375ebfc5c5a99eb6d88de11ac232c7 Mon Sep 17 00:00:00 2001 From: theMackabu Date: Sat, 4 Apr 2026 16:59:24 -0700 Subject: [PATCH] revise fs stream --- include/modules/buffer.h | 1 + src/modules/buffer.c | 13 +++++++++- src/modules/stream.c | 52 +++++++++++++++++++++++++++++++++------ tests/test_fs_streams.cjs | 37 +++++++++++++++++++--------- 4 files changed, 83 insertions(+), 20 deletions(-) diff --git a/include/modules/buffer.h b/include/modules/buffer.h index 2bb7e2f..2e51207 100644 --- a/include/modules/buffer.h +++ b/include/modules/buffer.h @@ -72,6 +72,7 @@ ant_value_t create_dataview_with_buffer( size_t buffer_get_external_memory(void); bool buffer_is_dataview(ant_value_t obj); +bool buffer_is_binary_source(ant_value_t value); bool buffer_source_get_bytes(ant_t *js, ant_value_t value, const uint8_t **out, size_t *len); #endif diff --git a/src/modules/buffer.c b/src/modules/buffer.c index c05fd8a..8af5920 100644 --- a/src/modules/buffer.c +++ b/src/modules/buffer.c @@ -39,10 +39,21 @@ bool buffer_is_dataview(ant_value_t obj) { return js_check_brand(obj, BRAND_DATAVIEW); } +bool buffer_is_binary_source(ant_value_t value) { + ant_value_t slot = 0; + + if (vtype(value) == T_TYPEDARRAY) return true; + if (!is_object_type(value)) return false; + if (buffer_is_dataview(value)) return true; + + slot = js_get_slot(value, SLOT_BUFFER); + return vtype(slot) == T_TYPEDARRAY || vtype(slot) == T_NUM; +} + bool buffer_source_get_bytes(ant_t *js, ant_value_t value, const uint8_t **out, size_t *len) { if (out) *out = NULL; if (len) *len = 0; - if (vtype(value) != T_TYPEDARRAY && !is_object_type(value)) return false; + if (!buffer_is_binary_source(value)) return false; ant_value_t slot = js_get_slot(value, SLOT_BUFFER); TypedArrayData *ta = (TypedArrayData *)js_gettypedarray(slot); diff --git a/src/modules/stream.c b/src/modules/stream.c index 07b5468..9af3d55 100644 --- a/src/modules/stream.c +++ b/src/modules/stream.c @@ -6,10 +6,10 @@ #include "internal.h" #include "silver/engine.h" #include "esm/loader.h" - #include "gc/roots.h" #include "modules/assert.h" +#include "modules/buffer.h" #include "modules/events.h" #include "modules/stream.h" #include "modules/symbol.h" @@ -34,6 +34,11 @@ static ant_value_t g_transform_ctor = 0; static ant_value_t g_passthrough_proto = 0; static ant_value_t g_passthrough_ctor = 0; +static ant_value_t stream_readable_maybe_read(ant_t *js, ant_value_t stream_obj); +static ant_value_t stream_readable_flush(ant_t *js, ant_value_t stream_obj); +static ant_value_t stream_readable_push_value(ant_t *js, ant_value_t stream_obj, ant_value_t chunk, ant_value_t encoding); +static ant_value_t stream_readable_continue_flowing(ant_t *js, ant_value_t *args, int nargs); + static ant_value_t stream_noop(ant_t *js, ant_value_t *args, int nargs) { return js_mkundef(); } @@ -134,13 +139,16 @@ static ant_value_t stream_normalize_chunk( ) { ant_value_t str_val = 0; - if (object_mode || is_null(chunk) || is_undefined(chunk) || vtype(chunk) == T_TYPEDARRAY) - return chunk; + if ( + object_mode || is_null(chunk) || is_undefined(chunk) || + vtype(chunk) == T_TYPEDARRAY || buffer_is_binary_source(chunk) + ) return chunk; if (vtype(chunk) == T_STR) return stream_make_buffer(js, chunk, encoding); str_val = js_tostring_val(js, chunk); if (is_err(str_val)) return str_val; + return stream_make_buffer(js, str_val, encoding); } @@ -285,6 +293,7 @@ static void stream_init_readable(ant_t *js, ant_value_t obj, ant_value_t raw_opt js_set(js, state, "ended", js_false); js_set(js, state, "endEmitted", js_false); js_set(js, state, "flowing", js_false); + 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, "buffer", js_mkarr(js)); @@ -341,9 +350,19 @@ static void stream_emit_error(ant_t *js, ant_value_t stream_obj, ant_value_t err eventemitter_emit_args(js, stream_obj, "error", args, 1); } -static ant_value_t stream_readable_maybe_read(ant_t *js, ant_value_t stream_obj); -static ant_value_t stream_readable_flush(ant_t *js, ant_value_t stream_obj); -static ant_value_t stream_readable_push_value(ant_t *js, ant_value_t stream_obj, ant_value_t chunk, ant_value_t encoding); +static void stream_readable_schedule_continue_flowing(ant_t *js, ant_value_t stream_obj) { + ant_value_t state = stream_readable_state(js, stream_obj); + + if (!is_object_type(state)) return; + if (!js_truthy(js, js_get(js, state, "flowing"))) return; + if (js_truthy(js, js_get(js, stream_obj, "destroyed"))) return; + if (js_truthy(js, js_get(js, state, "ended"))) return; + if (stream_readable_buffer_len(js, stream_obj) > 0) return; + if (js_truthy(js, js_get(js, state, "flowingReadScheduled"))) return; + + js_set(js, state, "flowingReadScheduled", js_true); + stream_schedule_microtask(js, stream_readable_continue_flowing, stream_obj); +} static ant_value_t js_stream_pause(ant_t *js, ant_value_t *args, int nargs) { ant_value_t stream_obj = stream_require_this(js, js_getthis(js), "stream"); @@ -598,13 +617,32 @@ static ant_value_t stream_readable_start_flowing(ant_t *js, ant_value_t *args, i return js_mkundef(); } +static ant_value_t stream_readable_continue_flowing(ant_t *js, ant_value_t *args, int nargs) { + ant_value_t stream_obj = js_get_slot(js_getcurrentfunc(js), SLOT_DATA); + ant_value_t state = stream_readable_state(js, stream_obj); + + if (!is_object_type(state)) return js_mkundef(); + js_set(js, state, "flowingReadScheduled", js_false); + + if (!js_truthy(js, js_get(js, state, "flowing"))) return js_mkundef(); + if (js_truthy(js, js_get(js, stream_obj, "destroyed"))) return js_mkundef(); + if (js_truthy(js, js_get(js, state, "ended"))) return js_mkundef(); + + stream_readable_maybe_read(js, stream_obj); + stream_readable_flush(js, stream_obj); + + return js_mkundef(); +} + static ant_value_t stream_readable_flush(ant_t *js, ant_value_t stream_obj) { ant_value_t state = stream_readable_state(js, stream_obj); + bool emitted_data = false; if (!is_object_type(state)) return js_mkundef(); while (js_truthy(js, js_get(js, state, "flowing")) && stream_readable_buffer_len(js, stream_obj) > 0) { ant_value_t chunk = stream_buffer_shift(js, stream_obj); + emitted_data = true; eventemitter_emit_args(js, stream_obj, "data", &chunk, 1); } @@ -617,7 +655,7 @@ static ant_value_t stream_readable_flush(ant_t *js, ant_value_t stream_obj) { js_set(js, stream_obj, "readableEnded", js_true); stream_emit_named(js, stream_obj, "end"); stream_emit_named(js, stream_obj, "close"); - } + } else if (emitted_data) stream_readable_schedule_continue_flowing(js, stream_obj); return js_mkundef(); } diff --git a/tests/test_fs_streams.cjs b/tests/test_fs_streams.cjs index 531de39..a0ddfec 100644 --- a/tests/test_fs_streams.cjs +++ b/tests/test_fs_streams.cjs @@ -1,7 +1,20 @@ const fs = require('node:fs'); -const sourcePath = 'tests/.fs_stream_source.txt'; -const copyPath = 'tests/.fs_stream_copy.txt'; -const content = 'hello from fs streams\nline two'; +const { Buffer } = require('node:buffer'); +const sourcePath = '/tmp/ant_fs_stream_seq_source.bin'; +const copyPath = '/tmp/ant_fs_stream_seq_copy.bin'; +const content = Buffer.alloc(70000); + +for (let i = 0; i < content.length; i++) { + content[i] = i & 255; +} + +function sameBuffer(left, right) { + if (!left || !right || left.length !== right.length) return false; + for (let i = 0; i < left.length; i++) { + if (left[i] !== right[i]) return false; + } + return true; +} try { fs.unlinkSync(sourcePath); @@ -19,27 +32,27 @@ const writer = fs.createWriteStream(sourcePath); if (!(writer instanceof fs.WriteStream)) throw new Error('createWriteStream() did not return fs.WriteStream'); writer.on('error', fail); writer.on('finish', () => { - const written = fs.readFileSync(sourcePath, 'utf8'); - if (written !== content) fail(new Error(`unexpected write content: ${written}`)); + const written = fs.readFileSync(sourcePath); + if (!sameBuffer(written, content)) fail(new Error(`unexpected write content length: ${written.length}`)); const reader = fs.createReadStream(sourcePath); if (!(reader instanceof fs.ReadStream)) fail(new Error('createReadStream() did not return fs.ReadStream')); - let readBack = ''; + const readBack = []; reader.on('error', fail); reader.on('data', chunk => { - readBack += chunk.toString(); + readBack.push(chunk); }); reader.on('end', () => { - if (readBack !== content) fail(new Error(`unexpected read content: ${readBack}`)); + if (!sameBuffer(Buffer.concat(readBack), content)) fail(new Error('unexpected read content')); const pipedReader = fs.createReadStream(sourcePath); const pipedWriter = fs.createWriteStream(copyPath); pipedReader.on('error', fail); pipedWriter.on('error', fail); pipedWriter.on('finish', () => { - const copied = fs.readFileSync(copyPath, 'utf8'); - if (copied !== content) fail(new Error(`unexpected piped content: ${copied}`)); + const copied = fs.readFileSync(copyPath); + if (!sameBuffer(copied, content)) fail(new Error(`unexpected piped content length: ${copied.length}`)); fs.unlinkSync(sourcePath); fs.unlinkSync(copyPath); @@ -49,5 +62,5 @@ writer.on('finish', () => { }); }); -writer.write('hello '); -writer.end('from fs streams\nline two'); +writer.write(content.subarray(0, 32768)); +writer.end(content.subarray(32768)); -- 2.51.2