diff --git a/examples/spec/blob.js b/examples/spec/blob.js index c94939d..22c5200 100644 --- a/examples/spec/blob.js +++ b/examples/spec/blob.js @@ -105,6 +105,32 @@ test('Blob not instanceof File', b1 instanceof File, false); const ft = await f1.text(); test('File.text() works', ft, 'content'); +console.log('\nBlob.stream()\n'); + +{ + const blob = new Blob(['hello']); + const stream = blob.stream(); + test('stream() returns ReadableStream', stream instanceof ReadableStream, true); + test('stream is not locked', stream.locked, false); + + const reader = stream.getReader(); + const { value, done } = await reader.read(); + test('stream chunk is Uint8Array', value instanceof Uint8Array, true); + test('stream chunk length', value.length, 5); + test('stream chunk byte 0', value[0], 104); + test('stream chunk byte 4', value[4], 111); + + const end = await reader.read(); + test('stream done after one chunk', end.done, true); +} + +{ + const blob = new Blob(); + const reader = blob.stream().getReader(); + const { done } = await reader.read(); + test('empty blob stream closes immediately', done, true); +} + console.log('\nSymbol.toStringTag\n'); test('Blob toStringTag', Object.prototype.toString.call(b1), '[object Blob]'); diff --git a/include/gc/modules.h b/include/gc/modules.h index f5d53ac..966a750 100644 --- a/include/gc/modules.h +++ b/include/gc/modules.h @@ -23,5 +23,6 @@ void gc_mark_abort(ant_t *js, gc_mark_fn mark); void gc_mark_domexception(ant_t *js, gc_mark_fn mark); void gc_mark_readable_streams(ant_t *js, gc_mark_fn mark); void gc_mark_writable_streams(ant_t *js, gc_mark_fn mark); +void gc_mark_pipes(ant_t *js, gc_mark_fn mark); #endif diff --git a/include/modules/assert.h b/include/modules/assert.h index 66c235e..f782884 100644 --- a/include/modules/assert.h +++ b/include/modules/assert.h @@ -1,8 +1,27 @@ #ifndef ANT_ASSERT_MODULE_H #define ANT_ASSERT_MODULE_H -#include "types.h" +#include "internal.h" +#include "silver/engine.h" ant_value_t assert_library(ant_t *js); +static inline bool promise_was_rejected(ant_value_t result) { + if (vtype(result) != T_PROMISE) return false; + ant_object_t *obj = js_obj_ptr(js_as_obj(result)); + return obj && obj->promise_state && obj->promise_state->state == 2; +} + +static inline void promise_mark_handled(ant_value_t v) { + if (vtype(v) != T_PROMISE) return; + ant_object_t *obj = js_obj_ptr(js_as_obj(v)); + if (obj && obj->promise_state) obj->promise_state->has_rejection_handler = true; +} + +static inline bool promise_was_fulfilled(ant_value_t result) { + if (vtype(result) != T_PROMISE) return false; + ant_object_t *obj = js_obj_ptr(js_as_obj(result)); + return obj && obj->promise_state && obj->promise_state->state == 1; +} + #endif diff --git a/include/streams/pipes.h b/include/streams/pipes.h new file mode 100644 index 0000000..b6c6849 --- /dev/null +++ b/include/streams/pipes.h @@ -0,0 +1,16 @@ +#ifndef STREAMS_PIPES_H +#define STREAMS_PIPES_H + +#include "types.h" +#include + +void init_pipes_proto(ant_t *js, ant_value_t rs_proto); +void gc_mark_pipes(ant_t *js, void (*mark)(ant_t *, ant_value_t)); + +ant_value_t readable_stream_pipe_to( + ant_t *js, ant_value_t source, ant_value_t dest, + bool prevent_close, bool prevent_abort, bool prevent_cancel, + ant_value_t signal +); + +#endif diff --git a/include/streams/readable.h b/include/streams/readable.h index e648740..49d0f7e 100644 --- a/include/streams/readable.h +++ b/include/streams/readable.h @@ -28,15 +28,42 @@ typedef struct { bool disturbed; } rs_stream_t; +extern ant_value_t g_rs_proto; +extern ant_value_t g_reader_proto; +extern ant_value_t g_controller_proto; + void init_readable_stream_module(void); void gc_mark_readable_streams(ant_t *js, void (*mark)(ant_t *, ant_value_t)); rs_stream_t *rs_get_stream(ant_value_t obj); -ant_value_t rs_stream_controller(ant_t *js, ant_value_t stream_obj); +rs_controller_t *rs_get_controller(ant_value_t obj); +ant_offset_t rs_ctrl_queue_len(ant_t *js, ant_value_t ctrl_obj); + +ant_value_t rs_ctrl_size(ant_value_t ctrl_obj); +ant_value_t rs_reader_reqs(ant_value_t reader_obj); +ant_value_t rs_stream_error(ant_value_t stream_obj); ant_value_t rs_stream_reader(ant_value_t stream_obj); +ant_value_t rs_reader_stream(ant_value_t reader_obj); +ant_value_t rs_reader_closed(ant_value_t reader_obj); + +ant_value_t rs_stream_controller(ant_t *js, ant_value_t stream_obj); +ant_value_t rs_default_reader_read(ant_t *js, ant_value_t reader_obj); +ant_value_t rs_cancel_reject(ant_t *js, ant_value_t *args, int nargs); +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_default_controller_clear_algorithms(ant_value_t ctrl_obj); +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); -ant_value_t readable_stream_cancel(ant_t *js, ant_value_t stream_obj, ant_value_t reason); #endif diff --git a/include/streams/writable.h b/include/streams/writable.h index 0750df2..af52404 100644 --- a/include/streams/writable.h +++ b/include/streams/writable.h @@ -33,9 +33,14 @@ void init_writable_stream_module(void); void gc_mark_writable_streams(ant_t *js, void (*mark)(ant_t *, ant_value_t)); ws_stream_t *ws_get_stream(ant_value_t obj); +ws_controller_t *ws_get_controller(ant_value_t obj); + ant_value_t ws_stream_controller(ant_value_t stream_obj); ant_value_t ws_stream_writer(ant_value_t stream_obj); +ant_value_t ws_acquire_writer(ant_t *js, ant_value_t stream_obj); +ant_value_t ws_writer_write(ant_t *js, ant_value_t writer_obj, ant_value_t chunk); + ant_value_t writable_stream_close(ant_t *js, ant_value_t stream_obj); ant_value_t writable_stream_abort(ant_t *js, ant_value_t stream_obj, ant_value_t reason); diff --git a/src/gc/objects.c b/src/gc/objects.c index 07efe3a..0ac4798 100644 --- a/src/gc/objects.c +++ b/src/gc/objects.c @@ -456,6 +456,7 @@ static void gc_mark_roots(ant_t *js) { gc_mark_domexception(js, gc_mark_value); gc_mark_readable_streams(js, gc_mark_value); gc_mark_writable_streams(js, gc_mark_value); + gc_mark_pipes(js, gc_mark_value); for ( ant_object_t *obj = g_pending_promises; diff --git a/src/modules/assert.c b/src/modules/assert.c index 1d99c47..8dcfae6 100644 --- a/src/modules/assert.c +++ b/src/modules/assert.c @@ -4,6 +4,8 @@ #include "ant.h" #include "errors.h" #include "internal.h" + +#include "modules/assert.h" #include "silver/engine.h" static ant_value_t assertion_error(ant_t *js, const char *msg, ant_value_t msg_val) { @@ -190,24 +192,6 @@ static ant_value_t assert_does_not_throw(ant_t *js, ant_value_t *args, int nargs return js_mkundef(); } -static bool promise_was_rejected(ant_value_t result) { - if (vtype(result) != T_PROMISE) return false; - ant_object_t *obj = js_obj_ptr(js_as_obj(result)); - return obj && obj->promise_state && obj->promise_state->state == 2; -} - -static void promise_mark_handled(ant_value_t result) { - if (vtype(result) != T_PROMISE) return; - ant_object_t *obj = js_obj_ptr(js_as_obj(result)); - if (obj && obj->promise_state) obj->promise_state->has_rejection_handler = true; -} - -static bool promise_was_fulfilled(ant_value_t result) { - if (vtype(result) != T_PROMISE) return false; - ant_object_t *obj = js_obj_ptr(js_as_obj(result)); - return obj && obj->promise_state && obj->promise_state->state == 1; -} - static ant_value_t assert_rejects(ant_t *js, ant_value_t *args, int nargs) { if (nargs < 1) return js_mkerr(js, "assert.rejects: first argument required"); ant_value_t promise = js_mkpromise(js); diff --git a/src/modules/blob.c b/src/modules/blob.c index 6553d46..80017f1 100644 --- a/src/modules/blob.c +++ b/src/modules/blob.c @@ -12,9 +12,11 @@ #include "internal.h" #include "descriptors.h" +#include "silver/engine.h" #include "modules/blob.h" #include "modules/buffer.h" #include "modules/symbol.h" +#include "streams/readable.h" ant_value_t g_blob_proto = 0; ant_value_t g_file_proto = 0; @@ -279,9 +281,32 @@ static ant_value_t js_blob_slice(ant_t *js, ant_value_t *args, int nargs) { return result; } +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(); + + if (bd && bd->size > 0 && bd->data) { + 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); + + return js_mkundef(); +} + static ant_value_t js_blob_stream(ant_t *js, ant_value_t *args, int nargs) { - (void)args; (void)nargs; - return js_mkerr_typed(js, JS_ERR_TYPE, "Blob.stream() is not yet implemented"); + 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); } static ant_value_t js_blob_ctor(ant_t *js, ant_value_t *args, int nargs) { diff --git a/src/modules/timer.c b/src/modules/timer.c index 246fa9f..af30ba2 100644 --- a/src/modules/timer.c +++ b/src/modules/timer.c @@ -347,23 +347,22 @@ void process_microtasks(ant_t *js) { } void process_immediates(ant_t *js) { - while (timer_state.immediates != NULL) { - immediate_entry_t *entry = timer_state.immediates; - timer_state.immediates = entry->next; - - if (timer_state.immediates == NULL) { - timer_state.immediates_tail = NULL; - } - - if (entry->active) { - ant_value_t args[0]; - sv_vm_call(js->vm, js, entry->callback, js_mkundef(), args, 0, NULL, false); - process_microtasks(js); - } - - free(entry); +while (timer_state.immediates != NULL) { + immediate_entry_t *entry = timer_state.immediates; + timer_state.immediates = entry->next; + + if (timer_state.immediates == NULL) { + timer_state.immediates_tail = NULL; } -} + + if (entry->active) { + ant_value_t args[0]; + sv_vm_call(js->vm, js, entry->callback, js_mkundef(), args, 0, NULL, false); + process_microtasks(js); + } + + free(entry); +}} int has_pending_immediates(void) { for ( diff --git a/src/streams/readable.c b/src/streams/readable.c index 86eff99..94beb31 100644 --- a/src/streams/readable.c +++ b/src/streams/readable.c @@ -9,11 +9,13 @@ #include "silver/engine.h" #include "modules/symbol.h" +#include "modules/assert.h" #include "streams/readable.h" +#include "streams/pipes.h" -static ant_value_t g_rs_proto; -static ant_value_t g_reader_proto; -static ant_value_t g_controller_proto; +ant_value_t g_rs_proto; +ant_value_t g_reader_proto; +ant_value_t g_controller_proto; rs_stream_t *rs_get_stream(ant_value_t obj) { ant_value_t s = js_get_slot(obj, SLOT_DATA); @@ -21,7 +23,7 @@ rs_stream_t *rs_get_stream(ant_value_t obj) { return (rs_stream_t *)(uintptr_t)(size_t)js_getnum(s); } -static rs_controller_t *rs_get_controller(ant_value_t obj) { +rs_controller_t *rs_get_controller(ant_value_t obj) { ant_value_t s = js_get_slot(obj, SLOT_DATA); if (vtype(s) != T_NUM) return NULL; return (rs_controller_t *)(uintptr_t)(size_t)js_getnum(s); @@ -57,7 +59,7 @@ ant_value_t rs_stream_reader(ant_value_t stream_obj) { return js_get_slot(stream_obj, SLOT_CTOR); } -static inline ant_value_t rs_stream_error(ant_value_t stream_obj) { +ant_value_t rs_stream_error(ant_value_t stream_obj) { return js_get_slot(stream_obj, SLOT_BUFFER); } @@ -73,7 +75,7 @@ static inline ant_value_t rs_ctrl_cancel(ant_value_t ctrl_obj) { return js_get_slot(ctrl_obj, SLOT_RS_CANCEL); } -static inline ant_value_t rs_ctrl_size(ant_value_t ctrl_obj) { +ant_value_t rs_ctrl_size(ant_value_t ctrl_obj) { return js_get_slot(ctrl_obj, SLOT_RS_SIZE); } @@ -81,24 +83,24 @@ static inline ant_value_t rs_ctrl_queue(ant_t *js, ant_value_t ctrl_obj) { return js_get_slot(ctrl_obj, SLOT_BUFFER); } -static inline ant_value_t rs_reader_stream(ant_value_t reader_obj) { +ant_value_t rs_reader_stream(ant_value_t reader_obj) { return js_get_slot(reader_obj, SLOT_ENTRIES); } -static inline ant_value_t rs_reader_closed(ant_value_t reader_obj) { +ant_value_t rs_reader_closed(ant_value_t reader_obj) { return js_get_slot(reader_obj, SLOT_RS_CLOSED); } -static inline ant_value_t rs_reader_reqs(ant_value_t reader_obj) { +ant_value_t rs_reader_reqs(ant_value_t reader_obj) { return js_get_slot(reader_obj, SLOT_BUFFER); } -static bool rs_reader_has_reqs(ant_t *js, ant_value_t reader_obj) { +bool rs_reader_has_reqs(ant_t *js, ant_value_t reader_obj) { ant_value_t arr = rs_reader_reqs(reader_obj); return vtype(arr) == T_ARR && js_arr_len(js, arr) > 0; } -static void rs_ctrl_queue_push(ant_t *js, ant_value_t ctrl_obj, ant_value_t value) { +void rs_ctrl_queue_push(ant_t *js, ant_value_t ctrl_obj, ant_value_t value) { ant_value_t arr = rs_ctrl_queue(js, ctrl_obj); if (vtype(arr) == T_ARR) js_arr_push(js, arr, value); } @@ -116,7 +118,7 @@ static ant_value_t rs_ctrl_queue_shift(ant_t *js, ant_value_t ctrl_obj) { return val; } -static ant_offset_t rs_ctrl_queue_len(ant_t *js, ant_value_t ctrl_obj) { +ant_offset_t rs_ctrl_queue_len(ant_t *js, ant_value_t ctrl_obj) { ant_value_t arr = rs_ctrl_queue(js, ctrl_obj); if (vtype(arr) != T_ARR) return 0; return js_arr_len(js, arr); @@ -140,10 +142,10 @@ static ant_value_t rs_reader_reqs_shift(ant_t *js, ant_value_t reader_obj) { return val; } -static void rs_default_controller_call_pull_if_needed(ant_t *js, ant_value_t controller_obj); -static bool rs_default_controller_can_close_or_enqueue(rs_controller_t *ctrl, rs_stream_t *stream); +void rs_default_controller_call_pull_if_needed(ant_t *js, ant_value_t controller_obj); +bool rs_default_controller_can_close_or_enqueue(rs_controller_t *ctrl, rs_stream_t *stream); -static void rs_default_controller_clear_algorithms(ant_value_t ctrl_obj) { +void rs_default_controller_clear_algorithms(ant_value_t ctrl_obj) { js_set_slot(ctrl_obj, SLOT_RS_PULL, js_mkundef()); js_set_slot(ctrl_obj, SLOT_RS_CANCEL, js_mkundef()); js_set_slot(ctrl_obj, SLOT_RS_SIZE, js_mkundef()); @@ -168,14 +170,14 @@ static bool rs_default_controller_should_call_pull(ant_t *js, rs_controller_t *c return desired > 0; } -static bool rs_default_controller_can_close_or_enqueue(rs_controller_t *ctrl, rs_stream_t *stream) { +bool rs_default_controller_can_close_or_enqueue(rs_controller_t *ctrl, rs_stream_t *stream) { if (!ctrl || !stream) return false; if (ctrl->close_requested) return false; if (stream->state != RS_STATE_READABLE) return false; return true; } -static void rs_fulfill_read_request(ant_t *js, ant_value_t stream_obj, ant_value_t chunk, bool done) { +void rs_fulfill_read_request(ant_t *js, ant_value_t stream_obj, ant_value_t chunk, bool done) { ant_value_t reader_obj = rs_stream_reader(stream_obj); if (!is_object_type(reader_obj)) return; ant_value_t promise = rs_reader_reqs_shift(js, reader_obj); @@ -184,7 +186,7 @@ static void rs_fulfill_read_request(ant_t *js, ant_value_t stream_obj, ant_value js_resolve_promise(js, promise, result); } -static void rs_default_reader_error_read_requests(ant_t *js, ant_value_t reader_obj, ant_value_t e) { +void rs_default_reader_error_read_requests(ant_t *js, ant_value_t reader_obj, ant_value_t e) { ant_value_t arr = rs_reader_reqs(reader_obj); if (vtype(arr) != T_ARR) return; ant_offset_t len = js_arr_len(js, arr); @@ -228,13 +230,13 @@ void readable_stream_error(ant_t *js, ant_value_t stream_obj, ant_value_t e) { } } -static ant_value_t rs_cancel_resolve(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 p = js_get_slot(js->current_func, SLOT_DATA); js_resolve_promise(js, p, js_mkundef()); return js_mkundef(); } -static ant_value_t rs_cancel_reject(ant_t *js, ant_value_t *args, int nargs) { +ant_value_t rs_cancel_reject(ant_t *js, ant_value_t *args, int nargs) { ant_value_t p = js_get_slot(js->current_func, SLOT_DATA); js_reject_promise(js, p, nargs > 0 ? args[0] : js_mkundef()); return js_mkundef(); @@ -308,15 +310,15 @@ static ant_value_t rs_pull_reject_handler(ant_t *js, ant_value_t *args, int narg return js_mkundef(); } -static void rs_default_controller_call_pull_if_needed(ant_t *js, ant_value_t controller_obj) { +void rs_default_controller_call_pull_if_needed(ant_t *js, ant_value_t controller_obj) { rs_controller_t *ctrl = rs_get_controller(controller_obj); if (!ctrl) return; + ant_value_t stream_obj = rs_ctrl_stream(controller_obj); rs_stream_t *stream = rs_get_stream(stream_obj); if (!stream) return; if (!rs_default_controller_should_call_pull(js, ctrl, stream, controller_obj)) return; - if (ctrl->pulling) { ctrl->pull_again = true; return; } ctrl->pulling = true; @@ -352,7 +354,7 @@ static void rs_default_controller_call_pull_if_needed(ant_t *js, ant_value_t con } else ctrl->pulling = false; } -static ant_value_t rs_default_reader_read(ant_t *js, ant_value_t reader_obj) { +ant_value_t rs_default_reader_read(ant_t *js, ant_value_t reader_obj) { ant_value_t stream_obj = rs_reader_stream(reader_obj); rs_stream_t *stream = rs_get_stream(stream_obj); if (!stream) return js_mkerr_typed(js, JS_ERR_TYPE, "Reader has no stream"); @@ -519,10 +521,14 @@ static ant_value_t js_rs_reader_release_lock(ant_t *js, ant_value_t *args, int n js_mkerr_typed(js, JS_ERR_TYPE, "Reader was released"); ant_value_t release_err = js->thrown_value; - if (stream->state == RS_STATE_READABLE) - js_reject_promise(js, rs_reader_closed(js->this_val), release_err); + if (stream->state == RS_STATE_READABLE) { + ant_value_t old_closed = rs_reader_closed(js->this_val); + js_reject_promise(js, old_closed, release_err); + promise_mark_handled(old_closed); + } js_reject_promise(js, new_closed, release_err); + promise_mark_handled(new_closed); js_set_slot(js->this_val, SLOT_RS_CLOSED, new_closed); js_set_slot(stream_obj, SLOT_CTOR, js_mkundef()); @@ -542,7 +548,7 @@ static ant_value_t js_rs_reader_cancel(ant_t *js, ant_value_t *args, int nargs) return readable_stream_cancel(js, stream_obj, reason); } -static ant_value_t js_rs_reader_ctor(ant_t *js, ant_value_t *args, int nargs) { +ant_value_t js_rs_reader_ctor(ant_t *js, ant_value_t *args, int nargs) { if (vtype(js->new_target) == T_UNDEF) return js_mkerr_typed(js, JS_ERR_TYPE, "ReadableStreamDefaultReader constructor requires 'new'"); if (nargs < 1) return js_mkerr_typed(js, JS_ERR_TYPE, "ReadableStreamDefaultReader requires a stream argument"); @@ -619,18 +625,6 @@ static ant_value_t js_rs_get_reader(ant_t *js, ant_value_t *args, int nargs) { return reader; } -static ant_value_t js_rs_tee(ant_t *js, ant_value_t *args, int nargs) { - return js_mkerr_typed(js, JS_ERR_TYPE, "ReadableStream.prototype.tee is not yet implemented"); -} - -static ant_value_t js_rs_pipe_through(ant_t *js, ant_value_t *args, int nargs) { - return js_mkerr_typed(js, JS_ERR_TYPE, "ReadableStream.prototype.pipeThrough is not yet implemented"); -} - -static ant_value_t js_rs_pipe_to(ant_t *js, ant_value_t *args, int nargs) { - return js_mkerr_typed(js, JS_ERR_TYPE, "ReadableStream.prototype.pipeTo is not yet implemented"); -} - static ant_value_t js_rs_values(ant_t *js, ant_value_t *args, int nargs) { return js_mkerr_typed(js, JS_ERR_TYPE, "ReadableStream async iteration is not yet implemented"); } @@ -809,6 +803,33 @@ static ant_value_t js_rs_ctor(ant_t *js, ant_value_t *args, int nargs) { return obj; } +ant_value_t rs_create_stream(ant_t *js, ant_value_t pull_fn, ant_value_t cancel_fn, double hwm) { + rs_stream_t *st = calloc(1, sizeof(rs_stream_t)); + if (!st) return js_mkerr(js, "out of memory"); + st->state = RS_STATE_READABLE; + + ant_value_t obj = js_mkobj(js); + js_set_proto_init(obj, g_rs_proto); + js_set_slot(obj, SLOT_DATA, ANT_PTR(st)); + js_set_finalizer(obj, rs_stream_finalize); + + ant_value_t ctrl_obj = setup_default_controller(js, obj, pull_fn, cancel_fn, js_mkundef(), hwm); + if (is_err(ctrl_obj)) return ctrl_obj; + + ant_value_t resolved = js_mkpromise(js); + js_resolve_promise(js, resolved, js_mkundef()); + ant_value_t res_fn = js_heavy_mkfun(js, rs_start_resolve_handler, ctrl_obj); + ant_value_t rej_fn = js_heavy_mkfun(js, rs_start_reject_handler, ctrl_obj); + ant_value_t then_fn = js_get(js, resolved, "then"); + + if (is_callable(then_fn)) { + ant_value_t then_args[2] = { res_fn, rej_fn }; + sv_vm_call(js->vm, js, then_fn, resolved, then_args, 2, NULL, false); + } + + return obj; +} + static ant_value_t js_rs_controller_ctor(ant_t *js, ant_value_t *args, int nargs) { return js_mkerr_typed(js, JS_ERR_TYPE, "ReadableStreamDefaultController cannot be constructed directly"); } @@ -824,56 +845,44 @@ void init_readable_stream_module(void) { ant_value_t g = js_glob(js); g_controller_proto = js_mkobj(js); - js_set_getter_desc(js, g_controller_proto, "desiredSize", 11, - js_mkfun(js_rs_controller_get_desired_size), JS_DESC_C); + js_set_getter_desc(js, g_controller_proto, "desiredSize", 11, js_mkfun(js_rs_controller_get_desired_size), JS_DESC_C); js_set(js, g_controller_proto, "close", js_mkfun(js_rs_controller_close)); js_set_descriptor(js, g_controller_proto, "close", 5, JS_DESC_W | JS_DESC_C); js_set(js, g_controller_proto, "enqueue", js_mkfun(js_rs_controller_enqueue)); js_set_descriptor(js, g_controller_proto, "enqueue", 7, JS_DESC_W | JS_DESC_C); js_set(js, g_controller_proto, "error", js_mkfun(js_rs_controller_error)); js_set_descriptor(js, g_controller_proto, "error", 5, JS_DESC_W | JS_DESC_C); - js_set_sym(js, g_controller_proto, get_toStringTag_sym(), - js_mkstr(js, "ReadableStreamDefaultController", 31)); + js_set_sym(js, g_controller_proto, get_toStringTag_sym(), js_mkstr(js, "ReadableStreamDefaultController", 31)); - ant_value_t ctrl_ctor = js_make_ctor(js, js_rs_controller_ctor, g_controller_proto, - "ReadableStreamDefaultController", 31); + ant_value_t ctrl_ctor = js_make_ctor(js, js_rs_controller_ctor, g_controller_proto, "ReadableStreamDefaultController", 31); js_set(js, g, "ReadableStreamDefaultController", ctrl_ctor); js_set_descriptor(js, g, "ReadableStreamDefaultController", 31, JS_DESC_W | JS_DESC_C); g_reader_proto = js_mkobj(js); - js_set_getter_desc(js, g_reader_proto, "closed", 6, - js_mkfun(js_rs_reader_get_closed), JS_DESC_C); + js_set_getter_desc(js, g_reader_proto, "closed", 6, js_mkfun(js_rs_reader_get_closed), JS_DESC_C); js_set(js, g_reader_proto, "read", js_mkfun(js_rs_reader_read)); js_set_descriptor(js, g_reader_proto, "read", 4, JS_DESC_W | JS_DESC_C); js_set(js, g_reader_proto, "releaseLock", js_mkfun(js_rs_reader_release_lock)); js_set_descriptor(js, g_reader_proto, "releaseLock", 11, JS_DESC_W | JS_DESC_C); js_set(js, g_reader_proto, "cancel", js_mkfun(js_rs_reader_cancel)); js_set_descriptor(js, g_reader_proto, "cancel", 6, JS_DESC_W | JS_DESC_C); - js_set_sym(js, g_reader_proto, get_toStringTag_sym(), - js_mkstr(js, "ReadableStreamDefaultReader", 27)); + js_set_sym(js, g_reader_proto, get_toStringTag_sym(), js_mkstr(js, "ReadableStreamDefaultReader", 27)); - ant_value_t reader_ctor = js_make_ctor(js, js_rs_reader_ctor, g_reader_proto, - "ReadableStreamDefaultReader", 27); + ant_value_t reader_ctor = js_make_ctor(js, js_rs_reader_ctor, g_reader_proto, "ReadableStreamDefaultReader", 27); js_set(js, g, "ReadableStreamDefaultReader", reader_ctor); js_set_descriptor(js, g, "ReadableStreamDefaultReader", 27, JS_DESC_W | JS_DESC_C); g_rs_proto = js_mkobj(js); - js_set_getter_desc(js, g_rs_proto, "locked", 6, - js_mkfun(js_rs_get_locked), JS_DESC_C); + js_set_getter_desc(js, g_rs_proto, "locked", 6, js_mkfun(js_rs_get_locked), JS_DESC_C); js_set(js, g_rs_proto, "cancel", js_mkfun(js_rs_cancel)); js_set_descriptor(js, g_rs_proto, "cancel", 6, JS_DESC_W | JS_DESC_C); js_set(js, g_rs_proto, "getReader", js_mkfun(js_rs_get_reader)); js_set_descriptor(js, g_rs_proto, "getReader", 9, JS_DESC_W | JS_DESC_C); - js_set(js, g_rs_proto, "tee", js_mkfun(js_rs_tee)); - js_set_descriptor(js, g_rs_proto, "tee", 3, JS_DESC_W | JS_DESC_C); - js_set(js, g_rs_proto, "pipeThrough", js_mkfun(js_rs_pipe_through)); - js_set_descriptor(js, g_rs_proto, "pipeThrough", 11, JS_DESC_W | JS_DESC_C); - js_set(js, g_rs_proto, "pipeTo", js_mkfun(js_rs_pipe_to)); - js_set_descriptor(js, g_rs_proto, "pipeTo", 6, JS_DESC_W | JS_DESC_C); + init_pipes_proto(js, g_rs_proto); + js_set(js, g_rs_proto, "values", js_mkfun(js_rs_values)); js_set_descriptor(js, g_rs_proto, "values", 6, JS_DESC_W | JS_DESC_C); - js_set_sym(js, g_rs_proto, get_toStringTag_sym(), - js_mkstr(js, "ReadableStream", 14)); + js_set_sym(js, g_rs_proto, get_toStringTag_sym(), js_mkstr(js, "ReadableStream", 14)); ant_value_t rs_ctor = js_make_ctor(js, js_rs_ctor, g_rs_proto, "ReadableStream", 14); js_set(js, g, "ReadableStream", rs_ctor); diff --git a/src/streams/writable.c b/src/streams/writable.c index fc2cf96..40317d7 100644 --- a/src/streams/writable.c +++ b/src/streams/writable.c @@ -9,6 +9,7 @@ #include "silver/engine.h" #include "modules/symbol.h" +#include "modules/assert.h" #include "modules/abort.h" #include "streams/writable.h" @@ -23,7 +24,7 @@ ws_stream_t *ws_get_stream(ant_value_t obj) { return (ws_stream_t *)(uintptr_t)(size_t)js_getnum(s); } -static ws_controller_t *ws_get_controller(ant_value_t obj) { +ws_controller_t *ws_get_controller(ant_value_t obj) { ant_value_t s = js_get_slot(obj, SLOT_DATA); if (vtype(s) != T_NUM) return NULL; return (ws_controller_t *)(uintptr_t)(size_t)js_getnum(s); @@ -238,12 +239,14 @@ static bool writable_stream_has_operation_in_flight(ant_value_t stream_obj) { static void ws_writer_ensure_ready_promise_rejected(ant_t *js, ant_value_t writer_obj, ant_value_t error) { ant_value_t ready = js_mkpromise(js); js_reject_promise(js, ready, error); + promise_mark_handled(ready); js_set_slot(writer_obj, SLOT_WS_READY, ready); } static void ws_writer_ensure_closed_promise_rejected(ant_t *js, ant_value_t writer_obj, ant_value_t error) { ant_value_t closed = js_mkpromise(js); js_reject_promise(js, closed, error); + promise_mark_handled(closed); js_set_slot(writer_obj, SLOT_RS_CLOSED, closed); } @@ -755,7 +758,7 @@ ant_value_t writable_stream_abort(ant_t *js, ant_value_t stream_obj, ant_value_t return promise; } -static ant_value_t ws_writer_write(ant_t *js, ant_value_t writer_obj, ant_value_t chunk) { +ant_value_t ws_writer_write(ant_t *js, ant_value_t writer_obj, ant_value_t chunk) { ant_value_t stream_obj = ws_writer_stream(writer_obj); if (!is_object_type(stream_obj)) { ant_value_t p = js_mkpromise(js); @@ -948,7 +951,7 @@ static ant_value_t js_ws_writer_write(ant_t *js, ant_value_t *args, int nargs) { return ws_writer_write(js, js->this_val, chunk); } -static ant_value_t js_ws_writer_ctor(ant_t *js, ant_value_t *args, int nargs) { +ant_value_t js_ws_writer_ctor(ant_t *js, ant_value_t *args, int nargs) { if (vtype(js->new_target) == T_UNDEF) return js_mkerr_typed(js, JS_ERR_TYPE, "WritableStreamDefaultWriter constructor requires 'new'"); if (nargs < 1) @@ -992,6 +995,17 @@ static ant_value_t js_ws_writer_ctor(ant_t *js, ant_value_t *args, int nargs) { return obj; } +ant_value_t ws_acquire_writer(ant_t *js, ant_value_t stream_obj) { + ant_value_t writer_args[1] = { stream_obj }; + ant_value_t saved = js->new_target; + + js->new_target = g_writer_proto; + ant_value_t writer = js_ws_writer_ctor(js, writer_args, 1); + js->new_target = saved; + + return writer; +} + static ant_value_t js_ws_get_locked(ant_t *js, ant_value_t *args, int nargs) { ws_stream_t *stream = ws_get_stream(js->this_val); if (!stream) return js_mkerr_typed(js, JS_ERR_TYPE, "Invalid WritableStream"); @@ -1233,25 +1247,19 @@ void init_writable_stream_module(void) { g_close_sentinel = js_mkobj(js); g_controller_proto = js_mkobj(js); - js_set_getter_desc(js, g_controller_proto, "signal", 6, - js_mkfun(js_ws_controller_get_signal), JS_DESC_C); + js_set_getter_desc(js, g_controller_proto, "signal", 6, js_mkfun(js_ws_controller_get_signal), JS_DESC_C); js_set(js, g_controller_proto, "error", js_mkfun(js_ws_controller_error)); js_set_descriptor(js, g_controller_proto, "error", 5, JS_DESC_W | JS_DESC_C); - js_set_sym(js, g_controller_proto, get_toStringTag_sym(), - js_mkstr(js, "WritableStreamDefaultController", 31)); + js_set_sym(js, g_controller_proto, get_toStringTag_sym(), js_mkstr(js, "WritableStreamDefaultController", 31)); - ant_value_t ctrl_ctor = js_make_ctor(js, js_ws_controller_ctor, g_controller_proto, - "WritableStreamDefaultController", 31); + ant_value_t ctrl_ctor = js_make_ctor(js, js_ws_controller_ctor, g_controller_proto, "WritableStreamDefaultController", 31); js_set(js, g, "WritableStreamDefaultController", ctrl_ctor); js_set_descriptor(js, g, "WritableStreamDefaultController", 31, JS_DESC_W | JS_DESC_C); g_writer_proto = js_mkobj(js); - js_set_getter_desc(js, g_writer_proto, "closed", 6, - js_mkfun(js_ws_writer_get_closed), JS_DESC_C); - js_set_getter_desc(js, g_writer_proto, "desiredSize", 11, - js_mkfun(js_ws_writer_get_desired_size), JS_DESC_C); - js_set_getter_desc(js, g_writer_proto, "ready", 5, - js_mkfun(js_ws_writer_get_ready), JS_DESC_C); + js_set_getter_desc(js, g_writer_proto, "closed", 6, js_mkfun(js_ws_writer_get_closed), JS_DESC_C); + js_set_getter_desc(js, g_writer_proto, "desiredSize", 11, js_mkfun(js_ws_writer_get_desired_size), JS_DESC_C); + js_set_getter_desc(js, g_writer_proto, "ready", 5, js_mkfun(js_ws_writer_get_ready), JS_DESC_C); js_set(js, g_writer_proto, "abort", js_mkfun(js_ws_writer_abort)); js_set_descriptor(js, g_writer_proto, "abort", 5, JS_DESC_W | JS_DESC_C); js_set(js, g_writer_proto, "close", js_mkfun(js_ws_writer_close)); @@ -1260,25 +1268,21 @@ void init_writable_stream_module(void) { js_set_descriptor(js, g_writer_proto, "releaseLock", 11, JS_DESC_W | JS_DESC_C); js_set(js, g_writer_proto, "write", js_mkfun(js_ws_writer_write)); js_set_descriptor(js, g_writer_proto, "write", 5, JS_DESC_W | JS_DESC_C); - js_set_sym(js, g_writer_proto, get_toStringTag_sym(), - js_mkstr(js, "WritableStreamDefaultWriter", 27)); + js_set_sym(js, g_writer_proto, get_toStringTag_sym(), js_mkstr(js, "WritableStreamDefaultWriter", 27)); - ant_value_t writer_ctor = js_make_ctor(js, js_ws_writer_ctor, g_writer_proto, - "WritableStreamDefaultWriter", 27); + ant_value_t writer_ctor = js_make_ctor(js, js_ws_writer_ctor, g_writer_proto, "WritableStreamDefaultWriter", 27); js_set(js, g, "WritableStreamDefaultWriter", writer_ctor); js_set_descriptor(js, g, "WritableStreamDefaultWriter", 27, JS_DESC_W | JS_DESC_C); g_ws_proto = js_mkobj(js); - js_set_getter_desc(js, g_ws_proto, "locked", 6, - js_mkfun(js_ws_get_locked), JS_DESC_C); + js_set_getter_desc(js, g_ws_proto, "locked", 6, js_mkfun(js_ws_get_locked), JS_DESC_C); js_set(js, g_ws_proto, "abort", js_mkfun(js_ws_abort)); js_set_descriptor(js, g_ws_proto, "abort", 5, JS_DESC_W | JS_DESC_C); js_set(js, g_ws_proto, "close", js_mkfun(js_ws_close)); js_set_descriptor(js, g_ws_proto, "close", 5, JS_DESC_W | JS_DESC_C); js_set(js, g_ws_proto, "getWriter", js_mkfun(js_ws_get_writer)); js_set_descriptor(js, g_ws_proto, "getWriter", 9, JS_DESC_W | JS_DESC_C); - js_set_sym(js, g_ws_proto, get_toStringTag_sym(), - js_mkstr(js, "WritableStream", 14)); + js_set_sym(js, g_ws_proto, get_toStringTag_sym(), js_mkstr(js, "WritableStream", 14)); ant_value_t ws_ctor = js_make_ctor(js, js_ws_ctor, g_ws_proto, "WritableStream", 14); js_set(js, g, "WritableStream", ws_ctor);