Timeout.cs 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409
  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. using System.Reactive.Disposables;
  6. using System.Threading;
  7. namespace System.Reactive.Linq.ObservableImpl
  8. {
  9. internal static class Timeout<TSource>
  10. {
  11. internal sealed class Relative : Producer<TSource, Relative._>
  12. {
  13. private readonly IObservable<TSource> _source;
  14. private readonly TimeSpan _dueTime;
  15. private readonly IObservable<TSource> _other;
  16. private readonly IScheduler _scheduler;
  17. public Relative(IObservable<TSource> source, TimeSpan dueTime, IObservable<TSource> other, IScheduler scheduler)
  18. {
  19. _source = source;
  20. _dueTime = dueTime;
  21. _other = other;
  22. _scheduler = scheduler;
  23. }
  24. protected override _ CreateSink(IObserver<TSource> observer, IDisposable cancel) => new _(this, observer, cancel);
  25. protected override IDisposable Run(_ sink) => sink.Run(_source);
  26. internal sealed class _ : IdentitySink<TSource>
  27. {
  28. private readonly TimeSpan _dueTime;
  29. private readonly IObservable<TSource> _other;
  30. private readonly IScheduler _scheduler;
  31. long _index;
  32. IDisposable _mainDisposable;
  33. IDisposable _otherDisposable;
  34. IDisposable _timerDisposable;
  35. public _(Relative parent, IObserver<TSource> observer, IDisposable cancel)
  36. : base(observer, cancel)
  37. {
  38. _dueTime = parent._dueTime;
  39. _other = parent._other;
  40. _scheduler = parent._scheduler;
  41. }
  42. public IDisposable Run(IObservable<TSource> source)
  43. {
  44. CreateTimer(0L);
  45. Disposable.SetSingle(ref _mainDisposable, source.SubscribeSafe(this));
  46. return this;
  47. }
  48. protected override void Dispose(bool disposing)
  49. {
  50. if (disposing)
  51. {
  52. Disposable.TryDispose(ref _mainDisposable);
  53. Disposable.TryDispose(ref _otherDisposable);
  54. Disposable.TryDispose(ref _timerDisposable);
  55. }
  56. base.Dispose(disposing);
  57. }
  58. private void CreateTimer(long idx)
  59. {
  60. if (Disposable.TrySetMultiple(ref _timerDisposable, null))
  61. {
  62. var d = _scheduler.Schedule((idx, instance: this), _dueTime, (_, state) => { state.instance.Timeout(state.idx); return Disposable.Empty; });
  63. Disposable.TrySetMultiple(ref _timerDisposable, d);
  64. }
  65. }
  66. private void Timeout(long idx)
  67. {
  68. if (Volatile.Read(ref _index) == idx && Interlocked.CompareExchange(ref _index, long.MaxValue, idx) == idx)
  69. {
  70. Disposable.TryDispose(ref _mainDisposable);
  71. var d = _other.Subscribe(GetForwarder());
  72. Disposable.SetSingle(ref _otherDisposable, d);
  73. }
  74. }
  75. public override void OnNext(TSource value)
  76. {
  77. var idx = Volatile.Read(ref _index);
  78. if (idx != long.MaxValue && Interlocked.CompareExchange(ref _index, idx + 1, idx) == idx)
  79. {
  80. // Do not swap in the BooleanDisposable.True here
  81. // As we'll need _timerDisposable to store the next timer
  82. // BD.True would cancel it immediately and break the operation
  83. Volatile.Read(ref _timerDisposable)?.Dispose();
  84. ForwardOnNext(value);
  85. CreateTimer(idx + 1);
  86. }
  87. }
  88. public override void OnError(Exception error)
  89. {
  90. if (Interlocked.Exchange(ref _index, long.MaxValue) != long.MaxValue)
  91. {
  92. Disposable.TryDispose(ref _timerDisposable);
  93. ForwardOnError(error);
  94. }
  95. }
  96. public override void OnCompleted()
  97. {
  98. if (Interlocked.Exchange(ref _index, long.MaxValue) != long.MaxValue)
  99. {
  100. Disposable.TryDispose(ref _timerDisposable);
  101. ForwardOnCompleted();
  102. }
  103. }
  104. }
  105. }
  106. internal sealed class Absolute : Producer<TSource, Absolute._>
  107. {
  108. private readonly IObservable<TSource> _source;
  109. private readonly DateTimeOffset _dueTime;
  110. private readonly IObservable<TSource> _other;
  111. private readonly IScheduler _scheduler;
  112. public Absolute(IObservable<TSource> source, DateTimeOffset dueTime, IObservable<TSource> other, IScheduler scheduler)
  113. {
  114. _source = source;
  115. _dueTime = dueTime;
  116. _other = other;
  117. _scheduler = scheduler;
  118. }
  119. protected override _ CreateSink(IObserver<TSource> observer, IDisposable cancel) => new _(_other, observer, cancel);
  120. protected override IDisposable Run(_ sink) => sink.Run(this);
  121. internal sealed class _ : IdentitySink<TSource>
  122. {
  123. private readonly IObservable<TSource> _other;
  124. private readonly object _gate = new object();
  125. private readonly SerialDisposable _subscription = new SerialDisposable();
  126. public _(IObservable<TSource> other, IObserver<TSource> observer, IDisposable cancel)
  127. : base(observer, cancel)
  128. {
  129. _other = other;
  130. }
  131. private bool _switched;
  132. public IDisposable Run(Absolute parent)
  133. {
  134. var original = new SingleAssignmentDisposable();
  135. _subscription.Disposable = original;
  136. _switched = false;
  137. var timer = parent._scheduler.Schedule(this, parent._dueTime, (_, state) => state.Timeout());
  138. original.Disposable = parent._source.SubscribeSafe(this);
  139. return StableCompositeDisposable.Create(_subscription, timer);
  140. }
  141. private IDisposable Timeout()
  142. {
  143. var timerWins = false;
  144. lock (_gate)
  145. {
  146. timerWins = !_switched;
  147. _switched = true;
  148. }
  149. if (timerWins)
  150. _subscription.Disposable = _other.SubscribeSafe(GetForwarder());
  151. return Disposable.Empty;
  152. }
  153. public override void OnNext(TSource value)
  154. {
  155. lock (_gate)
  156. {
  157. if (!_switched)
  158. ForwardOnNext(value);
  159. }
  160. }
  161. public override void OnError(Exception error)
  162. {
  163. var onErrorWins = false;
  164. lock (_gate)
  165. {
  166. onErrorWins = !_switched;
  167. _switched = true;
  168. }
  169. if (onErrorWins)
  170. {
  171. ForwardOnError(error);
  172. }
  173. }
  174. public override void OnCompleted()
  175. {
  176. var onCompletedWins = false;
  177. lock (_gate)
  178. {
  179. onCompletedWins = !_switched;
  180. _switched = true;
  181. }
  182. if (onCompletedWins)
  183. {
  184. ForwardOnCompleted();
  185. }
  186. }
  187. }
  188. }
  189. }
  190. internal sealed class Timeout<TSource, TTimeout> : Producer<TSource, Timeout<TSource, TTimeout>._>
  191. {
  192. private readonly IObservable<TSource> _source;
  193. private readonly IObservable<TTimeout> _firstTimeout;
  194. private readonly Func<TSource, IObservable<TTimeout>> _timeoutSelector;
  195. private readonly IObservable<TSource> _other;
  196. public Timeout(IObservable<TSource> source, IObservable<TTimeout> firstTimeout, Func<TSource, IObservable<TTimeout>> timeoutSelector, IObservable<TSource> other)
  197. {
  198. _source = source;
  199. _firstTimeout = firstTimeout;
  200. _timeoutSelector = timeoutSelector;
  201. _other = other;
  202. }
  203. protected override _ CreateSink(IObserver<TSource> observer, IDisposable cancel) => new _(this, observer, cancel);
  204. protected override IDisposable Run(_ sink) => sink.Run(this);
  205. internal sealed class _ : IdentitySink<TSource>
  206. {
  207. private readonly Func<TSource, IObservable<TTimeout>> _timeoutSelector;
  208. private readonly IObservable<TSource> _other;
  209. private readonly object _gate = new object();
  210. private readonly SerialDisposable _subscription = new SerialDisposable();
  211. private readonly SerialDisposable _timer = new SerialDisposable();
  212. public _(Timeout<TSource, TTimeout> parent, IObserver<TSource> observer, IDisposable cancel)
  213. : base(observer, cancel)
  214. {
  215. _timeoutSelector = parent._timeoutSelector;
  216. _other = parent._other;
  217. }
  218. private ulong _id;
  219. private bool _switched;
  220. public IDisposable Run(Timeout<TSource, TTimeout> parent)
  221. {
  222. var original = new SingleAssignmentDisposable();
  223. _subscription.Disposable = original;
  224. _id = 0UL;
  225. _switched = false;
  226. SetTimer(parent._firstTimeout);
  227. original.Disposable = parent._source.SubscribeSafe(this);
  228. return StableCompositeDisposable.Create(_subscription, _timer);
  229. }
  230. public override void OnNext(TSource value)
  231. {
  232. if (ObserverWins())
  233. {
  234. ForwardOnNext(value);
  235. var timeout = default(IObservable<TTimeout>);
  236. try
  237. {
  238. timeout = _timeoutSelector(value);
  239. }
  240. catch (Exception error)
  241. {
  242. ForwardOnError(error);
  243. return;
  244. }
  245. SetTimer(timeout);
  246. }
  247. }
  248. public override void OnError(Exception error)
  249. {
  250. if (ObserverWins())
  251. {
  252. ForwardOnError(error);
  253. }
  254. }
  255. public override void OnCompleted()
  256. {
  257. if (ObserverWins())
  258. {
  259. ForwardOnCompleted();
  260. }
  261. }
  262. private void SetTimer(IObservable<TTimeout> timeout)
  263. {
  264. var myid = _id;
  265. var d = new SingleAssignmentDisposable();
  266. _timer.Disposable = d;
  267. d.Disposable = timeout.SubscribeSafe(new TimeoutObserver(this, myid, d));
  268. }
  269. private sealed class TimeoutObserver : IObserver<TTimeout>
  270. {
  271. private readonly _ _parent;
  272. private readonly ulong _id;
  273. private readonly IDisposable _self;
  274. public TimeoutObserver(_ parent, ulong id, IDisposable self)
  275. {
  276. _parent = parent;
  277. _id = id;
  278. _self = self;
  279. }
  280. public void OnNext(TTimeout value)
  281. {
  282. if (TimerWins())
  283. _parent._subscription.Disposable = _parent._other.SubscribeSafe(_parent.GetForwarder());
  284. _self.Dispose();
  285. }
  286. public void OnError(Exception error)
  287. {
  288. if (TimerWins())
  289. {
  290. _parent.ForwardOnError(error);
  291. }
  292. }
  293. public void OnCompleted()
  294. {
  295. if (TimerWins())
  296. _parent._subscription.Disposable = _parent._other.SubscribeSafe(_parent.GetForwarder());
  297. }
  298. private bool TimerWins()
  299. {
  300. var res = false;
  301. lock (_parent._gate)
  302. {
  303. _parent._switched = (_parent._id == _id);
  304. res = _parent._switched;
  305. }
  306. return res;
  307. }
  308. }
  309. private bool ObserverWins()
  310. {
  311. var res = false;
  312. lock (_gate)
  313. {
  314. res = !_switched;
  315. if (res)
  316. {
  317. _id = unchecked(_id + 1);
  318. }
  319. }
  320. return res;
  321. }
  322. }
  323. }
  324. }