From 18f39e2db2560935dc476fb75221015a11fba170 Mon Sep 17 00:00:00 2001 From: theMackabu Date: Wed, 25 Mar 2026 19:17:57 -0700 Subject: [PATCH] improve blob streaming --- examples/demo/event_loop.js | 1 + include/streams/readable.h | 12 ++++++----- src/modules/blob.c | 22 ++++++------------- src/streams/readable.c | 43 +++++++++++++++++++++++++++++++++++++ 4 files changed, 58 insertions(+), 20 deletions(-) diff --git a/examples/demo/event_loop.js b/examples/demo/event_loop.js index 4da219e..c7e426d 100644 --- a/examples/demo/event_loop.js +++ b/examples/demo/event_loop.js @@ -19,4 +19,5 @@ process.once('beforeExit', () => { else formatted = String(Math.round(rate)); console.log(`\x1b[1;36m${formatted} event loop iterations/sec\x1b[0m \x1b[2m(${count.toLocaleString()} in ${(elapsed / 1000).toFixed(1)}s)\x1b[0m`); + if (typeof globalThis.Ant !== 'undefined') console.log('ant version:', Ant.version); }); diff --git a/include/streams/readable.h b/include/streams/readable.h index 49d0f7e..7b0a8e9 100644 --- a/include/streams/readable.h +++ b/include/streams/readable.h @@ -53,17 +53,19 @@ ant_value_t js_rs_reader_ctor(ant_t *js, ant_value_t *args, int nargs); ant_value_t rs_cancel_resolve(ant_t *js, ant_value_t *args, int nargs); ant_value_t readable_stream_cancel(ant_t *js, ant_value_t stream_obj, ant_value_t reason); ant_value_t rs_create_stream(ant_t *js, ant_value_t pull_fn, ant_value_t cancel_fn, double hwm); -ant_value_t rs_create_stream(ant_t *js, ant_value_t pull_fn, ant_value_t cancel_fn, double hwm); - -bool rs_reader_has_reqs(ant_t *js, ant_value_t reader_obj); -bool rs_default_controller_can_close_or_enqueue(rs_controller_t *ctrl, rs_stream_t *stream); +void rs_controller_close(ant_t *js, ant_value_t ctrl_obj); void rs_default_controller_clear_algorithms(ant_value_t ctrl_obj); +void rs_ctrl_queue_push(ant_t *js, ant_value_t ctrl_obj, ant_value_t value); +void rs_controller_enqueue(ant_t *js, ant_value_t ctrl_obj, ant_value_t chunk); void rs_default_controller_call_pull_if_needed(ant_t *js, ant_value_t controller_obj); void rs_default_reader_error_read_requests(ant_t *js, ant_value_t reader_obj, ant_value_t e); void rs_fulfill_read_request(ant_t *js, ant_value_t stream_obj, ant_value_t chunk, bool done); -void rs_ctrl_queue_push(ant_t *js, ant_value_t ctrl_obj, ant_value_t value); + void readable_stream_close(ant_t *js, ant_value_t stream_obj); void readable_stream_error(ant_t *js, ant_value_t stream_obj, ant_value_t e); +bool rs_reader_has_reqs(ant_t *js, ant_value_t reader_obj); +bool rs_default_controller_can_close_or_enqueue(rs_controller_t *ctrl, rs_stream_t *stream); + #endif diff --git a/src/modules/blob.c b/src/modules/blob.c index 80017f1..c89746f 100644 --- a/src/modules/blob.c +++ b/src/modules/blob.c @@ -12,7 +12,6 @@ #include "internal.h" #include "descriptors.h" -#include "silver/engine.h" #include "modules/blob.h" #include "modules/buffer.h" #include "modules/symbol.h" @@ -284,29 +283,22 @@ static ant_value_t js_blob_slice(ant_t *js, ant_value_t *args, int nargs) { static ant_value_t blob_stream_pull(ant_t *js, ant_value_t *args, int nargs) { ant_value_t blob_obj = js_get_slot(js->current_func, SLOT_DATA); blob_data_t *bd = blob_get_data(blob_obj); - ant_value_t controller = (nargs > 0) ? args[0] : js_mkundef(); + ant_value_t ctrl = (nargs > 0) ? args[0] : js_mkundef(); if (bd && bd->size > 0 && bd->data) { - ArrayBufferData *ab = create_array_buffer_data(bd->size); - if (ab) { + ArrayBufferData *ab = create_array_buffer_data(bd->size); + if (ab) { memcpy(ab->data, bd->data, bd->size); - ant_value_t chunk = create_typed_array(js, TYPED_ARRAY_UINT8, ab, 0, bd->size, "Uint8Array"); - ant_value_t enqueue_fn = js_get(js, controller, "enqueue"); - if (is_callable(enqueue_fn)) { - ant_value_t enq_args[1] = { chunk }; - sv_vm_call(js->vm, js, enqueue_fn, controller, enq_args, 1, NULL, false); - }} - } - - ant_value_t close_fn = js_get(js, controller, "close"); - if (is_callable(close_fn)) sv_vm_call(js->vm, js, close_fn, controller, NULL, 0, NULL, false); + rs_controller_enqueue(js, ctrl, create_typed_array(js, TYPED_ARRAY_UINT8, ab, 0, bd->size, "Uint8Array")); + }} + rs_controller_close(js, ctrl); return js_mkundef(); } static ant_value_t js_blob_stream(ant_t *js, ant_value_t *args, int nargs) { ant_value_t pull_fn = js_heavy_mkfun(js, blob_stream_pull, js->this_val); - return rs_create_stream(js, pull_fn, js_mkundef(), 0); + return rs_create_stream(js, pull_fn, js_mkundef(), 1); } static ant_value_t js_blob_ctor(ant_t *js, ant_value_t *args, int nargs) { diff --git a/src/streams/readable.c b/src/streams/readable.c index 94beb31..d9d71f1 100644 --- a/src/streams/readable.c +++ b/src/streams/readable.c @@ -495,6 +495,49 @@ static ant_value_t js_rs_controller_error(ant_t *js, ant_value_t *args, int narg return js_mkundef(); } +void rs_controller_enqueue(ant_t *js, ant_value_t ctrl_obj, ant_value_t chunk) { + rs_controller_t *ctrl = rs_get_controller(ctrl_obj); + if (!ctrl) return; + + ant_value_t stream_obj = rs_ctrl_stream(ctrl_obj); + rs_stream_t *stream = rs_get_stream(stream_obj); + if (!rs_default_controller_can_close_or_enqueue(ctrl, stream)) return; + + ant_value_t reader_obj = rs_stream_reader(stream_obj); + if (is_object_type(reader_obj) && rs_reader_has_reqs(js, reader_obj)) { + rs_fulfill_read_request(js, stream_obj, chunk, false); + rs_default_controller_call_pull_if_needed(js, ctrl_obj); + return; + } + + rs_ctrl_queue_push(js, ctrl_obj, chunk); + if (ctrl->queue_sizes_len >= ctrl->queue_sizes_cap) { + uint32_t new_cap = ctrl->queue_sizes_cap ? ctrl->queue_sizes_cap * 2 : 4; + double *new_sizes = realloc(ctrl->queue_sizes, new_cap * sizeof(double)); + if (new_sizes) { ctrl->queue_sizes = new_sizes; ctrl->queue_sizes_cap = new_cap; } + } + + if (ctrl->queue_sizes_len < ctrl->queue_sizes_cap) + ctrl->queue_sizes[ctrl->queue_sizes_len++] = 1; + ctrl->queue_total_size += 1; + rs_default_controller_call_pull_if_needed(js, ctrl_obj); +} + +void rs_controller_close(ant_t *js, ant_value_t ctrl_obj) { + rs_controller_t *ctrl = rs_get_controller(ctrl_obj); + if (!ctrl) return; + + ant_value_t stream_obj = rs_ctrl_stream(ctrl_obj); + rs_stream_t *stream = rs_get_stream(stream_obj); + if (!rs_default_controller_can_close_or_enqueue(ctrl, stream)) return; + ctrl->close_requested = true; + + if (rs_ctrl_queue_len(js, ctrl_obj) == 0) { + rs_default_controller_clear_algorithms(ctrl_obj); + readable_stream_close(js, stream_obj); + } +} + static ant_value_t js_rs_reader_get_closed(ant_t *js, ant_value_t *args, int nargs) { return rs_reader_closed(js->this_val); } -- 2.51.2