<PackageReference Include="System.Reactive" Version="4.1.3" />

SkipUntil<TSource>

sealed class SkipUntil<TSource> : Producer<TSource, _<TSource>>
using System.Reactive.Concurrency; using System.Reactive.Disposables; namespace System.Reactive.Linq.ObservableImpl { internal sealed class SkipUntil<TSource> : Producer<TSource, SkipUntil<TSource>._> { internal sealed class _ : IdentitySink<TSource> { private bool _open; private IDisposable _task; public _(IObserver<TSource> observer) : base(observer) { } public void Run(SkipUntil<TSource> parent) { Disposable.SetSingle(ref _task, Scheduler.ScheduleAction<_>(parent._scheduler, this, parent._startTime, (Action<_>)delegate(_ state) { state.Tick(); })); Run(parent._source); } protected override void Dispose(bool disposing) { if (disposing) Disposable.TryDispose(ref _task); base.Dispose(disposing); } private void Tick() { _open = true; } public override void OnNext(TSource value) { if (_open) ForwardOnNext(value); } } private readonly IObservable<TSource> _source; private readonly DateTimeOffset _startTime; internal readonly IScheduler _scheduler; public SkipUntil(IObservable<TSource> source, DateTimeOffset startTime, IScheduler scheduler) { _source = source; _startTime = startTime; _scheduler = scheduler; } public IObservable<TSource> Combine(DateTimeOffset startTime) { if (startTime <= _startTime) return this; return new SkipUntil<TSource>(_source, startTime, _scheduler); } protected override _ CreateSink(IObserver<TSource> observer) { return new _(observer); } protected override void Run(_ sink) { sink.Run(this); } } }