使用BehavourSubject<T> 支持属性
解决这个问题的一个有用的想法是使用BehaviorSubject<bool> 类型支持您的属性。这些有用地服务于 active 作为属性和该属性的值流的双重目的。
您可以将它们作为 observable 订阅,也可以通过 Value 属性访问它们的当前值。您可以通过OnNext 发送一个新值来更改它们。
例如,我们可以这样做:
private BehaviorSubject<bool> _isActive;
public bool IsActive
{
get { return _isActive.Value; }
set { _isActive.OnNext(value); }
}
在您的所有属性中都设置了这一点后,就可以非常简单地查看属性是否符合您陈述的复杂条件。假设 _status 和 _isBroken 是类似实现的支持主题,我们可以像这样设置订阅:
Observable.CombineLatest(_isActive,
_status,
(a,s) => a & !s).DistinctUntilChanged()
.Where(p => p)
.SelectMany(_ => Observable.Timer(TimeSpan.FromSeconds(3), scheduler)
.TakeUntil(_status.Where(st => st)))
.Subscribe(_ => _isBroken.OnNext(true));
部分行使用CombineLatest 并订阅_isActive 和_status 流。每当这些变化中的任何一个发生变化时,它都会发出 - 当_isActive 为真且_status 为假时,结果函数恰好设置了一个真值。 DistinctUntilChanged() 可防止将 _isActive 和 _status 设置为它们已有的值以启动新计时器。
然后我们使用Where 只过滤这个条件。
SelectMany 现在将采用真值并将每个值投射到 3 秒后发出的流中,使用 Timer - 但是我们使用 TakeUntil 压缩这个值_status 变为真的事件。 SelectMany 还将流流平展回单个布尔流。
这里不确定 - 你没有提到它,但你可能想考虑 _isActive 变为 false 是否也应该终止计时器。如果是这种情况,您可以使用Merge 在TakeUntil 中将对此的监视和_status 结合起来。
我们现在可以订阅这整个事情来设置_isBroken true 如果这个查询被触发,表明计时器已过期。
注意 Timer 的 scheduler 参数 - 这是存在的,因此我们可以传入测试调度程序。
我不确定我是否正确地捕捉到了您的所有逻辑 - 但如果不是,希望您能看到如何根据需要进行修改。
这是完整的示例。使用 nuget 包 rx-testing,这将在 LINQPad 中运行:
void Main()
{
var tests = new Tests();
tests.Test();
}
public class Foo
{
private BehaviorSubject<bool> _isActive;
private BehaviorSubject<bool> _isBroken;
private BehaviorSubject<bool> _status;
public bool IsActive
{
get { return _isActive.Value; }
set { _isActive.OnNext(value); }
}
public bool IsBroken { get { return _isBroken.Value; } }
public bool Status { get { return _status.Value; } }
public Foo(IObservable<MyObject> valueStream, IScheduler scheduler)
{
_isActive = new BehaviorSubject<bool>(false);
_isBroken = new BehaviorSubject<bool>(false);
_status = new BehaviorSubject<bool>(false);
// for debugging purposes
_isActive.Subscribe(a => Console.WriteLine(
"Time: " + scheduler.Now.TimeOfDay + " IsActive: " + a));
_isBroken.Subscribe(a => Console.WriteLine(
"Time: " + scheduler.Now.TimeOfDay + " IsBroken: " + a));
_status.Subscribe(a => Console.WriteLine(
"Time: " + scheduler.Now.TimeOfDay + " Status: " + a));
valueStream.Subscribe(UpdateValues);
Observable.CombineLatest(
_isActive,
_status,
(a,s) => a & !s).DistinctUntilChanged()
.Where(p => p)
.SelectMany(_ => Observable.Timer(TimeSpan.FromSeconds(3),
scheduler)
.TakeUntil(_status.Where(st => st)))
.Subscribe(_ => _isBroken.OnNext(true));
}
private void UpdateValues(MyObject obj)
{
_status.OnNext(obj.SpecialValue);
}
}
public class MyObject
{
public MyObject(bool specialValue)
{
SpecialValue = specialValue;
}
public bool SpecialValue { get; set; }
}
public class Tests : ReactiveTest
{
public void Test()
{
var testScheduler = new TestScheduler();
var statusStream = testScheduler.CreateColdObservable<bool>(
OnNext(TimeSpan.FromSeconds(1).Ticks, false),
OnNext(TimeSpan.FromSeconds(3).Ticks, true),
OnNext(TimeSpan.FromSeconds(5).Ticks, false));
var activeStream = testScheduler.CreateColdObservable<bool>(
OnNext(TimeSpan.FromSeconds(1).Ticks, false),
OnNext(TimeSpan.FromSeconds(6).Ticks, true));
var foo = new Foo(statusStream.Select(b => new MyObject(b)), testScheduler);
activeStream.Subscribe(b => foo.IsActive = b);
testScheduler.Start();
}
}
回复评论
如果你想让 isActive false 设置 isBroken false,那么我认为这加起来现在是这样说的:
isActive isStatus Action
T F Set Broken True after 3 seconds unless any other result occurs
T T Set Broken False immediately if not already false, cancel timer
F F Set Broken False immediately if not already false, cancel timer
F T Set Broken False immediately if not already false, cancel timer
在这种情况下,请使用以下查询:
Observable.CombineLatest(
_isActive,
_status,
(a,s) => a & !s).DistinctUntilChanged()
.Select(p => p ? Observable.Timer(TimeSpan.FromSeconds(3),
scheduler)
.Select(_ => true)
: Observable.Return(false))
.Switch()
.DistinctUntilChanged()
.Subscribe(res => _isBroken.OnNext(res));
注意变化:
-
SelectMany 现在是 Select,它将每个事件变成
- 3 秒后返回
true 的计时器
- 或直接
false
-
Select 的结果是一个布尔流:IObservable<IObservable<bool>>。我们希望任何出现的新流都取消任何先前的流。这就是Switch 将要做的——在过程中使结果变平。
- 我们现在应用第二个
DistinctUntilChanged(),因为取消的计时器可能会导致两个错误值连续出现在流中
- 最后我们将新出现的布尔值赋给
isBroken。