Multicast.cs 2.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475
  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.Reactive.Disposables;
  5. using System.Reactive.Subjects;
  6. namespace System.Reactive.Linq.ObservableImpl
  7. {
  8. internal sealed class Multicast<TSource, TIntermediate, TResult> : Producer<TResult, Multicast<TSource, TIntermediate, TResult>._>
  9. {
  10. private readonly IObservable<TSource> _source;
  11. private readonly Func<ISubject<TSource, TIntermediate>> _subjectSelector;
  12. private readonly Func<IObservable<TIntermediate>, IObservable<TResult>> _selector;
  13. public Multicast(IObservable<TSource> source, Func<ISubject<TSource, TIntermediate>> subjectSelector, Func<IObservable<TIntermediate>, IObservable<TResult>> selector)
  14. {
  15. _source = source;
  16. _subjectSelector = subjectSelector;
  17. _selector = selector;
  18. }
  19. protected override _ CreateSink(IObserver<TResult> observer, IDisposable cancel) => new _(observer, cancel);
  20. protected override IDisposable Run(_ sink) => sink.Run(this);
  21. internal sealed class _ : Sink<TResult>, IObserver<TResult>
  22. {
  23. public _(IObserver<TResult> observer, IDisposable cancel)
  24. : base(observer, cancel)
  25. {
  26. }
  27. public IDisposable Run(Multicast<TSource, TIntermediate, TResult> parent)
  28. {
  29. var observable = default(IObservable<TResult>);
  30. var connectable = default(IConnectableObservable<TIntermediate>);
  31. try
  32. {
  33. var subject =parent._subjectSelector();
  34. connectable = new ConnectableObservable<TSource, TIntermediate>(parent._source, subject);
  35. observable = parent._selector(connectable);
  36. }
  37. catch (Exception exception)
  38. {
  39. base._observer.OnError(exception);
  40. base.Dispose();
  41. return Disposable.Empty;
  42. }
  43. var subscription = observable.SubscribeSafe(this);
  44. var connection = connectable.Connect();
  45. return StableCompositeDisposable.Create(subscription, connection);
  46. }
  47. public void OnNext(TResult value)
  48. {
  49. base._observer.OnNext(value);
  50. }
  51. public void OnError(Exception error)
  52. {
  53. base._observer.OnError(error);
  54. base.Dispose();
  55. }
  56. public void OnCompleted()
  57. {
  58. base._observer.OnCompleted();
  59. base.Dispose();
  60. }
  61. }
  62. }
  63. }