Generate.cs 9.5 KB

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