Where.cs 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147
  1. using System;
  2. namespace UniRx.Operators
  3. {
  4. internal class WhereObservable<T> : OperatorObservableBase<T>
  5. {
  6. readonly IObservable<T> source;
  7. readonly Func<T, bool> predicate;
  8. readonly Func<T, int, bool> predicateWithIndex;
  9. public WhereObservable(IObservable<T> source, Func<T, bool> predicate)
  10. : base(source.IsRequiredSubscribeOnCurrentThread())
  11. {
  12. this.source = source;
  13. this.predicate = predicate;
  14. }
  15. public WhereObservable(IObservable<T> source, Func<T, int, bool> predicateWithIndex)
  16. : base(source.IsRequiredSubscribeOnCurrentThread())
  17. {
  18. this.source = source;
  19. this.predicateWithIndex = predicateWithIndex;
  20. }
  21. // Optimize for .Where().Where()
  22. public IObservable<T> CombinePredicate(Func<T, bool> combinePredicate)
  23. {
  24. if (this.predicate != null)
  25. {
  26. return new WhereObservable<T>(source, x => this.predicate(x) && combinePredicate(x));
  27. }
  28. else
  29. {
  30. return new WhereObservable<T>(this, combinePredicate);
  31. }
  32. }
  33. // Optimize for .Where().Select()
  34. public IObservable<TR> CombineSelector<TR>(Func<T, TR> selector)
  35. {
  36. if (this.predicate != null)
  37. {
  38. return new WhereSelectObservable<T, TR>(source, predicate, selector);
  39. }
  40. else
  41. {
  42. return new SelectObservable<T, TR>(this, selector); // can't combine
  43. }
  44. }
  45. protected override IDisposable SubscribeCore(IObserver<T> observer, IDisposable cancel)
  46. {
  47. if (predicate != null)
  48. {
  49. return source.Subscribe(new Where(this, observer, cancel));
  50. }
  51. else
  52. {
  53. return source.Subscribe(new Where_(this, observer, cancel));
  54. }
  55. }
  56. class Where : OperatorObserverBase<T, T>
  57. {
  58. readonly WhereObservable<T> parent;
  59. public Where(WhereObservable<T> parent, IObserver<T> observer, IDisposable cancel)
  60. : base(observer, cancel)
  61. {
  62. this.parent = parent;
  63. }
  64. public override void OnNext(T value)
  65. {
  66. var isPassed = false;
  67. try
  68. {
  69. isPassed = parent.predicate(value);
  70. }
  71. catch (Exception ex)
  72. {
  73. try { observer.OnError(ex); } finally { Dispose(); }
  74. return;
  75. }
  76. if (isPassed)
  77. {
  78. observer.OnNext(value);
  79. }
  80. }
  81. public override void OnError(Exception error)
  82. {
  83. try { observer.OnError(error); } finally { Dispose(); }
  84. }
  85. public override void OnCompleted()
  86. {
  87. try { observer.OnCompleted(); } finally { Dispose(); }
  88. }
  89. }
  90. class Where_ : OperatorObserverBase<T, T>
  91. {
  92. readonly WhereObservable<T> parent;
  93. int index;
  94. public Where_(WhereObservable<T> parent, IObserver<T> observer, IDisposable cancel)
  95. : base(observer, cancel)
  96. {
  97. this.parent = parent;
  98. this.index = 0;
  99. }
  100. public override void OnNext(T value)
  101. {
  102. var isPassed = false;
  103. try
  104. {
  105. isPassed = parent.predicateWithIndex(value, index++);
  106. }
  107. catch (Exception ex)
  108. {
  109. try { observer.OnError(ex); } finally { Dispose(); }
  110. return;
  111. }
  112. if (isPassed)
  113. {
  114. observer.OnNext(value);
  115. }
  116. }
  117. public override void OnError(Exception error)
  118. {
  119. try { observer.OnError(error); } finally { Dispose(); }
  120. }
  121. public override void OnCompleted()
  122. {
  123. try { observer.OnCompleted(); } finally { Dispose(); }
  124. }
  125. }
  126. }
  127. }