1
0

GroupedAsyncObservable.cs 1.4 KB

12345678910111213141516171819202122232425262728293031323334353637
  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.Reactive.Disposables;
  5. using System.Reactive.Subjects;
  6. using System.Threading.Tasks;
  7. namespace System.Reactive.Linq
  8. {
  9. internal sealed class GroupedAsyncObservable<TKey, TElement> : AsyncObservableBase<TElement>, IGroupedAsyncObservable<TKey, TElement>
  10. {
  11. private readonly IAsyncSubject<TElement> _subject;
  12. private readonly RefCountAsyncDisposable _disposable;
  13. public GroupedAsyncObservable(TKey key, IAsyncSubject<TElement> subject, RefCountAsyncDisposable disposable)
  14. {
  15. Key = key;
  16. _subject = subject;
  17. _disposable = disposable;
  18. }
  19. public TKey Key { get; }
  20. protected override async ValueTask<IAsyncDisposable> SubscribeAsyncCore(IAsyncObserver<TElement> observer)
  21. {
  22. if (_disposable != null)
  23. {
  24. var d = await _disposable.GetDisposableAsync().ConfigureAwait(false);
  25. var s = await _subject.SubscribeAsync(observer).ConfigureAwait(false);
  26. return StableCompositeAsyncDisposable.Create(d, s);
  27. }
  28. return await _subject.SubscribeAsync(observer).ConfigureAwait(false);
  29. }
  30. }
  31. }