CatchScheduler.cs 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168
  1. // Licensed to the .NET Foundation under one or more agreements.
  2. // The .NET Foundation licenses this file to you under the MIT License.
  3. // See the LICENSE file in the project root for more information.
  4. using System.Reactive.Disposables;
  5. using System.Runtime.CompilerServices;
  6. namespace System.Reactive.Concurrency
  7. {
  8. internal sealed class CatchScheduler<TException> : SchedulerWrapper
  9. where TException : Exception
  10. {
  11. private readonly Func<TException, bool> _handler;
  12. public CatchScheduler(IScheduler scheduler, Func<TException, bool> handler)
  13. : base(scheduler)
  14. {
  15. _handler = handler;
  16. }
  17. protected override Func<IScheduler, TState, IDisposable> Wrap<TState>(Func<IScheduler, TState, IDisposable> action)
  18. {
  19. return (self, state) =>
  20. {
  21. try
  22. {
  23. return action(GetRecursiveWrapper(self), state);
  24. }
  25. catch (TException exception) when (_handler(exception))
  26. {
  27. return Disposable.Empty;
  28. }
  29. };
  30. }
  31. public CatchScheduler(IScheduler scheduler, Func<TException, bool> handler, ConditionalWeakTable<IScheduler, IScheduler> cache)
  32. : base(scheduler, cache)
  33. {
  34. _handler = handler;
  35. }
  36. protected override SchedulerWrapper Clone(IScheduler scheduler, ConditionalWeakTable<IScheduler, IScheduler> cache)
  37. {
  38. return new CatchScheduler<TException>(scheduler, _handler, cache);
  39. }
  40. protected override bool TryGetService(IServiceProvider provider, Type serviceType, out object? service)
  41. {
  42. service = provider.GetService(serviceType);
  43. if (service != null)
  44. {
  45. if (serviceType == typeof(ISchedulerLongRunning))
  46. {
  47. service = new CatchSchedulerLongRunning((ISchedulerLongRunning)service, _handler);
  48. }
  49. else if (serviceType == typeof(ISchedulerPeriodic))
  50. {
  51. service = new CatchSchedulerPeriodic((ISchedulerPeriodic)service, _handler);
  52. }
  53. }
  54. return true;
  55. }
  56. private class CatchSchedulerLongRunning : ISchedulerLongRunning
  57. {
  58. private readonly ISchedulerLongRunning _scheduler;
  59. private readonly Func<TException, bool> _handler;
  60. public CatchSchedulerLongRunning(ISchedulerLongRunning scheduler, Func<TException, bool> handler)
  61. {
  62. _scheduler = scheduler;
  63. _handler = handler;
  64. }
  65. public IDisposable ScheduleLongRunning<TState>(TState state, Action<TState, ICancelable> action)
  66. {
  67. // Note that avoiding closure allocation here would introduce infinite generic recursion over the TState argument
  68. return _scheduler.ScheduleLongRunning(
  69. state,
  70. (state1, cancel) =>
  71. {
  72. try
  73. {
  74. action(state1, cancel);
  75. }
  76. catch (TException exception) when (_handler(exception))
  77. {
  78. }
  79. });
  80. }
  81. }
  82. private sealed class CatchSchedulerPeriodic : ISchedulerPeriodic
  83. {
  84. private sealed class PeriodicallyScheduledWorkItem<TState> : IDisposable
  85. {
  86. private IDisposable _cancel;
  87. private bool _failed;
  88. private readonly Func<TState, TState> _action;
  89. private readonly CatchSchedulerPeriodic _catchScheduler;
  90. public PeriodicallyScheduledWorkItem(CatchSchedulerPeriodic scheduler, TState state, TimeSpan period, Func<TState, TState> action)
  91. {
  92. _catchScheduler = scheduler;
  93. _action = action;
  94. // Note that avoiding closure allocation here would introduce infinite generic recursion over the TState argument
  95. Disposable.SetSingle(ref _cancel, scheduler._scheduler.SchedulePeriodic(state, period, state1 => Tick(state1).state));
  96. }
  97. public void Dispose()
  98. {
  99. Disposable.TryDispose(ref _cancel);
  100. }
  101. private (PeriodicallyScheduledWorkItem<TState> @this, TState state) Tick(TState state)
  102. {
  103. //
  104. // Cancellation may not be granted immediately; prevent from running user
  105. // code in that case. Periodic schedulers are assumed to introduce some
  106. // degree of concurrency, so we should return from the SchedulePeriodic
  107. // call eventually, allowing the d.Dispose() call in the catch block to
  108. // take effect.
  109. //
  110. if (_failed)
  111. {
  112. return default;
  113. }
  114. try
  115. {
  116. return (this, _action(state));
  117. }
  118. catch (TException exception)
  119. {
  120. _failed = true;
  121. if (!_catchScheduler._handler(exception))
  122. {
  123. throw;
  124. }
  125. Disposable.TryDispose(ref _cancel);
  126. return default;
  127. }
  128. }
  129. }
  130. private readonly ISchedulerPeriodic _scheduler;
  131. private readonly Func<TException, bool> _handler;
  132. public CatchSchedulerPeriodic(ISchedulerPeriodic scheduler, Func<TException, bool> handler)
  133. {
  134. _scheduler = scheduler;
  135. _handler = handler;
  136. }
  137. public IDisposable SchedulePeriodic<TState>(TState state, TimeSpan period, Func<TState, TState> action)
  138. {
  139. return new PeriodicallyScheduledWorkItem<TState>(this, state, period, action);
  140. }
  141. }
  142. }
  143. }