feat(dom): add serial Observable flatMap subscriptions

Queue source values until each mapped inner subscription completes, preserving
mapper indices, synchronous completion order, sharing and callback realms.
Trace the producer graph in V8 and release exhausted sources and queued values.

Finish shared abort algorithms before rethrowing the first IteratorClose
failure so cancellation always reaches both source and inner producers.
Cover Window/worker behavior, reentrancy, cancellation, realms and GC lifetimes.
This commit is contained in:
ldm0
2026-09-24 20:01:12 +08:00
parent c1f2c16d44
commit 1b05ab5f35
11 changed files with 751 additions and 44 deletions
@@ -4553,6 +4553,8 @@ dom/observable/tentative/observable-find.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-find.any.js?moli-wpt-any=window
dom/observable/tentative/observable-first.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-first.any.js?moli-wpt-any=window
dom/observable/tentative/observable-flatMap.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-flatMap.any.js?moli-wpt-any=window
dom/observable/tentative/observable-forEach.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-forEach.any.js?moli-wpt-any=window
dom/observable/tentative/observable-forEach.window.js?moli-wpt-script=window
+41 -1
View File
@@ -46,7 +46,7 @@ impl AbortAlgorithm {
/// Most internal abort algorithms cannot throw. Observable iterator closing
/// is an exception: its synchronous return() failure escapes AbortController.abort.
/// Keep this policy on the native callback, in the existing signal-owned list.
pub(crate) fn invoke_abort_algorithm<'s>(
fn invoke_abort_algorithm<'s>(
scope: &mut v8::PinScope<'s, '_>,
label: &str,
algorithm: v8::Local<'s, v8::Function>,
@@ -69,6 +69,46 @@ pub(crate) fn invoke_abort_algorithm<'s>(
}
}
/// A multi-source Observable shares a signal between its producers. A failing
/// IteratorClose must not leave later producers active. Finish the cancellation
/// snapshot before propagating its first exception to the abort caller.
pub(crate) fn invoke_abort_algorithms<'s>(
scope: &mut v8::PinScope<'s, '_>,
label: &str,
signal: v8::Local<'s, v8::Object>,
reason: v8::Local<'s, v8::Value>,
algorithms: Vec<AbortAlgorithm>,
) -> bool {
let mut first_error = None;
for algorithm in algorithms {
let Some(algorithm) = algorithm.prepare(scope) else {
continue;
};
let exception = {
v8::tc_scope!(let scope, scope);
if invoke_abort_algorithm(scope, label, algorithm, signal, reason) {
None
} else {
let Some(error) = scope.exception() else {
// Do not resume script execution after V8 termination.
return false;
};
scope.reset();
Some(error)
}
};
if first_error.is_none() {
first_error = exception;
}
}
if let Some(error) = first_error {
scope.throw_exception(error);
false
} else {
true
}
}
#[derive(Clone, Copy)]
enum AbortSignalOwner {
Window,
@@ -1,4 +1,4 @@
use crate::abort_signal_route::{AbortAlgorithm, invoke_abort_algorithm};
use crate::abort_signal_route::AbortAlgorithm;
pub(super) fn invoke_abort_algorithms<'s>(
scope: &mut v8::PinScope<'s, '_>,
@@ -7,21 +7,13 @@ pub(super) fn invoke_abort_algorithms<'s>(
abort_algorithms: Vec<AbortAlgorithm>,
) -> bool {
let signal = local_object_in_scope(scope, signal);
for algorithm in abort_algorithms {
let Some(algorithm) = algorithm.prepare(scope) else {
continue;
};
if !invoke_abort_algorithm(
scope,
"AbortSignal abort algorithm",
algorithm,
signal,
reason,
) {
return false;
}
}
true
crate::abort_signal_route::invoke_abort_algorithms(
scope,
"AbortSignal abort algorithm",
signal,
reason,
abort_algorithms,
)
}
pub(super) fn local_object_in_scope<'s>(
+4
View File
@@ -9,6 +9,7 @@ mod consume;
mod event_target;
mod finally;
mod first;
mod flat_map;
mod from;
mod inspect;
mod observer;
@@ -34,6 +35,8 @@ struct ObservablePrototype {
subscribe: (),
#[webapi(method, length = 1, callback = transform::map)]
map: (),
#[webapi(method = "flatMap", length = 1, callback = flat_map::flat_map)]
flat_map: (),
#[webapi(method, length = 1, callback = transform::filter)]
filter: (),
#[webapi(method, length = 1, callback = transform::take)]
@@ -318,6 +321,7 @@ fn subscribe_internal<'s>(
&& !until::subscribe(scope, observable, subscriber)
&& !inspect::subscribe(scope, observable, subscriber)
&& !finally::subscribe(scope, observable, subscriber)
&& !flat_map::subscribe(scope, observable, subscriber)
{
event_target::subscribe(scope, observable, subscriber);
}
+262
View File
@@ -0,0 +1,262 @@
//! Serial flattening keeps raw source values queued until the active inner
//! subscription completes. Queue cells and both observers are V8-traced.
use moli_webapi_declare::WebApiObject;
use super::{
callbacks, from,
observer::{self, Notification},
state::*,
subscribe_internal, subscriber_complete, subscriber_error, subscriber_next,
};
use crate::{
abort_signal_route::ResolvedAbortSignal,
util::{get_private_value, set_private_value},
webidl,
};
const SOURCE: &str = "__moliObservableFlatMapSource";
const MAPPER: &str = "__moliObservableFlatMapMapper";
const DOWNSTREAM: &str = "__moliFlatMapSubscriber";
const OWNER: &str = "__moliFlatMapSourceObserver";
const INNER: &str = "__moliFlatMapInnerObserver";
const BUSY: &str = "__moliFlatMapActiveInner";
const OUTER_COMPLETE: &str = "__moliFlatMapSourceCompleted";
const HEAD: &str = "__moliFlatMapQueueHead";
const TAIL: &str = "__moliFlatMapQueueTail";
const VALUE: &str = "__moliFlatMapQueuedValue";
const NEXT: &str = "__moliFlatMapQueueNext";
#[derive(webidl::WebIdlArgs)]
#[webidl(prefix = "Observable.flatMap")]
struct FlatMapArgs {
#[webidl(required, converter = "callback_function")]
mapper: webidl::WebIdlCallbackFunction,
}
pub(super) fn flat_map<'s>(
scope: &mut v8::PinScope<'s, '_>,
args: v8::FunctionCallbackArguments<'s>,
mut rv: v8::ReturnValue<'_, v8::Value>,
) {
let Some(parsed) = webidl::parse_args::<FlatMapArgs>(scope, &args) else {
return;
};
let Some(observable) = new_native_observable(scope, None) else {
return;
};
set_private_value(scope, observable, SOURCE, args.this().into());
set_callback(scope, observable, MAPPER, parsed.mapper);
rv.set(observable.into());
}
#[derive(WebApiObject)]
#[webapi(plain)]
struct SourceObserver<'scope> {
#[webapi(slot = observer::NATIVE_KIND)]
kind: i32,
#[webapi(slot = DOWNSTREAM)]
downstream: v8::Local<'scope, v8::Object>,
#[webapi(slot = MAPPER)]
mapper: v8::Local<'scope, v8::Object>,
}
#[derive(WebApiObject)]
#[webapi(plain)]
struct InnerObserver<'scope> {
#[webapi(slot = observer::NATIVE_KIND)]
kind: i32,
#[webapi(slot = OWNER)]
owner: v8::Local<'scope, v8::Object>,
}
#[derive(WebApiObject)]
#[webapi(plain)]
struct QueueEntry<'scope> {
#[webapi(slot = VALUE)]
value: v8::Local<'scope, v8::Value>,
}
pub(super) fn subscribe<'s>(
scope: &mut v8::PinScope<'s, '_>,
observable: v8::Local<'s, v8::Object>,
subscriber: v8::Local<'s, v8::Object>,
) -> bool {
let Some(source) = object_slot(scope, observable, SOURCE) else {
return false;
};
let mapper = object_slot(scope, observable, MAPPER).expect("flatMap mapper");
let Some(observer) = SourceObserver::new(observer::FLAT_MAP_SOURCE, subscriber, mapper)
.bind(scope)
.ok()
else {
return true;
};
if active(scope, subscriber) {
set_private_value(scope, subscriber, UPSTREAM_OBSERVER, observer.into());
}
let signal = signal(scope, subscriber);
subscribe_internal(scope, source, observer, Some(signal));
true
}
fn signal<'s>(
scope: &mut v8::PinScope<'s, '_>,
subscriber: v8::Local<'s, v8::Object>,
) -> ResolvedAbortSignal<'s> {
object_slot(scope, subscriber, SIGNAL)
.and_then(|signal| ResolvedAbortSignal::resolve(scope, signal))
.expect("flatMap Subscriber signal")
}
fn flag<'s>(
scope: &mut v8::PinScope<'s, '_>,
observer: v8::Local<'s, v8::Object>,
slot: &str,
) -> bool {
get_private_value(scope, observer, slot).is_some_and(|value| value.is_true())
}
fn enqueue<'s>(
scope: &mut v8::PinScope<'s, '_>,
observer: v8::Local<'s, v8::Object>,
value: v8::Local<'s, v8::Value>,
) {
let entry = QueueEntry::new(value)
.bind(scope)
.expect("flatMap queue entry");
if let Some(tail) = object_slot(scope, observer, TAIL) {
set_private_value(scope, tail, NEXT, entry.into());
} else {
set_private_value(scope, observer, HEAD, entry.into());
}
set_private_value(scope, observer, TAIL, entry.into());
}
fn dequeue<'s>(
scope: &mut v8::PinScope<'s, '_>,
observer: v8::Local<'s, v8::Object>,
) -> Option<v8::Local<'s, v8::Value>> {
let head = object_slot(scope, observer, HEAD)?;
let value = get_private_value(scope, head, VALUE).expect("flatMap queued value");
let next = object_slot(scope, head, NEXT);
let next_value = next.map_or_else(|| v8::undefined(scope).into(), Into::into);
set_private_value(scope, observer, HEAD, next_value);
if next.is_none() {
set_private_value(scope, observer, TAIL, v8::undefined(scope).into());
}
Some(value)
}
fn process_next<'s>(
scope: &mut v8::PinScope<'s, '_>,
observer: v8::Local<'s, v8::Object>,
value: v8::Local<'s, v8::Value>,
) {
let downstream = object_slot(scope, observer, DOWNSTREAM).expect("flatMap downstream");
let mapper = object_slot(scope, observer, MAPPER).expect("flatMap mapper");
let index = observer::index(scope, observer);
let index = v8::Number::new(scope, index as f64).into();
let mapped = match callbacks::invoke_value(scope, mapper, &[value, index]) {
Ok(mapped) => mapped,
Err(error) => {
subscriber_error(scope, downstream, error);
return;
}
};
observer::increment_index(scope, observer);
let (inner, error) = {
v8::tc_scope!(let scope, scope);
let inner = from::convert(scope, mapped);
let error = scope.exception();
scope.reset();
(inner, error)
};
if let Some(error) = error {
subscriber_error(scope, downstream, error);
return;
}
let Some(inner) = inner else {
return;
};
let Some(inner_observer) = InnerObserver::new(observer::FLAT_MAP_INNER, observer)
.bind(scope)
.ok()
else {
return;
};
set_private_value(scope, observer, INNER, inner_observer.into());
// Conversion may cancel downstream. The source still receives its fresh,
// inactive Subscriber, just as an explicitly pre-aborted subscribe does.
let signal = signal(scope, downstream);
subscribe_internal(scope, inner, inner_observer, Some(signal));
}
pub(super) fn notify<'s>(
scope: &mut v8::PinScope<'s, '_>,
observer: v8::Local<'s, v8::Object>,
notification: Notification<'s>,
kind: i32,
) {
let source_observer = if kind == observer::FLAT_MAP_INNER {
object_slot(scope, observer, OWNER).expect("flatMap source observer")
} else {
observer
};
let downstream = object_slot(scope, source_observer, DOWNSTREAM).expect("flatMap downstream");
match notification {
Notification::Next(value) if kind == observer::FLAT_MAP_SOURCE => {
if flag(scope, source_observer, BUSY) {
enqueue(scope, source_observer, value);
} else {
// Set before invoking mapper: reentrant source.next queues its
// raw value and does not overlap the current mapper or inner.
set_private_value(
scope,
source_observer,
BUSY,
v8::Boolean::new(scope, true).into(),
);
process_next(scope, source_observer, value);
}
}
Notification::Next(value) => subscriber_next(scope, downstream, value),
Notification::Error(error) => subscriber_error(scope, downstream, error),
Notification::Complete => {
set_private_value(
scope,
observer,
observer::SUBSCRIBER,
v8::undefined(scope).into(),
);
if kind == observer::FLAT_MAP_SOURCE {
set_private_value(
scope,
source_observer,
OUTER_COMPLETE,
v8::Boolean::new(scope, true).into(),
);
if !flag(scope, source_observer, BUSY) {
subscriber_complete(scope, downstream);
}
} else {
set_private_value(scope, source_observer, INNER, v8::undefined(scope).into());
if let Some(value) = dequeue(scope, source_observer) {
// Run before this complete() returns, including when the
// next inner completes synchronously and reenters here.
process_next(scope, source_observer, value);
} else {
set_private_value(
scope,
source_observer,
BUSY,
v8::Boolean::new(scope, false).into(),
);
if flag(scope, source_observer, OUTER_COMPLETE) {
subscriber_complete(scope, downstream);
}
}
}
}
}
}
+7 -2
View File
@@ -2,8 +2,8 @@
//! while script callbacks keep their typed Web IDL invocation boundary.
use super::{
callbacks, collect, consume, finally, first, inspect, invoke_and_report, state::*, transform,
until,
callbacks, collect, consume, finally, first, flat_map, inspect, invoke_and_report, state::*,
transform, until,
};
use crate::util::{get_private_value, set_private_value};
@@ -24,6 +24,8 @@ pub(super) const UNTIL_SOURCE: i32 = 13;
pub(super) const UNTIL_NOTIFIER: i32 = 14;
pub(super) const INSPECT: i32 = 15;
pub(super) const FINALLY: i32 = 16;
pub(super) const FLAT_MAP_SOURCE: i32 = 17;
pub(super) const FLAT_MAP_INNER: i32 = 18;
pub(super) const SUBSCRIBER: &str = "__moliObservableNativeSubscriber";
const INDEX: &str = "__moliObservableCallbackIndex";
@@ -93,6 +95,9 @@ pub(super) fn notify<'s>(
UNTIL_SOURCE | UNTIL_NOTIFIER => until::notify(scope, observer, notification, kind),
INSPECT => inspect::notify(scope, observer, notification),
FINALLY => finally::notify(scope, observer, notification),
FLAT_MAP_SOURCE | FLAT_MAP_INNER => {
flat_map::notify(scope, observer, notification, kind)
}
_ => unreachable!("unknown native Observable observer"),
}
let exception = scope.exception();
@@ -1,5 +1,109 @@
use super::*;
#[test]
fn observable_flat_map_preserves_serial_order_conversion_reentrancy_and_cancellation() {
let mut vm = new_storage_test_vm("https://observable-flat-map.test/");
vm.eval(&format!(
"({}).then(value => {{ globalThis.flatMapResult = JSON.stringify(value); }});",
include_str!("../../../tests/fixtures/observable-flat-map.js")
))
.expect("Observable.flatMap fixture should evaluate");
let result: serde_json::Value =
serde_json::from_str(&vm.eval("flatMapResult").unwrap()).unwrap();
assert_eq!(result["failures"], serde_json::json!([]), "{result}");
assert!(result["checks"].as_u64().unwrap() >= 120, "{result}");
}
#[test]
fn observable_flat_map_preserves_result_conversion_callback_and_cancellation_realms() {
let mut vm = new_storage_test_vm("https://observable-flat-map-realms.test/");
vm.eval("document.appendChild(document.createElement('iframe'))")
.unwrap();
materialize_single_child_default_realm_for_test(&mut vm, "Observable.flatMap realm");
vm.eval(&format!(
"({}).then(value => {{ globalThis.flatMapRealms = JSON.stringify(value); }});",
include_str!("../../../tests/fixtures/observable-flat-map-realms.js")
))
.expect("Observable.flatMap realms fixture should evaluate");
let result: serde_json::Value =
serde_json::from_str(&vm.eval("flatMapRealms").unwrap()).unwrap();
assert_eq!(result["failures"], serde_json::json!([]), "{result}");
assert_eq!(result["checks"], 31, "{result}");
}
#[test]
fn observable_flat_map_traces_both_producers_and_releases_queues_and_closed_graphs() {
let mut vm = new_storage_test_vm("https://observable-flat-map-gc.test/");
vm.eval(r#"
function makeFlatMapSource() {
let subscriber;
return {source: new Observable(s => { subscriber = s; }), get subscriber() { return subscriber; }};
}
function makeFlatMapper(inner) {
const token = {values: []};
return {token, callback: value => value === 0 ? inner : token.values};
}
globalThis.flatMapChains = [];
for (const mode of ['abandoned', 'complete', 'outer-error', 'inner-error', 'abort']) (() => {
const outer = makeFlatMapSource(), inner = makeFlatMapSource(), mapper = makeFlatMapper(inner.source);
const result = outer.source.flatMap(mapper.callback), ac = new AbortController(), queued = [{}, {}];
const promise = result.toArray(mode === 'abort' ? {signal: ac.signal} : undefined);
promise.catch(() => {});
outer.subscriber.next(0); inner.subscriber.next(1);
for (const value of queued) outer.subscriber.next(value);
const entry = {mode, templates: [outer.source, result].map(v => new WeakRef(v)),
outer: new WeakRef(outer.subscriber), inner: new WeakRef(inner.subscriber),
callbacks: [mapper.callback, mapper.token].map(v => new WeakRef(v)), queued: queued.map(v => new WeakRef(v))};
if (mode !== 'abandoned') entry.promise = promise;
if (mode === 'abort') entry.controller = ac;
flatMapChains.push(entry);
})();
"#).unwrap();
let collect = |vm: &mut StandaloneScriptVmHarness| {
vm.renderer_document_isolate
.clone()
.with_entered_renderer_document_isolate(|isolate| {
isolate.clear_kept_objects();
isolate.low_memory_notification();
Ok(())
})
.unwrap();
};
collect(&mut vm);
assert_eq!(vm.eval(r#"JSON.stringify([
flatMapChains.every(c => c.templates.every(ref => ref.deref() === undefined)),
flatMapChains.every(c => [c.outer, c.inner, ...c.callbacks, ...c.queued].every(ref => (ref.deref() !== undefined) === (c.mode !== 'abandoned')))
])"#).unwrap(), "[true,true]");
vm.eval(r#"
for (const c of flatMapChains.filter(c => c.promise)) {
c.promise.then(values => { c.correct = c.mode === 'complete' && JSON.stringify(values) === '[1,2]'; }, error => { c.correct = error === c.mode; });
const outer = c.outer.deref(), inner = c.inner.deref(); inner.next(2);
if (c.mode === 'complete') { outer.complete(); inner.complete(); }
else if (c.mode === 'outer-error') outer.error(c.mode);
else if (c.mode === 'inner-error') inner.error(c.mode);
else c.controller.abort(c.mode);
c.closed = [outer, inner];
}
"#).unwrap();
assert_eq!(vm.eval("flatMapChains.filter(c => c.promise).every(c => c.correct && c.closed.every(s => !s.active))").unwrap(), "true");
collect(&mut vm);
assert_eq!(vm.eval("flatMapChains.every(c => [...c.callbacks, ...c.queued].every(ref => ref.deref() === undefined))").unwrap(), "true");
vm.eval(r#"
globalThis.flatMapDrain = (() => {
const outer = makeFlatMapSource(), inner = makeFlatMapSource(), mapper = makeFlatMapper(inner.source), value = {};
const promise = outer.source.flatMap(mapper.callback).toArray();
outer.subscriber.next(0); outer.subscriber.next(value); outer.subscriber.complete();
return {promise, outer: new WeakRef(outer.subscriber), inner: new WeakRef(inner.subscriber),
value: new WeakRef(value), callback: new WeakRef(mapper.callback)};
})();
"#).unwrap();
collect(&mut vm);
assert_eq!(vm.eval("flatMapDrain.outer.deref() === undefined && flatMapDrain.inner.deref() !== undefined && flatMapDrain.value.deref() !== undefined && flatMapDrain.callback.deref() !== undefined").unwrap(), "true");
vm.eval("flatMapDrain.promise.then(v => { flatMapDrain.correct = v.length === 0; }); flatMapDrain.inner.deref().complete();").unwrap();
collect(&mut vm);
assert_eq!(vm.eval("flatMapDrain.correct && flatMapDrain.inner.deref() === undefined && flatMapDrain.value.deref() === undefined && flatMapDrain.callback.deref() === undefined").unwrap(), "true");
}
#[test]
fn observable_finally_preserves_teardown_order_sharing_and_reentrant_cancellation() {
let mut vm = new_storage_test_vm("https://observable-finally.test/");
+8 -25
View File
@@ -1,7 +1,7 @@
use std::collections::{HashMap, HashSet};
use std::{cell::RefCell, rc::Rc};
use crate::abort_signal_route::{AbortAlgorithm, invoke_abort_algorithm};
use crate::abort_signal_route::{AbortAlgorithm, invoke_abort_algorithms};
use crate::util::{get_private_value, set_private_value, v8str};
use crate::webidl;
@@ -297,7 +297,13 @@ fn run_worker_abort_steps<'s>(
std::mem::take(&mut state.abort_algorithms)
};
reject_worker_fetches_for_signal(scope, signal_id, reason);
if !invoke_worker_abort_algorithms(scope, signal, reason, abort_algorithms) {
if !invoke_abort_algorithms(
scope,
"Worker AbortSignal abort algorithm",
signal,
reason,
abort_algorithms,
) {
return false;
}
// Dispatch reads the shared listener registry after all abort algorithms.
@@ -305,29 +311,6 @@ fn run_worker_abort_steps<'s>(
true
}
fn invoke_worker_abort_algorithms<'s>(
scope: &mut v8::PinScope<'s, '_>,
signal: v8::Local<'s, v8::Object>,
reason: v8::Local<'s, v8::Value>,
abort_algorithms: Vec<AbortAlgorithm>,
) -> bool {
for algorithm in abort_algorithms {
let Some(algorithm) = algorithm.prepare(scope) else {
continue;
};
if !invoke_abort_algorithm(
scope,
"Worker AbortSignal abort algorithm",
algorithm,
signal,
reason,
) {
return false;
}
}
true
}
fn create_signal<'s>(
scope: &mut v8::PinScope<'s, '_>,
store: &mut WorkerAbortStore,
@@ -1,5 +1,22 @@
use super::*;
#[tokio::test]
async fn worker_observable_flat_map_preserves_serial_order_conversion_reentrancy_and_cancellation()
{
ensure_v8();
let mut handle = spawn_worker(
format!(
"({}).then(value => {{ postMessage(value); close(); }});",
include_str!("../../../../tests/fixtures/observable-flat-map.js")
),
"https://observable-flat-map.test/worker.js".into(),
);
let message = timeout(TIMEOUT, handle.recv()).await.unwrap().unwrap();
let result: serde_json::Value = serde_json::from_str(&expect_post_json(message)).unwrap();
assert_eq!(result["failures"], serde_json::json!([]), "{result}");
assert!(result["checks"].as_u64().unwrap() >= 120, "{result}");
}
#[tokio::test]
async fn worker_observable_finally_preserves_teardown_order_sharing_and_reentrant_cancellation() {
ensure_v8();
@@ -0,0 +1,58 @@
(async () => {
const failures = [];
let checks = 0;
const check = (value, label) => { checks++; if (!value) failures.push(label); };
const child = document.querySelector('iframe').contentWindow;
const method = child.Observable.prototype.flatMap, source = Observable.from([1, 2]);
child.mapperCalls = [];
const callback = child.Function('value', 'index', 'mapperCalls.push([value,index]); globalThis.mapperThis = this; globalThis.mapperArgc = arguments.length; return [value*10];');
const result = method.call(source, callback);
check(result instanceof child.Observable && !(result instanceof Observable), 'result uses callee realm');
check(Object.getPrototypeOf(result) === child.Observable.prototype, 'callee intrinsic prototype');
check(child.mapperCalls.length === 0, 'foreign mapper lazy');
const values = await result.toArray();
check(values instanceof child.Array && values.join(',') === '10,20', 'callee Array result');
check(child.mapperThis === child && child.mapperArgc === 2, 'callback own realm and argument count');
check(JSON.stringify(child.mapperCalls) === '[[1,0],[2,1]]', 'foreign mapper indices');
const local = Observable.prototype.flatMap.call(child.Observable.from([3]), value => [value]);
check(local instanceof Observable && !(local instanceof child.Observable), 'local method returns local Observable');
check(Object.getPrototypeOf(local) === Observable.prototype, 'local intrinsic prototype');
const localValues = await local.toArray();
check(localValues instanceof Array && localValues[0] === 3, 'local Array result');
const revoked = Proxy.revocable(source, {}); revoked.revoke();
for (const receiver of [{}, Object.create(source), new Proxy(source, {}), revoked.proxy]) {
let error; try { method.call(receiver, () => []); } catch (e) { error = e; }
check(error instanceof child.TypeError && !(error instanceof TypeError), 'receiver check uses callee TypeError');
}
for (const input of [[], [undefined], [null], [1], [{}]]) {
let error; try { method.apply(source, input); } catch (e) { error = e; }
check(error instanceof child.TypeError && !(error instanceof TypeError), 'callback conversion uses callee TypeError');
}
for (const value of [null, 1, {}, {then() {}}]) {
const error = await method.call(source, () => value).toArray().catch(e => e);
check(error instanceof child.TypeError && !(error instanceof TypeError), 'mapper-result conversion error uses callee TypeError');
}
const marker = new child.Error('mapper'); child.mapperError = marker;
for (const mapper of [child.Function('throw mapperError'), () => ({get [Symbol.iterator]() { throw marker; }}),
() => new child.Observable(s => s.error(marker)), () => child.Promise.reject(marker)]) {
const error = await method.call(source, mapper).toArray().catch(e => e);
check(error === marker && error instanceof child.Error, 'callback, conversion and inner errors retain identity');
}
const frameSource = new child.Observable(s => { child.pendingOuter = s; });
const inner = new child.Observable(s => { child.pendingInner = s; });
const ac = new AbortController(), reason = {}, order = [];
const pending = Observable.prototype.flatMap.call(frameSource, () => inner).toArray({signal: ac.signal});
const outcome = pending.catch(e => e);
child.pendingOuter.next(1);
child.pendingOuter.addTeardown(() => order.push('outer'));
child.pendingInner.addTeardown(() => order.push('inner'));
ac.abort(reason);
check(await outcome === reason, 'foreign pending graph preserves cancellation reason');
check(!child.pendingOuter.active && !child.pendingInner.active, 'cancellation reaches both foreign producers');
check(child.pendingOuter.signal.reason === reason && child.pendingInner.signal.reason === reason, 'foreign producer signals preserve reason');
check(order.join(',') === 'outer,inner', 'foreign cleanup order');
Object.setPrototypeOf(source, null);
const branded = method.call(source, value => [value]);
check(branded instanceof child.Observable && (await branded.toArray()).join(',') === '1,2', 'native brand survives prototype replacement');
return {checks, failures};
})()
+240
View File
@@ -0,0 +1,240 @@
(async () => {
'use strict';
const failures = [];
let checks = 0;
const check = (value, label) => { checks++; if (!value) failures.push(label); };
const same = (a, b, label) => check(JSON.stringify(a) === JSON.stringify(b), label);
const thrown = fn => { try { fn(); } catch (e) { return e; } };
const test = async (label, fn) => { try { await fn(); } catch (e) { check(false, label + ': ' + e); } };
const method = Observable.prototype.flatMap;
check(typeof method === 'function', 'flatMap exposed');
if (failures.length) return {checks, failures};
function subject() {
let subscriber, starts = 0;
return {source: new Observable(s => { subscriber = s; starts++; }),
get subscriber() { return subscriber; }, get starts() { return starts; }};
}
await test('Web IDL conversion and native operations', async () => {
const desc = Object.getOwnPropertyDescriptor(Observable.prototype, 'flatMap');
check(method.name === 'flatMap' && method.length === 1, 'name and length');
check(desc.enumerable && desc.writable && desc.configurable, 'descriptor');
check(thrown(() => new method(() => [])) instanceof TypeError, 'not a constructor');
const source = Observable.from([1]);
let traps = 0;
const revoked = Proxy.revocable(source, {}); revoked.revoke();
for (const receiver of [undefined, null, false, 1, Symbol(), {}, Object.create(source), Object.create(Observable.prototype),
new Proxy(source, {get() { traps++; }}), revoked.proxy]) {
check(thrown(() => method.call(receiver, () => [])) instanceof TypeError, 'invalid receiver rejected');
}
check(traps === 0, 'receiver check does not invoke Proxy traps');
check(thrown(() => source.flatMap()) instanceof TypeError, 'required mapper');
for (const callback of [undefined, null, false, 1, 1n, '', Symbol(), {}, [], {handleEvent() {}}]) {
check(thrown(() => source.flatMap(callback)) instanceof TypeError, 'non-callable mapper rejected');
}
let calls = 0;
const mapper = new Proxy(function(value, index) {
calls++;
check(this === undefined && arguments.length === 2, 'mapper receiver and arguments');
check(index === 0, 'fresh subscription resets index');
return [value, value + 1];
}, {get() { throw 'mapper property read'; }});
Object.defineProperty(source, 'constructor', {get() { throw 'constructor read'; }});
Object.setPrototypeOf(source, null);
const result = method.call(source, mapper, {get signal() { throw 'extra argument read'; }});
check(calls === 0, 'creation is lazy');
check(Object.getPrototypeOf(result) === Observable.prototype && Observable.from(result) === result, 'intrinsic branded result');
same(await result.toArray(), [1, 2], 'source prototype and species ignored');
same(await result.toArray(), [1, 2], 'reusable result');
check(calls === 2, 'one mapper call per source value');
class Derived extends Observable {}
check(!(new Derived(s => s.complete()).flatMap(mapper) instanceof Derived), 'does not use source species');
const saved = [];
for (const [object, key] of [[Observable, 'from'], [Observable.prototype, 'subscribe'],
[Subscriber.prototype, 'next'], [Subscriber.prototype, 'error'], [Subscriber.prototype, 'complete']]) {
saved.push([object, key, Object.getOwnPropertyDescriptor(object, key)]);
Object.defineProperty(object, key, {value() { throw key + ' invoked'; }, configurable: true});
}
try { same(await result.toArray(), [1, 2], 'internal conversion and notifications ignore public methods'); }
finally { for (const [object, key, descriptor] of saved) Object.defineProperty(object, key, descriptor); }
});
await test('conversion of mapper results', async () => {
for (const inner of [Observable.from([3]), [3], new Set([3]), Promise.resolve(3),
(async function* () { yield 3; })()]) {
same(await Observable.from([1]).flatMap(() => inner).toArray(), [3], 'maps convertible input');
}
for (const value of [undefined, null, false, 1, 1n, '', 'abc', Symbol(), {}, {then() { throw 'assimilated'; }}]) {
check(await Observable.from([1]).flatMap(() => value).toArray().catch(e => e) instanceof TypeError, 'invalid result rejects through observer');
}
const marker = {}, log = [];
const iterable = {get [Symbol.asyncIterator]() { log.push('async'); return undefined; },
get [Symbol.iterator]() { log.push('sync'); return function* () { log.push('open'); yield marker; }; }};
const values = await Observable.from([0]).flatMap(() => iterable).toArray();
check(values.length === 1 && values[0] === marker, 'inner value identity');
same(log, ['async', 'sync', 'sync', 'open'], 'conversion probes before subscription obtains iterator');
let thenReads = 0;
const preferred = {[Symbol.iterator]: function* () { yield 5; }, get then() { thenReads++; throw 'then read'; }};
same(await Observable.from([0]).flatMap(() => preferred).toArray(), [5], 'iterable result bypasses then property');
check(thenReads === 0, 'does not assimilate arbitrary thenables');
});
await test('raw queue, delayed mapper and synchronous completion order', () => {
const outer = subject(), inners = [], log = [], indices = [];
const result = outer.source.flatMap((value, index) => {
const id = value.id;
indices.push(index); log.push('map' + id);
return new Observable(s => {
inners.push(s); log.push('start' + id); s.addTeardown(() => log.push('cleanup' + id)); s.next(id);
if (id > 1) { s.complete(); log.push('after' + id); }
});
});
const values = [];
result.subscribe({next: v => values.push(v), complete: () => log.push('complete')});
const queued = {id: 20};
outer.subscriber.next({id: 1}); outer.subscriber.next(queued); outer.subscriber.next({id: 3});
same(log, ['map1', 'start1'], 'queued values do not invoke mapper early');
queued.id = 2; outer.subscriber.complete();
check(inners.length === 1 && inners[0].active, 'outer completion waits for active inner');
inners[0].complete(); log.push('after1');
same(values, [1, 2, 3], 'queue retains raw value identity');
same(indices, [0, 1, 2], 'serial mapper indices');
same(log, ['map1', 'start1', 'cleanup1', 'map2', 'start2', 'cleanup2', 'map3', 'start3', 'cleanup3',
'complete', 'after3', 'after2', 'after1'], 'drains next inner before previous complete returns');
});
await test('reentrant mapper, conversion and inner delivery', async () => {
const outer = subject(), mapped = [], values = [];
outer.source.flatMap((value, index) => {
mapped.push([value, index]);
if (value === 1) outer.subscriber.next(2);
return {[Symbol.iterator]() {
if (value === 1) outer.subscriber.next(3);
return [value][Symbol.iterator]();
}};
}).subscribe(value => { values.push(value); if (value === 1) outer.subscriber.next(4); });
outer.subscriber.next(1); outer.subscriber.complete();
same(mapped, [[1, 0], [2, 1], [3, 2], [4, 3]], 'reentrant pushes wait until mapper and inner finish');
same(values, [1, 2, 3, 4], 'reentrant output order');
const pending = subject(), gate = subject(), queued = [];
const result = pending.source.flatMap(value => { queued.push(value); return value === 0 ? gate.source : [value]; });
const promise = result.toArray();
for (let i = 0; i < 128; i++) pending.subscriber.next(i);
pending.subscriber.complete();
same(queued, [0], 'large queue maps lazily');
gate.subscriber.next(0); gate.subscriber.complete();
const collected = await promise;
check(collected.length === 128 && collected.every((v, i) => v === i), 'drains buffered synchronous inputs in FIFO order');
});
await test('sharing, distinct branches and last-consumer cancellation', () => {
const outer = subject(), inner = subject(), ac1 = new AbortController(), ac2 = new AbortController();
let maps = 0, otherMaps = 0;
const result = outer.source.flatMap(() => { maps++; return inner.source; });
const first = [], second = [];
result.subscribe(v => first.push(v), {signal: ac1.signal});
result.subscribe(v => second.push(v), {signal: ac2.signal});
outer.source.flatMap(() => { otherMaps++; return []; }).subscribe();
outer.subscriber.next(1); inner.subscriber.next(2); ac1.abort('first'); inner.subscriber.next(3);
check(maps === 1 && otherMaps === 1 && outer.starts === 1 && inner.starts === 1, 'shared result maps once and distinct branch separately');
check(outer.subscriber.active && inner.subscriber.active, 'first cancellation keeps both subscriptions');
same(first, [2], 'first consumer removed'); same(second, [2, 3], 'second consumer survives');
const reason = {}; ac2.abort(reason);
check(outer.subscriber.active && !inner.subscriber.active && inner.subscriber.signal.reason === reason, 'last result consumer only cancels its branch');
const oldInner = inner.subscriber;
result.subscribe(); outer.subscriber.next(2);
check(maps === 2 && otherMaps === 2 && inner.starts === 2 && inner.subscriber !== oldInner, 'resubscription uses fresh queue and inner');
inner.subscriber.complete(); outer.subscriber.complete();
});
await test('original errors, queue disposal and synchronous cancellation', () => {
const reports = [], onerror = e => { reports.push(e.error); e.preventDefault(); };
addEventListener('error', onerror);
try {
for (const mode of ['outer', 'inner', 'mapper', 'conversion', 'initializer']) {
for (const marker of [{}, null, undefined]) {
const outer = subject(), inner = subject(), log = [], errors = [];
let maps = 0, complete = 0;
outer.source.flatMap(() => {
maps++;
if (mode === 'mapper') throw marker;
if (mode === 'conversion') return {get [Symbol.asyncIterator]() { throw marker; }};
if (mode === 'initializer') return new Observable(() => { throw marker; });
return inner.source;
}).subscribe({error: e => { errors.push(e); log.push('error'); }, complete: () => complete++});
outer.subscriber.addTeardown(() => log.push('outer'));
outer.subscriber.next(1);
if (inner.subscriber) {
inner.subscriber.addTeardown(() => log.push('inner'));
outer.subscriber.next(2);
if (mode === 'outer') outer.subscriber.error(marker); else inner.subscriber.error(marker);
}
check(errors.length === 1 && errors[0] === marker && complete === 0, mode + ' error identity');
check(!outer.subscriber.active && (!inner.subscriber || !inner.subscriber.active), mode + ' closes both subscriptions');
check(maps === 1, mode + ' abandons queued values');
same(log, mode === 'outer' ? ['outer', 'inner', 'error'] : mode === 'inner' ? ['inner', 'outer', 'error'] : ['outer', 'error'], mode + ' cleanup precedes observer error');
}
}
same(reports, [], 'handled errors are not reported globally');
const outer = subject(), inner = subject(), ac = new AbortController(), reason = {}, cleanupError = {};
let maps = 0;
outer.source.flatMap(() => { maps++; return inner.source; }).subscribe({}, {signal: ac.signal});
outer.subscriber.next(1); outer.subscriber.next(2);
outer.subscriber.addTeardown(() => { throw cleanupError; });
ac.abort(reason);
check(!outer.subscriber.active && !inner.subscriber.active && maps === 1, 'explicit cancellation closes both and drops queue');
check(outer.subscriber.signal.reason === reason && inner.subscriber.signal.reason === reason, 'both receive same cancellation reason');
check(reports.length === 1 && reports[0] === cleanupError, 'teardown exception does not skip inner cancellation');
} finally { removeEventListener('error', onerror); }
});
await test('throwing IteratorClose still cancels both producers', () => {
for (const firstError of [{}, undefined]) {
const secondError = {}, ac = new AbortController(), log = [], reason = {};
let outerClosed = 0, innerClosed = 0, innerPulls = 0, caught, didThrow = false;
const outer = {[Symbol.iterator]() { return {
next() { return {value: 1}; },
return() { outerClosed++; log.push('outer return'); throw firstError; }
}; }};
const inner = {[Symbol.iterator]() { return {
next() { return ++innerPulls > 3 ? {done: true} : {value: innerPulls}; },
return() { innerClosed++; log.push('inner return'); throw secondError; }
}; }};
Observable.from(outer).flatMap(() => inner).finally(() => log.push('finally')).subscribe(() => {
try { ac.abort(reason); } catch (e) { didThrow = true; caught = e; }
log.push('after abort');
}, {signal: ac.signal});
check(didThrow && caught === firstError, 'first IteratorClose failure preserved, including undefined');
check(outerClosed === 1 && innerClosed === 1 && innerPulls === 1, 'both iterators close once before any extra pulls');
same(log, ['outer return', 'inner return', 'finally', 'after abort'], 'cleanup and finalizer finish before abort rethrows');
}
const ac = new AbortController(), marker = {}, first = subject(), second = subject(), log = [];
first.source.subscribe({}, {signal: ac.signal}); second.source.subscribe({}, {signal: ac.signal});
Observable.from({[Symbol.iterator]() { return {
next() { return {value: 1}; }, return() { throw marker; }
}; }}).subscribe(() => {
const third = new Observable(s => s.addTeardown(() => log.push('third')));
third.subscribe({}, {signal: ac.signal});
check(thrown(() => ac.abort()) === marker, 'shared direct signal preserves failing cancellation');
}, {signal: ac.signal});
check(!first.subscriber.active && !second.subscriber.active, 'earlier consumers also closed');
same(log, ['third'], 'later signal algorithm runs after failure without flatMap');
});
await test('pre-abort and cancellation inside mapper', () => {
const ac = new AbortController(), reason = {}, pre = subject(); let maps = 0, inactive;
pre.source.flatMap(() => { maps++; return []; }).subscribe({}, {signal: AbortSignal.abort(reason)});
check(pre.starts === 1 && !pre.subscriber.active && pre.subscriber.signal.reason === reason && maps === 0, 'pre-aborted source initialized inactive without mapping');
const outer = subject();
outer.source.flatMap(() => { ac.abort(reason); return new Observable(s => { inactive = s; }); }).subscribe({}, {signal: ac.signal});
outer.subscriber.next(1);
check(!outer.subscriber.active && inactive && !inactive.active && inactive.signal.reason === reason, 'mapper cancellation still initializes returned Observable inactive');
const captured = subject(), controller = new AbortController(); let calls = 0, inner;
captured.source.subscribe(() => controller.abort(reason));
captured.source.flatMap(() => { calls++; return new Observable(s => { inner = s; }); }).subscribe({}, {signal: controller.signal});
captured.subscriber.next(1);
check(calls === 1 && inner && !inner.active, 'captured next notification still maps after an earlier observer cancels');
captured.subscriber.complete();
});
return {checks, failures};
})()