完成阶段 4
This commit is contained in:
@@ -0,0 +1,14 @@
|
||||
using System;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Network.NetworkApplication
|
||||
{
|
||||
public interface INetworkMessageDispatcher
|
||||
{
|
||||
int PendingCount { get; }
|
||||
|
||||
void Enqueue(Func<Task> workItem);
|
||||
|
||||
Task<int> DrainAsync(int maxItems = int.MaxValue);
|
||||
}
|
||||
}
|
||||
+1
-1
@@ -1,5 +1,5 @@
|
||||
fileFormatVersion: 2
|
||||
guid: d26d19f5e4031fd4089d620dc62d5159
|
||||
guid: 688ae8436f5f43fa88382700ef9c5e58
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
@@ -0,0 +1,30 @@
|
||||
using System;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Network.NetworkApplication
|
||||
{
|
||||
public sealed class ImmediateNetworkMessageDispatcher : INetworkMessageDispatcher
|
||||
{
|
||||
public int PendingCount => 0;
|
||||
|
||||
public void Enqueue(Func<Task> workItem)
|
||||
{
|
||||
if (workItem == null)
|
||||
{
|
||||
throw new ArgumentNullException(nameof(workItem));
|
||||
}
|
||||
|
||||
workItem().GetAwaiter().GetResult();
|
||||
}
|
||||
|
||||
public Task<int> DrainAsync(int maxItems = int.MaxValue)
|
||||
{
|
||||
if (maxItems <= 0)
|
||||
{
|
||||
throw new ArgumentOutOfRangeException(nameof(maxItems));
|
||||
}
|
||||
|
||||
return Task.FromResult(0);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: c1b99c76091f41299af26864a9a25fb0
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,41 @@
|
||||
using System;
|
||||
using System.Collections.Concurrent;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Network.NetworkApplication
|
||||
{
|
||||
public sealed class MainThreadNetworkDispatcher : INetworkMessageDispatcher
|
||||
{
|
||||
private readonly ConcurrentQueue<Func<Task>> pendingWork = new();
|
||||
|
||||
public int PendingCount => pendingWork.Count;
|
||||
|
||||
public void Enqueue(Func<Task> workItem)
|
||||
{
|
||||
if (workItem == null)
|
||||
{
|
||||
throw new ArgumentNullException(nameof(workItem));
|
||||
}
|
||||
|
||||
pendingWork.Enqueue(workItem);
|
||||
}
|
||||
|
||||
public async Task<int> DrainAsync(int maxItems = int.MaxValue)
|
||||
{
|
||||
if (maxItems <= 0)
|
||||
{
|
||||
throw new ArgumentOutOfRangeException(nameof(maxItems));
|
||||
}
|
||||
|
||||
var processed = 0;
|
||||
|
||||
while (processed < maxItems && pendingWork.TryDequeue(out var workItem))
|
||||
{
|
||||
await workItem();
|
||||
processed++;
|
||||
}
|
||||
|
||||
return processed;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 3dc7b1ecbad541ea86a9d700dd5148e0
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -1,4 +1,4 @@
|
||||
using System;
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Net;
|
||||
using System.Threading.Tasks;
|
||||
@@ -11,16 +11,22 @@ namespace Network.NetworkApplication
|
||||
public class MessageManager
|
||||
{
|
||||
private readonly ITransport transport;
|
||||
private readonly INetworkMessageDispatcher dispatcher;
|
||||
|
||||
private readonly Dictionary<MessageType, Func<byte[], IPEndPoint, Task>> handlers =
|
||||
new();
|
||||
|
||||
public MessageManager(ITransport transport)
|
||||
public MessageManager(ITransport transport, INetworkMessageDispatcher dispatcher)
|
||||
{
|
||||
this.transport = transport ?? throw new ArgumentNullException(nameof(transport));
|
||||
this.transport.OnReceive += OnTransportReceiveAsync;
|
||||
this.dispatcher = dispatcher ?? throw new ArgumentNullException(nameof(dispatcher));
|
||||
this.transport.OnReceive += OnTransportReceive;
|
||||
}
|
||||
|
||||
public INetworkMessageDispatcher Dispatcher => dispatcher;
|
||||
|
||||
public int PendingMessageCount => dispatcher.PendingCount;
|
||||
|
||||
public void RegisterHandler(MessageType type, IMessageHandler handler)
|
||||
{
|
||||
if (handler == null)
|
||||
@@ -94,7 +100,12 @@ namespace Network.NetworkApplication
|
||||
transport.SendToAll(envelope.ToByteArray());
|
||||
}
|
||||
|
||||
private async void OnTransportReceiveAsync(byte[] data, IPEndPoint sender)
|
||||
public Task<int> DrainPendingMessagesAsync(int maxMessages = int.MaxValue)
|
||||
{
|
||||
return dispatcher.DrainAsync(maxMessages);
|
||||
}
|
||||
|
||||
private void OnTransportReceive(byte[] data, IPEndPoint sender)
|
||||
{
|
||||
try
|
||||
{
|
||||
@@ -104,7 +115,8 @@ namespace Network.NetworkApplication
|
||||
|
||||
if (handlers.TryGetValue(type, out var handler))
|
||||
{
|
||||
await handler(envelope.Payload.ToByteArray(), sender);
|
||||
var payload = envelope.Payload.ToByteArray();
|
||||
dispatcher.Enqueue(() => DispatchAsync(handler, payload, sender, type));
|
||||
}
|
||||
else
|
||||
{
|
||||
@@ -116,5 +128,21 @@ namespace Network.NetworkApplication
|
||||
Console.WriteLine($"[MessageManager] 消息处理错误:{ex.Message}");
|
||||
}
|
||||
}
|
||||
|
||||
private static async Task DispatchAsync(
|
||||
Func<byte[], IPEndPoint, Task> handler,
|
||||
byte[] payload,
|
||||
IPEndPoint sender,
|
||||
MessageType type)
|
||||
{
|
||||
try
|
||||
{
|
||||
await handler(payload, sender);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
Console.WriteLine($"[MessageManager] Handler 执行错误:{type} -> {ex.Message}");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
using System;
|
||||
using System.Threading.Tasks;
|
||||
using Network.NetworkTransport;
|
||||
|
||||
namespace Network.NetworkApplication
|
||||
{
|
||||
public sealed class SharedNetworkRuntime
|
||||
{
|
||||
public SharedNetworkRuntime(ITransport transport, INetworkMessageDispatcher dispatcher)
|
||||
{
|
||||
Transport = transport ?? throw new ArgumentNullException(nameof(transport));
|
||||
MessageManager = new MessageManager(transport, dispatcher ?? throw new ArgumentNullException(nameof(dispatcher)));
|
||||
}
|
||||
|
||||
public ITransport Transport { get; }
|
||||
|
||||
public MessageManager MessageManager { get; }
|
||||
|
||||
public Task StartAsync()
|
||||
{
|
||||
return Transport.StartAsync();
|
||||
}
|
||||
|
||||
public void Stop()
|
||||
{
|
||||
Transport.Stop();
|
||||
}
|
||||
|
||||
public Task<int> DrainPendingMessagesAsync(int maxMessages = int.MaxValue)
|
||||
{
|
||||
return MessageManager.DrainPendingMessagesAsync(maxMessages);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 314c8d46e1ae4d9eb4914d0cac7bb628
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,8 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 9cf2571f026e4872b3c07033fd0c21a9
|
||||
folderAsset: yes
|
||||
DefaultImporter:
|
||||
externalObjects: {}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,38 @@
|
||||
using System;
|
||||
using System.Threading.Tasks;
|
||||
using Network.NetworkApplication;
|
||||
using Network.NetworkTransport;
|
||||
|
||||
namespace Network.NetworkHost
|
||||
{
|
||||
public sealed class ServerNetworkHost
|
||||
{
|
||||
private readonly SharedNetworkRuntime runtime;
|
||||
|
||||
public ServerNetworkHost(ITransport transport, INetworkMessageDispatcher dispatcher = null)
|
||||
{
|
||||
runtime = new SharedNetworkRuntime(
|
||||
transport ?? throw new ArgumentNullException(nameof(transport)),
|
||||
dispatcher ?? new ImmediateNetworkMessageDispatcher());
|
||||
}
|
||||
|
||||
public MessageManager MessageManager => runtime.MessageManager;
|
||||
|
||||
public ITransport Transport => runtime.Transport;
|
||||
|
||||
public Task StartAsync()
|
||||
{
|
||||
return runtime.StartAsync();
|
||||
}
|
||||
|
||||
public void Stop()
|
||||
{
|
||||
runtime.Stop();
|
||||
}
|
||||
|
||||
public Task<int> DrainPendingMessagesAsync(int maxMessages = int.MaxValue)
|
||||
{
|
||||
return runtime.DrainPendingMessagesAsync(maxMessages);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 89f7ce53cdf54dc9ac44f14eaf11cf5d
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -1,158 +0,0 @@
|
||||
using System;
|
||||
using System.Collections.Concurrent;
|
||||
using System.Net;
|
||||
using System.Net.Sockets;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Network.NetworkTransport
|
||||
{
|
||||
public class ReliableUdpTransport : ITransport
|
||||
{
|
||||
private readonly UdpClient _client;
|
||||
private readonly IPEndPoint _defaultRemoteEndPoint;
|
||||
private readonly bool _isServer;
|
||||
|
||||
// Stage one keeps this class name for compatibility while collapsing it to plain UDP.
|
||||
private readonly ConcurrentDictionary<string, IPEndPoint> _knownRemoteEndPoints = new();
|
||||
|
||||
private volatile bool _isRunning;
|
||||
|
||||
public event Action<byte[], IPEndPoint> OnReceive;
|
||||
|
||||
private Task _receiveTask = Task.CompletedTask;
|
||||
|
||||
public ReliableUdpTransport(int listenPort)
|
||||
{
|
||||
_client = new UdpClient(listenPort);
|
||||
_isServer = true;
|
||||
Console.WriteLine($"[Transport] 服务端模式,监听端口: {listenPort}");
|
||||
}
|
||||
|
||||
public ReliableUdpTransport(string serverIP, int serverPort)
|
||||
{
|
||||
_client = new UdpClient(0);
|
||||
_defaultRemoteEndPoint = new IPEndPoint(IPAddress.Parse(serverIP), serverPort);
|
||||
_isServer = false;
|
||||
Console.WriteLine($"[Transport] 客户端模式,目标: {_defaultRemoteEndPoint}");
|
||||
}
|
||||
|
||||
public async Task StartAsync()
|
||||
{
|
||||
if (_isRunning)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
_knownRemoteEndPoints.Clear();
|
||||
_isRunning = true;
|
||||
Console.WriteLine("[Transport] 传输层启动");
|
||||
_receiveTask = ReceiveLoop();
|
||||
await Task.Yield();
|
||||
}
|
||||
|
||||
public void Stop()
|
||||
{
|
||||
if (!_isRunning)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
_isRunning = false;
|
||||
_client.Close();
|
||||
_knownRemoteEndPoints.Clear();
|
||||
Console.WriteLine("[Transport] 传输层停止");
|
||||
}
|
||||
|
||||
public void Send(byte[] data)
|
||||
{
|
||||
if (_defaultRemoteEndPoint == null)
|
||||
{
|
||||
throw new InvalidOperationException("Default remote endpoint is not configured.");
|
||||
}
|
||||
|
||||
SendTo(data, _defaultRemoteEndPoint);
|
||||
}
|
||||
|
||||
public void SendTo(byte[] data, IPEndPoint target)
|
||||
{
|
||||
if (data == null)
|
||||
{
|
||||
throw new ArgumentNullException(nameof(data));
|
||||
}
|
||||
|
||||
if (target == null)
|
||||
{
|
||||
throw new ArgumentNullException(nameof(target));
|
||||
}
|
||||
|
||||
EnsureRunning();
|
||||
RememberRemote(target);
|
||||
_client.Send(data, data.Length, target);
|
||||
Console.WriteLine($"[Transport] 发送数据到 {target}");
|
||||
}
|
||||
|
||||
public void SendToAll(byte[] data)
|
||||
{
|
||||
if (data == null)
|
||||
{
|
||||
throw new ArgumentNullException(nameof(data));
|
||||
}
|
||||
|
||||
EnsureRunning();
|
||||
|
||||
if (!_isServer)
|
||||
{
|
||||
throw new InvalidOperationException("SendToAll is only supported in server mode.");
|
||||
}
|
||||
|
||||
foreach (var remoteEndPoint in _knownRemoteEndPoints.Values)
|
||||
{
|
||||
_client.Send(data, data.Length, remoteEndPoint);
|
||||
Console.WriteLine($"[Transport] 广播数据到 {remoteEndPoint}");
|
||||
}
|
||||
}
|
||||
|
||||
private async Task ReceiveLoop()
|
||||
{
|
||||
while (_isRunning)
|
||||
{
|
||||
try
|
||||
{
|
||||
var result = await _client.ReceiveAsync();
|
||||
RememberRemote(result.RemoteEndPoint);
|
||||
OnReceive?.Invoke(result.Buffer, result.RemoteEndPoint);
|
||||
}
|
||||
catch (ObjectDisposedException) when (!_isRunning)
|
||||
{
|
||||
return;
|
||||
}
|
||||
catch (SocketException) when (!_isRunning)
|
||||
{
|
||||
return;
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
Console.WriteLine($"[Transport] 接收错误:{e.Message}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void EnsureRunning()
|
||||
{
|
||||
if (!_isRunning)
|
||||
{
|
||||
throw new InvalidOperationException("Transport has not been started.");
|
||||
}
|
||||
}
|
||||
|
||||
private void RememberRemote(IPEndPoint remoteEndPoint)
|
||||
{
|
||||
if (remoteEndPoint == null)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
_knownRemoteEndPoints[remoteEndPoint.ToString()] = remoteEndPoint;
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user