<PackageReference Include="System.Reactive" Version="4.0.0-preview.2.build.379" />

MaxDoubleNullable

sealed class MaxDoubleNullable : Producer<double?, _>
namespace System.Reactive.Linq.ObservableImpl { internal sealed class MaxDoubleNullable : Producer<double?, MaxDoubleNullable._> { internal sealed class _ : Sink<double?>, IObserver<double?> { private double? _lastValue; public _(IObserver<double?> observer, IDisposable cancel) : base(observer, cancel) { _lastValue = null; } public void OnNext(double? value) { if (value.HasValue) { if (_lastValue.HasValue) { if (value > _lastValue || double.IsNaN(value.Value)) _lastValue = value; } else _lastValue = value; } } public void OnError(Exception error) { _observer.OnError(error); base.Dispose(); } public void OnCompleted() { _observer.OnNext(_lastValue); _observer.OnCompleted(); base.Dispose(); } } private readonly IObservable<double?> _source; public MaxDoubleNullable(IObservable<double?> source) { _source = source; } protected override _ CreateSink(IObserver<double?> observer, IDisposable cancel) { return new _(observer, cancel); } protected override IDisposable Run(_ sink) { return _source.SubscribeSafe(sink); } } }