Buffer.cs 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526
  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.Concurrency;
  6. using System.Reactive.Disposables;
  7. using System.Threading;
  8. using System.Threading.Tasks;
  9. namespace System.Reactive.Linq
  10. {
  11. partial class AsyncObservable
  12. {
  13. public static IAsyncObservable<IList<TSource>> Buffer<TSource>(this IAsyncObservable<TSource> source, int count)
  14. {
  15. if (source == null)
  16. throw new ArgumentNullException(nameof(source));
  17. if (count <= 0)
  18. throw new ArgumentNullException(nameof(count));
  19. return Create<IList<TSource>>(observer => source.SubscribeSafeAsync(AsyncObserver.Buffer(observer, count)));
  20. }
  21. public static IAsyncObservable<IList<TSource>> Buffer<TSource>(this IAsyncObservable<TSource> source, int count, int skip)
  22. {
  23. if (source == null)
  24. throw new ArgumentNullException(nameof(source));
  25. if (count <= 0)
  26. throw new ArgumentNullException(nameof(count));
  27. if (skip <= 0)
  28. throw new ArgumentNullException(nameof(skip));
  29. return Create<IList<TSource>>(observer => source.SubscribeSafeAsync(AsyncObserver.Buffer(observer, count, skip)));
  30. }
  31. public static IAsyncObservable<IList<TSource>> Buffer<TSource>(this IAsyncObservable<TSource> source, TimeSpan timeSpan)
  32. {
  33. if (source == null)
  34. throw new ArgumentNullException(nameof(source));
  35. if (timeSpan < TimeSpan.Zero)
  36. throw new ArgumentNullException(nameof(timeSpan));
  37. return Create<IList<TSource>>(async observer =>
  38. {
  39. var (sink, timer) = await AsyncObserver.Buffer(observer, timeSpan).ConfigureAwait(false);
  40. var subscription = await source.SubscribeSafeAsync(sink).ConfigureAwait(false);
  41. return StableCompositeAsyncDisposable.Create(subscription, timer);
  42. });
  43. }
  44. public static IAsyncObservable<IList<TSource>> Buffer<TSource>(this IAsyncObservable<TSource> source, TimeSpan timeSpan, IAsyncScheduler scheduler)
  45. {
  46. if (source == null)
  47. throw new ArgumentNullException(nameof(source));
  48. if (timeSpan < TimeSpan.Zero)
  49. throw new ArgumentNullException(nameof(timeSpan));
  50. if (scheduler == null)
  51. throw new ArgumentNullException(nameof(scheduler));
  52. return Create<IList<TSource>>(async observer =>
  53. {
  54. var (sink, timer) = await AsyncObserver.Buffer(observer, timeSpan, scheduler).ConfigureAwait(false);
  55. var subscription = await source.SubscribeSafeAsync(sink).ConfigureAwait(false);
  56. return StableCompositeAsyncDisposable.Create(subscription, timer);
  57. });
  58. }
  59. public static IAsyncObservable<IList<TSource>> Buffer<TSource>(this IAsyncObservable<TSource> source, TimeSpan timeSpan, TimeSpan timeShift)
  60. {
  61. if (source == null)
  62. throw new ArgumentNullException(nameof(source));
  63. if (timeSpan < TimeSpan.Zero)
  64. throw new ArgumentNullException(nameof(timeSpan));
  65. if (timeShift < TimeSpan.Zero)
  66. throw new ArgumentNullException(nameof(timeShift));
  67. return Create<IList<TSource>>(async observer =>
  68. {
  69. var (sink, timer) = await AsyncObserver.Buffer(observer, timeSpan, timeShift).ConfigureAwait(false);
  70. var subscription = await source.SubscribeSafeAsync(sink).ConfigureAwait(false);
  71. return StableCompositeAsyncDisposable.Create(subscription, timer);
  72. });
  73. }
  74. public static IAsyncObservable<IList<TSource>> Buffer<TSource>(this IAsyncObservable<TSource> source, TimeSpan timeSpan, TimeSpan timeShift, IAsyncScheduler scheduler)
  75. {
  76. if (source == null)
  77. throw new ArgumentNullException(nameof(source));
  78. if (timeSpan < TimeSpan.Zero)
  79. throw new ArgumentNullException(nameof(timeSpan));
  80. if (timeShift < TimeSpan.Zero)
  81. throw new ArgumentNullException(nameof(timeShift));
  82. if (scheduler == null)
  83. throw new ArgumentNullException(nameof(scheduler));
  84. return Create<IList<TSource>>(async observer =>
  85. {
  86. var (sink, timer) = await AsyncObserver.Buffer(observer, timeSpan, timeShift, scheduler).ConfigureAwait(false);
  87. var subscription = await source.SubscribeSafeAsync(sink).ConfigureAwait(false);
  88. return StableCompositeAsyncDisposable.Create(subscription, timer);
  89. });
  90. }
  91. public static IAsyncObservable<IList<TSource>> Buffer<TSource>(this IAsyncObservable<TSource> source, TimeSpan timeSpan, int count)
  92. {
  93. if (source == null)
  94. throw new ArgumentNullException(nameof(source));
  95. if (timeSpan < TimeSpan.Zero)
  96. throw new ArgumentNullException(nameof(timeSpan));
  97. if (count <= 0)
  98. throw new ArgumentNullException(nameof(count));
  99. return Create<IList<TSource>>(async observer =>
  100. {
  101. var (sink, timer) = await AsyncObserver.Buffer(observer, timeSpan, count).ConfigureAwait(false);
  102. var subscription = await source.SubscribeSafeAsync(sink).ConfigureAwait(false);
  103. return StableCompositeAsyncDisposable.Create(subscription, timer);
  104. });
  105. }
  106. public static IAsyncObservable<IList<TSource>> Buffer<TSource>(this IAsyncObservable<TSource> source, TimeSpan timeSpan, int count, IAsyncScheduler scheduler)
  107. {
  108. if (source == null)
  109. throw new ArgumentNullException(nameof(source));
  110. if (timeSpan < TimeSpan.Zero)
  111. throw new ArgumentNullException(nameof(timeSpan));
  112. if (count <= 0)
  113. throw new ArgumentNullException(nameof(count));
  114. if (scheduler == null)
  115. throw new ArgumentNullException(nameof(scheduler));
  116. return Create<IList<TSource>>(async observer =>
  117. {
  118. var (sink, timer) = await AsyncObserver.Buffer(observer, timeSpan, count, scheduler).ConfigureAwait(false);
  119. var subscription = await source.SubscribeSafeAsync(sink).ConfigureAwait(false);
  120. return StableCompositeAsyncDisposable.Create(subscription, timer);
  121. });
  122. }
  123. }
  124. partial class AsyncObserver
  125. {
  126. public static IAsyncObserver<TSource> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, int count)
  127. {
  128. if (observer == null)
  129. throw new ArgumentNullException(nameof(observer));
  130. if (count <= 0)
  131. throw new ArgumentNullException(nameof(count));
  132. return Buffer(observer, count, count);
  133. }
  134. public static IAsyncObserver<TSource> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, int count, int skip)
  135. {
  136. if (observer == null)
  137. throw new ArgumentNullException(nameof(observer));
  138. if (count <= 0)
  139. throw new ArgumentNullException(nameof(count));
  140. if (skip <= 0)
  141. throw new ArgumentNullException(nameof(skip));
  142. var queue = new Queue<IList<TSource>>();
  143. var n = 0;
  144. void CreateBuffer() => queue.Enqueue(new List<TSource>());
  145. CreateBuffer();
  146. return Create<TSource>(
  147. async x =>
  148. {
  149. foreach (var buffer in queue)
  150. {
  151. buffer.Add(x);
  152. }
  153. var c = n - count + 1;
  154. if (c >= 0 && c % skip == 0)
  155. {
  156. var buffer = queue.Dequeue();
  157. if (buffer.Count > 0)
  158. {
  159. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  160. }
  161. }
  162. n++;
  163. if (n % skip == 0)
  164. {
  165. CreateBuffer();
  166. }
  167. },
  168. ex =>
  169. {
  170. while (queue.Count > 0)
  171. {
  172. queue.Dequeue().Clear();
  173. }
  174. return observer.OnErrorAsync(ex);
  175. },
  176. async () =>
  177. {
  178. while (queue.Count > 0)
  179. {
  180. var buffer = queue.Dequeue();
  181. if (buffer.Count > 0)
  182. {
  183. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  184. }
  185. }
  186. await observer.OnCompletedAsync().ConfigureAwait(false);
  187. }
  188. );
  189. }
  190. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan) => Buffer(observer, timeSpan, TaskPoolAsyncScheduler.Default);
  191. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan, IAsyncScheduler scheduler)
  192. {
  193. if (observer == null)
  194. throw new ArgumentNullException(nameof(observer));
  195. if (timeSpan < TimeSpan.Zero)
  196. throw new ArgumentNullException(nameof(timeSpan));
  197. if (scheduler == null)
  198. throw new ArgumentNullException(nameof(scheduler));
  199. return CoreAsync();
  200. async Task<(IAsyncObserver<TSource>, IAsyncDisposable)> CoreAsync()
  201. {
  202. var gate = new AsyncLock();
  203. var buffer = new List<TSource>();
  204. var sink = Create<TSource>(
  205. async x =>
  206. {
  207. using (await gate.LockAsync().ConfigureAwait(false))
  208. {
  209. buffer.Add(x);
  210. }
  211. },
  212. async ex =>
  213. {
  214. using (await gate.LockAsync().ConfigureAwait(false))
  215. {
  216. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  217. }
  218. },
  219. async () =>
  220. {
  221. using (await gate.LockAsync().ConfigureAwait(false))
  222. {
  223. if (buffer.Count > 0)
  224. {
  225. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  226. }
  227. await observer.OnCompletedAsync().ConfigureAwait(false);
  228. }
  229. }
  230. );
  231. var timer = await scheduler.ScheduleAsync(async ct =>
  232. {
  233. while (!ct.IsCancellationRequested)
  234. {
  235. using (await gate.LockAsync().ConfigureAwait(false))
  236. {
  237. if (buffer.Count > 0)
  238. {
  239. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  240. buffer = new List<TSource>();
  241. }
  242. }
  243. await scheduler.Delay(timeSpan, ct).RendezVous(scheduler);
  244. }
  245. }, timeSpan);
  246. return (sink, timer);
  247. };
  248. }
  249. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan, TimeSpan timeShift) => Buffer(observer, timeSpan, timeShift, TaskPoolAsyncScheduler.Default);
  250. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan, TimeSpan timeShift, IAsyncScheduler scheduler)
  251. {
  252. if (observer == null)
  253. throw new ArgumentNullException(nameof(observer));
  254. if (timeSpan < TimeSpan.Zero)
  255. throw new ArgumentNullException(nameof(timeSpan));
  256. if (timeShift < TimeSpan.Zero)
  257. throw new ArgumentNullException(nameof(timeShift));
  258. if (scheduler == null)
  259. throw new ArgumentNullException(nameof(scheduler));
  260. return CoreAsync();
  261. async Task<(IAsyncObserver<TSource>, IAsyncDisposable)> CoreAsync()
  262. {
  263. var gate = new AsyncLock();
  264. var queue = new Queue<List<TSource>>();
  265. queue.Enqueue(new List<TSource>());
  266. var sink = Create<TSource>(
  267. async x =>
  268. {
  269. using (await gate.LockAsync().ConfigureAwait(false))
  270. {
  271. foreach (var buffer in queue)
  272. {
  273. buffer.Add(x);
  274. }
  275. }
  276. },
  277. async ex =>
  278. {
  279. using (await gate.LockAsync().ConfigureAwait(false))
  280. {
  281. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  282. }
  283. },
  284. async () =>
  285. {
  286. using (await gate.LockAsync().ConfigureAwait(false))
  287. {
  288. while (queue.Count > 0)
  289. {
  290. var buffer = queue.Dequeue();
  291. if (buffer.Count > 0)
  292. {
  293. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  294. }
  295. }
  296. await observer.OnCompletedAsync().ConfigureAwait(false);
  297. }
  298. }
  299. );
  300. var nextOpen = timeShift;
  301. var nextClose = timeSpan;
  302. var totalTime = TimeSpan.Zero;
  303. var isOpen = false;
  304. var isClose = false;
  305. TimeSpan GetNextDue()
  306. {
  307. if (nextOpen == nextClose)
  308. {
  309. isOpen = isClose = true;
  310. }
  311. else if (nextClose < nextOpen)
  312. {
  313. isClose = true;
  314. isOpen = false;
  315. }
  316. else
  317. {
  318. isOpen = true;
  319. isClose = false;
  320. }
  321. var newTotalTime = isClose ? nextClose : nextOpen;
  322. var due = newTotalTime - totalTime;
  323. totalTime = newTotalTime;
  324. if (isOpen)
  325. {
  326. nextOpen += timeShift;
  327. }
  328. if (isClose)
  329. {
  330. nextClose += timeShift;
  331. }
  332. return due;
  333. }
  334. var timer = await scheduler.ScheduleAsync(async ct =>
  335. {
  336. while (!ct.IsCancellationRequested)
  337. {
  338. using (await gate.LockAsync().ConfigureAwait(false))
  339. {
  340. if (isClose)
  341. {
  342. var buffer = queue.Dequeue();
  343. if (buffer.Count > 0)
  344. {
  345. await observer.OnNextAsync(buffer).RendezVous(scheduler);
  346. }
  347. }
  348. if (isOpen)
  349. {
  350. queue.Enqueue(new List<TSource>());
  351. }
  352. }
  353. await scheduler.Delay(GetNextDue(), ct).RendezVous(scheduler);
  354. }
  355. }, GetNextDue());
  356. return (sink, timer);
  357. };
  358. }
  359. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan, int count) => Buffer(observer, timeSpan, count, TaskPoolAsyncScheduler.Default);
  360. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan, int count, IAsyncScheduler scheduler)
  361. {
  362. if (observer == null)
  363. throw new ArgumentNullException(nameof(observer));
  364. if (timeSpan < TimeSpan.Zero)
  365. throw new ArgumentNullException(nameof(timeSpan));
  366. if (count <= 0)
  367. throw new ArgumentNullException(nameof(count));
  368. if (scheduler == null)
  369. throw new ArgumentNullException(nameof(scheduler));
  370. return CoreAsync();
  371. async Task<(IAsyncObserver<TSource>, IAsyncDisposable)> CoreAsync()
  372. {
  373. var gate = new AsyncLock();
  374. var timer = new SerialAsyncDisposable();
  375. var buffer = new List<TSource>();
  376. var n = 0;
  377. var id = 0;
  378. async Task CreateTimerAsync(int timerId)
  379. {
  380. var d = await scheduler.ScheduleAsync(async ct =>
  381. {
  382. using (await gate.LockAsync().ConfigureAwait(false))
  383. {
  384. if (timerId == id)
  385. {
  386. await FlushAsync().ConfigureAwait(false);
  387. }
  388. }
  389. }, timeSpan);
  390. await timer.AssignAsync(d).ConfigureAwait(false);
  391. }
  392. async Task FlushAsync()
  393. {
  394. n = 0;
  395. ++id;
  396. var res = buffer;
  397. buffer = new List<TSource>();
  398. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  399. await CreateTimerAsync(id).ConfigureAwait(false);
  400. }
  401. var sink = Create<TSource>(
  402. async x =>
  403. {
  404. using (await gate.LockAsync().ConfigureAwait(false))
  405. {
  406. buffer.Add(x);
  407. if (++n == count)
  408. {
  409. await FlushAsync().ConfigureAwait(false);
  410. }
  411. }
  412. },
  413. async ex =>
  414. {
  415. using (await gate.LockAsync().ConfigureAwait(false))
  416. {
  417. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  418. }
  419. },
  420. async () =>
  421. {
  422. using (await gate.LockAsync().ConfigureAwait(false))
  423. {
  424. await observer.OnNextAsync(buffer).ConfigureAwait(false); // NB: We don't check for non-empty in sync Rx either.
  425. await observer.OnCompletedAsync().ConfigureAwait(false);
  426. }
  427. }
  428. );
  429. await CreateTimerAsync(0).ConfigureAwait(false);
  430. return (sink, timer);
  431. };
  432. }
  433. }
  434. }