Added behavior subject
This commit is contained in:
parent
9d05d76cfa
commit
d7d41b878c
23
src/cpl_reactive_extensions/behavior_subject.py
Normal file
23
src/cpl_reactive_extensions/behavior_subject.py
Normal file
@ -0,0 +1,23 @@
|
|||||||
|
from cpl_core.type import T
|
||||||
|
from cpl_reactive_extensions.observable import Observable
|
||||||
|
|
||||||
|
|
||||||
|
class BehaviorSubject(Observable):
|
||||||
|
def __init__(self, _t: type, value: T):
|
||||||
|
Observable.__init__(self, lambda x: x)
|
||||||
|
|
||||||
|
if not isinstance(value, _t):
|
||||||
|
raise TypeError(f"Expected {_t.__name__} not {type(value).__name__}")
|
||||||
|
|
||||||
|
self._t = _t
|
||||||
|
self._value = value
|
||||||
|
|
||||||
|
@property
|
||||||
|
def value(self) -> T:
|
||||||
|
return self._value
|
||||||
|
|
||||||
|
def next(self, value: T):
|
||||||
|
if not isinstance(value, self._t):
|
||||||
|
raise TypeError(f"Expected {self._t.__name__} not {type(value).__name__}")
|
||||||
|
|
||||||
|
self._value = value
|
@ -7,14 +7,3 @@ class Subject(Observable):
|
|||||||
Observable.__init__(self)
|
Observable.__init__(self)
|
||||||
|
|
||||||
self._t = _t
|
self._t = _t
|
||||||
self._value: T = None
|
|
||||||
|
|
||||||
@property
|
|
||||||
def value(self) -> T:
|
|
||||||
return self._value
|
|
||||||
|
|
||||||
def next(self, value: T):
|
|
||||||
if not isinstance(value, self._t):
|
|
||||||
raise TypeError(f"Expected {self._t.__name__} not {type(value).__name__}")
|
|
||||||
|
|
||||||
self._value = value
|
|
||||||
|
@ -4,6 +4,7 @@ import unittest
|
|||||||
from threading import Timer
|
from threading import Timer
|
||||||
|
|
||||||
from cpl_core.console import Console
|
from cpl_core.console import Console
|
||||||
|
from cpl_reactive_extensions.behavior_subject import BehaviorSubject
|
||||||
from cpl_reactive_extensions.observable import Observable
|
from cpl_reactive_extensions.observable import Observable
|
||||||
from cpl_reactive_extensions.observer import Observer
|
from cpl_reactive_extensions.observer import Observer
|
||||||
from cpl_reactive_extensions.subject import Subject
|
from cpl_reactive_extensions.subject import Subject
|
||||||
@ -114,3 +115,14 @@ class ReactiveTestCase(unittest.TestCase):
|
|||||||
observable.subscribe(subject, self._on_error)
|
observable.subscribe(subject, self._on_error)
|
||||||
|
|
||||||
self.assertFalse(self._error)
|
self.assertFalse(self._error)
|
||||||
|
|
||||||
|
def test_behavior_subject(self):
|
||||||
|
subject = BehaviorSubject(int, 0)
|
||||||
|
|
||||||
|
subject.subscribe(lambda x: Console.write_line("a", x))
|
||||||
|
|
||||||
|
subject.next(1)
|
||||||
|
subject.next(2)
|
||||||
|
|
||||||
|
subject.subscribe(lambda x: Console.write_line("b", x))
|
||||||
|
subject.next(3)
|
||||||
|
Loading…
Reference in New Issue
Block a user