Generate.cs 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375
  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. #nullable disable
  5. using System.Reactive.Concurrency;
  6. using System.Reactive.Disposables;
  7. namespace System.Reactive.Linq.ObservableImpl
  8. {
  9. internal static class Generate<TState, TResult>
  10. {
  11. internal sealed class NoTime : Producer<TResult, NoTime._>
  12. {
  13. private readonly TState _initialState;
  14. private readonly Func<TState, bool> _condition;
  15. private readonly Func<TState, TState> _iterate;
  16. private readonly Func<TState, TResult> _resultSelector;
  17. private readonly IScheduler _scheduler;
  18. public NoTime(TState initialState, Func<TState, bool> condition, Func<TState, TState> iterate, Func<TState, TResult> resultSelector, IScheduler scheduler)
  19. {
  20. _initialState = initialState;
  21. _condition = condition;
  22. _iterate = iterate;
  23. _resultSelector = resultSelector;
  24. _scheduler = scheduler;
  25. }
  26. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  27. protected override void Run(_ sink) => sink.Run(_scheduler);
  28. internal sealed class _ : IdentitySink<TResult>
  29. {
  30. private readonly Func<TState, bool> _condition;
  31. private readonly Func<TState, TState> _iterate;
  32. private readonly Func<TState, TResult> _resultSelector;
  33. public _(NoTime parent, IObserver<TResult> observer)
  34. : base(observer)
  35. {
  36. _condition = parent._condition;
  37. _iterate = parent._iterate;
  38. _resultSelector = parent._resultSelector;
  39. _state = parent._initialState;
  40. _first = true;
  41. }
  42. private TState _state;
  43. private bool _first;
  44. public void Run(IScheduler _scheduler)
  45. {
  46. var longRunning = _scheduler.AsLongRunning();
  47. if (longRunning != null)
  48. {
  49. SetUpstream(longRunning.ScheduleLongRunning(this, (@this, c) => @this.Loop(c)));
  50. }
  51. else
  52. {
  53. SetUpstream(_scheduler.Schedule(this, (@this, a) => @this.LoopRec(a)));
  54. }
  55. }
  56. private void Loop(ICancelable cancel)
  57. {
  58. while (!cancel.IsDisposed)
  59. {
  60. bool hasResult;
  61. var result = default(TResult);
  62. try
  63. {
  64. if (_first)
  65. {
  66. _first = false;
  67. }
  68. else
  69. {
  70. _state = _iterate(_state);
  71. }
  72. hasResult = _condition(_state);
  73. if (hasResult)
  74. {
  75. result = _resultSelector(_state);
  76. }
  77. }
  78. catch (Exception exception)
  79. {
  80. ForwardOnError(exception);
  81. return;
  82. }
  83. if (hasResult)
  84. {
  85. ForwardOnNext(result);
  86. }
  87. else
  88. {
  89. break;
  90. }
  91. }
  92. if (!cancel.IsDisposed)
  93. {
  94. ForwardOnCompleted();
  95. }
  96. }
  97. private void LoopRec(Action<_> recurse)
  98. {
  99. bool hasResult;
  100. var result = default(TResult);
  101. try
  102. {
  103. if (_first)
  104. {
  105. _first = false;
  106. }
  107. else
  108. {
  109. _state = _iterate(_state);
  110. }
  111. hasResult = _condition(_state);
  112. if (hasResult)
  113. {
  114. result = _resultSelector(_state);
  115. }
  116. }
  117. catch (Exception exception)
  118. {
  119. ForwardOnError(exception);
  120. return;
  121. }
  122. if (hasResult)
  123. {
  124. ForwardOnNext(result);
  125. recurse(this);
  126. }
  127. else
  128. {
  129. ForwardOnCompleted();
  130. }
  131. }
  132. }
  133. }
  134. internal sealed class Absolute : Producer<TResult, Absolute._>
  135. {
  136. private readonly TState _initialState;
  137. private readonly Func<TState, bool> _condition;
  138. private readonly Func<TState, TState> _iterate;
  139. private readonly Func<TState, TResult> _resultSelector;
  140. private readonly Func<TState, DateTimeOffset> _timeSelector;
  141. private readonly IScheduler _scheduler;
  142. public Absolute(TState initialState, Func<TState, bool> condition, Func<TState, TState> iterate, Func<TState, TResult> resultSelector, Func<TState, DateTimeOffset> timeSelector, IScheduler scheduler)
  143. {
  144. _initialState = initialState;
  145. _condition = condition;
  146. _iterate = iterate;
  147. _resultSelector = resultSelector;
  148. _timeSelector = timeSelector;
  149. _scheduler = scheduler;
  150. }
  151. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  152. protected override void Run(_ sink) => sink.Run(_scheduler, _initialState);
  153. internal sealed class _ : IdentitySink<TResult>
  154. {
  155. private readonly Func<TState, bool> _condition;
  156. private readonly Func<TState, TState> _iterate;
  157. private readonly Func<TState, TResult> _resultSelector;
  158. private readonly Func<TState, DateTimeOffset> _timeSelector;
  159. public _(Absolute parent, IObserver<TResult> observer)
  160. : base(observer)
  161. {
  162. _condition = parent._condition;
  163. _iterate = parent._iterate;
  164. _resultSelector = parent._resultSelector;
  165. _timeSelector = parent._timeSelector;
  166. _first = true;
  167. }
  168. private bool _first;
  169. private bool _hasResult;
  170. private TResult _result;
  171. private IDisposable _timerDisposable;
  172. public void Run(IScheduler outerScheduler, TState initialState)
  173. {
  174. var timer = new SingleAssignmentDisposable();
  175. Disposable.TrySetMultiple(ref _timerDisposable, timer);
  176. timer.Disposable = outerScheduler.Schedule((@this: this, initialState), (scheduler, tuple) => [email protected](scheduler, tuple.initialState));
  177. }
  178. protected override void Dispose(bool disposing)
  179. {
  180. Disposable.Dispose(ref _timerDisposable);
  181. base.Dispose(disposing);
  182. }
  183. private IDisposable InvokeRec(IScheduler self, TState state)
  184. {
  185. if (_hasResult)
  186. {
  187. ForwardOnNext(_result);
  188. }
  189. var time = default(DateTimeOffset);
  190. try
  191. {
  192. if (_first)
  193. {
  194. _first = false;
  195. }
  196. else
  197. {
  198. state = _iterate(state);
  199. }
  200. _hasResult = _condition(state);
  201. if (_hasResult)
  202. {
  203. _result = _resultSelector(state);
  204. time = _timeSelector(state);
  205. }
  206. }
  207. catch (Exception exception)
  208. {
  209. ForwardOnError(exception);
  210. return Disposable.Empty;
  211. }
  212. if (!_hasResult)
  213. {
  214. ForwardOnCompleted();
  215. return Disposable.Empty;
  216. }
  217. var timer = new SingleAssignmentDisposable();
  218. Disposable.TrySetMultiple(ref _timerDisposable, timer);
  219. timer.Disposable = self.Schedule((@this: this, state), time, (scheduler, tuple) => [email protected](scheduler, tuple.state));
  220. return Disposable.Empty;
  221. }
  222. }
  223. }
  224. internal sealed class Relative : Producer<TResult, Relative._>
  225. {
  226. private readonly TState _initialState;
  227. private readonly Func<TState, bool> _condition;
  228. private readonly Func<TState, TState> _iterate;
  229. private readonly Func<TState, TResult> _resultSelector;
  230. private readonly Func<TState, TimeSpan> _timeSelector;
  231. private readonly IScheduler _scheduler;
  232. public Relative(TState initialState, Func<TState, bool> condition, Func<TState, TState> iterate, Func<TState, TResult> resultSelector, Func<TState, TimeSpan> timeSelector, IScheduler scheduler)
  233. {
  234. _initialState = initialState;
  235. _condition = condition;
  236. _iterate = iterate;
  237. _resultSelector = resultSelector;
  238. _timeSelector = timeSelector;
  239. _scheduler = scheduler;
  240. }
  241. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  242. protected override void Run(_ sink) => sink.Run(_scheduler, _initialState);
  243. internal sealed class _ : IdentitySink<TResult>
  244. {
  245. private readonly Func<TState, bool> _condition;
  246. private readonly Func<TState, TState> _iterate;
  247. private readonly Func<TState, TResult> _resultSelector;
  248. private readonly Func<TState, TimeSpan> _timeSelector;
  249. public _(Relative parent, IObserver<TResult> observer)
  250. : base(observer)
  251. {
  252. _condition = parent._condition;
  253. _iterate = parent._iterate;
  254. _resultSelector = parent._resultSelector;
  255. _timeSelector = parent._timeSelector;
  256. _first = true;
  257. }
  258. private bool _first;
  259. private bool _hasResult;
  260. private TResult _result;
  261. private IDisposable _timerDisposable;
  262. public void Run(IScheduler outerScheduler, TState initialState)
  263. {
  264. var timer = new SingleAssignmentDisposable();
  265. Disposable.TrySetMultiple(ref _timerDisposable, timer);
  266. timer.Disposable = outerScheduler.Schedule((@this: this, initialState), (scheduler, tuple) => [email protected](scheduler, tuple.initialState));
  267. }
  268. protected override void Dispose(bool disposing)
  269. {
  270. Disposable.Dispose(ref _timerDisposable);
  271. base.Dispose(disposing);
  272. }
  273. private IDisposable InvokeRec(IScheduler self, TState state)
  274. {
  275. if (_hasResult)
  276. {
  277. ForwardOnNext(_result);
  278. }
  279. var time = default(TimeSpan);
  280. try
  281. {
  282. if (_first)
  283. {
  284. _first = false;
  285. }
  286. else
  287. {
  288. state = _iterate(state);
  289. }
  290. _hasResult = _condition(state);
  291. if (_hasResult)
  292. {
  293. _result = _resultSelector(state);
  294. time = _timeSelector(state);
  295. }
  296. }
  297. catch (Exception exception)
  298. {
  299. ForwardOnError(exception);
  300. return Disposable.Empty;
  301. }
  302. if (!_hasResult)
  303. {
  304. ForwardOnCompleted();
  305. return Disposable.Empty;
  306. }
  307. var timer = new SingleAssignmentDisposable();
  308. Disposable.TrySetMultiple(ref _timerDisposable, timer);
  309. timer.Disposable = self.Schedule((@this: this, state), time, (scheduler, tuple) => [email protected](scheduler, tuple.state));
  310. return Disposable.Empty;
  311. }
  312. }
  313. }
  314. }
  315. }