Repeat.cs 2.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899
  1. using System;
  2. namespace UniRx.Operators
  3. {
  4. internal class RepeatObservable<T> : OperatorObservableBase<T>
  5. {
  6. readonly T value;
  7. readonly int? repeatCount;
  8. readonly IScheduler scheduler;
  9. public RepeatObservable(T value, int? repeatCount, IScheduler scheduler)
  10. : base(scheduler == Scheduler.CurrentThread)
  11. {
  12. this.value = value;
  13. this.repeatCount = repeatCount;
  14. this.scheduler = scheduler;
  15. }
  16. protected override IDisposable SubscribeCore(IObserver<T> observer, IDisposable cancel)
  17. {
  18. observer = new Repeat(observer, cancel);
  19. if (repeatCount == null)
  20. {
  21. return scheduler.Schedule((Action self) =>
  22. {
  23. observer.OnNext(value);
  24. self();
  25. });
  26. }
  27. else
  28. {
  29. if (scheduler == Scheduler.Immediate)
  30. {
  31. var count = this.repeatCount.Value;
  32. for (int i = 0; i < count; i++)
  33. {
  34. observer.OnNext(value);
  35. }
  36. observer.OnCompleted();
  37. return Disposable.Empty;
  38. }
  39. else
  40. {
  41. var currentCount = this.repeatCount.Value;
  42. return scheduler.Schedule((Action self) =>
  43. {
  44. if (currentCount > 0)
  45. {
  46. observer.OnNext(value);
  47. currentCount--;
  48. }
  49. if (currentCount == 0)
  50. {
  51. observer.OnCompleted();
  52. return;
  53. }
  54. self();
  55. });
  56. }
  57. }
  58. }
  59. class Repeat : OperatorObserverBase<T, T>
  60. {
  61. public Repeat(IObserver<T> observer, IDisposable cancel)
  62. : base(observer, cancel)
  63. {
  64. }
  65. public override void OnNext(T value)
  66. {
  67. try
  68. {
  69. base.observer.OnNext(value);
  70. }
  71. catch
  72. {
  73. Dispose();
  74. throw;
  75. }
  76. }
  77. public override void OnError(Exception error)
  78. {
  79. try { observer.OnError(error); }
  80. finally { Dispose(); }
  81. }
  82. public override void OnCompleted()
  83. {
  84. try { observer.OnCompleted(); }
  85. finally { Dispose(); }
  86. }
  87. }
  88. }
  89. }