ReplayAsyncSubject.cs 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556
  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. public SequentialReplayAsyncSubject(TimeSpan window)
  31. : base(false, window)
  32. {
  33. }
  34. public SequentialReplayAsyncSubject(TimeSpan window, IAsyncScheduler scheduler)
  35. : base(false, window, scheduler)
  36. {
  37. }
  38. public SequentialReplayAsyncSubject(int bufferSize, TimeSpan window)
  39. : base(false, bufferSize, window)
  40. {
  41. }
  42. public SequentialReplayAsyncSubject(int bufferSize, TimeSpan window, IAsyncScheduler scheduler)
  43. : base(false, bufferSize, window, scheduler)
  44. {
  45. }
  46. }
  47. public sealed class ConcurrentReplayAsyncSubject<T> : ReplayAsyncSubject<T>
  48. {
  49. public ConcurrentReplayAsyncSubject()
  50. : base(true)
  51. {
  52. }
  53. public ConcurrentReplayAsyncSubject(int bufferSize)
  54. : base(true, bufferSize)
  55. {
  56. }
  57. public ConcurrentReplayAsyncSubject(IAsyncScheduler scheduler)
  58. : base(true, scheduler)
  59. {
  60. }
  61. public ConcurrentReplayAsyncSubject(int bufferSize, IAsyncScheduler scheduler)
  62. : base(true, bufferSize, scheduler)
  63. {
  64. }
  65. public ConcurrentReplayAsyncSubject(TimeSpan window)
  66. : base(false, window)
  67. {
  68. }
  69. public ConcurrentReplayAsyncSubject(TimeSpan window, IAsyncScheduler scheduler)
  70. : base(false, window, scheduler)
  71. {
  72. }
  73. public ConcurrentReplayAsyncSubject(int bufferSize, TimeSpan window)
  74. : base(false, bufferSize, window)
  75. {
  76. }
  77. public ConcurrentReplayAsyncSubject(int bufferSize, TimeSpan window, IAsyncScheduler scheduler)
  78. : base(false, bufferSize, window, scheduler)
  79. {
  80. }
  81. }
  82. public abstract class ReplayAsyncSubject<T> : IAsyncSubject<T>
  83. {
  84. protected readonly IAsyncSubject<T> _impl;
  85. public ReplayAsyncSubject(bool concurrent)
  86. : this(concurrent, int.MaxValue)
  87. {
  88. }
  89. public ReplayAsyncSubject(bool concurrent, int bufferSize)
  90. {
  91. if (bufferSize < 0)
  92. throw new ArgumentOutOfRangeException(nameof(bufferSize));
  93. if (bufferSize == 1)
  94. {
  95. _impl = new ReplayOne(concurrent, CreateImmediateObserver);
  96. }
  97. else if (bufferSize == int.MaxValue)
  98. {
  99. _impl = new ReplayAll(concurrent, CreateImmediateObserver);
  100. }
  101. else
  102. {
  103. _impl = new ReplayMany(concurrent, CreateImmediateObserver, bufferSize);
  104. }
  105. }
  106. public ReplayAsyncSubject(bool concurrent, IAsyncScheduler scheduler)
  107. : this(concurrent, int.MaxValue, scheduler)
  108. {
  109. }
  110. public ReplayAsyncSubject(bool concurrent, int bufferSize, IAsyncScheduler scheduler)
  111. {
  112. if (bufferSize < 0)
  113. throw new ArgumentOutOfRangeException(nameof(bufferSize));
  114. if (scheduler == null)
  115. throw new ArgumentNullException(nameof(scheduler));
  116. if (bufferSize == 1)
  117. {
  118. _impl = new ReplayOne(concurrent, o => CreateScheduledObserver(o, scheduler));
  119. }
  120. else if (bufferSize == int.MaxValue)
  121. {
  122. _impl = new ReplayAll(concurrent, o => CreateScheduledObserver(o, scheduler));
  123. }
  124. else
  125. {
  126. _impl = new ReplayMany(concurrent, o => CreateScheduledObserver(o, scheduler), bufferSize);
  127. }
  128. }
  129. public ReplayAsyncSubject(bool concurrent, TimeSpan window)
  130. {
  131. if (window < TimeSpan.Zero)
  132. throw new ArgumentOutOfRangeException(nameof(window));
  133. _impl = new ReplayTime(concurrent, ImmediateAsyncScheduler.Instance, int.MaxValue, window, CreateImmediateObserver);
  134. }
  135. public ReplayAsyncSubject(bool concurrent, TimeSpan window, IAsyncScheduler scheduler)
  136. {
  137. if (window < TimeSpan.Zero)
  138. throw new ArgumentOutOfRangeException(nameof(window));
  139. if (scheduler == null)
  140. throw new ArgumentNullException(nameof(scheduler));
  141. _impl = new ReplayTime(concurrent, scheduler, int.MaxValue, window, o => CreateScheduledObserver(o, scheduler));
  142. }
  143. public ReplayAsyncSubject(bool concurrent, int bufferSize, TimeSpan window)
  144. {
  145. if (window < TimeSpan.Zero)
  146. throw new ArgumentOutOfRangeException(nameof(window));
  147. if (bufferSize < 0)
  148. throw new ArgumentOutOfRangeException(nameof(bufferSize));
  149. _impl = new ReplayTime(concurrent, ImmediateAsyncScheduler.Instance, bufferSize, window, CreateImmediateObserver);
  150. }
  151. public ReplayAsyncSubject(bool concurrent, int bufferSize, TimeSpan window, IAsyncScheduler scheduler)
  152. {
  153. if (window < TimeSpan.Zero)
  154. throw new ArgumentOutOfRangeException(nameof(window));
  155. if (bufferSize < 0)
  156. throw new ArgumentOutOfRangeException(nameof(bufferSize));
  157. if (scheduler == null)
  158. throw new ArgumentNullException(nameof(scheduler));
  159. _impl = new ReplayTime(concurrent, scheduler, bufferSize, window, o => CreateScheduledObserver(o, scheduler));
  160. }
  161. private static IScheduledAsyncObserver<T> CreateImmediateObserver(IAsyncObserver<T> observer) => new FastImmediateAsyncObserver<T>(observer);
  162. private static IScheduledAsyncObserver<T> CreateScheduledObserver(IAsyncObserver<T> observer, IAsyncScheduler scheduler) => new ScheduledAsyncObserver<T>(observer, scheduler);
  163. public Task OnCompletedAsync() => _impl.OnCompletedAsync();
  164. public Task OnErrorAsync(Exception error) => _impl.OnErrorAsync(error ?? throw new ArgumentNullException(nameof(error)));
  165. public Task OnNextAsync(T value) => _impl.OnNextAsync(value);
  166. public Task<IAsyncDisposable> SubscribeAsync(IAsyncObserver<T> observer) => _impl.SubscribeAsync(observer ?? throw new ArgumentNullException(nameof(observer)));
  167. private abstract class ReplayBase : IAsyncSubject<T>
  168. {
  169. private readonly bool _concurrent;
  170. private readonly AsyncLock _lock = new AsyncLock();
  171. private readonly List<IScheduledAsyncObserver<T>> _observers = new List<IScheduledAsyncObserver<T>>(); // TODO: immutable array
  172. private bool _done;
  173. private Exception _error;
  174. public ReplayBase(bool concurrent)
  175. {
  176. _concurrent = concurrent;
  177. }
  178. public async Task OnCompletedAsync()
  179. {
  180. var observers = default(IScheduledAsyncObserver<T>[]);
  181. using (await _lock.LockAsync().ConfigureAwait(false))
  182. {
  183. if (!_done)
  184. {
  185. _done = true;
  186. Trim();
  187. observers = _observers.ToArray();
  188. if (_concurrent)
  189. {
  190. await Task.WhenAll(observers.Select(o => o.OnCompletedAsync())).ConfigureAwait(false);
  191. }
  192. else
  193. {
  194. foreach (var observer in observers)
  195. {
  196. await observer.OnCompletedAsync().ConfigureAwait(false);
  197. }
  198. }
  199. }
  200. }
  201. if (observers != null)
  202. {
  203. await EnsureActive(observers).ConfigureAwait(false);
  204. }
  205. }
  206. public async Task OnErrorAsync(Exception error)
  207. {
  208. var observers = default(IScheduledAsyncObserver<T>[]);
  209. using (await _lock.LockAsync().ConfigureAwait(false))
  210. {
  211. if (!_done)
  212. {
  213. _done = true;
  214. _error = error;
  215. Trim();
  216. observers = _observers.ToArray();
  217. if (_concurrent)
  218. {
  219. await Task.WhenAll(observers.Select(o => o.OnErrorAsync(error))).ConfigureAwait(false);
  220. }
  221. else
  222. {
  223. foreach (var observer in observers)
  224. {
  225. await observer.OnErrorAsync(error).ConfigureAwait(false);
  226. }
  227. }
  228. }
  229. }
  230. if (observers != null)
  231. {
  232. await EnsureActive(observers).ConfigureAwait(false);
  233. }
  234. }
  235. public async Task OnNextAsync(T value)
  236. {
  237. var observers = default(IScheduledAsyncObserver<T>[]);
  238. using (await _lock.LockAsync().ConfigureAwait(false))
  239. {
  240. if (!_done)
  241. {
  242. await NextAsync(value);
  243. Trim();
  244. observers = _observers.ToArray();
  245. if (_concurrent)
  246. {
  247. await Task.WhenAll(observers.Select(o => o.OnNextAsync(value))).ConfigureAwait(false);
  248. }
  249. else
  250. {
  251. foreach (var observer in observers)
  252. {
  253. await observer.OnNextAsync(value).ConfigureAwait(false);
  254. }
  255. }
  256. }
  257. }
  258. if (observers != null)
  259. {
  260. await EnsureActive(observers).ConfigureAwait(false);
  261. }
  262. }
  263. private async Task EnsureActive(IScheduledAsyncObserver<T>[] observers)
  264. {
  265. if (_concurrent)
  266. {
  267. await Task.WhenAll(observers.Select(o => o.EnsureActive())).ConfigureAwait(false);
  268. }
  269. else
  270. {
  271. foreach (var observer in observers)
  272. {
  273. await observer.EnsureActive().ConfigureAwait(false);
  274. }
  275. }
  276. }
  277. public async Task<IAsyncDisposable> SubscribeAsync(IAsyncObserver<T> observer)
  278. {
  279. var res = AsyncDisposable.Nop;
  280. var scheduled = CreateScheduledObserver(observer);
  281. var count = 0;
  282. using (await _lock.LockAsync().ConfigureAwait(false))
  283. {
  284. Trim();
  285. count = await ReplayAsync(scheduled).ConfigureAwait(false);
  286. if (_error != null)
  287. {
  288. count++;
  289. await scheduled.OnErrorAsync(_error).ConfigureAwait(false);
  290. }
  291. else if (_done)
  292. {
  293. count++;
  294. await scheduled.OnCompletedAsync().ConfigureAwait(false);
  295. }
  296. if (!_done)
  297. {
  298. _observers.Add(scheduled);
  299. res = new Subscription(this, scheduled);
  300. }
  301. }
  302. await scheduled.EnsureActive(count).ConfigureAwait(false);
  303. return res;
  304. }
  305. protected abstract IScheduledAsyncObserver<T> CreateScheduledObserver(IAsyncObserver<T> observer);
  306. protected abstract Task NextAsync(T value);
  307. protected abstract Task<int> ReplayAsync(IScheduledAsyncObserver<T> observer);
  308. protected abstract void Trim();
  309. private async Task UnsubscribeAsync(IScheduledAsyncObserver<T> observer)
  310. {
  311. using (await _lock.LockAsync().ConfigureAwait(false))
  312. {
  313. _observers.Remove(observer);
  314. }
  315. }
  316. private sealed class Subscription : IAsyncDisposable
  317. {
  318. private readonly ReplayBase _parent;
  319. private readonly IScheduledAsyncObserver<T> _scheduled;
  320. public Subscription(ReplayBase parent, IScheduledAsyncObserver<T> scheduled)
  321. {
  322. _parent = parent;
  323. _scheduled = scheduled;
  324. }
  325. public Task DisposeAsync() => _parent.UnsubscribeAsync(_scheduled);
  326. }
  327. }
  328. private abstract class ReplayBufferBase : ReplayBase
  329. {
  330. private readonly Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> _createObserver;
  331. public ReplayBufferBase(bool concurrent, Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> createObserver)
  332. : base(concurrent)
  333. {
  334. _createObserver = createObserver;
  335. }
  336. protected override IScheduledAsyncObserver<T> CreateScheduledObserver(IAsyncObserver<T> observer) => _createObserver(observer);
  337. }
  338. private sealed class ReplayOne : ReplayBufferBase
  339. {
  340. private bool _hasValue;
  341. private T _value;
  342. public ReplayOne(bool concurrent, Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> createObserver)
  343. : base(concurrent, createObserver)
  344. {
  345. }
  346. protected override Task NextAsync(T value)
  347. {
  348. _hasValue = true;
  349. _value = value;
  350. return Task.CompletedTask;
  351. }
  352. protected override async Task<int> ReplayAsync(IScheduledAsyncObserver<T> observer)
  353. {
  354. if (_hasValue)
  355. {
  356. await observer.OnNextAsync(_value).ConfigureAwait(false);
  357. return 1;
  358. }
  359. return 0;
  360. }
  361. protected override void Trim() { }
  362. }
  363. private abstract class ReplayManyBase : ReplayBufferBase
  364. {
  365. protected readonly Queue<T> _values = new Queue<T>();
  366. public ReplayManyBase(bool concurrent, Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> createObserver)
  367. : base(concurrent, createObserver)
  368. {
  369. }
  370. protected override Task NextAsync(T value)
  371. {
  372. _values.Enqueue(value);
  373. return Task.CompletedTask;
  374. }
  375. protected override async Task<int> ReplayAsync(IScheduledAsyncObserver<T> observer)
  376. {
  377. var count = _values.Count;
  378. foreach (var value in _values)
  379. {
  380. await observer.OnNextAsync(value).ConfigureAwait(false);
  381. }
  382. return count;
  383. }
  384. }
  385. private sealed class ReplayMany : ReplayManyBase
  386. {
  387. private readonly int _bufferSize;
  388. public ReplayMany(bool concurrent, Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> createObserver, int bufferSize)
  389. : base(concurrent, createObserver)
  390. {
  391. _bufferSize = bufferSize;
  392. }
  393. protected override void Trim()
  394. {
  395. while (_values.Count > _bufferSize)
  396. {
  397. _values.Dequeue();
  398. }
  399. }
  400. }
  401. private sealed class ReplayAll : ReplayManyBase
  402. {
  403. public ReplayAll(bool concurrent, Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> createObserver)
  404. : base(concurrent, createObserver)
  405. {
  406. }
  407. protected override void Trim() { }
  408. }
  409. private sealed class ReplayTime : ReplayBufferBase
  410. {
  411. private readonly IAsyncScheduler _scheduler;
  412. private readonly int _bufferSize;
  413. private readonly TimeSpan _window;
  414. private readonly Queue<Timestamped<T>> _values = new Queue<Timestamped<T>>();
  415. public ReplayTime(bool concurrent, IAsyncScheduler scheduler, int bufferSize, TimeSpan window, Func<IAsyncObserver<T>, IScheduledAsyncObserver<T>> createObserver)
  416. : base(concurrent, createObserver)
  417. {
  418. _scheduler = scheduler;
  419. _bufferSize = bufferSize;
  420. _window = window;
  421. }
  422. protected override Task NextAsync(T value)
  423. {
  424. _values.Enqueue(new Timestamped<T>(value, _scheduler.Now));
  425. return Task.CompletedTask;
  426. }
  427. protected override async Task<int> ReplayAsync(IScheduledAsyncObserver<T> observer)
  428. {
  429. var count = _values.Count;
  430. foreach (var value in _values)
  431. {
  432. await observer.OnNextAsync(value.Value).ConfigureAwait(false);
  433. }
  434. return count;
  435. }
  436. protected override void Trim()
  437. {
  438. while (_values.Count > _bufferSize)
  439. {
  440. _values.Dequeue();
  441. }
  442. var threshold = _scheduler.Now - _window;
  443. while (_values.Count > 0 && _values.Peek().Timestamp < threshold)
  444. {
  445. _values.Dequeue();
  446. }
  447. }
  448. }
  449. }
  450. }