| | | 1 | | // Copyright (c) ZeroC, Inc. |
| | | 2 | | |
| | | 3 | | using System.Diagnostics; |
| | | 4 | | |
| | | 5 | | namespace Ice.Internal; |
| | | 6 | | |
| | | 7 | | public class RetryTask : TimerTask, CancellationHandler |
| | | 8 | | { |
| | | 9 | | public RetryTask(Instance instance, RetryQueue retryQueue, ProxyOutgoingAsyncBase outAsync) |
| | | 10 | | { |
| | | 11 | | _instance = instance; |
| | | 12 | | _retryQueue = retryQueue; |
| | | 13 | | _outAsync = outAsync; |
| | | 14 | | } |
| | | 15 | | |
| | | 16 | | public void runTimerTask() |
| | | 17 | | { |
| | | 18 | | try |
| | | 19 | | { |
| | | 20 | | _outAsync.retry(); |
| | | 21 | | } |
| | | 22 | | finally |
| | | 23 | | { |
| | | 24 | | // |
| | | 25 | | // NOTE: this must be called last, destroy() blocks until all task |
| | | 26 | | // are removed to prevent the client thread pool to be destroyed |
| | | 27 | | // (we still need the client thread pool at this point to call |
| | | 28 | | // exception callbacks with CommunicatorDestroyedException). |
| | | 29 | | // |
| | | 30 | | _retryQueue.remove(this); |
| | | 31 | | } |
| | | 32 | | } |
| | | 33 | | |
| | | 34 | | public void asyncRequestCanceled(OutgoingAsyncBase outAsync, Ice.LocalException ex) |
| | | 35 | | { |
| | | 36 | | Debug.Assert(_outAsync == outAsync); |
| | | 37 | | if (_retryQueue.cancel(this)) |
| | | 38 | | { |
| | | 39 | | if (_instance.traceLevels().retry >= 1) |
| | | 40 | | { |
| | | 41 | | _instance.initializationData().logger.trace( |
| | | 42 | | _instance.traceLevels().retryCat, |
| | | 43 | | $"operation retry canceled\n{ex}"); |
| | | 44 | | } |
| | | 45 | | if (_outAsync.exception(ex)) |
| | | 46 | | { |
| | | 47 | | _outAsync.invokeExceptionAsync(); |
| | | 48 | | } |
| | | 49 | | } |
| | | 50 | | } |
| | | 51 | | |
| | | 52 | | public void destroy() |
| | | 53 | | { |
| | | 54 | | try |
| | | 55 | | { |
| | | 56 | | _outAsync.abort(new Ice.CommunicatorDestroyedException()); |
| | | 57 | | } |
| | | 58 | | catch (Ice.CommunicatorDestroyedException) |
| | | 59 | | { |
| | | 60 | | // Abort can throw if there's no callback, just ignore in this case |
| | | 61 | | } |
| | | 62 | | } |
| | | 63 | | |
| | | 64 | | private readonly Instance _instance; |
| | | 65 | | private readonly RetryQueue _retryQueue; |
| | | 66 | | private readonly ProxyOutgoingAsyncBase _outAsync; |
| | | 67 | | } |
| | | 68 | | |
| | | 69 | | public class RetryQueue |
| | | 70 | | { |
| | 1 | 71 | | public RetryQueue(Instance instance) => _instance = instance; |
| | | 72 | | |
| | | 73 | | public void add(ProxyOutgoingAsyncBase outAsync, int interval) |
| | | 74 | | { |
| | | 75 | | Debug.Assert(interval >= 0); |
| | 1 | 76 | | lock (_mutex) |
| | | 77 | | { |
| | 1 | 78 | | if (_instance == null) |
| | | 79 | | { |
| | 0 | 80 | | throw new Ice.CommunicatorDestroyedException(); |
| | | 81 | | } |
| | 1 | 82 | | var task = new RetryTask(_instance, this, outAsync); |
| | 1 | 83 | | outAsync.cancelable(task); // This will throw if the request is canceled. |
| | 1 | 84 | | _instance.timer().schedule(task, TimeSpan.FromMilliseconds(interval)); |
| | 1 | 85 | | _requests.Add(task, null); |
| | 1 | 86 | | } |
| | 1 | 87 | | } |
| | | 88 | | |
| | | 89 | | public void destroy() |
| | | 90 | | { |
| | 1 | 91 | | lock (_mutex) |
| | | 92 | | { |
| | 1 | 93 | | var keep = new Dictionary<RetryTask, object>(); |
| | 1 | 94 | | foreach (RetryTask task in _requests.Keys) |
| | | 95 | | { |
| | 0 | 96 | | if (_instance.timer().cancel(task)) |
| | | 97 | | { |
| | 0 | 98 | | task.destroy(); |
| | | 99 | | } |
| | | 100 | | else |
| | | 101 | | { |
| | 0 | 102 | | keep.Add(task, null); |
| | | 103 | | } |
| | | 104 | | } |
| | 1 | 105 | | _requests = keep; |
| | 1 | 106 | | _instance = null; |
| | 1 | 107 | | while (_requests.Count > 0) |
| | | 108 | | { |
| | 0 | 109 | | System.Threading.Monitor.Wait(_mutex); |
| | | 110 | | } |
| | 1 | 111 | | } |
| | 1 | 112 | | } |
| | | 113 | | |
| | | 114 | | public void remove(RetryTask task) |
| | | 115 | | { |
| | 1 | 116 | | lock (_mutex) |
| | | 117 | | { |
| | 1 | 118 | | if (_requests.Remove(task)) |
| | | 119 | | { |
| | 1 | 120 | | if (_instance == null && _requests.Count == 0) |
| | | 121 | | { |
| | | 122 | | // If we are destroying the queue, destroy is probably waiting on the queue to be empty. |
| | 0 | 123 | | System.Threading.Monitor.Pulse(_mutex); |
| | | 124 | | } |
| | | 125 | | } |
| | 1 | 126 | | } |
| | 1 | 127 | | } |
| | | 128 | | |
| | | 129 | | public bool cancel(RetryTask task) |
| | | 130 | | { |
| | 1 | 131 | | lock (_mutex) |
| | | 132 | | { |
| | | 133 | | // Only remove the task if we cancel it in the timer before it runs. If timer().cancel returns |
| | | 134 | | // false, the task is already executing and runTimerTask will call remove() to erase it; removing |
| | | 135 | | // it here would let destroy() observe an empty queue while the task is still running. When the |
| | | 136 | | // queue is being destroyed (_instance is null) every remaining task is running and is likewise |
| | | 137 | | // removed by remove(), which wakes destroy(). |
| | 1 | 138 | | if (_instance != null && _instance.timer().cancel(task)) |
| | | 139 | | { |
| | 1 | 140 | | return _requests.Remove(task); |
| | | 141 | | } |
| | 0 | 142 | | return false; |
| | | 143 | | } |
| | 1 | 144 | | } |
| | | 145 | | |
| | | 146 | | private Instance _instance; |
| | 1 | 147 | | private Dictionary<RetryTask, object> _requests = new(); |
| | 1 | 148 | | private readonly object _mutex = new(); |
| | | 149 | | } |