diff --git a/moli-benchmark/wpt-cross-current/passed-cases.txt b/moli-benchmark/wpt-cross-current/passed-cases.txt index b19d8f00d..e89c86477 100644 --- a/moli-benchmark/wpt-cross-current/passed-cases.txt +++ b/moli-benchmark/wpt-cross-current/passed-cases.txt @@ -4553,6 +4553,8 @@ 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 dom/observable/tentative/observable-first.any.js?moli-wpt-any=window +dom/observable/tentative/observable-flatMap.any.js?moli-wpt-any=dedicatedworker +dom/observable/tentative/observable-flatMap.any.js?moli-wpt-any=window dom/observable/tentative/observable-forEach.any.js?moli-wpt-any=dedicatedworker dom/observable/tentative/observable-forEach.any.js?moli-wpt-any=window dom/observable/tentative/observable-forEach.window.js?moli-wpt-script=window diff --git a/moli-renderer-v8/src/abort_signal_route.rs b/moli-renderer-v8/src/abort_signal_route.rs index 32b08cf92..22e7e25a9 100644 --- a/moli-renderer-v8/src/abort_signal_route.rs +++ b/moli-renderer-v8/src/abort_signal_route.rs @@ -46,7 +46,7 @@ impl AbortAlgorithm { /// Most internal abort algorithms cannot throw. Observable iterator closing /// is an exception: its synchronous return() failure escapes AbortController.abort. /// Keep this policy on the native callback, in the existing signal-owned list. -pub(crate) fn invoke_abort_algorithm<'s>( +fn invoke_abort_algorithm<'s>( scope: &mut v8::PinScope<'s, '_>, label: &str, algorithm: v8::Local<'s, v8::Function>, @@ -69,6 +69,46 @@ pub(crate) fn invoke_abort_algorithm<'s>( } } +/// A multi-source Observable shares a signal between its producers. A failing +/// IteratorClose must not leave later producers active. Finish the cancellation +/// snapshot before propagating its first exception to the abort caller. +pub(crate) fn invoke_abort_algorithms<'s>( + scope: &mut v8::PinScope<'s, '_>, + label: &str, + signal: v8::Local<'s, v8::Object>, + reason: v8::Local<'s, v8::Value>, + algorithms: Vec, +) -> bool { + let mut first_error = None; + for algorithm in algorithms { + let Some(algorithm) = algorithm.prepare(scope) else { + continue; + }; + let exception = { + v8::tc_scope!(let scope, scope); + if invoke_abort_algorithm(scope, label, algorithm, signal, reason) { + None + } else { + let Some(error) = scope.exception() else { + // Do not resume script execution after V8 termination. + return false; + }; + scope.reset(); + Some(error) + } + }; + if first_error.is_none() { + first_error = exception; + } + } + if let Some(error) = first_error { + scope.throw_exception(error); + false + } else { + true + } +} + #[derive(Clone, Copy)] enum AbortSignalOwner { Window, diff --git a/moli-renderer-v8/src/native_bridge/abort/event.rs b/moli-renderer-v8/src/native_bridge/abort/event.rs index 5f56a2c16..ab27d807c 100644 --- a/moli-renderer-v8/src/native_bridge/abort/event.rs +++ b/moli-renderer-v8/src/native_bridge/abort/event.rs @@ -1,4 +1,4 @@ -use crate::abort_signal_route::{AbortAlgorithm, invoke_abort_algorithm}; +use crate::abort_signal_route::AbortAlgorithm; pub(super) fn invoke_abort_algorithms<'s>( scope: &mut v8::PinScope<'s, '_>, @@ -7,21 +7,13 @@ pub(super) fn invoke_abort_algorithms<'s>( abort_algorithms: Vec, ) -> bool { let signal = local_object_in_scope(scope, signal); - for algorithm in abort_algorithms { - let Some(algorithm) = algorithm.prepare(scope) else { - continue; - }; - if !invoke_abort_algorithm( - scope, - "AbortSignal abort algorithm", - algorithm, - signal, - reason, - ) { - return false; - } - } - true + crate::abort_signal_route::invoke_abort_algorithms( + scope, + "AbortSignal abort algorithm", + signal, + reason, + abort_algorithms, + ) } pub(super) fn local_object_in_scope<'s>( diff --git a/moli-renderer-v8/src/observable.rs b/moli-renderer-v8/src/observable.rs index 975d8994f..834fdc0ce 100644 --- a/moli-renderer-v8/src/observable.rs +++ b/moli-renderer-v8/src/observable.rs @@ -9,6 +9,7 @@ mod consume; mod event_target; mod finally; mod first; +mod flat_map; mod from; mod inspect; mod observer; @@ -34,6 +35,8 @@ struct ObservablePrototype { subscribe: (), #[webapi(method, length = 1, callback = transform::map)] map: (), + #[webapi(method = "flatMap", length = 1, callback = flat_map::flat_map)] + flat_map: (), #[webapi(method, length = 1, callback = transform::filter)] filter: (), #[webapi(method, length = 1, callback = transform::take)] @@ -318,6 +321,7 @@ fn subscribe_internal<'s>( && !until::subscribe(scope, observable, subscriber) && !inspect::subscribe(scope, observable, subscriber) && !finally::subscribe(scope, observable, subscriber) + && !flat_map::subscribe(scope, observable, subscriber) { event_target::subscribe(scope, observable, subscriber); } diff --git a/moli-renderer-v8/src/observable/flat_map.rs b/moli-renderer-v8/src/observable/flat_map.rs new file mode 100644 index 000000000..0e69d07b1 --- /dev/null +++ b/moli-renderer-v8/src/observable/flat_map.rs @@ -0,0 +1,262 @@ +//! Serial flattening keeps raw source values queued until the active inner +//! subscription completes. Queue cells and both observers are V8-traced. + +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 = "__moliObservableFlatMapSource"; +const MAPPER: &str = "__moliObservableFlatMapMapper"; +const DOWNSTREAM: &str = "__moliFlatMapSubscriber"; +const OWNER: &str = "__moliFlatMapSourceObserver"; +const INNER: &str = "__moliFlatMapInnerObserver"; +const BUSY: &str = "__moliFlatMapActiveInner"; +const OUTER_COMPLETE: &str = "__moliFlatMapSourceCompleted"; +const HEAD: &str = "__moliFlatMapQueueHead"; +const TAIL: &str = "__moliFlatMapQueueTail"; +const VALUE: &str = "__moliFlatMapQueuedValue"; +const NEXT: &str = "__moliFlatMapQueueNext"; + +#[derive(webidl::WebIdlArgs)] +#[webidl(prefix = "Observable.flatMap")] +struct FlatMapArgs { + #[webidl(required, converter = "callback_function")] + mapper: webidl::WebIdlCallbackFunction, +} + +pub(super) fn flat_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>, +} + +#[derive(WebApiObject)] +#[webapi(plain)] +struct QueueEntry<'scope> { + #[webapi(slot = VALUE)] + value: v8::Local<'scope, v8::Value>, +} + +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("flatMap mapper"); + let Some(observer) = SourceObserver::new(observer::FLAT_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(scope, subscriber); + subscribe_internal(scope, source, observer, Some(signal)); + true +} + +fn signal<'s>( + scope: &mut v8::PinScope<'s, '_>, + subscriber: v8::Local<'s, v8::Object>, +) -> ResolvedAbortSignal<'s> { + object_slot(scope, subscriber, SIGNAL) + .and_then(|signal| ResolvedAbortSignal::resolve(scope, signal)) + .expect("flatMap Subscriber signal") +} + +fn flag<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, + slot: &str, +) -> bool { + get_private_value(scope, observer, slot).is_some_and(|value| value.is_true()) +} + +fn enqueue<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, + value: v8::Local<'s, v8::Value>, +) { + let entry = QueueEntry::new(value) + .bind(scope) + .expect("flatMap queue entry"); + if let Some(tail) = object_slot(scope, observer, TAIL) { + set_private_value(scope, tail, NEXT, entry.into()); + } else { + set_private_value(scope, observer, HEAD, entry.into()); + } + set_private_value(scope, observer, TAIL, entry.into()); +} + +fn dequeue<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, +) -> Option> { + let head = object_slot(scope, observer, HEAD)?; + let value = get_private_value(scope, head, VALUE).expect("flatMap queued value"); + let next = object_slot(scope, head, NEXT); + let next_value = next.map_or_else(|| v8::undefined(scope).into(), Into::into); + set_private_value(scope, observer, HEAD, next_value); + if next.is_none() { + set_private_value(scope, observer, TAIL, v8::undefined(scope).into()); + } + Some(value) +} + +fn process_next<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, + value: v8::Local<'s, v8::Value>, +) { + let downstream = object_slot(scope, observer, DOWNSTREAM).expect("flatMap downstream"); + let mapper = object_slot(scope, observer, MAPPER).expect("flatMap 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; + }; + let Some(inner_observer) = InnerObserver::new(observer::FLAT_MAP_INNER, observer) + .bind(scope) + .ok() + else { + return; + }; + set_private_value(scope, observer, INNER, inner_observer.into()); + // Conversion may cancel downstream. The source still receives its fresh, + // inactive Subscriber, just as an explicitly pre-aborted subscribe does. + let signal = signal(scope, downstream); + 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_observer = if kind == observer::FLAT_MAP_INNER { + object_slot(scope, observer, OWNER).expect("flatMap source observer") + } else { + observer + }; + let downstream = object_slot(scope, source_observer, DOWNSTREAM).expect("flatMap downstream"); + match notification { + Notification::Next(value) if kind == observer::FLAT_MAP_SOURCE => { + if flag(scope, source_observer, BUSY) { + enqueue(scope, source_observer, value); + } else { + // Set before invoking mapper: reentrant source.next queues its + // raw value and does not overlap the current mapper or inner. + set_private_value( + scope, + source_observer, + BUSY, + v8::Boolean::new(scope, true).into(), + ); + process_next(scope, source_observer, value); + } + } + Notification::Next(value) => subscriber_next(scope, downstream, value), + Notification::Error(error) => subscriber_error(scope, downstream, error), + Notification::Complete => { + set_private_value( + scope, + observer, + observer::SUBSCRIBER, + v8::undefined(scope).into(), + ); + if kind == observer::FLAT_MAP_SOURCE { + set_private_value( + scope, + source_observer, + OUTER_COMPLETE, + v8::Boolean::new(scope, true).into(), + ); + if !flag(scope, source_observer, BUSY) { + subscriber_complete(scope, downstream); + } + } else { + set_private_value(scope, source_observer, INNER, v8::undefined(scope).into()); + if let Some(value) = dequeue(scope, source_observer) { + // Run before this complete() returns, including when the + // next inner completes synchronously and reenters here. + process_next(scope, source_observer, value); + } else { + set_private_value( + scope, + source_observer, + BUSY, + v8::Boolean::new(scope, false).into(), + ); + if flag(scope, source_observer, OUTER_COMPLETE) { + subscriber_complete(scope, downstream); + } + } + } + } + } +} diff --git a/moli-renderer-v8/src/observable/observer.rs b/moli-renderer-v8/src/observable/observer.rs index 6bcca5d99..0dd2bbf63 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, inspect, invoke_and_report, state::*, transform, - until, + callbacks, collect, consume, finally, first, flat_map, inspect, invoke_and_report, state::*, + transform, until, }; use crate::util::{get_private_value, set_private_value}; @@ -24,6 +24,8 @@ pub(super) const UNTIL_SOURCE: i32 = 13; pub(super) const UNTIL_NOTIFIER: i32 = 14; 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 SUBSCRIBER: &str = "__moliObservableNativeSubscriber"; const INDEX: &str = "__moliObservableCallbackIndex"; @@ -93,6 +95,9 @@ pub(super) fn notify<'s>( UNTIL_SOURCE | UNTIL_NOTIFIER => until::notify(scope, observer, notification, kind), INSPECT => inspect::notify(scope, observer, notification), FINALLY => finally::notify(scope, observer, notification), + FLAT_MAP_SOURCE | FLAT_MAP_INNER => { + flat_map::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 6bcbe0448..4a6e33b3d 100644 --- a/moli-renderer-v8/src/script_vm/tests/observable.rs +++ b/moli-renderer-v8/src/script_vm/tests/observable.rs @@ -1,5 +1,109 @@ use super::*; +#[test] +fn observable_flat_map_preserves_serial_order_conversion_reentrancy_and_cancellation() { + let mut vm = new_storage_test_vm("https://observable-flat-map.test/"); + vm.eval(&format!( + "({}).then(value => {{ globalThis.flatMapResult = JSON.stringify(value); }});", + include_str!("../../../tests/fixtures/observable-flat-map.js") + )) + .expect("Observable.flatMap fixture should evaluate"); + let result: serde_json::Value = + serde_json::from_str(&vm.eval("flatMapResult").unwrap()).unwrap(); + assert_eq!(result["failures"], serde_json::json!([]), "{result}"); + assert!(result["checks"].as_u64().unwrap() >= 120, "{result}"); +} + +#[test] +fn observable_flat_map_preserves_result_conversion_callback_and_cancellation_realms() { + let mut vm = new_storage_test_vm("https://observable-flat-map-realms.test/"); + vm.eval("document.appendChild(document.createElement('iframe'))") + .unwrap(); + materialize_single_child_default_realm_for_test(&mut vm, "Observable.flatMap realm"); + vm.eval(&format!( + "({}).then(value => {{ globalThis.flatMapRealms = JSON.stringify(value); }});", + include_str!("../../../tests/fixtures/observable-flat-map-realms.js") + )) + .expect("Observable.flatMap realms fixture should evaluate"); + let result: serde_json::Value = + serde_json::from_str(&vm.eval("flatMapRealms").unwrap()).unwrap(); + assert_eq!(result["failures"], serde_json::json!([]), "{result}"); + assert_eq!(result["checks"], 31, "{result}"); +} + +#[test] +fn observable_flat_map_traces_both_producers_and_releases_queues_and_closed_graphs() { + let mut vm = new_storage_test_vm("https://observable-flat-map-gc.test/"); + vm.eval(r#" +function makeFlatMapSource() { + let subscriber; + return {source: new Observable(s => { subscriber = s; }), get subscriber() { return subscriber; }}; +} +function makeFlatMapper(inner) { + const token = {values: []}; + return {token, callback: value => value === 0 ? inner : token.values}; +} +globalThis.flatMapChains = []; +for (const mode of ['abandoned', 'complete', 'outer-error', 'inner-error', 'abort']) (() => { + const outer = makeFlatMapSource(), inner = makeFlatMapSource(), mapper = makeFlatMapper(inner.source); + const result = outer.source.flatMap(mapper.callback), ac = new AbortController(), queued = [{}, {}]; + const promise = result.toArray(mode === 'abort' ? {signal: ac.signal} : undefined); + promise.catch(() => {}); + outer.subscriber.next(0); inner.subscriber.next(1); + for (const value of queued) outer.subscriber.next(value); + const entry = {mode, templates: [outer.source, result].map(v => new WeakRef(v)), + outer: new WeakRef(outer.subscriber), inner: new WeakRef(inner.subscriber), + callbacks: [mapper.callback, mapper.token].map(v => new WeakRef(v)), queued: queued.map(v => new WeakRef(v))}; + if (mode !== 'abandoned') entry.promise = promise; + if (mode === 'abort') entry.controller = ac; + flatMapChains.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([ +flatMapChains.every(c => c.templates.every(ref => ref.deref() === undefined)), +flatMapChains.every(c => [c.outer, c.inner, ...c.callbacks, ...c.queued].every(ref => (ref.deref() !== undefined) === (c.mode !== 'abandoned'))) +])"#).unwrap(), "[true,true]"); + vm.eval(r#" +for (const c of flatMapChains.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(); inner.next(2); + 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("flatMapChains.filter(c => c.promise).every(c => c.correct && c.closed.every(s => !s.active))").unwrap(), "true"); + collect(&mut vm); + assert_eq!(vm.eval("flatMapChains.every(c => [...c.callbacks, ...c.queued].every(ref => ref.deref() === undefined))").unwrap(), "true"); + vm.eval(r#" +globalThis.flatMapDrain = (() => { + const outer = makeFlatMapSource(), inner = makeFlatMapSource(), mapper = makeFlatMapper(inner.source), value = {}; + const promise = outer.source.flatMap(mapper.callback).toArray(); + outer.subscriber.next(0); outer.subscriber.next(value); outer.subscriber.complete(); + return {promise, outer: new WeakRef(outer.subscriber), inner: new WeakRef(inner.subscriber), + value: new WeakRef(value), callback: new WeakRef(mapper.callback)}; +})(); +"#).unwrap(); + collect(&mut vm); + assert_eq!(vm.eval("flatMapDrain.outer.deref() === undefined && flatMapDrain.inner.deref() !== undefined && flatMapDrain.value.deref() !== undefined && flatMapDrain.callback.deref() !== undefined").unwrap(), "true"); + vm.eval("flatMapDrain.promise.then(v => { flatMapDrain.correct = v.length === 0; }); flatMapDrain.inner.deref().complete();").unwrap(); + collect(&mut vm); + assert_eq!(vm.eval("flatMapDrain.correct && flatMapDrain.inner.deref() === undefined && flatMapDrain.value.deref() === undefined && flatMapDrain.callback.deref() === undefined").unwrap(), "true"); +} + #[test] fn observable_finally_preserves_teardown_order_sharing_and_reentrant_cancellation() { let mut vm = new_storage_test_vm("https://observable-finally.test/"); diff --git a/moli-renderer-v8/src/worker/abort.rs b/moli-renderer-v8/src/worker/abort.rs index 3a82fd9bb..ee834aa4a 100644 --- a/moli-renderer-v8/src/worker/abort.rs +++ b/moli-renderer-v8/src/worker/abort.rs @@ -1,7 +1,7 @@ use std::collections::{HashMap, HashSet}; use std::{cell::RefCell, rc::Rc}; -use crate::abort_signal_route::{AbortAlgorithm, invoke_abort_algorithm}; +use crate::abort_signal_route::{AbortAlgorithm, invoke_abort_algorithms}; use crate::util::{get_private_value, set_private_value, v8str}; use crate::webidl; @@ -297,7 +297,13 @@ fn run_worker_abort_steps<'s>( std::mem::take(&mut state.abort_algorithms) }; reject_worker_fetches_for_signal(scope, signal_id, reason); - if !invoke_worker_abort_algorithms(scope, signal, reason, abort_algorithms) { + if !invoke_abort_algorithms( + scope, + "Worker AbortSignal abort algorithm", + signal, + reason, + abort_algorithms, + ) { return false; } // Dispatch reads the shared listener registry after all abort algorithms. @@ -305,29 +311,6 @@ fn run_worker_abort_steps<'s>( true } -fn invoke_worker_abort_algorithms<'s>( - scope: &mut v8::PinScope<'s, '_>, - signal: v8::Local<'s, v8::Object>, - reason: v8::Local<'s, v8::Value>, - abort_algorithms: Vec, -) -> bool { - for algorithm in abort_algorithms { - let Some(algorithm) = algorithm.prepare(scope) else { - continue; - }; - if !invoke_abort_algorithm( - scope, - "Worker AbortSignal abort algorithm", - algorithm, - signal, - reason, - ) { - return false; - } - } - true -} - fn create_signal<'s>( scope: &mut v8::PinScope<'s, '_>, store: &mut WorkerAbortStore, diff --git a/moli-renderer-v8/src/worker/thread/tests/postmessage.rs b/moli-renderer-v8/src/worker/thread/tests/postmessage.rs index 255f2f1bf..76308641d 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_flat_map_preserves_serial_order_conversion_reentrancy_and_cancellation() +{ + ensure_v8(); + let mut handle = spawn_worker( + format!( + "({}).then(value => {{ postMessage(value); close(); }});", + include_str!("../../../../tests/fixtures/observable-flat-map.js") + ), + "https://observable-flat-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() >= 120, "{result}"); +} + #[tokio::test] async fn worker_observable_finally_preserves_teardown_order_sharing_and_reentrant_cancellation() { ensure_v8(); diff --git a/moli-renderer-v8/tests/fixtures/observable-flat-map-realms.js b/moli-renderer-v8/tests/fixtures/observable-flat-map-realms.js new file mode 100644 index 000000000..238060ccc --- /dev/null +++ b/moli-renderer-v8/tests/fixtures/observable-flat-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.flatMap, 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.flatMap.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.flatMap.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-flat-map.js b/moli-renderer-v8/tests/fixtures/observable-flat-map.js new file mode 100644 index 000000000..5c651fce2 --- /dev/null +++ b/moli-renderer-v8/tests/fixtures/observable-flat-map.js @@ -0,0 +1,240 @@ +(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.flatMap; + check(typeof method === 'function', 'flatMap 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, 'flatMap'); + check(method.name === 'flatMap' && 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.flatMap()) instanceof TypeError, 'required mapper'); + for (const callback of [undefined, null, false, 1, 1n, '', Symbol(), {}, [], {handleEvent() {}}]) { + check(thrown(() => source.flatMap(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()).flatMap(mapper) instanceof Derived), 'does not use source species'); + const saved = []; + for (const [object, key] of [[Observable, 'from'], [Observable.prototype, 'subscribe'], + [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]).flatMap(() => 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]).flatMap(() => 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]).flatMap(() => 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]).flatMap(() => preferred).toArray(), [5], 'iterable result bypasses then property'); + check(thenReads === 0, 'does not assimilate arbitrary thenables'); + }); + + await test('raw queue, delayed mapper and synchronous completion order', () => { + const outer = subject(), inners = [], log = [], indices = []; + const result = outer.source.flatMap((value, index) => { + const id = value.id; + indices.push(index); log.push('map' + id); + return new Observable(s => { + inners.push(s); log.push('start' + id); s.addTeardown(() => log.push('cleanup' + id)); s.next(id); + if (id > 1) { s.complete(); log.push('after' + id); } + }); + }); + const values = []; + result.subscribe({next: v => values.push(v), complete: () => log.push('complete')}); + const queued = {id: 20}; + outer.subscriber.next({id: 1}); outer.subscriber.next(queued); outer.subscriber.next({id: 3}); + same(log, ['map1', 'start1'], 'queued values do not invoke mapper early'); + queued.id = 2; outer.subscriber.complete(); + check(inners.length === 1 && inners[0].active, 'outer completion waits for active inner'); + inners[0].complete(); log.push('after1'); + same(values, [1, 2, 3], 'queue retains raw value identity'); + same(indices, [0, 1, 2], 'serial mapper indices'); + same(log, ['map1', 'start1', 'cleanup1', 'map2', 'start2', 'cleanup2', 'map3', 'start3', 'cleanup3', + 'complete', 'after3', 'after2', 'after1'], 'drains next inner before previous complete returns'); + }); + + await test('reentrant mapper, conversion and inner delivery', async () => { + const outer = subject(), mapped = [], values = []; + outer.source.flatMap((value, index) => { + mapped.push([value, index]); + if (value === 1) outer.subscriber.next(2); + return {[Symbol.iterator]() { + if (value === 1) outer.subscriber.next(3); + return [value][Symbol.iterator](); + }}; + }).subscribe(value => { values.push(value); if (value === 1) outer.subscriber.next(4); }); + outer.subscriber.next(1); outer.subscriber.complete(); + same(mapped, [[1, 0], [2, 1], [3, 2], [4, 3]], 'reentrant pushes wait until mapper and inner finish'); + same(values, [1, 2, 3, 4], 'reentrant output order'); + const pending = subject(), gate = subject(), queued = []; + const result = pending.source.flatMap(value => { queued.push(value); return value === 0 ? gate.source : [value]; }); + const promise = result.toArray(); + for (let i = 0; i < 128; i++) pending.subscriber.next(i); + pending.subscriber.complete(); + same(queued, [0], 'large queue maps lazily'); + gate.subscriber.next(0); gate.subscriber.complete(); + const collected = await promise; + check(collected.length === 128 && collected.every((v, i) => v === i), 'drains buffered synchronous inputs in FIFO order'); + }); + + 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.flatMap(() => { 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.flatMap(() => { 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 queue and inner'); + inner.subscriber.complete(); outer.subscriber.complete(); + }); + + await test('original errors, queue 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.flatMap(() => { + 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')); + outer.subscriber.next(2); + 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 + ' abandons queued values'); + 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.flatMap(() => { maps++; return inner.source; }).subscribe({}, {signal: ac.signal}); + outer.subscriber.next(1); outer.subscriber.next(2); + outer.subscriber.addTeardown(() => { throw cleanupError; }); + ac.abort(reason); + check(!outer.subscriber.active && !inner.subscriber.active && maps === 1, 'explicit cancellation closes both and drops queue'); + 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).flatMap(() => 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 flatMap'); + }); + + await test('pre-abort and cancellation inside mapper', () => { + const ac = new AbortController(), reason = {}, pre = subject(); let maps = 0, inactive; + pre.source.flatMap(() => { 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.flatMap(() => { 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.flatMap(() => { 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}; +})()