CombineLatest.cs 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528
  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.Collections.Generic;
  5. using System.Collections.ObjectModel;
  6. using System.Linq;
  7. using System.Reactive.Disposables;
  8. namespace System.Reactive.Linq.ObservableImpl
  9. {
  10. #region Binary
  11. internal sealed class CombineLatest<TFirst, TSecond, TResult> : Producer<TResult, CombineLatest<TFirst, TSecond, TResult>._>
  12. {
  13. private readonly IObservable<TFirst> _first;
  14. private readonly IObservable<TSecond> _second;
  15. private readonly Func<TFirst, TSecond, TResult> _resultSelector;
  16. public CombineLatest(IObservable<TFirst> first, IObservable<TSecond> second, Func<TFirst, TSecond, TResult> resultSelector)
  17. {
  18. _first = first;
  19. _second = second;
  20. _resultSelector = resultSelector;
  21. }
  22. protected override _ CreateSink(IObserver<TResult> observer) => new _(_resultSelector, observer);
  23. protected override void Run(_ sink) => sink.Run(_first, _second);
  24. internal sealed class _ : IdentitySink<TResult>
  25. {
  26. private readonly Func<TFirst, TSecond, TResult> _resultSelector;
  27. public _(Func<TFirst, TSecond, TResult> resultSelector, IObserver<TResult> observer)
  28. : base(observer)
  29. {
  30. _resultSelector = resultSelector;
  31. }
  32. private object _gate;
  33. private IDisposable _firstDisposable;
  34. private IDisposable _secondDisposable;
  35. public void Run(IObservable<TFirst> first, IObservable<TSecond> second)
  36. {
  37. _gate = new object();
  38. var fstO = new FirstObserver(this);
  39. var sndO = new SecondObserver(this);
  40. fstO.Other = sndO;
  41. sndO.Other = fstO;
  42. Disposable.SetSingle(ref _firstDisposable, first.SubscribeSafe(fstO));
  43. Disposable.SetSingle(ref _secondDisposable, second.SubscribeSafe(sndO));
  44. }
  45. protected override void Dispose(bool disposing)
  46. {
  47. if (disposing)
  48. {
  49. Disposable.TryDispose(ref _firstDisposable);
  50. Disposable.TryDispose(ref _secondDisposable);
  51. }
  52. base.Dispose(disposing);
  53. }
  54. private sealed class FirstObserver : IObserver<TFirst>
  55. {
  56. private readonly _ _parent;
  57. private SecondObserver _other;
  58. public FirstObserver(_ parent)
  59. {
  60. _parent = parent;
  61. }
  62. public SecondObserver Other { set { _other = value; } }
  63. public bool HasValue { get; private set; }
  64. public TFirst Value { get; private set; }
  65. public bool Done { get; private set; }
  66. public void OnNext(TFirst value)
  67. {
  68. lock (_parent._gate)
  69. {
  70. HasValue = true;
  71. Value = value;
  72. if (_other.HasValue)
  73. {
  74. var res = default(TResult);
  75. try
  76. {
  77. res = _parent._resultSelector(value, _other.Value);
  78. }
  79. catch (Exception ex)
  80. {
  81. _parent.ForwardOnError(ex);
  82. return;
  83. }
  84. _parent.ForwardOnNext(res);
  85. }
  86. else if (_other.Done)
  87. {
  88. _parent.ForwardOnCompleted();
  89. }
  90. }
  91. }
  92. public void OnError(Exception error)
  93. {
  94. lock (_parent._gate)
  95. {
  96. _parent.ForwardOnError(error);
  97. }
  98. }
  99. public void OnCompleted()
  100. {
  101. lock (_parent._gate)
  102. {
  103. Done = true;
  104. if (_other.Done)
  105. {
  106. _parent.ForwardOnCompleted();
  107. }
  108. else
  109. {
  110. Disposable.TryDispose(ref _parent._firstDisposable);
  111. }
  112. }
  113. }
  114. }
  115. private sealed class SecondObserver : IObserver<TSecond>
  116. {
  117. private readonly _ _parent;
  118. private FirstObserver _other;
  119. public SecondObserver(_ parent)
  120. {
  121. _parent = parent;
  122. }
  123. public FirstObserver Other { set { _other = value; } }
  124. public bool HasValue { get; private set; }
  125. public TSecond Value { get; private set; }
  126. public bool Done { get; private set; }
  127. public void OnNext(TSecond value)
  128. {
  129. lock (_parent._gate)
  130. {
  131. HasValue = true;
  132. Value = value;
  133. if (_other.HasValue)
  134. {
  135. var res = default(TResult);
  136. try
  137. {
  138. res = _parent._resultSelector(_other.Value, value);
  139. }
  140. catch (Exception ex)
  141. {
  142. _parent.ForwardOnError(ex);
  143. return;
  144. }
  145. _parent.ForwardOnNext(res);
  146. }
  147. else if (_other.Done)
  148. {
  149. _parent.ForwardOnCompleted();
  150. }
  151. }
  152. }
  153. public void OnError(Exception error)
  154. {
  155. lock (_parent._gate)
  156. {
  157. _parent.ForwardOnError(error);
  158. }
  159. }
  160. public void OnCompleted()
  161. {
  162. lock (_parent._gate)
  163. {
  164. Done = true;
  165. if (_other.Done)
  166. {
  167. _parent.ForwardOnCompleted();
  168. }
  169. else
  170. {
  171. Disposable.TryDispose(ref _parent._secondDisposable);
  172. }
  173. }
  174. }
  175. }
  176. }
  177. }
  178. #endregion
  179. #region [3,16]-ary
  180. #region Helpers for n-ary overloads
  181. internal interface ICombineLatest
  182. {
  183. void Next(int index);
  184. void Fail(Exception error);
  185. void Done(int index);
  186. }
  187. internal abstract class CombineLatestSink<TResult> : IdentitySink<TResult>, ICombineLatest
  188. {
  189. protected readonly object _gate;
  190. private bool _hasValueAll;
  191. private readonly bool[] _hasValue;
  192. private readonly bool[] _isDone;
  193. protected CombineLatestSink(int arity, IObserver<TResult> observer)
  194. : base(observer)
  195. {
  196. _gate = new object();
  197. _hasValue = new bool[arity];
  198. _isDone = new bool[arity];
  199. }
  200. public void Next(int index)
  201. {
  202. if (!_hasValueAll)
  203. {
  204. _hasValue[index] = true;
  205. var hasValueAll = true;
  206. foreach (var hasValue in _hasValue)
  207. {
  208. if (!hasValue)
  209. {
  210. hasValueAll = false;
  211. break;
  212. }
  213. }
  214. _hasValueAll = hasValueAll;
  215. }
  216. if (_hasValueAll)
  217. {
  218. var res = default(TResult);
  219. try
  220. {
  221. res = GetResult();
  222. }
  223. catch (Exception ex)
  224. {
  225. ForwardOnError(ex);
  226. return;
  227. }
  228. ForwardOnNext(res);
  229. }
  230. else
  231. {
  232. var allOthersDone = true;
  233. for (var i = 0; i < _isDone.Length; i++)
  234. {
  235. if (i != index && !_isDone[i])
  236. {
  237. allOthersDone = false;
  238. break;
  239. }
  240. }
  241. if (allOthersDone)
  242. {
  243. ForwardOnCompleted();
  244. }
  245. }
  246. }
  247. protected abstract TResult GetResult();
  248. public void Fail(Exception error)
  249. {
  250. ForwardOnError(error);
  251. }
  252. public void Done(int index)
  253. {
  254. _isDone[index] = true;
  255. var allDone = true;
  256. foreach (var isDone in _isDone)
  257. {
  258. if (!isDone)
  259. {
  260. allDone = false;
  261. break;
  262. }
  263. }
  264. if (allDone)
  265. {
  266. ForwardOnCompleted();
  267. }
  268. }
  269. }
  270. internal sealed class CombineLatestObserver<T> : SafeObserver<T>
  271. {
  272. private readonly object _gate;
  273. private readonly ICombineLatest _parent;
  274. private readonly int _index;
  275. private T _value;
  276. public CombineLatestObserver(object gate, ICombineLatest parent, int index)
  277. {
  278. _gate = gate;
  279. _parent = parent;
  280. _index = index;
  281. }
  282. public T Value => _value;
  283. public override void OnNext(T value)
  284. {
  285. lock (_gate)
  286. {
  287. _value = value;
  288. _parent.Next(_index);
  289. }
  290. }
  291. public override void OnError(Exception error)
  292. {
  293. Dispose();
  294. lock (_gate)
  295. {
  296. _parent.Fail(error);
  297. }
  298. }
  299. public override void OnCompleted()
  300. {
  301. Dispose();
  302. lock (_gate)
  303. {
  304. _parent.Done(_index);
  305. }
  306. }
  307. }
  308. #endregion
  309. #endregion
  310. #region N-ary
  311. internal sealed class CombineLatest<TSource, TResult> : Producer<TResult, CombineLatest<TSource, TResult>._>
  312. {
  313. private readonly IEnumerable<IObservable<TSource>> _sources;
  314. private readonly Func<IList<TSource>, TResult> _resultSelector;
  315. public CombineLatest(IEnumerable<IObservable<TSource>> sources, Func<IList<TSource>, TResult> resultSelector)
  316. {
  317. _sources = sources;
  318. _resultSelector = resultSelector;
  319. }
  320. protected override _ CreateSink(IObserver<TResult> observer) => new _(_resultSelector, observer);
  321. protected override void Run(_ sink) => sink.Run(_sources);
  322. internal sealed class _ : IdentitySink<TResult>
  323. {
  324. private readonly Func<IList<TSource>, TResult> _resultSelector;
  325. public _(Func<IList<TSource>, TResult> resultSelector, IObserver<TResult> observer)
  326. : base(observer)
  327. {
  328. _resultSelector = resultSelector;
  329. }
  330. private object _gate;
  331. private bool[] _hasValue;
  332. private bool _hasValueAll;
  333. private List<TSource> _values;
  334. private bool[] _isDone;
  335. private IDisposable[] _subscriptions;
  336. public void Run(IEnumerable<IObservable<TSource>> sources)
  337. {
  338. var srcs = sources.ToArray();
  339. var N = srcs.Length;
  340. _hasValue = new bool[N];
  341. _hasValueAll = false;
  342. _values = new List<TSource>(N);
  343. for (var i = 0; i < N; i++)
  344. {
  345. _values.Add(default);
  346. }
  347. _isDone = new bool[N];
  348. _subscriptions = new IDisposable[N];
  349. _gate = new object();
  350. for (var i = 0; i < N; i++)
  351. {
  352. var j = i;
  353. var o = new SourceObserver(this, j);
  354. _subscriptions[j] = o;
  355. o.SetResource(srcs[j].SubscribeSafe(o));
  356. }
  357. SetUpstream(StableCompositeDisposable.CreateTrusted(_subscriptions));
  358. }
  359. private void OnNext(int index, TSource value)
  360. {
  361. lock (_gate)
  362. {
  363. _values[index] = value;
  364. _hasValue[index] = true;
  365. if (_hasValueAll || (_hasValueAll = _hasValue.All()))
  366. {
  367. var res = default(TResult);
  368. try
  369. {
  370. res = _resultSelector(new ReadOnlyCollection<TSource>(_values));
  371. }
  372. catch (Exception ex)
  373. {
  374. ForwardOnError(ex);
  375. return;
  376. }
  377. ForwardOnNext(res);
  378. }
  379. else if (_isDone.AllExcept(index))
  380. {
  381. ForwardOnCompleted();
  382. }
  383. }
  384. }
  385. private new void OnError(Exception error)
  386. {
  387. lock (_gate)
  388. {
  389. ForwardOnError(error);
  390. }
  391. }
  392. private void OnCompleted(int index)
  393. {
  394. lock (_gate)
  395. {
  396. _isDone[index] = true;
  397. if (_isDone.All())
  398. {
  399. ForwardOnCompleted();
  400. }
  401. else
  402. {
  403. _subscriptions[index].Dispose();
  404. }
  405. }
  406. }
  407. private sealed class SourceObserver : SafeObserver<TSource>
  408. {
  409. private readonly _ _parent;
  410. private readonly int _index;
  411. public SourceObserver(_ parent, int index)
  412. {
  413. _parent = parent;
  414. _index = index;
  415. }
  416. public override void OnNext(TSource value)
  417. {
  418. _parent.OnNext(_index, value);
  419. }
  420. public override void OnError(Exception error)
  421. {
  422. _parent.OnError(error);
  423. }
  424. public override void OnCompleted()
  425. {
  426. _parent.OnCompleted(_index);
  427. }
  428. }
  429. }
  430. }
  431. #endregion
  432. }