ReplayAsyncSubject.cs 16 KB

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