Join.cs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259
  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.Reactive.Disposables;
  6. using System.Threading;
  7. using System.Threading.Tasks;
  8. namespace System.Reactive.Linq
  9. {
  10. partial class AsyncObservable
  11. {
  12. public static IAsyncObservable<TResult> Join<TLeft, TRight, TLeftDuration, TRightDuration, TResult>(this IAsyncObservable<TLeft> left, IAsyncObservable<TRight> right, Func<TLeft, IAsyncObservable<TLeftDuration>> leftDurationSelector, Func<TRight, IAsyncObservable<TRightDuration>> rightDurationSelector, Func<TLeft, TRight, TResult> resultSelector)
  13. {
  14. if (left == null)
  15. throw new ArgumentNullException(nameof(left));
  16. if (right == null)
  17. throw new ArgumentNullException(nameof(right));
  18. if (leftDurationSelector == null)
  19. throw new ArgumentNullException(nameof(leftDurationSelector));
  20. if (rightDurationSelector == null)
  21. throw new ArgumentNullException(nameof(rightDurationSelector));
  22. if (resultSelector == null)
  23. throw new ArgumentNullException(nameof(resultSelector));
  24. return Create<TResult>(async observer =>
  25. {
  26. var subscriptions = new CompositeAsyncDisposable();
  27. var (leftObserver, rightObserver, disposable) = AsyncObserver.Join(observer, subscriptions, leftDurationSelector, rightDurationSelector, resultSelector);
  28. var leftSubscription = await left.SubscribeSafeAsync(leftObserver).ConfigureAwait(false);
  29. await subscriptions.AddAsync(leftSubscription).ConfigureAwait(false);
  30. var rightSubscription = await right.SubscribeSafeAsync(rightObserver).ConfigureAwait(false);
  31. await subscriptions.AddAsync(rightSubscription).ConfigureAwait(false);
  32. return disposable;
  33. });
  34. }
  35. }
  36. partial class AsyncObserver
  37. {
  38. public static (IAsyncObserver<TLeft>, IAsyncObserver<TRight>, IAsyncDisposable) Join<TLeft, TRight, TLeftDuration, TRightDuration, TResult>(IAsyncObserver<TResult> observer, IAsyncDisposable subscriptions, Func<TLeft, IAsyncObservable<TLeftDuration>> leftDurationSelector, Func<TRight, IAsyncObservable<TRightDuration>> rightDurationSelector, Func<TLeft, TRight, TResult> resultSelector)
  39. {
  40. if (observer == null)
  41. throw new ArgumentNullException(nameof(observer));
  42. if (subscriptions == null)
  43. throw new ArgumentNullException(nameof(subscriptions));
  44. if (leftDurationSelector == null)
  45. throw new ArgumentNullException(nameof(leftDurationSelector));
  46. if (rightDurationSelector == null)
  47. throw new ArgumentNullException(nameof(rightDurationSelector));
  48. if (resultSelector == null)
  49. throw new ArgumentNullException(nameof(resultSelector));
  50. var gate = new AsyncLock();
  51. var group = new CompositeAsyncDisposable(subscriptions);
  52. var leftMap = new SortedDictionary<int, TLeft>();
  53. var rightMap = new SortedDictionary<int, TRight>();
  54. var leftDone = false;
  55. var rightDone = false;
  56. var leftId = default(int);
  57. var rightId = default(int);
  58. async Task OnErrorAsync(Exception ex)
  59. {
  60. using (await gate.LockAsync().ConfigureAwait(false))
  61. {
  62. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  63. }
  64. }
  65. var leftObserver =
  66. Create<TLeft>(
  67. async x =>
  68. {
  69. var theLeftId = default(int);
  70. var theRightId = default(int);
  71. using (await gate.LockAsync().ConfigureAwait(false))
  72. {
  73. theLeftId = leftId++;
  74. theRightId = rightId;
  75. leftMap.Add(theLeftId, x);
  76. }
  77. var duration = default(IAsyncObservable<TLeftDuration>);
  78. try
  79. {
  80. duration = leftDurationSelector(x);
  81. }
  82. catch (Exception ex)
  83. {
  84. await OnErrorAsync(ex);
  85. return;
  86. }
  87. var sad = new SingleAssignmentAsyncDisposable();
  88. await group.AddAsync(sad).ConfigureAwait(false);
  89. var durationObserver =
  90. Create<TLeftDuration>(
  91. d => Task.CompletedTask,
  92. OnErrorAsync,
  93. async () =>
  94. {
  95. using (await gate.LockAsync().ConfigureAwait(false))
  96. {
  97. if (leftMap.Remove(theLeftId) && leftMap.Count == 0 && leftDone)
  98. {
  99. await observer.OnCompletedAsync().ConfigureAwait(false);
  100. }
  101. }
  102. await group.RemoveAsync(sad).ConfigureAwait(false);
  103. }
  104. );
  105. var durationSubscription = await duration.FirstOrDefault().SubscribeSafeAsync(durationObserver).ConfigureAwait(false);
  106. await sad.AssignAsync(durationSubscription).ConfigureAwait(false);
  107. using (await gate.LockAsync().ConfigureAwait(false))
  108. {
  109. foreach (var rightValue in rightMap)
  110. {
  111. if (rightValue.Key < theRightId)
  112. {
  113. var result = default(TResult);
  114. try
  115. {
  116. result = resultSelector(x, rightValue.Value);
  117. }
  118. catch (Exception ex)
  119. {
  120. await OnErrorAsync(ex);
  121. return;
  122. }
  123. await observer.OnNextAsync(result).ConfigureAwait(false);
  124. }
  125. }
  126. }
  127. },
  128. OnErrorAsync,
  129. async () =>
  130. {
  131. using (await gate.LockAsync().ConfigureAwait(false))
  132. {
  133. leftDone = true;
  134. if (rightDone || leftMap.Count == 0)
  135. {
  136. await observer.OnCompletedAsync().ConfigureAwait(false);
  137. }
  138. }
  139. }
  140. );
  141. var rightObserver =
  142. Create<TRight>(
  143. async x =>
  144. {
  145. var theLeftId = 0;
  146. var theRightId = 0;
  147. using (await gate.LockAsync().ConfigureAwait(false))
  148. {
  149. theRightId = rightId++;
  150. theLeftId = leftId;
  151. rightMap.Add(theRightId, x);
  152. }
  153. var duration = default(IAsyncObservable<TRightDuration>);
  154. try
  155. {
  156. duration = rightDurationSelector(x);
  157. }
  158. catch (Exception ex)
  159. {
  160. await OnErrorAsync(ex).ConfigureAwait(false);
  161. return;
  162. }
  163. var sad = new SingleAssignmentAsyncDisposable();
  164. await group.AddAsync(sad).ConfigureAwait(false);
  165. var durationObserver =
  166. Create<TRightDuration>(
  167. d => Task.CompletedTask,
  168. OnErrorAsync,
  169. async () =>
  170. {
  171. using (await gate.LockAsync().ConfigureAwait(false))
  172. {
  173. if (rightMap.Remove(theRightId) && rightMap.Count == 0 && rightDone)
  174. {
  175. await observer.OnCompletedAsync().ConfigureAwait(false);
  176. }
  177. }
  178. await group.RemoveAsync(sad).ConfigureAwait(false);
  179. }
  180. );
  181. var durationSubscription = await duration.FirstOrDefault().SubscribeSafeAsync(durationObserver).ConfigureAwait(false);
  182. await sad.AssignAsync(durationSubscription).ConfigureAwait(false);
  183. using (await gate.LockAsync().ConfigureAwait(false))
  184. {
  185. foreach (var leftValue in leftMap)
  186. {
  187. if (leftValue.Key < theLeftId)
  188. {
  189. var result = default(TResult);
  190. try
  191. {
  192. result = resultSelector(leftValue.Value, x);
  193. }
  194. catch (Exception ex)
  195. {
  196. await OnErrorAsync(ex);
  197. return;
  198. }
  199. await observer.OnNextAsync(result).ConfigureAwait(false);
  200. }
  201. }
  202. }
  203. },
  204. OnErrorAsync,
  205. async () =>
  206. {
  207. using (await gate.LockAsync().ConfigureAwait(false))
  208. {
  209. rightDone = true;
  210. if (leftDone || rightMap.Count == 0)
  211. {
  212. await observer.OnCompletedAsync().ConfigureAwait(false);
  213. }
  214. }
  215. }
  216. );
  217. return (leftObserver, rightObserver, group);
  218. }
  219. }
  220. }