feat(dom): add Observable take and drop operators

Reuse lazy native transform subscriptions for take/drop, preserving shared
counts, reentrant notification ordering, cancellation and early completion.
Keep each subscription's unsigned 64-bit count in a traced BigInt slot.

Fix unsigned-long-long conversion losing low bits for negative values by
applying the sign wrap in u64 after converting the magnitude. Add count,
conversion, Window/worker, realm, iterator-close and GC regression coverage.

Validation: fmt, all-target/all-feature workspace clippy, 172 targeted tests,
full nextest (19,297 passed; 13 skipped), 194 runtime checks, 40 realm checks,
and 43 WPT cases improving from 441/469 to 467/469 without regressions.
This commit is contained in:
ldm0
2026-09-22 21:47:16 +08:00
parent d8c7af63fb
commit c5970767ef
8 changed files with 516 additions and 19 deletions
@@ -4538,6 +4538,8 @@ dom/observable/tentative/idlharness.html
dom/observable/tentative/observable-constructor.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-constructor.any.js?moli-wpt-any=window
dom/observable/tentative/observable-constructor.window.js?moli-wpt-script=window
dom/observable/tentative/observable-drop.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-drop.any.js?moli-wpt-any=window
dom/observable/tentative/observable-event-target.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-event-target.any.js?moli-wpt-any=window
dom/observable/tentative/observable-event-target.window.js?moli-wpt-script=window
@@ -4561,6 +4563,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-take.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-take.any.js?moli-wpt-any=window
dom/observable/tentative/observable-toArray.any.js?moli-wpt-any=dedicatedworker
dom/observable/tentative/observable-toArray.any.js?moli-wpt-any=window
dom/ranges/Range-adopt-test.html
+4
View File
@@ -33,6 +33,10 @@ struct ObservablePrototype {
map: (),
#[webapi(method, length = 1, callback = transform::filter)]
filter: (),
#[webapi(method, length = 1, callback = transform::take)]
take: (),
#[webapi(method, length = 1, callback = transform::drop)]
drop: (),
#[webapi(method, length = 0, returns_promise, callback = first::first)]
first: (),
#[webapi(method, length = 0, returns_promise, callback = collect::last)]
+5 -1
View File
@@ -15,6 +15,8 @@ pub(super) const EVERY: i32 = 7;
pub(super) const FIND: i32 = 8;
pub(super) const MAP: i32 = 9;
pub(super) const FILTER: i32 = 10;
pub(super) const TAKE: i32 = 11;
pub(super) const DROP: i32 = 12;
pub(super) const SUBSCRIBER: &str = "__moliObservableNativeSubscriber";
const INDEX: &str = "__moliObservableCallbackIndex";
@@ -78,7 +80,9 @@ pub(super) fn notify<'s>(
FOR_EACH | REDUCE | SOME | EVERY | FIND => {
consume::notify(scope, observer, notification, kind);
}
MAP | FILTER => transform::notify(scope, observer, notification, kind),
MAP | FILTER | TAKE | DROP => {
transform::notify(scope, observer, notification, kind)
}
_ => unreachable!("unknown native Observable observer"),
}
let exception = scope.exception();
+106 -15
View File
@@ -1,5 +1,5 @@
//! Lazy map/filter producers. Each shared downstream subscription owns one
//! upstream observer, whose callback and index are traced by V8.
//! Lazy map/filter/take/drop producers. Each shared downstream subscription owns
//! one upstream observer, whose callback or remaining count is traced by V8.
use moli_webapi_declare::WebApiObject;
@@ -16,7 +16,7 @@ use crate::{
};
const SOURCE: &str = "__moliObservableTransformSource";
const CALLBACK: &str = "__moliObservableTransformCallback";
const PARAMETER: &str = "__moliObservableTransformParameter";
const KIND: &str = "__moliObservableTransformKind";
const DOWNSTREAM: &str = "__moliObservableTransformSubscriber";
@@ -27,12 +27,19 @@ struct TransformArgs {
callback: webidl::WebIdlCallbackFunction,
}
#[derive(webidl::WebIdlArgs)]
#[webidl(prefix = "Observable")]
struct CountArgs {
#[webidl(required, converter = "unsigned_long_long")]
amount: u64,
}
pub(super) fn map<'s>(
scope: &mut v8::PinScope<'s, '_>,
args: v8::FunctionCallbackArguments<'s>,
rv: v8::ReturnValue<'_, v8::Value>,
) {
create(scope, args, rv, observer::MAP);
create_callback(scope, args, rv, observer::MAP);
}
pub(super) fn filter<'s>(
@@ -40,10 +47,26 @@ pub(super) fn filter<'s>(
args: v8::FunctionCallbackArguments<'s>,
rv: v8::ReturnValue<'_, v8::Value>,
) {
create(scope, args, rv, observer::FILTER);
create_callback(scope, args, rv, observer::FILTER);
}
fn create<'s>(
pub(super) fn take<'s>(
scope: &mut v8::PinScope<'s, '_>,
args: v8::FunctionCallbackArguments<'s>,
rv: v8::ReturnValue<'_, v8::Value>,
) {
create_count(scope, args, rv, observer::TAKE);
}
pub(super) fn drop<'s>(
scope: &mut v8::PinScope<'s, '_>,
args: v8::FunctionCallbackArguments<'s>,
rv: v8::ReturnValue<'_, v8::Value>,
) {
create_count(scope, args, rv, observer::DROP);
}
fn create_callback<'s>(
scope: &mut v8::PinScope<'s, '_>,
args: v8::FunctionCallbackArguments<'s>,
mut rv: v8::ReturnValue<'_, v8::Value>,
@@ -52,18 +75,62 @@ fn create<'s>(
let Some(parsed) = webidl::parse_args::<TransformArgs>(scope, &args) else {
return;
};
let Some(observable) = new_native_observable(scope, None) else {
let Some(observable) = new_transform(scope, args.this(), kind) else {
return;
};
set_private_value(scope, observable, SOURCE, args.this().into());
set_callback(scope, observable, PARAMETER, parsed.callback);
rv.set(observable.into());
}
fn create_count<'s>(
scope: &mut v8::PinScope<'s, '_>,
args: v8::FunctionCallbackArguments<'s>,
mut rv: v8::ReturnValue<'_, v8::Value>,
kind: i32,
) {
let Some(parsed) = webidl::parse_args::<CountArgs>(scope, &args) else {
return;
};
let Some(observable) = new_transform(scope, args.this(), kind) else {
return;
};
set_count(scope, observable, parsed.amount);
rv.set(observable.into());
}
fn new_transform<'s>(
scope: &mut v8::PinScope<'s, '_>,
source: v8::Local<'s, v8::Object>,
kind: i32,
) -> Option<v8::Local<'s, v8::Object>> {
let observable = new_native_observable(scope, None)?;
set_private_value(scope, observable, SOURCE, source.into());
set_private_value(
scope,
observable,
KIND,
v8::Integer::new(scope, kind).into(),
);
set_callback(scope, observable, CALLBACK, parsed.callback);
rv.set(observable.into());
Some(observable)
}
fn count<'s>(scope: &mut v8::PinScope<'s, '_>, object: v8::Local<'s, v8::Object>) -> u64 {
get_private_value(scope, object, PARAMETER)
.and_then(|value| v8::Local::<v8::BigInt>::try_from(value).ok())
.expect("Observable transform count")
.u64_value()
.0
}
fn set_count<'s>(scope: &mut v8::PinScope<'s, '_>, object: v8::Local<'s, v8::Object>, count: u64) {
// Keep every bit after WebIDL conversion, including negative inputs modulo
// 2^64; storing a Number would round u64::MAX back up to 2^64.
set_private_value(
scope,
object,
PARAMETER,
v8::BigInt::new_from_u64(scope, count).into(),
);
}
#[derive(WebApiObject)]
@@ -73,8 +140,8 @@ struct TransformObserver<'scope> {
kind: i32,
#[webapi(slot = DOWNSTREAM)]
downstream: v8::Local<'scope, v8::Object>,
#[webapi(slot = CALLBACK)]
callback: v8::Local<'scope, v8::Object>,
#[webapi(slot = PARAMETER)]
parameter: v8::Local<'scope, v8::Value>,
}
pub(super) fn subscribe<'s>(
@@ -88,8 +155,14 @@ pub(super) fn subscribe<'s>(
let kind = get_private_value(scope, observable, KIND)
.and_then(|value| value.int32_value(scope))
.expect("Observable transform kind");
let callback = object_slot(scope, observable, CALLBACK).expect("Observable transform callback");
let Some(observer) = TransformObserver::new(kind, subscriber, callback)
if kind == observer::TAKE && count(scope, observable) == 0 {
// take(0) completes synchronously without ever subscribing upstream.
subscriber_complete(scope, subscriber);
return true;
}
let parameter = get_private_value(scope, observable, PARAMETER).expect("Transform parameter");
// Count state belongs to the subscription, not the reusable Observable.
let Some(observer) = TransformObserver::new(kind, subscriber, parameter)
.bind(scope)
.ok()
else {
@@ -118,11 +191,29 @@ pub(super) fn notify<'s>(
) {
let downstream = object_slot(scope, observer, DOWNSTREAM).expect("Transform downstream");
match notification {
Notification::Next(value) if kind == observer::TAKE => {
subscriber_next(scope, downstream, value);
// Delivery may reenter next() and consume more values. Read the
// current remaining count after delivery, as in the draft algorithm.
let remaining = count(scope, observer).wrapping_sub(1);
set_count(scope, observer, remaining);
if remaining == 0 {
subscriber_complete(scope, downstream);
}
}
Notification::Next(value) if kind == observer::DROP => {
let remaining = count(scope, observer);
if remaining > 0 {
set_count(scope, observer, remaining - 1);
} else {
subscriber_next(scope, downstream, value);
}
}
Notification::Next(value) => {
// A captured notification still invokes the callback after another
// observer cancels this branch. Subscriber.next filters delivery;
// Subscriber.error reports a callback failure after closure.
let callback = object_slot(scope, observer, CALLBACK).expect("Transform callback");
let callback = object_slot(scope, observer, PARAMETER).expect("Transform callback");
let index = observer::index(scope, observer);
let index = v8::Number::new(scope, index as f64).into();
match callbacks::invoke_value(scope, callback, &[value, index]) {
@@ -1,5 +1,138 @@
use super::*;
#[test]
fn observable_count_operators_preserve_conversion_sharing_reentrancy_and_cancellation() {
let mut vm = new_storage_test_vm("https://observable-count-operators.test/");
vm.eval(&format!(
"({}).then(value => {{ globalThis.countOperatorResult = JSON.stringify(value); }});",
include_str!("../../../tests/fixtures/observable-count-operators.js")
))
.expect("Observable count operator fixture should evaluate");
let result: serde_json::Value =
serde_json::from_str(&vm.eval("countOperatorResult").unwrap()).unwrap();
assert_eq!(result["failures"], serde_json::json!([]), "{result}");
assert!(result["checks"].as_u64().unwrap() >= 194, "{result}");
}
#[test]
fn observable_count_operators_preserve_conversion_result_and_exception_realms() {
let mut vm = new_storage_test_vm("https://observable-count-realms.test/");
vm.eval("document.appendChild(document.createElement('iframe'))")
.unwrap();
materialize_single_child_default_realm_for_test(&mut vm, "Observable count operator realm");
vm.eval(r#"
(async () => {
const child = document.querySelector('iframe').contentWindow, checks = [];
for (const name of ['take', 'drop']) {
const method = child.Observable.prototype[name];
child.countReads = 0;
const amount = new child.Object();
amount[Symbol.toPrimitive] = child.Function('globalThis.countThis = this; globalThis.countReads++; return 1;');
let starts = 0;
const source = new Observable(s => { starts++; s.next(1); s.next(2); s.complete(); });
const result = method.call(source, amount);
checks.push(result instanceof child.Observable, !(result instanceof Observable), Object.getPrototypeOf(result) === child.Observable.prototype);
checks.push(starts === 0, child.countReads === 1, child.countThis === amount);
const values = await result.toArray();
checks.push(values instanceof child.Array, values.length === 1 && values[0] === (name === 'take' ? 1 : 2));
const local = Observable.prototype[name].call(child.Observable.from([1, 2]), 1);
checks.push(local instanceof Observable, !(local instanceof child.Observable));
const localValues = await local.toArray();
checks.push(localValues instanceof Array, localValues.length === 1 && localValues[0] === (name === 'take' ? 1 : 2));
let reads = 0;
try { method.call({}, {valueOf() { reads++; return 1; }}); } catch (e) { checks.push(e instanceof child.TypeError, !(e instanceof TypeError)); }
checks.push(reads === 0);
try { method.call(source, 1n); } catch (e) { checks.push(e instanceof child.TypeError, !(e instanceof TypeError)); }
const marker = new child.RangeError('amount');
try { method.call(source, {valueOf() { throw marker; }}); } catch (e) { checks.push(e === marker, e instanceof child.RangeError); }
method.call(new Observable(s => s.error(marker)), 2).subscribe({error: e => checks.push(e === marker)});
}
globalThis.countOperatorRealms = JSON.stringify(checks);
})();
"#).unwrap();
let checks: Vec<bool> = serde_json::from_str(&vm.eval("countOperatorRealms").unwrap()).unwrap();
assert_eq!(checks.len(), 40);
assert!(checks.iter().all(|value| *value), "{checks:?}");
}
#[test]
fn observable_count_operator_chains_trace_pending_state_and_release_after_early_completion() {
let mut vm = new_storage_test_vm("https://observable-count-gc.test/");
vm.eval(r#"
globalThis.countChains = [];
for (const kept of [false, true]) (() => {
let subscriber;
const token = {}, callback = value => token && value;
const source = new Observable(s => { subscriber = s; });
const mapped = source.map(callback), dropped = mapped.drop(1), taken = dropped.take(2);
const promise = taken.toArray();
subscriber.next(1);
const entry = {kept, templates: [source, mapped, dropped, taken].map(value => new WeakRef(value)),
refs: {subscriber: new WeakRef(subscriber), callback: new WeakRef(callback), token: new WeakRef(token)}};
if (kept) entry.promise = promise;
countChains.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([
countChains.every(c => c.templates.every(ref => ref.deref() === undefined)),
countChains.every(c => Object.values(c.refs).every(ref => (ref.deref() !== undefined) === c.kept))
])"#
)
.unwrap(),
"[true,true]"
);
vm.eval(
r#"
const keptCountChain = countChains.find(c => c.kept);
keptCountChain.promise.then(values => { keptCountChain.values = values; });
{
const subscriber = keptCountChain.refs.subscriber.deref();
subscriber.next(2); subscriber.next(3);
globalThis.countSourceClosed = !subscriber.active;
}
"#,
)
.unwrap();
assert_eq!(
vm.eval("countSourceClosed && JSON.stringify(keptCountChain.values) === '[2,3]'")
.unwrap(),
"true"
);
collect(&mut vm);
assert_eq!(
vm.eval(
"countChains.every(c => Object.values(c.refs).every(ref => ref.deref() === undefined))"
)
.unwrap(),
"true"
);
vm.eval(r#"
globalThis.cancelledCountChain = (() => {
const ac = new AbortController(), token = {}, callback = value => token && value;
let subscriber;
const source = new Observable(s => { subscriber = s; });
const promise = source.map(callback).drop(1).take(2).toArray({signal: ac.signal});
promise.catch(() => {}); ac.abort();
return {promise, subscriber, signal: ac.signal, refs: [new WeakRef(callback), new WeakRef(token)]};
})();
"#).unwrap();
collect(&mut vm);
assert_eq!(vm.eval("!cancelledCountChain.subscriber.active && cancelledCountChain.refs.every(ref => ref.deref() === undefined)").unwrap(), "true");
}
#[test]
fn observable_transforms_preserve_lazy_sharing_cancellation_and_callback_semantics() {
let mut vm = new_storage_test_vm("https://observable-transforms.test/");
@@ -1,5 +1,22 @@
use super::*;
#[tokio::test]
async fn worker_observable_count_operators_preserve_conversion_sharing_reentrancy_and_cancellation()
{
ensure_v8();
let mut handle = spawn_worker(
format!(
"({}).then(value => {{ postMessage(value); close(); }});",
include_str!("../../../../tests/fixtures/observable-count-operators.js")
),
"https://observable-count-operators.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() >= 194, "{result}");
}
#[tokio::test]
async fn worker_observable_transforms_preserve_lazy_sharing_cancellation_and_callback_semantics() {
ensure_v8();
@@ -0,0 +1,219 @@
(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 names = ['take', 'drop'];
for (const name of names) check(typeof Observable.prototype[name] === 'function', name + ' exposed');
if (failures.length) return {checks, failures};
for (const name of names) {
const method = Observable.prototype[name];
await test(name + ' conversion and intrinsic creation', async () => {
const desc = Object.getOwnPropertyDescriptor(Observable.prototype, name);
check(method.name === name && method.length === 1, name + ' name and length');
check(desc.enumerable && desc.writable && desc.configurable, name + ' descriptor');
check(thrown(() => new method(1)) instanceof TypeError, name + ' non-constructible');
let starts = 0, conversions = 0, traps = 0;
const source = new Observable(s => { starts++; s.next(1); s.next(2); s.complete(); });
const amount = {[Symbol.toPrimitive](hint) { conversions++; check(hint === 'number', name + ' numeric conversion hint'); return 1; }};
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, amount)) instanceof TypeError, name + ' invalid receiver');
}
check(conversions === 0 && traps === 0, name + ' receiver validation precedes conversion and traps');
check(thrown(() => method.call(source)) instanceof TypeError, name + ' required argument');
for (const value of [1n, Symbol(), {[Symbol.toPrimitive]: () => 1n}, Object.create(null)]) {
check(thrown(() => method.call(source, value)) instanceof TypeError, name + ' invalid numeric conversion');
}
const marker = {};
check(thrown(() => method.call(source, {valueOf() { throw marker; }})) === marker, name + ' conversion exception identity');
check(starts === 0, name + ' invalid conversion does not subscribe');
const result = method.call(source, amount);
check(starts === 0 && conversions === 1, name + ' converts once at creation without subscribing');
check(result !== source && Object.getPrototypeOf(result) === Observable.prototype && Observable.from(result) === result, name + ' new branded intrinsic result');
same(await result.toArray(), name === 'take' ? [1] : [2], name + ' first result');
same(await result.toArray(), name === 'take' ? [1] : [2], name + ' count resets for next subscription');
check(starts === 2 && conversions === 1, name + ' subscriptions reuse converted amount');
const order = [];
const numeric = {valueOf() { order.push('valueOf'); return {}; }, toString() { order.push('toString'); return '1.9'; }};
same(await method.call(source, numeric).toArray(), name === 'take' ? [1] : [2], name + ' ordinary number conversion');
same(order, ['valueOf', 'toString'], name + ' conversion order');
let reads = 0;
Object.defineProperty(source, 'constructor', {get() { reads++; throw marker; }});
Object.setPrototypeOf(source, null);
same(await method.call(source, 1, {get signal() { reads++; throw marker; }}).toArray(), name === 'take' ? [1] : [2], name + ' genuine receiver ignores public prototype and extra options');
check(reads === 0, name + ' no species or options lookup');
class Subclass extends Observable {}
const derived = method.call(new Subclass(s => s.complete()), 1);
check(Object.getPrototypeOf(derived) === Observable.prototype && !(derived instanceof Subclass), name + ' result does not inherit source subclass');
});
await test(name + ' unsigned 64-bit amount', async () => {
const cases = [
[undefined, 0], [null, 0], [false, 0], [-0, 0], [-0.9, 0], [NaN, 0], [Infinity, 0], [-Infinity, 0],
['', 0], ['not numeric', 0], [2 ** 64, 0], [-(2 ** 64), 0],
[true, 1], [1.9, 1], ['1.9', 1], [new Number(1), 1], [2.9, 2], ['2', 2],
[-1, 3], [-2, 3], [-3.9, 3], [2 ** 32, 3], [2 ** 53, 3], [2 ** 64 - 2048, 3],
[2 ** 64 + 4096, 3], [-(2 ** 64) + 2048, 3],
];
for (const [amount, count] of cases) {
const values = [1, 2, 3];
same(await method.call(Observable.from(values), amount).toArray(),
name === 'take' ? values.slice(0, count) : values.slice(count), name + ' converts ' + String(amount));
}
});
await test(name + ' sharing, cancellation and fresh counts', () => {
let subscriber, starts = 0, teardowns = 0;
const source = new Observable(s => { starts++; subscriber = s; s.addTeardown(() => teardowns++); });
const keeper = new AbortController(); source.subscribe({}, {signal: keeper.signal});
const result = method.call(source, 2);
for (let round = 0; round < 2; round++) {
const left = [], right = [], a = new AbortController(), b = new AbortController();
let completions = 0;
result.subscribe(value => { left.push(value); a.abort(); }, {signal: a.signal});
result.subscribe({next: value => right.push(value), complete: () => completions++}, {signal: b.signal});
for (const value of [1, 2, 3, 4]) subscriber.next(value);
same(left, name === 'take' ? [1] : [3], name + ' cancelled observer stops receiving');
same(right, name === 'take' ? [1, 2] : [3, 4], name + ' observers share one count');
check(completions === (name === 'take' ? 1 : 0), name + ' completion only at take limit');
check(subscriber.active && starts === 1 && teardowns === 0, name + ' other source subscription stays alive');
b.abort();
}
keeper.abort(); check(!subscriber.active && teardowns === 1, name + ' last observer closes source');
});
await test(name + ' errors, completion and pre-abort', async () => {
for (const error of [{}, null, undefined]) {
const received = [], errors = [];
const source = new Observable(s => { s.next(1); s.error(error); s.next(2); });
method.call(source, 2).subscribe({next: value => received.push(value), error: e => errors.push(e), complete: () => errors.push('complete')});
same(received, name === 'take' ? [1] : [], name + ' error before count reached');
check(errors.length === 1 && errors[0] === error, name + ' source error identity');
}
const marker = {}, errors = [];
method.call(new Observable(() => { throw marker; }), 2).subscribe({error: e => errors.push(e)});
check(errors.length === 1 && errors[0] === marker, name + ' initializer exception forwarded');
same(await method.call(Observable.from([]), 2).toArray(), [], name + ' empty source');
same(await method.call(Observable.from([1]), 2).toArray(), name === 'take' ? [1] : [], name + ' source ends before amount reached');
for (const amount of [0, 2]) {
let starts = 0, notifications = 0, subscriber;
const source = new Observable(s => { starts++; subscriber = s; });
method.call(source, amount).subscribe({next() { notifications++; }, error() { notifications++; }, complete() { notifications++; }}, {signal: AbortSignal.abort(marker)});
check(starts === (name === 'take' && amount === 0 ? 0 : 1), name + ' pre-aborted source initialization');
check(!subscriber || !subscriber.active && subscriber.signal.reason === marker, name + ' pre-aborted Subscriber keeps reason');
check(notifications === 0, name + ' pre-aborted observer receives no notifications');
}
});
await test(name + ' iterator cancellation and raw values', async () => {
const ac = new AbortController(), marker = {}, values = [];
let pulls = 0, closes = 0;
const iterator = {next: () => ({value: ++pulls}), return() { closes++; return {}; }};
method.call(Observable.from({[Symbol.iterator]: () => iterator}), 2).subscribe({
next: value => { values.push(value); ac.abort(marker); }, error: () => values.push('error'), complete: () => values.push('complete'),
}, {signal: ac.signal});
same(values, name === 'take' ? [1] : [3], name + ' cancellation inside next suppresses terminal notifications');
check(pulls === (name === 'take' ? 1 : 3) && closes === 1, name + ' cancellation closes iterator exactly once');
let reads = 0;
const poison = {get then() { reads++; throw marker; }};
const rejected = Promise.reject(marker); rejected.catch(() => {});
const revoked = Proxy.revocable({}, {}); revoked.revoke();
const raw = [undefined, null, false, -0, NaN, 1n, Symbol(), poison, Promise.resolve(1), rejected, revoked.proxy];
const result = await method.call(Observable.from(raw), name === 'take' ? raw.length : 0).toArray();
check(result.length === raw.length && result.every((value, i) => Object.is(value, raw[i])), name + ' values are not converted or assimilated');
check(reads === 0, name + ' raw values do not read then');
});
}
await test('zero count and teardown ordering', () => {
const log = [];
const source = new Observable(s => {
log.push('start'); s.signal.addEventListener('abort', () => log.push('abort'));
s.addTeardown(() => log.push('teardown')); s.next(1); log.push('after 1'); s.next(2); log.push('after 2'); s.complete();
});
source.take(0).subscribe({complete: () => log.push('empty')});
same(log.splice(0), ['empty'], 'take(0) completes without starting source');
source.take(1).subscribe({next: value => log.push(value), complete: () => log.push('complete')});
same(log.splice(0), ['start', 1, 'abort', 'teardown', 'complete', 'after 1', 'after 2'], 'take closes upstream before downstream completion');
source.drop(0).subscribe({next: value => log.push(value), complete: () => log.push('complete')});
same(log.splice(0), ['start', 1, 'after 1', 2, 'after 2', 'abort', 'teardown', 'complete'], 'drop(0) mirrors full source');
let reads = 0;
const iterable = {get [Symbol.iterator]() { reads++; return () => { throw new Error('opened'); }; }};
Observable.from(iterable).take(0).subscribe();
check(reads === 1, 'take(0) does not open iterator after Observable.from conversion');
});
await test('reentrancy and notification snapshots', () => {
let subscriber;
const source = new Observable(s => { subscriber = s; });
const taken = source.take(1), log = [];
taken.subscribe({next: value => { log.push('a' + value); if (value === 1) subscriber.next(2); }, complete: () => log.push('a complete')});
taken.subscribe({next: value => log.push('b' + value), complete: () => log.push('b complete')});
subscriber.next(1); subscriber.next(3);
same(log, ['a1', 'a2', 'b2', 'a complete', 'b complete', 'b1'], 'take decrements after delivery and preserves captured downstream next');
const values = [];
source.take(2).subscribe({next: value => { values.push(value); if (value === 1) { subscriber.next(2); subscriber.next(3); } }, complete: () => values.push('complete')});
subscriber.next(1); subscriber.next(4);
same(values, [1, 2, 3, 'complete'], 'nested take decrements are read after reentrant delivery');
const dropped = [];
source.drop(1).subscribe(value => { dropped.push(value); if (value === 2) subscriber.next(3); });
subscriber.next(1); subscriber.next(2); subscriber.complete();
same(dropped, [2, 3], 'drop count stays zero during reentrant delivery');
for (const name of names) {
const ac = new AbortController(), received = [];
source.subscribe(() => ac.abort());
source[name](name === 'take' ? 1 : 0).subscribe({next: v => received.push(v), complete: () => received.push('complete')}, {signal: ac.signal});
subscriber.next(1); subscriber.complete();
same(received, [], name + ' cancelled branch rejects captured upstream notification');
}
});
await test('iterator close errors and async sources', async () => {
const closeError = {}, reports = [];
const onerror = e => { reports.push(e.error); e.preventDefault(); };
addEventListener('error', onerror);
try {
const iterator = {next: () => ({value: 7}), return() { throw closeError; }};
same(await Observable.from({[Symbol.iterator]: () => iterator}).take(1).toArray(), [7], 'take completion survives iterator close error');
check(reports.length === 1 && reports[0] === closeError, 'take reports close error once');
reports.length = 0;
for (const name of names) {
const ac = new AbortController(); let caught;
Observable.from({[Symbol.iterator]: () => iterator})[name](name === 'take' ? 2 : 0)
.subscribe(() => { caught = thrown(() => ac.abort()); }, {signal: ac.signal});
check(caught === closeError, name + ' explicit abort preserves close exception');
}
check(reports.length === 0, 'caught author abort errors are not reported');
} finally { removeEventListener('error', onerror); }
let pulls = 0, closes = 0;
const iterator = {next: () => Promise.resolve({value: ++pulls}), return() { closes++; return Promise.resolve({}); }};
same(await Observable.from({[Symbol.asyncIterator]: () => iterator}).drop(1).take(2).toArray(), [2, 3], 'async drop/take chain');
check(pulls === 3 && closes === 1, 'async chain cancels at last selected value');
async function* numbers() { yield 1; yield 2; yield 3; yield 4; }
const source = Observable.from(numbers());
const results = await Promise.all([source.take(1).toArray(), source.drop(1).take(2).toArray()]);
same(results, [[1], [2, 3]], 'shared async producer continues after shorter branch completes');
});
await test('intrinsic operations and mixed transforms', async () => {
same(await Observable.from([1, 2, 3, 4, 5]).map(v => v * 2).drop(1).filter(v => v % 4 === 0).take(2).toArray(), [4, 8], 'mixed callback and count operators');
const Constructor = Observable, take = Observable.prototype.take, drop = Observable.prototype.drop, subscribe = Observable.prototype.subscribe;
const originals = [globalThis.Observable, Observable.prototype.subscribe, Subscriber.prototype.next, Subscriber.prototype.error, Subscriber.prototype.complete];
const source = Observable.from([1, 2, 3, 4]), values = [];
const poison = () => { throw new Error('public implementation consulted'); };
let result;
try {
globalThis.Observable = Constructor.prototype.subscribe = Subscriber.prototype.next = Subscriber.prototype.error = Subscriber.prototype.complete = poison;
result = drop.call(take.call(source, 3), 1); subscribe.call(result, value => values.push(value));
} finally { [globalThis.Observable, Constructor.prototype.subscribe, Subscriber.prototype.next, Subscriber.prototype.error, Subscriber.prototype.complete] = originals; }
check(result instanceof Constructor, 'count operators use intrinsic Observable prototype');
same(values, [2, 3], 'count operators bypass public methods');
});
return {checks, failures};
})()
+28 -3
View File
@@ -1310,9 +1310,14 @@ fn unsigned_long_long(value: f64) -> u64 {
if !value.is_finite() || value == 0.0 {
return 0;
}
let integer = value.trunc();
let wrapped = integer.rem_euclid(2f64.powi(64));
wrapped as u64
// Adding 2^64 to a small negative remainder in f64 loses its low bits.
// Convert the magnitude first, then perform the sign wrap exactly in u64.
let magnitude = (value.trunc().abs() % 2f64.powi(64)) as u64;
if value.is_sign_negative() {
magnitude.wrapping_neg()
} else {
magnitude
}
}
fn enforce_range_unsigned_long_long(value: f64, context: Context) -> Result<u64, WebIdlError> {
@@ -1422,6 +1427,26 @@ mod tests {
assert_eq!(unsigned_long_long(1.9), 1);
}
#[test]
fn unsigned_long_long_preserves_low_bits_when_wrapping_negative_values() {
for (value, expected) in [
(-0.9, 0),
(-2.0, u64::MAX - 1),
(-3.9, u64::MAX - 2),
(-1025.0, u64::MAX - 1024),
(-2047.0, u64::MAX - 2046),
(-2f64.powi(63), 1 << 63),
(-2f64.powi(64) + 2048.0, 2048),
(-2f64.powi(64), 0),
(-2f64.powi(64) - 4096.0, u64::MAX - 4095),
(2f64.powi(64) - 2048.0, u64::MAX - 2047),
(2f64.powi(64), 0),
(2f64.powi(64) + 4096.0, 4096),
] {
assert_eq!(unsigned_long_long(value), expected, "{value}");
}
}
#[test]
fn enforce_range_unsigned_long_long_rejects_out_of_range_values() {
let context = Context::argument("IDBFactory.open", 2);