diff --git a/changes.md b/changes.md index b3dcb69ed..739473d66 100644 --- a/changes.md +++ b/changes.md @@ -2,11 +2,21 @@ ## Unreleased +<<<<<<< repo-assist/fix-issue-480-scheduler-forwarding-ea9c5c512fb11be9 +- Operators: Fixed scheduler forwarding in `pairwise`, `to_marbles`, and + `delay_with_mapper` (subscription-delay path). These operators now pass the + `scheduler` argument through to `source.subscribe(...)` and, in the case of + `delay_with_mapper`, to the subscription-delay observable, consistent with + all other pipeable operators. Closes #480 (partial — the operators listed + in the issue that were not yet fixed). + +======= - CI: Skip `tests/test_scheduler/test_mainloop/test_tkinterscheduler.py` on PyPy (all platforms). The module creates a Tk root at import time; PyPy's `_tkinter` finalizer calls `threading.notify_all` during interpreter shutdown and aborts the xdist worker, failing whichever unrelated test was running. Previously guarded only on macOS+PyPy. +>>>>>>> master - Fix: `reactivex.timer(duetime, period)` now correctly resets the initial delay on each resubscription (e.g. via `repeat()`). Previously `nonlocal duetime` in the subscribe closure mutated the shared outer variable so diff --git a/reactivex/operators/_delaywithmapper.py b/reactivex/operators/_delaywithmapper.py index 45c4626d1..cc7b60b47 100644 --- a/reactivex/operators/_delaywithmapper.py +++ b/reactivex/operators/_delaywithmapper.py @@ -95,7 +95,7 @@ def on_completed() -> None: start() else: subscription.disposable = sub_delay.subscribe( - lambda _: start(), observer.on_error, start + lambda _: start(), observer.on_error, start, scheduler=scheduler ) return CompositeDisposable(subscription, delays) diff --git a/reactivex/operators/_pairwise.py b/reactivex/operators/_pairwise.py index d1b4a9b0e..ec2e8d02f 100644 --- a/reactivex/operators/_pairwise.py +++ b/reactivex/operators/_pairwise.py @@ -49,7 +49,9 @@ def on_next(x: _T) -> None: if pair: observer.on_next(pair) - return source.subscribe(on_next, observer.on_error, observer.on_completed) + return source.subscribe( + on_next, observer.on_error, observer.on_completed, scheduler=scheduler + ) return Observable(subscribe) diff --git a/reactivex/operators/_tomarbles.py b/reactivex/operators/_tomarbles.py index faaef60a7..eab02106b 100644 --- a/reactivex/operators/_tomarbles.py +++ b/reactivex/operators/_tomarbles.py @@ -59,7 +59,9 @@ def on_completed(): observer.on_next("".join(n for n in result)) observer.on_completed() - return source.subscribe(on_next, on_error, on_completed) + return source.subscribe( + on_next, on_error, on_completed, scheduler=scheduler + ) return Observable(subscribe) diff --git a/tests/test_observable/test_delaywithmapper.py b/tests/test_observable/test_delaywithmapper.py index df14bd836..cfe8da95b 100644 --- a/tests/test_observable/test_delaywithmapper.py +++ b/tests/test_observable/test_delaywithmapper.py @@ -1,6 +1,9 @@ import unittest +from reactivex import Observable, abc from reactivex import operators as ops +from reactivex.disposable import Disposable +from reactivex.scheduler import ImmediateScheduler from reactivex.testing import ReactiveTest, TestScheduler on_next = ReactiveTest.on_next @@ -199,3 +202,37 @@ def create(): assert results.messages == [on_next(210 + 50, 2)] assert xs.subscriptions == [subscribe(200, 300)] assert ys.subscriptions == [subscribe(210, 260)] + + def test_delay_with_mapper_forwards_scheduler_to_subscription_delay( + self, + ) -> None: + captured: dict[str, abc.SchedulerBase | None] = {} + expected: abc.SchedulerBase = ImmediateScheduler() + + def sub_delay_subscribe( + observer: abc.ObserverBase[int], + scheduler: abc.SchedulerBase | None = None, + ) -> abc.DisposableBase: + captured["scheduler"] = scheduler + observer.on_completed() + return Disposable() + + sub_delay: Observable[int] = Observable(sub_delay_subscribe) + + def source_subscribe( + observer: abc.ObserverBase[int], + scheduler: abc.SchedulerBase | None = None, + ) -> abc.DisposableBase: + observer.on_completed() + return Disposable() + + source: Observable[int] = Observable(source_subscribe) + + def mapper(_: int) -> Observable[int]: + return source + + source.pipe(ops.delay_with_mapper(sub_delay, mapper)).subscribe( + scheduler=expected + ) + + assert captured["scheduler"] is expected diff --git a/tests/test_observable/test_marbles.py b/tests/test_observable/test_marbles.py index 81384d0e2..0350375d0 100644 --- a/tests/test_observable/test_marbles.py +++ b/tests/test_observable/test_marbles.py @@ -2,8 +2,11 @@ import unittest import reactivex -from reactivex import notification +from reactivex import Observable, abc, notification +from reactivex import operators as ops +from reactivex.disposable import Disposable from reactivex.observable.marbles import parse +from reactivex.scheduler import ImmediateScheduler from reactivex.testing import TestScheduler from reactivex.testing.reactivetest import ReactiveTest @@ -578,3 +581,22 @@ def test_hot_marble_with_timedelta(self): ReactiveTest.on_completed(300.7), ] assert results == expected + + +class TestToMarbles(unittest.TestCase): + def test_to_marbles_forwards_scheduler_to_source(self) -> None: + captured: dict[str, abc.SchedulerBase | None] = {} + expected: abc.SchedulerBase = ImmediateScheduler() + + def source_subscribe( + observer: abc.ObserverBase[int], + scheduler: abc.SchedulerBase | None = None, + ) -> abc.DisposableBase: + captured["scheduler"] = scheduler + observer.on_completed() + return Disposable() + + source: Observable[int] = Observable(source_subscribe) + source.pipe(ops.to_marbles()).subscribe(scheduler=expected) + + assert captured["scheduler"] is expected diff --git a/tests/test_observable/test_pairwise.py b/tests/test_observable/test_pairwise.py index 613c8fe25..55e0708c5 100644 --- a/tests/test_observable/test_pairwise.py +++ b/tests/test_observable/test_pairwise.py @@ -1,6 +1,9 @@ import unittest +from reactivex import Observable, abc from reactivex import operators as ops +from reactivex.disposable import Disposable +from reactivex.scheduler import ImmediateScheduler from reactivex.testing import ReactiveTest, TestScheduler on_next = ReactiveTest.on_next @@ -134,3 +137,20 @@ def create(): assert results.messages == [on_next(240, (4, 3))] assert xs.subscriptions == [subscribe(200, 280)] + + def test_pairwise_forwards_scheduler_to_source(self) -> None: + captured: dict[str, abc.SchedulerBase | None] = {} + expected: abc.SchedulerBase = ImmediateScheduler() + + def source_subscribe( + observer: abc.ObserverBase[int], + scheduler: abc.SchedulerBase | None = None, + ) -> abc.DisposableBase: + captured["scheduler"] = scheduler + observer.on_completed() + return Disposable() + + source: Observable[int] = Observable(source_subscribe) + source.pipe(ops.pairwise()).subscribe(scheduler=expected) + + assert captured["scheduler"] is expected