"""
测试 vools-reactive Connectable Observable
"""
import sys
import os
import asyncio
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from vools.reactive import Observable, Subject, BehaviorSubject, ReplaySubject, ops
from vools.reactive.connectable import (
ConnectableObservable,
publish,
share,
ref_count,
replay,
publish_replay,
auto_connect
)
def test_publish_basic():
"""测试 publish 操作符基础功能"""
print("=== publish 基础测试 ===")
result1 = []
result2 = []
source = Observable.from_iterable([1, 2, 3])
connectable = source.pipe(publish())
sub1 = connectable.subscribe(on_next=lambda x: result1.append(x))
assert result1 == [], f"Expected [], got {result1}"
print("[OK] subscribe before connect - no emission")
connection = connectable.connect()
assert result1 == [1, 2, 3], f"Expected [1,2,3], got {result1}"
print("[OK] connect triggers emission")
sub1.unsubscribe()
sub2 = connectable.subscribe(on_next=lambda x: result2.append(x))
print("[OK] publish basic")
def test_publish_multicast():
"""测试 publish 多播功能"""
print("\n=== publish 多播测试 ===")
results = []
source = Observable.from_iterable([1, 2, 3])
connectable = source.pipe(publish())
sub1_results = []
sub2_results = []
sub1 = connectable.subscribe(on_next=lambda x: sub1_results.append(x))
sub2 = connectable.subscribe(on_next=lambda x: sub2_results.append(x))
connectable.connect()
assert sub1_results == [1, 2, 3], f"Expected [1,2,3], got {sub1_results}"
assert sub2_results == [1, 2, 3], f"Expected [1,2,3], got {sub2_results}"
print("[OK] multicast to multiple subscribers")
sub1.unsubscribe()
sub2.unsubscribe()
def test_share():
"""测试 share 操作符 - 自动连接"""
print("\n=== share 测试 ===")
results = []
source = Observable.from_iterable([1, 2, 3])
shared = source.pipe(share())
sub1_results = []
shared.subscribe(on_next=lambda x: sub1_results.append(x))
assert sub1_results == [1, 2, 3], f"Expected [1,2,3], got {sub1_results}"
print("[OK] first subscriber triggers auto-connect")
sub2_results = []
shared.subscribe(on_next=lambda x: sub2_results.append(x))
assert sub2_results == [1, 2, 3], f"Expected [1,2,3], got {sub2_results}"
print("[OK] second subscriber receives replayed values")
def test_ref_count():
"""测试 ref_count 操作符"""
print("\n=== ref_count 测试 ===")
emit_count = [0]
subscribe_count = [0]
def make_counter_source():
"""创建一个记录发射次数的 source"""
return Observable.from_iterable([1, 2, 3]).pipe(
ops.do_on_next(lambda x: emit_count.__setitem__(0, emit_count[0] + 1))
)
source = make_counter_source().pipe(ref_count())
result1 = []
sub1 = source.subscribe(on_next=lambda x: result1.append(x))
assert emit_count[0] == 3, f"Expected 3 emissions (one per value), got {emit_count[0]}"
print("[OK] first subscriber triggers connection")
result2 = []
sub2 = source.subscribe(on_next=lambda x: result2.append(x))
assert emit_count[0] == 3, f"Expected still 3 emissions, got {emit_count[0]}"
print("[OK] second subscriber does not re-trigger connection")
assert result1 == [1, 2, 3], f"Expected [1,2,3], got {result1}"
assert result2 == [1, 2, 3], f"Expected [1,2,3], got {result2}"
sub1.unsubscribe()
sub2.unsubscribe()
print("[OK] ref_count")
def test_replay():
"""测试 replay 操作符"""
print("\n=== replay 测试 ===")
results = []
source = Observable.from_iterable([1, 2, 3])
connectable = source.pipe(replay())
sub1_results = []
sub1 = connectable.subscribe(on_next=lambda x: sub1_results.append(x))
connectable.connect()
assert sub1_results == [1, 2, 3], f"Expected [1,2,3], got {sub1_results}"
print("[OK] first subscriber receives all values")
sub2_results = []
sub2 = connectable.subscribe(on_next=lambda x: sub2_results.append(x))
assert sub2_results == [1, 2, 3], f"Expected [1,2,3], got {sub2_results}"
print("[OK] second subscriber receives replayed values")
sub1.unsubscribe()
sub2.unsubscribe()
def test_replay_with_buffer_size():
"""测试带缓冲大小的 replay"""
print("\n=== replay(buffer_size) 测试 ===")
source = Observable.from_iterable([1, 2, 3, 4, 5])
connectable = source.pipe(replay(buffer_size=2))
sub1 = connectable.subscribe()
connectable.connect()
sub2_results = []
sub2 = connectable.subscribe(on_next=lambda x: sub2_results.append(x))
assert sub2_results == [4, 5], f"Expected [4,5], got {sub2_results}"
print("[OK] replay with buffer_size=2")
sub1.unsubscribe()
sub2.unsubscribe()
def test_publish_replay():
"""测试 publish_replay 操作符"""
print("\n=== publish_replay 测试 ===")
results = []
source = Observable.from_iterable([1, 2, 3])
connectable = source.pipe(publish_replay())
sub1_results = []
sub1 = connectable.subscribe(on_next=lambda x: sub1_results.append(x))
connectable.connect()
assert sub1_results == [1, 2, 3], f"Expected [1,2,3], got {sub1_results}"
print("[OK] publish_replay first subscriber")
sub2_results = []
sub2 = connectable.subscribe(on_next=lambda x: sub2_results.append(x))
assert sub2_results == [1, 2, 3], f"Expected [1,2,3], got {sub2_results}"
print("[OK] publish_replay second subscriber")
sub1.unsubscribe()
sub2.unsubscribe()
def test_auto_connect():
"""测试 auto_connect 操作符"""
print("\n=== auto_connect 测试 ===")
emit_count = [0]
source = Observable.from_iterable([1, 2, 3]).pipe(
ops.do_on_next(lambda x: emit_count.__setitem__(0, emit_count[0] + 1))
)
connectable = source.pipe(publish())
auto = connectable.pipe(auto_connect(num_subscriptions=2))
result1 = []
sub1 = auto.subscribe(on_next=lambda x: result1.append(x))
assert emit_count[0] == 0, f"Expected 0, got {emit_count[0]}"
print("[OK] first subscriber does not trigger auto-connect")
result2 = []
sub2 = auto.subscribe(on_next=lambda x: result2.append(x))
assert emit_count[0] == 3, f"Expected 3, got {emit_count[0]}"
print("[OK] second subscriber triggers auto-connect")
assert result1 == [1, 2, 3], f"Expected [1,2,3], got {result1}"
assert result2 == [1, 2, 3], f"Expected [1,2,3], got {result2}"
sub1.unsubscribe()
sub2.unsubscribe()
def test_connectable_observable_connect_twice():
"""测试 ConnectableObservable 的 connect 方法只会连接一次"""
print("\n=== connect only once 测试 ===")
emission_count = [0]
source = Observable.from_iterable([1, 2, 3]).pipe(
ops.do_on_next(lambda x: emission_count.__setitem__(0, emission_count[0] + 1))
)
connectable = source.pipe(publish())
sub1 = connectable.subscribe()
connectable.connect()
connectable.connect()
assert emission_count[0] == 3, f"Expected 3 (one per value), got {emission_count[0]}"
print("[OK] connect() only triggers source subscription once")
sub1.unsubscribe()
def test_connectable_unsubscribe_before_connect():
"""测试在 connect 之前取消订阅"""
print("\n=== unsubscribe before connect 测试 ===")
results = []
source = Observable.from_iterable([1, 2, 3])
connectable = source.pipe(publish())
sub1 = connectable.subscribe(on_next=lambda x: results.append(x))
sub1.unsubscribe()
connectable.connect()
assert results == [], f"Expected [], got {results}"
print("[OK] unsubscribed before connect receives no values")
if __name__ == "__main__":
tests = [
test_publish_basic,
test_publish_multicast,
test_share,
test_ref_count,
test_replay,
test_replay_with_buffer_size,
test_publish_replay,
test_auto_connect,
test_connectable_observable_connect_twice,
test_connectable_unsubscribe_before_connect,
]
passed = 0
for test in tests:
try:
test()
passed += 1
except Exception as e:
print(f"[FAIL] {test.__name__}: {e}")
import traceback
traceback.print_exc()
print(f"\n{'='*50}")
print(f"Connectable Observable Tests: {passed}/{len(tests)} passed")
print(f"Coverage: {passed/len(tests)*100:.1f}%")