Merge.cs 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390
  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.ObservableImpl
  9. {
  10. internal sealed class Merge<TSource> : Producer<TSource>
  11. {
  12. private readonly IObservable<IObservable<TSource>> _sources;
  13. private readonly IObservable<Task<TSource>> _sourcesT;
  14. private readonly int _maxConcurrent;
  15. public Merge(IObservable<IObservable<TSource>> sources)
  16. {
  17. _sources = sources;
  18. }
  19. public Merge(IObservable<IObservable<TSource>> sources, int maxConcurrent)
  20. {
  21. _sources = sources;
  22. _maxConcurrent = maxConcurrent;
  23. }
  24. public Merge(IObservable<Task<TSource>> sources)
  25. {
  26. _sourcesT = sources;
  27. }
  28. protected override IDisposable Run(IObserver<TSource> observer, IDisposable cancel, Action<IDisposable> setSink)
  29. {
  30. if (_maxConcurrent > 0)
  31. {
  32. var sink = new MergeConcurrent(this, observer, cancel);
  33. setSink(sink);
  34. return sink.Run();
  35. }
  36. else if (_sourcesT != null)
  37. {
  38. var sink = new MergeImpl(this, observer, cancel);
  39. setSink(sink);
  40. return sink.Run();
  41. }
  42. else
  43. {
  44. var sink = new _(this, observer, cancel);
  45. setSink(sink);
  46. return sink.Run();
  47. }
  48. }
  49. class _ : Sink<TSource>, IObserver<IObservable<TSource>>
  50. {
  51. private readonly Merge<TSource> _parent;
  52. public _(Merge<TSource> parent, IObserver<TSource> observer, IDisposable cancel)
  53. : base(observer, cancel)
  54. {
  55. _parent = parent;
  56. }
  57. private object _gate;
  58. private bool _isStopped;
  59. private CompositeDisposable _group;
  60. private SingleAssignmentDisposable _sourceSubscription;
  61. public IDisposable Run()
  62. {
  63. _gate = new object();
  64. _isStopped = false;
  65. _group = new CompositeDisposable();
  66. _sourceSubscription = new SingleAssignmentDisposable();
  67. _group.Add(_sourceSubscription);
  68. _sourceSubscription.Disposable = _parent._sources.SubscribeSafe(this);
  69. return _group;
  70. }
  71. public void OnNext(IObservable<TSource> value)
  72. {
  73. var innerSubscription = new SingleAssignmentDisposable();
  74. _group.Add(innerSubscription);
  75. innerSubscription.Disposable = value.SubscribeSafe(new Iter(this, innerSubscription));
  76. }
  77. public void OnError(Exception error)
  78. {
  79. lock (_gate)
  80. {
  81. base._observer.OnError(error);
  82. base.Dispose();
  83. }
  84. }
  85. public void OnCompleted()
  86. {
  87. _isStopped = true;
  88. if (_group.Count == 1)
  89. {
  90. //
  91. // Notice there can be a race between OnCompleted of the source and any
  92. // of the inner sequences, where both see _group.Count == 1, and one is
  93. // waiting for the lock. There won't be a double OnCompleted observation
  94. // though, because the call to Dispose silences the observer by swapping
  95. // in a NopObserver<T>.
  96. //
  97. lock (_gate)
  98. {
  99. base._observer.OnCompleted();
  100. base.Dispose();
  101. }
  102. }
  103. else
  104. {
  105. _sourceSubscription.Dispose();
  106. }
  107. }
  108. class Iter : IObserver<TSource>
  109. {
  110. private readonly _ _parent;
  111. private readonly IDisposable _self;
  112. public Iter(_ parent, IDisposable self)
  113. {
  114. _parent = parent;
  115. _self = self;
  116. }
  117. public void OnNext(TSource value)
  118. {
  119. lock (_parent._gate)
  120. _parent._observer.OnNext(value);
  121. }
  122. public void OnError(Exception error)
  123. {
  124. lock (_parent._gate)
  125. {
  126. _parent._observer.OnError(error);
  127. _parent.Dispose();
  128. }
  129. }
  130. public void OnCompleted()
  131. {
  132. _parent._group.Remove(_self);
  133. if (_parent._isStopped && _parent._group.Count == 1)
  134. {
  135. //
  136. // Notice there can be a race between OnCompleted of the source and any
  137. // of the inner sequences, where both see _group.Count == 1, and one is
  138. // waiting for the lock. There won't be a double OnCompleted observation
  139. // though, because the call to Dispose silences the observer by swapping
  140. // in a NopObserver<T>.
  141. //
  142. lock (_parent._gate)
  143. {
  144. _parent._observer.OnCompleted();
  145. _parent.Dispose();
  146. }
  147. }
  148. }
  149. }
  150. }
  151. class MergeConcurrent : Sink<TSource>, IObserver<IObservable<TSource>>
  152. {
  153. private readonly Merge<TSource> _parent;
  154. public MergeConcurrent(Merge<TSource> parent, IObserver<TSource> observer, IDisposable cancel)
  155. : base(observer, cancel)
  156. {
  157. _parent = parent;
  158. }
  159. private object _gate;
  160. private Queue<IObservable<TSource>> _q;
  161. private bool _isStopped;
  162. private SingleAssignmentDisposable _sourceSubscription;
  163. private CompositeDisposable _group;
  164. private int _activeCount = 0;
  165. public IDisposable Run()
  166. {
  167. _gate = new object();
  168. _q = new Queue<IObservable<TSource>>();
  169. _isStopped = false;
  170. _activeCount = 0;
  171. _group = new CompositeDisposable();
  172. _sourceSubscription = new SingleAssignmentDisposable();
  173. _sourceSubscription.Disposable = _parent._sources.SubscribeSafe(this);
  174. _group.Add(_sourceSubscription);
  175. return _group;
  176. }
  177. public void OnNext(IObservable<TSource> value)
  178. {
  179. lock (_gate)
  180. {
  181. if (_activeCount < _parent._maxConcurrent)
  182. {
  183. _activeCount++;
  184. Subscribe(value);
  185. }
  186. else
  187. _q.Enqueue(value);
  188. }
  189. }
  190. public void OnError(Exception error)
  191. {
  192. lock (_gate)
  193. {
  194. base._observer.OnError(error);
  195. base.Dispose();
  196. }
  197. }
  198. public void OnCompleted()
  199. {
  200. lock (_gate)
  201. {
  202. _isStopped = true;
  203. if (_activeCount == 0)
  204. {
  205. base._observer.OnCompleted();
  206. base.Dispose();
  207. }
  208. else
  209. {
  210. _sourceSubscription.Dispose();
  211. }
  212. }
  213. }
  214. private void Subscribe(IObservable<TSource> innerSource)
  215. {
  216. var subscription = new SingleAssignmentDisposable();
  217. _group.Add(subscription);
  218. subscription.Disposable = innerSource.SubscribeSafe(new Iter(this, subscription));
  219. }
  220. class Iter : IObserver<TSource>
  221. {
  222. private readonly MergeConcurrent _parent;
  223. private readonly IDisposable _self;
  224. public Iter(MergeConcurrent parent, IDisposable self)
  225. {
  226. _parent = parent;
  227. _self = self;
  228. }
  229. public void OnNext(TSource value)
  230. {
  231. lock (_parent._gate)
  232. _parent._observer.OnNext(value);
  233. }
  234. public void OnError(Exception error)
  235. {
  236. lock (_parent._gate)
  237. {
  238. _parent._observer.OnError(error);
  239. _parent.Dispose();
  240. }
  241. }
  242. public void OnCompleted()
  243. {
  244. _parent._group.Remove(_self);
  245. lock (_parent._gate)
  246. {
  247. if (_parent._q.Count > 0)
  248. {
  249. var s = _parent._q.Dequeue();
  250. _parent.Subscribe(s);
  251. }
  252. else
  253. {
  254. _parent._activeCount--;
  255. if (_parent._isStopped && _parent._activeCount == 0)
  256. {
  257. _parent._observer.OnCompleted();
  258. _parent.Dispose();
  259. }
  260. }
  261. }
  262. }
  263. }
  264. }
  265. class MergeImpl : Sink<TSource>, IObserver<Task<TSource>>
  266. {
  267. private readonly Merge<TSource> _parent;
  268. public MergeImpl(Merge<TSource> parent, IObserver<TSource> observer, IDisposable cancel)
  269. : base(observer, cancel)
  270. {
  271. _parent = parent;
  272. }
  273. private object _gate;
  274. private volatile int _count;
  275. public IDisposable Run()
  276. {
  277. _gate = new object();
  278. _count = 1;
  279. return _parent._sourcesT.SubscribeSafe(this);
  280. }
  281. public void OnNext(Task<TSource> value)
  282. {
  283. Interlocked.Increment(ref _count);
  284. if (value.IsCompleted)
  285. {
  286. OnCompletedTask(value);
  287. }
  288. else
  289. {
  290. value.ContinueWith(OnCompletedTask);
  291. }
  292. }
  293. private void OnCompletedTask(Task<TSource> task)
  294. {
  295. switch (task.Status)
  296. {
  297. case TaskStatus.RanToCompletion:
  298. {
  299. lock (_gate)
  300. base._observer.OnNext(task.Result);
  301. OnCompleted();
  302. }
  303. break;
  304. case TaskStatus.Faulted:
  305. {
  306. lock (_gate)
  307. {
  308. base._observer.OnError(task.Exception.InnerException);
  309. base.Dispose();
  310. }
  311. }
  312. break;
  313. case TaskStatus.Canceled:
  314. {
  315. lock (_gate)
  316. {
  317. base._observer.OnError(new TaskCanceledException(task));
  318. base.Dispose();
  319. }
  320. }
  321. break;
  322. }
  323. }
  324. public void OnError(Exception error)
  325. {
  326. lock (_gate)
  327. {
  328. base._observer.OnError(error);
  329. base.Dispose();
  330. }
  331. }
  332. public void OnCompleted()
  333. {
  334. if (Interlocked.Decrement(ref _count) == 0)
  335. {
  336. lock (_gate)
  337. {
  338. base._observer.OnCompleted();
  339. base.Dispose();
  340. }
  341. }
  342. }
  343. }
  344. }
  345. }