forked from Rollocraft/CS2MultiplayerMod
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathMultiplayerSession.cs
More file actions
333 lines (230 loc) · 13.6 KB
/
Copy pathMultiplayerSession.cs
File metadata and controls
333 lines (230 loc) · 13.6 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
using System.Collections.Generic;
using System.Security.Cryptography.X509Certificates;
using CS2MPMod.Core.Diagnostics;
using CS2MPMod.Core.Networking;
using CS2MPMod.Core.Protocol;
using CS2MPMod.Core.Protocol.Messages;
namespace CS2MPMod.Core.Session
{
/// <summary>
/// The multiplayer core session manager. Owns transport, handshake, peer list,
/// keep-alives, and message routing. Host-authoritative: clients talk only to host,
/// which relays commands in canonical order. Challenge-response auth, rate budgeting,
/// and protocol violations trigger disconnects. All public methods run on game thread.
/// See <see cref="Update"/>, <see cref="HandshakeAuth"/>, <see cref="ITransport.Poll"/>.
/// </summary>
public sealed partial class MultiplayerSession
{
// A world transfer, autosave, or a long simulation frame can pause the game
// thread for several seconds. Ten seconds caused false host disconnects
// even though the underlying socket was still alive. Keep heartbeats
// frequent, but allow a generous grace window before declaring a peer dead.
private const int HeartbeatIntervalMs = 1000;
private const int PeerTimeoutMs = 30000;
private const int HandshakeTimeoutMs = 10000;
/// <summary>A join awaiting the host's manual approval is auto-declined after this
/// long, so an absent host never leaves the would-be player waiting forever and a
/// pre-handshake socket is never held open indefinitely.</summary>
private const int JoinApprovalTimeoutMs = 120000;
private const int HostPlayerId = 1;
/// <summary>Reassembling blobs allowed at once on a client.</summary>
private const int MaxActiveBlobs = 4;
/// <summary>A blob that receives no chunk for this long is abandoned.</summary>
private const int BlobStallTimeoutMs = 60000;
/// <summary>Minimum gap between accepted /sync requests - save+stream is expensive, kept short so post-join syncs aren't silently ignored.</summary>
private const long ResyncRequestCooldownMs = 5000;
private readonly IModLogger _log;
private readonly MessageCodec _codec;
private readonly List<ISessionObserver> _observers = new List<ISessionObserver>();
private readonly List<TransportEvent> _eventBuffer = new List<TransportEvent>();
private readonly Dictionary<int, Peer> _peers = new Dictionary<int, Peer>();
private readonly Dictionary<string, BlobReassembler> _blobs = new Dictionary<string, BlobReassembler>();
private readonly Dictionary<string, long> _blobTransferIds = new Dictionary<string, long>();
private readonly Dictionary<string, int> _allowedBlobChannels = new Dictionary<string, int>();
private readonly HashSet<ushort> _allowedCommandIds = new HashSet<ushort>();
private readonly HashSet<int> _administrativeRemovals = new HashSet<int>();
private readonly HashSet<string> _hostBannedAddresses = new HashSet<string>();
// Connections already told to go. The transport only removes a peer when its
// Disconnected event arrives, so without this every frame already queued behind a
// flood is dispatched - and logged - against a connection that is on its way out.
private readonly HashSet<int> _puntedConnections = new HashSet<int>();
private readonly FailedAuthTracker _failedAuth = new FailedAuthTracker();
private ITransport _transport;
private MultiplayerConfig _config;
private X509Certificate2 _certificate;
private PortForward _portForward;
private int _nextPlayerId = HostPlayerId + 1;
private long _lastHeartbeatMs;
private long _lastBlobSweepMs;
private long _lastAuthSweepMs;
private long _lastResyncAcceptedUnixMs;
private bool _challengeAnswered;
private bool _awaitingHostApproval;
private bool _worldSyncSuspended;
private long _worldSyncEpoch;
public MultiplayerSession(IModLogger log, MessageCodec codec = null)
{
_log = log ?? NullModLogger.Instance;
_codec = codec ?? MessageCodec.CreateDefault();
}
public SessionRole Role { get; private set; } = SessionRole.None;
public SessionStatus Status { get; private set; } = SessionStatus.Offline;
public int LocalPlayerId { get; private set; }
public string LocalPlayerName { get; private set; } = "Player";
/// <summary>True when the transport is actually running TLS.</summary>
public bool EncryptionActive { get; private set; }
/// <summary>True when this session requires a password.</summary>
public bool PasswordProtected => _config != null && _config.Password.Length > 0;
/// <summary>True when hosting beyond the local network (LAN filter off).</summary>
public bool PublicExposure => Role == SessionRole.Host && _config != null && !_config.LanOnly;
/// <summary>
/// Whether this session replicates the simulation's own decisions. The host answers from
/// its own config; a client answers with what the host announced when it was accepted,
/// never with its local setting - the two machines have to hold the same half of the
/// simulation or one of them waits forever for the other's messages.
/// </summary>
public bool SimulationSyncEnabled => Role == SessionRole.Client
? _hostSimulationSync
: _config == null || _config.SimulationSync;
/// <summary>Client-only: the host's answer, defaulted until the accept arrives.</summary>
private bool _hostSimulationSync = true;
/// <summary>How the active session reaches its peers (Direct before the first session).</summary>
public TransportMode Transport => _config != null ? _config.Transport : TransportMode.Direct;
/// <summary>True when the active session runs over a relay rather than a direct socket.</summary>
public bool UsesRelay => _config != null && _config.Transport == TransportMode.SteamRelay;
/// <summary>TCP port of the active session's config (0 before the first session).</summary>
public int Port => _config != null ? _config.Port : 0;
/// <summary>
/// What the router made of opening this host's port. Null whenever nothing was
/// asked: a client, a relay session, or a LAN-only host, none of which need one.
/// </summary>
public PortForwardState? PortForwardStatus =>
_portForward != null ? _portForward.State : (PortForwardState?)null;
/// <summary>The public address the router reported, or null if it never told us.</summary>
public string PortForwardAddress => _portForward != null ? _portForward.ExternalAddress : null;
/// <summary>Bytes queued in the transport but not yet on the wire (0 when idle).</summary>
public long PendingSendBytes => _transport != null ? _transport.PendingSendBytes : 0;
/// <summary>Epoch of the installed command baseline, for local diagnostic correlation.</summary>
public long CommandEpoch => _commandEpoch;
/// <summary>Channel of the blob currently being received, or null. For progress UX.</summary>
public string IncomingBlobChannel { get; private set; }
public int IncomingBlobReceived { get; private set; }
public int IncomingBlobTotal { get; private set; }
public long IncomingBlobTransferId { get; private set; }
/// <summary>
/// True between a world-sync Begin and its matching Resume/Abort. Gameplay traffic is
/// rejected at the session boundary during this interval, in addition to game-layer gates.
/// </summary>
public bool WorldSyncSuspended => _worldSyncSuspended;
// Host-side "Sending world %": a streamed blob is queued instantly (the send is
// non-blocking) and then drains off the transport's send thread; these track that
// drain so the host can show a progress bar instead of appearing frozen.
private bool _outgoingBlobActive;
private long _outgoingBlobTotal;
private long _outgoingBlobSent;
public bool OutgoingBlobActive => _outgoingBlobActive;
public long OutgoingBlobTotal => _outgoingBlobTotal;
public long OutgoingBlobSent => _outgoingBlobSent;
public IReadOnlyCollection<Peer> Peers => _peers.Values;
/// <summary>
/// Client-only: the host acknowledged the join and it is waiting for the host to
/// approve it by hand. True between the host's HandshakePending and its accept/reject.
/// </summary>
public bool AwaitingHostApproval => _awaitingHostApproval;
/// <summary>
/// Host-only: joins that passed every automatic check and are waiting for the host
/// to approve or decline them. Enumerated on the game thread alongside the pump.
/// </summary>
public IEnumerable<Peer> PendingJoins
{
get
{
foreach (var pair in _peers)
if (pair.Value.AwaitingApproval) yield return pair.Value;
}
}
public void AddObserver(ISessionObserver observer)
{
if (observer != null && !_observers.Contains(observer)) _observers.Add(observer);
}
public void RemoveObserver(ISessionObserver observer) => _observers.Remove(observer);
// ---- Authorization registries (filled by the game layer at startup) -----
/// <summary>
/// Declare a blob channel clients may receive with size ceiling. Unregistered blobs are dropped - secure by default.
/// </summary>
public void AllowBlobChannel(string channel, int maxBytes)
{
if (!string.IsNullOrEmpty(channel) && maxBytes > 0)
_allowedBlobChannels[channel] = maxBytes;
}
/// <summary>
/// Declare the simulation command ids peers are allowed to send. Once any id is
/// registered, a command outside the set disconnects its sender.
/// </summary>
public void AllowCommands(params ushort[] commandIds)
{
if (commandIds == null) return;
for (int i = 0; i < commandIds.Length; i++) _allowedCommandIds.Add(commandIds[i]);
}
// ---- Lifecycle --------------------------------------------------------
// ---- Per-tick pump ----------------------------------------------------
/// <summary>
/// Advance the session: drain transport events, dispatch messages, send
/// keep-alives, and reap timed-out peers. <paramref name="nowUnixMs"/> is the
/// caller's monotonic clock so the core stays free of <c>DateTime.Now</c>.
/// </summary>
public void Update(long nowUnixMs)
{
if (_transport == null) return;
_commandNowMs = nowUnixMs;
_eventBuffer.Clear();
_transport.Poll(_eventBuffer);
for (int i = 0; i < _eventBuffer.Count; i++)
HandleEvent(_eventBuffer[i], nowUnixMs);
if (Status == SessionStatus.Connected)
{
PumpHeartbeats(nowUnixMs);
ReapTimedOutPeers(nowUnixMs);
PumpCommandReplay(nowUnixMs);
SweepStalledBlobs(nowUnixMs);
PumpOutgoingBlobs();
UpdateOutgoingBlobProgress();
// The ban book only grows on failed auths, so a sparse sweep is plenty.
if (nowUnixMs - _lastAuthSweepMs >= 60000)
{
_lastAuthSweepMs = nowUnixMs;
_failedAuth.Prune(nowUnixMs);
}
}
}
/// <summary>Track how much of a streamed world has drained off the send thread.</summary>
private void UpdateOutgoingBlobProgress()
{
if (!_outgoingBlobActive || _transport == null) return;
long pending = _transport.PendingSendBytes;
long remaining = 0;
foreach (OutgoingBlob blob in _outgoingBlobs) remaining += blob.Data.Length - blob.Offset;
long sent = _outgoingBlobTotal - remaining - pending;
_outgoingBlobSent = sent < 0 ? 0 : (sent > _outgoingBlobTotal ? _outgoingBlobTotal : sent);
// Drained to a trickle (only small keep-alives/commands left): the world is sent.
// Gameplay traffic keeps flowing, so waiting for an exactly empty queue would leave
// the transfer reported as active for the rest of the session.
if (_outgoingBlobs.Count == 0 && pending < 65536)
{
_outgoingBlobSent = _outgoingBlobTotal;
_outgoingBlobActive = false;
}
}
// ---- Handshake --------------------------------------------------------
// ---- Keep-alive -------------------------------------------------------
// ---- Chat & commands --------------------------------------------------
// ---- On-demand world resync (/sync) ------------------------------------
// ---- Replicated state -------------------------------------------------
// ---- Player positions -------------------------------------------------
// ---- Large blobs (e.g. savegame for map sync) -------------------------
// ---- Send helpers -----------------------------------------------------
// ---- Observer fan-out -------------------------------------------------
// Every callback is isolated: one observer throwing must never kill the pump,
// stop the remaining observers, or take the session down.
}
}