From ac922d08e77beedb671c63e3cde767ae6701dffa Mon Sep 17 00:00:00 2001 From: ldm0 Date: Fri, 18 Sep 2026 18:09:41 +0800 Subject: [PATCH] feat(dom): add Observable error recovery Add catch(callback) with native receiver and callback bindings. Forward source values and completion, then convert the catcher's return value into one recovery subscription after a source error. Recovery errors go directly downstream. Replace the traced upstream observer when recovery starts so pending Promises retain the replacement producer without retaining the exhausted source and catcher. Reuse the existing signal-owned cancellation algorithms. Cover Window/Worker conversion, error identity, shared subscriptions, reentrant cancellation, teardown ordering, realms and GC. Verify all nine upstream catch subtests in both realms and add those cases to the pass catalogue. --- .../wpt-cross-current/passed-cases.txt | 2 + moli-renderer-v8/src/observable.rs | 4 + moli-renderer-v8/src/observable/catch.rs | 146 ++++++++++++ moli-renderer-v8/src/observable/observer.rs | 7 +- .../src/script_vm/tests/observable.rs | 125 ++++++++++ .../src/worker/thread/tests/postmessage.rs | 16 ++ .../tests/fixtures/observable-catch-realms.js | 58 +++++ .../tests/fixtures/observable-catch.js | 218 ++++++++++++++++++ 8 files changed, 574 insertions(+), 2 deletions(-) create mode 100644 moli-renderer-v8/src/observable/catch.rs create mode 100644 moli-renderer-v8/tests/fixtures/observable-catch-realms.js create mode 100644 moli-renderer-v8/tests/fixtures/observable-catch.js 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}; +})()