Init
This commit is contained in:
@@ -0,0 +1,8 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 895d4be03eb18e941b9a0b25d937da97
|
||||
folderAsset: yes
|
||||
DefaultImporter:
|
||||
externalObjects: {}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 13502c02ef6908e46b00ba64f6df5b77
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,38 @@
|
||||
namespace Network.Defines
|
||||
{
|
||||
public enum MessageType : byte
|
||||
{
|
||||
Unknow = 0,
|
||||
|
||||
// 游戏相关
|
||||
PlayerInput = 1,
|
||||
PlayerState = 2,
|
||||
PlayerAction = 3,
|
||||
GameState = 4,
|
||||
PlayerJoin = 5,
|
||||
PlayerLeave = 6,
|
||||
|
||||
// 聊天相关
|
||||
ChatMessage = 10,
|
||||
PrivateMessage = 11,
|
||||
SystemMessage = 12,
|
||||
|
||||
// 系统相关
|
||||
HeartBeat = 20,
|
||||
|
||||
LoginRequest = 21,
|
||||
LoginResponse = 22,
|
||||
|
||||
LogoutRequest = 23,
|
||||
|
||||
// 房间管理
|
||||
CreateRoom = 30,
|
||||
JoinRoom = 31,
|
||||
LeaveRoom = 32,
|
||||
RoomList = 33,
|
||||
|
||||
Heartbeat = 40,
|
||||
HeartbeatResponse = 41,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 1dd2a49fbda33064891c337fadc590ed
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,27 @@
|
||||
using UnityEngine;
|
||||
|
||||
namespace Network.Defines
|
||||
{
|
||||
public static class ProtoExtensions
|
||||
{
|
||||
public static UnityEngine.Vector3 ToVector3(this Vector3 vec)
|
||||
{
|
||||
return new UnityEngine.Vector3()
|
||||
{
|
||||
x = vec.X,
|
||||
y = vec.Y,
|
||||
z = vec.Z
|
||||
};
|
||||
}
|
||||
|
||||
public static Vector3 ToProtoVector3(UnityEngine.Vector3 vec)
|
||||
{
|
||||
return new Vector3()
|
||||
{
|
||||
X = vec.x,
|
||||
Y = vec.y,
|
||||
Z = vec.z
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: fefe99db6b80a43418075799ba64c176
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,11 @@
|
||||
using System;
|
||||
|
||||
namespace Network.Defines
|
||||
{
|
||||
public class SystemMessage
|
||||
{
|
||||
public string Content { get; set; }
|
||||
public DateTime Timestamp { get; set; } = DateTime.Now;
|
||||
public string Level { get; set; } = "info"; // "info", "warning", "error"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 9a9d4f06e72ce954598ae990364ac8ac
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,8 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 2abfb5f176e7f5e4c8bf371eb28ae588
|
||||
folderAsset: yes
|
||||
DefaultImporter:
|
||||
externalObjects: {}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,32 @@
|
||||
using System;
|
||||
using System.Net;
|
||||
using System.Threading.Tasks;
|
||||
using Google.Protobuf;
|
||||
using Network.Defines;
|
||||
|
||||
namespace Network.NetworkApplication
|
||||
{
|
||||
public class DelegateMessageHandler : IMessageHandler
|
||||
{
|
||||
private readonly Func<byte[], IPEndPoint, Task> handler;
|
||||
|
||||
public DelegateMessageHandler(Func<byte[], IPEndPoint, Task> handler)
|
||||
{
|
||||
this.handler = handler;
|
||||
}
|
||||
|
||||
public DelegateMessageHandler(Action<byte[], IPEndPoint> handler)
|
||||
{
|
||||
this.handler = (msg, sender) =>
|
||||
{
|
||||
handler(msg, sender);
|
||||
return Task.CompletedTask;
|
||||
};
|
||||
}
|
||||
|
||||
public Task HandleAsync(byte[] message, IPEndPoint sender)
|
||||
{
|
||||
return handler(message, sender);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: c940a6a8da48ccf429215f05e2d32154
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,10 @@
|
||||
using System.Net;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Network.NetworkApplication
|
||||
{
|
||||
public interface IMessageHandler
|
||||
{
|
||||
Task HandleAsync(byte[] message, IPEndPoint sender);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: c46d2128446560d409c48949e2636bf1
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,101 @@
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Net;
|
||||
using System.Threading.Tasks;
|
||||
using Google.Protobuf;
|
||||
using Network.Defines;
|
||||
using Network.NetworkTransport;
|
||||
|
||||
namespace Network.NetworkApplication
|
||||
{
|
||||
public class MessageManager
|
||||
{
|
||||
private readonly ITransport transport;
|
||||
|
||||
private readonly Dictionary<MessageType, Func<byte[], IPEndPoint, Task>> handlers =
|
||||
new Dictionary<MessageType, Func<byte[], IPEndPoint, Task>>();
|
||||
|
||||
public MessageManager(ITransport transport)
|
||||
{
|
||||
this.transport = transport;
|
||||
this.transport.OnReceive += OnTransportReceiveAsync;
|
||||
}
|
||||
|
||||
public void RegisterHandler(MessageType type, IMessageHandler handler)
|
||||
{
|
||||
handlers[type] = async (payload, sender) => { await handler.HandleAsync(payload, sender); };
|
||||
|
||||
Console.WriteLine($"[MessageManager] 注册处理器:{type}");
|
||||
}
|
||||
|
||||
public void RegisterHandler(MessageType type, Func<byte[], IPEndPoint, Task> handler)
|
||||
{
|
||||
var han = new DelegateMessageHandler(handler);
|
||||
RegisterHandler(type, new DelegateMessageHandler(handler));
|
||||
}
|
||||
|
||||
public void RegisterHandler(MessageType type, Action<byte[], IPEndPoint> handler)
|
||||
{
|
||||
var han = new DelegateMessageHandler((msg, sender) => { handler(msg, sender); });
|
||||
RegisterHandler(type, new DelegateMessageHandler((msg, sender) =>
|
||||
{
|
||||
handler(msg, sender);
|
||||
return Task.CompletedTask;
|
||||
}));
|
||||
}
|
||||
|
||||
public void SendMessage<T>(T message, MessageType type, IPEndPoint target = null) where T : IMessage
|
||||
{
|
||||
var envelope = new Envelope()
|
||||
{
|
||||
Type = (int)type,
|
||||
Payload = message.ToByteString()
|
||||
};
|
||||
|
||||
if (target != null)
|
||||
{
|
||||
transport.SendTo(envelope.ToByteArray(), target);
|
||||
}
|
||||
else
|
||||
{
|
||||
transport.Send(envelope.ToByteArray());
|
||||
}
|
||||
|
||||
Console.WriteLine($"[MessageManager] 发送消息:{type} -> {target?.ToString() ?? "default"}");
|
||||
}
|
||||
|
||||
public void BroadcastMessage<T>(T message, MessageType type) where T : IMessage
|
||||
{
|
||||
Console.WriteLine($"[MessageManager] 广播消息:{type}");
|
||||
var envelope = new Envelope()
|
||||
{
|
||||
Type = (int)type,
|
||||
Payload = message.ToByteString()
|
||||
};
|
||||
transport.SendToAll(envelope.ToByteArray());
|
||||
}
|
||||
|
||||
private async void OnTransportReceiveAsync(byte[] data, IPEndPoint sender)
|
||||
{
|
||||
try
|
||||
{
|
||||
var envelope = Envelope.Parser.ParseFrom(data);
|
||||
var type = (MessageType)envelope.Type;
|
||||
Console.WriteLine($"[MessageManager] 收到消息:{type} 来自 {sender}");
|
||||
|
||||
if (handlers.TryGetValue(type, out var handler))
|
||||
{
|
||||
await handler(envelope.Payload.ToByteArray(), sender);
|
||||
}
|
||||
else
|
||||
{
|
||||
Console.WriteLine($"[MessageManager] 警告:未注册的消息类型 {type}");
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
Console.WriteLine($"[MessageManager] 消息处理错误:{ex.Message}");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 75ac30aeadd168e44a8860b22bd467c8
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,8 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 32f1de5d4a6031049a2033d69eb21595
|
||||
folderAsset: yes
|
||||
DefaultImporter:
|
||||
externalObjects: {}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,193 @@
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Net;
|
||||
using System.Net.Sockets;
|
||||
|
||||
namespace Network.NetworkTransport
|
||||
{
|
||||
public class ClientSession
|
||||
{
|
||||
private ITransport _transport;
|
||||
|
||||
private IPEndPoint _remote;
|
||||
|
||||
public long LastActivityTs { get; private set; }
|
||||
|
||||
public uint SendSequenceNumber { get; private set; } = 0;
|
||||
|
||||
private int _currentTicks = 0;
|
||||
|
||||
private int _nextSendTicks = 0;
|
||||
|
||||
private int _sendInterval = 10;
|
||||
|
||||
// 重传时间 5s
|
||||
private long _retransmitTicks = 5000;
|
||||
|
||||
// 上层交付
|
||||
private readonly LinkedList<Packet> _sendQueue = new LinkedList<Packet>();
|
||||
|
||||
// 已发送但未确认
|
||||
private readonly LinkedList<Packet> _sendBuffer = new LinkedList<Packet>();
|
||||
|
||||
// 已收到但乱序
|
||||
private readonly LinkedList<Packet> _receiveBuffer = new LinkedList<Packet>();
|
||||
|
||||
// 已收到可交付
|
||||
private readonly LinkedList<Packet> _receiveQueue = new LinkedList<Packet>();
|
||||
|
||||
private bool _hasReceived = false;
|
||||
|
||||
private uint _expectedAck = 0;
|
||||
|
||||
private readonly object _lockObj = new object();
|
||||
|
||||
public ClientSession(ITransport transport, IPEndPoint remote)
|
||||
{
|
||||
_transport = transport;
|
||||
_remote = remote;
|
||||
LastActivityTs = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds();
|
||||
}
|
||||
|
||||
public uint GetExpectedAck() => _expectedAck;
|
||||
|
||||
public void SetSendInterval(int interval) => _sendInterval = interval;
|
||||
|
||||
private uint GetNextSendSequence()
|
||||
{
|
||||
lock (_lockObj)
|
||||
{
|
||||
return SendSequenceNumber++;
|
||||
}
|
||||
}
|
||||
|
||||
public void Tick(int currentTicks)
|
||||
{
|
||||
_currentTicks = currentTicks;
|
||||
if (_currentTicks >= _nextSendTicks)
|
||||
{
|
||||
_nextSendTicks = currentTicks + _sendInterval;
|
||||
SendPacketInternal();
|
||||
}
|
||||
}
|
||||
|
||||
public void SendPacket(byte[] data)
|
||||
{
|
||||
_sendQueue.AddLast(Packet.CreateDataPacket(GetNextSendSequence(), data));
|
||||
}
|
||||
|
||||
public List<Packet> ReceivePackets()
|
||||
{
|
||||
var list = new List<Packet>();
|
||||
lock (_lockObj)
|
||||
{
|
||||
while (_receiveQueue.Count > 0)
|
||||
{
|
||||
var packet = _receiveQueue.First.Value;
|
||||
|
||||
list.Add(packet);
|
||||
_receiveQueue.RemoveFirst();
|
||||
}
|
||||
}
|
||||
|
||||
return list;
|
||||
}
|
||||
|
||||
private void SendPacketInternal()
|
||||
{
|
||||
if (_hasReceived)
|
||||
{
|
||||
var packet = Packet.CreateAckPacket(_expectedAck);
|
||||
_sendBuffer.AddLast(packet);
|
||||
var bytes = packet.ToBytes();
|
||||
_transport.SendTo(bytes, _remote);
|
||||
_hasReceived = false;
|
||||
}
|
||||
|
||||
foreach (var packet in _receiveBuffer)
|
||||
{
|
||||
if (_currentTicks - packet.Timestamp > _retransmitTicks)
|
||||
{
|
||||
var bytes = packet.ToBytes();
|
||||
_transport.SendTo(bytes, _remote);
|
||||
}
|
||||
else break;
|
||||
}
|
||||
|
||||
while (_sendQueue.Count > 0)
|
||||
{
|
||||
var packet = _sendQueue.First.Value;
|
||||
_sendBuffer.AddLast(packet);
|
||||
var bytes = packet.ToBytes();
|
||||
_transport.SendTo(bytes, _remote);
|
||||
}
|
||||
}
|
||||
|
||||
public void ReceivePacketsInternal(Packet packet)
|
||||
{
|
||||
uint seq = packet.SequenceNumber;
|
||||
|
||||
// 是否是按序到达的包
|
||||
if (seq == _expectedAck)
|
||||
{
|
||||
_receiveQueue.AddLast(packet);
|
||||
while (_receiveBuffer.Count > 0)
|
||||
{
|
||||
var pendingPacket = _receiveBuffer.First.Value;
|
||||
if (seq != pendingPacket.SequenceNumber) break;
|
||||
seq++;
|
||||
_receiveQueue.AddLast(pendingPacket);
|
||||
_receiveBuffer.RemoveFirst();
|
||||
}
|
||||
|
||||
_expectedAck = seq + 1;
|
||||
_hasReceived = true;
|
||||
}
|
||||
// 将包按顺序追加在 receivingPackets 后面
|
||||
else
|
||||
{
|
||||
var firstNode = _receiveBuffer.First;
|
||||
while (firstNode.Next != null)
|
||||
{
|
||||
if (firstNode.Value.SequenceNumber > seq)
|
||||
{
|
||||
var node = new LinkedListNode<Packet>(packet);
|
||||
_receiveBuffer.AddBefore(firstNode, node);
|
||||
break;
|
||||
}
|
||||
|
||||
firstNode = firstNode.Next;
|
||||
}
|
||||
|
||||
if (firstNode == null) _receiveBuffer.AddLast(packet);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
public bool TryProcessReceiveSequence(uint sequenceNumber, out bool shouldDeliver)
|
||||
{
|
||||
lock (_lockObj)
|
||||
{
|
||||
LastActivityTs = DateTime.Now;
|
||||
|
||||
if (sequenceNumber == _expectedAck)
|
||||
{
|
||||
_expectedAck++;
|
||||
_receivedSequences.Add(sequenceNumber);
|
||||
shouldDeliver = true;
|
||||
return true;
|
||||
}
|
||||
else if (sequenceNumber < _expectedAck)
|
||||
{
|
||||
shouldDeliver = false;
|
||||
return _receivedSequences.Contains(sequenceNumber);
|
||||
}
|
||||
else
|
||||
{
|
||||
shouldDeliver = false;
|
||||
return false;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: ec6c25bc42967db499742dfa355380b7
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,15 @@
|
||||
using System;
|
||||
using System.Net;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Network.NetworkTransport
|
||||
{
|
||||
public interface ITransport
|
||||
{
|
||||
void SendTo(byte[] data, IPEndPoint target);
|
||||
void SendToAll(byte[] data);
|
||||
event Action<byte[], IPEndPoint> OnReceive;
|
||||
Task StartAsync();
|
||||
void Stop();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: dc400a702b75abc40bb454eb27a34249
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,63 @@
|
||||
using System;
|
||||
using System.Linq;
|
||||
using UnityEngine;
|
||||
|
||||
namespace Network.NetworkTransport
|
||||
{
|
||||
public enum PacketType : byte
|
||||
{
|
||||
Data = 1,
|
||||
Ack = 2,
|
||||
}
|
||||
|
||||
public struct Packet
|
||||
{
|
||||
public PacketType Type;
|
||||
public uint SequenceNumber;
|
||||
public byte[] Data;
|
||||
public long Timestamp;
|
||||
|
||||
public byte[] ToBytes()
|
||||
{
|
||||
var result = new byte[1 + 4 + 8 + Data.Length];
|
||||
result[0] = (byte)Type;
|
||||
BitConverter.GetBytes(SequenceNumber).CopyTo(result, 1);
|
||||
BitConverter.GetBytes(Timestamp).CopyTo(result, 5);
|
||||
Data.CopyTo(result, 13);
|
||||
return result;
|
||||
}
|
||||
|
||||
public static Packet FromBytes(byte[] data)
|
||||
{
|
||||
return new Packet
|
||||
{
|
||||
Type = (PacketType)data[0],
|
||||
SequenceNumber = BitConverter.ToUInt32(data, 1),
|
||||
Timestamp = BitConverter.ToInt64(data, 5),
|
||||
Data = new ArraySegment<byte>(data, 5, data.Length - 5).ToArray()
|
||||
};
|
||||
}
|
||||
|
||||
public static Packet CreateDataPacket(uint seqNum, byte[] data)
|
||||
{
|
||||
return new Packet
|
||||
{
|
||||
Type = PacketType.Data,
|
||||
SequenceNumber = seqNum,
|
||||
Data = data,
|
||||
Timestamp = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds()
|
||||
};
|
||||
}
|
||||
|
||||
public static Packet CreateAckPacket(uint seqNum)
|
||||
{
|
||||
return new Packet
|
||||
{
|
||||
Type = PacketType.Ack,
|
||||
SequenceNumber = seqNum,
|
||||
Data = Array.Empty<byte>(),
|
||||
Timestamp = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds()
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: b84a2cb7ffe3cd14180358559e526dbe
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,285 @@
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Net;
|
||||
using System.Net.Sockets;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Network.NetworkTransport
|
||||
{
|
||||
public class ReliableUdpTransport : ITransport
|
||||
{
|
||||
private readonly UdpClient _client;
|
||||
private readonly IPEndPoint _defaultRemoteEndPoint;
|
||||
private readonly bool _isServer;
|
||||
|
||||
private readonly List<ClientSession> _sessions = new();
|
||||
|
||||
private readonly Timer _retransmitTimer;
|
||||
private readonly Timer _cleanupTimer;
|
||||
|
||||
//TODO: volatile 关键字
|
||||
private volatile bool _isRunning;
|
||||
|
||||
// 配置参数
|
||||
private const int RetransmitTimeoutMs = 1000;
|
||||
private const int SessionTimeoutMs = 30000;
|
||||
private const int MaxRetransmitAttempts = 5;
|
||||
|
||||
public event Action<byte[], IPEndPoint> OnReceive;
|
||||
|
||||
private Task _receiveTask;
|
||||
|
||||
// 构造函数——服务端模式
|
||||
public ReliableUdpTransport(int listenPort)
|
||||
{
|
||||
_client = new UdpClient(listenPort);
|
||||
_isServer = true;
|
||||
_retransmitTimer = new Timer(CheckRetransmit, null, 100, 100);
|
||||
_cleanupTimer = new Timer(CleanupSessions, null, 5000, 5000);
|
||||
Console.WriteLine($"[Transport] 服务端模式,监听端口: {listenPort}");
|
||||
}
|
||||
|
||||
// 构造函数——客户端模式
|
||||
public ReliableUdpTransport(string serverIP, int serverPort)
|
||||
{
|
||||
_client = new UdpClient(0);
|
||||
_defaultRemoteEndPoint = new IPEndPoint(IPAddress.Parse(serverIP), serverPort);
|
||||
|
||||
_isServer = false;
|
||||
_retransmitTimer = new Timer(CheckRetransmit, null, 100, 100);
|
||||
_cleanupTimer = new Timer(CleanupSessions, null, 5000, 5000);
|
||||
Console.WriteLine($"[Transport] 客户端模式,目标: {_defaultRemoteEndPoint}");
|
||||
}
|
||||
|
||||
public async Task StartAsync()
|
||||
{
|
||||
_sessions.Clear();
|
||||
|
||||
_isRunning = true;
|
||||
Console.WriteLine("[Transport] 传输层启动");
|
||||
|
||||
// 开始接收数据
|
||||
_receiveTask = ReceiveLoop();
|
||||
await Task.Delay(100); // 给接收循环一点启动时间
|
||||
}
|
||||
|
||||
public void Tick()
|
||||
{
|
||||
foreach (var session in _sessions)
|
||||
{
|
||||
session.Tick(DateTime.UtcNow.Millisecond);
|
||||
}
|
||||
}
|
||||
|
||||
public void Stop()
|
||||
{
|
||||
_isRunning = false;
|
||||
_retransmitTimer.Dispose();
|
||||
_cleanupTimer.Dispose();
|
||||
_client.Close();
|
||||
_sessions.Clear();
|
||||
Console.WriteLine("[Transport] 传输层停止");
|
||||
}
|
||||
|
||||
public async void SendTo(Packet packet, IPEndPoint target)
|
||||
{
|
||||
if (!_isRunning)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
var bytes = packet.ToBytes();
|
||||
await _client.SendAsync(bytes, bytes.Length, target);
|
||||
|
||||
Console.WriteLine($"[Transport] 发送数据包到 {target}");
|
||||
}
|
||||
|
||||
public void SendToAll(byte[] data)
|
||||
{
|
||||
foreach (var session in _sessions)
|
||||
{
|
||||
session.SendPacket(data);
|
||||
}
|
||||
}
|
||||
|
||||
private async Task ReceiveLoop()
|
||||
{
|
||||
while (_isRunning)
|
||||
{
|
||||
try
|
||||
{
|
||||
var result = await _client.ReceiveAsync();
|
||||
var packet = Packet.FromBytes(result.Buffer);
|
||||
|
||||
if (packet.Type == PacketType.Data)
|
||||
{
|
||||
HandleDataPacket(packet, result.RemoteEndPoint);
|
||||
}
|
||||
else if (packet.Type == PacketType.Ack)
|
||||
{
|
||||
HandleAckPacket(packet, result.RemoteEndPoint);
|
||||
}
|
||||
}
|
||||
catch (ObjectDisposedException)
|
||||
{
|
||||
return; // 正常关闭
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
Console.WriteLine($"[Transport] 接收错误:{e.Message}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void HandleDataPacket(Packet packet, IPEndPoint senderEndPoint)
|
||||
{
|
||||
var session = GetOrCreateSession(senderEndPoint);
|
||||
|
||||
Console.WriteLine(
|
||||
$"[Transport] 收到数据包从{senderEndPoint} SeqNum={packet.SequenceNumber}, DataLen={packet.Data.Length}");
|
||||
|
||||
// 发送ACK
|
||||
var ackPacket = Packet.CreateAckPacket(packet.SequenceNumber);
|
||||
SendPacketTo(ackPacket, senderEndPoint);
|
||||
Console.WriteLine($"[Transport] 发送ACK 到 {senderEndPoint} SeqNum={packet.SequenceNumber}");
|
||||
|
||||
// 检查是否应该交付
|
||||
if (session.TryProcessReceiveSequence(packet.SequenceNumber, out bool shouldDeliver))
|
||||
{
|
||||
if (shouldDeliver)
|
||||
{
|
||||
OnReceive?.Invoke(packet.Data, senderEndPoint);
|
||||
Console.WriteLine($"[Transport] 交付数据包从 {senderEndPoint} SeqNum={packet.SequenceNumber}");
|
||||
}
|
||||
else
|
||||
{
|
||||
Console.WriteLine($"[Transport] 重复包从 {senderEndPoint} SeqNum={packet.SequenceNumber},忽略");
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
// 乱序到达,暂存(简化处理:直接丢弃,依赖重传)
|
||||
Console.WriteLine($"[Transport] 乱序包从 {senderEndPoint} SeqNum={packet.SequenceNumber},丢弃");
|
||||
}
|
||||
}
|
||||
|
||||
private void HandleAckPacket(Packet packet, IPEndPoint senderEndPoint)
|
||||
{
|
||||
var session = GetOrCreateSession(senderEndPoint);
|
||||
Console.WriteLine($"[Transport] 收到ACK从 {senderEndPoint} SeqNum={packet.SequenceNumber}");
|
||||
|
||||
if (session.PendingAcks.TryRemove(packet.SequenceNumber, out _))
|
||||
{
|
||||
Console.WriteLine($"[Transport] 确认包到 {senderEndPoint} SeqNum={packet.SequenceNumber}");
|
||||
}
|
||||
}
|
||||
|
||||
private ClientSession GetOrCreateSession(IPEndPoint endPoint)
|
||||
{
|
||||
string key = endPoint.ToString();
|
||||
return _sessions.GetOrAdd(key, _ =>
|
||||
{
|
||||
var session = new ClientSession(endPoint);
|
||||
Console.WriteLine($"创建新会话:{endPoint}");
|
||||
return session;
|
||||
});
|
||||
}
|
||||
|
||||
private void CheckRetransmit(object state)
|
||||
{
|
||||
if (!_isRunning)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
var now = DateTime.Now;
|
||||
var toRetransmit = new List<(IPEndPoint target, uint seqNum, Packet packet)>();
|
||||
|
||||
foreach (var sessionKvp in _sessions)
|
||||
{
|
||||
var session = sessionKvp.Value;
|
||||
foreach (var ackKvp in session.PendingAcks)
|
||||
{
|
||||
var timeSinceLastSend = now - ackKvp.Value.sendTime;
|
||||
if (timeSinceLastSend.TotalMilliseconds > RetransmitTimeoutMs)
|
||||
{
|
||||
toRetransmit.Add((session.EndPoint, ackKvp.Key, ackKvp.Value.packet));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
foreach (var (target, seqNum, packet) in toRetransmit)
|
||||
{
|
||||
var session = GetOrCreateSession(target);
|
||||
if (session.PendingAcks.ContainsKey(seqNum))
|
||||
{
|
||||
// 更新发送时间
|
||||
session.PendingAcks[seqNum] = (packet, now);
|
||||
SendPacketTo(packet, target);
|
||||
Console.WriteLine($"[Transport] 重传包到 {target} SeqNum={seqNum}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void CleanupSessions(object state)
|
||||
{
|
||||
if (!_isRunning)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
var now = DateTime.Now;
|
||||
var toRemove = new List<string>();
|
||||
|
||||
foreach (var sessionKvp in _sessions)
|
||||
{
|
||||
var session = sessionKvp.Value;
|
||||
var timeSinceLastActivity = now - session.LastActivity;
|
||||
|
||||
if (timeSinceLastActivity.TotalMilliseconds > SessionTimeoutMs)
|
||||
{
|
||||
toRemove.Add(sessionKvp.Key);
|
||||
}
|
||||
}
|
||||
|
||||
foreach (string key in toRemove)
|
||||
{
|
||||
if (_sessions.TryRemove(key, out var session))
|
||||
{
|
||||
Console.WriteLine($"[Transport] 清理超时会话:{session.EndPoint}");
|
||||
}
|
||||
}
|
||||
|
||||
if (_isServer)
|
||||
{
|
||||
PrintSessionInfo();
|
||||
}
|
||||
}
|
||||
|
||||
private async void SendPacketTo(Packet packet, IPEndPoint endPoint)
|
||||
{
|
||||
try
|
||||
{
|
||||
var data = packet.ToBytes();
|
||||
await _client.SendAsync(data, data.Length, endPoint);
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
Console.WriteLine($"[Transport] 发送错误:{e.Message}");
|
||||
}
|
||||
}
|
||||
|
||||
private void PrintSessionInfo()
|
||||
{
|
||||
Console.WriteLine($"当前活跃会话数:{_sessions.Count}");
|
||||
foreach (var sessionKvp in _sessions)
|
||||
{
|
||||
var session = sessionKvp.Value;
|
||||
Console.WriteLine(
|
||||
$" 会话:{session.EndPoint},发送SeqNum:{session.SendSequenceNumber},期望接收:{session.GetExpectedAck()},待确认: {session.PendingAcks.Count}");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: d26d19f5e4031fd4089d620dc62d5159
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
Reference in New Issue
Block a user