Buffer.cs 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732
  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. public static IAsyncObservable<IList<TSource>> Buffer<TSource, TBufferBoundary>(this IAsyncObservable<TSource> source, IAsyncObservable<TBufferBoundary> bufferBoundaries)
  124. {
  125. if (source == null)
  126. throw new ArgumentNullException(nameof(source));
  127. if (bufferBoundaries == null)
  128. throw new ArgumentNullException(nameof(bufferBoundaries));
  129. return Create<IList<TSource>>(async observer =>
  130. {
  131. var (sourceObserver, boundariesObserver) = AsyncObserver.Buffer<TSource, TBufferBoundary>(observer);
  132. var sourceSubscription = await source.SubscribeSafeAsync(sourceObserver).ConfigureAwait(false);
  133. var boundariesSubscription = await bufferBoundaries.SubscribeSafeAsync(boundariesObserver).ConfigureAwait(false);
  134. return StableCompositeAsyncDisposable.Create(sourceSubscription, boundariesSubscription);
  135. });
  136. }
  137. // REVIEW: This overload is inherited from Rx but arguably a bit esoteric as it doesn't provide context to the closing selector.
  138. public static IAsyncObservable<IList<TSource>> Buffer<TSource, TBufferClosing>(this IAsyncObservable<TSource> source, Func<IAsyncObservable<TBufferClosing>> bufferClosingSelector)
  139. {
  140. if (source == null)
  141. throw new ArgumentNullException(nameof(source));
  142. if (bufferClosingSelector == null)
  143. throw new ArgumentNullException(nameof(bufferClosingSelector));
  144. return Create<IList<TSource>>(async observer =>
  145. {
  146. var (sourceObserver, closingDisposable) = await AsyncObserver.Buffer<TSource, TBufferClosing>(observer, bufferClosingSelector).ConfigureAwait(false);
  147. var sourceSubscription = await source.SubscribeSafeAsync(sourceObserver).ConfigureAwait(false);
  148. return StableCompositeAsyncDisposable.Create(sourceSubscription, closingDisposable);
  149. });
  150. }
  151. }
  152. partial class AsyncObserver
  153. {
  154. public static IAsyncObserver<TSource> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, int count)
  155. {
  156. if (observer == null)
  157. throw new ArgumentNullException(nameof(observer));
  158. if (count <= 0)
  159. throw new ArgumentNullException(nameof(count));
  160. return Buffer(observer, count, count);
  161. }
  162. public static IAsyncObserver<TSource> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, int count, int skip)
  163. {
  164. if (observer == null)
  165. throw new ArgumentNullException(nameof(observer));
  166. if (count <= 0)
  167. throw new ArgumentNullException(nameof(count));
  168. if (skip <= 0)
  169. throw new ArgumentNullException(nameof(skip));
  170. var queue = new Queue<IList<TSource>>();
  171. var n = 0;
  172. void CreateBuffer() => queue.Enqueue(new List<TSource>());
  173. CreateBuffer();
  174. return Create<TSource>(
  175. async x =>
  176. {
  177. foreach (var buffer in queue)
  178. {
  179. buffer.Add(x);
  180. }
  181. var c = n - count + 1;
  182. if (c >= 0 && c % skip == 0)
  183. {
  184. var buffer = queue.Dequeue();
  185. if (buffer.Count > 0)
  186. {
  187. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  188. }
  189. }
  190. n++;
  191. if (n % skip == 0)
  192. {
  193. CreateBuffer();
  194. }
  195. },
  196. ex =>
  197. {
  198. while (queue.Count > 0)
  199. {
  200. queue.Dequeue().Clear();
  201. }
  202. return observer.OnErrorAsync(ex);
  203. },
  204. async () =>
  205. {
  206. while (queue.Count > 0)
  207. {
  208. var buffer = queue.Dequeue();
  209. if (buffer.Count > 0)
  210. {
  211. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  212. }
  213. }
  214. await observer.OnCompletedAsync().ConfigureAwait(false);
  215. }
  216. );
  217. }
  218. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan) => Buffer(observer, timeSpan, TaskPoolAsyncScheduler.Default);
  219. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan, IAsyncScheduler scheduler)
  220. {
  221. if (observer == null)
  222. throw new ArgumentNullException(nameof(observer));
  223. if (timeSpan < TimeSpan.Zero)
  224. throw new ArgumentNullException(nameof(timeSpan));
  225. if (scheduler == null)
  226. throw new ArgumentNullException(nameof(scheduler));
  227. return CoreAsync();
  228. async Task<(IAsyncObserver<TSource>, IAsyncDisposable)> CoreAsync()
  229. {
  230. var gate = new AsyncLock();
  231. var buffer = new List<TSource>();
  232. var sink = Create<TSource>(
  233. async x =>
  234. {
  235. using (await gate.LockAsync().ConfigureAwait(false))
  236. {
  237. buffer.Add(x);
  238. }
  239. },
  240. async ex =>
  241. {
  242. using (await gate.LockAsync().ConfigureAwait(false))
  243. {
  244. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  245. }
  246. },
  247. async () =>
  248. {
  249. using (await gate.LockAsync().ConfigureAwait(false))
  250. {
  251. if (buffer.Count > 0)
  252. {
  253. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  254. }
  255. await observer.OnCompletedAsync().ConfigureAwait(false);
  256. }
  257. }
  258. );
  259. var timer = await scheduler.ScheduleAsync(async ct =>
  260. {
  261. while (!ct.IsCancellationRequested)
  262. {
  263. using (await gate.LockAsync().ConfigureAwait(false))
  264. {
  265. if (buffer.Count > 0)
  266. {
  267. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  268. buffer = new List<TSource>();
  269. }
  270. }
  271. await scheduler.Delay(timeSpan, ct).RendezVous(scheduler);
  272. }
  273. }, timeSpan);
  274. return (sink, timer);
  275. };
  276. }
  277. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan, TimeSpan timeShift) => Buffer(observer, timeSpan, timeShift, TaskPoolAsyncScheduler.Default);
  278. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan, TimeSpan timeShift, IAsyncScheduler scheduler)
  279. {
  280. if (observer == null)
  281. throw new ArgumentNullException(nameof(observer));
  282. if (timeSpan < TimeSpan.Zero)
  283. throw new ArgumentNullException(nameof(timeSpan));
  284. if (timeShift < TimeSpan.Zero)
  285. throw new ArgumentNullException(nameof(timeShift));
  286. if (scheduler == null)
  287. throw new ArgumentNullException(nameof(scheduler));
  288. return CoreAsync();
  289. async Task<(IAsyncObserver<TSource>, IAsyncDisposable)> CoreAsync()
  290. {
  291. var gate = new AsyncLock();
  292. var queue = new Queue<List<TSource>>();
  293. queue.Enqueue(new List<TSource>());
  294. var sink = Create<TSource>(
  295. async x =>
  296. {
  297. using (await gate.LockAsync().ConfigureAwait(false))
  298. {
  299. foreach (var buffer in queue)
  300. {
  301. buffer.Add(x);
  302. }
  303. }
  304. },
  305. async ex =>
  306. {
  307. using (await gate.LockAsync().ConfigureAwait(false))
  308. {
  309. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  310. }
  311. },
  312. async () =>
  313. {
  314. using (await gate.LockAsync().ConfigureAwait(false))
  315. {
  316. while (queue.Count > 0)
  317. {
  318. var buffer = queue.Dequeue();
  319. if (buffer.Count > 0)
  320. {
  321. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  322. }
  323. }
  324. await observer.OnCompletedAsync().ConfigureAwait(false);
  325. }
  326. }
  327. );
  328. var nextOpen = timeShift;
  329. var nextClose = timeSpan;
  330. var totalTime = TimeSpan.Zero;
  331. var isOpen = false;
  332. var isClose = false;
  333. TimeSpan GetNextDue()
  334. {
  335. if (nextOpen == nextClose)
  336. {
  337. isOpen = isClose = true;
  338. }
  339. else if (nextClose < nextOpen)
  340. {
  341. isClose = true;
  342. isOpen = false;
  343. }
  344. else
  345. {
  346. isOpen = true;
  347. isClose = false;
  348. }
  349. var newTotalTime = isClose ? nextClose : nextOpen;
  350. var due = newTotalTime - totalTime;
  351. totalTime = newTotalTime;
  352. if (isOpen)
  353. {
  354. nextOpen += timeShift;
  355. }
  356. if (isClose)
  357. {
  358. nextClose += timeShift;
  359. }
  360. return due;
  361. }
  362. var timer = await scheduler.ScheduleAsync(async ct =>
  363. {
  364. while (!ct.IsCancellationRequested)
  365. {
  366. using (await gate.LockAsync().ConfigureAwait(false))
  367. {
  368. if (isClose)
  369. {
  370. var buffer = queue.Dequeue();
  371. if (buffer.Count > 0)
  372. {
  373. await observer.OnNextAsync(buffer).RendezVous(scheduler);
  374. }
  375. }
  376. if (isOpen)
  377. {
  378. queue.Enqueue(new List<TSource>());
  379. }
  380. }
  381. await scheduler.Delay(GetNextDue(), ct).RendezVous(scheduler);
  382. }
  383. }, GetNextDue());
  384. return (sink, timer);
  385. };
  386. }
  387. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan, int count) => Buffer(observer, timeSpan, count, TaskPoolAsyncScheduler.Default);
  388. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource>(IAsyncObserver<IList<TSource>> observer, TimeSpan timeSpan, int count, IAsyncScheduler scheduler)
  389. {
  390. if (observer == null)
  391. throw new ArgumentNullException(nameof(observer));
  392. if (timeSpan < TimeSpan.Zero)
  393. throw new ArgumentNullException(nameof(timeSpan));
  394. if (count <= 0)
  395. throw new ArgumentNullException(nameof(count));
  396. if (scheduler == null)
  397. throw new ArgumentNullException(nameof(scheduler));
  398. return CoreAsync();
  399. async Task<(IAsyncObserver<TSource>, IAsyncDisposable)> CoreAsync()
  400. {
  401. var gate = new AsyncLock();
  402. var timer = new SerialAsyncDisposable();
  403. var buffer = new List<TSource>();
  404. var n = 0;
  405. var id = 0;
  406. async Task CreateTimerAsync(int timerId)
  407. {
  408. var d = await scheduler.ScheduleAsync(async ct =>
  409. {
  410. using (await gate.LockAsync().ConfigureAwait(false))
  411. {
  412. if (timerId == id)
  413. {
  414. await FlushAsync().ConfigureAwait(false);
  415. }
  416. }
  417. }, timeSpan);
  418. await timer.AssignAsync(d).ConfigureAwait(false);
  419. }
  420. async Task FlushAsync()
  421. {
  422. n = 0;
  423. ++id;
  424. var res = buffer;
  425. buffer = new List<TSource>();
  426. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  427. await CreateTimerAsync(id).ConfigureAwait(false);
  428. }
  429. var sink = Create<TSource>(
  430. async x =>
  431. {
  432. using (await gate.LockAsync().ConfigureAwait(false))
  433. {
  434. buffer.Add(x);
  435. if (++n == count)
  436. {
  437. await FlushAsync().ConfigureAwait(false);
  438. }
  439. }
  440. },
  441. async ex =>
  442. {
  443. using (await gate.LockAsync().ConfigureAwait(false))
  444. {
  445. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  446. }
  447. },
  448. async () =>
  449. {
  450. using (await gate.LockAsync().ConfigureAwait(false))
  451. {
  452. await observer.OnNextAsync(buffer).ConfigureAwait(false); // NB: We don't check for non-empty in sync Rx either.
  453. await observer.OnCompletedAsync().ConfigureAwait(false);
  454. }
  455. }
  456. );
  457. await CreateTimerAsync(0).ConfigureAwait(false);
  458. return (sink, timer);
  459. };
  460. }
  461. public static (IAsyncObserver<TSource>, IAsyncObserver<TBufferBoundary>) Buffer<TSource, TBufferBoundary>(IAsyncObserver<IList<TSource>> observer)
  462. {
  463. if (observer == null)
  464. throw new ArgumentNullException(nameof(observer));
  465. var gate = new AsyncLock();
  466. var buffer = new List<TSource>();
  467. return
  468. (
  469. Create<TSource>(
  470. async x =>
  471. {
  472. using (await gate.LockAsync().ConfigureAwait(false))
  473. {
  474. buffer.Add(x);
  475. }
  476. },
  477. async ex =>
  478. {
  479. using (await gate.LockAsync().ConfigureAwait(false))
  480. {
  481. buffer.Clear();
  482. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  483. }
  484. },
  485. async () =>
  486. {
  487. using (await gate.LockAsync().ConfigureAwait(false))
  488. {
  489. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  490. await observer.OnCompletedAsync().ConfigureAwait(false);
  491. }
  492. }
  493. ),
  494. Create<TBufferBoundary>(
  495. async x =>
  496. {
  497. using (await gate.LockAsync().ConfigureAwait(false))
  498. {
  499. var oldBuffer = buffer;
  500. buffer = new List<TSource>();
  501. await observer.OnNextAsync(oldBuffer).ConfigureAwait(false);
  502. }
  503. },
  504. async ex =>
  505. {
  506. using (await gate.LockAsync().ConfigureAwait(false))
  507. {
  508. buffer.Clear();
  509. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  510. }
  511. },
  512. async () =>
  513. {
  514. using (await gate.LockAsync().ConfigureAwait(false))
  515. {
  516. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  517. await observer.OnCompletedAsync().ConfigureAwait(false);
  518. }
  519. }
  520. )
  521. );
  522. }
  523. public static Task<(IAsyncObserver<TSource>, IAsyncDisposable)> Buffer<TSource, TBufferClosing>(IAsyncObserver<IList<TSource>> observer, Func<IAsyncObservable<TBufferClosing>> bufferClosingSelector)
  524. {
  525. if (observer == null)
  526. throw new ArgumentNullException(nameof(observer));
  527. if (bufferClosingSelector == null)
  528. throw new ArgumentNullException(nameof(bufferClosingSelector));
  529. async Task<(IAsyncObserver<TSource>, IAsyncDisposable)> CoreAsync()
  530. {
  531. var closeSubscription = new SerialAsyncDisposable();
  532. var gate = new AsyncLock();
  533. var queueLock = new AsyncQueueLock();
  534. var buffer = new List<TSource>();
  535. async Task CreateBufferCloseAsync()
  536. {
  537. var closing = default(IAsyncObservable<TBufferClosing>);
  538. try
  539. {
  540. closing = bufferClosingSelector(); // REVIEW: Do we need an async variant?
  541. }
  542. catch (Exception ex)
  543. {
  544. using (await gate.LockAsync().ConfigureAwait(false))
  545. {
  546. buffer.Clear();
  547. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  548. }
  549. }
  550. var closingSubscription = new SingleAssignmentAsyncDisposable();
  551. await closeSubscription.AssignAsync(closingSubscription).ConfigureAwait(false);
  552. async Task CloseBufferAsync()
  553. {
  554. await closingSubscription.DisposeAsync().ConfigureAwait(false);
  555. using (await gate.LockAsync().ConfigureAwait(false))
  556. {
  557. var oldBuffer = buffer;
  558. buffer = new List<TSource>();
  559. await observer.OnNextAsync(oldBuffer).ConfigureAwait(false);
  560. }
  561. await queueLock.WaitAsync(CreateBufferCloseAsync).ConfigureAwait(false);
  562. }
  563. var closingObserver =
  564. Create<TBufferClosing>(
  565. x => CloseBufferAsync(),
  566. async ex =>
  567. {
  568. using (await gate.LockAsync().ConfigureAwait(false))
  569. {
  570. buffer.Clear();
  571. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  572. }
  573. },
  574. CloseBufferAsync
  575. );
  576. var closingSubscriptionInner = await closing.SubscribeAsync(closingObserver).ConfigureAwait(false);
  577. await closingSubscription.AssignAsync(closingSubscriptionInner).ConfigureAwait(false);
  578. }
  579. var sink =
  580. Create<TSource>(
  581. async x =>
  582. {
  583. using (await gate.LockAsync().ConfigureAwait(false))
  584. {
  585. buffer.Add(x);
  586. }
  587. },
  588. async ex =>
  589. {
  590. using (await gate.LockAsync().ConfigureAwait(false))
  591. {
  592. buffer.Clear();
  593. await observer.OnErrorAsync(ex).ConfigureAwait(false);
  594. }
  595. },
  596. async () =>
  597. {
  598. using (await gate.LockAsync().ConfigureAwait(false))
  599. {
  600. await observer.OnNextAsync(buffer).ConfigureAwait(false);
  601. await observer.OnCompletedAsync().ConfigureAwait(false);
  602. }
  603. }
  604. );
  605. await queueLock.WaitAsync(CreateBufferCloseAsync).ConfigureAwait(false);
  606. return (sink, closeSubscription);
  607. }
  608. return CoreAsync();
  609. }
  610. }
  611. }