Sink.cs 2.1 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970
  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. #if !NO_PERF
  5. using System.Threading;
  6. namespace System.Reactive
  7. {
  8. /// <summary>
  9. /// Base class for implementation of query operators, providing a lightweight sink that can be disposed to mute the outgoing observer.
  10. /// </summary>
  11. /// <typeparam name="TSource">Type of the resulting sequence's elements.</typeparam>
  12. /// <remarks>Implementations of sinks are responsible to enforce the message grammar on the associated observer. Upon sending a terminal message, a pairing Dispose call should be made to trigger cancellation of related resources and to mute the outgoing observer.</remarks>
  13. internal abstract class Sink<TSource> : IDisposable
  14. {
  15. protected internal volatile IObserver<TSource> _observer;
  16. private IDisposable _cancel;
  17. public Sink(IObserver<TSource> observer, IDisposable cancel)
  18. {
  19. _observer = observer;
  20. _cancel = cancel;
  21. }
  22. public virtual void Dispose()
  23. {
  24. _observer = NopObserver<TSource>.Instance;
  25. var cancel = Interlocked.Exchange(ref _cancel, null);
  26. if (cancel != null)
  27. {
  28. cancel.Dispose();
  29. }
  30. }
  31. public IObserver<TSource> GetForwarder()
  32. {
  33. return new _(this);
  34. }
  35. class _ : IObserver<TSource>
  36. {
  37. private readonly Sink<TSource> _forward;
  38. public _(Sink<TSource> forward)
  39. {
  40. _forward = forward;
  41. }
  42. public void OnNext(TSource value)
  43. {
  44. _forward._observer.OnNext(value);
  45. }
  46. public void OnError(Exception error)
  47. {
  48. _forward._observer.OnError(error);
  49. _forward.Dispose();
  50. }
  51. public void OnCompleted()
  52. {
  53. _forward._observer.OnCompleted();
  54. _forward.Dispose();
  55. }
  56. }
  57. }
  58. }
  59. #endif