Skip to content

Commit db91c6b

Browse files
committed
feat(sidecar): add routed data endpoints
- Add negotiated control, realtime, and bulk-stream profiles with bounded scheduling and backend-independent routing. - Stream large payloads with windowed acknowledgements, backpressure, retries, cancellation, timeouts, and end-to-end integrity checks. - Prefer confirmed routes for typed messages while retaining legacy delivery for older peers and canonicalizing relayed sender identities.
1 parent 7286987 commit db91c6b

33 files changed

Lines changed: 6590 additions & 57 deletions

src/Networking/Sidecar/Core/Handlers/RitsuLibSidecarBuiltInHandlers.cs

Lines changed: 2 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -5,15 +5,6 @@ namespace STS2RitsuLib.Networking.Sidecar
55
{
66
internal static class RitsuLibSidecarBuiltInHandlers
77
{
8-
private const RitsuLibSidecarPeerFeatures SupportedFeatures =
9-
RitsuLibSidecarPeerFeatures.ChunkedStreams |
10-
RitsuLibSidecarPeerFeatures.ManagedNetActions |
11-
RitsuLibSidecarPeerFeatures.BrotliPayloadCompression |
12-
RitsuLibSidecarPeerFeatures.ModelRightClickV2 |
13-
RitsuLibSidecarPeerFeatures.DeveloperActionsV1 |
14-
RitsuLibSidecarInternalPeerFeatures.MonsterIntentActionsV1 |
15-
RitsuLibSidecarInternalPeerFeatures.ExtendedDeveloperStateActionsV1;
16-
178
private static readonly RitsuLibSidecarChunkReassembly Chunks = new();
189

1910
internal static void Register()
@@ -145,7 +136,7 @@ private static void OnHandshake(RitsuLibSidecarDispatchContext ctx)
145136
buf.AsSpan(),
146137
selected,
147138
ok,
148-
SupportedFeatures);
139+
RitsuLibSidecarSupportedFeatures.All);
149140
var rm = RunManager.Instance;
150141
if (ctx.IsHostIngest)
151142
RitsuLibSidecarHighLevelSend.TrySendAsHostToPeer(
@@ -162,7 +153,7 @@ private static void OnHandshake(RitsuLibSidecarDispatchContext ctx)
162153
RitsuLibSidecarDeliverySemantics.StableSync);
163154

164155
RitsuLibFramework.Logger.Debug(
165-
$"[Sidecar] Handshake ack sent target={ctx.SenderNetId}, opcode={RitsuLibSidecarControlOpcodes.HandshakeAck}, payloadLen={buf.Length}, selectedWire={selected}, ok={ok}, senderFeatures={SupportedFeatures}");
156+
$"[Sidecar] Handshake ack sent target={ctx.SenderNetId}, opcode={RitsuLibSidecarControlOpcodes.HandshakeAck}, payloadLen={buf.Length}, selectedWire={selected}, ok={ok}, senderFeatures={RitsuLibSidecarSupportedFeatures.All}");
166157
}
167158

168159
private static void OnHandshakeAck(RitsuLibSidecarDispatchContext ctx)

src/Networking/Sidecar/Core/Messaging/RitsuLibSidecarTypedMessages.cs

