forked from Rollocraft/CS2MultiplayerMod
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathMessaging.cs
More file actions
363 lines (318 loc) · 17 KB
/
Copy pathMessaging.cs
File metadata and controls
363 lines (318 loc) · 17 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
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
using System;
using System.Collections.Generic;
using CS2MPMod.Core.Diagnostics;
using CS2MPMod.Core.Networking;
using CS2MPMod.Core.Protocol;
using CS2MPMod.Core.Protocol.Messages;
namespace CS2MPMod.Core.Session
{
public sealed partial class MultiplayerSession
{
private void HandleHeartbeat(ConnectionId connection, Peer peer, Heartbeat heartbeat, long nowUnixMs)
{
// An echo returns a timestamp WE sent, so now − echo is a true round-trip on
// our own clock — the peer's clock never enters the math (the two machines'
// clocks are unrelated: each side passes its own monotonic ms). Echoes are
// not echoed back, so a ping costs exactly one reply.
if (heartbeat.EchoOfMs > 0)
{
long rtt = nowUnixMs - heartbeat.EchoOfMs;
if (peer != null && rtt >= 0 && rtt < 60000) peer.LatencyMs = (int)rtt;
return;
}
// A ping: return the sender's timestamp so IT can measure its round-trip.
if (heartbeat.SentAtMs > 0)
SendTo(connection, new Heartbeat(nowUnixMs, heartbeat.SentAtMs));
}
private void PumpHeartbeats(long nowUnixMs)
{
if (nowUnixMs - _lastHeartbeatMs < HeartbeatIntervalMs) return;
_lastHeartbeatMs = nowUnixMs;
var beat = new Heartbeat(nowUnixMs);
if (Role == SessionRole.Host)
BroadcastToAll(beat, ConnectionId.None);
else
SendTo(ConnectionId.Server, beat);
}
private void ReapTimedOutPeers(long nowUnixMs)
{
List<Peer> dead = null; // dropped silently at the socket
List<Peer> unanswered = null; // approvals the host never answered
foreach (var pair in _peers)
{
Peer peer = pair.Value;
if (peer.Handshaked)
{
if (nowUnixMs - peer.LastSeenUnixMs > PeerTimeoutMs)
(dead ?? (dead = new List<Peer>())).Add(peer);
}
else if (peer.AwaitingApproval)
{
// The host never accepted or declined in time — auto-decline so the
// socket is freed and the waiting player is told why, instead of both
// sides hanging on an absent host.
if (nowUnixMs - peer.ConnectedAtUnixMs > JoinApprovalTimeoutMs)
(unanswered ?? (unanswered = new List<Peer>())).Add(peer);
}
// Pending sockets must finish the handshake promptly or make room.
else if (nowUnixMs - peer.ConnectedAtUnixMs > HandshakeTimeoutMs)
(dead ?? (dead = new List<Peer>())).Add(peer);
}
if (dead != null)
foreach (Peer peer in dead)
{
_log.Warn(LogTopic.Session,
(peer.Handshaked ? "Peer timed out: " : "Handshake timed out: ") + peer);
_transport.Disconnect(peer.Connection);
// The transport will also raise Disconnected; removal/notify happens there.
}
if (unanswered != null)
foreach (Peer peer in unanswered)
// Reject logs, delivers the reason, flushes, and removes the peer.
Reject(peer.Connection, "The host did not respond to your join request in time.");
}
/// <summary>
/// Send a chat line. On the host it is relayed to all clients. "/sync" is a
/// command, not a line: it asks the host for a fresh world stream instead.
/// </summary>
public void SendChat(string text)
{
if (Status != SessionStatus.Connected || string.IsNullOrEmpty(text)) return;
if (IsSyncCommand(text)) { RequestWorldSync(); return; }
text = WireGuard.SanitizeText(text, WireGuard.MaxChatLength);
if (text.Length == 0) return;
var message = new ChatMessage(LocalPlayerName, text);
if (Role == SessionRole.Host)
BroadcastToAll(message, ConnectionId.None);
else
SendTo(ConnectionId.Server, message);
}
private void HandleChat(ConnectionId from, Peer peer, ChatMessage chat, long nowUnixMs)
{
// A "/sync" line from a client (e.g. an older build that sends it as raw
// chat) is treated as the command it means.
if (Role == SessionRole.Host && IsSyncCommand(chat.Text))
{
HandleResyncRequest(from, peer, nowUnixMs);
return;
}
// Whatever arrives is displayed and logged — so control characters, fake
// newlines and kilometer-long lines are stripped before anything sees them.
chat.Text = WireGuard.SanitizeText(chat.Text, WireGuard.MaxChatLength);
// Never trust the sender's claimed name — display the one we authenticated.
if (Role == SessionRole.Host && peer != null)
chat.SenderName = peer.Name;
else
chat.SenderName = chat.SenderName == null
? null // system notice
: WireGuard.SanitizePlayerName(chat.SenderName);
NotifyChat(chat.SenderName, chat.Text);
// Host fans a client's message out to the other clients.
if (Role == SessionRole.Host)
BroadcastToAll(chat, from);
}
private static bool IsSyncCommand(string text) =>
text != null && text.Trim().Equals("/sync", StringComparison.OrdinalIgnoreCase);
/// <summary>
/// Player-initiated drift correction. A client asks the host to stream it the
/// current world; on the host it refreshes every client. The actual save+stream
/// is done by the observer (the game layer owns savegames).
/// </summary>
public void RequestWorldSync() => RequestWorldSync(ManualSyncReason, false);
/// <summary>Recovery initiated by the mod, independently of a player's sync button.</summary>
public void RequestAutomaticWorldSync(string reason) => RequestWorldSync(reason, true);
private void RequestWorldSync(string reason, bool automatic)
{
if (Status != SessionStatus.Connected) return;
reason = WireGuard.SanitizeText(reason, WireGuard.MaxResyncReasonLength);
if (reason.Length == 0) reason = automatic ? "automatic recovery" : ManualSyncReason;
if (_worldSyncSuspended)
{
_log.Detail(LogTopic.Session, "World sync request coalesced into active epoch " +
_worldSyncEpoch + " (" + reason + ").");
return;
}
if (Role == SessionRole.Client)
{
SendTo(ConnectionId.Server, new ResyncRequestMessage(LocalPlayerId, reason, automatic));
_log.Event(LogTopic.Session, "World sync request sent to host (" + reason + ").");
NotifyChat(null, automatic
? "The mod requested an automatic world sync - the host will stream you its city."
: "World sync requested - the host will stream you its city.");
}
else if (Role == SessionRole.Host)
{
_log.Event(LogTopic.Session,
(automatic ? "Automatic world recovery started on host (" :
"Host requested world sync for all clients (") + reason + ").");
string notice = automatic
? "The mod started an automatic world sync - streaming the city to all players."
: "World sync started - streaming the city to all players.";
BroadcastToAll(new ChatMessage(null, notice), ConnectionId.None);
NotifyResyncRequested(LocalPlayerId, ConnectionId.None);
NotifyChat(null, notice);
}
}
/// <summary>What a request carries when the player asked for it themselves.</summary>
internal const string ManualSyncReason = "requested by the player";
private void HandleResyncRequest(ConnectionId from, Peer peer, long nowUnixMs)
{
HandleResyncRequest(from, peer, nowUnixMs, ManualSyncReason, false);
}
private void HandleResyncRequest(ConnectionId from, Peer peer, long nowUnixMs, string reason,
bool automatic)
{
if (Role != SessionRole.Host) return;
reason = WireGuard.SanitizeText(reason, WireGuard.MaxResyncReasonLength);
if (reason.Length == 0) reason = automatic ? "automatic recovery" : ManualSyncReason;
// Rate limit: a misbehaving client spamming /sync would otherwise keep the
// host in a permanent save+stream loop. (Per-peer budgets run on top.)
if (nowUnixMs - _lastResyncAcceptedUnixMs < ResyncRequestCooldownMs)
{
_log.Warn(LogTopic.Session, "Ignoring world sync request from " +
(peer != null ? peer.ToString() : from.ToString()) + " (" + reason +
"): a world sync ran moments ago.");
return;
}
_lastResyncAcceptedUnixMs = nowUnixMs;
// The requester's identity comes from OUR peer table, never from the wire.
string name = peer != null && peer.Name != null ? peer.Name : from.ToString();
// Why the other machine gave up belongs in THIS log too. Without it the host's log
// reads "someone asked for a sync" for both a player pressing the button and a client
// pipeline that could not apply an edit - the two cases that need telling apart most.
_log.Event(LogTopic.Session, (automatic ? "Automatic world recovery requested by the mod on " :
"World sync requested by ") + name + ": " + reason + ".");
// Tell everyone the world is about to snap, then let the game layer stream it.
string notice = automatic
? "The mod triggered an automatic world sync to recover synchronization."
: name + " requested a world sync.";
BroadcastToAll(new ChatMessage(null, notice), ConnectionId.None);
NotifyChat(null, notice);
NotifyResyncRequested(peer != null ? peer.PlayerId : -1, from);
}
/// <summary>
/// Submit a simulation command for synchronization. The host applies it locally
/// and relays to clients; a client forwards it to the host, which then relays.
/// </summary>
public void SendCommand(long tick, ushort commandId, byte[] body)
{
if (Status != SessionStatus.Connected || _worldSyncSuspended || _commandRecoveryPending) return;
var message = new SimulationCommandMessage(LocalPlayerId, tick,
NextCommandSequence(), commandId, body, _commandEpoch);
if (Role == SessionRole.Host)
{
NotifyCommand(message); // apply on the host
JournalCommand(message);
BroadcastToAll(message, ConnectionId.None); // and to clients
}
else
{
SendTo(ConnectionId.Server, message); // host will echo/relay back
}
}
private void HandleCommand(ConnectionId from, Peer peer, SimulationCommandMessage command)
{
// Commands crossing the snapshot cut are deliberately rejected. Every participant
// installs the host snapshot before Resume, so applying only a suffix on one side would
// recreate the very drift this transaction is meant to repair.
if (_worldSyncSuspended) return;
if (command.Epoch != _commandEpoch) return;
// Only command ids the game layer registered are legitimate; anything else
// is a peer probing the surface.
if (_allowedCommandIds.Count > 0 && !_allowedCommandIds.Contains(command.CommandId))
{
Punt(from, peer, "unauthorized command id " + command.CommandId, "SimulationCommand");
return;
}
// The origin id drives every echo-skip; stamp it from OUR peer table so a
// client cannot impersonate another player (or the host) on the wire.
if (Role == SessionRole.Host && peer != null)
command.OriginPlayerId = peer.PlayerId;
// The host is the single ordering authority. Commands from clients arrive on
// separate sockets and therefore cannot safely use arrival order on receivers.
// Stamp before notifying the host or relaying so every remote peer sees the same
// canonical order. A non-zero sequence from an untrusted client is overwritten.
if (Role == SessionRole.Host)
command.Sequence = NextCommandSequence();
if (Role == SessionRole.Client)
{
if (command.Sequence <= 0 || command.Sequence == long.MaxValue)
{
RequestCommandRecovery("host command has an invalid sequence");
return;
}
ReceiveOrderedCommand(command);
return;
}
NotifyCommand(command);
if (Role == SessionRole.Host)
{
JournalCommand(command);
BroadcastToAll(command, ConnectionId.None); // sender must see the canonical sequence too
}
}
/// <summary>
/// Broadcast an authoritative state slice to all clients. Host-only: replicated
/// state flows one way, from the authority outward. A client call is ignored.
/// </summary>
public void SendState(byte channelId, byte[] data)
{
if (Role != SessionRole.Host || Status != SessionStatus.Connected || _worldSyncSuspended) return;
BroadcastToAll(new StateSnapshotMessage(channelId, data), ConnectionId.None);
}
private void HandleState(ConnectionId from, Peer peer, StateSnapshotMessage snapshot)
{
// Only clients apply replicated state; a client pushing "authoritative"
// state at the host is impersonating the authority.
if (Role == SessionRole.Host)
{
Punt(from, peer, "client sent a host-only state snapshot", "StateSnapshot");
return;
}
if (_worldSyncSuspended) return;
NotifyState(snapshot);
}
/// <summary>
/// Client -> host: submit an edit of a player-editable state channel (taxes, policies, ...).
/// Body uses channel's snapshot encoding; host applies it in next broadcast. Host-side edits
/// need no message - host's capture already picks them up.
/// </summary>
public void SendStateEdit(byte channelId, byte[] data)
{
if (Role != SessionRole.Client || Status != SessionStatus.Connected || _worldSyncSuspended) return;
SendTo(ConnectionId.Server, new StateEditMessage(LocalPlayerId, channelId, data));
}
private void HandleStateEdit(Peer peer, StateEditMessage edit)
{
// Only the host arbitrates edits; a client receiving one is a stray.
if (Role != SessionRole.Host) return;
if (_worldSyncSuspended) return;
if (peer != null) edit.OriginPlayerId = peer.PlayerId; // no impersonation
NotifyStateEdit(edit);
}
/// <summary>Publish the local player's camera and display-only hover outlines.</summary>
public void SendPlayerState(float x, float y, float z, float eyeX, float eyeY, float eyeZ, float yaw,
PlayerHoverShape[] hover = null)
{
if (Status != SessionStatus.Connected || _worldSyncSuspended) return;
// Presence is refreshed at 10 Hz. Skip samples under backpressure instead of
// adding stale hover traffic behind city updates on the reliable stream.
if (_transport == null || _transport.PendingSendBytes > 16 * 1024) return;
var message = new PlayerStateMessage(LocalPlayerId, x, y, z, eyeX, eyeY, eyeZ, yaw, hover);
if (Role == SessionRole.Host)
BroadcastToAll(message, ConnectionId.None);
else
SendTo(ConnectionId.Server, message);
}
private void HandlePlayerState(ConnectionId from, Peer peer, PlayerStateMessage state)
{
if (_worldSyncSuspended) return;
// Same anti-impersonation stamp as commands: positions are keyed by player id.
if (Role == SessionRole.Host && peer != null)
state.PlayerId = peer.PlayerId;
NotifyPlayerState(state);
if (Role == SessionRole.Host && _transport != null && _transport.PendingSendBytes <= 16 * 1024)
BroadcastToAll(state, from); // fan a client's position out to the others
}
}
}