Synchronize.cs 1.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657
  1. using System;
  2. using UniRx.Operators;
  3. namespace UniRx.Operators
  4. {
  5. internal class SynchronizeObservable<T> : OperatorObservableBase<T>
  6. {
  7. readonly IObservable<T> source;
  8. readonly object gate;
  9. public SynchronizeObservable(IObservable<T> source, object gate)
  10. : base(source.IsRequiredSubscribeOnCurrentThread())
  11. {
  12. this.source = source;
  13. this.gate = gate;
  14. }
  15. protected override IDisposable SubscribeCore(IObserver<T> observer, IDisposable cancel)
  16. {
  17. return source.Subscribe(new Synchronize(this, observer, cancel));
  18. }
  19. class Synchronize : OperatorObserverBase<T, T>
  20. {
  21. readonly SynchronizeObservable<T> parent;
  22. public Synchronize(SynchronizeObservable<T> parent, IObserver<T> observer, IDisposable cancel) : base(observer, cancel)
  23. {
  24. this.parent = parent;
  25. }
  26. public override void OnNext(T value)
  27. {
  28. lock (parent.gate)
  29. {
  30. base.observer.OnNext(value);
  31. }
  32. }
  33. public override void OnError(Exception error)
  34. {
  35. lock (parent.gate)
  36. {
  37. try { observer.OnError(error); } finally { Dispose(); };
  38. }
  39. }
  40. public override void OnCompleted()
  41. {
  42. lock (parent.gate)
  43. {
  44. try { observer.OnCompleted(); } finally { Dispose(); };
  45. }
  46. }
  47. }
  48. }
  49. }