调整嵌套类位置 + 制订阶段 5 task

This commit is contained in:
2026-03-27 16:37:43 +08:00
parent f21a4c51d2
commit ca26ab8e38
11 changed files with 446 additions and 213 deletions
@@ -0,0 +1,221 @@
using System;
using System.Net;
using System.Runtime.InteropServices;
using kcp;
namespace Network.NetworkTransport
{
public partial class KcpTransport
{
private sealed unsafe class KcpSession : IDisposable
{
private readonly KcpTransport _owner;
private readonly object _gate = new();
private readonly GCHandle _handle;
private IKCPCB* _kcp;
private bool _disposed;
private uint _nextUpdateAt;
public KcpSession(KcpTransport owner, IPEndPoint remoteEndPoint, uint conv)
{
_owner = owner ?? throw new ArgumentNullException(nameof(owner));
RemoteEndPoint = remoteEndPoint ?? throw new ArgumentNullException(nameof(remoteEndPoint));
Conv = conv;
LastActivityUtc = DateTime.UtcNow;
_handle = GCHandle.Alloc(this);
_kcp = KCP.ikcp_create(conv, (void*)GCHandle.ToIntPtr(_handle));
KCP.ikcp_setoutput(_kcp, &OutputCallback);
KCP.ikcp_nodelay(_kcp, DefaultNoDelay, DefaultInterval, DefaultResend, DefaultNc);
KCP.ikcp_wndsize(_kcp, DefaultSendWindow, DefaultReceiveWindow);
KCP.ikcp_setmtu(_kcp, DefaultMtu);
_nextUpdateAt = GetCurrentTimeMilliseconds();
}
public uint Conv { get; }
public IPEndPoint RemoteEndPoint { get; }
public DateTime LastActivityUtc { get; private set; }
public void Send(byte[] payload)
{
if (payload == null)
{
throw new ArgumentNullException(nameof(payload));
}
lock (_gate)
{
ThrowIfDisposed();
if (payload.Length == 0)
{
return;
}
fixed (byte* buffer = payload)
{
var result = KCP.ikcp_send(_kcp, buffer, payload.Length);
if (result < 0)
{
throw new InvalidOperationException($"KCP send failed with error code {result}.");
}
}
LastActivityUtc = DateTime.UtcNow;
UpdateNoLock(GetCurrentTimeMilliseconds());
}
}
public void Input(byte[] datagram)
{
if (datagram == null)
{
throw new ArgumentNullException(nameof(datagram));
}
if (datagram.Length == 0)
{
return;
}
lock (_gate)
{
ThrowIfDisposed();
fixed (byte* buffer = datagram)
{
var result = KCP.ikcp_input(_kcp, buffer, datagram.Length);
if (result < 0)
{
Console.WriteLine($"[KcpTransport] KCP input failed for {RemoteEndPoint}: {result}");
return;
}
}
LastActivityUtc = DateTime.UtcNow;
UpdateNoLock(GetCurrentTimeMilliseconds());
}
}
public bool TryReceive(out byte[] payload)
{
lock (_gate)
{
if (_disposed)
{
payload = null;
return false;
}
var size = KCP.ikcp_peeksize(_kcp);
if (size <= 0)
{
payload = null;
return false;
}
payload = new byte[size];
fixed (byte* buffer = payload)
{
var result = KCP.ikcp_recv(_kcp, buffer, payload.Length);
if (result < 0)
{
payload = null;
return false;
}
if (result != payload.Length)
{
Array.Resize(ref payload, result);
}
}
LastActivityUtc = DateTime.UtcNow;
return true;
}
}
public void UpdateIfDue(uint current)
{
lock (_gate)
{
if (_disposed)
{
return;
}
if (KCP._itimediff(current, _nextUpdateAt) < 0)
{
return;
}
UpdateNoLock(current);
}
}
public void Dispose()
{
lock (_gate)
{
if (_disposed)
{
return;
}
_disposed = true;
if (_kcp != null)
{
KCP.ikcp_release(_kcp);
_kcp = null;
}
if (_handle.IsAllocated)
{
_handle.Free();
}
}
}
private void UpdateNoLock(uint current)
{
KCP.ikcp_update(_kcp, current);
_nextUpdateAt = KCP.ikcp_check(_kcp, current);
}
private void ThrowIfDisposed()
{
if (_disposed || _kcp == null)
{
throw new ObjectDisposedException(nameof(KcpSession));
}
}
private int SendRaw(byte* buffer, int length)
{
return _owner.SendDatagram(buffer, length, RemoteEndPoint);
}
private static int OutputCallback(byte* buffer, int length, IKCPCB* kcp, void* user)
{
if (user == null)
{
return -1;
}
var handle = GCHandle.FromIntPtr((IntPtr)user);
if (handle.Target is not KcpSession session)
{
return -1;
}
return session.SendRaw(buffer, length);
}
}
}
}
@@ -0,0 +1,3 @@
fileFormatVersion: 2
guid: 63fb533d620e4ac0bf15ba0a9a71331e
timeCreated: 1774593005
@@ -9,7 +9,7 @@ using kcp;
namespace Network.NetworkTransport
{
public class KcpTransport : ITransport
public partial class KcpTransport : ITransport
{
private const uint DefaultConv = 1;
private const int DefaultNoDelay = 1;
@@ -127,7 +127,8 @@ namespace Network.NetworkTransport
if (!_isServer && !target.Equals(_defaultRemoteEndPoint))
{
throw new InvalidOperationException("Client mode only supports the configured default remote endpoint.");
throw new InvalidOperationException(
"Client mode only supports the configured default remote endpoint.");
}
var session = GetOrCreateSession(target, _defaultConv);
@@ -318,215 +319,6 @@ namespace Network.NetworkTransport
}
}
private unsafe sealed class KcpSession : IDisposable
{
private readonly KcpTransport _owner;
private readonly object _gate = new();
private readonly GCHandle _handle;
private IKCPCB* _kcp;
private bool _disposed;
private uint _nextUpdateAt;
public KcpSession(KcpTransport owner, IPEndPoint remoteEndPoint, uint conv)
{
_owner = owner ?? throw new ArgumentNullException(nameof(owner));
RemoteEndPoint = remoteEndPoint ?? throw new ArgumentNullException(nameof(remoteEndPoint));
Conv = conv;
LastActivityUtc = DateTime.UtcNow;
_handle = GCHandle.Alloc(this);
_kcp = KCP.ikcp_create(conv, (void*)GCHandle.ToIntPtr(_handle));
KCP.ikcp_setoutput(_kcp, &OutputCallback);
KCP.ikcp_nodelay(_kcp, DefaultNoDelay, DefaultInterval, DefaultResend, DefaultNc);
KCP.ikcp_wndsize(_kcp, DefaultSendWindow, DefaultReceiveWindow);
KCP.ikcp_setmtu(_kcp, DefaultMtu);
_nextUpdateAt = GetCurrentTimeMilliseconds();
}
public uint Conv { get; }
public IPEndPoint RemoteEndPoint { get; }
public DateTime LastActivityUtc { get; private set; }
public void Send(byte[] payload)
{
if (payload == null)
{
throw new ArgumentNullException(nameof(payload));
}
lock (_gate)
{
ThrowIfDisposed();
if (payload.Length == 0)
{
return;
}
fixed (byte* buffer = payload)
{
var result = KCP.ikcp_send(_kcp, buffer, payload.Length);
if (result < 0)
{
throw new InvalidOperationException($"KCP send failed with error code {result}.");
}
}
LastActivityUtc = DateTime.UtcNow;
UpdateNoLock(GetCurrentTimeMilliseconds());
}
}
public void Input(byte[] datagram)
{
if (datagram == null)
{
throw new ArgumentNullException(nameof(datagram));
}
if (datagram.Length == 0)
{
return;
}
lock (_gate)
{
ThrowIfDisposed();
fixed (byte* buffer = datagram)
{
var result = KCP.ikcp_input(_kcp, buffer, datagram.Length);
if (result < 0)
{
Console.WriteLine($"[KcpTransport] KCP input failed for {RemoteEndPoint}: {result}");
return;
}
}
LastActivityUtc = DateTime.UtcNow;
UpdateNoLock(GetCurrentTimeMilliseconds());
}
}
public bool TryReceive(out byte[] payload)
{
lock (_gate)
{
if (_disposed)
{
payload = null;
return false;
}
var size = KCP.ikcp_peeksize(_kcp);
if (size <= 0)
{
payload = null;
return false;
}
payload = new byte[size];
fixed (byte* buffer = payload)
{
var result = KCP.ikcp_recv(_kcp, buffer, payload.Length);
if (result < 0)
{
payload = null;
return false;
}
if (result != payload.Length)
{
Array.Resize(ref payload, result);
}
}
LastActivityUtc = DateTime.UtcNow;
return true;
}
}
public void UpdateIfDue(uint current)
{
lock (_gate)
{
if (_disposed)
{
return;
}
if (KCP._itimediff(current, _nextUpdateAt) < 0)
{
return;
}
UpdateNoLock(current);
}
}
public void Dispose()
{
lock (_gate)
{
if (_disposed)
{
return;
}
_disposed = true;
if (_kcp != null)
{
KCP.ikcp_release(_kcp);
_kcp = null;
}
if (_handle.IsAllocated)
{
_handle.Free();
}
}
}
private void UpdateNoLock(uint current)
{
KCP.ikcp_update(_kcp, current);
_nextUpdateAt = KCP.ikcp_check(_kcp, current);
}
private void ThrowIfDisposed()
{
if (_disposed || _kcp == null)
{
throw new ObjectDisposedException(nameof(KcpSession));
}
}
private int SendRaw(byte* buffer, int length)
{
return _owner.SendDatagram(buffer, length, RemoteEndPoint);
}
private static int OutputCallback(byte* buffer, int length, IKCPCB* kcp, void* user)
{
if (user == null)
{
return -1;
}
var handle = GCHandle.FromIntPtr((IntPtr)user);
if (handle.Target is not KcpSession session)
{
return -1;
}
return session.SendRaw(buffer, length);
}
}
}
}
}