Do.cs 9.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258
  1. using Cysharp.Threading.Tasks;
  2. using Cysharp.Threading.Tasks.Internal;
  3. using Cysharp.Threading.Tasks.Linq;
  4. using System;
  5. using System.Threading;
  6. using System.Threading.Tasks;
  7. namespace Cysharp.Threading.Tasks.Linq
  8. {
  9. public static partial class UniTaskAsyncEnumerable
  10. {
  11. public static IUniTaskAsyncEnumerable<TSource> Do<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Action<TSource> onNext)
  12. {
  13. Error.ThrowArgumentNullException(source, nameof(source));
  14. return source.Do(onNext, null, null);
  15. }
  16. public static IUniTaskAsyncEnumerable<TSource> Do<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Action<TSource> onNext, Action<Exception> onError)
  17. {
  18. Error.ThrowArgumentNullException(source, nameof(source));
  19. return source.Do(onNext, onError, null);
  20. }
  21. public static IUniTaskAsyncEnumerable<TSource> Do<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Action<TSource> onNext, Action onCompleted)
  22. {
  23. Error.ThrowArgumentNullException(source, nameof(source));
  24. return source.Do(onNext, null, onCompleted);
  25. }
  26. public static IUniTaskAsyncEnumerable<TSource> Do<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Action<TSource> onNext, Action<Exception> onError, Action onCompleted)
  27. {
  28. Error.ThrowArgumentNullException(source, nameof(source));
  29. return new Do<TSource>(source, onNext, onError, onCompleted);
  30. }
  31. public static IUniTaskAsyncEnumerable<TSource> Do<TSource>(this IUniTaskAsyncEnumerable<TSource> source, IObserver<TSource> observer)
  32. {
  33. Error.ThrowArgumentNullException(source, nameof(source));
  34. Error.ThrowArgumentNullException(observer, nameof(observer));
  35. return source.Do(observer.OnNext, observer.OnError, observer.OnCompleted); // alloc delegate.
  36. }
  37. // not yet impl.
  38. //public static IUniTaskAsyncEnumerable<TSource> DoAwait<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Func<TSource, UniTask> onNext)
  39. //{
  40. // throw new NotImplementedException();
  41. //}
  42. //public static IUniTaskAsyncEnumerable<TSource> DoAwait<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Func<TSource, UniTask> onNext, Func<Exception, UniTask> onError)
  43. //{
  44. // throw new NotImplementedException();
  45. //}
  46. //public static IUniTaskAsyncEnumerable<TSource> DoAwait<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Func<TSource, UniTask> onNext, Func<UniTask> onCompleted)
  47. //{
  48. // throw new NotImplementedException();
  49. //}
  50. //public static IUniTaskAsyncEnumerable<TSource> DoAwait<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Func<TSource, UniTask> onNext, Func<Exception, UniTask> onError, Func<UniTask> onCompleted)
  51. //{
  52. // throw new NotImplementedException();
  53. //}
  54. //public static IUniTaskAsyncEnumerable<TSource> DoAwaitWithCancellation<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Func<TSource, CancellationToken, UniTask> onNext)
  55. //{
  56. // throw new NotImplementedException();
  57. //}
  58. //public static IUniTaskAsyncEnumerable<TSource> DoAwaitWithCancellation<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Func<TSource, CancellationToken, UniTask> onNext, Func<Exception, CancellationToken, UniTask> onError)
  59. //{
  60. // throw new NotImplementedException();
  61. //}
  62. //public static IUniTaskAsyncEnumerable<TSource> DoAwaitWithCancellation<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Func<TSource, CancellationToken, UniTask> onNext, Func<CancellationToken, UniTask> onCompleted)
  63. //{
  64. // throw new NotImplementedException();
  65. //}
  66. //public static IUniTaskAsyncEnumerable<TSource> DoAwaitWithCancellation<TSource>(this IUniTaskAsyncEnumerable<TSource> source, Func<TSource, CancellationToken, UniTask> onNext, Func<Exception, CancellationToken, UniTask> onError, Func<CancellationToken, UniTask> onCompleted)
  67. //{
  68. // throw new NotImplementedException();
  69. //}
  70. }
  71. internal sealed class Do<TSource> : IUniTaskAsyncEnumerable<TSource>
  72. {
  73. readonly IUniTaskAsyncEnumerable<TSource> source;
  74. readonly Action<TSource> onNext;
  75. readonly Action<Exception> onError;
  76. readonly Action onCompleted;
  77. public Do(IUniTaskAsyncEnumerable<TSource> source, Action<TSource> onNext, Action<Exception> onError, Action onCompleted)
  78. {
  79. this.source = source;
  80. this.onNext = onNext;
  81. this.onError = onError;
  82. this.onCompleted = onCompleted;
  83. }
  84. public IUniTaskAsyncEnumerator<TSource> GetAsyncEnumerator(CancellationToken cancellationToken = default)
  85. {
  86. return new _Do(source, onNext, onError, onCompleted, cancellationToken);
  87. }
  88. sealed class _Do : MoveNextSource, IUniTaskAsyncEnumerator<TSource>
  89. {
  90. static readonly Action<object> MoveNextCoreDelegate = MoveNextCore;
  91. readonly IUniTaskAsyncEnumerable<TSource> source;
  92. readonly Action<TSource> onNext;
  93. readonly Action<Exception> onError;
  94. readonly Action onCompleted;
  95. CancellationToken cancellationToken;
  96. IUniTaskAsyncEnumerator<TSource> enumerator;
  97. UniTask<bool>.Awaiter awaiter;
  98. public _Do(IUniTaskAsyncEnumerable<TSource> source, Action<TSource> onNext, Action<Exception> onError, Action onCompleted, CancellationToken cancellationToken)
  99. {
  100. this.source = source;
  101. this.onNext = onNext;
  102. this.onError = onError;
  103. this.onCompleted = onCompleted;
  104. this.cancellationToken = cancellationToken;
  105. TaskTracker.TrackActiveTask(this, 3);
  106. }
  107. public TSource Current { get; private set; }
  108. public UniTask<bool> MoveNextAsync()
  109. {
  110. cancellationToken.ThrowIfCancellationRequested();
  111. completionSource.Reset();
  112. bool isCompleted = false;
  113. try
  114. {
  115. if (enumerator == null)
  116. {
  117. enumerator = source.GetAsyncEnumerator(cancellationToken);
  118. }
  119. awaiter = enumerator.MoveNextAsync().GetAwaiter();
  120. isCompleted = awaiter.IsCompleted;
  121. }
  122. catch (Exception ex)
  123. {
  124. CallTrySetExceptionAfterNotification(ex);
  125. return new UniTask<bool>(this, completionSource.Version);
  126. }
  127. if (isCompleted)
  128. {
  129. MoveNextCore(this);
  130. }
  131. else
  132. {
  133. awaiter.SourceOnCompleted(MoveNextCoreDelegate, this);
  134. }
  135. return new UniTask<bool>(this, completionSource.Version);
  136. }
  137. void CallTrySetExceptionAfterNotification(Exception ex)
  138. {
  139. if (onError != null)
  140. {
  141. try
  142. {
  143. onError(ex);
  144. }
  145. catch (Exception ex2)
  146. {
  147. completionSource.TrySetException(ex2);
  148. return;
  149. }
  150. }
  151. completionSource.TrySetException(ex);
  152. }
  153. bool TryGetResultWithNotification<T>(UniTask<T>.Awaiter awaiter, out T result)
  154. {
  155. try
  156. {
  157. result = awaiter.GetResult();
  158. return true;
  159. }
  160. catch (Exception ex)
  161. {
  162. CallTrySetExceptionAfterNotification(ex);
  163. result = default;
  164. return false;
  165. }
  166. }
  167. static void MoveNextCore(object state)
  168. {
  169. var self = (_Do)state;
  170. if (self.TryGetResultWithNotification(self.awaiter, out var result))
  171. {
  172. if (result)
  173. {
  174. var v = self.enumerator.Current;
  175. if (self.onNext != null)
  176. {
  177. try
  178. {
  179. self.onNext(v);
  180. }
  181. catch (Exception ex)
  182. {
  183. self.CallTrySetExceptionAfterNotification(ex);
  184. }
  185. }
  186. self.Current = v;
  187. self.completionSource.TrySetResult(true);
  188. }
  189. else
  190. {
  191. if (self.onCompleted != null)
  192. {
  193. try
  194. {
  195. self.onCompleted();
  196. }
  197. catch (Exception ex)
  198. {
  199. self.CallTrySetExceptionAfterNotification(ex);
  200. return;
  201. }
  202. }
  203. self.completionSource.TrySetResult(false);
  204. }
  205. }
  206. }
  207. public UniTask DisposeAsync()
  208. {
  209. TaskTracker.RemoveTracking(this);
  210. if (enumerator != null)
  211. {
  212. return enumerator.DisposeAsync();
  213. }
  214. return default;
  215. }
  216. }
  217. }
  218. }