process TODO.md step6
This commit is contained in:
@@ -1,4 +1,5 @@
|
||||
using System;
|
||||
using System.Threading.Tasks;
|
||||
using Network.NetworkHost;
|
||||
using Network.NetworkTransport;
|
||||
|
||||
@@ -84,6 +85,11 @@ namespace Network.NetworkApplication
|
||||
syncSequenceTracker);
|
||||
}
|
||||
|
||||
public static Task<ServerRuntimeHandle> StartServerRuntimeAsync(ServerRuntimeConfiguration configuration)
|
||||
{
|
||||
return ServerRuntimeEntryPoint.StartAsync(configuration);
|
||||
}
|
||||
|
||||
private static void ValidateDualPortConfiguration(int reliablePort, int? syncPort)
|
||||
{
|
||||
if (reliablePort <= 0)
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
using System;
|
||||
using Network.NetworkApplication;
|
||||
using Network.NetworkTransport;
|
||||
|
||||
namespace Network.NetworkHost
|
||||
{
|
||||
public sealed class ServerRuntimeConfiguration
|
||||
{
|
||||
public ServerRuntimeConfiguration(int reliablePort)
|
||||
{
|
||||
if (reliablePort <= 0)
|
||||
{
|
||||
throw new ArgumentOutOfRangeException(nameof(reliablePort), "Reliable port must be positive.");
|
||||
}
|
||||
|
||||
ReliablePort = reliablePort;
|
||||
}
|
||||
|
||||
public int ReliablePort { get; }
|
||||
|
||||
public int? SyncPort { get; set; }
|
||||
|
||||
public INetworkMessageDispatcher Dispatcher { get; set; }
|
||||
|
||||
public SessionReconnectPolicy ReconnectPolicy { get; set; }
|
||||
|
||||
public Func<DateTimeOffset> UtcNowProvider { get; set; }
|
||||
|
||||
public IMessageDeliveryPolicyResolver DeliveryPolicyResolver { get; set; }
|
||||
|
||||
public SyncSequenceTracker SyncSequenceTracker { get; set; }
|
||||
|
||||
public Func<int, ITransport> TransportFactory { get; set; }
|
||||
|
||||
internal void Validate()
|
||||
{
|
||||
if (ReliablePort <= 0)
|
||||
{
|
||||
throw new ArgumentOutOfRangeException(nameof(ReliablePort), "Reliable port must be positive.");
|
||||
}
|
||||
|
||||
if (!SyncPort.HasValue)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
if (SyncPort.Value <= 0)
|
||||
{
|
||||
throw new ArgumentOutOfRangeException(nameof(SyncPort), "Sync port must be positive.");
|
||||
}
|
||||
|
||||
if (SyncPort.Value == ReliablePort)
|
||||
{
|
||||
throw new ArgumentException("Sync port must differ from reliable port.", nameof(SyncPort));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: f6cf6d0542534955ad0dc92dd55a5429
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,40 @@
|
||||
using System;
|
||||
using System.Threading.Tasks;
|
||||
using Network.NetworkApplication;
|
||||
|
||||
namespace Network.NetworkHost
|
||||
{
|
||||
public static class ServerRuntimeEntryPoint
|
||||
{
|
||||
public static async Task<ServerRuntimeHandle> StartAsync(ServerRuntimeConfiguration configuration)
|
||||
{
|
||||
if (configuration == null)
|
||||
{
|
||||
throw new ArgumentNullException(nameof(configuration));
|
||||
}
|
||||
|
||||
configuration.Validate();
|
||||
|
||||
var host = NetworkIntegrationFactory.CreateServerHost(
|
||||
configuration.ReliablePort,
|
||||
configuration.SyncPort,
|
||||
configuration.Dispatcher,
|
||||
configuration.ReconnectPolicy,
|
||||
configuration.UtcNowProvider,
|
||||
configuration.DeliveryPolicyResolver,
|
||||
configuration.SyncSequenceTracker,
|
||||
configuration.TransportFactory);
|
||||
|
||||
try
|
||||
{
|
||||
await host.StartAsync();
|
||||
return new ServerRuntimeHandle(host);
|
||||
}
|
||||
catch
|
||||
{
|
||||
host.Stop();
|
||||
throw;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 4f683e87b65449fdb2ace6f6826dcc27
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,64 @@
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Net;
|
||||
using System.Threading.Tasks;
|
||||
using Network.NetworkApplication;
|
||||
|
||||
namespace Network.NetworkHost
|
||||
{
|
||||
public sealed class ServerRuntimeHandle : IDisposable
|
||||
{
|
||||
private readonly ServerNetworkHost host;
|
||||
private bool isStopped;
|
||||
|
||||
internal ServerRuntimeHandle(ServerNetworkHost host)
|
||||
{
|
||||
this.host = host ?? throw new ArgumentNullException(nameof(host));
|
||||
IsRunning = true;
|
||||
}
|
||||
|
||||
public ServerNetworkHost Host => host;
|
||||
|
||||
public bool IsRunning { get; private set; }
|
||||
|
||||
public IReadOnlyList<ManagedNetworkSession> ManagedSessions => host.ManagedSessions;
|
||||
|
||||
public event Action<MultiSessionLifecycleEvent> LifecycleChanged
|
||||
{
|
||||
add => host.LifecycleChanged += value;
|
||||
remove => host.LifecycleChanged -= value;
|
||||
}
|
||||
|
||||
public Task<int> DrainPendingMessagesAsync(int maxMessages = int.MaxValue)
|
||||
{
|
||||
return host.DrainPendingMessagesAsync(maxMessages);
|
||||
}
|
||||
|
||||
public void UpdateLifecycle()
|
||||
{
|
||||
host.UpdateLifecycle();
|
||||
}
|
||||
|
||||
public bool TryGetSession(IPEndPoint remoteEndPoint, out ManagedNetworkSession session)
|
||||
{
|
||||
return host.TryGetSession(remoteEndPoint, out session);
|
||||
}
|
||||
|
||||
public void Stop()
|
||||
{
|
||||
if (isStopped)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
isStopped = true;
|
||||
IsRunning = false;
|
||||
host.Stop();
|
||||
}
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
Stop();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 5580fbf7c9f748f3a33bc6043d14faea
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
@@ -0,0 +1,182 @@
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Net;
|
||||
using System.Threading.Tasks;
|
||||
using Google.Protobuf;
|
||||
using Network.Defines;
|
||||
using Network.NetworkApplication;
|
||||
using Network.NetworkHost;
|
||||
using Network.NetworkTransport;
|
||||
using NUnit.Framework;
|
||||
|
||||
namespace Tests.EditMode.Network
|
||||
{
|
||||
public class ServerRuntimeEntryPointTests
|
||||
{
|
||||
private static readonly IPEndPoint Peer = new(IPAddress.Loopback, 9100);
|
||||
|
||||
[Test]
|
||||
public void StartServerRuntimeAsync_ReliableOnly_StartsAndExposesServerHost()
|
||||
{
|
||||
var createdTransports = new Dictionary<int, FakeTransport>();
|
||||
var configuration = new ServerRuntimeConfiguration(9000)
|
||||
{
|
||||
TransportFactory = port => CreateTransport(createdTransports, port)
|
||||
};
|
||||
|
||||
using var runtime = NetworkIntegrationFactory.StartServerRuntimeAsync(configuration).GetAwaiter().GetResult();
|
||||
|
||||
Assert.That(createdTransports.Keys, Is.EquivalentTo(new[] { 9000 }));
|
||||
Assert.That(runtime.IsRunning, Is.True);
|
||||
Assert.That(runtime.Host.Transport, Is.SameAs(createdTransports[9000]));
|
||||
Assert.That(runtime.Host.SyncTransport, Is.Null);
|
||||
Assert.That(createdTransports[9000].StartCallCount, Is.EqualTo(1));
|
||||
}
|
||||
|
||||
[Test]
|
||||
public void StartServerRuntimeAsync_DualTransport_StartsBothConfiguredLanes()
|
||||
{
|
||||
var createdTransports = new Dictionary<int, FakeTransport>();
|
||||
var configuration = new ServerRuntimeConfiguration(9000)
|
||||
{
|
||||
SyncPort = 9001,
|
||||
TransportFactory = port => CreateTransport(createdTransports, port)
|
||||
};
|
||||
|
||||
using var runtime = ServerRuntimeEntryPoint.StartAsync(configuration).GetAwaiter().GetResult();
|
||||
|
||||
Assert.That(createdTransports.Keys, Is.EquivalentTo(new[] { 9000, 9001 }));
|
||||
Assert.That(runtime.Host.Transport, Is.SameAs(createdTransports[9000]));
|
||||
Assert.That(runtime.Host.SyncTransport, Is.SameAs(createdTransports[9001]));
|
||||
Assert.That(createdTransports[9000].StartCallCount, Is.EqualTo(1));
|
||||
Assert.That(createdTransports[9001].StartCallCount, Is.EqualTo(1));
|
||||
}
|
||||
|
||||
[Test]
|
||||
public void StartServerRuntimeAsync_SyncTransportStartFails_RollsBackStartedResources()
|
||||
{
|
||||
var createdTransports = new Dictionary<int, FakeTransport>();
|
||||
var configuration = new ServerRuntimeConfiguration(9000)
|
||||
{
|
||||
SyncPort = 9001,
|
||||
TransportFactory = port =>
|
||||
{
|
||||
var transport = CreateTransport(createdTransports, port);
|
||||
if (port == 9001)
|
||||
{
|
||||
transport.StartException = new InvalidOperationException("sync failed");
|
||||
}
|
||||
|
||||
return transport;
|
||||
}
|
||||
};
|
||||
|
||||
var exception = Assert.Throws<InvalidOperationException>(() =>
|
||||
ServerRuntimeEntryPoint.StartAsync(configuration).GetAwaiter().GetResult());
|
||||
|
||||
Assert.That(exception.Message, Is.EqualTo("sync failed"));
|
||||
Assert.That(createdTransports[9000].StartCallCount, Is.EqualTo(1));
|
||||
Assert.That(createdTransports[9000].StopCallCount, Is.EqualTo(1));
|
||||
Assert.That(createdTransports[9001].StartCallCount, Is.EqualTo(1));
|
||||
Assert.That(createdTransports[9001].StopCallCount, Is.EqualTo(1));
|
||||
}
|
||||
|
||||
[Test]
|
||||
public void ServerRuntimeHandle_DrainsMessages_ExposesManagedSessions_AndStopsIdempotently()
|
||||
{
|
||||
var createdTransports = new Dictionary<int, FakeTransport>();
|
||||
var configuration = new ServerRuntimeConfiguration(9000)
|
||||
{
|
||||
Dispatcher = new MainThreadNetworkDispatcher(),
|
||||
TransportFactory = port => CreateTransport(createdTransports, port)
|
||||
};
|
||||
var runtime = ServerRuntimeEntryPoint.StartAsync(configuration).GetAwaiter().GetResult();
|
||||
var handled = false;
|
||||
|
||||
runtime.Host.MessageManager.RegisterHandler(MessageType.Heartbeat, (payload, sender) =>
|
||||
{
|
||||
handled = true;
|
||||
});
|
||||
|
||||
createdTransports[9000].EmitReceive(CreateEnvelope(MessageType.Heartbeat, new Heartbeat()), Peer);
|
||||
|
||||
Assert.That(runtime.ManagedSessions.Count, Is.EqualTo(1));
|
||||
Assert.That(runtime.TryGetSession(Peer, out var session), Is.True);
|
||||
Assert.That(session.SessionManager.State, Is.EqualTo(ConnectionState.TransportConnected));
|
||||
Assert.That(runtime.Host.MessageManager.PendingMessageCount, Is.EqualTo(1));
|
||||
|
||||
runtime.DrainPendingMessagesAsync().GetAwaiter().GetResult();
|
||||
runtime.UpdateLifecycle();
|
||||
|
||||
Assert.That(handled, Is.True);
|
||||
Assert.That(runtime.Host.MessageManager.PendingMessageCount, Is.EqualTo(0));
|
||||
|
||||
runtime.Stop();
|
||||
runtime.Stop();
|
||||
|
||||
Assert.That(runtime.IsRunning, Is.False);
|
||||
Assert.That(runtime.ManagedSessions.Count, Is.EqualTo(0));
|
||||
Assert.That(createdTransports[9000].StopCallCount, Is.EqualTo(1));
|
||||
}
|
||||
|
||||
private static FakeTransport CreateTransport(IDictionary<int, FakeTransport> createdTransports, int port)
|
||||
{
|
||||
var transport = new FakeTransport();
|
||||
createdTransports.Add(port, transport);
|
||||
return transport;
|
||||
}
|
||||
|
||||
private static byte[] CreateEnvelope(MessageType type, IMessage payload)
|
||||
{
|
||||
return new Envelope
|
||||
{
|
||||
Type = (int)type,
|
||||
Payload = payload.ToByteString()
|
||||
}.ToByteArray();
|
||||
}
|
||||
|
||||
private sealed class FakeTransport : ITransport
|
||||
{
|
||||
public Exception StartException { get; set; }
|
||||
|
||||
public int StartCallCount { get; private set; }
|
||||
|
||||
public int StopCallCount { get; private set; }
|
||||
|
||||
public event Action<byte[], IPEndPoint> OnReceive;
|
||||
|
||||
public Task StartAsync()
|
||||
{
|
||||
StartCallCount++;
|
||||
if (StartException != null)
|
||||
{
|
||||
throw StartException;
|
||||
}
|
||||
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
public void Stop()
|
||||
{
|
||||
StopCallCount++;
|
||||
}
|
||||
|
||||
public void Send(byte[] data)
|
||||
{
|
||||
}
|
||||
|
||||
public void SendTo(byte[] data, IPEndPoint target)
|
||||
{
|
||||
}
|
||||
|
||||
public void SendToAll(byte[] data)
|
||||
{
|
||||
}
|
||||
|
||||
public void EmitReceive(byte[] data, IPEndPoint sender)
|
||||
{
|
||||
OnReceive?.Invoke(data, sender);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
fileFormatVersion: 2
|
||||
guid: 0a5c5d260e11429dafcee61f53a2f2d7
|
||||
MonoImporter:
|
||||
externalObjects: {}
|
||||
serializedVersion: 2
|
||||
defaultReferences: []
|
||||
executionOrder: 0
|
||||
icon: {instanceID: 0}
|
||||
userData:
|
||||
assetBundleName:
|
||||
assetBundleVariant:
|
||||
Reference in New Issue
Block a user