| | | 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.Collections.Generic; |
| | | 5 | | using System.Diagnostics; |
| | | 6 | | using System.Runtime.CompilerServices; |
| | | 7 | | using System.Runtime.InteropServices; |
| | | 8 | | using System.Threading; |
| | | 9 | | |
| | | 10 | | namespace System.Buffers |
| | | 11 | | { |
| | | 12 | | /// <summary> |
| | | 13 | | /// Provides an ArrayPool implementation meant to be used as the singleton returned from ArrayPool.Shared. |
| | | 14 | | /// </summary> |
| | | 15 | | /// <remarks> |
| | | 16 | | /// The implementation uses a tiered caching scheme, with a small per-thread cache for each array size, followed |
| | | 17 | | /// by a cache per array size shared by all threads, split into per-core stacks meant to be used by threads |
| | | 18 | | /// running on that core. Locks are used to protect each per-core stack, because a thread can migrate after |
| | | 19 | | /// checking its processor number, because multiple threads could interleave on the same core, and because |
| | | 20 | | /// a thread is allowed to check other core's buckets if its core's bucket is empty/full. |
| | | 21 | | /// </remarks> |
| | | 22 | | internal sealed partial class SharedArrayPool<T> : ArrayPool<T> |
| | | 23 | | { |
| | | 24 | | /// <summary>The number of buckets (array sizes) in the pool, one for each array length, starting from length 16 |
| | | 25 | | private const int NumBuckets = 27; // Utilities.SelectBucketIndex(1024 * 1024 * 1024 + 1) |
| | | 26 | | |
| | | 27 | | /// <summary>A per-thread array of arrays, to cache one array per array size per thread.</summary> |
| | | 28 | | [ThreadStatic] |
| | | 29 | | private static SharedArrayPoolThreadLocalArray[]? t_tlsBuckets; |
| | | 30 | | /// <summary>Used to keep track of all thread local buckets for trimming if needed.</summary> |
| | 4 | 31 | | private readonly ConditionalWeakTable<SharedArrayPoolThreadLocalArray[], object?> _allTlsBuckets = new Condition |
| | | 32 | | /// <summary> |
| | | 33 | | /// An array of per-core partitions. The slots are lazily initialized to avoid creating |
| | | 34 | | /// lots of overhead for unused array sizes. |
| | | 35 | | /// </summary> |
| | 4 | 36 | | private readonly SharedArrayPoolPartitions?[] _buckets = new SharedArrayPoolPartitions[NumBuckets]; |
| | | 37 | | /// <summary>Whether the callback to trim arrays in response to memory pressure has been created.</summary> |
| | | 38 | | private bool _trimCallbackCreated; |
| | | 39 | | |
| | | 40 | | /// <summary>Allocate a new <see cref="SharedArrayPoolPartitions"/> and try to store it into the <see cref="_buc |
| | | 41 | | private SharedArrayPoolPartitions CreatePerCorePartitions(int bucketIndex) |
| | | 42 | | { |
| | 0 | 43 | | var inst = new SharedArrayPoolPartitions(); |
| | 0 | 44 | | return Interlocked.CompareExchange(ref _buckets[bucketIndex], inst, null) ?? inst; |
| | | 45 | | } |
| | | 46 | | |
| | | 47 | | /// <summary>Gets an ID for the pool to use with events.</summary> |
| | 0 | 48 | | private int Id => GetHashCode(); |
| | | 49 | | |
| | | 50 | | [MethodImpl(MethodImplOptions.NoInlining)] |
| | | 51 | | public override T[] Rent(int minimumLength) |
| | | 52 | | { |
| | 69605 | 53 | | ArrayPoolEventSource log = ArrayPoolEventSource.Log; |
| | | 54 | | T[]? buffer; |
| | | 55 | | |
| | | 56 | | // Get the bucket number for the array length. The result may be out of range of buckets, |
| | | 57 | | // either for too large a value or for 0 and negative values. |
| | 69605 | 58 | | int bucketIndex = Utilities.SelectBucketIndex(minimumLength); |
| | | 59 | | |
| | | 60 | | // First, try to get an array from TLS if possible. |
| | 69605 | 61 | | SharedArrayPoolThreadLocalArray[]? tlsBuckets = t_tlsBuckets; |
| | 69605 | 62 | | if (tlsBuckets is not null && (uint)bucketIndex < (uint)tlsBuckets.Length) |
| | | 63 | | { |
| | 69600 | 64 | | buffer = Unsafe.As<T[]>(tlsBuckets[bucketIndex].Array); |
| | 69600 | 65 | | if (buffer is not null) |
| | | 66 | | { |
| | 69587 | 67 | | tlsBuckets[bucketIndex].Array = null; |
| | 69587 | 68 | | if (log.IsEnabled()) |
| | | 69 | | { |
| | 0 | 70 | | log.BufferRented(buffer.GetHashCode(), buffer.Length, Id, bucketIndex); |
| | | 71 | | } |
| | 69587 | 72 | | return buffer; |
| | | 73 | | } |
| | | 74 | | } |
| | | 75 | | |
| | | 76 | | // Next, try to get an array from one of the partitions. |
| | 18 | 77 | | SharedArrayPoolPartitions?[] perCoreBuckets = _buckets; |
| | 18 | 78 | | if ((uint)bucketIndex < (uint)perCoreBuckets.Length) |
| | | 79 | | { |
| | 18 | 80 | | SharedArrayPoolPartitions? b = perCoreBuckets[bucketIndex]; |
| | 18 | 81 | | if (b is not null) |
| | | 82 | | { |
| | 0 | 83 | | buffer = Unsafe.As<T[]>(b.TryPop()); |
| | 0 | 84 | | if (buffer is not null) |
| | | 85 | | { |
| | 0 | 86 | | if (log.IsEnabled()) |
| | | 87 | | { |
| | 0 | 88 | | log.BufferRented(buffer.GetHashCode(), buffer.Length, Id, bucketIndex); |
| | | 89 | | } |
| | 0 | 90 | | return buffer; |
| | | 91 | | } |
| | | 92 | | } |
| | | 93 | | |
| | | 94 | | // No buffer available. Ensure the length we'll allocate matches that of a bucket |
| | | 95 | | // so we can later return it. |
| | 18 | 96 | | minimumLength = Utilities.GetMaxSizeForBucket(bucketIndex); |
| | | 97 | | } |
| | 0 | 98 | | else if (minimumLength == 0) |
| | | 99 | | { |
| | | 100 | | // We allow requesting zero-length arrays (even though pooling such an array isn't valuable) |
| | | 101 | | // as it's a valid length array, and we want the pool to be usable in general instead of using |
| | | 102 | | // `new`, even for computed lengths. But, there's no need to log the empty array. Our pool is |
| | | 103 | | // effectively infinite for empty arrays and we'll never allocate for rents and never store for returns. |
| | 0 | 104 | | return []; |
| | | 105 | | } |
| | | 106 | | else |
| | | 107 | | { |
| | 0 | 108 | | ArgumentOutOfRangeException.ThrowIfNegative(minimumLength); |
| | | 109 | | } |
| | | 110 | | |
| | | 111 | | // For large arrays, we prefer to avoid the zero-initialization costs. However, as the resulting |
| | | 112 | | // arrays could end up containing arbitrary bit patterns, we only allow this for types for which |
| | | 113 | | // every possible bit pattern is valid. |
| | 18 | 114 | | buffer = typeof(T).IsPrimitive && typeof(T) != typeof(bool) ? |
| | 18 | 115 | | GC.AllocateUninitializedArray<T>(minimumLength) : |
| | 18 | 116 | | new T[minimumLength]; |
| | | 117 | | |
| | 18 | 118 | | if (log.IsEnabled()) |
| | | 119 | | { |
| | 0 | 120 | | int bufferId = buffer.GetHashCode(); |
| | 0 | 121 | | log.BufferRented(bufferId, buffer.Length, Id, ArrayPoolEventSource.NoBucketId); |
| | 0 | 122 | | log.BufferAllocated(bufferId, buffer.Length, Id, ArrayPoolEventSource.NoBucketId, bucketIndex >= _bucket |
| | 0 | 123 | | ArrayPoolEventSource.BufferAllocatedReason.OverMaximumSize : |
| | 0 | 124 | | ArrayPoolEventSource.BufferAllocatedReason.PoolExhausted); |
| | | 125 | | } |
| | 18 | 126 | | return buffer; |
| | | 127 | | } |
| | | 128 | | |
| | | 129 | | [MethodImpl(MethodImplOptions.NoInlining)] |
| | | 130 | | public override void Return(T[] array, bool clearArray = false) |
| | | 131 | | { |
| | 69605 | 132 | | if (array is null) |
| | | 133 | | { |
| | 0 | 134 | | ThrowHelper.ThrowArgumentNullException(ExceptionArgument.array); |
| | | 135 | | } |
| | | 136 | | |
| | | 137 | | // Determine with what bucket this array length is associated |
| | 69605 | 138 | | int bucketIndex = Utilities.SelectBucketIndex(array.Length); |
| | | 139 | | |
| | | 140 | | // Make sure our TLS buckets are initialized. Technically we could avoid doing |
| | | 141 | | // this if the array being returned is erroneous or too large for the pool, but the |
| | | 142 | | // former condition is an error we don't need to optimize for, and the latter is incredibly |
| | | 143 | | // rare, given a max size of 1B elements. |
| | 69605 | 144 | | SharedArrayPoolThreadLocalArray[] tlsBuckets = t_tlsBuckets ?? InitializeTlsBucketsAndTrimming(); |
| | | 145 | | |
| | 69605 | 146 | | bool haveBucket = false; |
| | 69605 | 147 | | bool returned = true; |
| | 69605 | 148 | | if ((uint)bucketIndex < (uint)tlsBuckets.Length) |
| | | 149 | | { |
| | 69605 | 150 | | haveBucket = true; |
| | | 151 | | |
| | | 152 | | // Clear the array if the user requested it. |
| | 69605 | 153 | | if (clearArray) |
| | | 154 | | { |
| | 0 | 155 | | Array.Clear(array); |
| | | 156 | | } |
| | | 157 | | |
| | | 158 | | // Check to see if the buffer is the correct size for this bucket. |
| | 69605 | 159 | | if (array.Length != Utilities.GetMaxSizeForBucket(bucketIndex)) |
| | | 160 | | { |
| | 0 | 161 | | throw new ArgumentException(SR.ArgumentException_BufferNotFromPool, nameof(array)); |
| | | 162 | | } |
| | | 163 | | |
| | | 164 | | // Store the array into the TLS bucket. If there's already an array in it, |
| | | 165 | | // push that array down into the partitions, preferring to keep the latest |
| | | 166 | | // one in TLS for better locality. |
| | 69605 | 167 | | ref SharedArrayPoolThreadLocalArray tla = ref tlsBuckets[bucketIndex]; |
| | 69605 | 168 | | Array? prev = tla.Array; |
| | 69605 | 169 | | tla = new SharedArrayPoolThreadLocalArray(array); |
| | 69605 | 170 | | if (prev is not null) |
| | | 171 | | { |
| | 0 | 172 | | SharedArrayPoolPartitions partitionsForArraySize = _buckets[bucketIndex] ?? CreatePerCorePartitions( |
| | 0 | 173 | | returned = partitionsForArraySize.TryPush(prev); |
| | | 174 | | } |
| | | 175 | | } |
| | | 176 | | |
| | | 177 | | // Log that the buffer was returned |
| | 69605 | 178 | | ArrayPoolEventSource log = ArrayPoolEventSource.Log; |
| | 69605 | 179 | | if (log.IsEnabled() && array.Length != 0) |
| | | 180 | | { |
| | 0 | 181 | | log.BufferReturned(array.GetHashCode(), array.Length, Id); |
| | 0 | 182 | | if (!(haveBucket & returned)) |
| | | 183 | | { |
| | 0 | 184 | | log.BufferDropped(array.GetHashCode(), array.Length, Id, |
| | 0 | 185 | | haveBucket ? bucketIndex : ArrayPoolEventSource.NoBucketId, |
| | 0 | 186 | | haveBucket ? ArrayPoolEventSource.BufferDroppedReason.Full : ArrayPoolEventSource.BufferDroppedR |
| | | 187 | | } |
| | | 188 | | } |
| | 69605 | 189 | | } |
| | | 190 | | |
| | | 191 | | public bool Trim() |
| | | 192 | | { |
| | 0 | 193 | | int currentMilliseconds = Environment.TickCount; |
| | 0 | 194 | | Utilities.MemoryPressure pressure = Utilities.GetMemoryPressure(); |
| | | 195 | | |
| | | 196 | | // Log that we're trimming. |
| | 0 | 197 | | ArrayPoolEventSource log = ArrayPoolEventSource.Log; |
| | 0 | 198 | | if (log.IsEnabled()) |
| | | 199 | | { |
| | 0 | 200 | | log.BufferTrimPoll(currentMilliseconds, (int)pressure); |
| | | 201 | | } |
| | | 202 | | |
| | | 203 | | // Trim each of the per-core buckets. |
| | 0 | 204 | | SharedArrayPoolPartitions?[] perCoreBuckets = _buckets; |
| | 0 | 205 | | for (int i = 0; i < perCoreBuckets.Length; i++) |
| | | 206 | | { |
| | 0 | 207 | | perCoreBuckets[i]?.Trim(currentMilliseconds, Id, pressure); |
| | | 208 | | } |
| | | 209 | | |
| | | 210 | | // Trim each of the TLS buckets. Note that threads may be modifying their TLS slots concurrently with |
| | | 211 | | // this trimming happening. We do not force synchronization with those operations, so we accept the fact |
| | | 212 | | // that we may end up firing a trimming event even if an array wasn't trimmed, and potentially |
| | | 213 | | // trim an array we didn't need to. Both of these should be rare occurrences. |
| | | 214 | | |
| | | 215 | | // Under high pressure, release all thread locals. |
| | 0 | 216 | | if (pressure == Utilities.MemoryPressure.High) |
| | | 217 | | { |
| | 0 | 218 | | if (!log.IsEnabled()) |
| | | 219 | | { |
| | 0 | 220 | | foreach (KeyValuePair<SharedArrayPoolThreadLocalArray[], object?> tlsBuckets in _allTlsBuckets) |
| | | 221 | | { |
| | 0 | 222 | | Array.Clear(tlsBuckets.Key); |
| | | 223 | | } |
| | | 224 | | } |
| | | 225 | | else |
| | | 226 | | { |
| | 0 | 227 | | foreach (KeyValuePair<SharedArrayPoolThreadLocalArray[], object?> tlsBuckets in _allTlsBuckets) |
| | | 228 | | { |
| | 0 | 229 | | SharedArrayPoolThreadLocalArray[] buckets = tlsBuckets.Key; |
| | 0 | 230 | | for (int i = 0; i < buckets.Length; i++) |
| | | 231 | | { |
| | 0 | 232 | | if (Interlocked.Exchange(ref buckets[i].Array, null) is T[] buffer) |
| | | 233 | | { |
| | 0 | 234 | | log.BufferTrimmed(buffer.GetHashCode(), buffer.Length, Id); |
| | | 235 | | } |
| | | 236 | | } |
| | | 237 | | } |
| | | 238 | | } |
| | | 239 | | } |
| | | 240 | | else |
| | | 241 | | { |
| | | 242 | | // Otherwise, release thread locals based on how long we've observed them to be stored. This time is |
| | | 243 | | // approximate, with the time set not when the array is stored but when we see it during a Trim, so it |
| | | 244 | | // takes at least two Trim calls (and thus two gen2 GCs) to drop an array, unless we're in high memory |
| | | 245 | | // pressure. These values have been set arbitrarily; we could tune them in the future. |
| | 0 | 246 | | uint millisecondsThreshold = pressure switch |
| | 0 | 247 | | { |
| | 0 | 248 | | Utilities.MemoryPressure.Medium => 15_000, |
| | 0 | 249 | | _ => 30_000, |
| | 0 | 250 | | }; |
| | | 251 | | |
| | 0 | 252 | | foreach (KeyValuePair<SharedArrayPoolThreadLocalArray[], object?> tlsBuckets in _allTlsBuckets) |
| | | 253 | | { |
| | 0 | 254 | | SharedArrayPoolThreadLocalArray[] buckets = tlsBuckets.Key; |
| | 0 | 255 | | for (int i = 0; i < buckets.Length; i++) |
| | | 256 | | { |
| | 0 | 257 | | if (buckets[i].Array is null) |
| | | 258 | | { |
| | | 259 | | continue; |
| | | 260 | | } |
| | | 261 | | |
| | | 262 | | // We treat 0 to mean it hasn't yet been seen in a Trim call. In the very rare case where Trim r |
| | | 263 | | // it'll take an extra Trim call to remove the array. |
| | 0 | 264 | | int lastSeen = buckets[i].MillisecondsTimeStamp; |
| | 0 | 265 | | if (lastSeen == 0) |
| | | 266 | | { |
| | 0 | 267 | | buckets[i].MillisecondsTimeStamp = currentMilliseconds; |
| | | 268 | | } |
| | 0 | 269 | | else if ((currentMilliseconds - lastSeen) >= millisecondsThreshold) |
| | | 270 | | { |
| | | 271 | | // Time noticeably wrapped, or we've surpassed the threshold. |
| | | 272 | | // Clear out the array, and log its being trimmed if desired. |
| | 0 | 273 | | if (Interlocked.Exchange(ref buckets[i].Array, null) is T[] buffer && |
| | 0 | 274 | | log.IsEnabled()) |
| | | 275 | | { |
| | 0 | 276 | | log.BufferTrimmed(buffer.GetHashCode(), buffer.Length, Id); |
| | | 277 | | } |
| | | 278 | | } |
| | | 279 | | } |
| | | 280 | | } |
| | | 281 | | } |
| | | 282 | | |
| | 0 | 283 | | return true; |
| | | 284 | | } |
| | | 285 | | |
| | | 286 | | private SharedArrayPoolThreadLocalArray[] InitializeTlsBucketsAndTrimming() |
| | | 287 | | { |
| | 4 | 288 | | Debug.Assert(t_tlsBuckets is null, $"Non-null {nameof(t_tlsBuckets)}"); |
| | | 289 | | |
| | 4 | 290 | | var tlsBuckets = new SharedArrayPoolThreadLocalArray[NumBuckets]; |
| | 4 | 291 | | t_tlsBuckets = tlsBuckets; |
| | | 292 | | |
| | 4 | 293 | | _allTlsBuckets.Add(tlsBuckets, null); |
| | 4 | 294 | | if (!Interlocked.Exchange(ref _trimCallbackCreated, true)) |
| | | 295 | | { |
| | 4 | 296 | | Gen2GcCallback.Register(s => ((SharedArrayPool<T>)s).Trim(), this); |
| | | 297 | | } |
| | | 298 | | |
| | 4 | 299 | | return tlsBuckets; |
| | | 300 | | } |
| | | 301 | | } |
| | | 302 | | |
| | | 303 | | // The following partition types are separated out of SharedArrayPool<T> to avoid |
| | | 304 | | // them being generic, in order to avoid the generic code size increase that would |
| | | 305 | | // result, in particular for Native AOT. The only thing that's necessary to actually |
| | | 306 | | // be generic is the return type of TryPop, and we can handle that at the access |
| | | 307 | | // site with a well-placed Unsafe.As. |
| | | 308 | | |
| | | 309 | | /// <summary>Wrapper for arrays stored in ThreadStatic buckets.</summary> |
| | | 310 | | internal struct SharedArrayPoolThreadLocalArray |
| | | 311 | | { |
| | | 312 | | /// <summary>The stored array.</summary> |
| | | 313 | | public Array? Array; |
| | | 314 | | /// <summary>Environment.TickCount timestamp for when this array was observed by Trim.</summary> |
| | | 315 | | public int MillisecondsTimeStamp; |
| | | 316 | | |
| | | 317 | | public SharedArrayPoolThreadLocalArray(Array array) |
| | | 318 | | { |
| | | 319 | | Array = array; |
| | | 320 | | MillisecondsTimeStamp = 0; |
| | | 321 | | } |
| | | 322 | | } |
| | | 323 | | |
| | | 324 | | /// <summary>Provides a collection of partitions, each of which is a pool of arrays.</summary> |
| | | 325 | | internal sealed class SharedArrayPoolPartitions |
| | | 326 | | { |
| | | 327 | | /// <summary>The partitions.</summary> |
| | | 328 | | private readonly Partition[] _partitions; |
| | | 329 | | |
| | | 330 | | /// <summary>Initializes the partitions.</summary> |
| | | 331 | | public SharedArrayPoolPartitions() |
| | | 332 | | { |
| | | 333 | | // Create the partitions. We create as many as there are processors, limited by our max. |
| | | 334 | | var partitions = new Partition[SharedArrayPoolStatics.s_partitionCount]; |
| | | 335 | | for (int i = 0; i < partitions.Length; i++) |
| | | 336 | | { |
| | | 337 | | partitions[i] = new Partition(); |
| | | 338 | | } |
| | | 339 | | _partitions = partitions; |
| | | 340 | | } |
| | | 341 | | |
| | | 342 | | /// <summary> |
| | | 343 | | /// Try to push the array into any partition with available space, starting with partition associated with the c |
| | | 344 | | /// If all partitions are full, the array will be dropped. |
| | | 345 | | /// </summary> |
| | | 346 | | [MethodImpl(MethodImplOptions.AggressiveInlining)] |
| | | 347 | | public bool TryPush(Array array) |
| | | 348 | | { |
| | | 349 | | // Try to push on to the associated partition first. If that fails, |
| | | 350 | | // round-robin through the other partitions. |
| | | 351 | | Partition[] partitions = _partitions; |
| | | 352 | | int index = (int)((uint)Thread.GetCurrentProcessorId() % (uint)SharedArrayPoolStatics.s_partitionCount); // |
| | | 353 | | for (int i = 0; i < partitions.Length; i++) |
| | | 354 | | { |
| | | 355 | | if (partitions[index].TryPush(array)) return true; |
| | | 356 | | if (++index == partitions.Length) index = 0; |
| | | 357 | | } |
| | | 358 | | |
| | | 359 | | return false; |
| | | 360 | | } |
| | | 361 | | |
| | | 362 | | /// <summary> |
| | | 363 | | /// Try to pop an array from any partition with available arrays, starting with partition associated with the cu |
| | | 364 | | /// If all partitions are empty, null is returned. |
| | | 365 | | /// </summary> |
| | | 366 | | [MethodImpl(MethodImplOptions.AggressiveInlining)] |
| | | 367 | | public Array? TryPop() |
| | | 368 | | { |
| | | 369 | | // Try to pop from the associated partition first. If that fails, round-robin through the other partitions. |
| | | 370 | | Array? arr; |
| | | 371 | | Partition[] partitions = _partitions; |
| | | 372 | | int index = (int)((uint)Thread.GetCurrentProcessorId() % (uint)SharedArrayPoolStatics.s_partitionCount); // |
| | | 373 | | for (int i = 0; i < partitions.Length; i++) |
| | | 374 | | { |
| | | 375 | | if ((arr = partitions[index].TryPop()) is not null) return arr; |
| | | 376 | | if (++index == partitions.Length) index = 0; |
| | | 377 | | } |
| | | 378 | | return null; |
| | | 379 | | } |
| | | 380 | | |
| | | 381 | | public void Trim(int currentMilliseconds, int id, Utilities.MemoryPressure pressure) |
| | | 382 | | { |
| | | 383 | | Partition[] partitions = _partitions; |
| | | 384 | | for (int i = 0; i < partitions.Length; i++) |
| | | 385 | | { |
| | | 386 | | partitions[i].Trim(currentMilliseconds, id, pressure); |
| | | 387 | | } |
| | | 388 | | } |
| | | 389 | | |
| | | 390 | | /// <summary>Provides a simple, bounded stack of arrays, protected by a lock.</summary> |
| | | 391 | | private sealed class Partition |
| | | 392 | | { |
| | | 393 | | /// <summary>The arrays in the partition.</summary> |
| | | 394 | | private readonly Array?[] _arrays = new Array[SharedArrayPoolStatics.s_maxArraysPerPartition]; |
| | | 395 | | /// <summary>Number of arrays stored in <see cref="_arrays"/>.</summary> |
| | | 396 | | private int _count; |
| | | 397 | | /// <summary>Timestamp set by Trim when it sees this as 0.</summary> |
| | | 398 | | private int _millisecondsTimestamp; |
| | | 399 | | |
| | | 400 | | [MethodImpl(MethodImplOptions.AggressiveInlining)] |
| | | 401 | | public bool TryPush(Array array) |
| | | 402 | | { |
| | | 403 | | bool enqueued = false; |
| | | 404 | | Monitor.Enter(this); |
| | | 405 | | Array?[] arrays = _arrays; |
| | | 406 | | int count = _count; |
| | | 407 | | if ((uint)count < (uint)arrays.Length) |
| | | 408 | | { |
| | | 409 | | if (count == 0) |
| | | 410 | | { |
| | | 411 | | // Reset the time stamp now that we're transitioning from empty to non-empty. |
| | | 412 | | // Trim will see this as 0 and initialize it to the current time when Trim is called. |
| | | 413 | | _millisecondsTimestamp = 0; |
| | | 414 | | } |
| | | 415 | | |
| | | 416 | | Unsafe.Add(ref MemoryMarshal.GetArrayDataReference(arrays), count) = array; // arrays[count] = array |
| | | 417 | | _count = count + 1; |
| | | 418 | | enqueued = true; |
| | | 419 | | } |
| | | 420 | | Monitor.Exit(this); |
| | | 421 | | return enqueued; |
| | | 422 | | } |
| | | 423 | | |
| | | 424 | | [MethodImpl(MethodImplOptions.AggressiveInlining)] |
| | | 425 | | public Array? TryPop() |
| | | 426 | | { |
| | | 427 | | Array? arr = null; |
| | | 428 | | Monitor.Enter(this); |
| | | 429 | | Array?[] arrays = _arrays; |
| | | 430 | | int count = _count - 1; |
| | | 431 | | if ((uint)count < (uint)arrays.Length) |
| | | 432 | | { |
| | | 433 | | arr = arrays[count]; |
| | | 434 | | arrays[count] = null; |
| | | 435 | | _count = count; |
| | | 436 | | } |
| | | 437 | | Monitor.Exit(this); |
| | | 438 | | return arr; |
| | | 439 | | } |
| | | 440 | | |
| | | 441 | | public void Trim(int currentMilliseconds, int id, Utilities.MemoryPressure pressure) |
| | | 442 | | { |
| | | 443 | | const int TrimAfterMS = 60 * 1000; // Trim after 60 seconds for low/mod |
| | | 444 | | const int HighTrimAfterMS = 10 * 1000; // Trim after 10 seconds for high pr |
| | | 445 | | |
| | | 446 | | if (_count == 0) |
| | | 447 | | { |
| | | 448 | | return; |
| | | 449 | | } |
| | | 450 | | |
| | | 451 | | int trimMilliseconds = pressure == Utilities.MemoryPressure.High ? HighTrimAfterMS : TrimAfterMS; |
| | | 452 | | |
| | | 453 | | lock (this) |
| | | 454 | | { |
| | | 455 | | if (_count == 0) |
| | | 456 | | { |
| | | 457 | | return; |
| | | 458 | | } |
| | | 459 | | |
| | | 460 | | if (_millisecondsTimestamp == 0) |
| | | 461 | | { |
| | | 462 | | _millisecondsTimestamp = currentMilliseconds; |
| | | 463 | | return; |
| | | 464 | | } |
| | | 465 | | |
| | | 466 | | if ((currentMilliseconds - _millisecondsTimestamp) <= trimMilliseconds) |
| | | 467 | | { |
| | | 468 | | return; |
| | | 469 | | } |
| | | 470 | | |
| | | 471 | | // We've elapsed enough time since the first item went into the partition. |
| | | 472 | | // Drop the top item(s) so they can be collected. |
| | | 473 | | |
| | | 474 | | int trimCount = pressure switch |
| | | 475 | | { |
| | | 476 | | Utilities.MemoryPressure.High => SharedArrayPoolStatics.s_maxArraysPerPartition, |
| | | 477 | | Utilities.MemoryPressure.Medium => 2, |
| | | 478 | | _ => 1, |
| | | 479 | | }; |
| | | 480 | | |
| | | 481 | | ArrayPoolEventSource log = ArrayPoolEventSource.Log; |
| | | 482 | | while (_count > 0 && trimCount-- > 0) |
| | | 483 | | { |
| | | 484 | | Array? array = _arrays[--_count]; |
| | | 485 | | Debug.Assert(array is not null, "No nulls should have been present in slots < _count."); |
| | | 486 | | _arrays[_count] = null; |
| | | 487 | | |
| | | 488 | | if (log.IsEnabled()) |
| | | 489 | | { |
| | | 490 | | log.BufferTrimmed(array.GetHashCode(), array.Length, id); |
| | | 491 | | } |
| | | 492 | | } |
| | | 493 | | |
| | | 494 | | _millisecondsTimestamp = _count > 0 ? |
| | | 495 | | _millisecondsTimestamp + (trimMilliseconds / 4) : // Give the remaining items a bit more time |
| | | 496 | | 0; |
| | | 497 | | } |
| | | 498 | | } |
| | | 499 | | } |
| | | 500 | | } |
| | | 501 | | |
| | | 502 | | internal static class SharedArrayPoolStatics |
| | | 503 | | { |
| | | 504 | | /// <summary>Number of partitions to employ.</summary> |
| | | 505 | | internal static readonly int s_partitionCount = GetPartitionCount(); |
| | | 506 | | /// <summary>The maximum number of arrays per array size to store per partition.</summary> |
| | | 507 | | internal static readonly int s_maxArraysPerPartition = GetMaxArraysPerPartition(); |
| | | 508 | | |
| | | 509 | | /// <summary>Gets the maximum number of partitions to shard arrays into.</summary> |
| | | 510 | | /// <remarks>Defaults to int.MaxValue. Whatever value is returned will end up being clamped to <see cref="Envir |
| | | 511 | | private static int GetPartitionCount() |
| | | 512 | | { |
| | | 513 | | int partitionCount = TryGetInt32EnvironmentVariable("DOTNET_SYSTEM_BUFFERS_SHAREDARRAYPOOL_MAXPARTITIONCOUNT |
| | | 514 | | result : |
| | | 515 | | int.MaxValue; // no limit other than processor count |
| | | 516 | | return Math.Min(partitionCount, Environment.ProcessorCount); |
| | | 517 | | } |
| | | 518 | | |
| | | 519 | | /// <summary>Gets the maximum number of arrays of a given size allowed to be cached per partition.</summary> |
| | | 520 | | /// <returns>Defaults to 32. This does not factor in or impact the number of arrays cached per thread in TLS (cu |
| | | 521 | | private static int GetMaxArraysPerPartition() |
| | | 522 | | { |
| | | 523 | | return TryGetInt32EnvironmentVariable("DOTNET_SYSTEM_BUFFERS_SHAREDARRAYPOOL_MAXARRAYSPERPARTITION", out int |
| | | 524 | | result : |
| | | 525 | | 32; // arbitrary limit |
| | | 526 | | } |
| | | 527 | | |
| | | 528 | | /// <summary>Look up an environment variable and try to parse it as an Int32.</summary> |
| | | 529 | | /// <remarks>This avoids using anything that might in turn recursively use the ArrayPool.</remarks> |
| | | 530 | | private static bool TryGetInt32EnvironmentVariable(string variable, out int result) |
| | | 531 | | { |
| | | 532 | | // Avoid globalization stack, as it might in turn be using ArrayPool. |
| | | 533 | | |
| | | 534 | | if (Environment.GetEnvironmentVariableCore_NoArrayPool(variable) is string envVar && |
| | | 535 | | envVar.Length is > 0 and <= 32) // arbitrary limit that allows for some spaces around the maximum length |
| | | 536 | | { |
| | | 537 | | ReadOnlySpan<char> value = envVar.AsSpan().Trim(' '); |
| | | 538 | | if (!value.IsEmpty && value.Length <= 10) |
| | | 539 | | { |
| | | 540 | | long tempResult = 0; |
| | | 541 | | foreach (char c in value) |
| | | 542 | | { |
| | | 543 | | uint digit = (uint)(c - '0'); |
| | | 544 | | if (digit > 9) |
| | | 545 | | { |
| | | 546 | | goto Fail; |
| | | 547 | | } |
| | | 548 | | |
| | | 549 | | tempResult = tempResult * 10 + digit; |
| | | 550 | | } |
| | | 551 | | |
| | | 552 | | if (tempResult is >= 0 and <= int.MaxValue) |
| | | 553 | | { |
| | | 554 | | result = (int)tempResult; |
| | | 555 | | return true; |
| | | 556 | | } |
| | | 557 | | } |
| | | 558 | | } |
| | | 559 | | |
| | | 560 | | Fail: |
| | | 561 | | result = 0; |
| | | 562 | | return false; |
| | | 563 | | } |
| | | 564 | | } |
| | | 565 | | } |
| | | 566 | | |