feat(dom): collect Observable values with last and toArray

Add completion-based Promise operators in Window and workers, preserving
raw values, shared subscriptions and direct AbortSignal ordering. Reuse
the native Promise observer lifetime and settlement logic with first.

Cover receiver conversion, callee realms, thenable reentrancy, inherited
array setters and value reachability. Related WPTs improve from 302/326
to 324/326; EventTarget regressions remain 150/150.
This commit is contained in:
ldm0
2026-09-24 20:01:11 +08:00
parent 999fcb83c3
commit 68ec15233d
9 changed files with 618 additions and 113 deletions
@@ -4541,6 +4541,10 @@ 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-last.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-last.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
dom/ranges/Range-attributes.html
dom/ranges/Range-cloneContents.html
+13 -1
View File
@@ -4,10 +4,12 @@
//! Window/worker owners.
mod callbacks;
mod collect;
mod event_target;
mod first;
mod from;
mod observer;
mod promise;
mod state;
pub(crate) use event_target::event_target_when;
@@ -27,6 +29,10 @@ struct ObservablePrototype {
subscribe: (),
#[webapi(method, length = 0, returns_promise, callback = first::first)]
first: (),
#[webapi(method, length = 0, returns_promise, callback = collect::last)]
last: (),
#[webapi(method = "toArray", length = 0, returns_promise, callback = collect::to_array)]
to_array: (),
}
#[derive(WebApiFunctionTemplate)]
@@ -220,6 +226,12 @@ fn subscribe_internal<'s>(
let mut observers = list(scope, subscriber, OBSERVERS);
observers.push(observer);
set_list(scope, subscriber, OBSERVERS, &observers);
let native = observer::is_native(scope, observer);
if native {
// Pending native Promise observers keep their producer reachable even
// when no cancellation callback provides the reference to Subscriber.
set_private_value(scope, observer, observer::SUBSCRIBER, subscriber.into());
}
if let Some(signal) = signal {
if signal.is_aborted(scope) {
if fresh {
@@ -239,7 +251,7 @@ fn subscribe_internal<'s>(
.expect("Observable abort algorithm should allocate");
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) {
if native {
// Native observers trace their private cancellation callback.
// The internal signal must not root an abandoned subscription.
signal.register_weak_rethrowing_algorithm(scope, algorithm);
+117
View File
@@ -0,0 +1,117 @@
//! Promise operators that retain values until the source completes.
use super::{
observer::{self, Notification},
promise, signal_arg,
state::object_slot,
subscribe_internal,
};
use crate::{
abort_signal_route::ResolvedAbortSignal,
util::{get_private_value, set_private_value, v8str},
webidl,
};
const HAS_VALUE: &str = "__moliObservableHasLastValue";
#[derive(webidl::WebIdlArgs)]
#[webidl(prefix = "Observable")]
struct CollectArgs<'scope> {
#[webidl(with = signal_arg)]
signal: Option<ResolvedAbortSignal<'scope>>,
}
pub(super) fn last<'s>(
scope: &mut v8::PinScope<'s, '_>,
args: v8::FunctionCallbackArguments<'s>,
rv: v8::ReturnValue<'_, v8::Value>,
) {
collect(scope, args, rv, observer::LAST);
}
pub(super) fn to_array<'s>(
scope: &mut v8::PinScope<'s, '_>,
args: v8::FunctionCallbackArguments<'s>,
rv: v8::ReturnValue<'_, v8::Value>,
) {
collect(scope, args, rv, observer::TO_ARRAY);
}
fn collect<'s>(
scope: &mut v8::PinScope<'s, '_>,
args: v8::FunctionCallbackArguments<'s>,
mut rv: v8::ReturnValue<'_, v8::Value>,
kind: i32,
) {
let Some(parsed) = webidl::parse_args::<CollectArgs<'s>>(scope, &args) else {
return;
};
let Some(resolver) = v8::PromiseResolver::new(scope) else {
return;
};
rv.set(resolver.get_promise(scope).into());
let Some(observer) = promise::new_observer(scope, kind, resolver, parsed.signal, true) else {
return;
};
if kind == observer::TO_ARRAY {
let values = v8::Array::new(scope, 0);
set_private_value(scope, observer, promise::VALUE, values.into());
}
// These operators do not cancel on next(), and subscribe directly with the
// caller's signal. Adding a dependent signal would change abort ordering.
subscribe_internal(scope, args.this(), observer, parsed.signal);
}
pub(super) fn notify<'s>(
scope: &mut v8::PinScope<'s, '_>,
observer: v8::Local<'s, v8::Object>,
notification: Notification<'s>,
kind: i32,
) {
if promise::is_settled(scope, observer) {
return;
}
match notification {
Notification::Next(value) if kind == observer::TO_ARRAY => {
let values = object_slot(scope, observer, promise::VALUE)
.and_then(|value| v8::Local::<v8::Array>::try_from(value).ok())
.expect("Observable values");
let Some(key) = v8::String::new(scope, &values.length().to_string()) else {
return;
};
// This private, dense array represents an Infra list until
// completion. Append without invoking inherited index setters.
values.create_data_property(scope, key.into(), value);
}
Notification::Next(value) => {
set_private_value(scope, observer, promise::VALUE, value);
set_private_value(
scope,
observer,
HAS_VALUE,
v8::Boolean::new(scope, true).into(),
);
}
Notification::Error(error) => {
if let Some(resolver) = promise::start_settlement(scope, observer) {
resolver.reject(scope, error);
}
}
Notification::Complete => {
let has_value = kind == observer::TO_ARRAY
|| get_private_value(scope, observer, HAS_VALUE)
.is_some_and(|value| value.is_true());
let value = get_private_value(scope, observer, promise::VALUE)
.unwrap_or_else(|| v8::undefined(scope).into());
if let Some(resolver) = promise::start_settlement(scope, observer) {
if has_value {
resolver.resolve(scope, value);
} else {
let error =
v8::Exception::range_error(scope, v8str(scope, "No values in Observable"));
resolver.reject(scope, error);
}
}
}
}
}
+17 -111
View File
@@ -1,24 +1,16 @@
use moli_webapi_declare::WebApiObject;
use super::{
callbacks,
observer::{self, NATIVE_KIND, Notification},
signal_arg,
observer::{self, Notification},
promise, signal_arg,
state::object_slot,
subscribe_internal,
};
use crate::{
abort_signal_route::ResolvedAbortSignal,
util::{get_private_value, set_private_value, v8str},
util::{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")]
@@ -27,21 +19,6 @@ struct FirstArgs<'scope> {
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>,
@@ -63,79 +40,22 @@ pub(super) fn first<'s>(
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(
let Some(observer) = promise::new_observer(
scope,
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.
resolver,
Some(signal),
parsed.signal.is_some(),
) else {
return;
};
set_private_value(
scope,
observer,
SETTLED,
v8::Boolean::new(scope, true).into(),
CONTROLLER_SIGNAL,
controller.value().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)
subscribe_internal(scope, args.this(), observer, Some(signal));
}
pub(super) fn notify<'s>(
@@ -145,7 +65,7 @@ pub(super) fn notify<'s>(
) {
match notification {
Notification::Next(value) => {
if let Some(resolver) = start_settlement(scope, observer) {
if let Some(resolver) = promise::start_settlement(scope, observer) {
resolver.resolve(scope, value);
}
// Abort after resolving, including when a then getter reenters
@@ -158,12 +78,12 @@ pub(super) fn notify<'s>(
}
}
Notification::Error(error) => {
if let Some(resolver) = start_settlement(scope, observer) {
if let Some(resolver) = promise::start_settlement(scope, observer) {
resolver.reject(scope, error);
}
}
Notification::Complete => {
if let Some(resolver) = start_settlement(scope, observer) {
if let Some(resolver) = promise::start_settlement(scope, observer) {
let error =
v8::Exception::range_error(scope, v8str(scope, "No values in Observable"));
resolver.reject(scope, error);
@@ -171,17 +91,3 @@ pub(super) fn notify<'s>(
}
}
}
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));
}
}
+5 -1
View File
@@ -1,11 +1,14 @@
//! 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 super::{callbacks, collect, 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 SUBSCRIBER: &str = "__moliObservableNativeSubscriber";
#[derive(Clone, Copy)]
pub(super) enum Notification<'s> {
@@ -41,6 +44,7 @@ pub(super) fn notify<'s>(
v8::tc_scope!(let scope, scope);
match kind {
FIRST => first::notify(scope, observer, notification),
LAST | TO_ARRAY => collect::notify(scope, observer, notification, kind),
_ => unreachable!("unknown native Observable observer"),
}
let exception = scope.exception();
+137
View File
@@ -0,0 +1,137 @@
//! Shared lifetime and single-settlement state for native Promise observers.
use moli_webapi_declare::WebApiObject;
use super::{callbacks, observer, state::object_slot};
use crate::{
abort_signal_route::ResolvedAbortSignal,
util::{get_private_value, set_private_value},
};
const RESOLVER: &str = "__moliObservablePromiseResolver";
const REJECTION_SIGNAL: &str = "__moliObservableRejectionSignal";
const REJECTION_ALGORITHM: &str = "__moliObservableRejectionAlgorithm";
const SETTLED: &str = "__moliObservablePromiseSettled";
const PROMISE_OBSERVER: &str = "__moliObservablePromiseObserver";
pub(super) const VALUE: &str = "__moliObservablePromiseValue";
#[derive(WebApiObject)]
#[webapi(plain)]
struct PromiseObserver<'scope> {
#[webapi(slot = observer::NATIVE_KIND)]
kind: i32,
#[webapi(slot = RESOLVER)]
resolver: v8::Local<'scope, v8::Object>,
#[webapi(slot = SETTLED)]
settled: bool,
}
pub(super) fn new_observer<'s>(
scope: &mut v8::PinScope<'s, '_>,
kind: i32,
resolver: v8::Local<'s, v8::PromiseResolver>,
signal: Option<ResolvedAbortSignal<'s>>,
caller_signal: bool,
) -> Option<v8::Local<'s, v8::Object>> {
if let Some(signal) = signal.filter(|signal| signal.is_aborted(scope)) {
let reason = signal.reason(scope);
resolver.reject(scope, reason);
return None;
}
let observer = PromiseObserver::new(kind, resolver.into(), false)
.bind(scope)
.ok()?;
let promise = resolver.get_promise(scope);
// 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());
if let Some(signal) = signal {
let algorithm = v8::Function::builder(aborted)
.data(observer.into())
.build(scope)
.expect("Observable Promise abort algorithm");
set_private_value(scope, observer, REJECTION_SIGNAL, signal.value().into());
set_private_value(scope, observer, REJECTION_ALGORITHM, algorithm.into());
if caller_signal {
// Like subscribe({signal}), caller-driven cancellation remains
// an owner while this operation is pending.
signal.register_algorithm(scope, algorithm);
} else {
signal.register_weak_rethrowing_algorithm(scope, algorithm);
}
}
Some(observer)
}
pub(super) fn is_settled<'s>(
scope: &mut v8::PinScope<'s, '_>,
observer: v8::Local<'s, v8::Object>,
) -> bool {
get_private_value(scope, observer, SETTLED).is_some_and(|value| value.is_true())
}
pub(super) fn start_settlement<'s>(
scope: &mut v8::PinScope<'s, '_>,
observer: v8::Local<'s, v8::Object>,
) -> Option<v8::Local<'s, v8::PromiseResolver>> {
if is_settled(scope, observer) {
return None;
}
// Resolve can execute a then getter that reenters the producer or aborts
// its signal. A Promise assimilating a thenable is still Pending, so its
// V8 state 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 Promise 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(),
);
set_private_value(
scope,
observer,
observer::SUBSCRIBER,
v8::undefined(scope).into(),
);
// Operators read their final value before settling. Release accumulated
// values on rejection or abort too, even while the result Promise survives.
set_private_value(scope, observer, VALUE, v8::undefined(scope).into());
Some(resolver)
}
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 Promise abort data");
if callbacks::is_current(scope, observer)
&& let Some(resolver) = start_settlement(scope, observer)
{
resolver.reject(scope, args.get(0));
}
}
@@ -1,5 +1,129 @@
use super::*;
#[test]
fn observable_collect_values_abort_order_and_native_promise_observers() {
let mut vm = new_storage_test_vm("https://observable-collect.test/");
vm.eval(&format!(
"({}).then(value => {{ globalThis.collectResult = JSON.stringify(value); }});",
include_str!("../../../tests/fixtures/observable-collect.js")
))
.expect("Observable collection fixture should evaluate");
let result = vm
.eval("collectResult")
.expect("Observable collection 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() >= 110, "{result}");
}
#[test]
fn observable_collect_uses_callee_promise_array_and_error_realms() {
let mut vm = new_storage_test_vm("https://observable-collect-realms.test/");
vm.eval("document.appendChild(document.createElement('iframe'))")
.unwrap();
materialize_single_child_default_realm_for_test(&mut vm, "Observable collection realm");
vm.eval(r#"
(async () => {
const child = document.querySelector('iframe').contentWindow, checks = [];
for (const name of ['last', 'toArray']) {
const method = child.Observable.prototype[name];
const promise = method.call(Observable.from([4, 5]));
checks.push(promise instanceof child.Promise, !(promise instanceof Promise));
const value = await promise;
checks.push(name === 'last' ? value === 5 : value instanceof child.Array && !(value instanceof Array) && value[1] === 5);
const invalid = method.call({});
checks.push(invalid instanceof child.Promise);
await invalid.catch(e => checks.push(e instanceof child.TypeError, !(e instanceof TypeError)));
const ac = new AbortController(), marker = {};
const pending = method.call(new Observable(() => {}), {signal: ac.signal});
ac.abort(marker);
await pending.catch(e => checks.push(e === marker));
const local = Observable.prototype[name].call(child.Observable.from([7]));
checks.push(local instanceof Promise, !(local instanceof child.Promise));
const result = await local;
checks.push(name === 'last' ? result === 7 : result instanceof Array && !(result instanceof child.Array) && result[0] === 7);
}
await child.Observable.prototype.last.call(Observable.from([])).catch(e => checks.push(e instanceof child.RangeError, !(e instanceof RangeError)));
globalThis.collectRealms = JSON.stringify(checks);
})();
"#).unwrap();
let result = vm.eval("collectRealms").unwrap();
let result: Vec<bool> = serde_json::from_str(&result).unwrap();
assert_eq!(result.len(), 22);
assert!(result.iter().all(|value| *value), "{result:?}");
}
#[test]
fn observable_collect_traces_pending_values_and_releases_them_after_completion_or_abort() {
let mut vm = new_storage_test_vm("https://observable-collect-gc.test/");
vm.eval(r#"
globalThis.collectCases = [];
globalThis.abandonedCollect = [];
for (const mode of ['last', 'toArray']) {
(() => {
let subscriber;
const source = new Observable(s => { subscriber = s; });
const promise = source[mode](), value = {};
subscriber.next(value);
abandonedCollect.push([new WeakRef(source), new WeakRef(subscriber), new WeakRef(promise), new WeakRef(value)]);
})();
(() => {
let subscriber;
const source = new Observable(s => { subscriber = s; });
const promise = source[mode](), first = {}, second = {};
subscriber.next(first); subscriber.next(second);
collectCases.push({mode, promise, source: new WeakRef(source), subscriber: new WeakRef(subscriber), first: new WeakRef(first), second: new WeakRef(second)});
})();
}
"#).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([
abandonedCollect.every(refs => refs.every(ref => ref.deref() === undefined)),
collectCases.every(c => c.source.deref() === undefined && c.subscriber.deref() !== undefined && c.second.deref() !== undefined),
collectCases[0].first.deref() === undefined, collectCases[1].first.deref() !== undefined
])"#).unwrap(), "[true,true,true,true]");
vm.eval(r#"
for (const c of collectCases) {
c.promise.then(value => { c.correct = c.mode === 'last' ? value === c.second.deref() : value[0] === c.first.deref() && value[1] === c.second.deref(); });
c.subscriber.deref().complete();
delete c.promise;
}
"#).unwrap();
assert_eq!(
vm.eval("collectCases.every(c => c.correct)").unwrap(),
"true"
);
collect(&mut vm);
assert_eq!(vm.eval("collectCases.every(c => c.subscriber.deref() === undefined && c.first.deref() === undefined && c.second.deref() === undefined)").unwrap(), "true");
vm.eval(r#"
globalThis.cancelledCollect = [];
for (const mode of ['last', 'toArray']) {
(() => {
const ac = new AbortController();
let subscriber;
const source = new Observable(s => { subscriber = s; });
const promise = source[mode]({signal: ac.signal}), value = {};
subscriber.next(value);
promise.catch(() => {});
ac.abort('cancelled');
cancelledCollect.push({promise, value: new WeakRef(value), subscriber: new WeakRef(subscriber), signal: ac.signal});
})();
}
"#).unwrap();
collect(&mut vm);
assert_eq!(vm.eval("cancelledCollect.every(c => c.value.deref() === undefined && c.subscriber.deref() === undefined)").unwrap(), "true");
}
#[test]
fn observable_first_promises_cancellation_reentrancy_and_native_observers() {
let mut vm = new_storage_test_vm("https://observable-first.test/");
@@ -1,5 +1,21 @@
use super::*;
#[tokio::test]
async fn worker_observable_collect_values_abort_order_and_native_promise_observers() {
ensure_v8();
let mut handle = spawn_worker(
format!(
"({}).then(value => {{ postMessage(value); close(); }});",
include_str!("../../../../tests/fixtures/observable-collect.js")
),
"https://observable-collect.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() >= 110, "{result}");
}
#[tokio::test]
async fn worker_observable_first_promises_cancellation_reentrancy_and_native_observers() {
ensure_v8();
+185
View File
@@ -0,0 +1,185 @@
(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 methods = ['last', 'toArray'];
for (const name of methods) check(typeof Observable.prototype[name] === 'function', name + ' exposed');
if (failures.length) return {checks, failures};
for (const name of methods) {
const method = Observable.prototype[name];
await test(name + ' conversion', async () => {
const descriptor = Object.getOwnPropertyDescriptor(Observable.prototype, name);
check(method.name === name && method.length === 0, name + ' name and length');
check(descriptor.enumerable && descriptor.writable && descriptor.configurable, name + ' descriptor');
check(thrown(() => new method()) instanceof TypeError, name + ' not constructible');
const source = Observable.from([1, 2]);
let conversions = 0, traps = 0;
const options = {get signal() { conversions++; }};
const revoked = Proxy.revocable(source, {}); revoked.revoke();
for (const receiver of [undefined, null, false, 1, 'x', Symbol(), {},
Object.create(Observable.prototype), Object.create(source),
new Proxy(source, {get() { traps++; throw 1; }}), revoked.proxy]) {
const promise = method.call(receiver, options);
check(promise instanceof Promise && await rejected(promise) instanceof TypeError, name + ' receiver rejection');
}
check(conversions === 0 && traps === 0, name + ' brand before conversions without Proxy traps');
for (const options of [true, 1, 'x', Symbol(), {signal: null}, {signal: {}},
{signal: new Proxy(new AbortController().signal, {})}]) {
check(await rejected(method.call(source, options)) instanceof TypeError, name + ' invalid options rejection');
}
const marker = {};
check(await rejected(method.call(source, {get signal() { throw marker; }})) === marker, name + ' getter error identity');
Object.setPrototypeOf(source, null);
for (const options of [undefined, null, {}, {signal: undefined}]) {
const value = await method.call(source, options);
check(name === 'last' ? value === 2 : value.length === 2 && value[1] === 2, name + ' genuine receiver with changed prototype');
}
await method.call(source, options);
check(conversions === 1, name + ' signal read once');
});
await test(name + ' lifecycle', async () => {
const log = [];
let subscriber, settled = false;
const source = new Observable(s => {
subscriber = s;
s.signal.addEventListener('abort', () => log.push('abort'));
s.addTeardown(() => log.push('teardown'));
});
const promise = source[name]().then(value => { settled = true; log.push('resolved'); return value; });
subscriber.next(1); subscriber.next(2);
await Promise.resolve();
check(!settled && subscriber.active, name + ' waits for completion');
subscriber.complete();
same(log, ['abort', 'teardown'], name + ' closes before resolving');
const value = await promise;
same(value, name === 'last' ? 2 : [1, 2], name + ' result');
same(log, ['abort', 'teardown', 'resolved'], name + ' reaction timing');
const error = {};
check(await rejected(new Observable(s => { s.next(1); s.error(error); })[name]()) === error, name + ' error overrides retained values');
check(await rejected(new Observable(() => { throw error; })[name]()) === error, name + ' initializer exception');
for (const terminal of ['complete', 'error']) {
const ac = new AbortController(), reason = {};
const source = new Observable(s => { s.next(7); s.addTeardown(() => ac.abort(reason)); s[terminal](error); });
check(await rejected(source[name]({signal: ac.signal})) === reason, name + ' teardown abort before ' + terminal);
}
});
await test(name + ' direct abort ordering', async () => {
let starts = 0, subscriber;
const reason = {}, source = new Observable(s => { starts++; subscriber = s; });
const log = [];
const preaborted = source[name]({signal: AbortSignal.abort(reason)}).catch(e => { log.push('reject'); return e; });
Promise.resolve().then(() => log.push('later'));
check(await preaborted === reason && starts === 0, name + ' pre-abort skips subscription');
same(log, ['reject', 'later'], name + ' pre-abort is immediately rejected');
const ac = new AbortController();
ac.signal.addEventListener('abort', () => { log.push('outer'); Promise.resolve().then(() => log.push('outer job')); });
const pending = source[name]({signal: ac.signal}).catch(e => { log.push('reject'); return e; });
subscriber.signal.addEventListener('abort', () => { log.push('inner'); Promise.resolve().then(() => log.push('inner job')); });
subscriber.addTeardown(() => { log.push('teardown'); Promise.resolve().then(() => log.push('teardown job')); });
log.length = 0;
ac.abort(reason);
same(log, ['inner', 'teardown', 'outer'], name + ' abort algorithms precede caller event');
check(await pending === reason && !subscriber.active && subscriber.signal.reason === reason, name + ' abort preserves reason');
same(log, ['inner', 'teardown', 'outer', 'reject', 'inner job', 'teardown job', 'outer job'], name + ' abort microtask order');
});
await test(name + ' snapshot reentrancy', async () => {
let subscriber;
const source = new Observable(s => { subscriber = s; });
source.subscribe(() => subscriber.complete());
const result = source[name]();
const handled = name === 'last' ? rejected(result) : result;
subscriber.next('after completion');
const value = await handled;
check(name === 'last' ? value instanceof RangeError : value.length === 0, name + ' stale next snapshot cannot change completed result');
});
}
await test('last values and thenables', async () => {
for (const value of [undefined, null, false, NaN, -0, 1n]) {
check(Object.is(await Observable.from([value]).last(), value), 'last distinguishes emitted value from empty');
}
check(await rejected(Observable.from([]).last()) instanceof RangeError, 'empty last RangeError');
const discarded = {get then() { throw new Error('discarded value inspected'); }};
check(await Observable.from([discarded, 3]).last() === 3, 'last never assimilates overwritten values');
let subscriber, reads = 0, resolveValue;
const ac = new AbortController(), source = new Observable(s => { subscriber = s; });
const promise = source.last({signal: ac.signal});
const value = {get then() { reads++; ac.abort('too late'); return resolve => { resolveValue = resolve; }; }};
subscriber.next(value);
check(reads === 0, 'last then getter waits until completion');
subscriber.complete();
check(reads === 1 && !subscriber.active, 'last assimilates after source closure');
await Promise.resolve(); resolveValue(41);
check(await promise === 41, 'last locks result before reentrant then getter abort');
});
await test('arrays preserve values and own properties', async () => {
const marker = {}, thenable = {get then() { throw marker; }}, promise = Promise.resolve(9), symbol = Symbol();
const values = [thenable, promise, undefined, null, symbol, 7n];
const result = await Observable.from(values).toArray();
check(result.length === values.length && result.every((value, i) => value === values[i]), 'toArray does not assimilate individual values');
check(Object.getPrototypeOf(result) === Array.prototype, 'toArray result uses intrinsic Array prototype');
const index = Object.getOwnPropertyDescriptor(result, '0');
check(index.writable && index.enumerable && index.configurable, 'result element is own data property');
const source = Observable.from([]), a = await source.toArray(), b = await source.toArray();
check(a.length === 0 && b.length === 0 && a !== b, 'each empty subscription returns a fresh array');
let subscriber, reads = 0;
const ac = new AbortController(), pending = new Observable(s => { subscriber = s; }).toArray({signal: ac.signal});
const saved = Object.getOwnPropertyDescriptor(Array.prototype, 'then');
let resultArray;
try {
Object.defineProperty(Array.prototype, 'then', {configurable: true, get() { reads++; resultArray = this; ac.abort('too late'); }});
subscriber.next(1); subscriber.next(2);
check(reads === 0, 'private collection does not read then');
subscriber.complete();
} finally {
if (saved) Object.defineProperty(Array.prototype, 'then', saved); else delete Array.prototype.then;
}
check(await pending === resultArray && reads === 1 && resultArray.length === 2, 'result array assimilation locks before caller abort');
});
await test('shared producer', async () => {
let starts = 0, subscriber;
const source = new Observable(s => { starts++; subscriber = s; });
const ac = new AbortController(), last = rejected(source.last({signal: ac.signal})), all = source.toArray(), first = source.first();
subscriber.next(1);
check(subscriber.active && starts === 1 && await first === 1, 'first leaves collecting observers active');
subscriber.next(2); ac.abort('cancel last');
check(await last === 'cancel last' && subscriber.active, 'last cancellation preserves other observers');
subscriber.next(3); subscriber.complete();
same(await all, [1, 2, 3], 'toArray keeps collecting shared values');
});
await test('native intrinsics', async () => {
for (const name of methods) {
const source = Observable.from([1, 2]), method = Observable.prototype[name];
const globals = ['Promise', 'Array', 'Observable'], saved = globals.map(key => globalThis[key]);
const targets = [[Observable.prototype, 'subscribe'], [Array.prototype, 'push'], [Array.prototype, '0']];
const descriptors = targets.map(([object, key]) => Object.getOwnPropertyDescriptor(object, key));
const poison = () => { throw new Error('author method or setter invoked'); };
let promise;
try {
globals.forEach(key => { globalThis[key] = poison; });
targets.forEach(([object, key]) => Object.defineProperty(object, key, {configurable: true, set: poison}));
promise = method.call(source);
} finally {
globals.forEach((key, i) => { globalThis[key] = saved[i]; });
targets.forEach(([object, key], i) => {
if (descriptors[i]) Object.defineProperty(object, key, descriptors[i]); else delete object[key];
});
}
check(promise instanceof Promise, name + ' intrinsic Promise');
same(await promise, name === 'last' ? 2 : [1, 2], name + ' bypasses public subscription and inherited setters');
}
});
return {checks, failures};
})()