First.cs 5.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166
  1. using System;
  2. using UniRx.Operators;
  3. namespace UniRx.Operators
  4. {
  5. internal class FirstObservable<T> : OperatorObservableBase<T>
  6. {
  7. readonly IObservable<T> source;
  8. readonly bool useDefault;
  9. readonly Func<T, bool> predicate;
  10. public FirstObservable(IObservable<T> source, bool useDefault)
  11. : base(source.IsRequiredSubscribeOnCurrentThread())
  12. {
  13. this.source = source;
  14. this.useDefault = useDefault;
  15. }
  16. public FirstObservable(IObservable<T> source, Func<T, bool> predicate, bool useDefault)
  17. : base(source.IsRequiredSubscribeOnCurrentThread())
  18. {
  19. this.source = source;
  20. this.predicate = predicate;
  21. this.useDefault = useDefault;
  22. }
  23. protected override IDisposable SubscribeCore(IObserver<T> observer, IDisposable cancel)
  24. {
  25. if (predicate == null)
  26. {
  27. return source.Subscribe(new First(this, observer, cancel));
  28. }
  29. else
  30. {
  31. return source.Subscribe(new First_(this, observer, cancel));
  32. }
  33. }
  34. class First : OperatorObserverBase<T, T>
  35. {
  36. readonly FirstObservable<T> parent;
  37. bool notPublished;
  38. public First(FirstObservable<T> parent, IObserver<T> observer, IDisposable cancel) : base(observer, cancel)
  39. {
  40. this.parent = parent;
  41. this.notPublished = true;
  42. }
  43. public override void OnNext(T value)
  44. {
  45. if (notPublished)
  46. {
  47. notPublished = false;
  48. observer.OnNext(value);
  49. try { observer.OnCompleted(); }
  50. finally { Dispose(); }
  51. return;
  52. }
  53. }
  54. public override void OnError(Exception error)
  55. {
  56. try { observer.OnError(error); }
  57. finally { Dispose(); }
  58. }
  59. public override void OnCompleted()
  60. {
  61. if (parent.useDefault)
  62. {
  63. if (notPublished)
  64. {
  65. observer.OnNext(default(T));
  66. }
  67. try { observer.OnCompleted(); }
  68. finally { Dispose(); }
  69. }
  70. else
  71. {
  72. if (notPublished)
  73. {
  74. try { observer.OnError(new InvalidOperationException("sequence is empty")); }
  75. finally { Dispose(); }
  76. }
  77. else
  78. {
  79. try { observer.OnCompleted(); }
  80. finally { Dispose(); }
  81. }
  82. }
  83. }
  84. }
  85. // with predicate
  86. class First_ : OperatorObserverBase<T, T>
  87. {
  88. readonly FirstObservable<T> parent;
  89. bool notPublished;
  90. public First_(FirstObservable<T> parent, IObserver<T> observer, IDisposable cancel) : base(observer, cancel)
  91. {
  92. this.parent = parent;
  93. this.notPublished = true;
  94. }
  95. public override void OnNext(T value)
  96. {
  97. if (notPublished)
  98. {
  99. bool isPassed;
  100. try
  101. {
  102. isPassed = parent.predicate(value);
  103. }
  104. catch (Exception ex)
  105. {
  106. try { observer.OnError(ex); }
  107. finally { Dispose(); }
  108. return;
  109. }
  110. if (isPassed)
  111. {
  112. notPublished = false;
  113. observer.OnNext(value);
  114. try { observer.OnCompleted(); }
  115. finally { Dispose(); }
  116. }
  117. }
  118. }
  119. public override void OnError(Exception error)
  120. {
  121. try { observer.OnError(error); }
  122. finally { Dispose(); }
  123. }
  124. public override void OnCompleted()
  125. {
  126. if (parent.useDefault)
  127. {
  128. if (notPublished)
  129. {
  130. observer.OnNext(default(T));
  131. }
  132. try { observer.OnCompleted(); }
  133. finally { Dispose(); }
  134. }
  135. else
  136. {
  137. if (notPublished)
  138. {
  139. try { observer.OnError(new InvalidOperationException("sequence is empty")); }
  140. finally { Dispose(); }
  141. }
  142. else
  143. {
  144. try { observer.OnCompleted(); }
  145. finally { Dispose(); }
  146. }
  147. }
  148. }
  149. }
  150. }
  151. }