From 33c4f582ec15d40a7a6dea119e8ee392f18f25d3 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Fri, 17 Apr 2026 17:57:29 +0000 Subject: [PATCH 1/2] fix(operators): forward scheduler in pairwise, to_marbles, delay_with_mapper Three operators were subscribing to their source without forwarding the scheduler received at subscription time, breaking scheduler propagation for downstream operators. * _pairwise.py: add scheduler=scheduler to source.subscribe() * _tomarbles.py: add scheduler=scheduler to source.subscribe() * _delaywithmapper.py: add scheduler=scheduler to sub_delay.subscribe() in the subscription-delay path (the main source path was already fixed) The other operators originally listed in issue #480 (delay_with_mapper main path, throttle_first, timeout_with_mapper) were already forwarding the scheduler correctly. Closes #480 (partial) Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- changes.md | 7 +++++++ reactivex/operators/_delaywithmapper.py | 2 +- reactivex/operators/_pairwise.py | 4 +++- reactivex/operators/_tomarbles.py | 4 +++- 4 files changed, 14 insertions(+), 3 deletions(-) diff --git a/changes.md b/changes.md index 9a97ee862..b13de7ea4 100644 --- a/changes.md +++ b/changes.md @@ -2,6 +2,13 @@ ## Unreleased +- 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). + - Testing: Fixed ruff lint issues in `tests/test_scheduler/` (import ordering, f-string upgrades, `object` base class removal) and removed it from the ruff exclude list. `tests/test_subject/` also removed from the ruff exclude 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) From f4ba4c63ff2fe3d3952a045069220ceaef46ad68 Mon Sep 17 00:00:00 2001 From: Dag Brattli Date: Fri, 17 Apr 2026 23:38:49 +0200 Subject: [PATCH 2/2] test(operators): add regression tests for scheduler propagation Cover pairwise, to_marbles, and delay_with_mapper (subscription-delay path): each test installs a source whose subscribe captures the scheduler it received, pipes through the operator, subscribes with a known scheduler, and asserts the captured scheduler matches. Verified to fail on the pre-fix code. Co-Authored-By: Claude Opus 4.7 (1M context) --- tests/test_observable/test_delaywithmapper.py | 37 +++++++++++++++++++ tests/test_observable/test_marbles.py | 24 +++++++++++- tests/test_observable/test_pairwise.py | 20 ++++++++++ 3 files changed, 80 insertions(+), 1 deletion(-) 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