<PackageReference Include="Relativity.Transfer.Client" Version="7.2.7" />

AverageDouble

sealed class AverageDouble : Producer<double, _>
namespace System.Reactive.Linq.ObservableImpl { internal sealed class AverageDouble : Producer<double, AverageDouble._> { internal sealed class _ : Sink<double>, IObserver<double> { private double _sum; private long _count; public _(IObserver<double> observer, IDisposable cancel) : base(observer, cancel) { _sum = 0; _count = 0; } public void OnNext(double value) { try { _sum += value; checked { _count++; } } catch (Exception error) { _observer.OnError(error); base.Dispose(); } } public void OnError(Exception error) { _observer.OnError(error); base.Dispose(); } public void OnCompleted() { if (_count > 0) { _observer.OnNext(_sum / (double)_count); _observer.OnCompleted(); } else _observer.OnError(new InvalidOperationException(Strings_Linq.NO_ELEMENTS)); base.Dispose(); } } private readonly IObservable<double> _source; public AverageDouble(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); } } }