From c5970767efe65c3fab4e0db99d3a6bbc08c0fcbd Mon Sep 17 00:00:00 2001 From: ldm0 Date: Fri, 18 Sep 2026 13:00:26 +0800 Subject: [PATCH] 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. --- .../wpt-cross-current/passed-cases.txt | 4 + moli-renderer-v8/src/observable.rs | 4 + moli-renderer-v8/src/observable/observer.rs | 6 +- moli-renderer-v8/src/observable/transform.rs | 121 ++++++++-- .../src/script_vm/tests/observable.rs | 133 +++++++++++ .../src/worker/thread/tests/postmessage.rs | 17 ++ .../fixtures/observable-count-operators.js | 219 ++++++++++++++++++ moli-webidl/src/convert.rs | 31 ++- 8 files changed, 516 insertions(+), 19 deletions(-) create mode 100644 moli-renderer-v8/tests/fixtures/observable-count-operators.js diff --git a/moli-benchmark/wpt-cross-current/passed-cases.txt b/moli-benchmark/wpt-cross-current/passed-cases.txt index 839ed4b85..0a5afa6e8 100644 --- a/moli-benchmark/wpt-cross-current/passed-cases.txt +++ b/moli-benchmark/wpt-cross-current/passed-cases.txt @@ -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 diff --git a/moli-renderer-v8/src/observable.rs b/moli-renderer-v8/src/observable.rs index 8b322c236..a97b683d1 100644 --- a/moli-renderer-v8/src/observable.rs +++ b/moli-renderer-v8/src/observable.rs @@ -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)] diff --git a/moli-renderer-v8/src/observable/observer.rs b/moli-renderer-v8/src/observable/observer.rs index 2be031adb..1d5eb4901 100644 --- a/moli-renderer-v8/src/observable/observer.rs +++ b/moli-renderer-v8/src/observable/observer.rs @@ -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(); diff --git a/moli-renderer-v8/src/observable/transform.rs b/moli-renderer-v8/src/observable/transform.rs index 1448ce8a1..183384353 100644 --- a/moli-renderer-v8/src/observable/transform.rs +++ b/moli-renderer-v8/src/observable/transform.rs @@ -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::(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::(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> { + 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::::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]) { diff --git a/moli-renderer-v8/src/script_vm/tests/observable.rs b/moli-renderer-v8/src/script_vm/tests/observable.rs index 354af2492..2997c5014 100644 --- a/moli-renderer-v8/src/script_vm/tests/observable.rs +++ b/moli-renderer-v8/src/script_vm/tests/observable.rs @@ -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 = 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/"); diff --git a/moli-renderer-v8/src/worker/thread/tests/postmessage.rs b/moli-renderer-v8/src/worker/thread/tests/postmessage.rs index 59acb56fb..2d2690d12 100644 --- a/moli-renderer-v8/src/worker/thread/tests/postmessage.rs +++ b/moli-renderer-v8/src/worker/thread/tests/postmessage.rs @@ -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(); diff --git a/moli-renderer-v8/tests/fixtures/observable-count-operators.js b/moli-renderer-v8/tests/fixtures/observable-count-operators.js new file mode 100644 index 000000000..e5c8b016d --- /dev/null +++ b/moli-renderer-v8/tests/fixtures/observable-count-operators.js @@ -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}; +})() diff --git a/moli-webidl/src/convert.rs b/moli-webidl/src/convert.rs index b09a1b15b..95d8e7e30 100644 --- a/moli-webidl/src/convert.rs +++ b/moli-webidl/src/convert.rs @@ -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 { @@ -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);