ReplayAsyncSubject.cs 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428
  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.Threading;
  9. using System.Threading.Tasks;
  10. namespace System.Reactive.Subjects
  11. {
  12. public sealed class SequentialReplayAsyncSubject<T> : ReplayAsyncSubject<T>
  13. {
  14. public SequentialReplayAsyncSubject()
  15. : base(false)
  16. {
  17. }
  18. public SequentialReplayAsyncSubject(int bufferSize)
  19. : base(false, bufferSize)
  20. {
  21. }
  22. public SequentialReplayAsyncSubject(IAsyncScheduler scheduler)
  23. : base(false, scheduler)
  24. {
  25. }
  26. public SequentialReplayAsyncSubject(int bufferSize, IAsyncScheduler scheduler)
  27. : base(false, bufferSize, scheduler)
  28. {
  29. }
  30. }
  31. public sealed class ConcurrentReplayAsyncSubject<T> : ReplayAsyncSubject<T>
  32. {
  33. public ConcurrentReplayAsyncSubject()
  34. : base(true)
  35. {
  36. }
  37. public ConcurrentReplayAsyncSubject(int bufferSize)
  38. : base(true, bufferSize)
  39. {
  40. }
  41. public ConcurrentReplayAsyncSubject(IAsyncScheduler scheduler)
  42. : base(true, scheduler)
  43. {
  44. }
  45. public ConcurrentReplayAsyncSubject(int bufferSize, IAsyncScheduler scheduler)
  46. : base(true, bufferSize, scheduler)
  47. {
  48. }
  49. }
  50. public abstract class ReplayAsyncSubject<T> : IAsyncSubject<T>
  51. {
  52. protected readonly IAsyncSubject<T> _impl;
  53. public ReplayAsyncSubject(bool concurrent)
  54. : this(concurrent, int.MaxValue)
  55. {
  56. }
  57. public ReplayAsyncSubject(bool concurrent, int bufferSize)
  58. {
  59. if (bufferSize < 0)
  60. throw new ArgumentNullException(nameof(bufferSize));
  61. if (bufferSize == 1)
  62. {
  63. _impl = new ReplayOne(concurrent, CreateImmediateObserver);
  64. }
  65. else if (bufferSize == int.MaxValue)
  66. {
  67. _impl = new ReplayAll(concurrent, CreateImmediateObserver);
  68. }
  69. else
  70. {
  71. _impl = new ReplayMany(concurrent, CreateImmediateObserver, bufferSize);
  72. }
  73. }
  74. public ReplayAsyncSubject(bool concurrent, IAsyncScheduler scheduler)
  75. : this(concurrent, int.MaxValue, scheduler)
  76. {
  77. }
  78. public ReplayAsyncSubject(bool concurrent, int bufferSize, IAsyncScheduler scheduler)
  79. {
  80. if (bufferSize < 0)
  81. throw new ArgumentNullException(nameof(bufferSize));
  82. if (scheduler == null)
  83. throw new ArgumentNullException(nameof(scheduler));
  84. if (bufferSize == 1)
  85. {
  86. _impl = new ReplayOne(concurrent, o => CreateScheduledObserver(o, scheduler));
  87. }
  88. else if (bufferSize == int.MaxValue)
  89. {
  90. _impl = new ReplayAll(concurrent, o => CreateScheduledObserver(o, scheduler));
  91. }
  92. else
  93. {
  94. _impl = new ReplayMany(concurrent, o => CreateScheduledObserver(o, scheduler), bufferSize);
  95. }
  96. }
  97. private static IScheduledAsyncObserver<T> CreateImmediateObserver(IAsyncObserver<T> observer) => new FastImmediateAsyncObserver<T>(observer);
  98. private static IScheduledAsyncObserver<T> CreateScheduledObserver(IAsyncObserver<T> observer, IAsyncScheduler scheduler) => new ScheduledAsyncObserver<T>(observer, scheduler);
  99. public Task OnCompletedAsync() => _impl.OnCompletedAsync();
  100. public Task OnErrorAsync(Exception error) => _impl.OnErrorAsync(error ?? throw new ArgumentNullException(nameof(error)));
  101. public Task OnNextAsync(T value) => _impl.OnNextAsync(value);
  102. public Task<IAsyncDisposable> SubscribeAsync(IAsyncObserver<T> observer) => _impl.SubscribeAsync(observer ?? throw new ArgumentNullException(nameof(observer)));
  103. private abstract class ReplayBase : IAsyncSubject<T>
  104. {
  105. private readonly bool _concurrent;
  106. private readonly AsyncLock _lock = new AsyncLock();
  107. private readonly List<IScheduledAsyncObserver<T>> _observers = new List<IScheduledAsyncObserver<T>>(); // TODO: immutable array
  108. private bool _done;
  109. private Exception _error;
  110. public ReplayBase(bool concurrent)
  111. {
  112. _concurrent = concurrent;
  113. }
  114. public async Task OnCompletedAsync()
  115. {
  116. var observers = default(IScheduledAsyncObserver<T>[]);
  117. using (await _lock.LockAsync().ConfigureAwait(false))
  118. {
  119. if (!_done)
  120. {
  121. _done = true;
  122. await TrimAsync().ConfigureAwait(false);
  123. observers = _observers.ToArray();
  124. if (_concurrent)
  125. {
  126. await Task.WhenAll(observers.Select(o => o.OnCompletedAsync())).ConfigureAwait(false);
  127. }
  128. else
  129. {
  130. foreach (var observer in observers)
  131. {
  132. await observer.OnCompletedAsync().ConfigureAwait(false);
  133. }
  134. }
  135. }
  136. }
  137. if (observers != null)
  138. {
  139. await EnsureActive(observers).ConfigureAwait(false);
  140. }
  141. }
  142. public async Task OnErrorAsync(Exception error)
  143. {
  144. var observers = default(IScheduledAsyncObserver<T>[]);
  145. using (await _lock.LockAsync().ConfigureAwait(false))
  146. {
  147. if (!_done)
  148. {
  149. _done = true;
  150. _error = error;
  151. await TrimAsync().ConfigureAwait(false);
  152. observers = _observers.ToArray();
  153. if (_concurrent)
  154. {
  155. await Task.WhenAll(observers.Select(o => o.OnErrorAsync(error))).ConfigureAwait(false);
  156. }
  157. else
  158. {
  159. foreach (var observer in observers)
  160. {
  161. await observer.OnErrorAsync(error).ConfigureAwait(false);
  162. }
  163. }
  164. }
  165. }
  166. if (observers != null)
  167. {
  168. await EnsureActive(observers).ConfigureAwait(false);
  169. }
  170. }
  171. public async Task OnNextAsync(T value)
  172. {
  173. var observers = default(IScheduledAsyncObserver<T>[]);
  174. using (await _lock.LockAsync().ConfigureAwait(false))
  175. {
  176. if (!_done)
  177. {
  178. await NextAsync(value);
  179. await TrimAsync().ConfigureAwait(false);
  180. observers = _observers.ToArray();
  181. if (_concurrent)
  182. {
  183. await Task.WhenAll(observers.Select(o => o.OnNextAsync(value))).ConfigureAwait(false);
  184. }
  185. else
  186. {
  187. foreach (var observer in observers)
  188. {
  189. await observer.OnNextAsync(value).ConfigureAwait(false);
  190. }
  191. }
  192. }
  193. }
  194. if (observers != null)
  195. {
  196. await EnsureActive(observers).ConfigureAwait(false);
  197. }
  198. }
  199. private async Task EnsureActive(IScheduledAsyncObserver<T>[] observers)
  200. {
  201. if (_concurrent)
  202. {
  203. await Task.WhenAll(observers.Select(o => o.EnsureActive())).ConfigureAwait(false);
  204. }
  205. else
  206. {
  207. foreach (var observer in observers)
  208. {
  209. await observer.EnsureActive().ConfigureAwait(false);
  210. }
  211. }
  212. }
  213. public async Task<IAsyncDisposable> SubscribeAsync(IAsyncObserver<T> observer)
  214. {
  215. var res = AsyncDisposable.Nop;
  216. var scheduled = CreateScheduledObserver(observer);
  217. var count = 0;
  218. using (await _lock.LockAsync().ConfigureAwait(false))
  219. {
  220. await TrimAsync().ConfigureAwait(false);
  221. count = await ReplayAsync(scheduled).ConfigureAwait(false);
  222. if (_error != null)
  223. {
  224. count++;
  225. await scheduled.OnErrorAsync(_error).ConfigureAwait(false);
  226. }
  227. else if (_done)
  228. {
  229. count++;
  230. await scheduled.OnCompletedAsync().ConfigureAwait(false);
  231. }
  232. if (!_done)
  233. {
  234. _observers.Add(scheduled);
  235. res = new Subscription(this, scheduled);
  236. }
  237. }
  238. await scheduled.EnsureActive(count).ConfigureAwait(false);
  239. return res;
  240. }
  241. protected abstract IScheduledAsyncObserver<T> CreateScheduledObserver(IAsyncObserver<T> observer);
  242. protected abstract Task NextAsync(T value);
  243. protected abstract Task<int> ReplayAsync(IScheduledAsyncObserver<T> observer);
  244. protected abstract Task TrimAsync();
  245. private async Task UnsubscribeAsync(IScheduledAsyncObserver<T> observer)
  246. {
  247. using (await _lock.LockAsync().ConfigureAwait(false))
  248. {
  249. _observers.Remove(observer);
  250. }
  251. }
  252. private sealed class Subscription : IAsyncDisposable
  253. {
  254. private readonly ReplayBase _parent;
  255. private readonly IScheduledAsyncObserver<T> _scheduled;
  256. public Subscription(ReplayBase parent, IScheduledAsyncObserver<T> scheduled)
  257. {
  258. _parent = parent;
  259. _scheduled = scheduled;
  260. }
  261. public Task DisposeAsync() => _parent.UnsubscribeAsync(_scheduled);
  262. }
  263. }
  264. private abstract class ReplayBufferBase : ReplayBase
  265. {
  266. private readonly Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> _createObserver;
  267. public ReplayBufferBase(bool concurrent, Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> createObserver)
  268. : base(concurrent)
  269. {
  270. _createObserver = createObserver;
  271. }
  272. protected override IScheduledAsyncObserver<T> CreateScheduledObserver(IAsyncObserver<T> observer) => _createObserver(observer);
  273. }
  274. private sealed class ReplayOne : ReplayBufferBase
  275. {
  276. private bool _hasValue;
  277. private T _value;
  278. public ReplayOne(bool concurrent, Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> createObserver)
  279. : base(concurrent, createObserver)
  280. {
  281. }
  282. protected override Task NextAsync(T value)
  283. {
  284. _hasValue = true;
  285. _value = value;
  286. return Task.CompletedTask;
  287. }
  288. protected override async Task<int> ReplayAsync(IScheduledAsyncObserver<T> observer)
  289. {
  290. if (_hasValue)
  291. {
  292. await observer.OnNextAsync(_value).ConfigureAwait(false);
  293. return 1;
  294. }
  295. return 0;
  296. }
  297. protected override Task TrimAsync() => Task.CompletedTask;
  298. }
  299. private abstract class ReplayManyBase : ReplayBufferBase
  300. {
  301. protected readonly Queue<T> _values = new Queue<T>();
  302. public ReplayManyBase(bool concurrent, Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> createObserver)
  303. : base(concurrent, createObserver)
  304. {
  305. }
  306. protected override Task NextAsync(T value)
  307. {
  308. _values.Enqueue(value);
  309. return Task.CompletedTask;
  310. }
  311. protected override async Task<int> ReplayAsync(IScheduledAsyncObserver<T> observer)
  312. {
  313. var count = _values.Count;
  314. foreach (var value in _values)
  315. {
  316. await observer.OnNextAsync(value).ConfigureAwait(false);
  317. }
  318. return count;
  319. }
  320. }
  321. private sealed class ReplayMany : ReplayManyBase
  322. {
  323. private readonly int _bufferSize;
  324. public ReplayMany(bool concurrent, Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> createObserver, int bufferSize)
  325. : base(concurrent, createObserver)
  326. {
  327. _bufferSize = bufferSize;
  328. }
  329. protected override Task TrimAsync()
  330. {
  331. while (_values.Count > _bufferSize)
  332. {
  333. _values.Dequeue();
  334. }
  335. return Task.CompletedTask;
  336. }
  337. }
  338. private sealed class ReplayAll : ReplayManyBase
  339. {
  340. public ReplayAll(bool concurrent, Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> createObserver)
  341. : base(concurrent, createObserver)
  342. {
  343. }
  344. protected override Task TrimAsync() => Task.CompletedTask;
  345. }
  346. }
  347. }