Network 层的迁移与强化已完成,TODO.md 与 MobaSyncMVP.md 给出后续方向

This commit is contained in:
SepComet
2026-03-27 22:14:33 +08:00
parent e361510100
commit 156d72bf4a
57 changed files with 2933 additions and 102 deletions
@@ -0,0 +1,50 @@
using System;
using System.Collections.Generic;
using System.Linq;
using Network.Defines;
namespace Network.NetworkApplication
{
public sealed class ClientPredictionBuffer
{
private readonly List<PlayerInput> pendingInputs = new();
public long? LastAuthoritativeTick { get; private set; }
public IReadOnlyList<PlayerInput> PendingInputs => pendingInputs;
public void Record(PlayerInput input)
{
if (input == null)
{
throw new ArgumentNullException(nameof(input));
}
if (pendingInputs.Count > 0 && pendingInputs[^1].Tick >= input.Tick)
{
return;
}
pendingInputs.Add(input);
}
public bool TryApplyAuthoritativeState(PlayerState state, out IReadOnlyList<PlayerInput> replayInputs)
{
if (state == null)
{
throw new ArgumentNullException(nameof(state));
}
if (LastAuthoritativeTick.HasValue && state.Tick <= LastAuthoritativeTick.Value)
{
replayInputs = Array.Empty<PlayerInput>();
return false;
}
LastAuthoritativeTick = state.Tick;
pendingInputs.RemoveAll(input => input.Tick <= state.Tick);
replayInputs = pendingInputs.ToArray();
return true;
}
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 0835bc8e8b5c4f14eab1d4c429ec6238
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,35 @@
using System;
namespace Network.NetworkApplication
{
public sealed class ClockSyncState
{
private readonly Func<DateTimeOffset> utcNowProvider;
public ClockSyncState(Func<DateTimeOffset> utcNowProvider = null)
{
this.utcNowProvider = utcNowProvider ?? (() => DateTimeOffset.UtcNow);
}
public long? CurrentServerTick { get; private set; }
public DateTimeOffset? LastSampleReceivedAtUtc { get; private set; }
public bool ObserveSample(long? serverTick)
{
if (!serverTick.HasValue)
{
return false;
}
if (CurrentServerTick.HasValue && serverTick.Value < CurrentServerTick.Value)
{
return false;
}
CurrentServerTick = serverTick.Value;
LastSampleReceivedAtUtc = utcNowProvider();
return true;
}
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 95a64b5509806e94aac0803ca2a9823f
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,22 @@
using System.Collections.Generic;
using Network.Defines;
namespace Network.NetworkApplication
{
public sealed class DefaultMessageDeliveryPolicyResolver : IMessageDeliveryPolicyResolver
{
private static readonly IReadOnlyDictionary<MessageType, DeliveryPolicy> DefaultPolicies =
new Dictionary<MessageType, DeliveryPolicy>
{
{ MessageType.PlayerInput, DeliveryPolicy.HighFrequencySync },
{ MessageType.PlayerState, DeliveryPolicy.HighFrequencySync }
};
public DeliveryPolicy Resolve(MessageType messageType)
{
return DefaultPolicies.TryGetValue(messageType, out var policy)
? policy
: DeliveryPolicy.ReliableOrdered;
}
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 8b39164ada6225c4db48ec637634ba86
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,8 @@
namespace Network.NetworkApplication
{
public enum DeliveryPolicy
{
ReliableOrdered = 0,
HighFrequencySync = 1
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 2e26e6d78859ced41a9e1241414e5953
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,9 @@
using Network.Defines;
namespace Network.NetworkApplication
{
public interface IMessageDeliveryPolicyResolver
{
DeliveryPolicy Resolve(MessageType messageType);
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 7732b6dea776deb43b39d6cb42874efc
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,16 @@
using System;
using System.Net;
namespace Network.NetworkApplication
{
public interface INetworkMessageLane
{
event Action<byte[], IPEndPoint> Received;
void Send(byte[] data);
void SendTo(byte[] data, IPEndPoint target);
void SendToAll(byte[] data);
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 910fad1550952a748bcbded789d8ae87
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -5,14 +5,20 @@ namespace Network.NetworkApplication
{
public sealed class ManagedNetworkSession
{
public ManagedNetworkSession(IPEndPoint remoteEndPoint, SessionManager sessionManager)
public ManagedNetworkSession(
IPEndPoint remoteEndPoint,
SessionManager sessionManager,
ClockSyncState clockSync)
{
RemoteEndPoint = remoteEndPoint ?? throw new ArgumentNullException(nameof(remoteEndPoint));
SessionManager = sessionManager ?? throw new ArgumentNullException(nameof(sessionManager));
ClockSync = clockSync ?? throw new ArgumentNullException(nameof(clockSync));
}
public IPEndPoint RemoteEndPoint { get; }
public SessionManager SessionManager { get; }
public ClockSyncState ClockSync { get; }
}
}
@@ -10,17 +10,58 @@ namespace Network.NetworkApplication
{
public class MessageManager
{
private readonly ITransport transport;
private readonly INetworkMessageLane reliableLane;
private readonly INetworkMessageLane syncLane;
private readonly INetworkMessageDispatcher dispatcher;
private readonly IMessageDeliveryPolicyResolver deliveryPolicyResolver;
private readonly SyncSequenceTracker syncSequenceTracker;
private readonly Dictionary<MessageType, Func<byte[], IPEndPoint, Task>> handlers =
new();
public MessageManager(ITransport transport, INetworkMessageDispatcher dispatcher)
: this(
CreateLane(transport),
dispatcher,
new DefaultMessageDeliveryPolicyResolver(),
null,
new SyncSequenceTracker())
{
this.transport = transport ?? throw new ArgumentNullException(nameof(transport));
}
public MessageManager(
ITransport reliableTransport,
INetworkMessageDispatcher dispatcher,
IMessageDeliveryPolicyResolver deliveryPolicyResolver,
ITransport syncTransport = null,
SyncSequenceTracker syncSequenceTracker = null)
: this(
CreateLane(reliableTransport),
dispatcher,
deliveryPolicyResolver,
CreateLaneIfDistinct(reliableTransport, syncTransport),
syncSequenceTracker)
{
}
public MessageManager(
INetworkMessageLane reliableLane,
INetworkMessageDispatcher dispatcher,
IMessageDeliveryPolicyResolver deliveryPolicyResolver = null,
INetworkMessageLane syncLane = null,
SyncSequenceTracker syncSequenceTracker = null)
{
this.reliableLane = reliableLane ?? throw new ArgumentNullException(nameof(reliableLane));
this.dispatcher = dispatcher ?? throw new ArgumentNullException(nameof(dispatcher));
this.transport.OnReceive += OnTransportReceive;
this.deliveryPolicyResolver = deliveryPolicyResolver ?? new DefaultMessageDeliveryPolicyResolver();
this.syncLane = syncLane;
this.syncSequenceTracker = syncSequenceTracker ?? new SyncSequenceTracker();
this.reliableLane.Received += OnTransportReceive;
if (this.syncLane != null && !ReferenceEquals(this.syncLane, this.reliableLane))
{
this.syncLane.Received += OnTransportReceive;
}
}
public INetworkMessageDispatcher Dispatcher => dispatcher;
@@ -71,14 +112,15 @@ namespace Network.NetworkApplication
Type = (int)type,
Payload = message.ToByteString()
};
var lane = ResolveLane(type);
if (target != null)
{
transport.SendTo(envelope.ToByteArray(), target);
lane.SendTo(envelope.ToByteArray(), target);
}
else
{
transport.Send(envelope.ToByteArray());
lane.Send(envelope.ToByteArray());
}
Console.WriteLine($"[MessageManager] 发送消息:{type} -> {target?.ToString() ?? "default"}");
@@ -97,7 +139,7 @@ namespace Network.NetworkApplication
Type = (int)type,
Payload = message.ToByteString()
};
transport.SendToAll(envelope.ToByteArray());
ResolveLane(type).SendToAll(envelope.ToByteArray());
}
public Task<int> DrainPendingMessagesAsync(int maxMessages = int.MaxValue)
@@ -112,10 +154,16 @@ namespace Network.NetworkApplication
var envelope = Envelope.Parser.ParseFrom(data);
var type = (MessageType)envelope.Type;
Console.WriteLine($"[MessageManager] 收到消息:{type} 来自 {sender}");
var payload = envelope.Payload.ToByteArray();
if (!syncSequenceTracker.ShouldAccept(type, payload, sender))
{
Console.WriteLine($"[MessageManager] 丢弃过期同步消息:{type} 来自 {sender}");
return;
}
if (handlers.TryGetValue(type, out var handler))
{
var payload = envelope.Payload.ToByteArray();
dispatcher.Enqueue(() => DispatchAsync(handler, payload, sender, type));
}
else
@@ -144,5 +192,30 @@ namespace Network.NetworkApplication
Console.WriteLine($"[MessageManager] Handler 执行错误:{type} -> {ex.Message}");
}
}
private INetworkMessageLane ResolveLane(MessageType type)
{
var policy = deliveryPolicyResolver.Resolve(type);
return policy == DeliveryPolicy.HighFrequencySync && syncLane != null
? syncLane
: reliableLane;
}
private static INetworkMessageLane CreateLane(ITransport transport)
{
return new TransportMessageLane(transport ?? throw new ArgumentNullException(nameof(transport)));
}
private static INetworkMessageLane CreateLaneIfDistinct(ITransport reliableTransport, ITransport syncTransport)
{
if (syncTransport == null)
{
return null;
}
return ReferenceEquals(reliableTransport, syncTransport)
? null
: CreateLane(syncTransport);
}
}
}
@@ -119,7 +119,9 @@ namespace Network.NetworkApplication
public void NotifyHeartbeatReceived(IPEndPoint remoteEndPoint, long? serverTick = null)
{
GetOrCreateSession(remoteEndPoint).SessionManager.NotifyHeartbeatReceived(serverTick);
var session = GetOrCreateSession(remoteEndPoint);
session.SessionManager.NotifyHeartbeatReceived();
session.ClockSync.ObserveSample(serverTick);
}
public void NotifyInboundActivity(IPEndPoint remoteEndPoint)
@@ -127,6 +129,11 @@ namespace Network.NetworkApplication
GetOrCreateSession(remoteEndPoint).SessionManager.NotifyInboundActivity();
}
public void ObserveAuthoritativeState(IPEndPoint remoteEndPoint, long? serverTick)
{
GetOrCreateSession(remoteEndPoint).ClockSync.ObserveSample(serverTick);
}
public bool RemoveSession(IPEndPoint remoteEndPoint, string reason = null)
{
SessionRegistration registration;
@@ -194,7 +201,8 @@ namespace Network.NetworkApplication
}
var sessionManager = new SessionManager(reconnectPolicy, utcNowProvider);
var session = new ManagedNetworkSession(normalizedEndPoint, sessionManager);
var clockSync = new ClockSyncState(utcNowProvider);
var session = new ManagedNetworkSession(normalizedEndPoint, sessionManager, clockSync);
Action<SessionLifecycleEvent> handler = lifecycleEvent =>
LifecycleChanged?.Invoke(new MultiSessionLifecycleEvent(session.RemoteEndPoint, session.SessionManager, lifecycleEvent));
@@ -235,4 +243,3 @@ namespace Network.NetworkApplication
}
}
}
@@ -32,8 +32,6 @@ namespace Network.NetworkApplication
public TimeSpan? LastRoundTripTime { get; private set; }
public long? LastServerTick { get; private set; }
public string LastFailureReason { get; private set; }
public bool CanSendHeartbeat => State == ConnectionState.LoggedIn;
@@ -104,11 +102,10 @@ namespace Network.NetworkApplication
RaiseEvent(SessionEventKind.HeartbeatSent, State, State, lastHeartbeatSentUtc.Value);
}
public void NotifyHeartbeatReceived(long? serverTick = null)
public void NotifyHeartbeatReceived()
{
var now = utcNowProvider();
lastLivenessUtc = now;
LastServerTick = serverTick;
if (lastHeartbeatSentUtc.HasValue)
{
LastRoundTripTime = now - lastHeartbeatSentUtc.Value;
@@ -10,19 +10,34 @@ namespace Network.NetworkApplication
ITransport transport,
INetworkMessageDispatcher dispatcher,
SessionReconnectPolicy reconnectPolicy = null,
Func<DateTimeOffset> utcNowProvider = null)
Func<DateTimeOffset> utcNowProvider = null,
ITransport syncTransport = null,
IMessageDeliveryPolicyResolver deliveryPolicyResolver = null,
SyncSequenceTracker syncSequenceTracker = null,
ClockSyncState clockSync = null)
{
Transport = transport ?? throw new ArgumentNullException(nameof(transport));
SyncTransport = syncTransport;
SessionManager = new SessionManager(reconnectPolicy, utcNowProvider);
MessageManager = new MessageManager(transport, dispatcher ?? throw new ArgumentNullException(nameof(dispatcher)));
ClockSync = clockSync ?? new ClockSyncState(utcNowProvider);
MessageManager = new MessageManager(
transport,
dispatcher ?? throw new ArgumentNullException(nameof(dispatcher)),
deliveryPolicyResolver ?? new DefaultMessageDeliveryPolicyResolver(),
syncTransport,
syncSequenceTracker ?? new SyncSequenceTracker());
}
public ITransport Transport { get; }
public ITransport SyncTransport { get; }
public MessageManager MessageManager { get; }
public SessionManager SessionManager { get; }
public ClockSyncState ClockSync { get; }
public event Action<SessionLifecycleEvent> LifecycleChanged
{
add => SessionManager.LifecycleChanged += value;
@@ -32,13 +47,26 @@ namespace Network.NetworkApplication
public async Task StartAsync()
{
await Transport.StartAsync();
if (SyncTransport != null && !ReferenceEquals(SyncTransport, Transport))
{
await SyncTransport.StartAsync();
}
SessionManager.NotifyTransportConnected();
PublishMetricsSessionSnapshot();
}
public void Stop()
{
Transport.Stop();
if (SyncTransport != null && !ReferenceEquals(SyncTransport, Transport))
{
SyncTransport.Stop();
}
SessionManager.NotifyTransportDisconnected("Transport stopped");
PublishMetricsSessionSnapshot();
}
public Task<int> DrainPendingMessagesAsync(int maxMessages = int.MaxValue)
@@ -49,36 +77,90 @@ namespace Network.NetworkApplication
public void NotifyLoginStarted()
{
SessionManager.NotifyLoginStarted();
PublishMetricsSessionSnapshot();
}
public void NotifyLoginSucceeded()
{
SessionManager.NotifyLoginSucceeded();
PublishMetricsSessionSnapshot();
}
public void NotifyLoginFailed(string reason = null)
{
SessionManager.NotifyLoginFailed(reason);
PublishMetricsSessionSnapshot();
}
public void NotifyHeartbeatSent()
{
SessionManager.NotifyHeartbeatSent();
PublishMetricsSessionSnapshot();
}
public void NotifyHeartbeatReceived(long? serverTick = null)
{
SessionManager.NotifyHeartbeatReceived(serverTick);
SessionManager.NotifyHeartbeatReceived();
ClockSync.ObserveSample(serverTick);
PublishMetricsSessionSnapshot();
}
public void NotifyInboundActivity()
{
SessionManager.NotifyInboundActivity();
PublishMetricsSessionSnapshot();
}
public void UpdateLifecycle()
{
SessionManager.Evaluate();
PublishMetricsSessionSnapshot();
}
public void ObserveAuthoritativeState(long? serverTick)
{
ClockSync.ObserveSample(serverTick);
PublishMetricsSessionSnapshot();
}
private void PublishMetricsSessionSnapshot()
{
RecordMetricsSessionSnapshot(Transport, "shared-runtime", SessionManager, ClockSync, remoteEndPoint: null);
if (SyncTransport != null && !ReferenceEquals(SyncTransport, Transport))
{
RecordMetricsSessionSnapshot(SyncTransport, "shared-runtime-sync", SessionManager, ClockSync, remoteEndPoint: null);
}
}
private static void RecordMetricsSessionSnapshot(
ITransport transport,
string scope,
SessionManager sessionManager,
ClockSyncState clockSync,
System.Net.IPEndPoint remoteEndPoint)
{
if (transport is not ITransportMetricsSink metricsSink || sessionManager == null)
{
return;
}
metricsSink.RecordApplicationSessionSnapshot(new TransportApplicationSessionSnapshot
{
Scope = scope,
RemoteEndPoint = remoteEndPoint?.ToString(),
ConnectionState = sessionManager.State.ToString(),
CanSendHeartbeat = sessionManager.CanSendHeartbeat,
LastRoundTripTimeMs = sessionManager.LastRoundTripTime.HasValue
? (long?)System.Math.Max(0d, sessionManager.LastRoundTripTime.Value.TotalMilliseconds)
: null,
LastFailureReason = sessionManager.LastFailureReason,
LastLivenessUtc = sessionManager.LastLivenessUtc,
LastHeartbeatSentUtc = sessionManager.LastHeartbeatSentUtc,
NextReconnectAtUtc = sessionManager.NextReconnectAtUtc,
CurrentServerTick = clockSync?.CurrentServerTick,
ObservedAtUtc = DateTimeOffset.UtcNow
});
}
}
}
@@ -0,0 +1,70 @@
using System;
using System.Collections.Generic;
using System.Net;
using Network.Defines;
namespace Network.NetworkApplication
{
public sealed class SyncSequenceTracker
{
private readonly object gate = new();
private readonly Dictionary<string, long> latestSequenceByStream = new();
public bool ShouldAccept(MessageType messageType, byte[] payload, IPEndPoint sender)
{
if (!TryResolveSequence(messageType, payload, sender, out var streamKey, out var sequence))
{
return true;
}
lock (gate)
{
if (latestSequenceByStream.TryGetValue(streamKey, out var latestSequence) &&
sequence < latestSequence)
{
return false;
}
latestSequenceByStream[streamKey] = sequence;
return true;
}
}
private static bool TryResolveSequence(
MessageType messageType,
byte[] payload,
IPEndPoint sender,
out string streamKey,
out long sequence)
{
switch (messageType)
{
case MessageType.PlayerInput:
{
var input = PlayerInput.Parser.ParseFrom(payload);
streamKey = $"input:{Normalize(sender)}:{input.PlayerId}";
sequence = input.Tick;
return true;
}
case MessageType.PlayerState:
{
var state = PlayerState.Parser.ParseFrom(payload);
streamKey = $"state:{state.PlayerId}";
sequence = state.Tick;
return true;
}
default:
streamKey = null;
sequence = 0;
return false;
}
}
private static string Normalize(IPEndPoint sender)
{
return sender == null ? "unknown" : sender.ToString();
}
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 377bd3719da11b346926ee82dd56bc45
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant:
@@ -0,0 +1,39 @@
using System;
using System.Net;
using Network.NetworkTransport;
namespace Network.NetworkApplication
{
public sealed class TransportMessageLane : INetworkMessageLane
{
private readonly ITransport transport;
public TransportMessageLane(ITransport transport)
{
this.transport = transport ?? throw new ArgumentNullException(nameof(transport));
this.transport.OnReceive += HandleReceive;
}
public event Action<byte[], IPEndPoint> Received;
public void Send(byte[] data)
{
transport.Send(data);
}
public void SendTo(byte[] data, IPEndPoint target)
{
transport.SendTo(data, target);
}
public void SendToAll(byte[] data)
{
transport.SendToAll(data);
}
private void HandleReceive(byte[] data, IPEndPoint sender)
{
Received?.Invoke(data, sender);
}
}
}
@@ -0,0 +1,11 @@
fileFormatVersion: 2
guid: 4c7a1935b45cd2e418814d4b6a5c0b45
MonoImporter:
externalObjects: {}
serializedVersion: 2
defaultReferences: []
executionOrder: 0
icon: {instanceID: 0}
userData:
assetBundleName:
assetBundleVariant: