feat(dom): consume Observables with forEach and reduce

Implement native callback-driven Promise consumers for Window and workers,
using traced WebIDL callbacks and shared cancellation/settlement helpers.
Preserve exception identity, raw reducer values, reentrant notification
ordering, callback realms, and independent subscriptions.

Add conversion, cancellation, cross-realm, GC, and worker coverage. Record
five newly passing WPT cases: the audited suites improve from 330/372 to
364/372 subtests, with EventTarget regressions remaining 150/150.

Validation: cargo fmt --all; workspace all-targets/all-features clippy with
-D warnings; CLI build; 156 focused tests; full nextest (19,283 passed,
13 skipped). The supplemental fixture passes 132 checks in both globals.
This commit is contained in:
ldm0
2026-09-23 00:14:08 +08:00
parent c894362a2c
commit df54ffc84f
10 changed files with 616 additions and 41 deletions
@@ -4541,8 +4541,13 @@ 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/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
dom/observable/tentative/observable-last.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-last.any.js?moli-wpt-any=window
dom/observable/tentative/observable-reduce.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-reduce.any.js?moli-wpt-any=window
dom/observable/tentative/observable-toArray.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-toArray.any.js?moli-wpt-any=window
dom/ranges/Range-adopt-test.html
+5
View File
@@ -5,6 +5,7 @@
mod callbacks;
mod collect;
mod consume;
mod event_target;
mod first;
mod from;
@@ -33,6 +34,10 @@ struct ObservablePrototype {
last: (),
#[webapi(method = "toArray", length = 0, returns_promise, callback = collect::to_array)]
to_array: (),
#[webapi(method = "forEach", length = 1, returns_promise, callback = consume::for_each)]
for_each: (),
#[webapi(method, length = 1, returns_promise, callback = consume::reduce)]
reduce: (),
}
#[derive(WebApiFunctionTemplate)]
+19 -3
View File
@@ -42,17 +42,33 @@ pub(super) fn invoke<'s>(
carrier: v8::Local<'s, v8::Object>,
arguments: &[v8::Local<'s, v8::Value>],
) -> Option<v8::Local<'s, v8::Value>> {
invoke_value(scope, carrier, arguments).err()
}
/// The same typed boundary, retaining the raw result for any-returning
/// callbacks such as reducers. This does not assimilate returned thenables.
pub(super) fn invoke_value<'s>(
scope: &mut v8::PinScope<'s, '_>,
carrier: v8::Local<'s, v8::Object>,
arguments: &[v8::Local<'s, v8::Value>],
) -> Result<v8::Local<'s, v8::Value>, v8::Local<'s, v8::Value>> {
let callback = V8TracedWebIdlCallbackFunction::from_object(carrier).prepare(scope);
let context = callback.relevant_context(scope);
if !context_is_current(scope, context) {
return None;
return Ok(v8::undefined(scope).into());
}
v8::tc_scope!(let scope, scope);
let receiver = v8::undefined(scope).into();
invoke_synchronous_webidl_callback_function(scope, &callback, receiver, arguments);
let result = invoke_synchronous_webidl_callback_function(scope, &callback, receiver, arguments);
let exception = scope.exception();
scope.reset();
exception
if let Some(exception) = exception {
Err(exception)
} else {
Ok(result
.map(|value| v8::Local::new(scope, value))
.unwrap_or_else(|| v8::undefined(scope).into()))
}
}
pub(super) fn report_callback_exception<'s>(
+212
View File
@@ -0,0 +1,212 @@
//! Callback-driven Promise consumers use traced Web IDL callbacks and cancel
//! only their own observer when a visitor or reducer throws.
use super::{
callbacks,
observer::{self, Notification},
promise, signal_arg,
state::{object_slot, set_callback},
subscribe_internal,
};
use crate::{
abort_signal_route::ResolvedAbortSignal,
util::{get_private_value, set_private_value, v8str},
webidl,
};
const CALLBACK: &str = "__moliObservableConsumerCallback";
const INDEX: &str = "__moliObservableConsumerIndex";
const HAS_ACCUMULATOR: &str = "__moliObservableHasAccumulator";
// Unlike terminal collection state, a reducer's accumulator and callback remain
// observable through a next() snapshot already being dispatched at cancellation.
// The observer retains them until that snapshot is released, without Rust roots.
const ACCUMULATOR: &str = "__moliObservableAccumulator";
#[derive(webidl::WebIdlArgs)]
#[webidl(prefix = "Observable.forEach")]
struct ForEachArgs<'scope> {
#[webidl(required, converter = "callback_function")]
callback: webidl::WebIdlCallbackFunction,
#[webidl(with = signal_arg)]
signal: Option<ResolvedAbortSignal<'scope>>,
}
#[derive(webidl::WebIdlArgs)]
#[webidl(prefix = "Observable.reduce")]
struct ReduceArgs<'scope> {
#[webidl(required, converter = "callback_function")]
callback: webidl::WebIdlCallbackFunction,
#[webidl(converter = "raw")]
initial: Option<v8::Local<'scope, v8::Value>>,
#[webidl(with = signal_arg)]
signal: Option<ResolvedAbortSignal<'scope>>,
}
pub(super) fn for_each<'s>(
scope: &mut v8::PinScope<'s, '_>,
args: v8::FunctionCallbackArguments<'s>,
mut rv: v8::ReturnValue<'_, v8::Value>,
) {
let Some(parsed) = webidl::parse_args::<ForEachArgs<'s>>(scope, &args) else {
return;
};
if let Some(promise) = consume(
scope,
args.this(),
parsed.callback,
None,
parsed.signal,
observer::FOR_EACH,
) {
rv.set(promise.into());
}
}
pub(super) fn reduce<'s>(
scope: &mut v8::PinScope<'s, '_>,
args: v8::FunctionCallbackArguments<'s>,
mut rv: v8::ReturnValue<'_, v8::Value>,
) {
let Some(parsed) = webidl::parse_args::<ReduceArgs<'s>>(scope, &args) else {
return;
};
// Web IDL treats undefined for an optional argument without a default as
// missing. Option<raw> preserves null and other actual seed values.
if let Some(promise) = consume(
scope,
args.this(),
parsed.callback,
parsed.initial,
parsed.signal,
observer::REDUCE,
) {
rv.set(promise.into());
}
}
fn consume<'s>(
scope: &mut v8::PinScope<'s, '_>,
source: v8::Local<'s, v8::Object>,
callback: webidl::WebIdlCallbackFunction,
initial: Option<v8::Local<'s, v8::Value>>,
signal: Option<ResolvedAbortSignal<'s>>,
kind: i32,
) -> Option<v8::Local<'s, v8::Promise>> {
let resolver = v8::PromiseResolver::new(scope)?;
let promise = resolver.get_promise(scope);
let Some((observer, signal)) = promise::new_controlled_observer(scope, kind, resolver, signal)
else {
return Some(promise);
};
set_callback(scope, observer, CALLBACK, callback);
if let Some(initial) = initial {
set_accumulator(scope, observer, initial);
}
subscribe_internal(scope, source, observer, Some(signal));
Some(promise)
}
fn index<'s>(scope: &mut v8::PinScope<'s, '_>, observer: v8::Local<'s, v8::Object>) -> u64 {
get_private_value(scope, observer, INDEX)
.and_then(|value| v8::Local::<v8::BigInt>::try_from(value).ok())
.map_or(0, |value| value.u64_value().0)
}
fn increment_index<'s>(scope: &mut v8::PinScope<'s, '_>, observer: v8::Local<'s, v8::Object>) {
let next = index(scope, observer).wrapping_add(1);
set_private_value(
scope,
observer,
INDEX,
v8::BigInt::new_from_u64(scope, next).into(),
);
}
fn set_accumulator<'s>(
scope: &mut v8::PinScope<'s, '_>,
observer: v8::Local<'s, v8::Object>,
value: v8::Local<'s, v8::Value>,
) {
set_private_value(scope, observer, ACCUMULATOR, value);
set_private_value(
scope,
observer,
HAS_ACCUMULATOR,
v8::Boolean::new(scope, true).into(),
);
}
pub(super) fn notify<'s>(
scope: &mut v8::PinScope<'s, '_>,
observer: v8::Local<'s, v8::Object>,
notification: Notification<'s>,
kind: i32,
) {
let has_accumulator =
get_private_value(scope, observer, HAS_ACCUMULATOR).is_some_and(|value| value.is_true());
match notification {
Notification::Next(value) => {
// Do not skip a snapshot notification after settlement: a prior
// observer can cancel/complete the source during this same next().
if kind == observer::REDUCE && !has_accumulator {
set_accumulator(scope, observer, value);
increment_index(scope, observer);
return;
}
let callback =
object_slot(scope, observer, CALLBACK).expect("Observable consumer callback");
let accumulator = get_private_value(scope, observer, ACCUMULATOR)
.unwrap_or_else(|| v8::undefined(scope).into());
let idx = index(scope, observer);
let arguments = [
accumulator,
value,
v8::Number::new(scope, idx as f64).into(),
];
let arguments = if kind == observer::REDUCE {
&arguments[..]
} else {
&arguments[1..]
};
let result = callbacks::invoke_value(scope, callback, arguments);
if let Err(error) = result {
if let Some(resolver) = promise::start_settlement(scope, observer) {
resolver.reject(scope, error);
}
promise::abort_controller(scope, observer, error);
}
// The draft increments after invocation; a reentrant next() sees
// the current index. Read it again to preserve nested increments.
increment_index(scope, observer);
if kind == observer::REDUCE
&& let Ok(value) = result
{
set_accumulator(scope, observer, value);
}
}
Notification::Error(error) => {
if let Some(resolver) = promise::start_settlement(scope, observer) {
resolver.reject(scope, error);
}
}
Notification::Complete => {
let value = if kind == observer::REDUCE {
get_private_value(scope, observer, ACCUMULATOR)
.unwrap_or_else(|| v8::undefined(scope).into())
} else {
v8::undefined(scope).into()
};
if let Some(resolver) = promise::start_settlement(scope, observer) {
if kind == observer::REDUCE && !has_accumulator {
let error = v8::Exception::type_error(
scope,
v8str(scope, "Reduce of empty Observable with no initial value"),
);
resolver.reject(scope, error);
} else {
resolver.resolve(scope, value);
}
}
}
}
}
+6 -37
View File
@@ -1,16 +1,8 @@
use super::{
observer::{self, Notification},
promise, signal_arg,
state::object_slot,
subscribe_internal,
promise, signal_arg, subscribe_internal,
};
use crate::{
abort_signal_route::ResolvedAbortSignal,
util::{set_private_value, v8str},
webidl,
};
const CONTROLLER_SIGNAL: &str = "__moliObservableFirstControllerSignal";
use crate::{abort_signal_route::ResolvedAbortSignal, util::v8str, webidl};
#[derive(webidl::WebIdlArgs)]
#[webidl(prefix = "Observable.first")]
@@ -32,29 +24,11 @@ pub(super) fn first<'s>(
};
let promise = resolver.get_promise(scope);
rv.set(promise.into());
let Some(controller) = ResolvedAbortSignal::new(scope) else {
let Some((observer, signal)) =
promise::new_controlled_observer(scope, observer::FIRST, resolver, parsed.signal)
else {
return;
};
let mut sources = vec![controller];
sources.extend(parsed.signal);
let Some(signal) = ResolvedAbortSignal::dependent(scope, &sources) else {
return;
};
let Some(observer) = promise::new_observer(
scope,
observer::FIRST,
resolver,
Some(signal),
parsed.signal.is_some(),
) else {
return;
};
set_private_value(
scope,
observer,
CONTROLLER_SIGNAL,
controller.value().into(),
);
subscribe_internal(scope, args.this(), observer, Some(signal));
}
@@ -70,12 +44,7 @@ pub(super) fn notify<'s>(
}
// 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);
}
promise::abort_controller(scope, observer, v8::undefined(scope).into());
}
Notification::Error(error) => {
if let Some(resolver) = promise::start_settlement(scope, observer) {
+4 -1
View File
@@ -1,13 +1,15 @@
//! Internal observer steps share notification ordering with script observers,
//! while script callbacks keep their typed Web IDL invocation boundary.
use super::{callbacks, collect, first, invoke_and_report, state::*};
use super::{callbacks, collect, consume, first, invoke_and_report, state::*};
use crate::util::get_private_value;
pub(super) const NATIVE_KIND: &str = "__moliObservableNativeObserver";
pub(super) const FIRST: i32 = 1;
pub(super) const LAST: i32 = 2;
pub(super) const TO_ARRAY: i32 = 3;
pub(super) const FOR_EACH: i32 = 4;
pub(super) const REDUCE: i32 = 5;
pub(super) const SUBSCRIBER: &str = "__moliObservableNativeSubscriber";
#[derive(Clone, Copy)]
@@ -45,6 +47,7 @@ pub(super) fn notify<'s>(
match kind {
FIRST => first::notify(scope, observer, notification),
LAST | TO_ARRAY => collect::notify(scope, observer, notification, kind),
FOR_EACH | REDUCE => consume::notify(scope, observer, notification, kind),
_ => unreachable!("unknown native Observable observer"),
}
let exception = scope.exception();
@@ -13,6 +13,7 @@ const REJECTION_SIGNAL: &str = "__moliObservableRejectionSignal";
const REJECTION_ALGORITHM: &str = "__moliObservableRejectionAlgorithm";
const SETTLED: &str = "__moliObservablePromiseSettled";
const PROMISE_OBSERVER: &str = "__moliObservablePromiseObserver";
const CONTROLLER_SIGNAL: &str = "__moliObservableOperatorController";
pub(super) const VALUE: &str = "__moliObservablePromiseValue";
#[derive(WebApiObject)]
@@ -26,6 +27,46 @@ struct PromiseObserver<'scope> {
settled: bool,
}
/// Operators that may stop their own subscription use a private controller
/// and the existing AbortSignal dependency graph, preserving caller abort order.
pub(super) fn new_controlled_observer<'s>(
scope: &mut v8::PinScope<'s, '_>,
kind: i32,
resolver: v8::Local<'s, v8::PromiseResolver>,
caller_signal: Option<ResolvedAbortSignal<'s>>,
) -> Option<(v8::Local<'s, v8::Object>, ResolvedAbortSignal<'s>)> {
let controller = ResolvedAbortSignal::new(scope)?;
let mut sources = vec![controller];
sources.extend(caller_signal);
let signal = ResolvedAbortSignal::dependent(scope, &sources)?;
let observer = new_observer(scope, kind, resolver, Some(signal), caller_signal.is_some())?;
set_private_value(
scope,
observer,
CONTROLLER_SIGNAL,
controller.value().into(),
);
Some((observer, signal))
}
pub(super) fn abort_controller<'s>(
scope: &mut v8::PinScope<'s, '_>,
observer: v8::Local<'s, v8::Object>,
reason: v8::Local<'s, v8::Value>,
) {
if let Some(signal) = object_slot(scope, observer, CONTROLLER_SIGNAL)
.and_then(|signal| ResolvedAbortSignal::resolve(scope, signal))
{
// Signal-abort defaults an undefined reason, including `throw undefined`.
let reason = if reason.is_undefined() {
crate::native_bridge::abort::abort_error_value(scope)
} else {
reason
};
signal.abort(scope, reason);
}
}
pub(super) fn new_observer<'s>(
scope: &mut v8::PinScope<'s, '_>,
kind: i32,
@@ -1,5 +1,119 @@
use super::*;
#[test]
fn observable_callback_consumers_conversion_reentrancy_cancellation_and_exception_identity() {
let mut vm = new_storage_test_vm("https://observable-consumers.test/");
vm.eval(&format!(
"({}).then(value => {{ globalThis.consumerResult = JSON.stringify(value); }});",
include_str!("../../../tests/fixtures/observable-callback-consumers.js")
))
.expect("Observable callback consumers fixture should evaluate");
let result = vm
.eval("consumerResult")
.expect("Observable consumers 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() >= 132, "{result}");
}
#[test]
fn observable_callback_consumers_preserve_callback_realms_and_callee_promises() {
let mut vm = new_storage_test_vm("https://observable-consumers-realms.test/");
vm.eval("document.appendChild(document.createElement('iframe'))")
.unwrap();
materialize_single_child_default_realm_for_test(&mut vm, "Observable callback consumers realm");
vm.eval(r#"
(async () => {
const child = document.querySelector('iframe').contentWindow, checks = [];
const callback = child.Function('a', 'b', 'globalThis.consumerThis = this; return a + b;');
for (const name of ['forEach', 'reduce']) {
const method = child.Observable.prototype[name];
const call = (method, source, cb) => Reflect.apply(method, source, name === 'reduce' ? [cb, 10] : [cb]);
const promise = call(method, Observable.from([5]), callback);
checks.push(promise instanceof child.Promise, !(promise instanceof Promise));
checks.push(await promise === (name === 'reduce' ? 15 : undefined), child.consumerThis === child);
const local = call(Observable.prototype[name], child.Observable.from([7]), (a, b) => a + b);
checks.push(local instanceof Promise, !(local instanceof child.Promise), await local === (name === 'reduce' ? 17 : undefined));
const invalid = call(method, {}, () => {});
checks.push(invalid instanceof child.Promise);
await invalid.catch(e => checks.push(e instanceof child.TypeError, !(e instanceof TypeError)));
const marker = new child.Error('callback error');
await call(method, Observable.from([1]), () => { throw marker; }).catch(e => checks.push(e === marker, e instanceof child.Error));
}
await child.Observable.prototype.reduce.call(Observable.from([]), () => {}).catch(e => checks.push(e instanceof child.TypeError, !(e instanceof TypeError)));
globalThis.consumerRealms = JSON.stringify(checks);
})();
"#).unwrap();
let result: Vec<bool> = serde_json::from_str(&vm.eval("consumerRealms").unwrap()).unwrap();
assert_eq!(result.len(), 26);
assert!(result.iter().all(|value| *value), "{result:?}");
}
#[test]
fn observable_callback_consumers_trace_callbacks_and_release_abandoned_and_cancelled_state() {
let mut vm = new_storage_test_vm("https://observable-consumers-gc.test/");
vm.eval(r#"
globalThis.consumerCases = [];
globalThis.abandonedConsumers = [];
for (const name of ['forEach', 'reduce']) {
for (const kept of [false, true]) (() => {
let subscriber;
const source = new Observable(s => { subscriber = s; });
const token = {}, seed = {}, callback = () => token;
const promise = name === 'reduce' ? source.reduce(callback, seed) : source.forEach(callback);
const refs = {source: new WeakRef(source), subscriber: new WeakRef(subscriber), callback: new WeakRef(callback), token: new WeakRef(token), seed: new WeakRef(seed)};
if (kept) consumerCases.push({name, promise, refs});
else abandonedConsumers.push(refs);
})();
}
"#).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([
abandonedConsumers.every(refs => Object.values(refs).every(ref => ref.deref() === undefined)),
consumerCases.every(c => c.refs.source.deref() === undefined && c.refs.subscriber.deref() !== undefined && c.refs.callback.deref() !== undefined && c.refs.token.deref() !== undefined),
consumerCases[0].refs.seed.deref() === undefined, consumerCases[1].refs.seed.deref() !== undefined
])"#).unwrap(), "[true,true,true,true]");
vm.eval(r#"
for (const c of consumerCases) {
c.promise.then(value => { c.correct = value === (c.name === 'reduce' ? c.refs.seed.deref() : undefined); });
c.refs.subscriber.deref().complete(); delete c.promise;
}
"#).unwrap();
assert_eq!(
vm.eval("consumerCases.every(c => c.correct)").unwrap(),
"true"
);
collect(&mut vm);
assert_eq!(vm.eval("consumerCases.every(c => Object.values(c.refs).every(ref => ref.deref() === undefined))").unwrap(), "true");
vm.eval(r#"
globalThis.cancelledConsumers = [];
for (const name of ['forEach', 'reduce']) (() => {
const ac = new AbortController(), token = {}, seed = {}, callback = () => token;
let subscriber;
const source = new Observable(s => { subscriber = s; });
const promise = name === 'reduce' ? source.reduce(callback, seed, {signal: ac.signal}) : source.forEach(callback, {signal: ac.signal});
promise.catch(() => {}); ac.abort('cancelled');
cancelledConsumers.push({promise, signal: ac.signal, refs: [new WeakRef(token), new WeakRef(seed), new WeakRef(callback), new WeakRef(subscriber)]});
})();
"#).unwrap();
collect(&mut vm);
assert_eq!(
vm.eval("cancelledConsumers.every(c => c.refs.every(ref => ref.deref() === undefined))")
.unwrap(),
"true"
);
}
#[test]
fn observable_collect_values_abort_order_and_native_promise_observers() {
let mut vm = new_storage_test_vm("https://observable-collect.test/");
@@ -1,5 +1,22 @@
use super::*;
#[tokio::test]
async fn worker_observable_callback_consumers_conversion_reentrancy_cancellation_and_exception_identity()
{
ensure_v8();
let mut handle = spawn_worker(
format!(
"({}).then(value => {{ postMessage(value); close(); }});",
include_str!("../../../../tests/fixtures/observable-callback-consumers.js")
),
"https://observable-consumers.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() >= 132, "{result}");
}
#[tokio::test]
async fn worker_observable_collect_values_abort_order_and_native_promise_observers() {
ensure_v8();
@@ -0,0 +1,193 @@
(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 rejected = p => p.then(() => { throw new Error('expected rejection'); }, e => e);
const thrown = fn => { try { fn(); } catch (e) { return e; } };
const test = async (label, fn) => { try { await fn(); } catch (e) { check(false, label + ': ' + e); } };
const names = ['forEach', 'reduce'];
for (const name of names) check(typeof Observable.prototype[name] === 'function', name + ' exposed');
if (failures.length) return {checks, failures};
for (const name of names) {
const method = Observable.prototype[name];
const call = (source, callback, options) => Reflect.apply(method, source,
name === 'reduce' ? [callback, 10, options] : [callback, options]);
await test(name + ' conversion', async () => {
const desc = Object.getOwnPropertyDescriptor(Observable.prototype, name);
check(method.name === name && method.length === 1, name + ' name and length');
check(desc.enumerable && desc.writable && desc.configurable, name + ' descriptor');
check(thrown(() => new method(() => {})) instanceof TypeError, name + ' non-constructor');
const source = Observable.from([1]);
let reads = 0, traps = 0;
const options = {get signal() { reads++; }};
const revoked = Proxy.revocable(source, {}); revoked.revoke();
for (const value of [undefined, null, false, 1, Symbol(), {}, Object.create(source),
Object.create(Observable.prototype), new Proxy(source, {get() { traps++; }}), revoked.proxy]) {
const promise = call(value, () => {}, options);
check(promise instanceof Promise && await rejected(promise) instanceof TypeError, name + ' receiver rejects Promise');
}
for (const callback of [undefined, null, {}, {handleEvent() {}}, 1, Object.create(Function.prototype)]) {
check(await rejected(call(source, callback, options)) instanceof TypeError, name + ' invalid callback');
}
check(await rejected(method.call(source)) instanceof TypeError, name + ' required callback');
check(reads === 0 && traps === 0, name + ' receiver and callback validation precede options');
for (const options of [1, true, {signal: null}, {signal: {}}, {signal: new Proxy(new AbortController().signal, {})}]) {
check(await rejected(call(source, () => {}, options)) instanceof TypeError, name + ' invalid options');
}
const marker = {};
check(await rejected(call(source, () => {}, {get signal() { throw marker; }})) === marker, name + ' getter error identity');
Object.setPrototypeOf(source, null);
await call(source, () => {}, options);
check(reads === 1, name + ' genuine brand with changed prototype and one signal read');
const proxy = Proxy.revocable(() => {}, {});
const promise = call(Observable.from([1]), proxy.proxy, {get signal() { proxy.revoke(); }});
check(await rejected(promise) instanceof TypeError, name + ' callback revoked after conversion rejects');
check(await rejected(call(Observable.from([1]), class Visitor {})) instanceof TypeError, name + ' class callback fails at invocation');
});
await test(name + ' invocation', async () => {
const log = [], marker = {}, returned = {get then() { throw marker; }};
let subscriber;
const source = new Observable(s => { subscriber = s; });
const callback = new Proxy(function (...args) {
check(this === undefined, name + ' strict callback this is undefined');
log.push(args);
return name === 'reduce' ? args[0] + args[1] : returned;
}, {});
const promise = call(source, callback);
subscriber.next(1); subscriber.next(2);
same(log, name === 'reduce' ? [[10, 1, 0], [11, 2, 1]] : [[1, 0], [2, 1]], name + ' synchronous arguments');
check(subscriber.active, name + ' success keeps producer active');
subscriber.complete();
check(await promise === (name === 'reduce' ? 13 : undefined), name + ' completion result');
check(await rejected(call(new Observable(s => s.error(marker)), () => {})) === marker, name + ' source error identity');
check(await rejected(call(new Observable(() => { throw marker; }), () => {})) === marker, name + ' initializer error identity');
});
await test(name + ' callback errors', async () => {
for (const marker of [{}, null, undefined]) {
let subscriber, visits = 0, teardowns = 0;
const source = new Observable(s => {
subscriber = s; s.addTeardown(() => teardowns++);
s.next(1); s.next(2); s.complete();
});
const error = await rejected(call(source, () => { visits++; throw marker; }));
check(error === marker && visits === 1 && teardowns === 1 && !subscriber.active, name + ' callback throw rejects and cancels once');
check(marker === undefined ? subscriber.signal.reason.name === 'AbortError' : subscriber.signal.reason === marker, name + ' callback throw abort reason');
}
const ac = new AbortController(), reason = {}, later = {};
const promise = call(Observable.from([1]), () => { ac.abort(reason); throw later; }, {signal: ac.signal});
check(await rejected(promise) === reason, name + ' reentrant caller abort wins before callback throw');
});
await test(name + ' cancellation timing', async () => {
const reason = {}, preaborted = AbortSignal.abort(reason);
let starts = 0, calls = 0, subscriber;
const source = new Observable(s => { starts++; subscriber = s; });
check(await rejected(call(source, () => calls++, {signal: preaborted})) === reason && starts === 0 && calls === 0, name + ' pre-abort skips producer');
const ac = new AbortController(), log = [];
ac.signal.addEventListener('abort', () => { log.push('outer'); queueMicrotask(() => log.push('outer job')); });
const promise = call(source, () => {}, {signal: ac.signal}).catch(e => { log.push('reject'); return e; });
subscriber.signal.addEventListener('abort', () => { log.push('inner'); queueMicrotask(() => log.push('inner job')); });
subscriber.addTeardown(() => log.push('teardown'));
ac.abort(reason);
same(log, ['outer', 'inner', 'teardown'], name + ' dependent signal abort order');
check(await promise === reason, name + ' caller reason identity');
same(log, ['outer', 'inner', 'teardown', 'outer job', 'reject', 'inner job'], name + ' dependent signal microtask order');
});
await test(name + ' reentrant dispatch', async () => {
let subscriber;
const source = new Observable(s => { subscriber = s; }), indices = [], values = [];
const promise = call(source, (...args) => {
const value = args[name === 'reduce' ? 1 : 0], idx = args.at(-1);
indices.push(idx); values.push(value);
if (value === 1) subscriber.next(2);
return name === 'reduce' ? args[0] + value : undefined;
});
subscriber.next(1); subscriber.next(3); subscriber.complete();
same(values, [1, 2, 3], name + ' nested emission order');
// The draft increments after callback invocation; Chromium currently
// increments before it and returns [0, 1, 2] in this case.
same(indices, [0, 0, 2], name + ' draft reentrant index order');
check(await promise === (name === 'reduce' ? 14 : undefined), name + ' outer reducer result replaces nested result');
for (const terminal of ['abort', 'complete']) {
const ac = new AbortController(), marker = {}, received = [];
const source = new Observable(s => { subscriber = s; });
source.subscribe(() => { if (terminal === 'abort') ac.abort(marker); else subscriber.complete(); });
const promise = call(source, (...args) => { received.push(args); return 99; }, {signal: ac.signal});
const handled = terminal === 'abort' ? rejected(promise) : promise;
subscriber.next(1);
same(received, name === 'reduce' ? [[10, 1, 0]] : [[1, 0]], name + ' already captured next runs after ' + terminal);
check(await handled === (terminal === 'abort' ? marker : name === 'reduce' ? 10 : undefined), name + ' snapshot cannot replace settled result');
subscriber.complete();
}
});
}
await test('reduce seeds and raw results', async () => {
const empty = Observable.from([]);
check(await rejected(empty.reduce(() => {})) instanceof TypeError, 'empty reduction without seed rejects');
check(await rejected(empty.reduce(() => {}, undefined)) instanceof TypeError, 'optional undefined seed is missing under Web IDL');
check(await Observable.from([undefined]).reduce(() => { throw 1; }) === undefined, 'emitted undefined is a real accumulator');
check(await empty.reduce(() => {}, null) === null, 'null is a real seed');
const reason = {};
check(await rejected(empty.reduce(() => {}, undefined, {signal: AbortSignal.abort(reason)})) === reason, 'optional missing seed still converts later options');
const seed = {}, values = [];
check(await empty.reduce(() => { throw 1; }, seed) === seed, 'empty reduction preserves seed identity');
const sum = await Observable.from([1, 2, 3]).reduce((a, v, i) => { values.push([a, v, i]); return a + v; });
check(sum === 6, 'seedless reduction result');
same(values, [[1, 2, 1], [3, 3, 2]], 'first emitted value seeds without callback');
const thenable = {get then() { throw new Error('intermediate result assimilated'); }};
let calls = 0;
check(await Observable.from([1, 2]).reduce((a, v) => { calls++; if (v === 1) return thenable; check(a === thenable, 'raw intermediate accumulator'); return 42; }, 0) === 42 && calls === 2, 'intermediate thenables not assimilated');
let subscriber, reads = 0, resolveValue;
const ac = new AbortController(), value = {get then() { reads++; ac.abort('late'); return resolve => { resolveValue = resolve; }; }};
const promise = new Observable(s => { subscriber = s; }).reduce(() => value, 0, {signal: ac.signal});
subscriber.next(1); check(reads === 0, 'final thenable not read before completion');
subscriber.complete(); await Promise.resolve(); resolveValue(8);
check(await promise === 8 && reads === 1, 'final assimilation locks before then getter abort');
const ignored = Promise.reject('visitor promise'); ignored.catch(() => {});
check(await Observable.from([1]).forEach(() => ignored) === undefined, 'visitor return Promise is ignored');
});
await test('shared subscription and iterator close', async () => {
let subscriber;
const source = new Observable(s => { subscriber = s; }), marker = {};
const all = source.toArray(), completion = rejected(source.forEach(() => { throw marker; }));
subscriber.next(1); check(subscriber.active && await completion === marker, 'callback error removes only its own observer');
subscriber.next(2); subscriber.complete();
same(await all, [1, 2], 'other observer continues after visitor error');
const closeError = {}, errors = [];
const onerror = e => { if (e.error === closeError) { errors.push(e.error); e.preventDefault(); } };
addEventListener('error', onerror);
try {
for (const name of names) {
const iterator = {next: () => ({value: 1}), return() { throw closeError; }};
const source = Observable.from({[Symbol.iterator]: () => iterator});
const promise = name === 'reduce' ? source.reduce(() => { throw marker; }, 0) : source.forEach(() => { throw marker; });
check(await rejected(promise) === marker, name + ' close error cannot replace callback error');
}
check(errors.length === 2, 'iterator close errors reported once per consumer');
} finally { removeEventListener('error', onerror); }
});
await test('intrinsic operations', async () => {
const source = Observable.from([1, 2]), forEach = Observable.prototype.forEach, reduce = Observable.prototype.reduce;
const saved = [globalThis.Promise, globalThis.AbortController, Observable.prototype.subscribe, AbortSignal.any];
const poison = () => { throw new Error('public native implementation consulted'); };
let each, reduced;
try {
globalThis.Promise = globalThis.AbortController = Observable.prototype.subscribe = AbortSignal.any = poison;
each = forEach.call(source, () => {}); reduced = reduce.call(source, (a, v) => a + v, 0);
} finally {
[globalThis.Promise, globalThis.AbortController, Observable.prototype.subscribe, AbortSignal.any] = saved;
}
check(each instanceof Promise && await each === undefined, 'forEach bypasses mutable globals and subscribe');
check(reduced instanceof Promise && await reduced === 3, 'reduce bypasses mutable globals and subscribe');
});
return {checks, failures};
})()