using System.Net; using System.Net.Sockets; using FinalFactory.Rendezvous.Contracts; using LiteNetLib; namespace FinalFactory.Rendezvous.Client; public enum RendezvousHostState { Active = 1, ManagerStopped = 2, Disposed = 3, } public sealed class RendezvousHostAttemptCompletedEventArgs : EventArgs { [Obsolete("Completion events now expose a typed Outcome. Construct these arguments only for legacy test doubles.")] public RendezvousHostAttemptCompletedEventArgs( JoinAttemptId attemptId, RendezvousConnectionState state, NetPeer? peer) : this(attemptId, state, RendezvousCompletionInvariant.FromLegacy(state, peer)) { } internal RendezvousHostAttemptCompletedEventArgs( JoinAttemptId attemptId, RendezvousConnectionState state, RendezvousConnectionOutcome outcome) { if (attemptId.Value == Guid.Empty) { throw new ArgumentException("The completed attempt ID is invalid.", nameof(attemptId)); } RendezvousCompletionInvariant.Validate(state, outcome); AttemptId = attemptId; State = state; Outcome = outcome; } public JoinAttemptId AttemptId { get; } public RendezvousConnectionState State { get; } public RendezvousConnectionOutcome Outcome { get; } public NetPeer? Peer => Outcome.Peer; } public sealed class RendezvousHostCoordinator : IDisposable { private readonly NetManager _manager; private readonly RendezvousNetListener _networkEvents; private readonly EventBasedNatPunchListener _punchEvents; private readonly IPEndPoint _mediator; private readonly PublishedSession _session; private readonly IRendezvousJoinClient _joinClient; private readonly RendezvousCoordinatorOptions _options; private readonly IRendezvousCoordinatorClock _clock; private readonly ConnectionTicketValidator _tickets; private readonly Dictionary _attempts = []; private readonly Dictionary _acceptedPeers = []; private readonly Dictionary _deferredRequests = []; private readonly Dictionary _terminalAttempts = []; private readonly Queue _attemptSchedule = []; private readonly SortedDictionary> _deadlines = []; private readonly List _cleanupScratch = []; private HostJoinAttempt[]? _latestSnapshot; private DateTimeOffset _nextPresenceAt = DateTimeOffset.MinValue; private DateTimeOffset _nextTerminalCleanupAt = DateTimeOffset.MinValue; private int _refreshing; private int _polling; private bool _subscriptionsReleased; private int _disposed; public RendezvousHostCoordinator( NetManager manager, RendezvousNetListener networkEvents, IPEndPoint mediator, PublishedSession session, IRendezvousJoinClient joinClient, RendezvousCoordinatorOptions? options = null) : this( manager, networkEvents, mediator, session, joinClient, options, new SystemRendezvousCoordinatorClock(), null) { } internal RendezvousHostCoordinator( NetManager manager, RendezvousNetListener networkEvents, IPEndPoint mediator, PublishedSession session, IRendezvousJoinClient joinClient, RendezvousCoordinatorOptions? options, IRendezvousCoordinatorClock clock, ConnectionTicketValidator? tickets) { _manager = manager ?? throw new ArgumentNullException(nameof(manager)); _networkEvents = networkEvents ?? throw new ArgumentNullException(nameof(networkEvents)); _punchEvents = _networkEvents.PunchEvents; _mediator = mediator ?? throw new ArgumentNullException(nameof(mediator)); _session = session ?? throw new ArgumentNullException(nameof(session)); _joinClient = joinClient ?? throw new ArgumentNullException(nameof(joinClient)); _options = (options ?? new RendezvousCoordinatorOptions()).CopyAndValidate(); _clock = clock ?? throw new ArgumentNullException(nameof(clock)); _tickets = tickets ?? new ConnectionTicketValidator(); RendezvousManagerGuard.Validate(_manager, _networkEvents); ValidateInputs(); _networkEvents.RendezvousConnectionRequest += OnConnectionRequest; _networkEvents.RendezvousPeerConnected += OnPeerConnected; _networkEvents.RendezvousPeerDisconnected += OnPeerDisconnected; _networkEvents.RendezvousNetworkError += OnNetworkError; _punchEvents.NatIntroductionSuccess += OnNatIntroductionSuccess; } public event EventHandler? AttemptCompleted; public RendezvousHostState State { get; private set; } = RendezvousHostState.Active; public int PendingAttemptCount => _attempts.Count; internal int DeferredRequestCount => _deferredRequests.Count; public async Task> RefreshJoinAttemptsAsync( CancellationToken cancellationToken = default) { ThrowIfDisposed(); if (Interlocked.Exchange(ref _refreshing, 1) != 0) { throw new InvalidOperationException("A host invitation refresh is already running."); } try { RendezvousClientResult> result = await _joinClient.BrowseAllForHostAsync( _session, cancellationToken: cancellationToken).ConfigureAwait(false); if (!result.IsSuccess || result.Value is null) { return RendezvousClientResult.Failure( result.Error, result.Message, result.RetryAfterSeconds); } HostJoinAttempt[] snapshot = result.Value.Select(CopyAttempt).ToArray(); if (Volatile.Read(ref _disposed) != 0) { throw new ObjectDisposedException(nameof(RendezvousHostCoordinator)); } Interlocked.Exchange(ref _latestSnapshot, snapshot); if (Volatile.Read(ref _disposed) != 0) { Interlocked.Exchange(ref _latestSnapshot, null); throw new ObjectDisposedException(nameof(RendezvousHostCoordinator)); } return RendezvousClientResult.Success(snapshot.Length); } finally { Volatile.Write(ref _refreshing, 0); } } public void Poll() { ThrowIfDisposed(); if (State != RendezvousHostState.Active) { return; } if (Interlocked.Exchange(ref _polling, 1) != 0) { throw new InvalidOperationException("The Rendezvous coordinator cannot be polled concurrently or recursively."); } try { ApplySnapshots(); if (!_manager.IsRunning) { Stop( RendezvousHostState.ManagerStopped, RendezvousConnectionState.ManagerStopped, ConnectionOutcomeKind.ManagerStopped); return; } _manager.NatPunchModule.PollEvents(); _manager.PollEvents(); _manager.NatPunchModule.PollEvents(); if (State != RendezvousHostState.Active) { return; } DateTimeOffset now = _clock.UtcNow; TimeSpan elapsed = _clock.Elapsed; if (!_manager.IsRunning) { Stop( RendezvousHostState.ManagerStopped, RendezvousConnectionState.ManagerStopped, ConnectionOutcomeKind.ManagerStopped); return; } RefreshPresence(now); ProcessDueDeadlines(elapsed); if (State != RendezvousHostState.Active) { return; } int checks = Math.Min( _attemptSchedule.Count, _options.MaximumAttemptChecksPerPoll); for (int index = 0; index < checks; index++) { JoinAttemptId attemptId = _attemptSchedule.Dequeue(); if (!_attempts.TryGetValue(attemptId, out PendingHostAttempt? attempt)) { continue; } if (attempt.State != RendezvousConnectionState.Punching) { continue; } if (attempt.Retry.IsDue(elapsed)) { if (attempt.Retry.IsExhausted) { CompleteAttempt( attemptId, RendezvousConnectionState.TimedOut, ConnectionOutcomeKind.PunchTimedOut, RendezvousConnectionOutcomeSource.LocalTraversal, RendezvousConnectionFailureCategory.NatTraversal, RendezvousConnectionPhase.NatTraversal); continue; } _manager.NatPunchModule.SendNatIntroduceRequest( _mediator, NatPunchRequestTokenCodec.Encode( NatPunchPeerRole.Host, attempt.Invitation.MediationHandle, attempt.Invitation.HostPunchCapability)); attempt.Retry.RecordRequest(); } _attemptSchedule.Enqueue(attemptId); } if (now >= _nextTerminalCleanupAt) { _cleanupScratch.Clear(); foreach (KeyValuePair terminal in _terminalAttempts) { if (terminal.Value <= now) { _cleanupScratch.Add(terminal.Key); } } foreach (JoinAttemptId attemptId in _cleanupScratch) { _terminalAttempts.Remove(attemptId); } _nextTerminalCleanupAt = now + TimeSpan.FromSeconds(1); } } finally { Volatile.Write(ref _polling, 0); } } public void Dispose() { if (Interlocked.Exchange(ref _disposed, 1) != 0) { return; } Stop( RendezvousHostState.Disposed, RendezvousConnectionState.Disposed, ConnectionOutcomeKind.Disposed); Interlocked.Exchange(ref _latestSnapshot, null); _attemptSchedule.Clear(); _deadlines.Clear(); _terminalAttempts.Clear(); _cleanupScratch.Clear(); _tickets.Dispose(); } public override string ToString() => $"[RendezvousHostCoordinator {_session.ListingId}; credentials redacted]"; private void ApplySnapshots() { HostJoinAttempt[]? latest = Interlocked.Exchange(ref _latestSnapshot, null); if (latest is null) { return; } DateTimeOffset now = _clock.UtcNow; TimeSpan elapsed = _clock.Elapsed; foreach (HostJoinAttempt invitation in latest) { if (invitation.AttemptId.Value == Guid.Empty || invitation.MediationHandle.Value == Guid.Empty || !ContractValidation.IsCapabilityValid(invitation.HostPunchCapability) || !ContractValidation.IsConnectionTicketValid( invitation.ConnectionTicketDigest)) { continue; } if (invitation.IsCancelled) { if (_attempts.ContainsKey(invitation.AttemptId)) { CompleteAttempt( invitation.AttemptId, RendezvousConnectionState.Cancelled, ConnectionOutcomeKind.Cancelled, RendezvousConnectionOutcomeSource.RendezvousService, RendezvousConnectionFailureCategory.Lifecycle, RendezvousConnectionPhase.Authorization); } _terminalAttempts[invitation.AttemptId] = invitation.ExpiresAt; continue; } if (invitation.ExpiresAt <= now || _attempts.ContainsKey(invitation.AttemptId) || _terminalAttempts.ContainsKey(invitation.AttemptId)) { continue; } TimeSpan attemptDeadline = elapsed + (invitation.ExpiresAt - now); TimeSpan punchDeadline = Min( attemptDeadline, elapsed + _options.PunchTimeout); _attempts.Add( invitation.AttemptId, new PendingHostAttempt( CopyAttempt(invitation), new RendezvousPunchRetrySchedule(_options, _clock), elapsed, attemptDeadline, punchDeadline)); EnqueueDeadline( new HostAttemptDeadline( invitation.AttemptId, RendezvousConnectionState.Punching, punchDeadline)); _attemptSchedule.Enqueue(invitation.AttemptId); } } private void RefreshPresence(DateTimeOffset now) { if (now < _nextPresenceAt || now >= _session.ExpiresAt) { return; } _manager.NatPunchModule.SendNatIntroduceRequest( _mediator, NatPunchRequestTokenCodec.Encode( NatPunchPeerRole.HostPresence, _session.HostPresenceHandle, _session.HostPresenceCapability)); _nextPresenceAt = now + TimeSpan.FromSeconds(_session.HostPresenceRefreshAfterSeconds); } private void OnNatIntroductionSuccess( IPEndPoint target, NatAddressType addressType, string encodedIntroduction) { _ = target; _ = addressType; if (!NatIntroductionTokenCodec.TryDecode( encodedIntroduction, out NatIntroductionToken? introduction) || introduction is null || !_attempts.TryGetValue(introduction.AttemptId, out PendingHostAttempt? attempt) || !NatIntroductionTokenCodec.MatchesDigest( introduction.ConnectionTicket, attempt.Invitation.ConnectionTicketDigest) || !_tickets.TryAuthorize( introduction.AttemptId, introduction.ConnectionTicket, Min( attempt.Invitation.ExpiresAt, _clock.UtcNow + _options.ConnectionTicketLifetime))) { return; } attempt.State = RendezvousConnectionState.Connecting; attempt.DirectDeadline = Min( attempt.AttemptDeadline, _clock.Elapsed + _options.DirectConnectTimeout); EnqueueDeadline(new HostAttemptDeadline( introduction.AttemptId, RendezvousConnectionState.Connecting, attempt.DirectDeadline.Value)); if (_deferredRequests.Remove( introduction.AttemptId, out DeferredConnectionRequest? deferred)) { AcceptAuthorizedRequest( introduction.AttemptId, attempt, deferred.Request, deferred.ConnectionTicket); } } private void OnConnectionRequest(ConnectionRequest request) { ReadOnlySpan data = request.Data.GetRemainingBytesSpan(); if (!DirectConnectionRequestCodec.IsRendezvousRequest(data)) { return; } if (!DirectConnectionRequestCodec.TryDecode(data, out DirectConnectionRequest? connection) || connection is null || !_attempts.TryGetValue(connection.AttemptId, out PendingHostAttempt? attempt) || !NatIntroductionTokenCodec.MatchesDigest( connection.ConnectionTicket, attempt.Invitation.ConnectionTicketDigest)) { request.RejectForce([]); return; } if (attempt.State == RendezvousConnectionState.Punching) { _deferredRequests[connection.AttemptId] = new( request, connection.ConnectionTicket); return; } if (attempt.State != RendezvousConnectionState.Connecting) { request.RejectForce([]); return; } AcceptAuthorizedRequest( connection.AttemptId, attempt, request, connection.ConnectionTicket); } private void OnPeerConnected(NetPeer peer) { if (_acceptedPeers.TryGetValue(peer, out JoinAttemptId attemptId)) { CompleteAttempt( attemptId, RendezvousConnectionState.Connected, ConnectionOutcomeKind.Connected, RendezvousConnectionOutcomeSource.LocalTraversal, RendezvousConnectionFailureCategory.None, RendezvousConnectionPhase.Complete, peer); } } private void OnPeerDisconnected(NetPeer peer, DisconnectInfo disconnectInfo) { _ = disconnectInfo; if (_acceptedPeers.TryGetValue(peer, out JoinAttemptId attemptId)) { ConnectionOutcomeKind kind = disconnectInfo.Reason == DisconnectReason.Timeout ? ConnectionOutcomeKind.DirectConnectTimedOut : ConnectionOutcomeKind.TransportError; CompleteAttempt( attemptId, kind == ConnectionOutcomeKind.DirectConnectTimedOut ? RendezvousConnectionState.TimedOut : RendezvousConnectionState.Rejected, kind, RendezvousConnectionOutcomeSource.LocalTraversal, RendezvousConnectionFailureCategory.DirectConnection, RendezvousConnectionPhase.DirectConnection); } } private void OnNetworkError(IPEndPoint endpoint, SocketError socketError) { _ = socketError; if (!endpoint.Equals(_mediator)) { return; } foreach (JoinAttemptId attemptId in _attempts .Where(static item => item.Value.State == RendezvousConnectionState.Punching) .Select(static item => item.Key) .ToArray()) { CompleteAttempt( attemptId, RendezvousConnectionState.Rejected, ConnectionOutcomeKind.MediatorUnavailable, RendezvousConnectionOutcomeSource.LocalTraversal, RendezvousConnectionFailureCategory.Mediation, RendezvousConnectionPhase.Mediation); } } private void CompleteAttempt( JoinAttemptId attemptId, RendezvousConnectionState state, ConnectionOutcomeKind kind, RendezvousConnectionOutcomeSource source, RendezvousConnectionFailureCategory category, RendezvousConnectionPhase phase, NetPeer? peer = null) { if (TryCompleteAttempt( attemptId, state, kind, source, category, phase, peer, out RendezvousHostAttemptCompletedEventArgs? completion)) { AttemptCompleted?.Invoke(this, completion!); } } private bool TryCompleteAttempt( JoinAttemptId attemptId, RendezvousConnectionState state, ConnectionOutcomeKind kind, RendezvousConnectionOutcomeSource source, RendezvousConnectionFailureCategory category, RendezvousConnectionPhase phase, NetPeer? peer, out RendezvousHostAttemptCompletedEventArgs? completion) { completion = null; if (!_attempts.Remove(attemptId, out PendingHostAttempt? attempt)) { return false; } if (attempt.AcceptedPeer is not null) { _acceptedPeers.Remove(attempt.AcceptedPeer); if (kind != ConnectionOutcomeKind.Connected) { attempt.AcceptedPeer.Disconnect(); } } if (_deferredRequests.Remove(attemptId, out DeferredConnectionRequest? deferred)) { deferred.Request.RejectForce([]); } _tickets.Revoke(attemptId); _terminalAttempts[attemptId] = attempt.Invitation.ExpiresAt; RendezvousConnectionOutcome outcome = RendezvousConnectionOutcome.Create( kind, source, category, phase, _clock.Elapsed - attempt.StartedAt, peer: peer); completion = new(attemptId, state, outcome); return true; } private void Stop( RendezvousHostState hostState, RendezvousConnectionState attemptState, ConnectionOutcomeKind outcomeKind) { if (State != RendezvousHostState.Active) { return; } State = hostState; List completions = []; foreach (JoinAttemptId attemptId in _attempts.Keys.ToArray()) { RendezvousConnectionPhase phase = _attempts[attemptId].State == RendezvousConnectionState.Connecting ? RendezvousConnectionPhase.DirectConnection : RendezvousConnectionPhase.NatTraversal; if (TryCompleteAttempt( attemptId, attemptState, outcomeKind, RendezvousConnectionOutcomeSource.Lifecycle, RendezvousConnectionFailureCategory.Lifecycle, phase, null, out RendezvousHostAttemptCompletedEventArgs? completion)) { completions.Add(completion!); } } ReleaseSubscriptions(); foreach (RendezvousHostAttemptCompletedEventArgs completion in completions) { AttemptCompleted?.Invoke(this, completion); } } private void ReleaseSubscriptions() { if (_subscriptionsReleased) { return; } _networkEvents.RendezvousConnectionRequest -= OnConnectionRequest; _networkEvents.RendezvousPeerConnected -= OnPeerConnected; _networkEvents.RendezvousPeerDisconnected -= OnPeerDisconnected; _networkEvents.RendezvousNetworkError -= OnNetworkError; _punchEvents.NatIntroductionSuccess -= OnNatIntroductionSuccess; _subscriptionsReleased = true; } private void ValidateInputs() { if (_mediator.Port is < 1 or > 65_535 || _session.HostPresenceHandle.Value == Guid.Empty || !ContractValidation.IsCapabilityValid(_session.HostPresenceCapability) || _session.HostPresenceRefreshAfterSeconds < 1 || _session.ExpiresAt <= _clock.UtcNow) { throw new ArgumentException("The host traversal inputs are invalid."); } } private static HostJoinAttempt CopyAttempt(HostJoinAttempt attempt) => new() { AttemptId = attempt.AttemptId, MediationHandle = attempt.MediationHandle, HostPunchCapability = attempt.HostPunchCapability, ConnectionTicketDigest = attempt.ConnectionTicketDigest, IsCancelled = attempt.IsCancelled, ExpiresAt = attempt.ExpiresAt, }; private static TimeSpan Min(TimeSpan left, TimeSpan right) => left <= right ? left : right; private static DateTimeOffset Min(DateTimeOffset left, DateTimeOffset right) => left <= right ? left : right; private void EnqueueDeadline(HostAttemptDeadline deadline) { if (!_deadlines.TryGetValue(deadline.Deadline.Ticks, out Queue? bucket)) { bucket = new Queue(); _deadlines.Add(deadline.Deadline.Ticks, bucket); } bucket.Enqueue(deadline); } private void ProcessDueDeadlines(TimeSpan elapsed) { while (_deadlines.Count > 0) { KeyValuePair> first = _deadlines.First(); if (first.Key > elapsed.Ticks) { return; } HostAttemptDeadline deadline = first.Value.Dequeue(); if (first.Value.Count == 0) { _deadlines.Remove(first.Key); } if (!_attempts.TryGetValue(deadline.AttemptId, out PendingHostAttempt? attempt) || attempt.State != deadline.ExpectedState || (deadline.ExpectedState == RendezvousConnectionState.Punching ? attempt.PunchDeadline : attempt.DirectDeadline) != deadline.Deadline) { continue; } bool expired = elapsed >= attempt.AttemptDeadline; CompleteAttempt( deadline.AttemptId, RendezvousConnectionState.TimedOut, expired ? ConnectionOutcomeKind.AttemptExpired : deadline.ExpectedState == RendezvousConnectionState.Punching ? ConnectionOutcomeKind.PunchTimedOut : ConnectionOutcomeKind.DirectConnectTimedOut, expired ? RendezvousConnectionOutcomeSource.RendezvousService : RendezvousConnectionOutcomeSource.LocalTraversal, expired ? RendezvousConnectionFailureCategory.Authorization : deadline.ExpectedState == RendezvousConnectionState.Punching ? RendezvousConnectionFailureCategory.NatTraversal : RendezvousConnectionFailureCategory.DirectConnection, expired ? RendezvousConnectionPhase.Authorization : deadline.ExpectedState == RendezvousConnectionState.Punching ? RendezvousConnectionPhase.NatTraversal : RendezvousConnectionPhase.DirectConnection); if (State != RendezvousHostState.Active) { return; } } } private void AcceptAuthorizedRequest( JoinAttemptId attemptId, PendingHostAttempt attempt, ConnectionRequest request, string connectionTicket) { ConnectionTicketConsumptionResult consumption = _tickets.Consume( attemptId, connectionTicket); if (consumption != ConnectionTicketConsumptionResult.Accepted) { request.RejectForce([]); return; } NetPeer peer = request.Accept(); attempt.AcceptedPeer = peer; _acceptedPeers[peer] = attemptId; } private void ThrowIfDisposed() { if (Volatile.Read(ref _disposed) != 0) { throw new ObjectDisposedException(nameof(RendezvousHostCoordinator)); } } private sealed class PendingHostAttempt( HostJoinAttempt invitation, RendezvousPunchRetrySchedule retry, TimeSpan startedAt, TimeSpan attemptDeadline, TimeSpan punchDeadline) { internal HostJoinAttempt Invitation { get; } = invitation; internal RendezvousPunchRetrySchedule Retry { get; } = retry; internal TimeSpan StartedAt { get; } = startedAt; internal TimeSpan AttemptDeadline { get; } = attemptDeadline; internal TimeSpan PunchDeadline { get; } = punchDeadline; internal TimeSpan? DirectDeadline { get; set; } internal RendezvousConnectionState State { get; set; } = RendezvousConnectionState.Punching; internal NetPeer? AcceptedPeer { get; set; } } private sealed class HostAttemptDeadline( JoinAttemptId attemptId, RendezvousConnectionState expectedState, TimeSpan deadline) { internal JoinAttemptId AttemptId { get; } = attemptId; internal RendezvousConnectionState ExpectedState { get; } = expectedState; internal TimeSpan Deadline { get; } = deadline; } private sealed class DeferredConnectionRequest( ConnectionRequest request, string connectionTicket) { internal ConnectionRequest Request { get; } = request; internal string ConnectionTicket { get; } = connectionTicket; } }