Lines changed: 201 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
using System.Text.Json;
2+
using MegaCrit.Sts2.Core.Multiplayer;
23
using MegaCrit.Sts2.Core.Multiplayer.Game;
34
using MegaCrit.Sts2.Core.Multiplayer.Transport;
45
using MegaCrit.Sts2.Core.Runs;
@@ -72,10 +73,24 @@ public readonly record struct SidecarTypedMessageReceivedEvent(
7273
/// </para>
7374
/// <para xml:lang="zh-CN">用于类型化 Sidecar 描述符注册、冲突检查、订阅及便捷发送的注册表。</para>
7475
/// </summary>
76+
/// <remarks>
77+
/// <para xml:lang="en">
78+
/// Registered descriptors receive a bounded routed-endpoint compatibility bridge when local capacity permits.
79+
/// Convenience sends prefer the routed path for each compatible, route-confirmed recipient and use the legacy
80+
/// opcode path only for recipients or payloads that cannot use it. Split broadcasts never send both paths to
81+
/// the same recipient.
82+
/// </para>
83+
/// <para xml:lang="zh-CN">
84+
/// 本地容量允许时,已注册描述符会获得有界的路由端点兼容桥。便捷发送会针对每个兼容且已确认路由的
85+
/// 接收方优先使用新路径,仅在接收方或载荷不适用时回退旧操作码路径;拆分广播绝不会向同一接收方
86+
/// 同时发送两条路径。
87+
/// </para>
88+
/// </remarks>
7589
public static class RitsuLibSidecarTypedMessageRegistry
7690
{
7791
private static readonly Lock Gate = new();
7892
private static readonly Dictionary<ulong, RegistrationBase> Registrations = [];
93+
private static int _migrationEndpointCount;
7994

8095
/// <summary>
8196
/// <para xml:lang="en">Raised after any typed message is successfully deserialized and dispatched.</para>
@@ -100,6 +115,7 @@ public static ulong Register<T>(RitsuLibSidecarMessageDescriptor<T> descriptor)
100115
ArgumentException.ThrowIfNullOrEmpty(descriptor.MessageKey);
101116
ArgumentNullException.ThrowIfNull(descriptor.Serialize);
102117
ArgumentNullException.ThrowIfNull(descriptor.Deserialize);
118+
RitsuLibSidecarProtocol.EnsureDefaultHandlers();
103119

104120
var opcode = RitsuLibSidecarOpcodes.For(descriptor.ModuleId, descriptor.MessageKey);
105121
lock (Gate)
@@ -121,6 +137,7 @@ public static ulong Register<T>(RitsuLibSidecarMessageDescriptor<T> descriptor)
121137
descriptor.Serialize,
122138
descriptor.Deserialize,
123139
descriptor.Delivery);
140+
TryAttachMigrationEndpoint(opcode, reg);
124141
Registrations[opcode] = reg;
125142
RitsuLibSidecarBus.RegisterHandler(opcode, ctx => HandleDispatch(opcode, reg, in ctx));
126143
}
@@ -170,8 +187,18 @@ public static bool SendToHost<T>(INetGameService? netService, RitsuLibSidecarMes
170187
T message)
171188
{
172189
var opcode = Register(descriptor);
173-
var payload = descriptor.Serialize(message);
174-
return RitsuLibSidecarHighLevelSend.TrySendAsClient(netService, opcode, payload, descriptor.Delivery);
190+
var registration = GetRegistration<T>(opcode);
191+
var payload = SerializePayload(registration, message);
192+
if (CanUseRoutedPath(netService) && registration.EndpointHandle is { } endpoint)
193+
{
194+
var result = endpoint.SendToHost(payload);
195+
if (result.IsAccepted)
196+
return true;
197+
if (!ShouldFallbackToLegacy(result.Status))
198+
return false;
199+
}
200+
201+
return RitsuLibSidecarHighLevelSend.TrySendAsClient(netService, opcode, payload, registration.Delivery);
175202
}
176203

177204
/// <summary>
@@ -181,9 +208,7 @@ public static bool SendToHost<T>(INetGameService? netService, RitsuLibSidecarMes
181208
public static bool SendToHost<T>(RunManager? runManager, RitsuLibSidecarMessageDescriptor<T> descriptor,
182209
T message)
183210
{
184-
var opcode = Register(descriptor);
185-
var payload = descriptor.Serialize(message);
186-
return RitsuLibSidecarHighLevelSend.TrySendAsClient(runManager, opcode, payload, descriptor.Delivery);
211+
return SendToHost(runManager?.NetService, descriptor, message);
187212
}
188213

189214
/// <summary>
@@ -194,9 +219,19 @@ public static bool SendToPeer<T>(INetGameService? netService, ulong peerNetId,
194219
RitsuLibSidecarMessageDescriptor<T> descriptor, T message)
195220
{
196221
var opcode = Register(descriptor);
197-
var payload = descriptor.Serialize(message);
222+
var registration = GetRegistration<T>(opcode);
223+
var payload = SerializePayload(registration, message);
224+
if (CanUseRoutedPath(netService) && registration.EndpointHandle is { } endpoint)
225+
{
226+
var result = endpoint.SendToPeer(peerNetId, payload);
227+
if (result.IsAccepted)
228+
return true;
229+
if (!ShouldFallbackToLegacy(result.Status))
230+
return false;
231+
}
232+
198233
return RitsuLibSidecarHighLevelSend.TrySendAsHostToPeer(netService, peerNetId, opcode, payload,
199-
descriptor.Delivery);
234+
registration.Delivery);
200235
}
201236

