Generate.cs 9.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289
  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. using System.Reactive.Concurrency;
  5. using System.Reactive.Disposables;
  6. namespace System.Reactive.Linq.ObservableImpl
  7. {
  8. internal sealed class Generate<TState, TResult> : Producer<TResult>
  9. {
  10. private readonly TState _initialState;
  11. private readonly Func<TState, bool> _condition;
  12. private readonly Func<TState, TState> _iterate;
  13. private readonly Func<TState, TResult> _resultSelector;
  14. private readonly Func<TState, DateTimeOffset> _timeSelectorA;
  15. private readonly Func<TState, TimeSpan> _timeSelectorR;
  16. private readonly IScheduler _scheduler;
  17. public Generate(TState initialState, Func<TState, bool> condition, Func<TState, TState> iterate, Func<TState, TResult> resultSelector, IScheduler scheduler)
  18. {
  19. _initialState = initialState;
  20. _condition = condition;
  21. _iterate = iterate;
  22. _resultSelector = resultSelector;
  23. _scheduler = scheduler;
  24. }
  25. public Generate(TState initialState, Func<TState, bool> condition, Func<TState, TState> iterate, Func<TState, TResult> resultSelector, Func<TState, DateTimeOffset> timeSelector, IScheduler scheduler)
  26. {
  27. _initialState = initialState;
  28. _condition = condition;
  29. _iterate = iterate;
  30. _resultSelector = resultSelector;
  31. _timeSelectorA = timeSelector;
  32. _scheduler = scheduler;
  33. }
  34. public Generate(TState initialState, Func<TState, bool> condition, Func<TState, TState> iterate, Func<TState, TResult> resultSelector, Func<TState, TimeSpan> timeSelector, IScheduler scheduler)
  35. {
  36. _initialState = initialState;
  37. _condition = condition;
  38. _iterate = iterate;
  39. _resultSelector = resultSelector;
  40. _timeSelectorR = timeSelector;
  41. _scheduler = scheduler;
  42. }
  43. protected override IDisposable Run(IObserver<TResult> observer, IDisposable cancel, Action<IDisposable> setSink)
  44. {
  45. if (_timeSelectorA != null)
  46. {
  47. var sink = new SelectorA(this, observer, cancel);
  48. setSink(sink);
  49. return sink.Run();
  50. }
  51. else if (_timeSelectorR != null)
  52. {
  53. var sink = new Delta(this, observer, cancel);
  54. setSink(sink);
  55. return sink.Run();
  56. }
  57. else
  58. {
  59. var sink = new _(this, observer, cancel);
  60. setSink(sink);
  61. return sink.Run();
  62. }
  63. }
  64. class SelectorA : Sink<TResult>
  65. {
  66. private readonly Generate<TState, TResult> _parent;
  67. public SelectorA(Generate<TState, TResult> parent, IObserver<TResult> observer, IDisposable cancel)
  68. : base(observer, cancel)
  69. {
  70. _parent = parent;
  71. }
  72. private bool _first;
  73. private bool _hasResult;
  74. private TResult _result;
  75. public IDisposable Run()
  76. {
  77. _first = true;
  78. _hasResult = false;
  79. _result = default(TResult);
  80. return _parent._scheduler.Schedule(_parent._initialState, InvokeRec);
  81. }
  82. private IDisposable InvokeRec(IScheduler self, TState state)
  83. {
  84. var time = default(DateTimeOffset);
  85. if (_hasResult)
  86. base._observer.OnNext(_result);
  87. try
  88. {
  89. if (_first)
  90. _first = false;
  91. else
  92. state = _parent._iterate(state);
  93. _hasResult = _parent._condition(state);
  94. if (_hasResult)
  95. {
  96. _result = _parent._resultSelector(state);
  97. time = _parent._timeSelectorA(state);
  98. }
  99. }
  100. catch (Exception exception)
  101. {
  102. base._observer.OnError(exception);
  103. base.Dispose();
  104. return Disposable.Empty;
  105. }
  106. if (!_hasResult)
  107. {
  108. base._observer.OnCompleted();
  109. base.Dispose();
  110. return Disposable.Empty;
  111. }
  112. return self.Schedule(state, time, InvokeRec);
  113. }
  114. }
  115. class Delta : Sink<TResult>
  116. {
  117. private readonly Generate<TState, TResult> _parent;
  118. public Delta(Generate<TState, TResult> parent, IObserver<TResult> observer, IDisposable cancel)
  119. : base(observer, cancel)
  120. {
  121. _parent = parent;
  122. }
  123. private bool _first;
  124. private bool _hasResult;
  125. private TResult _result;
  126. public IDisposable Run()
  127. {
  128. _first = true;
  129. _hasResult = false;
  130. _result = default(TResult);
  131. return _parent._scheduler.Schedule(_parent._initialState, InvokeRec);
  132. }
  133. private IDisposable InvokeRec(IScheduler self, TState state)
  134. {
  135. var time = default(TimeSpan);
  136. if (_hasResult)
  137. base._observer.OnNext(_result);
  138. try
  139. {
  140. if (_first)
  141. _first = false;
  142. else
  143. state = _parent._iterate(state);
  144. _hasResult = _parent._condition(state);
  145. if (_hasResult)
  146. {
  147. _result = _parent._resultSelector(state);
  148. time = _parent._timeSelectorR(state);
  149. }
  150. }
  151. catch (Exception exception)
  152. {
  153. base._observer.OnError(exception);
  154. base.Dispose();
  155. return Disposable.Empty;
  156. }
  157. if (!_hasResult)
  158. {
  159. base._observer.OnCompleted();
  160. base.Dispose();
  161. return Disposable.Empty;
  162. }
  163. return self.Schedule(state, time, InvokeRec);
  164. }
  165. }
  166. class _ : Sink<TResult>
  167. {
  168. private readonly Generate<TState, TResult> _parent;
  169. public _(Generate<TState, TResult> parent, IObserver<TResult> observer, IDisposable cancel)
  170. : base(observer, cancel)
  171. {
  172. _parent = parent;
  173. }
  174. private TState _state;
  175. private bool _first;
  176. public IDisposable Run()
  177. {
  178. _state = _parent._initialState;
  179. _first = true;
  180. var longRunning = _parent._scheduler.AsLongRunning();
  181. if (longRunning != null)
  182. {
  183. return longRunning.ScheduleLongRunning(Loop);
  184. }
  185. else
  186. {
  187. return _parent._scheduler.Schedule(LoopRec);
  188. }
  189. }
  190. private void Loop(ICancelable cancel)
  191. {
  192. while (!cancel.IsDisposed)
  193. {
  194. var hasResult = false;
  195. var result = default(TResult);
  196. try
  197. {
  198. if (_first)
  199. _first = false;
  200. else
  201. _state = _parent._iterate(_state);
  202. hasResult = _parent._condition(_state);
  203. if (hasResult)
  204. result = _parent._resultSelector(_state);
  205. }
  206. catch (Exception exception)
  207. {
  208. base._observer.OnError(exception);
  209. base.Dispose();
  210. return;
  211. }
  212. if (hasResult)
  213. base._observer.OnNext(result);
  214. else
  215. break;
  216. }
  217. if (!cancel.IsDisposed)
  218. base._observer.OnCompleted();
  219. base.Dispose();
  220. }
  221. private void LoopRec(Action recurse)
  222. {
  223. var hasResult = false;
  224. var result = default(TResult);
  225. try
  226. {
  227. if (_first)
  228. _first = false;
  229. else
  230. _state = _parent._iterate(_state);
  231. hasResult = _parent._condition(_state);
  232. if (hasResult)
  233. result = _parent._resultSelector(_state);
  234. }
  235. catch (Exception exception)
  236. {
  237. base._observer.OnError(exception);
  238. base.Dispose();
  239. return;
  240. }
  241. if (hasResult)
  242. {
  243. base._observer.OnNext(result);
  244. recurse();
  245. }
  246. else
  247. {
  248. base._observer.OnCompleted();
  249. base.Dispose();
  250. }
  251. }
  252. }
  253. }
  254. }