Select.cs 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190
  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. using System.Threading.Tasks;
  5. namespace System.Reactive.Linq
  6. {
  7. public partial class AsyncObservable
  8. {
  9. public static IAsyncObservable<TResult> Select<TSource, TResult>(this IAsyncObservable<TSource> source, Func<TSource, TResult> selector)
  10. {
  11. if (source == null)
  12. throw new ArgumentNullException(nameof(source));
  13. if (selector == null)
  14. throw new ArgumentNullException(nameof(selector));
  15. return Create(
  16. source,
  17. selector,
  18. default(TResult),
  19. (source, selector, observer) => source.SubscribeSafeAsync(AsyncObserver.Select(observer, selector)));
  20. }
  21. public static IAsyncObservable<TResult> Select<TSource, TResult>(this IAsyncObservable<TSource> source, Func<TSource, ValueTask<TResult>> selector)
  22. {
  23. if (source == null)
  24. throw new ArgumentNullException(nameof(source));
  25. if (selector == null)
  26. throw new ArgumentNullException(nameof(selector));
  27. return Create(
  28. source,
  29. selector,
  30. default(TResult),
  31. (source, selector, observer) => source.SubscribeSafeAsync(AsyncObserver.Select(observer, selector)));
  32. }
  33. public static IAsyncObservable<TResult> Select<TSource, TResult>(this IAsyncObservable<TSource> source, Func<TSource, int, TResult> selector)
  34. {
  35. if (source == null)
  36. throw new ArgumentNullException(nameof(source));
  37. if (selector == null)
  38. throw new ArgumentNullException(nameof(selector));
  39. return Create(
  40. source,
  41. selector,
  42. default(TResult),
  43. (source, selector, observer) => source.SubscribeSafeAsync(AsyncObserver.Select(observer, selector)));
  44. }
  45. public static IAsyncObservable<TResult> Select<TSource, TResult>(this IAsyncObservable<TSource> source, Func<TSource, int, ValueTask<TResult>> selector)
  46. {
  47. if (source == null)
  48. throw new ArgumentNullException(nameof(source));
  49. if (selector == null)
  50. throw new ArgumentNullException(nameof(selector));
  51. return Create(
  52. source,
  53. selector,
  54. default(TResult),
  55. (source, selector, observer) => source.SubscribeSafeAsync(AsyncObserver.Select(observer, selector)));
  56. }
  57. }
  58. public partial class AsyncObserver
  59. {
  60. public static IAsyncObserver<TSource> Select<TSource, TResult>(IAsyncObserver<TResult> observer, Func<TSource, TResult> selector)
  61. {
  62. if (observer == null)
  63. throw new ArgumentNullException(nameof(observer));
  64. if (selector == null)
  65. throw new ArgumentNullException(nameof(selector));
  66. return Create<TSource>(
  67. async x =>
  68. {
  69. TResult res;
  70. try
  71. {
  72. res = selector(x);
  73. }
  74. catch (Exception ex)
  75. {
  76. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  77. return;
  78. }
  79. await observer.OnNextAsync(res).ConfigureAwait(false);
  80. },
  81. observer.OnErrorAsync,
  82. observer.OnCompletedAsync
  83. );
  84. }
  85. public static IAsyncObserver<TSource> Select<TSource, TResult>(IAsyncObserver<TResult> observer, Func<TSource, ValueTask<TResult>> selector)
  86. {
  87. if (observer == null)
  88. throw new ArgumentNullException(nameof(observer));
  89. if (selector == null)
  90. throw new ArgumentNullException(nameof(selector));
  91. return Create<TSource>(
  92. async x =>
  93. {
  94. TResult res;
  95. try
  96. {
  97. res = await selector(x).ConfigureAwait(false);
  98. }
  99. catch (Exception ex)
  100. {
  101. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  102. return;
  103. }
  104. await observer.OnNextAsync(res).ConfigureAwait(false);
  105. },
  106. observer.OnErrorAsync,
  107. observer.OnCompletedAsync
  108. );
  109. }
  110. public static IAsyncObserver<TSource> Select<TSource, TResult>(IAsyncObserver<TResult> observer, Func<TSource, int, TResult> selector)
  111. {
  112. if (observer == null)
  113. throw new ArgumentNullException(nameof(observer));
  114. if (selector == null)
  115. throw new ArgumentNullException(nameof(selector));
  116. var i = 0;
  117. return Create<TSource>(
  118. async x =>
  119. {
  120. TResult res;
  121. try
  122. {
  123. res = selector(x, checked(i++));
  124. }
  125. catch (Exception ex)
  126. {
  127. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  128. return;
  129. }
  130. await observer.OnNextAsync(res).ConfigureAwait(false);
  131. },
  132. observer.OnErrorAsync,
  133. observer.OnCompletedAsync
  134. );
  135. }
  136. public static IAsyncObserver<TSource> Select<TSource, TResult>(IAsyncObserver<TResult> observer, Func<TSource, int, ValueTask<TResult>> selector)
  137. {
  138. if (observer == null)
  139. throw new ArgumentNullException(nameof(observer));
  140. if (selector == null)
  141. throw new ArgumentNullException(nameof(selector));
  142. var i = 0;
  143. return Create<TSource>(
  144. async x =>
  145. {
  146. TResult res;
  147. try
  148. {
  149. res = await selector(x, checked(i++)).ConfigureAwait(false);
  150. }
  151. catch (Exception ex)
  152. {
  153. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  154. return;
  155. }
  156. await observer.OnNextAsync(res).ConfigureAwait(false);
  157. },
  158. observer.OnErrorAsync,
  159. observer.OnCompletedAsync
  160. );
  161. }
  162. }
  163. }