feat(network)!: add v2 framing, pooled buffers and bounded dispatch
This commit is contained in:
@@ -54,7 +54,7 @@ namespace ShrinkNetwork
|
||||
Serializer = serializer ?? throw new ArgumentNullException(nameof(serializer));
|
||||
MessageRegistry = messageRegistry ?? throw new ArgumentNullException(nameof(messageRegistry));
|
||||
Router = router ?? throw new ArgumentNullException(nameof(router));
|
||||
_dispatchScheduler = dispatchScheduler ?? ShrinkNetworkDispatchSchedulers.Inline;
|
||||
_dispatchScheduler = dispatchScheduler ?? new ShrinkNetworkInlineDispatchScheduler();
|
||||
}
|
||||
|
||||
public IShrinkNetworkSerializer Serializer { get; }
|
||||
@@ -74,6 +74,8 @@ namespace ShrinkNetwork
|
||||
set => _dispatchScheduler = value ?? throw new ArgumentNullException(nameof(value));
|
||||
}
|
||||
|
||||
/// <summary>Reliable receive work could not be processed. Session-control transports also disconnect the peer.</summary>
|
||||
public event Action<long>? OnDispatchRejected;
|
||||
public event Action<ShrinkNetworkSession>? OnSessionConnected;
|
||||
public event Action<ShrinkNetworkSession>? OnSessionDisconnected;
|
||||
|
||||
@@ -273,11 +275,11 @@ namespace ShrinkNetwork
|
||||
if (messageType == null)
|
||||
throw new ArgumentNullException(nameof(messageType));
|
||||
|
||||
return SendPacketAsync(session, messageType, kind, requestToken, route, Serializer.Serialize(message));
|
||||
return SendPacketAsync(session, messageType, kind, requestToken, route, Array.Empty<byte>(), message);
|
||||
}
|
||||
|
||||
private async UniTask SendPacketAsync(ShrinkNetworkSession session, Type messageType,
|
||||
ShrinkNetworkPacketKind kind, ShrinkRequestToken requestToken, string? route, byte[] payload)
|
||||
ShrinkNetworkPacketKind kind, ShrinkRequestToken requestToken, string? route, byte[] payload, object? message = null)
|
||||
{
|
||||
if (session == null)
|
||||
throw new ArgumentNullException(nameof(session));
|
||||
@@ -299,16 +301,22 @@ namespace ShrinkNetwork
|
||||
Payload = payload
|
||||
};
|
||||
|
||||
var packetData = Serializer.Serialize(packet);
|
||||
using var encoded = ShrinkPacketCodec.Encode(packet, Serializer, message);
|
||||
var packetData = encoded.WrittenMemory;
|
||||
Interlocked.Increment(ref _packetsSent);
|
||||
Interlocked.Add(ref _bytesSent, packetData.Length);
|
||||
if (transport is IShrinkNetworkMemoryTransport memoryTransport)
|
||||
{
|
||||
await memoryTransport.SendAsync(session.SessionId, packetData);
|
||||
return;
|
||||
}
|
||||
if (transport is IShrinkNetworkAsyncTransport asyncTransport)
|
||||
{
|
||||
await asyncTransport.SendAsync(session.SessionId, packetData);
|
||||
await asyncTransport.SendAsync(session.SessionId, packetData.ToArray());
|
||||
return;
|
||||
}
|
||||
|
||||
transport.Send(session.SessionId, packetData);
|
||||
transport.Send(session.SessionId, packetData.ToArray());
|
||||
}
|
||||
|
||||
private void OnTransportEvent(ShrinkNetworkTransportEvent evt)
|
||||
@@ -330,6 +338,12 @@ namespace ShrinkNetwork
|
||||
return;
|
||||
}
|
||||
|
||||
// Responses complete pending RPCs independently of the serial handler queue, including nested calls.
|
||||
if (evt.Type == ShrinkNetworkTransportEventType.Packet && ShrinkPacketCodec.IsResponse(evt.PacketData))
|
||||
{
|
||||
HandlePacketAsync(evt).Forget();
|
||||
return;
|
||||
}
|
||||
ScheduleTransportEventAsync(evt).Forget();
|
||||
}
|
||||
|
||||
@@ -337,9 +351,23 @@ namespace ShrinkNetwork
|
||||
{
|
||||
try
|
||||
{
|
||||
var scheduled = await DispatchScheduler.ScheduleAsync(() => HandleTransportEventAsync(evt));
|
||||
var scheduler = DispatchScheduler;
|
||||
_sessions.TryGetValue(evt.SessionId, out var receivedSession);
|
||||
UniTask DispatchCurrent()
|
||||
{
|
||||
return receivedSession != null && _sessions.TryGetValue(evt.SessionId, out var current) && ReferenceEquals(current, receivedSession)
|
||||
? HandleTransportEventAsync(evt) : UniTask.CompletedTask;
|
||||
}
|
||||
var scheduled = scheduler is IShrinkNetworkPacketDispatchScheduler queue
|
||||
? await queue.ScheduleAsync(evt.SessionId, evt.PacketData?.Length ?? 0, DispatchCurrent)
|
||||
: await scheduler.ScheduleAsync(DispatchCurrent);
|
||||
if (!scheduled)
|
||||
{
|
||||
Interlocked.Increment(ref _dispatchQueueRejectedCount);
|
||||
OnDispatchRejected?.Invoke(evt.SessionId);
|
||||
if (_transport is IShrinkNetworkSessionControlTransport control)
|
||||
control.DisconnectSession(evt.SessionId, "SHRINK-NET-CONGESTION: reliable receive queue rejected work.");
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
@@ -395,12 +423,12 @@ namespace ShrinkNetwork
|
||||
{
|
||||
Interlocked.Increment(ref _packetsReceived);
|
||||
Interlocked.Add(ref _bytesReceived, evt.PacketData.Length);
|
||||
var packet = Serializer.Deserialize<ShrinkNetworkPacket>(evt.PacketData);
|
||||
var packet = ShrinkPacketCodec.Decode(evt.PacketData);
|
||||
|
||||
if (!_sessions.TryGetValue(evt.SessionId, out var session))
|
||||
{
|
||||
session = _sessions.GetOrAdd(evt.SessionId,
|
||||
id => new ShrinkNetworkSession(id, evt.RemoteAddress, this));
|
||||
// Queued packets from a disconnected session cannot resurrect its state.
|
||||
return;
|
||||
}
|
||||
|
||||
if (!ValidatePacketCompatibility(packet, evt.SessionId))
|
||||
@@ -410,7 +438,7 @@ namespace ShrinkNetwork
|
||||
|
||||
if (packet.Kind == ShrinkNetworkPacketKind.Response)
|
||||
{
|
||||
HandleResponse(packet);
|
||||
HandleResponse(session, packet);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -425,7 +453,7 @@ namespace ShrinkNetwork
|
||||
}
|
||||
|
||||
var resolvedMeta = meta!;
|
||||
var message = Serializer.Deserialize(packet.Payload, resolvedMeta.MessageType);
|
||||
var message = DeserializePayload(packet.Payload, resolvedMeta.MessageType);
|
||||
if (message == null)
|
||||
{
|
||||
ShrinkNetworkLogger.Warn($"[ShrinkNetwork] Failed to deserialize message for opcode {packet.Opcode}.");
|
||||
@@ -442,6 +470,12 @@ namespace ShrinkNetwork
|
||||
ShrinkNetworkLogger.Warn($"[ShrinkNetwork] No handler found for {resolvedMeta.MessageType.FullName}");
|
||||
}
|
||||
}
|
||||
catch (ShrinkProtocolException ex)
|
||||
{
|
||||
Interlocked.Increment(ref _protocolViolations);
|
||||
if (DisconnectOnProtocolViolation && _transport is IShrinkNetworkSessionControlTransport control)
|
||||
control.DisconnectSession(evt.SessionId, ex.Message);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
Interlocked.Increment(ref _serializationErrorCount);
|
||||
@@ -450,9 +484,10 @@ namespace ShrinkNetwork
|
||||
}
|
||||
}
|
||||
|
||||
private void HandleResponse(ShrinkNetworkPacket packet)
|
||||
private void HandleResponse(ShrinkNetworkSession session, ShrinkNetworkPacket packet)
|
||||
{
|
||||
if (!_pendingRequests.TryRemove(packet.RequestToken, out var pending))
|
||||
if (!_pendingRequests.TryGetValue(packet.RequestToken, out var expected) || expected.SessionId != session.SessionId ||
|
||||
!_pendingRequests.TryRemove(packet.RequestToken, out var pending))
|
||||
{
|
||||
ShrinkNetworkLogger.Warn($"[ShrinkNetwork] Pending request not found. RequestToken={packet.RequestToken}");
|
||||
return;
|
||||
@@ -460,7 +495,7 @@ namespace ShrinkNetwork
|
||||
|
||||
try
|
||||
{
|
||||
var response = Serializer.Deserialize(packet.Payload, pending.ResponseType);
|
||||
var response = DeserializePayload(packet.Payload, pending.ResponseType);
|
||||
pending.CompletionSource.TrySetResult(response);
|
||||
}
|
||||
catch (Exception ex)
|
||||
@@ -469,6 +504,10 @@ namespace ShrinkNetwork
|
||||
}
|
||||
}
|
||||
|
||||
private object DeserializePayload(ReadOnlyMemory<byte> payload, Type type) =>
|
||||
Serializer is IShrinkNetworkBufferSerializer buffered
|
||||
? buffered.Deserialize(payload, type) : Serializer.Deserialize(payload.ToArray(), type);
|
||||
|
||||
private async UniTask<TResponse> WaitForPendingResponse<TResponse>(ShrinkRequestToken requestToken, PendingRequest pending,
|
||||
ShrinkRpcCallOptions? options)
|
||||
where TResponse : class, IShrinkNetworkResponse
|
||||
@@ -627,20 +666,42 @@ namespace ShrinkNetwork
|
||||
|
||||
public static class ShrinkNetworkDispatchSchedulers
|
||||
{
|
||||
public static IShrinkNetworkDispatchScheduler Inline { get; } =
|
||||
new ShrinkNetworkInlineDispatchScheduler();
|
||||
public static IShrinkNetworkDispatchScheduler Inline => new ShrinkNetworkInlineDispatchScheduler();
|
||||
}
|
||||
|
||||
public sealed class ShrinkNetworkInlineDispatchScheduler : IShrinkNetworkDispatchScheduler
|
||||
public interface IShrinkNetworkPacketDispatchScheduler : IShrinkNetworkDispatchScheduler
|
||||
{
|
||||
public async UniTask<bool> ScheduleAsync(Func<UniTask> callback)
|
||||
{
|
||||
if (callback == null)
|
||||
throw new ArgumentNullException(nameof(callback));
|
||||
UniTask<bool> ScheduleAsync(long sessionId, int bytes, Func<UniTask> callback);
|
||||
}
|
||||
|
||||
await callback();
|
||||
return true;
|
||||
/// <summary>Automatically drains a bounded serial queue. Callbacks start on the initiating transport/continuation thread.</summary>
|
||||
public sealed class ShrinkNetworkInlineDispatchScheduler : IShrinkNetworkPacketDispatchScheduler, IDisposable
|
||||
{
|
||||
private readonly object _gate = new();
|
||||
private readonly ShrinkNetworkWorkQueue _queue = new(4096, 64 * 1024 * 1024);
|
||||
private bool _draining;
|
||||
public ShrinkNetworkQueueDiagnostics CaptureDiagnostics() => _queue.CaptureDiagnostics();
|
||||
public UniTask<bool> ScheduleAsync(Func<UniTask> callback) => ScheduleAsync(0, 0, callback);
|
||||
public async UniTask<bool> ScheduleAsync(long sessionId, int bytes, Func<UniTask> callback)
|
||||
{
|
||||
var completion = _queue.EnqueueAsync(sessionId, "receive", bytes, callback);
|
||||
var start = false;
|
||||
lock (_gate) { if (!_draining) { _draining = true; start = true; } }
|
||||
if (start) DrainAsync().Forget();
|
||||
return await completion == ShrinkNetworkQueueResult.Completed;
|
||||
}
|
||||
private async UniTask DrainAsync()
|
||||
{
|
||||
while (true)
|
||||
{
|
||||
await _queue.PumpAsync(int.MaxValue);
|
||||
lock (_gate)
|
||||
{
|
||||
if (_queue.PendingCount == 0) { _draining = false; return; }
|
||||
}
|
||||
}
|
||||
}
|
||||
public void Dispose() => _queue.Dispose();
|
||||
}
|
||||
|
||||
public enum ShrinkNetworkDispatchOverflowPolicy
|
||||
@@ -650,163 +711,30 @@ namespace ShrinkNetwork
|
||||
DropOldest = 2
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// A caller-pumped, bounded dispatch queue. Unity can pump it from Update
|
||||
/// while a dedicated server can keep the default inline scheduler.
|
||||
/// </summary>
|
||||
public sealed class ShrinkNetworkDispatchQueue : IShrinkNetworkDispatchScheduler, IDisposable
|
||||
/// <summary>Caller-pumped serial receive queue; rejects overflow without silently dropping reliable packets.</summary>
|
||||
public sealed class ShrinkNetworkDispatchQueue : IShrinkNetworkPacketDispatchScheduler, IDisposable
|
||||
{
|
||||
private sealed class WorkItem
|
||||
{
|
||||
public Func<UniTask> Callback = null!;
|
||||
public UniTaskCompletionSource<bool> Completion = null!;
|
||||
}
|
||||
|
||||
private readonly ConcurrentQueue<WorkItem> _queue = new();
|
||||
private readonly object _lifecycleLock = new();
|
||||
private readonly int _capacity;
|
||||
private readonly ShrinkNetworkDispatchOverflowPolicy _overflowPolicy;
|
||||
private int _queuedCount;
|
||||
private int _pumping;
|
||||
private int _disposed;
|
||||
private long _rejectedCount;
|
||||
private long _droppedCount;
|
||||
|
||||
private readonly ShrinkNetworkWorkQueue _queue;
|
||||
public ShrinkNetworkDispatchQueue(int capacity,
|
||||
ShrinkNetworkDispatchOverflowPolicy overflowPolicy = ShrinkNetworkDispatchOverflowPolicy.Reject)
|
||||
ShrinkNetworkDispatchOverflowPolicy overflowPolicy = ShrinkNetworkDispatchOverflowPolicy.Reject,
|
||||
long byteCapacity = 64 * 1024 * 1024, int perSessionCapacity = int.MaxValue)
|
||||
{
|
||||
if (capacity <= 0)
|
||||
throw new ArgumentOutOfRangeException(nameof(capacity));
|
||||
|
||||
_capacity = capacity;
|
||||
_overflowPolicy = overflowPolicy;
|
||||
}
|
||||
|
||||
public int Capacity => _capacity;
|
||||
public int PendingCount => Volatile.Read(ref _queuedCount);
|
||||
public long RejectedCount => Volatile.Read(ref _rejectedCount);
|
||||
public long DroppedCount => Volatile.Read(ref _droppedCount);
|
||||
|
||||
public UniTask<bool> ScheduleAsync(Func<UniTask> callback)
|
||||
{
|
||||
if (callback == null)
|
||||
throw new ArgumentNullException(nameof(callback));
|
||||
if (Volatile.Read(ref _disposed) != 0)
|
||||
return UniTask.FromException<bool>(new ObjectDisposedException(nameof(ShrinkNetworkDispatchQueue)));
|
||||
|
||||
var item = new WorkItem
|
||||
{
|
||||
Callback = callback,
|
||||
Completion = new UniTaskCompletionSource<bool>()
|
||||
};
|
||||
|
||||
while (true)
|
||||
{
|
||||
if (Volatile.Read(ref _queuedCount) >= _capacity)
|
||||
{
|
||||
switch (_overflowPolicy)
|
||||
{
|
||||
case ShrinkNetworkDispatchOverflowPolicy.Reject:
|
||||
Interlocked.Increment(ref _rejectedCount);
|
||||
return UniTask.FromResult(false);
|
||||
case ShrinkNetworkDispatchOverflowPolicy.DropNewest:
|
||||
Interlocked.Increment(ref _droppedCount);
|
||||
return UniTask.FromResult(false);
|
||||
case ShrinkNetworkDispatchOverflowPolicy.DropOldest:
|
||||
if (_queue.TryDequeue(out var dropped))
|
||||
{
|
||||
Interlocked.Decrement(ref _queuedCount);
|
||||
Interlocked.Increment(ref _droppedCount);
|
||||
dropped.Completion.TrySetResult(false);
|
||||
continue;
|
||||
}
|
||||
|
||||
Thread.Yield();
|
||||
continue;
|
||||
default:
|
||||
throw new ArgumentOutOfRangeException();
|
||||
}
|
||||
}
|
||||
|
||||
var currentCount = Volatile.Read(ref _queuedCount);
|
||||
if (currentCount >= _capacity ||
|
||||
Interlocked.CompareExchange(ref _queuedCount, currentCount + 1, currentCount) != currentCount)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
lock (_lifecycleLock)
|
||||
{
|
||||
if (Volatile.Read(ref _disposed) != 0)
|
||||
{
|
||||
Interlocked.Decrement(ref _queuedCount);
|
||||
item.Completion.TrySetResult(false);
|
||||
return item.Completion.Task;
|
||||
}
|
||||
|
||||
_queue.Enqueue(item);
|
||||
return item.Completion.Task;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public UniTask<int> PumpAsync(int maxItems)
|
||||
{
|
||||
if (maxItems <= 0)
|
||||
throw new ArgumentOutOfRangeException(nameof(maxItems));
|
||||
if (Interlocked.Exchange(ref _pumping, 1) == 1)
|
||||
return UniTask.FromResult(0);
|
||||
|
||||
return PumpCoreAsync(maxItems);
|
||||
}
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
lock (_lifecycleLock)
|
||||
{
|
||||
if (Interlocked.Exchange(ref _disposed, 1) != 0)
|
||||
return;
|
||||
|
||||
while (_queue.TryDequeue(out var item))
|
||||
{
|
||||
Interlocked.Decrement(ref _queuedCount);
|
||||
item.Completion.TrySetResult(false);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private async UniTask<int> PumpCoreAsync(int maxItems)
|
||||
{
|
||||
var processed = 0;
|
||||
try
|
||||
{
|
||||
while (processed < maxItems && _queue.TryDequeue(out var item))
|
||||
{
|
||||
Interlocked.Decrement(ref _queuedCount);
|
||||
await ExecuteItemAsync(item);
|
||||
processed++;
|
||||
}
|
||||
|
||||
return processed;
|
||||
}
|
||||
finally
|
||||
{
|
||||
Volatile.Write(ref _pumping, 0);
|
||||
}
|
||||
}
|
||||
|
||||
private static async UniTask ExecuteItemAsync(WorkItem item)
|
||||
{
|
||||
try
|
||||
{
|
||||
await item.Callback();
|
||||
item.Completion.TrySetResult(true);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
item.Completion.TrySetException(ex);
|
||||
}
|
||||
if (overflowPolicy != ShrinkNetworkDispatchOverflowPolicy.Reject)
|
||||
throw new ArgumentException("Protocol v2 reliable dispatch requires Reject. Use an explicit state key on ShrinkNetworkWorkQueue for replaceable state.", nameof(overflowPolicy));
|
||||
Capacity = capacity;
|
||||
_queue = new ShrinkNetworkWorkQueue(capacity, byteCapacity, perSessionCapacity);
|
||||
}
|
||||
public int Capacity { get; }
|
||||
public int PendingCount => _queue.PendingCount;
|
||||
public long RejectedCount => CaptureDiagnostics().Rejected;
|
||||
public long DroppedCount => 0;
|
||||
public ShrinkNetworkQueueDiagnostics CaptureDiagnostics() => _queue.CaptureDiagnostics();
|
||||
public UniTask<bool> ScheduleAsync(Func<UniTask> callback) => ScheduleAsync(0, 0, callback);
|
||||
public async UniTask<bool> ScheduleAsync(long sessionId, int bytes, Func<UniTask> callback) =>
|
||||
await _queue.EnqueueAsync(sessionId, "receive", bytes, callback) == ShrinkNetworkQueueResult.Completed;
|
||||
public UniTask<int> PumpAsync(int maxItems) => _queue.PumpAsync(maxItems);
|
||||
public UniTask<int> PumpAsync(int maxItems, long maxBytes, TimeSpan timeBudget) => _queue.PumpAsync(maxItems, maxBytes, timeBudget);
|
||||
public void Dispose() => _queue.Dispose();
|
||||
}
|
||||
|
||||
public sealed class ShrinkNetworkServiceDiagnosticsSnapshot
|
||||
|
||||
@@ -0,0 +1,159 @@
|
||||
#nullable enable
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Diagnostics;
|
||||
using System.Threading;
|
||||
using Cysharp.Threading.Tasks;
|
||||
|
||||
namespace ShrinkNetwork
|
||||
{
|
||||
public enum ShrinkNetworkQueueResult { Completed, Rejected, Replaced, Canceled }
|
||||
|
||||
public sealed class ShrinkNetworkQueueDiagnostics
|
||||
{
|
||||
public int PendingCount { get; internal set; }
|
||||
public long PendingBytes { get; internal set; }
|
||||
public long Rejected { get; internal set; }
|
||||
public long Replaced { get; internal set; }
|
||||
public long Completed { get; internal set; }
|
||||
public double OldestWaitMilliseconds { get; internal set; }
|
||||
public double LastWaitMilliseconds { get; internal set; }
|
||||
}
|
||||
|
||||
/// <summary>Caller-pumped serial execution, fair across session/channel partitions. Only queued state is replaceable.</summary>
|
||||
public sealed class ShrinkNetworkWorkQueue : IDisposable
|
||||
{
|
||||
private sealed class Work
|
||||
{
|
||||
public Func<UniTask> Callback = null!;
|
||||
public UniTaskCompletionSource<ShrinkNetworkQueueResult> Completion = new();
|
||||
public string? StateKey;
|
||||
public int Bytes;
|
||||
public long Enqueued = Stopwatch.GetTimestamp();
|
||||
public CancellationToken Cancellation;
|
||||
}
|
||||
private sealed class Partition
|
||||
{
|
||||
public readonly LinkedList<Work> Items = new();
|
||||
public readonly Dictionary<string, LinkedListNode<Work>> States = new(StringComparer.Ordinal);
|
||||
}
|
||||
private readonly object _gate = new();
|
||||
private readonly Dictionary<(long, string), Partition> _partitions = new();
|
||||
private readonly Queue<(long, string)> _ready = new();
|
||||
private readonly int _capacity;
|
||||
private readonly long _byteCapacity;
|
||||
private readonly int _partitionCapacity;
|
||||
private int _count, _pumping;
|
||||
private long _bytes, _rejected, _replaced, _completed;
|
||||
private double _lastWait;
|
||||
private bool _disposed;
|
||||
|
||||
public int PendingCount { get { lock (_gate) return _count; } }
|
||||
|
||||
public ShrinkNetworkWorkQueue(int capacity, long byteCapacity, int perPartitionCapacity = int.MaxValue)
|
||||
{
|
||||
if (capacity <= 0 || byteCapacity <= 0 || perPartitionCapacity <= 0) throw new ArgumentOutOfRangeException(nameof(capacity));
|
||||
_capacity = capacity; _byteCapacity = byteCapacity; _partitionCapacity = perPartitionCapacity;
|
||||
}
|
||||
|
||||
/// <summary>A nonempty stateKey explicitly permits replacing an unsent state in this session/channel. Never use for RPC or snapshot fragments.</summary>
|
||||
public UniTask<ShrinkNetworkQueueResult> EnqueueAsync(long sessionId, string channel, int byteCount,
|
||||
Func<UniTask> callback, string? stateKey = null, CancellationToken cancellationToken = default)
|
||||
{
|
||||
if (callback == null) throw new ArgumentNullException(nameof(callback));
|
||||
if (channel == null) throw new ArgumentNullException(nameof(channel));
|
||||
if (byteCount < 0) throw new ArgumentOutOfRangeException(nameof(byteCount));
|
||||
Work? replaced = null;
|
||||
Work item;
|
||||
lock (_gate)
|
||||
{
|
||||
if (_disposed || cancellationToken.IsCancellationRequested) return UniTask.FromResult(ShrinkNetworkQueueResult.Canceled);
|
||||
var key = (sessionId, channel);
|
||||
_partitions.TryGetValue(key, out var partition);
|
||||
LinkedListNode<Work>? old = null;
|
||||
if (!string.IsNullOrEmpty(stateKey)) partition?.States.TryGetValue(stateKey!, out old);
|
||||
var nextBytes = _bytes - (old?.Value.Bytes ?? 0) + byteCount;
|
||||
if (nextBytes > _byteCapacity || (old == null && (_count >= _capacity || (partition?.Items.Count ?? 0) >= _partitionCapacity)))
|
||||
{ _rejected++; return UniTask.FromResult(ShrinkNetworkQueueResult.Rejected); }
|
||||
item = new Work { Callback = callback, Bytes = byteCount, StateKey = string.IsNullOrEmpty(stateKey) ? null : stateKey, Cancellation = cancellationToken };
|
||||
if (partition == null) { partition = new Partition(); _partitions.Add(key, partition); _ready.Enqueue(key); }
|
||||
if (old != null)
|
||||
{
|
||||
replaced = old.Value;
|
||||
// Move a replacement to the tail: later state must not jump ahead of intervening reliable operations.
|
||||
partition.Items.Remove(old); _replaced++;
|
||||
}
|
||||
else _count++;
|
||||
var node = partition.Items.AddLast(item);
|
||||
if (item.StateKey != null) partition.States[item.StateKey] = node;
|
||||
_bytes = nextBytes;
|
||||
}
|
||||
replaced?.Completion.TrySetResult(ShrinkNetworkQueueResult.Replaced);
|
||||
return item.Completion.Task;
|
||||
}
|
||||
|
||||
public async UniTask<int> PumpAsync(int maxItems, long maxBytes = long.MaxValue, TimeSpan? timeBudget = null)
|
||||
{
|
||||
if (maxItems <= 0 || maxBytes <= 0 || (timeBudget.HasValue && timeBudget.Value <= TimeSpan.Zero)) throw new ArgumentOutOfRangeException(nameof(maxItems));
|
||||
if (Interlocked.Exchange(ref _pumping, 1) != 0) return 0;
|
||||
var started = Stopwatch.GetTimestamp();
|
||||
var processed = 0;
|
||||
long bytes = 0;
|
||||
try
|
||||
{
|
||||
while (processed < maxItems && (!timeBudget.HasValue || Elapsed(started) < timeBudget.Value.TotalMilliseconds))
|
||||
{
|
||||
Work item;
|
||||
lock (_gate)
|
||||
{
|
||||
if (_disposed || _ready.Count == 0) break;
|
||||
var key = _ready.Peek();
|
||||
var partition = _partitions[key];
|
||||
item = partition.Items.First!.Value;
|
||||
// Allow one oversized item so a byte budget cannot permanently starve a valid packet.
|
||||
if (processed > 0 && item.Bytes > maxBytes - bytes) break;
|
||||
_ready.Dequeue(); partition.Items.RemoveFirst();
|
||||
if (item.StateKey != null) partition.States.Remove(item.StateKey);
|
||||
if (partition.Items.Count == 0) _partitions.Remove(key); else _ready.Enqueue(key);
|
||||
_count--; _bytes -= item.Bytes; _lastWait = Elapsed(item.Enqueued);
|
||||
}
|
||||
try
|
||||
{
|
||||
if (item.Cancellation.IsCancellationRequested) item.Completion.TrySetResult(ShrinkNetworkQueueResult.Canceled);
|
||||
else { await item.Callback(); item.Completion.TrySetResult(ShrinkNetworkQueueResult.Completed); lock (_gate) _completed++; }
|
||||
}
|
||||
catch (OperationCanceledException) { item.Completion.TrySetResult(ShrinkNetworkQueueResult.Canceled); }
|
||||
catch (Exception ex) { item.Completion.TrySetException(ex); }
|
||||
processed++; bytes += item.Bytes;
|
||||
}
|
||||
return processed;
|
||||
}
|
||||
finally { Volatile.Write(ref _pumping, 0); }
|
||||
}
|
||||
|
||||
public ShrinkNetworkQueueDiagnostics CaptureDiagnostics()
|
||||
{
|
||||
lock (_gate)
|
||||
{
|
||||
long oldest = Stopwatch.GetTimestamp();
|
||||
foreach (var partition in _partitions.Values)
|
||||
if (partition.Items.First != null) oldest = Math.Min(oldest, partition.Items.First.Value.Enqueued);
|
||||
return new ShrinkNetworkQueueDiagnostics { PendingCount = _count, PendingBytes = _bytes, Rejected = _rejected, Replaced = _replaced,
|
||||
Completed = _completed, LastWaitMilliseconds = _lastWait, OldestWaitMilliseconds = _count == 0 ? 0 : Elapsed(oldest) };
|
||||
}
|
||||
}
|
||||
private static double Elapsed(long start) => (Stopwatch.GetTimestamp() - start) * 1000d / Stopwatch.Frequency;
|
||||
public void Dispose()
|
||||
{
|
||||
List<Work> canceled = new();
|
||||
lock (_gate)
|
||||
{
|
||||
if (_disposed) return;
|
||||
_disposed = true;
|
||||
foreach (var partition in _partitions.Values) canceled.AddRange(partition.Items);
|
||||
_partitions.Clear(); _ready.Clear(); _count = 0; _bytes = 0;
|
||||
}
|
||||
foreach (var item in canceled) item.Completion.TrySetResult(ShrinkNetworkQueueResult.Canceled);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 3b35a44b897396549b5a1dec58ae05c0
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
Reference in New Issue
Block a user