mirror of
https://github.com/lexmount/moli.git
synced 2026-09-24 16:01:31 +00:00
feat(dom): add Observable switchMap subscriptions
Cancel the previous inner before mapping each source value, convert mapper results through the native Observable path, and wait for active inner completion after the source completes. Preserve the draft's shared controller reference under reentrant mapper, conversion and teardown callbacks. Trace each live inner and release replaced observers without persistent roots outside the existing signal stores. Finish Window and Worker dependent abort dispatch after IteratorClose failures, then propagate the first exception. Cover Window and Worker conversion, switching, cancellation, cross-realm receivers and callbacks, GC, and reentrant synchronous completion. Verify 61 related WPT cases: 576/592 to 590/592 subtests, no regressions. Add switchMap and finally Window/Worker cases to the verified pass catalogue.
This commit is contained in:
@@ -4549,6 +4549,8 @@ dom/observable/tentative/observable-every.any.js?moli-wpt-any=dedicatedworker
|
||||
dom/observable/tentative/observable-every.any.js?moli-wpt-any=window
|
||||
dom/observable/tentative/observable-filter.any.js?moli-wpt-any=dedicatedworker
|
||||
dom/observable/tentative/observable-filter.any.js?moli-wpt-any=window
|
||||
dom/observable/tentative/observable-finally.any.js?moli-wpt-any=dedicatedworker
|
||||
dom/observable/tentative/observable-finally.any.js?moli-wpt-any=window
|
||||
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
|
||||
@@ -4569,6 +4571,8 @@ 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-some.any.js?moli-wpt-any=dedicatedworker
|
||||
dom/observable/tentative/observable-some.any.js?moli-wpt-any=window
|
||||
dom/observable/tentative/observable-switchMap.any.js?moli-wpt-any=dedicatedworker
|
||||
dom/observable/tentative/observable-switchMap.any.js?moli-wpt-any=window
|
||||
dom/observable/tentative/observable-take.any.js?moli-wpt-any=dedicatedworker
|
||||
dom/observable/tentative/observable-take.any.js?moli-wpt-any=window
|
||||
dom/observable/tentative/observable-takeUntil.any.js?moli-wpt-any=dedicatedworker
|
||||
|
||||
@@ -300,11 +300,29 @@ impl AbortStore {
|
||||
signals_to_abort.push((dependent_signal_id, v8::Local::new(scope, signal)));
|
||||
}
|
||||
}
|
||||
// Complete the dependency snapshot even if a source's IteratorClose
|
||||
// throws. Dependents were already marked aborted and cannot be retried.
|
||||
let mut first_error = None;
|
||||
for (signal_id, signal) in signals_to_abort {
|
||||
if !self.run_abort_steps(scope, host, signal, signal_id, reason) {
|
||||
return;
|
||||
let error = {
|
||||
v8::tc_scope!(let scope, scope);
|
||||
if self.run_abort_steps(scope, host, signal, signal_id, reason) {
|
||||
None
|
||||
} else {
|
||||
let Some(error) = scope.exception() else {
|
||||
return;
|
||||
};
|
||||
scope.reset();
|
||||
Some(error)
|
||||
}
|
||||
};
|
||||
if first_error.is_none() {
|
||||
first_error = error;
|
||||
}
|
||||
}
|
||||
if let Some(error) = first_error {
|
||||
scope.throw_exception(error);
|
||||
}
|
||||
}
|
||||
|
||||
fn run_abort_steps<'s>(
|
||||
|
||||
@@ -15,6 +15,7 @@ mod inspect;
|
||||
mod observer;
|
||||
mod promise;
|
||||
mod state;
|
||||
mod switch_map;
|
||||
mod transform;
|
||||
mod until;
|
||||
|
||||
@@ -37,6 +38,8 @@ struct ObservablePrototype {
|
||||
map: (),
|
||||
#[webapi(method = "flatMap", length = 1, callback = flat_map::flat_map)]
|
||||
flat_map: (),
|
||||
#[webapi(method = "switchMap", length = 1, callback = switch_map::switch_map)]
|
||||
switch_map: (),
|
||||
#[webapi(method, length = 1, callback = transform::filter)]
|
||||
filter: (),
|
||||
#[webapi(method, length = 1, callback = transform::take)]
|
||||
@@ -322,6 +325,7 @@ fn subscribe_internal<'s>(
|
||||
&& !inspect::subscribe(scope, observable, subscriber)
|
||||
&& !finally::subscribe(scope, observable, subscriber)
|
||||
&& !flat_map::subscribe(scope, observable, subscriber)
|
||||
&& !switch_map::subscribe(scope, observable, subscriber)
|
||||
{
|
||||
event_target::subscribe(scope, observable, subscriber);
|
||||
}
|
||||
|
||||
@@ -3,7 +3,7 @@
|
||||
|
||||
use super::{
|
||||
callbacks, collect, consume, finally, first, flat_map, inspect, invoke_and_report, state::*,
|
||||
transform, until,
|
||||
switch_map, transform, until,
|
||||
};
|
||||
use crate::util::{get_private_value, set_private_value};
|
||||
|
||||
@@ -26,6 +26,8 @@ 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 SWITCH_MAP_SOURCE: i32 = 19;
|
||||
pub(super) const SWITCH_MAP_INNER: i32 = 20;
|
||||
pub(super) const SUBSCRIBER: &str = "__moliObservableNativeSubscriber";
|
||||
const INDEX: &str = "__moliObservableCallbackIndex";
|
||||
|
||||
@@ -98,6 +100,9 @@ pub(super) fn notify<'s>(
|
||||
FLAT_MAP_SOURCE | FLAT_MAP_INNER => {
|
||||
flat_map::notify(scope, observer, notification, kind)
|
||||
}
|
||||
SWITCH_MAP_SOURCE | SWITCH_MAP_INNER => {
|
||||
switch_map::notify(scope, observer, notification, kind)
|
||||
}
|
||||
_ => unreachable!("unknown native Observable observer"),
|
||||
}
|
||||
let exception = scope.exception();
|
||||
|
||||
@@ -0,0 +1,268 @@
|
||||
//! Switch inner subscriptions through the existing AbortSignal dependency graph.
|
||||
//! Reentrant mapper/teardown calls can leave multiple live inner observers, so
|
||||
//! trace each until its own subscription completes or is cancelled.
|
||||
|
||||
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 = "__moliObservableSwitchMapSource";
|
||||
const MAPPER: &str = "__moliObservableSwitchMapMapper";
|
||||
const DOWNSTREAM: &str = "__moliSwitchMapSubscriber";
|
||||
const OWNER: &str = "__moliSwitchMapSourceObserver";
|
||||
const INNERS: &str = "__moliSwitchMapInnerObservers";
|
||||
const CURRENT_SIGNAL: &str = "__moliSwitchMapCurrentSignal";
|
||||
const OUTER_COMPLETE: &str = "__moliSwitchMapSourceCompleted";
|
||||
const INNER_SIGNAL: &str = "__moliSwitchMapInnerSignal";
|
||||
const CLEANUP: &str = "__moliSwitchMapInnerCleanup";
|
||||
|
||||
#[derive(webidl::WebIdlArgs)]
|
||||
#[webidl(prefix = "Observable.switchMap")]
|
||||
struct SwitchMapArgs {
|
||||
#[webidl(required, converter = "callback_function")]
|
||||
mapper: webidl::WebIdlCallbackFunction,
|
||||
}
|
||||
|
||||
pub(super) fn switch_map<'s>(
|
||||
scope: &mut v8::PinScope<'s, '_>,
|
||||
args: v8::FunctionCallbackArguments<'s>,
|
||||
mut rv: v8::ReturnValue<'_, v8::Value>,
|
||||
) {
|
||||
let Some(parsed) = webidl::parse_args::<SwitchMapArgs>(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>,
|
||||
#[webapi(slot = INNER_SIGNAL)]
|
||||
signal: v8::Local<'scope, v8::Object>,
|
||||
}
|
||||
|
||||
fn signal_slot<'s>(
|
||||
scope: &mut v8::PinScope<'s, '_>,
|
||||
object: v8::Local<'s, v8::Object>,
|
||||
slot: &str,
|
||||
) -> Option<ResolvedAbortSignal<'s>> {
|
||||
object_slot(scope, object, slot).and_then(|signal| ResolvedAbortSignal::resolve(scope, signal))
|
||||
}
|
||||
|
||||
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("switchMap mapper");
|
||||
let Some(observer) = SourceObserver::new(observer::SWITCH_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_slot(scope, subscriber, SIGNAL).expect("switchMap Subscriber signal");
|
||||
subscribe_internal(scope, source, observer, Some(signal));
|
||||
true
|
||||
}
|
||||
|
||||
fn release_inner<'s>(scope: &mut v8::PinScope<'s, '_>, inner: v8::Local<'s, v8::Object>) {
|
||||
let owner = object_slot(scope, inner, OWNER).expect("switchMap owner");
|
||||
let mut inners = list(scope, owner, INNERS);
|
||||
inners.retain(|entry| *entry != inner);
|
||||
set_list(scope, owner, INNERS, &inners);
|
||||
if let Some(signal) = signal_slot(scope, inner, INNER_SIGNAL)
|
||||
&& let Some(callback) = object_slot(scope, inner, CLEANUP)
|
||||
.and_then(|value| v8::Local::<v8::Function>::try_from(value).ok())
|
||||
{
|
||||
signal.unregister_algorithm(scope, callback);
|
||||
}
|
||||
// Root the producer through the cancellation snapshot before dropping its
|
||||
// persistent edge. Its native cancellation callback is deliberately weak.
|
||||
let _producer = object_slot(scope, inner, observer::SUBSCRIBER);
|
||||
for slot in [INNER_SIGNAL, CLEANUP, observer::SUBSCRIBER] {
|
||||
set_private_value(scope, inner, slot, v8::undefined(scope).into());
|
||||
}
|
||||
}
|
||||
|
||||
fn cancelled<'s>(
|
||||
scope: &mut v8::PinScope<'s, '_>,
|
||||
args: v8::FunctionCallbackArguments<'s>,
|
||||
_rv: v8::ReturnValue<'_, v8::Value>,
|
||||
) {
|
||||
let inner = v8::Local::<v8::Object>::try_from(args.data()).expect("switchMap cleanup data");
|
||||
release_inner(scope, inner);
|
||||
}
|
||||
|
||||
fn process_next<'s>(
|
||||
scope: &mut v8::PinScope<'s, '_>,
|
||||
observer: v8::Local<'s, v8::Object>,
|
||||
value: v8::Local<'s, v8::Value>,
|
||||
) {
|
||||
if let Some(signal) = signal_slot(scope, observer, CURRENT_SIGNAL) {
|
||||
let error = {
|
||||
v8::tc_scope!(let scope, scope);
|
||||
let reason = crate::native_bridge::abort::abort_error_value(scope);
|
||||
signal.abort(scope, reason);
|
||||
let error = scope.exception();
|
||||
scope.reset();
|
||||
error
|
||||
};
|
||||
// Internal switching must finish even if IteratorClose fails. Explicit
|
||||
// author cancellation still uses the AbortController rethrow boundary.
|
||||
if let Some(error) = error {
|
||||
callbacks::report(scope, observer, error);
|
||||
}
|
||||
}
|
||||
let Some(initial_signal) = ResolvedAbortSignal::new(scope) else {
|
||||
return;
|
||||
};
|
||||
set_private_value(
|
||||
scope,
|
||||
observer,
|
||||
CURRENT_SIGNAL,
|
||||
initial_signal.value().into(),
|
||||
);
|
||||
let downstream = object_slot(scope, observer, DOWNSTREAM).expect("switchMap downstream");
|
||||
let mapper = object_slot(scope, observer, MAPPER).expect("switchMap 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;
|
||||
};
|
||||
// The draft passes the current controller by reference: mapper/conversion
|
||||
// reentrancy can replace it before subscription. If a reentrant inner also
|
||||
// completed, it cleared that reference; the original (now aborted) signal
|
||||
// safely initializes this superseded inner as inactive instead.
|
||||
let current = signal_slot(scope, observer, CURRENT_SIGNAL).unwrap_or(initial_signal);
|
||||
let downstream_signal = signal_slot(scope, downstream, SIGNAL).expect("switchMap signal");
|
||||
let Some(signal) = ResolvedAbortSignal::dependent(scope, &[current, downstream_signal]) else {
|
||||
return;
|
||||
};
|
||||
let Some(inner_observer) =
|
||||
InnerObserver::new(observer::SWITCH_MAP_INNER, observer, signal.value())
|
||||
.bind(scope)
|
||||
.ok()
|
||||
else {
|
||||
return;
|
||||
};
|
||||
if !signal.is_aborted(scope) {
|
||||
let mut inners = list(scope, observer, INNERS);
|
||||
inners.push(inner_observer);
|
||||
set_list(scope, observer, INNERS, &inners);
|
||||
let cleanup = v8::Function::builder(cancelled)
|
||||
.data(inner_observer.into())
|
||||
.build(scope)
|
||||
.expect("switchMap inner cleanup");
|
||||
set_private_value(scope, inner_observer, CLEANUP, cleanup.into());
|
||||
signal.register_weak_algorithm(scope, cleanup);
|
||||
}
|
||||
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 = if kind == observer::SWITCH_MAP_INNER {
|
||||
object_slot(scope, observer, OWNER).expect("switchMap source observer")
|
||||
} else {
|
||||
observer
|
||||
};
|
||||
let downstream = object_slot(scope, source, DOWNSTREAM).expect("switchMap downstream");
|
||||
match notification {
|
||||
Notification::Next(value) if kind == observer::SWITCH_MAP_SOURCE => {
|
||||
process_next(scope, source, value);
|
||||
}
|
||||
Notification::Next(value) => subscriber_next(scope, downstream, value),
|
||||
Notification::Error(error) => {
|
||||
if kind == observer::SWITCH_MAP_INNER {
|
||||
release_inner(scope, observer);
|
||||
}
|
||||
subscriber_error(scope, downstream, error);
|
||||
}
|
||||
Notification::Complete => {
|
||||
if kind == observer::SWITCH_MAP_SOURCE {
|
||||
set_private_value(
|
||||
scope,
|
||||
observer,
|
||||
observer::SUBSCRIBER,
|
||||
v8::undefined(scope).into(),
|
||||
);
|
||||
set_private_value(
|
||||
scope,
|
||||
source,
|
||||
OUTER_COMPLETE,
|
||||
v8::Boolean::new(scope, true).into(),
|
||||
);
|
||||
if signal_slot(scope, source, CURRENT_SIGNAL).is_none() {
|
||||
subscriber_complete(scope, downstream);
|
||||
}
|
||||
} else {
|
||||
release_inner(scope, observer);
|
||||
if get_private_value(scope, source, OUTER_COMPLETE).is_some_and(|v| v.is_true()) {
|
||||
subscriber_complete(scope, downstream);
|
||||
} else {
|
||||
set_private_value(scope, source, CURRENT_SIGNAL, v8::undefined(scope).into());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,152 @@
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn observable_switch_map_preserves_switch_order_conversion_reentrancy_and_cancellation() {
|
||||
let mut vm = new_storage_test_vm("https://observable-switch-map.test/");
|
||||
vm.eval(&format!(
|
||||
"({}).then(value => {{ globalThis.switchMapResult = JSON.stringify(value); }});",
|
||||
include_str!("../../../tests/fixtures/observable-switch-map.js")
|
||||
))
|
||||
.expect("Observable.switchMap fixture should evaluate");
|
||||
let result: serde_json::Value =
|
||||
serde_json::from_str(&vm.eval("switchMapResult").unwrap()).unwrap();
|
||||
assert_eq!(result["failures"], serde_json::json!([]), "{result}");
|
||||
assert!(result["checks"].as_u64().unwrap() >= 170, "{result}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn observable_switch_map_preserves_result_conversion_callback_and_cancellation_realms() {
|
||||
let mut vm = new_storage_test_vm("https://observable-switch-map-realms.test/");
|
||||
vm.eval("document.appendChild(document.createElement('iframe'))")
|
||||
.unwrap();
|
||||
materialize_single_child_default_realm_for_test(&mut vm, "Observable.switchMap realm");
|
||||
vm.eval(&format!(
|
||||
"({}).then(value => {{ globalThis.switchMapRealms = JSON.stringify(value); }});",
|
||||
include_str!("../../../tests/fixtures/observable-switch-map-realms.js")
|
||||
))
|
||||
.expect("Observable.switchMap realms fixture should evaluate");
|
||||
let result: serde_json::Value =
|
||||
serde_json::from_str(&vm.eval("switchMapRealms").unwrap()).unwrap();
|
||||
assert_eq!(result["failures"], serde_json::json!([]), "{result}");
|
||||
assert_eq!(result["checks"], 31, "{result}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn observable_switch_map_traces_reentrant_inners_and_releases_replaced_and_closed_graphs() {
|
||||
let mut vm = new_storage_test_vm("https://observable-switch-map-gc.test/");
|
||||
vm.eval(r#"
|
||||
function makeSwitchSource() {
|
||||
let subscriber;
|
||||
return {source: new Observable(s => { subscriber = s; }), get subscriber() { return subscriber; }};
|
||||
}
|
||||
function makeSwitchMapper() {
|
||||
const token = {calls: 0};
|
||||
return {token, callback: value => { token.calls++; return value; }};
|
||||
}
|
||||
globalThis.switchMapChains = [];
|
||||
for (const mode of ['abandoned', 'complete', 'outer-error', 'inner-error', 'abort']) (() => {
|
||||
const outer = makeSwitchSource(), old = makeSwitchSource(), inner = makeSwitchSource(), mapper = makeSwitchMapper();
|
||||
const result = outer.source.switchMap(mapper.callback), ac = new AbortController();
|
||||
const promise = result.toArray(mode === 'abort' ? {signal: ac.signal} : undefined);
|
||||
promise.catch(() => {});
|
||||
outer.subscriber.next(old.source); old.subscriber.next(1);
|
||||
outer.subscriber.next(inner.source); inner.subscriber.next(2);
|
||||
const entry = {mode, templates: [outer.source, old.source, inner.source, result].map(v => new WeakRef(v)),
|
||||
old: new WeakRef(old.subscriber), outer: new WeakRef(outer.subscriber), inner: new WeakRef(inner.subscriber),
|
||||
callbacks: [mapper.callback, mapper.token].map(v => new WeakRef(v))};
|
||||
if (mode !== 'abandoned') entry.promise = promise;
|
||||
if (mode === 'abort') entry.controller = ac;
|
||||
switchMapChains.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([
|
||||
switchMapChains.every(c => c.templates.every(ref => ref.deref() === undefined)),
|
||||
switchMapChains.every(c => c.old.deref() === undefined),
|
||||
switchMapChains.every(c => [c.outer, c.inner, ...c.callbacks].every(ref => (ref.deref() !== undefined) === (c.mode !== 'abandoned')))
|
||||
])"#).unwrap(), "[true,true,true]");
|
||||
vm.eval(r#"
|
||||
for (const c of switchMapChains.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();
|
||||
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("switchMapChains.filter(c => c.promise).every(c => c.correct && c.closed.every(s => !s.active))").unwrap(), "true");
|
||||
collect(&mut vm);
|
||||
assert_eq!(
|
||||
vm.eval("switchMapChains.every(c => c.callbacks.every(ref => ref.deref() === undefined))")
|
||||
.unwrap(),
|
||||
"true"
|
||||
);
|
||||
vm.eval(r#"
|
||||
globalThis.switchMapDrain = (() => {
|
||||
const outer = makeSwitchSource(), inner = makeSwitchSource(), mapper = makeSwitchMapper();
|
||||
const promise = outer.source.switchMap(mapper.callback).toArray();
|
||||
outer.subscriber.next(inner.source); outer.subscriber.complete();
|
||||
return {promise, outer: new WeakRef(outer.subscriber), inner: new WeakRef(inner.subscriber), callback: new WeakRef(mapper.callback)};
|
||||
})();
|
||||
"#).unwrap();
|
||||
collect(&mut vm);
|
||||
assert_eq!(vm.eval("switchMapDrain.outer.deref() === undefined && switchMapDrain.inner.deref() !== undefined && switchMapDrain.callback.deref() !== undefined").unwrap(), "true");
|
||||
vm.eval("switchMapDrain.promise.then(v => { switchMapDrain.correct = v.length === 0; }); switchMapDrain.inner.deref().complete();").unwrap();
|
||||
collect(&mut vm);
|
||||
assert_eq!(vm.eval("switchMapDrain.correct && switchMapDrain.inner.deref() === undefined && switchMapDrain.callback.deref() === undefined").unwrap(), "true");
|
||||
vm.eval(r#"
|
||||
function makeReentrantSwitchMapper(outer) {
|
||||
return value => {
|
||||
if (value.next) outer.subscriber.next({current: value.next});
|
||||
return value.current;
|
||||
};
|
||||
}
|
||||
globalThis.reentrantSwitch = (() => {
|
||||
const outer = makeSwitchSource(), first = makeSwitchSource(), second = makeSwitchSource();
|
||||
const promise = outer.source.switchMap(makeReentrantSwitchMapper(outer)).toArray();
|
||||
outer.subscriber.next({current: first.source, next: second.source});
|
||||
return {promise, source: new WeakRef(outer.subscriber), first: new WeakRef(first.subscriber), second: new WeakRef(second.subscriber)};
|
||||
})();
|
||||
"#).unwrap();
|
||||
collect(&mut vm);
|
||||
assert_eq!(vm.eval("[reentrantSwitch.source, reentrantSwitch.first, reentrantSwitch.second].every(ref => ref.deref() !== undefined)").unwrap(), "true");
|
||||
vm.eval("reentrantSwitch.promise.then(v => { reentrantSwitch.correct = JSON.stringify(v) === '[1,2]'; }); reentrantSwitch.first.deref().next(1); reentrantSwitch.second.deref().next(2); reentrantSwitch.source.deref().complete(); reentrantSwitch.first.deref().complete();").unwrap();
|
||||
collect(&mut vm);
|
||||
assert_eq!(vm.eval("reentrantSwitch.correct && [reentrantSwitch.source, reentrantSwitch.first, reentrantSwitch.second].every(ref => ref.deref() === undefined)").unwrap(), "true");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn observable_switch_map_handles_completion_during_reentrant_mapping() {
|
||||
let mut vm = new_storage_test_vm("https://observable-switch-map-reentrant-complete.test/");
|
||||
assert_eq!(vm.eval(r#"
|
||||
JSON.stringify(['mapper', 'conversion'].map(mode => {
|
||||
let source, superseded, completions = 0;
|
||||
new Observable(s => { source = s; }).switchMap(value => {
|
||||
if (value === 2) return [];
|
||||
if (mode === 'mapper') source.next(2);
|
||||
return mode === 'mapper' ? new Observable(s => { superseded = s; }) : {
|
||||
get [Symbol.asyncIterator]() { source.next(2); return undefined; },
|
||||
[Symbol.iterator]() { throw 'inactive conversion must not obtain an iterator'; }
|
||||
};
|
||||
}).subscribe({complete: () => completions++});
|
||||
source.next(1); source.complete();
|
||||
return completions === 1 && (mode === 'conversion' || (!superseded.active && superseded.signal.aborted));
|
||||
}))
|
||||
"#).unwrap(), "[true,true]");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn observable_flat_map_preserves_serial_order_conversion_reentrancy_and_cancellation() {
|
||||
let mut vm = new_storage_test_vm("https://observable-flat-map.test/");
|
||||
|
||||
@@ -275,11 +275,29 @@ fn abort_worker_signal<'s>(
|
||||
}
|
||||
signals_to_abort
|
||||
};
|
||||
// Complete the dependency snapshot even if a source's IteratorClose
|
||||
// throws. Dependents were already marked aborted and cannot be retried.
|
||||
let mut first_error = None;
|
||||
for (signal_id, signal) in signals_to_abort {
|
||||
if !run_worker_abort_steps(store, scope, signal, signal_id, reason) {
|
||||
return;
|
||||
let error = {
|
||||
v8::tc_scope!(let scope, scope);
|
||||
if run_worker_abort_steps(store, scope, signal, signal_id, reason) {
|
||||
None
|
||||
} else {
|
||||
let Some(error) = scope.exception() else {
|
||||
return;
|
||||
};
|
||||
scope.reset();
|
||||
Some(error)
|
||||
}
|
||||
};
|
||||
if first_error.is_none() {
|
||||
first_error = error;
|
||||
}
|
||||
}
|
||||
if let Some(error) = first_error {
|
||||
scope.throw_exception(error);
|
||||
}
|
||||
}
|
||||
|
||||
fn run_worker_abort_steps<'s>(
|
||||
|
||||
@@ -1,5 +1,22 @@
|
||||
use super::*;
|
||||
|
||||
#[tokio::test]
|
||||
async fn worker_observable_switch_map_preserves_switch_order_conversion_reentrancy_and_cancellation()
|
||||
{
|
||||
ensure_v8();
|
||||
let mut handle = spawn_worker(
|
||||
format!(
|
||||
"({}).then(value => {{ postMessage(value); close(); }});",
|
||||
include_str!("../../../../tests/fixtures/observable-switch-map.js")
|
||||
),
|
||||
"https://observable-switch-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() >= 170, "{result}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn worker_observable_flat_map_preserves_serial_order_conversion_reentrancy_and_cancellation()
|
||||
{
|
||||
|
||||
@@ -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.switchMap, 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.switchMap.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.switchMap.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};
|
||||
})()
|
||||
@@ -0,0 +1,339 @@
|
||||
(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.switchMap;
|
||||
check(typeof method === 'function', 'switchMap 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, 'switchMap');
|
||||
check(method.name === 'switchMap' && 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.switchMap()) instanceof TypeError, 'required mapper');
|
||||
for (const callback of [undefined, null, false, 1, 1n, '', Symbol(), {}, [], {handleEvent() {}}]) {
|
||||
check(thrown(() => source.switchMap(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()).switchMap(mapper) instanceof Derived), 'does not use source species');
|
||||
const saved = [];
|
||||
for (const [object, key] of [[Observable, 'from'], [Observable.prototype, 'subscribe'],
|
||||
[AbortSignal, 'any'], [AbortController.prototype, 'abort'], [globalThis, 'AbortController'],
|
||||
[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]).switchMap(() => 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]).switchMap(() => 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]).switchMap(() => 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]).switchMap(() => preferred).toArray(), [5], 'iterable result bypasses then property');
|
||||
check(thenReads === 0, 'does not assimilate arbitrary thenables');
|
||||
});
|
||||
|
||||
await test('switch cancellation and completion ordering', () => {
|
||||
const outer = subject(), inners = [], order = [], values = [], indices = [];
|
||||
const result = outer.source.switchMap((value, index) => {
|
||||
indices.push(index); order.push('map' + value);
|
||||
if (inners.length) check(!inners.at(-1).active, 'old inner closes before next mapper');
|
||||
return new Observable(s => {
|
||||
inners.push(s); order.push('start' + value);
|
||||
s.signal.addEventListener('abort', () => order.push('abort' + value));
|
||||
s.addTeardown(() => order.push('cleanup' + value)); s.next(value);
|
||||
});
|
||||
});
|
||||
result.subscribe({next: v => values.push(v), complete: () => order.push('complete')});
|
||||
outer.subscriber.next(1); outer.subscriber.next(2); inners[0].next('stale'); inners[0].complete();
|
||||
same(indices, [0, 1], 'indices begin at zero per subscription');
|
||||
same(values, [1, 2], 'cancelled inner no longer emits');
|
||||
check(inners[0].signal.reason instanceof DOMException && inners[0].signal.reason.name === 'AbortError', 'switch reason defaults to AbortError');
|
||||
same(order, ['map1', 'start1', 'abort1', 'cleanup1', 'map2', 'start2'], 'old abort and cleanup precede mapper');
|
||||
outer.subscriber.complete();
|
||||
check(inners[1].active && !order.includes('complete'), 'outer completion waits for current inner');
|
||||
inners[1].next(3); inners[1].complete();
|
||||
same(values, [1, 2, 3], 'last inner survives outer completion');
|
||||
same(order.slice(-3), ['abort2', 'cleanup2', 'complete'], 'final inner cleanup precedes downstream complete');
|
||||
const empty = subject(); let completes = 0;
|
||||
empty.source.switchMap(() => []).subscribe({complete: () => completes++});
|
||||
empty.subscriber.complete(); check(completes === 1, 'empty outer completes immediately');
|
||||
const sync = subject(); let syncCompletes = 0;
|
||||
sync.source.switchMap(v => [v]).subscribe({complete: () => syncCompletes++});
|
||||
sync.subscriber.next(1); check(syncCompletes === 0, 'completed inner waits for outer');
|
||||
sync.subscriber.complete(); check(syncCompletes === 1, 'completed inner releases completion gate');
|
||||
});
|
||||
|
||||
await test('reentrant mapper and conversion use shared controller reference', () => {
|
||||
for (const mode of ['mapper', 'conversion']) {
|
||||
const outer = subject(), inners = {}, calls = [], values = [];
|
||||
outer.source.switchMap((value, index) => {
|
||||
calls.push([value, index]);
|
||||
const inner = new Observable(s => { inners[value] = s; });
|
||||
if (value === 1 && mode === 'mapper') outer.subscriber.next(2);
|
||||
if (value === 1 && mode === 'conversion') return {
|
||||
get [Symbol.asyncIterator]() { outer.subscriber.next(2); return undefined; },
|
||||
[Symbol.iterator]() { return [value][Symbol.iterator](); }
|
||||
};
|
||||
return inner;
|
||||
}).subscribe({next: v => values.push(v), complete: () => values.push('complete')});
|
||||
outer.subscriber.next(1);
|
||||
same(calls, mode === 'mapper' ? [[1, 0], [2, 0]] : [[1, 0], [2, 1]], mode + ' invocation/index update order');
|
||||
if (mode === 'mapper') {
|
||||
check(inners[1].active && inners[2].active, 'mapper reentrancy retains both inner observers');
|
||||
inners[1].next('one'); inners[2].next('two');
|
||||
outer.subscriber.next(3);
|
||||
check(!inners[1].active && !inners[2].active && inners[3].active, 'next switch cancels both observers sharing current controller');
|
||||
inners[3].complete();
|
||||
} else {
|
||||
check(inners[2].active, 'conversion reentrancy keeps earlier inner active');
|
||||
inners[2].next('two');
|
||||
}
|
||||
outer.subscriber.complete();
|
||||
same(values, mode === 'mapper' ? ['one', 'two', 'complete'] : [1, 'two', 'complete'], mode + ' output and completion');
|
||||
check(Object.values(inners).every(s => !s.active), mode + ' leaves no live inner after completion');
|
||||
}
|
||||
});
|
||||
|
||||
await test('teardown and delivery reentrancy preserve all subscriptions', () => {
|
||||
const outer = subject(), inners = {}, log = [], values = [];
|
||||
outer.source.switchMap(v => new Observable(s => {
|
||||
inners[v] = s; log.push('start' + v);
|
||||
s.addTeardown(() => { log.push('close' + v); if (v === 1) outer.subscriber.next(3); });
|
||||
})).subscribe({next: v => values.push(v), complete: () => values.push('complete')});
|
||||
outer.subscriber.next(1); outer.subscriber.next(2);
|
||||
same(log, ['start1', 'close1', 'start3', 'start2'], 'reentrant teardown initializes before interrupted switch resumes');
|
||||
check(inners[2].active && inners[3].active, 'both reentrant inners stay live');
|
||||
inners[2].next(2); inners[3].next(3); outer.subscriber.complete();
|
||||
check(inners[2].active && inners[3].active, 'source completion still waits');
|
||||
inners[2].complete();
|
||||
same(values, [2, 3, 'complete'], 'inner completion closes remaining reentrant inner');
|
||||
check(!inners[3].active, 'dependent signal cancels reentrant inner');
|
||||
const delivery = subject(), seen = [], created = [];
|
||||
delivery.source.switchMap(v => new Observable(s => {
|
||||
created.push(s); s.next(v); if (v === 1) s.next('stale');
|
||||
})).subscribe(v => { seen.push(v); if (v === 1) delivery.subscriber.next(2); });
|
||||
delivery.subscriber.next(1);
|
||||
same(seen, [1, 2], 'delivery switches synchronously and suppresses later old notifications');
|
||||
check(!created[0].active && created[1].active, 'delivery cancels old inner before initializer resumes');
|
||||
delivery.subscriber.complete(); created[1].complete();
|
||||
});
|
||||
|
||||
await test('asynchronous switching ignores superseded resolutions', async () => {
|
||||
const outer = subject(), resolvers = [], values = [];
|
||||
const result = outer.source.switchMap(() => new Promise(resolve => resolvers.push(resolve)));
|
||||
const promise = result.toArray();
|
||||
outer.subscriber.next(1); outer.subscriber.next(2); outer.subscriber.complete();
|
||||
resolvers[0]('old'); await Promise.resolve();
|
||||
promise.then(v => values.push(...v));
|
||||
check(values.length === 0, 'old promise does not complete result');
|
||||
resolvers[1]('new'); same(await promise, ['new'], 'only active Promise delivers');
|
||||
const asyncOuter = subject(), asyncLog = [], gates = [];
|
||||
const asyncResult = asyncOuter.source.switchMap(v => ({[Symbol.asyncIterator]() { return {
|
||||
next() { return new Promise(resolve => gates.push(resolve)); },
|
||||
return(reason) { asyncLog.push([v, reason.name]); return Promise.resolve({done: true}); }
|
||||
}; }}));
|
||||
const ac = new AbortController(); asyncResult.subscribe({}, {signal: ac.signal});
|
||||
asyncOuter.subscriber.next(1); asyncOuter.subscriber.next(2);
|
||||
check(asyncLog.length === 1 && asyncLog[0][0] === 1, 'switch closes async iterator');
|
||||
gates[0]({value: 'old'}); await Promise.resolve(); await Promise.resolve();
|
||||
check(gates.length === 2, 'cancelled async iterator is not pulled again');
|
||||
ac.abort(); check(asyncLog.length === 2 && asyncLog[1][0] === 2, 'cancel closes latest async iterator');
|
||||
});
|
||||
|
||||
await test('switch only removes its own observer from a shared inner', () => {
|
||||
const outer = subject(), shared = subject(), keep = new AbortController(), values = [], other = [];
|
||||
shared.source.subscribe(v => other.push(v), {signal: keep.signal});
|
||||
outer.source.switchMap(v => v === 1 ? shared.source : [v]).subscribe(v => values.push(v));
|
||||
outer.subscriber.next(1); shared.subscriber.next('a'); outer.subscriber.next(2);
|
||||
check(shared.subscriber.active, 'other consumer keeps replaced producer alive');
|
||||
shared.subscriber.next('b'); same(values, ['a', 2], 'switched observer detached'); same(other, ['a', 'b'], 'other observer unaffected');
|
||||
keep.abort(); check(!shared.subscriber.active, 'last independent consumer closes producer');
|
||||
outer.subscriber.complete();
|
||||
});
|
||||
|
||||
|
||||
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.switchMap(() => { 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.switchMap(() => { 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 inner');
|
||||
inner.subscriber.complete(); outer.subscriber.complete();
|
||||
});
|
||||
|
||||
await test('original errors, inner 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.switchMap(() => {
|
||||
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'));
|
||||
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 + ' maps once');
|
||||
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.switchMap(() => { maps++; return inner.source; }).subscribe({}, {signal: ac.signal});
|
||||
outer.subscriber.next(1);
|
||||
outer.subscriber.addTeardown(() => { throw cleanupError; });
|
||||
ac.abort(reason);
|
||||
check(!outer.subscriber.active && !inner.subscriber.active && maps === 1, 'explicit cancellation closes both subscriptions');
|
||||
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).switchMap(() => 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 switchMap');
|
||||
});
|
||||
|
||||
await test('dependent cancellation finishes after source IteratorClose failure', () => {
|
||||
for (const marker of [{}, undefined]) {
|
||||
const ac = new AbortController(), log = [];
|
||||
let inner, caught, didThrow = false, pulls = 0;
|
||||
Observable.from({[Symbol.iterator]() { return {
|
||||
next() { return ++pulls === 1 ? {value: 1} : {done: true}; },
|
||||
return() { log.push('outer'); throw marker; }
|
||||
}; }}).subscribe(() => {
|
||||
new Observable(s => { inner = s; s.addTeardown(() => log.push('inner')); })
|
||||
.subscribe({}, {signal: AbortSignal.any([ac.signal])});
|
||||
try { ac.abort('stop'); } catch (e) { didThrow = true; caught = e; }
|
||||
}, {signal: ac.signal});
|
||||
check(didThrow && caught === marker, 'dependent dispatch preserves first exception including undefined');
|
||||
check(!inner.active && inner.signal.reason === 'stop', 'dependent observer is cancelled despite earlier failure');
|
||||
same(log, ['outer', 'inner'], 'dependent cleanup finishes before rethrow');
|
||||
}
|
||||
});
|
||||
|
||||
await test('captured inner notification survives observer removal', () => {
|
||||
const outer = subject(), inner = subject(), keep = new AbortController(), seen = [];
|
||||
inner.source.subscribe(() => outer.subscriber.next(2), {signal: keep.signal});
|
||||
outer.source.switchMap(value => value === 1 ? inner.source : [value]).subscribe(v => seen.push(v));
|
||||
outer.subscriber.next(1); inner.subscriber.next('captured');
|
||||
same(seen, [2, 'captured'], 'already captured next steps still run after reentrant switch');
|
||||
keep.abort(); outer.subscriber.complete();
|
||||
});
|
||||
|
||||
await test('pre-abort and cancellation inside mapper', () => {
|
||||
const ac = new AbortController(), reason = {}, pre = subject(); let maps = 0, inactive;
|
||||
pre.source.switchMap(() => { 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.switchMap(() => { 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.switchMap(() => { 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};
|
||||
})()
|
||||
Reference in New Issue
Block a user