Tikfinity/Runtime/TikfinityHub.cs
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
using Sandbox;
namespace TikfinitySbox;
/// <summary>
/// Place one in the active game scene. The host connects to TikFinity and owns all routing.
/// </summary>
public sealed class TikfinityHub : Component
{
[Property, Group( "Setup" )] public TikfinitySettings Settings { get; set; }
[Property, Group( "Setup" )] public bool AutoConnect { get; set; } = true;
[Property, ReadOnly, Group( "Runtime" )] public TikfinityConnectionState ConnectionState { get; private set; } = TikfinityConnectionState.Disabled;
[Property, ReadOnly, Group( "Runtime" )] public int PendingMessages => _pendingMessages.Count;
[Property, ReadOnly, Group( "Runtime" )] public string LastStatus { get; private set; } = "Not started";
readonly Queue<string> _pendingMessages = new();
readonly TikfinityDeduplicator _deduplicator = new();
readonly TikfinityRateLimiter _rateLimiter = new();
readonly Dictionary<string, double> _cooldowns = new( StringComparer.OrdinalIgnoreCase );
readonly List<TikfinityEvent> _parseBuffer = new();
TikfinitySettings _fallbackSettings;
WebSocket _socket;
bool _connectionLoopRunning;
bool _stopRequested;
int _droppedMessages;
double _lastCooldownPrune;
protected override void OnStart()
{
var primaryHub = Scene.GetAllComponents<TikfinityHub>().Where( hub => hub.Enabled ).FirstOrDefault();
if ( primaryHub != this )
{
ConnectionState = TikfinityConnectionState.Disabled;
LastStatus = "Duplicate hub disabled";
Log.Warning( "TikFinity: only one TikfinityHub may be active in a scene." );
Enabled = false;
return;
}
if ( !AutoConnect )
{
ConnectionState = TikfinityConnectionState.Disabled;
LastStatus = "Auto-connect disabled";
return;
}
EnsureConnectionLoop();
}
protected override void OnUpdate()
{
if ( !Networking.IsHost )
{
if ( _socket is not null ) StopConnection();
ConnectionState = TikfinityConnectionState.WaitingForHost;
LastStatus = "Only the session host connects to TikFinity";
return;
}
if ( AutoConnect ) EnsureConnectionLoop();
DrainMessages();
}
protected override void OnDestroy()
{
StopConnection();
}
public void Connect()
{
AutoConnect = true;
_stopRequested = false;
EnsureConnectionLoop();
}
public void Disconnect()
{
AutoConnect = false;
StopConnection();
ConnectionState = TikfinityConnectionState.Stopped;
LastStatus = "Disconnected manually";
}
public void DispatchTestInteraction( string route, int amount = 1, string nickname = "Local Tester" )
{
if ( !Networking.IsHost ) return;
if ( string.IsNullOrWhiteSpace( route ) ) return;
var maximum = ActiveSettings.MaxAmountPerInteraction;
var safeAmount = Clamp( amount, 1, maximum < 1 ? 1 : maximum );
DispatchInteraction( new TikfinityInteraction
{
Route = route.Trim(),
Amount = safeAmount,
Username = "local_test",
Nickname = Limit( nickname, TikfinityEventParser.MaxDisplayLength ),
GiftName = "Test Gift",
IsTest = true
} );
}
void EnsureConnectionLoop()
{
if ( _connectionLoopRunning || !Networking.IsHost || !AutoConnect ) return;
_stopRequested = false;
_ = RunConnectionLoop();
}
async Task RunConnectionLoop()
{
_connectionLoopRunning = true;
_stopRequested = false;
var attempt = 0;
try
{
while ( GameObject.IsValid() && Enabled && AutoConnect && Networking.IsHost && !_stopRequested )
{
var settings = ActiveSettings;
if ( !TryValidateEndpoint( settings.Endpoint, out var endpoint, out var validationError ) )
{
ConnectionState = TikfinityConnectionState.BlockedEndpoint;
LastStatus = validationError;
Log.Warning( $"TikFinity: {validationError}" );
return;
}
CleanupSocket();
var messageSize = Clamp( settings.MaxMessageBytes, 1024, 262144 );
_socket = new WebSocket( messageSize );
_socket.OnMessageReceived += HandleMessage;
try
{
ConnectionState = attempt == 0 ? TikfinityConnectionState.Connecting : TikfinityConnectionState.Reconnecting;
LastStatus = $"Connecting to {endpoint}";
await _socket.Connect( endpoint );
if ( !GameObject.IsValid() || _stopRequested ) return;
attempt = 0;
ConnectionState = TikfinityConnectionState.Connected;
LastStatus = "Connected to TikFinity";
Log.Info( "TikFinity: connected to local event server." );
while ( GameObject.IsValid() && !_stopRequested && AutoConnect && Networking.IsHost && _socket is not null && _socket.IsConnected )
await Task.DelayRealtimeSeconds( 0.25f );
}
catch ( Exception ex )
{
LastStatus = $"Connection failed: {Limit( ex.Message, 160 )}";
if ( settings.DebugLogging ) Log.Warning( $"TikFinity: {LastStatus}" );
}
finally
{
CleanupSocket();
}
if ( _stopRequested || !AutoConnect || !Networking.IsHost ) break;
attempt++;
ConnectionState = TikfinityConnectionState.Reconnecting;
var initialDelay = settings.InitialReconnectSeconds < 0.25f ? 0.25f : settings.InitialReconnectSeconds;
var maximumDelay = settings.MaxReconnectSeconds < 1f ? 1f : settings.MaxReconnectSeconds;
var exponent = attempt > 8 ? 8 : attempt;
var delay = MathF.Min( maximumDelay, initialDelay * MathF.Pow( 2f, exponent - 1 ) );
LastStatus = $"Retrying in {delay:0.0}s";
await Task.DelayRealtimeSeconds( delay );
}
}
finally
{
CleanupSocket();
_connectionLoopRunning = false;
if ( ConnectionState != TikfinityConnectionState.BlockedEndpoint )
ConnectionState = AutoConnect ? TikfinityConnectionState.WaitingForHost : TikfinityConnectionState.Stopped;
}
}
void HandleMessage( string message )
{
var settings = ActiveSettings;
if ( string.IsNullOrWhiteSpace( message ) ) return;
if ( message.Length > Clamp( settings.MaxMessageBytes, 1024, 262144 ) )
{
_droppedMessages++;
return;
}
var maximum = Clamp( settings.MaxPendingMessages, 1, 4096 );
if ( _pendingMessages.Count >= maximum )
{
_droppedMessages++;
return;
}
_pendingMessages.Enqueue( message );
}
void DrainMessages()
{
if ( !Networking.IsHost ) return;
var settings = ActiveSettings;
var budget = Clamp( settings.MaxMessagesPerFrame, 1, 128 );
for ( var i = 0; i < budget && _pendingMessages.Count > 0; i++ )
{
var message = _pendingMessages.Dequeue();
_parseBuffer.Clear();
if ( !TikfinityEventParser.TryParseMany( message, _parseBuffer, out var error ) )
{
if ( settings.DebugLogging ) Log.Warning( $"TikFinity: discarded invalid message ({error})." );
continue;
}
foreach ( var evt in _parseBuffer ) ProcessEvent( evt, settings );
}
if ( _droppedMessages > 0 && settings.DebugLogging )
{
Log.Warning( $"TikFinity: dropped {_droppedMessages} over-limit messages." );
_droppedMessages = 0;
}
}
void ProcessEvent( TikfinityEvent evt, TikfinitySettings settings )
{
if ( evt is null ) return;
var now = (double)Time.Now;
PruneCooldowns( now );
// TikFinity can emit many intermediate updates for one streak. Do not let
// those updates consume gameplay rate-limit budget or execute any route.
if ( evt.Kind == TikfinityEventKind.Gift && settings.WaitForGiftStreakEnd && evt.HasRepeatEnd && !evt.RepeatEnd )
return;
if ( !_deduplicator.TryRemember( evt.Fingerprint(), now, settings.DuplicateWindowSeconds, settings.MaxDuplicateEntries ) )
return;
var userKey = !string.IsNullOrWhiteSpace( evt.UserId ) ? evt.UserId : evt.Username;
if ( !_rateLimiter.TryTake( evt.Kind.ToString(), userKey, now, settings.GlobalEventsPerSecond, settings.UserEventsPerSecond ) )
return;
ITikfinityEvents.Post( listener => listener.OnEvent( evt ) );
if ( evt.Kind != TikfinityEventKind.Gift ) return;
var rules = settings.GiftRules ?? new List<TikfinityGiftRule>();
foreach ( var rule in rules.Where( rule => rule is not null && rule.RuleEnabled ).OrderByDescending( rule => rule.Priority ) )
{
if ( !rule.Matches( evt ) ) continue;
if ( IsCoolingDown( rule, evt, now ) ) continue;
var interaction = new TikfinityInteraction
{
Route = rule.Route.Trim(),
Amount = rule.CalculateAmount( evt, settings.MaxAmountPerInteraction ),
UserId = evt.UserId,
Username = evt.Username,
Nickname = DisplayNameFor( evt ),
GiftId = evt.GiftId,
GiftName = evt.GiftName,
GiftCoins = evt.GiftCoins,
TotalCoins = evt.TotalCoins,
RepeatCount = evt.RepeatCount
};
StartCooldown( rule, evt, now );
DispatchInteraction( interaction );
}
}
void DispatchInteraction( TikfinityInteraction interaction )
{
if ( interaction is null || !Networking.IsHost ) return;
ITikfinityEvents.Post( listener => listener.OnInteraction( interaction ) );
if ( ActiveSettings.BroadcastPresentationToClients )
BroadcastPresentation( interaction.Route, interaction.Nickname, interaction.GiftName, interaction.Amount, interaction.IsTest );
else
ITikfinityEvents.Post( listener => listener.OnPresentation( interaction.Route, interaction.Nickname, interaction.GiftName, interaction.Amount, interaction.IsTest ) );
}
[Rpc.Broadcast( NetFlags.HostOnly | NetFlags.Reliable )]
void BroadcastPresentation( string route, string nickname, string giftName, int amount, bool isTest )
{
if ( Rpc.Calling && !Rpc.Caller.IsHost ) return;
ITikfinityEvents.Post( listener => listener.OnPresentation(
Limit( route, 96 ),
Limit( nickname, TikfinityEventParser.MaxDisplayLength ),
Limit( giftName, TikfinityEventParser.MaxDisplayLength ),
Clamp( amount, 1, ActiveSettings.MaxAmountPerInteraction ),
isTest ) );
}
bool IsCoolingDown( TikfinityGiftRule rule, TikfinityEvent evt, double now )
{
if ( rule.CooldownSeconds <= 0f ) return false;
var key = CooldownKey( rule, evt );
return _cooldowns.TryGetValue( key, out var expiry ) && expiry > now;
}
void StartCooldown( TikfinityGiftRule rule, TikfinityEvent evt, double now )
{
if ( rule.CooldownSeconds <= 0f ) return;
_cooldowns[CooldownKey( rule, evt )] = now + rule.CooldownSeconds;
}
void PruneCooldowns( double now )
{
if ( now - _lastCooldownPrune < 30.0 ) return;
_lastCooldownPrune = now;
var remove = new List<string>();
foreach ( var pair in _cooldowns )
{
if ( pair.Value <= now ) remove.Add( pair.Key );
}
foreach ( var key in remove ) _cooldowns.Remove( key );
}
static string CooldownKey( TikfinityGiftRule rule, TikfinityEvent evt )
{
var key = rule.Route.Trim();
if ( !rule.CooldownPerUser ) return key;
var user = !string.IsNullOrWhiteSpace( evt.UserId ) ? evt.UserId : evt.Username;
return $"{key}|{user}";
}
void StopConnection()
{
_stopRequested = true;
CleanupSocket();
_pendingMessages.Clear();
}
void CleanupSocket()
{
if ( _socket is null ) return;
_socket.OnMessageReceived -= HandleMessage;
_socket.Dispose();
_socket = null;
}
TikfinitySettings ActiveSettings => Settings ?? (_fallbackSettings ??= new TikfinitySettings());
static bool TryValidateEndpoint( string endpoint, out string normalized, out string error )
{
normalized = "";
error = "";
if ( string.IsNullOrWhiteSpace( endpoint ) )
{
error = "The TikFinity endpoint is empty.";
return false;
}
var value = endpoint.Trim();
var lower = value.ToLowerInvariant();
var prefixes = new[]
{
"ws://127.0.0.1:",
"ws://localhost:",
"ws://[::1]:"
};
var prefixLength = 0;
foreach ( var prefix in prefixes )
{
if ( !lower.StartsWith( prefix, StringComparison.Ordinal ) ) continue;
prefixLength = prefix.Length;
break;
}
if ( prefixLength == 0 )
{
error = "The endpoint must start with ws://127.0.0.1:, ws://localhost:, or ws://[::1]:.";
return false;
}
var pathIndex = value.IndexOf( '/', prefixLength );
var portText = pathIndex < 0
? value.Substring( prefixLength )
: value.Substring( prefixLength, pathIndex - prefixLength );
if ( portText.IndexOfAny( new[] { '?', '#', '@' } ) >= 0 || !int.TryParse( portText, out var port ) )
{
error = "The TikFinity endpoint must contain a numeric loopback port.";
return false;
}
if ( port < 1 || port > 65535 )
{
error = "The TikFinity endpoint port is invalid.";
return false;
}
normalized = value;
return true;
}
static string DisplayNameFor( TikfinityEvent evt )
{
if ( !string.IsNullOrWhiteSpace( evt.Nickname ) ) return evt.Nickname;
if ( !string.IsNullOrWhiteSpace( evt.Username ) ) return evt.Username;
return "TikTok Viewer";
}
static int Clamp( int value, int minimum, int maximum )
{
if ( maximum < minimum ) maximum = minimum;
return value < minimum ? minimum : value > maximum ? maximum : value;
}
static string Limit( string value, int maximum )
{
if ( string.IsNullOrEmpty( value ) || maximum < 1 ) return "";
return value.Length <= maximum ? value : value.Substring( 0, maximum );
}
}