| | | 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 | | |
| | | 4 | | using System.Diagnostics; |
| | | 5 | | using System.Net; |
| | | 6 | | using System.Threading.Tasks; |
| | | 7 | | |
| | | 8 | | namespace System.Threading |
| | | 9 | | { |
| | | 10 | | /// <summary>Provides an async mutex.</summary> |
| | | 11 | | /// <remarks> |
| | | 12 | | /// This could be achieved with a <see cref="SemaphoreSlim"/> constructed with an initial |
| | | 13 | | /// and max limit of 1. However, this implementation is optimized to the needs of ManagedWebSocket, |
| | | 14 | | /// which is that we expect zero contention in typical use cases. |
| | | 15 | | /// </remarks> |
| | | 16 | | internal sealed class AsyncMutex |
| | | 17 | | { |
| | | 18 | | /// <summary>Fast-path gate count tracking access to the mutex.</summary> |
| | | 19 | | /// <remarks> |
| | | 20 | | /// If the value is 1, the mutex can be entered atomically with an interlocked operation. |
| | | 21 | | /// If the value is less than or equal to 0, the mutex is held and requires fallback to enter it. |
| | | 22 | | /// </remarks> |
| | 17284 | 23 | | private int _gate = 1; |
| | | 24 | | /// <summary>Secondary check guarded by the lock to indicate whether the mutex is acquired.</summary> |
| | | 25 | | /// <remarks> |
| | | 26 | | /// This is only meaningful after having updated <see cref="_gate"/> via interlockeds and taken the appropriate |
| | | 27 | | /// If after decrementing <see cref="_gate"/> we end up with a negative count, the mutex is contended, hence |
| | | 28 | | /// <see cref="_lockedSemaphoreFull"/> starting as <c>true</c>. The primary purpose of this field |
| | | 29 | | /// is to handle the race condition between one thread acquiring the mutex, then another thread trying to acquir |
| | | 30 | | /// and getting as far as completing the interlocked operation, and then the original thread releasing; at that |
| | | 31 | | /// it'll hit the lock and we need to store that the mutex is available to enter. If we instead used a |
| | | 32 | | /// SemaphoreSlim as the fallback from the interlockeds, this would have been its count, and it would have start |
| | | 33 | | /// with an initial count of 0. |
| | | 34 | | /// </remarks> |
| | 17284 | 35 | | private bool _lockedSemaphoreFull = true; |
| | | 36 | | /// <summary>The tail of the double-linked circular waiting queue.</summary> |
| | | 37 | | /// <remarks> |
| | | 38 | | /// Waiters are added at the tail. |
| | | 39 | | /// Items are dequeued from the head (tail.Prev). |
| | | 40 | | /// </remarks> |
| | | 41 | | private Waiter? _waitersTail; |
| | | 42 | | |
| | | 43 | | /// <summary>Gets whether the mutex is currently held by some operation (not necessarily the caller).</summary> |
| | | 44 | | /// <remarks>This should be used only for asserts and debugging.</remarks> |
| | 79530 | 45 | | public bool IsHeld => _gate != 1; |
| | | 46 | | |
| | | 47 | | /// <summary>Gets the object used to synchronize contended operations.</summary> |
| | 646364 | 48 | | private object SyncObj => this; |
| | | 49 | | |
| | | 50 | | /// <summary>Asynchronously waits to enter the mutex.</summary> |
| | | 51 | | /// <param name="cancellationToken">The CancellationToken token to observe.</param> |
| | | 52 | | /// <returns>A task that will complete when the mutex has been entered or the enter canceled.</returns> |
| | | 53 | | public Task EnterAsync(CancellationToken cancellationToken) |
| | 708392 | 54 | | { |
| | | 55 | | // If cancellation was requested, bail immediately. |
| | | 56 | | // If the mutex is not currently held nor contended, enter immediately. |
| | | 57 | | // Otherwise, fall back to a more expensive likely-asynchronous wait. |
| | | 58 | | |
| | 708392 | 59 | | if (cancellationToken.IsCancellationRequested) |
| | 0 | 60 | | { |
| | 0 | 61 | | return Task.FromCanceled(cancellationToken); |
| | | 62 | | } |
| | | 63 | | |
| | 708392 | 64 | | int gate = Interlocked.Decrement(ref _gate); |
| | 708392 | 65 | | if (gate >= 0) |
| | 385210 | 66 | | { |
| | 385210 | 67 | | return Task.CompletedTask; |
| | | 68 | | } |
| | | 69 | | |
| | 323182 | 70 | | if (NetEventSource.Log.IsEnabled()) NetEventSource.Trace(this, $"Waiting to enter, queue length {-gate}"); |
| | | 71 | | |
| | 323182 | 72 | | return Contended(cancellationToken); |
| | | 73 | | |
| | | 74 | | // Everything that follows is the equivalent of: |
| | | 75 | | // return _sem.WaitAsync(cancellationToken); |
| | | 76 | | // if _sem were to be constructed as `new SemaphoreSlim(0)`. |
| | | 77 | | |
| | | 78 | | Task Contended(CancellationToken cancellationToken) |
| | 323182 | 79 | | { |
| | 323182 | 80 | | var w = new Waiter(this); |
| | | 81 | | |
| | | 82 | | // We need to register for cancellation before storing the waiter into the list. |
| | | 83 | | // If we registered after, we might leak a registration if the mutex was exited and the waiter |
| | | 84 | | // removed from the list prior to CancellationRegistration being properly assigned. By registering befor |
| | | 85 | | // there's a different race condition, that of cancellation being requested prior to storing the waiter |
| | | 86 | | // the list; if that happens, we could end up adding the waiter and have it still stored in the list eve |
| | | 87 | | // though OnCancellation was called. So once we hold the lock, which OnCancellation also needs to take, |
| | | 88 | | // check again whether cancellation has been requested,and avoid storing the waiter if it has. |
| | 323182 | 89 | | w.CancellationRegistration = cancellationToken.UnsafeRegister((s, token) => OnCancellation(s, token), w) |
| | | 90 | | |
| | 323182 | 91 | | lock (SyncObj) |
| | 323182 | 92 | | { |
| | | 93 | | // Now that we're holding the lock, check to see whether the async lock is acquirable. |
| | 323182 | 94 | | if (!_lockedSemaphoreFull) |
| | 0 | 95 | | { |
| | | 96 | | // If we are able to acquire the lock, we're done; we just need to clean up after the registrati |
| | 0 | 97 | | w.CancellationRegistration.Unregister(); |
| | 0 | 98 | | _lockedSemaphoreFull = true; |
| | 0 | 99 | | return Task.CompletedTask; |
| | | 100 | | } |
| | | 101 | | |
| | | 102 | | // Now that we're holding the lock and thus synchronized with OnCancellation, check to see |
| | | 103 | | // if cancellation has been requested. |
| | 323182 | 104 | | if (cancellationToken.IsCancellationRequested) |
| | 0 | 105 | | { |
| | 0 | 106 | | w.TrySetCanceled(cancellationToken); |
| | 0 | 107 | | return w.Task; |
| | | 108 | | } |
| | | 109 | | |
| | | 110 | | // The lock couldn't be acquired. |
| | | 111 | | // Add the waiter to the linked list of waiters. |
| | 323182 | 112 | | if (_waitersTail is null) |
| | 323182 | 113 | | { |
| | 323182 | 114 | | w.Next = w.Prev = w; |
| | 323182 | 115 | | } |
| | | 116 | | else |
| | 0 | 117 | | { |
| | 0 | 118 | | Debug.Assert(_waitersTail.Next != null && _waitersTail.Prev != null); |
| | 0 | 119 | | w.Next = _waitersTail; |
| | 0 | 120 | | w.Prev = _waitersTail.Prev; |
| | 0 | 121 | | w.Prev.Next = w.Next.Prev = w; |
| | 0 | 122 | | } |
| | 323182 | 123 | | _waitersTail = w; |
| | 323182 | 124 | | } |
| | | 125 | | |
| | | 126 | | // Return the waiter as a value task. |
| | 323182 | 127 | | return w.Task; |
| | | 128 | | |
| | | 129 | | // Cancels the specified waiter if it's still in the list. |
| | | 130 | | static void OnCancellation(object? state, CancellationToken cancellationToken) |
| | 0 | 131 | | { |
| | 0 | 132 | | Waiter? w = (Waiter)state!; |
| | 0 | 133 | | AsyncMutex m = w.Owner; |
| | | 134 | | |
| | 0 | 135 | | lock (m.SyncObj) |
| | 0 | 136 | | { |
| | 0 | 137 | | bool inList = w.Next != null; |
| | 0 | 138 | | if (inList) |
| | 0 | 139 | | { |
| | | 140 | | // The waiter is in the list. |
| | 0 | 141 | | Debug.Assert(w.Prev != null); |
| | | 142 | | |
| | | 143 | | // The gate counter was decremented when this waiter was added. We need |
| | | 144 | | // to undo that. Since the waiter is still in the list, the lock must |
| | | 145 | | // still be held by someone, which means we don't need to do anything with |
| | | 146 | | // the result of this increment. If it increments to < 1, then there are |
| | | 147 | | // still other waiters. If it increments to 1, we're in a rare race condition |
| | | 148 | | // where there are no other waiters and the owner just incremented the gate |
| | | 149 | | // count; they would have seen it be < 1, so they will proceed to take the |
| | | 150 | | // contended code path and synchronize on the lock we're holding... once we |
| | | 151 | | // release it, they will appropriately update state. |
| | 0 | 152 | | Interlocked.Increment(ref m._gate); |
| | | 153 | | |
| | 0 | 154 | | if (w.Next == w) |
| | 0 | 155 | | { |
| | 0 | 156 | | Debug.Assert(m._waitersTail == w); |
| | 0 | 157 | | m._waitersTail = null; |
| | 0 | 158 | | } |
| | | 159 | | else |
| | 0 | 160 | | { |
| | 0 | 161 | | w.Next!.Prev = w.Prev; |
| | 0 | 162 | | w.Prev.Next = w.Next; |
| | 0 | 163 | | if (m._waitersTail == w) |
| | 0 | 164 | | { |
| | 0 | 165 | | m._waitersTail = w.Next; |
| | 0 | 166 | | } |
| | 0 | 167 | | } |
| | | 168 | | |
| | | 169 | | // Remove it from the list. |
| | 0 | 170 | | w.Next = w.Prev = null; |
| | 0 | 171 | | } |
| | | 172 | | else |
| | 0 | 173 | | { |
| | | 174 | | // The waiter was no longer in the list. We must not cancel it. |
| | 0 | 175 | | w = null; |
| | 0 | 176 | | } |
| | 0 | 177 | | } |
| | | 178 | | |
| | | 179 | | // If the waiter was in the list, we removed it under the lock and thus own |
| | | 180 | | // the ability to cancel it. Do so. |
| | 0 | 181 | | w?.TrySetCanceled(cancellationToken); |
| | 0 | 182 | | } |
| | 323182 | 183 | | } |
| | 708392 | 184 | | } |
| | | 185 | | |
| | | 186 | | /// <summary>Releases the mutex.</summary> |
| | | 187 | | /// <remarks>The caller must logically own the mutex. This is not validated.</remarks> |
| | | 188 | | public void Exit() |
| | 708392 | 189 | | { |
| | | 190 | | // This is the equivalent of: |
| | | 191 | | // _sem.Release(); |
| | | 192 | | // if _sem were to be constructed as `new SemaphoreSlim(0)`. |
| | 708392 | 193 | | int gate = Interlocked.Increment(ref _gate); |
| | 708392 | 194 | | if (gate < 1) |
| | 323182 | 195 | | { |
| | 323182 | 196 | | if (NetEventSource.Log.IsEnabled()) NetEventSource.Trace(this, $"Unblocking next waiter on exit, remaini |
| | 323182 | 197 | | Contended(); |
| | 323182 | 198 | | } |
| | | 199 | | |
| | | 200 | | void Contended() |
| | 323182 | 201 | | { |
| | | 202 | | Waiter? w; |
| | 323182 | 203 | | lock (SyncObj) |
| | 323182 | 204 | | { |
| | 323182 | 205 | | Debug.Assert(_lockedSemaphoreFull); |
| | | 206 | | |
| | 323182 | 207 | | w = _waitersTail; |
| | 323182 | 208 | | if (w is null) |
| | 0 | 209 | | { |
| | 0 | 210 | | _lockedSemaphoreFull = false; |
| | 0 | 211 | | } |
| | | 212 | | else |
| | 323182 | 213 | | { |
| | 323182 | 214 | | Debug.Assert(w.Next != null && w.Prev != null); |
| | 323182 | 215 | | Debug.Assert(w.Next != w || w.Prev == w); |
| | 323182 | 216 | | Debug.Assert(w.Prev != w || w.Next == w); |
| | | 217 | | |
| | 323182 | 218 | | if (w.Next == w) |
| | 323182 | 219 | | { |
| | 323182 | 220 | | _waitersTail = null; |
| | 323182 | 221 | | } |
| | | 222 | | else |
| | 0 | 223 | | { |
| | 0 | 224 | | w = w.Prev; // get the head |
| | 0 | 225 | | Debug.Assert(w.Next != null && w.Prev != null); |
| | 0 | 226 | | Debug.Assert(w.Next != w && w.Prev != w); |
| | | 227 | | |
| | 0 | 228 | | w.Next.Prev = w.Prev; |
| | 0 | 229 | | w.Prev.Next = w.Next; |
| | 0 | 230 | | } |
| | | 231 | | |
| | 323182 | 232 | | w.Next = w.Prev = null; |
| | 323182 | 233 | | } |
| | 323182 | 234 | | } |
| | | 235 | | |
| | | 236 | | // Either there wasn't a waiter, or we got one and successfully removed it from the list, |
| | | 237 | | // at which point we own the ability to complete it. Do so. |
| | 323182 | 238 | | if (w is not null) |
| | 323182 | 239 | | { |
| | 323182 | 240 | | w.CancellationRegistration.Unregister(); |
| | 323182 | 241 | | w.TrySetResult(); |
| | 323182 | 242 | | } |
| | 323182 | 243 | | } |
| | 708392 | 244 | | } |
| | | 245 | | |
| | | 246 | | /// <summary>Represents a waiter for the mutex.</summary> |
| | | 247 | | private sealed class Waiter : TaskCompletionSource |
| | | 248 | | { |
| | 646364 | 249 | | public Waiter(AsyncMutex owner) : base(TaskCreationOptions.RunContinuationsAsynchronously) => Owner = owner; |
| | 0 | 250 | | public AsyncMutex Owner { get; } |
| | 646364 | 251 | | public CancellationTokenRegistration CancellationRegistration { get; set; } |
| | 1939092 | 252 | | public Waiter? Next { get; set; } |
| | 1615910 | 253 | | public Waiter? Prev { get; set; } |
| | | 254 | | } |
| | | 255 | | } |
| | | 256 | | } |
| | | 257 | | |