ToObservable.cs 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208
  1. // Licensed to the .NET Foundation under one or more agreements.
  2. // The .NET Foundation licenses this file to you under the MIT License.
  3. // See the LICENSE file in the project root for more information.
  4. using System.Collections.Generic;
  5. using System.Reactive.Concurrency;
  6. using System.Reactive.Disposables;
  7. namespace System.Reactive.Linq.ObservableImpl
  8. {
  9. internal sealed class ToObservableRecursive<TSource> : Producer<TSource, ToObservableRecursive<TSource>._>
  10. {
  11. private readonly IEnumerable<TSource> _source;
  12. private readonly IScheduler _scheduler;
  13. public ToObservableRecursive(IEnumerable<TSource> source, IScheduler scheduler)
  14. {
  15. _source = source;
  16. _scheduler = scheduler;
  17. }
  18. protected override _ CreateSink(IObserver<TSource> observer) => new _(observer);
  19. protected override void Run(_ sink) => sink.Run(_source, _scheduler);
  20. internal sealed class _ : IdentitySink<TSource>
  21. {
  22. private IEnumerator<TSource> _enumerator;
  23. private volatile bool _disposed;
  24. public _(IObserver<TSource> observer)
  25. : base(observer)
  26. {
  27. }
  28. public void Run(IEnumerable<TSource> source, IScheduler scheduler)
  29. {
  30. try
  31. {
  32. _enumerator = source.GetEnumerator();
  33. }
  34. catch (Exception exception)
  35. {
  36. ForwardOnError(exception);
  37. return;
  38. }
  39. //
  40. // We never allow the scheduled work to be cancelled. Instead, the _disposed flag
  41. // is used to have LoopRec bail out and perform proper clean-up of the
  42. // enumerator.
  43. //
  44. scheduler.Schedule(this, (innerScheduler, @this) => @this.LoopRec(innerScheduler));
  45. }
  46. protected override void Dispose(bool disposing)
  47. {
  48. base.Dispose(disposing);
  49. if (disposing)
  50. {
  51. _disposed = true;
  52. }
  53. }
  54. private IDisposable LoopRec(IScheduler scheduler)
  55. {
  56. var hasNext = false;
  57. var ex = default(Exception);
  58. var current = default(TSource);
  59. var enumerator = _enumerator;
  60. if (_disposed)
  61. {
  62. _enumerator.Dispose();
  63. _enumerator = null;
  64. return Disposable.Empty;
  65. }
  66. try
  67. {
  68. hasNext = enumerator.MoveNext();
  69. if (hasNext)
  70. {
  71. current = enumerator.Current;
  72. }
  73. }
  74. catch (Exception exception)
  75. {
  76. ex = exception;
  77. }
  78. if (ex != null)
  79. {
  80. enumerator.Dispose();
  81. _enumerator = null;
  82. ForwardOnError(ex);
  83. return Disposable.Empty;
  84. }
  85. if (!hasNext)
  86. {
  87. enumerator.Dispose();
  88. _enumerator = null;
  89. ForwardOnCompleted();
  90. return Disposable.Empty;
  91. }
  92. ForwardOnNext(current);
  93. //
  94. // We never allow the scheduled work to be cancelled. Instead, the _disposed flag
  95. // is used to have LoopRec bail out and perform proper clean-up of the
  96. // enumerator.
  97. //
  98. scheduler.Schedule(this, (innerScheduler, @this) => @this.LoopRec(innerScheduler));
  99. return Disposable.Empty;
  100. }
  101. }
  102. }
  103. internal sealed class ToObservableLongRunning<TSource> : Producer<TSource, ToObservableLongRunning<TSource>._>
  104. {
  105. private readonly IEnumerable<TSource> _source;
  106. private readonly ISchedulerLongRunning _scheduler;
  107. public ToObservableLongRunning(IEnumerable<TSource> source, ISchedulerLongRunning scheduler)
  108. {
  109. _source = source;
  110. _scheduler = scheduler;
  111. }
  112. protected override _ CreateSink(IObserver<TSource> observer) => new _(observer);
  113. protected override void Run(_ sink) => sink.Run(_source, _scheduler);
  114. internal sealed class _ : IdentitySink<TSource>
  115. {
  116. public _(IObserver<TSource> observer)
  117. : base(observer)
  118. {
  119. }
  120. public void Run(IEnumerable<TSource> source, ISchedulerLongRunning scheduler)
  121. {
  122. IEnumerator<TSource> e;
  123. try
  124. {
  125. e = source.GetEnumerator();
  126. }
  127. catch (Exception exception)
  128. {
  129. ForwardOnError(exception);
  130. return;
  131. }
  132. SetUpstream(scheduler.ScheduleLongRunning((@this: this, e), (tuple, cancelable) => [email protected](tuple.e, cancelable)));
  133. }
  134. private void Loop(IEnumerator<TSource> enumerator, ICancelable cancel)
  135. {
  136. while (!cancel.IsDisposed)
  137. {
  138. var hasNext = false;
  139. var ex = default(Exception);
  140. var current = default(TSource);
  141. try
  142. {
  143. hasNext = enumerator.MoveNext();
  144. if (hasNext)
  145. {
  146. current = enumerator.Current;
  147. }
  148. }
  149. catch (Exception exception)
  150. {
  151. ex = exception;
  152. }
  153. if (ex != null)
  154. {
  155. ForwardOnError(ex);
  156. break;
  157. }
  158. if (!hasNext)
  159. {
  160. ForwardOnCompleted();
  161. break;
  162. }
  163. ForwardOnNext(current);
  164. }
  165. enumerator.Dispose();
  166. Dispose();
  167. }
  168. }
  169. }
  170. }