diff --git a/moli-benchmark/wpt-cross-current/passed-cases.txt b/moli-benchmark/wpt-cross-current/passed-cases.txt index e89c86477..8d94d4c2a 100644 --- a/moli-benchmark/wpt-cross-current/passed-cases.txt +++ b/moli-benchmark/wpt-cross-current/passed-cases.txt @@ -4549,6 +4549,8 @@ dom/observable/tentative/observable-every.any.js?moli-wpt-any=dedicatedworker dom/observable/tentative/observable-every.any.js?moli-wpt-any=window dom/observable/tentative/observable-filter.any.js?moli-wpt-any=dedicatedworker dom/observable/tentative/observable-filter.any.js?moli-wpt-any=window +dom/observable/tentative/observable-finally.any.js?moli-wpt-any=dedicatedworker +dom/observable/tentative/observable-finally.any.js?moli-wpt-any=window dom/observable/tentative/observable-find.any.js?moli-wpt-any=dedicatedworker dom/observable/tentative/observable-find.any.js?moli-wpt-any=window dom/observable/tentative/observable-first.any.js?moli-wpt-any=dedicatedworker @@ -4569,6 +4571,8 @@ dom/observable/tentative/observable-reduce.any.js?moli-wpt-any=dedicatedworker dom/observable/tentative/observable-reduce.any.js?moli-wpt-any=window dom/observable/tentative/observable-some.any.js?moli-wpt-any=dedicatedworker dom/observable/tentative/observable-some.any.js?moli-wpt-any=window +dom/observable/tentative/observable-switchMap.any.js?moli-wpt-any=dedicatedworker +dom/observable/tentative/observable-switchMap.any.js?moli-wpt-any=window dom/observable/tentative/observable-take.any.js?moli-wpt-any=dedicatedworker dom/observable/tentative/observable-take.any.js?moli-wpt-any=window dom/observable/tentative/observable-takeUntil.any.js?moli-wpt-any=dedicatedworker diff --git a/moli-renderer-v8/src/native_bridge/abort.rs b/moli-renderer-v8/src/native_bridge/abort.rs index d5fa58b33..b7e3843a1 100644 --- a/moli-renderer-v8/src/native_bridge/abort.rs +++ b/moli-renderer-v8/src/native_bridge/abort.rs @@ -300,11 +300,29 @@ impl AbortStore { signals_to_abort.push((dependent_signal_id, v8::Local::new(scope, signal))); } } + // Complete the dependency snapshot even if a source's IteratorClose + // throws. Dependents were already marked aborted and cannot be retried. + let mut first_error = None; for (signal_id, signal) in signals_to_abort { - if !self.run_abort_steps(scope, host, signal, signal_id, reason) { - return; + let error = { + v8::tc_scope!(let scope, scope); + if self.run_abort_steps(scope, host, signal, signal_id, reason) { + None + } else { + let Some(error) = scope.exception() else { + return; + }; + scope.reset(); + Some(error) + } + }; + if first_error.is_none() { + first_error = error; } } + if let Some(error) = first_error { + scope.throw_exception(error); + } } fn run_abort_steps<'s>( diff --git a/moli-renderer-v8/src/observable.rs b/moli-renderer-v8/src/observable.rs index 834fdc0ce..fd8019e42 100644 --- a/moli-renderer-v8/src/observable.rs +++ b/moli-renderer-v8/src/observable.rs @@ -15,6 +15,7 @@ mod inspect; mod observer; mod promise; mod state; +mod switch_map; mod transform; mod until; @@ -37,6 +38,8 @@ struct ObservablePrototype { map: (), #[webapi(method = "flatMap", length = 1, callback = flat_map::flat_map)] flat_map: (), + #[webapi(method = "switchMap", length = 1, callback = switch_map::switch_map)] + switch_map: (), #[webapi(method, length = 1, callback = transform::filter)] filter: (), #[webapi(method, length = 1, callback = transform::take)] @@ -322,6 +325,7 @@ fn subscribe_internal<'s>( && !inspect::subscribe(scope, observable, subscriber) && !finally::subscribe(scope, observable, subscriber) && !flat_map::subscribe(scope, observable, subscriber) + && !switch_map::subscribe(scope, observable, subscriber) { event_target::subscribe(scope, observable, subscriber); } diff --git a/moli-renderer-v8/src/observable/observer.rs b/moli-renderer-v8/src/observable/observer.rs index 0dd2bbf63..b9594a9ad 100644 --- a/moli-renderer-v8/src/observable/observer.rs +++ b/moli-renderer-v8/src/observable/observer.rs @@ -3,7 +3,7 @@ use super::{ callbacks, collect, consume, finally, first, flat_map, inspect, invoke_and_report, state::*, - transform, until, + switch_map, transform, until, }; use crate::util::{get_private_value, set_private_value}; @@ -26,6 +26,8 @@ pub(super) const INSPECT: i32 = 15; pub(super) const FINALLY: i32 = 16; 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 SUBSCRIBER: &str = "__moliObservableNativeSubscriber"; const INDEX: &str = "__moliObservableCallbackIndex"; @@ -98,6 +100,9 @@ pub(super) fn notify<'s>( FLAT_MAP_SOURCE | FLAT_MAP_INNER => { flat_map::notify(scope, observer, notification, kind) } + SWITCH_MAP_SOURCE | SWITCH_MAP_INNER => { + switch_map::notify(scope, observer, notification, kind) + } _ => unreachable!("unknown native Observable observer"), } let exception = scope.exception(); diff --git a/moli-renderer-v8/src/observable/switch_map.rs b/moli-renderer-v8/src/observable/switch_map.rs new file mode 100644 index 000000000..7aade4dc2 --- /dev/null +++ b/moli-renderer-v8/src/observable/switch_map.rs @@ -0,0 +1,268 @@ +//! Switch inner subscriptions through the existing AbortSignal dependency graph. +//! Reentrant mapper/teardown calls can leave multiple live inner observers, so +//! trace each until its own subscription completes or is cancelled. + +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::{get_private_value, set_private_value}, + webidl, +}; + +const SOURCE: &str = "__moliObservableSwitchMapSource"; +const MAPPER: &str = "__moliObservableSwitchMapMapper"; +const DOWNSTREAM: &str = "__moliSwitchMapSubscriber"; +const OWNER: &str = "__moliSwitchMapSourceObserver"; +const INNERS: &str = "__moliSwitchMapInnerObservers"; +const CURRENT_SIGNAL: &str = "__moliSwitchMapCurrentSignal"; +const OUTER_COMPLETE: &str = "__moliSwitchMapSourceCompleted"; +const INNER_SIGNAL: &str = "__moliSwitchMapInnerSignal"; +const CLEANUP: &str = "__moliSwitchMapInnerCleanup"; + +#[derive(webidl::WebIdlArgs)] +#[webidl(prefix = "Observable.switchMap")] +struct SwitchMapArgs { + #[webidl(required, converter = "callback_function")] + mapper: webidl::WebIdlCallbackFunction, +} + +pub(super) fn switch_map<'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, MAPPER, parsed.mapper); + rv.set(observable.into()); +} + +#[derive(WebApiObject)] +#[webapi(plain)] +struct SourceObserver<'scope> { + #[webapi(slot = observer::NATIVE_KIND)] + kind: i32, + #[webapi(slot = DOWNSTREAM)] + downstream: v8::Local<'scope, v8::Object>, + #[webapi(slot = MAPPER)] + mapper: v8::Local<'scope, v8::Object>, +} + +#[derive(WebApiObject)] +#[webapi(plain)] +struct InnerObserver<'scope> { + #[webapi(slot = observer::NATIVE_KIND)] + kind: i32, + #[webapi(slot = OWNER)] + owner: v8::Local<'scope, v8::Object>, + #[webapi(slot = INNER_SIGNAL)] + signal: v8::Local<'scope, v8::Object>, +} + +fn signal_slot<'s>( + scope: &mut v8::PinScope<'s, '_>, + object: v8::Local<'s, v8::Object>, + slot: &str, +) -> Option> { + object_slot(scope, object, slot).and_then(|signal| ResolvedAbortSignal::resolve(scope, 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 mapper = object_slot(scope, observable, MAPPER).expect("switchMap mapper"); + let Some(observer) = SourceObserver::new(observer::SWITCH_MAP_SOURCE, subscriber, mapper) + .bind(scope) + .ok() + else { + return true; + }; + if active(scope, subscriber) { + set_private_value(scope, subscriber, UPSTREAM_OBSERVER, observer.into()); + } + let signal = signal_slot(scope, subscriber, SIGNAL).expect("switchMap Subscriber signal"); + subscribe_internal(scope, source, observer, Some(signal)); + true +} + +fn release_inner<'s>(scope: &mut v8::PinScope<'s, '_>, inner: v8::Local<'s, v8::Object>) { + let owner = object_slot(scope, inner, OWNER).expect("switchMap owner"); + let mut inners = list(scope, owner, INNERS); + inners.retain(|entry| *entry != inner); + set_list(scope, owner, INNERS, &inners); + if let Some(signal) = signal_slot(scope, inner, INNER_SIGNAL) + && let Some(callback) = object_slot(scope, inner, CLEANUP) + .and_then(|value| v8::Local::::try_from(value).ok()) + { + signal.unregister_algorithm(scope, callback); + } + // Root the producer through the cancellation snapshot before dropping its + // persistent edge. Its native cancellation callback is deliberately weak. + let _producer = object_slot(scope, inner, observer::SUBSCRIBER); + for slot in [INNER_SIGNAL, CLEANUP, observer::SUBSCRIBER] { + set_private_value(scope, inner, slot, v8::undefined(scope).into()); + } +} + +fn cancelled<'s>( + scope: &mut v8::PinScope<'s, '_>, + args: v8::FunctionCallbackArguments<'s>, + _rv: v8::ReturnValue<'_, v8::Value>, +) { + let inner = v8::Local::::try_from(args.data()).expect("switchMap cleanup data"); + release_inner(scope, inner); +} + +fn process_next<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, + value: v8::Local<'s, v8::Value>, +) { + if let Some(signal) = signal_slot(scope, observer, CURRENT_SIGNAL) { + let error = { + v8::tc_scope!(let scope, scope); + let reason = crate::native_bridge::abort::abort_error_value(scope); + signal.abort(scope, reason); + let error = scope.exception(); + scope.reset(); + error + }; + // Internal switching must finish even if IteratorClose fails. Explicit + // author cancellation still uses the AbortController rethrow boundary. + if let Some(error) = error { + callbacks::report(scope, observer, error); + } + } + let Some(initial_signal) = ResolvedAbortSignal::new(scope) else { + return; + }; + set_private_value( + scope, + observer, + CURRENT_SIGNAL, + initial_signal.value().into(), + ); + let downstream = object_slot(scope, observer, DOWNSTREAM).expect("switchMap downstream"); + let mapper = object_slot(scope, observer, MAPPER).expect("switchMap mapper"); + let index = observer::index(scope, observer); + let index = v8::Number::new(scope, index as f64).into(); + let mapped = match callbacks::invoke_value(scope, mapper, &[value, index]) { + Ok(mapped) => mapped, + Err(error) => { + subscriber_error(scope, downstream, error); + return; + } + }; + observer::increment_index(scope, observer); + let (inner, error) = { + v8::tc_scope!(let scope, scope); + let inner = from::convert(scope, mapped); + 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; + }; + // The draft passes the current controller by reference: mapper/conversion + // reentrancy can replace it before subscription. If a reentrant inner also + // completed, it cleared that reference; the original (now aborted) signal + // safely initializes this superseded inner as inactive instead. + let current = signal_slot(scope, observer, CURRENT_SIGNAL).unwrap_or(initial_signal); + let downstream_signal = signal_slot(scope, downstream, SIGNAL).expect("switchMap signal"); + let Some(signal) = ResolvedAbortSignal::dependent(scope, &[current, downstream_signal]) else { + return; + }; + let Some(inner_observer) = + InnerObserver::new(observer::SWITCH_MAP_INNER, observer, signal.value()) + .bind(scope) + .ok() + else { + return; + }; + if !signal.is_aborted(scope) { + let mut inners = list(scope, observer, INNERS); + inners.push(inner_observer); + set_list(scope, observer, INNERS, &inners); + let cleanup = v8::Function::builder(cancelled) + .data(inner_observer.into()) + .build(scope) + .expect("switchMap inner cleanup"); + set_private_value(scope, inner_observer, CLEANUP, cleanup.into()); + signal.register_weak_algorithm(scope, cleanup); + } + subscribe_internal(scope, inner, inner_observer, Some(signal)); +} + +pub(super) fn notify<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, + notification: Notification<'s>, + kind: i32, +) { + let source = if kind == observer::SWITCH_MAP_INNER { + object_slot(scope, observer, OWNER).expect("switchMap source observer") + } else { + observer + }; + let downstream = object_slot(scope, source, DOWNSTREAM).expect("switchMap downstream"); + match notification { + Notification::Next(value) if kind == observer::SWITCH_MAP_SOURCE => { + process_next(scope, source, value); + } + Notification::Next(value) => subscriber_next(scope, downstream, value), + Notification::Error(error) => { + if kind == observer::SWITCH_MAP_INNER { + release_inner(scope, observer); + } + subscriber_error(scope, downstream, error); + } + Notification::Complete => { + if kind == observer::SWITCH_MAP_SOURCE { + set_private_value( + scope, + observer, + observer::SUBSCRIBER, + v8::undefined(scope).into(), + ); + set_private_value( + scope, + source, + OUTER_COMPLETE, + v8::Boolean::new(scope, true).into(), + ); + if signal_slot(scope, source, CURRENT_SIGNAL).is_none() { + subscriber_complete(scope, downstream); + } + } else { + release_inner(scope, observer); + if get_private_value(scope, source, OUTER_COMPLETE).is_some_and(|v| v.is_true()) { + subscriber_complete(scope, downstream); + } else { + set_private_value(scope, source, CURRENT_SIGNAL, v8::undefined(scope).into()); + } + } + } + } +} diff --git a/moli-renderer-v8/src/script_vm/tests/observable.rs b/moli-renderer-v8/src/script_vm/tests/observable.rs index 4a6e33b3d..5040accdd 100644 --- a/moli-renderer-v8/src/script_vm/tests/observable.rs +++ b/moli-renderer-v8/src/script_vm/tests/observable.rs @@ -1,5 +1,152 @@ use super::*; +#[test] +fn observable_switch_map_preserves_switch_order_conversion_reentrancy_and_cancellation() { + let mut vm = new_storage_test_vm("https://observable-switch-map.test/"); + vm.eval(&format!( + "({}).then(value => {{ globalThis.switchMapResult = JSON.stringify(value); }});", + include_str!("../../../tests/fixtures/observable-switch-map.js") + )) + .expect("Observable.switchMap fixture should evaluate"); + let result: serde_json::Value = + serde_json::from_str(&vm.eval("switchMapResult").unwrap()).unwrap(); + assert_eq!(result["failures"], serde_json::json!([]), "{result}"); + assert!(result["checks"].as_u64().unwrap() >= 170, "{result}"); +} + +#[test] +fn observable_switch_map_preserves_result_conversion_callback_and_cancellation_realms() { + let mut vm = new_storage_test_vm("https://observable-switch-map-realms.test/"); + vm.eval("document.appendChild(document.createElement('iframe'))") + .unwrap(); + materialize_single_child_default_realm_for_test(&mut vm, "Observable.switchMap realm"); + vm.eval(&format!( + "({}).then(value => {{ globalThis.switchMapRealms = JSON.stringify(value); }});", + include_str!("../../../tests/fixtures/observable-switch-map-realms.js") + )) + .expect("Observable.switchMap realms fixture should evaluate"); + let result: serde_json::Value = + serde_json::from_str(&vm.eval("switchMapRealms").unwrap()).unwrap(); + assert_eq!(result["failures"], serde_json::json!([]), "{result}"); + assert_eq!(result["checks"], 31, "{result}"); +} + +#[test] +fn observable_switch_map_traces_reentrant_inners_and_releases_replaced_and_closed_graphs() { + let mut vm = new_storage_test_vm("https://observable-switch-map-gc.test/"); + vm.eval(r#" +function makeSwitchSource() { + let subscriber; + return {source: new Observable(s => { subscriber = s; }), get subscriber() { return subscriber; }}; +} +function makeSwitchMapper() { + const token = {calls: 0}; + return {token, callback: value => { token.calls++; return value; }}; +} +globalThis.switchMapChains = []; +for (const mode of ['abandoned', 'complete', 'outer-error', 'inner-error', 'abort']) (() => { + const outer = makeSwitchSource(), old = makeSwitchSource(), inner = makeSwitchSource(), mapper = makeSwitchMapper(); + const result = outer.source.switchMap(mapper.callback), ac = new AbortController(); + const promise = result.toArray(mode === 'abort' ? {signal: ac.signal} : undefined); + promise.catch(() => {}); + outer.subscriber.next(old.source); old.subscriber.next(1); + outer.subscriber.next(inner.source); inner.subscriber.next(2); + const entry = {mode, templates: [outer.source, old.source, inner.source, result].map(v => new WeakRef(v)), + old: new WeakRef(old.subscriber), outer: new WeakRef(outer.subscriber), inner: new WeakRef(inner.subscriber), + callbacks: [mapper.callback, mapper.token].map(v => new WeakRef(v))}; + if (mode !== 'abandoned') entry.promise = promise; + if (mode === 'abort') entry.controller = ac; + switchMapChains.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([ +switchMapChains.every(c => c.templates.every(ref => ref.deref() === undefined)), +switchMapChains.every(c => c.old.deref() === undefined), +switchMapChains.every(c => [c.outer, c.inner, ...c.callbacks].every(ref => (ref.deref() !== undefined) === (c.mode !== 'abandoned'))) +])"#).unwrap(), "[true,true,true]"); + vm.eval(r#" +for (const c of switchMapChains.filter(c => c.promise)) { + c.promise.then(values => { c.correct = c.mode === 'complete' && JSON.stringify(values) === '[1,2]'; }, error => { c.correct = error === c.mode; }); + const outer = c.outer.deref(), inner = c.inner.deref(); + if (c.mode === 'complete') { outer.complete(); inner.complete(); } + else if (c.mode === 'outer-error') outer.error(c.mode); + else if (c.mode === 'inner-error') inner.error(c.mode); + else c.controller.abort(c.mode); + c.closed = [outer, inner]; +} +"#).unwrap(); + assert_eq!(vm.eval("switchMapChains.filter(c => c.promise).every(c => c.correct && c.closed.every(s => !s.active))").unwrap(), "true"); + collect(&mut vm); + assert_eq!( + vm.eval("switchMapChains.every(c => c.callbacks.every(ref => ref.deref() === undefined))") + .unwrap(), + "true" + ); + vm.eval(r#" +globalThis.switchMapDrain = (() => { + const outer = makeSwitchSource(), inner = makeSwitchSource(), mapper = makeSwitchMapper(); + const promise = outer.source.switchMap(mapper.callback).toArray(); + outer.subscriber.next(inner.source); outer.subscriber.complete(); + return {promise, outer: new WeakRef(outer.subscriber), inner: new WeakRef(inner.subscriber), callback: new WeakRef(mapper.callback)}; +})(); +"#).unwrap(); + collect(&mut vm); + assert_eq!(vm.eval("switchMapDrain.outer.deref() === undefined && switchMapDrain.inner.deref() !== undefined && switchMapDrain.callback.deref() !== undefined").unwrap(), "true"); + vm.eval("switchMapDrain.promise.then(v => { switchMapDrain.correct = v.length === 0; }); switchMapDrain.inner.deref().complete();").unwrap(); + collect(&mut vm); + assert_eq!(vm.eval("switchMapDrain.correct && switchMapDrain.inner.deref() === undefined && switchMapDrain.callback.deref() === undefined").unwrap(), "true"); + vm.eval(r#" +function makeReentrantSwitchMapper(outer) { + return value => { + if (value.next) outer.subscriber.next({current: value.next}); + return value.current; + }; +} +globalThis.reentrantSwitch = (() => { + const outer = makeSwitchSource(), first = makeSwitchSource(), second = makeSwitchSource(); + const promise = outer.source.switchMap(makeReentrantSwitchMapper(outer)).toArray(); + outer.subscriber.next({current: first.source, next: second.source}); + return {promise, source: new WeakRef(outer.subscriber), first: new WeakRef(first.subscriber), second: new WeakRef(second.subscriber)}; +})(); +"#).unwrap(); + collect(&mut vm); + assert_eq!(vm.eval("[reentrantSwitch.source, reentrantSwitch.first, reentrantSwitch.second].every(ref => ref.deref() !== undefined)").unwrap(), "true"); + vm.eval("reentrantSwitch.promise.then(v => { reentrantSwitch.correct = JSON.stringify(v) === '[1,2]'; }); reentrantSwitch.first.deref().next(1); reentrantSwitch.second.deref().next(2); reentrantSwitch.source.deref().complete(); reentrantSwitch.first.deref().complete();").unwrap(); + collect(&mut vm); + assert_eq!(vm.eval("reentrantSwitch.correct && [reentrantSwitch.source, reentrantSwitch.first, reentrantSwitch.second].every(ref => ref.deref() === undefined)").unwrap(), "true"); +} + +#[test] +fn observable_switch_map_handles_completion_during_reentrant_mapping() { + let mut vm = new_storage_test_vm("https://observable-switch-map-reentrant-complete.test/"); + assert_eq!(vm.eval(r#" +JSON.stringify(['mapper', 'conversion'].map(mode => { + let source, superseded, completions = 0; + new Observable(s => { source = s; }).switchMap(value => { + if (value === 2) return []; + if (mode === 'mapper') source.next(2); + return mode === 'mapper' ? new Observable(s => { superseded = s; }) : { + get [Symbol.asyncIterator]() { source.next(2); return undefined; }, + [Symbol.iterator]() { throw 'inactive conversion must not obtain an iterator'; } + }; + }).subscribe({complete: () => completions++}); + source.next(1); source.complete(); + return completions === 1 && (mode === 'conversion' || (!superseded.active && superseded.signal.aborted)); +})) +"#).unwrap(), "[true,true]"); +} + #[test] fn observable_flat_map_preserves_serial_order_conversion_reentrancy_and_cancellation() { let mut vm = new_storage_test_vm("https://observable-flat-map.test/"); diff --git a/moli-renderer-v8/src/worker/abort.rs b/moli-renderer-v8/src/worker/abort.rs index ee834aa4a..c1892c971 100644 --- a/moli-renderer-v8/src/worker/abort.rs +++ b/moli-renderer-v8/src/worker/abort.rs @@ -275,11 +275,29 @@ fn abort_worker_signal<'s>( } signals_to_abort }; + // Complete the dependency snapshot even if a source's IteratorClose + // throws. Dependents were already marked aborted and cannot be retried. + let mut first_error = None; for (signal_id, signal) in signals_to_abort { - if !run_worker_abort_steps(store, scope, signal, signal_id, reason) { - return; + let error = { + v8::tc_scope!(let scope, scope); + if run_worker_abort_steps(store, scope, signal, signal_id, reason) { + None + } else { + let Some(error) = scope.exception() else { + return; + }; + scope.reset(); + Some(error) + } + }; + if first_error.is_none() { + first_error = error; } } + if let Some(error) = first_error { + scope.throw_exception(error); + } } fn run_worker_abort_steps<'s>( diff --git a/moli-renderer-v8/src/worker/thread/tests/postmessage.rs b/moli-renderer-v8/src/worker/thread/tests/postmessage.rs index 76308641d..d3e6e6cdf 100644 --- a/moli-renderer-v8/src/worker/thread/tests/postmessage.rs +++ b/moli-renderer-v8/src/worker/thread/tests/postmessage.rs @@ -1,5 +1,22 @@ use super::*; +#[tokio::test] +async fn worker_observable_switch_map_preserves_switch_order_conversion_reentrancy_and_cancellation() + { + ensure_v8(); + let mut handle = spawn_worker( + format!( + "({}).then(value => {{ postMessage(value); close(); }});", + include_str!("../../../../tests/fixtures/observable-switch-map.js") + ), + "https://observable-switch-map.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() >= 170, "{result}"); +} + #[tokio::test] async fn worker_observable_flat_map_preserves_serial_order_conversion_reentrancy_and_cancellation() { diff --git a/moli-renderer-v8/tests/fixtures/observable-switch-map-realms.js b/moli-renderer-v8/tests/fixtures/observable-switch-map-realms.js new file mode 100644 index 000000000..e56c53151 --- /dev/null +++ b/moli-renderer-v8/tests/fixtures/observable-switch-map-realms.js @@ -0,0 +1,58 @@ +(async () => { + const failures = []; + let checks = 0; + const check = (value, label) => { checks++; if (!value) failures.push(label); }; + const child = document.querySelector('iframe').contentWindow; + const method = child.Observable.prototype.switchMap, source = Observable.from([1, 2]); + child.mapperCalls = []; + const callback = child.Function('value', 'index', 'mapperCalls.push([value,index]); globalThis.mapperThis = this; globalThis.mapperArgc = arguments.length; return [value*10];'); + 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.mapperCalls.length === 0, 'foreign mapper lazy'); + const values = await result.toArray(); + check(values instanceof child.Array && values.join(',') === '10,20', 'callee Array result'); + check(child.mapperThis === child && child.mapperArgc === 2, 'callback own realm and argument count'); + check(JSON.stringify(child.mapperCalls) === '[[1,0],[2,1]]', 'foreign mapper indices'); + const local = Observable.prototype.switchMap.call(child.Observable.from([3]), value => [value]); + 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 value of [null, 1, {}, {then() {}}]) { + const error = await method.call(source, () => value).toArray().catch(e => e); + check(error instanceof child.TypeError && !(error instanceof TypeError), 'mapper-result conversion error uses callee TypeError'); + } + const marker = new child.Error('mapper'); child.mapperError = marker; + for (const mapper of [child.Function('throw mapperError'), () => ({get [Symbol.iterator]() { throw marker; }}), + () => new child.Observable(s => s.error(marker)), () => child.Promise.reject(marker)]) { + const error = await method.call(source, mapper).toArray().catch(e => e); + check(error === marker && error instanceof child.Error, 'callback, conversion and inner errors retain identity'); + } + const frameSource = new child.Observable(s => { child.pendingOuter = s; }); + const inner = new child.Observable(s => { child.pendingInner = s; }); + const ac = new AbortController(), reason = {}, order = []; + const pending = Observable.prototype.switchMap.call(frameSource, () => inner).toArray({signal: ac.signal}); + const outcome = pending.catch(e => e); + child.pendingOuter.next(1); + child.pendingOuter.addTeardown(() => order.push('outer')); + child.pendingInner.addTeardown(() => order.push('inner')); + ac.abort(reason); + check(await outcome === reason, 'foreign pending graph preserves cancellation reason'); + check(!child.pendingOuter.active && !child.pendingInner.active, 'cancellation reaches both foreign producers'); + check(child.pendingOuter.signal.reason === reason && child.pendingInner.signal.reason === reason, 'foreign producer signals preserve reason'); + check(order.join(',') === 'outer,inner', 'foreign cleanup order'); + Object.setPrototypeOf(source, null); + const branded = method.call(source, value => [value]); + check(branded instanceof child.Observable && (await branded.toArray()).join(',') === '1,2', 'native brand survives prototype replacement'); + return {checks, failures}; +})() diff --git a/moli-renderer-v8/tests/fixtures/observable-switch-map.js b/moli-renderer-v8/tests/fixtures/observable-switch-map.js new file mode 100644 index 000000000..c02b76237 --- /dev/null +++ b/moli-renderer-v8/tests/fixtures/observable-switch-map.js @@ -0,0 +1,339 @@ +(async () => { + 'use strict'; + const failures = []; + let checks = 0; + const check = (value, label) => { checks++; if (!value) failures.push(label); }; + const same = (a, b, label) => check(JSON.stringify(a) === JSON.stringify(b), label); + const thrown = fn => { try { fn(); } catch (e) { return e; } }; + const test = async (label, fn) => { try { await fn(); } catch (e) { check(false, label + ': ' + e); } }; + const method = Observable.prototype.switchMap; + check(typeof method === 'function', 'switchMap exposed'); + if (failures.length) return {checks, failures}; + function subject() { + let subscriber, starts = 0; + return {source: new Observable(s => { subscriber = s; starts++; }), + get subscriber() { return subscriber; }, get starts() { return starts; }}; + } + + await test('Web IDL conversion and native operations', async () => { + const desc = Object.getOwnPropertyDescriptor(Observable.prototype, 'switchMap'); + check(method.name === 'switchMap' && method.length === 1, 'name and length'); + check(desc.enumerable && desc.writable && desc.configurable, 'descriptor'); + check(thrown(() => new method(() => [])) instanceof TypeError, 'not a constructor'); + const source = Observable.from([1]); + let traps = 0; + const revoked = Proxy.revocable(source, {}); revoked.revoke(); + 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 rejected'); + } + check(traps === 0, 'receiver check does not invoke Proxy traps'); + check(thrown(() => source.switchMap()) instanceof TypeError, 'required mapper'); + for (const callback of [undefined, null, false, 1, 1n, '', Symbol(), {}, [], {handleEvent() {}}]) { + check(thrown(() => source.switchMap(callback)) instanceof TypeError, 'non-callable mapper rejected'); + } + let calls = 0; + const mapper = new Proxy(function(value, index) { + calls++; + check(this === undefined && arguments.length === 2, 'mapper receiver and arguments'); + check(index === 0, 'fresh subscription resets index'); + return [value, value + 1]; + }, {get() { throw 'mapper property read'; }}); + Object.defineProperty(source, 'constructor', {get() { throw 'constructor read'; }}); + Object.setPrototypeOf(source, null); + const result = method.call(source, mapper, {get signal() { throw 'extra argument read'; }}); + check(calls === 0, 'creation is lazy'); + check(Object.getPrototypeOf(result) === Observable.prototype && Observable.from(result) === result, 'intrinsic branded result'); + same(await result.toArray(), [1, 2], 'source prototype and species ignored'); + same(await result.toArray(), [1, 2], 'reusable result'); + check(calls === 2, 'one mapper call per source value'); + class Derived extends Observable {} + check(!(new Derived(s => s.complete()).switchMap(mapper) instanceof Derived), 'does not use source species'); + const saved = []; + for (const [object, key] of [[Observable, 'from'], [Observable.prototype, 'subscribe'], + [AbortSignal, 'any'], [AbortController.prototype, 'abort'], [globalThis, 'AbortController'], + [Subscriber.prototype, 'next'], [Subscriber.prototype, 'error'], [Subscriber.prototype, 'complete']]) { + saved.push([object, key, Object.getOwnPropertyDescriptor(object, key)]); + Object.defineProperty(object, key, {value() { throw key + ' invoked'; }, configurable: true}); + } + try { same(await result.toArray(), [1, 2], 'internal conversion and notifications ignore public methods'); } + finally { for (const [object, key, descriptor] of saved) Object.defineProperty(object, key, descriptor); } + }); + + await test('conversion of mapper results', async () => { + for (const inner of [Observable.from([3]), [3], new Set([3]), Promise.resolve(3), + (async function* () { yield 3; })()]) { + same(await Observable.from([1]).switchMap(() => inner).toArray(), [3], 'maps convertible input'); + } + for (const value of [undefined, null, false, 1, 1n, '', 'abc', Symbol(), {}, {then() { throw 'assimilated'; }}]) { + check(await Observable.from([1]).switchMap(() => value).toArray().catch(e => e) instanceof TypeError, 'invalid result rejects through observer'); + } + const marker = {}, log = []; + const iterable = {get [Symbol.asyncIterator]() { log.push('async'); return undefined; }, + get [Symbol.iterator]() { log.push('sync'); return function* () { log.push('open'); yield marker; }; }}; + const values = await Observable.from([0]).switchMap(() => iterable).toArray(); + check(values.length === 1 && values[0] === marker, 'inner value identity'); + same(log, ['async', 'sync', 'sync', 'open'], 'conversion probes before subscription obtains iterator'); + let thenReads = 0; + const preferred = {[Symbol.iterator]: function* () { yield 5; }, get then() { thenReads++; throw 'then read'; }}; + same(await Observable.from([0]).switchMap(() => preferred).toArray(), [5], 'iterable result bypasses then property'); + check(thenReads === 0, 'does not assimilate arbitrary thenables'); + }); + + await test('switch cancellation and completion ordering', () => { + const outer = subject(), inners = [], order = [], values = [], indices = []; + const result = outer.source.switchMap((value, index) => { + indices.push(index); order.push('map' + value); + if (inners.length) check(!inners.at(-1).active, 'old inner closes before next mapper'); + return new Observable(s => { + inners.push(s); order.push('start' + value); + s.signal.addEventListener('abort', () => order.push('abort' + value)); + s.addTeardown(() => order.push('cleanup' + value)); s.next(value); + }); + }); + result.subscribe({next: v => values.push(v), complete: () => order.push('complete')}); + outer.subscriber.next(1); outer.subscriber.next(2); inners[0].next('stale'); inners[0].complete(); + same(indices, [0, 1], 'indices begin at zero per subscription'); + same(values, [1, 2], 'cancelled inner no longer emits'); + check(inners[0].signal.reason instanceof DOMException && inners[0].signal.reason.name === 'AbortError', 'switch reason defaults to AbortError'); + same(order, ['map1', 'start1', 'abort1', 'cleanup1', 'map2', 'start2'], 'old abort and cleanup precede mapper'); + outer.subscriber.complete(); + check(inners[1].active && !order.includes('complete'), 'outer completion waits for current inner'); + inners[1].next(3); inners[1].complete(); + same(values, [1, 2, 3], 'last inner survives outer completion'); + same(order.slice(-3), ['abort2', 'cleanup2', 'complete'], 'final inner cleanup precedes downstream complete'); + const empty = subject(); let completes = 0; + empty.source.switchMap(() => []).subscribe({complete: () => completes++}); + empty.subscriber.complete(); check(completes === 1, 'empty outer completes immediately'); + const sync = subject(); let syncCompletes = 0; + sync.source.switchMap(v => [v]).subscribe({complete: () => syncCompletes++}); + sync.subscriber.next(1); check(syncCompletes === 0, 'completed inner waits for outer'); + sync.subscriber.complete(); check(syncCompletes === 1, 'completed inner releases completion gate'); + }); + + await test('reentrant mapper and conversion use shared controller reference', () => { + for (const mode of ['mapper', 'conversion']) { + const outer = subject(), inners = {}, calls = [], values = []; + outer.source.switchMap((value, index) => { + calls.push([value, index]); + const inner = new Observable(s => { inners[value] = s; }); + if (value === 1 && mode === 'mapper') outer.subscriber.next(2); + if (value === 1 && mode === 'conversion') return { + get [Symbol.asyncIterator]() { outer.subscriber.next(2); return undefined; }, + [Symbol.iterator]() { return [value][Symbol.iterator](); } + }; + return inner; + }).subscribe({next: v => values.push(v), complete: () => values.push('complete')}); + outer.subscriber.next(1); + same(calls, mode === 'mapper' ? [[1, 0], [2, 0]] : [[1, 0], [2, 1]], mode + ' invocation/index update order'); + if (mode === 'mapper') { + check(inners[1].active && inners[2].active, 'mapper reentrancy retains both inner observers'); + inners[1].next('one'); inners[2].next('two'); + outer.subscriber.next(3); + check(!inners[1].active && !inners[2].active && inners[3].active, 'next switch cancels both observers sharing current controller'); + inners[3].complete(); + } else { + check(inners[2].active, 'conversion reentrancy keeps earlier inner active'); + inners[2].next('two'); + } + outer.subscriber.complete(); + same(values, mode === 'mapper' ? ['one', 'two', 'complete'] : [1, 'two', 'complete'], mode + ' output and completion'); + check(Object.values(inners).every(s => !s.active), mode + ' leaves no live inner after completion'); + } + }); + + await test('teardown and delivery reentrancy preserve all subscriptions', () => { + const outer = subject(), inners = {}, log = [], values = []; + outer.source.switchMap(v => new Observable(s => { + inners[v] = s; log.push('start' + v); + s.addTeardown(() => { log.push('close' + v); if (v === 1) outer.subscriber.next(3); }); + })).subscribe({next: v => values.push(v), complete: () => values.push('complete')}); + outer.subscriber.next(1); outer.subscriber.next(2); + same(log, ['start1', 'close1', 'start3', 'start2'], 'reentrant teardown initializes before interrupted switch resumes'); + check(inners[2].active && inners[3].active, 'both reentrant inners stay live'); + inners[2].next(2); inners[3].next(3); outer.subscriber.complete(); + check(inners[2].active && inners[3].active, 'source completion still waits'); + inners[2].complete(); + same(values, [2, 3, 'complete'], 'inner completion closes remaining reentrant inner'); + check(!inners[3].active, 'dependent signal cancels reentrant inner'); + const delivery = subject(), seen = [], created = []; + delivery.source.switchMap(v => new Observable(s => { + created.push(s); s.next(v); if (v === 1) s.next('stale'); + })).subscribe(v => { seen.push(v); if (v === 1) delivery.subscriber.next(2); }); + delivery.subscriber.next(1); + same(seen, [1, 2], 'delivery switches synchronously and suppresses later old notifications'); + check(!created[0].active && created[1].active, 'delivery cancels old inner before initializer resumes'); + delivery.subscriber.complete(); created[1].complete(); + }); + + await test('asynchronous switching ignores superseded resolutions', async () => { + const outer = subject(), resolvers = [], values = []; + const result = outer.source.switchMap(() => new Promise(resolve => resolvers.push(resolve))); + const promise = result.toArray(); + outer.subscriber.next(1); outer.subscriber.next(2); outer.subscriber.complete(); + resolvers[0]('old'); await Promise.resolve(); + promise.then(v => values.push(...v)); + check(values.length === 0, 'old promise does not complete result'); + resolvers[1]('new'); same(await promise, ['new'], 'only active Promise delivers'); + const asyncOuter = subject(), asyncLog = [], gates = []; + const asyncResult = asyncOuter.source.switchMap(v => ({[Symbol.asyncIterator]() { return { + next() { return new Promise(resolve => gates.push(resolve)); }, + return(reason) { asyncLog.push([v, reason.name]); return Promise.resolve({done: true}); } + }; }})); + const ac = new AbortController(); asyncResult.subscribe({}, {signal: ac.signal}); + asyncOuter.subscriber.next(1); asyncOuter.subscriber.next(2); + check(asyncLog.length === 1 && asyncLog[0][0] === 1, 'switch closes async iterator'); + gates[0]({value: 'old'}); await Promise.resolve(); await Promise.resolve(); + check(gates.length === 2, 'cancelled async iterator is not pulled again'); + ac.abort(); check(asyncLog.length === 2 && asyncLog[1][0] === 2, 'cancel closes latest async iterator'); + }); + + await test('switch only removes its own observer from a shared inner', () => { + const outer = subject(), shared = subject(), keep = new AbortController(), values = [], other = []; + shared.source.subscribe(v => other.push(v), {signal: keep.signal}); + outer.source.switchMap(v => v === 1 ? shared.source : [v]).subscribe(v => values.push(v)); + outer.subscriber.next(1); shared.subscriber.next('a'); outer.subscriber.next(2); + check(shared.subscriber.active, 'other consumer keeps replaced producer alive'); + shared.subscriber.next('b'); same(values, ['a', 2], 'switched observer detached'); same(other, ['a', 'b'], 'other observer unaffected'); + keep.abort(); check(!shared.subscriber.active, 'last independent consumer closes producer'); + outer.subscriber.complete(); + }); + + + await test('sharing, distinct branches and last-consumer cancellation', () => { + const outer = subject(), inner = subject(), ac1 = new AbortController(), ac2 = new AbortController(); + let maps = 0, otherMaps = 0; + const result = outer.source.switchMap(() => { maps++; return inner.source; }); + const first = [], second = []; + result.subscribe(v => first.push(v), {signal: ac1.signal}); + result.subscribe(v => second.push(v), {signal: ac2.signal}); + outer.source.switchMap(() => { otherMaps++; return []; }).subscribe(); + outer.subscriber.next(1); inner.subscriber.next(2); ac1.abort('first'); inner.subscriber.next(3); + check(maps === 1 && otherMaps === 1 && outer.starts === 1 && inner.starts === 1, 'shared result maps once and distinct branch separately'); + check(outer.subscriber.active && inner.subscriber.active, 'first cancellation keeps both subscriptions'); + same(first, [2], 'first consumer removed'); same(second, [2, 3], 'second consumer survives'); + const reason = {}; ac2.abort(reason); + check(outer.subscriber.active && !inner.subscriber.active && inner.subscriber.signal.reason === reason, 'last result consumer only cancels its branch'); + const oldInner = inner.subscriber; + result.subscribe(); outer.subscriber.next(2); + check(maps === 2 && otherMaps === 2 && inner.starts === 2 && inner.subscriber !== oldInner, 'resubscription uses fresh inner'); + inner.subscriber.complete(); outer.subscriber.complete(); + }); + + await test('original errors, inner disposal and synchronous cancellation', () => { + const reports = [], onerror = e => { reports.push(e.error); e.preventDefault(); }; + addEventListener('error', onerror); + try { + for (const mode of ['outer', 'inner', 'mapper', 'conversion', 'initializer']) { + for (const marker of [{}, null, undefined]) { + const outer = subject(), inner = subject(), log = [], errors = []; + let maps = 0, complete = 0; + outer.source.switchMap(() => { + maps++; + if (mode === 'mapper') throw marker; + if (mode === 'conversion') return {get [Symbol.asyncIterator]() { throw marker; }}; + if (mode === 'initializer') return new Observable(() => { throw marker; }); + return inner.source; + }).subscribe({error: e => { errors.push(e); log.push('error'); }, complete: () => complete++}); + outer.subscriber.addTeardown(() => log.push('outer')); + outer.subscriber.next(1); + if (inner.subscriber) { + inner.subscriber.addTeardown(() => log.push('inner')); + if (mode === 'outer') outer.subscriber.error(marker); else inner.subscriber.error(marker); + } + check(errors.length === 1 && errors[0] === marker && complete === 0, mode + ' error identity'); + check(!outer.subscriber.active && (!inner.subscriber || !inner.subscriber.active), mode + ' closes both subscriptions'); + check(maps === 1, mode + ' maps once'); + same(log, mode === 'outer' ? ['outer', 'inner', 'error'] : mode === 'inner' ? ['inner', 'outer', 'error'] : ['outer', 'error'], mode + ' cleanup precedes observer error'); + } + } + same(reports, [], 'handled errors are not reported globally'); + const outer = subject(), inner = subject(), ac = new AbortController(), reason = {}, cleanupError = {}; + let maps = 0; + outer.source.switchMap(() => { maps++; return inner.source; }).subscribe({}, {signal: ac.signal}); + outer.subscriber.next(1); + outer.subscriber.addTeardown(() => { throw cleanupError; }); + ac.abort(reason); + check(!outer.subscriber.active && !inner.subscriber.active && maps === 1, 'explicit cancellation closes both subscriptions'); + check(outer.subscriber.signal.reason === reason && inner.subscriber.signal.reason === reason, 'both receive same cancellation reason'); + check(reports.length === 1 && reports[0] === cleanupError, 'teardown exception does not skip inner cancellation'); + } finally { removeEventListener('error', onerror); } + }); + + await test('throwing IteratorClose still cancels both producers', () => { + for (const firstError of [{}, undefined]) { + const secondError = {}, ac = new AbortController(), log = [], reason = {}; + let outerClosed = 0, innerClosed = 0, innerPulls = 0, caught, didThrow = false; + const outer = {[Symbol.iterator]() { return { + next() { return {value: 1}; }, + return() { outerClosed++; log.push('outer return'); throw firstError; } + }; }}; + const inner = {[Symbol.iterator]() { return { + next() { return ++innerPulls > 3 ? {done: true} : {value: innerPulls}; }, + return() { innerClosed++; log.push('inner return'); throw secondError; } + }; }}; + Observable.from(outer).switchMap(() => inner).finally(() => log.push('finally')).subscribe(() => { + try { ac.abort(reason); } catch (e) { didThrow = true; caught = e; } + log.push('after abort'); + }, {signal: ac.signal}); + check(didThrow && caught === firstError, 'first IteratorClose failure preserved, including undefined'); + check(outerClosed === 1 && innerClosed === 1 && innerPulls === 1, 'both iterators close once before any extra pulls'); + same(log, ['outer return', 'inner return', 'finally', 'after abort'], 'cleanup and finalizer finish before abort rethrows'); + } + const ac = new AbortController(), marker = {}, first = subject(), second = subject(), log = []; + first.source.subscribe({}, {signal: ac.signal}); second.source.subscribe({}, {signal: ac.signal}); + Observable.from({[Symbol.iterator]() { return { + next() { return {value: 1}; }, return() { throw marker; } + }; }}).subscribe(() => { + const third = new Observable(s => s.addTeardown(() => log.push('third'))); + third.subscribe({}, {signal: ac.signal}); + check(thrown(() => ac.abort()) === marker, 'shared direct signal preserves failing cancellation'); + }, {signal: ac.signal}); + check(!first.subscriber.active && !second.subscriber.active, 'earlier consumers also closed'); + same(log, ['third'], 'later signal algorithm runs after failure without switchMap'); + }); + + await test('dependent cancellation finishes after source IteratorClose failure', () => { + for (const marker of [{}, undefined]) { + const ac = new AbortController(), log = []; + let inner, caught, didThrow = false, pulls = 0; + Observable.from({[Symbol.iterator]() { return { + next() { return ++pulls === 1 ? {value: 1} : {done: true}; }, + return() { log.push('outer'); throw marker; } + }; }}).subscribe(() => { + new Observable(s => { inner = s; s.addTeardown(() => log.push('inner')); }) + .subscribe({}, {signal: AbortSignal.any([ac.signal])}); + try { ac.abort('stop'); } catch (e) { didThrow = true; caught = e; } + }, {signal: ac.signal}); + check(didThrow && caught === marker, 'dependent dispatch preserves first exception including undefined'); + check(!inner.active && inner.signal.reason === 'stop', 'dependent observer is cancelled despite earlier failure'); + same(log, ['outer', 'inner'], 'dependent cleanup finishes before rethrow'); + } + }); + + await test('captured inner notification survives observer removal', () => { + const outer = subject(), inner = subject(), keep = new AbortController(), seen = []; + inner.source.subscribe(() => outer.subscriber.next(2), {signal: keep.signal}); + outer.source.switchMap(value => value === 1 ? inner.source : [value]).subscribe(v => seen.push(v)); + outer.subscriber.next(1); inner.subscriber.next('captured'); + same(seen, [2, 'captured'], 'already captured next steps still run after reentrant switch'); + keep.abort(); outer.subscriber.complete(); + }); + + await test('pre-abort and cancellation inside mapper', () => { + const ac = new AbortController(), reason = {}, pre = subject(); let maps = 0, inactive; + pre.source.switchMap(() => { maps++; return []; }).subscribe({}, {signal: AbortSignal.abort(reason)}); + check(pre.starts === 1 && !pre.subscriber.active && pre.subscriber.signal.reason === reason && maps === 0, 'pre-aborted source initialized inactive without mapping'); + const outer = subject(); + outer.source.switchMap(() => { ac.abort(reason); return new Observable(s => { inactive = s; }); }).subscribe({}, {signal: ac.signal}); + outer.subscriber.next(1); + check(!outer.subscriber.active && inactive && !inactive.active && inactive.signal.reason === reason, 'mapper cancellation still initializes returned Observable inactive'); + const captured = subject(), controller = new AbortController(); let calls = 0, inner; + captured.source.subscribe(() => controller.abort(reason)); + captured.source.switchMap(() => { calls++; return new Observable(s => { inner = s; }); }).subscribe({}, {signal: controller.signal}); + captured.subscriber.next(1); + check(calls === 1 && inner && !inner.active, 'captured next notification still maps after an earlier observer cancels'); + captured.subscriber.complete(); + }); + return {checks, failures}; +})()