DelaySubscription.cs 2.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990
  1. // Licensed to the .NET Foundation under one or more agreements.
  2. // The .NET Foundation licenses this file to you under the Apache 2.0 License.
  3. // See the LICENSE file in the project root for more information.
  4. using System.Reactive.Concurrency;
  5. namespace System.Reactive.Linq.ObservableImpl
  6. {
  7. internal abstract class DelaySubscription<TSource> : Producer<TSource>
  8. {
  9. private readonly IObservable<TSource> _source;
  10. private readonly IScheduler _scheduler;
  11. public DelaySubscription(IObservable<TSource> source, IScheduler scheduler)
  12. {
  13. _source = source;
  14. _scheduler = scheduler;
  15. }
  16. internal sealed class Relative : DelaySubscription<TSource>
  17. {
  18. private readonly TimeSpan _dueTime;
  19. public Relative(IObservable<TSource> source, TimeSpan dueTime, IScheduler scheduler)
  20. : base(source, scheduler)
  21. {
  22. _dueTime = dueTime;
  23. }
  24. protected override IDisposable CreateSink(IObserver<TSource> observer, IDisposable cancel) => new _(observer, cancel);
  25. protected override IDisposable Run(IObserver<TSource> observer, IDisposable cancel, Action<IDisposable> setSink)
  26. {
  27. var sink = new _(observer, cancel);
  28. setSink(sink);
  29. return _scheduler.Schedule(sink, _dueTime, Subscribe);
  30. }
  31. }
  32. internal sealed class Absolute : DelaySubscription<TSource>
  33. {
  34. private readonly DateTimeOffset _dueTime;
  35. public Absolute(IObservable<TSource> source, DateTimeOffset dueTime, IScheduler scheduler)
  36. : base(source, scheduler)
  37. {
  38. _dueTime = dueTime;
  39. }
  40. protected override IDisposable CreateSink(IObserver<TSource> observer, IDisposable cancel) => new _(observer, cancel);
  41. protected override IDisposable Run(IObserver<TSource> observer, IDisposable cancel, Action<IDisposable> setSink)
  42. {
  43. var sink = new _(observer, cancel);
  44. setSink(sink);
  45. return _scheduler.Schedule(sink, _dueTime, Subscribe);
  46. }
  47. }
  48. private IDisposable Subscribe(IScheduler _, _ sink)
  49. {
  50. return _source.SubscribeSafe(sink);
  51. }
  52. private sealed class _ : Sink<TSource>, IObserver<TSource>
  53. {
  54. public _(IObserver<TSource> observer, IDisposable cancel)
  55. : base(observer, cancel)
  56. {
  57. }
  58. public void OnNext(TSource value)
  59. {
  60. base._observer.OnNext(value);
  61. }
  62. public void OnError(Exception error)
  63. {
  64. base._observer.OnError(error);
  65. base.Dispose();
  66. }
  67. public void OnCompleted()
  68. {
  69. base._observer.OnCompleted();
  70. base.Dispose();
  71. }
  72. }
  73. }
  74. }