Code/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 );
    }
}