202237
/// <summary>
@@ -207,9 +242,40 @@ public static bool Broadcast<T>(INetGameService? netService, RitsuLibSidecarMess
207242
T message)
208243
{
209244
var opcode = Register(descriptor);
210-
var payload = descriptor.Serialize(message);
211-
return RitsuLibSidecarHighLevelSend.TrySendAsHostBroadcast(netService, opcode, payload,
212-
descriptor.Delivery);
245+
var registration = GetRegistration<T>(opcode);
246+
var payload = SerializePayload(registration, message);
247+
if (!CanUseRoutedPath(netService) ||
248+
netService is not NetHostGameService { IsConnected: true } host ||
249+
registration.EndpointHandle is not { } endpoint)
250+
return RitsuLibSidecarHighLevelSend.TrySendAsHostBroadcast(
251+
netService,
252+
opcode,
253+
payload,
254+
registration.Delivery);
255+
256+
var routedParticipants = endpoint.GetParticipantsSnapshot().ToHashSet();
257+
foreach (var peer in host.ConnectedPeers)
258+
{
259+
if (!peer.readyForBroadcasting ||
260+
!RitsuLibSidecarSessionManager.CanSendToPeer(peer.peerId))
261+
continue;
262+
263+
if (routedParticipants.Contains(peer.peerId))
264+
{
265+
var result = endpoint.SendToPeer(peer.peerId, payload);
266+
if (result.IsAccepted || !ShouldFallbackToLegacy(result.Status))
267+
continue;
268+
}
269+
270+
RitsuLibSidecarHighLevelSend.TrySendAsHostToPeer(
271+
netService,
272+
peer.peerId,
273+
opcode,
274+
payload,
275+
registration.Delivery);
276+
}
277+
278+
return true;
213279
}
214280

215281
/// <summary>
@@ -219,24 +285,93 @@ public static bool Broadcast<T>(INetGameService? netService, RitsuLibSidecarMess
219285
public static bool Broadcast<T>(RunManager? runManager, RitsuLibSidecarMessageDescriptor<T> descriptor,
220286
T message)
221287
{
222-
var opcode = Register(descriptor);
223-
var payload = descriptor.Serialize(message);
224-
return RitsuLibSidecarHighLevelSend.TrySendAsHostBroadcast(runManager, opcode, payload,
225-
descriptor.Delivery);
288+
return Broadcast(runManager?.NetService, descriptor, message);
289+
}
290+
291+
private static void TryAttachMigrationEndpoint<T>(ulong opcode, Registration<T> registration)
292+
{
293+
if (_migrationEndpointCount >= RitsuLibSidecarEndpointPolicy.MaxLegacyTypedMigrationEndpoints)
294+
return;
295+
var deliveryProfile = ResolveDeliveryProfile(registration.Delivery);
296+
var maxPayloadBytes = deliveryProfile == RitsuLibSidecarDeliveryProfile.Control
297+
? RitsuLibSidecarEndpointPolicy.MaxControlPayloadBytes
298+
: RitsuLibSidecarEndpointPolicy.MaxRealtimePayloadBytes;
299+
try
300+
{
301+
var endpoint = RitsuLibSidecarEndpoints.Register(
302+
new(
303+
"ritsulib.typed",
304+
$"message/{opcode:x16}",
305+
1,
306+
1,
307+
deliveryProfile,
308+
RitsuLibSidecarEndpointTopology.HostAuthority,
309+
maxPayloadBytes,
310+
RitsuLibSidecarEndpointDispatchMode.ReceiveThread),
311+
message => HandleEndpointDispatch(opcode, registration, message));
312+
registration.AttachEndpoint(endpoint);
313+
_migrationEndpointCount++;
314+
}
315+
catch (InvalidOperationException ex)
316+
{
317+
RitsuLibSidecarRepeatedWarningLog.Warn(
318+
$"typed-endpoint-unavailable:opcode={opcode}:{ex.Message}",
319+
$"[Sidecar] Routed typed-message endpoint unavailable opcode={opcode}; legacy path remains active: {ex.Message}");
320+
}
321+
}
322+
323+
private static void HandleEndpointDispatch<T>(
324+
ulong opcode,
325+
Registration<T> registration,
326+
RitsuLibSidecarEndpointMessage endpointMessage)
327+
{
328+
var deliveryProfile = ResolveDeliveryProfile(registration.Delivery);
329+
if (!RitsuLibSidecarEndpointTransport.TryGetNetworkParameters(
330+
deliveryProfile,
331+
out var transferMode,
332+
out var channel))
333+
return;
334+
HandlePayload(
335+
opcode,
336+
registration,
337+
endpointMessage.Payload.Span,
338+
endpointMessage.OriginalSenderNetId,
339+
transferMode,
340+
channel,
341+
RitsuLibSidecarSessionManager.CurrentNetService is NetHostGameService);
226342
}
227343

228344
private static void HandleDispatch<T>(ulong opcode, Registration<T> registration,
229345
in RitsuLibSidecarDispatchContext context)
346+
{
347+
HandlePayload(
348+
opcode,
349+
registration,
350+
context.Payload.Span,
351+
context.SenderNetId,
352+
context.TransferMode,
353+
context.Channel,
354+
context.IsHostIngest);
355+
}
356+
357+
private static void HandlePayload<T>(
358+
ulong opcode,
359+
Registration<T> registration,
360+
ReadOnlySpan<byte> payload,
361+
ulong senderNetId,
362+
NetTransferMode transferMode,
363+
int channel,
364+
bool isHostIngest)
230365
{
231366
T message;
232367
try
233368
{
234-
message = registration.Deserialize(context.Payload.Span);
369+
message = registration.Deserialize(payload);
235370
}
236371
catch (Exception ex) when (RitsuLibExceptionPolicy.IsRecoverable(ex))
237372
{
238373
RitsuLibSidecarRepeatedWarningLog.Warn(
239-
$"typed-deserialize:opcode={opcode}:sender={context.SenderNetId}:{ex.GetType().FullName}:{ex.Message}",
374+
$"typed-deserialize:opcode={opcode}:sender={senderNetId}:{ex.GetType().FullName}:{ex.Message}",
240375
$"[Sidecar] Typed message deserialize failed opcode={opcode}: {ex.Message}");
241376
return;
242377
}
@@ -249,15 +384,54 @@ private static void HandleDispatch<T>(ulong opcode, Registration<T> registration
249384

250385
var typedContext = new RitsuLibSidecarTypedDispatchContext<T>(
251386
message,
252-
context.SenderNetId,
253-
context.TransferMode,
254-
context.Channel,
255-
context.IsHostIngest);
387+
senderNetId,
388+
transferMode,
389+
channel,
390+
isHostIngest);
256391
foreach (var handler in handlers)
257392
handler(typedContext);
258393

259394
TypedMessageReceived?.Invoke(
260-
new(opcode, registration.ModuleId, registration.MessageKey, context.SenderNetId));
395+
new(opcode, registration.ModuleId, registration.MessageKey, senderNetId));
396+
}
397+
398+
private static Registration<T> GetRegistration<T>(ulong opcode)
399+
{
400+
lock (Gate)
401+
{
402+
return Registrations.TryGetValue(opcode, out var registration) &&
403+
registration is Registration<T> typed
404+
? typed
405+
: throw new InvalidOperationException(
406+
"Typed descriptor registered with incompatible payload type.");
407+
}
408+
}
409+
410+
private static byte[] SerializePayload<T>(Registration<T> registration, T message)
411+
{
412+
return registration.Serialize(message)
413+
?? throw new InvalidOperationException("Typed message serializer returned null.");
414+
}
415+
416+
private static bool CanUseRoutedPath(INetGameService? netService)
417+
{
418+
return ReferenceEquals(netService, RitsuLibSidecarSessionManager.CurrentNetService);
419+
}
420+
421+
private static bool ShouldFallbackToLegacy(RitsuLibSidecarSendStatus status)
422+
{
423+
return status is RitsuLibSidecarSendStatus.RouteUnavailable
424+
or RitsuLibSidecarSendStatus.ProfileUnsupported
425+
or RitsuLibSidecarSendStatus.PayloadTooLarge
426+
or RitsuLibSidecarSendStatus.DestinationUnavailable;
427+
}
428+
429+
private static RitsuLibSidecarDeliveryProfile ResolveDeliveryProfile(
430+
RitsuLibSidecarDeliverySemantics delivery)
431+
{
432+
return delivery == RitsuLibSidecarDeliverySemantics.BestEffort
433+
? RitsuLibSidecarDeliveryProfile.RealtimeDatagram
434+
: RitsuLibSidecarDeliveryProfile.Control;
261435
}
262436

