From b86a30705e30175295a5cea96d48749122b30d9c Mon Sep 17 00:00:00 2001 From: theMackabu Date: Tue, 20 Jan 2026 17:09:54 -0800 Subject: [PATCH] add Observable constructor and method test suite unify JS type constants into enum and flag macros --- examples/spec/observable.js | 422 ++++++++++++++++++++++++++++++++++++ include/internal.h | 21 +- src/ant.c | 11 +- src/gc.c | 42 ++-- src/modules/observable.c | 150 ++++++------- 5 files changed, 517 insertions(+), 129 deletions(-) create mode 100644 examples/spec/observable.js diff --git a/examples/spec/observable.js b/examples/spec/observable.js new file mode 100644 index 0000000..538cec7 --- /dev/null +++ b/examples/spec/observable.js @@ -0,0 +1,422 @@ +import { test, testThrows, testDeep, summary } from './helpers.js'; + +console.log('Observable Constructor Tests\n'); + +test('Observable is defined', typeof Observable, 'function'); +test('Observable.prototype exists', typeof Observable.prototype, 'object'); +test('Observable.prototype.constructor is Observable', Observable.prototype.constructor, Observable); + +testThrows('Constructor throws without arguments', () => new Observable()); +testThrows('Constructor throws with non-function', () => new Observable({})); +testThrows('Constructor throws with null', () => new Observable(null)); +testThrows('Constructor throws with number', () => new Observable(1)); + +let subscriberCalled = false; +new Observable(() => { + subscriberCalled = true; +}); +test('Subscriber not called by constructor', subscriberCalled, false); + +console.log('\nObservable.prototype.subscribe Tests\n'); + +test('subscribe method exists', typeof Observable.prototype.subscribe, 'function'); + +const obs1 = new Observable(sink => null); +let noThrow = true; +try { + obs1.subscribe(null); + obs1.subscribe(undefined); + obs1.subscribe(1); + obs1.subscribe({}); + obs1.subscribe(() => {}); +} catch (e) { + noThrow = false; +} +test('subscribe accepts any observer type', noThrow, true); + +const list1 = []; +const error1 = new Error('test error'); +new Observable(s => { + s.next(1); + s.error(error1); +}).subscribe( + x => list1.push('next:' + x), + e => list1.push(e), + () => list1.push('complete') +); +test('First subscribe arg is next callback', list1[0], 'next:1'); +test('Second subscribe arg is error callback', list1[1], error1); + +const list2 = []; +new Observable(s => { + s.complete(); +}).subscribe( + x => list2.push('next:' + x), + e => list2.push(e), + () => list2.push('complete') +); +test('Third subscribe arg is complete callback', list2[0], 'complete'); + +let observer1 = null; +new Observable(x => { + observer1 = x; +}).subscribe({}); +test('Subscriber receives observer object', typeof observer1, 'object'); +test('Observer has next method', typeof observer1.next, 'function'); +test('Observer has error method', typeof observer1.error, 'function'); +test('Observer has complete method', typeof observer1.complete, 'function'); + +console.log('\nSubscription Tests\n'); + +let cleanupCalled = 0; +const sub1 = new Observable(observer => { + return () => { + cleanupCalled++; + }; +}).subscribe({}); + +test('subscribe returns subscription object', typeof sub1, 'object'); +test('subscription has unsubscribe method', typeof sub1.unsubscribe, 'function'); +test('subscription.closed is false before unsubscribe', sub1.closed, false); +test('unsubscribe returns undefined', sub1.unsubscribe(), undefined); +test('cleanup function called on unsubscribe', cleanupCalled, 1); +test('subscription.closed is true after unsubscribe', sub1.closed, true); + +sub1.unsubscribe(); +test('cleanup not called again on second unsubscribe', cleanupCalled, 1); + +let cleanupOnError = 0; +new Observable(sink => { + sink.error(1); + return () => { + cleanupOnError++; + }; +}).subscribe({ error: () => {} }); +test('cleanup called on error', cleanupOnError, 1); + +let cleanupOnComplete = 0; +new Observable(sink => { + sink.complete(); + return () => { + cleanupOnComplete++; + }; +}).subscribe({}); +test('cleanup called on complete', cleanupOnComplete, 1); + +let unsubscribeCalled = 0; +const sub2 = new Observable(sink => { + return { + unsubscribe: () => { + unsubscribeCalled++; + } + }; +}).subscribe({}); +sub2.unsubscribe(); +test('subscription object with unsubscribe is valid return', unsubscribeCalled, 1); + +console.log('\nSubscriber Error Handling Tests\n'); + +const thrownError = new Error('subscriber error'); +let caughtError = null; +new Observable(() => { + throw thrownError; +}).subscribe({ + error: e => { + caughtError = e; + } +}); +test('subscriber errors sent to observer.error', caughtError, thrownError); + +console.log('\nObserver.start Tests\n'); + +const events1 = []; +const obs2 = new Observable(observer => { + events1.push('subscriber'); + observer.complete(); +}); + +let startSubscription = null; +let startThisVal = null; +const observer2 = { + start(subscription) { + events1.push('start'); + startSubscription = subscription; + startThisVal = this; + } +}; +const sub3 = obs2.subscribe(observer2); + +test('start called before subscriber', events1[0], 'start'); +test('subscriber called after start', events1[1], 'subscriber'); +test('start receives subscription', startSubscription, sub3); +test('start called with observer as this', startThisVal, observer2); + +const events2 = []; +const obs3 = new Observable(() => { + events2.push('subscriber'); +}); +obs3.subscribe({ + start(subscription) { + events2.push('start'); + subscription.unsubscribe(); + } +}); +test('unsubscribe in start prevents subscriber call', events2.length, 1); +test('only start was called', events2[0], 'start'); + +console.log('\nSubscriptionObserver.next Tests\n'); + +const token1 = {}; +let nextReceived = null; +let nextArgs = []; +new Observable(observer => { + observer.next(token1, 'extra'); +}).subscribe({ + next(value, ...args) { + nextReceived = value; + nextArgs = args; + } +}); +test('next forwards value to observer', nextReceived, token1); +test('next does not forward extra arguments', nextArgs.length, 0); + +let nextReturnVal = null; +new Observable(observer => { + nextReturnVal = observer.next(); + observer.complete(); +}).subscribe({ next: () => 'return value' }); +test('next suppresses observer return value', nextReturnVal, undefined); + +let nextAfterClose = 'not called'; +new Observable(observer => { + observer.complete(); + nextReturnVal = observer.next(); +}).subscribe({ + next: () => { + nextAfterClose = 'called'; + } +}); +test('next returns undefined when closed', nextReturnVal, undefined); +test('next does not call observer when closed', nextAfterClose, 'not called'); + +console.log('\nSubscriptionObserver.error Tests\n'); + +let errorReceived = null; +const errorToken = new Error('error token'); +new Observable(observer => { + observer.error(errorToken); +}).subscribe({ + error: e => { + errorReceived = e; + } +}); +test('error forwards exception to observer', errorReceived, errorToken); + +let errorAfterClose = 'not called'; +new Observable(observer => { + observer.complete(); + observer.error(new Error()); +}).subscribe({ + error: () => { + errorAfterClose = 'called'; + } +}); +test('error does not call observer when closed', errorAfterClose, 'not called'); + +let nextAfterError = 'not called'; +new Observable(observer => { + observer.error(new Error()); + observer.next(1); +}).subscribe({ + next: () => { + nextAfterError = 'called'; + }, + error: () => {} +}); +test('next not called after error', nextAfterError, 'not called'); + +console.log('\nSubscriptionObserver.complete Tests\n'); + +let completeReceived = false; +new Observable(observer => { + observer.complete(); +}).subscribe({ + complete: () => { + completeReceived = true; + } +}); +test('complete calls observer.complete', completeReceived, true); + +let completeAfterClose = 'not called'; +new Observable(observer => { + observer.complete(); + observer.complete(); +}).subscribe({ + complete: () => { + completeAfterClose = 'called'; + } +}); +test('complete only called once', completeAfterClose, 'called'); + +let nextAfterComplete = 'not called'; +new Observable(observer => { + observer.complete(); + observer.next(1); +}).subscribe({ + next: () => { + nextAfterComplete = 'called'; + } +}); +test('next not called after complete', nextAfterComplete, 'not called'); + +console.log('\nObservable.of Tests\n'); + +test('Observable.of exists', typeof Observable.of, 'function'); + +const ofValues = []; +Observable.of(1, 2, 3, 4).subscribe({ + next: v => ofValues.push(v), + complete: () => ofValues.push('done') +}); +testDeep('Observable.of delivers all values', ofValues, [1, 2, 3, 4, 'done']); + +const ofEmpty = []; +Observable.of().subscribe({ + next: v => ofEmpty.push(v), + complete: () => ofEmpty.push('done') +}); +testDeep('Observable.of with no args calls complete', ofEmpty, ['done']); + +console.log('\nObservable.from Tests\n'); + +test('Observable.from exists', typeof Observable.from, 'function'); + +testThrows('Observable.from throws with null', () => Observable.from(null)); +testThrows('Observable.from throws with undefined', () => Observable.from(undefined)); +testThrows('Observable.from throws with no args', () => Observable.from()); + +const fromArray = []; +Observable.from([1, 2, 3]).subscribe({ + next: v => fromArray.push(v), + complete: () => fromArray.push('done') +}); +testDeep('Observable.from works with arrays', fromArray, [1, 2, 3, 'done']); + +const customObservable = { + [Symbol.observable]() { + return new Observable(observer => { + observer.next('custom'); + observer.complete(); + }); + } +}; +const fromCustom = []; +Observable.from(customObservable).subscribe({ + next: v => fromCustom.push(v), + complete: () => fromCustom.push('done') +}); +testDeep('Observable.from works with @@observable', fromCustom, ['custom', 'done']); + +const obs4 = new Observable(o => { + o.next(1); + o.complete(); +}); +const fromObs = Observable.from(obs4); +test('Observable.from returns same observable if already Observable', fromObs, obs4); + +testThrows('Observable.from throws with non-iterable object', () => { + Observable.from({}).subscribe({}); +}); + +console.log('\nSymbol.observable Tests\n'); + +test('Symbol.observable exists', typeof Symbol.observable, 'symbol'); + +const obs5 = new Observable(() => {}); +test('Observable has @@observable method', typeof obs5[Symbol.observable], 'function'); +test('@@observable returns this', obs5[Symbol.observable](), obs5); + +console.log('\nEarly Unsubscribe Tests\n'); + +const earlyValues = []; +let earlyCleanup = false; +let sub4; +new Observable(observer => { + observer.next(1); + observer.next(2); + observer.next(3); + observer.complete(); + return () => { + earlyCleanup = true; + }; +}).subscribe({ + start: subscription => { sub4 = subscription; }, + next: v => { + earlyValues.push(v); + if (v === 2) sub4.unsubscribe(); + } +}); +test('early unsubscribe stops delivery', earlyValues.length <= 3, true); +test('cleanup called on early unsubscribe', earlyCleanup, true); + +console.log('\nIterator Closing Tests\n'); + +let iteratorClosed = false; +const customIterable = { + [Symbol.iterator]() { + let i = 0; + return { + next() { + if (i < 5) { + const val = i; + i++; + return { value: val, done: false }; + } + return { done: true }; + }, + return() { + iteratorClosed = true; + return { done: true }; + } + }; + } +}; + +const iterValues = []; +let sub5; +new Observable(observer => { + const iter = customIterable[Symbol.iterator](); + let result = iter.next(); + while (!result.done) { + observer.next(result.value); + if (observer.closed) { + if (iter.return) iter.return(); + break; + } + result = iter.next(); + } + if (!observer.closed) observer.complete(); +}).subscribe({ + start: subscription => { sub5 = subscription; }, + next: v => { + iterValues.push(v); + if (v === 2) sub5.unsubscribe(); + } +}); + +test('iterator early unsubscribe stops at value 2', iterValues.length <= 3, true); +test('iterator return() called on early unsubscribe', iteratorClosed, true); + +console.log('\nObserver.closed Property Tests\n'); + +let closedBefore = null; +let closedAfter = null; +new Observable(observer => { + closedBefore = observer.closed; + observer.complete(); + closedAfter = observer.closed; +}).subscribe({}); +test('observer.closed is false before complete', closedBefore, false); +test('observer.closed is true after complete', closedAfter, true); + +summary(); diff --git a/include/internal.h b/include/internal.h index 262afff..6b7ada3 100644 --- a/include/internal.h +++ b/include/internal.h @@ -49,18 +49,11 @@ struct js { int parse_depth; // recursion depth of parser (for stack overflow protection) }; -#define JS_T_OBJ 0 -#define JS_T_PROP 1 -#define JS_T_STR 2 - -#define JS_V_OBJ 0 -#define JS_V_PROP 1 -#define JS_V_STR 2 -#define JS_V_FUNC 7 -#define JS_V_ARR 11 -#define JS_V_PROMISE 12 -#define JS_V_BIGINT 14 -#define JS_V_GENERATOR 17 +enum { + T_OBJ, T_PROP, T_STR, T_UNDEF, T_NULL, T_NUM, T_BOOL, T_FUNC, + T_CODEREF, T_CFUNC, T_ERR, T_ARR, T_PROMISE, T_TYPEDARRAY, + T_BIGINT, T_PROPREF, T_SYMBOL, T_GENERATOR, T_FFI +}; #define JS_HASH_SIZE 512 #define JS_MAX_PARSE_DEPTH (1024 * 2) @@ -73,10 +66,12 @@ struct js { #define NANBOX_DATA_MASK 0x0000FFFFFFFFFFFFULL #define TYPE_FLAG(t) (1u << (t)) +#define T_NEEDS_PROTO_FALLBACK (TYPE_FLAG(T_FUNC) | TYPE_FLAG(T_ARR) | TYPE_FLAG(T_PROMISE)) #define T_NON_NUMERIC_MASK (TYPE_FLAG(T_STR) | TYPE_FLAG(T_ARR) | TYPE_FLAG(T_FUNC) | TYPE_FLAG(T_CFUNC) | TYPE_FLAG(T_OBJ)) -#define is_non_numeric(v) ((1u << vtype(v)) & T_NON_NUMERIC_MASK) void js_gc_update_roots(GC_UPDATE_ARGS); bool js_has_pending_coroutines(void); +#define is_non_numeric(v) ((1u << vtype(v)) & T_NON_NUMERIC_MASK) + #endif diff --git a/src/ant.c b/src/ant.c index 8a2d9d0..222c713 100644 --- a/src/ant.c +++ b/src/ant.c @@ -449,15 +449,6 @@ static const uint8_t body_end_tok[TOK_MAX] = { [TOK_SEMICOLON] = 1, [TOK_COMMA] = 1, [TOK_EOF] = 1, }; -enum { - T_OBJ, T_PROP, T_STR, T_UNDEF, T_NULL, T_NUM, T_BOOL, T_FUNC, - T_CODEREF, T_CFUNC, T_ERR, T_ARR, T_PROMISE, T_TYPEDARRAY, - T_BIGINT, T_PROPREF, T_SYMBOL, T_GENERATOR, T_FFI -}; - -#define TYPE_MASK(t) (1u << (t)) -#define T_NEEDS_PROTO_FALLBACK (TYPE_MASK(T_FUNC) | TYPE_MASK(T_ARR) | TYPE_MASK(T_PROMISE)) - static const char *typestr_raw(uint8_t t) { const char *names[] = { "object", "prop", "string", "undefined", "null", "number", @@ -5112,7 +5103,7 @@ static jsoff_t lkp_proto(struct js *js, jsval_t obj, const char *key, size_t len uint8_t pt = vtype(proto); if (pt == T_NULL) break; if (pt != T_OBJ && pt != T_ARR && pt != T_FUNC) { - if (TYPE_MASK(t) & T_NEEDS_PROTO_FALLBACK) { + if (TYPE_FLAG(t) & T_NEEDS_PROTO_FALLBACK) { cur = get_prototype_for_type(js, t); t = vtype(cur); if (t == T_NULL || t == T_UNDEF) break; diff --git a/src/gc.c b/src/gc.c index 26767b5..0ca2d35 100644 --- a/src/gc.c +++ b/src/gc.c @@ -97,10 +97,10 @@ static inline void gc_saveval(uint8_t *mem, jsoff_t off, jsval_t val) { static jsoff_t gc_esize(jsoff_t w) { jsoff_t cleaned = w & ~FLAGMASK; switch (cleaned & 3U) { - case JS_T_OBJ: return (jsoff_t)(sizeof(jsoff_t) + sizeof(jsoff_t)); - case JS_T_PROP: return (jsoff_t)(sizeof(jsoff_t) + sizeof(jsoff_t) + sizeof(jsval_t)); - case JS_T_STR: return (jsoff_t)(sizeof(jsoff_t) + ((cleaned >> 2) + 3) / 4 * 4); - default: return (jsoff_t)~0U; + case T_OBJ: return (jsoff_t)(sizeof(jsoff_t) + sizeof(jsoff_t)); + case T_PROP: return (jsoff_t)(sizeof(jsoff_t) + sizeof(jsoff_t) + sizeof(jsval_t)); + case T_STR: return (jsoff_t)(sizeof(jsoff_t) + ((cleaned >> 2) + 3) / 4 * 4); + default: return (jsoff_t)~0U; } } @@ -137,7 +137,7 @@ static jsoff_t gc_copy_string(gc_ctx_t *ctx, jsoff_t old_off) { if (new_off != (jsoff_t)~0) return new_off; jsoff_t header = gc_loadoff(ctx->js->mem, old_off); - if ((header & 3) != JS_T_STR) return old_off; + if ((header & 3) != T_STR) return old_off; jsoff_t size = gc_esize(header); if (size == (jsoff_t)~0) return old_off; @@ -160,7 +160,7 @@ static jsoff_t gc_copy_prop(gc_ctx_t *ctx, jsoff_t old_off) { if (new_off != (jsoff_t)~0) return new_off; jsoff_t header = gc_loadoff(ctx->js->mem, old_off); - if ((header & 3) != JS_T_PROP) { + if ((header & 3) != T_PROP) { return old_off; } @@ -210,7 +210,7 @@ static jsoff_t gc_copy_object(gc_ctx_t *ctx, jsoff_t old_off) { if (new_off != (jsoff_t)~0) return new_off; jsoff_t header = gc_loadoff(ctx->js->mem, old_off); - if ((header & 3) != JS_T_OBJ) return old_off; + if ((header & 3) != T_OBJ) return old_off; jsoff_t size = gc_esize(header); if (size == (jsoff_t)~0) return old_off; @@ -244,10 +244,10 @@ static jsoff_t gc_copy_entity(gc_ctx_t *ctx, jsoff_t old_off) { jsoff_t header = gc_loadoff(ctx->js->mem, old_off); switch (header & 3) { - case JS_T_OBJ: return gc_copy_object(ctx, old_off); - case JS_T_PROP: return gc_copy_prop(ctx, old_off); - case JS_T_STR: return gc_copy_string(ctx, old_off); - default: return old_off; + case T_OBJ: return gc_copy_object(ctx, old_off); + case T_PROP: return gc_copy_prop(ctx, old_off); + case T_STR: return gc_copy_string(ctx, old_off); + default: return old_off; } } @@ -258,17 +258,17 @@ static jsval_t gc_update_val(gc_ctx_t *ctx, jsval_t val) { jsoff_t old_off = (jsoff_t)gc_vdata(val); switch (type) { - case JS_V_OBJ: - case JS_V_FUNC: - case JS_V_ARR: - case JS_V_PROMISE: - case JS_V_GENERATOR: { + case T_OBJ: + case T_FUNC: + case T_ARR: + case T_PROMISE: + case T_GENERATOR: { if (old_off >= ctx->js->brk) return val; jsoff_t new_off = gc_copy_object(ctx, old_off); if (new_off != (jsoff_t)~0) return gc_mkval(type, new_off); break; } - case JS_V_STR: { + case T_STR: { if (old_off >= ctx->js->brk) return val; jsoff_t new_off = gc_copy_string(ctx, old_off); if (new_off != (jsoff_t)~0) { @@ -276,7 +276,7 @@ static jsval_t gc_update_val(gc_ctx_t *ctx, jsval_t val) { } break; } - case JS_V_PROP: { + case T_PROP: { if (old_off >= ctx->js->brk) return val; jsoff_t new_off = gc_copy_prop(ctx, old_off); if (new_off != (jsoff_t)~0) { @@ -284,7 +284,7 @@ static jsval_t gc_update_val(gc_ctx_t *ctx, jsval_t val) { } break; } - case JS_V_BIGINT: { + case T_BIGINT: { if (old_off >= ctx->js->brk) return val; jsoff_t new_off = fwd_lookup(&ctx->fwd, old_off); if (new_off != (jsoff_t)~0) { @@ -357,13 +357,13 @@ size_t js_gc_compact(struct js *js) { if (js->brk > 0) { jsoff_t header_at_0 = gc_loadoff(js->mem, 0); - if ((header_at_0 & 3) == JS_T_OBJ) gc_copy_object(&ctx, 0); + if ((header_at_0 & 3) == T_OBJ) gc_copy_object(&ctx, 0); } jsoff_t scope_off = (jsoff_t)gc_vdata(js->scope); if (scope_off < js->brk) { jsoff_t new_scope = gc_copy_object(&ctx, scope_off); - js->scope = gc_mkval(JS_V_OBJ, new_scope); + js->scope = gc_mkval(T_OBJ, new_scope); } js->this_val = gc_update_val(&ctx, js->this_val); diff --git a/src/modules/observable.c b/src/modules/observable.c index e6e0ea9..362dd13 100644 --- a/src/modules/observable.c +++ b/src/modules/observable.c @@ -2,12 +2,17 @@ #include #include -#include "ant.h" +#include "internal.h" #include "runtime.h" #include "modules/symbol.h" #include "modules/observable.h" +static inline bool is_callable(jsval_t val) { + uint8_t t = vtype(val); + return t == T_FUNC || t == T_CFUNC; +} + static bool subscription_closed(struct js *js, jsval_t subscription) { jsval_t observer = js_get_slot(js, subscription, SLOT_SUBSCRIPTION_OBSERVER); return js_type(observer) == JS_UNDEF; @@ -16,14 +21,12 @@ static bool subscription_closed(struct js *js, jsval_t subscription) { static void cleanup_subscription(struct js *js, jsval_t subscription) { jsval_t cleanup = js_get_slot(js, subscription, SLOT_SUBSCRIPTION_CLEANUP); if (js_type(cleanup) == JS_UNDEF) return; - if (js_type(cleanup) != JS_FUNC) return; + if (!is_callable(cleanup)) return; js_set_slot(js, subscription, SLOT_SUBSCRIPTION_CLEANUP, js_mkundef()); jsval_t result = js_call(js, cleanup, NULL, 0); - if (js_type(result) == JS_ERR) { - fprintf(stderr, "Error in subscription cleanup: %s\n", js_str(js, result)); - } + if (js_type(result) == JS_ERR) fprintf(stderr, "Error in subscription cleanup: %s\n", js_str(js, result)); } static jsval_t create_subscription(struct js *js, jsval_t observer) { @@ -38,10 +41,7 @@ static jsval_t js_subscription_get_closed(struct js *js, jsval_t *args, int narg (void)args; (void)nargs; jsval_t subscription = js_getthis(js); - if (js_type(subscription) != JS_OBJ) { - return js_mkerr_typed(js, JS_ERR_TYPE, "Subscription.closed getter called on non-object"); - } - + if (js_type(subscription) != JS_OBJ) return js_mkerr_typed(js, JS_ERR_TYPE, "Subscription.closed getter called on non-object"); return subscription_closed(js, subscription) ? js_mktrue() : js_mkfalse(); } @@ -101,13 +101,11 @@ static jsval_t js_subobs_next(struct js *js, jsval_t *args, int nargs) { if (js_type(observer) != JS_OBJ) return js_mkundef(); jsval_t nextMethod = js_get(js, observer, "next"); - if (js_type(nextMethod) == JS_FUNC) { + if (is_callable(nextMethod)) { jsval_t value = (nargs > 0) ? args[0] : js_mkundef(); jsval_t call_args[1] = {value}; jsval_t result = js_call_with_this(js, nextMethod, observer, call_args, 1); - if (js_type(result) == JS_ERR) { - fprintf(stderr, "Error in observer.next: %s\n", js_str(js, result)); - } + if (js_type(result) == JS_ERR) fprintf(stderr, "Error in observer.next: %s\n", js_str(js, result)); } return js_mkundef(); @@ -132,13 +130,11 @@ static jsval_t js_subobs_error(struct js *js, jsval_t *args, int nargs) { if (js_type(observer) == JS_OBJ) { jsval_t errorMethod = js_get(js, observer, "error"); - if (js_type(errorMethod) == JS_FUNC) { + if (is_callable(errorMethod)) { jsval_t exception = (nargs > 0) ? args[0] : js_mkundef(); jsval_t call_args[1] = {exception}; jsval_t result = js_call_with_this(js, errorMethod, observer, call_args, 1); - if (js_type(result) == JS_ERR) { - fprintf(stderr, "Error in observer.error: %s\n", js_str(js, result)); - } + if (js_type(result) == JS_ERR) fprintf(stderr, "Error in observer.error: %s\n", js_str(js, result)); } } @@ -166,11 +162,9 @@ static jsval_t js_subobs_complete(struct js *js, jsval_t *args, int nargs) { if (js_type(observer) == JS_OBJ) { jsval_t completeMethod = js_get(js, observer, "complete"); - if (js_type(completeMethod) == JS_FUNC) { + if (is_callable(completeMethod)) { jsval_t result = js_call_with_this(js, completeMethod, observer, NULL, 0); - if (js_type(result) == JS_ERR) { - fprintf(stderr, "Error in observer.complete: %s\n", js_str(js, result)); - } + if (js_type(result) == JS_ERR) fprintf(stderr, "Error in observer.complete: %s\n", js_str(js, result)); } } @@ -180,13 +174,16 @@ static jsval_t js_subobs_complete(struct js *js, jsval_t *args, int nargs) { static jsval_t create_subscription_observer(struct js *js, jsval_t subscription) { jsval_t subobs = js_mkobj(js); + js_set_slot(js, subobs, SLOT_DATA, subscription); js_set(js, subobs, "next", js_mkfun(js_subobs_next)); js_set(js, subobs, "error", js_mkfun(js_subobs_error)); js_set(js, subobs, "complete", js_mkfun(js_subobs_complete)); js_set(js, subobs, get_toStringTag_sym_key(), js_mkstr(js, "SubscriptionObserver", 20)); + jsval_t closed_getter = js_mkfun(js_subobs_get_closed); js_set_getter_desc(js, subobs, "closed", 6, closed_getter, JS_DESC_E | JS_DESC_C); + return subobs; } @@ -198,7 +195,7 @@ static jsval_t js_cleanup_fn(struct js *js, jsval_t *args, int nargs) { if (js_type(subscription) != JS_OBJ) return js_mkundef(); jsval_t unsubscribe = js_get(js, subscription, "unsubscribe"); - if (js_type(unsubscribe) == JS_FUNC) { + if (is_callable(unsubscribe)) { return js_call_with_this(js, unsubscribe, subscription, NULL, 0); } @@ -211,7 +208,7 @@ static jsval_t execute_subscriber(struct js *js, jsval_t subscriber, jsval_t obs if (js_type(subscriberResult) == JS_ERR) return subscriberResult; if (js_type(subscriberResult) == JS_NULL || js_type(subscriberResult) == JS_UNDEF) return js_mkundef(); - if (js_type(subscriberResult) == JS_FUNC) return subscriberResult; + if (is_callable(subscriberResult)) return subscriberResult; if (js_type(subscriberResult) == JS_OBJ) { jsval_t result = js_get(js, subscriberResult, "unsubscribe"); @@ -236,13 +233,13 @@ static jsval_t js_observable_subscribe(struct js *js, jsval_t *args, int nargs) } jsval_t subscriber = js_get_slot(js, O, SLOT_OBSERVABLE_SUBSCRIBER); - if (js_type(subscriber) != JS_FUNC) { + if (!is_callable(subscriber)) { return js_mkerr_typed(js, JS_ERR_TYPE, "Observable has no [[Subscriber]] internal slot"); } jsval_t observer; - if (nargs > 0 && js_type(args[0]) == JS_FUNC) { + if (nargs > 0 && is_callable(args[0])) { jsval_t nextCallback = args[0]; jsval_t errorCallback = (nargs > 1) ? args[1] : js_mkundef(); jsval_t completeCallback = (nargs > 2) ? args[2] : js_mkundef(); @@ -253,15 +250,13 @@ static jsval_t js_observable_subscribe(struct js *js, jsval_t *args, int nargs) js_set(js, observer, "complete", completeCallback); } else if (nargs > 0 && js_type(args[0]) == JS_OBJ) { observer = args[0]; - } else { - observer = js_mkobj(js); - } + } else observer = js_mkobj(js); jsval_t subscription = create_subscription(js, observer); setup_subscription_methods(js, subscription); jsval_t start = js_get(js, observer, "start"); - if (js_type(start) == JS_FUNC) { + if (is_callable(start)) { jsval_t start_args[1] = {subscription}; jsval_t result = js_call_with_this(js, start, observer, start_args, 1); if (js_type(result) == JS_ERR) { @@ -274,18 +269,16 @@ static jsval_t js_observable_subscribe(struct js *js, jsval_t *args, int nargs) jsval_t subscriberResult = execute_subscriber(js, subscriber, subscriptionObserver); if (js_type(subscriberResult) == JS_ERR) { - jsval_t error_args[1] = {subscriberResult}; + jsval_t thrown_error = js->thrown_value; + js->thrown_value = js_mkundef(); + js->flags &= ~F_THROW; + + jsval_t error_args[1] = {thrown_error}; jsval_t error_method = js_get(js, subscriptionObserver, "error"); - if (js_type(error_method) == JS_FUNC) { - js_call_with_this(js, error_method, subscriptionObserver, error_args, 1); - } - } else { - js_set_slot(js, subscription, SLOT_SUBSCRIPTION_CLEANUP, subscriberResult); - } + if (is_callable(error_method)) js_call_with_this(js, error_method, subscriptionObserver, error_args, 1); + } else js_set_slot(js, subscription, SLOT_SUBSCRIPTION_CLEANUP, subscriberResult); - if (subscription_closed(js, subscription)) { - cleanup_subscription(js, subscription); - } + if (subscription_closed(js, subscription)) cleanup_subscription(js, subscription); return subscription; } @@ -301,7 +294,7 @@ static jsval_t js_observable_constructor(struct js *js, jsval_t *args, int nargs } jsval_t subscriber = args[0]; - if (js_type(subscriber) != JS_FUNC) { + if (!is_callable(subscriber)) { return js_mkerr_typed(js, JS_ERR_TYPE, "Observable subscriber must be a function"); } @@ -332,47 +325,38 @@ static jsval_t js_of_subscriber(struct js *js, jsval_t *args, int nargs) { jsval_t value = js_get(js, items, key); jsval_t next = js_get(js, observer, "next"); - if (js_type(next) == JS_FUNC) { + if (is_callable(next)) { jsval_t next_args[1] = {value}; js_call_with_this(js, next, observer, next_args, 1); } - if (js_type(subscription) == JS_OBJ && subscription_closed(js, subscription)) { - return js_mkundef(); - } + if (js_type(subscription) == JS_OBJ && subscription_closed(js, subscription)) return js_mkundef(); } jsval_t complete = js_get(js, observer, "complete"); - if (js_type(complete) == JS_FUNC) { - js_call_with_this(js, complete, observer, NULL, 0); - } + if (is_callable(complete)) js_call_with_this(js, complete, observer, NULL, 0); return js_mkundef(); } static jsval_t js_observable_of(struct js *js, jsval_t *args, int nargs) { jsval_t items = js_mkarr(js); - for (int i = 0; i < nargs; i++) { - js_arr_push(js, items, args[i]); - } - - jsval_t subscriber = js_mkobj(js); - js_set_slot(js, subscriber, SLOT_DATA, items); - js_set_slot(js, subscriber, SLOT_CFUNC, js_mkfun(js_of_subscriber)); - jsval_t subscriber_func = js_obj_to_func(subscriber); + for (int i = 0; i < nargs; i++) js_arr_push(js, items, args[i]); + jsval_t subscriber_func = js_heavy_mkfun(js, js_of_subscriber, items); jsval_t ctor_args[1] = {subscriber_func}; + return js_observable_constructor(js, ctor_args, 1); } static jsval_t js_from_delegating(struct js *js, jsval_t *args, int nargs) { jsval_t F = js_getcurrentfunc(js); - jsval_t observable = js_get_slot(js, F, SLOT_DATA); + jsval_t observable = js_get_slot(js, F, SLOT_DATA); if (js_type(observable) != JS_OBJ) return js_mkundef(); jsval_t subscribe = js_get(js, observable, "subscribe"); - if (js_type(subscribe) == JS_FUNC) { + if (is_callable(subscribe)) { return js_call_with_this(js, subscribe, observable, args, nargs); } @@ -391,7 +375,7 @@ static jsval_t js_from_iteration(struct js *js, jsval_t *args, int nargs) { jsval_t observer = args[0]; jsval_t subscription = js_get_slot(js, observer, SLOT_DATA); - if (js_type(iteratorMethod) != JS_FUNC) { + if (!is_callable(iteratorMethod)) { return js_mkerr_typed(js, JS_ERR_TYPE, "Object is not iterable"); } @@ -401,7 +385,7 @@ static jsval_t js_from_iteration(struct js *js, jsval_t *args, int nargs) { } jsval_t nextMethod = js_get(js, iterator, "next"); - if (js_type(nextMethod) != JS_FUNC) { + if (!is_callable(nextMethod)) { return js_mkerr_typed(js, JS_ERR_TYPE, "Iterator must have a next method"); } @@ -412,61 +396,60 @@ static jsval_t js_from_iteration(struct js *js, jsval_t *args, int nargs) { jsval_t done = js_get(js, next, "done"); if (js_truthy(js, done)) { jsval_t complete = js_get(js, observer, "complete"); - if (js_type(complete) == JS_FUNC) { - js_call_with_this(js, complete, observer, NULL, 0); - } + if (is_callable(complete)) js_call_with_this(js, complete, observer, NULL, 0); return js_mkundef(); } jsval_t nextValue = js_get(js, next, "value"); jsval_t obs_next = js_get(js, observer, "next"); - if (js_type(obs_next) == JS_FUNC) { + if (is_callable(obs_next)) { jsval_t next_args[1] = {nextValue}; js_call_with_this(js, obs_next, observer, next_args, 1); } if (js_type(subscription) == JS_OBJ && subscription_closed(js, subscription)) { jsval_t returnMethod = js_get(js, iterator, "return"); - if (js_type(returnMethod) == JS_FUNC) { - js_call_with_this(js, returnMethod, iterator, NULL, 0); - } + if (is_callable(returnMethod)) js_call_with_this(js, returnMethod, iterator, NULL, 0); return js_mkundef(); } } } static jsval_t js_observable_from(struct js *js, jsval_t *args, int nargs) { - if (nargs < 1) { - return js_mkerr_typed(js, JS_ERR_TYPE, "Observable.from requires an argument"); + if (nargs < 1) return js_mkerr_typed(js, JS_ERR_TYPE, "Observable.from requires an argument"); + jsval_t x = args[0]; + + if (js_type(x) == JS_NULL || js_type(x) == JS_UNDEF) { + return js_mkerr_typed(js, JS_ERR_TYPE, "Cannot convert null or undefined to observable"); } - jsval_t x = args[0]; jsval_t observableMethod = js_get(js, x, get_observable_sym_key()); - if (js_type(observableMethod) == JS_FUNC) { + if (is_callable(observableMethod)) { jsval_t observable = js_call_with_this(js, observableMethod, x, NULL, 0); if (js_type(observable) != JS_OBJ) { return js_mkerr_typed(js, JS_ERR_TYPE, "@@observable must return an object"); } - jsval_t constructor = js_get(js, observable, "constructor"); - jsval_t C = js_get(js, js_glob(js), "Observable"); - - if (constructor == C) return observable; - - jsval_t subscriber = js_mkobj(js); - js_set_slot(js, subscriber, SLOT_DATA, observable); - js_set_slot(js, subscriber, SLOT_CFUNC, js_mkfun(js_from_delegating)); - jsval_t subscriber_func = js_obj_to_func(subscriber); + jsval_t existing_subscriber = js_get_slot(js, observable, SLOT_OBSERVABLE_SUBSCRIBER); + if (is_callable(existing_subscriber)) return observable; + jsval_t subscriber_func = js_heavy_mkfun(js, js_from_delegating, observable); jsval_t ctor_args[1] = {subscriber_func}; + return js_observable_constructor(js, ctor_args, 1); } jsval_t iteratorMethod = js_get(js, x, get_iterator_sym_key()); - if (js_type(iteratorMethod) != JS_FUNC) { + if (!is_callable(iteratorMethod) && vtype(x) == T_ARR) { + jsval_t array_ctor = js_get(js, js_glob(js), "Array"); + jsval_t array_proto = js_get(js, array_ctor, "prototype"); + iteratorMethod = js_get(js, array_proto, get_iterator_sym_key()); + } + + if (!is_callable(iteratorMethod)) { return js_mkerr_typed(js, JS_ERR_TYPE, "Object is not observable or iterable"); } @@ -474,12 +457,9 @@ static jsval_t js_observable_from(struct js *js, jsval_t *args, int nargs) { js_set(js, data, "iterable", x); js_set(js, data, "iteratorMethod", iteratorMethod); - jsval_t subscriber = js_mkobj(js); - js_set_slot(js, subscriber, SLOT_DATA, data); - js_set_slot(js, subscriber, SLOT_CFUNC, js_mkfun(js_from_iteration)); - jsval_t subscriber_func = js_obj_to_func(subscriber); - + jsval_t subscriber_func = js_heavy_mkfun(js, js_from_iteration, data); jsval_t ctor_args[1] = {subscriber_func}; + return js_observable_constructor(js, ctor_args, 1); } -- 2.51.2