diff --git a/moli-benchmark/wpt-cross-current/passed-cases.txt b/moli-benchmark/wpt-cross-current/passed-cases.txt index 4c2202281..e358f2dcf 100644 --- a/moli-benchmark/wpt-cross-current/passed-cases.txt +++ b/moli-benchmark/wpt-cross-current/passed-cases.txt @@ -4539,6 +4539,8 @@ dom/observable/tentative/observable-constructor.window.js?moli-wpt-script=window dom/observable/tentative/observable-event-target.any.js?moli-wpt-any=dedicatedworker dom/observable/tentative/observable-event-target.any.js?moli-wpt-any=window dom/observable/tentative/observable-event-target.window.js?moli-wpt-script=window +dom/observable/tentative/observable-first.any.js?moli-wpt-any=dedicatedworker +dom/observable/tentative/observable-first.any.js?moli-wpt-any=window dom/ranges/Range-adopt-test.html dom/ranges/Range-attributes.html dom/ranges/Range-cloneContents.html diff --git a/moli-renderer-v8/src/abort_signal_route.rs b/moli-renderer-v8/src/abort_signal_route.rs index 91150359b..3744bcb00 100644 --- a/moli-renderer-v8/src/abort_signal_route.rs +++ b/moli-renderer-v8/src/abort_signal_route.rs @@ -87,6 +87,18 @@ pub(crate) struct ResolvedAbortSignal<'s> { } impl<'s> ResolvedAbortSignal<'s> { + /// Uses the same dependency graph as AbortSignal.any, without calling a + /// mutable JavaScript static method or relaying through author event listeners. + pub(crate) fn dependent(scope: &mut v8::PinScope<'s, '_>, sources: &[Self]) -> Option { + let sources: Vec<_> = sources.iter().map(|source| source.signal).collect(); + let signal = if context_host_ptr_from_global_bridge(scope).is_some() { + crate::native_bridge::abort::new_dependent_abort_signal(scope, &sources)? + } else { + crate::worker::abort::new_worker_dependent_abort_signal(scope, &sources)? + }; + Self::resolve(scope, signal) + } + /// Creates a signal in the current realm without consulting author-visible /// constructors or maintaining another store for native algorithms. pub(crate) fn new(scope: &mut v8::PinScope<'s, '_>) -> Option { diff --git a/moli-renderer-v8/src/native_bridge/abort.rs b/moli-renderer-v8/src/native_bridge/abort.rs index de2eae9a1..fefeaa4b9 100644 --- a/moli-renderer-v8/src/native_bridge/abort.rs +++ b/moli-renderer-v8/src/native_bridge/abort.rs @@ -20,6 +20,7 @@ pub(crate) use signal::{ }; pub(crate) use statics::{ abort_signal_any_callback, abort_signal_static_abort_callback, abort_signal_timeout_callback, + new_dependent_abort_signal, }; const ABORT_SIGNAL_ID_SLOT: &str = "__lmAbortSignalId"; diff --git a/moli-renderer-v8/src/native_bridge/abort/statics.rs b/moli-renderer-v8/src/native_bridge/abort/statics.rs index fd33f60d8..255e6959e 100644 --- a/moli-renderer-v8/src/native_bridge/abort/statics.rs +++ b/moli-renderer-v8/src/native_bridge/abort/statics.rs @@ -74,23 +74,22 @@ pub(crate) fn abort_signal_any_callback<'s>( let Some(parsed) = webidl::parse_args::>(scope, &args) else { return; }; - let Some(host_ptr) = context_host_ptr_from_global_bridge(scope) else { + if let Some(signal) = new_dependent_abort_signal(scope, &parsed.signals) { + rv.set(signal.into()); + } else { rv.set_null(); - return; - }; + } +} + +pub(crate) fn new_dependent_abort_signal<'s>( + scope: &mut v8::PinScope<'s, '_>, + signals: &[v8::Local<'s, v8::Object>], +) -> Option> { + let host_ptr = context_host_ptr_from_global_bridge(scope)?; let host = unsafe { &mut *host_ptr }; - let Some(signal) = create_signal(scope, host, false, None) else { - rv.set_null(); - return; - }; - let Some(composite_signal_id) = AbortStore::signal_id_from_object(scope, signal) else { - rv.set_null(); - return; - }; - - let signals = parsed.signals; - - for source_signal in &signals { + let signal = create_signal(scope, host, false, None)?; + let composite_signal_id = AbortStore::signal_id_from_object(scope, signal)?; + for source_signal in signals { let Some(source_signal_id) = AbortStore::signal_id_from_object(scope, *source_signal) else { continue; @@ -106,18 +105,16 @@ pub(crate) fn abort_signal_any_callback<'s>( continue; }; host.abort_signal(scope, signal, reason); - rv.set(signal.into()); - return; + return Some(signal); } host.native_bridge_mut().abort.set_signal_sources( composite_signal_id, signals - .into_iter() - .filter_map(|signal| AbortStore::signal_id_from_object(scope, signal)), + .iter() + .filter_map(|signal| AbortStore::signal_id_from_object(scope, *signal)), ); - - rv.set(signal.into()); + Some(signal) } fn abort_signal_timeout_fire_native_callback( diff --git a/moli-renderer-v8/src/observable.rs b/moli-renderer-v8/src/observable.rs index 6149ab707..c5b300f8d 100644 --- a/moli-renderer-v8/src/observable.rs +++ b/moli-renderer-v8/src/observable.rs @@ -5,7 +5,9 @@ mod callbacks; mod event_target; +mod first; mod from; +mod observer; mod state; pub(crate) use event_target::event_target_when; @@ -23,6 +25,8 @@ use state::*; struct ObservablePrototype { #[webapi(method, length = 0, callback = subscribe)] subscribe: (), + #[webapi(method, length = 0, returns_promise, callback = first::first)] + first: (), } #[derive(WebApiFunctionTemplate)] @@ -189,7 +193,15 @@ fn subscribe<'s>( let Some(parsed) = webidl::parse_args::>(scope, &args) else { return; }; - let observable = args.this(); + subscribe_internal(scope, args.this(), parsed.observer, parsed.signal); +} + +fn subscribe_internal<'s>( + scope: &mut v8::PinScope<'s, '_>, + observable: v8::Local<'s, v8::Object>, + observer: v8::Local<'s, v8::Object>, + signal: Option>, +) { if !is_current(scope, observable) { return; } @@ -206,9 +218,9 @@ fn subscribe<'s>( } }; let mut observers = list(scope, subscriber, OBSERVERS); - observers.push(parsed.observer); + observers.push(observer); set_list(scope, subscriber, OBSERVERS, &observers); - if let Some(signal) = parsed.signal { + if let Some(signal) = signal { if signal.is_aborted(scope) { if fresh { let reason = signal.reason(scope); @@ -220,15 +232,20 @@ fn subscribe<'s>( set_list(scope, subscriber, OBSERVERS, &observers); } } else { - let data = - v8::Array::new_with_elements(scope, &[subscriber.into(), parsed.observer.into()]); + let data = v8::Array::new_with_elements(scope, &[subscriber.into(), observer.into()]); let algorithm = v8::Function::builder(cancel_observer) .data(data.into()) .build(scope) .expect("Observable abort algorithm should allocate"); - set_private_value(scope, parsed.observer, INPUT_SIGNAL, signal.value().into()); - set_private_value(scope, parsed.observer, ABORT_ALGORITHM, algorithm.into()); - signal.register_rethrowing_algorithm(scope, algorithm); + set_private_value(scope, observer, INPUT_SIGNAL, signal.value().into()); + set_private_value(scope, observer, ABORT_ALGORITHM, algorithm.into()); + if observer::is_native(scope, observer) { + // Native observers trace their private cancellation callback. + // The internal signal must not root an abandoned subscription. + signal.register_weak_rethrowing_algorithm(scope, algorithm); + } else { + signal.register_rethrowing_algorithm(scope, algorithm); + } } } if fresh { @@ -354,9 +371,7 @@ fn subscriber_next<'s>( } // Reentrant subscribe/cancel must not change this notification's snapshot. for observer in list(scope, subscriber, OBSERVERS) { - if let Some(callback) = object_slot(scope, observer, NEXT) { - invoke_and_report(scope, callback, &[value]); - } + observer::notify(scope, observer, observer::Notification::Next(value)); } } @@ -378,11 +393,7 @@ fn subscriber_error<'s>( let observers = list(scope, subscriber, OBSERVERS); set_list(scope, subscriber, OBSERVERS, &[]); for observer in observers { - if let Some(callback) = object_slot(scope, observer, ERROR) { - invoke_and_report(scope, callback, &[error]); - } else { - callbacks::report_default_error(scope, error); - } + observer::notify(scope, observer, observer::Notification::Error(error)); } } @@ -418,9 +429,7 @@ fn subscriber_complete<'s>( let observers = list(scope, subscriber, OBSERVERS); set_list(scope, subscriber, OBSERVERS, &[]); for observer in observers { - if let Some(callback) = object_slot(scope, observer, COMPLETE) { - invoke_and_report(scope, callback, &[]); - } + observer::notify(scope, observer, observer::Notification::Complete); } } diff --git a/moli-renderer-v8/src/observable/first.rs b/moli-renderer-v8/src/observable/first.rs new file mode 100644 index 000000000..fba653870 --- /dev/null +++ b/moli-renderer-v8/src/observable/first.rs @@ -0,0 +1,187 @@ +use moli_webapi_declare::WebApiObject; + +use super::{ + callbacks, + observer::{self, NATIVE_KIND, Notification}, + signal_arg, + state::object_slot, + subscribe_internal, +}; +use crate::{ + abort_signal_route::ResolvedAbortSignal, + util::{get_private_value, set_private_value, v8str}, + webidl, +}; + +const RESOLVER: &str = "__moliObservableFirstResolver"; +const CONTROLLER_SIGNAL: &str = "__moliObservableFirstControllerSignal"; +const REJECTION_SIGNAL: &str = "__moliObservableFirstRejectionSignal"; +const REJECTION_ALGORITHM: &str = "__moliObservableFirstRejectionAlgorithm"; +const SETTLED: &str = "__moliObservableFirstSettled"; +const PROMISE_OBSERVER: &str = "__moliObservablePromiseObserver"; + +#[derive(webidl::WebIdlArgs)] +#[webidl(prefix = "Observable.first")] +struct FirstArgs<'scope> { + #[webidl(with = signal_arg)] + signal: Option>, +} + +#[derive(WebApiObject)] +#[webapi(plain)] +struct FirstObserver<'scope> { + #[webapi(slot = NATIVE_KIND)] + kind: i32, + #[webapi(slot = RESOLVER)] + resolver: v8::Local<'scope, v8::Object>, + #[webapi(slot = CONTROLLER_SIGNAL)] + controller_signal: v8::Local<'scope, v8::Object>, + #[webapi(slot = REJECTION_SIGNAL)] + rejection_signal: v8::Local<'scope, v8::Object>, + #[webapi(slot = SETTLED)] + settled: bool, +} + +pub(super) fn first<'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(resolver) = v8::PromiseResolver::new(scope) else { + return; + }; + let promise = resolver.get_promise(scope); + rv.set(promise.into()); + let Some(controller) = ResolvedAbortSignal::new(scope) else { + return; + }; + let mut sources = vec![controller]; + sources.extend(parsed.signal); + let Some(signal) = ResolvedAbortSignal::dependent(scope, &sources) else { + return; + }; + if signal.is_aborted(scope) { + let reason = signal.reason(scope); + resolver.reject(scope, reason); + return; + } + let observer = FirstObserver::new( + observer::FIRST, + resolver.into(), + controller.value(), + signal.value(), + false, + ) + .bind(scope) + .expect("Observable.first observer"); + // A reachable pending Promise retains its observer. Without an external + // signal or producer, this cycle remains entirely V8-traced and collectible. + set_private_value(scope, promise.into(), PROMISE_OBSERVER, observer.into()); + let algorithm = v8::Function::builder(aborted) + .data(observer.into()) + .build(scope) + .expect("Observable.first abort algorithm"); + set_private_value(scope, observer, REJECTION_ALGORITHM, algorithm.into()); + if parsed.signal.is_some() { + // As with an ordinary subscribe({signal}), caller-driven cancellation + // remains an owner while its observer is pending. + signal.register_algorithm(scope, algorithm); + } else { + signal.register_weak_rethrowing_algorithm(scope, algorithm); + } + subscribe_internal(scope, args.this(), observer, Some(signal)); +} + +fn start_settlement<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, +) -> Option> { + if get_private_value(scope, observer, SETTLED).is_some_and(|value| value.is_true()) { + return None; + } + // Lock the result before Resolve can run a then getter and reenter the + // producer or abort its signal. Promise::state can still be Pending while + // assimilating the first value, so it is not an already-resolved flag. + set_private_value( + scope, + observer, + SETTLED, + v8::Boolean::new(scope, true).into(), + ); + if let Some(algorithm) = object_slot(scope, observer, REJECTION_ALGORITHM) + .and_then(|value| v8::Local::::try_from(value).ok()) + { + let signal = object_slot(scope, observer, REJECTION_SIGNAL) + .and_then(|signal| ResolvedAbortSignal::resolve(scope, signal))?; + signal.unregister_algorithm(scope, algorithm); + set_private_value( + scope, + observer, + REJECTION_ALGORITHM, + v8::undefined(scope).into(), + ); + } + let resolver = object_slot(scope, observer, RESOLVER).expect("Observable.first resolver"); + // SAFETY: this private slot is populated only with PromiseResolver::new; + // it is never read from a public property or supplied by script. + let resolver = unsafe { v8::Local::::cast_unchecked(resolver) }; + let promise = resolver.get_promise(scope); + set_private_value( + scope, + promise.into(), + PROMISE_OBSERVER, + v8::undefined(scope).into(), + ); + Some(resolver) +} + +pub(super) fn notify<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, + notification: Notification<'s>, +) { + match notification { + Notification::Next(value) => { + if let Some(resolver) = start_settlement(scope, observer) { + resolver.resolve(scope, value); + } + // Abort after resolving, including when a then getter reenters + // next(). This removes just this observer from a shared producer. + if let Some(signal) = object_slot(scope, observer, CONTROLLER_SIGNAL) + .and_then(|signal| ResolvedAbortSignal::resolve(scope, signal)) + { + let reason = crate::native_bridge::abort::abort_error_value(scope); + signal.abort(scope, reason); + } + } + Notification::Error(error) => { + if let Some(resolver) = start_settlement(scope, observer) { + resolver.reject(scope, error); + } + } + Notification::Complete => { + if let Some(resolver) = start_settlement(scope, observer) { + let error = + v8::Exception::range_error(scope, v8str(scope, "No values in Observable")); + resolver.reject(scope, error); + } + } + } +} + +fn aborted<'s>( + scope: &mut v8::PinScope<'s, '_>, + args: v8::FunctionCallbackArguments<'s>, + _rv: v8::ReturnValue<'_, v8::Value>, +) { + let observer = + v8::Local::::try_from(args.data()).expect("Observable.first abort data"); + if callbacks::is_current(scope, observer) + && let Some(resolver) = start_settlement(scope, observer) + { + resolver.reject(scope, args.get(0)); + } +} diff --git a/moli-renderer-v8/src/observable/observer.rs b/moli-renderer-v8/src/observable/observer.rs new file mode 100644 index 000000000..7e13a5aa1 --- /dev/null +++ b/moli-renderer-v8/src/observable/observer.rs @@ -0,0 +1,77 @@ +//! Internal observer steps share notification ordering with script observers, +//! while script callbacks keep their typed Web IDL invocation boundary. + +use super::{callbacks, first, invoke_and_report, state::*}; +use crate::util::get_private_value; + +pub(super) const NATIVE_KIND: &str = "__moliObservableNativeObserver"; +pub(super) const FIRST: i32 = 1; + +#[derive(Clone, Copy)] +pub(super) enum Notification<'s> { + Next(v8::Local<'s, v8::Value>), + Error(v8::Local<'s, v8::Value>), + Complete, +} + +pub(super) fn is_native<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, +) -> bool { + get_private_value(scope, observer, NATIVE_KIND).is_some_and(|value| value.is_int32()) +} + +pub(super) fn notify<'s>( + scope: &mut v8::PinScope<'s, '_>, + observer: v8::Local<'s, v8::Object>, + notification: Notification<'s>, +) { + if let Some(kind) = get_private_value(scope, observer, NATIVE_KIND) + .filter(|value| value.is_int32()) + .and_then(|value| value.int32_value(scope)) + { + if !callbacks::is_current(scope, observer) { + return; + } + let Some(context) = observer.get_creation_context(scope) else { + return; + }; + let scope = &mut v8::ContextScope::new(scope, context); + let exception = { + v8::tc_scope!(let scope, scope); + match kind { + FIRST => first::notify(scope, observer, notification), + _ => unreachable!("unknown native Observable observer"), + } + let exception = scope.exception(); + scope.reset(); + exception + }; + // Internal observer steps cannot throw through Subscriber.next/error/ + // complete. In particular, first() has already resolved its Promise + // before a throwing iterator return() is encountered during cancellation. + if let Some(exception) = exception { + callbacks::report(scope, observer, exception); + } + return; + } + match notification { + Notification::Next(value) => { + if let Some(callback) = object_slot(scope, observer, NEXT) { + invoke_and_report(scope, callback, &[value]); + } + } + Notification::Error(error) => { + if let Some(callback) = object_slot(scope, observer, ERROR) { + invoke_and_report(scope, callback, &[error]); + } else { + callbacks::report_default_error(scope, error); + } + } + Notification::Complete => { + if let Some(callback) = object_slot(scope, observer, COMPLETE) { + invoke_and_report(scope, callback, &[]); + } + } + } +} diff --git a/moli-renderer-v8/src/script_vm/tests/observable.rs b/moli-renderer-v8/src/script_vm/tests/observable.rs index 1c03369fe..1d0ab11b0 100644 --- a/moli-renderer-v8/src/script_vm/tests/observable.rs +++ b/moli-renderer-v8/src/script_vm/tests/observable.rs @@ -1,5 +1,141 @@ use super::*; +#[test] +fn observable_first_promises_cancellation_reentrancy_and_native_observers() { + let mut vm = new_storage_test_vm("https://observable-first.test/"); + vm.eval(&format!( + "({}).then(value => {{ globalThis.firstResult = JSON.stringify(value); }});", + include_str!("../../../tests/fixtures/observable-first.js") + )) + .expect("Observable.first fixture should evaluate"); + let result = vm + .eval("firstResult") + .expect("Observable.first fixture should settle"); + let result: serde_json::Value = serde_json::from_str(&result).unwrap(); + assert_eq!(result["failures"], serde_json::json!([]), "{result}"); + assert!(result["checks"].as_u64().unwrap() >= 72, "{result}"); +} + +#[test] +fn observable_first_pending_promises_trace_observers_without_rooting_abandoned_cycles() { + let mut vm = new_storage_test_vm("https://observable-first-gc.test/"); + vm.eval( + r#" +(() => { + const captured = {}; + const source = new Observable(s => { + globalThis.weakFirstSubscriber = new WeakRef(s); + s.addTeardown(() => captured); + }); + globalThis.weakFirstCapture = new WeakRef(captured); + globalThis.weakFirstSource = new WeakRef(source); + globalThis.weakFirstPromise = new WeakRef(source.first()); +})(); +(() => { + const iterator = {next: () => new Promise(() => {})}; + globalThis.weakFirstIterator = new WeakRef(iterator); + Observable.from({[Symbol.asyncIterator]: () => iterator}).first(); +})(); +(() => { + const source = new Observable(s => { + globalThis.weakKeptSubscriber = new WeakRef(s); + globalThis.deliverFirst = s.next.bind(s); + }); + globalThis.weakKeptSource = new WeakRef(source); + globalThis.keptFirstPromise = source.first(); +})(); +(() => { + const ac = new AbortController(); + globalThis.cancelFirst = ac.abort.bind(ac); + new Observable(s => { globalThis.weakAbortFirstSubscriber = new WeakRef(s); }) + .first({signal: ac.signal}).catch(reason => { globalThis.firstAbortReason = reason; }); +})(); +"#, + ) + .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([ +weakFirstSubscriber.deref() === undefined, weakFirstCapture.deref() === undefined, +weakFirstSource.deref() === undefined, weakFirstPromise.deref() === undefined, +weakFirstIterator.deref() === undefined, weakKeptSource.deref() === undefined, +weakKeptSubscriber.deref() !== undefined, weakAbortFirstSubscriber.deref() !== undefined +])"# + ) + .unwrap(), + "[true,true,true,true,true,true,true,true]" + ); + vm.eval("delete globalThis.deliverFirst;").unwrap(); + collect(&mut vm); + assert_eq!( + vm.eval("weakKeptSubscriber.deref() !== undefined").unwrap(), + "true" + ); + vm.eval( + r#" +keptFirstPromise.then(value => { globalThis.firstDelivered = value; }); +weakKeptSubscriber.deref().next(31); +cancelFirst('cancelled'); +delete globalThis.keptFirstPromise; +delete globalThis.cancelFirst; +"#, + ) + .unwrap(); + assert_eq!( + vm.eval("JSON.stringify([firstDelivered, firstAbortReason])") + .unwrap(), + "[31,\"cancelled\"]" + ); + collect(&mut vm); + assert_eq!(vm.eval("JSON.stringify([weakKeptSubscriber.deref() === undefined, weakAbortFirstSubscriber.deref() === undefined])").unwrap(), "[true,true]"); +} + +#[test] +fn observable_first_uses_callee_promise_and_error_realms_with_foreign_sources() { + let mut vm = new_storage_test_vm("https://observable-first-realms.test/"); + vm.eval("document.appendChild(document.createElement('iframe'))") + .unwrap(); + materialize_single_child_default_realm_for_test(&mut vm, "Observable.first realm"); + vm.eval(r#" +(async () => { + const child = document.querySelector('iframe').contentWindow, checks = []; + const first = child.Observable.prototype.first; + const source = new Observable(s => s.next(5)); + const promise = first.call(source); + checks.push(promise instanceof child.Promise, !(promise instanceof Promise), await promise === 5); + let conversions = 0; + const invalid = first.call({}, {get signal() { conversions++; }}); + checks.push(invalid instanceof child.Promise); + await invalid.catch(e => checks.push(e instanceof child.TypeError, !(e instanceof TypeError))); + checks.push(conversions === 0); + await first.call(new Observable(s => s.complete())).catch(e => checks.push(e instanceof child.RangeError, !(e instanceof RangeError))); + const foreign = new child.Observable(s => s.next(6)); + const local = Observable.prototype.first.call(foreign); + checks.push(local instanceof Promise, !(local instanceof child.Promise), await local === 6); + const marker = {}, ac = new AbortController(); + const pending = first.call(new Observable(() => {}), {signal: ac.signal}); + ac.abort(marker); + await pending.catch(e => checks.push(e === marker)); + globalThis.firstRealms = JSON.stringify(checks); +})(); +"#).unwrap(); + assert_eq!( + vm.eval("firstRealms").unwrap(), + "[true,true,true,true,true,true,true,true,true,true,true,true,true]" + ); +} + #[test] fn observable_from_iterables_promises_cancellation_and_exception_timing() { let mut vm = new_storage_test_vm("https://observable-from.test/"); diff --git a/moli-renderer-v8/src/worker/abort.rs b/moli-renderer-v8/src/worker/abort.rs index 7cbc9728e..3a82fd9bb 100644 --- a/moli-renderer-v8/src/worker/abort.rs +++ b/moli-renderer-v8/src/worker/abort.rs @@ -595,20 +595,21 @@ pub(crate) fn worker_abort_signal_any_callback<'s>( let Some(parsed) = webidl::parse_args::>(scope, &args) else { return; }; - let Some(store) = worker_abort_store(scope) else { + if let Some(signal) = new_worker_dependent_abort_signal(scope, &parsed.signals) { + rv.set(signal.into()); + } else { rv.set_null(); - return; - }; - let Some(signal) = create_signal(scope, &mut store.borrow_mut(), false, None) else { - rv.set_null(); - return; - }; - let Some(composite_signal_id) = WorkerAbortStore::signal_id_from_object(scope, signal) else { - rv.set_null(); - return; - }; - let signals = parsed.signals; - for source_signal in &signals { + } +} + +pub(crate) fn new_worker_dependent_abort_signal<'s>( + scope: &mut v8::PinScope<'s, '_>, + signals: &[v8::Local<'s, v8::Object>], +) -> Option> { + let store = worker_abort_store(scope)?; + let signal = create_signal(scope, &mut store.borrow_mut(), false, None)?; + let composite_signal_id = WorkerAbortStore::signal_id_from_object(scope, signal)?; + for source_signal in signals { let Some(source_signal_id) = WorkerAbortStore::signal_id_from_object(scope, *source_signal) else { continue; @@ -623,16 +624,15 @@ pub(crate) fn worker_abort_signal_any_callback<'s>( continue; }; abort_worker_signal(&store, scope, signal, reason); - rv.set(signal.into()); - return; + return Some(signal); } store.borrow_mut().set_signal_sources( composite_signal_id, signals - .into_iter() - .filter_map(|signal| WorkerAbortStore::signal_id_from_object(scope, signal)), + .iter() + .filter_map(|signal| WorkerAbortStore::signal_id_from_object(scope, *signal)), ); - rv.set(signal.into()); + Some(signal) } pub(crate) fn worker_abort_signal_aborted_getter_function<'s>( diff --git a/moli-renderer-v8/src/worker/thread/tests/postmessage.rs b/moli-renderer-v8/src/worker/thread/tests/postmessage.rs index f1459ea80..a0281c041 100644 --- a/moli-renderer-v8/src/worker/thread/tests/postmessage.rs +++ b/moli-renderer-v8/src/worker/thread/tests/postmessage.rs @@ -1,5 +1,21 @@ use super::*; +#[tokio::test] +async fn worker_observable_first_promises_cancellation_reentrancy_and_native_observers() { + ensure_v8(); + let mut handle = spawn_worker( + format!( + "({}).then(value => {{ postMessage(value); close(); }});", + include_str!("../../../../tests/fixtures/observable-first.js") + ), + "https://observable-first.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() >= 72, "{result}"); +} + #[tokio::test] async fn worker_observable_from_iterables_promises_cancellation_and_exception_timing() { ensure_v8(); diff --git a/moli-renderer-v8/tests/fixtures/observable-first.js b/moli-renderer-v8/tests/fixtures/observable-first.js new file mode 100644 index 000000000..fa98477b9 --- /dev/null +++ b/moli-renderer-v8/tests/fixtures/observable-first.js @@ -0,0 +1,188 @@ +(async () => { + 'use strict'; + const failures = []; + let checks = 0; + const check = (value, label) => { checks++; if (!value) failures.push(label); }; + const same = (actual, expected, label) => check(JSON.stringify(actual) === JSON.stringify(expected), label); + const thrown = fn => { try { fn(); } catch (error) { return error; } }; + const rejected = promise => promise.then(() => { throw new Error('expected rejection'); }, error => error); + const test = async (label, fn) => { try { await fn(); } catch (error) { check(false, label + ': ' + error); } }; + const first = Observable.prototype.first; + if (typeof first !== 'function') { + check(false, 'Observable.first is exposed'); + return {checks, failures}; + } + + await test('receiver and conversion', async () => { + const descriptor = Object.getOwnPropertyDescriptor(Observable.prototype, 'first'); + check(first.name === 'first' && first.length === 0, 'name and length'); + check(descriptor.enumerable && descriptor.writable && descriptor.configurable, 'method descriptor'); + check(thrown(() => new first()) instanceof TypeError, 'not constructible'); + const source = new Observable(subscriber => subscriber.next(7)); + let reads = 0, traps = 0; + const options = {get signal() { reads++; return undefined; }}; + const revoked = Proxy.revocable(source, {}); revoked.revoke(); + const forged = [undefined, null, false, 1, 'x', Symbol(), {}, Object.create(Observable.prototype), + Object.create(source), new Proxy(source, {get() { traps++; throw 1; }}), revoked.proxy]; + for (const receiver of forged) { + const promise = first.call(receiver, options); + check(promise instanceof Promise && await rejected(promise) instanceof TypeError, 'invalid receiver rejects a Promise'); + } + check(reads === 0 && traps === 0, 'brand check precedes option conversion and does not inspect proxies'); + Object.setPrototypeOf(source, null); + check(await first.call(source) === 7, 'native brand survives prototype replacement'); + for (const value of [1, true, 'x', Symbol(), 1n]) { + check(await rejected(first.call(source, value)) instanceof TypeError, 'non-object options reject'); + } + for (const value of [null, undefined, {}, {signal: undefined}]) { + check(await first.call(source, value) === 7, 'empty dictionary accepted'); + } + const ac = new AbortController(); + for (const signal of [null, {}, Object.create(AbortSignal.prototype), new Proxy(ac.signal, {})]) { + check(await rejected(first.call(source, {signal})) instanceof TypeError, 'invalid signal rejects'); + } + const marker = {}; + check(await rejected(first.call(source, {get signal() { throw marker; }})) === marker, 'option getter rejection identity'); + check(await first.call(source, options) === 7 && reads === 1, 'signal read exactly once'); + }); + + await test('lifecycle', async () => { + const log = [], value = {}; + let subscriber; + const source = new Observable(s => { + subscriber = s; + s.signal.addEventListener('abort', () => log.push('abort')); + s.addTeardown(() => log.push('teardown')); + log.push('before'); s.next(value); + log.push(s.active ? 'active' : 'inactive'); + s.next('ignored'); s.complete(); + }); + const promise = source.first().then(result => { log.push('resolved'); return result; }); + same(log, ['before', 'abort', 'teardown', 'inactive'], 'synchronous cancellation order'); + check(subscriber.signal.aborted && subscriber.signal.reason.name === 'AbortError', 'upstream receives default abort reason'); + check(await promise === value, 'first value identity'); + same(log, ['before', 'abort', 'teardown', 'inactive', 'resolved'], 'Promise reaction runs after teardown'); + check(await rejected(new Observable(s => s.complete()).first()) instanceof RangeError, 'empty source rejects with RangeError'); + const marker = {}; + check(await rejected(new Observable(s => s.error(marker)).first()) === marker, 'source error identity'); + check(await rejected(new Observable(() => { throw marker; }).first()) === marker, 'initializer exception identity'); + }); + + await test('abort and sharing', async () => { + const marker = {}, ac = new AbortController(), log = []; + let subscriber, starts = 0; + const source = new Observable(s => { + starts++; subscriber = s; + s.addTeardown(() => log.push('teardown')); + }); + const aborted = AbortSignal.abort(marker); + check(await rejected(source.first({signal: aborted})) === marker && starts === 0, 'pre-abort skips initializer'); + const promise = source.first({signal: ac.signal}); + const rejection = rejected(promise); + ac.abort(marker); + check(!subscriber.active && subscriber.signal.reason === marker, 'input abort closes upstream with same reason'); + same(log, ['teardown'], 'input abort runs teardown once'); + check(await rejection === marker, 'input abort rejects with reason identity'); + const values = [], shared = new AbortController(); + source.subscribe(value => values.push(value), {signal: shared.signal}); + check(await rejected(source.first({signal: aborted})) === marker && subscriber.active, 'pre-aborted first leaves shared producer active'); + const firstValue = source.first(); + subscriber.next(1); + check(subscriber.active && starts === 2, 'first only removes its observer from a shared subscription'); + subscriber.next(2); + check(await firstValue === 1, 'shared first value'); + same(values, [1, 2], 'other observer continues receiving'); + shared.abort(); + check(!subscriber.active && log.length === 2, 'last observer cancellation closes shared producer'); + for (const terminal of ['complete', 'error']) { + const duringTeardown = new AbortController(), reason = {}; + const source = new Observable(s => { + s.addTeardown(() => duringTeardown.abort(reason)); + s[terminal]('source error'); + }); + check(await rejected(source.first({signal: duringTeardown.signal})) === reason, 'teardown abort wins before ' + terminal + ' notification'); + } + }); + + await test('thenable reentrancy', async () => { + for (const reenter of ['next', 'complete', 'abort']) { + const ac = new AbortController(), log = [], marker = {}; + let subscriber, thenReads = 0, secondReads = 0; + const source = new Observable(s => { subscriber = s; s.addTeardown(() => log.push('teardown')); }); + const promise = source.first({signal: ac.signal}); + const value = {get then() { + thenReads++; log.push('get then'); + if (reenter === 'next') subscriber.next({get then() { secondReads++; }}); + if (reenter === 'complete') subscriber.complete(); + if (reenter === 'abort') ac.abort(marker); + check(!subscriber.active, 'reentrant ' + reenter + ' cancels synchronously'); + return resolve => { log.push('then'); resolve(42); }; + }}; + subscriber.next(value); + check(thenReads === 1 && secondReads === 0, 'first resolve locks before then getter reentrancy'); + check(await promise === 42, 'thenable result survives reentrant ' + reenter); + same(log, ['get then', 'teardown', 'then'], 'thenable cancellation order for ' + reenter); + } + let subscriber, resolveValue; + const ac = new AbortController(), value = new Promise(resolve => { resolveValue = resolve; }); + const promise = new Observable(s => { subscriber = s; }).first({signal: ac.signal}); + subscriber.next(value); + ac.abort('too late'); resolveValue(9); + check(await promise === 9, 'late input abort cannot replace pending assimilation'); + const marker = {}; + let closed = false; + const throwing = new Observable(s => { + s.addTeardown(() => { closed = true; }); + s.next({get then() { throw marker; }}); + }); + check(await rejected(throwing.first()) === marker && closed, 'throwing then getter still cancels source'); + }); + + await test('dependent signal ordering', async () => { + const ac = new AbortController(), reason = {}, value = {}; + let subscriber; + const promise = new Observable(s => { subscriber = s; }).first({signal: ac.signal}); + ac.signal.addEventListener('abort', () => subscriber.next(value)); + ac.abort(reason); + check(await promise === value, 'source abort event precedes dependent signal algorithms'); + check(!subscriber.active, 'dependent abort still removes observer'); + }); + + await test('iterable cancellation', async () => { + for (const symbol of [Symbol.iterator, Symbol.asyncIterator]) { + let pulls = 0, returns = 0; + const iterator = {next() { pulls++; return {value: 8}; }, return() { returns++; return {}; }}; + check(await Observable.from({[symbol]: () => iterator}).first() === 8, 'first from iterable ' + String(symbol)); + check(pulls === 1 && returns === 1, 'one pull and one close ' + String(symbol)); + } + const marker = {}, errors = []; + const onerror = event => { if (event.error === marker) { errors.push(event.error); event.preventDefault(); } }; + addEventListener('error', onerror); + try { + const iterator = {next: () => ({value: 5}), return() { throw marker; }}; + check(await Observable.from({[Symbol.iterator]: () => iterator}).first() === 5, 'close exception cannot replace resolved value'); + check(errors.length === 1 && errors[0] === marker, 'close exception reported globally once'); + } finally { removeEventListener('error', onerror); } + }); + + await test('intrinsics and internal subscription', async () => { + const source = Observable.from([17]); + const globals = ['Observable', 'Promise', 'AbortController']; + const saved = globals.map(name => globalThis[name]); + const methods = [[Observable.prototype, 'subscribe'], [AbortSignal, 'any'], + [Subscriber.prototype, 'next'], [Subscriber.prototype, 'error'], [Subscriber.prototype, 'complete']]; + const descriptors = methods.map(([object, name]) => Object.getOwnPropertyDescriptor(object, name)); + const poison = () => { throw new Error('public implementation consulted'); }; + let promise; + try { + globals.forEach(name => { globalThis[name] = poison; }); + methods.forEach(([object, name]) => { object[name] = poison; }); + promise = first.call(source); + } finally { + globals.forEach((name, index) => { globalThis[name] = saved[index]; }); + methods.forEach(([object, name], index) => Object.defineProperty(object, name, descriptors[index])); + } + check(promise instanceof Promise && await promise === 17, 'native first ignores replaced public constructors and methods'); + }); + return {checks, failures}; +})()