SkipUntil.cs 6.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212
  1. // Licensed to the .NET Foundation under one or more agreements.
  2. // The .NET Foundation licenses this file to you under the Apache 2.0 License.
  3. // See the LICENSE file in the project root for more information.
  4. #if !NO_PERF
  5. using System;
  6. using System.Reactive.Concurrency;
  7. using System.Reactive.Disposables;
  8. using System.Threading;
  9. namespace System.Reactive.Linq.ObservableImpl
  10. {
  11. class SkipUntil<TSource, TOther> : Producer<TSource>
  12. {
  13. private readonly IObservable<TSource> _source;
  14. private readonly IObservable<TOther> _other;
  15. public SkipUntil(IObservable<TSource> source, IObservable<TOther> other)
  16. {
  17. _source = source;
  18. _other = other;
  19. }
  20. protected override IDisposable Run(IObserver<TSource> observer, IDisposable cancel, Action<IDisposable> setSink)
  21. {
  22. var sink = new _(this, observer, cancel);
  23. setSink(sink);
  24. return sink.Run();
  25. }
  26. class _ : Sink<TSource>
  27. {
  28. private readonly SkipUntil<TSource, TOther> _parent;
  29. public _(SkipUntil<TSource, TOther> parent, IObserver<TSource> observer, IDisposable cancel)
  30. : base(observer, cancel)
  31. {
  32. _parent = parent;
  33. }
  34. public IDisposable Run()
  35. {
  36. var sourceObserver = new T(this);
  37. var otherObserver = new O(this, sourceObserver);
  38. var sourceSubscription = _parent._source.SubscribeSafe(sourceObserver);
  39. var otherSubscription = _parent._other.SubscribeSafe(otherObserver);
  40. sourceObserver.Disposable = sourceSubscription;
  41. otherObserver.Disposable = otherSubscription;
  42. return StableCompositeDisposable.Create(
  43. sourceSubscription,
  44. otherSubscription
  45. );
  46. }
  47. class T : IObserver<TSource>
  48. {
  49. private readonly _ _parent;
  50. public volatile IObserver<TSource> _observer;
  51. private readonly SingleAssignmentDisposable _subscription;
  52. public T(_ parent)
  53. {
  54. _parent = parent;
  55. _observer = NopObserver<TSource>.Instance;
  56. _subscription = new SingleAssignmentDisposable();
  57. }
  58. public IDisposable Disposable
  59. {
  60. set { _subscription.Disposable = value; }
  61. }
  62. public void OnNext(TSource value)
  63. {
  64. _observer.OnNext(value);
  65. }
  66. public void OnError(Exception error)
  67. {
  68. _parent._observer.OnError(error);
  69. _parent.Dispose();
  70. }
  71. public void OnCompleted()
  72. {
  73. _observer.OnCompleted();
  74. _subscription.Dispose(); // We can't cancel the other stream yet, it may be on its way to dispatch an OnError message and we don't want to have a race.
  75. }
  76. }
  77. class O : IObserver<TOther>
  78. {
  79. private readonly _ _parent;
  80. private readonly T _sourceObserver;
  81. private readonly SingleAssignmentDisposable _subscription;
  82. public O(_ parent, T sourceObserver)
  83. {
  84. _parent = parent;
  85. _sourceObserver = sourceObserver;
  86. _subscription = new SingleAssignmentDisposable();
  87. }
  88. public IDisposable Disposable
  89. {
  90. set { _subscription.Disposable = value; }
  91. }
  92. public void OnNext(TOther value)
  93. {
  94. _sourceObserver._observer = _parent._observer;
  95. _subscription.Dispose();
  96. }
  97. public void OnError(Exception error)
  98. {
  99. _parent._observer.OnError(error);
  100. _parent.Dispose();
  101. }
  102. public void OnCompleted()
  103. {
  104. _subscription.Dispose();
  105. }
  106. }
  107. }
  108. }
  109. class SkipUntil<TSource> : Producer<TSource>
  110. {
  111. private readonly IObservable<TSource> _source;
  112. private readonly DateTimeOffset _startTime;
  113. internal readonly IScheduler _scheduler;
  114. public SkipUntil(IObservable<TSource> source, DateTimeOffset startTime, IScheduler scheduler)
  115. {
  116. _source = source;
  117. _startTime = startTime;
  118. _scheduler = scheduler;
  119. }
  120. public IObservable<TSource> Omega(DateTimeOffset startTime)
  121. {
  122. //
  123. // Maximum semantics:
  124. //
  125. // t 0--1--2--3--4--5--6--7-> t 0--1--2--3--4--5--6--7->
  126. //
  127. // xs --o--o--o--o--o--o--| xs --o--o--o--o--o--o--|
  128. // xs.SU(5AM) xxxxxxxxxxxxxxxx-o--| xs.SU(3AM) xxxxxxxxxx-o--o--o--|
  129. // xs.SU(5AM).SU(3AM) xxxxxxxxx--------o--| xs.SU(3AM).SU(5AM) xxxxxxxxxxxxxxxx-o--|
  130. //
  131. if (startTime <= _startTime)
  132. return this;
  133. else
  134. return new SkipUntil<TSource>(_source, startTime, _scheduler);
  135. }
  136. protected override IDisposable Run(IObserver<TSource> observer, IDisposable cancel, Action<IDisposable> setSink)
  137. {
  138. var sink = new _(this, observer, cancel);
  139. setSink(sink);
  140. return sink.Run();
  141. }
  142. class _ : Sink<TSource>, IObserver<TSource>
  143. {
  144. private readonly SkipUntil<TSource> _parent;
  145. private volatile bool _open;
  146. public _(SkipUntil<TSource> parent, IObserver<TSource> observer, IDisposable cancel)
  147. : base(observer, cancel)
  148. {
  149. _parent = parent;
  150. }
  151. public IDisposable Run()
  152. {
  153. var t = _parent._scheduler.Schedule(_parent._startTime, Tick);
  154. var d = _parent._source.SubscribeSafe(this);
  155. return StableCompositeDisposable.Create(t, d);
  156. }
  157. private void Tick()
  158. {
  159. _open = true;
  160. }
  161. public void OnNext(TSource value)
  162. {
  163. if (_open)
  164. base._observer.OnNext(value);
  165. }
  166. public void OnError(Exception error)
  167. {
  168. base._observer.OnError(error);
  169. base.Dispose();
  170. }
  171. public void OnCompleted()
  172. {
  173. base._observer.OnCompleted();
  174. base.Dispose();
  175. }
  176. }
  177. }
  178. }
  179. #endif