This commit is contained in:
SepComet
2026-03-27 13:27:14 +08:00
parent ff9ee1291f
commit e5851795c7
43 changed files with 1544 additions and 35 deletions
+3
View File
@@ -0,0 +1,3 @@
fileFormatVersion: 2
guid: 2f6e8be9b1a348c388f73388a3821247
timeCreated: 1774579341
@@ -24,4 +24,4 @@ namespace Network.Defines
};
}
}
}
}
@@ -0,0 +1,16 @@
using System;
namespace Network.NetworkApplication
{
public enum ConnectionState
{
Disconnected = 0,
TransportConnected = 1,
LoginPending = 2,
LoggedIn = 3,
LoginFailed = 4,
TimedOut = 5,
ReconnectPending = 6,
Reconnecting = 7,
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 4db1cb66c3f181d43b233641c5c42b0d
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,18 @@
using System;
using System.Net;
namespace Network.NetworkApplication
{
public sealed class ManagedNetworkSession
{
public ManagedNetworkSession(IPEndPoint remoteEndPoint, SessionManager sessionManager)
{
RemoteEndPoint = remoteEndPoint ?? throw new ArgumentNullException(nameof(remoteEndPoint));
SessionManager = sessionManager ?? throw new ArgumentNullException(nameof(sessionManager));
}
public IPEndPoint RemoteEndPoint { get; }
public SessionManager SessionManager { get; }
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 41dcf02596d3304418a07f265ae8f525
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,24 @@
using System;
using System.Net;
namespace Network.NetworkApplication
{
public sealed class MultiSessionLifecycleEvent
{
public MultiSessionLifecycleEvent(
IPEndPoint remoteEndPoint,
SessionManager sessionManager,
SessionLifecycleEvent lifecycleEvent)
{
RemoteEndPoint = remoteEndPoint ?? throw new ArgumentNullException(nameof(remoteEndPoint));
SessionManager = sessionManager ?? throw new ArgumentNullException(nameof(sessionManager));
LifecycleEvent = lifecycleEvent ?? throw new ArgumentNullException(nameof(lifecycleEvent));
}
public IPEndPoint RemoteEndPoint { get; }
public SessionManager SessionManager { get; }
public SessionLifecycleEvent LifecycleEvent { get; }
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 82ee5b8b86486b04eab6d76cc6859cc5
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,238 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Net;
namespace Network.NetworkApplication
{
public sealed class MultiSessionManager
{
private readonly object gate = new();
private readonly Dictionary<string, SessionRegistration> sessions = new();
private readonly SessionReconnectPolicy reconnectPolicy;
private readonly Func<DateTimeOffset> utcNowProvider;
public MultiSessionManager(
SessionReconnectPolicy reconnectPolicy = null,
Func<DateTimeOffset> utcNowProvider = null)
{
this.reconnectPolicy = reconnectPolicy ?? SessionReconnectPolicy.Default;
this.utcNowProvider = utcNowProvider ?? (() => DateTimeOffset.UtcNow);
}
public event Action<MultiSessionLifecycleEvent> LifecycleChanged;
public int SessionCount
{
get
{
lock (gate)
{
return sessions.Count;
}
}
}
public IReadOnlyList<ManagedNetworkSession> Sessions
{
get
{
lock (gate)
{
return sessions.Values
.Select(registration => registration.Session)
.ToArray();
}
}
}
public ManagedNetworkSession GetOrCreateSession(IPEndPoint remoteEndPoint)
{
return GetOrCreateRegistration(remoteEndPoint).Session;
}
public bool TryGetSession(IPEndPoint remoteEndPoint, out ManagedNetworkSession session)
{
var key = BuildKey(remoteEndPoint);
lock (gate)
{
if (sessions.TryGetValue(key, out var registration))
{
session = registration.Session;
return true;
}
}
session = null;
return false;
}
public bool TryGetSessionManager(IPEndPoint remoteEndPoint, out SessionManager sessionManager)
{
if (TryGetSession(remoteEndPoint, out var session))
{
sessionManager = session.SessionManager;
return true;
}
sessionManager = null;
return false;
}
public void ObserveTransportActivity(IPEndPoint remoteEndPoint)
{
var sessionManager = GetOrCreateSession(remoteEndPoint).SessionManager;
if (sessionManager.State == ConnectionState.Disconnected)
{
sessionManager.NotifyTransportConnected();
}
sessionManager.NotifyInboundActivity();
}
public void NotifyTransportConnected(IPEndPoint remoteEndPoint)
{
GetOrCreateSession(remoteEndPoint).SessionManager.NotifyTransportConnected();
}
public void NotifyLoginStarted(IPEndPoint remoteEndPoint)
{
GetOrCreateSession(remoteEndPoint).SessionManager.NotifyLoginStarted();
}
public void NotifyLoginSucceeded(IPEndPoint remoteEndPoint)
{
GetOrCreateSession(remoteEndPoint).SessionManager.NotifyLoginSucceeded();
}
public void NotifyLoginFailed(IPEndPoint remoteEndPoint, string reason = null)
{
GetOrCreateSession(remoteEndPoint).SessionManager.NotifyLoginFailed(reason);
}
public void NotifyHeartbeatSent(IPEndPoint remoteEndPoint)
{
GetOrCreateSession(remoteEndPoint).SessionManager.NotifyHeartbeatSent();
}
public void NotifyHeartbeatReceived(IPEndPoint remoteEndPoint, long? serverTick = null)
{
GetOrCreateSession(remoteEndPoint).SessionManager.NotifyHeartbeatReceived(serverTick);
}
public void NotifyInboundActivity(IPEndPoint remoteEndPoint)
{
GetOrCreateSession(remoteEndPoint).SessionManager.NotifyInboundActivity();
}
public bool RemoveSession(IPEndPoint remoteEndPoint, string reason = null)
{
SessionRegistration registration;
var key = BuildKey(remoteEndPoint);
lock (gate)
{
if (!sessions.TryGetValue(key, out registration))
{
return false;
}
sessions.Remove(key);
}
registration.Session.SessionManager.NotifyTransportDisconnected(reason);
registration.Session.SessionManager.LifecycleChanged -= registration.Handler;
return true;
}
public void RemoveAllSessions(string reason = null)
{
SessionRegistration[] registrations;
lock (gate)
{
registrations = sessions.Values.ToArray();
sessions.Clear();
}
foreach (var registration in registrations)
{
registration.Session.SessionManager.NotifyTransportDisconnected(reason);
registration.Session.SessionManager.LifecycleChanged -= registration.Handler;
}
}
public void UpdateLifecycle()
{
SessionManager[] activeSessions;
lock (gate)
{
activeSessions = sessions.Values
.Select(registration => registration.Session.SessionManager)
.ToArray();
}
foreach (var session in activeSessions)
{
session.Evaluate();
}
}
private SessionRegistration GetOrCreateRegistration(IPEndPoint remoteEndPoint)
{
var normalizedEndPoint = Normalize(remoteEndPoint);
var key = normalizedEndPoint.ToString();
lock (gate)
{
if (sessions.TryGetValue(key, out var registration))
{
return registration;
}
var sessionManager = new SessionManager(reconnectPolicy, utcNowProvider);
var session = new ManagedNetworkSession(normalizedEndPoint, sessionManager);
Action<SessionLifecycleEvent> handler = lifecycleEvent =>
LifecycleChanged?.Invoke(new MultiSessionLifecycleEvent(session.RemoteEndPoint, session.SessionManager, lifecycleEvent));
sessionManager.LifecycleChanged += handler;
registration = new SessionRegistration(session, handler);
sessions.Add(key, registration);
return registration;
}
}
private static string BuildKey(IPEndPoint remoteEndPoint)
{
return Normalize(remoteEndPoint).ToString();
}
private static IPEndPoint Normalize(IPEndPoint remoteEndPoint)
{
if (remoteEndPoint == null)
{
throw new ArgumentNullException(nameof(remoteEndPoint));
}
return new IPEndPoint(remoteEndPoint.Address, remoteEndPoint.Port);
}
private sealed class SessionRegistration
{
public SessionRegistration(ManagedNetworkSession session, Action<SessionLifecycleEvent> handler)
{
Session = session;
Handler = handler;
}
public ManagedNetworkSession Session { get; }
public Action<SessionLifecycleEvent> Handler { get; }
}
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 607451432d1e802409107e2a816b3587
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,16 @@
namespace Network.NetworkApplication
{
public enum SessionEventKind
{
TransportConnected = 0,
LoginStarted = 1,
LoginSucceeded = 2,
LoginFailed = 3,
HeartbeatSent = 4,
HeartbeatReceived = 5,
TimedOut = 6,
ReconnectScheduled = 7,
ReconnectStarted = 8,
Disconnected = 9,
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: f9f3ec0c5a33d23478a9b1848bd05431
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,31 @@
using System;
namespace Network.NetworkApplication
{
public sealed class SessionLifecycleEvent
{
public SessionLifecycleEvent(
SessionEventKind kind,
ConnectionState previousState,
ConnectionState currentState,
DateTimeOffset occurredAtUtc,
string reason = null)
{
Kind = kind;
PreviousState = previousState;
CurrentState = currentState;
OccurredAtUtc = occurredAtUtc;
Reason = reason;
}
public SessionEventKind Kind { get; }
public ConnectionState PreviousState { get; }
public ConnectionState CurrentState { get; }
public DateTimeOffset OccurredAtUtc { get; }
public string Reason { get; }
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 549538bb90b911142b9e627bd38df213
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,198 @@
using System;
namespace Network.NetworkApplication
{
public sealed class SessionManager
{
private readonly Func<DateTimeOffset> utcNowProvider;
private DateTimeOffset? lastLivenessUtc;
private DateTimeOffset? lastHeartbeatSentUtc;
private DateTimeOffset? nextReconnectAtUtc;
public SessionManager(
SessionReconnectPolicy reconnectPolicy = null,
Func<DateTimeOffset> utcNowProvider = null)
{
ReconnectPolicy = reconnectPolicy ?? SessionReconnectPolicy.Default;
this.utcNowProvider = utcNowProvider ?? (() => DateTimeOffset.UtcNow);
State = ConnectionState.Disconnected;
}
public event Action<SessionLifecycleEvent> LifecycleChanged;
public ConnectionState State { get; private set; }
public SessionReconnectPolicy ReconnectPolicy { get; }
public DateTimeOffset? LastLivenessUtc => lastLivenessUtc;
public DateTimeOffset? LastHeartbeatSentUtc => lastHeartbeatSentUtc;
public DateTimeOffset? NextReconnectAtUtc => nextReconnectAtUtc;
public TimeSpan? LastRoundTripTime { get; private set; }
public long? LastServerTick { get; private set; }
public string LastFailureReason { get; private set; }
public bool CanSendHeartbeat => State == ConnectionState.LoggedIn;
public bool IsReconnectDue
{
get
{
if (State != ConnectionState.ReconnectPending || nextReconnectAtUtc == null)
{
return false;
}
return utcNowProvider() >= nextReconnectAtUtc.Value;
}
}
public bool IsHeartbeatDue
{
get
{
if (!CanSendHeartbeat)
{
return false;
}
if (lastHeartbeatSentUtc == null)
{
return true;
}
return utcNowProvider() - lastHeartbeatSentUtc.Value >= ReconnectPolicy.HeartbeatInterval;
}
}
public void NotifyTransportConnected()
{
var now = utcNowProvider();
lastLivenessUtc = now;
lastHeartbeatSentUtc = null;
nextReconnectAtUtc = null;
LastFailureReason = null;
TransitionTo(ConnectionState.TransportConnected, SessionEventKind.TransportConnected, now);
}
public void NotifyLoginStarted()
{
TransitionTo(ConnectionState.LoginPending, SessionEventKind.LoginStarted, utcNowProvider());
}
public void NotifyLoginSucceeded()
{
var now = utcNowProvider();
lastLivenessUtc = now;
LastFailureReason = null;
TransitionTo(ConnectionState.LoggedIn, SessionEventKind.LoginSucceeded, now);
}
public void NotifyLoginFailed(string reason = null)
{
LastFailureReason = reason;
TransitionTo(ConnectionState.LoginFailed, SessionEventKind.LoginFailed, utcNowProvider(), reason);
}
public void NotifyHeartbeatSent()
{
lastHeartbeatSentUtc = utcNowProvider();
RaiseEvent(SessionEventKind.HeartbeatSent, State, State, lastHeartbeatSentUtc.Value);
}
public void NotifyHeartbeatReceived(long? serverTick = null)
{
var now = utcNowProvider();
lastLivenessUtc = now;
LastServerTick = serverTick;
if (lastHeartbeatSentUtc.HasValue)
{
LastRoundTripTime = now - lastHeartbeatSentUtc.Value;
}
RaiseEvent(SessionEventKind.HeartbeatReceived, State, State, now);
}
public void NotifyInboundActivity()
{
lastLivenessUtc = utcNowProvider();
}
public void NotifyTransportDisconnected(string reason = null)
{
LastFailureReason = reason;
nextReconnectAtUtc = null;
TransitionTo(ConnectionState.Disconnected, SessionEventKind.Disconnected, utcNowProvider(), reason);
}
public void Evaluate()
{
var now = utcNowProvider();
if (ShouldTimeout(now))
{
LastFailureReason = "Heartbeat timeout";
TransitionTo(ConnectionState.TimedOut, SessionEventKind.TimedOut, now, LastFailureReason);
if (ReconnectPolicy.AutoReconnect)
{
nextReconnectAtUtc = now + ReconnectPolicy.ReconnectDelay;
TransitionTo(ConnectionState.ReconnectPending, SessionEventKind.ReconnectScheduled, now, LastFailureReason);
}
return;
}
if (State == ConnectionState.ReconnectPending && nextReconnectAtUtc.HasValue && now >= nextReconnectAtUtc.Value)
{
TransitionTo(ConnectionState.Reconnecting, SessionEventKind.ReconnectStarted, now, LastFailureReason);
}
}
private bool ShouldTimeout(DateTimeOffset now)
{
if (State != ConnectionState.TransportConnected && State != ConnectionState.LoginPending && State != ConnectionState.LoggedIn)
{
return false;
}
if (!lastLivenessUtc.HasValue)
{
return false;
}
return now - lastLivenessUtc.Value >= ReconnectPolicy.HeartbeatTimeout;
}
private void TransitionTo(
ConnectionState newState,
SessionEventKind eventKind,
DateTimeOffset occurredAtUtc,
string reason = null)
{
if (State == newState)
{
RaiseEvent(eventKind, State, State, occurredAtUtc, reason);
return;
}
var previousState = State;
State = newState;
RaiseEvent(eventKind, previousState, newState, occurredAtUtc, reason);
}
private void RaiseEvent(
SessionEventKind kind,
ConnectionState previousState,
ConnectionState currentState,
DateTimeOffset occurredAtUtc,
string reason = null)
{
LifecycleChanged?.Invoke(new SessionLifecycleEvent(kind, previousState, currentState, occurredAtUtc, reason));
}
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: ce1fd028b318b084e8dc5a3084f9e9c1
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,48 @@
using System;
namespace Network.NetworkApplication
{
public sealed class SessionReconnectPolicy
{
public static SessionReconnectPolicy Default { get; } = new(
heartbeatInterval: TimeSpan.FromSeconds(2),
heartbeatTimeout: TimeSpan.FromSeconds(6),
reconnectDelay: TimeSpan.FromSeconds(1),
autoReconnect: true);
public SessionReconnectPolicy(
TimeSpan heartbeatInterval,
TimeSpan heartbeatTimeout,
TimeSpan reconnectDelay,
bool autoReconnect)
{
if (heartbeatInterval <= TimeSpan.Zero)
{
throw new ArgumentOutOfRangeException(nameof(heartbeatInterval));
}
if (heartbeatTimeout <= TimeSpan.Zero)
{
throw new ArgumentOutOfRangeException(nameof(heartbeatTimeout));
}
if (reconnectDelay < TimeSpan.Zero)
{
throw new ArgumentOutOfRangeException(nameof(reconnectDelay));
}
HeartbeatInterval = heartbeatInterval;
HeartbeatTimeout = heartbeatTimeout;
ReconnectDelay = reconnectDelay;
AutoReconnect = autoReconnect;
}
public TimeSpan HeartbeatInterval { get; }
public TimeSpan HeartbeatTimeout { get; }
public TimeSpan ReconnectDelay { get; }
public bool AutoReconnect { get; }
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 1dbefe0de75d3904ca75f2dbf73b61d1
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -1,4 +1,4 @@
using System;
using System;
using System.Threading.Tasks;
using Network.NetworkTransport;
@@ -6,9 +6,14 @@ namespace Network.NetworkApplication
{
public sealed class SharedNetworkRuntime
{
public SharedNetworkRuntime(ITransport transport, INetworkMessageDispatcher dispatcher)
public SharedNetworkRuntime(
ITransport transport,
INetworkMessageDispatcher dispatcher,
SessionReconnectPolicy reconnectPolicy = null,
Func<DateTimeOffset> utcNowProvider = null)
{
Transport = transport ?? throw new ArgumentNullException(nameof(transport));
SessionManager = new SessionManager(reconnectPolicy, utcNowProvider);
MessageManager = new MessageManager(transport, dispatcher ?? throw new ArgumentNullException(nameof(dispatcher)));
}
@@ -16,19 +21,64 @@ namespace Network.NetworkApplication
public MessageManager MessageManager { get; }
public Task StartAsync()
public SessionManager SessionManager { get; }
public event Action<SessionLifecycleEvent> LifecycleChanged
{
return Transport.StartAsync();
add => SessionManager.LifecycleChanged += value;
remove => SessionManager.LifecycleChanged -= value;
}
public async Task StartAsync()
{
await Transport.StartAsync();
SessionManager.NotifyTransportConnected();
}
public void Stop()
{
Transport.Stop();
SessionManager.NotifyTransportDisconnected("Transport stopped");
}
public Task<int> DrainPendingMessagesAsync(int maxMessages = int.MaxValue)
{
return MessageManager.DrainPendingMessagesAsync(maxMessages);
}
public void NotifyLoginStarted()
{
SessionManager.NotifyLoginStarted();
}
public void NotifyLoginSucceeded()
{
SessionManager.NotifyLoginSucceeded();
}
public void NotifyLoginFailed(string reason = null)
{
SessionManager.NotifyLoginFailed(reason);
}
public void NotifyHeartbeatSent()
{
SessionManager.NotifyHeartbeatSent();
}
public void NotifyHeartbeatReceived(long? serverTick = null)
{
SessionManager.NotifyHeartbeatReceived(serverTick);
}
public void NotifyInboundActivity()
{
SessionManager.NotifyInboundActivity();
}
public void UpdateLifecycle()
{
SessionManager.Evaluate();
}
}
}
@@ -1,4 +1,6 @@
using System;
using System;
using System.Collections.Generic;
using System.Net;
using System.Threading.Tasks;
using Network.NetworkApplication;
using Network.NetworkTransport;
@@ -7,32 +9,101 @@ namespace Network.NetworkHost
{
public sealed class ServerNetworkHost
{
private readonly SharedNetworkRuntime runtime;
private readonly ITransport transport;
private readonly MessageManager messageManager;
public ServerNetworkHost(ITransport transport, INetworkMessageDispatcher dispatcher = null)
public ServerNetworkHost(
ITransport transport,
INetworkMessageDispatcher dispatcher = null,
SessionReconnectPolicy reconnectPolicy = null,
Func<DateTimeOffset> utcNowProvider = null)
{
runtime = new SharedNetworkRuntime(
transport ?? throw new ArgumentNullException(nameof(transport)),
dispatcher ?? new ImmediateNetworkMessageDispatcher());
this.transport = transport ?? throw new ArgumentNullException(nameof(transport));
SessionCoordinator = new MultiSessionManager(reconnectPolicy, utcNowProvider);
this.transport.OnReceive += HandleTransportReceive;
messageManager = new MessageManager(this.transport, dispatcher ?? new ImmediateNetworkMessageDispatcher());
}
public MessageManager MessageManager => runtime.MessageManager;
public MessageManager MessageManager => messageManager;
public ITransport Transport => runtime.Transport;
public ITransport Transport => transport;
// Server-side lifecycle entry point: inspect and control per-peer session state here.
public MultiSessionManager SessionCoordinator { get; }
public IReadOnlyList<ManagedNetworkSession> ManagedSessions => SessionCoordinator.Sessions;
public event Action<MultiSessionLifecycleEvent> LifecycleChanged
{
add => SessionCoordinator.LifecycleChanged += value;
remove => SessionCoordinator.LifecycleChanged -= value;
}
public Task StartAsync()
{
return runtime.StartAsync();
return transport.StartAsync();
}
public void Stop()
{
runtime.Stop();
transport.Stop();
SessionCoordinator.RemoveAllSessions("Transport stopped");
}
public Task<int> DrainPendingMessagesAsync(int maxMessages = int.MaxValue)
{
return runtime.DrainPendingMessagesAsync(maxMessages);
return messageManager.DrainPendingMessagesAsync(maxMessages);
}
public void UpdateLifecycle()
{
SessionCoordinator.UpdateLifecycle();
}
public bool TryGetSession(IPEndPoint remoteEndPoint, out ManagedNetworkSession session)
{
return SessionCoordinator.TryGetSession(remoteEndPoint, out session);
}
public void NotifyLoginStarted(IPEndPoint remoteEndPoint)
{
SessionCoordinator.NotifyLoginStarted(remoteEndPoint);
}
public void NotifyLoginSucceeded(IPEndPoint remoteEndPoint)
{
SessionCoordinator.NotifyLoginSucceeded(remoteEndPoint);
}
public void NotifyLoginFailed(IPEndPoint remoteEndPoint, string reason = null)
{
SessionCoordinator.NotifyLoginFailed(remoteEndPoint, reason);
}
public void NotifyHeartbeatSent(IPEndPoint remoteEndPoint)
{
SessionCoordinator.NotifyHeartbeatSent(remoteEndPoint);
}
public void NotifyHeartbeatReceived(IPEndPoint remoteEndPoint, long? serverTick = null)
{
SessionCoordinator.NotifyHeartbeatReceived(remoteEndPoint, serverTick);
}
public void NotifyInboundActivity(IPEndPoint remoteEndPoint)
{
SessionCoordinator.NotifyInboundActivity(remoteEndPoint);
}
public bool RemoveSession(IPEndPoint remoteEndPoint, string reason = null)
{
return SessionCoordinator.RemoveSession(remoteEndPoint, reason);
}
private void HandleTransportReceive(byte[] _, IPEndPoint sender)
{
SessionCoordinator.ObserveTransportActivity(sender);
}
}
}
+27 -3
View File
@@ -29,6 +29,7 @@ public class NetworkManager : MonoBehaviour
var transport = new KcpTransport("127.0.0.1", 8080);
var dispatcher = new MainThreadNetworkDispatcher();
_networkRuntime = new SharedNetworkRuntime(transport, dispatcher);
_networkRuntime.LifecycleChanged += HandleLifecycleChanged;
var startTask = _networkRuntime.StartAsync();
yield return new WaitUntil(() => startTask.IsCompleted);
@@ -50,6 +51,8 @@ public class NetworkManager : MonoBehaviour
return;
}
_networkRuntime.UpdateLifecycle();
if (!_networkDrainTask.IsCompleted)
{
return;
@@ -65,7 +68,12 @@ public class NetworkManager : MonoBehaviour
private void OnDestroy()
{
_networkRuntime?.Stop();
if (_networkRuntime != null)
{
_networkRuntime.LifecycleChanged -= HandleLifecycleChanged;
_networkRuntime.Stop();
}
if (Instance == this)
{
Instance = null;
@@ -76,13 +84,16 @@ public class NetworkManager : MonoBehaviour
{
while (true)
{
if (_serverPoint != null)
if (_networkRuntime != null
&& _serverPoint != null
&& _networkRuntime.SessionManager.IsHeartbeatDue)
{
var heartbeat = new Heartbeat();
_networkRuntime.MessageManager.SendMessage(heartbeat, MessageType.Heartbeat, _serverPoint);
_networkRuntime.NotifyHeartbeatSent();
}
yield return new WaitForSeconds(2.0f);
yield return new WaitForSeconds(0.25f);
}
}
@@ -98,13 +109,16 @@ public class NetworkManager : MonoBehaviour
private void HandleLoginResponse(byte[] data, IPEndPoint sender)
{
var response = LoginResponse.Parser.ParseFrom(data);
_networkRuntime.NotifyInboundActivity();
_serverPoint = sender;
if (response.Result)
{
_networkRuntime.NotifyLoginSucceeded();
MasterManager.Instance.InitPlayersState(response);
}
else
{
_networkRuntime.NotifyLoginFailed("UserId already exists");
_wrongWindow.SetActive(true);
Debug.LogError("UserId 已经存在");
}
@@ -112,6 +126,7 @@ public class NetworkManager : MonoBehaviour
private void HandlePlayerState(byte[] data, IPEndPoint sender)
{
_networkRuntime.NotifyInboundActivity();
var message = PlayerState.Parser.ParseFrom(data);
MasterManager.Instance.MovePlayer(message.PlayerId, message);
Debug.Log($"收到PlayerState::PlayerID={message.PlayerId},Position=" + message.Position.ToVector3().ToString());
@@ -120,6 +135,7 @@ public class NetworkManager : MonoBehaviour
private void HandleHeartbeatResponse(byte[] data, IPEndPoint sender)
{
var response = HeartbeatResponse.Parser.ParseFrom(data);
_networkRuntime.NotifyHeartbeatReceived(response.ServerTick);
var player = MasterManager.Instance.GetCurrentPlayer();
if (player != null)
{
@@ -129,17 +145,24 @@ public class NetworkManager : MonoBehaviour
private void HandleLogoutRequest(byte[] data, IPEndPoint sender)
{
_networkRuntime.NotifyInboundActivity();
var request = LogoutRequest.Parser.ParseFrom(data);
MasterManager.Instance.UnregisterPlayer(request.PlayerId);
}
private void HandlePlayerJoin(byte[] data, IPEndPoint sender)
{
_networkRuntime.NotifyInboundActivity();
var playerJoin = PlayerJoin.Parser.ParseFrom(data);
if (MasterManager.Instance.LocalPlayerId == playerJoin.PlayerId) return;
MasterManager.Instance.RegisterRemotePlayer(playerJoin.PlayerId, playerJoin.Position.ToVector3());
}
private void HandleLifecycleChanged(SessionLifecycleEvent lifecycleEvent)
{
Debug.Log($"[NetworkManager] Session {lifecycleEvent.PreviousState} -> {lifecycleEvent.CurrentState} ({lifecycleEvent.Kind}) {lifecycleEvent.Reason}");
}
public void SendPlayerInput(string playerId, Vector3 input)
{
var message = new PlayerInput()
@@ -164,6 +187,7 @@ public class NetworkManager : MonoBehaviour
PlayerId = playerId,
Speed = speed
};
_networkRuntime.NotifyLoginStarted();
_networkRuntime.MessageManager.SendMessage(request, MessageType.LoginRequest);
Debug.Log($"Sent login request to player {playerId}");
}
@@ -0,0 +1,226 @@
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 SessionLifecycleTests
{
[Test]
public void SharedNetworkRuntime_StartAsync_TransitionsToTransportConnectedButNotLoggedIn()
{
var transport = new FakeTransport();
var runtime = new SharedNetworkRuntime(transport, new ImmediateNetworkMessageDispatcher());
runtime.StartAsync().GetAwaiter().GetResult();
Assert.That(runtime.SessionManager.State, Is.EqualTo(ConnectionState.TransportConnected));
Assert.That(runtime.SessionManager.CanSendHeartbeat, Is.False);
}
[Test]
public void LoginFailure_IsDistinctFromTransportConnectedState()
{
var transport = new FakeTransport();
var runtime = new SharedNetworkRuntime(transport, new ImmediateNetworkMessageDispatcher());
runtime.StartAsync().GetAwaiter().GetResult();
runtime.NotifyLoginStarted();
runtime.NotifyLoginFailed("bad credentials");
Assert.That(runtime.SessionManager.State, Is.EqualTo(ConnectionState.LoginFailed));
Assert.That(runtime.SessionManager.LastFailureReason, Is.EqualTo("bad credentials"));
}
[Test]
public void HeartbeatTimeout_SchedulesAndStartsReconnect()
{
var clock = new MutableClock(new DateTimeOffset(2026, 3, 27, 0, 0, 0, TimeSpan.Zero));
var policy = new SessionReconnectPolicy(
heartbeatInterval: TimeSpan.FromSeconds(2),
heartbeatTimeout: TimeSpan.FromSeconds(5),
reconnectDelay: TimeSpan.FromSeconds(3),
autoReconnect: true);
var transport = new FakeTransport();
var runtime = new SharedNetworkRuntime(transport, new ImmediateNetworkMessageDispatcher(), policy, clock.UtcNow);
var events = new List<SessionEventKind>();
runtime.LifecycleChanged += lifecycleEvent => events.Add(lifecycleEvent.Kind);
runtime.StartAsync().GetAwaiter().GetResult();
runtime.NotifyLoginStarted();
runtime.NotifyLoginSucceeded();
clock.Advance(TimeSpan.FromSeconds(6));
runtime.UpdateLifecycle();
Assert.That(runtime.SessionManager.State, Is.EqualTo(ConnectionState.ReconnectPending));
Assert.That(events, Does.Contain(SessionEventKind.TimedOut));
Assert.That(events, Does.Contain(SessionEventKind.ReconnectScheduled));
clock.Advance(TimeSpan.FromSeconds(3));
runtime.UpdateLifecycle();
Assert.That(runtime.SessionManager.State, Is.EqualTo(ConnectionState.Reconnecting));
Assert.That(events, Does.Contain(SessionEventKind.ReconnectStarted));
}
[Test]
public void HeartbeatResponse_UpdatesRttAndServerTick_WithoutChangingLoggedInState()
{
var clock = new MutableClock(new DateTimeOffset(2026, 3, 27, 0, 0, 0, TimeSpan.Zero));
var transport = new FakeTransport();
var runtime = new SharedNetworkRuntime(transport, new ImmediateNetworkMessageDispatcher(), utcNowProvider: clock.UtcNow);
runtime.StartAsync().GetAwaiter().GetResult();
runtime.NotifyLoginStarted();
runtime.NotifyLoginSucceeded();
runtime.NotifyHeartbeatSent();
clock.Advance(TimeSpan.FromMilliseconds(120));
runtime.NotifyHeartbeatReceived(321);
Assert.That(runtime.SessionManager.State, Is.EqualTo(ConnectionState.LoggedIn));
Assert.That(runtime.SessionManager.LastRoundTripTime, Is.EqualTo(TimeSpan.FromMilliseconds(120)));
Assert.That(runtime.SessionManager.LastServerTick, Is.EqualTo(321));
}
[Test]
public void ServerNetworkHost_TracksMultipleSessionsIndependently()
{
var clock = new MutableClock(new DateTimeOffset(2026, 3, 27, 0, 0, 0, TimeSpan.Zero));
var policy = new SessionReconnectPolicy(
heartbeatInterval: TimeSpan.FromSeconds(2),
heartbeatTimeout: TimeSpan.FromSeconds(5),
reconnectDelay: TimeSpan.FromSeconds(3),
autoReconnect: true);
var transport = new FakeTransport();
var host = new ServerNetworkHost(transport, reconnectPolicy: policy, utcNowProvider: clock.UtcNow);
var peerA = new IPEndPoint(IPAddress.Parse("127.0.0.1"), 5001);
var peerB = new IPEndPoint(IPAddress.Parse("127.0.0.1"), 5002);
host.StartAsync().GetAwaiter().GetResult();
transport.EmitReceive(CreateEnvelope(MessageType.Heartbeat), peerA);
transport.EmitReceive(CreateEnvelope(MessageType.Heartbeat), peerB);
host.NotifyLoginStarted(peerA);
host.NotifyLoginSucceeded(peerA);
host.NotifyLoginStarted(peerB);
host.NotifyLoginSucceeded(peerB);
clock.Advance(TimeSpan.FromSeconds(6));
host.NotifyHeartbeatReceived(peerB, 99);
host.UpdateLifecycle();
Assert.That(host.ManagedSessions.Count, Is.EqualTo(2));
Assert.That(host.TryGetSession(peerA, out var sessionA), Is.True);
Assert.That(host.TryGetSession(peerB, out var sessionB), Is.True);
Assert.That(sessionA.SessionManager.State, Is.EqualTo(ConnectionState.ReconnectPending));
Assert.That(sessionB.SessionManager.State, Is.EqualTo(ConnectionState.LoggedIn));
Assert.That(sessionB.SessionManager.LastServerTick, Is.EqualTo(99));
}
[Test]
public void ServerNetworkHost_RemoveSession_DoesNotDisturbOtherPeers()
{
var transport = new FakeTransport();
var host = new ServerNetworkHost(transport);
var peerA = new IPEndPoint(IPAddress.Parse("127.0.0.1"), 5001);
var peerB = new IPEndPoint(IPAddress.Parse("127.0.0.1"), 5002);
host.StartAsync().GetAwaiter().GetResult();
transport.EmitReceive(CreateEnvelope(MessageType.Heartbeat), peerA);
transport.EmitReceive(CreateEnvelope(MessageType.Heartbeat), peerB);
var removed = host.RemoveSession(peerA, "peer closed");
Assert.That(removed, Is.True);
Assert.That(host.ManagedSessions.Count, Is.EqualTo(1));
Assert.That(host.TryGetSession(peerA, out _), Is.False);
Assert.That(host.TryGetSession(peerB, out var sessionB), Is.True);
Assert.That(sessionB.SessionManager.State, Is.EqualTo(ConnectionState.TransportConnected));
}
[Test]
public void ServerNetworkHost_LifecycleEventsIncludeRemotePeerIdentity()
{
var transport = new FakeTransport();
var host = new ServerNetworkHost(transport);
var peer = new IPEndPoint(IPAddress.Parse("127.0.0.1"), 5001);
MultiSessionLifecycleEvent receivedEvent = null;
host.LifecycleChanged += lifecycleEvent => receivedEvent = lifecycleEvent;
host.StartAsync().GetAwaiter().GetResult();
transport.EmitReceive(CreateEnvelope(MessageType.Heartbeat), peer);
Assert.That(receivedEvent, Is.Not.Null);
Assert.That(receivedEvent.RemoteEndPoint, Is.EqualTo(peer));
Assert.That(receivedEvent.LifecycleEvent.CurrentState, Is.EqualTo(ConnectionState.TransportConnected));
}
private static byte[] CreateEnvelope(MessageType type)
{
return new Envelope
{
Type = (int)type
}.ToByteArray();
}
private sealed class MutableClock
{
public MutableClock(DateTimeOffset now)
{
Now = now;
}
public DateTimeOffset Now { get; private set; }
public DateTimeOffset UtcNow()
{
return Now;
}
public void Advance(TimeSpan delta)
{
Now = Now.Add(delta);
}
}
private sealed class FakeTransport : ITransport
{
public event Action<byte[], IPEndPoint> OnReceive;
public Task StartAsync()
{
return Task.CompletedTask;
}
public void Stop()
{
}
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: 44395d472582e2c41a8c31b766751035
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant: