feat(dom): resolve first Observable values with native observers

Add Observable.first in Window and workers through the shared WebIDL
receiver and Promise adapters. Native observers use internal subscription
and AbortSignal dependency graphs, preserving shared producers and
reentrant Promise resolution without consulting public JS methods.

Cover thenable reentrancy, abort/teardown order, iterable closing,
callee realms, and observer reachability. Related WPTs improve from
290/304 to 302/304; EventTarget regressions remain 150/150.
This commit is contained in:
ldm0
2026-09-24 20:01:11 +08:00
parent 006df888d9
commit 999fcb83c3
11 changed files with 683 additions and 58 deletions
@@ -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
@@ -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<Self> {
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<Self> {
@@ -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";
@@ -74,23 +74,22 @@ pub(crate) fn abort_signal_any_callback<'s>(
let Some(parsed) = webidl::parse_args::<abort_signal::AnyArgs<'s>>(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<v8::Local<'s, v8::Object>> {
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(
+28 -19
View File
@@ -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::<SubscribeArgs<'s>>(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<ResolvedAbortSignal<'s>>,
) {
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);
}
}
+187
View File
@@ -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<ResolvedAbortSignal<'scope>>,
}
#[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::<FirstArgs<'s>>(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<v8::Local<'s, v8::PromiseResolver>> {
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::<v8::Function>::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::<v8::PromiseResolver>::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::<v8::Object>::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));
}
}
@@ -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, &[]);
}
}
}
}
@@ -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/");
+18 -18
View File
@@ -595,20 +595,21 @@ pub(crate) fn worker_abort_signal_any_callback<'s>(
let Some(parsed) = webidl::parse_args::<abort_signal::AnyArgs<'s>>(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<v8::Local<'s, v8::Object>> {
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>(
@@ -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();
+188
View File
@@ -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};
})()