263437
private abstract class RegistrationBase(string moduleId, string messageKey)
@@ -278,6 +452,12 @@ private sealed class Registration<T>(
278452
public Func<ReadOnlySpan<byte>, T> Deserialize { get; } = deserialize;
279453
public RitsuLibSidecarDeliverySemantics Delivery { get; } = delivery;
280454
public List<Action<RitsuLibSidecarTypedDispatchContext<T>>> Handlers { get; } = [];
455+
public RitsuLibSidecarEndpointHandle? EndpointHandle { get; private set; }
456+
457+
public void AttachEndpoint(RitsuLibSidecarEndpointHandle endpoint)
458+
{
459+
EndpointHandle = endpoint;
460+
}
281461
}
282462

283463
private sealed class Subscription(Action dispose) : IDisposable

src/Networking/Sidecar/Core/RitsuLibSidecarProtocol.cs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ public static void EnsureDefaultHandlers()
5252
RitsuLibSidecarSessionManager.EnsureProvidersBootstrapped();
5353
RitsuLibSidecarBuiltInHandlers.Register();
5454
RitsuLibSidecarSyncMessages.RegisterBuiltInHandler();
55+
RitsuLibSidecarEndpointProtocol.EnsureRegistered();
5556
ModRightClickRegistry.RegisterBuiltInSyncDescriptors();
5657
RitsuDebugActionProtocol.EnsureHandlersRegistered();
5758
RitsuLibSidecarNetworkingLifecycle.EnsureHooksInstalled();

0 commit comments

Comments
 (0)