diff --git a/Basis Server/BasisNetworkCore/Protocol/BasisNetworkCommons.cs b/Basis Server/BasisNetworkCore/Protocol/BasisNetworkCommons.cs index 00ee264f4f..da6c55eac7 100644 --- a/Basis Server/BasisNetworkCore/Protocol/BasisNetworkCommons.cs +++ b/Basis Server/BasisNetworkCore/Protocol/BasisNetworkCommons.cs @@ -1083,9 +1083,21 @@ public static int DecodeAvatarIntervalMs(byte encoded, int baseIntervalMs) // ── Server-bound ───────────────────────────────────────────────────── /// Developer hook — data only delivered to the server public const byte ServerBoundChannel = 31; + + // ── Custom Server Data Pub/Sub channel ────────────────────────────────────────────────── + /// Custom Server Data Pub/Sub channel. + public const byte CustomServerDataChannel = 32; + /// Client subscribes to a PubSub channel. + public const byte CustomServerData_Subscribe = 1; + /// Client unsubscribes from a PubSub channel. + public const byte CustomServerData_Unsubscribe = 2; + /// Server sends a PubSub message to the Client. + public const byte CustomServerData_Message = 3; + /// Server sends a PubSub initial state to the Client. + public const byte CustomServerData_InitialState = 4; // ── Admin ──────────────────────────────────────────────────────────── - // Channels 32 & 33 are free (held the removed server-side database). + // Channel 33 is free (32 and 33 previously held the removed server-side database). /// Admin messages from client public const byte AdminChannel = 34; diff --git a/Basis Server/BasisNetworkCore/Serializable/Protocol/BasisCustomServerDataMessages.cs b/Basis Server/BasisNetworkCore/Serializable/Protocol/BasisCustomServerDataMessages.cs new file mode 100644 index 0000000000..e01ffc6e35 --- /dev/null +++ b/Basis Server/BasisNetworkCore/Serializable/Protocol/BasisCustomServerDataMessages.cs @@ -0,0 +1,97 @@ +using System; +using Basis.Network.Core; + +public static partial class SerializableBasis +{ + [Serializable] + public struct CustomServerDataSubscribeRequest + { + public string ChannelName; + public Guid RequestID; + + public void Serialize(NetDataWriter writer) + { + writer.Put(ChannelName); + writer.Put(RequestID); + } + + public bool Deserialize(NetDataReader reader) + { + if (reader.TryGetString(out ChannelName) && reader.AvailableBytes >= 16) + { + RequestID = reader.GetGuid(); + return true; + } + + return false; + } + } + + [Serializable] + public struct CustomServerDataUnsubscribeRequest + { + public string ChannelName; + public Guid RequestID; + + public void Serialize(NetDataWriter writer) + { + writer.Put(ChannelName); + writer.Put(RequestID); + } + + public bool Deserialize(NetDataReader reader) + { + if (reader.TryGetString(out ChannelName) && reader.AvailableBytes >= 16) + { + RequestID = reader.GetGuid(); + return true; + } + + return false; + } + } + + [Serializable] + public struct CustomServerDataMessage + { + public string ChannelName; + public byte[] Data; + + public void Serialize(NetDataWriter writer) + { + writer.Put(ChannelName); + writer.PutBytesWithLength(Data); + } + + public bool Deserialize(NetDataReader reader) + { + return reader.TryGetString(out ChannelName) && reader.TryGetBytesWithLength(out Data); + } + } + + [Serializable] + public struct CustomServerDataInitialState + { + public string ChannelName; + public byte[] Data; + public Guid RequestID; + + public void Serialize(NetDataWriter writer) + { + writer.Put(ChannelName); + writer.PutBytesWithLength(Data); + writer.Put(RequestID); + } + + public bool Deserialize(NetDataReader reader) + { + if (reader.TryGetString(out ChannelName) && reader.TryGetBytesWithLength(out Data) && reader.AvailableBytes >= 16) + { + RequestID = reader.GetGuid(); + return true; + } + + return false; + } + } +} \ No newline at end of file diff --git a/Basis Server/BasisNetworkServer/Core/BasisServerHandleEvents.cs b/Basis Server/BasisNetworkServer/Core/BasisServerHandleEvents.cs index 3004d4b17c..8d0f80664a 100644 --- a/Basis Server/BasisNetworkServer/Core/BasisServerHandleEvents.cs +++ b/Basis Server/BasisNetworkServer/Core/BasisServerHandleEvents.cs @@ -408,6 +408,7 @@ private static bool CleanupPeerSubsystems(NetPeer peer, int id) BasisNetworkPIPCamera.RemovePlayer(id); BasisNetworkContentShare.RemovePlayerSpheres(id); BasisNetworkImageCache.RemovePlayerImages(id); + BasisNetworkHandleCustomServerData.RemovePlayerSubscriptions(id); // Drops this peer's egress bucket and any replay still queued for it. Without this a // recycled player id would inherit the previous holder's spent budget. BasisImageBandwidthGovernor.RemovePeer(id); diff --git a/Basis Server/BasisNetworkServer/Handlers/BasisNetworkHandleCustomServerData.cs b/Basis Server/BasisNetworkServer/Handlers/BasisNetworkHandleCustomServerData.cs new file mode 100644 index 0000000000..c0b16a983c --- /dev/null +++ b/Basis Server/BasisNetworkServer/Handlers/BasisNetworkHandleCustomServerData.cs @@ -0,0 +1,181 @@ +using Basis.Network.Core; +using System; +using System.Collections.Concurrent; +using System.Collections.Generic; +using System.Linq; + +namespace BasisNetworkServer +{ + public interface IBasisCustomServerDataPublisher + { + /// + /// Generates an initial state message for a new subscriber. + /// If it returns an empty list, no initial state will be sent. + /// + List GetInitialState(); + } + + /// + /// Provides a custom server data PubSub service typically for use by props, so that they may arbitrary live information from modified servers. + /// + public static class BasisNetworkHandleCustomServerData + { + private sealed class ChannelState + { + public readonly string Name; + public readonly IBasisCustomServerDataPublisher Publisher; + + public readonly ConcurrentDictionary> PeerToSubscriptionsDict = new(); + public readonly object Lock = new object(); + + public ChannelState(string name, IBasisCustomServerDataPublisher publisher) + { + Name = name; + Publisher = publisher; + } + } + + private static readonly ConcurrentDictionary Channels = new(); + + public static void HandleEvent(NetPeer peer, NetPacketReader reader) + { + if (!reader.TryGetByte(out byte sub)) { reader.Recycle(); return; } + + if (sub == BasisNetworkCommons.CustomServerData_Subscribe) + { + var req = new SerializableBasis.CustomServerDataSubscribeRequest(); + if (req.Deserialize(reader)) + HandleSubscribeRequest(peer, req); + } + else if (sub == BasisNetworkCommons.CustomServerData_Unsubscribe) + { + var req = new SerializableBasis.CustomServerDataUnsubscribeRequest(); + if (req.Deserialize(reader)) + HandleUnsubscribeRequest(peer, req); + } + reader.Recycle(); + } + + public static void RegisterChannel(string name, IBasisCustomServerDataPublisher publisher) + { + if (string.IsNullOrEmpty(name)) throw new ArgumentException("Channel name cannot be empty", nameof(name)); + if (publisher == null) throw new ArgumentNullException(nameof(publisher)); + + if (!Channels.TryAdd(name, new ChannelState(name, publisher))) + { + throw new InvalidOperationException($"Channel '{name}' is already registered."); + } + } + + public static bool UnregisterChannel(string name) + { + return Channels.TryRemove(name, out _); + } + + public static void HandleSubscribeRequest(NetPeer peer, SerializableBasis.CustomServerDataSubscribeRequest request) + { + if (!Channels.TryGetValue(request.ChannelName, out var channel)) + { + BNL.LogWarning($"Peer {peer.Id} tried to subscribe to non-existent channel: {request.ChannelName}"); + return; + } + + List initialStateMessages = null; + lock (channel.Lock) + { + var peerToSubscription = channel.PeerToSubscriptionsDict.GetOrAdd(peer.Id, _ => new HashSet()); + peerToSubscription.Add(request.RequestID); + initialStateMessages = channel.Publisher.GetInitialState(); + } + + foreach (byte[] initialState in initialStateMessages) + { + var initial = new SerializableBasis.CustomServerDataInitialState + { + ChannelName = request.ChannelName, + Data = initialState, + RequestID = request.RequestID + }; + SendMessageToSpecificPeer(peer, BasisNetworkCommons.CustomServerData_InitialState, initial); + } + } + + public static void HandleUnsubscribeRequest(NetPeer peer, SerializableBasis.CustomServerDataUnsubscribeRequest request) + { + if (!Channels.TryGetValue(request.ChannelName, out var channel)) + { + return; + } + + lock (channel.Lock) + { + if (channel.PeerToSubscriptionsDict.TryGetValue(peer.Id, out var peerToSubscription)) + { + peerToSubscription.Remove(request.RequestID); + if (peerToSubscription.Count == 0) + { + channel.PeerToSubscriptionsDict.TryRemove(peer.Id, out _); + } + } + } + } + + public static void Publish(string channelName, byte[] data) + { + if (!Channels.TryGetValue(channelName, out var channel)) + { + return; + } + + int[] targets; + lock (channel.Lock) + { + targets = channel.PeerToSubscriptionsDict.Keys.ToArray(); + } + + if (targets.Length == 0) return; + + var update = new SerializableBasis.CustomServerDataMessage + { + ChannelName = channelName, + Data = data + }; + + NetDataWriter writer = NetworkServer.RentWriter(); + writer.Put(BasisNetworkCommons.CustomServerData_Message); + update.Serialize(writer); + + foreach (var peerId in targets) + { + if (NetworkServer.AuthenticatedPeers.TryGetValue(peerId, out var peer)) + { + peer.Send(writer, BasisNetworkCommons.CustomServerDataChannel, DeliveryMethod.ReliableOrdered); + } + } + + NetworkServer.ReturnWriter(writer); + } + + public static void RemovePlayerSubscriptions(int peerId) + { + foreach (var channel in Channels.Values) + { + lock (channel.Lock) + { + channel.PeerToSubscriptionsDict.TryRemove(peerId, out _); + } + } + } + + private static void SendMessageToSpecificPeer(NetPeer peer, byte subType, SerializableBasis.CustomServerDataInitialState message) + { + NetDataWriter writer = NetworkServer.RentWriter(); + writer.Put(subType); + + message.Serialize(writer); + + peer.Send(writer, BasisNetworkCommons.CustomServerDataChannel, DeliveryMethod.ReliableOrdered); + NetworkServer.ReturnWriter(writer); + } + } +} diff --git a/Basis Server/BasisNetworkServer/Messaging/BasisServerMessageRegistry.cs b/Basis Server/BasisNetworkServer/Messaging/BasisServerMessageRegistry.cs index 82e429b350..361eb307e8 100644 --- a/Basis Server/BasisNetworkServer/Messaging/BasisServerMessageRegistry.cs +++ b/Basis Server/BasisNetworkServer/Messaging/BasisServerMessageRegistry.cs @@ -375,6 +375,9 @@ private static void RegisterCoreHandlers() RegisterCore(BasisNetworkCommons.P2PChannel, (peer, reader, channel, dm) => BasisServerP2PBroker.HandleP2PMessage(reader, peer)); // reads sub-type byte, routes, recycles inside + RegisterCore(BasisNetworkCommons.CustomServerDataChannel, (peer, reader, channel, dm) => + BasisNetworkHandleCustomServerData.HandleEvent(peer, reader)); // recycles inside + RegisterCore(BasisNetworkCommons.RegistryControlChannel, (peer, reader, channel, dm) => { if (reader.TryGetByte(out byte sub) && sub == BasisNetworkCommons.RegistrySub_Subscribe) diff --git a/Basis Server/BasisServerTests/Networking/CustomServerDataTests.cs b/Basis Server/BasisServerTests/Networking/CustomServerDataTests.cs new file mode 100644 index 0000000000..f0b4bebff3 --- /dev/null +++ b/Basis Server/BasisServerTests/Networking/CustomServerDataTests.cs @@ -0,0 +1,154 @@ +using Basis.Network.Core; +using BasisNetworkServer; +using Xunit; +using static SerializableBasis; + +namespace BasisServerTests; + +[Collection("BasisServer shared network statics")] +public class CustomServerDataTests +{ + private class TestDataPublisher : IBasisCustomServerDataPublisher + { + public byte[] InitialState; + public List GetInitialState() => InitialState == null ? new List() : new List { InitialState }; + } + + [Fact] + public void RegisterChannel_SucceedsAndCanReRegister() + { + TestDataPublisher publisher = new TestDataPublisher(); + string channelName = "test.channel.reg"; + + // We can't easily check existence anymore via API, but we can verify registration doesn't throw + BasisNetworkHandleCustomServerData.RegisterChannel(channelName, publisher); + + // Duplicate should still throw + Assert.Throws(() => BasisNetworkHandleCustomServerData.RegisterChannel(channelName, publisher)); + + BasisNetworkHandleCustomServerData.UnregisterChannel(channelName); + + // Should be able to register again after unregister + BasisNetworkHandleCustomServerData.RegisterChannel(channelName, publisher); + BasisNetworkHandleCustomServerData.UnregisterChannel(channelName); + } + + [Fact] + public void RegisterChannel_DuplicateName_ThrowsInvalidOperationException() + { + TestDataPublisher publisher = new TestDataPublisher(); + string channelName = "test.channel.dup"; + BasisNetworkHandleCustomServerData.RegisterChannel(channelName, publisher); + + Assert.Throws(() => BasisNetworkHandleCustomServerData.RegisterChannel(channelName, publisher)); + + BasisNetworkHandleCustomServerData.UnregisterChannel(channelName); + } + + [Fact] + public void HandleSubscribeRequest_SendsInitialStateAndAddsSubscriber() + { + using var scope = new ServerStaticsScope(); + string channelName = "test.subscribe"; + byte[] initialState = new byte[] { 1, 2, 3 }; + var provider = new TestDataPublisher { InitialState = initialState }; + BasisNetworkHandleCustomServerData.RegisterChannel(channelName, provider); + + var peer = new FakeNetPeer(1, "127.0.0.1"); + var requestId = Guid.NewGuid(); + var request = new CustomServerDataSubscribeRequest + { + ChannelName = channelName, + RequestID = requestId + }; + + BasisNetworkHandleCustomServerData.HandleSubscribeRequest(peer, request); + + // Verify initial state sent + var sentMessage = peer.Sent.FirstOrDefault(s => s.Channel == BasisNetworkCommons.CustomServerDataChannel); + Assert.NotNull(sentMessage); + + var reader = new NetDataReader(sentMessage.Data); + Assert.Equal(BasisNetworkCommons.CustomServerData_InitialState, reader.GetByte()); + + var response = new CustomServerDataInitialState(); + Assert.True(response.Deserialize(reader)); + Assert.Equal(channelName, response.ChannelName); + Assert.Equal(initialState, response.Data); + Assert.Equal(requestId, response.RequestID); + + BasisNetworkHandleCustomServerData.UnregisterChannel(channelName); + } + + [Fact] + public void HandleUnsubscribeRequest_RemovesSpecificSubscription() + { + using var scope = new ServerStaticsScope(); + string channelName = "test.unsubscribe"; + var provider = new TestDataPublisher(); + BasisNetworkHandleCustomServerData.RegisterChannel(channelName, provider); + + var peer = new FakeNetPeer(1, "127.0.0.1"); + NetworkServer.AuthenticatedPeers[1] = peer; + NetworkServer.RebuildPeerSnapshot(); + + var id1 = Guid.NewGuid(); + var id2 = Guid.NewGuid(); + + BasisNetworkHandleCustomServerData.HandleSubscribeRequest(peer, new CustomServerDataSubscribeRequest { ChannelName = channelName, RequestID = id1 }); + BasisNetworkHandleCustomServerData.HandleSubscribeRequest(peer, new CustomServerDataSubscribeRequest { ChannelName = channelName, RequestID = id2 }); + + // Both subscribed, verify publish reaches peer + byte[] updateData = new byte[] { 42 }; + BasisNetworkHandleCustomServerData.Publish(channelName, updateData); + + // One update should be sent (peer is unique subscriber) + Assert.Single(peer.Sent.Where(s => s.Data[0] == BasisNetworkCommons.CustomServerData_Message)); + peer.Sent.Clear(); + + // Unsubscribe one + BasisNetworkHandleCustomServerData.HandleUnsubscribeRequest(peer, new CustomServerDataUnsubscribeRequest { ChannelName = channelName, RequestID = id1 }); + + // Still one subscription left, should still receive updates + BasisNetworkHandleCustomServerData.Publish(channelName, updateData); + Assert.Single(peer.Sent.Where(s => s.Data[0] == BasisNetworkCommons.CustomServerData_Message)); + peer.Sent.Clear(); + + // Unsubscribe second + BasisNetworkHandleCustomServerData.HandleUnsubscribeRequest(peer, new CustomServerDataUnsubscribeRequest { ChannelName = channelName, RequestID = id2 }); + + // No more subscriptions, should not receive updates + BasisNetworkHandleCustomServerData.Publish(channelName, updateData); + Assert.Empty(peer.Sent.Where(s => s.Data[0] == BasisNetworkCommons.CustomServerData_Message)); + + BasisNetworkHandleCustomServerData.UnregisterChannel(channelName); + } + + [Fact] + public void RemovePlayerSubscriptions_CleansUpAllChannels() + { + using var scope = new ServerStaticsScope(); + string chan1 = "test.cleanup.1"; + string chan2 = "test.cleanup.2"; + BasisNetworkHandleCustomServerData.RegisterChannel(chan1, new TestDataPublisher()); + BasisNetworkHandleCustomServerData.RegisterChannel(chan2, new TestDataPublisher()); + + var peer = new FakeNetPeer(1, "127.0.0.1"); + NetworkServer.AuthenticatedPeers[1] = peer; + NetworkServer.RebuildPeerSnapshot(); + + BasisNetworkHandleCustomServerData.HandleSubscribeRequest(peer, new CustomServerDataSubscribeRequest { ChannelName = chan1, RequestID = Guid.NewGuid() }); + BasisNetworkHandleCustomServerData.HandleSubscribeRequest(peer, new CustomServerDataSubscribeRequest { ChannelName = chan2, RequestID = Guid.NewGuid() }); + + BasisNetworkHandleCustomServerData.RemovePlayerSubscriptions(peer.Id); + + byte[] updateData = new byte[] { 0 }; + BasisNetworkHandleCustomServerData.Publish(chan1, updateData); + BasisNetworkHandleCustomServerData.Publish(chan2, updateData); + + Assert.Empty(peer.Sent.Where(s => s.Data[0] == BasisNetworkCommons.CustomServerData_Message)); + + BasisNetworkHandleCustomServerData.UnregisterChannel(chan1); + BasisNetworkHandleCustomServerData.UnregisterChannel(chan2); + } +} diff --git a/Basis/Packages/com.basis.framework/Networking/BasisNetworkEvents.cs b/Basis/Packages/com.basis.framework/Networking/BasisNetworkEvents.cs index 76e3f0558a..8dbe47475c 100644 --- a/Basis/Packages/com.basis.framework/Networking/BasisNetworkEvents.cs +++ b/Basis/Packages/com.basis.framework/Networking/BasisNetworkEvents.cs @@ -629,6 +629,16 @@ private static void RegisterCoreHandlers() BasisP2PManager.HandleServerMessage(Reader); }); + BasisClientMessageRegistry.RegisterCore(BasisNetworkCommons.CustomServerDataChannel, (peer, Reader, channel, deliveryMethod) => + { + if (ValidateSize(Reader, peer, channel) == false) + { + Reader.Recycle(); + return; + } + BasisNetworkHandleCustomServerData.HandleMessage(Reader, deliveryMethod); + }); + BasisClientMessageRegistry.RegisterCore(BasisNetworkCommons.EventsChannel, (peer, Reader, channel, deliveryMethod) => { if (ValidateSize(Reader, peer, channel) == false) diff --git a/Basis/Packages/com.basis.framework/Networking/Handles/BasisNetworkHandleCustomServerData.cs b/Basis/Packages/com.basis.framework/Networking/Handles/BasisNetworkHandleCustomServerData.cs new file mode 100644 index 0000000000..5c4f88bff6 --- /dev/null +++ b/Basis/Packages/com.basis.framework/Networking/Handles/BasisNetworkHandleCustomServerData.cs @@ -0,0 +1,25 @@ +using Basis.Network.Core; + +namespace Basis.Scripts.Networking +{ + public static class BasisNetworkHandleCustomServerData + { + public delegate void CustomServerDataMessageDelegate(byte[] buffer, DeliveryMethod deliveryMethod); + public static event CustomServerDataMessageDelegate OnCustomServerDataMessageReceived; + + public static void HandleMessage(NetPacketReader reader, DeliveryMethod deliveryMethod) + { + if (OnCustomServerDataMessageReceived == null) return; + + try + { + byte[] data = reader.GetRemainingBytes(); + OnCustomServerDataMessageReceived?.Invoke(data, deliveryMethod); + } + finally + { + reader.Recycle(); + } + } + } +} diff --git a/Basis/Packages/com.basis.framework/Networking/Handles/BasisNetworkHandleCustomServerData.cs.meta b/Basis/Packages/com.basis.framework/Networking/Handles/BasisNetworkHandleCustomServerData.cs.meta new file mode 100644 index 0000000000..7f434ff021 --- /dev/null +++ b/Basis/Packages/com.basis.framework/Networking/Handles/BasisNetworkHandleCustomServerData.cs.meta @@ -0,0 +1,3 @@ +fileFormatVersion: 2 +guid: 2d81c3a04322405da46396d16a0f3f7b +timeCreated: 1788768535 \ No newline at end of file diff --git a/Basis/Packages/com.basis.server/BasisNetworkCore/Protocol/BasisNetworkCommons.cs b/Basis/Packages/com.basis.server/BasisNetworkCore/Protocol/BasisNetworkCommons.cs index 00ee264f4f..da6c55eac7 100644 --- a/Basis/Packages/com.basis.server/BasisNetworkCore/Protocol/BasisNetworkCommons.cs +++ b/Basis/Packages/com.basis.server/BasisNetworkCore/Protocol/BasisNetworkCommons.cs @@ -1083,9 +1083,21 @@ public static int DecodeAvatarIntervalMs(byte encoded, int baseIntervalMs) // ── Server-bound ───────────────────────────────────────────────────── /// Developer hook — data only delivered to the server public const byte ServerBoundChannel = 31; + + // ── Custom Server Data Pub/Sub channel ────────────────────────────────────────────────── + /// Custom Server Data Pub/Sub channel. + public const byte CustomServerDataChannel = 32; + /// Client subscribes to a PubSub channel. + public const byte CustomServerData_Subscribe = 1; + /// Client unsubscribes from a PubSub channel. + public const byte CustomServerData_Unsubscribe = 2; + /// Server sends a PubSub message to the Client. + public const byte CustomServerData_Message = 3; + /// Server sends a PubSub initial state to the Client. + public const byte CustomServerData_InitialState = 4; // ── Admin ──────────────────────────────────────────────────────────── - // Channels 32 & 33 are free (held the removed server-side database). + // Channel 33 is free (32 and 33 previously held the removed server-side database). /// Admin messages from client public const byte AdminChannel = 34; diff --git a/Basis/Packages/com.basis.server/BasisNetworkCore/Serializable/Protocol/BasisCustomServerDataMessages.cs b/Basis/Packages/com.basis.server/BasisNetworkCore/Serializable/Protocol/BasisCustomServerDataMessages.cs new file mode 100644 index 0000000000..e01ffc6e35 --- /dev/null +++ b/Basis/Packages/com.basis.server/BasisNetworkCore/Serializable/Protocol/BasisCustomServerDataMessages.cs @@ -0,0 +1,97 @@ +using System; +using Basis.Network.Core; + +public static partial class SerializableBasis +{ + [Serializable] + public struct CustomServerDataSubscribeRequest + { + public string ChannelName; + public Guid RequestID; + + public void Serialize(NetDataWriter writer) + { + writer.Put(ChannelName); + writer.Put(RequestID); + } + + public bool Deserialize(NetDataReader reader) + { + if (reader.TryGetString(out ChannelName) && reader.AvailableBytes >= 16) + { + RequestID = reader.GetGuid(); + return true; + } + + return false; + } + } + + [Serializable] + public struct CustomServerDataUnsubscribeRequest + { + public string ChannelName; + public Guid RequestID; + + public void Serialize(NetDataWriter writer) + { + writer.Put(ChannelName); + writer.Put(RequestID); + } + + public bool Deserialize(NetDataReader reader) + { + if (reader.TryGetString(out ChannelName) && reader.AvailableBytes >= 16) + { + RequestID = reader.GetGuid(); + return true; + } + + return false; + } + } + + [Serializable] + public struct CustomServerDataMessage + { + public string ChannelName; + public byte[] Data; + + public void Serialize(NetDataWriter writer) + { + writer.Put(ChannelName); + writer.PutBytesWithLength(Data); + } + + public bool Deserialize(NetDataReader reader) + { + return reader.TryGetString(out ChannelName) && reader.TryGetBytesWithLength(out Data); + } + } + + [Serializable] + public struct CustomServerDataInitialState + { + public string ChannelName; + public byte[] Data; + public Guid RequestID; + + public void Serialize(NetDataWriter writer) + { + writer.Put(ChannelName); + writer.PutBytesWithLength(Data); + writer.Put(RequestID); + } + + public bool Deserialize(NetDataReader reader) + { + if (reader.TryGetString(out ChannelName) && reader.TryGetBytesWithLength(out Data) && reader.AvailableBytes >= 16) + { + RequestID = reader.GetGuid(); + return true; + } + + return false; + } + } +} \ No newline at end of file diff --git a/Basis/Packages/com.basis.server/BasisNetworkCore/Serializable/Protocol/BasisCustomServerDataMessages.cs.meta b/Basis/Packages/com.basis.server/BasisNetworkCore/Serializable/Protocol/BasisCustomServerDataMessages.cs.meta new file mode 100644 index 0000000000..8526b4e515 --- /dev/null +++ b/Basis/Packages/com.basis.server/BasisNetworkCore/Serializable/Protocol/BasisCustomServerDataMessages.cs.meta @@ -0,0 +1,3 @@ +fileFormatVersion: 2 +guid: b87b926cddf341f58a9eb3db157668d8 +timeCreated: 1788773252 \ No newline at end of file diff --git a/Basis/Packages/com.basis.server/BasisNetworkServer/Core/BasisServerHandleEvents.cs b/Basis/Packages/com.basis.server/BasisNetworkServer/Core/BasisServerHandleEvents.cs index 3733e6b98b..d617b4eb49 100644 --- a/Basis/Packages/com.basis.server/BasisNetworkServer/Core/BasisServerHandleEvents.cs +++ b/Basis/Packages/com.basis.server/BasisNetworkServer/Core/BasisServerHandleEvents.cs @@ -387,6 +387,7 @@ private static bool CleanupPeerSubsystems(NetPeer peer, int id) BasisNetworkPIPCamera.RemovePlayer(id); BasisNetworkContentShare.RemovePlayerSpheres(id); BasisNetworkImageCache.RemovePlayerImages(id); + BasisNetworkHandleCustomServerData.RemovePlayerSubscriptions(id); // Drops this peer's egress bucket and any replay still queued for it. Without this a // recycled player id would inherit the previous holder's spent budget. BasisImageBandwidthGovernor.RemovePeer(id); diff --git a/Basis/Packages/com.basis.server/BasisNetworkServer/Handlers/BasisNetworkHandlePubSub.cs b/Basis/Packages/com.basis.server/BasisNetworkServer/Handlers/BasisNetworkHandlePubSub.cs new file mode 100644 index 0000000000..1c2eecca21 --- /dev/null +++ b/Basis/Packages/com.basis.server/BasisNetworkServer/Handlers/BasisNetworkHandlePubSub.cs @@ -0,0 +1,181 @@ +using Basis.Network.Core; +using System; +using System.Collections.Concurrent; +using System.Collections.Generic; +using System.Linq; + +namespace BasisNetworkServer +{ + public interface IPubSubDataProvider + { + /// + /// Generates an initial state message for a new subscriber. + /// If it returns an empty list, no initial state will be sent. + /// + List GetInitialState(); + } + + /// + /// Provides a PubSub service typically for use by props, so that they may arbitrary live information from modified servers. + /// + public static class BasisNetworkHandlePubSub + { + private sealed class ChannelState + { + public readonly string Name; + public readonly IPubSubDataProvider Provider; + + public readonly ConcurrentDictionary> PeerToSubscriptionsDict = new(); + public readonly object Lock = new object(); + + public ChannelState(string name, IPubSubDataProvider provider) + { + Name = name; + Provider = provider; + } + } + + private static readonly ConcurrentDictionary Channels = new(); + + public static void HandleEvent(NetPeer peer, NetPacketReader reader) + { + if (!reader.TryGetByte(out byte sub)) { reader.Recycle(); return; } + + if (sub == BasisNetworkCommons.PubSub_Subscribe) + { + var req = new SerializableBasis.PubSubSubscribeRequest(); + if (req.Deserialize(reader)) + HandleSubscribeRequest(peer, req); + } + else if (sub == BasisNetworkCommons.PubSub_Unsubscribe) + { + var req = new SerializableBasis.PubSubUnsubscribeRequest(); + if (req.Deserialize(reader)) + HandleUnsubscribeRequest(peer, req); + } + reader.Recycle(); + } + + public static void RegisterChannel(string name, IPubSubDataProvider provider) + { + if (string.IsNullOrEmpty(name)) throw new ArgumentException("Channel name cannot be empty", nameof(name)); + if (provider == null) throw new ArgumentNullException(nameof(provider)); + + if (!Channels.TryAdd(name, new ChannelState(name, provider))) + { + throw new InvalidOperationException($"Channel '{name}' is already registered."); + } + } + + public static bool UnregisterChannel(string name) + { + return Channels.TryRemove(name, out _); + } + + public static void HandleSubscribeRequest(NetPeer peer, SerializableBasis.PubSubSubscribeRequest request) + { + if (!Channels.TryGetValue(request.ChannelName, out var channel)) + { + BNL.LogWarning($"Peer {peer.Id} tried to subscribe to non-existent channel: {request.ChannelName}"); + return; + } + + List initialStateMessages = null; + lock (channel.Lock) + { + var peerToSubscription = channel.PeerToSubscriptionsDict.GetOrAdd(peer.Id, _ => new HashSet()); + peerToSubscription.Add(request.RequestID); + initialStateMessages = channel.Provider.GetInitialState(); + } + + foreach (byte[] initialState in initialStateMessages) + { + var initial = new SerializableBasis.PubSubInitialState + { + ChannelName = request.ChannelName, + Data = initialState, + RequestID = request.RequestID + }; + SendMessageToSpecificPeer(peer, BasisNetworkCommons.PubSub_Initial, initial); + } + } + + public static void HandleUnsubscribeRequest(NetPeer peer, SerializableBasis.PubSubUnsubscribeRequest request) + { + if (!Channels.TryGetValue(request.ChannelName, out var channel)) + { + return; + } + + lock (channel.Lock) + { + if (channel.PeerToSubscriptionsDict.TryGetValue(peer.Id, out var peerToSubscription)) + { + peerToSubscription.Remove(request.RequestID); + if (peerToSubscription.Count == 0) + { + channel.PeerToSubscriptionsDict.TryRemove(peer.Id, out _); + } + } + } + } + + public static void Publish(string channelName, byte[] data) + { + if (!Channels.TryGetValue(channelName, out var channel)) + { + return; + } + + int[] targets; + lock (channel.Lock) + { + targets = channel.PeerToSubscriptionsDict.Keys.ToArray(); + } + + if (targets.Length == 0) return; + + var update = new SerializableBasis.PubSubMessage + { + ChannelName = channelName, + Data = data + }; + + NetDataWriter writer = NetworkServer.RentWriter(); + writer.Put(BasisNetworkCommons.PubSub_Message); + update.Serialize(writer); + + foreach (var peerId in targets) + { + if (NetworkServer.AuthenticatedPeers.TryGetValue(peerId, out var peer)) + { + peer.Send(writer, BasisNetworkCommons.PubSubChannel, DeliveryMethod.ReliableOrdered); + } + } + + NetworkServer.ReturnWriter(writer); + } + + public static void RemovePlayerSubscriptions(int peerId) + { + foreach (var channel in Channels.Values) + { + lock (channel.Lock) + { + channel.PeerToSubscriptionsDict.TryRemove(peerId, out _); + } + } + } + + private static void SendMessageToSpecificPeer(NetPeer peer, byte subType, SerializableBasis.PubSubInitialState message) + { + NetDataWriter writer = NetworkServer.RentWriter(); + writer.Put(subType); + + message.Serialize(writer); + + peer.Send(writer, BasisNetworkCommons.PubSubChannel, DeliveryMethod.ReliableOrdered); + NetworkServer.ReturnWriter(writer); + } + } +} diff --git a/Basis/Packages/com.basis.server/BasisNetworkServer/Handlers/BasisNetworkHandlePubSub.cs.meta b/Basis/Packages/com.basis.server/BasisNetworkServer/Handlers/BasisNetworkHandlePubSub.cs.meta new file mode 100644 index 0000000000..551f2483ce --- /dev/null +++ b/Basis/Packages/com.basis.server/BasisNetworkServer/Handlers/BasisNetworkHandlePubSub.cs.meta @@ -0,0 +1,2 @@ +fileFormatVersion: 2 +guid: 7c70ae3d16e6a634da04add2af5a64f1 \ No newline at end of file diff --git a/Basis/Packages/com.basis.server/BasisNetworkServer/Messaging/BasisServerMessageRegistry.cs b/Basis/Packages/com.basis.server/BasisNetworkServer/Messaging/BasisServerMessageRegistry.cs index 82e429b350..361eb307e8 100644 --- a/Basis/Packages/com.basis.server/BasisNetworkServer/Messaging/BasisServerMessageRegistry.cs +++ b/Basis/Packages/com.basis.server/BasisNetworkServer/Messaging/BasisServerMessageRegistry.cs @@ -375,6 +375,9 @@ private static void RegisterCoreHandlers() RegisterCore(BasisNetworkCommons.P2PChannel, (peer, reader, channel, dm) => BasisServerP2PBroker.HandleP2PMessage(reader, peer)); // reads sub-type byte, routes, recycles inside + RegisterCore(BasisNetworkCommons.CustomServerDataChannel, (peer, reader, channel, dm) => + BasisNetworkHandleCustomServerData.HandleEvent(peer, reader)); // recycles inside + RegisterCore(BasisNetworkCommons.RegistryControlChannel, (peer, reader, channel, dm) => { if (reader.TryGetByte(out byte sub) && sub == BasisNetworkCommons.RegistrySub_Subscribe) diff --git a/Basis/Packages/com.basis.shim/Shims/BasisCustomServerDataSubscriberShim.cs b/Basis/Packages/com.basis.shim/Shims/BasisCustomServerDataSubscriberShim.cs new file mode 100644 index 0000000000..e1c6be6358 --- /dev/null +++ b/Basis/Packages/com.basis.shim/Shims/BasisCustomServerDataSubscriberShim.cs @@ -0,0 +1,158 @@ +using System; +using System.Collections.Generic; +using Basis.Network.Core; +using Basis.Scripts.Networking; + +namespace Basis.Shims +{ + /// + /// Provides access to the custom data CustomServerData channel of the currently connected server.
+ ///
+ /// The custom server data channels are meant for the server to provide custom live data to props + /// or other content that requests it, such as a list of players who recently joined the server, + /// weather forecasts, integration with Discord, etc.
+ ///
+ /// The available capabilities are entirely dependent on the server.
+ ///
+ /// You should call UnsubscribeAll() in your OnDestroy() method. + ///
+ /// The custom server data channels are NOT designed for props or other content to communicate with each other: + /// the clients cannot publish messages to the channels; only the server can send messages to the clients. + ///
+ public class BasisCustomServerDataSubscriberShim + { + private readonly Dictionary _subscriptions = new(); + private bool _isHooked; + + public delegate void InitialStateReceivedDelegate(string channelName, byte[] data); + public delegate void MessageReceivedDelegate(string channelName, byte[] data); + + public event InitialStateReceivedDelegate InitialStateReceived; + public event MessageReceivedDelegate MessageReceived; + + ~BasisCustomServerDataSubscriberShim() + { + BasisNetworkHandleCustomServerData.OnCustomServerDataMessageReceived -= OnCustomServerDataMessageReceived; + } + + private void OnCustomServerDataMessageReceived(byte[] buffer, DeliveryMethod deliveryMethod) + { + try + { + if (buffer == null || buffer.Length == 0) return; + + var reader = new NetDataReader(buffer); + if (!reader.TryGetByte(out byte subType)) return; + + if (subType == BasisNetworkCommons.CustomServerData_Message) + { + var msg = new SerializableBasis.CustomServerDataMessage(); + if (msg.Deserialize(reader)) + { + if (_subscriptions.ContainsKey(msg.ChannelName)) + { + MessageReceived?.Invoke(msg.ChannelName, msg.Data); + } + } + } + else if (subType == BasisNetworkCommons.CustomServerData_InitialState) + { + var initialState = new SerializableBasis.CustomServerDataInitialState(); + if (initialState.Deserialize(reader)) + { + if (_subscriptions.ContainsKey(initialState.ChannelName)) + { + InitialStateReceived?.Invoke(initialState.ChannelName, initialState.Data); + } + } + } + } + catch (Exception e) + { + BasisDebug.LogError($"[BasisCustomServerDataShim] Error while processing CustomServerData message, this exception will not be re-thrown: {e}"); + // Do not throw the exception, as we want to avoid Cilbox props disrupting the network message processing. + } + } + + /// + /// Subscribes to a CustomServerData channel.
+ ///
+ /// Channels can only be subscribed to once per instance of the shim.
+ /// When subscribing, the InitialStateReceived will trigger if that channel provides an initial state.
+ /// Multiple props can subscribe to the same channel:
+ /// - Each prop may receive a different initial state, depending on when that prop subscribes.
+ /// - Non-initial state messages are sent from the server to the user once, and then dispatched to all the shims that require it. + ///
+ /// Name of the CustomServerData channel + public void Subscribe(string channelName) + { + if (string.IsNullOrEmpty(channelName)) return; + if (_subscriptions.TryGetValue(channelName, out _)) return; + + if (!_isHooked) + { + _isHooked = true; + BasisNetworkHandleCustomServerData.OnCustomServerDataMessageReceived += OnCustomServerDataMessageReceived; + } + + Guid requestID = Guid.NewGuid(); + _subscriptions[channelName] = requestID; + + var request = new SerializableBasis.CustomServerDataSubscribeRequest + { + ChannelName = channelName, + RequestID = requestID + }; + + NetDataWriter writer = new NetDataWriter(); + writer.Put(BasisNetworkCommons.CustomServerData_Subscribe); + request.Serialize(writer); + + BasisNetworkConnection.LocalPlayerPeer?.Send(writer, BasisNetworkCommons.CustomServerDataChannel, DeliveryMethod.ReliableOrdered); + } + + /// + /// Unsubscribes from a CustomServerData channel. + /// + /// + public void Unsubscribe(string channelName) + { + if (string.IsNullOrEmpty(channelName)) return; + if (!_subscriptions.TryGetValue(channelName, out Guid requestID)) return; + + var request = new SerializableBasis.CustomServerDataUnsubscribeRequest + { + ChannelName = channelName, + RequestID = requestID + }; + + NetDataWriter writer = new NetDataWriter(); + writer.Put(BasisNetworkCommons.CustomServerData_Unsubscribe); + request.Serialize(writer); + + BasisNetworkConnection.LocalPlayerPeer?.Send(writer, BasisNetworkCommons.CustomServerDataChannel, DeliveryMethod.ReliableOrdered); + _subscriptions.Remove(channelName); + + if (_subscriptions.Count == 0) + { + BasisNetworkHandleCustomServerData.OnCustomServerDataMessageReceived -= OnCustomServerDataMessageReceived; + _isHooked = false; + } + } + + /// + /// Unsubscribes from all CustomServerData channels. + /// + public void UnsubscribeAll() + { + List channels = new List(_subscriptions.Keys); + foreach (var channel in channels) + { + Unsubscribe(channel); + } + + BasisNetworkHandleCustomServerData.OnCustomServerDataMessageReceived -= OnCustomServerDataMessageReceived; + _isHooked = false; + } + } +} diff --git a/Basis/Packages/com.basis.shim/Shims/BasisCustomServerDataSubscriberShim.cs.meta b/Basis/Packages/com.basis.shim/Shims/BasisCustomServerDataSubscriberShim.cs.meta new file mode 100644 index 0000000000..d58327193d --- /dev/null +++ b/Basis/Packages/com.basis.shim/Shims/BasisCustomServerDataSubscriberShim.cs.meta @@ -0,0 +1,3 @@ +fileFormatVersion: 2 +guid: b44694c8f1124dc693b03822876a5451 +timeCreated: 1788764341 \ No newline at end of file