QueryLanguage.Creation.cs 9.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265
  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.Linq;
  6. using System.Reactive.Concurrency;
  7. using System.Reactive.Disposables;
  8. using System.Reactive.Threading.Tasks;
  9. using System.Threading;
  10. using System.Threading.Tasks;
  11. namespace System.Reactive.Linq
  12. {
  13. using ObservableImpl;
  14. internal partial class QueryLanguage
  15. {
  16. #region - Create -
  17. public virtual IObservable<TSource> Create<TSource>(Func<IObserver<TSource>, IDisposable> subscribe)
  18. {
  19. return new AnonymousObservable<TSource>(subscribe);
  20. }
  21. public virtual IObservable<TSource> Create<TSource>(Func<IObserver<TSource>, Action> subscribe)
  22. {
  23. return new AnonymousObservable<TSource>(o =>
  24. {
  25. var a = subscribe(o);
  26. return a != null ? Disposable.Create(a) : Disposable.Empty;
  27. });
  28. }
  29. #endregion
  30. #region - CreateAsync -
  31. public virtual IObservable<TResult> Create<TResult>(Func<IObserver<TResult>, CancellationToken, Task> subscribeAsync)
  32. {
  33. return new AnonymousObservable<TResult>(observer =>
  34. {
  35. var cancellable = new CancellationDisposable();
  36. var taskObservable = subscribeAsync(observer, cancellable.Token).ToObservable();
  37. var taskCompletionObserver = new AnonymousObserver<Unit>(Stubs<Unit>.Ignore, observer.OnError, observer.OnCompleted);
  38. var subscription = taskObservable.Subscribe(taskCompletionObserver);
  39. return StableCompositeDisposable.Create(cancellable, subscription);
  40. });
  41. }
  42. public virtual IObservable<TResult> Create<TResult>(Func<IObserver<TResult>, Task> subscribeAsync)
  43. {
  44. return Create<TResult>((observer, token) => subscribeAsync(observer));
  45. }
  46. public virtual IObservable<TResult> Create<TResult>(Func<IObserver<TResult>, CancellationToken, Task<IDisposable>> subscribeAsync)
  47. {
  48. return new AnonymousObservable<TResult>(observer =>
  49. {
  50. var subscription = new SingleAssignmentDisposable();
  51. var cancellable = new CancellationDisposable();
  52. var taskObservable = subscribeAsync(observer, cancellable.Token).ToObservable();
  53. var taskCompletionObserver = new AnonymousObserver<IDisposable>(d => subscription.Disposable = d ?? Disposable.Empty, observer.OnError, Stubs.Nop);
  54. //
  55. // We don't cancel the subscription below *ever* and want to make sure the returned resource gets disposed eventually.
  56. // Notice because we're using the AnonymousObservable<T> type, we get auto-detach behavior for free.
  57. //
  58. taskObservable.Subscribe(taskCompletionObserver);
  59. return StableCompositeDisposable.Create(cancellable, subscription);
  60. });
  61. }
  62. public virtual IObservable<TResult> Create<TResult>(Func<IObserver<TResult>, Task<IDisposable>> subscribeAsync)
  63. {
  64. return Create<TResult>((observer, token) => subscribeAsync(observer));
  65. }
  66. public virtual IObservable<TResult> Create<TResult>(Func<IObserver<TResult>, CancellationToken, Task<Action>> subscribeAsync)
  67. {
  68. return new AnonymousObservable<TResult>(observer =>
  69. {
  70. var subscription = new SingleAssignmentDisposable();
  71. var cancellable = new CancellationDisposable();
  72. var taskObservable = subscribeAsync(observer, cancellable.Token).ToObservable();
  73. var taskCompletionObserver = new AnonymousObserver<Action>(a => subscription.Disposable = a != null ? Disposable.Create(a) : Disposable.Empty, observer.OnError, Stubs.Nop);
  74. //
  75. // We don't cancel the subscription below *ever* and want to make sure the returned resource eventually gets disposed.
  76. // Notice because we're using the AnonymousObservable<T> type, we get auto-detach behavior for free.
  77. //
  78. taskObservable.Subscribe(taskCompletionObserver);
  79. return StableCompositeDisposable.Create(cancellable, subscription);
  80. });
  81. }
  82. public virtual IObservable<TResult> Create<TResult>(Func<IObserver<TResult>, Task<Action>> subscribeAsync)
  83. {
  84. return Create<TResult>((observer, token) => subscribeAsync(observer));
  85. }
  86. #endregion
  87. #region + Defer +
  88. public virtual IObservable<TValue> Defer<TValue>(Func<IObservable<TValue>> observableFactory)
  89. {
  90. return new Defer<TValue>(observableFactory);
  91. }
  92. #endregion
  93. #region + DeferAsync +
  94. public virtual IObservable<TValue> Defer<TValue>(Func<Task<IObservable<TValue>>> observableFactoryAsync)
  95. {
  96. return Defer(() => StartAsync(observableFactoryAsync).Merge());
  97. }
  98. public virtual IObservable<TValue> Defer<TValue>(Func<CancellationToken, Task<IObservable<TValue>>> observableFactoryAsync)
  99. {
  100. return Defer(() => StartAsync(observableFactoryAsync).Merge());
  101. }
  102. #endregion
  103. #region + Empty +
  104. public virtual IObservable<TResult> Empty<TResult>()
  105. {
  106. return new Empty<TResult>(SchedulerDefaults.ConstantTimeOperations);
  107. }
  108. public virtual IObservable<TResult> Empty<TResult>(IScheduler scheduler)
  109. {
  110. return new Empty<TResult>(scheduler);
  111. }
  112. #endregion
  113. #region + Generate +
  114. public virtual IObservable<TResult> Generate<TState, TResult>(TState initialState, Func<TState, bool> condition, Func<TState, TState> iterate, Func<TState, TResult> resultSelector)
  115. {
  116. return new Generate<TState, TResult>.NoTime(initialState, condition, iterate, resultSelector, SchedulerDefaults.Iteration);
  117. }
  118. public virtual IObservable<TResult> Generate<TState, TResult>(TState initialState, Func<TState, bool> condition, Func<TState, TState> iterate, Func<TState, TResult> resultSelector, IScheduler scheduler)
  119. {
  120. return new Generate<TState, TResult>.NoTime(initialState, condition, iterate, resultSelector, scheduler);
  121. }
  122. #endregion
  123. #region + Never +
  124. public virtual IObservable<TResult> Never<TResult>()
  125. {
  126. return new Never<TResult>();
  127. }
  128. #endregion
  129. #region + Range +
  130. public virtual IObservable<int> Range(int start, int count)
  131. {
  132. return Range_(start, count, SchedulerDefaults.Iteration);
  133. }
  134. public virtual IObservable<int> Range(int start, int count, IScheduler scheduler)
  135. {
  136. return Range_(start, count, scheduler);
  137. }
  138. private static IObservable<int> Range_(int start, int count, IScheduler scheduler)
  139. {
  140. return new Range(start, count, scheduler);
  141. }
  142. #endregion
  143. #region + Repeat +
  144. public virtual IObservable<TResult> Repeat<TResult>(TResult value)
  145. {
  146. return new Repeat<TResult>.Forever(value, SchedulerDefaults.Iteration);
  147. }
  148. public virtual IObservable<TResult> Repeat<TResult>(TResult value, IScheduler scheduler)
  149. {
  150. return new Repeat<TResult>.Forever(value, scheduler);
  151. }
  152. public virtual IObservable<TResult> Repeat<TResult>(TResult value, int repeatCount)
  153. {
  154. return new Repeat<TResult>.Count(value, repeatCount, SchedulerDefaults.Iteration);
  155. }
  156. public virtual IObservable<TResult> Repeat<TResult>(TResult value, int repeatCount, IScheduler scheduler)
  157. {
  158. return new Repeat<TResult>.Count(value, repeatCount, scheduler);
  159. }
  160. #endregion
  161. #region + Return +
  162. public virtual IObservable<TResult> Return<TResult>(TResult value)
  163. {
  164. return new Return<TResult>(value, SchedulerDefaults.ConstantTimeOperations);
  165. }
  166. public virtual IObservable<TResult> Return<TResult>(TResult value, IScheduler scheduler)
  167. {
  168. return new Return<TResult>(value, scheduler);
  169. }
  170. #endregion
  171. #region + Throw +
  172. public virtual IObservable<TResult> Throw<TResult>(Exception exception)
  173. {
  174. return new Throw<TResult>(exception, SchedulerDefaults.ConstantTimeOperations);
  175. }
  176. public virtual IObservable<TResult> Throw<TResult>(Exception exception, IScheduler scheduler)
  177. {
  178. return new Throw<TResult>(exception, scheduler);
  179. }
  180. #endregion
  181. #region + Using +
  182. public virtual IObservable<TSource> Using<TSource, TResource>(Func<TResource> resourceFactory, Func<TResource, IObservable<TSource>> observableFactory) where TResource : IDisposable
  183. {
  184. return new Using<TSource, TResource>(resourceFactory, observableFactory);
  185. }
  186. #endregion
  187. #region - UsingAsync -
  188. public virtual IObservable<TSource> Using<TSource, TResource>(Func<CancellationToken, Task<TResource>> resourceFactoryAsync, Func<TResource, CancellationToken, Task<IObservable<TSource>>> observableFactoryAsync) where TResource : IDisposable
  189. {
  190. return Observable.FromAsync<TResource>(resourceFactoryAsync)
  191. .SelectMany(resource =>
  192. Observable.Using<TSource, TResource>(
  193. () => resource,
  194. resource_ => Observable.FromAsync<IObservable<TSource>>(ct => observableFactoryAsync(resource_, ct)).Merge()
  195. )
  196. );
  197. }
  198. #endregion
  199. }
  200. }