RepeatSafe.cs 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138
  1. using System;
  2. using System.Collections;
  3. using System.Collections.Generic;
  4. namespace UniRx.Operators
  5. {
  6. internal class RepeatSafeObservable<T> : OperatorObservableBase<T>
  7. {
  8. readonly IEnumerable<IObservable<T>> sources;
  9. public RepeatSafeObservable(IEnumerable<IObservable<T>> sources, bool isRequiredSubscribeOnCurrentThread)
  10. : base(isRequiredSubscribeOnCurrentThread)
  11. {
  12. this.sources = sources;
  13. }
  14. protected override IDisposable SubscribeCore(IObserver<T> observer, IDisposable cancel)
  15. {
  16. return new RepeatSafe(this, observer, cancel).Run();
  17. }
  18. class RepeatSafe : OperatorObserverBase<T, T>
  19. {
  20. readonly RepeatSafeObservable<T> parent;
  21. readonly object gate = new object();
  22. IEnumerator<IObservable<T>> e;
  23. SerialDisposable subscription;
  24. Action nextSelf;
  25. bool isDisposed;
  26. bool isRunNext;
  27. public RepeatSafe(RepeatSafeObservable<T> parent, IObserver<T> observer, IDisposable cancel) : base(observer, cancel)
  28. {
  29. this.parent = parent;
  30. }
  31. public IDisposable Run()
  32. {
  33. isDisposed = false;
  34. isRunNext = false;
  35. e = parent.sources.GetEnumerator();
  36. subscription = new SerialDisposable();
  37. var schedule = Scheduler.DefaultSchedulers.TailRecursion.Schedule(RecursiveRun);
  38. return StableCompositeDisposable.Create(schedule, subscription, Disposable.Create(() =>
  39. {
  40. lock (gate)
  41. {
  42. isDisposed = true;
  43. e.Dispose();
  44. }
  45. }));
  46. }
  47. void RecursiveRun(Action self)
  48. {
  49. lock (gate)
  50. {
  51. this.nextSelf = self;
  52. if (isDisposed) return;
  53. var current = default(IObservable<T>);
  54. var hasNext = false;
  55. var ex = default(Exception);
  56. try
  57. {
  58. hasNext = e.MoveNext();
  59. if (hasNext)
  60. {
  61. current = e.Current;
  62. if (current == null) throw new InvalidOperationException("sequence is null.");
  63. }
  64. else
  65. {
  66. e.Dispose();
  67. }
  68. }
  69. catch (Exception exception)
  70. {
  71. ex = exception;
  72. e.Dispose();
  73. }
  74. if (ex != null)
  75. {
  76. try { observer.OnError(ex); }
  77. finally { Dispose(); }
  78. return;
  79. }
  80. if (!hasNext)
  81. {
  82. try { observer.OnCompleted(); }
  83. finally { Dispose(); }
  84. return;
  85. }
  86. var source = e.Current;
  87. var d = new SingleAssignmentDisposable();
  88. subscription.Disposable = d;
  89. d.Disposable = source.Subscribe(this);
  90. }
  91. }
  92. public override void OnNext(T value)
  93. {
  94. isRunNext = true;
  95. base.observer.OnNext(value);
  96. }
  97. public override void OnError(Exception error)
  98. {
  99. try { observer.OnError(error); }
  100. finally { Dispose(); }
  101. }
  102. public override void OnCompleted()
  103. {
  104. if (isRunNext && !isDisposed)
  105. {
  106. isRunNext = false;
  107. this.nextSelf();
  108. }
  109. else
  110. {
  111. e.Dispose();
  112. if (!isDisposed)
  113. {
  114. try { observer.OnCompleted(); }
  115. finally { Dispose(); }
  116. }
  117. }
  118. }
  119. }
  120. }
  121. }