diff --git a/moli-benchmark/wpt-cross-current/passed-cases.txt b/moli-benchmark/wpt-cross-current/passed-cases.txt index 8d94d4c2a..1fc940fe9 100644 --- a/moli-benchmark/wpt-cross-current/passed-cases.txt +++ b/moli-benchmark/wpt-cross-current/passed-cases.txt @@ -4537,6 +4537,8 @@ dom/observable/tentative/crashtests/observable-gc.any.js?moli-wpt-any=window dom/observable/tentative/crashtests/observable-takeUntil-toArray.any.js?moli-wpt-any=dedicatedworker dom/observable/tentative/crashtests/observable-takeUntil-toArray.any.js?moli-wpt-any=window dom/observable/tentative/idlharness.html +dom/observable/tentative/observable-catch.any.js?moli-wpt-any=dedicatedworker +dom/observable/tentative/observable-catch.any.js?moli-wpt-any=window dom/observable/tentative/observable-constructor.any.js?moli-wpt-any=dedicatedworker dom/observable/tentative/observable-constructor.any.js?moli-wpt-any=window dom/observable/tentative/observable-constructor.window.js?moli-wpt-script=window diff --git a/moli-renderer-v8/src/observable.rs b/moli-renderer-v8/src/observable.rs index fd8019e42..b2b361788 100644 --- a/moli-renderer-v8/src/observable.rs +++ b/moli-renderer-v8/src/observable.rs @@ -4,6 +4,7 @@ //! Window/worker owners. mod callbacks; +mod catch; mod collect; mod consume; mod event_target; @@ -52,6 +53,8 @@ struct ObservablePrototype { inspect: (), #[webapi(method, length = 1, callback = finally::finally)] finally: (), + #[webapi(method, length = 1, callback = catch::catch)] + catch: (), #[webapi(method, length = 0, returns_promise, callback = first::first)] first: (), #[webapi(method, length = 0, returns_promise, callback = collect::last)] @@ -326,6 +329,7 @@ fn subscribe_internal<'s>( && !finally::subscribe(scope, observable, subscriber) && !flat_map::subscribe(scope, observable, subscriber) && !switch_map::subscribe(scope, observable, subscriber) + && !catch::subscribe(scope, observable, subscriber) { event_target::subscribe(scope, observable, subscriber); } diff --git a/moli-renderer-v8/src/observable/catch.rs b/moli-renderer-v8/src/observable/catch.rs new file mode 100644 index 000000000..f1389c4de --- /dev/null +++ b/moli-renderer-v8/src/observable/catch.rs @@ -0,0 +1,146 @@ +//! Recover once from the source's error. Replacement errors pass downstream; +//! the recovery subscription does not retain the exhausted source or catcher. + +use moli_webapi_declare::WebApiObject; + +use super::{ + callbacks, from, + observer::{self, Notification}, + state::*, + subscribe_internal, subscriber_complete, subscriber_error, subscriber_next, +}; +use crate::{abort_signal_route::ResolvedAbortSignal, util::set_private_value, webidl}; + +const SOURCE: &str = "__moliObservableCatchSource"; +const CALLBACK: &str = "__moliObservableCatchCallback"; +const DOWNSTREAM: &str = "__moliCatchSubscriber"; + +#[derive(webidl::WebIdlArgs)] +#[webidl(prefix = "Observable.catch")] +struct CatchArgs { + #[webidl(required, converter = "callback_function")] + callback: webidl::WebIdlCallbackFunction, +} + +pub(super) fn catch<'s>( + scope: &mut v8::PinScope<'s, '_>, + args: v8::FunctionCallbackArguments<'s>, + mut rv: v8::ReturnValue<'_, v8::Value>, +) { + let Some(parsed) = webidl::parse_args::(scope, &args) else { + return; + }; + let Some(observable) = new_native_observable(scope, None) else { + return; + }; + set_private_value(scope, observable, SOURCE, args.this().into()); + set_callback(scope, observable, CALLBACK, parsed.callback); + rv.set(observable.into()); +} + +#[derive(WebApiObject)] +#[webapi(plain)] +struct CatchObserver<'scope> { + #[webapi(slot = observer::NATIVE_KIND)] + kind: i32, + #[webapi(slot = DOWNSTREAM)] + downstream: v8::Local<'scope, v8::Object>, +} + +fn subscribe_upstream<'s>( + scope: &mut v8::PinScope<'s, '_>, + observable: v8::Local<'s, v8::Object>, + subscriber: v8::Local<'s, v8::Object>, + observer: v8::Local<'s, v8::Object>, +) { + if active(scope, subscriber) { + // Replace this edge when recovery starts. A pending terminal Promise + // keeps the replacement producer alive without retaining the old one. + set_private_value(scope, subscriber, UPSTREAM_OBSERVER, observer.into()); + } + let signal = object_slot(scope, subscriber, SIGNAL) + .and_then(|signal| ResolvedAbortSignal::resolve(scope, signal)) + .expect("catch Subscriber signal"); + subscribe_internal(scope, observable, observer, Some(signal)); +} + +pub(super) fn subscribe<'s>( + scope: &mut v8::PinScope<'s, '_>, + observable: v8::Local<'s, v8::Object>, + subscriber: v8::Local<'s, v8::Object>, +) -> bool { + let Some(source) = object_slot(scope, observable, SOURCE) else { + return false; + }; + let callback = object_slot(scope, observable, CALLBACK).expect("Observable catcher"); + let Some(observer) = CatchObserver::new(observer::CATCH_SOURCE, subscriber) + .bind(scope) + .ok() + else { + return true; + }; + set_private_value(scope, observer, CALLBACK, callback.into()); + subscribe_upstream(scope, source, subscriber, observer); + true +} + +fn recover<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, + downstream: v8::Local<'s, v8::Object>, + error: v8::Local<'s, v8::Value>, +) { + let callback = object_slot(scope, observer, CALLBACK).expect("catch callback"); + // Source.error already closed its Subscriber before notifying us. Neither + // its producer nor the one-shot callback belongs to the replacement graph. + for slot in [observer::SUBSCRIBER, CALLBACK] { + set_private_value(scope, observer, slot, v8::undefined(scope).into()); + } + let result = match callbacks::invoke_value(scope, callback, &[error]) { + Ok(result) => result, + Err(error) => { + subscriber_error(scope, downstream, error); + return; + } + }; + let (inner, error) = { + v8::tc_scope!(let scope, scope); + let inner = from::convert(scope, result); + let error = scope.exception(); + scope.reset(); + (inner, error) + }; + if let Some(error) = error { + subscriber_error(scope, downstream, error); + return; + } + let Some(inner) = inner else { + return; + }; + let Some(inner_observer) = CatchObserver::new(observer::CATCH_INNER, downstream) + .bind(scope) + .ok() + else { + return; + }; + // Cancellation inside the callback/conversion still initializes a script + // replacement with an inactive Subscriber, following pre-aborted subscribe. + subscribe_upstream(scope, inner, downstream, inner_observer); +} + +pub(super) fn notify<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, + notification: Notification<'s>, + kind: i32, +) { + let downstream = object_slot(scope, observer, DOWNSTREAM).expect("catch downstream"); + match notification { + Notification::Next(value) => subscriber_next(scope, downstream, value), + Notification::Error(error) if kind == observer::CATCH_SOURCE => { + recover(scope, observer, downstream, error); + } + Notification::Error(error) => subscriber_error(scope, downstream, error), + Notification::Complete => subscriber_complete(scope, downstream), + } +} diff --git a/moli-renderer-v8/src/observable/observer.rs b/moli-renderer-v8/src/observable/observer.rs index b9594a9ad..62d6fab38 100644 --- a/moli-renderer-v8/src/observable/observer.rs +++ b/moli-renderer-v8/src/observable/observer.rs @@ -2,8 +2,8 @@ //! while script callbacks keep their typed Web IDL invocation boundary. use super::{ - callbacks, collect, consume, finally, first, flat_map, inspect, invoke_and_report, state::*, - switch_map, transform, until, + callbacks, catch, collect, consume, finally, first, flat_map, inspect, invoke_and_report, + state::*, switch_map, transform, until, }; use crate::util::{get_private_value, set_private_value}; @@ -28,6 +28,8 @@ pub(super) const FLAT_MAP_SOURCE: i32 = 17; pub(super) const FLAT_MAP_INNER: i32 = 18; pub(super) const SWITCH_MAP_SOURCE: i32 = 19; pub(super) const SWITCH_MAP_INNER: i32 = 20; +pub(super) const CATCH_SOURCE: i32 = 21; +pub(super) const CATCH_INNER: i32 = 22; pub(super) const SUBSCRIBER: &str = "__moliObservableNativeSubscriber"; const INDEX: &str = "__moliObservableCallbackIndex"; @@ -103,6 +105,7 @@ pub(super) fn notify<'s>( SWITCH_MAP_SOURCE | SWITCH_MAP_INNER => { switch_map::notify(scope, observer, notification, kind) } + CATCH_SOURCE | CATCH_INNER => catch::notify(scope, observer, notification, kind), _ => unreachable!("unknown native Observable observer"), } let exception = scope.exception(); diff --git a/moli-renderer-v8/src/script_vm/tests/observable.rs b/moli-renderer-v8/src/script_vm/tests/observable.rs index 5040accdd..a7944848f 100644 --- a/moli-renderer-v8/src/script_vm/tests/observable.rs +++ b/moli-renderer-v8/src/script_vm/tests/observable.rs @@ -1,5 +1,130 @@ use super::*; +#[test] +fn observable_catch_preserves_recovery_order_conversion_reentrancy_and_cancellation() { + let mut vm = new_storage_test_vm("https://observable-catch.test/"); + vm.eval(&format!( + "({}).then(value => {{ globalThis.catchResult = JSON.stringify(value); }});", + include_str!("../../../tests/fixtures/observable-catch.js") + )) + .expect("Observable.catch fixture should evaluate"); + let result: serde_json::Value = serde_json::from_str(&vm.eval("catchResult").unwrap()).unwrap(); + assert_eq!(result["failures"], serde_json::json!([]), "{result}"); + assert!(result["checks"].as_u64().unwrap() >= 130, "{result}"); +} + +#[test] +fn observable_catch_preserves_result_conversion_callback_and_cancellation_realms() { + let mut vm = new_storage_test_vm("https://observable-catch-realms.test/"); + vm.eval("document.appendChild(document.createElement('iframe'))") + .unwrap(); + materialize_single_child_default_realm_for_test(&mut vm, "Observable.catch realm"); + vm.eval(&format!( + "({}).then(value => {{ globalThis.catchRealms = JSON.stringify(value); }});", + include_str!("../../../tests/fixtures/observable-catch-realms.js") + )) + .expect("Observable.catch realms fixture should evaluate"); + let result: serde_json::Value = serde_json::from_str(&vm.eval("catchRealms").unwrap()).unwrap(); + assert_eq!(result["failures"], serde_json::json!([]), "{result}"); + assert_eq!(result["checks"], 35, "{result}"); +} + +#[test] +fn observable_catch_releases_exhausted_sources_and_callbacks_while_tracing_recovery() { + let mut vm = new_storage_test_vm("https://observable-catch-gc.test/"); + vm.eval(r#" +function makeCatchSource() { + let subscriber; + return {source: new Observable(s => { subscriber = s; }), get subscriber() { return subscriber; }}; +} +function makeCatchInner(entry) { + return new Observable(s => { entry.inner = new WeakRef(s); }); +} +function makeCatcher(value) { + const token = {value}; + return {token, callback: () => token.value}; +} +globalThis.catchChains = []; +for (const mode of ['abandoned', 'complete', 'source-abort', 'recover-complete', 'recover-error', 'recover-abort']) (() => { + const outer = makeCatchSource(), entry = {mode}, inner = makeCatchInner(entry), catcher = makeCatcher(inner); + const result = outer.source.catch(catcher.callback), controller = new AbortController(); + const promise = result.toArray(mode.endsWith('abort') ? {signal: controller.signal} : undefined); + promise.catch(() => {}); outer.subscriber.next(1); + Object.assign(entry, {templates: [outer.source, result].map(v => new WeakRef(v)), + source: new WeakRef(outer.subscriber), innerTemplate: new WeakRef(inner), + callbacks: [catcher.callback, catcher.token].map(v => new WeakRef(v))}); + if (mode !== 'abandoned') entry.promise = promise; + if (mode.endsWith('abort')) entry.controller = controller; + catchChains.push(entry); +})(); +"#).unwrap(); + let collect = |vm: &mut StandaloneScriptVmHarness| { + vm.renderer_document_isolate + .clone() + .with_entered_renderer_document_isolate(|isolate| { + isolate.clear_kept_objects(); + isolate.low_memory_notification(); + Ok(()) + }) + .unwrap(); + }; + collect(&mut vm); + assert_eq!(vm.eval(r#"JSON.stringify([ +catchChains.every(c => c.templates.every(r => r.deref() === undefined)), +catchChains.every(c => [c.source, c.innerTemplate, ...c.callbacks].every(r => (r.deref() !== undefined) === (c.mode !== 'abandoned'))) +])"#).unwrap(), "[true,true]"); + vm.eval( + r#" +for (const c of catchChains.filter(c => c.promise)) { + c.promise.then(v => { c.outcome = v; }, e => { c.outcome = e; }); + const source = c.source.deref(); + if (c.mode === 'complete') source.complete(); + else if (c.mode === 'source-abort') c.controller.abort(c.mode); + else { source.error('original'); c.inner.deref().next(2); } +} +"#, + ) + .unwrap(); + collect(&mut vm); + assert_eq!( + vm.eval( + r#"JSON.stringify([ +catchChains.every(c => c.source.deref() === undefined), +catchChains.every(c => [c.innerTemplate, ...c.callbacks].every(r => r.deref() === undefined)), +catchChains.filter(c => c.mode.startsWith('recover')).every(c => c.inner.deref() !== undefined), +catchChains.find(c => c.mode === 'complete').outcome.join(',') === '1', +catchChains.find(c => c.mode === 'source-abort').outcome === 'source-abort' +])"# + ) + .unwrap(), + "[true,true,true,true,true]" + ); + vm.eval( + r#" +for (const c of catchChains.filter(c => c.mode.startsWith('recover'))) { + const inner = c.inner.deref(); + if (c.mode === 'recover-complete') inner.complete(); + else if (c.mode === 'recover-error') inner.error(c.mode); + else c.controller.abort(c.mode); +} +"#, + ) + .unwrap(); + collect(&mut vm); + assert_eq!( + vm.eval( + r#"JSON.stringify([ +catchChains.filter(c => c.inner).every(c => c.inner.deref() === undefined), +catchChains.find(c => c.mode === 'recover-complete').outcome.join(',') === '1,2', +catchChains.find(c => c.mode === 'recover-error').outcome === 'recover-error', +catchChains.find(c => c.mode === 'recover-abort').outcome === 'recover-abort' +])"# + ) + .unwrap(), + "[true,true,true,true]" + ); +} + #[test] fn observable_switch_map_preserves_switch_order_conversion_reentrancy_and_cancellation() { let mut vm = new_storage_test_vm("https://observable-switch-map.test/"); diff --git a/moli-renderer-v8/src/worker/thread/tests/postmessage.rs b/moli-renderer-v8/src/worker/thread/tests/postmessage.rs index d3e6e6cdf..b4697eb1f 100644 --- a/moli-renderer-v8/src/worker/thread/tests/postmessage.rs +++ b/moli-renderer-v8/src/worker/thread/tests/postmessage.rs @@ -1,5 +1,21 @@ use super::*; +#[tokio::test] +async fn worker_observable_catch_preserves_recovery_order_conversion_reentrancy_and_cancellation() { + ensure_v8(); + let mut handle = spawn_worker( + format!( + "({}).then(value => {{ postMessage(value); close(); }});", + include_str!("../../../../tests/fixtures/observable-catch.js") + ), + "https://observable-catch.test/worker.js".into(), + ); + let message = timeout(TIMEOUT, handle.recv()).await.unwrap().unwrap(); + let result: serde_json::Value = serde_json::from_str(&expect_post_json(message)).unwrap(); + assert_eq!(result["failures"], serde_json::json!([]), "{result}"); + assert!(result["checks"].as_u64().unwrap() >= 130, "{result}"); +} + #[tokio::test] async fn worker_observable_switch_map_preserves_switch_order_conversion_reentrancy_and_cancellation() { diff --git a/moli-renderer-v8/tests/fixtures/observable-catch-realms.js b/moli-renderer-v8/tests/fixtures/observable-catch-realms.js new file mode 100644 index 000000000..c001b3292 --- /dev/null +++ b/moli-renderer-v8/tests/fixtures/observable-catch-realms.js @@ -0,0 +1,58 @@ +(async () => { + const failures = []; let checks = 0; + const check = (value, name) => { checks++; if (!value) failures.push(name); }; + const child = document.querySelector('iframe').contentWindow; + const method = child.Observable.prototype.catch, marker = {}, source = new Observable(s => s.error(marker)); + child.recoveryValue = 10; + const callback = child.Function('error', 'globalThis.catcherThis = this; globalThis.catcherArgc = arguments.length; globalThis.catcherError = error; return [recoveryValue];'); + const result = method.call(source, callback); + check(result instanceof child.Observable && !(result instanceof Observable), 'result uses callee realm'); + check(Object.getPrototypeOf(result) === child.Observable.prototype, 'callee intrinsic prototype'); + check(child.catcherArgc === undefined, 'foreign callback lazy'); + const values = await result.toArray(); + check(values instanceof child.Array && values[0] === 10, 'callee Array result'); + check(child.catcherThis === child && child.catcherArgc === 1, 'callback own realm and argument count'); + check(child.catcherError === marker, 'foreign callback receives original error'); + const local = Observable.prototype.catch.call(new child.Observable(s => s.error(3)), e => [e]); + check(local instanceof Observable && !(local instanceof child.Observable), 'local method returns local Observable'); + check(Object.getPrototypeOf(local) === Observable.prototype, 'local intrinsic prototype'); + const localValues = await local.toArray(); + check(localValues instanceof Array && localValues[0] === 3, 'local Array result'); + const revoked = Proxy.revocable(source, {}); revoked.revoke(); + for (const receiver of [{}, Object.create(source), new Proxy(source, {}), revoked.proxy]) { + let error; try { method.call(receiver, () => []); } catch (e) { error = e; } + check(error instanceof child.TypeError && !(error instanceof TypeError), 'receiver check uses callee TypeError'); + } + for (const input of [[], [undefined], [null], [1], [{}]]) { + let error; try { method.apply(source, input); } catch (e) { error = e; } + check(error instanceof child.TypeError && !(error instanceof TypeError), 'callback conversion uses callee TypeError'); + } + for (const input of [null, 1, {}, {then() {}}]) { + const error = await method.call(source, () => input).toArray().catch(e => e); + check(error instanceof child.TypeError && !(error instanceof TypeError), 'recovery conversion uses callee TypeError'); + } + const error = new child.Error('recovery'); child.recoveryError = error; + for (const callback of [child.Function('throw recoveryError'), () => ({get [Symbol.iterator]() { throw error; }}), + () => new child.Observable(s => s.error(error)), () => child.Promise.reject(error)]) { + check(await method.call(source, callback).toArray().catch(e => e) === error, 'callback, conversion and replacement error identity'); + } + for (const phase of ['source', 'recovery']) { + let outer, inner; + const foreignSource = new child.Observable(s => { outer = s; }); + const foreignInner = new child.Observable(s => { inner = s; }); + const ac = new AbortController(), reason = {}, log = []; + const pending = Observable.prototype.catch.call(foreignSource, () => foreignInner).toArray({signal: ac.signal}); + const outcome = pending.catch(e => e); + outer.addTeardown(() => log.push('source')); + if (phase === 'recovery') { outer.error(marker); inner.addTeardown(() => log.push('inner')); } + ac.abort(reason); + check(await outcome === reason, phase + ' cancellation rejection identity'); + check(!outer.active && (!inner || !inner.active), phase + ' foreign producers close'); + check((inner || outer).signal.reason === reason, phase + ' foreign signal reason identity'); + check(log.join(',') === (phase === 'source' ? 'source' : 'source,inner'), phase + ' foreign cleanup order'); + } + Object.setPrototypeOf(source, null); + const branded = method.call(source, () => [1]); + check(branded instanceof child.Observable && (await branded.toArray())[0] === 1, 'native brand survives prototype replacement'); + return {checks, failures}; +})() diff --git a/moli-renderer-v8/tests/fixtures/observable-catch.js b/moli-renderer-v8/tests/fixtures/observable-catch.js new file mode 100644 index 000000000..0d12432d9 --- /dev/null +++ b/moli-renderer-v8/tests/fixtures/observable-catch.js @@ -0,0 +1,218 @@ +(async () => { + 'use strict'; + const failures = []; let checks = 0; + const check = (v, name) => { checks++; if (!v) failures.push(name); }; + const same = (a, b, name) => check(JSON.stringify(a) === JSON.stringify(b), name); + const thrown = fn => { try { fn(); } catch (e) { return e; } }; + const test = async (name, fn) => { try { await fn(); } catch (e) { check(false, name + ': ' + e); } }; + const method = Observable.prototype.catch; + check(typeof method === 'function', 'catch exposed'); + if (failures.length) return {checks, failures}; + const fail = error => new Observable(s => s.error(error)); + function subject() { + let subscriber, starts = 0; + return {source: new Observable(s => { subscriber = s; starts++; }), + get subscriber() { return subscriber; }, get starts() { return starts; }}; + } + + await test('binding and intrinsic operations', async () => { + const desc = Object.getOwnPropertyDescriptor(Observable.prototype, 'catch'); + check(method.length === 1 && method.name === 'catch', 'name and length'); + check(desc.enumerable && desc.writable && desc.configurable, 'descriptor'); + check(thrown(() => new method(() => [])) instanceof TypeError, 'not constructible'); + const source = fail(1), revoked = Proxy.revocable(source, {}); revoked.revoke(); let traps = 0; + for (const receiver of [undefined, null, false, 1, Symbol(), {}, Object.create(source), Object.create(Observable.prototype), + new Proxy(source, {get() { traps++; }}), revoked.proxy]) { + check(thrown(() => method.call(receiver, () => [])) instanceof TypeError, 'invalid receiver'); + } + check(traps === 0, 'native receiver check bypasses author traps'); + check(thrown(() => source.catch()) instanceof TypeError, 'required callback'); + for (const callback of [undefined, null, false, 1, 1n, '', Symbol(), {}, [], {handleEvent() {}}]) { + check(thrown(() => source.catch(callback)) instanceof TypeError, 'non-callable callback'); + } + let calls = 0; + const callback = new Proxy(function(error) { + calls++; check(this === undefined && arguments.length === 1 && error === 1, 'callback receiver and exact argument'); + return [2]; + }, {get() { throw 'callback property read'; }}); + Object.defineProperty(source, 'constructor', {get() { throw 'source constructor'; }}); + Object.setPrototypeOf(source, null); + const result = method.call(source, callback, {get signal() { throw 'extra argument'; }}); + check(calls === 0, 'creation is lazy'); + check(Object.getPrototypeOf(result) === Observable.prototype && Observable.from(result) === result, 'intrinsic branded result'); + same(await result.toArray(), [2], 'native source ignores prototype and constructor'); + same(await result.toArray(), [2], 'fresh subscription can recover again'); + check(calls === 2, 'one catch per subscription'); + class Derived extends Observable {} + check(!(new Derived(s => s.complete()).catch(callback) instanceof Derived), 'source species ignored'); + const saved = []; + for (const [object, key] of [[Observable, 'from'], [Observable.prototype, 'subscribe'], [Subscriber.prototype, 'next'], + [Subscriber.prototype, 'complete'], [Subscriber.prototype, 'error']]) { + saved.push([object, key, Object.getOwnPropertyDescriptor(object, key)]); + Object.defineProperty(object, key, {value() { throw key; }, configurable: true}); + } + // A retained source Subscriber emits through its saved intrinsic method. + const nativeError = saved.find(([object, key]) => object === Subscriber.prototype && key === 'error')[2].value; + const primitiveSource = new Observable(s => nativeError.call(s, 1)); + try { same(await primitiveSource.catch(callback).toArray(), [2], 'internal conversion and subscription ignore public methods'); } + finally { for (const [object, key, descriptor] of saved) Object.defineProperty(object, key, descriptor); } + }); + + await test('pass through and replacement conversion', async () => { + let calls = 0; + same(await Observable.from([1, 2]).catch(() => { calls++; return []; }).toArray(), [1, 2], 'normal values pass through'); + check(calls === 0, 'completion does not call catcher'); + const marker = {}; + for (const input of [Observable.from([marker]), [marker], new Set([marker]), Promise.resolve(marker), + (async function* () { yield marker; })()]) { + const values = await fail('source').catch(() => input).toArray(); + check(values.length === 1 && values[0] === marker, 'converts recovery result and preserves value identity'); + } + for (const input of [undefined, null, false, 1, 1n, '', 'abc', Symbol(), {}, {then() { throw 'then'; }}]) { + check(await fail('source').catch(() => input).toArray().catch(e => e) instanceof TypeError, 'invalid recovery result rejects downstream'); + } + const log = []; + same(await fail(1).catch(() => ({ + get [Symbol.asyncIterator]() { log.push('async'); return undefined; }, + get [Symbol.iterator]() { log.push('sync'); return function* () { log.push('open'); yield 2; }; }, + get then() { throw 'then getter'; } + })).toArray(), [2], 'iterable conversion bypasses then property'); + same(log, ['async', 'sync', 'sync', 'open'], 'conversion probes precede obtaining iterator'); + const revoked = Proxy.revocable(() => [], {}); revoked.revoke(); + check(await fail(1).catch(revoked.proxy).toArray().catch(e => e) instanceof TypeError, 'revoked callable throws at invocation'); + }); + + await test('cleanup precedes recovery and recovery errors are not caught twice', async () => { + const marker = {}, log = [], outer = subject(), inner = subject(); let calls = 0; + const promise = outer.source.finally(() => log.push('source finally')).catch(error => { + calls++; check(error === marker && !outer.subscriber.active, 'catch sees original error after source closure'); + log.push('catch'); return inner.source; + }).finally(() => log.push('result finally')).toArray(); + outer.subscriber.signal.addEventListener('abort', () => log.push('source abort')); + outer.subscriber.addTeardown(() => log.push('source teardown')); + outer.subscriber.next(1); outer.subscriber.error(marker); + same(log, ['source abort', 'source teardown', 'source finally', 'catch'], 'source cleanup finishes before recovery callback'); + check(inner.subscriber.active && calls === 1, 'replacement remains active'); + inner.subscriber.addTeardown(() => log.push('inner teardown')); inner.subscriber.next(2); inner.subscriber.complete(); + same(await promise, [1, 2], 'preserves prefix and recovery values'); + same(log.slice(-2), ['inner teardown', 'result finally'], 'recovery cleanup precedes result finalizer'); + for (const mode of ['callback', 'conversion', 'initializer', 'inner']) for (const error of [{}, null, undefined]) { + let recovered = 0; const errors = [], values = []; + fail('original').catch(() => { + recovered++; + if (mode === 'callback') throw error; + if (mode === 'conversion') return {get [Symbol.asyncIterator]() { throw error; }}; + if (mode === 'initializer') return new Observable(() => { throw error; }); + return new Observable(s => { s.next(1); s.error(error); }); + }).subscribe({next: v => values.push(v), error: e => errors.push(e), complete: () => values.push('complete')}); + check(errors.length === 1 && errors[0] === error, mode + ' preserves replacement error identity'); + check(recovered === 1, mode + ' does not invoke catcher recursively'); + same(values, mode === 'inner' ? [1] : [], mode + ' does not complete after error'); + } + for (const error of [{}, null, undefined]) { + let seen, count = 0; + same(await new Observable(() => { throw error; }).catch(e => { seen = e; count++; return [3]; }).toArray(), [3], 'initializer errors recover'); + check(count === 1 && seen === error, 'initializer error identity reaches catcher'); + } + let starts = 0, recovered = 0; + const retry = new Observable(s => { starts++; s.next(starts); s.error(starts); }); + const events = []; + retry.catch(() => { recovered++; return retry; }).subscribe({next: v => events.push(v), error: e => events.push('error' + e)}); + same(events, [1, 2, 'error2'], 'returning source retries once then forwards second error'); + check(starts === 2 && recovered === 1, 'no implicit retry loop'); + }); + + await test('sharing, branches and later subscriptions', () => { + const outer = subject(), inner = subject(), ac1 = new AbortController(), ac2 = new AbortController(); + const a = [], b = [], branch = []; let calls = 0, branchCalls = 0; + const result = outer.source.catch(() => { calls++; return inner.source; }); + result.subscribe(v => a.push(v), {signal: ac1.signal}); result.subscribe(v => b.push(v), {signal: ac2.signal}); + outer.source.catch(() => { branchCalls++; return [9]; }).subscribe(v => branch.push(v)); + outer.subscriber.next(1); outer.subscriber.error('source'); + check(outer.starts === 1 && inner.starts === 1 && calls === 1 && branchCalls === 1, 'shared source and distinct catch branches'); + inner.subscriber.next(2); ac1.abort(); inner.subscriber.next(3); + check(inner.subscriber.active, 'first consumer removal keeps recovery alive'); + same(a, [1, 2], 'first consumer removed'); same(b, [1, 2, 3], 'second consumer survives'); same(branch, [1, 9], 'other branch recovers separately'); + const reason = {}; ac2.abort(reason); + check(!inner.subscriber.active && inner.subscriber.signal.reason === reason, 'last consumer cancels recovery with same reason'); + result.subscribe(); outer.subscriber.error('again'); + check(outer.starts === 2 && inner.starts === 2 && calls === 2, 'resubscription uses fresh source and replacement'); + inner.subscriber.complete(); + }); + + await test('cancellation before recovery, during callback, and in conversion', () => { + const pre = subject(); let calls = 0; + const reason = {}; + pre.source.catch(() => { calls++; return []; }).subscribe({}, {signal: AbortSignal.abort(reason)}); + check(pre.starts === 1 && !pre.subscriber.active && pre.subscriber.signal.reason === reason && calls === 0, 'pre-aborted source initialized inactive'); + for (const phase of ['before', 'callback', 'conversion']) { + const outer = subject(), ac = new AbortController(), log = []; let inactive, count = 0, iteratorCalls = 0; + outer.source.catch(() => { + count++; + if (phase === 'callback') { ac.abort(reason); return new Observable(s => { inactive = s; }); } + return {get [Symbol.asyncIterator]() { ac.abort(reason); return undefined; }, + [Symbol.iterator]() { iteratorCalls++; return [1][Symbol.iterator](); }}; + }).subscribe(v => log.push(v), {signal: ac.signal}); + if (phase === 'before') ac.abort(reason); else outer.subscriber.error('source'); + check(!outer.subscriber.active && log.length === 0, phase + ' cancellation suppresses output'); + check(count === (phase === 'before' ? 0 : 1), phase + ' catcher count'); + if (phase === 'callback') check(inactive && !inactive.active && inactive.signal.reason === reason, 'cancelled callback still initializes inactive replacement'); + if (phase === 'conversion') check(iteratorCalls === 0, 'cancelled conversion does not obtain an iterator'); + } + const captured = subject(), ac = new AbortController(); let count = 0, inner; + captured.source.subscribe({error: () => ac.abort(reason)}); + captured.source.catch(() => { count++; return new Observable(s => { inner = s; }); }).subscribe({}, {signal: ac.signal}); + captured.subscriber.error('captured'); + check(count === 1 && inner && !inner.active, 'captured error still calls catcher after earlier observer cancels'); + }); + + await test('cancellation order and IteratorClose failures', () => { + const log = [], ac = new AbortController(), reason = {}; let inner; + ac.signal.addEventListener('abort', () => log.push('consumer abort')); + fail('source').catch(() => new Observable(s => { + inner = s; s.signal.addEventListener('abort', () => log.push('inner abort')); s.addTeardown(() => log.push('inner teardown')); + })).finally(() => log.push('finally')).subscribe({}, {signal: ac.signal}); + ac.abort(reason); + same(log, ['inner abort', 'inner teardown', 'finally', 'consumer abort'], 'nested cancellation algorithms precede outer abort event'); + check(inner.signal.reason === reason, 'replacement receives cancellation identity'); + for (const marker of [{}, undefined]) { + const controller = new AbortController(), log = []; let pulls = 0, returns = 0, caught, didThrow = false; + fail(1).catch(() => ({[Symbol.iterator]() { return { + next() { return ++pulls < 3 ? {value: 1} : {done: true}; }, return() { returns++; log.push('return'); throw marker; } + }; }})).finally(() => log.push('finally')).subscribe(() => { + try { controller.abort(); } catch (e) { caught = e; didThrow = true; } + log.push('after abort'); + }, {signal: controller.signal}); + check(didThrow && caught === marker, 'recovery IteratorClose exception propagates including undefined'); + check(pulls === 1 && returns === 1, 'cancelled recovery closes once without extra pulls'); + same(log, ['return', 'finally', 'after abort'], 'finalizer still runs before abort rethrows'); + } + }); + + await test('reentrancy and compositions', async () => { + const outer = subject(), order = []; let calls = 0; + const result = outer.source.catch(() => { calls++; order.push('catch'); return [2]; }); + const values = []; result.subscribe(v => values.push(v)); + outer.subscriber.addTeardown(() => { + order.push('teardown'); result.subscribe(v => order.push('late' + v)); + }); + outer.subscriber.error('source'); + same(order, ['teardown', 'catch', 'late2'], 'reentrant result subscription joins recovery'); + same(values, [2], 'original consumer receives recovery'); check(calls === 1, 'reentrant join does not invoke callback twice'); + same(await Observable.from([1, 2, 3]).flatMap(v => v === 2 ? fail('two').catch(() => []) : [v]).toArray(), [1, 3], 'flatMap continues after recovered inner'); + same(await fail('a').catch(() => fail('b')).catch(e => [e]).toArray(), ['b'], 'outer catch can recover replacement errors'); + let resolve; const pending = new Promise(r => { resolve = r; }); const ac = new AbortController(); + const promise = fail('source').catch(() => pending).toArray({signal: ac.signal}); const outcome = promise.catch(e => e); + const reason = {}; ac.abort(reason); resolve('late'); + check(await outcome === reason, 'promise recovery cancellation keeps consumer rejection'); + let source; const reports = [], cleanupError = {}, observerError = {}, onerror = e => { reports.push(e.error); e.preventDefault(); }; let caught = 0; + addEventListener('error', onerror); + try { + new Observable(s => { source = s; }).catch(() => { caught++; return []; }).subscribe(() => { throw observerError; }); + source.addTeardown(() => { throw cleanupError; }); source.next(1); source.complete(); + check(caught === 0, 'observer and teardown exceptions do not trigger catch'); + check(reports.length === 2 && reports[0] === observerError && reports[1] === cleanupError, 'observer and teardown exceptions report separately'); + } finally { removeEventListener('error', onerror); } + }); + return {checks, failures}; +})()